From 5c536b0295437adfe2bf9bd092fa7075a9d65ee1 Mon Sep 17 00:00:00 2001 From: macodev00 <273427913+macodev00@users.noreply.github.com> Date: Thu, 24 Sep 2026 08:55:57 +0000 Subject: [PATCH] fix(web): dispatch queued follow-ups after leaving the thread Queue behavior held follow-ups in the open chat view, so leaving the thread delayed thread.turn.start until it was selected again. An app-wide coordinator now sends them when a tool finishes or the turn ends. Stopping a claimed background send restores that prompt to the composer instead of parking it as a held queue row. Co-authored-by: maco --- .../components/BackgroundQueuedMessages.tsx | 135 +++++++ apps/web/src/components/ChatView.tsx | 29 +- .../lib/backgroundQueueActiveThread.test.ts | 31 ++ .../src/lib/backgroundQueueActiveThread.ts | 17 + .../lib/sendBackgroundQueuedMessage.test.ts | 275 ++++++++++++++ .../src/lib/sendBackgroundQueuedMessage.ts | 357 ++++++++++++++++++ apps/web/src/queuedMessageStore.test.ts | 32 +- apps/web/src/queuedMessageStore.ts | 58 ++- apps/web/src/routes/__root.tsx | 2 + 9 files changed, 927 insertions(+), 9 deletions(-) create mode 100644 apps/web/src/components/BackgroundQueuedMessages.tsx create mode 100644 apps/web/src/lib/backgroundQueueActiveThread.test.ts create mode 100644 apps/web/src/lib/backgroundQueueActiveThread.ts create mode 100644 apps/web/src/lib/sendBackgroundQueuedMessage.test.ts create mode 100644 apps/web/src/lib/sendBackgroundQueuedMessage.ts diff --git a/apps/web/src/components/BackgroundQueuedMessages.tsx b/apps/web/src/components/BackgroundQueuedMessages.tsx new file mode 100644 index 000000000000..4f6ee65f9fb3 --- /dev/null +++ b/apps/web/src/components/BackgroundQueuedMessages.tsx @@ -0,0 +1,135 @@ +import { useEffect, useLayoutEffect, useMemo, useRef } from "react"; +import { useParams } from "@tanstack/react-router"; +import { useShallow } from "zustand/react/shallow"; +import { parseScopedThreadKey } from "@t3tools/client-runtime/environment"; +import { derivePendingRequests } from "@t3tools/client-runtime/pending-requests"; + +import { useComposerDraftStore } from "../composerDraftStore"; +import { useClientSettingsHydrated } from "../hooks/useSettings"; +import { activeQueueThreadKey } from "../lib/backgroundQueueActiveThread"; +import { sendBackgroundQueuedMessage } from "../lib/sendBackgroundQueuedMessage"; +import { + isQueuedMessageDue, + latestCompletedToolActivityId, + useQueuedMessageStore, + useQueuedMessages, +} from "../queuedMessageStore"; +import { derivePhase } from "../session-logic"; +import { useEnvironments } from "../state/environments"; +import { useServerConfigs, useThread, useThreadShell } from "../state/entities"; +import { useEnvironmentThread } from "../state/threads"; +import { resolveThreadRouteTarget } from "../threadRoutes"; + +export function BackgroundQueueCoordinator() { + const target = useParams({ strict: false, select: resolveThreadRouteTarget }); + const draftId = target?.kind === "draft" ? target.draftId : null; + const draftThread = useComposerDraftStore((store) => + draftId ? store.getDraftSession(draftId) : null, + ); + return ( + + ); +} + +/** Watches every thread that has a queue, including ones the chat view is not showing. */ +export function BackgroundQueuedMessages({ activeThreadKey }: { activeThreadKey: string | null }) { + const keys = useQueuedMessageStore( + useShallow((state) => [ + ...new Set([ + ...Object.keys(state.queuesByThreadKey), + ...Object.keys(state.backgroundSendsByThreadKey), + ]), + ]), + ); + return keys.map((threadKey) => ( + + )); +} + +function BackgroundThreadQueue({ threadKey, active }: { threadKey: string; active: boolean }) { + const ref = useMemo(() => parseScopedThreadKey(threadKey), [threadKey]); + const thread = useThread(ref); + const shell = useThreadShell(ref); + const sending = useQueuedMessageStore( + (state) => state.backgroundSendsByThreadKey[threadKey] !== undefined, + ); + const detail = useEnvironmentThread(ref?.environmentId ?? null, ref?.threadId ?? null); + const { environments } = useEnvironments(); + const configs = useServerConfigs(); + const hydrated = useClientSettingsHydrated(); + const message = useQueuedMessages(threadKey)[0]; + const rewinding = useComposerDraftStore((state) => state.rewindingThreadKeys.has(threadKey)); + const config = ref ? configs.get(ref.environmentId) : undefined; + const phase = derivePhase(thread?.session ?? null); + const latestToolActivityId = latestCompletedToolActivityId(thread?.activities ?? []); + const pending = derivePendingRequests(thread?.activities ?? []); + const connected = environments.some( + (environment) => + environment.environmentId === ref?.environmentId && + environment.connection.phase === "connected", + ); + const provider = config?.providers.find( + (entry) => entry.instanceId === message?.sendOptions?.modelSelection.instanceId, + ); + const providerReady = + provider?.enabled === true && + provider.installed === true && + provider.availability !== "unavailable" && + provider.status === "ready"; + const blocked = + (active && !sending) || + rewinding || + !connected || + !hydrated || + !thread || + !shell || + detail.status !== "live" || + pending.approvals.length > 0 || + pending.userInputs.length > 0 || + !providerReady; + const attempted = useRef(null); + const live = useRef({ blocked, phase, latestToolActivityId }); + useLayoutEffect(() => { + live.current = { blocked, phase, latestToolActivityId }; + }, [blocked, phase, latestToolActivityId]); + + useEffect(() => { + if (!ref || !message?.sendOptions || blocked) return; + if ( + !isQueuedMessageDue({ + message, + phase, + latestToolActivityId, + }) + ) { + return; + } + const boundary = `${message.id}\0${phase}\0${latestToolActivityId ?? ""}\0${thread?.latestTurn?.turnId ?? ""}`; + if (attempted.current === boundary) return; + attempted.current = boundary; + const options = message.sendOptions; + void sendBackgroundQueuedMessage( + ref, + message, + options, + () => { + const current = live.current; + return ( + !current.blocked && + isQueuedMessageDue({ + message, + phase: current.phase, + latestToolActivityId: current.latestToolActivityId, + }) + ); + }, + () => live.current.latestToolActivityId, + ); + }, [blocked, latestToolActivityId, message, phase, ref, thread?.latestTurn?.turnId]); + + return null; +} diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 9af410f5722e..012217fe5afd 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -3201,7 +3201,7 @@ export default function ChatView(props: ChatViewProps) { localDispatchStartedAt, latestUserMessageAt, isPreparingWorktree: isLocallyPreparingWorktree, - isSendBusy, + isSendBusy: isLocalSendBusy, backgroundSubmissionPending, } = useLocalDispatchState({ activeThread, @@ -3211,6 +3211,11 @@ export default function ChatView(props: ChatViewProps) { activePendingUserInput: activePendingUserInput?.requestId ?? null, threadError, }); + const isBackgroundQueueSending = useQueuedMessageStore( + (state) => + activeThreadKey !== null && state.backgroundSendsByThreadKey[activeThreadKey] !== undefined, + ); + const isSendBusy = isLocalSendBusy || isBackgroundQueueSending; const optimisticCompactionMessage = optimisticUserMessages.at(-1); const pendingCompactionMessage = isSendBusy && @@ -7657,6 +7662,22 @@ export default function ChatView(props: ChatViewProps) { previewAnnotations: [...composerPreviewAnnotations], reviewComments: [...composerReviewComments], submissionIntent, + sendOptions: { + modelSelection: ctxSelectedModelSelection, + runtimeMode, + interactionMode: sendInteractionMode, + promptEffort: resolvePromptInjectedEffort( + getProviderModelCapabilities( + ctxSelectedProviderModels, + ctxSelectedModel, + ctxSelectedProvider, + ), + ctxSelectedPromptEffort, + ), + ...(localCheckoutBranchMismatch + ? { branch: localCheckoutBranchMismatch.currentBranch } + : {}), + }, queuedAfterToolActivityId: latestCompletedToolActivityId(threadActivities), createdAt: new Date().toISOString(), }); @@ -8619,9 +8640,9 @@ export default function ChatView(props: ChatViewProps) { } }; - // Sends the oldest queued message once it is due: a tool call finished - // after it was queued, or the turn ended. Only one leaves per boundary; the - // take inside onSend re-anchors the rest. + // Sends the oldest queued message for the open thread. Other threads are + // dispatched by BackgroundQueueCoordinator, which stays mounted across navigation. + // Only one leaves per boundary; the take inside onSend re-anchors the rest. const sendQueuedMessage = useEffectEvent((message: QueuedComposerMessage) => { void onSend(undefined, message.submissionIntent, undefined, message); }); diff --git a/apps/web/src/lib/backgroundQueueActiveThread.test.ts b/apps/web/src/lib/backgroundQueueActiveThread.test.ts new file mode 100644 index 000000000000..a7ba2fc7c351 --- /dev/null +++ b/apps/web/src/lib/backgroundQueueActiveThread.test.ts @@ -0,0 +1,31 @@ +import { scopeThreadRef } from "@t3tools/client-runtime/environment"; +import { ThreadId } from "@t3tools/contracts"; +import { describe, expect, it } from "vite-plus/test"; + +import { DraftId } from "../composerDraftStore"; +import { resolveThreadRouteTarget } from "../threadRoutes"; +import { activeQueueThreadKey } from "./backgroundQueueActiveThread"; + +describe("activeQueueThreadKey", () => { + it("uses the promoted server thread while the route is still a draft", () => { + const promotedTo = scopeThreadRef("env-2" as never, ThreadId.make("server-thread")); + + expect( + activeQueueThreadKey(resolveThreadRouteTarget({ draftId: DraftId.make("draft-1") }), { + environmentId: "env-1" as never, + threadId: ThreadId.make("draft-thread"), + promotedTo, + }), + ).toBe("env-2:server-thread"); + }); + + it("leaves an unpromoted draft to the chat view", () => { + expect( + activeQueueThreadKey(resolveThreadRouteTarget({ draftId: DraftId.make("draft-1") }), { + environmentId: "env-1" as never, + threadId: ThreadId.make("draft-thread"), + promotedTo: null, + }), + ).toBeNull(); + }); +}); diff --git a/apps/web/src/lib/backgroundQueueActiveThread.ts b/apps/web/src/lib/backgroundQueueActiveThread.ts new file mode 100644 index 000000000000..ce6112efa7f7 --- /dev/null +++ b/apps/web/src/lib/backgroundQueueActiveThread.ts @@ -0,0 +1,17 @@ +import { scopedThreadKey } from "@t3tools/client-runtime/environment"; + +import { resolveActiveThreadRouteRef, type ThreadRouteTarget } from "../threadRoutes"; + +type DraftRouteState = Parameters[1]; + +/** + * The thread whose queue the open chat owns. A draft route that has already + * promoted still shows that server thread, so Stop must stay on the chat path. + */ +export function activeQueueThreadKey( + target: ThreadRouteTarget | null, + draftThread: DraftRouteState, +): string | null { + const ref = resolveActiveThreadRouteRef(target, draftThread); + return ref ? scopedThreadKey(ref) : null; +} diff --git a/apps/web/src/lib/sendBackgroundQueuedMessage.test.ts b/apps/web/src/lib/sendBackgroundQueuedMessage.test.ts new file mode 100644 index 000000000000..bbb02099ead8 --- /dev/null +++ b/apps/web/src/lib/sendBackgroundQueuedMessage.test.ts @@ -0,0 +1,275 @@ +import { EnvironmentId, ProviderInstanceId, ThreadId } from "@t3tools/contracts"; +import { scopeThreadRef } from "@t3tools/client-runtime/environment"; +import { beforeEach, describe, expect, it, vi } from "vite-plus/test"; + +import { useComposerDraftStore } from "../composerDraftStore"; +import { useQueuedMessageStore, type QueuedMessageSendOptions } from "../queuedMessageStore"; +import { sendBackgroundQueuedMessage } from "./sendBackgroundQueuedMessage"; + +const mocks = vi.hoisted(() => ({ + runAtomCommand: vi.fn(), + readThreadShell: vi.fn(), + config: null as unknown, + toast: vi.fn(), + startTurn: Symbol("startTurn"), + updateMetadata: Symbol("updateMetadata"), + setRuntimeMode: Symbol("setRuntimeMode"), + setInteractionMode: Symbol("setInteractionMode"), +})); + +vi.mock("@t3tools/client-runtime/state/runtime", () => ({ + runAtomCommand: mocks.runAtomCommand, + squashAtomCommandFailure: (result: { readonly error: unknown }) => result.error, +})); + +vi.mock("../rpc/atomRegistry", () => ({ + appAtomRegistry: { + get: () => ({ + get: () => mocks.config, + }), + }, +})); + +vi.mock("../state/server", () => ({ + environmentServerConfigsAtom: Symbol("configs"), +})); + +vi.mock("../state/threads", () => ({ + threadEnvironment: { + startTurn: mocks.startTurn, + updateMetadata: mocks.updateMetadata, + setRuntimeMode: mocks.setRuntimeMode, + setInteractionMode: mocks.setInteractionMode, + }, +})); + +vi.mock("../state/entities", () => ({ + readThreadShell: mocks.readThreadShell, +})); + +vi.mock("../components/ui/toast", () => ({ + toastManager: { add: mocks.toast }, +})); + +vi.mock("./attachmentUploadQueue", () => ({ + awaitAttachmentUploads: vi.fn(), + getUploadedAttachments: vi.fn(), + releaseDraftAttachments: vi.fn(), + startAttachmentUpload: vi.fn(), +})); + +const environmentId = EnvironmentId.make("env-1"); +const threadRef = scopeThreadRef(environmentId, ThreadId.make("thread-1")); +const threadKey = `${environmentId}:${ThreadId.make("thread-1")}`; + +const sendOptions: QueuedMessageSendOptions = { + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5" }, + runtimeMode: "full-access", + interactionMode: "default", + promptEffort: null, +}; + +function enqueue(prompt: string, options: QueuedMessageSendOptions = sendOptions) { + return useQueuedMessageStore.getState().enqueue(threadKey, { + prompt, + images: [], + files: [], + terminalContexts: [], + previewAnnotations: [], + reviewComments: [], + submissionIntent: "foreground", + sendOptions: options, + queuedAfterToolActivityId: "tool-1", + createdAt: "2026-09-24T00:00:00.000Z", + }); +} + +describe("sendBackgroundQueuedMessage", () => { + beforeEach(() => { + useQueuedMessageStore.setState({ + queuesByThreadKey: {}, + backgroundSendsByThreadKey: {}, + drainGeneration: 0, + }); + useComposerDraftStore.setState({ + draftsByThreadKey: {}, + draftThreadsByThreadKey: {}, + }); + mocks.runAtomCommand.mockReset(); + mocks.runAtomCommand.mockResolvedValue({ _tag: "Success", value: undefined }); + mocks.toast.mockReset(); + mocks.config = { + providers: [ + { + instanceId: sendOptions.modelSelection.instanceId, + driver: "codex", + enabled: true, + installed: true, + availability: "available", + status: "ready", + }, + ], + environment: { + capabilities: { + attachmentUploads: false, + inlineMessageContext: true, + }, + }, + }; + mocks.readThreadShell.mockReset(); + mocks.readThreadShell.mockReturnValue({ + modelSelection: sendOptions.modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: "main", + }); + }); + + it("starts the turn when the queued thread is no longer selected", async () => { + const message = enqueue("please continue"); + + const sent = await sendBackgroundQueuedMessage( + threadRef, + message, + sendOptions, + () => true, + () => "tool-1", + ); + + expect(sent).toBe(true); + expect(mocks.runAtomCommand).toHaveBeenCalledWith( + expect.anything(), + mocks.startTurn, + expect.objectContaining({ + environmentId, + input: expect.objectContaining({ + threadId: threadRef.threadId, + message: expect.objectContaining({ text: "please continue" }), + }), + }), + { reportFailure: false }, + ); + expect(useQueuedMessageStore.getState().queuesByThreadKey[threadKey]).toBeUndefined(); + expect(useQueuedMessageStore.getState().backgroundSendsByThreadKey[threadKey]).toBeUndefined(); + }); + + it("persists the checkout branch captured when the message was queued", async () => { + const options = { ...sendOptions, branch: "feature" }; + const message = enqueue("please continue", options); + + const sent = await sendBackgroundQueuedMessage( + threadRef, + message, + options, + () => true, + () => "tool-1", + ); + + expect(sent).toBe(true); + expect(mocks.runAtomCommand).toHaveBeenCalledWith( + expect.anything(), + mocks.updateMetadata, + expect.objectContaining({ + input: expect.objectContaining({ branch: "feature", worktreePath: null }), + }), + { reportFailure: false }, + ); + }); + + it("leaves the message queued when the thread is not eligible yet", async () => { + const message = enqueue("please continue"); + + const sent = await sendBackgroundQueuedMessage( + threadRef, + message, + sendOptions, + () => false, + () => null, + ); + + expect(sent).toBe(false); + expect(mocks.runAtomCommand).not.toHaveBeenCalled(); + expect(useQueuedMessageStore.getState().queuesByThreadKey[threadKey]?.[0]?.id).toBe(message.id); + }); + + it("holds a failed send for the user instead of retrying it", async () => { + const message = enqueue("please continue"); + mocks.runAtomCommand.mockResolvedValue({ _tag: "Failure", error: new Error("provider down") }); + + const sent = await sendBackgroundQueuedMessage( + threadRef, + message, + sendOptions, + () => true, + () => null, + ); + + expect(sent).toBe(false); + const held = useQueuedMessageStore.getState().queuesByThreadKey[threadKey]?.[0]; + expect(held?.id).toBe(message.id); + expect(held?.holdUntilUserAction).toBe(true); + expect(mocks.toast).toHaveBeenCalled(); + }); + + it("does not start a turn when Stop drains the send", async () => { + const message = enqueue("please continue"); + mocks.runAtomCommand.mockImplementation(async (_registry, command) => { + if (command === mocks.updateMetadata) { + useQueuedMessageStore.getState().drain(threadKey); + } + return { _tag: "Success", value: undefined }; + }); + mocks.readThreadShell.mockReturnValue({ + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "other" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + }); + + const sent = await sendBackgroundQueuedMessage( + threadRef, + message, + sendOptions, + () => true, + () => null, + ); + + expect(sent).toBe(false); + expect(mocks.runAtomCommand.mock.calls.some((call) => call[1] === mocks.startTurn)).toBe(false); + expect(useQueuedMessageStore.getState().queuesByThreadKey[threadKey]).toBeUndefined(); + expect(useComposerDraftStore.getState().getComposerDraft(threadRef)?.prompt).toBe( + "please continue", + ); + }); + + it("appends a cancelled send to an existing composer draft", async () => { + const message = enqueue("please continue"); + useComposerDraftStore.getState().setPrompt(threadRef, "kept"); + mocks.runAtomCommand.mockImplementation(async (_registry, command) => { + if (command === mocks.setRuntimeMode) { + useQueuedMessageStore.getState().drain(threadKey); + } + return { _tag: "Success", value: undefined }; + }); + mocks.readThreadShell.mockReturnValue({ + modelSelection: sendOptions.modelSelection, + runtimeMode: "approval-required", + interactionMode: "default", + branch: "main", + }); + + const sent = await sendBackgroundQueuedMessage( + threadRef, + message, + sendOptions, + () => true, + () => null, + ); + + expect(sent).toBe(false); + expect(useComposerDraftStore.getState().getComposerDraft(threadRef)?.prompt).toBe( + "kept\n\nplease continue", + ); + expect(useQueuedMessageStore.getState().queuesByThreadKey[threadKey]).toBeUndefined(); + }); +}); diff --git a/apps/web/src/lib/sendBackgroundQueuedMessage.ts b/apps/web/src/lib/sendBackgroundQueuedMessage.ts new file mode 100644 index 000000000000..f3e7266783ac --- /dev/null +++ b/apps/web/src/lib/sendBackgroundQueuedMessage.ts @@ -0,0 +1,357 @@ +import { PROVIDER_SEND_TURN_MAX_ATTACHMENTS, type ScopedThreadRef } from "@t3tools/contracts"; +import { scopedThreadKey } from "@t3tools/client-runtime/environment"; +import { runAtomCommand, squashAtomCommandFailure } from "@t3tools/client-runtime/state/runtime"; +import { applyClaudePromptEffortPrefix } from "@t3tools/shared/model"; +import { serializeLegacyContextMessage } from "@t3tools/shared/composerContextLegacySend"; +import { appAtomRegistry } from "../rpc/atomRegistry"; +import { environmentServerConfigsAtom } from "../state/server"; +import { threadEnvironment } from "../state/threads"; +import { readThreadShell } from "../state/entities"; +import { + useQueuedMessageStore, + type QueuedComposerMessage, + type QueuedMessageSendOptions, +} from "../queuedMessageStore"; +import { + deriveComposerSendState, + getAntigravitySendBlockReason, + readFileAsDataUrl, + resolveThreadMetadataUpdateForNextTurn, +} from "../components/ChatView.logic"; +import { getComposerSubmissionValidationMessage } from "../components/chat/composerSubmission"; +import { ATTACHMENT_ONLY_BOOTSTRAP_PROMPT } from "../components/chat/composerPromptHistory"; +import { fileAttachmentCapabilityBlockReason } from "../components/chat/composerAttachmentFiles"; +import { toastManager } from "../components/ui/toast"; +import { stackedThreadToast } from "../components/ui/toastHelpers"; +import { useComposerDraftStore } from "../composerDraftStore"; +import { buildMessageContext, terminalContextReference } from "./composerContextRecords"; +import { removeInlineContextReference } from "./composerContextReferences"; +import { + awaitAttachmentUploads, + getUploadedAttachments, + releaseDraftAttachments, + startAttachmentUpload, +} from "./attachmentUploadQueue"; +import { newMessageId } from "./utils"; + +/** + * Dispatches one queued follow-up for a thread that is not on screen. + * The open thread keeps ChatView's send path. + */ +export async function sendBackgroundQueuedMessage( + ref: ScopedThreadRef, + message: QueuedComposerMessage, + options: QueuedMessageSendOptions, + canSend: () => boolean, + latestToolActivityId: () => string | null, +): Promise { + const key = scopedThreadKey(ref); + const stillQueued = () => + useQueuedMessageStore + .getState() + .queuesByThreadKey[key]?.some((entry) => entry.id === message.id) === true; + const readConfig = () => appAtomRegistry.get(environmentServerConfigsAtom).get(ref.environmentId); + const attachments = [...message.images, ...message.files]; + let taken = false; + + const cancelled = () => + useQueuedMessageStore.getState().backgroundSendsByThreadKey[key]?.cancelled === true; + + try { + if (!canSend() || !stillQueued()) return false; + const config = readConfig(); + const provider = config?.providers.find( + (entry) => entry.instanceId === options.modelSelection.instanceId, + ); + if ( + !config || + !provider?.enabled || + !provider.installed || + provider.availability === "unavailable" || + provider.status !== "ready" + ) { + return false; + } + const providerBlockReason = getAntigravitySendBlockReason( + provider, + options.modelSelection.model, + ); + if (providerBlockReason) throw new Error(providerBlockReason); + + const { sendableTerminalContexts, hasSendableContent } = deriveComposerSendState({ + prompt: message.prompt, + imageCount: attachments.length, + terminalContexts: message.terminalContexts, + elementContextCount: message.previewAnnotations.length + message.reviewComments.length, + }); + if (!hasSendableContent) { + useQueuedMessageStore.getState().remove(key, message.id); + return false; + } + + const prompt = message.terminalContexts + .filter((context) => !sendableTerminalContexts.includes(context)) + .reduce( + (text, context) => + removeInlineContextReference(text, terminalContextReference(context).contextId).prompt, + message.prompt, + ) + .trim(); + const text = applyClaudePromptEffortPrefix( + prompt || ATTACHMENT_ONLY_BOOTSTRAP_PROMPT, + options.promptEffort, + ); + const validation = getComposerSubmissionValidationMessage({ + prompt: message.prompt, + providerInput: text, + submissionTarget: "provider-turn", + }); + if (validation) throw new Error(validation); + + const checkFiles = () => { + const current = readConfig(); + const reason = fileAttachmentCapabilityBlockReason({ + files: message.files, + attachmentUploadsCapabilityKnown: current !== undefined, + supportsAttachmentUploads: current?.environment.capabilities.attachmentUploads === true, + maxFileAttachmentBytes: + current?.environment.capabilities.fileAttachments?.maxUploadBytes ?? null, + }); + if (reason) throw new Error(reason); + }; + checkFiles(); + + const upload = config.environment.capabilities.attachmentUploads === true; + if (upload && attachments.length > 0) { + for (const attachment of attachments) { + startAttachmentUpload({ + environmentId: ref.environmentId, + image: attachment, + draftTarget: ref, + }); + } + await awaitAttachmentUploads(attachments.map((attachment) => attachment.id)); + if (!canSend() || cancelled()) return false; + } + + const wireAttachments = await Promise.all( + attachments.map(async (attachment) => { + if (upload) { + const uploaded = getUploadedAttachments({ + environmentId: ref.environmentId, + images: [attachment], + })?.[0]; + if (!uploaded) throw new Error("Retry or remove failed uploads before sending."); + return uploaded; + } + if (attachment.type !== "image") { + throw new Error("This server does not support file attachments."); + } + return { + type: "image" as const, + id: attachment.id, + name: attachment.name, + mimeType: attachment.mimeType, + sizeBytes: attachment.sizeBytes, + dataUrl: await readFileAsDataUrl(attachment.file), + ...(attachment.source ? { source: attachment.source } : {}), + }; + }), + ); + checkFiles(); + if (!canSend() || !stillQueued()) return false; + + const shell = readThreadShell(ref); + if (!shell) return false; + if (!useQueuedMessageStore.getState().take(key, message.id, latestToolActivityId(), true)) { + return false; + } + taken = true; + + const giveBack = () => { + if (cancelled()) { + restoreCancelledQueuedMessageToComposer(ref, message); + return; + } + useQueuedMessageStore.getState().holdAtFront(key, message); + }; + if (cancelled() || !canSend()) { + giveBack(); + return false; + } + + const createdAt = new Date().toISOString(); + const metadataUpdate = resolveThreadMetadataUpdateForNextTurn({ + currentModelSelection: shell.modelSelection, + nextModelSelection: options.modelSelection, + currentBranch: shell.branch, + ...(options.branch !== undefined ? { nextBranch: options.branch } : {}), + }); + if (metadataUpdate) { + const result = await runAtomCommand( + appAtomRegistry, + threadEnvironment.updateMetadata, + { + environmentId: ref.environmentId, + input: { threadId: ref.threadId, ...metadataUpdate }, + }, + { reportFailure: false }, + ); + if (result._tag === "Failure") throw squashAtomCommandFailure(result); + } + if (cancelled() || !canSend()) { + giveBack(); + return false; + } + if (options.runtimeMode !== shell.runtimeMode) { + const result = await runAtomCommand( + appAtomRegistry, + threadEnvironment.setRuntimeMode, + { + environmentId: ref.environmentId, + input: { threadId: ref.threadId, runtimeMode: options.runtimeMode, createdAt }, + }, + { reportFailure: false }, + ); + if (result._tag === "Failure") throw squashAtomCommandFailure(result); + } + if (cancelled() || !canSend()) { + giveBack(); + return false; + } + if (options.interactionMode !== shell.interactionMode) { + const result = await runAtomCommand( + appAtomRegistry, + threadEnvironment.setInteractionMode, + { + environmentId: ref.environmentId, + input: { + threadId: ref.threadId, + interactionMode: options.interactionMode, + createdAt, + }, + }, + { reportFailure: false }, + ); + if (result._tag === "Failure") throw squashAtomCommandFailure(result); + } + if (cancelled() || !canSend()) { + giveBack(); + return false; + } + + const context = buildMessageContext({ + terminalContexts: sendableTerminalContexts, + previewAnnotations: message.previewAnnotations, + reviewComments: message.reviewComments, + attachments: attachments.map((attachment, index) => ({ + attachment, + attachmentId: wireAttachments[index]?.id ?? attachment.id, + })), + }); + const inlineContext = readConfig()?.environment.capabilities.inlineMessageContext === true; + const result = await runAtomCommand( + appAtomRegistry, + threadEnvironment.startTurn, + { + environmentId: ref.environmentId, + input: { + threadId: ref.threadId, + message: { + messageId: newMessageId(), + role: "user", + text: + context && !inlineContext + ? serializeLegacyContextMessage({ text, records: context.records }) + : text, + attachments: wireAttachments, + ...(context && inlineContext ? { context } : {}), + }, + modelSelection: options.modelSelection, + runtimeMode: options.runtimeMode, + interactionMode: options.interactionMode, + createdAt, + }, + }, + { reportFailure: false }, + ); + if (result._tag === "Failure") throw squashAtomCommandFailure(result); + if (upload) releaseDraftAttachments(attachments); + return true; + } catch (error) { + if (cancelled()) { + if (taken) restoreCancelledQueuedMessageToComposer(ref, message); + return false; + } + if (taken || stillQueued()) { + useQueuedMessageStore.getState().holdAtFront(key, message); + toastManager.add({ + type: "error", + title: "Queued message not sent", + description: + error instanceof Error ? error.message : "Open the thread to retry the queued message.", + }); + } + return false; + } finally { + if (taken) useQueuedMessageStore.getState().finishBackgroundSend(key); + } +} + +/** + * Stop already drained the rest of the queue into the composer. A claimed + * send is no longer in that drain, so cancellation writes it back here + * instead of parking it as a held queue row. + */ +function restoreCancelledQueuedMessageToComposer( + ref: ScopedThreadRef, + message: QueuedComposerMessage, +) { + const store = useComposerDraftStore.getState(); + const draft = store.getComposerDraft(ref); + const prompts = [draft?.prompt ?? "", message.prompt] + .map((prompt) => prompt.trim()) + .filter((prompt) => prompt.length > 0); + store.setPrompt(ref, prompts.join("\n\n")); + + const room = Math.max( + 0, + PROVIDER_SEND_TURN_MAX_ATTACHMENTS - (draft?.images.length ?? 0) - (draft?.files.length ?? 0), + ); + const attachments = [...message.images, ...message.files]; + const restored = attachments.slice(0, room); + const overflow = attachments.slice(room); + const images = restored.filter((attachment) => attachment.type === "image"); + const files = restored.filter((attachment) => attachment.type === "file"); + if (images.length > 0) store.addImages(ref, images); + if (files.length > 0) store.addFiles(ref, files); + if (overflow.length > 0) { + useQueuedMessageStore.getState().enqueue(scopedThreadKey(ref), { + prompt: "", + images: overflow.filter((attachment) => attachment.type === "image"), + files: overflow.filter((attachment) => attachment.type === "file"), + terminalContexts: [], + previewAnnotations: [], + reviewComments: [], + submissionIntent: "foreground", + queuedAfterToolActivityId: message.queuedAfterToolActivityId, + holdUntilUserAction: true, + createdAt: new Date().toISOString(), + }); + toastManager.add( + stackedThreadToast({ + type: "info", + title: "Some attachments stayed queued", + description: `A message holds at most ${PROVIDER_SEND_TURN_MAX_ATTACHMENTS} attachments. Use Send now on the queued row when you want the rest to go.`, + }), + ); + } + + store.setTerminalContexts(ref, [...(draft?.terminalContexts ?? []), ...message.terminalContexts]); + const next = store.getComposerDraft(ref); + store.setPreviewAnnotations(ref, [ + ...(next?.previewAnnotations ?? []), + ...message.previewAnnotations, + ]); + store.setReviewComments(ref, [...(next?.reviewComments ?? []), ...message.reviewComments]); +} diff --git a/apps/web/src/queuedMessageStore.test.ts b/apps/web/src/queuedMessageStore.test.ts index 33869d89da62..cd477a8a9577 100644 --- a/apps/web/src/queuedMessageStore.test.ts +++ b/apps/web/src/queuedMessageStore.test.ts @@ -23,7 +23,11 @@ function makeMessage(prompt: string): Omit { describe("queuedMessageStore", () => { beforeEach(() => { - useQueuedMessageStore.setState({ queuesByThreadKey: {}, drainGeneration: 0 }); + useQueuedMessageStore.setState({ + queuesByThreadKey: {}, + backgroundSendsByThreadKey: {}, + drainGeneration: 0, + }); }); it("keeps messages in submission order per thread", () => { @@ -98,6 +102,32 @@ describe("queuedMessageStore", () => { expect(useQueuedMessageStore.getState().drainGeneration).toBe(1); expect(useQueuedMessageStore.getState().queuesByThreadKey["thread-b"]).toHaveLength(1); }); + + it("lets one background send claim a message and ignores a second take", () => { + const { enqueue, take, finishBackgroundSend } = useQueuedMessageStore.getState(); + const entry = enqueue("thread-a", makeMessage("first")); + + expect(take("thread-a", entry.id, null, true)?.prompt).toBe("first"); + expect(take("thread-a", entry.id, null)).toBeNull(); + expect(useQueuedMessageStore.getState().backgroundSendsByThreadKey["thread-a"]).toEqual({ + cancelled: false, + }); + + finishBackgroundSend("thread-a"); + expect(useQueuedMessageStore.getState().backgroundSendsByThreadKey["thread-a"]).toBeUndefined(); + }); + + it("marks an in-flight background send cancelled when the thread is drained", () => { + const { enqueue, take, drain } = useQueuedMessageStore.getState(); + const entry = enqueue("thread-a", makeMessage("first")); + take("thread-a", entry.id, null, true); + enqueue("thread-a", makeMessage("second")); + + expect(drain("thread-a").map((message) => message.prompt)).toEqual(["second"]); + expect(useQueuedMessageStore.getState().backgroundSendsByThreadKey["thread-a"]).toEqual({ + cancelled: true, + }); + }); }); describe("queued message dispatch timing", () => { diff --git a/apps/web/src/queuedMessageStore.ts b/apps/web/src/queuedMessageStore.ts index b342a1380f7a..024ef15e652c 100644 --- a/apps/web/src/queuedMessageStore.ts +++ b/apps/web/src/queuedMessageStore.ts @@ -1,4 +1,10 @@ -import type { PreviewAnnotationPayload } from "@t3tools/contracts"; +import type { + ModelSelection, + PreviewAnnotationPayload, + ProviderInteractionMode, + RuntimeMode, +} from "@t3tools/contracts"; +import type { resolvePromptInjectedEffort } from "@t3tools/shared/model"; import { create } from "zustand"; import type { ComposerSubmissionIntent } from "./composer-logic"; @@ -12,6 +18,19 @@ import type { ReviewCommentContext } from "./reviewCommentContext"; * carries the full draft snapshot so the send path can dispatch it later with * the same text, attachments, and contexts the user pressed Enter on. */ +/** Composer choices captured at queue time, so a later send does not read another thread. */ +export interface QueuedMessageSendOptions { + modelSelection: ModelSelection; + runtimeMode: RuntimeMode; + interactionMode: ProviderInteractionMode; + promptEffort: ReturnType; + /** + * Checkout the user had selected when they queued. Dispatch persists it so + * the follow-up does not stay on the thread's old branch. + */ + branch?: string; +} + export interface QueuedComposerMessage { id: string; prompt: string; @@ -21,6 +40,7 @@ export interface QueuedComposerMessage { previewAnnotations: PreviewAnnotationPayload[]; reviewComments: ReviewCommentContext[]; submissionIntent: ComposerSubmissionIntent; + sendOptions?: QueuedMessageSendOptions; /** * The newest completed tool activity at queue time. A different id later * means a tool call finished after the user queued, which is the boundary @@ -37,6 +57,9 @@ export interface QueuedComposerMessage { interface QueuedMessageStoreState { queuesByThreadKey: Record; + /** Set while a non-selected thread is sending, so Stop and the open chat cannot both dispatch. */ + backgroundSendsByThreadKey: Record; + finishBackgroundSend: (threadKey: string) => void; /** * Bumped by `drain`. A send that took a message before a drain and finishes * its upload after it compares this to the value it captured and gives up, @@ -53,6 +76,7 @@ interface QueuedMessageStoreState { threadKey: string, id: string, toolActivityId: string | null, + background?: boolean, ) => QueuedComposerMessage | null; /** Removes one message without touching the others' anchors. Null when already gone. */ remove: (threadKey: string, id: string) => QueuedComposerMessage | null; @@ -70,6 +94,14 @@ const EMPTY_QUEUE: QueuedComposerMessage[] = []; /** In-memory only: a queued message is a live intent, not a draft worth persisting. */ export const useQueuedMessageStore = create()((set, get) => ({ queuesByThreadKey: {}, + backgroundSendsByThreadKey: {}, + finishBackgroundSend: (threadKey) => + set((state) => { + if (state.backgroundSendsByThreadKey[threadKey] === undefined) return state; + const backgroundSendsByThreadKey = { ...state.backgroundSendsByThreadKey }; + delete backgroundSendsByThreadKey[threadKey]; + return { backgroundSendsByThreadKey }; + }), drainGeneration: 0, enqueue: (threadKey, message) => { const entry: QueuedComposerMessage = { ...message, id: randomUUID() }; @@ -81,10 +113,10 @@ export const useQueuedMessageStore = create()((set, get })); return entry; }, - take: (threadKey, id, toolActivityId) => { + take: (threadKey, id, toolActivityId, background = false) => { const queue = get().queuesByThreadKey[threadKey]; const entry = queue?.find((message) => message.id === id); - if (!queue || !entry) { + if (!queue || !entry || get().backgroundSendsByThreadKey[threadKey]) { return null; } set((state) => { @@ -101,7 +133,17 @@ export const useQueuedMessageStore = create()((set, get } else { queuesByThreadKey[threadKey] = remaining; } - return { queuesByThreadKey }; + return { + queuesByThreadKey, + ...(background + ? { + backgroundSendsByThreadKey: { + ...state.backgroundSendsByThreadKey, + [threadKey]: { cancelled: false }, + }, + } + : {}), + }; }); return entry; }, @@ -139,6 +181,14 @@ export const useQueuedMessageStore = create()((set, get }); }, drain: (threadKey) => { + if (get().backgroundSendsByThreadKey[threadKey]) { + set((state) => ({ + backgroundSendsByThreadKey: { + ...state.backgroundSendsByThreadKey, + [threadKey]: { cancelled: true }, + }, + })); + } const queue = get().queuesByThreadKey[threadKey]; if (!queue || queue.length === 0) { return EMPTY_QUEUE; diff --git a/apps/web/src/routes/__root.tsx b/apps/web/src/routes/__root.tsx index eee703ac3f47..0e30d7736d99 100644 --- a/apps/web/src/routes/__root.tsx +++ b/apps/web/src/routes/__root.tsx @@ -73,6 +73,7 @@ import { type KeybindingsUpdateToastController, } from "../components/KeybindingsUpdateToast.logic"; +import { BackgroundQueueCoordinator } from "../components/BackgroundQueuedMessages"; import { getDesktopSnapShotBridge } from "../lib/desktopSnapShot"; import { installDesktopPasteAsText } from "../lib/desktopPasteAsText"; import { shouldResumeSnapShotSetupOnStartup } from "../lib/snapShotSetupResume"; @@ -224,6 +225,7 @@ function RootRouteView() { +