Arbiter is a lightweight, distributed task scheduling engine designed specifically for ML workloads. It accepts computational DAGs (Directed Acyclic Graphs), parses dependencies, and executes tasks in parallel across a cluster of GPU-aware worker nodes.
Built entirely from scratch using asyncio, gRPC, and Protocol Buffers. There are zero heavy framework dependencies—no Celery, no Ray, no Langchain. Just pure, rigorously typed Python.
The system follows a centralized master-worker architecture. The Scheduler maintains the global state of the DAGs, while Workers execute tasks and report results over persistent, bidirectional gRPC streams.
graph TB
Client[Client SDK<br/>Python] -->|gRPC SubmitDAG| Scheduler
subgraph Scheduler["Scheduler (Master Node)"]
DAGParser[DAG Parser &<br/>Topological Sorter]
TaskQueue[Async Task Queue<br/>Dependency Tracker]
ResourceManager[Resource Manager<br/>GPU/VRAM Tracker]
StateMachine[Task State Machine<br/>PENDING → RUNNING → COMPLETED]
end
Scheduler <-->|gRPC Bi-Di Stream<br/>TaskStream| Worker1
Scheduler <-->|gRPC Bi-Di Stream<br/>TaskStream| Worker2
Scheduler <-->|gRPC Bi-Di Stream<br/>TaskStream| WorkerN
subgraph Worker1["Worker Node 1"]
W1Exec[Task Executor]
W1GPU[GPU Resource Probe]
end
subgraph Worker2["Worker Node 2"]
W2Exec[Task Executor]
W2GPU[GPU Resource Probe]
end
- DAG Parsing & Validation: Uses Kahn's algorithm (O(V+E)) for topological sorting and strict cycle detection. Tasks with missing dependencies are rejected at the boundary.
- Dependency-Aware Scheduling: Tasks remain PENDING until all upstream dependencies are marked COMPLETED. The scheduler dynamically unlocks dependents.
- GPU-Aware Resource Management: Workers register available VRAM. The scheduler tracks allocation and ensures tasks are only assigned to workers with sufficient compute resources.
- gRPC Bidirectional Streaming: Persistent, low-latency connections for task assignment and result reporting. Eliminates HTTP polling overhead.
- Strict State Transitions: A built-in state machine enforces valid task lifecycles (PENDING → SCHEDULED → RUNNING → COMPLETED/FAILED), preventing race conditions.
- Uncompromising Type Safety: 100% mypy --strict compliance and ruff linting. No untyped functions, no dynamic dictionaries in hot paths.
arbiter/
├── proto/ # gRPC protocol definitions
│ └── arbiter.proto
├── src/
│ └── arbiter/
│ ├── core/ # Shared internals
│ │ ├── exceptions.py # Custom exception hierarchy
│ │ └── models.py # Typed dataclasses & Enums
│ ├── generated/ # Auto-generated gRPC stubs (ignored by linters)
│ ├── scheduler/ # Master node logic
│ │ ├── dag.py # Kahn's algorithm & DAG validation
│ │ ├── resource_manager.py # Concurrent GPU/VRAM tracking
│ │ ├── server.py # gRPC servicer & stream handler
│ │ ├── state_machine.py # Task lifecycle enforcement
│ │ ├── task_queue.py # Async-safe dependency queue
│ │ └── __main__.py # Scheduler entry point
│ ├── worker/ # Worker node logic
│ │ ├── agent.py # gRPC client & stream manager
│ │ ├── executor.py # Task execution (Mock for now)
│ │ └── __main__.py # Worker entry point
│ └── __init__.py
├── tests/
│ ├── e2e/ # End-to-end distributed tests
│ ├── integration/ # gRPC server integration tests
│ └── unit/ # DAG parser and queue unit tests
├── scripts/
│ └── submit_dag.py # Client script to submit test workloads
├── docker-compose.yml # 1 Scheduler + 3 Workers cluster
├── Dockerfile # Multi-stage Python 3.11-slim build
└── pyproject.toml # Hatchling build, ruff, mypy config
Why gRPC over REST/HTTP? Internal ML infrastructure requires low-latency, typed contracts. gRPC's bidirectional streaming is essential for real-time task assignment without polling overhead. Protobufs enforce schema evolution.
Why custom asyncio Queue over Celery/Redis? Introducing a message broker adds network hops and serialization overhead. For tightly coupled, in-memory scheduling, an asyncio.Lock protected queue is orders of magnitude faster and eliminates external state management.
Why dataclasses over Pydantic? Pydantic is excellent for API boundaries, but its runtime validation overhead is unnecessary for internal state. Standard dataclasses with frozen=True provide immutability and are significantly faster in hot paths.
You must have Docker and Docker Compose installed.
1. Build and launch the cluster: Spins up 1 Scheduler and 3 Worker containers.
docker compose up --build -d2. Submit a DAG workload: Submit a test DAG (A → B, A → C, B&C → D) to the exposed scheduler port.
python scripts/submit_dag.py3. Observe distributed execution: Filter the logs to watch the scheduler assign tasks and workers execute them in parallel.
docker compose logs -f | findstr /C:"assigning_task" /C:"task_completed"Expected output: Task A completes, then Tasks B and C are assigned to different workers simultaneously.
4. Tear down:
docker compose downRequires Python 3.11+.
Install dependencies:
python -m venv venv
source venv/bin/activate # Windows: venv\Scripts\activate
pip install -e ".[dev]"Regenerate gRPC stubs (if proto changes):
python -m grpc_tools.protoc -I proto --python_out=src/arbiter/generated --grpc_python_out=src/arbiter/generated proto/arbiter.protoRun strict quality gates:
ruff check .
mypy src/
pytest tests/