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

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.

Scheduler ▶ Decision Engine

Runs a pure-Python scheduling pipeline in run_scheduling_decision within 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 (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:

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.

Database Insertion and Configuration

The submit function calls insert_job_and_config in 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 (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 (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

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.

Programmatic Submission

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

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, 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, 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. 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 (lines 88-99). The scheduler uses these bands to determine task ordering and enforce band-specific capacity limits during the assignment phase.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →