1import {2 createComposerAttachmentFile,3 prepareComposerMessageForWorkspace,4} from "../../lib/composerAttachments";5import * as desktopCommands from "../../lib/desktopCommands";6import { type NewChatLandingTarget, resolveDefaultNewChatTarget } from "../../lib/newChatLanding"
;
7import { buildAttachmentSignature, buildUserInputDisplayText } from "../attachmentInputs";
8import {
9 type ComposerDraft,
10 type ComposerDraftAttachment,
11 type ComposerDraftRevision,
12 clearComposerDraftRevision,
13 composerDraftKeyForNewChatTarget,
14 composerDraftKeyForThread,
15 createComposerDraftAttachment,
16 createEmptyComposerDraft,
17 getComposerDraftAttachmentValidationMessage,
18 hasComposerDraftState,
19 pruneComposerDrafts as pruneComposerDraftEntries,
20 resolveActiveComposerDraftKey,
21 revokeComposerDraftAttachmentPreviews,
22} from "../composerDrafts";
23import {
24 type ComposerSubmission,
25 type ComposerSubmissionRequest,
26 cloneComposerDraftForSubmission,
27 composerSubmissionErrorMessage,
28 findComposerSubmissionById,
29 isComposerSubmissionInFlight,
30} from "../composerSubmission";
31import { isInteractionThreadVisible } from "../interactionVisibility";
32import {
33 googleProviderOptionsForReasoningEffort,
34 isGoogleReasoningEffortValue,
35 isOpenAiReasoningEffortValue,
36} from "../openaiCompatibleProviderOptions";
37import {
38 type AbortableActionOptions,
39 type AppStoreActions,
40 type AppStoreState,
41 appendThreadTranscript,
42 beginThreadSelectionRequest,
43 buildContextPreamble,
44 bumpWorkspaceJsonRpcSocketGeneration,
45 bumpWorkspaceStartGeneration,
46 clearPendingThreadSteers,
47 clearThreadSelectionRequest,
48 clearWorkspaceJsonRpcSocketGeneration,
49 clearWorkspaceStartState,
50 disposeWorkspaceJsonRpcState,
51 ensureControlSocket,
52 ensureServerRunning,
53 ensureThreadRuntime,
54 ensureThreadSocket,
55 ensureWorkspaceRuntime,
56 extractUsageStateFromTranscript,
57 getEffectiveThreadLastEventSeq,
58 isCurrentThreadSelectionRequest,
59 makeId,
60 nowIso,
61 persistNow,
62 pushNotification,
63 queuePendingThreadMessage,
64 RUNTIME,
65 requestJsonRpcControlEvent,
66 requestSessionSnapshot,
67 type StoreGet,
68 type StoreSet,
69 sendThread,
70 sendUserMessageToThread,
71 syncDesktopStateCache,
72 truncateTitle,
73} from "../store.helpers";
74import type { FileAttachmentInput } from "../store.helpers/jsonRpcSocket";
75import { requestJsonRpc } from "../store.helpers/jsonRpcSocket";
76import { createOneOffWorkspaceRecord } from "../store.helpers/oneOffWorkspaceRecord";
77import {
78 beginCreationOperationIntent,
79 type CreationOperationControl,
80 invalidateNavigationIntent,
81 isCreationNavigationIntentCurrent,
82 isOperationAbortError,
83 recordThreadNavigationIntent,
84 waitForOperation,
85} from "../store.helpers/operationIntent";
86import { operationKey, runAcknowledgedOperation } from "../store.helpers/operations";
87import { waitForNextPaintOrTimeout } from "../store.helpers/paintScheduling";
88import { persist } from "../store.helpers/persistence";
89import { MAX_FEED_ITEMS } from "../store.helpers/threadEventReducerContext";
90import { isStandardChatThread } from "../threadFilters";
91import { hydrateTranscriptSnapshot } from "../transcriptHydration";
92import {
93 type ChatInteraction,
94 type FeedItem,
95 isOneOffChatWorkspace,
96 type SessionSnapshot,
97 type SessionSnapshotFingerprint,
98 type ThreadBusyPolicy,
99 type ThreadRecord,
100 type TranscriptEvent,
101} from "../types";
102
103const RECONNECT_OUTCOME_TIMEOUT_MS = 15_000;
104const RECONNECT_OUTCOME_POLL_MS = 50;
105const composerSubmissionCreationControl = new Map<string, CreationOperationControl>();
106
107type HydrateThreadSelectionOptions = {
108 preserveView?: boolean;
109 reconnectAfterHydration?: boolean;
110 skipWorkspaceSelectOnReconnect?: boolean;
111} & AbortableActionOptions;
112
113function createEmptyComposerDraftForState(state: AppStoreState, key: string): ComposerDraft {
114 const draft = createEmptyComposerDraft(nowIso());
115 const revisionFloor = state.composerDraftRevisionFloorByKey[key];
116 return revisionFloor
117 ? {
118 ...draft,
119 revision: revisionFloor.revision,
120 generation: revisionFloor.generation,
121 }
122 : draft;
123}
124
125/**
126 * Queue the first message of a brand-new chat AND render its user bubble
127 * immediately, before the workspace socket/thread session exist. The
128 * pre-generated clientMessageId travels with the queued message so the
129 * eventual `turn/start` reuses it and the server echo dedups against the
130 * optimistic bubble instead of rendering a duplicate.
131 */
132function queueOptimisticFirstThreadMessage(
133 set: StoreSet,
134 threadId: string,
135 text: string,
136 attachments?: FileAttachmentInput[],
137 references?: import("../../lib/wsProtocol").TurnReference[],
138 draftSubmission?: ComposerDraftRevision,
139 presetClientMessageId?: string,
140): void {
141 const trimmed = text.trim();
142 const hasAttachments = (attachments?.length ?? 0) > 0;
143 if (!trimmed && !hasAttachments) return;
144
145 const clientMessageId = presetClientMessageId ?? makeId();
146 queuePendingThreadMessage(
147 threadId,
148 trimmed,
149 attachments,
150 references,
151 clientMessageId,
152 draftSubmission,
153 );
154
155 const optimisticSeen = RUNTIME.optimisticUserMessageIds.get(threadId) ?? new Set<string>();
156 optimisticSeen.add(clientMessageId);
157 RUNTIME.optimisticUserMessageIds.set(threadId, optimisticSeen);
158
159 const bubble: FeedItem = {
160 id: clientMessageId,
161 kind: "message",
162 role: "user",
163 ts: nowIso(),
164 text: buildUserInputDisplayText(trimmed, attachments),
165 };
166 set((s) => {
167 const rt = s.threadRuntimeById[threadId];
168 if (!rt) return {};
169 return {
170 threadRuntimeById: {
171 ...s.threadRuntimeById,
172 [threadId]: {
173 ...rt,
174 feed: [...rt.feed, bubble].slice(-MAX_FEED_ITEMS),
175 pendingTurnStart: {
176 clientMessageId,
177 text: trimmed,
178 attachmentSignature: buildAttachmentSignature(attachments),
179 status: "sending",
180 },
181 },
182 },
183 };
184 });
185}
186
187function updateInteraction(
188 set: StoreSet,
189 threadId: string,
190 requestId: string,
191 update: (interaction: ChatInteraction) => ChatInteraction,
192): void {
193 set((state) => {
194 const interactions = state.interactionsByThread[threadId];
195 if (!interactions?.some((interaction) => interaction.requestId === requestId)) {
196 return {};
197 }
198 return {
199 interactionsByThread: {
200 ...state.interactionsByThread,
201 [threadId]: interactions.map((interaction) =>
202 interaction.requestId === requestId ? update(interaction) : interaction,
203 ),
204 },
205 };
206 });
207}
208
209function findLatestVisibleSandboxInteraction(
210 state: AppStoreState,
211): { threadId: string; interaction: ChatInteraction } | null {
212 const eligible = (interaction: ChatInteraction) =>
213 interaction.kind === "approval" &&
214 interaction.approvalKind === "sandbox" &&
215 (interaction.status === "pending" || interaction.status === "failed");
216 const selectedThreadId = state.selectedThreadId;
217 const selectedInteraction = selectedThreadId
218 ? state.interactionsByThread[selectedThreadId]?.filter(eligible).at(-1)
219 : undefined;
220 if (
221 selectedThreadId &&
222 selectedInteraction &&
223 isInteractionThreadVisible(state, selectedThreadId)
224 ) {
225 return { threadId: selectedThreadId, interaction: selectedInteraction };
226 }
227
228 let latest: { threadId: string; interaction: ChatInteraction } | null = null;
229 for (const [threadId, interactions] of Object.entries(state.interactionsByThread)) {
230 if (!isInteractionThreadVisible(state, threadId)) continue;
231 for (const interaction of interactions) {
232 if (!eligible(interaction)) continue;
233 if (!latest || interaction.receivedSequence > latest.interaction.receivedSequence) {
234 latest = { threadId, interaction };
235 }
236 }
237 }
238 return latest;
239}
240
241export async function hydrateThreadSelection(
242 get: StoreGet,
243 set: StoreSet,
244 threadId: string,
245 options: HydrateThreadSelectionOptions = {},
246): Promise<void> {
247 const isOperationCurrent = () => options.signal?.aborted !== true;
248 const isSelectionCurrent = (requestId: number) =>
249 isOperationCurrent() &&
250 get().selectedThreadId === threadId &&
251 isCurrentThreadSelectionRequest(threadId, requestId);
252
253 const clearThreadHydrationIfCurrent = (requestId: number) => {
254 if (!isOperationCurrent() || !isCurrentThreadSelectionRequest(threadId, requestId)) {
255 return;
256 }
257 set((state) => {
258 const rt = state.threadRuntimeById[threadId];
259 if (!rt) return {};
260 return {
261 threadRuntimeById: {
262 ...state.threadRuntimeById,
263 [threadId]: {
264 ...rt,
265 hydrating: false,
266 transcriptOnly: false,
267 },
268 },
269 };
270 });
271 clearThreadSelectionRequest(threadId, requestId);
272 };
273
274 if (!isOperationCurrent()) return;
275 const thread = get().threads.find((candidate) => candidate.id === threadId);
276 if (!thread) return;
277 const selectedTaskIdForThread = (candidate: ThreadRecord): string | null => {
278 if (isStandardChatThread(candidate, { includeDrafts: true })) return null;
279 return typeof candidate.taskId === "string" && candidate.taskId.trim().length > 0
280 ? candidate.taskId
281 : null;
282 };
283
284 const threadFingerprint = (candidate: ThreadRecord): SessionSnapshotFingerprint => ({
285 updatedAt: candidate.lastMessageAt,
286 messageCount: candidate.messageCount,
287 lastEventSeq: getEffectiveThreadLastEventSeq(get(), candidate.id),
288 });
289
290 const fingerprintMatches = (
291 left: SessionSnapshotFingerprint,
292 right: SessionSnapshotFingerprint,
293 ): boolean =>
294 left.updatedAt === right.updatedAt &&
295 left.messageCount === right.messageCount &&
296 left.lastEventSeq === right.lastEventSeq;
297
298 const cacheSessionSnapshot = (snapshot: SessionSnapshot) => {
299 if (!isOperationCurrent()) return;
300 RUNTIME.sessionSnapshots.set(snapshot.sessionId, {
301 fingerprint: {
302 updatedAt: snapshot.updatedAt,
303 messageCount: snapshot.messageCount,
304 lastEventSeq: snapshot.lastEventSeq,
305 },
306 snapshot,
307 });
308 syncDesktopStateCache(get);
309 };
310
311 const applySessionSnapshot = (
312 selectedThreadId: string,
313 sessionId: string,
314 snapshot: SessionSnapshot,
315 ) => {
316 if (!isOperationCurrent()) return;
317 set((state) => {
318 const nextThreads = state.threads.map((candidate) =>
319 candidate.id === selectedThreadId
320 ? {
321 ...candidate,
322 title: snapshot.title,
323 titleSource: snapshot.titleSource,
324 lastMessageAt: snapshot.updatedAt,
325 sessionId,
326 messageCount: snapshot.messageCount,
327 lastEventSeq: snapshot.lastEventSeq,
328 }
329 : candidate,
330 );
331 const currentRuntime = state.threadRuntimeById[selectedThreadId];
332 const currentThread = state.threads.find((candidate) => candidate.id === selectedThreadId);
333 const snapshotWorkspace = currentThread
334 ? state.workspaces.find((workspace) => workspace.id === currentThread.workspaceId)
335 : null;
336 const snapshotConfig =
337 currentRuntime?.config ??
338 (snapshotWorkspace
339 ? {
340 provider: snapshot.provider,
341 model: snapshot.model,
342 workingDirectory: snapshotWorkspace.path,
343 }
344 : null);
345 return {
346 threads: nextThreads,
347 threadRuntimeById: {
348 ...state.threadRuntimeById,
349 [selectedThreadId]: {
350 ...currentRuntime,
351 sessionId,
352 lastEventSeq: snapshot.lastEventSeq,
353 sessionKind: snapshot.sessionKind,
354 parentSessionId: snapshot.parentSessionId,
355 role: snapshot.role,
356 mode: snapshot.mode,
357 depth: snapshot.depth ?? 0,
358 nickname: snapshot.nickname,
359 requestedModel: snapshot.requestedModel,
360 effectiveModel: snapshot.effectiveModel,
361 requestedReasoningEffort: snapshot.requestedReasoningEffort,
362 effectiveReasoningEffort: snapshot.effectiveReasoningEffort,
363 executionState: snapshot.executionState,
364 lastMessagePreview: snapshot.lastMessagePreview,
365 agents: snapshot.agents,
366 sessionUsage: snapshot.sessionUsage,
367 lastTurnUsage: snapshot.lastTurnUsage,
368 feed: snapshot.feed,
369 hydrating: false,
370 transcriptOnly: false,
371 connected: currentRuntime?.connected ?? false,
372 config: snapshotConfig,
373 sessionConfig: currentRuntime?.sessionConfig ?? null,
374 enableMcp: currentRuntime?.enableMcp ?? null,
375 busy: currentRuntime?.busy ?? false,
376 busySince: currentRuntime?.busySince ?? null,
377 activeTurnId: currentRuntime?.activeTurnId ?? null,
378 pendingSteer: currentRuntime?.pendingSteer ?? null,
379 wsUrl: currentRuntime?.wsUrl ?? null,
380 },
381 },
382 };
383 });
384 };
385
386 const transcriptIdsForThread = (
387 candidate: Pick<ThreadRecord, "id" | "sessionId" | "legacyTranscriptId">,
388 ): string[] => {
389 const ids = [candidate.legacyTranscriptId ?? null, candidate.sessionId ?? null, candidate.id];
390 return [
391 ...new Set(
392 ids.filter(
393 (value): value is string => typeof value === "string" && value.trim().length > 0,
394 ),
395 ),
396 ];
397 };
398
399 const hydrateLegacyTranscript = async (candidate: ThreadRecord) => {
400 if (!isOperationCurrent()) return null;
401 const transcriptIds = transcriptIdsForThread(candidate);
402 if (transcriptIds.length === 0) return null;
403
404 if (transcriptIds.length === 1 && typeof desktopCommands.hydrateTranscript === "function") {
405 const firstTranscriptId = transcriptIds[0];
406 if (!firstTranscriptId) return null;
407 const snapshot = await desktopCommands.hydrateTranscript({ threadId: firstTranscriptId });
408 return isOperationCurrent() ? snapshot : null;
409 }
410
411 const transcript: TranscriptEvent[] = [];
412 let successfulReads = 0;
413 let firstError: unknown = null;
414 for (const transcriptId of transcriptIds) {
415 if (!isOperationCurrent()) return null;
416 try {
417 transcript.push(...(await desktopCommands.readTranscript({ threadId: transcriptId })));
418 if (!isOperationCurrent()) return null;
419 successfulReads += 1;
420 } catch (error) {
421 firstError ??= error;
422 }
423 }
424
425 if (successfulReads === 0 && firstError) {
426 throw firstError;
427 }
428
429 transcript.sort((left, right) => left.ts.localeCompare(right.ts));
430 return hydrateTranscriptSnapshot(transcript);
431 };
432
433 if (!isOperationCurrent()) return;
434 ensureThreadRuntime(get, set, threadId);
435 if (thread.draft) {
436 if (!isOperationCurrent()) return;
437 const selectedTaskId = selectedTaskIdForThread(thread);
438 set((state) => {
439 return {
440 selectedThreadId: threadId,
441 selectedWorkspaceId: thread.workspaceId,
442 selectedTaskId,
443 view: options.preserveView ? state.view : "chat",
444 threadRuntimeById: {
445 ...state.threadRuntimeById,
446 [threadId]: {
447 ...state.threadRuntimeById[threadId],
448 hydrating: false,
449 transcriptOnly: false,
450 },
451 },
452 };
453 });
454 syncDesktopStateCache(get);
455 return;
456 }
457
458 const rt = get().threadRuntimeById[threadId];
459 if (get().selectedThreadId === threadId && RUNTIME.threadSelectionRequests.has(threadId)) {
460 if (!isOperationCurrent()) return;
461 const selectedTaskId = selectedTaskIdForThread(thread);
462 set((state) => ({
463 selectedWorkspaceId: thread.workspaceId,
464 selectedTaskId,
465 view: options.preserveView ? state.view : "chat",
466 }));
467 syncDesktopStateCache(get);
468 return;
469 }
470 if (get().selectedThreadId === threadId && rt?.connected) {
471 if (!isOperationCurrent()) return;
472 const selectedTaskId = selectedTaskIdForThread(thread);
473 set((state) => ({
474 selectedWorkspaceId: thread.workspaceId,
475 selectedTaskId,
476 view: options.preserveView ? state.view : "chat",
477 }));
478 syncDesktopStateCache(get);
479 if (options.reconnectAfterHydration) {
480 const reconnect = get().reconnectThread(threadId, undefined, {
481 skipWorkspaceSelect: options.skipWorkspaceSelectOnReconnect,
482 refreshSnapshot: false,
483 signal: options.signal,
484 });
485 if (options.signal) {
486 await reconnect.catch(() => {
487 // The next socket reconnect or thread selection will re-assert the live subscription.
488 });
489 } else {
490 void reconnect.catch(() => {
491 // The next socket reconnect or thread selection will re-assert the live subscription.
492 });
493 }
494 }
495 return;
496 }
497
498 const alreadyLoaded = rt?.feed && rt.feed.length > 0;
499 const sessionId = rt?.sessionId ?? thread.sessionId;
500 const expectedFingerprint = threadFingerprint(thread);
501 const cachedSnapshot = sessionId ? RUNTIME.sessionSnapshots.get(sessionId) : null;
502 const matchingCachedSnapshot =
503 sessionId &&
504 cachedSnapshot &&
505 fingerprintMatches(cachedSnapshot.fingerprint, expectedFingerprint)
506 ? cachedSnapshot.snapshot
507 : null;
508
509 if (sessionId && cachedSnapshot && !matchingCachedSnapshot) {
510 console.debug(
511 `[selectThread] Cache fingerprint mismatch for session ${sessionId}: cached ${JSON.stringify(cachedSnapshot.fingerprint)} vs expected ${JSON.stringify(expectedFingerprint)}`,
512 );
513 }
514
515 const skipHarnessSnapshotFetch = Boolean(alreadyLoaded && matchingCachedSnapshot);
516 const shouldFetchHarnessSnapshot =
517 Boolean(sessionId) &&
518 !skipHarnessSnapshotFetch &&
519 (thread.messageCount > 0 ||
520 getEffectiveThreadLastEventSeq(get(), thread.id) > 0 ||
521 Boolean(thread.legacyTranscriptId));
522
523 if (!isOperationCurrent()) return;
524 const requestId = beginThreadSelectionRequest(threadId);
525 const selectedTaskId = selectedTaskIdForThread(thread);
526 if (!isOperationCurrent()) return;
527 set((state) => {
528 if (!isOperationCurrent()) return {};
529 return {
530 selectedThreadId: threadId,
531 selectedWorkspaceId: thread.workspaceId,
532 selectedTaskId,
533 view: options.preserveView ? state.view : "chat",
534 threadRuntimeById: {
535 ...state.threadRuntimeById,
536 [threadId]: {
537 ...state.threadRuntimeById[threadId],
538 hydrating: !alreadyLoaded && !matchingCachedSnapshot,
539 transcriptOnly: false,
540 },
541 },
542 };
543 });
544 if (!isOperationCurrent()) return;
545 syncDesktopStateCache(get);
546
547 let appliedCachedSnapshot = false;
548 if (matchingCachedSnapshot && sessionId) {
549 if (!isOperationCurrent()) return;
550 applySessionSnapshot(threadId, sessionId, matchingCachedSnapshot);
551 appliedCachedSnapshot = true;
552 }
553
554 await waitForNextPaintOrTimeout();
555 if (!isSelectionCurrent(requestId)) {
556 clearThreadHydrationIfCurrent(requestId);
557 return;
558 }
559
560 if (!isSelectionCurrent(requestId)) {
561 if (appliedCachedSnapshot) {
562 clearThreadHydrationIfCurrent(requestId);
563 }
564 return;
565 }
566
567 let stayTranscriptOnly = false;
568 if (!alreadyLoaded || matchingCachedSnapshot) {
569 try {
570 let loadedFromHarness = false;
571 if (sessionId && shouldFetchHarnessSnapshot) {
572 await ensureServerRunning(get, set, thread.workspaceId, { signal: options.signal });
573 if (!isSelectionCurrent(requestId)) {
574 clearThreadHydrationIfCurrent(requestId);
575 return;
576 }
577 ensureControlSocket(get, set, thread.workspaceId);
578 if (!isSelectionCurrent(requestId)) {
579 clearThreadHydrationIfCurrent(requestId);
580 return;
581 }
582 const snapshot = await requestSessionSnapshot(get, set, thread.workspaceId, sessionId);
583 if (!isSelectionCurrent(requestId)) {
584 clearThreadHydrationIfCurrent(requestId);
585 return;
586 }
587 if (snapshot) {
588 if (!isSelectionCurrent(requestId)) return;
589 applySessionSnapshot(threadId, sessionId, snapshot);
590 cacheSessionSnapshot(snapshot);
591 loadedFromHarness = true;
592 } else if (matchingCachedSnapshot) {
593 applySessionSnapshot(threadId, sessionId, matchingCachedSnapshot);
594 } else {
595 stayTranscriptOnly = true;
596 }
597 }
598
599 if (!loadedFromHarness && !matchingCachedSnapshot && !alreadyLoaded) {
600 const snapshot = await hydrateLegacyTranscript(thread);
601 if (!snapshot) {
602 throw new Error("No harness snapshot or legacy transcript cache was available.");
603 }
604 if (!isSelectionCurrent(requestId)) {
605 clearThreadHydrationIfCurrent(requestId);
606 return;
607 }
608 set((state) => {
609 const currentRuntime = state.threadRuntimeById[threadId];
610 return {
611 threadRuntimeById: {
612 ...state.threadRuntimeById,
613 [threadId]: {
614 ...currentRuntime,
615 sessionUsage: snapshot.sessionUsage,
616 lastTurnUsage: snapshot.lastTurnUsage,
617 agents: snapshot.agents,
618 feed: snapshot.feed,
619 hydrating: false,
620 transcriptOnly: true,
621 },
622 },
623 };
624 });
625 }
626 } catch (error) {
627 if (!isSelectionCurrent(requestId)) {
628 clearThreadHydrationIfCurrent(requestId);
629 return;
630 }
631
632 const detail = error instanceof Error ? error.message : String(error);
633 set((state) => ({
634 notifications: pushNotification(state.notifications, {
635 id: makeId(),
636 ts: nowIso(),
637 kind: "error",
638 title: "Transcript load failed",
639 detail,
640 }),
641 }));
642 clearThreadHydrationIfCurrent(requestId);
643 return;
644 }
645 }
646
647 if (!isSelectionCurrent(requestId)) {
648 clearThreadHydrationIfCurrent(requestId);
649 return;
650 }
651
652 if (stayTranscriptOnly) {
653 if (!isOperationCurrent()) return;
654 set((state) => ({
655 threadRuntimeById: {
656 ...state.threadRuntimeById,
657 [threadId]: {
658 ...state.threadRuntimeById[threadId],
659 hydrating: false,
660 transcriptOnly: true,
661 },
662 },
663 }));
664 clearThreadSelectionRequest(threadId, requestId);
665 return;
666 }
667
668 if (!isOperationCurrent()) return;
669 set((state) => ({
670 threadRuntimeById: {
671 ...state.threadRuntimeById,
672 [threadId]: { ...state.threadRuntimeById[threadId], hydrating: false, transcriptOnly: false },
673 },
674 }));
675
676 if (options.reconnectAfterHydration) {
677 await get().reconnectThread(threadId, undefined, {
678 selectionRequestId: requestId,
679 skipWorkspaceSelect: options.skipWorkspaceSelectOnReconnect,
680 signal: options.signal,
681 });
682 }
683 if (!isOperationCurrent()) return;
684 clearThreadSelectionRequest(threadId, requestId);
685}
686
687export function createThreadActions(
688 set: StoreSet,
689 get: StoreGet,
690): Pick<
691 AppStoreActions,
692 | "removeThread"
693 | "archiveThread"
694 | "restoreThread"
695 | "deleteThreadHistory"
696 | "renameThread"
697 | "newThread"
698 | "openNewChatLanding"
699 | "setNewChatLandingTarget"
700 | "selectThread"
701 | "reconnectThread"
702 | "reconnectThreadWithFeedback"
703 | "sendMessage"
704 | "submitComposerDraft"
705 | "retryComposerSubmission"
706 | "cancelComposerSubmission"
707 | "editAcceptedComposerSubmission"
708 | "dismissComposerSubmission"
709 | "completeComposerSubmission"
710 | "failComposerSubmission"
711 | "cancelThread"
712 | "clearThreadUsageHardCap"
713 | "setThreadModel"
714 | "setThreadReasoningEffort"
715 | "setComposerText"
716 | "addComposerAttachments"
717 | "removeComposerAttachment"
718 | "setComposerDraftModel"
719 | "setComposerDraftReasoningEffort"
720 | "clearComposerDraft"
721 | "discardComposerDraft"
722 | "pruneComposerDrafts"
723 | "setInjectContext"
724 | "answerAsk"
725 | "answerApproval"
726 | "dismissPrompt"
727 | "retryInteractionResponse"
728 | "loadAllThreadUsage"
729> {
730 const closeThreadSession = (threadId: string) => {
731 sendThread(get, threadId, (sessionId) => ({ type: "session_close", sessionId }));
732 };
733
734 const transcriptIdsForThread = (
735 thread: Pick<ThreadRecord, "id" | "sessionId" | "legacyTranscriptId">,
736 ): string[] => {
737 const ids = [thread.legacyTranscriptId ?? null, thread.sessionId ?? null, thread.id];
738 return [
739 ...new Set(
740 ids.filter(
741 (value): value is string => typeof value === "string" && value.trim().length > 0,
742 ),
743 ),
744 ];
745 };
746
747 const sessionSnapshotIdsForThread = (
748 thread: Pick<ThreadRecord, "sessionId">,
749 runtimeSessionId?: string | null,
750 ): string[] => {
751 const ids = [runtimeSessionId ?? null, thread.sessionId ?? null];
752 return [
753 ...new Set(
754 ids.filter(
755 (value): value is string => typeof value === "string" && value.trim().length > 0,
756 ),
757 ),
758 ];
759 };
760
761 const readTranscriptEvents = async (
762 thread: Pick<ThreadRecord, "id" | "sessionId" | "legacyTranscriptId">,
763 ): Promise<TranscriptEvent[] | null> => {
764 const transcriptIds = transcriptIdsForThread(thread);
765 if (transcriptIds.length === 0) return null;
766
767 const transcripts: TranscriptEvent[][] = [];
768 let successfulReads = 0;
769 let firstError: unknown = null;
770
771 for (const transcriptId of transcriptIds) {
772 try {
773 const events = await desktopCommands.readTranscript({ threadId: transcriptId });
774 transcripts.push(events);
775 successfulReads += 1;
776 } catch (error) {
777 firstError ??= error;
778 }
779 }
780
781 if (successfulReads === 0 && firstError) {
782 throw firstError;
783 }
784
785 return transcripts.flat().sort((left, right) => left.ts.localeCompare(right.ts));
786 };
787
788 const projectWorkspaces = () =>
789 get().workspaces.filter((workspace) => !isOneOffChatWorkspace(workspace));
790
791 const cleanupRemovedWorkspaceRuntime = async (workspaceId: string): Promise<void> => {
792 bumpWorkspaceStartGeneration(workspaceId);
793 bumpWorkspaceJsonRpcSocketGeneration(workspaceId);
794 const jsonRpcSocket = RUNTIME.jsonRpcSockets.get(workspaceId);
795 try {
796 jsonRpcSocket?.close();
797 } catch {
798 // ignore
799 }
800 RUNTIME.jsonRpcSockets.delete(workspaceId);
801 clearWorkspaceJsonRpcSocketGeneration(workspaceId);
802
803 try {
804 await desktopCommands.stopWorkspaceServer({ workspaceId });
805 } catch {
806 // ignore
807 } finally {
808 disposeWorkspaceJsonRpcState(get, workspaceId);
809 clearWorkspaceStartState(workspaceId);
810 }
811 };
812
813 const discardCancelledOneOffWorkspace = async (workspace: {
814 id: string;
815 path: string;
816 }): Promise<void> => {
817 const workspaceId = workspace.id;
818 const workspaceRecord = get().workspaces.find((entry) => entry.id === workspaceId);
819 if (workspaceRecord) {
820 await cleanupRemovedWorkspaceRuntime(workspaceId);
821 set((state) => {
822 const remainingWorkspaces = state.workspaces.filter((entry) => entry.id !== workspaceId);
823 return {
824 workspaces: remainingWorkspaces,
825 quickChatPreparedWorkspaceId:
826 state.quickChatPreparedWorkspaceId === workspaceId
827 ? null
828 : state.quickChatPreparedWorkspaceId,
829 selectedWorkspaceId:
830 state.selectedWorkspaceId === workspaceId
831 ? (remainingWorkspaces[0]?.id ?? null)
832 : state.selectedWorkspaceId,
833 };
834 });
835 await persistNow(get);
836 }
837 try {
838 await desktopCommands.trashPath({ path: workspace.path });
839 } catch {
840 // Best-effort cleanup; the path remains confined to Cowork's one-off chat root.
841 }
842 };
843
844 const submissionEntry = (
845 submissionId: string,
846 ): { key: string; submission: ComposerSubmission } | null => {
847 for (const [key, submission] of Object.entries(get().composerSubmissionsByKey)) {
848 if (submission.id === submissionId) return { key, submission };
849 }
850 return null;
851 };
852
853 const updateSubmission = (
854 submissionId: string,
855 update: (submission: ComposerSubmission) => ComposerSubmission,
856 ): boolean => {
857 let updated = false;
858 set((state) => {
859 const submission = findComposerSubmissionById(state.composerSubmissionsByKey, submissionId);
860 if (!submission) return {};
861 const key = Object.keys(state.composerSubmissionsByKey).find(
862 (candidate) => state.composerSubmissionsByKey[candidate]?.id === submissionId,
863 );
864 if (!key) return {};
865 updated = true;
866 return {
867 composerSubmissionsByKey: {
868 ...state.composerSubmissionsByKey,
869 [key]: update(submission),
870 },
871 };
872 });
873 return updated;
874 };
875
876 const markSubmissionSending = (
877 submissionId: string | undefined,
878 delivery?: ComposerSubmission["delivery"],
879 ): void => {
880 if (!submissionId) return;
881 updateSubmission(submissionId, (submission) => ({
882 ...submission,
883 phase: "sending",
884 delivery: delivery ?? submission.delivery,
885 error: null,
886 }));
887 };
888
889 const failSubmission = (submissionId: string, error: unknown): void => {
890 updateSubmission(submissionId, (submission) => ({
891 ...submission,
892 phase: "failed",
893 error: composerSubmissionErrorMessage(error),
894 }));
895 };
896
897 const removeSubmission = (submissionId: string): void => {
898 set((state) => {
899 const key = Object.keys(state.composerSubmissionsByKey).find(
900 (candidate) => state.composerSubmissionsByKey[candidate]?.id === submissionId,
901 );
902 if (!key) return {};
903 const next = { ...state.composerSubmissionsByKey };
904 delete next[key];
905 return { composerSubmissionsByKey: next };
906 });
907 };
908
909 const revokeAttachmentPreviewsIfUnreferenced = (
910 attachments: readonly ComposerDraftAttachment[],
911 ): void => {
912 const state = get();
913 const retainedPreviewUrls = new Set([
914 ...Object.values(state.composerDraftsByKey).flatMap((draft) =>
915 draft.attachments
916 .map((attachment) => attachment.previewUrl)
917 .filter((previewUrl): previewUrl is string => Boolean(previewUrl)),
918 ),
919 ...Object.values(state.composerSubmissionsByKey).flatMap((submission) =>
920 submission.draft.attachments
921 .map((attachment) => attachment.previewUrl)
922 .filter((previewUrl): previewUrl is string => Boolean(previewUrl)),
923 ),
924 ]);
925 revokeComposerDraftAttachmentPreviews(
926 attachments.filter(
927 (attachment) =>
928 typeof attachment.previewUrl === "string" &&
929 !retainedPreviewUrls.has(attachment.previewUrl),
930 ),
931 );
932 };
933
934 const currentSubmissionDelivery = (
935 request: ComposerSubmissionRequest,
936 ): ComposerSubmission["delivery"] =>
937 request.kind === "thread" && get().threadRuntimeById[request.threadId]?.busy ? "steer" : "send";
938
939 const executeComposerSubmission = async (submissionId: string): Promise<void> => {
940 const entry = submissionEntry(submissionId);
941 if (entry?.submission.phase !== "preparing") return;
942 const { submission } = entry;
943
944 try {
945 if (submission.request.kind === "newChat") {
946 const { target, provider, model, reasoningEffort } = submission.request;
947 const creationControl = composerSubmissionCreationControl.get(submissionId);
948 const started = await get().newThread({
949 scope: target.kind === "project" ? "project" : "oneOff",
950 ...(target.kind === "project" ? { workspaceId: target.workspaceId } : {}),
951 firstMessage: submission.draft.text,
952 titleHint:
953 submission.draft.text.trim() || submission.draft.attachments[0]?.filename || "New chat",
954 mode: "session",
955 draftAttachments: submission.draft.attachments,
956 references:
957 submission.draft.references.length > 0 ? submission.draft.references : undefined,
958 provider,
959 model,
960 reasoningEffort: reasoningEffort ?? undefined,
961 draftSubmission: submission.owner,
962 clientMessageId: submission.clientMessageId,
963 ...creationControl,
964 });
965 if (!started) {
966 if (creationControl?.signal?.aborted) {
967 removeSubmission(submissionId);
968 } else {
969 failSubmission(submissionId, new Error("The new chat could not be started."));
970 }
971 }
972 return;
973 }
974
975 const targetThreadId = submission.request.threadId;
976 const thread = get().threads.find((candidate) => candidate.id === targetThreadId);
977 if (!thread) {
978 failSubmission(submissionId, new Error("The target chat is no longer available."));
979 return;
980 }
981 const prepared =
982 submission.prepared ??
983 (await prepareComposerMessageForWorkspace(
984 get,
985 set,
986 thread.workspaceId,
987 submission.draft.text,
988 submission.draft.attachments,
989 { threadId: thread.id },
990 ));
991 if (!submission.prepared) {
992 updateSubmission(submissionId, (current) =>
993 current.phase === "preparing" ? { ...current, prepared } : current,
994 );
995 }
996 const latest = submissionEntry(submissionId);
997 if (latest?.submission.phase !== "preparing") return;
998 const exactPrepared = latest.submission.prepared ?? prepared;
999 if (!exactPrepared.text.trim() && !exactPrepared.attachments?.length) {
1000 failSubmission(submissionId, new Error("The message has no sendable content."));
1001 return;
1002 }
1003
1004 const delivery = currentSubmissionDelivery(latest.submission.request);
1005 markSubmissionSending(submissionId, delivery);
1006 const accepted = await get().sendMessage(
1007 exactPrepared.text,
1008 delivery === "steer" ? "steer" : "reject",
1009 exactPrepared.attachments,
1010 latest.submission.draft.references.length > 0
1011 ? latest.submission.draft.references
1012 : undefined,
1013 {
1014 targetThreadId: thread.id,
1015 draftSubmission: latest.submission.owner,
1016 clientMessageId: latest.submission.clientMessageId,
1017 },
1018 );
1019 if (!accepted) {
1020 failSubmission(submissionId, new Error("The message was not accepted by the chat."));
1021 }
1022 } catch (error) {
1023 if (composerSubmissionCreationControl.get(submissionId)?.signal?.aborted) {
1024 removeSubmission(submissionId);
1025 } else {
1026 failSubmission(submissionId, error);
1027 }
1028 } finally {
1029 composerSubmissionCreationControl.delete(submissionId);
1030 }
1031 };
1032
1033 return {
1034 archiveThread: async (threadId: string) => {
1035 set((s) => ({
1036 threads: s.threads.map((t) =>
1037 t.id === threadId ? { ...t, archived: true, archivedAt: nowIso() } : t,
1038 ),
1039 selectedThreadId: s.selectedThreadId === threadId ? null : s.selectedThreadId,
1040 }));
1041 await persistNow(get);
1042 },
1043
1044 restoreThread: async (threadId: string) => {
1045 set((s) => ({
1046 threads: s.threads.map((t) =>
1047 t.id === threadId ? { ...t, archived: false, archivedAt: undefined } : t,
1048 ),
1049 }));
1050 await persistNow(get);
1051 },
1052
1053 removeThread: async (threadId: string) => {
1054 const thread = get().threads.find((t) => t.id === threadId);
1055 get().discardComposerDraft(composerDraftKeyForThread(threadId));
1056 const runtimeSessionId = get().threadRuntimeById[threadId]?.sessionId ?? null;
1057 const sessionSnapshotIds = thread
1058 ? sessionSnapshotIdsForThread(thread, runtimeSessionId)
1059 : runtimeSessionId
1060 ? [runtimeSessionId]
1061 : [];
1062 closeThreadSession(threadId);
1063 RUNTIME.optimisticUserMessageIds.delete(threadId);
1064 RUNTIME.pendingThreadMessages.delete(threadId);
1065 RUNTIME.pendingThreadAttachments.delete(threadId);
1066 RUNTIME.pendingThreadReferences.delete(threadId);
1067 RUNTIME.pendingWorkspaceDefaultApplyByThread.delete(threadId);
1068 RUNTIME.modelStreamByThread.delete(threadId);
1069 RUNTIME.threadSelectionRequests.delete(threadId);
1070 clearPendingThreadSteers(threadId);
1071
1072 for (const sessionId of sessionSnapshotIds) {
1073 RUNTIME.sessionSnapshots.delete(sessionId);
1074 }
1075
1076 const threadWorkspace = thread
1077 ? (get().workspaces.find((workspace) => workspace.id === thread.workspaceId) ?? null)
1078 : null;
1079 const removeOneOffWorkspace =
1080 Boolean(thread && isOneOffChatWorkspace(threadWorkspace)) &&
1081 !get().threads.some(
1082 (candidate) => candidate.id !== threadId && candidate.workspaceId === thread?.workspaceId,
1083 );
1084 const workspaceIdToRemove = removeOneOffWorkspace && thread ? thread.workspaceId : null;
1085
1086 set((s) => {
1087 const remainingThreads = s.threads.filter((t) => t.id !== threadId);
1088 const selectedThreadId = s.selectedThreadId === threadId ? null : s.selectedThreadId;
1089 const remainingWorkspaces = workspaceIdToRemove
1090 ? s.workspaces.filter((workspace) => workspace.id !== workspaceIdToRemove)
1091 : s.workspaces;
1092 const fallbackWorkspaceId =
1093 remainingWorkspaces.find((workspace) => !isOneOffChatWorkspace(workspace))?.id ??
1094 remainingWorkspaces[0]?.id ??
1095 null;
1096
1097 const nextThreadRuntimeById = { ...s.threadRuntimeById };
1098 delete nextThreadRuntimeById[threadId];
1099
1100 const nextInteractionsByThread = { ...s.interactionsByThread };
1101 delete nextInteractionsByThread[threadId];
1102
1103 return {
1104 workspaces: remainingWorkspaces,
1105 threads: remainingThreads,
1106 selectedThreadId,
1107 interactionsByThread: nextInteractionsByThread,
1108 threadRuntimeById: nextThreadRuntimeById,
1109 selectedWorkspaceId:
1110 s.selectedWorkspaceId === workspaceIdToRemove
1111 ? fallbackWorkspaceId
1112 : s.selectedWorkspaceId,
1113 };
1114 });
1115
1116 if (workspaceIdToRemove) {
1117 await cleanupRemovedWorkspaceRuntime(workspaceIdToRemove);
1118 }
1119
1120 if (thread) {
1121 for (const transcriptId of transcriptIdsForThread(thread)) {
1122 try {
1123 await desktopCommands.deleteTranscript({ threadId: transcriptId });
1124 } catch {
1125 // ignore
1126 }
1127 }
1128 }
1129
1130 await persistNow(get);
1131 },
1132
1133 deleteThreadHistory: async (threadId: string) => {
1134 const thread = get().threads.find((t) => t.id === threadId);
1135 if (!thread) return;
1136 const targetSessionId = get().threadRuntimeById[threadId]?.sessionId ?? thread.sessionId;
1137
1138 let deleteOk = false;
1139 if (targetSessionId) {
1140 await ensureServerRunning(get, set, thread.workspaceId);
1141 ensureControlSocket(get, set, thread.workspaceId);
1142 deleteOk = await requestJsonRpcControlEvent(
1143 get,
1144 set,
1145 thread.workspaceId,
1146 "cowork/session/delete",
1147 {
1148 cwd: get().workspaces.find((workspace) => workspace.id === thread.workspaceId)?.path,
1149 targetSessionId,
1150 },
1151 );
1152 }
1153
1154 await get().removeThread(threadId);
1155
1156 if (!targetSessionId) return;
1157
1158 set((s) => ({
1159 notifications: pushNotification(s.notifications, {
1160 id: makeId(),
1161 ts: nowIso(),
1162 kind: deleteOk ? "info" : "error",
1163 title: deleteOk ? "Session history deleted" : "Delete session history failed",
1164 detail: deleteOk ? targetSessionId : "Control session is unavailable.",
1165 }),
1166 }));
1167 },
1168
1169 renameThread: (threadId: string, newTitle: string) => {
1170 const trimmed = newTitle.trim();
1171 if (!trimmed) return;
1172
1173 set((s) => ({
1174 threads: s.threads.map((t) =>
1175 t.id === threadId ? { ...t, title: trimmed, titleSource: "manual" } : t,
1176 ),
1177 }));
1178 void persistNow(get);
1179
1180 sendThread(get, threadId, (sessionId) => ({
1181 type: "set_session_title",
1182 sessionId,
1183 title: trimmed,
1184 }));
1185 },
1186
1187 newThread: async (opts) => {
1188 const operationIntent = opts?.intent ?? beginCreationOperationIntent();
1189 const canNavigate = () => isCreationNavigationIntentCurrent(operationIntent);
1190 const reportPhase = opts?.onPhase ?? (() => {});
1191 reportPhase("preparing");
1192 if (opts?.signal?.aborted) return false;
1193 const explicitWorkspace = opts?.workspaceId
1194 ? (get().workspaces.find((workspace) => workspace.id === opts.workspaceId) ?? null)
1195 : null;
1196 const scope =
1197 opts?.scope ??
1198 (explicitWorkspace && !isOneOffChatWorkspace(explicitWorkspace) ? "project" : "oneOff");
1199 const hasQueuedAttachments =
1200 (opts?.attachments && opts.attachments.length > 0) ||
1201 (opts?.attachmentFiles && opts.attachmentFiles.length > 0) ||
1202 (opts?.draftAttachments && opts.draftAttachments.length > 0);
1203 const createSessionImmediately =
1204 opts?.mode === "session" || Boolean(opts?.firstMessage?.trim()) || hasQueuedAttachments;
1205
1206 let workspaceId: string | null = null;
1207 let createdOneOffWorkspace: Awaited<ReturnType<typeof createOneOffWorkspaceRecord>> | null =
1208 null;
1209 if (scope === "oneOff") {
1210 const preparedWorkspaceId = get().quickChatPreparedWorkspaceId;
1211 const preparedWorkspace = preparedWorkspaceId
1212 ? (get().workspaces.find((workspace) => workspace.id === preparedWorkspaceId) ?? null)
1213 : null;
1214 if (preparedWorkspace) {
1215 createdOneOffWorkspace = preparedWorkspace;
1216 workspaceId = preparedWorkspace.id;
1217 set({
1218 quickChatPreparedWorkspaceId: null,
1219 ...(canNavigate() ? { selectedWorkspaceId: preparedWorkspace.id } : {}),
1220 });
1221 ensureWorkspaceRuntime(get, set, preparedWorkspace.id);
1222 } else {
1223 if (preparedWorkspaceId) {
1224 set({ quickChatPreparedWorkspaceId: null });
1225 }
1226 const oneOffWorkspacePromise = createOneOffWorkspaceRecord(
1227 get,
1228 opts?.titleHint ?? opts?.firstMessage,
1229 );
1230 try {
1231 const oneOffWorkspace = await waitForOperation(oneOffWorkspacePromise, opts?.signal);
1232 createdOneOffWorkspace = oneOffWorkspace;
1233 workspaceId = oneOffWorkspace.id;
1234 set((s) => {
1235 const next = {
1236 workspaces: [oneOffWorkspace, ...s.workspaces],
1237 };
1238 return canNavigate() ? { ...next, selectedWorkspaceId: oneOffWorkspace.id } : next;
1239 });
1240 ensureWorkspaceRuntime(get, set, oneOffWorkspace.id);
1241 } catch (error) {
1242 if (isOperationAbortError(error)) {
1243 const oneOffWorkspace = await oneOffWorkspacePromise.catch(() => null);
1244 if (oneOffWorkspace) {
1245 await discardCancelledOneOffWorkspace(oneOffWorkspace);
1246 }
1247 return false;
1248 }
1249 set((s) => ({
1250 notifications: pushNotification(s.notifications, {
1251 id: makeId(),
1252 ts: nowIso(),
1253 kind: "error",
1254 title: "Unable to create chat",
1255 detail: error instanceof Error ? error.message : String(error),
1256 }),
1257 }));
1258 return false;
1259 }
1260 }
1261 } else {
1262 workspaceId =
1263 opts?.workspaceId ??
1264 (get().selectedWorkspaceId &&
1265 !isOneOffChatWorkspace(
1266 get().workspaces.find((workspace) => workspace.id === get().selectedWorkspaceId),
1267 )
1268 ? get().selectedWorkspaceId
1269 : null) ??
1270 projectWorkspaces()[0]?.id ??
1271 null;
1272
1273 if (!workspaceId) {
1274 if (get().desktopFeatureFlags.workspaceLifecycle === false) {
1275 set((s) => ({
1276 notifications: pushNotification(s.notifications, {
1277 id: makeId(),
1278 ts: nowIso(),
1279 kind: "info",
1280 title: "Workspace management is disabled",
1281 detail:
1282 "Enable Workspace lifecycle actions in Settings -> Feature Flags to add a project workspace.",
1283 }),
1284 }));
1285 return false;
1286 }
1287 await get().addWorkspace({ intent: operationIntent });
1288 workspaceId = projectWorkspaces()[0]?.id ?? null;
1289 if (!workspaceId) return false;
1290 }
1291
1292 if (canNavigate() && get().selectedWorkspaceId !== workspaceId) {
1293 set({ selectedWorkspaceId: workspaceId });
1294 }
1295
1296 if (!createSessionImmediately) {
1297 const existingDraft = get().threads.find(
1298 (thread) => thread.workspaceId === workspaceId && thread.draft === true,
1299 );
1300 if (existingDraft) {
1301 if (canNavigate()) {
1302 set({
1303 selectedThreadId: existingDraft.id,
1304 selectedTaskId: null,
1305 view: "chat",
1306 newChatLandingTarget: null,
1307 });
1308 }
1309 ensureThreadRuntime(get, set, existingDraft.id);
1310 await persistNow(get);
1311 return true;
1312 }
1313 }
1314 }
1315
1316 if (!workspaceId) return false;
1317
1318 const threadId = makeId();
1319 let firstMessage = opts?.firstMessage ?? "";
1320 let resolvedAttachments = opts?.attachments;
1321 const draftAttachments =
1322 opts?.draftAttachments ?? opts?.attachmentFiles?.map(createComposerAttachmentFile) ?? [];
1323 // Attachment preparation may need the workspace server (large files fall
1324 // back to a JSON-RPC upload), so only text-only first messages can be
1325 // rendered optimistically before the server start wait.
1326 const needsAttachmentPreparation = draftAttachments.length > 0;
1327 const previousLandingTarget = get().newChatLandingTarget;
1328 let queuedDraftSubmission = opts?.draftSubmission;
1329 let rollbackDraftRekey: {
1330 previousKey: string;
1331 draft: ComposerDraft;
1332 submission?: ComposerSubmission;
1333 } | null = null;
1334
1335 // The thread record, navigation, and draft re-key are all local store
1336 // operations — no server RPC — so they run before the (potentially slow)
1337 // server start to give immediate feedback. Any later failure rolls the
1338 // local state back to the exact pre-send draft.
1339 const createLocalThreadRecord = async (): Promise<void> => {
1340 reportPhase("creating");
1341 const createdAt = nowIso();
1342 const title = opts?.titleHint ? truncateTitle(opts.titleHint) : "New chat";
1343
1344 const thread: ThreadRecord = {
1345 id: threadId,
1346 workspaceId,
1347 title,
1348 titleSource: "default",
1349 createdAt,
1350 lastMessageAt: createdAt,
1351 status: "active",
1352 sessionId: null,
1353 messageCount: 0,
1354 lastEventSeq: 0,
1355 draft: !createSessionImmediately,
1356 ...(opts?.reasoningEffort ? { reasoningEffort: opts.reasoningEffort } : {}),
1357 };
1358
1359 set((s) => {
1360 let composerDraftsByKey = s.composerDraftsByKey;
1361 let composerSubmissionsByKey = s.composerSubmissionsByKey;
1362 if (queuedDraftSubmission) {
1363 const previousKey = queuedDraftSubmission.key;
1364 const submittedDraft = composerDraftsByKey[previousKey];
1365 if (submittedDraft?.revision === queuedDraftSubmission.revision) {
1366 const nextKey = composerDraftKeyForThread(threadId);
1367 const nextDrafts = { ...composerDraftsByKey };
1368 delete nextDrafts[previousKey];
1369 nextDrafts[nextKey] = submittedDraft;
1370 composerDraftsByKey = nextDrafts;
1371 queuedDraftSubmission = {
1372 key: nextKey,
1373 revision: queuedDraftSubmission.revision,
1374 ...(queuedDraftSubmission.submissionId
1375 ? { submissionId: queuedDraftSubmission.submissionId }
1376 : {}),
1377 };
1378 const submission = s.composerSubmissionsByKey[previousKey];
1379 rollbackDraftRekey = { previousKey, draft: submittedDraft };
1380 if (submission?.id === queuedDraftSubmission.submissionId) {
1381 rollbackDraftRekey = { previousKey, draft: submittedDraft, submission };
1382 composerSubmissionsByKey = { ...s.composerSubmissionsByKey };
1383 delete composerSubmissionsByKey[previousKey];
1384 composerSubmissionsByKey[nextKey] = {
1385 ...submission,
1386 owner: queuedDraftSubmission,
1387 request: { kind: "thread", threadId },
1388 };
1389 }
1390 }
1391 }
1392 const next = {
1393 threads: [thread, ...s.threads],
1394 composerDraftsByKey,
1395 composerSubmissionsByKey,
1396 };
1397 return canNavigate()
1398 ? {
1399 ...next,
1400 selectedWorkspaceId: workspaceId,
1401 selectedThreadId: threadId,
1402 selectedTaskId: null,
1403 view: "chat" as const,
1404 newChatLandingTarget: null,
1405 }
1406 : next;
1407 });
1408 ensureThreadRuntime(get, set, threadId);
1409 set((s) => ({
1410 threadRuntimeById: {
1411 ...s.threadRuntimeById,
1412 [threadId]: {
1413 ...s.threadRuntimeById[threadId],
1414 transcriptOnly: false,
1415 draftComposerProvider: opts?.provider ?? null,
1416 draftComposerModel: opts?.model?.trim() || null,
1417 composerReasoningEffort: opts?.reasoningEffort ?? null,
1418 },
1419 },
1420 }));
1421 await persistNow(get);
1422 };
1423
1424 const rollbackCreatedThread = async (): Promise<void> => {
1425 if (!get().threads.some((candidate) => candidate.id === threadId)) return;
1426 RUNTIME.optimisticUserMessageIds.delete(threadId);
1427 RUNTIME.pendingThreadMessages.delete(threadId);
1428 RUNTIME.pendingThreadAttachments.delete(threadId);
1429 RUNTIME.pendingThreadReferences.delete(threadId);
1430 RUNTIME.pendingWorkspaceDefaultApplyByThread.delete(threadId);
1431 RUNTIME.modelStreamByThread.delete(threadId);
1432 const rekey = rollbackDraftRekey;
1433 rollbackDraftRekey = null;
1434 set((s) => {
1435 const threadRuntimeById = { ...s.threadRuntimeById };
1436 delete threadRuntimeById[threadId];
1437 let composerDraftsByKey = s.composerDraftsByKey;
1438 let composerSubmissionsByKey = s.composerSubmissionsByKey;
1439 if (rekey) {
1440 const threadDraftKey = composerDraftKeyForThread(threadId);
1441 const nextDrafts = { ...composerDraftsByKey };
1442 delete nextDrafts[threadDraftKey];
1443 // Keep any edits made while the send was in flight; otherwise this
1444 // is the exact submitted draft object moved back untouched.
1445 nextDrafts[rekey.previousKey] = composerDraftsByKey[threadDraftKey] ?? rekey.draft;
1446 composerDraftsByKey = nextDrafts;
1447 if (rekey.submission) {
1448 const nextSubmissions = { ...composerSubmissionsByKey };
1449 delete nextSubmissions[threadDraftKey];
1450 nextSubmissions[rekey.previousKey] = rekey.submission;
1451 composerSubmissionsByKey = nextSubmissions;
1452 }
1453 }
1454 return {
1455 threads: s.threads.filter((candidate) => candidate.id !== threadId),
1456 threadRuntimeById,
1457 composerDraftsByKey,
1458 composerSubmissionsByKey,
1459 ...(s.selectedThreadId === threadId
1460 ? {
1461 selectedThreadId: null,
1462 selectedTaskId: null,
1463 newChatLandingTarget: s.newChatLandingTarget ?? previousLandingTarget,
1464 }
1465 : {}),
1466 };
1467 });
1468 await persistNow(get);
1469 };
1470
1471 const discardFailedThreadAndWorkspace = async (): Promise<void> => {
1472 await rollbackCreatedThread();
1473 if (createdOneOffWorkspace) {
1474 await discardCancelledOneOffWorkspace(createdOneOffWorkspace);
1475 }
1476 };
1477
1478 const queueFirstMessageOptimistically = (): void => {
1479 if (queuedDraftSubmission?.submissionId) {
1480 const prepared = {
1481 text: firstMessage,
1482 attachments:
1483 resolvedAttachments && resolvedAttachments.length > 0
1484 ? resolvedAttachments
1485 : undefined,
1486 };
1487 updateSubmission(queuedDraftSubmission.submissionId, (submission) => ({
1488 ...submission,
1489 prepared: submission.prepared ?? prepared,
1490 }));
1491 }
1492 markSubmissionSending(queuedDraftSubmission?.submissionId);
1493 const hasFirstMessage = Boolean(firstMessage.trim());
1494 const hasResolvedAttachments = Boolean(
1495 resolvedAttachments && resolvedAttachments.length > 0,
1496 );
1497 if (hasFirstMessage || hasResolvedAttachments) {
1498 queueOptimisticFirstThreadMessage(
1499 set,
1500 threadId,
1501 firstMessage,
1502 resolvedAttachments,
1503 opts?.references,
1504 queuedDraftSubmission,
1505 opts?.clientMessageId,
1506 );
1507 recordThreadNavigationIntent(threadId, operationIntent);
1508 }
1509 };
1510
1511 if (createSessionImmediately) {
1512 await createLocalThreadRecord();
1513 if (!needsAttachmentPreparation) {
1514 queueFirstMessageOptimistically();
1515 }
1516 }
1517
1518 let url: string | null = null;
1519 if (createSessionImmediately) {
1520 reportPhase("starting-server");
1521 try {
1522 await ensureServerRunning(get, set, workspaceId, { signal: opts?.signal });
1523 } catch (error) {
1524 await discardFailedThreadAndWorkspace();
1525 if (isOperationAbortError(error)) return false;
1526 throw error;
1527 }
1528 ensureControlSocket(get, set, workspaceId);
1529
1530 const wsRt = get().workspaceRuntimeById[workspaceId];
1531 url = wsRt?.serverUrl ?? null;
1532 if (!url) {
1533 set((s) => ({
1534 notifications: pushNotification(s.notifications, {
1535 id: makeId(),
1536 ts: nowIso(),
1537 kind: "error",
1538 title: "Unable to create session",
1539 detail: wsRt?.error ?? "Workspace server is not ready.",
1540 }),
1541 }));
1542 await discardFailedThreadAndWorkspace();
1543 return false;
1544 }
1545 }
1546
1547 if (needsAttachmentPreparation) {
1548 reportPhase("processing-attachments");
1549 try {
1550 const prepared = await prepareComposerMessageForWorkspace(
1551 get,
1552 set,
1553 workspaceId,
1554 firstMessage,
1555 draftAttachments,
1556 { threadId, signal: opts?.signal },
1557 );
1558 resolvedAttachments = prepared.attachments;
1559 firstMessage = prepared.text;
1560 } catch (error) {
1561 await discardFailedThreadAndWorkspace();
1562 if (isOperationAbortError(error)) return false;
1563 throw error;
1564 }
1565 }
1566 if (opts?.signal?.aborted) {
1567 await discardFailedThreadAndWorkspace();
1568 return false;
1569 }
1570
1571 if (!createSessionImmediately) {
1572 await createLocalThreadRecord();
1573 return true;
1574 }
1575
1576 if (!url) {
1577 return false;
1578 }
1579
1580 if (needsAttachmentPreparation) {
1581 queueFirstMessageOptimistically();
1582 }
1583 ensureThreadSocket(
1584 get,
1585 set,
1586 threadId,
1587 url,
1588 firstMessage,
1589 Boolean(firstMessage.trim()),
1590 resolvedAttachments,
1591 );
1592 return true;
1593 },
1594
1595 openNewChatLanding: async (opts?: {
1596 defaultTargetKind?: "project" | "oneOff";
1597 target?: NewChatLandingTarget;
1598 }) => {
1599 invalidateNavigationIntent();
1600 const state = get();
1601 const landingTarget: NewChatLandingTarget =
1602 opts?.target ??
1603 (opts?.defaultTargetKind === "oneOff"
1604 ? { kind: "oneOff" }
1605 : resolveDefaultNewChatTarget(state.workspaces, state.selectedWorkspaceId));
1606 set({
1607 selectedThreadId: null,
1608 selectedTaskId: null,
1609 view: "chat",
1610 newChatLandingTarget: landingTarget,
1611 });
1612 syncDesktopStateCache(get);
1613 await persistNow(get);
1614 },
1615
1616 setNewChatLandingTarget: (target) => {
1617 invalidateNavigationIntent();
1618 set({ newChatLandingTarget: target });
1619 syncDesktopStateCache(get);
1620 },
1621
1622 selectThread: async (threadId: string, options = {}) => {
1623 if (options.signal?.aborted) return;
1624 invalidateNavigationIntent();
1625 set({ newChatLandingTarget: null });
1626 if (options.signal?.aborted) return;
1627 await hydrateThreadSelection(get, set, threadId, {
1628 reconnectAfterHydration: true,
1629 skipWorkspaceSelectOnReconnect: true,
1630 signal: options.signal,
1631 });
1632 },
1633
1634 reconnectThread: async (
1635 threadId: string,
1636 firstMessage?: string,
1637 opts?: {
1638 selectionRequestId?: number;
1639 skipWorkspaceSelect?: boolean;
1640 attachments?: import("../store.helpers/jsonRpcSocket").FileAttachmentInput[];
1641 references?: import("../../lib/wsProtocol").TurnReference[];
1642 refreshSnapshot?: boolean;
1643 signal?: AbortSignal;
1644 draftSubmission?: ComposerDraftRevision;
1645 clientMessageId?: string;
1646 },
1647 ) => {
1648 const isReconnectCurrent = () =>
1649 opts?.signal?.aborted !== true &&
1650 (opts?.selectionRequestId === undefined ||
1651 (get().selectedThreadId === threadId &&
1652 isCurrentThreadSelectionRequest(threadId, opts.selectionRequestId)));
1653
1654 if (!isReconnectCurrent()) return false;
1655 ensureThreadRuntime(get, set, threadId);
1656
1657 const thread = get().threads.find((t) => t.id === threadId);
1658 if (!thread) return false;
1659
1660 const hasQueuedAttachments = opts?.attachments && opts.attachments.length > 0;
1661 if (thread.draft && !firstMessage?.trim() && !hasQueuedAttachments) {
1662 return false;
1663 }
1664
1665 if (!opts?.skipWorkspaceSelect) {
1666 await get().selectWorkspace(thread.workspaceId, { signal: opts?.signal });
1667 if (!isReconnectCurrent()) return false;
1668 }
1669 await ensureServerRunning(get, set, thread.workspaceId, { signal: opts?.signal });
1670 if (!isReconnectCurrent()) return false;
1671 ensureControlSocket(get, set, thread.workspaceId);
1672
1673 const url = get().workspaceRuntimeById[thread.workspaceId]?.serverUrl;
1674 if (!url) {
1675 if (!isReconnectCurrent()) return false;
1676 set((s) => ({
1677 notifications: pushNotification(s.notifications, {
1678 id: makeId(),
1679 ts: nowIso(),
1680 kind: "error",
1681 title: "Workspace server unavailable",
1682 detail: "Workspace server is not ready.",
1683 }),
1684 }));
1685 return false;
1686 }
1687 if (!isReconnectCurrent()) return false;
1688
1689 const hasFirstMessage = firstMessage?.trim();
1690 if (hasFirstMessage || hasQueuedAttachments) {
1691 if (!isReconnectCurrent()) return false;
1692 queuePendingThreadMessage(
1693 threadId,
1694 firstMessage ?? "",
1695 opts?.attachments,
1696 opts?.references,
1697 opts?.clientMessageId,
1698 opts?.draftSubmission,
1699 );
1700 }
1701 if (!isReconnectCurrent()) return false;
1702 ensureThreadSocket(
1703 get,
1704 set,
1705 threadId,
1706 url,
1707 firstMessage,
1708 Boolean(firstMessage?.trim()),
1709 opts?.attachments,
1710 opts?.refreshSnapshot !== undefined ? { refreshSnapshot: opts.refreshSnapshot } : undefined,
1711 );
1712 return true;
1713 },
1714
1715 reconnectThreadWithFeedback: async (threadId: string) =>
1716 await runAcknowledgedOperation(get, set, {
1717 key: operationKey("thread-reconnect", threadId),
1718 label: "Reconnect chat",
1719 errorTitle: "Chat did not reconnect",
1720 errorMessage: "Cowork could not reconnect this chat.",
1721 repairAction: "Your draft is safe. Check the workspace connection and retry.",
1722 execute: async () => {
1723 const accepted = await get().reconnectThread(threadId);
1724 if (!accepted) {
1725 throw new Error("Cowork could not start a reconnect attempt for this chat.");
1726 }
1727 const startedAt = Date.now();
1728 while (Date.now() - startedAt < RECONNECT_OUTCOME_TIMEOUT_MS) {
1729 const state = get();
1730 if (!state.threads.some((thread) => thread.id === threadId)) {
1731 throw new Error("This chat is no longer available.");
1732 }
1733 if (state.threadRuntimeById[threadId]?.connected === true) {
1734 return;
1735 }
1736 await new Promise<void>((resolve) => {
1737 setTimeout(resolve, RECONNECT_OUTCOME_POLL_MS);
1738 });
1739 }
1740 throw new Error("Cowork could not reconnect this chat before the request timed out.");
1741 },
1742 }),
1743
1744 sendMessage: async (
1745 text: string,
1746 busyPolicy: ThreadBusyPolicy = "reject",
1747 attachments?: import("../store.helpers/jsonRpcSocket").FileAttachmentInput[],
1748 references?: import("../../lib/wsProtocol").TurnReference[],
1749 options?: {
1750 targetThreadId?: string;
1751 draftSubmission?: ComposerDraftRevision;
1752 clientMessageId?: string;
1753 retryToolItemIds?: string[];
1754 },
1755 ): Promise<boolean> => {
1756 const activeThreadId = options?.targetThreadId ?? get().selectedThreadId;
1757 if (!activeThreadId) return false;
1758
1759 const thread = get().threads.find((t) => t.id === activeThreadId);
1760 if (!thread) return false;
1761 const threadDraftKey = composerDraftKeyForThread(activeThreadId);
1762 const currentDraft = get().composerDraftsByKey[threadDraftKey];
1763 const draftSubmission =
1764 options?.draftSubmission ??
1765 (currentDraft && currentDraft.text.trim() === text.trim()
1766 ? { key: threadDraftKey, revision: currentDraft.revision }
1767 : undefined);
1768
1769 if (!(thread.workspaceId in get().taskSummariesByWorkspaceId)) {
1770 // Warm task summaries without blocking the send; the server enforces
1771 // task locks authoritatively and the UI catches up when this resolves.
1772 void get().refreshTasks(thread.workspaceId);
1773 }
1774
1775 const rt = get().threadRuntimeById[activeThreadId];
1776 const trimmed = text.trim();
1777 const hasAttachments = attachments && attachments.length > 0;
1778 if (!trimmed && !hasAttachments) return false;
1779
1780 const taskCommand = !hasAttachments ? trimmed.match(/^\/task(?:\s+([\s\S]*))?$/i) : null;
1781 if (taskCommand) {
1782 try {
1783 recordThreadNavigationIntent(activeThreadId);
1784 await ensureServerRunning(get, set, thread.workspaceId);
1785 ensureControlSocket(get, set, thread.workspaceId);
1786 await requestJsonRpc(get, set, thread.workspaceId, "command/execute", {
1787 threadId: activeThreadId,
1788 name: "task",
1789 arguments: taskCommand[1]?.trim() ?? "",
1790 clientMessageId: options?.clientMessageId ?? makeId(),
1791 });
1792 if (draftSubmission) get().completeComposerSubmission(draftSubmission);
1793 return true;
1794 } catch (error) {
1795 const detail = error instanceof Error ? error.message : String(error);
1796 set((state) => ({
1797 notifications: pushNotification(state.notifications, {
1798 id: makeId(),
1799 ts: nowIso(),
1800 kind: "error",
1801 title: "Unable to start task mode",
1802 detail,
1803 }),
1804 }));
1805 return false;
1806 }
1807 }
1808
1809 if (rt?.transcriptOnly) {
1810 const preamble = get().injectContext ? buildContextPreamble(rt?.feed ?? []) : "";
1811 const firstMessage = preamble ? `${preamble}${trimmed}` : trimmed;
1812 const workspace = get().workspaces.find((candidate) => candidate.id === thread.workspaceId);
1813 const started = await get().newThread({
1814 workspaceId: thread.workspaceId,
1815 scope: isOneOffChatWorkspace(workspace) ? "oneOff" : "project",
1816 titleHint: thread.title,
1817 firstMessage,
1818 attachments,
1819 references,
1820 draftSubmission,
1821 clientMessageId: options?.clientMessageId,
1822 });
1823 if (!started) return false;
1824 return true;
1825 }
1826
1827 if (thread.status !== "active" || !rt?.sessionId) {
1828 const preamble = get().injectContext ? buildContextPreamble(rt?.feed ?? []) : "";
1829 const firstMessage = preamble ? `${preamble}${trimmed}` : trimmed;
1830 const reconnected = await get().reconnectThread(activeThreadId, firstMessage, {
1831 attachments,
1832 references,
1833 draftSubmission,
1834 clientMessageId: options?.clientMessageId,
1835 });
1836 if (!reconnected) return false;
1837 recordThreadNavigationIntent(activeThreadId);
1838 return true;
1839 }
1840
1841 const accepted = sendUserMessageToThread(
1842 get,
1843 set,
1844 activeThreadId,
1845 trimmed,
1846 busyPolicy,
1847 attachments,
1848 references,
1849 options?.clientMessageId,
1850 draftSubmission,
1851 options?.retryToolItemIds,
1852 );
1853 if (!accepted) return false;
1854 recordThreadNavigationIntent(activeThreadId);
1855 return true;
1856 },
1857
1858 submitComposerDraft: (request, control) => {
1859 const key =
1860 request.kind === "thread"
1861 ? composerDraftKeyForThread(request.threadId)
1862 : composerDraftKeyForNewChatTarget(request.target);
1863 const state = get();
1864 if ((state.composerAttachmentIngestionCountByKey[key] ?? 0) > 0) return false;
1865 const draft = state.composerDraftsByKey[key];
1866 if (!draft || (!draft.text.trim() && draft.attachments.length === 0)) return false;
1867 const existing = state.composerSubmissionsByKey[key];
1868 if (isComposerSubmissionInFlight(existing)) return false;
1869
1870 const submissionId = makeId();
1871 const owner: ComposerDraftRevision = {
1872 key,
1873 revision: draft.revision,
1874 submissionId,
1875 };
1876 const submission: ComposerSubmission = {
1877 id: submissionId,
1878 clientMessageId: makeId(),
1879 owner,
1880 request,
1881 draft: cloneComposerDraftForSubmission(draft),
1882 prepared: null,
1883 phase: "preparing",
1884 delivery: currentSubmissionDelivery(request),
1885 error: null,
1886 };
1887 let claimed = false;
1888 set((current) => {
1889 if ((current.composerAttachmentIngestionCountByKey[key] ?? 0) > 0) return {};
1890 if (isComposerSubmissionInFlight(current.composerSubmissionsByKey[key])) return {};
1891 const currentDraft = current.composerDraftsByKey[key];
1892 if (!currentDraft || currentDraft.revision !== owner.revision) return {};
1893 claimed = true;
1894 return {
1895 composerSubmissionsByKey: {
1896 ...current.composerSubmissionsByKey,
1897 [key]: submission,
1898 },
1899 };
1900 });
1901 if (!claimed) return false;
1902 if (existing) revokeAttachmentPreviewsIfUnreferenced(existing.draft.attachments);
1903 if (control) {
1904 composerSubmissionCreationControl.set(submissionId, control);
1905 }
1906 void executeComposerSubmission(submissionId);
1907 return true;
1908 },
1909
1910 retryComposerSubmission: (key, control) => {
1911 const ownerKey = key ?? resolveActiveComposerDraftKey(get());
1912 const submission = get().composerSubmissionsByKey[ownerKey];
1913 if (submission?.phase !== "failed") return false;
1914 let claimed = false;
1915 set((state) => {
1916 const current = state.composerSubmissionsByKey[ownerKey];
1917 if (!current || current.id !== submission.id || current.phase !== "failed") return {};
1918 claimed = true;
1919 return {
1920 composerSubmissionsByKey: {
1921 ...state.composerSubmissionsByKey,
1922 [ownerKey]: {
1923 ...current,
1924 phase: "preparing",
1925 delivery: currentSubmissionDelivery(current.request),
1926 error: null,
1927 },
1928 },
1929 };
1930 });
1931 if (!claimed) return false;
1932 if (control) {
1933 composerSubmissionCreationControl.set(submission.id, control);
1934 }
1935 void executeComposerSubmission(submission.id);
1936 return true;
1937 },
1938
1939 cancelComposerSubmission: (key) => {
1940 const ownerKey = key ?? resolveActiveComposerDraftKey(get());
1941 const submission = get().composerSubmissionsByKey[ownerKey];
1942 if (!isComposerSubmissionInFlight(submission)) return false;
1943 let cancelled = false;
1944 set((state) => {
1945 const current = state.composerSubmissionsByKey[ownerKey];
1946 if (!current || current.id !== submission.id || !isComposerSubmissionInFlight(current)) {
1947 return {};
1948 }
1949 const next = { ...state.composerSubmissionsByKey };
1950 delete next[ownerKey];
1951 cancelled = true;
1952 return { composerSubmissionsByKey: next };
1953 });
1954 return cancelled;
1955 },
1956
1957 editAcceptedComposerSubmission: (key) => {
1958 const ownerKey = key ?? resolveActiveComposerDraftKey(get());
1959 const submission = get().composerSubmissionsByKey[ownerKey];
1960 if (submission?.phase !== "accepted") return false;
1961 let restored = false;
1962 set((state) => {
1963 const currentSubmission = state.composerSubmissionsByKey[ownerKey];
1964 if (
1965 !currentSubmission ||
1966 currentSubmission.id !== submission.id ||
1967 currentSubmission.phase !== "accepted"
1968 ) {
1969 return {};
1970 }
1971 const current = state.composerDraftsByKey[ownerKey];
1972 if (hasComposerDraftState(current)) {
1973 return {};
1974 }
1975 const nextSubmissions = { ...state.composerSubmissionsByKey };
1976 delete nextSubmissions[ownerKey];
1977 restored = true;
1978 return {
1979 composerDraftsByKey: {
1980 ...state.composerDraftsByKey,
1981 [ownerKey]: {
1982 ...cloneComposerDraftForSubmission(currentSubmission.draft),
1983 revision: Math.max(current?.revision ?? 0, currentSubmission.owner.revision) + 1,
1984 generation: Math.max(current?.generation ?? 0, currentSubmission.draft.generation),
1985 updatedAt: nowIso(),
1986 },
1987 },
1988 composerSubmissionsByKey: nextSubmissions,
1989 };
1990 });
1991 if (restored) persist(get);
1992 return restored;
1993 },
1994
1995 dismissComposerSubmission: (key) => {
1996 const ownerKey = key ?? resolveActiveComposerDraftKey(get());
1997 const submission = get().composerSubmissionsByKey[ownerKey];
1998 if (!submission || isComposerSubmissionInFlight(submission)) return;
1999 set((state) => {
2000 if (state.composerSubmissionsByKey[ownerKey]?.id !== submission.id) return {};
2001 const next = { ...state.composerSubmissionsByKey };
2002 delete next[ownerKey];
2003 return { composerSubmissionsByKey: next };
2004 });
2005 revokeAttachmentPreviewsIfUnreferenced(submission.draft.attachments);
2006 },
2007
2008 completeComposerSubmission: (owner) => {
2009 if (!owner.submissionId) {
2010 get().clearComposerDraft(owner);
2011 return;
2012 }
2013 const entry = submissionEntry(owner.submissionId);
2014 if (!entry) return;
2015 const { key, submission } = entry;
2016 if (submission.delivery === "steer") {
2017 set((state) => {
2018 const current = state.composerSubmissionsByKey[key];
2019 if (!current || current.id !== submission.id) return {};
2020 const cleared = clearComposerDraftRevision(state.composerDraftsByKey, current.owner);
2021 return {
2022 ...(cleared.cleared ? { composerDraftsByKey: cleared.drafts } : {}),
2023 composerSubmissionsByKey: {
2024 ...state.composerSubmissionsByKey,
2025 [key]: {
2026 ...current,
2027 phase: "accepted",
2028 error: null,
2029 },
2030 },
2031 };
2032 });
2033 persist(get);
2034 return;
2035 }
2036
2037 set((state) => {
2038 const current = state.composerSubmissionsByKey[key];
2039 if (!current || current.id !== submission.id) return {};
2040 const next = { ...state.composerSubmissionsByKey };
2041 delete next[key];
2042 return { composerSubmissionsByKey: next };
2043 });
2044 const cleared = get().clearComposerDraft(submission.owner);
2045 if (!cleared) revokeAttachmentPreviewsIfUnreferenced(submission.draft.attachments);
2046 },
2047
2048 failComposerSubmission: (submissionId, error) => {
2049 failSubmission(submissionId, error);
2050 },
2051
2052 cancelThread: (threadId: string, opts?: { includeSubagents?: boolean }) => {
2053 let claimed = false;
2054 set((state) => {
2055 const runtime = state.threadRuntimeById[threadId];
2056 if (!runtime?.busy || runtime.interruptPending) return {};
2057 claimed = true;
2058 return {
2059 threadRuntimeById: {
2060 ...state.threadRuntimeById,
2061 [threadId]: {
2062 ...runtime,
2063 interruptPending: true,
2064 },
2065 },
2066 };
2067 });
2068 if (!claimed) return false;
2069
2070 const ok = sendThread(
2071 get,
2072 threadId,
2073 (sid) => ({
2074 type: "cancel",
2075 sessionId: sid,
2076 ...(opts?.includeSubagents !== undefined
2077 ? { includeSubagents: opts.includeSubagents }
2078 : {}),
2079 }),
2080 {
2081 onSettled: (error) => {
2082 if (!error) return;
2083 set((state) => {
2084 const runtime = state.threadRuntimeById[threadId];
2085 if (!runtime?.interruptPending) return {};
2086 return {
2087 threadRuntimeById: {
2088 ...state.threadRuntimeById,
2089 [threadId]: {
2090 ...runtime,
2091 interruptPending: false,
2092 },
2093 },
2094 notifications: pushNotification(state.notifications, {
2095 id: makeId(),
2096 ts: nowIso(),
2097 kind: "error",
2098 title: "Unable to stop response",
2099 detail: composerSubmissionErrorMessage(error),
2100 }),
2101 };
2102 });
2103 },
2104 },
2105 );
2106 if (!ok) {
2107 set((s) => ({
2108 threadRuntimeById: {
2109 ...s.threadRuntimeById,
2110 [threadId]: {
2111 ...s.threadRuntimeById[threadId],
2112 interruptPending: false,
2113 },
2114 },
2115 notifications: pushNotification(s.notifications, {
2116 id: makeId(),
2117 ts: nowIso(),
2118 kind: "error",
2119 title: "Not connected",
2120 detail: "Unable to cancel this run.",
2121 }),
2122 }));
2123 return false;
2124 }
2125 return true;
2126 },
2127
2128 clearThreadUsageHardCap: (threadId: string) => {
2129 const ok = sendThread(get, threadId, (sessionId) => ({
2130 type: "set_session_usage_budget",
2131 sessionId,
2132 stopAtUsd: null,
2133 }));
2134 if (!ok) {
2135 set((s) => ({
2136 notifications: pushNotification(s.notifications, {
2137 id: makeId(),
2138 ts: nowIso(),
2139 kind: "error",
2140 title: "Not connected",
2141 detail: "Unable to clear the session hard cap.",
2142 }),
2143 }));
2144 return;
2145 }
2146
2147 appendThreadTranscript(threadId, "client", {
2148 type: "set_session_usage_budget",
2149 sessionId: get().threadRuntimeById[threadId]?.sessionId,
2150 stopAtUsd: null,
2151 });
2152 },
2153
2154 setThreadModel: (threadId, provider, model) => {
2155 const thread = get().threads.find((t) => t.id === threadId);
2156 if (!thread) return;
2157
2158 if (thread.draft) {
2159 set((s) => ({
2160 threads: s.threads.map((candidate) =>
2161 candidate.id === threadId ? { ...candidate, reasoningEffort: undefined } : candidate,
2162 ),
2163 threadRuntimeById: {
2164 ...s.threadRuntimeById,
2165 [threadId]: {
2166 ...s.threadRuntimeById[threadId],
2167 draftComposerProvider: provider,
2168 draftComposerModel: model,
2169 composerReasoningEffort: null,
2170 },
2171 },
2172 }));
2173 persist(get);
2174 return;
2175 }
2176
2177 const rt = get().threadRuntimeById[threadId];
2178 if (!rt?.sessionId) return;
2179 set((state) => ({
2180 threads: state.threads.map((candidate) =>
2181 candidate.id === threadId ? { ...candidate, reasoningEffort: undefined } : candidate,
2182 ),
2183 threadRuntimeById: {
2184 ...state.threadRuntimeById,
2185 [threadId]: {
2186 ...state.threadRuntimeById[threadId],
2187 composerReasoningEffort: null,
2188 },
2189 },
2190 }));
2191 persist(get);
2192 const pendingApply = RUNTIME.pendingWorkspaceDefaultApplyByThread.get(threadId);
2193 if (pendingApply?.draftModelSelection) {
2194 RUNTIME.pendingWorkspaceDefaultApplyByThread.set(threadId, {
2195 ...pendingApply,
2196 draftModelSelection: null,
2197 });
2198 }
2199 const ok = sendThread(get, threadId, (sessionId) => ({
2200 type: "set_model",
2201 sessionId,
2202 provider,
2203 model,
2204 }));
2205 if (ok) {
2206 appendThreadTranscript(threadId, "client", {
2207 type: "set_model",
2208 sessionId: rt.sessionId,
2209 provider,
2210 model,
2211 });
2212 }
2213 },
2214
2215 setThreadReasoningEffort: (threadId, provider, effort) => {
2216 const providerConfig =
2217 (provider === "openai" || provider === "codex-cli") && isOpenAiReasoningEffortValue(effort)
2218 ? { [provider]: { reasoningEffort: effort } }
2219 : provider === "google" && isGoogleReasoningEffortValue(effort)
2220 ? { google: googleProviderOptionsForReasoningEffort(effort) }
2221 : null;
2222 if (!providerConfig) return;
2223 const thread = get().threads.find((candidate) => candidate.id === threadId);
2224 if (!thread) return;
2225 ensureThreadRuntime(get, set, threadId);
2226 const currentRuntime = get().threadRuntimeById[threadId];
2227 if (!thread.draft && currentRuntime?.busy) return;
2228
2229 set((state) => ({
2230 threads: state.threads.map((candidate) =>
2231 candidate.id === threadId ? { ...candidate, reasoningEffort: effort } : candidate,
2232 ),
2233 threadRuntimeById: {
2234 ...state.threadRuntimeById,
2235 [threadId]: {
2236 ...state.threadRuntimeById[threadId],
2237 composerReasoningEffort: effort,
2238 },
2239 },
2240 }));
2241 persist(get);
2242
2243 if (thread.draft) return;
2244 const rt = currentRuntime;
2245 if (!rt?.sessionId) return;
2246 const config = {
2247 providerOptions: providerConfig,
2248 };
2249 const ok = sendThread(get, threadId, (sessionId) => ({
2250 type: "set_config",
2251 sessionId,
2252 config,
2253 }));
2254 if (ok) {
2255 appendThreadTranscript(threadId, "client", {
2256 type: "set_config",
2257 sessionId: rt.sessionId,
2258 config,
2259 });
2260 } else {
2261 // The change never left the client, so no session_config ack will
2262 // arrive to clear the optimistic value — revert it now so the selector
2263 // does not stay stuck on an effort the session never received.
2264 set((state) => {
2265 const current = state.threadRuntimeById[threadId];
2266 if (!current || current.composerReasoningEffort !== effort) return {};
2267 return {
2268 threadRuntimeById: {
2269 ...state.threadRuntimeById,
2270 [threadId]: { ...current, composerReasoningEffort: null },
2271 },
2272 };
2273 });
2274 }
2275 },
2276
2277 setComposerText: (text, references = []) => {
2278 set((state) => {
2279 const key = resolveActiveComposerDraftKey(state);
2280 const current =
2281 state.composerDraftsByKey[key] ?? createEmptyComposerDraftForState(state, key);
2282 return {
2283 composerDraftsByKey: {
2284 ...state.composerDraftsByKey,
2285 [key]: {
2286 ...current,
2287 revision: current.revision + 1,
2288 updatedAt: nowIso(),
2289 text,
2290 references: references.map((reference) => ({ ...reference })),
2291 },
2292 },
2293 };
2294 });
2295 get().pruneComposerDrafts(undefined, Number.POSITIVE_INFINITY);
2296 persist(get);
2297 },
2298
2299 addComposerAttachments: async (files) => {
2300 if (files.length === 0) return;
2301 const ownerKey = resolveActiveComposerDraftKey(get());
2302 let ownerGeneration = 1;
2303 set((state) => {
2304 const existing = state.composerDraftsByKey[ownerKey];
2305 const current = existing ?? createEmptyComposerDraftForState(state, ownerKey);
2306 ownerGeneration = Math.max(1, current.generation);
2307 const pendingCount = (state.composerAttachmentIngestionCountByKey[ownerKey] ?? 0) + 1;
2308 return {
2309 composerDraftsByKey: {
2310 ...state.composerDraftsByKey,
2311 [ownerKey]: {
2312 ...current,
2313 generation: ownerGeneration,
2314 },
2315 },
2316 composerAttachmentIngestionCountByKey: {
2317 ...state.composerAttachmentIngestionCountByKey,
2318 [ownerKey]: pendingCount,
2319 },
2320 };
2321 });
2322 const previousIngestion = RUNTIME.composerAttachmentIngestionTail ?? Promise.resolve();
2323 const ingestion = previousIngestion
2324 .catch(() => undefined)
2325 .then(async () => {
2326 const current = get().composerDraftsByKey[ownerKey];
2327 if (current?.generation !== ownerGeneration) return;
2328 const validationMessage = getComposerDraftAttachmentValidationMessage(
2329 get().composerDraftsByKey,
2330 ownerKey,
2331 files,
2332 );
2333 if (validationMessage) {
2334 throw new Error(validationMessage);
2335 }
2336
2337 const attachments: ComposerDraftAttachment[] = [];
2338 try {
2339 for (const file of files) {
2340 attachments.push(await createComposerDraftAttachment(file));
2341 }
2342 } catch (error) {
2343 revokeComposerDraftAttachmentPreviews(attachments);
2344 throw error;
2345 }
2346
2347 let accepted = false;
2348 set((state) => {
2349 const draft = state.composerDraftsByKey[ownerKey];
2350 if (draft?.generation !== ownerGeneration) return {};
2351 accepted = true;
2352 return {
2353 composerDraftsByKey: {
2354 ...state.composerDraftsByKey,
2355 [ownerKey]: {
2356 ...draft,
2357 revision: draft.revision + 1,
2358 updatedAt: nowIso(),
2359 attachments: [...draft.attachments, ...attachments],
2360 },
2361 },
2362 };
2363 });
2364 if (!accepted) {
2365 revokeComposerDraftAttachmentPreviews(attachments);
2366 return;
2367 }
2368 persist(get);
2369 });
2370 let trackedIngestion: Promise<void>;
2371 trackedIngestion = ingestion.finally(() => {
2372 set((state) => {
2373 const pendingCount = (state.composerAttachmentIngestionCountByKey[ownerKey] ?? 1) - 1;
2374 const nextPendingCounts = {
2375 ...state.composerAttachmentIngestionCountByKey,
2376 };
2377 if (pendingCount > 0) nextPendingCounts[ownerKey] = pendingCount;
2378 else delete nextPendingCounts[ownerKey];
2379 return { composerAttachmentIngestionCountByKey: nextPendingCounts };
2380 });
2381 if (RUNTIME.composerAttachmentIngestionTail === trackedIngestion) {
2382 RUNTIME.composerAttachmentIngestionTail = null;
2383 }
2384 get().pruneComposerDrafts(undefined, Number.POSITIVE_INFINITY);
2385 });
2386 RUNTIME.composerAttachmentIngestionTail = trackedIngestion;
2387 await trackedIngestion;
2388 },
2389
2390 removeComposerAttachment: (index) => {
2391 let removed: ComposerDraftAttachment[] = [];
2392 set((state) => {
2393 const key = resolveActiveComposerDraftKey(state);
2394 const current = state.composerDraftsByKey[key];
2395 if (!current?.attachments[index]) return {};
2396 removed = [current.attachments[index]];
2397 return {
2398 composerDraftsByKey: {
2399 ...state.composerDraftsByKey,
2400 [key]: {
2401 ...current,
2402 revision: current.revision + 1,
2403 updatedAt: nowIso(),
2404 attachments: current.attachments.filter(
2405 (_attachment, currentIndex) => currentIndex !== index,
2406 ),
2407 },
2408 },
2409 };
2410 });
2411 revokeAttachmentPreviewsIfUnreferenced(removed);
2412 get().pruneComposerDrafts(undefined, Number.POSITIVE_INFINITY);
2413 persist(get);
2414 },
2415
2416 setComposerDraftModel: (provider, model) => {
2417 const normalizedModel = model.trim();
2418 if (!normalizedModel) return;
2419 set((state) => {
2420 const key = resolveActiveComposerDraftKey(state);
2421 const current =
2422 state.composerDraftsByKey[key] ?? createEmptyComposerDraftForState(state, key);
2423 if (current.provider === provider && current.model === normalizedModel) return {};
2424 return {
2425 composerDraftsByKey: {
2426 ...state.composerDraftsByKey,
2427 [key]: {
2428 ...current,
2429 revision: current.revision + 1,
2430 updatedAt: nowIso(),
2431 provider,
2432 model: normalizedModel,
2433 reasoningEffort: null,
2434 },
2435 },
2436 };
2437 });
2438 get().pruneComposerDrafts(undefined, Number.POSITIVE_INFINITY);
2439 persist(get);
2440 },
2441
2442 setComposerDraftReasoningEffort: (effort) => {
2443 set((state) => {
2444 const key = resolveActiveComposerDraftKey(state);
2445 const current =
2446 state.composerDraftsByKey[key] ?? createEmptyComposerDraftForState(state, key);
2447 if (current.reasoningEffort === effort) return {};
2448 return {
2449 composerDraftsByKey: {
2450 ...state.composerDraftsByKey,
2451 [key]: {
2452 ...current,
2453 revision: current.revision + 1,
2454 updatedAt: nowIso(),
2455 reasoningEffort: effort,
2456 },
2457 },
2458 };
2459 });
2460 get().pruneComposerDrafts(undefined, Number.POSITIVE_INFINITY);
2461 persist(get);
2462 },
2463
2464 clearComposerDraft: (owner) => {
2465 let cleared = false;
2466 let removedAttachments: ReturnType<typeof clearComposerDraftRevision>["removedAttachments"] =
2467 [];
2468 set((state) => {
2469 const result = clearComposerDraftRevision(state.composerDraftsByKey, owner);
2470 cleared = result.cleared;
2471 removedAttachments = result.removedAttachments;
2472 return result.cleared ? { composerDraftsByKey: result.drafts } : {};
2473 });
2474 if (cleared) {
2475 revokeAttachmentPreviewsIfUnreferenced(removedAttachments);
2476 get().pruneComposerDrafts(undefined, Number.POSITIVE_INFINITY);
2477 persist(get);
2478 }
2479 return cleared;
2480 },
2481
2482 discardComposerDraft: (key) => {
2483 const ownerKey = key ?? resolveActiveComposerDraftKey(get());
2484 const current = get().composerDraftsByKey[ownerKey];
2485 if (!current) return false;
2486 const owner = { key: ownerKey, revision: current.revision };
2487 return get().clearComposerDraft(owner);
2488 },
2489
2490 pruneComposerDrafts: (nowMs, maxAgeMs) => {
2491 let removedAttachments: ReturnType<typeof pruneComposerDraftEntries>["removedAttachments"] =
2492 [];
2493 let changed = false;
2494 set((state) => {
2495 const result = pruneComposerDraftEntries(state.composerDraftsByKey, {
2496 nowMs,
2497 validThreadIds: new Set(state.threads.map((thread) => thread.id)),
2498 validProjectWorkspaceIds: new Set(
2499 state.workspaces
2500 .filter((workspace) => !isOneOffChatWorkspace(workspace))
2501 .map((workspace) => workspace.id),
2502 ),
2503 activeKey: resolveActiveComposerDraftKey(state),
2504 maxAgeMs,
2505 protectedKeys: new Set([
2506 ...Object.entries(state.composerAttachmentIngestionCountByKey)
2507 .filter(([, count]) => count > 0)
2508 .map(([key]) => key),
2509 ...Object.keys(state.composerSubmissionsByKey),
2510 ]),
2511 });
2512 removedAttachments = result.removedAttachments;
2513 changed = result.removedKeys.length > 0;
2514 if (result.removedKeys.length === 0) return {};
2515 const composerDraftRevisionFloorByKey = {
2516 ...state.composerDraftRevisionFloorByKey,
2517 };
2518 for (const key of result.removedKeys) {
2519 const removedDraft = state.composerDraftsByKey[key];
2520 if (!removedDraft) continue;
2521 const currentFloor = composerDraftRevisionFloorByKey[key];
2522 composerDraftRevisionFloorByKey[key] = {
2523 revision: Math.max(currentFloor?.revision ?? 0, removedDraft.revision),
2524 generation: Math.max(currentFloor?.generation ?? 0, removedDraft.generation),
2525 };
2526 }
2527 return {
2528 composerDraftsByKey: result.drafts,
2529 composerDraftRevisionFloorByKey,
2530 };
2531 });
2532 revokeAttachmentPreviewsIfUnreferenced(removedAttachments);
2533 if (changed) persist(get);
2534 },
2535
2536 setInjectContext: (v) => set({ injectContext: v }),
2537
2538 answerAsk: (threadId, requestId, answer) => {
2539 const interaction = get().interactionsByThread[threadId]?.find(
2540 (candidate) => candidate.requestId === requestId,
2541 );
2542 if (
2543 interaction?.kind !== "ask" ||
2544 (interaction.status !== "pending" && interaction.status !== "failed")
2545 ) {
2546 return false;
2547 }
2548 updateInteraction(set, threadId, requestId, (current) => {
2549 if (current.kind !== "ask") return current;
2550 const { error: _error, ...rest } = current;
2551 return { ...rest, status: "responding", response: answer };
2552 });
2553 const sent = sendThread(get, threadId, (sessionId) => ({
2554 type: "ask_response",
2555 sessionId,
2556 requestId,
2557 answer,
2558 }));
2559 if (!sent) {
2560 updateInteraction(set, threadId, requestId, (current) => ({
2561 ...current,
2562 status: "failed",
2563 error: "The response could not be sent. Reconnect and retry.",
2564 }));
2565 return false;
2566 }
2567 appendThreadTranscript(threadId, "client", {
2568 type: "ask_response",
2569 sessionId: get().threadRuntimeById[threadId]?.sessionId,
2570 requestId,
2571 answer,
2572 });
2573 return true;
2574 },
2575
2576 answerApproval: (threadId, requestId, approved) => {
2577 const interaction = get().interactionsByThread[threadId]?.find(
2578 (candidate) => candidate.requestId === requestId,
2579 );
2580 if (
2581 interaction?.kind !== "approval" ||
2582 (interaction.status !== "pending" && interaction.status !== "failed")
2583 ) {
2584 return false;
2585 }
2586 updateInteraction(set, threadId, requestId, (current) => {
2587 if (current.kind !== "approval") return current;
2588 const { error: _error, ...rest } = current;
2589 return { ...rest, status: "responding", response: approved };
2590 });
2591 const sent = sendThread(get, threadId, (sessionId) => ({
2592 type: "approval_response",
2593 sessionId,
2594 requestId,
2595 approved,
2596 }));
2597 if (!sent) {
2598 updateInteraction(set, threadId, requestId, (current) => ({
2599 ...current,
2600 status: "failed",
2601 error: "The response could not be sent. Reconnect and retry.",
2602 }));
2603 return false;
2604 }
2605 appendThreadTranscript(threadId, "client", {
2606 type: "approval_response",
2607 sessionId: get().threadRuntimeById[threadId]?.sessionId,
2608 requestId,
2609 approved,
2610 });
2611 return true;
2612 },
2613
2614 dismissPrompt: () => {
2615 const pending = findLatestVisibleSandboxInteraction(get());
2616 if (pending?.interaction.kind !== "approval") return;
2617 get().answerApproval(pending.threadId, pending.interaction.requestId, false);
2618 },
2619
2620 retryInteractionResponse: (threadId, requestId) => {
2621 const interaction = get().interactionsByThread[threadId]?.find(
2622 (candidate) => candidate.requestId === requestId,
2623 );
2624 if (interaction?.status !== "failed") return false;
2625 if (interaction.kind === "ask" && interaction.response !== undefined) {
2626 return get().answerAsk(threadId, requestId, interaction.response);
2627 }
2628 if (interaction.kind === "approval" && interaction.response !== undefined) {
2629 return get().answerApproval(threadId, requestId, interaction.response);
2630 }
2631 return false;
2632 },
2633
2634 loadAllThreadUsage: async () => {
2635 const threads = get().threads;
2636 const existing = get().threadRuntimeById;
2637
2638 await Promise.all(
2639 threads.map(async (thread) => {
2640 // Skip threads that already have usage loaded
2641 const rt = existing[thread.id];
2642 if (rt?.sessionUsage !== undefined && rt.sessionUsage !== null) return;
2643
2644 try {
2645 const transcript = await readTranscriptEvents(thread);
2646 if (!transcript) return; // No transcript available
2647 const usageState = extractUsageStateFromTranscript(transcript);
2648 if (!usageState.sessionUsage) return; // No usage in this thread
2649
2650 ensureThreadRuntime(get, set, thread.id);
2651 set((s) => ({
2652 threadRuntimeById: {
2653 ...s.threadRuntimeById,
2654 [thread.id]: {
2655 ...s.threadRuntimeById[thread.id],
2656 ...usageState,
2657 },
2658 },
2659 }));
2660 } catch {
2661 // Skip threads whose transcripts can't be read
2662 }
2663 }),
2664 );
2665 },
2666 };
2667}
2668