# How Marin-Fray Provides a Distributed Execution Substrate

> Marin-Fray provides a distributed execution substrate by unifying Iris job orchestration. Learn how it translates specifications, supports actor hosting, coscheduling, and dynamic discovery.

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

---

**Marin-Fray delivers a distributed execution substrate by wrapping the Iris job orchestration system behind a unified Client interface, translating high-level resource and job specifications into low-level Iris protobuf messages while supporting actor hosting, coscheduling, and dynamic discovery.**

The `marin-community/marin` repository implements Fray as a lightweight, extensible layer that turns Iris-managed clusters into a portable distributed execution substrate. By abstracting backend specifics behind a `Client` protocol, Fray enables users to submit jobs, host actors, and manage replicated actor groups without interacting directly with Iris RPC details.

## Backend Abstraction Layer

At the core of Fray's distributed execution substrate is the **`FrayIrisClient`** class defined in [`lib/fray/src/fray/iris_backend.py`](https://github.com/marin-community/marin/blob/main/lib/fray/src/fray/iris_backend.py). This class implements the `Client` protocol, acting as a bridge between Fray's high-level APIs and Iris's job orchestration primitives.

The abstraction allows alternative backends (such as local thread pools or custom cloud schedulers) to be swapped by implementing the same protocol, making the substrate portable across different infrastructure targets.

### Resource Conversion Logic

Before submitting work to Iris, Fray must translate its native dataclasses into Iris protobuf structures. The conversion functions handle three critical mappings:

- **`convert_resources`** – Transforms Fray's `ResourceConfig` into Iris `ResourceSpec` objects, mapping device types like GPU or TPU to their Iris equivalents.
- **`convert_constraints`** – Converts Fray `Constraint` objects into Iris placement constraints.
- **`convert_entrypoint`** – Packages Fray `Entrypoint` definitions into Iris `Entrypoint` structs.

These conversion utilities ensure that high-level user specifications in [`fray/types.py`](https://github.com/marin-community/marin/blob/main/fray/types.py) align with the low-level requirements of the Iris backend.

### Coscheduling Configuration

For multi-replica workloads requiring gang scheduling, **`resolve_coscheduling`** selects the appropriate Iris `CoschedulingConfig` based on device type and replica count. When running on TPUs or multi-GPU nodes, this function configures the backend to reserve resources across machines simultaneously, preventing partial allocation deadlocks.

## Job Submission Pipeline

The **`submit`** method in `FrayIrisClient` orchestrates the transition from Fray `JobRequest` to Iris `JobRequest`. This method performs several critical operations:

1. Builds the Iris protobuf request using the conversion utilities
2. Applies optional wrappers for multiprocess execution or Nsight profiling
3. Invokes `iris.submit` to enqueue the job
4. Wraps the returned Iris `Job` in an **`IrisJobHandle`**

Error handling includes specific detection of duplicate job names, raising `FrayJobAlreadyExists` when Iris reports naming conflicts.

### Job Lifecycle Management

**`IrisJobHandle`** provides a thin abstraction over Iris job state, exposing methods like `status()`, `wait()`, `logs()`, and `terminate()`. This class maps Iris's detailed state machine to Fray's simplified `JobStatus` enum, offering a consistent interface for monitoring distributed work regardless of backend specifics.

## Actor Hosting and Group Management

Fray's substrate supports fine-grained distributed computing through an actor model layered atop Iris jobs.

### Hosting Actors Within Jobs

The **`_host_actor`** function enables an actor class to run inside an Iris job replica. When invoked, it:

- Creates an `ActorServer` instance
- Registers a unique endpoint following the pattern `<job_id>/<name>-<task_index>`
- Blocks on a shutdown event to keep the job alive

This mechanism allows driver-side services and parameter servers to coexist within the same Iris job that runs compute workers, reducing scheduling overhead.

### Dynamic Actor Discovery

**`IrisActorGroup`** implements scalable actor farms by polling the Iris resolver for endpoints matching a specific prefix. This class tracks discovered actors across job replicas and provides helper methods:

- **`discover_new`** – Identifies recently registered actors
- **`wait_ready`** – Blocks until a specified count of actors are available
- **`shutdown`** – Coordinates graceful termination across the group

By leveraging Iris's endpoint registry rather than spawning separate jobs per replica, `IrisActorGroup` enables elastic, fault-tolerant collections of workers without overwhelming the scheduler.

## Code Examples

Below are practical patterns for utilizing Fray's distributed execution substrate.

### Submitting a Distributed Job

```python
from fray import FrayIrisClient, JobRequest, ResourceConfig, GpuConfig, FrayEntrypoint

# Initialize the client against an Iris controller

client = FrayIrisClient(controller_address="iris.example.com:1234")

def train_step():
    print("Executing distributed training step")

# Configure resources for 2 GPU replicas

request = JobRequest(
    name="distributed_training",
    entrypoint=FrayEntrypoint.from_callable(train_step),
    resources=ResourceConfig(device=GpuConfig(variant="v100", count=2)),
    replicas=2,
)

# Submit and monitor

handle = client.submit(request)
print(f"Job ID: {handle.job_id}")
print(f"Status: {handle.status()}")
handle.wait()

```

### Creating an Actor Group

```python
from fray import ActorConfig

class ParameterServer:
    def __init__(self, initial_weights):
        self.weights = initial_weights
    
    def update(self, gradients):
        for key, grad in gradients.items():
            self.weights[key] -= 0.01 * grad
        return self.weights

# Launch 4 replicas as a single Iris job with gang scheduling

ps_group = client.create_actor_group(
    actor_class=ParameterServer,
    name="parameter_server",
    count=4,
    resources=ResourceConfig(device=GpuConfig(variant="v100", count=1)),
    actor_config=ActorConfig(priority=1),
    initial_weights={"layer1": 0.0, "layer2": 0.0},
)

# Wait for at least one replica to be ready

ps_handle = ps_group.wait_ready(count=1)[0]

# Execute remote method call

future = ps_handle.update.remote({"layer1": 0.5, "layer2": 0.3})
updated_weights = future.result()

```

### Hosting a Service Actor

```python
from fray import host_actor

class Logger:
    def log(self, message: str) -> None:
        print(f"[LOG] {message}")

# Host within the current job context for other tasks to discover

logger = host_actor(Logger, name="experiment_logger")

# Other job replicas can resolve this endpoint via the Iris resolver

```

## Summary

- **Fray provides a distributed execution substrate by wrapping Iris** behind a unified `Client` interface, isolating users from protobuf and RPC complexities.
- **Resource conversion functions** in [`iris_backend.py`](https://github.com/marin-community/marin/blob/main/iris_backend.py) map Fray's `ResourceConfig` and `Entrypoint` objects to Iris `ResourceSpec` and job definitions.
- **Coscheduling support** enables gang-scheduled multi-replica jobs through the `resolve_coscheduling` utility, critical for TPU and multi-GPU workloads.
- **Job submission** occurs through `FrayIrisClient.submit`, which returns `IrisJobHandle` instances for lifecycle management and status monitoring.
- **Actor hosting** via `_host_actor` runs actor classes inside job replicas, registering unique endpoints in Iris's resolver.
- **Dynamic discovery** through `IrisActorGroup` polls the Iris resolver to build elastic, fault-tolerant actor collections without per-actor job overhead.

## Frequently Asked Questions

### How does Marin-Fray translate high-level job specifications into Iris-specific calls?

Fray utilizes conversion functions located in [`lib/fray/src/fray/iris_backend.py`](https://github.com/marin-community/marin/blob/main/lib/fray/src/fray/iris_backend.py) to transform Fray dataclasses into Iris protobuf messages. The `convert_resources` function maps `ResourceConfig` objects to Iris `ResourceSpec`, while `convert_entrypoint` packages Python callables into Iris job entrypoints. This translation layer ensures that user code remains backend-agnostic while correctly interfacing with Iris's RPC requirements.

### What mechanism enables gang scheduling for multi-replica workloads?

The `resolve_coscheduling` function in [`iris_backend.py`](https://github.com/marin-community/marin/blob/main/iris_backend.py) analyzes the device type and replica count to configure an Iris `CoschedulingConfig`. For TPU workloads or GPU gangs, this ensures all replicas are allocated simultaneously across machines, preventing partial resource acquisition that could deadlock distributed training jobs. The configuration is embedded in the `JobRequest` before submission to the Iris controller.

### How does Fray handle actor discovery across distributed job replicas?

Fray implements dynamic actor discovery through the `IrisActorGroup` class, which polls the Iris resolver for endpoints matching a job-specific prefix. Rather than spawning separate Iris jobs for each actor replica, `IrisActorGroup` tracks endpoints formatted as `<job_id>/<name>-<task_index>` within a single job. This approach reduces scheduler load while enabling elastic scaling and fault tolerance through methods like `discover_new` and `wait_ready`.

### Can Fray support execution backends other than Iris?

Yes. Fray defines a `Client` protocol that `FrayIrisClient` implements for the Iris backend. Alternative substrates such as local thread pools, Kubernetes operators, or custom cloud schedulers can be integrated by implementing the same protocol methods (`submit`, `create_actor`, `create_actor_group`, etc.). This abstraction makes the distributed execution substrate portable across different infrastructure targets without modifying user code.