Skip to content

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

Capacity(pools: dict[str, int] = dict())

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_inputs raises TransferIntegrityError before 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. ExecutionBackendContract rejects 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=None means 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

PathResolver(mount_root: str = '')

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

resolve(logical_path: str) -> str

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

ResourceRequest(
    gpu: str | None = None,
    count: int = 1,
    memory_gb: float | None = None,
)

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 code n (also set in exit_code).
  • timeout: the task ran longer than ContainerTaskSpec.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 match Transfer.digest when the task checked it (backends whose input_digest_verification is "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.

deterministic_failure property

deterministic_failure: bool

Whether the failure will recur identically on retry. False when the task did not fail.

error property

error: str | None

The failure's human-readable detail, or None when the task did not fail.

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 Transfer means: make uri readable at logical_path before the container runs;
  • an output Transfer means: save logical_path to uri after 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

K8sJobExistsError(
    name: str, *, namespace: str, detail: str | None = None
)

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 $ALTAR_K8S_NAMESPACE, else default.

None
runtime KubernetesRuntime | None

Kubernetes API facade (see KubernetesRuntime). Defaults to _K8sApiRuntime, which imports kubernetes_asyncio lazily.

None
mount_root str

pod-internal mount path for the logical layout. Defaults to /mnt/volume.

CONTAINER_MOUNT_ROOT
pvc_name str | None

PersistentVolumeClaim mounted at mount_root in every pod. Defaults to $ALTAR_K8S_PVC; if unset, no volume is mounted (a cluster whose Storage is an object store with a CSI or gcsfuse mount needs none, and staging happens elsewhere).

None
image_pull_policy str

container imagePullPolicy. Defaults to IfNotPresent.

'IfNotPresent'
backoff_limit int

Job backoffLimit, the pod-level retries K8s does before marking the Job failed. The default of 0 leaves retries to the platform's reconcile or recover step.

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 nodeSelector labels, so a GPU ResourceRequest lands on the right nodes (for example {"a100": {"cloud.google.com/gke-accelerator": "nvidia-tesla-a100"}}). A pool with no entry just gets the nvidia.com/gpu count and no selector.

None
max_slots dict[str, int] | None

per-pool ceiling for capacity()'s free-slot math, as {pool: max}. If None, capacity() returns an empty advertisement.

None
digest_verifier_image str

image for the verify-inputs init container, which needs the shell tools verify_files_script lists. Defaults to DEFAULT_DIGEST_VERIFIER_IMAGE, busybox:1.36 pinned by digest; point it at a mirror on a cluster without registry access, and pin the mirror by digest too.

DEFAULT_DIGEST_VERIFIER_IMAGE
trust_verification_records bool

whether verify-inputs accepts an input from a current verification record instead of hashing it. Defaults to True. Set it to False when untrusted parties can write the volume, because they could also write a record.

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.

capacity async

capacity() -> Capacity

Free slots per pool: max_slots - active, where active is the count of non-terminal Jobs. Returns an empty advertisement when max_slots is unset.

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_job(
    manifest: dict[str, Any], *, namespace: str
) -> str

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_job(
    name: str, *, namespace: str, uid: str | None = None
) -> None

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

diagnose(
    name: str, *, namespace: str
) -> tuple[bool, str | None]

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 /mnt/volume in each container. On a cluster backend this role is filled by a volume or PVC. Defaults to $ALTAR_LOCAL_DATA_ROOT, otherwise <cwd>/.altar-local-data.

None
slots int

Number of free slots that capacity() reports for the single "normal" pool. A local host has no GPU pools. Core does not read this value, though a driver may.

4
storage Storage | None

Storage adapter used to move a Transfer.uri in and out of the host data_root. Defaults to LocalStorage, which uses plain host paths.

None
container_user str | None

Optional Docker uid:gid identity. Set this to the owner of data_root when the bind mount is backed by root-squashed NFS. Defaults to the image's configured user.

None
run_command CommandRunner | None

Injectable async docker runner (see CommandRunner). Defaults to the real CLI.

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 $ALTAR_MODAL_APP, else "altar-exec".

None
volume_name str | None

Modal Volume mounted at /mnt/volume in every sandbox; it holds the shared file layout. Defaults to $ALTAR_MODAL_VOLUME, else "altar-data".

None
storage Storage | None

Storage used to move a Transfer.uri in and out of a host temp file before it is uploaded to, or after it is downloaded from, the Modal Volume. Defaults to LocalStorage.

None
runtime ModalRuntime | None

Modal SDK facade (see ModalRuntime). Defaults to _ModalSdkRuntime, which imports modal lazily.

None
pool_accelerators dict[str, str] | None

Maps a pool name to a Modal accelerator. Defaults to DEFAULT_POOL_ACCELERATORS so the shipped plugins work out of the box; a deployment with different hardware overrides it (for example, map GpuPool.NORMAL_GPU to "A100").

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 spec.timeout_s is longer than Modal allows a sandbox to run.

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

poll_sandbox(sandbox_id: str) -> int | None

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

upload_to_volume(
    volume_name: str, local_path: str, remote_path: str
) -> None

Copy a host file into the named Modal Volume at remote_path, relative to the volume root.

download_from_volume abstractmethod async

download_from_volume(
    volume_name: str, remote_path: str, local_path: str
) -> None

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

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

ReconcileStalled(groups: list[str])

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(kind: str, parameters: dict[str, Any]) -> str

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

LocalStorage(logger: Logger | None = None)

Bases: Storage

A Storage implementation backed by the local filesystem.

upload_from_local_file async

upload_from_local_file(
    local_path: str, dest_path: str
) -> None

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_csv(
    path: str, rows: list[Any], delimiter: str = "\t"
) -> None

Write the rows directly to disk, skipping the in-memory text round-trip.

as_local_file async

as_local_file(path: str) -> AsyncIterator[str]

Yield the path unchanged, since the file is already on the local disk.

Storage

Storage(logger: Logger | None = None)

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_text(path: str) -> str

Read a text file and return its contents as a string.

read_binary abstractmethod async

read_binary(path: str) -> bytes

Read a file and return its contents as bytes.

write_text abstractmethod async

write_text(path: str, content: str) -> None

Write a string to a text file.

write_binary abstractmethod async

write_binary(path: str, content: bytes) -> None

Write bytes to a file.

exists abstractmethod async

exists(path: str) -> bool

Return whether a file exists at the given path.

makedirs abstractmethod async

makedirs(path: str, exist_ok: bool = True) -> None

Create a directory and any missing parent directories.

delete abstractmethod async

delete(path: str) -> None

Delete the file at the given path.

upload_from_local_file abstractmethod async

upload_from_local_file(
    local_path: str, dest_path: str
) -> None

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

as_local_file(path: str) -> AsyncIterator[str]

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_csv(
    path: str, rows: list[Any], delimiter: str = "\t"
) -> None

Write a list of rows to a delimited file.

open_text_reader async

open_text_reader(path: str) -> StringIO

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

has_verification_record(
    path: str | PathLike[str], digest: str
) -> bool

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

verify_files_script(*, trust_records: bool = True) -> str

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

all_expected_folds(
    count_label: str = "num_folds",
    default: int | None = None,
) -> ReadyPredicate

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

all_folds(n: int) -> ReadyPredicate

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

all_label_values(
    label: str, expected: set[str]
) -> ReadyPredicate

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

both_alleles

both_alleles() -> ReadyPredicate

Ready when both the ref and alt shards are present.

count_at_least

count_at_least(n: int) -> ReadyPredicate

Ready when at least n shards have succeeded. A single-shard kind uses count_at_least(1).

group_by

group_by(*label_keys: str) -> GroupKeyFn

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.

  1. Poll every open handle from the ledger, in bulk and without blocking.
  2. Write each observed state back to the ledger.
  3. Re-derive groups from the ledger's succeeded-and-unconsumed shards, bucketed by each kind's group_key. For any group whose ready_when now passes, call on_ready once and mark that group's shards consumed, so it never fires twice. Callbacks run sequentially by default; callers may set max_concurrent_callbacks to process independent groups concurrently with a fixed upper bound.
  4. Surface any group blocked by a terminally-failed shard, which can never reach ready_when. Dead-letter it via on_failed if the spec supplies one, otherwise report it in stalled_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

single() -> ReadyPredicate

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.