Skip to content

runner

tit.jobs.runner

Runner interface + the local subprocess implementation (TODO.md §2.3).

Runner is deliberately narrow (spawn only) so a future SbatchRunner (3.1, HPC) can drop in: everything else — waiting for exit, killing a process tree — works on a bare pid and doesn't care how the process was started, which is also what makes cancelling a re-attached job (one this server process never spawned) possible.

RunRequest dataclass

RunRequest(job_id: str, argv: list[str], cwd: str, env: dict[str, str] = dict(), stdout_path: str = devnull)

Everything :meth:Runner.spawn needs to start one job's process.

Runner

Bases: ABC

Starts a job's OS process. Waiting/cancelling operate on the resulting pid directly.

spawn abstractmethod async

spawn(request: RunRequest) -> Process

Start the process; must not block until exit (fire-and-return).

Source code in tit/jobs/runner.py
@abc.abstractmethod
async def spawn(self, request: RunRequest) -> asyncio.subprocess.Process:
    """Start the process; must not block until exit (fire-and-return)."""

LocalPopenRunner

Bases: Runner

asyncio.create_subprocess_exec with the safe-spawn options from TODO.md §2.3.

Its own process group (:func:tit.jobs.processes.spawn_kwargs -- POSIX start_new_session, Windows CREATE_NEW_PROCESS_GROUP -- so a cancel's SIGTERM/SIGKILL (or their Windows equivalents) never hits the server itself), stdin DEVNULL, stdout appended to stdout_path with stderr merged in, close_fds=True. Never preexec_fn (unsafe with threads) and never a pipe (a dead server would BrokenPipe the runner instead of leaving it to finish on its own).

ResourceSampler

ResourceSampler(pid: int)

Per-job CPU%/RSS sampler over the job's whole process tree, kept for the job's lifetime.

psutil.Process.cpu_percent(interval=None) is a delta against the previous call on the same object, so a fresh Process(pid) per poll always reads 0.0. One sampler therefore owns one Process per pid it has seen (the root and every descendant -- SimNIBS/FastSurfer/PARDISO spawn workers), created on first sight and dropped when it exits, and every :meth:sample sums CPU% and RSS across the live tree.

Memory is the tree's PSS (proportional set size) where the platform reports it (Linux, /proc/<pid>/smaps_rollup, ~0.06 ms per process), so the leadfield a dozen joblib workers share copy-on-write is counted once, not twelve times -- summing plain RSS over an ex-search tree read 54 GB where PSS reads 9 GB. Elsewhere it falls back to RSS. The field keeps the name rss on the wire.

Running statistics: cpu_peak/rss_peak are the maxima; cpu_avg/rss_avg are simple means over samples (samples are taken at a fixed cadence by the manager, so this is time-weighted to within one interval). The very first CPU reading of every process is 0.0 by psutil convention and is not counted, so a one-sample job does not report 0 %.

Source code in tit/jobs/runner.py
def __init__(self, pid: int) -> None:
    self.pid = pid
    self._procs: dict[int, psutil.Process] = {}
    self._primed: set[int] = set()
    self.n_samples = 0
    self.cpu_peak: float | None = None
    self.cpu_avg: float | None = None
    self.rss_peak: int | None = None
    self.rss_avg: int | None = None
    self._cpu_sum = 0.0
    self._rss_sum = 0
    try:
        self._procs[pid] = psutil.Process(pid)
    except (psutil.NoSuchProcess, psutil.AccessDenied, ValueError):
        pass

sample

sample() -> tuple[float | None, int | None]

One reading of (CPU %, RSS bytes) summed over the tree; (None, None) if the root is gone. Updates the running peak/average.

Source code in tit/jobs/runner.py
def sample(self) -> tuple[float | None, int | None]:
    """One reading of (CPU %, RSS bytes) summed over the tree; ``(None, None)`` if the
    root is gone. Updates the running peak/average."""
    procs = self._tree()
    if not procs:
        return None, None
    cpu = 0.0
    rss = 0
    counted = False
    seen_any = False
    for proc in procs:
        try:
            c = round(proc.cpu_percent(interval=None), 1)
            rss += _resident_bytes(proc)
        except (psutil.NoSuchProcess, psutil.AccessDenied, ValueError):
            continue
        seen_any = True
        if proc.pid in self._primed:
            cpu += c
            counted = True
        else:
            self._primed.add(proc.pid)
    if not seen_any:
        return None, None
    cpu = round(cpu, 1)
    self.rss_peak = rss if self.rss_peak is None else max(self.rss_peak, rss)
    if counted:
        self.n_samples += 1
        self._cpu_sum += cpu
        self._rss_sum += rss
        self.cpu_peak = cpu if self.cpu_peak is None else max(self.cpu_peak, cpu)
        self.cpu_avg = round(self._cpu_sum / self.n_samples, 1)
        self.rss_avg = int(self._rss_sum / self.n_samples)
    return (cpu if counted else None), rss

runner_env

runner_env(job_id: str, events_file: str, *, interface: str = 'api', base_env: dict[str, str] | None = None, cpus: float | None = None, kind: str | None = None, subject_ids: list[str] | tuple[str, ...] | None = None) -> dict[str, str]

The child process environment (TODO.md §2.3): unbuffered/faulthandler Python, job identity, thread knobs derived from the budget, and TIT_SERVER_TOKEN / TIT_SERVER_SETTINGS_FILE scrubbed so a job can never read the server's own auth secret directly, or the 0600 settings file a --reload parent hands its child (which itself carries the token in plain JSON) -- see tit/server/__main__.py's module docstring, which promises both are scrubbed (ra_14 finding #10: only the first one actually was).

Also adds -no_signal_handler to PETSC_OPTIONS (see :data:PETSC_NO_SIGNAL_HANDLER) so cancelling a job does not end its log in a PETSc/MPI crash block; a value the caller already set is kept and appended to.

Source code in tit/jobs/runner.py
def runner_env(
    job_id: str,
    events_file: str,
    *,
    interface: str = "api",
    base_env: dict[str, str] | None = None,
    cpus: float | None = None,
    kind: str | None = None,
    subject_ids: list[str] | tuple[str, ...] | None = None,
) -> dict[str, str]:
    """The child process environment (TODO.md §2.3): unbuffered/faulthandler Python, job
    identity, thread knobs derived from the budget, and ``TIT_SERVER_TOKEN`` /
    ``TIT_SERVER_SETTINGS_FILE`` scrubbed so a job can never read the server's own auth secret
    directly, or the 0600 settings file a ``--reload`` parent hands its child (which itself
    carries the token in plain JSON) -- see ``tit/server/__main__.py``'s module docstring,
    which promises both are scrubbed (ra_14 finding #10: only the first one actually was).

    Also adds ``-no_signal_handler`` to ``PETSC_OPTIONS`` (see
    :data:`PETSC_NO_SIGNAL_HANDLER`) so cancelling a job does not end its log in a PETSc/MPI
    crash block; a value the caller already set is kept and appended to.
    """
    env = dict(base_env if base_env is not None else os.environ)
    env.pop("TIT_SERVER_TOKEN", None)
    env.pop("TIT_SERVER_SETTINGS_FILE", None)
    existing_pythonpath = env.get("PYTHONPATH", "")
    env["PYTHONPATH"] = (
        os.pathsep.join([_TIT_ROOT, existing_pythonpath])
        if existing_pythonpath
        else _TIT_ROOT
    )
    env["PYTHONUNBUFFERED"] = "1"
    env["PYTHONFAULTHANDLER"] = "1"
    existing_petsc = env.get(PETSC_OPTIONS_ENV, "")
    if PETSC_NO_SIGNAL_HANDLER not in existing_petsc.split():
        env[PETSC_OPTIONS_ENV] = (
            f"{existing_petsc} {PETSC_NO_SIGNAL_HANDLER}".strip()
            if existing_petsc
            else PETSC_NO_SIGNAL_HANDLER
        )
    env[ENV_JOB_ID] = job_id
    if kind:
        env[ENV_JOB_KIND] = kind
    if subject_ids:
        env[ENV_JOB_SUBJECT_IDS] = ",".join(str(sid) for sid in subject_ids)
    env[ENV_EVENTS_FILE] = events_file
    env[ENV_INTERFACE] = interface
    if cpus is not None and cpus > 0:
        threads = str(max(1, int(cpus)))
        for name in ("OMP_NUM_THREADS", "MKL_NUM_THREADS", "NUMBA_NUM_THREADS"):
            env[name] = threads
        env["TI_NIFTI_WORKERS"] = threads
        # The same number, named rather than inferred, for code that picks a *worker count*
        # rather than a thread count (`tit.opt.ex.parallel.resolve_n_jobs`,
        # `tit.opt.flex`): a solver's "use every core" default must not exceed -- or fall
        # short of -- the CPUs the plan showed the user. See `tit.cpu.job_cpus`.
        env[JOB_CPUS_ENV] = threads
    return env

is_alive

is_alive(pid: int, create_time: float | None) -> bool

True if pid is a live, non-zombie process — and, when create_time is known, it's still the same process (not pid reuse after the original exited). Never true for pid 1 or this server's own pid, regardless of what create_time claims (see :func:_is_untouchable_pid).

Source code in tit/jobs/runner.py
def is_alive(pid: int, create_time: float | None) -> bool:
    """True if *pid* is a live, non-zombie process — and, when *create_time* is known, it's
    still the same process (not pid reuse after the original exited). Never true for pid 1 or
    this server's own pid, regardless of what ``create_time`` claims (see
    :func:`_is_untouchable_pid`)."""
    if _is_untouchable_pid(pid):
        return False
    try:
        proc = psutil.Process(pid)
        if proc.status() == psutil.STATUS_ZOMBIE:
            return False
        if (
            create_time is not None
            and abs(proc.create_time() - create_time) > PID_REUSE_TOLERANCE_S
        ):
            return False
        return True
    except (psutil.NoSuchProcess, psutil.AccessDenied, ValueError):
        return False

terminate_tree async

terminate_tree(pid: int, create_time: float | None = None, *, grace_s: float = DEFAULT_GRACE_S) -> None

Snapshot descendants, SIGTERM the tree, wait up to grace_s, SIGKILL survivors.

Uses asyncio.sleep for the grace period so the caller's event loop keeps servicing other jobs while this one shuts down. Never raises: a pid that's already gone is a no-op. Also a no-op for pid 1 or this server's own pid (ra_14 finding #5, see :func:_is_untouchable_pid) -- a re-attached "running" job can only ever point there via a corrupted or crafted status.json, never a job this server actually spawned itself.

Source code in tit/jobs/runner.py
async def terminate_tree(
    pid: int, create_time: float | None = None, *, grace_s: float = DEFAULT_GRACE_S
) -> None:
    """Snapshot descendants, SIGTERM the tree, wait up to *grace_s*, SIGKILL survivors.

    Uses ``asyncio.sleep`` for the grace period so the caller's event loop keeps servicing other
    jobs while this one shuts down. Never raises: a pid that's already gone is a no-op. Also a
    no-op for pid 1 or this server's own pid (ra_14 finding #5, see :func:`_is_untouchable_pid`)
    -- a re-attached "running" job can only ever point there via a corrupted or crafted
    ``status.json``, never a job this server actually spawned itself.
    """
    if _is_untouchable_pid(pid):
        logger.warning("terminate_tree refused to signal untouchable pid %s", pid)
        return
    try:
        root = psutil.Process(pid)
        if create_time is not None and abs(root.create_time() - create_time) > 2.0:
            return  # pid was reused; nothing to do
    except (psutil.NoSuchProcess, psutil.AccessDenied, ValueError):
        return

    procs = [root] + root.children(recursive=True)
    for proc in procs:
        with _ignore_gone():
            send_terminate(proc)

    deadline = asyncio.get_event_loop().time() + grace_s
    while asyncio.get_event_loop().time() < deadline:
        if not any(_still_running(p) for p in procs):
            return
        await asyncio.sleep(0.2)

    for proc in procs:
        if _still_running(proc):
            with _ignore_gone():
                send_kill(proc)

stop_docker_siblings async

stop_docker_siblings(job_id: str, timeout_s: float = DOCKER_CLI_TIMEOUT_S) -> None

Best-effort docker stop for any container labelled tit.job_id=<id>.

QSIPrep/QSIRecon's DooD builders add this label (TODO.md §2.3); a job that never spawned a sibling container, or a host without a docker CLI, makes this a silent no-op. Every call is bounded by timeout_s so an unresponsive daemon degrades to a logged warning.

Source code in tit/jobs/runner.py
async def stop_docker_siblings(
    job_id: str, timeout_s: float = DOCKER_CLI_TIMEOUT_S
) -> None:
    """Best-effort ``docker stop`` for any container labelled ``tit.job_id=<id>``.

    QSIPrep/QSIRecon's DooD builders add this label (TODO.md §2.3); a job that never spawned a
    sibling container, or a host without a docker CLI, makes this a silent no-op. Every call is
    bounded by *timeout_s* so an unresponsive daemon degrades to a logged warning.
    """
    out = await _run_docker(
        ["ps", "-q", "--filter", f"label=tit.job_id={job_id}"], timeout_s
    )
    if not out:
        return
    ids = out.decode().split()
    if not ids:
        return
    await _run_docker(["stop", *ids], timeout_s)

stop_docker_siblings_via_engine async

stop_docker_siblings_via_engine(job_id: str, timeout_s: float = DOCKER_CLI_TIMEOUT_S) -> int

Stop this job's sibling containers through the bounded Engine-API client.

Used by startup reconciliation (JobManager._reconcile_all), which must not shell out: at startup a docker CLI may not be on PATH at all, and every step before the server accepts requests has to be hard-bounded. :func:stop_docker_siblings (the CLI form) stays for the cancel path, whose behaviour and tests predate this. Never raises — an unreachable or wedged daemon degrades to a logged warning and 0.

Source code in tit/jobs/runner.py
async def stop_docker_siblings_via_engine(
    job_id: str, timeout_s: float = DOCKER_CLI_TIMEOUT_S
) -> int:
    """Stop this job's sibling containers through the bounded Engine-API client.

    Used by startup reconciliation (``JobManager._reconcile_all``), which must not shell out:
    at startup a ``docker`` CLI may not be on PATH at all, and every step before the server
    accepts requests has to be hard-bounded. :func:`stop_docker_siblings` (the CLI form) stays
    for the cancel path, whose behaviour and tests predate this. Never raises — an unreachable
    or wedged daemon degrades to a logged warning and ``0``.
    """
    try:
        return await asyncio.wait_for(
            asyncio.to_thread(_stop_siblings_via_engine_blocking, job_id, timeout_s),
            timeout=timeout_s,
        )
    except (asyncio.TimeoutError, TimeoutError):
        logger.warning(
            "job %s: Docker did not answer within %.1fs; leaving any sibling containers alone",
            job_id,
            timeout_s,
        )
        return 0
    except Exception as exc:  # daemon absent, socket refused, API error
        logger.debug("job %s: could not stop sibling containers: %s", job_id, exc)
        return 0