Skip to content

GatewayBrowserClient leaves streams pending after an error-only frame #506

Description

@MicroMilo

Summary

For an error-only response, GatewayBrowserClient remains connected, retains the stream entry, and leaves next() pending. Late events can still be delivered after the error; only a later close clears and ends the stream.

Expected behavior

A transport error must reject or close all pending streams and clean connection state, or enter a bounded recovery state.

Actual behavior

For an error-only response, GatewayBrowserClient remains connected, retains the stream entry, and leaves next() pending. Late events can still be delivered after the error; only a later close clears and ends the stream.

Impact

Browser callers can wait indefinitely after a transport error and may receive stale events from an already failed connection.

Reproduction

In the browser client, start a stream and deliver a transport response containing only an error, without a normal close. Inspect connection state, the stream registry, and the pending next() promise, then deliver a late event. The error should reject/close the stream and clean state; the observed result is connected=true with the stream still pending and late events still delivered until a later close.

Minimal reproduction script

From the repository root, save this as repro_browser_error_only_stream.mts and run:

pnpm install --frozen-lockfile
pnpm exec tsx repro_browser_error_only_stream.mts
import {
  GatewayBrowserClient,
  type WebSocketLike,
} from "./src/web/client/GatewayBrowserClient.ts";

type Listener = (event: { data?: unknown; code?: number; reason?: string }) => void;

function delay(ms: number): Promise<void> {
  return new Promise((resolve) => setTimeout(resolve, ms));
}

async function observe<T>(promise: Promise<T>, timeoutMs = 25): Promise<Record<string, unknown>> {
  let settled = false;
  const outcome = await Promise.race([
    promise.then((value) => {
      settled = true;
      return { state: "resolved", value };
    }).catch((error: unknown) => {
      settled = true;
      return { state: "rejected", error: error instanceof Error ? error.message : String(error) };
    }),
    delay(timeoutMs).then(() => ({ state: "pending_after_timeout" })),
  ]);
  return { ...outcome, settledAtObservation: settled };
}

class BrowserSocket implements WebSocketLike {
  readonly readyState = 1;
  readonly sentFrames: Array<{ type?: string; id?: string; method?: string }> = [];
  private readonly listeners = new Map<string, Array<{ listener: Listener; once: boolean }>>();

  constructor() {
    queueMicrotask(() => this.emit("open", {}));
  }

  addEventListener(
    type: "open" | "message" | "close" | "error",
    listener: Listener,
    options?: { once?: boolean },
  ): void {
    const values = this.listeners.get(type) ?? [];
    values.push({ listener, once: options?.once === true });
    this.listeners.set(type, values);
  }

  send(data: string): void {
    const frame = JSON.parse(data) as { type?: string; id?: string; method?: string };
    this.sentFrames.push(frame);
    if (frame.type === "hello") {
      queueMicrotask(() => this.emit("message", {
        data: JSON.stringify({
          type: "hello_ok",
          protocolVersion: "1.0",
          serverVersion: "fixture",
          serverInfo: { mode: "remote", protocolVersion: "1.0", sessionCount: 0 },
        }),
      }));
    }
  }

  close(): void {
    this.emitClose();
  }

  emitError(): void {
    this.emit("error", {});
  }

  emitClose(): void {
    this.emit("close", { code: 1006, reason: "fixture-transport-close" });
  }

  emitEvent(id: string, final = false): void {
    this.emit("message", {
      data: JSON.stringify({
        type: "event",
        id,
        seq: 0,
        final,
        event: { type: "assistant_text_delta", text: "late-fixture-event" },
      }),
    });
  }

  emitResponseError(id: string): void {
    this.emit("message", {
      data: JSON.stringify({
        type: "response",
        id,
        ok: false,
        error: { code: "gateway_request_failed", message: "fixture stream failure" },
      }),
    });
  }

  private emit(type: string, event: { data?: unknown; code?: number; reason?: string }): void {
    const values = this.listeners.get(type) ?? [];
    for (const entry of [...values]) {
      entry.listener(event);
      if (entry.once) {
        const current = this.listeners.get(type) ?? [];
        const index = current.indexOf(entry);
        if (index >= 0) current.splice(index, 1);
      }
    }
  }
}

type PrivateState = {
  streams: Map<string, unknown>;
  pending: Map<string, unknown>;
  closed: boolean;
  connectError?: Error;
};

function stateOf(client: GatewayBrowserClient): Record<string, unknown> {
  const state = client as unknown as PrivateState;
  return {
    connected: client.connected,
    closed: state.closed,
    streamMapSize: state.streams.size,
    pendingMapSize: state.pending.size,
    connectError: state.connectError?.message ?? null,
  };
}

async function makeClient(ids: string[]): Promise<{ client: GatewayBrowserClient; socket: BrowserSocket }> {
  let socket: BrowserSocket | undefined;
  const client = new GatewayBrowserClient({
    url: "ws://fixture.invalid",
    token: "fixture-token",
    clientName: "test",
    newId: () => ids.shift() ?? "fixture-id-fallback",
    webSocketFactory: () => {
      socket = new BrowserSocket();
      return socket;
    },
  });
  await client.connect();
  return { client, socket: socket! };
}

function submit(client: GatewayBrowserClient): { id: string; iterator: AsyncIterator<unknown> } {
  const stream = client.submitTurn({
    sessionKey: "fixture-session",
    channelKey: "web",
    message: "fixture",
  });
  const iterator = stream[Symbol.asyncIterator]();
  return { id: ((client as unknown as { streams: Map<string, unknown> }).streams.keys().next().value as string), iterator };
}

