From 94a979ef0eaacda936e7e903e16d9600090e1a0b Mon Sep 17 00:00:00 2001 From: testikun <320479488+testikun@users.noreply.github.com> Date: Fri, 4 Sep 2026 15:22:37 +0800 Subject: [PATCH 1/2] fix(web): recover from quiet or stalled connections --- tests/web/app-render.test.ts | 55 +++++++++++++++ tests/web/web-host.test.ts | 40 +++++++++++ web/host/web-host.ts | 47 +++++++++++-- web/ui/app.js | 133 +++++++++++++++++++++++++++++------ web/ui/styles.css | 4 ++ 5 files changed, 254 insertions(+), 25 deletions(-) diff --git a/tests/web/app-render.test.ts b/tests/web/app-render.test.ts index c35a3f88..4e59a666 100644 --- a/tests/web/app-render.test.ts +++ b/tests/web/app-render.test.ts @@ -348,6 +348,7 @@ async function renderApp( URL, TextDecoder, TextEncoder, + AbortController, Element: class Element {}, }; context.window = { @@ -401,6 +402,17 @@ async function renderApp( "sendPrompt", context as vm.Context, ) as () => Promise, + api: vm.runInContext("api", context as vm.Context) as ( + path: string, + options?: Record, + ) => Promise, + readEventChunk: vm.runInContext( + "readEventChunk", + context as vm.Context, + ) as ( + reader: { read(): Promise }, + timeoutMs?: number, + ) => Promise, updateComposer: vm.runInContext( "updateComposer", context as vm.Context, @@ -463,6 +475,49 @@ test("app.js resnapshots on an SSE cursor gap", async () => { assert.ok(app.readerCancellations() >= 1); }); +test("app.js uses quiet-stream heartbeats for bounded snapshot recovery", async () => { + const heartbeat = ": heartbeat\n\n"; + const app = await renderApp({ + eventRecords: [heartbeat.repeat(4)], + }); + assert.equal(app.eventFetches(), 1); + assert.ok(app.snapshotFetches() >= 2); + assert.equal(app.state.cursor, SNAPSHOT.cursor); +}); + +test("app.js bounds API waits and explains duplicate prompt admission", async () => { + const app = await renderApp(); + app.context.fetch = async ( + _url: unknown, + options?: { signal?: AbortSignal }, + ) => + new Promise((_resolve, reject) => { + options?.signal?.addEventListener("abort", () => { + const error = new Error("aborted"); + error.name = "AbortError"; + reject(error); + }); + }); + await assert.rejects( + app.api("/api/stuck", { timeoutMs: 5, timeoutMessage: "bounded timeout" }), + /bounded timeout/u, + ); + await assert.rejects( + app.readEventChunk({ read: () => new Promise(() => {}) }, 5), + /event stream stalled/u, + ); + + app.state.promptAdmissionPending = true; + const input = app.elements.get("prompt-input"); + assert.ok(input); + input.value = "another message"; + await app.sendPrompt(); + assert.equal( + app.elements.get("composer-hint")?.textContent, + "OpenPI is still accepting the previous message.", + ); +}); + test("app.js invalidates snapshots for cross-tab session metadata events", async () => { const app = await renderApp({ eventRecords: [ diff --git a/tests/web/web-host.test.ts b/tests/web/web-host.test.ts index d80df4fa..1f0ff162 100644 --- a/tests/web/web-host.test.ts +++ b/tests/web/web-host.test.ts @@ -697,6 +697,46 @@ async function startTestHost(runtime: WebRuntimeController) { return { host, launched, headers }; } +test("quiet SSE clients receive heartbeats without advancing the event cursor", async () => { + const cwd = await mkdtemp(join(tmpdir(), "openpi-web-heartbeat-")); + const host = new WebHost({ + runtime: testRuntime(cwd), + sseHeartbeatMs: 10, + }); + try { + await host.start(); + const launched = new URL(host.url); + const token = new URLSearchParams(launched.hash.slice(1)).get("token"); + assert.ok(token); + const headers = { Authorization: `Bearer ${token}` }; + const before = (await ( + await fetch(`${launched.origin}/api/snapshot`, { headers }) + ).json()) as { cursor: number }; + const response = await fetch( + `${launched.origin}/events?cursor=${before.cursor}`, + { headers }, + ); + assert.equal(response.status, 200); + assert.ok(response.body); + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let received = ""; + while (!received.includes(": heartbeat\n\n")) { + const chunk = await reader.read(); + assert.equal(chunk.done, false); + received += decoder.decode(chunk.value, { stream: true }); + } + const after = (await ( + await fetch(`${launched.origin}/api/snapshot`, { headers }) + ).json()) as { cursor: number }; + assert.equal(after.cursor, before.cursor); + await reader.cancel(); + } finally { + await host.stop(); + await rm(cwd, { recursive: true, force: true }); + } +}); + test("adapter initialization fails before the Host starts listening", async () => { const cwd = await mkdtemp(join(tmpdir(), "openpi-web-startup-failure-")); const runtime = testRuntime(cwd); diff --git a/web/host/web-host.ts b/web/host/web-host.ts index 6203e647..d9a59cf3 100644 --- a/web/host/web-host.ts +++ b/web/host/web-host.ts @@ -33,6 +33,7 @@ const MAX_COMMAND_BYTES = 16 * 1024; const MAX_SSE_CLIENTS = 8; const MAX_SSE_BUFFER_BYTES = 256 * 1024; const MAX_SSE_REPLAY_BYTES = MAX_SSE_BUFFER_BYTES; +const DEFAULT_SSE_HEARTBEAT_MS = 15_000; const SERVER_CLOSE_DRAIN_MS = 500; const DEFAULT_SHUTDOWN_TIMEOUT_MS = 5_000; const execFileAsync = promisify(execFile); @@ -45,6 +46,7 @@ export interface WebHostOptions { allowedOrigins?: readonly string[]; directoryChooser?: (signal: AbortSignal) => Promise; shutdownTimeoutMs?: number; + sseHeartbeatMs?: number; } export class WebHost { @@ -52,6 +54,10 @@ export class WebHost { private readonly token: Buffer; private readonly adapter: PiWebAdapter; private readonly clients = new Set(); + private readonly clientHeartbeats = new Map< + ServerResponse, + ReturnType + >(); private readonly events: WebEvent[] = []; private sequence = 0; private port = 0; @@ -63,6 +69,7 @@ export class WebHost { WebHostOptions["directoryChooser"] >; private readonly shutdownTimeoutMs: number; + private readonly sseHeartbeatMs: number; private readonly unsubscribeCapabilities: () => void; private readonly unsubscribeRuntime: () => void; private readonly chooserAbort = new AbortController(); @@ -90,6 +97,14 @@ export class WebHost { ) { throw new Error("Web host shutdown timeout must be a positive integer"); } + this.sseHeartbeatMs = + options.sseHeartbeatMs ?? DEFAULT_SSE_HEARTBEAT_MS; + if ( + !Number.isSafeInteger(this.sseHeartbeatMs) || + this.sseHeartbeatMs <= 0 + ) { + throw new Error("SSE heartbeat interval must be a positive integer"); + } this.adapter = new PiWebAdapter(options.runtime); this.onEvent = options.onEvent; this.unsubscribeCapabilities = subscribeWebCapabilities((scope) => { @@ -185,8 +200,7 @@ export class WebHost { client.writableLength > MAX_SSE_BUFFER_BYTES || !client.write(record) ) { - this.clients.delete(client); - client.destroy(); + this.removeSseClient(client, "destroy"); } } this.onEvent?.(event.type, event.detail); @@ -204,8 +218,7 @@ export class WebHost { this.unsubscribeCapabilities(); this.unsubscribeRuntime(); this.chooserAbort.abort(); - for (const client of this.clients) client.end(); - this.clients.clear(); + for (const client of [...this.clients]) this.removeSseClient(client, "end"); const closeServer = this.server.listening ? new Promise((resolve) => { const forceClose = setTimeout( @@ -736,7 +749,31 @@ export class WebHost { // ordering without treating normal backpressure as a broken client. for (const record of replay) response.write(record); this.clients.add(response); - response.on("close", () => this.clients.delete(response)); + const heartbeat = setInterval(() => { + if ( + response.destroyed || + response.writableEnded || + response.writableLength > MAX_SSE_BUFFER_BYTES || + !response.write(": heartbeat\n\n") + ) { + this.removeSseClient(response, "destroy"); + } + }, this.sseHeartbeatMs); + heartbeat.unref(); + this.clientHeartbeats.set(response, heartbeat); + response.on("close", () => this.removeSseClient(response)); + } + + private removeSseClient( + response: ServerResponse, + close?: "destroy" | "end", + ) { + this.clients.delete(response); + const heartbeat = this.clientHeartbeats.get(response); + if (heartbeat) clearInterval(heartbeat); + this.clientHeartbeats.delete(response); + if (close === "destroy" && !response.destroyed) response.destroy(); + else if (close === "end" && !response.writableEnded) response.end(); } private parseCursor(value: string | undefined | null) { diff --git a/web/ui/app.js b/web/ui/app.js index bcff449d..38ff60b9 100644 --- a/web/ui/app.js +++ b/web/ui/app.js @@ -30,6 +30,7 @@ const state = { snapshotGeneration: 0, livePhase: "idle", liveRetry: null, + composerFeedback: null, query: "", selectedWorkspace: null, language: navigator.language?.toLowerCase().startsWith("zh") ? "zh" : "en", @@ -67,6 +68,10 @@ const translations = { enterHint: "Enter to send, Shift+Enter for a new line.", activeOnlyHint: "Only the active Web session accepts messages.", acceptedHint: "Message accepted by OpenPI Web.", + admissionPendingHint: "OpenPI is still accepting the previous message.", + admissionTimeout: "Prompt admission timed out. Your draft was restored; try again.", + requestTimeout: "The Web request timed out.", + reconnectingHint: "Live updates were interrupted. Reconnecting and checking canonical state...", modelRunning: "Working...", modelPreparing: "Preparing task...", modelRetrying: "Retrying model request...", @@ -103,6 +108,10 @@ const translations = { enterHint: "按 Enter 发送,Shift+Enter 换行。", activeOnlyHint: "只有当前 Web 会话可以接收消息。", acceptedHint: "OpenPI Web 已接收消息。", + admissionPendingHint: "OpenPI 仍在接收上一条消息,请稍候。", + admissionTimeout: "消息接收超时,草稿已恢复,请重试。", + requestTimeout: "Web 请求已超时。", + reconnectingHint: "实时更新已中断,正在重连并核对权威状态……", modelRunning: "正在运行...", modelPreparing: "正在准备任务...", modelRetrying: "模型请求重试中...", @@ -127,6 +136,10 @@ function applyLanguage() { const $ = (id) => document.getElementById(id); const tokenStorageKey = "openpi.web.token"; +const DEFAULT_API_TIMEOUT_MS = 15_000; +const PROMPT_ADMISSION_TIMEOUT_MS = 30_000; +const SSE_STALE_TIMEOUT_MS = 45_000; +const HEARTBEATS_PER_SNAPSHOT = 4; const fragmentToken = new URLSearchParams(location.hash.slice(1)).get("token"); let token = fragmentToken; if (fragmentToken) { @@ -139,6 +152,21 @@ const headers = (json = false) => ({ Authorization: `Bearer ${token}`, ...(json ? { "Content-Type": "application/json" } : {}), }); + +function setComposerFeedback(message, kind = "status") { + state.composerFeedback = message ? { message, kind } : null; + const hint = $("composer-hint"); + if (!hint) return; + hint.textContent = message || ""; + for (const candidate of ["status", "error", "connection"]) { + hint.classList.toggle(candidate, candidate === kind && Boolean(message)); + } +} + +function clearComposerFeedback(kind) { + if (kind && state.composerFeedback?.kind !== kind) return; + setComposerFeedback(""); +} const escapeHtml = (value) => String(value ?? "").replace( /[&<>"']/g, @@ -196,13 +224,37 @@ function renderMarkdown(value) { } async function api(path, options = {}) { - const response = await fetch(path, { - ...options, - headers: { ...headers(Boolean(options.body)), ...options.headers }, - }); - const body = await response.json().catch(() => ({})); - if (!response.ok) throw new Error(body.error || `Request failed (${response.status})`); - return body; + const { + timeoutMs = DEFAULT_API_TIMEOUT_MS, + timeoutMessage = t("requestTimeout"), + ...requestOptions + } = options; + const controller = new AbortController(); + let timedOut = false; + const timer = window.setTimeout(() => { + timedOut = true; + controller.abort(); + }, timeoutMs); + timer.unref?.(); + try { + const response = await fetch(path, { + ...requestOptions, + signal: controller.signal, + headers: { + ...headers(Boolean(requestOptions.body)), + ...requestOptions.headers, + }, + }); + const body = await response.json().catch(() => ({})); + if (!response.ok) + throw new Error(body.error || `Request failed (${response.status})`); + return body; + } catch (error) { + if (timedOut) throw new Error(timeoutMessage); + throw error; + } finally { + window.clearTimeout(timer); + } } function sessionTitle(session) { @@ -527,11 +579,17 @@ function updateComposer() { : active ? t("promptMessage") : t("promptReadonly"); - $("composer-hint").textContent = canCompose + const defaultHint = canCompose ? state.snapshot.runtime.status === "running" || state.liveRunning ? t("queuedHint") : t("enterHint") : t("activeOnlyHint"); + const feedback = state.composerFeedback; + const hint = $("composer-hint"); + hint.textContent = feedback?.message || defaultHint; + for (const candidate of ["status", "error", "connection"]) { + hint.classList.toggle(candidate, feedback?.kind === candidate); + } } async function selectModel(value) { @@ -640,8 +698,7 @@ async function refreshSnapshot({ ) return false; $("connection-state").textContent = "Unavailable"; $("connection-state").classList.add("reconnecting"); - $("composer-hint").textContent = error.message; - $("composer-hint").classList.add("error"); + setComposerFeedback(error.message, "error"); return false; } } @@ -718,11 +775,14 @@ async function sendPrompt() { await chooseWorkspace(); } const content = $("prompt-input").value.trim(); + if (state.promptAdmissionPending) { + setComposerFeedback(t("admissionPendingHint")); + return; + } if ( !content || !state.selectedWorkspace || - state.sessionSwitching || - state.promptAdmissionPending + state.sessionSwitching ) return; if (!state.snapshot?.selectedSession?.id) { await createSession(state.selectedWorkspace); @@ -738,12 +798,14 @@ async function sendPrompt() { ].slice(-8); state.promptAdmissionPending = true; state.promptAdmissionToken = admissionToken; + clearComposerFeedback(); renderConversation(); - $("composer-hint").classList.remove("error"); try { const receipt = await api("/api/prompt", { method: "POST", body: JSON.stringify({ sessionId, content }), + timeoutMs: PROMPT_ADMISSION_TIMEOUT_MS, + timeoutMessage: t("admissionTimeout"), }); if (epoch !== state.sessionEpoch || state.promptAdmissionToken !== admissionToken) return; const alreadySettled = state.terminalPromptIds.has(receipt.id); @@ -751,7 +813,7 @@ async function sendPrompt() { state.livePhase = alreadySettled ? "idle" : "preparing"; $("prompt-input").value = ""; resizePrompt(); - $("composer-hint").textContent = t("acceptedHint"); + setComposerFeedback(t("acceptedHint")); scheduleSnapshotRefresh(120); } catch (error) { if (epoch !== state.sessionEpoch || state.promptAdmissionToken !== admissionToken) return; @@ -761,8 +823,7 @@ async function sendPrompt() { state.liveMessages = state.liveMessages.filter( (entry) => entry.key !== optimisticKey, ); - $("composer-hint").textContent = error.message; - $("composer-hint").classList.add("error"); + setComposerFeedback(error.message, "error"); } finally { if (epoch === state.sessionEpoch && state.promptAdmissionToken === admissionToken) { state.promptAdmissionPending = false; @@ -897,8 +958,7 @@ async function createSession(workspacePath) { } function showNotice(message) { - $("composer-hint").textContent = message; - $("composer-hint").classList.add("error"); + setComposerFeedback(message, "error"); } function openWorkspaceMenu(path, anchor) { @@ -1193,17 +1253,27 @@ async function connectEvents() { if (!response.ok || !response.body) throw new Error("event connection failed"); $("connection-state").textContent = "Connected"; $("connection-state").classList.remove("reconnecting"); + clearComposerFeedback("connection"); reconnectDelay = 500; reader = response.body.getReader(); const decoder = new TextDecoder(); let buffer = ""; + let heartbeatCount = 0; while (true) { - const { done, value } = await reader.read(); + const { done, value } = await readEventChunk(reader); if (done) throw new Error("event connection closed"); buffer += decoder.decode(value, { stream: true }); const records = buffer.split("\n\n"); buffer = records.pop() || ""; for (const record of records) { + if (record.split("\n").some((item) => item === ": heartbeat")) { + heartbeatCount++; + if (heartbeatCount >= HEARTBEATS_PER_SNAPSHOT) { + heartbeatCount = 0; + scheduleSnapshotRefresh(0); + } + continue; + } const line = record.split("\n").find((item) => item.startsWith("data: ")); if (!line) continue; const event = JSON.parse(line.slice(6)); @@ -1224,6 +1294,7 @@ async function connectEvents() { ? false : await refreshSnapshot({ resetCursor: true }); if (!recovered) resetLiveState(); + setComposerFeedback(t("reconnectingHint"), "connection"); await new Promise((resolve) => setTimeout(resolve, reconnectDelay)); reconnectDelay = Math.min(reconnectDelay * 2, 5_000); } finally { @@ -1232,6 +1303,24 @@ async function connectEvents() { } } +async function readEventChunk(reader, timeoutMs = SSE_STALE_TIMEOUT_MS) { + let timer; + try { + return await Promise.race([ + reader.read(), + new Promise((_, reject) => { + timer = window.setTimeout( + () => reject(new Error("event stream stalled")), + timeoutMs, + ); + timer.unref?.(); + }), + ]); + } finally { + if (timer !== undefined) window.clearTimeout(timer); + } +} + const collapseButton = $("collapse-sidebar"); const collapsedStorageKey = "openpi.sidebar-collapsed"; const setSidebarCollapsed = (collapsed) => { @@ -1384,7 +1473,11 @@ $("composer")?.addEventListener("submit", (event) => { if (state.selectedWorkspace) void sendPrompt(); else void chooseWorkspace(); }); -$("prompt-input")?.addEventListener("input", resizePrompt); +$("prompt-input")?.addEventListener("input", () => { + if (!state.promptAdmissionPending) clearComposerFeedback(); + resizePrompt(); + updateComposer(); +}); $("prompt-input")?.addEventListener("keydown", (event) => { if (event.isComposing || event.keyCode === 229) return; if (event.key === "Enter" && !event.shiftKey) { diff --git a/web/ui/styles.css b/web/ui/styles.css index 0270ed9a..54d4a402 100644 --- a/web/ui/styles.css +++ b/web/ui/styles.css @@ -507,7 +507,9 @@ body.sidebar-collapsed .icon-button { width: 36px; height: 36px; } .send-button:disabled { background: var(--warm-accent-disabled); color: var(--subtle); } .send-button svg { width: 17px; height: 17px; stroke-width: 2; } .composer-hint { display: none; } +.composer-hint.status, .composer-hint.error, .composer-hint.connection { display: block; margin: 7px 2px 0; color: var(--muted); font-size: 11px; line-height: 1.35; } .composer-hint.error { color: var(--error); } +.composer-hint.connection { color: #9a6700; } .sidebar-scrim { display: none; } @media (max-width: 760px) { @@ -551,6 +553,8 @@ body.sidebar-collapsed .icon-button { width: 36px; height: 36px; } .conversation { scroll-behavior: auto; } } +.composer-hint.status, .composer-hint.error, .composer-hint.connection { display: block; } + /* The conversation view should start with the actual messages. On narrow screens keep only the sidebar trigger as an overlay so navigation remains available without restoring the session header. */ From 72b334ded7b3dc0c77de066930546ef004fd797c Mon Sep 17 00:00:00 2001 From: testikun <320479488+testikun@users.noreply.github.com> Date: Fri, 4 Sep 2026 15:32:58 +0800 Subject: [PATCH 2/2] test(web): keep browser watchdog timers referenced --- web/ui/app.js | 2 -- 1 file changed, 2 deletions(-) diff --git a/web/ui/app.js b/web/ui/app.js index 38ff60b9..770f7ec7 100644 --- a/web/ui/app.js +++ b/web/ui/app.js @@ -235,7 +235,6 @@ async function api(path, options = {}) { timedOut = true; controller.abort(); }, timeoutMs); - timer.unref?.(); try { const response = await fetch(path, { ...requestOptions, @@ -1313,7 +1312,6 @@ async function readEventChunk(reader, timeoutMs = SSE_STALE_TIMEOUT_MS) { () => reject(new Error("event stream stalled")), timeoutMs, ); - timer.unref?.(); }), ]); } finally {