| 24 | |
| 25 | |
| 26 | class ClientConnection: |
| 27 | def __init__(self, directory: Optional[str] = None): |
| 28 | self.ctx: HqClientContext = ffi.connect_to_server(directory) |
| 29 | |
| 30 | def submit_job(self, job_description: JobDescription) -> JobId: |
| 31 | return ffi.submit_job(self.ctx, job_description) |
| 32 | |
| 33 | def wait_for_jobs(self, job_ids: Sequence[JobId], callback) -> List[JobId]: |
| 34 | """Blocks until jobs are finished. Returns the number of failed tasks""" |
| 35 | return ffi.wait_for_jobs(self.ctx, job_ids, callback) |
| 36 | |
| 37 | def stop_server(self): |
| 38 | return ffi.stop_server(self.ctx) |
| 39 | |
| 40 | def get_failed_tasks(self, job_ids: Sequence[JobId]) -> TaskFailureMap: |
| 41 | jobs = ffi.get_failed_tasks(self.ctx, job_ids) |
| 42 | return { |
| 43 | job_id: {task_id: FailedTaskContext(**data) for (task_id, data) in task_data.items()} |
| 44 | for (job_id, task_data) in jobs.items() |
| 45 | } |
| 46 | |
| 47 | def forget_job(self, job_id: JobId): |
| 48 | return ffi.forget_job(self.ctx, job_id) |
no outgoing calls