返回 DeepSeek-Reasonix
useRemoteSession.ts
根目录 / desktop / frontend / src / lib / useRemoteSession.ts
1 import { useCallback, useEffect, useMemo, useRef, useState, useSyncExternalStore } from "react";
2 import { useRuntimeSession } from "./useRuntimeState";
3 import { createLegacyRemotePolicyNoticeTracker } from "./legacyRemotePolicyNotice";
4 import { app, onRemoteTabEvent, onRemoteTabState } from "./bridge";
5 import { queuedFollowupOutcome } from "./pendingFollowup";
6 import { onRemoteTabUpdated } from "./remoteTabEvents";
7 import { hydrateRemoteTelemetry, loadRemoteStatusSnapshot } from "./remoteTelemetry";
8 import { remoteStatusTakenOver, remoteStatusToAction } from "./remoteStatus";
9 import { isRemoteTakeoverError } from "./remoteErrors";
10 import { useRemoteForkTurn } from "./remoteForkTurn";
11 import { useT } from "./i18n";
12 import type { CancelOutcome } from "./inboxCancel";
13 import type { HistoryMessage } from "./types";
14 import { createTurnSubmissionId, initialState, reducer, type ControllerLiveStore, type HistoryLoadOutcome, type HistoryLoadTrigger, type State } from "./useController";
15 import { isUnknownSubmissionError } from "./localSubmissionState";
16 import { TranscriptSessionFollower } from "./transcriptSessionFollower";
17 import { getTranscriptStore } from "./transcriptStore";
18 import type { NavigateToTurn } from "./historyTurnNavigation";
19 import { historyReplaceAction } from "./sessionTranscriptMode";
20 import { isAuthoritativeRemoteStatus, remoteCheckpoints, remoteComposerState, remoteGoalRuntime, remoteGoalView } from "./remoteStatus";
21 import type { CollaborationMode, CommandInfo, EffortInfo, GoalLifecycleView, GoalRuntime, GoalStatus, QualityFloor, RemoteTabStateValue, TabMeta, ToolApprovalMode, WireEvent } from "./types";
22 import type { RemoteAskAnswer } from "./remoteTypes";
23 import type { ForkTargetView } from "./forkTargets";
24
25
26 // The remote session reuses the local transcript pipeline end to end: serve
27 // frames share the agent event wire form, so they run through the same
28 // reducer that drives local tabs, and /history hydrates through the same
29 // history action. The surface and composer therefore consume exactly the
30 // shapes the local UI consumes.
31
32 // RemoteSessionApi is the surface-facing contract of useRemoteSession.
33 // Last session path mounted per remote tab. Switching sessions inside one tab
34 // keeps the tab-keyed reducer state, so a resident paint must know the mounted
35 // content belongs to a different session before replacing it.
36 const lastPaintedSessionPathByTab = new Map<string, string>();
37
38 export interface RemoteSessionApi {
39 state: RemoteTabStateValue;
40 error: string;
41 transcript: State;
42 liveStore: ControllerLiveStore;
43 hydrated: boolean;
44 /** Content is painted from the resident cache; the authoritative hydrate is still running. */
45 revalidating: boolean;
46 syncMode?: "v2";
47 loadOlderHistory?: (targetTurn?: number, trigger?: HistoryLoadTrigger) => Promise<HistoryLoadOutcome>;
48 loadNewerHistory?: (latest?: boolean, current?: () => boolean) => Promise<HistoryLoadOutcome>;
49 navigateToTurn?: NavigateToTurn;
50 running: boolean;
51 /** The serve's label for the active model, for the composer capsule. */
52 modelLabel: string;
53 commands: CommandInfo[];
54 composerProfile?: {
55 collaborationMode: CollaborationMode;
56 toolApprovalMode: ToolApprovalMode;
57 goal: string;
58 goalStatus?: GoalStatus;
59 qualityFloor: QualityFloor;
60 };
61 goalRuntime?: GoalRuntime;
62 goalView?: GoalLifecycleView;
63 effort?: EffortInfo;
64 /** Changes whenever the tab adopts a new/reconnected Serve session snapshot. */
65 surfaceGeneration: number;
66 promptError: string;
67 submit: (text: string, displayText?: string, choice?: import("./modelApplication").ModelApplicationChoice) => Promise<void>;
68 runManagementCommand: (text: string, rehydrate?: boolean) => Promise<void>;
69 compact: (instructions: string) => Promise<void>;
70 cancelTurn: () => Promise<void>;
71 approve: (callId: string, decision: string) => Promise<void>;
72 resolvePlanDecision: (callId: string, action: "start_execution" | "revise_plan" | "exit_plan", feedback?: string) => Promise<void>;
73 answer: (callId: string, answers: RemoteAskAnswer[]) => Promise<void>;
74 clearExtensionForm: (pluginId: string, surfaceId: string, formInstanceId?: string) => void;
75 rewind: (turn: number, scope: string) => Promise<void>;
76 /** Creates the child session for one turn; returns its id, or undefined with the reason in promptError. */
77 forkTurn: (target: ForkTargetView) => Promise<{ sessionId: string; operationId: string } | undefined>;
78 acknowledgeFork: (operationId: string) => Promise<void>;
79 setModel: (ref: string) => Promise<void>;
80 setEffort: (level: string) => Promise<void>;
81 setQualityFloor: (floor: QualityFloor) => Promise<void>;
82 pauseGoal: () => Promise<void>;
83 resumeGoal: () => Promise<void>;
84 editGoal: (objective: string, maxGoalRounds: number | null) => Promise<void>;
85 steer: (input: string) => Promise<void>;
86 cancelJob: (jobId: string) => Promise<boolean>;
87 drainApprovals: (ids: string[]) => void;
88 retryHydration: () => Promise<void>;
89 }
90
91 export function useRemoteComposer(
92 session: RemoteSessionApi,
93 showToast: (message: string, level: "warn" | "error") => void,
94 ) {
95 const onSend = useCallback(async (displayText: string, submitText = displayText, choice?: import("./modelApplication").ModelApplicationChoice) => {
96 const text = (submitText || displayText).trim();
97 if (!text) return;
98 await session.submit(text, displayText, choice);
99 }, [session, showToast]);
100 const onCancel = useCallback(async (_queuedItemIDs?: string[]): Promise<CancelOutcome> => {
101 void session.cancelTurn().catch((error) => {
102 showToast(error instanceof Error ? error.message : String(error), "error");
103 });
104 return { discardedItemIds: [] };
105 }, [session, showToast]);
106 return { onSend, onCancel };
107 }
108
109 export function useActiveRemoteSession(
110 activeTab: TabMeta | undefined,
111 showToast: (message: string, level: "warn" | "error") => void,
112 ) {
113 const t = useT();
114 const active = Boolean(activeTab?.remote);
115 const session = useRemoteSession(active && activeTab ? activeTab.id : undefined, activeTab?.remoteState, activeTab?.sessionPath);
116 const composer = useRemoteComposer(session, showToast);
117 useEffect(() => {
118 if (!activeTab?.remote || !activeTab.id) return;
119 const legacyQuality = session.transcript.items.some(item => item.kind === "notice" && (item.code === "final_readiness" || item.variant === "delivery"));
120 const key = `${activeTab.id}\u0000${activeTab.sessionPath ?? ""}`;
121 const notice = legacyRemotePolicyNotice(key, session.composerProfile?.qualityFloor, session.goalRuntime?.stopCause, legacyQuality);
122 if (notice) showToast(t(notice), "warn");
123 }, [activeTab?.remote, activeTab?.id, activeTab?.sessionPath, session.composerProfile?.qualityFloor, session.transcript.items, session.goalRuntime?.stopCause, showToast, t]);
124 return { active, session, ready: active && session.state === "ready" && session.hydrated && Boolean(session.composerProfile), ...composer };
125 }
126
127 const legacyRemotePolicyNotice = createLegacyRemotePolicyNoticeTracker();
128
129 export function useRemoteSession(tabId: string | undefined, initial?: RemoteTabStateValue, sessionPath?: string): RemoteSessionApi {
130 const runtimeState = useRuntimeSession(tabId, sessionPath);
131 const [state, setState] = useState<RemoteTabStateValue>(initial === "disconnected" ? "connecting" : (initial ?? "connecting"));
132 const [error, setError] = useState("");
133 const transcript = useSyncExternalStore(
134 useCallback(listener => tabId ? getTranscriptStore().subscribeState(tabId, listener) : () => {}, [tabId]),
135 useCallback(() => (tabId ? getTranscriptStore().states.get(tabId) : undefined) ?? initialState, [tabId]),
136 );
137 const [modelLabel, setModelLabel] = useState("");
138 const [commands, setCommands] = useState<CommandInfo[]>([]);
139 const [composerProfile, setComposerProfile] = useState<RemoteSessionApi["composerProfile"]>();
140 const [goalRuntime, setGoalRuntime] = useState<GoalRuntime>();
141 const [goalView, setGoalView] = useState<GoalLifecycleView>();
142 const [effort, setEffortInfo] = useState<EffortInfo>();
143 const [surfaceGeneration, setSurfaceGeneration] = useState(0);
144 const [promptError, setPromptError] = useState("");
145 const [hydrated, setHydrated] = useState(false);
146 // A resident-cache paint shows the session instantly while the authoritative
147 // hydrate reconciles in the background.
148 const [revalidating, setRevalidating] = useState(false);
149 const cachePaintedRef = useRef(false);
150 const olderRef = useRef<((trigger?: HistoryLoadTrigger) => Promise<HistoryLoadOutcome>) | undefined>(undefined);
151 const navigateRef = useRef<NavigateToTurn | undefined>(undefined);
152 const newerRef = useRef<((latest?: boolean, current?: () => boolean) => Promise<HistoryLoadOutcome>) | undefined>(undefined);
153 const transcriptRef = useRef(transcript);
154 const submitBindingRef = useRef<object>({});
155 const setTranscript = useCallback((update: State | ((state: State) => State)) => {
156 const owned = (tabId ? getTranscriptStore().states.get(tabId) : undefined) ?? initialState;
157 const next = typeof update === "function" ? update(owned) : update;
158 transcriptRef.current = next;
159 if (tabId) getTranscriptStore().setState(tabId, next);
160 }, [tabId]);
161 const { forkTurn, acknowledgeFork, forkTargetsRefreshRef } = useRemoteForkTurn(app, tabId, sessionPath, setTranscript, setPromptError);
162 const liveListenersRef = useRef(new Set<() => void>());
163 const hydratedRef = useRef(false);
164 const hydratingRef = useRef(false);
165 const bufferedEventsRef = useRef<WireEvent[]>([]);
166 // The pre-activation history prime and the follower are never both allowed
167 // to write the resident store. "retired" is terminal for the mounted
168 // identity: it is set the moment the follower publishes its first cut (or
169 // the legacy fallback installs), after which every prime attempt is a no-op.
170 const primeRef = useRef<"idle" | "loading" | "primed" | "retired">("idle");
171 const hydrateRef = useRef<{ tabId: string; run: (force?: boolean) => Promise<void> } | null>(null);
172 const refreshStatusRef = useRef<{ tabId: string; run: () => Promise<void> } | null>(null);
173 const reconcileHistoryRef = useRef<(() => Promise<void>) | null>(null);
174 const activityRevisionRef = useRef(0);
175 const eventTurnIdRef = useRef<string | undefined>(undefined);
176 const runtimeAtActivityRef = useRef(runtimeState.state);
177 // True while the serve reports a local runtime on the host owns this
178 // session; a spectator surface idles with no other status polling, so this
179 // flag also drives a slow reconcile loop below. The ref mirrors the last
180 // observation so the owner-return transition can be detected at the source.
181 const [spectator, setSpectator] = useState(false);
182 const spectatorRef = useRef(false);
183 const noteOwnership = useCallback((takenOver: boolean) => {
184 const wasSpectator = spectatorRef.current;
185 spectatorRef.current = takenOver;
186 setSpectator(takenOver);
187 // Ownership returning is the only path back to the live transcript: the
188 // legacy fallback installed a protocol-1 view whose submit() refuses to
189 // send, and neither the state channel nor the status poll re-attaches the
190 // follower on its own.
191 if (wasSpectator && !takenOver) void hydrateRef.current?.run().catch(() => undefined);
192 }, []);
193
194 useEffect(() => {
195 for (const listener of liveListenersRef.current) listener();
196 }, [transcript]);
197
198 const liveStore = useMemo<ControllerLiveStore>(() => ({
199 subscribe(requestedTabId, listener) {
200 if (!tabId || requestedTabId !== tabId) return () => undefined;
201 liveListenersRef.current.add(listener);
202 return () => liveListenersRef.current.delete(listener);
203 },
204 getSnapshot(requestedTabId) {
205 return requestedTabId === tabId ? transcriptRef.current.live : undefined;
206 },
207 getModelActiveAt(requestedTabId) {
208 return requestedTabId === tabId ? transcriptRef.current.turnModelActiveAt : undefined;
209 },
210 getRateOutputQuarters(requestedTabId) {
211 return requestedTabId === tabId ? transcriptRef.current.turnRateSample?.outputQuarters : undefined;
212 },
213 }), [tabId]);
214
215 const applyRemoteStatus = useCallback((status: unknown) => {
216 if (!isAuthoritativeRemoteStatus(status)) return;
217 const next = remoteComposerState(status);
218 setModelLabel(next.modelLabel);
219 setComposerProfile(next.composerProfile);
220 setGoalRuntime(remoteGoalRuntime(status));
221 setGoalView(remoteGoalView(status));
222 setEffortInfo(next.effort);
223 noteOwnership(remoteStatusTakenOver(status));
224 }, [noteOwnership]);
225
226 useEffect(() => {
227 if (!tabId) return;
228 transcriptRef.current = getTranscriptStore().states.get(tabId) ?? initialState;
229 submitBindingRef.current = {};
230 // Restored shells arrive as disconnected shells. Activation must kick the
231 // backend revive (SetActiveTab → bootstrap) and never park the UI on a
232 // reconnect placeholder — treat them as connecting until ready/error.
233 const revivedFromShell = initial === "disconnected";
234 const mountedState = revivedFromShell ? "connecting" : (initial ?? "connecting");
235 setState(mountedState);
236 setError("");
237 setPromptError("");
238 // The store owns mounted content across reconnects and tab switches.
239 eventTurnIdRef.current = undefined;
240 setModelLabel("");
241 setCommands([]);
242 setComposerProfile(undefined);
243 setGoalRuntime(undefined);
244 setGoalView(undefined);
245 setEffortInfo(undefined);
246 hydratedRef.current = false;
247 hydratingRef.current = false;
248 bufferedEventsRef.current = [];
249 primeRef.current = "idle";
250 spectatorRef.current = false;
251 setSpectator(false);
252 setHydrated(false);
253 cachePaintedRef.current = false;
254 setRevalidating(false);
255 let cancelled = false;
256 let generation = 0;
257 let follower: TranscriptSessionFollower | undefined;
258 const dispatch = (action: import("./useController").Action) => {
259 if (!cancelled) setTranscript(current => reducer(current, action));
260 };
261 // History-first: the persisted canonical window is readable before the
262 // serve activates the runtime, so publish a one-shot durable baseline as
263 // soon as the tab's identity resolves. This is deliberately not a second
264 // live transcript owner — once the runtime is ready, hydrate()'s follower
265 // installs the authoritative protocol-v2 cut and owns everything after it.
266 // Mirrors the local primeReadableHistoryForTab contract.
267 const transcriptHasContent = () => {
268 const mounted = transcriptRef.current;
269 return mounted.items.length > 0 || Boolean(mounted.live?.text || mounted.live?.reasoning);
270 };
271 // Local sessions paint a resident session synchronously before any I/O
272 // (useController's peek path). The remote surface had no such path: every
273 // switch paid the network round trips even for a session whose transcript
274 // was already resident. Paint it, then let hydrate() reconcile.
275 const paintResidentCache = () => {
276 if (!tabId || !sessionPath) return false;
277 const previousPath = lastPaintedSessionPathByTab.get(tabId);
278 lastPaintedSessionPathByTab.set(tabId, sessionPath);
279 // Same session (reconnect) or a restored mount already showing content:
280 // the mounted state is this session's and may be fresher than the cut.
281 if (previousPath === sessionPath || (previousPath === undefined && transcriptHasContent())) return false;
282 const resident = getTranscriptStore().peek(tabId, sessionPath);
283 if (!resident) return false;
284 cachePaintedRef.current = true;
285 setTranscript(current => reducer(current, historyReplaceAction(resident)));
286 hydratedRef.current = true;
287 setHydrated(true);
288 setRevalidating(true);
289 return true;
290 };
291 const primeEarlyHistory = async () => {
292 if (primeRef.current !== "idle") return;
293 // Only a blank transcript may be primed, and the check has to run before
294 // loadLatest: that call bumps the resident session generation (retiring
295 // the follower's in-flight reads) and replaces records before it reads
296 // `current`, so a transcript that already has content must never reach
297 // it. history_replace has no revision guard either — a stale or empty
298 // window landing after the follower's cut would wipe the conversation.
299 if (transcriptHasContent()) { primeRef.current = "retired"; return; }
300 primeRef.current = "loading";
301 // The prime is scoped to this mounted identity and to store ownership,
302 // not to a hydrate generation: hydrate() bumps the generation the moment
303 // it starts, and a follower that then fails or stalls must not have
304 // discarded the only baseline the tab could show.
305 const current = () => !cancelled && primeRef.current === "loading";
306 try {
307 const projection = await getTranscriptStore().loadLatest(tabId, sessionPath ?? "", { current });
308 if (!projection || !current()) return;
309 if (transcriptHasContent()) { primeRef.current = "retired"; return; }
310 primeRef.current = "primed";
311 setTranscript(current => reducer(current, historyReplaceAction(projection)));
312 } catch {
313 // Before the attach handshake lands (or on a legacy serve) the window
314 // read is unavailable. A miss stays non-fatal: the next attach
315 // publication retries, and the ready-time hydration takes over.
316 } finally {
317 if (primeRef.current === "loading") primeRef.current = "idle";
318 }
319 };
320 const refreshStatus = async () => {
321 const ticket = generation;
322 const status = await app.RemoteTabStatus(tabId);
323 if (cancelled || ticket !== generation) return;
324 applyRemoteStatus(status);
325 setTranscript(current => hydrateRemoteTelemetry(current, status));
326 };
327 const hydrate = async () => {
328 const ticket = ++generation;
329 follower?.stop();
330 follower = new TranscriptSessionFollower(tabId, sessionPath ?? "", true, action => {
331 if (cancelled || ticket !== generation) return;
332 // The follower's install cut makes it the store owner (connection
333 // status frames precede it and own nothing); an early history prime
334 // still in flight must not land after that cut.
335 if (action.type === "transcript_v2_snapshot") primeRef.current = "retired";
336 dispatch(action);
337 });
338 if (!cachePaintedRef.current) setHydrated(false);
339 try {
340 await follower.start();
341 const loaded = await loadRemoteStatusSnapshot(tabId, mountedState === "ready" ? 3 : 60,
342 () => cancelled || ticket !== generation, isAuthoritativeRemoteStatus, true);
343 if (!loaded || cancelled || ticket !== generation) return;
344 const [snapshot, status] = loaded;
345 applyRemoteStatus(status);
346 setCommands(Array.isArray(snapshot.commands) ? snapshot.commands as CommandInfo[] : []);
347 setTranscript(current => hydrateRemoteTelemetry(reducer(current,
348 { type: "checkpoints", checkpoints: remoteCheckpoints(snapshot.checkpoints) }), status));
349 hydratedRef.current = true;
350 setState("ready");
351 setHydrated(true);
352 setError("");
353 cachePaintedRef.current = false;
354 setRevalidating(false);
355 setSurfaceGeneration(value => value + 1);
356 void forkTargetsRefreshRef.current?.();
357 } catch (error) {
358 if (cancelled || ticket !== generation) return;
359 // The transcript protocol requires the live runtime that owns the
360 // session. Only a session taken over by a local runtime on the serve
361 // host (the Follow request answers 409, or status reports the
362 // take-over) may fall back to the identity/legacy history view; every
363 // other failure keeps its error and waits for the next ready
364 // publication or an explicit retry.
365 const takenOver = isRemoteTakeoverError(error) || await app.RemoteTabStatus(tabId).then(remoteStatusTakenOver, () => false);
366 if (cancelled || ticket !== generation) return;
367 if (!takenOver) { setError(String(error)); return; }
368 try {
369 const legacyLoaded = await loadRemoteStatusSnapshot(tabId, mountedState === "ready" ? 3 : 60,
370 () => cancelled || ticket !== generation, isAuthoritativeRemoteStatus, false);
371 if (!legacyLoaded || cancelled || ticket !== generation) return;
372 const [snapshot, status] = legacyLoaded;
373 const messages = Array.isArray(snapshot.history) ? snapshot.history as HistoryMessage[] : [];
374 const checkpoints = remoteCheckpoints(snapshot.checkpoints);
375 primeRef.current = "retired";
376 applyRemoteStatus(status);
377 setCommands(Array.isArray(snapshot.commands) ? snapshot.commands as CommandInfo[] : []);
378 setTranscript(current => {
379 let next = reducer(current, { type: "history", messages, remote: true });
380 next = reducer(next, { type: "checkpoints", checkpoints });
381 next = reducer(next, remoteStatusToAction(status, Date.now(), next.running));
382 return hydrateRemoteTelemetry(next, status);
383 });
384 hydratedRef.current = true;
385 setState("ready");
386 setHydrated(true);
387 setError("");
388 cachePaintedRef.current = false;
389 setRevalidating(false);
390 setSurfaceGeneration(value => value + 1);
391 void forkTargetsRefreshRef.current?.();
392 } catch (fallbackError) {
393 if (!cancelled && ticket === generation) setError(String(fallbackError));
394 }
395 }
396 };
397 const offContent = getTranscriptStore().subscribe(tabId, change => change.projection
398 ? dispatch({ type: "transcript_records", projection: change.projection, confirmedUsers: [] })
399 : dispatch({ type: "history_items_patch", patches: change.patches, expected: change.expected }));
400 navigateRef.current = async (target, current) => (await import("./historyTurnNavigation")).navigateHistoryTurn(tabId, transcriptRef.current, () => cancelled ? undefined : transcriptRef.current, dispatch, target, () => !cancelled && current());
401 olderRef.current = async () => {
402 if (transcriptRef.current.historyOlderLoading) return "empty";
403 dispatch({ type: "history_older_start" });
404 try {
405 const page = await getTranscriptStore().loadOlder(tabId, sessionPath ?? "");
406 if (!page || cancelled) return "empty";
407 if (page.kind === "reload") { await hydrate(); return "loaded"; }
408 if (transcriptRef.current.transcriptProtocol === 2) dispatch({ type: "transcript_records", projection: page, confirmedUsers: [] });
409 else dispatch({ type: "history_prepend", items: page.prependItems, removeIds: page.removeIds,
410 startTurn: page.startTurn, endTurn: page.endTurn, totalTurns: page.totalTurns,
411 hasOlder: page.hasOlder, hasNewer: page.hasNewer, revision: page.revision, digest: page.digest });
412 return "loaded";
413 } catch (error) {
414 dispatch({ type: "history_older_error", error: String(error) });
415 return "empty";
416 }
417 };
418 hydrateRef.current = { tabId, run: hydrate };
419 newerRef.current = async (latest = false, current = () => true) => {
420 if (transcriptRef.current.historyNewerLoading) return "empty";
421 dispatch({ type: "history_newer_start" });
422 try {
423 if (latest) {
424 const projection = await getTranscriptStore().loadLatest(tabId, sessionPath ?? "", { current: () => !cancelled && current() });
425 if (!projection || cancelled) { dispatch({ type: "history_newer_error", error: "" }); return "empty"; }
426 dispatch({ type: "transcript_records", projection: { ...projection, removeIds: [] }, confirmedUsers: [] });
427 return "loaded";
428 }
429 const page = await getTranscriptStore().loadNewer(tabId, sessionPath ?? "");
430 if (!page || cancelled) { dispatch({ type: "history_newer_error", error: "" }); return "empty"; }
431 if (page.kind === "stale") { await hydrate(); return "loaded"; }
432 if (transcriptRef.current.transcriptProtocol === 2) dispatch({ type: "transcript_records", projection: page, confirmedUsers: [] });
433 else dispatch({ type: "history_append", items: page.items,
434 startTurn: page.startTurn, endTurn: page.endTurn, totalTurns: page.totalTurns,
435 hasOlder: page.hasOlder, hasNewer: page.hasNewer, revision: page.revision, digest: page.digest });
436 return "loaded";
437 } catch (error) {
438 dispatch({ type: "history_newer_error", error: String(error) });
439 return "empty";
440 }
441 };
442 refreshStatusRef.current = { tabId, run: refreshStatus };
443 reconcileHistoryRef.current = hydrate;
444 const offState = onRemoteTabState(tabId, next => {
445 if (cancelled) return;
446 setState(next.state);
447 setError(next.error ?? "");
448 if (next.state === "ready") void hydrate();
449 else if (next.state === "disconnected") {
450 setHydrated(false);
451 setRevalidating(false);
452 dispatch({ type: "transcript_connection", status: "disconnected" });
453 }
454 });
455 // Ownership flips arrive as tab meta updates (an explicit reclaim clears
456 // the pin there before any status poll runs); mirror them into the
457 // spectator flag that drives the reconcile loop below.
458 const offMeta = onRemoteTabUpdated(meta => {
459 if (cancelled || meta?.id !== tabId) return;
460 noteOwnership(Boolean(meta.takenOver));
461 // The attach publication is the reliable "identity live, activation
462 // still in flight" signal — retry the early history read there.
463 void primeEarlyHistory();
464 });
465 // The legacy event channel carries ancillary invalidations only.
466 const offEvent = onRemoteTabEvent(tabId, raw => {
467 const event = raw as WireEvent;
468 if (event.kind === "turn_done") {
469 void refreshStatus().catch(() => undefined);
470 void forkTargetsRefreshRef.current?.();
471 }
472 });
473 if (revivedFromShell) void app.SetActiveTab(tabId).catch(() => undefined);
474 paintResidentCache();
475 void primeEarlyHistory();
476 void hydrate();
477 return () => {
478 submitBindingRef.current = {};
479 cancelled = true;
480 generation++;
481 follower?.stop();
482 offContent();
483 offState();
484 offMeta();
485 offEvent();
486 olderRef.current = undefined;
487 navigateRef.current = undefined;
488 newerRef.current = undefined;
489 hydrateRef.current = null;
490 refreshStatusRef.current = null;
491 reconcileHistoryRef.current = null;
492 };
493 }, [applyRemoteStatus, noteOwnership, tabId, sessionPath, setTranscript]);
494
495 // A spectator surface idles with no status traffic: the running watchdog
496 // only reconciles turns, and the read-only composer blocks the sends that
497 // would otherwise refresh status. A stale ownership observation could pin
498 // the takeover banner forever, so poll at a slow cadence until the serve
499 // reports the session free again.
500 useEffect(() => {
501 if (!tabId || state !== "ready" || !spectator) return;
502 const timer = window.setInterval(() => {
503 void refreshStatusRef.current?.run().catch(() => undefined);
504 }, 5_000);
505 return () => window.clearInterval(timer);
506 }, [tabId, state, spectator]);
507
508 const submit = useCallback(async (text: string, displayText = text, choice?: import("./modelApplication").ModelApplicationChoice) => {
509 if (!tabId) return;
510 if (transcriptRef.current.transcriptProtocol !== 2 || transcriptRef.current.transcriptConnection !== "connected") {
511 throw new Error("Transcript v2 is not synchronized. Upgrade Desktop and Serve together, or reconnect.");
512 }
513 const trimmed = text.trim();
514 if (!trimmed) return;
515 // Optimistic user bubble, exactly like the local send path. seq rides
516 // the reducer's counter; the submission id only needs uniqueness.
517 const before = getTranscriptStore().states.get(tabId) ?? initialState;
518 const binding = submitBindingRef.current;
519 const current = () => submitBindingRef.current === binding && getTranscriptStore().states.get(tabId)?.sessionGen === before.sessionGen;
520 const submissions = Object.values(before.localSubmissions);
521 if (submissions.some(item => item.status === "sending" && (item.submitText ?? item.text).trim() === trimmed)) return;
522 const unresolved = submissions.find(item => item.status === "unknown");
523 if(unresolved && (choice || (unresolved.submitText ?? unresolved.text).trim()!==trimmed)) {
524 throw Object.assign(new Error("Confirm the previous submission before sending another message"),{data:{submissionOutcome:"unknown"}});
525 }
526 const submissionId = unresolved?.submissionId ?? createTurnSubmissionId(tabId, before.sessionGen, before.seq, before.meta?.runtime?.epoch);
527 activityRevisionRef.current += 1;
528 runtimeAtActivityRef.current = runtimeState.state;
529 if(!unresolved) setTranscript((s) => reducer(s, { type: "user", text: displayText.trim(), submitText:trimmed, seq: s.seq, submissionId }));
530 try {
531 if (choice) {
532 if (!app.SubmitRemoteTabWithModelApplication) throw new Error("Upgrade Desktop to use model application recovery");
533 await app.SubmitRemoteTabWithModelApplication(tabId, trimmed, submissionId, choice);
534 } else if (app.SubmitRemoteTabWithSubmission) await app.SubmitRemoteTabWithSubmission(tabId, trimmed, submissionId);
535 else await app.SubmitRemoteTab(tabId, trimmed);
536 if (current()) setTranscript(s => reducer(s, { type: "send_confirmed", submissionId }));
537 } catch (e) {
538 // A busy-window submit is durably queued by the desktop. The queue strip
539 // is its representation, so retract the optimistic bubble instead of
540 // reporting a failure; the receipt rides the error for the composer.
541 if (queuedFollowupOutcome(e) !== undefined) {
542 if (current()) setTranscript((s) => reducer(s, { type: "send_queued", submissionId }));
543 throw e;
544 }
545 // Roll the optimistic running flag back — a refused/failed submit must
546 // never leave the pill spinning (same contract as the local send path).
547 const error = `Send failed: ${e instanceof Error ? e.message : String(e)}`;
548 if (current()) setTranscript((s) => reducer(s, { type: isUnknownSubmissionError(e) ? "turn_submit_unknown" : "send_failed", submissionId, error }));
549 throw e;
550 }
551 }, [tabId, runtimeState.state, setTranscript]);
552
553 const runManagementCommand = useCallback(async (text: string, rehydrate = false) => {
554 if (!tabId) return;
555 const trimmed = text.trim();
556 if (!trimmed) return;
557 // Management verbs produce notices/state changes rather than a model
558 // turn, so do not create the optimistic conversational bubble used by
559 // submit(). Refresh the authoritative profile after the command settles.
560 await app.SubmitRemoteTab(tabId, trimmed);
561 if (rehydrate) {
562 const hydration = hydrateRef.current;
563 if (hydration?.tabId === tabId) await hydration.run(true);
564 return;
565 }
566 const current = refreshStatusRef.current;
567 if (current?.tabId === tabId) await current.run();
568 }, [tabId]);
569
570 const cancelTurn = useCallback(async () => {
571 if (!tabId) return;
572 await app.CancelRemoteTab(tabId);
573 }, [tabId]);
574
575 const approve = useCallback(async (callId: string, decision: string) => {
576 if (!tabId) return;
577 setPromptError("");
578 try {
579 await app.ApproveRemoteTab(tabId, callId, decision);
580 setTranscript((s) => s.approval?.id === callId ? { ...s, approval: undefined } : s);
581 } catch (error) {
582 setPromptError(error instanceof Error ? error.message : String(error));
583 throw error;
584 }
585 }, [tabId]);
586
587 const resolvePlanDecision = useCallback(async (
588 callId: string,
589 action: "start_execution" | "revise_plan" | "exit_plan",
590 feedback = "",
591 ) => {
592 if (!tabId) return;
593 setPromptError("");
594 try {
595 await app.ResolveRemoteTabPlanDecision(tabId, callId, action, feedback);
596 setTranscript((s) => s.approval?.id === callId ? { ...s, approval: undefined } : s);
597 } catch (error) {
598 setPromptError(error instanceof Error ? error.message : String(error));
599 throw error;
600 }
601 }, [tabId]);
602
603 const answer = useCallback(async (callId: string, answers: RemoteAskAnswer[]) => {
604 if (!tabId) return;
605 setPromptError("");
606 try {
607 await app.AnswerRemoteTab(tabId, callId, answers);
608 setTranscript((s) => s.ask?.id === callId ? { ...s, ask: undefined } : s);
609 } catch (error) {
610 setPromptError(error instanceof Error ? error.message : String(error));
611 throw error;
612 }
613 }, [tabId]);
614
615 const clearExtensionForm = useCallback((pluginId: string, surfaceId: string, formInstanceId?: string) => {
616 setTranscript((s) => s.extensionForm?.pluginId === pluginId && s.extensionForm.surfaceId === surfaceId &&
617 (!formInstanceId || s.extensionForm.formInstanceId === formInstanceId)
618 ? reducer(s, { type: "clearExtensionForm" }) : s);
619 }, []);
620
621 const retryHydration = useCallback((): Promise<void> => {
622 setError("");
623 const current = hydrateRef.current;
624 if (!current || current.tabId !== tabId) return Promise.resolve();
625 return current.run(true);
626 }, [tabId]);
627
628 const compact = useCallback(async (instructions: string) => {
629 if (!tabId) return;
630 await app.CompactRemoteTab(tabId, instructions);
631 await retryHydration();
632 }, [retryHydration, tabId]);
633
634 const refreshStatus = useCallback((): Promise<void> => {
635 const current = refreshStatusRef.current;
636 if (!current || current.tabId !== tabId) return Promise.resolve();
637 return current.run();
638 }, [tabId]);
639
640 const cancelJob = useCallback(async (jobId: string) => {
641 if (!tabId) return false;
642 try {
643 await app.CancelRemoteTabJobs(tabId, [jobId]);
644 await refreshStatus();
645 return true;
646 } catch (error) {
647 setPromptError(String(error));
648 return false;
649 }
650 }, [refreshStatus, tabId]);
651
652 const rewind = useCallback(async (turn: number, scope: string) => {
653 if (!tabId) return;
654 setPromptError("");
655 try {
656 switch (scope) {
657 // No fork scope: the serve's /fork switches the parent session; forkTurn creates a child instead.
658 case "summ-from":
659 await app.SummarizeRemoteTab(tabId, turn, "from");
660 break;
661 case "summ-upto":
662 await app.SummarizeRemoteTab(tabId, turn, "upto");
663 break;
664 case "code":
665 case "conversation":
666 case "both":
667 await app.RewindRemoteTab(tabId, String(turn), scope);
668 break;
669 default:
670 throw new Error(`Unsupported remote rewind scope: ${scope}`);
671 }
672 await retryHydration();
673 } catch (error) {
674 setPromptError(error instanceof Error ? error.message : String(error));
675 throw error;
676 }
677 }, [retryHydration, tabId]);
678
679 const setEffort = useCallback(async (level: string) => {
680 if (!tabId) return;
681 await app.SetRemoteTabEffort(tabId, level);
682 await refreshStatus();
683 }, [refreshStatus, tabId]);
684
685 const setModel = useCallback(async (ref: string) => {
686 if (!tabId) return;
687 await app.SetRemoteTabModel(tabId, ref);
688 await refreshStatus();
689 }, [refreshStatus, tabId]);
690
691 const setQualityFloor = useCallback(async (floor: QualityFloor) => {
692 if (!tabId) return;
693 // Compatibility only. The new client never asks an old server to change
694 // policy behind the user's back; its next status remains authoritative.
695 if (floor !== "standard" && floor !== "delivery") throw new Error(`Unknown retired execution setting: ${floor}`);
696 await app.SetRemoteTabQualityFloor(tabId, floor);
697 }, [tabId]);
698
699 const pauseGoal = useCallback(async () => {
700 if (!tabId) return;
701 await app.PauseRemoteTabGoal(tabId);
702 await refreshStatus();
703 }, [refreshStatus, tabId]);
704
705 const resumeGoal = useCallback(async () => {
706 if (!tabId) return;
707 await app.ResumeRemoteTabGoal(tabId);
708 await refreshStatus();
709 }, [refreshStatus, tabId]);
710
711 const editGoal = useCallback(async (objective: string, maxGoalRounds: number | null) => {
712 if (!tabId) return;
713 await app.EditRemoteTabGoal(tabId, objective, maxGoalRounds);
714 await refreshStatus();
715 }, [refreshStatus, tabId]);
716
717 const steer = useCallback(async (input: string) => {
718 if (!tabId) return;
719 await app.SteerRemoteTab(tabId, input);
720 }, [tabId]);
721
722 const drainApprovals = useCallback((ids: string[]) => {
723 setTranscript((current) => reducer(current, { type: "approval_drained", ids, epoch: current.promptEpoch }));
724 }, []);
725
726 return {
727 state, error, transcript, liveStore, hydrated, revalidating, syncMode: "v2", navigateToTurn: (target, current) => navigateRef.current?.(target, current) ?? Promise.resolve("cancelled"), loadOlderHistory: (_targetTurn?: number, trigger?: HistoryLoadTrigger) => olderRef.current?.(trigger) ?? Promise.resolve("empty"), loadNewerHistory: (latest = false, current) => newerRef.current?.(latest, current) ?? Promise.resolve("empty"), running: transcript.running, modelLabel, commands,
728 composerProfile, goalRuntime, goalView, effort, surfaceGeneration, promptError, submit, runManagementCommand, compact, cancelTurn,
729 approve, resolvePlanDecision, answer, clearExtensionForm, rewind, forkTurn, acknowledgeFork, setModel, setEffort, setQualityFloor, pauseGoal, resumeGoal, editGoal, steer, cancelJob,
730 drainApprovals, retryHydration,
731 };
732 }
733
733 lines TYPESCRIPT