Spaces:
Running on Zero
Running on Zero
| # ===================================================== | |
| # Apckeyl Framework | |
| # Version 1.0 | |
| # task_manager.py | |
| # ===================================================== | |
| """ | |
| Apckeyl Task Manager. | |
| Version 1.0 | |
| Создаёт и хранит задачи Control Plane. | |
| На этом этапе Task Manager работает только | |
| в памяти процесса. | |
| Постоянное хранение, очередь и внешняя база | |
| будут добавлены позже. | |
| """ | |
| from datetime import datetime, timezone | |
| from itertools import count | |
| from control_plane import control_plane | |
| # ===================================================== | |
| # Task Manager | |
| # ===================================================== | |
| class TaskManager: | |
| def __init__(self, control_plane_instance=None): | |
| if control_plane_instance is None: | |
| control_plane_instance = control_plane | |
| self.control_plane = control_plane_instance | |
| self._tasks = {} | |
| self._counter = count(1) | |
| # ================================================= | |
| # Create Task | |
| # ================================================= | |
| def create_task( | |
| self, | |
| task_type, | |
| payload=None, | |
| ): | |
| if not task_type: | |
| raise ValueError( | |
| "task_type is required" | |
| ) | |
| module = self.control_plane.resolve_task( | |
| task_type | |
| ) | |
| number = next( | |
| self._counter | |
| ) | |
| task_id = ( | |
| f"task_{number:06d}" | |
| ) | |
| task = { | |
| "task_id": task_id, | |
| "task_type": task_type, | |
| "status": "created", | |
| "module_id": module[ | |
| "module_id" | |
| ], | |
| "module_name": module[ | |
| "module_name" | |
| ], | |
| "payload": payload, | |
| "created_at": ( | |
| datetime.now( | |
| timezone.utc | |
| ).isoformat() | |
| ), | |
| } | |
| self._tasks[ | |
| task_id | |
| ] = task | |
| return dict(task) | |
| # ================================================= | |
| # Get Task | |
| # ================================================= | |
| def get_task( | |
| self, | |
| task_id, | |
| ): | |
| task = self._tasks.get( | |
| task_id | |
| ) | |
| if task is None: | |
| return None | |
| return dict(task) | |
| # ================================================= | |
| # Update Status | |
| # ================================================= | |
| def update_status( | |
| self, | |
| task_id, | |
| status, | |
| ): | |
| task = self._tasks.get( | |
| task_id | |
| ) | |
| if task is None: | |
| raise KeyError( | |
| f"Unknown task: {task_id}" | |
| ) | |
| task["status"] = status | |
| return dict(task) | |
| # ================================================= | |
| # List Tasks | |
| # ================================================= | |
| def list_tasks(self): | |
| return [ | |
| dict(task) | |
| for task | |
| in self._tasks.values() | |
| ] | |
| # ================================================= | |
| # Delete Task | |
| # ================================================= | |
| def delete_task( | |
| self, | |
| task_id, | |
| ): | |
| if task_id in self._tasks: | |
| del self._tasks[ | |
| task_id | |
| ] | |
| # ===================================================== | |
| # Default Task Manager | |
| # ===================================================== | |
| task_manager = TaskManager() |