# How Marin-Iris Manages Distributed Job Execution: Architecture and Implementation

> Learn how Marin Iris manages distributed job execution across local and cloud GPUs and Kubernetes. Discover its controller-database-scheduler architecture for consistent state.

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

---

**Marin-Iris orchestrates distributed training and inference workloads across local GPU, Cloud GPU/TPU, and federated Kubernetes clusters through a controller-database-scheduler architecture that guarantees state consistency by preventing backends from writing directly to the job database.**

Marin-Iris serves as the job-orchestration layer within the `marin-community/marin` repository, driving distributed execution through a deterministic scheduling pipeline. Understanding how Marin-Iris manages distributed job execution reveals a design that separates state management from compute provisioning, ensuring reliable coordination across heterogeneous clusters.

## Architecture Overview

Marin-Iris consists of three tightly-coupled components that coordinate distributed workloads:

**Controller ▶ Database**

Persists job metadata, configuration, and task state in a SQL database. All scheduling decisions derive from a consistent snapshot of this database, with schema definitions for `jobs_table` and `job_workdir_files_table` located in [`iris/cluster/schema.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/schema.py).

**Scheduler ▶ Decision Engine**

Runs a pure-Python scheduling pipeline in `run_scheduling_decision` within [`iris/cluster/controller/backend.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/backend.py) (lines 47-65). This engine expands database snapshots into placement decisions through availability enrichment, gating, ordering, assignment, and pre-emption phases.

**Backend ▶ Workers**

Implements the `TaskBackend` protocol defined in [`iris/cluster/controller/backend.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/backend.py) (lines 67-74). Backends such as native Docker, Kubernetes, Kueue, or Slurm expose capacity and execute tasks via a thin RPC layer, but never write directly to the database.

## Job Submission Flow

### CLI and Client Requests

Users initiate jobs via the command line:

```bash
uv run iris --cluster=marin job run \
    --cpu=4 --memory=8G \
    --extra=cpu \
    --region=us-east5 \
    python train.py

```

The CLI marshals a `controller_pb2.Controller.LaunchJobRequest` and transmits it to the controller, as documented in [`docs/tutorials/train-an-lm.md`](https://github.com/marin-community/marin/blob/main/docs/tutorials/train-an-lm.md).

### Database Insertion and Configuration

The `submit` function calls `insert_job_and_config` in [`iris/cluster/controller/ops/job.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/ops/job.py) (lines 57-73 and 84-92). This transaction writes the `jobs` row and associated configuration, resolves the priority band, and records the authenticated submitter for root jobs.

### Task Materialization

Following successful insertion, the controller invokes `_materialize_tasks` in [`iris/cluster/controller/ops/job.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/ops/job.py) (lines 64-85). This creates a contiguous block of **PENDING** task rows (one per replica) by reserving a priority-insertion base and bulk-inserting tasks.

### Audit Logging

The insertion registers an audit event (`job_submitted`) emitted once the transaction commits (lines 48-53).

## The Deterministic Scheduling Pipeline

Every controller tick builds an in-memory snapshot of the current database state and passes it to the scheduler. The pipeline executes entirely without I/O:

1. **Capacity Enrichment**: `run_scheduling_decision` enriches worker attributes with `availability:<variant>` markers if the backend supplies per-zone accelerator capability information (lines 68-78).

2. **Gating**: `apply_scheduling_gates` filters tasks blocked by resource caps, quotas, or expiration (lines 84-88).

3. **Ordering**: `compute_scheduling_order` produces a deterministic priority sequence for remaining tasks (lines 108-109).

4. **Assignments**: `scheduler.find_assignments` executes the core placement algorithm, matching tasks to workers based on constraints, priority band, and capacity (lines 60-62). The `apply_placements` function collects results (lines 48-64).

5. **Pre-emption**: `apply_preemptions` generates pre-empt plans for higher-priority jobs requiring slots, committed atomically with assignments (lines 110-120).

6. **Result**: Returns a `ScheduleResult` containing assignments, preemptions, unschedulable tasks, diagnostics, and enriched context (lines 19-36).

## Backend Interaction Protocol

The controller communicates with compute providers through the `TaskBackend` protocol:

- **`observe`**: Publishes provider status including worker liveness, capacity, and pending hints as a `BackendObservation` (lines 91-93).

- **`schedule`**: Invokes the pure decision pipeline for backends performing their own placement (e.g., Kueue). For backends allowing Iris to schedule tasks directly, this returns empty and the controller drives placement.

- **`runtime_image`**: Resolves container images with backend-specific fallback defaults (lines 82-88).

This architecture guarantees state consistency because all mutations funnel through controller transactions; backends remain read-only regarding database state.

## Federated Execution and Hierarchical Jobs

Iris supports federated clusters where parent jobs launch child jobs on remote clusters. The `cluster` argument in `insert_job_and_config` routes to the child cluster's backend while the parent cluster materializes tasks locally. The controller mirrors child state back into the parent database, enabling unified visibility across clusters (lines 68-71).

## Priority Bands and Resource Management

Jobs specify a `PRIORITY_BAND` (interactive, batch, or `INHERIT`). The `resolve_priority_band` function normalizes this at ingestion in [`iris/cluster/controller/ops/job.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/ops/job.py) (lines 88-99), ensuring stored rows contain concrete bands. The scheduler uses these bands to order tasks and enforce band-specific capacity limits.

