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
15 changes: 14 additions & 1 deletion packages/webui/src/client/ConnectionStatus.tsx
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import type { ReactElement } from "react";
import { useSessionRuntimeState } from "./session-runtime-store.js";
import { useWebuiEventChannelDegraded } from "./connection-health.js";
import type { WebuiStreamState } from "./stream.js";

/**
Expand Down Expand Up @@ -88,7 +89,19 @@ export function ConnectionStatus({
retryLabel = "重试连接",
}: ConnectionStatusProps): ReactElement | null {
const { state } = useSessionRuntimeState(sessionId);
const connection = projectWebuiConnectionState(state.stream.phase);
// Two independent signals, merged: the selected session's stream phase (a
// turn that failed mid-flight) and the event channel's health (the
// always-on watcher, which is the only live link while the page is idle).
// Either one saying "trouble" shows the region; a terminal stream failure
// outranks a degraded-but-retrying channel.
const streamConnection = projectWebuiConnectionState(state.stream.phase);
const channelDegraded = useWebuiEventChannelDegraded();
const connection =
streamConnection === "failed"
? "failed"
: streamConnection === "reconnecting" || channelDegraded
? "reconnecting"
: "connected";
if (hideWhenConnected && connection === "connected") return null;
const copy = COPY[connection];
// A failure usually carries the server's own reason (`refusal`), and `status`
Expand Down
98 changes: 98 additions & 0 deletions packages/webui/src/client/connection-health.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
// Event-channel connection health — the transport's always-on watcher,
// surfaced.
//
// The WebUI holds no single "connection" to lose: every request opens its own
// WebSocket. The one socket that IS long-lived is the `watchEvents`
// subscription — the shell keeps one from boot, and while no turn is running
// it is the only live link to the server. When it drops between turns, the
// session stream's phase stays `idle` (which the connection banner reads as
// "connected"), the transport's reconnect loop retries silently every 250ms,
// and the page shows nothing at all. That silent gap is what this module
// closes: `watchEvents` reports its socket's health here, and the banner
// merges the signal with the per-session stream phase.
//
// Aggregation rule — a watcher is *degraded* only once it has been healthy
// and then went down. A watcher that has never been accepted (boot in
// progress, or a host that never answers events) is "unknown", not degraded:
// flagging it would flash 重连 at every page load, and the honest statement
// about a link that never existed is nothing, not "reconnecting". With
// several watchers mounted (the shell's plus the workspace panels'), any
// once-healthy watcher being down degrades the channel; the transport owns
// the reporting, so every current and future `watchEvents` caller is covered
// without each of them wiring callbacks.

import { useEffect, useState } from "react";

interface WebuiWatcherHealth {
readonly token: number;
healthy: boolean;
everHealthy: boolean;
}

const watchers = new Map<number, WebuiWatcherHealth>();
const listeners = new Set<() => void>();
let nextToken = 0;

function degraded(): boolean {
for (const watcher of watchers.values())
if (watcher.everHealthy && !watcher.healthy) return true;
return false;
}

function notifyIfChanged(before: boolean): void {
if (degraded() !== before) for (const listener of listeners) listener();
}

/** Registers one `watchEvents` subscription. Called by the transport only. */
export function registerWebuiEventWatcher(): number {
nextToken += 1;
watchers.set(nextToken, { token: nextToken, healthy: false, everHealthy: false });
return nextToken;
}

/** The watcher's `watchEvents` request was accepted — its socket is live. */
export function markWebuiEventWatcherHealthy(token: number): void {
const watcher = watchers.get(token);
if (!watcher) return;
const before = degraded();
watcher.healthy = true;
watcher.everHealthy = true;
notifyIfChanged(before);
}

/** The watcher's socket closed — it will reconnect on the transport's own loop. */
export function markWebuiEventWatcherDown(token: number): void {
const watcher = watchers.get(token);
if (!watcher) return;
const before = degraded();
watcher.healthy = false;
notifyIfChanged(before);
}

/** The subscription unsubscribed (component unmounted); it no longer counts. */
export function unregisterWebuiEventWatcher(token: number): void {
const before = degraded();
watchers.delete(token);
notifyIfChanged(before);
}

/**
* Whether the event channel has lost a link it previously held. Subscribe
* through the hook; read imperatively only for assertions in tests.
*/
export function isWebuiEventChannelDegraded(): boolean {
return degraded();
}

