"""File-based task registry for background dispatch tasks. The dispatch MCP server runs in a separate process, so we use a JSONL event file as the communication channel. The main agent process tails this file and pushes events to the room (UI) and speaks results. Event file: /tmp/tasks/events.jsonl Task state: /tmp/tasks/{task_id}.json """ from __future__ import annotations import asyncio import json import os import time import uuid from dataclasses import dataclass, field from enum import Enum from pathlib import Path from typing import Any, Callable TASKS_DIR = Path(os.environ.get("TASKS_DIR", "/tmp/tasks")) EVENTS_FILE = TASKS_DIR / "events.jsonl" class TaskStatus(str, Enum): PENDING = "pending" RUNNING = "running" COMPLETED = "completed" FAILED = "failed" @dataclass class TaskStep: role: str content: str timestamp: float = field(default_factory=time.time) @dataclass class Task: id: str description: str status: TaskStatus = TaskStatus.PENDING steps: list[TaskStep] = field(default_factory=list) result: str | None = None error: str | None = None created_at: float = field(default_factory=time.time) completed_at: float | None = None class TaskRegistry: """File-backed registry. Writers (MCP process) append to events.jsonl. Readers (main agent) tail the file and dispatch to subscribers.""" def __init__(self): TASKS_DIR.mkdir(parents=True, exist_ok=True) if not EVENTS_FILE.exists(): EVENTS_FILE.touch() self._tasks: dict[str, Task] = {} self._listeners: list[Callable[[str, dict], Any]] = [] self._tail_pos = 0 def create(self, description: str) -> Task: task = Task(id=uuid.uuid4().hex[:8], description=description) self._tasks[task.id] = task self._write_event(task.id, {"event": "created", "task": self._task_dict(task)}) return task def add_step(self, task_id: str, role: str, content: str): task = self._tasks.get(task_id) if task: task.steps.append(TaskStep(role=role, content=content)) if task.status == TaskStatus.PENDING: task.status = TaskStatus.RUNNING self._save_task(task) self._write_event(task_id, {"event": "step", "role": role, "content": content}) def complete(self, task_id: str, result: str): task = self._tasks.get(task_id) if task: task.status = TaskStatus.COMPLETED task.result = result task.completed_at = time.time() self._save_task(task) self._write_event(task_id, {"event": "completed", "result": result}) def fail(self, task_id: str, error: str): task = self._tasks.get(task_id) if task: task.status = TaskStatus.FAILED task.error = error task.completed_at = time.time() self._save_task(task) self._write_event(task_id, {"event": "failed", "error": error}) def get(self, task_id: str) -> dict | None: t = self._tasks.get(task_id) return self._task_dict(t) if t else None def list_all(self) -> list[dict]: return [self._task_dict(t) for t in self._tasks.values()] def subscribe(self, callback: Callable[[str, dict], Any]): self._listeners.append(callback) async def poll_events(self) -> list[tuple[str, dict]]: """Read new events from the file. Call periodically from the agent.""" events = [] try: with open(EVENTS_FILE, "r") as f: f.seek(self._tail_pos) for line in f: line = line.strip() if not line: continue try: ev = json.loads(line) events.append((ev.get("task_id", ""), ev)) except json.JSONDecodeError: pass self._tail_pos = f.tell() except (OSError, IOError): pass return events async def notify_subscribers(self, task_id: str, event: dict): for cb in self._listeners: try: result = cb(task_id, event) if asyncio.iscoroutine(result): await result except Exception: pass def _write_event(self, task_id: str, event: dict): payload = json.dumps({"task_id": task_id, **event}) with open(EVENTS_FILE, "a") as f: f.write(payload + "\n") def _save_task(self, task: Task): path = TASKS_DIR / f"{task.id}.json" path.write_text(json.dumps(self._task_dict(task))) @staticmethod def _task_dict(task: Task) -> dict: return { "id": task.id, "description": task.description, "status": task.status.value, "steps": [ {"role": s.role, "content": s.content, "ts": s.timestamp} for s in task.steps ], "result": task.result, "error": task.error, "created_at": task.created_at, "completed_at": task.completed_at, } # Global singleton — one per process registry = TaskRegistry()