## Practical Code Examples

### Submitting via CLI

```bash
uv run iris --cluster=marin job run \
    --cpu=4 --memory=8G \
    --extra=cpu \
    --region=us-east5 \
    python train.py

```

This creates a `LaunchJobRequest` handled by the `submit` path in [`iris/cluster/controller/ops/job.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/ops/job.py).

### Programmatic Submission

```python
from iris.client import IrisClient
from iris.rpc import controller_pb2

client = IrisClient(cluster="marin")
req = controller_pb2.Controller.LaunchJobRequest(
    job_name="example",
    resources=controller_pb2.ResourceSpecProto(
        cpu_millicores=4000, 
        memory_bytes=8*1024**3
    ),
    command=["python", "train.py"],
)

job_id = client.submit_job(req)  # Calls submit → insert_job_and_config

print(f"Submitted job {job_id}")

```

### Inspecting Job Status

```python
info = client.get_job_info(job_id)
print(f"State: {info.state}")   # TASK_STATE_PENDING, RUNNING, etc.

print(f"Tasks: {info.num_tasks}")

```

## Summary

- Marin-Iris separates concerns into a database-backed controller, pure-Python scheduler, and pluggable backends to manage distributed job execution.
- Job submission flows through `insert_job_and_config` and `_materialize_tasks` in [`iris/cluster/controller/ops/job.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/ops/job.py), creating atomic database transactions with audit logging.
- The scheduling pipeline executes six deterministic phases—enrichment, gating, ordering, assignment, pre-emption, and result compilation—without database I/O via `run_scheduling_decision`.
- The `TaskBackend` protocol abstracts Docker, Kubernetes, Kueue, and Slurm while enforcing a strict boundary preventing backends from writing to the controller database.
- Federated execution enables hierarchical jobs across clusters with unified state mirroring in the parent database.
- Priority bands (interactive, batch, inherit) resolve at ingestion to enforce deterministic scheduling order and capacity limits.

## Frequently Asked Questions

### How does Marin-Iris ensure database consistency across distributed workers?

The controller maintains exclusive write access to the SQL database. Backends implement the `TaskBackend` protocol and communicate via RPC, but they never write directly to the database. All state changes flow through the controller's transaction logic in [`iris/cluster/controller/ops/job.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/ops/job.py), ensuring a single source of truth for job metadata and task states.

### What scheduling algorithm does Marin-Iris use for task placement?

Marin-Iris uses a multi-phase deterministic algorithm implemented in `run_scheduling_decision` within [`iris/cluster/controller/backend.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/backend.py). The pipeline enriches capacity data, applies scheduling gates for resource constraints, computes priority ordering, finds assignments matching tasks to workers, and applies preemptions when higher-priority jobs require resources.

### Can Marin-Iris schedule jobs across multiple Kubernetes clusters?

Yes. Marin-Iris supports federated execution where parent jobs launch child jobs on remote clusters using the `cluster` argument in `insert_job_and_config`. The controller mirrors child cluster state back into the parent database, providing a unified view of tasks across local and remote infrastructure.

### How are job priorities handled in the Marin-Iris scheduler?

Jobs specify a `PRIORITY_BAND` (interactive, batch, or `INHERIT`) which `resolve_priority_band` normalizes to a concrete band during ingestion in [`iris/cluster/controller/ops/job.py`](https://github.com/marin-community/marin/blob/main/iris/cluster/controller/ops/job.py) (lines 88-99). The scheduler uses these bands to determine task ordering and enforce band-specific capacity limits during the assignment phase.