| 103 | worker.terminate() |
| 104 | |
| 105 | class JobManager: |
| 106 | def __init__(self, args): |
| 107 | self.past_jobs = [] |
| 108 | self.work_queue = multiprocessing.Queue() |
| 109 | self.result_queue = multiprocessing.Queue() |
| 110 | self.next_job_id = 0 |
| 111 | self.args = args |
| 112 | self.workers = JobWorkers(args) |
| 113 | self.workers.launch(self.work_queue, self.result_queue) |
| 114 | |
| 115 | def _add_job(self, task_op, code = None, test = None): |
| 116 | job_id = self.next_job_id |
| 117 | job = (task_op, job_id, code, test) |
| 118 | self.past_jobs.append(job) |
| 119 | self.next_job_id += 1 |
| 120 | |
| 121 | self.work_queue.put(job) |
| 122 | |
| 123 | return job_id |
| 124 | |
| 125 | def add_code_job(self): |
| 126 | return self._add_job("code") |
| 127 | |
| 128 | def add_test_job(self): |
| 129 | return self._add_job("test") |
| 130 | |
| 131 | def add_improve_code_job(self, code): |
| 132 | return self._add_job("improve_code", code=code) |
| 133 | |
| 134 | def add_improve_test_job(self, test): |
| 135 | return self._add_job("improve_test", test=test) |
| 136 | |
| 137 | def add_judge_pair_job(self, code, test): |
| 138 | return self._add_job("judge_pair", code=code, test=test) |
| 139 | |
| 140 | def get_results(self, timeout: float = 1.0) -> List[Tuple[str, int, Any, Any]]: |
| 141 | results = [] |
| 142 | try: |
| 143 | results.append(self.result_queue.get(timeout=timeout)) |
| 144 | while not self.result_queue.empty(): |
| 145 | results.append(self.result_queue.get(timeout=timeout)) |
| 146 | except Empty: |
| 147 | pass |
| 148 | return results |
| 149 | |
| 150 | def active_workers(self): |
| 151 | return self.workers.active_workers.value |
| 152 | |
| 153 | def approx_queue_depth(self): |
| 154 | return self.work_queue.qsize() |
| 155 | |
| 156 | def terminate(self): |
| 157 | self.workers.terminate() |