async function errorOnly(): Promise<Record<string, unknown>> {
  const { client, socket } = await makeClient(["stream-error-only", "request-after-error"]);
  const { id, iterator } = submit(client);
  const next = iterator.next();
  socket.emitError();
  const duringError = await observe(next);
  const stateAfterError = stateOf(client);
  socket.emitError();
  const stateAfterRepeatedError = stateOf(client);
  const connectAfterError = await observe(client.connect());
  const request = client.request("list_sessions", {});
  const requestDuringError = await observe(request);
  const stateWithNewRequest = stateOf(client);
  socket.emitResponseError(id);
  const streamAfterResponseError = await observe(next);
  const stateAfterResponseForStream = stateOf(client);
  socket.emitClose();
  const afterClose = await observe(next);
  const requestAfterClose = await observe(request);
  const finalState = stateOf(client);
  client.close();
  return {
    duringError,
    stateAfterError,
    stateAfterRepeatedError,
    connectAfterError,
    requestDuringError,
    stateWithNewRequest,
    streamAfterResponseError,
    stateAfterResponseForStream,
    afterClose,
    requestAfterClose,
    finalState,
  };
}

async function errorThenLateMessage(): Promise<Record<string, unknown>> {
  const { client, socket } = await makeClient(["stream-late-message"]);
  const { id, iterator } = submit(client);
  const first = iterator.next();
  socket.emitError();
  await delay(30);
  const stateAfterError = stateOf(client);
  socket.emitEvent(id);
  const lateEvent = await observe(first);
  const stateAfterLateEvent = stateOf(client);
  const second = iterator.next();
  socket.emitClose();
  const afterClose = await observe(second);
  const finalState = stateOf(client);
  client.close();
  return { stateAfterError, lateEvent, stateAfterLateEvent, afterClose, finalState };
}

async function errorOnlyForAwait(): Promise<Record<string, unknown>> {
  const { client, socket } = await makeClient(["stream-for-await"]);
  const stream = client.submitTurn({
    sessionKey: "fixture-session",
    channelKey: "web",
    message: "fixture",
  });
  const trace: string[] = [];
  const consumer = (async () => {
    try {
      for await (const event of stream) {
        trace.push(event.type);
      }
      trace.push("completed");
    } catch {
      trace.push("caught");
    } finally {
      trace.push("finally");
    }
  })();
  socket.emitError();
  const afterError = await observe(consumer);
  const traceAfterError = [...trace];
  const stateAfterError = stateOf(client);
  socket.emitClose();
  const afterClose = await observe(consumer);
  const traceAfterClose = [...trace];
  const finalState = stateOf(client);
  client.close();
  return { afterError, traceAfterError, stateAfterError, afterClose, traceAfterClose, finalState };
}

async function errorThenLateFinal(): Promise<Record<string, unknown>> {
  const { client, socket } = await makeClient(["stream-late-final"]);
  const { id, iterator } = submit(client);
  const next = iterator.next();
  socket.emitError();
  await delay(30);
  socket.emitEvent(id, true);
  const afterLateFinal = await observe(next);
  const stateAfterLateFinal = stateOf(client);
  client.close();
  return { afterLateFinal, stateAfterLateFinal };
}

async function twoStreamsErrorOnly(): Promise<Record<string, unknown>> {
  const { client, socket } = await makeClient(["stream-one", "stream-two"]);
  const first = submit(client);
  const second = submit(client);
  const firstNext = first.iterator.next();
  const secondNext = second.iterator.next();
  socket.emitError();
  const observations = {
    first: await observe(firstNext),
    second: await observe(secondNext),
    stateAfterError: stateOf(client),
  };
  socket.emitClose();
  observations.firstAfterClose = await observe(firstNext);
  observations.secondAfterClose = await observe(secondNext);
  observations.finalState = stateOf(client);
  client.close();
  return observations;
}

async function closeThenError(): Promise<Record<string, unknown>> {
  const { client, socket } = await makeClient(["stream-close-then-error"]);
  const { iterator } = submit(client);
  const next = iterator.next();
  socket.emitClose();
  const afterClose = await observe(next);
  const stateAfterClose = stateOf(client);
  socket.emitError();
  const stateAfterLateError = stateOf(client);
  client.close();
  return { afterClose, stateAfterClose, stateAfterLateError };
}

const result = {
  candidateId: "cand-browser-stream-error-only-liveness",
  source: "src/web/client/GatewayBrowserClient.ts",
  transport: "sanitized in-memory WebSocketLike; no credentials, network, or provider",
  cases: {
    errorOnly: await errorOnly(),
    errorThenLateMessage: await errorThenLateMessage(),
    errorOnlyForAwait: await errorOnlyForAwait(),
    errorThenLateFinal: await errorThenLateFinal(),
    twoStreamsErrorOnly: await twoStreamsErrorOnly(),
    closeThenError: await closeThenError(),
  },
};

console.log(JSON.stringify(result, null, 2));

Relevant source locations

  • src/web/client/GatewayBrowserClient.ts:90-98
  • src/web/client/GatewayBrowserClient.ts:314-342

Suggested direction

Make the external-input path establish one durable, identity-bound state/receipt before returning success; propagate explicit terminal outcomes to every channel and client; and add a regression test for the reproduced boundary.

This report is about functional behavior, not security. The reproduction uses deterministic in-memory or isolated fixtures and contains no credentials or private data.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions