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 withload() -> 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.
TaskLoads
module-attribute
¶
TaskLoads = typing.Annotated[
LoadableTask[LoadedT_co], Polymorphic()
]
TaskStruct
module-attribute
¶
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),versionand 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
¶
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
¶
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__
¶
Validate that subclasses implement either run() or run_aio().
Also wraps run() and run_aio() methods with precheck validation.
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
¶
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
¶
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
¶
Asynchronously return artifacts to be stored in the registry after task completion.
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
¶
Override PolymorphicRoot.resolve to handle AliasTask deserialization.
from_registry
classmethod
¶
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:
|
registry
|
An optional registry instance to use for loading metadata. If not
provided, the default registry from
TYPE:
|
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).
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.
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>
__map_generic_args_to_ancestor__
classmethod
¶
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:
|
args
|
The generic args of this class (e.g., (str,) for Task[str])
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
tuple | None
|
The mapped args for the ancestor, or None if mapping is not applicable. |
task
¶
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). |
task_executor
|
Where tasks run (default:
:class:
TYPE:
|
fail_mode
|
Stop at the first failure (
TYPE:
|
registry
|
Default: the configured registry (none: no plan, no claims — the build runs purely locally).
TYPE:
|
max_concurrent_discover
|
Completion checks in flight while walking.
TYPE:
|
resume_build_id
|
Resume this build: its plan for this scope is reused, observations re-sent, and failed members reset.
TYPE:
|
register_all
|
Expand complete tasks too, so every edge is recorded.
TYPE:
|
on_registry_failure
|
TYPE:
|
concurrency_config
|
Build-local limits (how many tasks run at once in this process).
TYPE:
|
concurrency_limiter
|
A pre-built limiter to use instead of
TYPE:
|
claim_config
|
How claims are waited on and renewed.
TYPE:
|
settings
|
Environment variables applied for the build's duration
(the scope's second half;
TYPE:
|
limit_key_selector
|
The registry concurrency-limit keys a task runs under, sent with its claim.
TYPE:
|
description
|
A description for a new build.
TYPE:
|
raise_on_failure
|
In
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
BuildSummary
|
BuildSummary with status, task counts and build id. |
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). |
registry
|
Default: the configured registry (none: no plan, no claims).
TYPE:
|
fail_mode
|
TYPE:
|
sync_run_default
|
How a sync-only task runs:
TYPE:
|
resume_build_id
|
Resume this build (its plan for this scope is reused, observations re-sent, failed members reset).
TYPE:
|
register_all
|
Expand complete tasks too.
TYPE:
|
on_registry_failure
|
TYPE:
|
claim_config
|
Claim waiting and renewal.
TYPE:
|
settings
|
Environment variables applied for the build's duration.
Omitted on a resume, the build's active plan's settings are
reused;
TYPE:
|
limit_key_selector
|
Registry concurrency-limit keys per task, sent with its claim.
TYPE:
|
description
|
A description for a new build.
TYPE:
|
max_concurrent_discover
|
Completion checks in flight while walking.
TYPE:
|
raise_on_failure
|
In
TYPE:
|
namespace
¶
Set the task namespace for the module and any submodules.
| PARAMETER | DESCRIPTION |
|---|---|
namespace
|
The namespace to set for the module.
TYPE:
|
scope
|
The module scope, typically passed as
TYPE:
|
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
¶
Set the task namespace for the module to the module import path.
| PARAMETER | DESCRIPTION |
|---|---|
scope
|
The module scope, typically passed as
TYPE:
|
Usage:
get_file_target
¶
Get a file target for the given relative path.
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
¶
Summary of a build execution.
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
__repr__
¶
Return a human-readable summary of the build.
Source code in stardag/build/_base.py
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:
|
max_async_workers
|
Maximum concurrent async tasks (semaphore-based).
TYPE:
|
max_thread_workers
|
Maximum concurrent thread pool workers.
TYPE:
|
max_process_workers
|
Maximum concurrent process pool workers.
TYPE:
|
Source code in stardag/build/_concurrent.py
setup
async
¶
Initialize worker pools.
Source code in stardag/build/_concurrent.py
teardown
async
¶
Shutdown worker pools.
Source code in stardag/build/_concurrent.py
get_executor_details
async
¶
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
submit
async
¶
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
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
¶
Execute a task in-process (or block on a remote one).
| RETURNS | DESCRIPTION |
|---|---|
None | TaskStruct | TaskExecutionError
|
|
None | TaskStruct | TaskExecutionError
|
|
None | TaskStruct | TaskExecutionError
|
|
Source code in stardag/build/_base.py
setup
abstractmethod
async
¶
teardown
abstractmethod
async
¶
cancel
async
¶
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
get_executor_metadata
async
¶
Descriptive executor metadata for executions of task, without
starting anything (stamped on the claiming start). Best-effort.
get_executor_details
async
¶
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
execution_timeout_seconds
¶
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
reports_lifecycle
¶
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
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
supports_detached
¶
Whether this executor can run task as a detached execution
(one that survives the orchestrator). Default: False.
submit_detached
async
¶
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
cancel_detached
async
¶
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
can_spawn_scheduler_ticks
¶
Whether :meth:spawn_scheduler_tick reaches a deployed tick
(a resident build then drains the registry's wake candidates).
spawn_scheduler_tick
¶
Spawn a reactive scheduler tick for build_id on app_name.
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
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). |
task_executor
|
Where tasks run (default:
:class:
TYPE:
|
fail_mode
|
Stop at the first failure (
TYPE:
|
registry
|
Default: the configured registry (none: no plan, no claims — the build runs purely locally).
TYPE:
|
max_concurrent_discover
|
Completion checks in flight while walking.
TYPE:
|
resume_build_id
|
Resume this build: its plan for this scope is reused, observations re-sent, and failed members reset.
TYPE:
|
register_all
|
Expand complete tasks too, so every edge is recorded.
TYPE:
|
on_registry_failure
|
TYPE:
|
concurrency_config
|
Build-local limits (how many tasks run at once in this process).
TYPE:
|
concurrency_limiter
|
A pre-built limiter to use instead of
TYPE:
|
claim_config
|
How claims are waited on and renewed.
TYPE:
|
settings
|
Environment variables applied for the build's duration
(the scope's second half;
TYPE:
|
limit_key_selector
|
The registry concurrency-limit keys a task runs under, sent with its claim.
TYPE:
|
description
|
A description for a new build.
TYPE:
|
raise_on_failure
|
In
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
BuildSummary
|
BuildSummary with status, task counts and build id. |
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
106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 | |
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
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). |
registry
|
Default: the configured registry (none: no plan, no claims).
TYPE:
|
fail_mode
|
TYPE:
|
sync_run_default
|
How a sync-only task runs:
TYPE:
|
resume_build_id
|
Resume this build (its plan for this scope is reused, observations re-sent, failed members reset).
TYPE:
|
register_all
|
Expand complete tasks too.
TYPE:
|
on_registry_failure
|
TYPE:
|
claim_config
|
Claim waiting and renewal.
TYPE:
|
settings
|
Environment variables applied for the build's duration.
Omitted on a resume, the build's active plan's settings are
reused;
TYPE:
|
limit_key_selector
|
Registry concurrency-limit keys per task, sent with its claim.
TYPE:
|
description
|
A description for a new build.
TYPE:
|
max_concurrent_discover
|
Completion checks in flight while walking.
TYPE:
|
raise_on_failure
|
In
TYPE:
|
Source code in stardag/build/_sequential.py
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
¶
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.
DirectoryTarget
¶
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
exists_aio
async
¶
mark_done_aio
async
¶
Async version of mark_done().
Source code in stardag/target/_base.py
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
¶
TargetFactory
¶
Source code in stardag/target/_factory.py
get_file_target
¶
Get a file target.
| PARAMETER | DESCRIPTION |
|---|---|
relpath
|
The path to the target, relative to the configured root path for
TYPE:
|
target_root_key
|
The key to the target root to use.
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
FileTarget
|
A file target. |
Source code in stardag/target/_factory.py
get_directory_target
¶
Get a directory target.
| PARAMETER | DESCRIPTION |
|---|---|
relpath
|
The path to the target, relative to the configured root path for
TYPE:
|
target_root_key
|
The key to the target root to use.
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
DirectoryTarget
|
A directory target. |
Source code in stardag/target/_factory.py
get_path
¶
Get the full (/"absolute") path (/"URI") to the target.
Source code in stardag/target/_factory.py
Registry Module¶
registry
¶
The task registry (v2).
- :class:
RegistryABC: the interface every engine and integration uses. - :class:
APIRegistry: its implementation over the/api/v2HTTP 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
¶
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
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
build_resume
¶
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
build_complete
¶
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
build_exit_early
¶
POST /builds/{id}/exit-early: the resident driver stops;
nothing is released (its in-flight executions keep reporting).
build_list_running
¶
RUNNING builds, most recently active first (the watchdog sweep).
plan_create
¶
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
plan_register_members
¶
POST /plans/{id}/members: one chunk (at most 1000 items), in
one transaction.
plan_seal
¶
POST /plans/{id}/seal: verify the static phase and seal (a
replacement activates here).
plan_roots_info
¶
plan_roots
¶
build_skip_blocked
¶
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
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
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
member_retry
¶
Reset to PENDING (the fail mode's retry). Idempotent by state; refused 409 on COMPLETED and on a live claim.
member_interrupt
¶
The platform ended the execution and nothing will restart it: INTERRUPTED (actionable), claim released.
Source code in stardag/registry/_base.py
member_preempt
¶
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
member_cancel
¶
One task's cancel, by the build holding its claim (409
not_claim_holder otherwise).
member_skip
¶
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
member_exclude
¶
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
member_discovery_failed
¶
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
claim_renew
¶
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
build_list_executions
¶
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
execution_report_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
task_get
¶
GET /tasks/{task_id}: a completion's identity and state, with
its instances (each a construction under one scope), newest first.
task_list_artifacts
¶
task_upload_artifacts
¶
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
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
deployment_activate
¶
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
deployment_list
¶
Newest first. current=True keeps one row per app: its
activated deployment with the highest generation.
Source code in stardag/registry/_base.py
concurrency_limit_set
¶
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
concurrency_limit_delete
¶
concurrency_limit_list
¶
concurrency_limit_list_detailed
¶
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
build_set_reactive_meta
¶
Mark the build reactively scheduled by app_name; None
tick_kwargs keeps the stored configuration.
Source code in stardag/registry/_base.py
scheduler_lease_release
¶
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
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 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:
|
| RETURNS | DESCRIPTION |
|---|---|
StardagConfig
|
Fully resolved StardagConfig (actual type is StardagConfig). |
Source code in stardag/config/loader.py
104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 | |