Merge pull request #315 from djtzemx/fix/dsh-plugin-observation-stream-reconnect

fix(dsh-plugin): reconnect the observation stream after a fatal failure
This commit is contained in:
Zhang GH
2026-09-27 00:25:59 +08:00
committed by GitHub
5 changed files with 281 additions and 13 deletions
@@ -239,14 +239,16 @@ function StopSessionAction(props: {
);
}
type ChromeState = "active" | "idle" | "error" | "reconnecting";
/** Flat status dot, specced after the BSK popup's ConnectionStatusIndicator. */
function StatusDot({ state }: { state: "active" | "idle" | "error" | "dead" }) {
function StatusDot({ state }: { state: ChromeState | "dead" }) {
const color =
state === "active"
? "bg-emerald-500"
: state === "error"
? "bg-red-500"
: state === "dead"
: state === "dead" || state === "reconnecting"
? "bg-amber-500"
: "bg-muted-foreground/40";
return (
@@ -325,6 +327,8 @@ export function OverlayBody(props: {
focus: SessionObservation | undefined;
sessions: readonly SessionObservation[];
available: boolean;
/** The live feed dropped; sessions shown are the last ones received. */
reconnecting: boolean;
pinnedId: string | null;
onTogglePin: (sessionId: string) => void;
now: number;
@@ -342,6 +346,7 @@ export function OverlayBody(props: {
focus,
sessions,
available,
reconnecting,
pinnedId,
onTogglePin,
now,
@@ -380,12 +385,20 @@ export function OverlayBody(props: {
void store.interrupt(focus.sessionId).finally(() => setInterrupting(false));
};
const statusText = !available
? "browser unavailable"
: focus === undefined
? "no session"
: `${focus.sessionId} · ${focus.action === "idle" ? "idle" : focus.action} · ${formatElapsed(focus.since, now)}`;
const state = !available ? "error" : focus !== undefined ? statusOf(focus) : "idle";
const statusText = reconnecting
? "reconnecting…"
: !available
? "browser unavailable"
: focus === undefined
? "no session"
: `${focus.sessionId} · ${focus.action === "idle" ? "idle" : focus.action} · ${formatElapsed(focus.since, now)}`;
const state: ChromeState = reconnecting
? "reconnecting"
: !available
? "error"
: focus !== undefined
? statusOf(focus)
: "idle";
return (
<div
@@ -401,7 +414,7 @@ export function OverlayBody(props: {
onPointerDown={onHeaderPointerDown}
role="presentation"
>
<StatusDot state={state === "error" ? "error" : state === "active" ? "active" : "idle"} />
<StatusDot state={state} />
<span className={css["status-text"]}>{statusText}</span>
{onUseFloating !== undefined ? (
<IconAction
@@ -620,6 +633,7 @@ export function ObservationOverlay({ store }: { store: ObservationClientStore })
focus={focus}
sessions={snapshot.sessions}
available={snapshot.available}
reconnecting={snapshot.reconnecting}
pinnedId={pinnedId}
onTogglePin={onTogglePin}
now={now}
@@ -654,21 +668,28 @@ export function ObservationOverlay({ store }: { store: ObservationClientStore })
if (snapshot.sessions.length === 0) return null;
if (collapsed) {
const state = focus !== undefined ? statusOf(focus) : "idle";
const state: ChromeState = snapshot.reconnecting
? "reconnecting"
: focus !== undefined
? statusOf(focus)
: "idle";
return (
<button
type="button"
className={cn(css.capsule, "bsk-obs")}
data-state={state}
data-testid="obs-capsule"
aria-label="Expand browser observation overlay"
onClick={() => setCollapsed(false)}
>
<StatusDot state={state} />
<span className={css["capsule-text"]}>
{snapshot.sessions.length} session{snapshot.sessions.length === 1 ? "" : "s"}
{focus !== undefined && focus.action !== "idle"
? ` · ${focus.action} · ${formatElapsed(focus.since, now)}`
: ""}
{snapshot.reconnecting
? " · reconnecting…"
: focus !== undefined && focus.action !== "idle"
? ` · ${focus.action} · ${formatElapsed(focus.since, now)}`
: ""}
</span>
</button>
);
@@ -100,6 +100,7 @@ export function ObservationSidebarTab({
focus={focus}
sessions={snapshot.sessions}
available={snapshot.available}
reconnecting={snapshot.reconnecting}
pinnedId={pinnedId}
onTogglePin={onTogglePin}
now={now}
@@ -12,6 +12,10 @@ import { ObservationPresentation } from "./observation-presentation";
export interface EventSourceLike {
onmessage: ((event: { data: string }) => void) | null;
/** Present on a real EventSource; test doubles may omit it. */
onerror?: ((event: unknown) => void) | null;
/** 0 CONNECTING, 1 OPEN, 2 CLOSED; undefined on test doubles. */
readyState?: number;
close(): void;
}
@@ -46,12 +50,25 @@ export interface OverlaySnapshot {
readonly displayFrames: Readonly<Record<string, ThumbnailState>>;
/** False when the host reports the browser/daemon as unreachable. */
readonly available: boolean;
/**
* True from a stream error until the stream delivers a frame again;
* meanwhile sessions and availability are the last ones received.
*/
readonly reconnecting: boolean;
}
const EVENTS_URL = "/bsk-observation/events";
const INTERRUPT_URL = "/bsk-observation/interrupt";
const STOP_URL = "/bsk-observation/stop";
const THUMBNAIL_RETRY_DELAYS_MS = [1000, 3000];
/**
* Backoff for re-creating the live event stream after a *fatal* failure.
* A non-200 response (for example the observation route answering 404 while
* the plugin is still starting, or right after a plugin reload) is terminal
* for an EventSource: the browser never retries it. Without this the view
* stays empty until the page is reloaded.
*/
const EVENT_STREAM_RETRY_DELAYS_MS = [1000, 3000, 10000, 30000];
function revoke(url: string | undefined): void {
if (url !== undefined && typeof URL.revokeObjectURL === "function") {
@@ -71,14 +88,19 @@ export class ObservationClientStore {
private readonly retryCounts = new Map<string, number>();
private listeners = new Set<() => void>();
private events: EventSourceLike | undefined;
/** Pending re-creation of the event stream after a fatal failure. */
private reconnectTimer: ReturnType<typeof setTimeout> | undefined;
private reconnectAttempts = 0;
private snapshot: OverlaySnapshot = {
sessions: [],
subscribed: false,
thumbnails: {},
displayFrames: {},
available: true,
reconnecting: false,
};
private available = true;
private reconnecting = false;
private started = false;
/** Refcount of mounted consumers (overlay card, sidebar tab, sidebar fiber). */
private consumers = 0;
@@ -100,6 +122,7 @@ export class ObservationClientStore {
thumbnails: Object.fromEntries(this.thumbs),
displayFrames: this.buildDisplayFrames(),
available: this.available,
reconnecting: this.reconnecting,
};
for (const listener of [...this.listeners]) listener();
}
@@ -136,6 +159,8 @@ export class ObservationClientStore {
}
private connectEvents(): void {
clearTimeout(this.reconnectTimer);
this.reconnectTimer = undefined;
const previous = this.events;
this.events = undefined;
previous?.close();
@@ -151,8 +176,36 @@ export class ObservationClientStore {
} catch {
return;
}
// Healthy traffic resets the backoff and ends the outage.
this.reconnectAttempts = 0;
this.reconnecting = false;
this.apply(event);
};
// A non-200 response is fatal for an EventSource: the browser fails the
// connection permanently and never retries it. Transient drops keep
// readyState 0/1 and are still retried by the EventSource itself, so only
// a CLOSED (2) readyState asks us to rebuild the stream.
events.onerror = () => {
if (!this.started || this.events !== events) return;
if (events.readyState === undefined || events.readyState === 2) this.scheduleReconnect();
if (this.reconnecting) return;
this.reconnecting = true;
this.publish();
};
}
/** Re-create the stream after a fatal failure, with a bounded backoff. */
private scheduleReconnect(): void {
if (this.reconnectTimer !== undefined) return;
const index = Math.min(this.reconnectAttempts, EVENT_STREAM_RETRY_DELAYS_MS.length - 1);
const delay = EVENT_STREAM_RETRY_DELAYS_MS[index];
if (delay === undefined) return;
this.reconnectAttempts += 1;
this.reconnectTimer = setTimeout(() => {
this.reconnectTimer = undefined;
if (!this.started) return;
this.connectEvents();
}, delay);
}
/** Every initial connection and reconnect starts with the server's snapshot. */
@@ -168,6 +221,10 @@ export class ObservationClientStore {
this.events = undefined;
this.started = false;
previous?.close();
clearTimeout(this.reconnectTimer);
this.reconnectTimer = undefined;
this.reconnectAttempts = 0;
this.reconnecting = false;
this.clearThumbnails();
this.sessions.clear();
this.available = true;
@@ -197,6 +197,68 @@ describe("ObservationOverlay", () => {
expect(screen.getByTestId("obs-card")).toBeTruthy();
});
it("never commits the card without sessions, even while the feed is down", async () => {
const h = makeHarness([]);
const inserted: Node[] = [];
const collect = (records: MutationRecord[]) => {
for (const record of records) inserted.push(...record.addedNodes);
};
const observer = new MutationObserver(collect);
observer.observe(document.body, { childList: true, subtree: true });
render(<ObservationOverlay store={h.store} />);
await waitFor(() => expect(h.es).not.toThrow());
act(() => {
h.es().readyState = 2;
h.es().onerror?.({});
});
collect(observer.takeRecords());
observer.disconnect();
expect(h.store.getSnapshot().reconnecting).toBe(true);
const cardCommitted = inserted.some(
(node) => node instanceof Element && node.closest("[data-obs-card]") !== null,
);
expect(cardCommitted).toBe(false);
});
it("marks a live session as reconnecting until the feed recovers", async () => {
const h = makeHarness([BUSY]);
render(<ObservationOverlay store={h.store} />);
await screen.findByText(/s1 · clicking/);
act(() => {
h.es().readyState = 0;
h.es().onerror?.({});
});
const header = screen.getByTestId("obs-header");
expect(header.textContent).toContain("reconnecting…");
expect(header.querySelector("[data-state]")?.getAttribute("data-state")).toBe("reconnecting");
expect(screen.queryByText(/s1 · clicking/)).toBeNull();
act(() => h.emitRaw({ type: "snapshot", sessions: [BUSY], available: true }));
expect(header.textContent).not.toContain("reconnecting");
expect(header.querySelector("[data-state]")?.getAttribute("data-state")).toBe("active");
expect(screen.getByText(/s1 · clicking/)).toBeTruthy();
});
it("collapses to a capsule that drops stale action timing while reconnecting", async () => {
const h = makeHarness([BUSY]);
render(<ObservationOverlay store={h.store} />);
await screen.findByText(/s1 · clicking/);
fireEvent.click(screen.getByRole("button", { name: "Collapse" }));
const capsule = await screen.findByTestId("obs-capsule");
expect(capsule.textContent).toContain("clicking");
expect(capsule.getAttribute("data-state")).toBe("active");
act(() => {
h.es().readyState = 0;
h.es().onerror?.({});
});
expect(capsule.textContent).toContain("reconnecting…");
expect(capsule.textContent).not.toContain("clicking");
expect(capsule.getAttribute("data-state")).toBe("reconnecting");
expect(capsule.querySelector("[data-state]")?.getAttribute("data-state")).toBe("reconnecting");
act(() => h.emitRaw({ type: "snapshot", sessions: [BUSY], available: true }));
expect(capsule.textContent).toContain("clicking");
expect(capsule.getAttribute("data-state")).toBe("active");
});
it("shows the status row and a placeholder without a thumbnail", async () => {
const h = makeHarness([BUSY]);
render(<ObservationOverlay store={h.store} />);
@@ -457,3 +457,130 @@ describe("thumbnail viewer leases", () => {
expect(esInstances).toHaveLength(2);
});
});
describe("event stream recovery", () => {
function streamAt(sources: EventSourceLike[], index: number): EventSourceLike {
const source = sources[index];
if (source === undefined) throw new Error(`no event stream at index ${index}`);
return source;
}
it("recreates a fatally failed stream after the backoff", async () => {
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
const { store, esInstances, eventUrls } = harness([OBS_IDLE]);
store.start();
await Promise.resolve();
expect(store.getSnapshot().reconnecting).toBe(false);
// A 404 from the observation route is fatal: CLOSED, never retried by the browser.
const failed = streamAt(esInstances, 0);
failed.readyState = 2;
failed.onerror?.({});
expect(esInstances).toHaveLength(1);
expect(store.getSnapshot()).toMatchObject({ reconnecting: true, sessions: [OBS_IDLE] });
await vi.advanceTimersByTimeAsync(1000);
expect(esInstances).toHaveLength(2);
expect(eventUrls).toEqual([
"/bsk-observation/events?thumbnails=0",
"/bsk-observation/events?thumbnails=0",
]);
expect(failed.close).toHaveBeenCalledOnce();
expect(store.getSnapshot()).toMatchObject({ reconnecting: false, sessions: [OBS_IDLE] });
});
it("stays reconnecting until a rebuilt stream delivers a frame", async () => {
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
const { store, esInstances } = harness([OBS_IDLE]);
store.start();
await Promise.resolve();
const failed = streamAt(esInstances, 0);
failed.readyState = 2;
failed.onerror?.({});
await vi.advanceTimersByTimeAsync(999);
expect(esInstances).toHaveLength(1);
// The rebuilt stream exists before its snapshot frame arrives.
vi.advanceTimersByTime(1);
expect(esInstances).toHaveLength(2);
expect(store.getSnapshot().reconnecting).toBe(true);
await Promise.resolve();
expect(store.getSnapshot().reconnecting).toBe(false);
});
it("leaves a transient drop to the EventSource's own retry", async () => {
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
const { store, esInstances } = harness([OBS_IDLE]);
store.start();
await Promise.resolve();
const dropped = streamAt(esInstances, 0);
dropped.readyState = 0;
dropped.onerror?.({});
expect(store.getSnapshot().reconnecting).toBe(true);
await vi.advanceTimersByTimeAsync(60_000);
expect(esInstances).toHaveLength(1);
// The browser's own reconnect replays the snapshot on the same stream.
emit(dropped, { type: "snapshot", sessions: [OBS_BUSY], available: true });
expect(store.getSnapshot()).toMatchObject({ reconnecting: false, sessions: [OBS_BUSY] });
});
it("lets a thumbnail switch take over a pending reconnect", async () => {
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
const { store, esInstances, eventUrls } = harness([OBS_IDLE]);
store.start();
await Promise.resolve();
const failed = streamAt(esInstances, 0);
failed.readyState = 2;
failed.onerror?.({});
store.watchThumbnails();
await Promise.resolve();
expect(store.getSnapshot().reconnecting).toBe(false);
// The stale backoff timer must not tear down the healthy replacement.
await vi.advanceTimersByTimeAsync(60_000);
expect(eventUrls).toEqual([
"/bsk-observation/events?thumbnails=0",
"/bsk-observation/events?thumbnails=1",
]);
expect(streamAt(esInstances, 1).close).not.toHaveBeenCalled();
});
it("resets the backoff after healthy traffic", async () => {
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
const { store, esInstances } = harness([OBS_IDLE]);
store.start();
await Promise.resolve();
const first = streamAt(esInstances, 0);
first.readyState = 2;
first.onerror?.({});
await vi.advanceTimersByTimeAsync(1000);
expect(esInstances).toHaveLength(2);
// Any frame proves the new stream is healthy, so the next failure waits 1s again.
const second = streamAt(esInstances, 1);
emit(second, { type: "upsert", session: OBS_BUSY });
second.readyState = 2;
second.onerror?.({});
await vi.advanceTimersByTimeAsync(1000);
expect(esInstances).toHaveLength(3);
});
it("cancels a pending reconnect on stop", async () => {
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
const { store, esInstances } = harness([OBS_IDLE]);
store.start();
await Promise.resolve();
const failed = streamAt(esInstances, 0);
failed.readyState = 2;
failed.onerror?.({});
store.stop();
expect(store.getSnapshot().reconnecting).toBe(false);
await vi.advanceTimersByTimeAsync(60_000);
expect(esInstances).toHaveLength(1);
});
});