Skip to content

manager

tit.jobs.manager

JobManager — the in-process job engine (TODO.md §2.3–2.5).

Owns a private background thread running its own asyncio event loop (independent of whatever loop tit.server itself runs requests on), so spawning/awaiting subprocesses never depends on the request lifecycle of any particular HTTP call. Public methods are synchronous and thread-safe (a single RLock guards the shared job table); everything async (subprocess spawn/wait, event tailing, termination) happens on the manager's own loop.

Testing this directly (construct a JobManager pointed at a tmp_path with a fake runner, call its synchronous methods, poll get() for state changes, shutdown() at teardown) needs no asyncio test plugin — see tests/test_jobs_manager.py / tests/fake_runner.py.

JobManager

JobManager(project_dir: str, *, runner: Runner | None = None, runner_cwd: str | None = None, poll_interval: float = 0.25, budget: Cost | None = None, command_for: Any = None)
Source code in tit/jobs/manager.py
def __init__(
    self,
    project_dir: str,
    *,
    runner: Runner | None = None,
    runner_cwd: str | None = None,
    poll_interval: float = 0.25,
    budget: Cost | None = None,
    command_for: Any = None,
) -> None:
    self.project_dir = project_dir
    self.runner = runner or LocalPopenRunner()
    self.runner_cwd = runner_cwd or project_dir
    # Bound to this manager's project root so a `tools` job's path arguments can be
    # jailed to it (tit.jobs.kinds.check_tool_args). A caller-supplied builder (tests,
    # and only tests) keeps the plain three-argument signature.
    self._command_for = command_for or functools.partial(
        kinds.command_for, project_dir=project_dir
    )
    self.poll_interval = poll_interval
    self.registry = JobRegistry(project_dir)
    # An injected budget (tests) is fixed; otherwise the CPU half is re-read from the user's
    # global CPU limit on every admission pass (see `_admission_budget`).
    self._fixed_budget = budget
    self._budget = budget or scheduler.discover_budget()

    self._lock = threading.RLock()
    self._specs: dict[str, JobSpec] = {}
    self._status: dict[str, JobStatus] = {}
    self._processes: dict[str, asyncio.subprocess.Process] = {}
    self._tailers: dict[str, EventTailer] = {}
    self._last_event_ts: dict[str, float] = {}
    self._last_exit_code: dict[str, int] = {}
    #: One sampler per running job, created when its pid appears, dropped at finalize.
    self._samplers: dict[str, ResourceSampler] = {}
    self._last_sample_ts: dict[str, float] = {}
    self._cancelled: set[str] = set()
    #: Non-terminal jobs found in the store at start(), i.e. left over from a previous
    #: server life -- reconciled once, before the first tick (see _reconcile_all).
    self._stranded: list[str] = []
    #: Set once _reconcile_all() has finished; start() waits on it (see start()).
    self._reconciled = threading.Event()

    self._status_subs: list[queue.Queue[dict[str, Any]]] = []
    self._event_subs: dict[str, list[queue.Queue[dict[str, Any]]]] = {}

    self._loop: asyncio.AbstractEventLoop | None = None
    self._main_task: asyncio.Task | None = None
    self._thread: threading.Thread | None = None
    self._started = False

submit_plan

submit_plan(planned: list[PlannedJob], *, created_by: str = 'api') -> dict[str, Any]

Submit a labelled DAG of :class:PlannedJob (tit.jobs.plans.plan_preprocessing) under one shared group_id, resolving after_labels to real job ids.

The scheduler runs one job per product at a time (:func:tit.jobs.scheduler.evaluate), so a group's jobs run one after another.

Source code in tit/jobs/manager.py
def submit_plan(
    self,
    planned: list[PlannedJob],
    *,
    created_by: str = "api",
) -> dict[str, Any]:
    """Submit a labelled DAG of :class:`PlannedJob` (``tit.jobs.plans.plan_preprocessing``)
    under one shared ``group_id``, resolving ``after_labels`` to real job ids.

    The scheduler runs one job per product at a time (:func:`tit.jobs.scheduler.evaluate`),
    so a group's jobs run one after another.
    """
    group_id = new_job_id()
    label_to_id: dict[str, str] = {}
    submitted: list[dict[str, Any]] = []
    # Labels must be defined before they're referenced (plan_preprocessing emits them in
    # dependency order); resolve what we can, ignore an unknown label defensively.
    for job in planned:
        after_ids = [
            label_to_id[lbl] for lbl in job.after_labels if lbl in label_to_id
        ]
        status = self.submit(
            job.kind,
            job.config,
            job.subject_ids,
            after=after_ids,
            tags=job.tags,
            created_by=created_by,
            group_id=group_id,
            overwrite=job.overwrite,
        )
        label_to_id[job.label] = status["id"]
        submitted.append(status)
    return {"group_id": group_id, "jobs": submitted}

get_outputs

