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