How the Socket.IO Event Bridge Synchronizes Live Agent Progress to the OpenHuman Frontend

OpenHuman's Rust core publishes typed domain events to a global event bus, which a dedicated Socket.IO bridge forwards to connected clients via room-targeted emissions—using explicit deduplication logic to prevent duplicate deliveries while streaming live agent and sub-agent progress to the React frontend.

The OpenHuman project implements a real-time architecture where a Rust backend streams execution updates to a React frontend via Socket.IO. This article examines the Socket.IO event bridge implemented in src/core/socketio.rs, tracing how internal domain events flow from the core's broadcast bus through deduplicated room-based delivery to browser-based Redux stores.

Event Emission Inside the Rust Core

The OpenHuman core operates as a Rust process that publishes a rich set of domain events on a global event bus (crate::core::bus::BUS). Agents, sub-agents, the web-chat dispatcher, dictation hot-key listener, overlay, notifications, transcription service, memory sync, orchestration, MCP setup, and session-expiry logic all publish typed events implementing Serialize.

Each event type represents a plain-old-data structure. For web-channel communication, the system specifically uses WebChannelEvent, which carries fields such as event, client_id, thread_id, request_id, and an optional subagent: Option<SubagentProgressDetail> field for nested progress tracking.

Bridge Creation and Event Forwarding

The attach_socketio function in src/core/socketio.rs initializes the Socket.IO layer on the HTTP server and spawns background bridges that consume internal events. The spawn_web_channel_bridge(io) function creates a Tokio task that subscribes to the event stream and forwards payloads to connected clients.

let mut rx = crate::openhuman::web_chat::subscribe_web_channel_events();
loop {
    let event = match rx.recv().await { … } ;
    emit_web_channel_event(&io_web, event);
}

The subscribe_web_channel_events function returns a broadcast receiver subscribed to the WebChannelEvent stream. When an event arrives, emit_web_channel_event serializes the payload using serde_json::to_value and emits it to targeted Socket.IO rooms.

Room-Based Delivery and Deduplication

Every Socket.IO client joins two distinct rooms upon connection:

  • Client-ID room – uniquely identified by the socket's client_id
  • Thread room – formatted as thread:<thread_id>, joined when the client emits thread:subscribe

The emit_web_channel_event function implements sophisticated routing logic to ensure exactly-once delivery. Because socketioxide does not automatically deduplicate emissions when a socket belongs to multiple rooms, the bridge explicitly excludes the primary room from the thread room broadcast:

let primary = event.client_id.clone();
let thread_room = (event.client_id != "system" && !event.thread_id.is_empty())
    .then(|| format!("thread:{}", event.thread_id));

io.to(primary.clone()).emit(&name, &payload);
if let Some(alias) = event_alias(&name) { io.to(primary.clone()).emit(alias, &payload); }

if let Some(tr) = thread_room {
    io.to(tr.clone())
        .except(primary.clone())   // dedup: prevent double send
        .emit(&name, &payload);
    if let Some(alias) = event_alias(&name) {
        io.to(tr).except(primary).emit(alias, &payload);
    }
}

The except(primary) call prevents the "double thinking" bug that previously caused streaming deltas to render twice when a client qualified for both its private room and the shared thread room.

Sub-Agent Progress Serialization

WebChannelEvent carries an optional subagent: Option<SubagentProgressDetail> field that enables granular progress tracking for spawned sub-agents. The core populates this structure for lifecycle events including subagent_spawned, subagent_iteration_start, subagent_tool_call, subagent_tool_result, and subagent_completed.

The bridge forwards the complete JSON payload without transformation, allowing the frontend to render live sub-agent rows—including progress bars, iteration counters, tool names, and elapsed time—without overloading the top-level event fields.

use crate::openhuman::web_chat::{emit_web_channel_event, WebChannelEvent, SubagentProgressDetail};

