Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
66 changes: 66 additions & 0 deletions src/proxy/websocket-client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,73 @@ function identityHeaders(): Headers {
});
}

/**
* A stream whose readable never produces a frame, so the client's read stays
* pending until something settles it. `close()` resolves `closed`, which is all
* a real `WebSocketStream` does: it does not settle a read already awaiting the
* next frame.
*/
function silentStream(): {
stream: UpstreamWebSocketStream;
readCancelled: () => boolean;
} {
let cancelled = false;
let resolveClosed: (info: { closeCode?: number; reason?: string }) => void = () => {};
const closed = new Promise<{ closeCode?: number; reason?: string }>((resolve) => {
resolveClosed = resolve;
});

const readable = new ReadableStream<string | Uint8Array>({
cancel() {
cancelled = true;
},
});

return {
stream: {
opened: Promise.resolve({
readable,
writable: new WritableStream<string | Uint8Array>(),
}),
closed,
close(closeInfo) {
resolveClosed(closeInfo ?? {});
},
},
readCancelled: () => cancelled,
};
}

describe("upstream WebSocket client", () => {
it("settles a pending read when the socket closes", async () => {
// Regression guard for a leak that surfaced as an intermittent CI failure on
// unrelated pull requests: closing left `read()` awaiting the next frame, so
// `--trace-leaks` reported an async operation outliving the test.
const { stream, readCancelled } = silentStream();
const socket = new UpstreamWebSocket(
"ws://upstream.test/_ws",
identityHeaders(),
() => stream,
);

await new Promise<void>((resolve) => {
socket.onopen = () => resolve();
});
assertEquals(readCancelled(), false, "the read is pending while the socket is open");

const closed = new Promise<void>((resolve) => {
socket.onclose = () => resolve();
});
socket.close(1000, "done");
await closed;

assertEquals(
readCancelled(),
true,
"closing must settle the read, or the operation outlives the connection",
);
});

it("presents the proxy identity headers on the handshake", async () => {
const server = startUpstreamServer();
try {
Expand Down
8 changes: 8 additions & 0 deletions src/proxy/websocket-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@ async function toChunk(
export class UpstreamWebSocket {
#stream: UpstreamWebSocketStream;
#writer: WritableStreamDefaultWriter<string | Uint8Array> | null = null;
#reader: ReadableStreamDefaultReader<string | Uint8Array> | null = null;
#readyState: number = WebSocket.CONNECTING;
#writes: Promise<void> = Promise.resolve();
#settled = false;
Expand Down Expand Up @@ -137,6 +138,7 @@ export class UpstreamWebSocket {

async #pump(readable: ReadableStream<string | Uint8Array>): Promise<void> {
const reader = readable.getReader();
this.#reader = reader;
try {
while (true) {
const { done, value } = await reader.read();
Expand All @@ -146,6 +148,7 @@ export class UpstreamWebSocket {
} catch (error) {
this.#fail(error);
} finally {
this.#reader = null;
reader.releaseLock();
}
}
Expand All @@ -161,6 +164,11 @@ export class UpstreamWebSocket {
if (this.#settled) return;
this.#settled = true;
this.#readyState = WebSocket.CLOSED;
// Settle the read still pending on the socket. Closing the stream resolves
// `closed` and fires this handler, but it does not settle a `read()` that is
// already awaiting the next frame, so the operation would outlive the
// connection it belongs to.
this.#reader?.cancel().catch(() => {});
this.onclose?.(new CloseEvent("close", { code, reason, wasClean }));
}
}
Expand Down
Loading