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 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
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 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
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
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
|