// Inside sub-agent execution:
let event = WebChannelEvent {
    event: "subagent_iteration_start".into(),
    client_id: client_id.clone(),
    thread_id: thread_id.clone(),
    request_id: request_id.clone(),
    subagent: Some(SubagentProgressDetail {
        mode: Some("typed".into()),
        child_iteration: Some(iteration),
        ..Default::default()
    }),
    ..Default::default()
};
emit_web_channel_event(&io, event);

Frontend Consumption and Reconnection Handling

The frontend socketService (app/src/services/socketService.ts) establishes the Socket.IO connection, passing a per-process bearer token during the handshake. Upon connect, the service automatically re-joins all active thread rooms to ensure in-flight turn streaming continues seamlessly after network interruptions:

this.socket.on('connect', () => {
  const threadState = store.getState().thread;
  const roomThreadIds = new Set<string>(Object.keys(threadState?.activeThreadIds ?? {}));
  if (threadState?.selectedThreadId) roomThreadIds.add(threadState.selectedThreadId);
  for (const threadId of roomThreadIds) {
    this.socket?.emit('thread:subscribe', { thread_id: threadId });
  }
});

The service registers listeners for event families including chat_*, dictation:*, memory:*, flow:*, and auth:session_expired. Each listener dispatches browser CustomEvent instances (e.g., openhuman:memory-sync-stage) or updates Redux state directly. For live agent progress, the frontend listens to chat_message, tool_call, tool_result, and subagent_* events, updating the timeline UI by reading the nested subagent object within the JSON payload.

import { socketService } from './socketService';

socketService.on('subagent_tool_call', (payload) => {
  window.dispatchEvent(
    new CustomEvent('openhuman:subagent-progress', { detail: payload })
  );
});

Summary

  • Internal Bus Architecture: The Rust core publishes typed WebChannelEvent structures to crate::core::bus::BUS, which the Socket.IO bridge consumes via subscribe_web_channel_events.
  • Bridge Implementation: src/core/socketio.rs contains attach_socketio and spawn_web_channel_bridge, which spawn Tokio tasks that serialize events using serde_json::to_value and emit them via Socket.IO.
  • Deduplication Logic: The emit_web_channel_event function targets both client-ID and thread rooms, using except(primary) to prevent duplicate delivery when sockets belong to both rooms.
  • Sub-Agent Granularity: Events include an optional SubagentProgressDetail structure that exposes iteration counts, tool calls, and completion states without polluting the main event namespace.
  • Resilient Frontend: app/src/services/socketService.ts handles automatic room re-subscription on reconnect and dispatches events to Redux or DOM CustomEvents for UI rendering.

Frequently Asked Questions

How does the Socket.IO event bridge prevent duplicate events when targeting multiple rooms?

The bridge explicitly calls except(primary.clone()) when emitting to the thread room, excluding the client-ID room from the broadcast. Because socketioxide does not deduplicate emissions when a socket belongs to both the individual client room and the shared thread room, this exclusion guarantees each client receives exactly one copy of every event, eliminating the "double thinking" bug that previously duplicated streaming deltas.

What happens to live agent progress updates when the frontend reconnects?

Upon connect, the frontend socketService queries the Redux store for active thread IDs and the currently selected thread, then emits thread:subscribe for each room. This ensures that any in-flight agent turns or sub-agent executions continue streaming progress updates immediately after the connection restores, without requiring a full page refresh or manual re-subscription.

How are sub-agent progress details structured differently from main agent events?

Rather than creating separate top-level event types for sub-agents, the system embeds a subagent: Option<SubagentProgressDetail> field within the standard WebChannelEvent structure. This optional payload contains sub-agent-specific metadata including child_iteration, mode, and tool call details, allowing the frontend to render hierarchical progress indicators while maintaining a consistent event schema across the bridge.

Which Rust module initializes the Socket.IO server attachment?

The src/core/socketio.rs module exposes attach_socketio, which integrates the Socket.IO layer with the HTTP server and spawns the background bridge tasks including spawn_web_channel_bridge. This module serves as the sole integration point between the internal event bus and external Socket.IO clients according to the OpenHuman source code architecture.

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 →