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

> Discover how the Socket.IO event bridge synchronizes live agent progress in OpenHuman. Learn about typed domain events, room-targeted emissions, and deduplication logic for seamless frontend updates.

- Repository: [Tiny Humans/openhuman](https://github.com/tinyhumansai/openhuman)
- Tags: how-to-guide
- Published: 2026-08-27

---

**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`](https://github.com/tinyhumansai/openhuman/blob/main/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`](https://github.com/tinyhumansai/openhuman/blob/main/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.

```rust
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:

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

```rust
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`](https://github.com/tinyhumansai/openhuman/blob/main/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:

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

```typescript
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`](https://github.com/tinyhumansai/openhuman/blob/main/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`](https://github.com/tinyhumansai/openhuman/blob/main/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`](https://github.com/tinyhumansai/openhuman/blob/main/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.