How Iris Federation Routing Works Across Cross-Cluster Storage

Iris federation routing uses the PeerRouter class to classify jobs at submit time as LOCAL, QUEUE, or REJECT based on explicit cluster pins, local feasibility, and peer capability advertisements.

Iris federation routing enables the marin-community/marin project to distribute workloads across independent clusters while maintaining a consistent view of job state. The system decouples routing decisions from execution by storing all cross-cluster metadata in a shared relational database. This design allows jobs to be handed off between clusters with durable state tracking and automatic recovery from network partitions.

The PeerRouter Classification Logic

The core routing decision happens in lib/iris/src/iris/cluster/federation/router.py within the PeerRouter.classify() method. This method evaluates every job submission against a strict priority order to determine where computation should occur.

The classify() method inspects four specific conditions in sequence:

  1. Explicit cluster pins – If the job specifies cluster=<peer>, it is forced into the federation queue for that specific peer.
  2. Local feasibility – When the local backend can satisfy the job constraints, the router returns LOCAL disposition.
  3. Peer capability match – Each reachable peer periodically sends a heartbeat describing available backends (CPU, GPU, TPU). If any peer advertises a matching backend, the job receives QUEUE disposition.
  4. Default rejection – If no local or peer backend satisfies the constraints, the router returns REJECT.

The implementation follows this exact logic:

def classify(self, request: RoutingRequest) -> SubmitPlan:
    if request.cluster_pin:
        return SubmitPlan(SubmitDisposition.QUEUE, pinned_peer_id=request.cluster_pin)
    if request.local_feasible:
        return SubmitPlan(SubmitDisposition.LOCAL)
    if any(_peer_can_host(peer, request.constraints) for peer in self._peers.values()):
        return SubmitPlan(SubmitDisposition.QUEUE)
    return SubmitPlan(SubmitDisposition.REJECT)

Cross-Cluster Storage Architecture

Iris federation routing does not move data directly between clusters. Instead, both parent and peer clusters read and write to a shared relational store (such as Cloud SQL) that acts as the source of truth for all federated job metadata.

Shared Database Tables

The storage schema is defined through migration files in the controller module. Three primary tables enable cross-cluster coordination:

  • federation_changelog – Records every hand-off event, including timestamps and which peer currently owns a job. Created by migration 0035_federation_unify in lib/iris/src/iris/cluster/controller/migrations/0035_federation_unify.py.
  • federation_sync_state – Tracks the last successful synchronization timestamp for each job, enabling incremental sync between clusters.
  • peer_heartbeat – Implicitly stores the most recent backend capability advertisement from each reachable peer.

Migration 0037_federation_fixup in lib/iris/src/iris/cluster/controller/migrations/0037_federation_unify.py demonstrates schema evolution by dropping obsolete columns, showing how the storage contract matures over time.

State Synchronization and Recovery

Because job state lives in a central database rather than in-memory, both clusters maintain a consistent view of metadata even during network partitions. If a peer becomes unreachable, the parent cluster retains the last known state from the federation_changelog and federation_sync_state tables. When the peer rejoins, it reads the same tables to resume synchronization.

This behavior is validated in lib/iris/tests/cluster/test_federation.py, which exercises peer outage scenarios and verifies that hand-off recovery works correctly using the shared storage layer.

End-to-End Routing Flow

The complete lifecycle of a federated job involves distinct phases across the routing and storage layers:

  1. The submitter constructs a RoutingRequest containing job constraints and a flag indicating whether the local backend can host the job.
  2. PeerRouter.classify() evaluates the request, consulting cached heartbeats from each FederationPeer instance.
  3. For QUEUE dispositions, the job writes an entry to the federation queue in the shared database.
  4. The federation tick (executed by the control plane) scans queued jobs, selects a reachable peer with matching advertised capabilities, and updates federation_changelog with the assignment.
  5. The peer cluster reads its assignment from the shared tables, launches the job on its local backend, and periodically writes status updates back to the same storage.
  6. The parent cluster polls the database to aggregate peer progress and surfaces the unified view to users.

Practical Implementation Example

You can interact with the federation routing logic directly using the PeerRouter class:

from iris.cluster.federation.router import PeerRouter, RoutingRequest, SubmitDisposition
from iris.cluster.federation.peer import build_peers
from iris.cluster.constraints import Constraint, ConstraintOp

# Build peer objects from configuration

peers = build_peers([...])  # Returns list of FederationPeer instances

router = PeerRouter(peers)

# Create routing request for GPU-only job without explicit pin

constraints = [Constraint(key="device_type", op=ConstraintOp.EQ, value="gpu")]
request = RoutingRequest(
    constraints=constraints,
    local_feasible=False,
    cluster_pin=""  # Empty string indicates no explicit pin

)

plan = router.classify(request)

if plan.disposition == SubmitDisposition.QUEUE:
    print(f"Job queued for peer '{plan.pinned_peer_id or 'auto-selected'}'")
elif plan.disposition == SubmitDisposition.LOCAL:
    print("Job scheduled locally")
else:
    print("Job rejected: no suitable backend available")

This example demonstrates how to programmatically evaluate where a job will execute based on实时 peer capabilities and local resource availability.

Summary

  • Iris federation routing classifies every job at submit time using the PeerRouter.classify() method in lib/iris/src/iris/cluster/federation/router.py.
  • Routing decisions prioritize explicit cluster pins, then local feasibility, then peer capability advertisements, finally rejecting unsatisfiable requests.
  • Cross-cluster storage relies on shared relational tables (federation_changelog, federation_sync_state) defined in migrations 0035 and 0037 to persist job state independently of cluster availability.
  • The decoupled architecture separates routing classification from authorization and execution, enabling durable hand-offs and automatic recovery from peer outages as tested in lib/iris/tests/cluster/test_federation.py.

Frequently Asked Questions

How does PeerRouter decide between local and remote execution?

The PeerRouter.classify() method checks request.local_feasible immediately after handling explicit cluster pins. If the local backend can satisfy the job constraints and no pin forces remote execution, it returns SubmitDisposition.LOCAL. Otherwise, it evaluates peer capabilities to determine if remote queuing is possible.

What happens when a peer cluster becomes unreachable?

When a peer stops sending heartbeats, the PeerRouter no longer considers it in the any(_peer_can_host(...)) check. Jobs already queued for that peer remain in the shared database tables (federation_changelog and federation_sync_state), allowing the parent cluster to retain state until the peer rejoins and resumes synchronization.

Where is federation state persisted?

All cross-cluster state lives in a central relational database accessible to both parent and peer clusters. Migration file 0035_federation_unify.py creates the federation_changelog table, while 0037_federation_fixup.py manages schema evolution. This shared storage ensures consistency without direct cluster-to-cluster communication.

How does the classify method handle explicit cluster pins?

If request.cluster_pin contains a non-empty value (such as cluster=us-west-peer), the classify() method immediately returns SubmitDisposition.QUEUE with pinned_peer_id set to that value. This bypasses both local feasibility and peer capability checks, forcing the job into the federation queue for the specified cluster regardless of current availability.

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 →