# How Marin's StepRunner Handles DAG Execution: A Deep Dive into the Core Scheduler

> Discover how Marin's StepRunner handles DAG execution. Learn about its eager pulling, dependency expansion, and parallel execution with thread pools and caching.

- Repository: [The Marin Project/marin](https://github.com/marin-community/marin)
- Tags: deep-dive
- Published: 2026-09-10

---

**Marin's StepRunner executes directed acyclic graphs (DAGs) by eagerly pulling steps from an iterable, expanding transitive dependencies via post-order traversal, and launching steps as soon as all parent dependencies complete, while supporting parallel execution through a thread pool and intelligent caching.**

Marin is an open-source pipeline orchestration framework maintained by `marin-community/marin`. At its heart lies the **StepRunner**, a deterministic DAG scheduler that orchestrates `StepSpec` objects into a coherent execution flow. Understanding how this component handles dependency resolution, concurrency, and caching is essential for building efficient, reproducible data pipelines with Marin.

## What is StepSpec? The Building Block of Marin's DAG

Before execution begins, Marin represents each unit of work as a **`StepSpec`**—an immutable data object defined in [`lib/marin/src/marin/execution/step_spec.py`](https://github.com/marin-community/marin/blob/main/lib/marin/src/marin/execution/step_spec.py) (lines 25–102). Unlike imperative task definitions, a `StepSpec` describes **what** to run, not **how** to run it.

Key fields include:

- **`name`** – A human-readable identifier for the step.
- **`hash_attrs`** – JSON-serializable configuration used for content hashing and cache invalidation.
- **`deps`** – A list of other `StepSpec` objects representing upstream dependencies.
- **`fn`** – The callable (or `RemoteCallable`) that performs the actual computation.
- **`resources`** – An optional `ResourceConfig` that triggers remote execution via Fray.
- **`writes_record`** – A boolean indicating if the step writes its own artifact record for lazy artifacts.

The **output path** of a step—used for status files, cache checks, and as a unique execution key—is derived from a prefix plus the step's `name_with_hash`. This path serves as the canonical identity for the step throughout the DAG execution process.

## The Execution Loop: How StepRunner Schedules Steps

The `StepRunner.run()` method in [`lib/marin/src/marin/execution/step_runner.py`](https://github.com/marin-community/marin/blob/main/lib/marin/src/marin/execution/step_runner.py) (lines 27–44) implements the core scheduling logic. It maintains several bookkeeping sets to track execution state:

- **`completed`** – Paths of steps that finished successfully.
- **`failed`** – Paths of steps that errored or had failed dependencies.
- **`running`** – Mapping of `output_path → JobHandle` for active executions.
- **`waiting`** – List of steps whose dependencies are not yet satisfied.
- **`pruned`** – Paths of cache-only leaves skipped due to existing outputs.
- **`scheduled`** – Global set of all paths already seen for deduplication.

The runner follows a four-phase cycle for each incoming step: expand the dependency graph, harvest completed jobs, flush waiting steps with satisfied dependencies, and determine the fate of the current step based on its state and cache status.

### Expanding the DAG with _expand_unseen

When `StepRunner.run()` receives a step, it invokes `_expand_unseen(step, seen, is_built, pruned)` (lines 39–56) to traverse the dependency graph **post-order**. This function walks the step's transitive dependencies and returns a list of steps that have not yet been scheduled.

If `_expand_unseen` encounters a cycle, it raises `ValueError`, though cycles should never appear because the graph is explicitly built as a DAG. When `is_built(node)` returns `True`—indicating the step is already cached—the node's output path is added to `pruned`, and its children are **not** traversed. This optimization prunes entire sub-trees of already-successful work from the schedule.

### The Core Scheduling Algorithm

After expansion, the runner processes each step through a decision tree:

1. **Cache confirmation** – If the path exists in `pruned`, call `_confirm_pruned()` to verify the cached artifact.
2. **Failure propagation** – If any dependency path exists in `failed`, mark the step as failed.
3. **Immediate launch** – If all dependencies are in `completed`, submit the step via `_do_launch()`.
4. **Defer execution** – Otherwise, append the step to `waiting` for later processing.

The loop continues until both `running` and `waiting` collections are empty. If any steps failed during execution, the runner raises a combined exception containing all failure details.

## Parallel Execution and Concurrency Control

Marin utilizes a **`ThreadPoolExecutor`** to enable parallel step execution. The executor is initialized with `max_concurrent` workers (defaulting to 8) and persists for the duration of the `run()` call.

Each ready step is submitted to the pool via `_launch_step`, which wraps the step's `run_step` call in a worker function. This wrapper also propagates the current Fray client context, ensuring that remote jobs inherit the same authentication and configuration as the parent process. The executor guarantees that no more than `max_concurrent` steps execute simultaneously, regardless of how many steps are ready in the `waiting` queue.

## Caching, Drift Detection, and Pruning

Before submitting a step to the thread pool, `_launch_step` (lines 65–78) performs several validation checks to avoid redundant work:

- **Dry-run guard** – When `dry_run=True`, the runner skips all I/O operations and simulates execution.
- **Drift check** – The `check_drift` function determines if mutable (development) artifacts require rebuilding regardless of cache status.
- **Cache validation** – The runner reads the step's status file; if it contains `STATUS_SUCCESS` and the step is immutable, the execution is skipped entirely.

These mechanisms ensure that Marin's DAG execution is **idempotent**—re-running a pipeline will automatically skip steps that completed successfully in previous runs, while still detecting configuration changes that necessitate rebuilding.

## Remote Execution via Fray Integration

When a `StepSpec` defines a `resources` field, the StepRunner automatically dispatches the step as a **Fray job** rather than executing it locally. This remote execution flow (lines 155–215) branches based on the callable type:

- **Standard functions** with resources trigger `_run_iris_job()`.
- **`RemoteCallable` instances** invoke `_run_remote_step()`.

Both paths ultimately call `_submit_iris_job()`, which constructs a `JobRequest`, submits it to the current Fray client, and blocks until completion. This allows Marin pipelines to seamlessly scale from local development to distributed computing without changing the DAG structure.

## Practical Examples

### Building a Simple Dependency DAG

```python
from marin.execution.step_spec import StepSpec
from marin.execution.step_runner import StepRunner

# Leaf step – tokenizes some data

tokenize = StepSpec(
    name="tokenize",
    fn=lambda out: open(out, "w").write("tokens"),
    hash_attrs={"vocab": "v1"},
)

# Dependent step – builds a vocab from the tokenized output

build_vocab = StepSpec(
    name="build_vocab",
    deps=[tokenize],
    fn=lambda out: open(out, "w").write("vocab"),
    hash_attrs={"size": 1000},
)

# Top step – trains a model that depends on both previous steps

train = StepSpec(
    name="train",
    deps=[tokenize, build_vocab],
    fn=lambda out: open(out, "w").write("model"),
    hash_attrs={"layers": 12},
)

# Run the DAG

StepRunner().run([train])

```

The runner expands the transitive dependencies (`tokenize`, `build_vocab`) and executes them in topological order before launching the `train` step.

### Parallel Execution with Concurrency Limits

```python

# Create many independent steps

steps = [
    StepSpec(name=f"step_{i}", fn=lambda out, i=i: open(out, "w").write(str(i)))
    for i in range(20)
]

# Run up to 5 steps concurrently

StepRunner().run(steps, max_concurrent=5)

```

This configuration processes all 20 steps using only 5 concurrent workers, respecting resource constraints while maximizing throughput.

### Remote Execution with Resource Requirements

```python
from fray.types import ResourceConfig

remote = StepSpec(
    name="remote_job",
    fn=lambda out: open(out, "w").write("remote result"),
    resources=ResourceConfig(cpu=4, memory_gb=16),  # forces Fray submission

)

StepRunner().run([remote])

```

When `resources` is specified, the StepRunner automatically submits the step to the Fray cluster rather than executing it in the local process.

## Summary

- **StepSpec** objects in [`lib/marin/src/marin/execution/step_spec.py`](https://github.com/marin-community/marin/blob/main/lib/marin/src/marin/execution/step_spec.py) serve as immutable DAG nodes containing dependency references, hashable configuration, and execution logic.
- **Post-order expansion** via `_expand_unseen` prunes cached sub-trees and detects cycles before scheduling begins.
- **State-machine scheduling** tracks completed, failed, running, and waiting steps to ensure dependencies complete before downstream execution.
- **ThreadPoolExecutor** enables parallel execution with configurable concurrency limits (`max_concurrent`).
- **Intelligent caching** checks status files and drift conditions to skip redundant work, ensuring idempotent pipeline runs.
- **Fray integration** automatically promotes steps with `ResourceConfig` to remote jobs without altering the DAG structure.

## Frequently Asked Questions

### How does StepRunner handle circular dependencies in the DAG?

The `_expand_unseen` method raises a `ValueError` if it encounters a cycle during the post-order traversal of the dependency graph. However, cycles should never appear in practice because Marin's `StepSpec` graphs are explicitly constructed as directed acyclic graphs by design.

### What happens when a step fails during DAG execution?

When a step fails, its output path is added to the `failed` set. The StepRunner then marks all downstream steps as failed during the scheduling phase, preventing wasted computation on dependent tasks. Once the execution loop completes, the runner raises a combined exception containing details of all failures encountered during the run.

### Can StepRunner resume a partially completed pipeline?

Yes. The runner automatically prunes steps where `is_built(node)` returns `True` by adding them to the `pruned` set and skipping their children. When re-running a pipeline, successfully completed steps with valid status files and matching hashes are skipped entirely, allowing the pipeline to resume from the point of failure or interruption.

### What is the difference between local and remote step execution?

Local execution runs the step's `fn` callable directly within the ThreadPoolExecutor. Remote execution occurs when a `StepSpec` includes a `resources` field or contains a `RemoteCallable`; in these cases, the StepRunner submits the work to a Fray cluster via `_submit_iris_job()`, blocking until the remote job completes while still respecting the local concurrency limits for job submission.