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);