Skip to content

API Reference

Auto-generated API documentation from source code.

Core Module

Stardag: Declarative and composable DAG framework for Python.

Stardag provides a clean Python API for representing persistently stored assets as a declarative Directed Acyclic Graph (DAG).

Basic usage::

import stardag as sd

@sd.task
def get_range(limit: int) -> list[int]:
    return list(range(limit))

@sd.task
def get_sum(integers: sd.Depends[list[int]]) -> int:
    return sum(integers)

task = get_sum(integers=get_range(limit=10))
sd.build(task)
print(task.target().load())  # 45

Core components:

  • :func:task - Decorator for creating tasks from functions
  • :class:Task - Task with automatic serialization and filesystem targets
  • :class:LoadableTask - Abstract base for tasks with load() -> T
  • :class:TargetTask - Base class for tasks with typed target outputs
  • :class:Depends - Dependency injection type annotation
  • :func:build - Execute task and its dependencies

See https://docs.stardag.com for full documentation.

TODO: Expand docstrings for all public API components.

Depends module-attribute

Depends = typing.Annotated[_DependsT, _DependsOnMarker]

TaskLoads module-attribute

TaskLoads = typing.Annotated[
    LoadableTask[LoadedT_co], Polymorphic()
]

TaskStruct module-attribute

TaskStruct = Union[
    "BaseTask",
    Sequence["TaskStruct"],
    Mapping[str, "TaskStruct"],
]

target_factory_provider module-attribute

target_factory_provider = resource_provider(
    type_=TargetFactory, default_factory=TargetFactory
)

BaseTask

Bases: PolymorphicRoot

The base of every task class.

A subclass declares a task's parameters as pydantic fields. An object of it is a task object; what the registry stores is an instance — a row holding one construction of a task under a deterministic scope. One task object planned under two scopes is two instances; one instance rehydrates into any number of task objects.

A task has two identities:

  • :attr:id, the task id: a hash of the class (namespace, name), version and every significant field (StardagField(significant=True), the default). It is the promise about output — completion and the execution claim are global on it — and it names the target.
  • :attr:instance_hash: the hash of the canonical instance body, i.e. of all fields, defaults included. It answers "how exactly was this task constructed", and is meaningful only together with the scope it is registered under.

Two task objects with one task id and different instance hashes are two ways of asking for one completion; a plan may hold only one of them.

id cached property

id

The task id: the completion identity.

uuid5 over the canonical hash-mode dump — the class discriminators, version and every significant field; a field whose value equals its compat_default is dropped, and a nested task contributes its own task id. Custom serializers may special-case the "hash" mode: that is the user's control over completion identity.

instance_hash cached property

instance_hash

The hash of the canonical instance body (all parameters).

Not an identifier on its own: an instance is a registry row keyed by its scope and this hash, so the same hash under two scopes is two instances. Use :attr:id for "which task", and the scope together with this for "which instance". Has no hash mode of its own; custom serializers affect it only through ordinary serialization, which must be stable (see stardag._core.instance.check_serialization_stability).

__init_subclass__

__init_subclass__(**kwargs)

Validate that subclasses implement either run() or run_aio().

Also wraps run() and run_aio() methods with precheck validation.

complete abstractmethod

complete()

Declare if the task is complete.

complete_aio async

complete_aio()

Asynchronously declare if the task is complete.

run

run()

Execute the task logic (sync).

Override this method for synchronous tasks. If you only override run_aio(), this method will automatically run it via asyncio.run().

RETURNS DESCRIPTION
None | Generator[TaskStruct, None, None]

None for simple tasks, or a Generator yielding TaskStruct for

None | Generator[TaskStruct, None, None]

tasks with dynamic dependencies (See Dynamic Dependencies Contract below).

RAISES DESCRIPTION
RuntimeError

If called from within an existing event loop when only run_aio() is implemented. In that case, call run_aio() directly instead.

NotImplementedError

If run_aio() is an async generator (dynamic deps). Async generators cannot be automatically converted to sync generators.

Dynamic Dependencies Contract: When a task yields dynamic dependencies via a generator, the BUILD SYSTEM guarantees that ALL yielded tasks are COMPLETE before the generator is resumed. The task can rely on this contract:

def run(self):
    # Do some initial work to get info about what additional dependencies
    # are needed
    initial_data = "..."

    # Yield deps we need to be built first
    task_a = TaskA(input=initial_data)
    task_b = TaskB(input=initial_data)
    yield [task_a, task_b]

    # CONTRACT: When we reach here, ALL deps are complete.
    # We can safely access their outputs.
    result_a = task_a.target().load()
    result_b = task_b.target().load()

    # Yield more deps if needed
    task_c = TaskC(input=result_a)
    yield task_c

    # Again, TaskC is complete when we reach here
    self.target().save(task_c.target().load() + result_b)

This contract is essential for correctness - tasks can depend on previously yielded tasks being complete before continuing execution.

run_aio async

run_aio()

Execute the task logic (async).

Override this method for asynchronous tasks. If you only override run(), this method will automatically delegate to it.

For dynamic dependencies, you can use 'yield' which makes this an async generator. Note that async generator methods have different type signatures that may require type: ignore comments.

RETURNS DESCRIPTION
None | Generator[TaskStruct, None, None]

None for simple tasks, or a Generator/AsyncGenerator for

None | Generator[TaskStruct, None, None]

tasks with dynamic dependencies.

Dynamic Dependencies Contract

Same as run() - the build system guarantees that ALL yielded tasks are COMPLETE before the generator is resumed. See run() docstring for detailed documentation and examples.

artifacts

artifacts()

Return artifacts to be stored in the registry after task completion.

Override this method to expose rich outputs (reports, summaries, structured data) that will be viewable in the registry UI.

This method is called after the task completes successfully. It should be stateless - loading any required data from the task's target rather than relying on in-memory state.

RETURNS DESCRIPTION
Sequence[Artifact]

Sequence of artifacts (MarkdownArtifact, JSONArtifact, etc.)

artifacts_aio async

artifacts_aio()

Asynchronously return artifacts to be stored in the registry after task completion.

instance_body

instance_body()

The instance body: the registry-mode dump of every field, defaults included, nested tasks as their full bodies — parsed from the same canonical JSON that :attr:instance_hash hashes, so what is sent is exactly what was hashed. A fresh dict on every call.

resolve classmethod

resolve(namespace, name, extra)

Override PolymorphicRoot.resolve to handle AliasTask deserialization.

from_registry classmethod

from_registry(id, registry=None)

Instantiate the task from the registry.

A task id may have several instances — constructions under different scopes — any of which rehydrates into this completion. A task has no parameters, an instance does, so without a scope to narrow by, this takes the newest instance's body (TaskInfo orders them newest first).

Validated in compat mode, same as task_from_registry_data: the recomputed task id is checked against the requested one, so a removed or renamed significant field — dropped by compat mode's lenient rules for non-significant fields — cannot silently return a task with a different completion identity than the one asked for.

PARAMETER DESCRIPTION
id

The UUID (or string representation) of the task to load.

TYPE: UUID | str

registry

An optional registry instance to use for loading metadata. If not provided, the default registry from registry_provider will be used.

TYPE: Union[RegistryABC, None] DEFAULT: None

Returns: The reconstructed task object.

RAISES DESCRIPTION
TaskRehydrationError

The recomputed task id does not match the requested one.

LoadableTask

Bases: BaseTask, ABC, Generic[LoadedT_co]

A task that can load its output as a typed value.

This is the minimal interface required by :class:~stardag.TaskLoads: any BaseTask subclass that implements load() -> T is compatible with TaskLoads[T].

Both :class:~stardag.Task (via diamond inheritance) and bare subclasses of LoadableTask satisfy TaskLoads[T].

Subclasses must implement at least one of load() or load_aio(). The missing method will delegate to the other automatically (mirroring the run/run_aio pattern on BaseTask).

load

load()

Load the output of this task (sync).

If only load_aio() is implemented, this delegates via asyncio.run(). Raises RuntimeError if called from within an existing event loop.

load_aio async

load_aio()

Asynchronously load the output of this task.

If only load() is implemented, this delegates to it.

TargetTask

Bases: BaseTask, Generic[TargetType]

Base class for tasks that produce a target output.

Extends BaseTask with a typed target() method and a default complete() implementation that checks whether the target exists.

Most users should subclass :class:~stardag.Task (which extends this class with automatic serialization and filesystem target management) rather than using TargetTask directly.

complete

complete()

Check if the task is complete.

complete_aio async

complete_aio()

Asynchronously check if the task is complete.

target abstractmethod

target()

The task output target.

Task

Bases: TargetTask[LoadableSaveableFileSystemTarget[LoadedT]], LoadableTask[LoadedT], ABC, Generic[LoadedT]

A base class for tasks with automatic serialization and filesystem targets.

The target of a Task is a LoadableSaveableFileSystemTarget that uses a serializer inferred from the generic type parameter LoadedT.

The target file path is automatically constructed based on the task's namespace, name, version, and unique ID and has the following structure:

[<relpath_base>/][<namespace>/]<name>/v<version>/[<relpath_extra>/]
<id>[:2]/<id>[2:4]/<id>[/<relpath_filename>].<relpath_extension>

You can override the following properties to customize the target path: _relpath_base, _relpath_extra, _relpath_filename, and _relpath_extension.

See stardag.target.serialize.get_serializer for details on how the serializer is inferred from the generic type parameter, and how to customize it.

Example:

import stardag as sd

class MyTask(sd.Task[dict[str, int]]):
    def run(self):
        self._save({"a": 1, "b": 2})

my_task = MyTask()

print(my_task.target())
# FileSerializable(../MyTask/03/6f/036f6e71-1b3c-54b8-aec1-182359f1e09a.json)

print(my_task.target().serializer)
# <stardag.target.serialize.JSONSerializer at 0x1064e4710>

serializer property

serializer

The serializer used for this task's target.

__map_generic_args_to_ancestor__ classmethod

__map_generic_args_to_ancestor__(ancestor_origin, args)

Map generic args from Task to how they appear on an ancestor class.

This enables type compatibility checking when using Task with polymorphic annotations like TaskLoads[T] and SubClass[TargetTask[LoadableTarget[T]]].

PARAMETER DESCRIPTION
ancestor_origin

The ancestor class to map args to

TYPE: type

args

The generic args of this class (e.g., (str,) for Task[str])

TYPE: tuple

RETURNS DESCRIPTION
tuple | None

The mapped args for the ancestor, or None if mapping is not applicable.

load

load()

Load the task target and run any LoadValidators.

load_aio async

load_aio()

Async load — delegates to the target's load_aio and validates.

TaskRef dataclass

TaskRef(name, version, id)

task

task(
    _func: Callable[_PWrapped, _FuncReturnT],
    *,
    name: str | None = None,
    version: str = "",
    relpath: RelpathSettings
    | _RelpathOverride
    | None = None,
    target_root_key: str | None = None,
) -> Type[_FunctionTask[_FuncReturnT, _PWrapped]]
task(
    *,
    name: str | None = None,
    version: str = "",
    relpath: RelpathSettings
    | _RelpathOverride
    | None = None,
    target_root_key: str | None = None,
) -> _TaskWrapper
task(
    _func=None,
    *,
    name=None,
    version="",
    relpath=None,
    target_root_key=None,
)

build_aio async

build_aio(
    tasks,
    task_executor=None,
    fail_mode=FAIL_FAST,
    registry=None,
    max_concurrent_discover=50,
    resume_build_id=None,
    register_all=False,
    on_registry_failure="raise",
    concurrency_config=None,
    concurrency_limiter=None,
    claim_config=None,
    settings=None,
    limit_key_selector=None,
    description=None,
    raise_on_failure=True,
)