get_outputs(job_id: str) -> dict[str, Any] | None

What the job's output folder holds on disk -- {"folder", "files"}.

The folder is the one the job's registered artifacts point at (a dir artifact first, else the directory of its first file artifact); the files are read from the filesystem now, so the Artifacts tab shows every output of a run rather than only the handful a runner registered. folder is None (and files empty) for a job that registered nothing -- one that failed before writing, or project_init.

Source code in tit/jobs/manager.py
def get_outputs(self, job_id: str) -> dict[str, Any] | None:
    """What the job's output folder holds *on disk* -- ``{"folder", "files"}``.

    The folder is the one the job's registered artifacts point at (a ``dir`` artifact
    first, else the directory of its first file artifact); the files are read from the
    filesystem now, so the Artifacts tab shows every output of a run rather than only the
    handful a runner registered. ``folder`` is ``None`` (and ``files`` empty) for a job
    that registered nothing -- one that failed before writing, or ``project_init``.
    """
    from tit.catalog import output_files

    with self._lock:
        status = self._status.get(job_id)
        if status is None:
            return None
        artifacts = list(status.artifacts)
    folder = None
    for a in artifacts:
        if a.kind == "dir" and a.path:
            folder = a.path
            break
    if folder is None:
        first = next((a.path for a in artifacts if a.path), None)
        if first:
            folder = os.path.dirname(first)
    files = output_files(folder, project_root=self.project_dir) if folder else []
    return {"folder": folder, "files": files}

force

force(job_id: str) -> dict[str, Any] | None

Force a stuck/lost job straight to a terminal state without waiting on its process (the OpenAPI contract's semantics for POST /api/jobs/{id}/force — a recovery escape hatch, distinct from "skip lock waits").

Source code in tit/jobs/manager.py
def force(self, job_id: str) -> dict[str, Any] | None:
    """Force a stuck/lost job straight to a terminal state without waiting on its process
    (the OpenAPI contract's semantics for ``POST /api/jobs/{id}/force`` — a recovery escape
    hatch, distinct from "skip lock waits").
    """
    with self._lock:
        status = self._status.get(job_id)
        if status is None:
            return None
        if status.state not in ("running", "lost", "queued"):
            return status.to_api(self.project_dir)
        pid = status.pid
        self._cancelled.add(job_id)
        state = "lost" if status.state != "queued" else "cancelled"
        self._finalize_locked(
            status,
            state=state,
            exit_code=None,
            error=JobError(
                type="forced",
                message="forced to a terminal state by the user without waiting for the process",
                last_lines=self._log_tail(job_id),
            ),
        )
    locks.release_job(self.project_dir, job_id)
    if pid is not None and self._loop is not None:
        # Best-effort, fire-and-forget: still try to actually stop it in the background.
        asyncio.run_coroutine_threadsafe(terminate_tree(pid), self._loop)
    with self._lock:
        return self._status[job_id].to_api(self.project_dir)

delete

delete(job_id: str) -> str

Returns "deleted", "not_found", or "not_terminal".

Source code in tit/jobs/manager.py
def delete(self, job_id: str) -> str:
    """Returns ``"deleted"``, ``"not_found"``, or ``"not_terminal"``."""
    from tit.jobs.spec import TERMINAL_STATES

    with self._lock:
        status = self._status.get(job_id)
        if status is None:
            return "not_found"
        if status.state not in TERMINAL_STATES:
            return "not_terminal"
        # Keep the in-memory record if persistent deletion fails, so it can be retried.
        self.registry.delete(job_id)
        del self._status[job_id]
        self._specs.pop(job_id, None)
        self._tailers.pop(job_id, None)
        self._event_subs.pop(job_id, None)
    return "deleted"

subscribe_events

subscribe_events(job_id: str, since: int = 0) -> Queue[dict[str, Any]]

Backfill and stream a registered job; raise ValueError for unknown ids.

Source code in tit/jobs/manager.py
def subscribe_events(
    self, job_id: str, since: int = 0
) -> queue.Queue[dict[str, Any]]:
    """Backfill and stream a registered job; raise ``ValueError`` for unknown ids."""
    with self._lock:
        # WebSocket subscription keys are client input, unlike scheduler-owned ids.
        # Match REST backfill's registration check before a key becomes a path.
        if job_id not in self._specs:
            raise ValueError(f"unknown job: {job_id}")
        backlog = read_events(events_path(self.project_dir, job_id), since=since)
        # No consumer exists until this method returns. A fixed-size blocking backfill
        # deadlocked the manager (and its status updates) for logs over 10,000 events.
        q: queue.Queue[dict[str, Any]] = queue.Queue(
            maxsize=len(backlog) + SUBSCRIBER_QUEUE_MAXSIZE
        )
        for event in backlog:
            q.put_nowait(event)
        self._event_subs.setdefault(job_id, []).append(q)
    return q