Skip to content

GatewayBrowserClient reports malformed transport close as clean stream completion #505

Description

@MicroMilo

Summary

For close code 1002, the browser request rejects but the first pending stream.next() resolves done=true; a later next() rejects. After one partial assistant event, browser for-await completes normally without a final event while the Node client rejects the corresponding waiter.

Expected behavior

A malformed transport close must reject pending and partial stream consumers rather than appear as normal iterator completion, with Node and browser behavior aligned.

Actual behavior

For close code 1002, the browser request rejects but the first pending stream.next() resolves done=true; a later next() rejects. After one partial assistant event, browser for-await completes normally without a final event while the Node client rejects the corresponding waiter.

Impact

Browser consumers can skip error/retry handling and treat an incomplete turn as normally finished.

Reproduction

In the browser client, begin a stream, close the transport with a malformed protocol close such as code 1002, and observe the pending next() and a stream that has already received one assistant event. Compare the same sequence with the Node client. The browser path should reject incomplete consumers; the observed result is done=true/normal for the first waiter or for-await loop while the Node control rejects.

Minimal reproduction script

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

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

type Frame = {
  type?: string;
  id?: string;
  method?: string;
  final?: boolean;
  seq?: number;
  event?: unknown;
};

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

const TOKEN = "fixture-token";
const HELLO_OK = {
  type: "hello_ok",
  protocolVersion: "1.0",
  serverVersion: "fixture",
  serverInfo: { mode: "remote", protocolVersion: "1.0", sessionCount: 0 },
};

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

async function observe<T>(promise: Promise<T>, timeoutMs = 40): Promise<Record<string, unknown>> {
  return Promise.race([
    promise.then((result) => ({ state: "resolved", result })).catch((error: unknown) => ({
      state: "rejected",
      error: error instanceof Error ? error.message : String(error),
    })),
    delay(timeoutMs).then(() => ({ state: `pending_after_${timeoutMs}ms` })),
  ]);
}

