How Marin-Fray Provides a Distributed Execution Substrate

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. 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 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

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

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

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 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 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 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.

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 →