# How Iris Federation Routing Works Across Cross-Cluster Storage

> Learn how Iris federation routing classifies jobs as LOCAL QUEUE or REJECT using PeerRouter and cluster pins. Optimize cross-cluster storage efficiency.

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

---

**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`](https://github.com/marin-community/marin/blob/main/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:

```python
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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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:

```python
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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/0035_federation_unify.py)** creates the `federation_changelog` table, while **[`0037_federation_fixup.py`](https://github.com/marin-community/marin/blob/main/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.