Execution and task state¶
Backend authors and orchestration code should import from altar.execution.
JobQueue and get_job_queue_registry are provisional: Altar ships no
job-queue adapter or consumer, so they may change in a minor release.
execution
¶
Stable compute, reconciliation, and object-storage contracts, plus the provisional job-queue contract.
These names are the supported boundary for execution-backend authors and orchestration code. The module forwards to existing implementations so established deep imports keep working without becoming public API.
JobQueue and get_job_queue_registry are provisional and listed in __provisional__: they remain
importable, but they are excluded from the stable-API commitment until the contract has a public adapter, an
in-repository consumer, documentation, and conformance coverage of that adapter.
Capacity
dataclass
¶
A backend's optional report of free slots per named pool, e.g. {"normal": 87, "a100": 4}.
It is optional: autoscaling backends report {}. altar never reads this — only a deployment's
submission driver does, for its own slot math. It is part of the interface so that driver can ask, not
because altar schedules.
free_for
¶
free_for(resources: ResourceRequest) -> int
Return the free slots in the pool that would serve resources, or 0 if the pool is unknown.
gpu=None maps to the CPU/default pool key "normal"; any GpuPool value looks up that pool name
directly (GpuPool.NORMAL_GPU looks up "gpu"). The pool keys a backend advertises in pools must
match the gpu names its tasks request, or a GPU task reads 0.
ContainerTaskSpec
dataclass
¶
ContainerTaskSpec(
kind: str,
image: str,
command: list[str],
labels: dict[str, str] = dict(),
inputs: list[Transfer] = list(),
outputs: list[Transfer] = list(),
resources: ResourceRequest = ResourceRequest(),
timeout_s: int | None = None,
avoid_hosts: tuple[str, ...] = (),
plugin_identity: PluginRunIdentity | None = None,
)
What to run, described without reference to any particular backend.
labels is the structured identity used downstream for grouping and readiness: {job_id, model_id, fold,
allele, ...}. The backend copies it onto the returned TaskHandle and stores it in the ledger, so
grouping and readiness read it directly instead of parsing the task's name. command is the container
argv. It carries model knowledge and is produced by the model plugin's build_scoring_plan (for
ChromBPNet, a model runtime CLI).
timeout_s is the task's wall-clock limit in seconds, counted from when the backend starts it; None
means no limit. Every reference backend enforces it: Kubernetes as the Job's activeDeadlineSeconds
(which also counts time spent scheduling the pod and pulling its image), Modal as the sandbox timeout, and
the local backend by killing the container on the first poll after the deadline. A platform that caps task
lifetime runs a None task for its maximum instead of a shorter default, and rejects a longer explicit
limit (Modal: 24 hours). A task that runs out of time fails with reason "timeout".
avoid_hosts names cluster hosts the task must not be placed on. A driver that has observed a host
failing tasks (a wedged accelerator, a broken image cache) passes it here so the submission lands
elsewhere. It is advisory scheduling input, not identity, so it is deliberately kept out of labels. No
reference backend acts on it yet; they accept and ignore it.
ExecutionBackend
¶
Bases: ABC
Runs a ContainerTaskSpec to completion and reports its terminal state.
Adapters register under altar.execution_backends and are looked up by name; name is the registry key
("local", "modal", "kubernetes"). Only submit and poll are required. Capacity, staging,
cancellation, and recovery have working defaults, so a managed backend implements almost nothing.
input_digest_verification declares where the backend checks an input's Transfer.digest:
"staging":stage_inputsraisesTransferIntegrityErrorbefore any task is submitted."task": each task checks its own inputs before its command starts and fails with reason"input_digest_mismatch". A backend whose inputs are already on a shared volume uses this, because the bytes the task reads are only reachable from inside the task.None(the default): the backend does not verify digests.ExecutionBackendContractrejects this, because content-addressed model weights must be verified on every backend.
capacity
async
¶
capacity() -> Capacity
Report free slots per pool for a deployment's submission driver to use in its slot math. altar
never reads this, so the default returns an empty report. A backend overrides it only to expose fixed
pools to a driver; autoscaling backends keep the default (Capacity()).
stage_inputs
async
¶
stage_inputs(transfers: list[Transfer]) -> None
Make each input readable at its logical_path before the container runs. The default does nothing:
managed backends satisfy inputs through submit, as function arguments or a mount. Backends that stage
files override this.
collect_outputs
async
¶
collect_outputs(transfers: list[Transfer]) -> None
Save each output's logical_path to its uri after the container succeeds. The default does
nothing.
path_resolver
¶
path_resolver() -> PathResolver
Return the backend's logical-to-physical path mapping. The default is identity (empty mount root); a backend that mounts platform storage under a root overrides this.
submit
abstractmethod
async
¶
submit(spec: ContainerTaskSpec) -> TaskHandle
Launch spec and return an opaque handle carrying spec.labels. This returns once the task is
accepted, not when it finishes. A backend must store a structured task record (the ledger) rather than
rely on the handle string, because the external id may not carry the full identity.
Submitting a spec that was submitted before must not raise, because a driver retries a failed or interrupted run by submitting the same specs again. A backend that mints a fresh id per task simply launches another task. A backend that names tasks by their spec, as the Kubernetes backend does, returns a handle to the earlier task while that task is pending or running, so the work never runs twice at once, and replaces it with a new run once it has finished, whether it succeeded or failed. A finished task is never taken as the new run's result: a spec names its inputs by path, and their content may have changed since.
poll
abstractmethod
async
¶
poll(handles: list[TaskHandle]) -> list[TaskStatus]
Probe handles and return their current states without blocking. This is all the reconcile driver
and the readiness checks need. The kubernetes backend lists terminal states in bulk; the modal backend
uses FunctionCall.get(timeout=0); the local backend uses docker inspect.
await_completion
async
¶
await_completion(
handles: list[TaskHandle],
) -> list[TaskStatus]
Block until every handle reaches a terminal state, then return the final statuses.
The base class does not implement this. How often to poll is a driver's choice, so altar does not
bake a busy-wait loop in here: a deployment that needs blocking completion supplies its own loop, or a
backend with a native wait (Modal's future.get()) overrides this. A backend that does neither inherits
the CapabilityNotSupported error below.
cancel
async
¶
cancel(handle: TaskHandle) -> None
Best-effort stop or delete the task the handle refers to, addressed by its external_id. A
completion driver calls this to free a finished job's backend resources once its group has fanned in.
It is idempotent: a task that is already gone counts as cancelled, so a missing task never raises.
altar never calls this; a driver does. Whether and when to reclaim resources is a driver's choice,
like recover. The default does nothing, which suits a managed backend whose tasks are auto-reaped or
one with no delete primitive (the modal backend keeps it). The kubernetes backend overrides it to
delete the Job by external_id, and the local backend to docker rm --force its container (which
otherwise lingers, since local containers run without --rm).
recover
async
¶
recover(
failures: list[TaskStatus], *, banned: set[str]
) -> list[TaskHandle]
Restart tasks that failed transiently — for example reschedule them off banned bad hosts — and
return the new handles. The default does nothing. None of the reference backends (local, modal,
kubernetes) overrides it, so a driver that needs restarts resubmits the task's spec itself.
altar never calls this; a driver does. The recovery policy — which failures to retry, the restart
budget, when to ban a host — is a driver's choice. The reconcile driver dead-letters a failed group
instead of retrying it, whereas a driver that wants restarts calls recover and reads
deterministic_failure to tell an application error, which will fail again, from transient
infrastructure trouble, which is worth a restart. A backend implements this for such drivers; a driver
that does not retry leaves the default alone.
GpuPool
¶
Bases: StrEnum
Names for the GPU pools a task can ask for in ResourceRequest.gpu.
ResourceRequest.gpu holds a pool name, not a hardware model. The name identifies a class of machine
that the backend maps to real hardware: the kubernetes backend maps it to a node pool, and the modal
backend maps it to an accelerator type. That mapping is the backend's job, so a plugin only names the
pool it needs. Because StrEnum members are strings, ResourceRequest.gpu stays typed as str | None
and a third-party backend can define its own pool names; these members give the shipped plugins and
backends one agreed spelling.
NORMAL_GPU("gpu"): the default GPU pool.gpu=Nonemeans CPU-only, which is different from the normal pool.A100("a100"): the A100 pool, a higher-end class a driver can switch to under load.
PathResolver
¶
Turns a logical path into the physical path a backend's container reads or writes.
A logical path is the platform's storage-agnostic file layout: jobs/$id/..., models/$id/...,
genomes/$g/.... That layout is the same for every backend. The only per-backend difference is the mount
root, so this base resolver just prepends it. The kubernetes backend prepends /mnt/volume; the local
backend prepends its bind-mount root; the default (empty root) returns the path unchanged, which a backend
that addresses storage by URI can inherit. Override resolve only when the mapping is not a simple
prefix.
resolve
¶
Return the physical path for logical_path under this backend's mount root.
With an empty mount root the path is returned unchanged. Otherwise the joined result is run through
posixpath.normpath so it is canonical: any ./ segment a logical path carries (for example from
os.path.join(Path(""), ...)) is collapsed. Without this the container argv would read
/mnt/volume/./jobs/..., which is the same file but not the same string as /mnt/volume/jobs/....
The empty-root case is left alone because a backend may pass a gs://... URI through here, and
normpath would collapse the // in the scheme.
ResourceRequest
dataclass
¶
The compute a task needs. A backend uses it to pick a pool or template.
gpu names a pool or accelerator class the backend understands: None means CPU-only, and any other
value is a GpuPool name such as "gpu" (the normal GPU pool) or "a100" (the A100 pool). See
GpuPool for why this stays an open str rather than a closed Literal. count is how many devices of
that pool one task needs. It must be at least 1, and every reference backend requests exactly count
devices; there is no "all devices" value, because a task that silently received a different device count
on each backend would not be portable. count is ignored for a CPU-only task. This type is small on
purpose; a backend reads only the fields it needs.
TaskFailure
dataclass
¶
TaskFailure(
reason: str,
detail: str | None = None,
deterministic: bool = False,
exit_code: int | None = None,
)
Why a task failed, described so a driver can act on it and an operator can read it.
reason is a short stable code meant for grouping and alerting; detail is the platform's own message,
meant for a human. Keeping them apart matters because the two are used differently: a driver branches on
reason, a person reads detail. The reference backends report the same code for the same event:
exit_<n>: the task's command exited with coden(also set inexit_code).timeout: the task ran longer thanContainerTaskSpec.timeout_s.oom_killed: the platform killed the task for exceeding its memory limit.input_digest_mismatch: a content-addressed input was missing or did not matchTransfer.digestwhen the task checked it (backends whoseinput_digest_verificationis"task").container_missing: the backend can no longer find the task, for example because it was deleted.job_failed: the platform reported failure without enough detail for a finer code.
A third-party backend may add its own codes (for example "admission_rejected" or "image_pull") but
should reuse these for the same events.
deterministic separates an application error, which will fail again on retry, from a transient
infrastructure failure such as a bad host or a wedged accelerator, which is worth a recover. It is the
single most consequential field: a transient failure misread as deterministic is abandoned after one
attempt, and a deterministic one misread as transient burns the whole restart budget rediscovering the
same error. When the evidence is ambiguous, False is the safe answer — a wasted retry costs compute, a
wrongly-abandoned shard costs the user their result.
A backend fills this in from whatever failure detail its platform exposes. The thresholds — how many restarts, when to quarantine a host, when to tell the user — are a driver's policy, not the backend's.
TaskHandle
dataclass
¶
TaskHandle(
backend: str,
external_id: str,
labels: dict[str, str] = dict(),
plugin_identity: PluginRunIdentity | None = None,
)
An opaque backend reference to a submitted task, plus its structured identity.
external_id is whatever the backend uses to find the task again: a Kubernetes job name, a Modal
FunctionCall id, a container id. This module treats it as opaque. The task's identity lives in labels,
copied from the spec when it is submitted, so the completion and grouping code never parses external_id.
A model task's exact plugin ABI identity is copied separately in plugin_identity; it is not flattened into
scheduler labels.
TaskState
¶
Bases: StrEnum
The state of a submitted task, normalized to the same four values across backends.
Each member is its own lowercase string, so a state serializes to its plain value without reading
.value.
TaskStatus
dataclass
¶
TaskStatus(
handle: TaskHandle,
state: TaskState,
failure: TaskFailure | None = None,
)
The result of probing one TaskHandle. failure is set only when state is FAILED.
Transfer
dataclass
¶
Transfer(
uri: str,
logical_path: str,
locality_key: str | None = None,
digest: str | None = None,
)
A file that must move between platform storage and the task's filesystem.
uri is a storage location the platform Storage understands (gs://..., file://...,
http(s)://...). logical_path is where the container reads or writes the file; the backend's
PathResolver turns it into a physical path. Direction comes from which list on ContainerTaskSpec the
Transfer sits in:
- an input
Transfermeans: makeurireadable atlogical_pathbefore the container runs; - an output
Transfermeans: savelogical_pathtouriafter the container succeeds.
The requirement is about the state at run time, not about which method does the work. Backends that stage
files do it in stage_inputs and collect_outputs. Managed backends do it through submit, as function
arguments or a mount, and it costs nothing when the backend's own storage is already the platform
Storage.
digest is the expected SHA-256 content identity for an input: sha256: plus the hex SHA-256 of the
bytes of one regular file, the same value sha256sum prints. It never covers a directory. A resource made
of several files is shipped as one archive (a tar of a model directory, for example) whose bytes carry the
digest, as ResourceReference requires a bundle digest for multi-file selections. A digest-bearing input
that resolves to a directory is rejected the same way as a mismatch. Every reference backend verifies it
before the task's command runs; ExecutionBackend.input_digest_verification says where. It is
independent of uri, so an input mapper may relocate the same resource without changing scientific
identity. Output transfers normally leave it unset.
locality_key is an optional co-scheduling hint. Transfers that share a key should, where a backend can
arrange it, run on the same worker so a warm local cache is reused across them — for example a cluster
backend can route them to one worker so a model's tar and peaks files stay cached. Without a key, sharding
is based on the path. It is only a hint: autoscaling backends ignore it and correctness never depends on
it.
TransferIntegrityError
¶
Bases: DigestMismatchError
Staged bytes do not match the content identity declared by a transfer.
It subclasses altar_identity.DigestMismatchError, which model runtimes raise for the same failure, so one
except DigestMismatchError catches a mismatch wherever it was detected. Both are ValueErrors.
K8sJobExistsError
¶
Bases: Exception
Raised by KubernetesRuntime.create_job when a Job with the manifest's name already exists in the
namespace (the API server's 409 AlreadyExists).
KubernetesExecutionBackend.submit catches it to adopt or replace the existing Job, so a runtime must
raise this, not its SDK's own error, for a name conflict. submit raises it too when an active Job of that
name was rendered from a different manifest, or when the conflict does not settle within about a minute.
K8sJobInfo
dataclass
¶
K8sJobInfo(
name: str,
state: TaskState,
labels: dict[str, str] = dict(),
annotations: dict[str, str] = dict(),
uid: str | None = None,
)
The parts of a K8s Job that the backend needs: its name, normalized state, labels, annotations, and uid.
Returned by list_jobs and read_job so the backend never handles a raw V1Job; the SDK type stays behind
the runtime boundary. state is normalized to a TaskState by the runtime, read from the Job's
conditions. labels are the pod-template and Job labels that carry identity. annotations are the Job's
annotations; submit reads MANIFEST_DIGEST_ANNOTATION from them before it adopts a Job. uid is the
server-assigned id of this particular Job object, which submit uses so that replacing a finished Job never
deletes a newer Job of the same name. A runtime that does not report it leaves it None.
KubernetesExecutionBackend
¶
KubernetesExecutionBackend(
*,
namespace: str | None = None,
runtime: KubernetesRuntime | None = None,
mount_root: str = CONTAINER_MOUNT_ROOT,
pvc_name: str | None = None,
image_pull_policy: str = "IfNotPresent",
backoff_limit: int = 0,
run_as_user: int | None = None,
run_as_group: int | None = None,
fs_group: int | None = None,
pool_node_selectors: dict[str, dict[str, str]]
| None = None,
max_slots: dict[str, int] | None = None,
digest_verifier_image: str = DEFAULT_DIGEST_VERIFIER_IMAGE,
trust_verification_records: bool = True,
)
Bases: ExecutionBackend
Runs a ContainerTaskSpec as a Kubernetes Job on any cluster.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
namespace
|
str | None
|
cluster namespace Jobs are created in. Defaults to |
None
|
runtime
|
KubernetesRuntime | None
|
Kubernetes API facade (see |
None
|
mount_root
|
str
|
pod-internal mount path for the logical layout. Defaults to |
CONTAINER_MOUNT_ROOT
|
pvc_name
|
str | None
|
PersistentVolumeClaim mounted at |
None
|
image_pull_policy
|
str
|
container |
'IfNotPresent'
|
backoff_limit
|
int
|
Job |
0
|
run_as_user
|
int | None
|
Optional pod UID for writable-volume compatibility. |
None
|
run_as_group
|
int | None
|
Optional pod primary GID. |
None
|
fs_group
|
int | None
|
Optional supplemental volume GID. Set these explicitly for root-squashed or restricted PVCs. |
None
|
pool_node_selectors
|
dict[str, dict[str, str]] | None
|
maps a pool name to |
None
|
max_slots
|
dict[str, int] | None
|
per-pool ceiling for |
None
|
digest_verifier_image
|
str
|
image for the |
DEFAULT_DIGEST_VERIFIER_IMAGE
|
trust_verification_records
|
bool
|
whether |
True
|
path_resolver
¶
path_resolver() -> PathResolver
The pod's view of a logical path: <mount_root>/<logical>, the mounted volume. Every backend's
resolver returns this same shape, so a spec's command runs unchanged on any of them.
submit
async
¶
submit(spec: ContainerTaskSpec) -> TaskHandle
Render spec to a Job manifest and create it without waiting. The created Job name is the opaque
external_id; identity lives in labels, stamped as K8s labels, never parsed from the name.
The verify-inputs init container is added here, after build_manifest, so it mounts every volume
the task's container mounts, including volumes a subclass adds in its own build_manifest.
The manifest's hash is stamped last, in the MANIFEST_DIGEST_ANNOTATION annotation, so it covers
everything the Job will run with.
When a Job with the same name already exists, which by job_name means the same spec, submit does not
raise. It returns a handle to that Job while it is pending or running and its manifest hash matches, so
the same work never runs twice at once. It raises K8sJobExistsError if an active Job's manifest hash
differs. A finished Job, succeeded or failed, is deleted with its pods and created again, because a
spec names inputs by path and their content may have changed since.
job_name
¶
job_name(spec: ContainerTaskSpec) -> str
Candidate Job name for spec: <kind>-<short spec hash>, safe under DNS-1123 and at most 63
characters.
The hash covers everything spec_to_dict records: image, command, labels, transfers, resources,
timeout, and plugin identity. Resubmitting the same spec therefore meets the Job the first submission
created, which submit adopts while it is active or replaces once it has finished, so the same work
never runs twice at once; a spec whose image or configuration changed gets a new Job. avoid_hosts is left out, as it is from
spec_to_dict, because it is placement rather than identity. A subclass can override this with its own
naming convention; submit treats two specs that share a name as the same work.
build_manifest
¶
build_manifest(spec: ContainerTaskSpec) -> dict[str, Any]
Build a standard batch/v1 Job manifest from spec: image, argv, resources, labels, deadline, and
the shared volume mount. This is the main hook to override: a subclass extends it to add
cluster-specific pod config such as node affinity, tolerations, or sidecars.
poll
async
¶
poll(handles: list[TaskHandle]) -> list[TaskStatus]
Map each handle's Job name (external_id) to its state via one bulk namespace list. Identity is in
handle.labels; the list only supplies the state. A failed Job is explained by the runtime's
describe_failure. A handle whose Job is absent from the list is a terminal failure with reason
"container_missing": a created Job is listed at once, so an absent one was deleted and can never
finish.
cancel
async
¶
cancel(handle: TaskHandle) -> None
Delete the Job named by handle.external_id and its pods. Idempotent: delete_job is best-effort,
so a Job that is already gone is a no-op rather than an error.
KubernetesRuntime
¶
Bases: ABC
The part of the Kubernetes API this backend uses, defined as an interface so the backend can be tested
without kubernetes_asyncio or a cluster. The default _K8sApiRuntime implements it against the real SDK;
tests pass a fake that records calls and returns programmed job lists.
Every method takes and returns plain values (a manifest dict, job names, a namespace, a label selector).
The real implementation builds the BatchV1Api and CoreV1Api clients internally, so no Kubernetes type
crosses this boundary.
create_job
abstractmethod
async
¶
Create the Job described by manifest in namespace and return its server-visible name. This equals
manifest["metadata"]["name"] unless the server generates one.
The Job must keep manifest["metadata"]["annotations"], which list_jobs and read_job report back.
Raises K8sJobExistsError when a Job with that name already exists.
list_jobs
abstractmethod
async
¶
list_jobs(
*, namespace: str, label_selector: str | None = None
) -> list[K8sJobInfo]
List all Jobs in namespace without blocking, optionally filtered by label_selector, each with its
normalized TaskState, labels, annotations, and, where the runtime can report it, uid. poll and
capacity both build on this one primitive.
read_job
async
¶
read_job(name: str, *, namespace: str) -> K8sJobInfo | None
Return the Job name in namespace with its normalized state, annotations, and uid, or None when no
such Job exists or it is already being deleted, since such a Job can be neither adopted nor replaced yet.
submit calls this after create_job reports a name conflict. The default finds the Job in
list_jobs, so a runtime written before this method existed keeps working; the real runtime reads the
one Job directly and also reports a Job that is being deleted as None.
delete_job
abstractmethod
async
¶
Delete the Job name and its pods in namespace. Best-effort: a Job that is already gone is not an
error. Used by cancellation, by submit to replace a finished Job, and by recovery and cleanup.
With uid, delete the Job only if it is still the object with that uid (a delete precondition), and do
nothing if the name now belongs to a newer Job. submit passes uid only when read_job reported one,
so a runtime that never reports a uid need not accept it.
describe_failure
async
¶
describe_failure(
name: str, *, namespace: str
) -> TaskFailure
Explain why the failed Job name failed, using the stable reason codes listed on TaskFailure.
The default wraps the coarser diagnose and reports reason "job_failed", so a runtime that only
implements diagnose keeps working. The real runtime reads the Job's failure condition and its pods'
container states to report timeout, input_digest_mismatch, oom_killed, or exit_<n>.
diagnose
async
¶
Classify a failed Job as (deterministic_failure, error): an application error that won't pass on
retry, versus transient infrastructure trouble (OOM, eviction, preemption).
Only the default describe_failure calls this. The default returns (True, None), treating every
failure as deterministic. Prefer overriding describe_failure, which can also report a reason code.
LocalDockerExecutionBackend
¶
LocalDockerExecutionBackend(
*,
data_root: str | None = None,
slots: int = 4,
storage: Storage | None = None,
container_user: str | None = None,
run_command: CommandRunner | None = None,
)
Bases: ExecutionBackend
Runs a ContainerTaskSpec as a detached Docker container on the local host.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
data_root
|
str | None
|
Host directory bind-mounted as |
None
|
slots
|
int
|
Number of free slots that |
4
|
storage
|
Storage | None
|
|
None
|
container_user
|
str | None
|
Optional Docker |
None
|
run_command
|
CommandRunner | None
|
Injectable async |
None
|
capacity
async
¶
capacity() -> Capacity
Report a single "normal" pool with slots free slots. A local host has no GPU pools.
path_resolver
¶
path_resolver() -> PathResolver
Return the container view of a logical path: /mnt/volume/<logical>. Every backend uses the same
root, so a spec's command refers to the same paths on any of them.
stage_inputs
async
¶
stage_inputs(transfers: list[Transfer]) -> None
Copy each input into the host data_root at its mirrored location, so it is readable at its
logical path inside the container. Transfers run one at a time.
An input with a digest is verified after it is copied, so the check covers the bytes the container
will read. A copy that does not match is deleted before TransferIntegrityError is raised, so no
later task can read it. A copy that matches gets a verification record (see verify_file_digest).
An input whose destination already holds an unchanged, recorded copy of the same digest is not copied
or hashed again.
collect_outputs
async
¶
collect_outputs(transfers: list[Transfer]) -> None
Copy each output from the host data_root back to its uri through the Storage adapter.
submit
async
¶
submit(spec: ContainerTaskSpec) -> TaskHandle
Start the spec with docker run -d and return a TaskHandle carrying its labels. The container
id becomes the external_id. A task's identity lives in labels, written as --label pairs, not
in a container name.
poll
async
¶
poll(handles: list[TaskHandle]) -> list[TaskStatus]
Check each container's state with docker inspect and return one TaskStatus per handle. This
does not wait for the containers. A container past its timeout_s is killed here.
cancel
async
¶
cancel(handle: TaskHandle) -> None
Reap the container named by handle.external_id with docker rm --force.
Local containers run without --rm — deliberately, so poll can read a finished container's exit
code with docker inspect after it exits. The cost is that an exited container lingers on the host
until something removes it. A completion driver calls cancel once a task's group has fanned in, to
reclaim it; without this override the base no-op would leak every container. docker rm --force
both stops a still-running container and removes an exited one, so one call covers either state.
Idempotent: a container that is already gone makes docker rm exit non-zero, which is ignored so a
missing task never raises (matching the base contract).
ModalExecutionBackend
¶
ModalExecutionBackend(
*,
app_name: str | None = None,
volume_name: str | None = None,
storage: Storage | None = None,
runtime: ModalRuntime | None = None,
pool_accelerators: dict[str, str] | None = None,
)
Bases: ExecutionBackend
Runs a ContainerTaskSpec as a Modal Sandbox on Modal's managed serverless platform.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
app_name
|
str | None
|
Modal App the sandboxes attach to, looked up or created on first use. Defaults to
|
None
|
volume_name
|
str | None
|
Modal Volume mounted at |
None
|
storage
|
Storage | None
|
|
None
|
runtime
|
ModalRuntime | None
|
Modal SDK facade (see |
None
|
pool_accelerators
|
dict[str, str] | None
|
Maps a pool name to a Modal accelerator. Defaults to |
None
|
capacity
async
¶
capacity() -> Capacity
Report no fixed capacity. Modal autoscales, so there are no fixed slots to ration. A backend that
sees an empty Capacity submits without gating.
path_resolver
¶
path_resolver() -> PathResolver
The sandbox's view of a logical path: /mnt/volume/<logical>, the Modal Volume mount. Every backend
returns this same shape, so a spec's command runs unchanged on any of them.
stage_inputs
async
¶
stage_inputs(transfers: list[Transfer]) -> None
Make each input readable at its logical path inside the sandbox. For each transfer, pull its uri
from storage to a host temp file, then upload that into the Modal Volume at the mirrored path.
Transfers run one at a time.
collect_outputs
async
¶
collect_outputs(transfers: list[Transfer]) -> None
Write each output from the Modal Volume back to its uri. For each transfer, download it to a host
temp file, then hand that file to storage.
submit
async
¶
submit(spec: ContainerTaskSpec) -> TaskHandle
Launch spec as a Modal Sandbox and return a handle carrying its labels. The sandbox id is the
opaque external_id; identity lives in labels, stamped as Modal tags, never in a name.
A spec without a timeout_s gets Modal's maximum sandbox lifetime, MODAL_MAX_SANDBOX_TIMEOUT_S.
Raises:
| Type | Description |
|---|---|
ValueError
|
If |
poll
async
¶
poll(handles: list[TaskHandle]) -> list[TaskStatus]
Check each sandbox's state with Modal's non-blocking Sandbox.poll().
ModalRuntime
¶
Bases: ABC
The part of the Modal SDK this backend uses, defined as an interface so the backend can be tested
without the modal package or an account. The default _ModalSdkRuntime implements it against the real
SDK; tests pass a fake that records calls.
Every method takes and returns plain values (image refs, argv, string paths, the app and volume names).
The real implementation builds the Modal App, Image, Volume, and Sandbox objects internally, so
no Modal type crosses this boundary.
run_sandbox
abstractmethod
async
¶
run_sandbox(
*,
app_name: str,
volume_name: str,
mount_root: str,
image: str,
command: list[str],
gpu: str | None,
memory_mb: int | None,
timeout_s: int,
labels: dict[str, str],
) -> str
Create a detached sandbox running image with command, with volume_name mounted at
mount_root, and return its opaque sandbox id. timeout_s is the sandbox's effective lifetime, always
set: the backend has already mapped a spec without a limit to MODAL_MAX_SANDBOX_TIMEOUT_S.
poll_sandbox
abstractmethod
async
¶
Probe a sandbox without blocking: None if it is still running, else its integer exit code.
A sandbox that ran past its timeout reports TIMEOUT_EXIT_CODE, as the Modal SDK does. Raise
LookupError when no sandbox has this id.
upload_to_volume
abstractmethod
async
¶
Copy a host file into the named Modal Volume at remote_path, relative to the volume root.
download_from_volume
abstractmethod
async
¶
Copy remote_path (relative to the volume root) out of the named Modal Volume to a host file.
GroupSpec
dataclass
¶
GroupSpec(
kind: str,
group_key: GroupKeyFn,
ready_when: ReadyPredicate,
on_ready: OnReady,
on_failed: OnFailed | None = None,
)
How one kind of shard groups and when a group is complete.
group_key buckets shards into groups. ready_when decides when a group has fanned in. on_ready is
the next-stage callback, called once per group with (group_key, succeeded_handles). All of these read
from handle.labels.
on_failed is the optional dead-letter path. When a group is blocked by a terminally-failed shard, so
it can never reach ready_when, reconcile_once calls on_failed(group_key, failed_handles). Returning
True means the group was dead-lettered: it is consumed, and the pipeline drains instead of stalling.
Returning False defers it — the group is left untouched for a later tick. That distinction matters
because "this shard failed" and "this shard is beyond saving" are not the same claim. A driver that can
still relaunch a failed shard, but was unable to this tick (no capacity, a dependency not ready), must be
able to say so; without that, a transient shortage silently becomes permanent abandonment of the group
and of every sibling that had already succeeded.
Leave it None to let the reference loop escalate the stall itself with ReconcileStalled. A driver that
surfaces failures its own way, such as through its own restart step, also leaves it None and reads
stalled_groups instead.
InMemoryTaskLedger
¶
Bases: TaskLedger
In-process reference implementation of TaskLedger for development and tests.
Records are keyed by (backend, external_id), the ledger's unique constraint.
ReconcileResult
dataclass
¶
ReconcileResult(
polled: int = 0,
succeeded: int = 0,
failed: int = 0,
fired_groups: list[str] = list(),
dead_lettered_groups: list[str] = list(),
deferred_groups: list[str] = list(),
stalled_groups: list[str] = list(),
)
Summary of one reconcile_once tick, for a driver to log and to decide whether to keep looping.
ReconcileStalled
¶
Bases: RuntimeError
Raised when a group is blocked by a terminally-failed shard and has no on_failed handler.
Without this, the reference run_reconcile_loop would drain and exit while the group's succeeded
siblings stay stuck. groups carries the stalled "kind:group_key" strings. Direct reconcile_once
callers never see this exception; they get the stalled groups in stalled_groups, and only the
reference loop raises.
TaskLedger
¶
Bases: ABC
Durable record of every submitted task and its latest state.
This module defines the interface and an in-memory reference; a deployment backs it with a database
table. reconcile_once reads and writes only through this interface, so the same driver works over a
fresh in-memory ledger in tests or a persistent database in production.
record
abstractmethod
async
¶
record(
handle: TaskHandle,
kind: str,
spec: dict[str, Any] | None = None,
) -> None
Persist a freshly-submitted task.
Idempotent per handle while that handle's record is still open, so a submit retried after a partial failure never resets a task that is already in flight.
A consumed record is the exception, and the distinction matters. Backends generally derive a
task's name from its identity and take the first free suffix, so work re-submitted after an earlier
attempt was settled and cleaned up arrives under the exact same external_id. Skipping it as
"already recorded" leaves the new task running with nothing in the ledger pointing at it: its group
stays permanently a shard short, and any recovery that looks for units with no open work re-submits
it into the same gap on every pass. So an implementation MUST re-open a consumed record for the new
task — resetting its state, restart count, failure detail, and spec — rather than leave it settled.
open_handles
abstractmethod
async
¶
open_handles() -> list[TaskHandle]
Return handles that are not yet terminal (PENDING or RUNNING) and not consumed — the ones to poll.
set_state
abstractmethod
async
¶
set_state(
handle: TaskHandle,
state: TaskState,
failure: TaskFailure | None = None,
) -> None
Record the latest observed state for a task, and why it failed when it did.
failure is what the backend learned about a FAILED task. Persisting it is the difference between
an operator being able to answer "why did this shard die" from the ledger and having to go ask the
cluster — which, by the time anyone looks, has usually garbage-collected the evidence.
active_succeeded
abstractmethod
async
¶
active_succeeded() -> list[TaskRecord]
Return succeeded, not-yet-consumed records — the candidates reconciliation buckets into groups.
active_failed
abstractmethod
async
¶
active_failed() -> list[TaskRecord]
Return failed, not-yet-consumed records.
These are the mirror of active_succeeded. Reconciliation uses them to detect a group blocked by a
terminally-failed shard, so it can dead-letter the group rather than stall.
consume
abstractmethod
async
¶
consume(handles: list[TaskHandle]) -> None
Mark handles consumed once their group has fired, so the group fires exactly once.
TaskRecord
dataclass
¶
TaskRecord(
handle: TaskHandle,
kind: str,
state: TaskState = PENDING,
restart_count: int = 0,
spec: dict[str, Any] | None = None,
consumed: bool = False,
failure: TaskFailure | None = None,
)
One row of the execution_tasks ledger: the durable record submit() writes.
The handle alone is an opaque backend id and does not carry identity, so the record also persists the
shard's labels and kind, and completion reads full identity off the record. consumed marks a shard
whose group has already fired, so it is not counted again on the next tick.
JobQueue
¶
Bases: ABC
Hands a job off to the compute layer. Each adapter implements one transport (workflow orchestrator, Pub/Sub, local queue).
Provisional: exported for existing host applications, but excluded from the stable-API commitment until a public adapter, an in-repository consumer, and conformance coverage exist (see the module docstring).
name is the key the registry uses to identify an adapter — the altar.job_queues entry-point name such
as "local". Each adapter sets it.
enqueue
abstractmethod
async
¶
Enqueue a job of kind with parameters and return an opaque run handle id.
Returns as soon as the job is accepted, not when it finishes. kind is an open, deployment-owned
identifier (see the module docstring). The adapter resolves it to its own transport target and
should raise on a kind it does not recognize.
LocalStorage
¶
Bases: Storage
A Storage implementation backed by the local filesystem.
upload_from_local_file
async
¶
Copy a local file to the destination path.
When the source and destination resolve to the same file, this returns without doing anything.
That case arises with the local backend when storage and the compute volume are the same
directory, so a file is already where it needs to be. shutil.copy2 raises SameFileError on
such a copy, so this checks for it first rather than failing.
read_csv
async
¶
read_csv(
path: str,
delimiter: str = "\t",
fieldnames: list[str] | None = None,
) -> list[dict[str, Any]]
Read the delimited file directly from disk, skipping the in-memory text round-trip.
write_csv
async
¶
Write the rows directly to disk, skipping the in-memory text round-trip.
as_local_file
async
¶
Yield the path unchanged, since the file is already on the local disk.
Storage
¶
Bases: ABC
A backend-independent interface for reading and writing files.
Subclasses implement the file operations against a specific backend, such as the local disk or GCS.
Each subclass sets name to its registry key ("local", "gcs", and so on), which is how a caller
looks the adapter up in the altar.storage registry.
read_text
abstractmethod
async
¶
Read a text file and return its contents as a string.
read_binary
abstractmethod
async
¶
Read a file and return its contents as bytes.
write_text
abstractmethod
async
¶
Write a string to a text file.
write_binary
abstractmethod
async
¶
Write bytes to a file.
exists
abstractmethod
async
¶
Return whether a file exists at the given path.
makedirs
abstractmethod
async
¶
Create a directory and any missing parent directories.
upload_from_local_file
abstractmethod
async
¶
Copy a file from the local disk into storage by streaming it.
Use this for large files. It streams the upload rather than reading the whole file into memory.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
local_path
|
str
|
Path to the local file to upload. |
required |
dest_path
|
str
|
Destination path in storage. |
required |
as_local_file
abstractmethod
async
¶
Yield a local filesystem path for a stored file, for the duration of the context.
Local storage yields the path unchanged. Remote storage downloads the file to a temporary path and yields that; the temporary file is removed when the context exits.
Usage
async with storage.as_local_file(path) as local_path: with open(local_path) as f: # Process file
read_csv
async
¶
read_csv(
path: str,
delimiter: str = "\t",
fieldnames: list[str] | None = None,
) -> list[dict[str, Any]]
Read a delimited file and return one dict per row, keyed by column name.
write_csv
async
¶
Write a list of rows to a delimited file.
open_text_reader
async
¶
Read a text file and return its contents as an in-memory file object.
open_text_writer
¶
open_text_writer(path: str) -> TextWriter
Return a context manager for writing text to a file (see TextWriter).
TextWriter
¶
TextWriter(storage: Storage, path: str)
An async context manager for writing a text file through a Storage.
Writes made inside the block go to an in-memory buffer. When the block exits without an error, the
buffer is flushed to the file with Storage.write_text. Because that flush is awaited, the writer
only supports async with.
get_execution_backend_registry
¶
get_execution_backend_registry() -> PluginRegistry[
type[ExecutionBackend]
]
Return the shared altar.execution_backends registry.
has_verification_record
¶
Return whether path is a regular file with a current verification record for digest.
The record must name digest and match the file's current size, modification time, status-change time,
and inode number. A missing, unreadable, or stale record returns False. Checking a record reads two
small pieces of metadata and never reads the file's bytes. Symlinks are followed.
spec_from_dict
¶
spec_from_dict(data: dict[str, Any]) -> ContainerTaskSpec
Rebuild a ContainerTaskSpec from spec_to_dict output. Unknown keys are ignored, so a spec stored by
an older version stays loadable after a field is added.
spec_to_dict
¶
spec_to_dict(spec: ContainerTaskSpec) -> dict[str, Any]
Serialize a ContainerTaskSpec to plain JSON-able values, for storage in a ledger row.
Paired with spec_from_dict. The point of persisting a spec is that a task can be re-launched from the
ledger alone, long after the in-memory object that produced it is gone — which is what recovery from a
task that vanished without finishing requires.
verify_file_digest
¶
verify_file_digest(
path: str | PathLike[str],
digest: str,
*,
trust_record: bool = True,
record: bool = True,
) -> bool
Check that path is a regular file whose SHA-256 is digest, reading its bytes only when needed.
Model runtimes and task-side verifiers call this before they use a staged resource. With trust_record,
a current verification record (see has_verification_record) accepts the file without hashing it. A file
staged once and read by many tasks is then hashed once instead of once per task. Otherwise the whole file
is hashed. With record, a successful hash writes a record so later checks can skip the hash.
A record means only that the bytes matched when it was written and that the file has not visibly changed
since. Anyone who can write the storage can also write a record, so a deployment that shares storage with
untrusted writers passes trust_record=False.
Returns True when the file was hashed and False when a record accepted it. Raises
TransferIntegrityError when the path is not a regular file or its bytes do not match. Raises
ValueError when digest is not sha256: plus 64 lowercase hex characters.
This is altar_identity.verify_file, which runtimes call directly, with Altar's error type. Unlike
verify_file, it trusts and writes records by default, because it is the task-side check of files a staging
step already hashed.
verify_files_script
¶
Return a POSIX shell program that verifies files the way verify_file_digest does.
Run it as sh -c <program> <name> <line>..., where each line is <64 hex> <path>: the digest's hex,
two spaces, and the path, as sha256sum -c reads them. For each line the program requires a regular file
(symlinks are followed). It then accepts a current verification record when trust_records is true, or
else hashes the file and writes a record after a match. It checks every line, prints one result per file,
and exits non-zero if any file is missing, is not a regular file, or does not match. A record that cannot
be written never fails the check.
The Kubernetes backend runs this program in its input-verification init container. A host that stages files onto shared storage through a shell can run it after staging so the tasks that read those files accept them from the record.
verify_transfer_digest
¶
verify_transfer_digest(
transfer: Transfer, path: str, *, record: bool = False
) -> None
Verify path against an input transfer's optional SHA-256 content identity.
A digest covers one regular file, so a path that is missing or is a directory raises
TransferIntegrityError, the same error as a mismatch. Symlinks are followed. The bytes are always hashed:
a staging backend calls this at the point where bytes enter the storage its tasks read.
With record=True, a successful check also writes a verification record beside path (see
verify_file_digest). Later checks that trust records then accept the file without reading it again.
Pass it only for a copy the backend owns, never for a caller's source file.
all_expected_folds
¶
Ready when folds 0..N-1 are all present, where N is read from the shards themselves.
Each fold shard carries its fold count in the count_label label, so N comes from the shard rather than
from a constant here. A ModelPlugin stamps its own num_folds onto every fold shard it emits. One
shared GroupSpec then completes a 5-fold model's group at 5 and a 1-shard model's group at 1, even
within the same job. This generalizes all_folds from a fixed n to a per-group n read off the
group's own labels.
default is the fold count used when no shard in the group carries a valid count_label. It is None
by default: core does not guess one architecture's fold count. An unlabeled group without a default
raises ValueError from the predicate instead of pending forever, and reconcile_once propagates it
before firing or consuming any group in that tick. Every in-repository binding stamps the label. When a
group's shards report several different counts, the largest wins, so a partially-labeled group still
waits for its full width.
all_folds
¶
Ready when folds 0..n-1 are all present.
Pass the plugin's own num_folds; there is no default fold count. Prefer all_expected_folds for a
completion rule shared across models. A fixed n here serves only one fold count, so a job mixing a
5-fold model with a 1-shard model would stall or fire early.
all_label_values
¶
Ready when the observed shards cover every expected value of label.
This is the general form of both fan-in rules. A three-member ensemble uses
all_label_values("fold", {"0", "1", "2"}); a per-allele stage uses
all_label_values("allele", {"ref", "alt"}). Label values are strings, so a fold is "3".
count_at_least
¶
Ready when at least n shards have succeeded. A single-shard kind uses count_at_least(1).
group_by
¶
Build a group-key function that joins the given label values with :.
For example, group_by("job_id", "model_id") yields the key "<job_id>:<model_id>", read from a
shard's labels. A missing label contributes an empty segment, so a shard with incomplete labels groups
predictably instead of crashing.
reconcile_once
async
¶
reconcile_once(
backend: ExecutionBackend,
ledger: TaskLedger,
group_specs: dict[str, GroupSpec] | list[GroupSpec],
*,
max_concurrent_callbacks: int | None = None,
) -> ReconcileResult
Run one non-blocking reconciliation tick.
- Poll every open handle from the ledger, in bulk and without blocking.
- Write each observed state back to the ledger.
- Re-derive groups from the ledger's succeeded-and-unconsumed shards, bucketed by each kind's
group_key. For any group whoseready_whennow passes, callon_readyonce and mark that group's shards consumed, so it never fires twice. Callbacks run sequentially by default; callers may setmax_concurrent_callbacksto process independent groups concurrently with a fixed upper bound. - Surface any group blocked by a terminally-failed shard, which can never reach
ready_when. Dead-letter it viaon_failedif the spec supplies one, otherwise report it installed_groups. Its succeeded siblings are never left silently unconsumed.
Readiness predicates are evaluated for every succeeded group before any callback fires. A predicate that
raises, such as all_expected_folds() over shards with no fold-count label, propagates out of this call
after the polled states are written but before any group fires or is consumed, so the misconfiguration
surfaces immediately and the tick can be retried once it is fixed. run_reconcile_loop propagates it too.
Groups are re-derived from the ledger every tick rather than accumulated in memory. A fold that finished
three ticks ago is still counted because its SUCCEEDED record persists until the group fires. Apart
from the ledger, the function keeps no state across calls, so a driver can call it on any cadence and
survive restarts. Retrying a transient failure is a driver's job (see the module docstring): a driver
that retries resubmits a fresh shard and consumes the failed record before the next tick, so a
still-unconsumed failed shard seen here is genuinely terminal and its group cannot complete.
run_reconcile_loop
async
¶
run_reconcile_loop(
backend: ExecutionBackend,
ledger: TaskLedger,
group_specs: dict[str, GroupSpec] | list[GroupSpec],
*,
interval_s: float = 5.0,
max_ticks: int | None = None,
stop_when_idle: bool = True,
) -> list[ReconcileResult]
Reference driver that calls reconcile_once every interval_s.
The loop stops when no open handles remain (if stop_when_idle) or when max_ticks is reached. It is a
simple built-in loop; a deployment that wants its own scheduling, such as an external cron, can ignore it
and call reconcile_once directly. Returns each tick's result.
stop_when_idle ends the loop once a tick sees zero open handles and fires no group, meaning the
pipeline has drained. max_ticks is a safety bound for callers that don't want an unbounded loop.
Raises ReconcileStalled if a tick reports a group blocked by a failed shard with no on_failed
handler. Such a group can never drain, so the loop escalates rather than looping forever or exiting as if
idle. Give the blocking kind's GroupSpec an on_failed dead-letter handler to keep the loop running
past it.
single
¶
Ready as soon as a one-shard group's single shard succeeds. Same as count_at_least(1).
get_job_queue_registry
¶
get_job_queue_registry() -> PluginRegistry[type[JobQueue]]
Return the shared altar.job_queues registry.
Provisional, like the JobQueue contract it resolves; see the module docstring.
get_storage_registry
¶
get_storage_registry() -> PluginRegistry[type[Storage]]
Return the shared registry of storage adapters.