(self, message: Any)
| 71 | self.failed_tasks = 0 |
| 72 | self._emit_queue_metrics_unlocked() |
| 73 | |
| 74 | def create_task(self, message: Any) -> str: |
| 75 | with self._task_available: |
| 76 | if hasattr(message, "task_id") and message.task_id in self._tasks: |
| 77 | raise RuntimeError(f"Task ID {message.task_id} already exists") |
| 78 | |
| 79 | active_tasks = sum(1 for t in self._tasks.values() if t.status in [TaskStatus.PENDING, TaskStatus.PROCESSING]) |
| 80 | if active_tasks >= self.max_queue_size: |
| 81 | raise RuntimeError(f"Task queue is full (max {self.max_queue_size} tasks)") |
| 82 | |
| 83 | task_id = getattr(message, "task_id", str(uuid.uuid4())) |
| 84 | task_info = TaskInfo(task_id=task_id, status=TaskStatus.PENDING, message=message, save_result_path=getattr(message, "save_result_path", None)) |
| 85 | |
| 86 | self._tasks[task_id] = task_info |
| 87 | self.total_tasks += 1 |
| 88 | |
| 89 | self._cleanup_old_tasks() |
| 90 | self._emit_queue_metrics_unlocked() |
| 91 | self._task_available.notify() |
| 92 | |
| 93 | return task_id |
| 94 | |
| 95 | def start_task(self, task_id: str) -> TaskInfo: |
no test coverage detected