Build tasks concurrently using hybrid async/thread/process execution.

Walks the DAG from the roots (stopping at complete tasks), registers the plan with the registry — roots first, the rest in post-order, then sealed — and runs every task whose dependencies are met on task_executor. With a registry every execution claims, so two builds (or a build and a reactive one) never run a task twice at once.

PARAMETER DESCRIPTION
tasks

The root task(s).

TYPE: Sequence[BaseTask] | BaseTask

task_executor

Where tasks run (default: :class:HybridConcurrentTaskExecutor). An executor running tasks on a deployed Modal app makes the build plan under that app's current deployment (D13).

TYPE: TaskExecutorABC | None DEFAULT: None

fail_mode

Stop at the first failure (FAIL_FAST) or run everything whose dependencies are met (CONTINUE).

TYPE: FailMode DEFAULT: FAIL_FAST

registry

Default: the configured registry (none: no plan, no claims — the build runs purely locally).

TYPE: RegistryABC | None DEFAULT: None

max_concurrent_discover

Completion checks in flight while walking.

TYPE: int DEFAULT: 50

resume_build_id

Resume this build: its plan for this scope is reused, observations re-sent, and failed members reset.

TYPE: UUID | None DEFAULT: None

register_all

Expand complete tasks too, so every edge is recorded.

TYPE: bool DEFAULT: False

on_registry_failure

"raise" (default) or "warn" — carry on through a registry outage; a refusal always raises.

TYPE: OnRegistryFailure DEFAULT: 'raise'

concurrency_config

Build-local limits (how many tasks run at once in this process).

TYPE: ConcurrencyConfig | None DEFAULT: None

concurrency_limiter

A pre-built limiter to use instead of concurrency_config.

TYPE: ConcurrencyLimiter | None DEFAULT: None

claim_config

How claims are waited on and renewed.

TYPE: ClaimConfig | None DEFAULT: None

settings

Environment variables applied for the build's duration (the scope's second half; STARDAG_* / MODAL_* refused). Omitted on a resume, the build's active plan's settings are reused; {} explicitly means none. Two concurrent builds in one process with different settings are refused.

TYPE: Mapping[str, str] | None DEFAULT: None

limit_key_selector

The registry concurrency-limit keys a task runs under, sent with its claim.

TYPE: LimitKeySelector | None DEFAULT: None

description

A description for a new build.

TYPE: str | None DEFAULT: None

raise_on_failure

In FAIL_FAST mode, re-raise the failure's exception (default). False returns the FAILURE summary instead — build id, failed task, error — after the build has stopped at the first failure the same way.

TYPE: bool DEFAULT: True

RETURNS DESCRIPTION
BuildSummary

BuildSummary with status, task counts and build id. STOPPED

BuildSummary

when the registry stopped the build (an operator cancelled it):

BuildSummary

the engine then stops without writing a build status.

build_sequential

build_sequential(
    tasks,
    registry=None,
    fail_mode=FAIL_FAST,
    dual_run_default="sync",
    resume_build_id=None,
    register_all=False,
    on_registry_failure="raise",
    claim_config=None,
    settings=None,
    limit_key_selector=None,
    description=None,
    max_concurrent_discover=16,
    raise_on_failure=True,
)

Build tasks sequentially, from synchronous code (for debugging).

Task code runs on the calling thread: sync tasks via run(), async-only tasks via asyncio.run(run_aio()) (so not from inside a running event loop), dual tasks by dual_run_default. See :func:build_sequential_aio for the other arguments.

build_sequential_aio async

build_sequential_aio(
    tasks,
    registry=None,
    fail_mode=FAIL_FAST,
    sync_run_default="blocking",
    resume_build_id=None,
    register_all=False,
    on_registry_failure="raise",
    claim_config=None,
    settings=None,
    limit_key_selector=None,
    description=None,
    max_concurrent_discover=16,
    raise_on_failure=True,
)

Build tasks sequentially from async code (for debugging).

PARAMETER DESCRIPTION
tasks

The root task(s).

TYPE: Sequence[BaseTask] | BaseTask

registry

Default: the configured registry (none: no plan, no claims).

TYPE: RegistryABC | None DEFAULT: None

fail_mode

FAIL_FAST (default) or CONTINUE.

TYPE: FailMode DEFAULT: FAIL_FAST

sync_run_default

How a sync-only task runs: "blocking" on the loop, or "thread".

TYPE: Literal['thread', 'blocking'] DEFAULT: 'blocking'

resume_build_id

Resume this build (its plan for this scope is reused, observations re-sent, failed members reset).

TYPE: UUID | None DEFAULT: None

register_all

Expand complete tasks too.

TYPE: bool DEFAULT: False

on_registry_failure

"raise" or "warn" (outages only).

TYPE: OnRegistryFailure DEFAULT: 'raise'

claim_config

Claim waiting and renewal.

TYPE: ClaimConfig | None DEFAULT: None

settings

Environment variables applied for the build's duration. Omitted on a resume, the build's active plan's settings are reused; {} explicitly means none.

TYPE: Mapping[str, str] | None DEFAULT: None

limit_key_selector

Registry concurrency-limit keys per task, sent with its claim.

TYPE: LimitKeySelector | None DEFAULT: None

description

A description for a new build.

TYPE: str | None DEFAULT: None

max_concurrent_discover

Completion checks in flight while walking.

TYPE: int DEFAULT: 16

raise_on_failure

In FAIL_FAST mode, re-raise the failure's exception (default); False returns the FAILURE summary.

TYPE: bool DEFAULT: True

namespace

namespace(namespace, scope)

Set the task namespace for the module and any submodules.

PARAMETER DESCRIPTION
namespace

The namespace to set for the module.

TYPE: str

scope

The module scope, typically passed as __name__.

TYPE: str

Usage:

