Orchestrator / task_manager.py
evgeniy778's picture
Add Task Manager Version 1.0
4595394 verified
Raw
History Blame Contribute Delete
3.6 kB
# =====================================================
# 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()