diff --git a/packages/dsh-plugin-browserskill/src/client/ObservationOverlay.tsx b/packages/dsh-plugin-browserskill/src/client/ObservationOverlay.tsx
index c1f30ec3..e75c07fd 100644
--- a/packages/dsh-plugin-browserskill/src/client/ObservationOverlay.tsx
+++ b/packages/dsh-plugin-browserskill/src/client/ObservationOverlay.tsx
@@ -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 (
-
+
{statusText}
{onUseFloating !== undefined ? (
setCollapsed(false)}
>
{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)}`
+ : ""}
);
diff --git a/packages/dsh-plugin-browserskill/src/client/observation-sidebar.tsx b/packages/dsh-plugin-browserskill/src/client/observation-sidebar.tsx
index 6894c809..89e7d028 100644
--- a/packages/dsh-plugin-browserskill/src/client/observation-sidebar.tsx
+++ b/packages/dsh-plugin-browserskill/src/client/observation-sidebar.tsx
@@ -100,6 +100,7 @@ export function ObservationSidebarTab({
focus={focus}
sessions={snapshot.sessions}
available={snapshot.available}
+ reconnecting={snapshot.reconnecting}
pinnedId={pinnedId}
onTogglePin={onTogglePin}
now={now}
diff --git a/packages/dsh-plugin-browserskill/src/client/observation-store.ts b/packages/dsh-plugin-browserskill/src/client/observation-store.ts
index 4bbb7aa4..0c0e86be 100644
--- a/packages/dsh-plugin-browserskill/src/client/observation-store.ts
+++ b/packages/dsh-plugin-browserskill/src/client/observation-store.ts
@@ -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>;
/** 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();
private listeners = new Set<() => void>();
private events: EventSourceLike | undefined;
+ /** Pending re-creation of the event stream after a fatal failure. */
+ private reconnectTimer: ReturnType | 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;
diff --git a/packages/dsh-plugin-browserskill/tests/client/observation-overlay.test.tsx b/packages/dsh-plugin-browserskill/tests/client/observation-overlay.test.tsx
index e2573971..0a1a8399 100644
--- a/packages/dsh-plugin-browserskill/tests/client/observation-overlay.test.tsx
+++ b/packages/dsh-plugin-browserskill/tests/client/observation-overlay.test.tsx
@@ -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();
+ 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();
+ 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();
+ 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();
diff --git a/packages/dsh-plugin-browserskill/tests/client/observation-store.test.ts b/packages/dsh-plugin-browserskill/tests/client/observation-store.test.ts
index aab72298..8fd225ba 100644
--- a/packages/dsh-plugin-browserskill/tests/client/observation-store.test.ts
+++ b/packages/dsh-plugin-browserskill/tests/client/observation-store.test.ts
@@ -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);
+ });
+});