```python import stardag as sd sd.namespace("my_custom_namespace", name)

class MyNamespacedTask(sd.Task[int]): a: int

def run(self):
    self._save(self.a)

assert MyNamespacedTask.get_namespace() == "my_custom_namespace"

auto_namespace

auto_namespace(scope)

Set the task namespace for the module to the module import path.

PARAMETER DESCRIPTION
scope

The module scope, typically passed as __name__.

TYPE: str

Usage:

import stardag as sd

sd.auto_namespace(__name__)

class MyAutoNamespacedTask(sd.Task[int]):
    a: int

    def run(self):
        self._save(self.a)

assert MyAutoNamespacedTask.get_namespace() == __name__

get_file_target

get_file_target(
    relpath, target_root_key=DEFAULT_TARGET_ROOT_KEY
)

Get a file target for the given relative path.

get_directory_target

get_directory_target(
    relpath, target_root_key=DEFAULT_TARGET_ROOT_KEY
)

Build Module

build

Build module for stardag.

Primary build functions: - build(): Concurrent build, recommended for real workloads from a sync context - build_aio(): Async concurrent build, for an async context or a running loop - build_sequential(): Sync sequential build (for debugging) - build_sequential_aio(): Async sequential build (for debugging)

With a registry, every build plans under a scope (deployment, settings) and every execution claims (see docs/design/registry-v2/design.md); without one, it runs purely locally.

Task executor: - HybridConcurrentTaskExecutor: Routes tasks to async/thread/process pools - RoutedTaskExecutor: Routes tasks between executors

Interfaces: - TaskExecutorABC: Abstract base class for custom task executors - ExecutionModeSelector: Protocol for custom execution mode selection

Concurrency limiting: - ConcurrencyConfig: Build-local overall and named limits - ConcurrencyLimiter: Protocol for custom limiters

Reactive scheduling (ticks) and task-module declaration: - run_tick_aio / TickConfig / TickSummary - expand_task_module_patterns() / import_task_modules() / plan_rehydration()

BuildSummary dataclass

BuildSummary(
    status,
    task_count,
    build_id=None,
    error=None,
    failed_task=None,
)

Summary of a build execution.

raise_on_failure

raise_on_failure()

Raise :class:BuildFailed if the build status is FAILURE or STOPPED (the build ended without producing its roots).

Source code in stardag/build/_base.py
def raise_on_failure(self) -> None:
    """Raise :class:`BuildFailed` if the build status is ``FAILURE`` or
    ``STOPPED`` (the build ended without producing its roots)."""
    if self.status in (BuildExitStatus.FAILURE, BuildExitStatus.STOPPED):
        raise BuildFailed(self)

__repr__

__repr__()

Return a human-readable summary of the build.

Source code in stardag/build/_base.py
def __repr__(self) -> str:
    """Return a human-readable summary of the build."""
    tc = self.task_count
    status_icon = "✓" if self.status == BuildExitStatus.SUCCESS else "✗"
    lines = [
        f"Build {self.status.value.upper()} {status_icon}",
    ]
    if self.build_id:
        lines.append(f"  Build ID: {self.build_id}")
    lines.extend(
        [
            f"  Discovered: {tc.discovered}",
            f"  Previously completed: {tc.previously_completed}",
            f"  Succeeded: {tc.succeeded}",
            f"  Failed: {tc.failed}",
        ]
    )
    if tc.cancelled > 0:
        lines.append(f"  Cancelled: {tc.cancelled}")
    if tc.skipped > 0:
        lines.append(f"  Skipped: {tc.skipped}")
    if tc.pending > 0:
        lines.append(f"  Pending: {tc.pending}")
    if self.failed_task is not None:
        lines.append(f"  Failed task: {describe_task(self.failed_task)}")
    if self.error:
        lines.append(f"  Error: {self.error}")
    return "\n".join(lines)

BuildExitStatus

Bases: StrEnum

FailMode

Bases: StrEnum

How to handle task failures during build.

ATTRIBUTE DESCRIPTION
FAIL_FAST

Stop the build at the first task failure.

CONTINUE

Continue executing all tasks whose dependencies are met, even if some tasks have failed.

HybridConcurrentTaskExecutor

HybridConcurrentTaskExecutor(
    execution_mode_selector=None,
    max_async_workers=10,
    max_thread_workers=10,
    max_process_workers=None,
)

Bases: TaskExecutorABC

Task executor with async, thread, and process pools.

Routes tasks to appropriate execution context based on ExecutionModeSelector. Handles generator suspension for dynamic dependencies.

Note: This executor does not handle registry calls - those are managed by the build() function. The executor only executes tasks and returns results.

For routing tasks to different executors (e.g., some to Modal, some local), use RoutedTaskExecutor to compose multiple executors.

Alternative: For fully async multiprocessing without thread pools, one could implement an AIOMultiprocessingTaskExecutor using libraries like aiomultiprocess.

PARAMETER DESCRIPTION
execution_mode_selector

Callable to select execution mode per task.

TYPE: ExecutionModeSelector | None DEFAULT: None

max_async_workers

Maximum concurrent async tasks (semaphore-based).

TYPE: int DEFAULT: 10

max_thread_workers

Maximum concurrent thread pool workers.

TYPE: int DEFAULT: 10

max_process_workers

Maximum concurrent process pool workers.

TYPE: int | None DEFAULT: None

Source code in stardag/build/_concurrent.py
def __init__(
    self,
    execution_mode_selector: ExecutionModeSelector | None = None,
    max_async_workers: int = 10,
    max_thread_workers: int = 10,
    max_process_workers: int | None = None,
) -> None:
    self.execution_mode_selector = (
        execution_mode_selector or DefaultExecutionModeSelector()
    )
    self.max_async_workers = max_async_workers
    self.max_thread_workers = max_thread_workers
    self.max_process_workers = max_process_workers

    # Pools - initialized in setup()
    self._async_semaphore: asyncio.Semaphore | None = None
    self._thread_pool: ThreadPoolExecutor | None = None
    self._process_pool: ProcessPoolExecutor | None = None

    # Track suspended generators (task_id -> sync or async generator)
    # For in-process execution where we can suspend and resume
    self._suspended_generators: dict[
        UUID,
        Union[
            Generator[TaskStruct, None, None],
            AsyncGenerator[TaskStruct, None],
        ],
    ] = {}

    # Track tasks pending re-execution (task_id -> True)
    # For cross-process/remote execution: when task yields incomplete deps,
    # it's re-executed from scratch after deps complete (idempotent re-execution)
    self._pending_reexecution: set[UUID] = set()

setup async

setup()

Initialize worker pools.

Source code in stardag/build/_concurrent.py
async def setup(self) -> None:
    """Initialize worker pools."""
    import multiprocessing as mp

    self._async_semaphore = asyncio.Semaphore(self.max_async_workers)
    self._thread_pool = ThreadPoolExecutor(max_workers=self.max_thread_workers)
    if self.max_process_workers:
        # Use 'spawn' explicitly for cross-platform compatibility.
        # Python 3.14 changed the default from 'fork' to 'forkserver' on Linux,
        # which can cause issues with environment variable inheritance.
        # 'spawn' is the safest option and works consistently across platforms.
        self._process_pool = ProcessPoolExecutor(
            max_workers=self.max_process_workers,
            mp_context=mp.get_context("spawn"),
        )

teardown async

teardown()

Shutdown worker pools.

Source code in stardag/build/_concurrent.py
async def teardown(self) -> None:
    """Shutdown worker pools."""
    if self._thread_pool:
        self._thread_pool.shutdown(wait=True)
        self._thread_pool = None
    if self._process_pool:
        self._process_pool.shutdown(wait=True)
        self._process_pool = None
    self._async_semaphore = None
    self._suspended_generators.clear()
    self._pending_reexecution.clear()

get_executor_details async

get_executor_details(task)

This process, and the execution mode task runs in (the pool a process-mode task lands in is not known before it runs).

Source code in stardag/build/_concurrent.py
async def get_executor_details(self, task: BaseTask) -> ExecutorDetails:
    """This process, and the execution mode ``task`` runs in (the pool
    a process-mode task lands in is not known before it runs)."""
    return in_process_executor_details(self.execution_mode_selector(task).value)

submit async

submit(task)

Execute a task and return result.

Note: This method does not make any registry calls. The build function is responsible for calling start_task, complete_task, and fail_task.

Source code in stardag/build/_concurrent.py
async def submit(self, task: BaseTask) -> None | TaskStruct | TaskExecutionError:
    """Execute a task and return result.

    Note: This method does not make any registry calls. The build function
    is responsible for calling start_task, complete_task, and fail_task.
    """
    # Check if we're resuming a suspended generator (in-process dynamic deps)
    if task.id in self._suspended_generators:
        gen = self._suspended_generators[task.id]
        if hasattr(gen, "__anext__"):
            return await self._resume_generator_aio(task)
        return self._resume_generator(task)

    # Check if task is pending re-execution (cross-process dynamic deps)
    # Task yielded incomplete deps, deps are now built, re-execute task
    if task.id in self._pending_reexecution:
        self._pending_reexecution.discard(task.id)

    mode = self.execution_mode_selector(task)

    try:
        result = await self._execute_task(task, mode)
        return await self._handle_result(task, result)
    except Exception as e:
        return TaskExecutionError(
            exception=e,
            traceback="".join(tb_module.format_exception(e)),
        )

TaskExecutorABC

Bases: ABC

Abstract base for task executors.

Receives tasks and executes them according to some policy. The executor is responsible for executing tasks in the appropriate context (async/thread/process/remote) and for generator suspension. Dependency resolution and the registry are the build engine's.

submit abstractmethod async

submit(task)

Execute a task in-process (or block on a remote one).

RETURNS DESCRIPTION
None | TaskStruct | TaskExecutionError
  • None: completed, no dynamic dependencies.
None | TaskStruct | TaskExecutionError
  • TaskStruct: suspended on dynamic dependencies.
None | TaskStruct | TaskExecutionError
  • TaskExecutionError: failed, with the traceback captured where it happened.
Source code in stardag/build/_base.py
@abstractmethod
async def submit(self, task: BaseTask) -> None | TaskStruct | TaskExecutionError:
    """Execute a task in-process (or block on a remote one).

    Returns:
        - None: completed, no dynamic dependencies.
        - TaskStruct: suspended on dynamic dependencies.
        - TaskExecutionError: failed, with the traceback captured where
            it happened.
    """
    ...

setup abstractmethod async

setup()

Setup any resources needed for the task runner (pools, etc.).

Source code in stardag/build/_base.py
@abstractmethod
async def setup(self) -> None:
    """Setup any resources needed for the task runner (pools, etc.)."""
    ...

teardown abstractmethod async

teardown()

Teardown any resources used by the task executor.

Source code in stardag/build/_base.py
@abstractmethod
async def teardown(self) -> None:
    """Teardown any resources used by the task executor."""
    ...

cancel async

cancel(task)

Best-effort cancel an in-flight task.

Default: no-op; the engine also cancels the asyncio future wrapping submit(), which propagates into cooperative awaitables. Executors of detached executions must override this to stop the remote work — cancelling wait() does not.

Source code in stardag/build/_base.py
async def cancel(self, task: BaseTask) -> None:
    """Best-effort cancel an in-flight task.

    Default: no-op; the engine also cancels the asyncio future wrapping
    ``submit()``, which propagates into cooperative awaitables. Executors
    of detached executions must override this to stop the remote work —
    cancelling ``wait()`` does not.
    """
    pass

get_executor_metadata async

get_executor_metadata(task)

Descriptive executor metadata for executions of task, without starting anything (stamped on the claiming start). Best-effort.

Source code in stardag/build/_base.py
async def get_executor_metadata(self, task: BaseTask) -> dict[str, Any] | None:
    """Descriptive executor metadata for executions of ``task``, without
    starting anything (stamped on the claiming start). Best-effort."""
    return None

get_executor_details async

get_executor_details(task)

What the claiming start records for an execution of task (see :class:ExecutorDetails). Default: the metadata only; an in-process executor also names itself and this process.

Source code in stardag/build/_base.py
async def get_executor_details(self, task: BaseTask) -> ExecutorDetails:
    """What the claiming start records for an execution of ``task``
    (see :class:`ExecutorDetails`). Default: the metadata only; an
    in-process executor also names itself and this process."""
    return ExecutorDetails(executor_metadata=await self.get_executor_metadata(task))

execution_timeout_seconds

execution_timeout_seconds(task)

Wall-clock limit this backend enforces on an execution of task — a fact (e.g. Modal's per-function timeout), from which a detached execution's claim TTL is derived. None when there is none or it cannot be resolved. Must not raise or do I/O.

Source code in stardag/build/_base.py
def execution_timeout_seconds(self, task: BaseTask) -> float | None:
    """Wall-clock limit this backend enforces on an execution of
    ``task`` — a *fact* (e.g. Modal's per-function ``timeout``), from
    which a detached execution's claim TTL is derived. None when there
    is none or it cannot be resolved. Must not raise or do I/O."""
    return None

reports_lifecycle

reports_lifecycle(task)

Whether the execution side reports this task's lifecycle.

When True the worker executing it reports its own start (with its executor ref), completion, failure and yields — naming the execution the engine claimed — and the engine reports none of them. Default: False, the engine reports everything.

Source code in stardag/build/_base.py
def reports_lifecycle(self, task: BaseTask) -> bool:
    """Whether the *execution side* reports this task's lifecycle.

    When True the worker executing it reports its own start (with its
    executor ref), completion, failure and yields — naming the execution
    the engine claimed — and the engine reports none of them. Default:
    False, the engine reports everything.
    """
    return False

deployment_app_name

deployment_app_name()

The deployed app this executor runs tasks on, if any.

A driver whose tasks run on a deployed app plans under that app's current deployment (D13) — its workers yield into the plan, and the registry refuses a yield from another deployment. Default: None (a pure local executor; the build plans under a local deployment).

Source code in stardag/build/_base.py
def deployment_app_name(self) -> str | None:
    """The deployed app this executor runs tasks on, if any.

    A driver whose tasks run on a deployed app plans under that app's
    current deployment (D13) — its workers yield into the plan, and the
    registry refuses a yield from another deployment. Default: None (a
    pure local executor; the build plans under a local deployment).
    """
    return None

supports_detached

supports_detached(task)

Whether this executor can run task as a detached execution (one that survives the orchestrator). Default: False.

Source code in stardag/build/_base.py
def supports_detached(self, task: BaseTask) -> bool:
    """Whether this executor can run ``task`` as a detached execution
    (one that survives the orchestrator). Default: False."""
    return False

submit_detached async

submit_detached(task, *, execution_id)

Start a detached execution of task and return its handle.

execution_id is the execution the engine claimed. An executor whose workers report their own lifecycle must forward it into the execution: every report names it, and the registry applies a report only while that execution holds the task's claim.

RAISES DESCRIPTION
Exception

the execution could not be started (the engine records the failure against the execution).

Source code in stardag/build/_base.py
async def submit_detached(
    self, task: BaseTask, *, execution_id: UUID
) -> DetachedHandle:
    """Start a detached execution of ``task`` and return its handle.

    ``execution_id`` is the execution the engine claimed. An executor
    whose workers report their own lifecycle must forward it into the
    execution: every report names it, and the registry applies a report
    only while that execution holds the task's claim.

    Raises:
        Exception: the execution could not be started (the engine
            records the failure against the execution).
    """
    raise NotImplementedError(
        f"{type(self).__name__} does not support detached execution"
    )

cancel_detached async

cancel_detached(task, executor, ref)

Best-effort stop of a detached execution by its reference.

Called only for an execution this engine spawned itself whose start the registry then refused — an orphan nothing else can find. Default: no-op.

Source code in stardag/build/_base.py
async def cancel_detached(self, task: BaseTask, executor: str, ref: str) -> None:
    """Best-effort stop of a detached execution by its reference.

    Called only for an execution this engine spawned itself whose start
    the registry then refused — an orphan nothing else can find.
    Default: no-op.
    """
    pass

can_spawn_scheduler_ticks

can_spawn_scheduler_ticks()

Whether :meth:spawn_scheduler_tick reaches a deployed tick (a resident build then drains the registry's wake candidates).

Source code in stardag/build/_base.py
def can_spawn_scheduler_ticks(self) -> bool:
    """Whether :meth:`spawn_scheduler_tick` reaches a deployed ``tick``
    (a resident build then drains the registry's wake candidates)."""
    return False

spawn_scheduler_tick

spawn_scheduler_tick(build_id, app_name)

Spawn a reactive scheduler tick for build_id on app_name.

Source code in stardag/build/_base.py
def spawn_scheduler_tick(self, build_id: UUID, app_name: str) -> None:
    """Spawn a reactive scheduler tick for ``build_id`` on ``app_name``."""
    raise NotImplementedError(f"{type(self).__name__} cannot spawn scheduler ticks")

build

build(
    tasks,
    task_executor=None,
    fail_mode=FAIL_FAST,
    registry=None,
    max_concurrent_discover=50,
    resume_build_id=None,
    register_all=False,
    on_registry_failure="raise",
    concurrency_config=None,
    concurrency_limiter=None,
    claim_config=None,
    settings=None,
    limit_key_selector=None,
    description=None,
    raise_on_failure=True,
)

Build tasks concurrently (sync wrapper for build_aio).

This is the recommended entry point for building tasks from synchronous code. Wraps the async build_aio() function; see it for the arguments.

Note

This function cannot be called from within an already running event loop. If you're in an async context (e.g., inside an async function, or using frameworks like Playwright, FastAPI, etc.), use await build_aio() instead.

Source code in stardag/build/_concurrent.py
def build(
    tasks: Sequence[BaseTask] | BaseTask,
    task_executor: TaskExecutorABC | None = None,
    fail_mode: FailMode = FailMode.FAIL_FAST,
    registry: RegistryABC | None = None,
    max_concurrent_discover: int = 50,
    resume_build_id: UUID | None = None,
    register_all: bool = False,
    on_registry_failure: OnRegistryFailure = "raise",
    concurrency_config: ConcurrencyConfig | None = None,
    concurrency_limiter: ConcurrencyLimiter | None = None,
    claim_config: ClaimConfig | None = None,
    settings: Mapping[str, str] | None = None,
    limit_key_selector: LimitKeySelector | None = None,
    description: str | None = None,
    raise_on_failure: bool = True,
) -> BuildSummary:
    """Build tasks concurrently (sync wrapper for build_aio).

    This is the recommended entry point for building tasks from synchronous code.
    Wraps the async build_aio() function; see it for the arguments.

    Note:
        This function cannot be called from within an already running event loop.
        If you're in an async context (e.g., inside an async function, or using
        frameworks like Playwright, FastAPI, etc.), use `await build_aio()` instead.
    """
    try:
        return asyncio.run(
            build_aio(
                tasks,
                task_executor=task_executor,
                fail_mode=fail_mode,
                registry=registry,
                max_concurrent_discover=max_concurrent_discover,
                resume_build_id=resume_build_id,
                register_all=register_all,
                on_registry_failure=on_registry_failure,
                concurrency_config=concurrency_config,
                concurrency_limiter=concurrency_limiter,
                claim_config=claim_config,
                settings=settings,
                limit_key_selector=limit_key_selector,
                description=description,
                raise_on_failure=raise_on_failure,
            )
        )
    except RuntimeError as e:
        if "cannot be called from a running event loop" in str(e):
            raise RuntimeError(
                "build() cannot be used from within an already running event loop. "
                "Use 'await build_aio()' instead, or 'build_sequential()' if you "
                "need synchronous execution without an event loop."
            ) from e
        raise

build_aio async

build_aio(
    tasks,
    task_executor=None,
    fail_mode=FAIL_FAST,
    registry=None,
    max_concurrent_discover=50,
    resume_build_id=None,
    register_all=False,
    on_registry_failure="raise",
    concurrency_config=None,
    concurrency_limiter=None,
    claim_config=None,
    settings=None,
    limit_key_selector=None,
    description=None,
    raise_on_failure=True,
)

Build tasks concurrently using hybrid async/thread/process execution.

Walks the DAG from the roots (stopping at complete tasks), registers the plan with the registry — roots first, the rest in post-order, then sealed — and runs every task whose dependencies are met on task_executor. With a registry every execution claims, so two builds (or a build and a reactive one) never run a task twice at once.

PARAMETER DESCRIPTION
tasks

The root task(s).

TYPE: Sequence[BaseTask] | BaseTask

task_executor

Where tasks run (default: :class:HybridConcurrentTaskExecutor). An executor running tasks on a deployed Modal app makes the build plan under that app's current deployment (D13).

TYPE: TaskExecutorABC | None DEFAULT: None

fail_mode

Stop at the first failure (FAIL_FAST) or run everything whose dependencies are met (CONTINUE).

TYPE: FailMode DEFAULT: FAIL_FAST

registry

Default: the configured registry (none: no plan, no claims — the build runs purely locally).

TYPE: RegistryABC | None DEFAULT: None

max_concurrent_discover

Completion checks in flight while walking.

TYPE: int DEFAULT: 50

resume_build_id

Resume this build: its plan for this scope is reused, observations re-sent, and failed members reset.

TYPE: UUID | None DEFAULT: None

register_all

Expand complete tasks too, so every edge is recorded.

TYPE: bool DEFAULT: False

on_registry_failure

"raise" (default) or "warn" — carry on through a registry outage; a refusal always raises.

TYPE: OnRegistryFailure DEFAULT: 'raise'

concurrency_config

Build-local limits (how many tasks run at once in this process).

TYPE: ConcurrencyConfig | None DEFAULT: None

concurrency_limiter

A pre-built limiter to use instead of concurrency_config.

TYPE: ConcurrencyLimiter | None DEFAULT: None

claim_config

How claims are waited on and renewed.

TYPE: ClaimConfig | None DEFAULT: None

settings

Environment variables applied for the build's duration (the scope's second half; STARDAG_* / MODAL_* refused). Omitted on a resume, the build's active plan's settings are reused; {} explicitly means none. Two concurrent builds in one process with different settings are refused.

TYPE: Mapping[str, str] | None DEFAULT: None

limit_key_selector

The registry concurrency-limit keys a task runs under, sent with its claim.

TYPE: LimitKeySelector | None DEFAULT: None

description

A description for a new build.

TYPE: str | None DEFAULT: None

raise_on_failure

In FAIL_FAST mode, re-raise the failure's exception (default). False returns the FAILURE summary instead — build id, failed task, error — after the build has stopped at the first failure the same way.

TYPE: bool DEFAULT: True

RETURNS DESCRIPTION
BuildSummary

BuildSummary with status, task counts and build id. STOPPED

BuildSummary

when the registry stopped the build (an operator cancelled it):

BuildSummary

the engine then stops without writing a build status.

Source code in stardag/build/_resident.py
async def build_aio(
    tasks: Sequence[BaseTask] | BaseTask,
    task_executor: TaskExecutorABC | None = None,
    fail_mode: FailMode = FailMode.FAIL_FAST,
    registry: RegistryABC | None = None,
    max_concurrent_discover: int = 50,
    resume_build_id: UUID | None = None,
    register_all: bool = False,
    on_registry_failure: OnRegistryFailure = "raise",
    concurrency_config: ConcurrencyConfig | None = None,
    concurrency_limiter: ConcurrencyLimiter | None = None,
    claim_config: ClaimConfig | None = None,
    settings: Mapping[str, str] | None = None,
    limit_key_selector: LimitKeySelector | None = None,
    description: str | None = None,
    raise_on_failure: bool = True,
) -> BuildSummary:
    """Build tasks concurrently using hybrid async/thread/process execution.

    Walks the DAG from the roots (stopping at complete tasks), registers the
    plan with the registry — roots first, the rest in post-order, then
    sealed — and runs every task whose dependencies are met on
    ``task_executor``. With a registry every execution claims, so two
    builds (or a build and a reactive one) never run a task twice at once.

    Args:
        tasks: The root task(s).
        task_executor: Where tasks run (default:
            :class:`HybridConcurrentTaskExecutor`). An executor running tasks
            on a deployed Modal app makes the build plan under that app's
            current deployment (D13).
        fail_mode: Stop at the first failure (``FAIL_FAST``) or run
            everything whose dependencies are met (``CONTINUE``).
        registry: Default: the configured registry (none: no plan, no
            claims — the build runs purely locally).
        max_concurrent_discover: Completion checks in flight while walking.
        resume_build_id: Resume this build: its plan for this scope is
            reused, observations re-sent, and failed members reset.
        register_all: Expand complete tasks too, so every edge is recorded.
        on_registry_failure: ``"raise"`` (default) or ``"warn"`` — carry on
            through a registry *outage*; a refusal always raises.
        concurrency_config: Build-local limits (how many tasks run at
            once in this process).
        concurrency_limiter: A pre-built limiter to use instead of
            ``concurrency_config``.
        claim_config: How claims are waited on and renewed.
        settings: Environment variables applied for the build's duration
            (the scope's second half; ``STARDAG_*`` / ``MODAL_*`` refused).
            Omitted on a resume, the build's active plan's settings are
            reused; ``{}`` explicitly means none.
            Two concurrent builds in one process with different settings
            are refused.
        limit_key_selector: The registry concurrency-limit keys a task runs
            under, sent with its claim.
        description: A description for a new build.
        raise_on_failure: In ``FAIL_FAST`` mode, re-raise the failure's
            exception (default). False returns the ``FAILURE`` summary
            instead — build id, failed task, error — after the build has
            stopped at the first failure the same way.

    Returns:
        BuildSummary with status, task counts and build id. ``STOPPED``
        when the registry stopped the build (an operator cancelled it):
        the engine then stops without writing a build status.
    """
    roots = [tasks] if isinstance(tasks, BaseTask) else list(tasks)
    for index, task in enumerate(roots):
        if not isinstance(task, BaseTask):
            raise ValueError(
                f"Invalid task at index {index}: {task} (must be BaseTask)"
            )
    registry = registry if registry is not None else registry_provider.get()
    try:
        checked_settings = await resolve_settings_aio(
            registry, resume_build_id, settings
        )
    except SettingsError:
        raise
    except Exception as e:
        handle_registry_error(
            e, "Could not read the resumed build's settings", on_registry_failure
        )
        checked_settings = {}
    with resident_settings(checked_settings):
        engine = _ResidentEngine(
            roots,
            task_executor=task_executor or _default_executor(),
            fail_mode=fail_mode,
            session=ResidentSession(
                registry,
                on_registry_failure=on_registry_failure,
                claim_config=claim_config,
                settings=checked_settings,
                limit_key_selector=limit_key_selector,
            ),
            max_concurrent_discover=max_concurrent_discover,
            register_all=register_all,
            limiter=build_concurrency_limiter(concurrency_config, concurrency_limiter),
            raise_on_failure=raise_on_failure,
        )
        return await engine.run(
            resume_build_id=resume_build_id, description=description
        )

build_sequential

build_sequential(
    tasks,
    registry=None,
    fail_mode=FAIL_FAST,
    dual_run_default="sync",
    resume_build_id=None,
    register_all=False,
    on_registry_failure="raise",
    claim_config=None,
    settings=None,
    limit_key_selector=None,
    description=None,
    max_concurrent_discover=16,
    raise_on_failure=True,
)

Build tasks sequentially, from synchronous code (for debugging).

Task code runs on the calling thread: sync tasks via run(), async-only tasks via asyncio.run(run_aio()) (so not from inside a running event loop), dual tasks by dual_run_default. See :func:build_sequential_aio for the other arguments.

Source code in stardag/build/_sequential.py
def build_sequential(
    tasks: Sequence[BaseTask] | BaseTask,
    registry: RegistryABC | None = None,
    fail_mode: FailMode = FailMode.FAIL_FAST,
    dual_run_default: Literal["sync", "async"] = "sync",
    resume_build_id: UUID | None = None,
    register_all: bool = False,
    on_registry_failure: OnRegistryFailure = "raise",
    claim_config: ClaimConfig | None = None,
    settings: Mapping[str, str] | None = None,
    limit_key_selector: LimitKeySelector | None = None,
    description: str | None = None,
    max_concurrent_discover: int = 16,
    raise_on_failure: bool = True,
) -> BuildSummary:
    """Build tasks sequentially, from synchronous code (for debugging).

    Task code runs on the calling thread: sync tasks via ``run()``,
    async-only tasks via ``asyncio.run(run_aio())`` (so not from inside a
    running event loop), dual tasks by ``dual_run_default``. See
    :func:`build_sequential_aio` for the other arguments.
    """
    roots = _roots(tasks)
    registry = registry if registry is not None else registry_provider.get()
    try:
        checked = resolve_settings(registry, resume_build_id, settings)
    except SettingsError:
        raise
    except Exception as e:
        handle_registry_error(
            e, "Could not read the resumed build's settings", on_registry_failure
        )
        checked = {}
    caller = _CallerThread()
    engine = _SequentialEngine(
        roots,
        session=_session(
            registry,
            on_registry_failure=on_registry_failure,
            claim_config=claim_config,
            settings=checked,
            limit_key_selector=limit_key_selector,
        ),
        runner=_SyncRunner(caller, dual_run_default),
        fail_mode=fail_mode,
        register_all=register_all,
        max_concurrent_discover=max_concurrent_discover,
        raise_on_failure=raise_on_failure,
    )
    with resident_settings(checked):
        loop = asyncio.new_event_loop()
        thread = threading.Thread(
            target=loop.run_forever, name="stardag-sequential", daemon=True
        )
        thread.start()
        try:
            done = asyncio.run_coroutine_threadsafe(
                engine.run(resume_build_id=resume_build_id, description=description),
                loop,
            )
            try:
                caller.serve_until(done)
            except BaseException:
                done.cancel()
                raise
            return done.result()
        finally:
            loop.call_soon_threadsafe(loop.stop)
            thread.join()
            loop.close()

build_sequential_aio async

build_sequential_aio(
    tasks,
    registry=None,
    fail_mode=FAIL_FAST,
    sync_run_default="blocking",
    resume_build_id=None,
    register_all=False,
    on_registry_failure="raise",
    claim_config=None,
    settings=None,
    limit_key_selector=None,
    description=None,
    max_concurrent_discover=16,
    raise_on_failure=True,
)

Build tasks sequentially from async code (for debugging).

PARAMETER DESCRIPTION
tasks

The root task(s).

TYPE: Sequence[BaseTask] | BaseTask

registry

Default: the configured registry (none: no plan, no claims).

TYPE: RegistryABC | None DEFAULT: None

fail_mode

FAIL_FAST (default) or CONTINUE.

TYPE: FailMode DEFAULT: FAIL_FAST

sync_run_default

How a sync-only task runs: "blocking" on the loop, or "thread".

TYPE: Literal['thread', 'blocking'] DEFAULT: 'blocking'

resume_build_id

Resume this build (its plan for this scope is reused, observations re-sent, failed members reset).

TYPE: UUID | None DEFAULT: None

register_all

Expand complete tasks too.

TYPE: bool DEFAULT: False

on_registry_failure

"raise" or "warn" (outages only).

TYPE: OnRegistryFailure DEFAULT: 'raise'

claim_config

Claim waiting and renewal.

TYPE: ClaimConfig | None DEFAULT: None

settings

Environment variables applied for the build's duration. Omitted on a resume, the build's active plan's settings are reused; {} explicitly means none.

TYPE: Mapping[str, str] | None DEFAULT: None

limit_key_selector

Registry concurrency-limit keys per task, sent with its claim.

TYPE: LimitKeySelector | None DEFAULT: None

description

A description for a new build.

TYPE: str | None DEFAULT: None

max_concurrent_discover

Completion checks in flight while walking.

TYPE: int DEFAULT: 16

raise_on_failure

In FAIL_FAST mode, re-raise the failure's exception (default); False returns the FAILURE summary.

TYPE: bool DEFAULT: True

Source code in stardag/build/_sequential.py
async def build_sequential_aio(
    tasks: Sequence[BaseTask] | BaseTask,
    registry: RegistryABC | None = None,
    fail_mode: FailMode = FailMode.FAIL_FAST,
    sync_run_default: Literal["thread", "blocking"] = "blocking",
    resume_build_id: UUID | None = None,
    register_all: bool = False,
    on_registry_failure: OnRegistryFailure = "raise",
    claim_config: ClaimConfig | None = None,
    settings: Mapping[str, str] | None = None,
    limit_key_selector: LimitKeySelector | None = None,
    description: str | None = None,
    max_concurrent_discover: int = 16,
    raise_on_failure: bool = True,
) -> BuildSummary:
    """Build tasks sequentially from async code (for debugging).

    Args:
        tasks: The root task(s).
        registry: Default: the configured registry (none: no plan, no
            claims).
        fail_mode: ``FAIL_FAST`` (default) or ``CONTINUE``.
        sync_run_default: How a sync-only task runs: ``"blocking"`` on the
            loop, or ``"thread"``.
        resume_build_id: Resume this build (its plan for this scope is
            reused, observations re-sent, failed members reset).
        register_all: Expand complete tasks too.
        on_registry_failure: ``"raise"`` or ``"warn"`` (outages only).
        claim_config: Claim waiting and renewal.
        settings: Environment variables applied for the build's duration.
            Omitted on a resume, the build's active plan's settings are
            reused; ``{}`` explicitly means none.
        limit_key_selector: Registry concurrency-limit keys per task, sent
            with its claim.
        description: A description for a new build.
        max_concurrent_discover: Completion checks in flight while walking.
        raise_on_failure: In ``FAIL_FAST`` mode, re-raise the failure's
            exception (default); False returns the ``FAILURE`` summary.
    """
    roots = _roots(tasks)
    registry = registry if registry is not None else registry_provider.get()
    try:
        checked = await resolve_settings_aio(registry, resume_build_id, settings)
    except SettingsError:
        raise
    except Exception as e:
        handle_registry_error(
            e, "Could not read the resumed build's settings", on_registry_failure
        )
        checked = {}
    engine = _SequentialEngine(
        roots,
        session=_session(
            registry,
            on_registry_failure=on_registry_failure,
            claim_config=claim_config,
            settings=checked,
            limit_key_selector=limit_key_selector,
        ),
        runner=_AsyncRunner(sync_run_default),
        fail_mode=fail_mode,
        register_all=register_all,
        max_concurrent_discover=max_concurrent_discover,
        raise_on_failure=raise_on_failure,
    )
    with resident_settings(checked):
        return await engine.run(
            resume_build_id=resume_build_id, description=description
        )

Target Module

target

target_factory_provider module-attribute

target_factory_provider = resource_provider(
    type_=TargetFactory, default_factory=TargetFactory
)

FileSystemTarget

Bases: Target, Protocol

Minimal base protocol for filesystem-backed targets.

Both FileTarget (file-oriented) and DirectoryTarget (directory-oriented) implement this protocol.

FileTarget

FileTarget(uri)

Bases: _FileTargetGeneric[bytes], Protocol

A file-oriented filesystem target with open/read/write capabilities.

Inherits all file I/O methods from _FileTargetGeneric: open(), proxy_path(), exists(), and their async variants. Concrete implementations: LocalFileTarget, RemoteFileTarget, InMemoryFileTarget.

Source code in stardag/target/_base.py
def __init__(self, uri: str) -> None:
    self.uri = uri

DirectoryTarget

DirectoryTarget(uri, prototype)

Bases: FileSystemTarget

A target representing a directory of file targets.

Manages a collection of sub-targets (files) under a common URI prefix, with a flag file to track completion. Sub-targets are created via get_sub_target() or the / operator.

Source code in stardag/target/_base.py
def __init__(
    self,
    uri: str,
    prototype: typing.Type[FileTarget] | typing.Callable[[str], FileTarget],
) -> None:
    self.uri = uri.removesuffix("/") + "/"
    self.prototype = prototype
    self._flag_target = prototype(self.uri[:-1] + "._DONE")
    self._sub_keys: set[str] = set()

exists_aio async

exists_aio()

Async check if directory is marked as done.

Source code in stardag/target/_base.py
async def exists_aio(self) -> bool:
    """Async check if directory is marked as done."""
    return await self._flag_target.exists_aio()

mark_done_aio async

mark_done_aio()

Async version of mark_done().

Source code in stardag/target/_base.py
async def mark_done_aio(self) -> None:
    """Async version of mark_done()."""
    async with self.sub_keys_target().proxy_path_aio("w") as path:
        async with aiofiles.open(path, "w") as f:
            await f.write("\n".join(sorted(self._sub_keys)))
    async with self._flag_target.proxy_path_aio("w") as path:
        async with aiofiles.open(path, "w") as f:
            await f.write("")  # empty file

LoadableSaveableFileSystemTarget

Bases: LoadableSaveableTarget[LoadedT], FileSystemTarget, Generic[LoadedT], Protocol

A filesystem target (file or directory) that supports load/save.

This is the return type of Task.target(). It provides: - load() -> LoadedT and save(obj: LoadedT) (from LoadableSaveableTarget) - uri: str and exists() -> bool (from FileSystemTarget)

LocalFileTarget

LocalFileTarget(uri)

Bases: FileTarget

TODO use luigi-style atomic writes.

Source code in stardag/target/_base.py
def __init__(self, uri: str) -> None:
    # Expand ~ to user home directory
    self.uri = os.path.expanduser(uri)

exists_aio async

exists_aio()

Asynchronously check if the local file exists.

Source code in stardag/target/_base.py
async def exists_aio(self) -> bool:
    """Asynchronously check if the local file exists."""
    return await aiofiles.os.path.exists(self.path)

TargetFactory

TargetFactory(
    target_roots=None, prefix_to_target_prototype=None
)
Source code in stardag/target/_factory.py
def __init__(
    self,
    target_roots: dict[str, str] | None = None,
    prefix_to_target_prototype: PrefixToFileTargetPrototype | None = None,
) -> None:
    # If no target_roots provided, get from central config
    if target_roots is None:
        target_roots = config_provider.get().target.roots

    self.target_roots = {
        key: value.removesuffix("/") + "/" for key, value in target_roots.items()
    }
    self.prefix_to_target_prototype = (
        prefix_to_target_prototype or get_default_prefix_to_target_prototype()
    )

get_file_target

get_file_target(
    relpath, target_root_key=DEFAULT_TARGET_ROOT_KEY
)

Get a file target.

PARAMETER DESCRIPTION
relpath

The path to the target, relative to the configured root path for target_root_key.

TYPE: str

target_root_key

The key to the target root to use.

TYPE: str DEFAULT: DEFAULT_TARGET_ROOT_KEY

RETURNS DESCRIPTION
FileTarget

A file target.

Source code in stardag/target/_factory.py
def get_file_target(
    self,
    relpath: str,
    target_root_key: str = DEFAULT_TARGET_ROOT_KEY,
) -> FileTarget:
    """Get a file target.

    Args:
        relpath: The path to the target, relative to the configured root path for
            `target_root_key`.
        target_root_key: The key to the target root to use.

    Returns:
        A file target.
    """
    if self._is_full_path(relpath):
        path = relpath
    else:
        path = self.get_path(relpath, target_root_key)
    target_prototype = self._get_target_prototype(path)
    return target_prototype(path)

get_directory_target

get_directory_target(
    relpath, target_root_key=DEFAULT_TARGET_ROOT_KEY
)

Get a directory target.

PARAMETER DESCRIPTION
relpath

The path to the target, relative to the configured root path for target_root_key.

TYPE: str

target_root_key

The key to the target root to use.

TYPE: str DEFAULT: DEFAULT_TARGET_ROOT_KEY

RETURNS DESCRIPTION
DirectoryTarget

A directory target.

Source code in stardag/target/_factory.py
def get_directory_target(
    self,
    relpath: str,
    target_root_key: str = DEFAULT_TARGET_ROOT_KEY,
) -> DirectoryTarget:
    """Get a directory target.

    Args:
        relpath: The path to the target, relative to the configured root path for
            `target_root_key`.
        target_root_key: The key to the target root to use.

    Returns:
        A directory target.
    """
    if self._is_full_path(relpath):
        path = relpath
    else:
        path = self.get_path(relpath, target_root_key)
    target_prototype = self._get_target_prototype(path)
    return DirectoryTarget(path, target_prototype)

get_path

get_path(relpath, target_root_key=DEFAULT_TARGET_ROOT_KEY)

Get the full (/"absolute") path (/"URI") to the target.

Source code in stardag/target/_factory.py
def get_path(
    self, relpath: str, target_root_key: str = DEFAULT_TARGET_ROOT_KEY
) -> str:
    """Get the full (/"absolute") path (/"URI") to the target."""
    target_root = self.target_roots.get(target_root_key)
    if target_root is None:
        example_json = json.dumps({target_root_key: "...", "default": "..."})
        raise ValueError(
            f"No target root is configured for key: '{target_root_key}'. "
            f"Available keys are: {list(self.target_roots.keys())}. Set the missing "
            "target root in your registry config or via environment variable: "
            f"`STARDAG_TARGET_ROOTS='{example_json}'`."
        )

    return f"{target_root}{relpath}"

Registry Module

registry

The task registry (v2).

  • :class:RegistryABC: the interface every engine and integration uses.
  • :class:APIRegistry: its implementation over the /api/v2 HTTP API.
  • :class:NoOpRegistry: the default when no registry is configured; the engines make no registry call against it.
  • :data:registry_provider: the configured registry for this process.

The response and registration models are re-exported here; see :mod:stardag.registry._models for the vocabulary (an instance is a registry row under a scope; the Python object is a task object).

registry_provider module-attribute

registry_provider = resource_provider(
    RegistryABC, init_registry
)

APIRegistry

APIRegistry(
    api_url=None,
    timeout=None,
    environment_id=None,
    api_key=None,
)

Bases: APIRegistryReads, RegistryABC

The v2 registry over HTTP.

Stateless with respect to builds (every call names its build, plan or task), so one instance serves many builds — it is a process-wide singleton through registry_provider.

Authentication: an API key (explicit, or STARDAG_API_KEY) or the browser-login JWT of the active profile. The inspecting reads are in :class:~stardag.registry._api_reads.APIRegistryReads.

Source code in stardag/registry/_api_http.py
def __init__(
    self,
    api_url: str | None = None,
    timeout: float | None = None,
    environment_id: str | None = None,
    api_key: str | None = None,
):
    config = config_provider.get()
    reg = config.registry
    resolved_url = api_url or (reg.url if reg else None)
    if not resolved_url:
        raise ValueError(
            "APIRegistry requires a registry URL. "
            "Set STARDAG_API_URL or configure a profile with a registry."
        )
    self.api_url = resolved_url.rstrip("/")
    self.timeout = (
        timeout
        if timeout is not None
        else (reg.timeout if reg else DEFAULT_API_TIMEOUT)
    )
    self.environment_id = environment_id or (reg.environment_id if reg else None)

    resolved_api_key = api_key or (
        reg.auth.api_key.get_secret_value() if reg and reg.auth.api_key else None
    )
    self._auth: httpx.Auth | None
    if resolved_api_key:
        self._auth = StardagAPIKeyAuth(resolved_api_key)
    elif reg and reg.auth.access_token:
        self._auth = StardagTokenAuth(
            access_token=reg.auth.access_token.get_secret_value(),
            workspace_id=reg.workspace_id,
            user_email=reg.auth.user_email,
            registry_url=reg.url,
            registry_name=config.context.registry_name,
        )
        if not self.environment_id:
            logger.warning(
                "APIRegistry: JWT auth requires environment_id. "
                "Run 'stardag config set environment <id>' to set it."
            )
    else:
        self._auth = None
        logger.warning(
            "APIRegistry initialized without authentication. "
            "Run 'stardag auth login' or set STARDAG_API_KEY env var."
        )
    self._client: httpx.Client | None = None
    self._async_client: httpx.AsyncClient | None = None
    self._async_client_loop: asyncio.AbstractEventLoop | None = None

RegistryABC

Bases: RegistryReadsABC

The v2 registry client interface. See the module docstring. The inspecting reads (a plan, paged lists, a task's executions and events, a deployment) are in :class:~stardag.registry._base_reads.RegistryReadsABC.

build_create

build_create(
    *,
    root_task_ids,
    build_id=None,
    name=None,
    description=None,
    executor_metadata=None,
)

POST /builds: a RUNNING build requesting root_task_ids. Idempotent on a client-minted build_id.

Source code in stardag/registry/_base.py
def build_create(
    self,
    *,
    root_task_ids: Sequence[str],
    build_id: UUID | None = None,
    name: str | None = None,
    description: str | None = None,
    executor_metadata: dict[str, Any] | None = None,
) -> BuildInfo:
    """``POST /builds``: a RUNNING build requesting ``root_task_ids``.
    Idempotent on a client-minted ``build_id``."""
    raise _missing(self, "build_create")

build_resume

build_resume(
    build_id,
    *,
    deployment_id=None,
    settings=None,
    executor_metadata=None,
)

POST /builds/{id}/resume: make the build RUNNING again and, when the caller names its scope, reuse or reactivate the plan for it.

Source code in stardag/registry/_base.py
def build_resume(
    self,
    build_id: UUID,
    *,
    deployment_id: UUID | None = None,
    settings: Mapping[str, str] | None = None,
    executor_metadata: dict[str, Any] | None = None,
) -> ResumeResult:
    """``POST /builds/{id}/resume``: make the build RUNNING again and,
    when the caller names its scope, reuse or reactivate the plan for
    it."""
    raise _missing(self, "build_resume")

build_complete

build_complete(build_id, *, force=False)

POST /builds/{id}/complete, refused (409 plan_incomplete) unless the active plan is sealed and every non-excluded member is COMPLETED; force overrides outstanding members only.

Source code in stardag/registry/_base.py
def build_complete(self, build_id: UUID, *, force: bool = False) -> BuildInfo:
    """``POST /builds/{id}/complete``, refused (409 ``plan_incomplete``)
    unless the active plan is sealed and every non-excluded member is
    COMPLETED; ``force`` overrides outstanding members only."""
    raise _missing(self, "build_complete")

build_exit_early

build_exit_early(build_id)

POST /builds/{id}/exit-early: the resident driver stops; nothing is released (its in-flight executions keep reporting).

Source code in stardag/registry/_base.py
def build_exit_early(self, build_id: UUID) -> BuildInfo:
    """``POST /builds/{id}/exit-early``: the resident driver stops;
    nothing is released (its in-flight executions keep reporting)."""
    raise _missing(self, "build_exit_early")

build_list_running

build_list_running(*, reactive_app_name=None, limit=100)

RUNNING builds, most recently active first (the watchdog sweep).

Source code in stardag/registry/_base.py
def build_list_running(
    self, *, reactive_app_name: str | None = None, limit: int = 100
) -> list[UUID]:
    """RUNNING builds, most recently active first (the watchdog sweep)."""
    raise _missing(self, "build_list_running")

plan_create

plan_create(
    build_id, *, plan_id, deployment_id, settings, roots
)

POST /builds/{id}/plans: look up or create the build's plan for (deployment, settings), admitting roots first and unexpanded. An existing plan is returned as it is (its own id, not plan_id).

Source code in stardag/registry/_base.py
def plan_create(
    self,
    build_id: UUID,
    *,
    plan_id: UUID,
    deployment_id: UUID,
    settings: Mapping[str, str],
    roots: Sequence[RegistrationItem],
) -> PlanInfo:
    """``POST /builds/{id}/plans``: look up or create the build's plan
    for ``(deployment, settings)``, admitting ``roots`` first and
    unexpanded. An existing plan is returned as it is (its own id, not
    ``plan_id``)."""
    raise _missing(self, "plan_create")

plan_register_members

plan_register_members(plan_id, items)

POST /plans/{id}/members: one chunk (at most 1000 items), in one transaction.

Source code in stardag/registry/_base.py
def plan_register_members(
    self, plan_id: UUID, items: Sequence[RegistrationItem]
) -> MembersResult:
    """``POST /plans/{id}/members``: one chunk (at most 1000 items), in
    one transaction."""
    raise _missing(self, "plan_register_members")

plan_seal

plan_seal(plan_id)

POST /plans/{id}/seal: verify the static phase and seal (a replacement activates here).

Source code in stardag/registry/_base.py
def plan_seal(self, plan_id: UUID) -> PlanInfo:
    """``POST /plans/{id}/seal``: verify the static phase and seal (a
    replacement activates here)."""
    raise _missing(self, "plan_seal")

plan_roots_info

plan_roots_info(plan_id)

GET /plans/{id}/roots: the plan's scope and root members.

Source code in stardag/registry/_base.py
def plan_roots_info(self, plan_id: UUID) -> PlanRoots:
    """``GET /plans/{id}/roots``: the plan's scope and root members."""
    raise _missing(self, "plan_roots_info")

plan_roots

plan_roots(plan_id)

The plan's root members with their instance bodies (rollover).

Source code in stardag/registry/_base.py
def plan_roots(self, plan_id: UUID) -> list[FrontierMember]:
    """The plan's root members with their instance bodies (rollover)."""
    raise _missing(self, "plan_roots")

build_skip_blocked

build_skip_blocked(build_id)

POST /builds/{id}/skip-blocked: mark the active plan's members transitively blocked by a failed, cancelled or skipped upstream SKIPPED; returns their task ids.

Source code in stardag/registry/_base.py
def build_skip_blocked(self, build_id: UUID) -> list[str]:
    """``POST /builds/{id}/skip-blocked``: mark the active plan's members
    transitively blocked by a failed, cancelled or skipped upstream
    SKIPPED; returns their task ids."""
    raise _missing(self, "build_skip_blocked")

member_start

member_start(
    plan_id,
    task_id,
    *,
    execution_id,
    claim=True,
    claim_ttl_seconds=None,
    executor=None,
    executor_ref=None,
    executor_metadata=None,
    limit_keys=(),
)

POST /plans/{id}/members/{task_id}/start.

A claiming start takes the claim for execution_id (minted by the caller before the spawn) and is the decision: refused 409 with task_already_completed, task_already_running, upstream_incomplete, member_excluded, plan_superseded or concurrency_limit_reached. A retried granted start is a no-op. A non-claiming start is the holder's own "I am running" report, with the executor details the claim could not know.

Source code in stardag/registry/_base.py
def member_start(
    self,
    plan_id: UUID,
    task_id: str,
    *,
    execution_id: UUID,
    claim: bool = True,
    claim_ttl_seconds: int | None = None,
    executor: str | None = None,
    executor_ref: str | None = None,
    executor_metadata: dict[str, Any] | None = None,
    limit_keys: Sequence[str] = (),
) -> TransitionResult:
    """``POST /plans/{id}/members/{task_id}/start``.

    A **claiming** start takes the claim for ``execution_id`` (minted by
    the caller before the spawn) and is the decision: refused 409 with
    ``task_already_completed``, ``task_already_running``,
    ``upstream_incomplete``, ``member_excluded``, ``plan_superseded`` or
    ``concurrency_limit_reached``. A retried granted start is a no-op. A
    **non-claiming** start is the holder's own "I am running" report,
    with the executor details the claim could not know."""
    raise _missing(self, "member_start")

member_yield

member_yield(
    plan_id,
    task_id,
    *,
    execution_id,
    deployment_id,
    batch_id,
    items,
    yielded,
    suspend,
)

POST /plans/{id}/members/{task_id}/yield: one yield batch in one transaction — the children and their static closure land, the parent gets dynamic edges to yielded (instance hashes), and with suspend the parent is SUSPENDED and its claim released. A batch re-delivered under the same batch_id is replayed. Refused 409 deployment_mismatch when deployment_id is not the plan's.

Source code in stardag/registry/_base.py
def member_yield(
    self,
    plan_id: UUID,
    task_id: str,
    *,
    execution_id: UUID,
    deployment_id: UUID,
    batch_id: UUID,
    items: Sequence[RegistrationItem],
    yielded: Sequence[str],
    suspend: bool,
) -> YieldResult:
    """``POST /plans/{id}/members/{task_id}/yield``: one yield batch in
    one transaction — the children and their static closure land, the
    parent gets dynamic edges to ``yielded`` (instance hashes), and with
    ``suspend`` the parent is SUSPENDED and its claim released. A batch
    re-delivered under the same ``batch_id`` is replayed. Refused 409
    ``deployment_mismatch`` when ``deployment_id`` is not the plan's."""
    raise _missing(self, "member_yield")

member_retry

member_retry(plan_id, task_id)

Reset to PENDING (the fail mode's retry). Idempotent by state; refused 409 on COMPLETED and on a live claim.

Source code in stardag/registry/_base.py
def member_retry(self, plan_id: UUID, task_id: str) -> TransitionResult:
    """Reset to PENDING (the fail mode's retry). Idempotent by state;
    refused 409 on COMPLETED and on a live claim."""
    raise _missing(self, "member_retry")

member_interrupt

member_interrupt(
    plan_id, task_id, *, execution_id, error_message=None
)

The platform ended the execution and nothing will restart it: INTERRUPTED (actionable), claim released.

Source code in stardag/registry/_base.py
def member_interrupt(
    self,
    plan_id: UUID,
    task_id: str,
    *,
    execution_id: UUID,
    error_message: str | None = None,
) -> TransitionResult:
    """The platform ended the execution and nothing will restart it:
    INTERRUPTED (actionable), claim released."""
    raise _missing(self, "member_interrupt")

member_preempt

member_preempt(plan_id, task_id, *, execution_id)

The backend restarts the execution itself: no status change, the claim kept but due to lapse soon. Not an end: the restart reports under the same execution id, and its non-claiming start restores the claim's TTL.

Source code in stardag/registry/_base.py
def member_preempt(
    self, plan_id: UUID, task_id: str, *, execution_id: UUID
) -> TransitionResult:
    """The backend restarts the execution itself: no status change, the
    claim kept but due to lapse soon. Not an end: the restart reports
    under the same execution id, and its non-claiming start restores
    the claim's TTL."""
    raise _missing(self, "member_preempt")

member_cancel

member_cancel(plan_id, task_id)

One task's cancel, by the build holding its claim (409 not_claim_holder otherwise).

Source code in stardag/registry/_base.py
def member_cancel(self, plan_id: UUID, task_id: str) -> TransitionResult:
    """One task's cancel, by the build holding its claim (409
    ``not_claim_holder`` otherwise)."""
    raise _missing(self, "member_cancel")

member_skip

member_skip(plan_id, task_id)

A member that cannot run because an upstream failed: SKIPPED. A scheduling decision naming no execution; 409 task_not_skippable on FAILED / CANCELLED.

Source code in stardag/registry/_base.py
def member_skip(self, plan_id: UUID, task_id: str) -> TransitionResult:
    """A member that cannot run because an upstream failed: SKIPPED. A
    scheduling decision naming no execution; 409 ``task_not_skippable``
    on FAILED / CANCELLED."""
    raise _missing(self, "member_skip")

member_exclude

member_exclude(plan_id, task_id, *, reason=None)

An operator gives up on a member in this plan. Cascades to its downstream closure; an excluded root fails the build. The global status is untouched.

Source code in stardag/registry/_base.py
def member_exclude(
    self, plan_id: UUID, task_id: str, *, reason: str | None = None
) -> ExclusionResult:
    """An operator gives up on a member in this plan. Cascades to its
    downstream closure; an excluded root fails the build. The global
    status is untouched."""
    raise _missing(self, "member_exclude")

member_discovery_failed

member_discovery_failed(plan_id, task_id, *, error)

The tick could not discover a member (its requires() raised, its body did not rehydrate): excluded as discovery_failed, with the same cascade as :meth:member_exclude.

Source code in stardag/registry/_base.py
def member_discovery_failed(
    self, plan_id: UUID, task_id: str, *, error: str
) -> ExclusionResult:
    """The tick could not discover a member (its ``requires()`` raised,
    its body did not rehydrate): excluded as ``discovery_failed``, with
    the same cascade as :meth:`member_exclude`."""
    raise _missing(self, "member_discovery_failed")

claim_renew

claim_renew(
    task_id, *, execution_id, claim_ttl_seconds=None
)

POST /tasks/{task_id}/claim/renew: extend an in-process execution's claim (D11). Refused 409 claim_not_held unless execution_id holds the live claim.

Source code in stardag/registry/_base.py
def claim_renew(
    self,
    task_id: str,
    *,
    execution_id: UUID,
    claim_ttl_seconds: int | None = None,
) -> TransitionResult:
    """``POST /tasks/{task_id}/claim/renew``: extend an in-process
    execution's claim (D11). Refused 409 ``claim_not_held`` unless
    ``execution_id`` holds the live claim."""
    raise _missing(self, "claim_renew")

build_list_executions

build_list_executions(
    build_id,
    *,
    not_in_current_plan=False,
    include_ended=False,
)

GET /builds/{id}/executions: the build's executions with no end reported (builds stop, a worker's cancellation checkpoint); not_in_current_plan keeps the orphans; include_ended lists the whole ledger (every execution granted, ended or not).

Source code in stardag/registry/_base.py
def build_list_executions(
    self,
    build_id: UUID,
    *,
    not_in_current_plan: bool = False,
    include_ended: bool = False,
) -> list[ExecutionInfo]:
    """``GET /builds/{id}/executions``: the build's executions with no end
    reported (``builds stop``, a worker's cancellation checkpoint);
    ``not_in_current_plan`` keeps the orphans; ``include_ended`` lists
    the whole ledger (every execution granted, ended or not)."""
    raise _missing(self, "build_list_executions")

execution_report_stopped

execution_report_stopped(
    execution_id, *, outcome="stopped"
)

POST /executions/{id}/stopped: an operator ends an execution — stopped (its call was cancelled) or lost (it could not be, and no report of it will ever be applied). If it is the task's current execution with its claim unreleased, the claim is released cancelled and the task is CANCELLED (a revocation is not a result).

Source code in stardag/registry/_base.py
def execution_report_stopped(
    self, execution_id: UUID, *, outcome: StopOutcome = "stopped"
) -> TransitionResult:
    """``POST /executions/{id}/stopped``: an operator ends an execution —
    ``stopped`` (its call was cancelled) or ``lost`` (it could not be,
    and no report of it will ever be applied). If it is the task's
    current execution with its claim unreleased, the claim is released
    ``cancelled`` and the task is CANCELLED (a revocation is not a
    result)."""
    raise _missing(self, "execution_report_stopped")

task_get

task_get(task_id)

GET /tasks/{task_id}: a completion's identity and state, with its instances (each a construction under one scope), newest first.

Source code in stardag/registry/_base.py
def task_get(self, task_id: str) -> TaskInfo:
    """``GET /tasks/{task_id}``: a completion's identity and state, with
    its instances (each a construction under one scope), newest first."""
    raise _missing(self, "task_get")

task_list_artifacts

task_list_artifacts(task_id)

GET /tasks/{task_id}/artifacts.

Source code in stardag/registry/_base.py
def task_list_artifacts(self, task_id: str) -> list[TaskArtifactInfo]:
    """``GET /tasks/{task_id}/artifacts``."""
    raise _missing(self, "task_list_artifacts")

task_upload_artifacts

task_upload_artifacts(
    plan_id, task_id, artifacts, *, execution_id=None
)

POST /plans/{plan_id}/members/{task_id}/artifacts: upsert artifacts onto the task named by its membership of plan_id (404 not_a_member otherwise). Artifacts belong to the task once uploaded, not the plan or execution — execution_id is informational only.

Source code in stardag/registry/_base.py
def task_upload_artifacts(
    self,
    plan_id: UUID,
    task_id: str,
    artifacts: "Sequence[Artifact]",
    *,
    execution_id: UUID | None = None,
) -> None:
    """``POST /plans/{plan_id}/members/{task_id}/artifacts``: upsert
    artifacts onto the task named by its membership of ``plan_id`` (404
    ``not_a_member`` otherwise). Artifacts belong to the task once
    uploaded, not the plan or execution — ``execution_id`` is
    informational only."""
    raise _missing(self, "task_upload_artifacts")

deployment_create

deployment_create(
    *,
    kind,
    code_id,
    deployment_id=None,
    app_name=None,
    image_id=None,
    modal_app_id=None,
)

POST /deployments. A Modal deployment is created before its deploy (the server assigns generation) and activated after; a local one is looked up or created by code_id and is born activated.

Source code in stardag/registry/_base.py
def deployment_create(
    self,
    *,
    kind: DeploymentKind,
    code_id: str,
    deployment_id: UUID | None = None,
    app_name: str | None = None,
    image_id: str | None = None,
    modal_app_id: str | None = None,
) -> DeploymentInfo:
    """``POST /deployments``. A Modal deployment is created **before**
    its deploy (the server assigns ``generation``) and activated after;
    a local one is looked up or created by ``code_id`` and is born
    activated."""
    raise _missing(self, "deployment_create")

deployment_activate

deployment_activate(
    deployment_id, *, modal_app_id=None, image_id=None
)

POST /deployments/{id}/activate, with what only the finished deploy knows (a given value fills a NULL or must match).

Source code in stardag/registry/_base.py
def deployment_activate(
    self,
    deployment_id: UUID,
    *,
    modal_app_id: str | None = None,
    image_id: str | None = None,
) -> DeploymentInfo:
    """``POST /deployments/{id}/activate``, with what only the finished
    deploy knows (a given value fills a NULL or must match)."""
    raise _missing(self, "deployment_activate")

deployment_list

deployment_list(
    *, kind=None, app_name=None, current=False, limit=100
)

Newest first. current=True keeps one row per app: its activated deployment with the highest generation.

Source code in stardag/registry/_base.py
def deployment_list(
    self,
    *,
    kind: DeploymentKind | None = None,
    app_name: str | None = None,
    current: bool = False,
    limit: int = 100,
) -> list[DeploymentInfo]:
    """Newest first. ``current=True`` keeps one row per app: its
    activated deployment with the highest generation."""
    raise _missing(self, "deployment_list")

concurrency_limit_set

concurrency_limit_set(key, max_concurrent)

PUT /concurrency-limits/{key}: create or replace the cap on how many tasks carrying key may hold a live claim at once.

Source code in stardag/registry/_base.py
def concurrency_limit_set(self, key: str, max_concurrent: int) -> None:
    """``PUT /concurrency-limits/{key}``: create or replace the cap on
    how many tasks carrying ``key`` may hold a live claim at once."""
    raise _missing(self, "concurrency_limit_set")

concurrency_limit_delete

concurrency_limit_delete(key)

DELETE /concurrency-limits/{key} (404 unknown_limit).

Source code in stardag/registry/_base.py
def concurrency_limit_delete(self, key: str) -> None:
    """``DELETE /concurrency-limits/{key}`` (404 ``unknown_limit``)."""
    raise _missing(self, "concurrency_limit_delete")

concurrency_limit_list

concurrency_limit_list()

GET /concurrency-limits: key -> max_concurrent.

Source code in stardag/registry/_base.py
def concurrency_limit_list(self) -> dict[str, int]:
    """``GET /concurrency-limits``: key -> max_concurrent."""
    raise _missing(self, "concurrency_limit_list")

concurrency_limit_list_detailed

concurrency_limit_list_detailed(*, include_holders=False)

GET /concurrency-limits, parsed in full: each key's cap, how many slots are in use, and — with include_holders — by which tasks. What stardag concurrency-limits list/holders render; :meth:concurrency_limit_list stays the plain key -> cap mapping for callers that only want the configuration.

Source code in stardag/registry/_base.py
def concurrency_limit_list_detailed(
    self, *, include_holders: bool = False
) -> list[ConcurrencyLimitInfo]:
    """``GET /concurrency-limits``, parsed in full: each key's cap, how
    many slots are in use, and — with ``include_holders`` — by which
    tasks. What ``stardag concurrency-limits list``/``holders`` render;
    :meth:`concurrency_limit_list` stays the plain key -> cap mapping
    for callers that only want the configuration."""
    raise _missing(self, "concurrency_limit_list_detailed")

build_set_reactive_meta

build_set_reactive_meta(
    build_id, *, app_name, tick_kwargs=None
)

Mark the build reactively scheduled by app_name; None tick_kwargs keeps the stored configuration.

Source code in stardag/registry/_base.py
def build_set_reactive_meta(
    self,
    build_id: UUID,
    *,
    app_name: str,
    tick_kwargs: dict[str, Any] | None = None,
) -> BuildInfo:
    """Mark the build reactively scheduled by ``app_name``; ``None``
    ``tick_kwargs`` keeps the stored configuration."""
    raise _missing(self, "build_set_reactive_meta")

scheduler_lease_release

scheduler_lease_release(build_id, *, owner_id)

Drop the lease if owner_id still holds it; held reports whether it did (a lost tick cannot clear its successor's lease).

Source code in stardag/registry/_base.py
def scheduler_lease_release(
    self, build_id: UUID, *, owner_id: str
) -> SchedulerLeaseResult:
    """Drop the lease if ``owner_id`` still holds it; ``held`` reports
    whether it did (a lost tick cannot clear its successor's lease)."""
    raise _missing(self, "scheduler_lease_release")

close

close()

Release connections (a no-op for registries that hold none).

Source code in stardag/registry/_base.py
def close(self) -> None:
    """Release connections (a no-op for registries that hold none)."""

NoOpRegistry

Bases: RegistryABC

No registry configured.

The engines recognise it by exact type and make no registry call at all (D11: a single-process build without a registry has no plan and no claims). Its methods raise, so a code path that reaches one by mistake fails loudly rather than pretending to have recorded something.

Configuration

config

Centralized configuration for Stardag SDK.

This module provides a unified configuration system that consolidates: - Target factory settings (target roots) - Registry settings (URL, workspace, environment, auth, timeout) - Config context (provenance: which profile/registry name was used)

Configuration is loaded from multiple sources with the following priority: 1. Environment variables (STARDAG_*) 2. Project config (.stardag/config.toml in working directory or parents) 3. User config (~/.stardag/config.toml) 4. Defaults

Usage

from stardag.config import get_config

config = get_config() if config.registry: print(config.registry.url) print(config.target.roots)

Environment Variables (highest priority): STARDAG_PROFILE - Profile name to use (looks up in config.toml) STARDAG_API_URL - Registry API URL override STARDAG_REGISTRY_URL - Deprecated alias for STARDAG_API_URL STARDAG_WORKSPACE_ID - Direct workspace ID override STARDAG_ENVIRONMENT_ID - Direct environment ID override STARDAG_API_KEY - API key for authentication STARDAG_TARGET_ROOTS - JSON dict of target roots (override) STARDAG_NO_REGISTRY - Set to 1/true to force offline/local mode

load_config

load_config(use_project_config=True)

Load configuration from all sources.

Priority (highest to lowest): 1. Environment variables (STARDAG_*) 2. Project config (.stardag/config.toml in repo) 3. User config (~/.stardag/config.toml) 4. Defaults

PARAMETER DESCRIPTION
use_project_config

Whether to load .stardag/config.toml from project.

TYPE: bool DEFAULT: True

RETURNS DESCRIPTION
StardagConfig

Fully resolved StardagConfig (actual type is StardagConfig).

Source code in stardag/config/loader.py
def load_config(
    use_project_config: bool = True,
) -> StardagConfig:
    """Load configuration from all sources.

    Priority (highest to lowest):
    1. Environment variables (STARDAG_*)
    2. Project config (.stardag/config.toml in repo)
    3. User config (~/.stardag/config.toml)
    4. Defaults

    Args:
        use_project_config: Whether to load .stardag/config.toml from project.

    Returns:
        Fully resolved StardagConfig (actual type is StardagConfig).
    """
    # 1. Load env vars first (highest priority)
    env_settings = StardagSettings()

    # Short-circuit: STARDAG_NO_REGISTRY forces offline/local mode
    if env_settings.no_registry:
        env_target_roots = _parse_target_roots_from_env()
        target_roots = env_target_roots or {
            DEFAULT_TARGET_ROOT_KEY: DEFAULT_TARGET_ROOT
        }
        return StardagConfig(
            registry=None,
            target=TargetConfig(roots=target_roots),
        )

    # 2. Load user and project TOML configs
    user_toml = load_toml_file(get_user_config_path())
    project_toml = {}
    if use_project_config:
        project_path = find_project_config()
        if project_path:
            project_toml = load_toml_file(project_path)

    # Merge configs (project overrides user)
    toml_config = _merge_toml_configs(user_toml, project_toml)

    # 3. Resolve profile -> (registry, user, workspace, environment)
    profile_name: str | None = None
    registry_name: str | None = None
    registry_url: str | None = None
    user: str | None = None
    workspace_id: str | None = None
    environment_id: str | None = None

    # Resolve API URL: STARDAG_API_URL (canonical) or STARDAG_REGISTRY_URL (deprecated)
    explicit_url = env_settings.api_url
    if not explicit_url:
        legacy_url = os.environ.get("STARDAG_REGISTRY_URL")
        if legacy_url:
            import warnings

            warnings.warn(
                "STARDAG_REGISTRY_URL is deprecated, use STARDAG_API_URL instead.",
                DeprecationWarning,
                stacklevel=2,
            )
            explicit_url = legacy_url

    # Check for direct env var overrides first
    if explicit_url:
        registry_url = explicit_url
        workspace_id = env_settings.workspace_id
        environment_id = env_settings.environment_id
        # Even with direct env var overrides, try to inherit user/registry_name
        # from the active profile so that token auth (OIDC refresh) still works.
        _profile_name = env_settings.profile or toml_config.default.get("profile")
        if _profile_name:
            _profile = toml_config.profile.get(_profile_name)
            if _profile:
                profile_name = _profile_name
                registry_name = _profile.registry
                user = _profile.user
    # Then check for profile-based config
    elif env_settings.profile:
        profile_name = env_settings.profile
    # Fall back to default profile from config
    elif toml_config.default.get("profile"):
        profile_name = toml_config.default["profile"]

    # If we have a profile, look it up
    if profile_name and not registry_url:
        profile = toml_config.profile.get(profile_name)
        if profile:
            registry_name = profile.registry
            user = profile.user  # Optional user for multi-user support
            workspace_value = profile.workspace  # Could be slug or ID
            environment_value = profile.environment  # Could be slug or ID

            # Look up registry URL from registry name
            registry_url_from_toml = toml_config.registry.get(registry_name)
            if registry_url_from_toml:
                registry_url = registry_url_from_toml
            else:
                logger.warning(
                    f"Profile '{profile_name}' references unknown registry '{registry_name}'"
                )

            # Resolve workspace slug to ID if needed
            if _looks_like_uuid(workspace_value):
                workspace_id = workspace_value
            else:
                # Try to resolve from cache
                cached_workspace_id = get_cached_workspace_id(
                    registry_name, workspace_value
                )
                if cached_workspace_id:
                    workspace_id = cached_workspace_id
                else:
                    # Store the slug - will need to be resolved at runtime
                    workspace_id = workspace_value
                    logger.debug(
                        f"Workspace '{workspace_value}' is a slug, not cached. "
                        "Run 'stardag auth refresh' to resolve."
                    )

            # Resolve environment slug to ID if needed
            if _looks_like_uuid(environment_value):
                environment_id = environment_value
            elif workspace_id and _looks_like_uuid(workspace_id):
                # Can only resolve environment if we have a resolved workspace ID
                cached_env_id = get_cached_environment_id(
                    registry_name, workspace_id, environment_value
                )
                if cached_env_id:
                    environment_id = cached_env_id
                else:
                    # Store the slug - will need to be resolved at runtime
                    environment_id = environment_value
                    logger.debug(
                        f"Environment '{environment_value}' is a slug, not cached. "
                        "Run 'stardag auth refresh' to resolve."
                    )
            else:
                # Workspace is not resolved, can't resolve environment either
                environment_id = environment_value
        else:
            logger.warning(f"Profile '{profile_name}' not found in config")

    # 4. Resolve target roots
    # Priority: env > cached > default
    target_roots: dict[str, str]
    env_target_roots = _parse_target_roots_from_env()
    if env_target_roots:
        target_roots = env_target_roots
    elif registry_url and workspace_id and environment_id:
        cached_roots = get_cached_target_roots(
            registry_url, workspace_id, environment_id
        )
        if cached_roots:
            target_roots = cached_roots
        else:
            target_roots = {DEFAULT_TARGET_ROOT_KEY: DEFAULT_TARGET_ROOT}
    else:
        target_roots = {DEFAULT_TARGET_ROOT_KEY: DEFAULT_TARGET_ROOT}

    # 5. Load access token from cache (if we have profile info)
    # If token is expired, try to refresh it automatically
    access_token: str | None = None
    if registry_name and workspace_id and user:
        token_cache_path = get_access_token_cache_path(
            registry_name, workspace_id, user
        )
        if token_cache_path.exists():
            token_data = load_json_file(token_cache_path)
            # Check if token is still valid
            import time

            expires_at = token_data.get("expires_at", 0)
            if expires_at > time.time():
                access_token = token_data.get("access_token")

        # If no valid token in cache, try to refresh it
        if not access_token:
            try:
                from stardag.registry._auth import (
                    ensure_access_token as _ensure_token,
                )

                access_token = _ensure_token(
                    registry_name, workspace_id, user, registry_url=registry_url
                )
            except Exception:
                # Silently fail - user can manually refresh with `stardag auth refresh`
                pass

    # 6. Get API key from env
    api_key_raw = os.environ.get("STARDAG_API_KEY")
    api_key: SecretStr | None = env_settings.api_key or (
        SecretStr(api_key_raw) if api_key_raw else None
    )

    # 7. Build canonical RegistryConfig (or None for offline mode)
    registry_cfg: RegistryConfig | None = None
    if registry_url:
        registry_cfg = RegistryConfig(
            url=registry_url,
            workspace_id=workspace_id or "",
            environment_id=environment_id or "",
            auth=RegistryAuth(
                api_key=api_key,
                user_email=user,
                access_token=SecretStr(access_token) if access_token else None,
            ),
            timeout=env_settings.api_timeout or DEFAULT_API_TIMEOUT,
        )

    return StardagConfig(
        registry=registry_cfg,
        target=TargetConfig(roots=target_roots),
        context=ConfigContext(
            profile=profile_name,
            registry_name=registry_name,
        ),
    )