From c65fc9e3c81c41aafe421aa90a00514b84343285 Mon Sep 17 00:00:00 2001 From: Devin Foley Date: Sun, 20 Sep 2026 12:43:25 -0700 Subject: [PATCH] fix: recover authentication and browser connection failures (#13724) fix: recover authentication and browser connection failures Include connect timeouts in the bounded retry policy for idempotent actor synchronization. Handle WebSocket constructor failures through existing reconnect paths and preserve HTTP polling while realtime is unavailable. Refresh visible company queries until the socket recovers and clear all fallback timers on hiding or unmount. Verify 172 focused tests, server/UI typechecks, UI build, and design token gates. Full workspace build/typecheck require the unavailable Rust toolchain; the full test run is tracked separately. Co-Authored-By: Paperclip --- doc/DATABASE.md | 6 + doc/DEVELOPING.md | 9 ++ .../cloud-tenant-transient-db-retry.test.ts | 22 ++-- server/src/middleware/auth.ts | 12 +- .../transcript/useLiveRunTranscripts.test.tsx | 31 +++++ .../transcript/useLiveRunTranscripts.ts | 13 +- .../LiveUpdatesProvider.recovery.test.tsx | 123 ++++++++++++++++++ ui/src/context/LiveUpdatesProvider.tsx | 37 +++++- ui/src/lib/websocket.test.ts | 28 ++++ ui/src/lib/websocket.ts | 9 ++ .../pages/AgentDetail.log-visibility.test.tsx | 37 +++++- ui/src/pages/AgentDetail.production.tsx | 7 +- ui/src/pages/AgentDetail.tsx | 7 +- 13 files changed, 311 insertions(+), 30 deletions(-) create mode 100644 ui/src/context/LiveUpdatesProvider.recovery.test.tsx create mode 100644 ui/src/lib/websocket.test.ts create mode 100644 ui/src/lib/websocket.ts diff --git a/doc/DATABASE.md b/doc/DATABASE.md index 64c8875fd2..fe3f2e0dc2 100644 --- a/doc/DATABASE.md +++ b/doc/DATABASE.md @@ -159,6 +159,12 @@ PostgreSQL backends and checks rejection, pool recovery, and transaction isolati Remove the patch when an upstream release passes these tests. Installs of the unmodified `postgres` package outside this workspace do not include the patch. +Trusted-header actor synchronization retries transient connection failures, +including `CONNECT_TIMEOUT`, at most twice. This retry applies only to the +idempotent actor synchronization operations, not arbitrary transactions. A +persistent outage still fails the request after the bounded retries; each +connection attempt remains subject to the configured database connect timeout. + ## Switching between modes The database mode is controlled by `DATABASE_URL`: diff --git a/doc/DEVELOPING.md b/doc/DEVELOPING.md index dd0698608c..2fc0f508f0 100644 --- a/doc/DEVELOPING.md +++ b/doc/DEVELOPING.md @@ -1575,3 +1575,12 @@ from stored configuration problems. Verify connection transport and endpoint fields before disabling a connection. Verify workspace ownership, active runs, Git state, and runtime-service readiness before closing a workspace. A missing URL or old workspace timestamp alone does not prove that a row is disposable. + +### Browser realtime connection recovery + +If a browser cannot construct a WebSocket, the live-update and run transcript +clients use their disconnected retry paths. While the company event stream is +disconnected, visible active queries refresh every 15 seconds. This fallback +stops when the socket opens, the tab is hidden, or the provider unmounts. A +reconnected socket also refreshes visible queries to recover missed events. +Run log views retain their existing HTTP polling fallback. diff --git a/server/src/__tests__/cloud-tenant-transient-db-retry.test.ts b/server/src/__tests__/cloud-tenant-transient-db-retry.test.ts index 94b9c5a8ef..a72e8a3c54 100644 --- a/server/src/__tests__/cloud-tenant-transient-db-retry.test.ts +++ b/server/src/__tests__/cloud-tenant-transient-db-retry.test.ts @@ -6,16 +6,17 @@ import { } from "../middleware/auth.ts"; /** The shape drizzle produces: a wrapper whose `cause` is the driver error. */ -function driverClosedError(code: string): Error { +function driverConnectionError(code: string): Error { const driver = Object.assign(new Error(`write ${code} db.example.internal:5432`), { code }); return new Error("Failed query: insert into \"companies\" (…)", { cause: driver }); } describe("isTransientDbConnectionError", () => { - it("detects a closed-connection code anywhere on the cause chain", () => { - expect(isTransientDbConnectionError(driverClosedError("CONNECTION_CLOSED"))).toBe(true); - expect(isTransientDbConnectionError(driverClosedError("CONNECTION_ENDED"))).toBe(true); - expect(isTransientDbConnectionError(driverClosedError("CONNECTION_DESTROYED"))).toBe(true); + it("detects a transient connection code anywhere on the cause chain", () => { + expect(isTransientDbConnectionError(driverConnectionError("CONNECT_TIMEOUT"))).toBe(true); + expect(isTransientDbConnectionError(driverConnectionError("CONNECTION_CLOSED"))).toBe(true); + expect(isTransientDbConnectionError(driverConnectionError("CONNECTION_ENDED"))).toBe(true); + expect(isTransientDbConnectionError(driverConnectionError("CONNECTION_DESTROYED"))).toBe(true); const bare = Object.assign(new Error("write CONNECTION_CLOSED host:5432"), { code: "CONNECTION_CLOSED", }); @@ -26,6 +27,7 @@ describe("isTransientDbConnectionError", () => { expect(isTransientDbConnectionError(new Error("boom"))).toBe(false); const unique = Object.assign(new Error("duplicate key"), { code: "23505" }); expect(isTransientDbConnectionError(unique)).toBe(false); + expect(isTransientDbConnectionError(driverConnectionError("28P01"))).toBe(false); expect(isTransientDbConnectionError(new Error("outer", { cause: unique }))).toBe(false); expect(isTransientDbConnectionError("CONNECTION_CLOSED")).toBe(false); expect(isTransientDbConnectionError(undefined)).toBe(false); @@ -33,11 +35,11 @@ describe("isTransientDbConnectionError", () => { }); describe("retryOnTransientDbConnectionError", () => { - it("retries after a transient closed connection", async () => { + it.each(["CONNECT_TIMEOUT", "CONNECTION_CLOSED"])("retries after %s", async (code) => { let calls = 0; const result = await retryOnTransientDbConnectionError(async () => { calls += 1; - if (calls === 1) throw driverClosedError("CONNECTION_CLOSED"); + if (calls === 1) throw driverConnectionError(code); return "ok"; }); expect(result).toBe("ok"); @@ -50,7 +52,7 @@ describe("retryOnTransientDbConnectionError", () => { let calls = 0; const result = await retryOnTransientDbConnectionError(async () => { calls += 1; - if (calls <= 2) throw driverClosedError("CONNECTION_CLOSED"); + if (calls <= 2) throw driverConnectionError("CONNECTION_CLOSED"); return "ok"; }); expect(result).toBe("ok"); @@ -68,12 +70,12 @@ describe("retryOnTransientDbConnectionError", () => { expect(calls).toBe(1); }); - it("propagates the failure once the replay budget is spent", async () => { + it.each(["CONNECT_TIMEOUT", "CONNECTION_CLOSED"])("propagates %s once the replay budget is spent", async (code) => { let calls = 0; await expect( retryOnTransientDbConnectionError(async () => { calls += 1; - throw driverClosedError("CONNECTION_CLOSED"); + throw driverConnectionError(code); }), ).rejects.toThrow("Failed query"); expect(calls).toBe(3); diff --git a/server/src/middleware/auth.ts b/server/src/middleware/auth.ts index e873aa309e..22286b7b3f 100644 --- a/server/src/middleware/auth.ts +++ b/server/src/middleware/auth.ts @@ -542,13 +542,15 @@ export function cloudActorHeaderSourceFromHeaders( } /** - * postgres.js codes for a connection the server side closed out from under - * an in-flight query — a pooled Postgres endpoint recycling or suspending + * postgres.js codes for connection establishment timing out or for a + * connection the server side closed out from under an in-flight query — + * a pooled Postgres endpoint recycling or suspending * (observed 2026-09-03 with a managed pooler closing the socket mid-INSERT). * The driver reconnects transparently on the next query; only the statement * that was on the wire is lost. */ const transientDbConnectionCodes = new Set([ + "CONNECT_TIMEOUT", "CONNECTION_CLOSED", "CONNECTION_ENDED", "CONNECTION_DESTROYED", @@ -556,7 +558,7 @@ const transientDbConnectionCodes = new Set([ /** * True when the error chain (drizzle wraps the driver error as `cause`) - * carries a postgres.js closed-connection code. Exported for tests. + * carries a postgres.js transient connection code. Exported for tests. */ export function isTransientDbConnectionError(error: unknown): boolean { for (let current: unknown = error; current instanceof Error; current = current.cause) { @@ -568,7 +570,7 @@ export function isTransientDbConnectionError(error: unknown): boolean { /** * Runs `run` and retries it up to twice when it fails on a transient - * closed-connection error. Two replays, not one: when a pooled endpoint + * connection error. Two replays, not one: when a pooled endpoint * suspends or recycles, EVERY pooled socket is dead at once, so the first * replay can draw another stale socket from the pool and fail identically * (observed 2026-09-12: retried actor resolution still surfacing @@ -587,7 +589,7 @@ export async function retryOnTransientDbConnectionError(run: () => Promise } /** - * Trusted-header actor resolution with a single transient-connection retry. + * Trusted-header actor resolution with bounded transient-connection retries. * The tenant sync inside is idempotent end to end — every write is an * upsert/on-conflict/delete and the write debounce records only after the * whole sync succeeds — so replaying it after a dropped connection is safe, diff --git a/ui/src/components/transcript/useLiveRunTranscripts.test.tsx b/ui/src/components/transcript/useLiveRunTranscripts.test.tsx index e82c282df4..b54d92b372 100644 --- a/ui/src/components/transcript/useLiveRunTranscripts.test.tsx +++ b/ui/src/components/transcript/useLiveRunTranscripts.test.tsx @@ -736,6 +736,37 @@ describe("useLiveRunTranscripts", () => { } }); + it("keeps HTTP log polling when the socket constructor fails and retries later", async () => { + vi.useFakeTimers(); + globalThis.WebSocket = class { + constructor() { throw new TypeError("WebSocket is not a constructor"); } + } as unknown as typeof WebSocket; + const runs = [{ id: "run-1", status: "running", adapterType: "codex_local" }]; + function Harness() { + useLiveRunTranscripts({ companyId: "company-1", runs }); + return null; + } + const root = createRoot(document.createElement("div")); + try { + await act(async () => root.render()); + expect(logMock).toHaveBeenCalledTimes(1); + await act(async () => vi.advanceTimersByTimeAsync(30_000)); + expect(logMock.mock.calls.length).toBeGreaterThan(1); + globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket; + await act(async () => vi.advanceTimersByTimeAsync(15_000)); + expect(FakeWebSocket.instances).toHaveLength(1); + await act(async () => FakeWebSocket.instances[0].triggerOpen()); + // Cleanup must not depend on the global constructor remaining available. + globalThis.WebSocket = undefined as unknown as typeof WebSocket; + } finally { + await act(async () => root.unmount()); + } + const calls = logMock.mock.calls.length; + await act(async () => vi.advanceTimersByTimeAsync(30_000)); + expect(logMock).toHaveBeenCalledTimes(calls); + expect(FakeWebSocket.instances[0].closeCalls).toHaveLength(1); + }); + it("backs off exponentially when the live event socket keeps failing", async () => { vi.useFakeTimers(); try { diff --git a/ui/src/components/transcript/useLiveRunTranscripts.ts b/ui/src/components/transcript/useLiveRunTranscripts.ts index cce877edb7..fa6a7f0695 100644 --- a/ui/src/components/transcript/useLiveRunTranscripts.ts +++ b/ui/src/components/transcript/useLiveRunTranscripts.ts @@ -9,6 +9,7 @@ import { heartbeatsApi } from "../../api/heartbeats"; import { buildTranscript, getUIAdapter, onAdapterChange, type RunLogChunk, type TranscriptEntry } from "../../adapters"; import { queryKeys } from "../../lib/queryKeys"; import { buildSameOriginWebSocketUrl } from "../../lib/websocket-url"; +import { tryCreateWebSocket } from "../../lib/websocket"; import { mergeRunLogChunks, parsePersistedLogContent, @@ -20,6 +21,8 @@ import { // durable fix is server push (SSE/websocket) for transcript deltas so idle tabs // do no periodic work at all; the constants below only reduce the churn of the // current polling approach. +const SOCKET_CONNECTING = 0; +const SOCKET_OPEN = 1; const LOG_POLL_INTERVAL_MS = 2000; const LOG_READ_LIMIT_BYTES = 256_000; // When realtime websocket updates are enabled, the frequent log poll is @@ -420,7 +423,11 @@ export function useLiveRunTranscripts({ const url = buildSameOriginWebSocketUrl( `/api/companies/${encodeURIComponent(companyId)}/events/ws`, ); - socket = new WebSocket(url); + socket = tryCreateWebSocket(url); + if (!socket) { + scheduleReconnect(); + return; + } socket.onopen = () => { if (closed) return; @@ -506,14 +513,14 @@ export function useLiveRunTranscripts({ socket.onmessage = null; socket.onerror = null; socket.onclose = null; - if (socket.readyState === WebSocket.CONNECTING) { + if (socket.readyState === SOCKET_CONNECTING) { // Defer the close until the handshake completes so the browser // does not emit a noisy "closed before the connection is established" // warning during rapid run teardown. socket.onopen = () => { socket?.close(1000, "live_run_transcripts_unmount"); }; - } else if (socket.readyState === WebSocket.OPEN) { + } else if (socket.readyState === SOCKET_OPEN) { socket.close(1000, "live_run_transcripts_unmount"); } } diff --git a/ui/src/context/LiveUpdatesProvider.recovery.test.tsx b/ui/src/context/LiveUpdatesProvider.recovery.test.tsx new file mode 100644 index 0000000000..6dcf566c00 --- /dev/null +++ b/ui/src/context/LiveUpdatesProvider.recovery.test.tsx @@ -0,0 +1,123 @@ +// @vitest-environment jsdom +import { act } from "react"; +import { createRoot, type Root } from "react-dom/client"; +import { QueryClient, QueryClientProvider, useQuery } from "@tanstack/react-query"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { queryKeys } from "../lib/queryKeys"; + +const { pushToast } = vi.hoisted(() => ({ pushToast: vi.fn() })); +vi.mock("./CompanyContext", () => ({ + useCompany: () => ({ selectedCompanyId: "company-1", selectedCompany: { id: "company-1" } }), +})); +vi.mock("./ToastContext", () => ({ useToastActions: () => ({ pushToast }) })); +vi.mock("../lib/router", () => ({ useLocation: () => ({ pathname: "/tasks" }) })); +vi.mock("../api/auth", () => ({ + authApi: { getSession: async () => ({ user: { id: "viewer" }, session: { id: "session", userId: "viewer" } }) }, +})); +vi.mock("../api/health", () => ({ healthApi: { get: async () => ({ deploymentMode: "authenticated" }) } })); +import { LiveUpdatesProvider } from "./LiveUpdatesProvider"; + +(globalThis as typeof globalThis & { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true; + +class Socket { + static instances: Socket[] = []; + readyState = 1; + onopen: (() => void) | null = null; + onclose: (() => void) | null = null; + onmessage = null; + onerror = null; + constructor() { Socket.instances.push(this); } + close() { this.readyState = 3; } +} + +describe("LiveUpdatesProvider connection recovery", () => { + let root: Root; + let container: HTMLDivElement; + let client: QueryClient; + const read = vi.fn(async () => "current task data"); + + beforeEach(() => { + vi.useFakeTimers(); + vi.spyOn(document, "visibilityState", "get").mockReturnValue("visible"); + read.mockClear(); + Socket.instances = []; + client = new QueryClient({ defaultOptions: { queries: { retry: false, gcTime: Infinity, staleTime: Infinity } } }); + client.setQueryData(queryKeys.auth.session, { user: { id: "viewer" }, session: { id: "session", userId: "viewer" } }); + client.setQueryData(queryKeys.health, { deploymentMode: "authenticated" }); + container = document.createElement("div"); + root = createRoot(container); + }); + + afterEach(async () => { + await act(async () => root.unmount()); + client.clear(); + vi.restoreAllMocks(); + vi.unstubAllGlobals(); + vi.useRealTimers(); + }); + + async function render() { + function Tasks() { + const { data } = useQuery({ queryKey: ["company-1", "tasks"], queryFn: read }); + return {data}; + } + await act(async () => root.render( + , + )); + await act(async () => vi.advanceTimersByTimeAsync(1)); + } + + it.each(["missing", "throws"])("polls visible data and resumes realtime when the constructor %s", async (failure) => { + vi.stubGlobal("WebSocket", failure === "missing" ? undefined : class { + constructor() { throw new DOMException("Blocked", "SecurityError"); } + }); + await render(); + expect(container.textContent).toBe("current task data"); + expect(read).toHaveBeenCalledTimes(1); + await act(async () => vi.advanceTimersByTimeAsync(15_000)); + expect(read).toHaveBeenCalledTimes(2); + + vi.stubGlobal("WebSocket", Socket); + await act(async () => vi.advanceTimersByTimeAsync(15_000)); + expect(Socket.instances).toHaveLength(1); + await act(async () => Socket.instances[0].onopen?.()); + const readsAfterRecovery = read.mock.calls.length; + await act(async () => vi.advanceTimersByTimeAsync(30_000)); + expect(read).toHaveBeenCalledTimes(readsAfterRecovery); + expect(Socket.instances).toHaveLength(1); + + // An ordinary disconnect must re-enable the same fallback. + await act(async () => Socket.instances[0].onclose?.()); + await act(async () => vi.advanceTimersByTimeAsync(15_000)); + expect(read.mock.calls.length).toBeGreaterThan(readsAfterRecovery); + }); + + it("stops fallback reads and connection retries while hidden and after unmount", async () => { + const construct = vi.fn(); + vi.stubGlobal("WebSocket", class { + constructor() { construct(); throw new Error("Unavailable"); } + }); + await render(); + const visibility = vi.spyOn(document, "visibilityState", "get"); + await act(async () => { + visibility.mockReturnValue("hidden"); + document.dispatchEvent(new Event("visibilitychange")); + }); + const attempts = construct.mock.calls.length; + await act(async () => vi.advanceTimersByTimeAsync(30_000)); + expect(read).toHaveBeenCalledTimes(1); + expect(construct).toHaveBeenCalledTimes(attempts); + await act(async () => { + visibility.mockReturnValue("visible"); + document.dispatchEvent(new Event("visibilitychange")); + }); + await act(async () => vi.advanceTimersByTimeAsync(15_001)); + expect(read.mock.calls.length).toBeGreaterThan(1); + await act(async () => root.unmount()); + const readsAtUnmount = read.mock.calls.length; + const attemptsAtUnmount = construct.mock.calls.length; + await act(async () => vi.advanceTimersByTimeAsync(30_000)); + expect(read).toHaveBeenCalledTimes(readsAtUnmount); + expect(construct).toHaveBeenCalledTimes(attemptsAtUnmount); + }); +}); diff --git a/ui/src/context/LiveUpdatesProvider.tsx b/ui/src/context/LiveUpdatesProvider.tsx index d6f365e9c5..8ea8c42dd0 100644 --- a/ui/src/context/LiveUpdatesProvider.tsx +++ b/ui/src/context/LiveUpdatesProvider.tsx @@ -49,10 +49,12 @@ import { extractCompanyPrefixFromPath, toCompanyRelativePath } from "../lib/comp import { useLocation } from "../lib/router"; import { agentRouteRef } from "../lib/utils"; import { buildSameOriginWebSocketUrl } from "../lib/websocket-url"; +import { tryCreateWebSocket } from "../lib/websocket"; const TOAST_COOLDOWN_WINDOW_MS = 10_000; const TOAST_COOLDOWN_MAX = 3; const RECONNECT_SUPPRESS_MS = 2000; +const DISCONNECTED_POLL_INTERVAL_MS = 15_000; const SOCKET_CONNECTING = 0; const SOCKET_OPEN = 1; const TERMINAL_RUN_STATUSES = new Set([ @@ -1903,6 +1905,7 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) { let closed = false; let reconnectAttempt = 0; let reconnectTimer: number | null = null; + let pollTimer: number | null = null; let socket: WebSocket | null = null; const clearReconnect = () => { @@ -1912,8 +1915,23 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) { } }; + const stopPolling = () => { + if (pollTimer !== null) { + window.clearInterval(pollTimer); + pollTimer = null; + } + }; + + const startPolling = () => { + if (closed || pollTimer !== null) return; + // Visible queries still need fresh state when realtime is unavailable. + pollTimer = window.setInterval(() => { + void queryClient.invalidateQueries({ type: "active" }, { cancelRefetch: false }); + }, DISCONNECTED_POLL_INTERVAL_MS); + }; + const scheduleReconnect = () => { - if (closed) return; + if (closed || reconnectTimer !== null) return; reconnectAttempt += 1; const delayMs = Math.min( 15000, @@ -1930,7 +1948,12 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) { const url = buildSameOriginWebSocketUrl( `/api/companies/${encodeURIComponent(liveCompanyId)}/events/ws`, ); - const nextSocket = new WebSocket(url); + const nextSocket = tryCreateWebSocket(url); + if (!nextSocket) { + startPolling(); + scheduleReconnect(); + return; + } socket = nextSocket; nextSocket.onopen = () => { @@ -1938,13 +1961,11 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) { closeSocketQuietly(nextSocket, "stale_connection"); return; } + stopPolling(); if (reconnectAttempt > 0) { gateRef.current.suppressUntil = Date.now() + RECONNECT_SUPPRESS_MS; - // Reconcile after a gap: events missed while disconnected can't be - // replayed yet, so refetch the event-sourced live-runs list once. - queryClient.invalidateQueries({ - queryKey: queryKeys.liveRuns(liveCompanyId), - }); + // Reconcile all visible data after a gap: missed events cannot be replayed. + void queryClient.invalidateQueries({ type: "active" }, { cancelRefetch: false }); } reconnectAttempt = 0; }; @@ -1989,6 +2010,7 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) { if (socket !== nextSocket) return; socket = null; if (closed) return; + startPolling(); scheduleReconnect(); }; }; @@ -2002,6 +2024,7 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) { closed = true; window.clearTimeout(connectTimer); clearReconnect(); + stopPolling(); const activeSocket = socket; socket = null; closeSocketQuietly(activeSocket, "provider_unmount"); diff --git a/ui/src/lib/websocket.test.ts b/ui/src/lib/websocket.test.ts new file mode 100644 index 0000000000..1054a0f547 --- /dev/null +++ b/ui/src/lib/websocket.test.ts @@ -0,0 +1,28 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { tryCreateWebSocket } from "./websocket"; + +afterEach(() => vi.unstubAllGlobals()); + +describe("tryCreateWebSocket", () => { + it.each([undefined, {}, () => undefined])("returns null for an unavailable constructor: %s", (value) => { + vi.stubGlobal("WebSocket", value); + expect(tryCreateWebSocket("ws://localhost/events")).toBeNull(); + }); + + it("returns null when browser policy prevents construction", () => { + vi.stubGlobal("WebSocket", class { + constructor() { throw new DOMException("Blocked", "SecurityError"); } + }); + expect(tryCreateWebSocket("ws://localhost/events")).toBeNull(); + }); + + it("uses the current constructor so a later retry can recover", () => { + vi.stubGlobal("WebSocket", undefined); + expect(tryCreateWebSocket("ws://localhost/events")).toBeNull(); + class Socket { + constructor(readonly url: string) {} + } + vi.stubGlobal("WebSocket", Socket); + expect(tryCreateWebSocket("ws://localhost/events")).toEqual(new Socket("ws://localhost/events")); + }); +}); diff --git a/ui/src/lib/websocket.ts b/ui/src/lib/websocket.ts new file mode 100644 index 0000000000..46633142cf --- /dev/null +++ b/ui/src/lib/websocket.ts @@ -0,0 +1,9 @@ +/** Constructor failures do not dispatch socket events. Callers must use their + * disconnected fallback and retry path when this returns null. */ +export function tryCreateWebSocket(url: string): WebSocket | null { + try { + return new WebSocket(url); + } catch { + return null; + } +} diff --git a/ui/src/pages/AgentDetail.log-visibility.test.tsx b/ui/src/pages/AgentDetail.log-visibility.test.tsx index b2dca3a615..d0a96b68a6 100644 --- a/ui/src/pages/AgentDetail.log-visibility.test.tsx +++ b/ui/src/pages/AgentDetail.log-visibility.test.tsx @@ -6,8 +6,8 @@ import { afterEach, expect, it, vi } from "vitest"; import { LogViewer } from "./AgentDetail"; import { LogViewer as ProductionLogViewer } from "./AgentDetail.production"; -const { log, empty } = vi.hoisted(() => ({ log: vi.fn(), empty: [] })); -vi.mock("../api/heartbeats", () => ({ heartbeatsApi: { log } })); +const { log, events, empty } = vi.hoisted(() => ({ log: vi.fn(), events: vi.fn(async () => []), empty: [] })); +vi.mock("../api/heartbeats", () => ({ heartbeatsApi: { log, events } })); vi.mock("@tanstack/react-query", async (original) => ({ ...await original(), useQuery: () => ({ data: empty }), @@ -21,7 +21,7 @@ vi.mock("../components/transcript/RunTranscriptView", () => ({ RunTranscriptView: ({ entries }: { entries: Array<{ chunk: string }> }) =>
{entries.map(line => line.chunk).join(" ")}
, })); (globalThis as typeof globalThis & { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true; -afterEach(() => { vi.restoreAllMocks(); log.mockReset(); }); +afterEach(() => { vi.restoreAllMocks(); vi.unstubAllGlobals(); vi.useRealTimers(); log.mockReset(); events.mockClear(); }); it.each([LogViewer, ProductionLogViewer])("retains legacy history and reads only the next offset on visibility recovery (%#)", async (Viewer) => { const visibility = vi.spyOn(document, "visibilityState", "get").mockReturnValue("visible"); @@ -61,3 +61,34 @@ it.each([LogViewer, ProductionLogViewer])("retains legacy history and reads only container.remove(); } }); + +it.each([LogViewer, ProductionLogViewer])("polls logs when WebSocket construction fails and recovers on retry (%#)", async (Viewer) => { + vi.useFakeTimers(); + vi.spyOn(document, "visibilityState", "get").mockReturnValue("visible"); + vi.stubGlobal("WebSocket", undefined); + log.mockResolvedValue({ content: "", nextOffset: 0 }); + const run = { id: "run-1", companyId: "company-1", agentId: "agent-1", status: "running", logRef: "log" } as HeartbeatRun; + const root = createRoot(document.createElement("div")); + const sockets: Array<{ onopen: (() => void) | null; close: () => void }> = []; + try { + await act(async () => root.render()); + const initialReads = log.mock.calls.length; + await act(async () => vi.advanceTimersByTimeAsync(2000)); + expect(log.mock.calls.length).toBeGreaterThan(initialReads); + expect(events).toHaveBeenCalled(); + vi.stubGlobal("WebSocket", class { + onopen = null; + close = vi.fn(); + constructor() { sockets.push(this); } + }); + await act(async () => vi.advanceTimersByTimeAsync(1000)); + expect(sockets).toHaveLength(1); + await act(async () => sockets[0].onopen?.()); + const readsWhenConnected = log.mock.calls.length; + await act(async () => vi.advanceTimersByTimeAsync(4000)); + expect(log).toHaveBeenCalledTimes(readsWhenConnected); + } finally { + await act(async () => root.unmount()); + } + expect(sockets[0].close).toHaveBeenCalled(); +}); diff --git a/ui/src/pages/AgentDetail.production.tsx b/ui/src/pages/AgentDetail.production.tsx index a64210b408..91efc76433 100644 --- a/ui/src/pages/AgentDetail.production.tsx +++ b/ui/src/pages/AgentDetail.production.tsx @@ -63,6 +63,7 @@ import { SourceResolvedFoldCallout } from "../components/SourceResolvedFoldCallo import { SourceResolvedFoldBadge } from "../components/SourceResolvedFoldBadge"; import { readSourceResolvedWatchdogFold } from "../lib/source-resolved-watchdog-fold"; import { buildSameOriginWebSocketUrl } from "../lib/websocket-url"; +import { tryCreateWebSocket } from "../lib/websocket"; import { formatCents, formatDate, relativeTime, formatTokens, visibleRunCostUsd } from "../lib/utils"; import { cn } from "../lib/utils"; import { describeRunRetryState } from "../lib/runRetryState"; @@ -4037,7 +4038,11 @@ export function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType const url = buildSameOriginWebSocketUrl( `/api/companies/${encodeURIComponent(run.companyId)}/events/ws`, ); - socket = new WebSocket(url); + socket = tryCreateWebSocket(url); + if (!socket) { + scheduleReconnect(); + return; + } socket.onopen = () => { setIsStreamingConnected(true); diff --git a/ui/src/pages/AgentDetail.tsx b/ui/src/pages/AgentDetail.tsx index fc833c2aa0..b6896c6f72 100644 --- a/ui/src/pages/AgentDetail.tsx +++ b/ui/src/pages/AgentDetail.tsx @@ -58,6 +58,7 @@ import { SourceResolvedFoldCallout } from "../components/SourceResolvedFoldCallo import { SourceResolvedFoldBadge } from "../components/SourceResolvedFoldBadge"; import { readSourceResolvedWatchdogFold } from "../lib/source-resolved-watchdog-fold"; import { buildSameOriginWebSocketUrl } from "../lib/websocket-url"; +import { tryCreateWebSocket } from "../lib/websocket"; import { formatDate, relativeTime, formatTokens, visibleRunCostUsd } from "../lib/utils"; import { cn } from "../lib/utils"; import { describeRunRetryState } from "../lib/runRetryState"; @@ -4119,7 +4120,11 @@ export function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType const url = buildSameOriginWebSocketUrl( `/api/companies/${encodeURIComponent(run.companyId)}/events/ws`, ); - socket = new WebSocket(url); + socket = tryCreateWebSocket(url); + if (!socket) { + scheduleReconnect(); + return; + } socket.onopen = () => { setIsStreamingConnected(true);