How Chat Streaming Works in Open Agents: Server-to-Client Architecture Explained

Open Agents implements chat streaming by wrapping durable workflow output in a cancel-aware ReadableStream, serializing chunks as NDJSON, and using React hooks with recovery policies to handle network interruptions and tab visibility changes.

The vercel-labs/open-agents repository delivers real-time AI responses using a resilient streaming architecture that bridges durable workflows with browser-based React hooks. Unlike simple SSE implementations, this system handles abrupt client disconnections, tab switching, and server-side workflow recovery through a carefully orchestrated pipeline of stream utilities and state machines.

Server-Side Chat Streaming Architecture

The server-side implementation centers on transforming workflow execution into an HTTP-streamed response that survives network volatility.

Starting the Workflow Stream

When a client initiates a chat, the POST /api/chat route in apps/web/app/api/chat/route.ts launches a durable workflow using runAgentWorkflow and immediately claims a streaming slot to prevent duplicate executions.

// apps/web/app/api/chat/route.ts (lines 20-30)
const run = await start(runAgentWorkflow, [
  { messages, chatId, sessionId, userId, modelId: mainModelSelection.id },
]);

// Atomically reserve the stream ID to prevent duplicate workflows
await compareAndSetChatActiveStreamId(chatId, null, run.runId);

The workflow returns a run object containing run.getReadable(), which provides a raw ReadableStream of WebAgentUIMessageChunk objects. Before sending this to the client, the system wraps it with cancellation logic.

Wrapping with Cancellation Support

The createCancelableReadableStream utility in apps/web/lib/chat/create-cancelable-readable-stream.ts (lines 16-71) intercepts the workflow stream to handle browser-initiated aborts gracefully. It creates a new ReadableStream that listens for cancellation signals and properly releases the underlying workflow reader.

// apps/web/lib/chat/create-cancelable-readable-stream.ts
export function createCancelableReadableStream<T>(source: ReadableStream<T>) {
  const reader = source.getReader();
  let isCancelled = false;

  const cancelReader = async () => {
    if (isCancelled) return;
    isCancelled = true;
    try { await reader.cancel(); } catch {}
    try { reader.releaseLock(); } catch {}
  };

  return new ReadableStream<T>({
    async pull(controller) {
      try {
        const { done, value } = await reader.read();
        if (done) { controller.close(); return; }
        controller.enqueue(value);
      } catch (e) {
        if (isCancelled || isAbortLikeError(e)) {
          controller.close();
        } else {
          controller.error(e);
        }
      }
    },
    async cancel() { await cancelReader(); },
  });
}

This wrapper ensures that when a client closes their tab, the workflow does not continue consuming resources indefinitely. It specifically catches AbortError, ResponseAborted, and 404 "not ok" messages to prevent resource leaks.

Returning the NDJSON Stream

Finally, the wrapped stream returns to the client via createUIMessageStreamResponse, which sets the x-workflow-run-id header for potential reconnection and serializes each chunk as newline-delimited JSON (NDJSON).

// apps/web/app/api/chat/route.ts (lines 77-84)
const stream = createCancelableReadableStream(
  run.getReadable<WebAgentUIMessageChunk>(),
);

return createUIMessageStreamResponse({
  stream,
  headers: { "x-workflow-run-id": run.runId },
});

Client-Side Stream Reliability

While the server produces the stream, the client implements sophisticated recovery mechanisms to handle stalled connections and interrupted workflows.

Monitoring with useStreamRecovery

The useStreamRecovery hook in apps/web/app/sessions/[sessionId]/chats/[chatId]/hooks/use-stream-recovery.ts (lines 24-68) registers event listeners for visibilitychange, focus, and online events to detect when the user returns to a tab or regains connectivity.

export function useStreamRecovery({
  sessionId, chatId, status, isChatInFlight,
  hasAssistantRenderableContent, retryChatStream,
}: UseStreamRecoveryParams) {
  const inFlightStartedAtRef = useRef<number | null>(null);
  const lastRecoveryAtRef = useRef(0);
  const probeInFlightRef = useRef(false);

  const maybeRecover = useCallback(() => {
    const decision = getStreamRecoveryDecision({
      now: Date.now(),
      lastRecoveryAt: lastRecoveryAtRef.current,
      status,
      hasAssistantRenderableContent,
      inFlightStartedAt: inFlightStartedAtRef.current,
      isProbeInFlight: probeInFlightRef.current,
    });

    if (decision === "retry-error") {
      retryChatStream({ auto: true });
    } else if (decision === "probe") {
      // Probe the server to check workflow status
    }
  }, [status, hasAssistantRenderableContent, retryChatStream]);
}