function mapSizes(client: unknown): Record<string, number> {
  const value = client as {
    pending?: Map<unknown, unknown>;
    streams?: Map<unknown, unknown>;
  };
  return {
    pending: value.pending?.size ?? -1,
    streams: value.streams?.size ?? -1,
  };
}

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

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

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

  send(data: string): void {
    const frame = JSON.parse(data) as Frame;
    this.sentFrames.push(frame);
    if (frame.type === "hello") {
      queueMicrotask(() => this.emit("message", { data: JSON.stringify(HELLO_OK) }));
    }
  }

  close(): void {
    this.emitClose(1000, "fixture-client-close");
  }

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

  emitClose(code: number, reason: string): void {
    this.emit("close", { code, reason });
  }

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

  private emit(type: string, event: SocketEvent): void {
    const entries = [...(this.listeners.get(type) ?? [])];
    for (const entry of entries) {
      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 BrowserCase = {
  client: GatewayBrowserClient;
  socket: BrowserSocket;
  streamId: string;
  stream: AsyncIterable<unknown>;
};

async function makeBrowserCase(id: string): Promise<BrowserCase> {
  let socket: BrowserSocket | undefined;
  const client = new GatewayBrowserClient({
    url: "ws://fixture.invalid",
    token: TOKEN,
    clientName: "test",
    newId: () => id,
    webSocketFactory: () => {
      socket = new BrowserSocket();
      return socket;
    },
  });
  await client.connect();
  const stream = client.submitTurn({
    sessionKey: "fixture-session",
    channelKey: "web",
    message: "fixture-turn",
  });
  const request = socket?.sentFrames.find((frame) => frame.type === "request");
  if (!request?.id || !socket) {
    throw new Error("fixture failed to capture submit_turn request");
  }
  return { client, socket, streamId: request.id, stream };
}

async function consumeWithForAwait(stream: AsyncIterable<unknown>): Promise<Record<string, unknown>> {
  const eventTypes: string[] = [];
  try {
    for await (const event of stream) {
      eventTypes.push((event as { type?: string }).type ?? "unknown");
    }
    return { outcome: "completed_normally", eventTypes };
  } catch (error: unknown) {
    return {
      outcome: "caught_error",
      eventTypes,
      error: error instanceof Error ? error.message : String(error),
    };
  }
}

async function browserPendingClose(code: number, reason: string): Promise<Record<string, unknown>> {
  const { client, socket, streamId, stream } = await makeBrowserCase(`browser-pending-${code}`);
  const request = observe(client.request("describe_server", {}));
  const consumer = consumeWithForAwait(stream);
  await delay(0);
  socket.emitClose(code, reason);
  const result = await observe(consumer);
  const requestResult = await observe(request);
  const iterator = stream[Symbol.asyncIterator]();
  const followUpOne = await observe(iterator.next());
  const followUpTwo = await observe(iterator.next());
  const afterClose = { connected: client.connected, maps: mapSizes(client) };
  // A repeated close and a post-close event must not create another terminal
  // result or deliver data after the first close.
  socket.emitClose(1011, "late-close");
  socket.emitEvent(streamId, "after-close");
  client.close();
  return {
    close: { code, reason },
    requestResult,
    forAwait: result,
    followUpOne,
    followUpTwo,
    afterClose,
  };
}

async function browserPartialThenClose(): Promise<Record<string, unknown>> {
  const { client, socket, streamId, stream } = await makeBrowserCase("browser-partial");
  const consumer = consumeWithForAwait(stream);
  await delay(0);
  socket.emitEvent(streamId, "partial");
  await delay(0);
  socket.emitClose(1002, "malformed-frame");
  const result = await observe(consumer);
  const afterClose = { connected: client.connected, maps: mapSizes(client) };
  client.close();
  return { forAwait: result, afterClose };
}

async function browserQueuedThenClose(): Promise<Record<string, unknown>> {
  const { client, socket, streamId, stream } = await makeBrowserCase("browser-queued");
  socket.emitEvent(streamId, "queued-before-close");
  socket.emitClose(1002, "malformed-frame");
  const result = await observe(consumeWithForAwait(stream));
  const afterClose = { connected: client.connected, maps: mapSizes(client) };
  client.close();
  return { forAwait: result, afterClose };
}

class NodeSocket {
  static readonly OPEN = 1;
  readonly readyState = 1;
  private readonly listeners = new Map<
    string,
    Array<{ listener: SocketListener; once: boolean }>
  >();

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

  addEventListener(type: string, listener: SocketListener, options?: { once?: boolean }): void {
    const entries = this.listeners.get(type) ?? [];
    entries.push({ listener, once: options?.once === true });
    this.listeners.set(type, entries);
  }

  removeEventListener(type: string, listener: SocketListener): void {
    const entries = this.listeners.get(type) ?? [];
    this.listeners.set(type, entries.filter((entry) => entry.listener !== listener));
  }

  send(data: string): void {
    const frame = JSON.parse(data) as Frame;
    if (frame.type === "hello") {
      queueMicrotask(() => this.emit("message", { data: JSON.stringify(HELLO_OK) }));
    }
  }

  close(): void {
    this.emit("close", { code: 1000, reason: "fixture-client-close" });
  }

  emitClose(code: number, reason: string): void {
    this.emit("close", { code, reason });
  }

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

async function nodePendingClose(): Promise<Record<string, unknown>> {
  const originalWebSocket = (globalThis as { WebSocket?: unknown }).WebSocket;
  let socket: NodeSocket | undefined;
  (globalThis as unknown as { WebSocket: typeof NodeSocket }).WebSocket = class extends NodeSocket {
    constructor() {
      super();
      socket = this;
    }
  } as typeof NodeSocket;
  try {
    const client = new GatewayWsClient({
      url: "ws://fixture.invalid",
      token: TOKEN,
      clientName: "test",
    });
    await client.connect();
    const request = observe(client.request("describe_server", {}));
    const stream = client.stream("submit_turn", {
      sessionKey: "fixture-session",
      message: "fixture-turn",
    });
    const consumer = consumeWithForAwait(stream);
    await delay(0);
    socket?.emitClose(1002, "malformed-frame");
    const result = await observe(consumer);
    const requestResult = await observe(request);
    const afterClose = { maps: mapSizes(client) };
    client.close();
    return { requestResult, forAwait: result, afterClose };
  } finally {
    if (originalWebSocket === undefined) {
      delete (globalThis as { WebSocket?: unknown }).WebSocket;
    } else {
      (globalThis as { WebSocket?: unknown }).WebSocket = originalWebSocket;
    }
  }
}

const browser1000 = await browserPendingClose(1000, "normal-close");
const browser1002 = await browserPendingClose(1002, "malformed-frame");
const browser1011 = await browserPendingClose(1011, "server-error");
const browserPartial = await browserPartialThenClose();
const browserQueued = await browserQueuedThenClose();
const node = await nodePendingClose();

console.log(JSON.stringify({ browser1000, browser1002, browser1011, browserPartial, browserQueued, node }, null, 2));

Relevant source locations

  • src/web/client/GatewayBrowserClient.ts:327-368
  • src/web/client/GatewayBrowserClient.ts:400-450
  • src/gateway/client/GatewayWsClient.ts:186-195
  • src/gateway/client/GatewayWsClient.ts:223-248

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.

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