Skip to content

scheduler

tit.jobs.scheduler

Admission logic (TODO.md §2.4): pure decision function plus budget discovery.

evaluate() takes a snapshot of the world (the candidate job, every other job the manager knows about, currently held locks, and the running/budget cost totals) and returns one :class:Decision — never touches the filesystem or a clock other than what it's handed, so the "budget fan-out of 4 fake jobs on an 8-cpu budget" and "lock waits with waiting_on" scheduler tests can drive it directly with synthetic state.

:class:tit.jobs.manager.JobManager owns the loop that calls this once per queued job per tick and acts on the result (spawn, mark skipped, or leave queued with the reported waiting_on/ budget_wait).

One job per product (maintainer rule, 2026-09-22): at most one job of each product in :data:PRODUCT_OF runs at a time -- one Preprocess job, one Simulator job, one Optimizer job, one Analyzer job -- whether it came from one multi-subject group or from separate submissions. The count is read off each job's live status in jobs, so it survives a server restart. Parallelism inside a job is the job's own architecture, under the global CPU limit (:func:discover_budget).

evaluate

evaluate(job: JobSpec, jobs: dict[str, JobStatus], current_holders: list[dict[str, Any]], *, running_cost: Cost, budget: Cost) -> Decision

Admission decision for one queued job (does not mutate anything).

Source code in tit/jobs/scheduler.py
def evaluate(
    job: JobSpec,
    jobs: dict[str, JobStatus],
    current_holders: list[dict[str, Any]],
    *,
    running_cost: Cost,
    budget: Cost,
) -> Decision:
    """Admission decision for one queued *job* (does not mutate anything)."""
    dep_decision = _dependency_state(job, jobs)
    if dep_decision is not None:
        return dep_decision

    requests = locks.keys_for(job.kind, job.subject_ids, job.config)
    blocking = locks.match_conflicts(requests, current_holders, self_job_id=job.id)
    if blocking:
        waiting_on = [
            WaitingOn(
                key=holder.get("key", holder.get("resource", "?")),
                job_id=holder["job_id"],
            )
            for holder in blocking
        ]
        return Decision(waiting_on=waiting_on)

    product = PRODUCT_OF.get(job.kind)
    if product is not None:
        busy = next(
            (
                st
                for st in jobs.values()
                if st.state == "running"
                and st.id != job.id
                and PRODUCT_OF.get(st.kind) == product
            ),
            None,
        )
        if busy is not None:
            return Decision(
                budget_wait=(
                    f"waiting for {product}: one {product} job runs at a time "
                    f"({busy.id} is running)"
                )
            )

    projected = running_cost + job.cost
    if job.kind != "viewer" and not projected.fits(budget):
        return Decision(
            budget_wait=(
                f"waiting for budget: need {job.cost.cpus:g} cpu / {job.cost.mem_gb:g} GB, "
                f"{running_cost.cpus:g} cpu / {running_cost.mem_gb:g} GB already running, "
                f"budget is {budget.cpus:g} cpu / {budget.mem_gb:g} GB"
            )
        )

    return Decision(admit=True)

cascade_skip

cascade_skip(job_id: str, jobs: dict[str, JobStatus], edges: dict[str, list[str]]) -> list[str]

Every job (transitively) depending on job_id that is still queued, for cascading a skip: edges maps a job id to the ids of jobs that list it in after.

Source code in tit/jobs/scheduler.py
def cascade_skip(
    job_id: str, jobs: dict[str, JobStatus], edges: dict[str, list[str]]
) -> list[str]:
    """Every job (transitively) depending on *job_id* that is still queued, for cascading a
    skip: *edges* maps a job id to the ids of jobs that list it in ``after``.
    """
    skipped: list[str] = []
    frontier = [job_id]
    seen = {job_id}
    while frontier:
        current = frontier.pop()
        for dependant in edges.get(current, []):
            if dependant in seen:
                continue
            seen.add(dependant)
            status = jobs.get(dependant)
            if status is not None and status.state == "queued":
                skipped.append(dependant)
            frontier.append(dependant)
    return skipped

build_after_edges

build_after_edges(specs: dict[str, JobSpec]) -> dict[str, list[str]]

job_id -> [ids of jobs that name it in their own after list] (reverse of after).

Source code in tit/jobs/scheduler.py
def build_after_edges(specs: dict[str, JobSpec]) -> dict[str, list[str]]:
    """``job_id -> [ids of jobs that name it in their own after list]`` (reverse of ``after``)."""
    edges: dict[str, list[str]] = {}
    for job_id, spec in specs.items():
        for dep_id in spec.after:
            edges.setdefault(dep_id, []).append(job_id)
    return edges

discover_budget

discover_budget() -> Cost

The pool every admitted job shares (TODO.md §2.4).

CPUs are the user's global CPU limit (:func:tit.cpu.cpu_limit, 70 % of the container's cores by default, set on the Settings page), never the whole container: TI-Toolbox shares the machine with the user's own work. RAM is the container's cgroup limit, clamped by currently-available memory.

Source code in tit/jobs/scheduler.py
def discover_budget() -> Cost:
    """The pool every admitted job shares (TODO.md §2.4).

    CPUs are the user's global CPU limit (:func:`tit.cpu.cpu_limit`, 70 % of the container's
    cores by default, set on the Settings page), never the whole container: TI-Toolbox shares
    the machine with the user's own work. RAM is the container's cgroup limit, clamped by
    currently-available memory.
    """
    mem_gb: float
    try:
        from tit.pre.qsi.utils import get_inherited_dood_resources

        _, cgroup_mem_gb = get_inherited_dood_resources()
        mem_gb = float(cgroup_mem_gb)
    except (
        Exception
    ):  # pragma: no cover - defensive; qsi utils is another lane's module
        mem_gb = 8.0
    try:
        available_gb = psutil.virtual_memory().available / (1024**3)
        mem_gb = min(mem_gb, available_gb)
    except (psutil.Error, OSError):
        pass
    return Cost(cpus=float(cpu_limit()), mem_gb=max(mem_gb, 1.0))