Files
fusion/packages/dashboard/app/hooks/useLiveTranscript.ts
gsxdsm f26cbedf4f fix(dashboard): close the code-review findings on the mobile tab-discard work
An 11-reviewer pass over f157bf7460..f5163d8351 found defects in the mobile
tab-discard change set itself. This fixes them.

Silent data loss (the recurring defect class):
- AgentDetailView reconnect refetched limit:100 and replaced wholesale, so 380
  displayed lines vanished with no "Load older" and no indicator; it now
  reconciles through the shared logStreamReconcile helper.
- useActivityLog.loadMore past the cap discarded the page it had just fetched
  while advancing the cursor and leaving hasMore true, so the feed silently
  stopped paginating behind a live-looking button.
- useAgentLogs: loadMore and resyncFromServer had no mutual exclusion, a
  no-overlap resync discarded explicitly paged-back history, a resync outliving
  the reconnect delay left an unmarked gap, and the live-tail trim could evict
  the gap marker itself.
- useLiveTranscript's resync overwrote live entries that raced the refetch.

The premise itself was not fully delivered:
- useProjects, useNodes, and useMeshState never called clearInterval, so they
  polled the whole time the tab was hidden. useProjects is mounted for the
  entire session, so the page never went idle -- the primary mechanism this
  work depends on. All three now use the shared visibility gate.
- sse-bus fired onReconnect twice per reconnect cycle and fanned out ~28
  subscribers in one tick, against a ~6-connection-per-origin cap on a waking
  radio. The successful open is now the single authority, and the fan-out uses
  the same exported stagger primitive as the polling path rather than a second
  copy of the slot formula.
- A channel first subscribed during the hidden window opened a live EventSource
  and keepalive; suspension is now a module-level condition openChannel
  consults, and a channel opened inside the grace window re-arms it.

