1 Orchestrator and Control Plane
jfabian edited this page 2026-08-13 11:15:47 -03:00

Orchestrator & Control Plane

Agent Lifecycle

An agent's state machine is defined in vox/agents/lifecycle.py via the AgentState enum. Unlike the conceptual documentation which mentions only four states (BOOTING, ACTIVE, PAUSED, TERMINATED), the actual implementation contains nine states with explicitly validated transitions:

class AgentState(Enum):
    BOOTING = auto()   # Capability initialization
    IDLE = auto()      # Ready but not active
    ACTIVE = auto()    # Operational, processes events
    PAUSING = auto()   # Pause transition
    PAUSED = auto()    # Suspended, events queued
    RESUMING = auto()  # Resume transition
    STOPPING = auto()  # Stop transition
    STOPPED = auto()   # Terminated (replaces TERMINATED)
    FAILED = auto()    # Irreversible boot error

Allowed Transitions

BOOTING ──→ IDLE ──→ ACTIVE ──→ PAUSING ──→ PAUSED ──→ RESUMING ──→ ACTIVE
  │                    │                                                    │
  │                    ▼                                                    ▼
  └──→ FAILED      STOPPING ──────────────────────────────────────────→ STOPPED
                        ↑
  (any state can go to STOPPING)

can_transition_to() implements validation:

_TRANSITIONS: dict[AgentState, set[AgentState]] = {
    AgentState.BOOTING:   {AgentState.IDLE, AgentState.ACTIVE, AgentState.FAILED, AgentState.STOPPING},
    AgentState.IDLE:      {AgentState.BOOTING, AgentState.STOPPING},
    AgentState.ACTIVE:    {AgentState.PAUSING, AgentState.STOPPING},
    AgentState.PAUSING:   {AgentState.PAUSED, AgentState.STOPPING},
    AgentState.PAUSED:    {AgentState.RESUMING, AgentState.STOPPING},
    AgentState.RESUMING:  {AgentState.ACTIVE, AgentState.STOPPING},
    AgentState.STOPPING:  {AgentState.STOPPED},
    AgentState.STOPPED:   {AgentState.BOOTING},
    AgentState.FAILED:    {AgentState.BOOTING},
}

Any invalid transition attempt raises RuntimeError.

Detailed Boot Flow

  1. VOXOrchestrator calls _hire_agent() which instantiates VOXAgent(folder, orchestrator, logger)
  2. The VOXAgent constructor runs _bootstrap():
    • _load_manifest(): loads and validates agent.yml against MANIFEST_SCHEMA
    • _load_env(): optionally loads agent's .env
    • _validate_identity(): verifies name and id exist
    • _mount_capabilities(): iterates declared capabilities, requests instances from the Orchestrator via orchestrator.get_capability_instance(cap_id), and mounts them as VOXBoundCapability
    • _load_roles(): scans roles/ in the agent directory, dynamically imports Python modules, finds VOXRole subclasses, and registers handlers
    • _validate_role_dependencies(): verifies each role has its required capabilities mounted
  3. If autostart is true, the Orchestrator calls agent.boot():
    • Transitions to BOOTING
    • Invokes cap.boot() on each mounted capability
    • If any fails, transitions to FAILED
    • If all OK, transitions to ACTIVE and emits on_boot event

Control Plane over Unix Domain Socket

The control plane is implemented in vox/runtime/control_plane.py. It exposes a JSON-RPC server over a Unix socket at /tmp/vox.sock (default).

Implementation

async def start_control_plane(runtime, logger):
    if os.path.exists(runtime.config.uds_path):
        os.remove(runtime.config.uds_path)

    server = await asyncio.start_unix_server(
        lambda r, w: handle_control_command(r, w, runtime, logger),
        path=str(runtime.config.uds_path),
    )
    os.chmod(runtime.config.uds_path, 0o600)  # Owner only

Supported Commands

Command Handler Description
status _handle_status Full fleet snapshot
list _handle_list Agent list
stop <name> _handle_stop Stop an agent
start <name> _handle_start Start an inactive agent
restart <name> _handle_restart Restart an agent
pause <name> _handle_pause Pause an agent
resume <name> _handle_resume Resume an agent

UDS Client

The CLI (vox.cli._send_uds_command) connects to the socket and sends commands:

async def _send_uds_command(cmd: str, args: list[str]) -> dict:
    uds_path = os.environ.get("VOX_UDS_PATH", _DEFAULT_UDS_PATH)
    reader, writer = await asyncio.open_unix_connection(uds_path)
    payload = json.dumps({"cmd": cmd, "args": args}) + "\n"
    writer.write(payload.encode())
    await writer.drain()
    response = await reader.readline()
    return json.loads(response.decode().strip())

REST API (aiohttp)

In parallel with the UDS control plane, VOXAPIServer (vox/api_server.py) exposes a REST API on port 8000:

Method Path Action
GET / System summary
GET /fleet Fleet snapshot
GET /agents Agent list
GET /agents/{id} Agent detail
GET /agents/{id}/commands Agent commands
GET /capabilities Registered capabilities
POST /agents/{id}/pause Pause agent
POST /agents/{id}/resume Resume agent
POST /agents/{id}/stop Stop agent
POST /agents/{id}/start Start agent
POST /agents/{id}/restart Restart agent

Event Buffering during PAUSED

EventQueue (vox/agents/lifecycle.py:54) is a fixed-capacity circular buffer (256 events by default) that retains messages when the agent is not in ACTIVE state.

Mechanism

class EventQueue:
    def __init__(self, capacity: int = 256):
        self._queue: deque[PendingEvent] = deque()
        self._dropped_count: int = 0

    def enqueue(self, event_name: str, kwargs: Dict[str, Any]) -> bool:
        if len(self._queue) >= self._capacity:
            self._dropped_count += 1
            return False
        self._queue.append(PendingEvent(event_name, kwargs))
        return True

    def drain(self) -> List[PendingEvent]:
        events = list(self._queue)
        self._queue.clear()
        return events

Integration in VOXAgent.emit()

async def emit(self, event_name: str, **kwargs) -> None:
    if self._state in (AgentState.PAUSED, AgentState.PAUSING):
        self._event_queue.enqueue(event_name, kwargs)
        return
    if self._state != AgentState.ACTIVE:
        return
    targets = self.event_router.get(event_name, [])
    for role in targets:
        await role.handle_event(event_name, **kwargs)

Resume Behavior

async def resume(self) -> None:
    self._state = AgentState.RESUMING
    pending = self._event_queue.drain()
    self._state = AgentState.ACTIVE
    self.logger.ok(f"Resumed with {len(pending)} queued events")
    for evt in pending:
        await self.emit(evt.event_name, **evt.kwargs)

Design implications:

  • No event loss during PAUSED (up to capacity limit)
  • Events exceeding capacity are dropped with a counter
  • FIFO order is preserved during drain
  • Events during BOOTING/STOPPING/FAILED are silently ignored (neither queued nor processed)