OpenBMB / OpenBMB/PilotDeck

GatewayBrowserClient reports malformed transport close as clean stream completion

Open
#505 0 comments 0 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

Dominant language
TypeScript
Stars
4k
Forks
453
Avg merge
12h 30m
Merged PRs (30d)
46

Description

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.

Contributor guide

No contributing guide indexed for this repository

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Start with src/web/client/GatewayBrowserClient.ts:327-368 and 400-450, then compare the corresponding handling in src/gateway/client/GatewayWsClient.ts:186-195 and 223-248. Run the provided repro_browser_stream_close.mts script with pnpm install --frozen-lockfile and pnpm exec tsx repro_browser_stream_close.mts. Done means malformed close code 1002 rejects pending and partial browser consumers, aligns with Node behavior, and has a regression test for the reproduced boundary.

Written by the indexing model from the issue text.

Assessment

Tech stack
typescript
Domain
api, testing-qa
Issue type
Bug
Difficulty
3/5
Estimated time
1-2 days
Activity status
Quiet
Clarity
Clearly specified
Newbie friendliness
72/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.