Credentials and correctness:
- The service worker persisted every GET /api/* to durable Cache Storage,
  including /api/settings with daemonToken, githubAuthToken, gitlabAuthToken
  and ntfyAccessToken in plaintext, with no exclusion and no purge path --
  "Clear all cached data" only walked localStorage. Now gated, bounded, and
  genuinely purgeable.
- useTasks cleared its own snapshot when the mount revalidation failed on a
  waking radio, so the board blanked and the next restore was empty too.
  Suspension-class failures no longer destroy the cache.
- A single-row SSE update reset lastFetchTimeMs to now while an hours-old
  hydrated snapshot was on screen, re-marking every in-progress card stuck.
- ListView's "Select all visible tasks" acted on the full filtered set while
  only 50 rows rendered, so a bulk delete reached rows the operator could not
  see. Column's search window reset keyed on a boolean, so refining a query
  kept the expanded window.

Tests that could not fail:
- App.test.tsx mocked TerminalModal as isOpen ? <div/> : null, making the
  unmount-on-close invariant unobservable; MockEventSource kept its listeners
  after close(), so cases passed with their onReconnect handlers deleted.
- The SSE resync ratchet scanned only hooks/, exempting ~13 component call
  sites -- the exact regression it exists to prevent.
- MissionControlPanel's bespoke poll and the xterm scrollback constants and
  WebGL disposal had no coverage at all.

Verified: tsc -p tsconfig.app.json clean, pnpm lint clean, pnpm
check:changesets clean, 877 tests passing across 36 scoped files.
Known unrelated red: MailboxView.test.tsx's FN-8407 CSS guard fails at HEAD
too -- this diff adds no @media rule and no .mailbox-view--mobile selector,
the only two things that assertion inspects. Left alone deliberately.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-26 11:17:52 -07:00

228 lines
9.2 KiB
TypeScript

import { useState, useEffect, useRef } from "react";
import type { AgentLogEntry } from "@fusion/core";
import { fetchAgentLogsWithMeta } from "../api";
import { subscribeSse } from "../sse-bus";
import { appendWithoutDuplicates } from "./logStreamReconcile";
import { createResyncRetryRunner } from "./resyncRetry";
// Render shows only the first 20 entries; keep a small buffer above that for
// hot-reconnect/scrollback but never let the array grow unbounded — long-lived
// active agents would otherwise leak hundreds of MB into React state.
const MAX_TRANSCRIPT_ENTRIES = 200;
/**
* Log entry from an agent's execution stream.
*
* Note: SSE payloads from `/api/tasks/:id/logs/stream` contain `text` field
* (matching `AgentLogEntry` from `@fusion/core`). This interface normalizes
* to `text` for rendering. Legacy payloads with `content` are also supported
* for backward compatibility.
*/
export interface TranscriptEntry {
type: string;
/** Canonical text content — matches `AgentLogEntry.text` */
text: string;
timestamp?: string;
/** Legacy field — normalized to `text` if present */
content?: string;
}
/** A live/persisted log row in the shared oldest-first reconcile shape, keeping the legacy field. */
type ReconcilableEntry = AgentLogEntry & { content?: string };
/** Oldest-first reconcile shape -> rendered newest-first, bounded transcript. */
function toTranscriptEntries(entries: ReconcilableEntry[]): TranscriptEntry[] {
return entries
.slice(-MAX_TRANSCRIPT_ENTRIES)
.reverse()
.map((raw) => ({
type: raw.type ?? "text",
text: raw.text ?? "",
timestamp: raw.timestamp || undefined,
...(raw.content !== undefined ? { content: raw.content } : {}),
}));
}
/**
* Hook that manages live transcript streaming for a task.
*
* Features project-context isolation to prevent cross-project transcript bleed:
* - Tracks project context version to detect stale events after project switches
* - Resets entries and connection state immediately on context change
* - Rejects stale SSE events from previous EventSource instances
*
* When `taskId` changes, a new SSE connection is opened for the new task.
* When `projectId` changes, all state is reset and a new connection is opened.
*
* **Reconnect semantics**: the stream replays nothing on open, so every reconnect refetches the
* persisted tail and merges it with events that raced the fetch (see resyncTranscript).
*/
export function useLiveTranscript(taskId: string | undefined, projectId?: string) {
const [entries, setEntries] = useState<TranscriptEntry[]>([]);
const [isConnected, setIsConnected] = useState(false);
// Refs for state that needs to survive re-renders
const unsubscribeRef = useRef<(() => void) | null>(null);
/*
FNXC:TaskTranscript 2026-07-26-18:40:
A reconnect refetch is in flight; live `agent:log` events are parked here instead of being
prepended. Ported from useAgentLogs: the refetched page is a snapshot from BEFORE those events, so
writing it over a buffer that moved meanwhile deleted the newest lines with no replay path (the
stream never resends). Parked entries are merged after the page, deduped by content, so an entry
present in both is rendered once.
*/
const resyncInFlightRef = useRef(false);
const pendingLiveRef = useRef<ReconcilableEntry[]>([]);
// Track the project context version to detect stale events after project switches.
// Incremented whenever projectId changes, invalidating any in-flight SSE handlers.
const projectContextVersionRef = useRef(0);
// Track previous values to detect context changes
const previousTaskIdRef = useRef<string | undefined>(taskId);
const previousProjectIdRef = useRef<string | undefined>(projectId);
// Detect context changes and reset state immediately
const contextChanged =
previousTaskIdRef.current !== taskId ||
previousProjectIdRef.current !== projectId;
if (contextChanged) {
previousTaskIdRef.current = taskId;
previousProjectIdRef.current = projectId;
projectContextVersionRef.current++;
// Drop existing subscription
if (unsubscribeRef.current) {
unsubscribeRef.current();
unsubscribeRef.current = null;
}
// Reset state immediately to prevent stale data visibility
resyncInFlightRef.current = false;
pendingLiveRef.current = [];
setEntries([]);
setIsConnected(false);
}
useEffect(() => {
if (!taskId) {
setEntries([]);
setIsConnected(false);
return;
}
// Capture context version at effect start - stale events will be rejected
const contextVersionAtStart = projectContextVersionRef.current;
// Build stream URL with optional projectId for multi-project support
let url = `/api/tasks/${encodeURIComponent(taskId)}/logs/stream`;
if (projectId) {
url += `?projectId=${encodeURIComponent(projectId)}`;
}
/*
FNXC:TaskTranscript 2026-07-26-14:50:
Missed-event recovery. `/api/tasks/:id/logs/stream` replays nothing on connect — it only forwards
entries emitted while the socket is live — so every SSE gap (error reconnect, or the mobile
hidden-tab suspend) punched a permanent hole in the transcript, silently mixing entries from
before and after the gap with no marker. On reopen, converge on the persisted tail from
GET /tasks/:id/logs, which is the same data the stream would have delivered, merged with
anything that streamed in during the fetch. The API returns oldest-first; this hook renders
newest-first.
*/
const isStale = () => projectContextVersionRef.current !== contextVersionAtStart;
/*
FNXC:TaskTranscript 2026-07-26-18:44:
CORRECTION to the previous resync, which called setEntries(server page) unconditionally: entries
that arrived while the fetch was in flight were older than the state write and were DELETED, and
this stream replays nothing, so they were unrecoverable. The fetch window is now a parking window
(pendingLiveRef) and the parked entries are merged onto the page with the shared
`appendWithoutDuplicates`, the same de-dup useAgentLogs uses, so a line that is in both the page
and the live stream renders exactly once.
The old catch comment ("the next reopen retries") was ALSO wrong — a healthy connection may not
reopen again for hours. Failures go through the shared bounded ladder, and parked entries are
flushed rather than dropped when the page cannot be fetched.
*/
const resyncTranscript = async () => {
if (resyncInFlightRef.current || isStale()) return;
resyncInFlightRef.current = true;
pendingLiveRef.current = [];
let reconciled = false;
try {
const result = await fetchAgentLogsWithMeta(taskId, projectId, { limit: MAX_TRANSCRIPT_ENTRIES });
if (isStale()) return;
const pending = pendingLiveRef.current;
pendingLiveRef.current = [];
reconciled = true;
setEntries(toTranscriptEntries(appendWithoutDuplicates(result.entries, pending) as ReconcilableEntry[]));
} finally {
resyncInFlightRef.current = false;
if (!reconciled && !isStale()) {
const pending = pendingLiveRef.current;
pendingLiveRef.current = [];
if (pending.length > 0) {
// Failed/stale refetch: a degraded live tail still beats silently dropping lines.
setEntries((prev) => [...toTranscriptEntries(pending), ...prev].slice(0, MAX_TRANSCRIPT_ENTRIES));
}
}
}
};
const transcriptResync = createResyncRetryRunner({ run: resyncTranscript });
const unsubscribe = subscribeSse(url, {
onReconnect: () => transcriptResync.trigger(),
events: {
"agent:log": (event) => {
if (isStale()) return;
try {
const raw = JSON.parse(event.data) as Partial<TranscriptEntry>;
// Normalize: canonical `text` field, with legacy `content` fallback
const entry: ReconcilableEntry = {
type: (raw.type ?? "text") as AgentLogEntry["type"],
text: raw.text ?? raw.content ?? "",
timestamp: raw.timestamp ?? "",
taskId,
...(raw.content !== undefined ? { content: raw.content } : {}),
};
if (resyncInFlightRef.current) {
pendingLiveRef.current.push(entry);
return;
}
setEntries(prev => {
const next = [...toTranscriptEntries([entry]), ...prev];
return next.length > MAX_TRANSCRIPT_ENTRIES
? next.slice(0, MAX_TRANSCRIPT_ENTRIES)
: next;
});
} catch { /* skip malformed events */ }
},
},
onOpen: () => {
if (projectContextVersionRef.current === contextVersionAtStart) {
setIsConnected(true);
}
},
onError: () => {
if (projectContextVersionRef.current === contextVersionAtStart) {
setIsConnected(false);
}
},
});
unsubscribeRef.current = unsubscribe;
return () => {
transcriptResync.dispose();
unsubscribe();
unsubscribeRef.current = null;
if (projectContextVersionRef.current === contextVersionAtStart) {
setIsConnected(false);
}
};
}, [taskId, projectId]);
return { entries, isConnected };
}