29 lines
1.1 KiB
Python
29 lines
1.1 KiB
Python
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
|