export function useWebuiEventChannelDegraded(): boolean {
const [value, setValue] = useState(degraded);
useEffect(() => {
setValue(degraded());
const listener = () => setValue(degraded());
listeners.add(listener);
return () => {
listeners.delete(listener);
};
}, []);
return value;
}
15 changes: 15 additions & 0 deletions packages/webui/src/client/transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,12 @@ import type {
WebuiGoalPatchRequest,
WebuiGoalEnabledResult,
} from "../server/port.js";
import {
markWebuiEventWatcherDown,
markWebuiEventWatcherHealthy,
registerWebuiEventWatcher,
unregisterWebuiEventWatcher,
} from "./connection-health.js";

declare const document: {
readonly visibilityState: string;
Expand Down Expand Up @@ -248,6 +254,12 @@ export function createWebuiTransport({
*/
onReconnect?: () => void,
): () => void {
// This socket is the page's only long-lived link while no turn runs, so
// its health is the user's "am I still connected" signal between turns.
// Reported from inside the transport so every watcher — the shell's and
// any panel's — counts without each call site wiring callbacks; the
// store's ever-healthy rule keeps boot quiet.
const watcherToken = registerWebuiEventWatcher();
let stopped = false;
let socket: WebuiSocket | undefined;
let reconnectTimer: ReturnType<typeof setTimeout> | undefined;
Expand Down Expand Up @@ -287,6 +299,7 @@ export function createWebuiTransport({
if (frame.requestId !== requestId) return;
if (frame.kind === "error") return;
if (frame.kind === "response") {
markWebuiEventWatcherHealthy(watcherToken);
if (acknowledged) return;
acknowledged = true;
// Fires once per connection: the first ack is the initial
Expand All @@ -302,6 +315,7 @@ export function createWebuiTransport({
ws.addEventListener("close", () => {
if (stopped || socket !== ws) return;
socket = undefined;
markWebuiEventWatcherDown(watcherToken);
reconnectTimer = setTimeout(connect, 250);
});
ws.addEventListener("error", () => undefined);
Expand Down Expand Up @@ -332,6 +346,7 @@ export function createWebuiTransport({
if (typeof window !== "undefined")
window.removeEventListener("online", reconnectWhenAvailable);
socket?.close();
unregisterWebuiEventWatcher(watcherToken);
};
}

Expand Down
200 changes: 200 additions & 0 deletions packages/webui/test/unit/connection-health.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,200 @@
// Unit tests for the event-channel health store and the transport's
// reporting into it.
//
// Q-b's defect, reproduced first against the built client: while the page is
// idle (no turn running), killing every socket and gating new ones produced
// NO visible signal at all — the connection banner reads the session
// stream's phase, which stays `idle` ("connected") between turns, and the
// transport's 250ms silent reconnect loop told nobody. The watcher socket is
// the only long-lived link in that state, so its health is the signal.
//
// Two layers here:
// * the store's aggregation rule — degraded means "lost a link it
// previously held"; a never-accepted watcher is unknown, not degraded,
// so a booting page never flashes a reconnecting banner;
// * the transport wiring — `watchEvents` registers, marks healthy on the
// server's acceptance, marks down on socket close, and unregisters on
// unsubscribe, so every current and future watcher counts without its
// call site changing.
//
// The store is module-level (the runtime-store precedent), so every case
// unregisters in `afterEach` — a leaked once-healthy watcher would degrade
// every later suite in this worker.

import { afterEach, describe, expect, it, vi } from "vitest";

import {
isWebuiEventChannelDegraded,
markWebuiEventWatcherDown,
markWebuiEventWatcherHealthy,
registerWebuiEventWatcher,
unregisterWebuiEventWatcher,
} from "../../src/client/connection-health.js";
import { createWebuiTransport } from "../../src/client/transport.js";

const tracked: number[] = [];

function register(): number {
const token = registerWebuiEventWatcher();
tracked.push(token);
return token;
}

afterEach(() => {
for (const token of tracked.splice(0)) unregisterWebuiEventWatcher(token);
expect(isWebuiEventChannelDegraded()).toBe(false);
});

describe("the watcher health store", () => {
it("treats a never-accepted watcher as unknown, not degraded", () => {
// Boot shape: the socket exists but the server has not answered yet.
// Flagging this would flash 正在重连 on every page load.
register();
expect(isWebuiEventChannelDegraded()).toBe(false);
});

it("degrades only after a healthy watcher goes down", () => {
const token = register();
markWebuiEventWatcherHealthy(token);
expect(isWebuiEventChannelDegraded()).toBe(false);
markWebuiEventWatcherDown(token);
expect(isWebuiEventChannelDegraded()).toBe(true);
markWebuiEventWatcherHealthy(token);
expect(isWebuiEventChannelDegraded()).toBe(false);
});

it("aggregates: any once-healthy watcher down degrades the channel", () => {
// The shell keeps one watcher from boot; the workspace panels open more.
// One dying link is a degraded channel even while the others stream.
const shell = register();
const panel = register();
markWebuiEventWatcherHealthy(shell);
markWebuiEventWatcherHealthy(panel);
markWebuiEventWatcherDown(panel);
expect(isWebuiEventChannelDegraded()).toBe(true);
markWebuiEventWatcherHealthy(panel);
expect(isWebuiEventChannelDegraded()).toBe(false);
});

it("unregistering a down watcher stops it from counting", () => {
// A panel unmounting its watcher mid-outage must not pin the channel
// degraded forever — the shell's link is the one that matters.
const shell = register();
const panel = register();
markWebuiEventWatcherHealthy(shell);
markWebuiEventWatcherHealthy(panel);
markWebuiEventWatcherDown(panel);
expect(isWebuiEventChannelDegraded()).toBe(true);
unregisterWebuiEventWatcher(panel);
expect(isWebuiEventChannelDegraded()).toBe(false);
});

it("marks on unknown tokens are ignored", () => {
markWebuiEventWatcherHealthy(999_999);
markWebuiEventWatcherDown(999_999);
unregisterWebuiEventWatcher(999_999);
expect(isWebuiEventChannelDegraded()).toBe(false);
});
});

/** Minimal WebuiSocket double: enough open/message/close/error for watchEvents. */
class FakeSocket {
readonly listeners = new Map<string, Array<(event: { data?: unknown }) => void>>();
readonly sent: string[] = [];
closed = false;

addEventListener(type: "open" | "message" | "error" | "close", listener: (event: { data?: unknown }) => void): void {
const list = this.listeners.get(type) ?? [];
list.push(listener);
this.listeners.set(type, list);
}

emit(type: "open" | "message" | "error" | "close", event: { data?: unknown } = {}): void {
for (const listener of this.listeners.get(type) ?? []) listener(event);
}

send(data: string): void {
this.sent.push(data);
}

close(): void {
if (this.closed) return;
this.closed = true;
this.emit("close");
}
}

function watchEventsHarness() {
const sockets: FakeSocket[] = [];
class Socket extends FakeSocket {
constructor() {
super();
sockets.push(this);
}
}
const transport = createWebuiTransport({
websocketUrl: "ws://harness.invalid",
token: "t",
webSocket: Socket as unknown as new (url: string) => FakeSocket,
});
return { transport, sockets };
}

/** Drives one watcher socket through acceptance. */
function accept(socket: FakeSocket): void {
socket.emit("open");
const request = JSON.parse(socket.sent[0] ?? "{}") as { requestId?: string };
socket.emit("message", {
data: JSON.stringify({ protocolVersion: 1, kind: "response", requestId: request.requestId, body: { ok: true } }),
});
}

describe("the transport's watcher reporting", () => {
it("marks healthy on the server's acceptance and down on socket close", () => {
const { transport, sockets } = watchEventsHarness();
const unsubscribe = transport.watchEvents(() => undefined);
expect(isWebuiEventChannelDegraded()).toBe(false);

const socket = sockets[0]!;
accept(socket);
expect(isWebuiEventChannelDegraded()).toBe(false);

socket.close();
expect(isWebuiEventChannelDegraded()).toBe(true);
unsubscribe();
});

it("clears the degradation when the reconnect is accepted", async () => {
const { transport, sockets } = watchEventsHarness();
const onReconnect = vi.fn();
const unsubscribe = transport.watchEvents(() => undefined, onReconnect);
accept(sockets[0]!);
sockets[0]!.close();
expect(isWebuiEventChannelDegraded()).toBe(true);

// The 250ms reconnect timer opens a fresh socket; accepting it restores
// the channel. Advance the timer by waiting it out for real — the delay
// is short, and fake timers would also have to drive the microtasks the
// acceptance path relies on.
await new Promise((resolve) => setTimeout(resolve, 320));
const reopened = sockets.at(-1)!;
expect(reopened).not.toBe(sockets[0]!);
accept(reopened);
expect(isWebuiEventChannelDegraded()).toBe(false);
// onReconnect fires per acceptance — the first ack plus this one.
expect(onReconnect).toHaveBeenCalledTimes(2);
unsubscribe();
});

it("unsubscribing unregisters the watcher", () => {
const { transport, sockets } = watchEventsHarness();
const unsubscribe = transport.watchEvents(() => undefined);
accept(sockets[0]!);
sockets[0]!.close();
expect(isWebuiEventChannelDegraded()).toBe(true);
unsubscribe();
// The watcher no longer counts, so the channel is not degraded even
// though its socket never came back before the unsubscribe.
expect(isWebuiEventChannelDegraded()).toBe(false);
});
});
Loading
Loading