from __future__ import annotations from django.db import transaction from control_plane.events.bus import EventBus from control_plane.events.models import EventType from control_plane.projects.models import Task, TaskStatus class TaskScheduler: def __init__(self, bus: EventBus | None = None) -> None: self.bus = bus or EventBus() def claim_next_ready_task(self) -> Task | None: with transaction.atomic(): candidates = ( Task.objects.select_for_update(skip_locked=True) .filter(status=TaskStatus.READY) .order_by("-priority", "created_at") ) for task in candidates: blocking_dependencies = task.dependency_edges.exclude(depends_on__status=TaskStatus.COMPLETE) if blocking_dependencies.exists(): continue task.status = TaskStatus.RUNNING task.save(update_fields=["status", "updated_at"]) self.bus.publish(EventType.TASK_STARTED, project=task.project, task=task, payload={"task_id": str(task.id)}) return task return None