Recovery Policies and Decision Logic

The actual decision-making logic lives in stream-recovery-policy.ts (lines 19-73), which exports getStreamRecoveryDecision to determine whether the client should retry, probe, or wait.

export function getStreamRecoveryDecision({
  now, lastRecoveryAt, status,
  hasAssistantRenderableContent, inFlightStartedAt,
  isProbeInFlight, isVisibilityRecovery = false,
  minIntervalMs = STREAM_RECOVERY_MIN_INTERVAL_MS,
  stallMs = STREAM_RECOVERY_STALL_MS,
}: Options): StreamRecoveryDecision {
  // Prevent aggressive recovery attempts
  if (now - lastRecoveryAt < minIntervalMs) return "none";
  if (status === "error") return "retry-error";

  // When tab becomes visible and stream is ready, probe first
  if (isVisibilityRecovery && status === "ready") {
    return isProbeInFlight ? "none" : "probe";
  }

  // Only probe if we've been waiting longer than 4000ms (STREAM_RECOVERY_STALL_MS)
  if (status !== "submitted" || hasAssistantRenderableContent) return "none";
  if (inFlightStartedAt === null || now - inFlightStartedAt < stallMs) return "none";
  return isProbeInFlight ? "none" : "probe";
}

Key recovery constants defined in the policy file include:

  • STREAM_RECOVERY_STALL_MS: 4000ms threshold before considering a stream stalled
  • STREAM_RECOVERY_MIN_INTERVAL_MS: Minimum time between recovery attempts to prevent spam

Handling Tab Visibility and Network Changes

When the hook detects a stalled stream (status === "submitted" with no assistant content for over 4000ms), it schedules a recovery timeout. If the user switches tabs and returns, the hook sends a lightweight probe request to GET /api/sessions/:sessionId/chats to verify if the server-side workflow is still processing. If the server confirms active streaming, the hook triggers a soft retry via retryChatStream({auto:true, strategy:"soft"}), which reopens the connection without resetting workflow state.

Summary

  • Durable workflows generate chat streams through runAgentWorkflow, which executes agent logic and exposes a ReadableStream via run.getReadable().
  • Cancellation safety is enforced by createCancelableReadableStream in apps/web/lib/chat/create-cancelable-readable-stream.ts, which gracefully handles browser disconnects by calling reader.cancel() and releasing locks.
  • NDJSON serialization delivers message chunks to the browser with the x-workflow-run-id header enabling stream recovery and reconnection.
  • Automatic recovery is managed by the useStreamRecovery hook and stream-recovery-policy.ts, which monitor for 4000ms stalls and tab visibility changes to trigger probes or soft retries without duplicating workflow executions.

Frequently Asked Questions

How does Open Agents handle chat streaming when a user closes their browser tab?

According to the source code in apps/web/lib/chat/create-cancelable-readable-stream.ts, when a browser closes a connection, it triggers the wrapper stream's cancel() method. This method calls reader.cancel() on the underlying workflow reader and releases the lock, ensuring the durable workflow does not continue running indefinitely after the client disconnects. The wrapper specifically catches AbortError and ResponseAborted exceptions to prevent these from propagating as unhandled errors.

What triggers a stream recovery attempt in the Open Agents client?

The useStreamRecovery hook schedules a recovery attempt when three conditions align: the chat status is "submitted", no assistant content has rendered yet, and the stream has been in flight for longer than STREAM_RECOVERY_STALL_MS (4000ms). Additionally, if the user switches tabs and returns (visibilitychange event) while the status is "ready", the hook triggers a probe to verify if the server is still streaming, as implemented in apps/web/app/sessions/[sessionId]/chats/[chatId]/stream-recovery-policy.ts.

What is the difference between a probe and a retry in the chat streaming recovery system?

A probe is a lightweight GET request to check the current chat status on the server, used when the tab becomes visible again or when the stream appears stalled. A retry (specifically a "soft retry") reopens a new streaming connection to continue receiving messages without resetting the server-side workflow state. The getStreamRecoveryDecision function returns "probe" for status checks and "retry-error" only when the status explicitly equals "error", ensuring retries happen only when necessary rather than on every reconnection.

Where does the chat stream originate in the Open Agents architecture?

The stream originates in apps/web/app/api/chat/route.ts, where the start(runAgentWorkflow, ...) function launches a durable workflow. This workflow produces a ReadableStream of UI message chunks that the API wraps with createCancelableReadableStream before returning it via createUIMessageStreamResponse. The workflow itself runs the agent logic and streams generated content back through the Vercel Workflow runtime, as evidenced by the run.getReadable<WebAgentUIMessageChunk>() call.

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 →