From ebdde89ea1d3a0967f040ea99f1d68064abf9a4e Mon Sep 17 00:00:00 2001 From: John Paul Soliva Date: Sat, 26 Sep 2026 09:52:03 +0900 Subject: [PATCH] fix(desktop): send a queued prompt once when several windows drain it Every desktop window auto-drains the shared composer queue: ChatBar drains the selected session and useBackgroundQueueDrain drains every other idle one. Since #122953 each window also sees the other windows' entries. The drain guard is a renderer-local ref, so two windows can pick the same entry and both submit it before either removes it. The gateway does not reject the second copy: a queued submit to a busy session is queued server-side and runs as its own turn. Both drain paths now run inside withQueueDrainClaim, a per-session Web Lock that the browser arbitrates across windows and frees when the holder closes. The entry is picked inside the claim from the live persisted queue, so a window that waited finds an entry the holder already sent gone and sends nothing. ChatBar auto-drain re-locates its entry by id and does not count "already sent elsewhere" as a failed attempt, so it can no longer run up the stuck-queue toast. Co-authored-by: furancis --- .../chat/composer/hooks/use-composer-queue.ts | 76 ++++++++++--------- .../hooks/use-background-queue-drain.test.tsx | 45 +++++++++++ .../hooks/use-background-queue-drain.ts | 59 +++++++------- apps/desktop/src/store/composer-queue.ts | 16 ++++ 4 files changed, 132 insertions(+), 64 deletions(-) diff --git a/apps/desktop/src/app/chat/composer/hooks/use-composer-queue.ts b/apps/desktop/src/app/chat/composer/hooks/use-composer-queue.ts index 027aae1cd8..ad8edfb8b2 100644 --- a/apps/desktop/src/app/chat/composer/hooks/use-composer-queue.ts +++ b/apps/desktop/src/app/chat/composer/hooks/use-composer-queue.ts @@ -20,7 +20,8 @@ import { removeQueuedPrompt, shouldAutoDrain, unparkQueuedPrompts, - updateQueuedPrompt + updateQueuedPrompt, + withQueueDrainClaim } from '@/store/composer-queue' import { notify } from '@/store/notifications' import { $sessionsLoading } from '@/store/session' @@ -203,50 +204,55 @@ export function useComposerQueue({ }, [activeQueueSessionKey, attachments, clearDraft, draftRef, scope.attachments]) // All queue drain paths share one lock + send-then-remove sequence. - // `pickEntry` lets each caller choose head, by-id, or skip-edited. + // `pickEntry` lets each caller choose head, by-id, or skip-edited, from the + // queue as it stands inside the cross-window claim. Resolves null when there + // is nothing to send: another window already sent the picked entry. const runDrain = useCallback( - async (pickEntry: (entries: QueuedPromptEntry[]) => QueuedPromptEntry | undefined): Promise => { + async (pickEntry: (entries: QueuedPromptEntry[]) => QueuedPromptEntry | undefined): Promise => { if (drainingQueueRef.current || !activeQueueSessionKey) { return false } const drainQueueSessionKey = activeQueueSessionKey const drainRuntimeSessionId = sessionId ?? null - const entry = pickEntry(getQueuedPrompts(drainQueueSessionKey)) - - if (!entry) { - return false - } drainingQueueRef.current = true try { - const accepted = await Promise.resolve( - onSubmit(entry.text, { - attachments: entry.attachments, - ...(entry.displayText ? { displayText: entry.displayText } : {}), - ...(entry.displayKind ? { displayKind: entry.displayKind } : {}), - fromQueue: true, - sessionId: drainRuntimeSessionId, - storedSessionId: drainQueueSessionKey - }) - ) + return await withQueueDrainClaim(drainQueueSessionKey, async queue => { + const entry = pickEntry(queue) - if (accepted === false) { - return false - } + if (!entry) { + return null + } - drainFailuresRef.current.delete(entry.id) - // Submit now owns the blob: previews (optimistic bubble); do not revoke. - removeQueuedPrompt(drainQueueSessionKey, entry.id, { retainPreviewUrls: true }) - resetBrowseState(drainRuntimeSessionId) - // A successful drain means the queue is flowing again — lift any park - // so the remaining entries follow. Manual drains (Enter on an empty - // composer, the per-row send arrow) are exactly the resume gestures a - // parked queue waits for; the auto path only reaches here unparked. - unparkQueuedPrompts(drainQueueSessionKey) + const accepted = await Promise.resolve( + onSubmit(entry.text, { + attachments: entry.attachments, + ...(entry.displayText ? { displayText: entry.displayText } : {}), + ...(entry.displayKind ? { displayKind: entry.displayKind } : {}), + fromQueue: true, + sessionId: drainRuntimeSessionId, + storedSessionId: drainQueueSessionKey + }) + ) - return true + if (accepted === false) { + return false + } + + drainFailuresRef.current.delete(entry.id) + // Submit now owns the blob: previews (optimistic bubble); do not revoke. + removeQueuedPrompt(drainQueueSessionKey, entry.id, { retainPreviewUrls: true }) + resetBrowseState(drainRuntimeSessionId) + // A successful drain means the queue is flowing again — lift any park + // so the remaining entries follow. Manual drains (Enter on an empty + // composer, the per-row send arrow) are exactly the resume gestures a + // parked queue waits for; the auto path only reaches here unparked. + unparkQueuedPrompts(drainQueueSessionKey) + + return true + }) } finally { drainingQueueRef.current = false } @@ -263,7 +269,7 @@ export function useComposerQueue({ [queueEditRef] // reads the edit id off a ref so the lock-holder always sees the latest ) - const drainNextQueued = useCallback(() => runDrain(pickDrainHead), [pickDrainHead, runDrain]) + const drainNextQueued = useCallback(async () => (await runDrain(pickDrainHead)) === true, [pickDrainHead, runDrain]) const sendQueuedNow = useCallback( (id: string) => { @@ -375,9 +381,11 @@ export function useComposerQueue({ } } - void runDrain(() => entry) + // By id: inside the claim the head may already be gone — sent by another + // window — which is not a failed send. + void runDrain(entries => entries.find(e => e.id === entry.id)) .then(sent => { - if (!sent) { + if (sent === false) { onFail() } }) diff --git a/apps/desktop/src/app/session/hooks/use-background-queue-drain.test.tsx b/apps/desktop/src/app/session/hooks/use-background-queue-drain.test.tsx index ea18717dab..f4241b5bbb 100644 --- a/apps/desktop/src/app/session/hooks/use-background-queue-drain.test.tsx +++ b/apps/desktop/src/app/session/hooks/use-background-queue-drain.test.tsx @@ -114,6 +114,51 @@ describe('useBackgroundQueueDrain', () => { await waitFor(() => expect(getQueuedPrompts('stored-session-a')).toHaveLength(0)) }) + it('submits a queued entry once when two idle windows both drain it', async () => { + // Web Locks are arbitrated across windows by the browser; jsdom has none. + // A FIFO mutex per lock name stands in for it. + const tails = new Map>() + + const request = (name: string, callback: () => Promise) => { + const granted = (tails.get(name) ?? Promise.resolve()).then(callback) + tails.set( + name, + granted.catch(() => {}) + ) + + return granted + } + + Object.defineProperty(window.navigator, 'locks', { configurable: true, value: { request } }) + + try { + const runtimeMap = { current: new Map([['stored-session-a', 'rt-session-a']]) } + let accept!: (accepted: boolean) => void + const submitText = vi.fn(() => new Promise(resolve => (accept = resolve))) + + enqueueQueuedPrompt('stored-session-a', { text: 'send me once', attachments: [] }) + clearAllSessionStates() + + // Two windows, each viewing another chat, share session A's queue. + render() + render() + + await waitFor(() => expect(submitText).toHaveBeenCalled()) + await new Promise(resolve => window.setTimeout(resolve, 0)) + await act(async () => accept(true)) + + await waitFor(() => expect(getQueuedPrompts('stored-session-a')).toHaveLength(0)) + // Let the window that waited on the claim finish its turn at the queue. + await act(async () => { + await Promise.all(tails.values()) + }) + + expect(submitText).toHaveBeenCalledTimes(1) + } finally { + delete (window.navigator as { locks?: unknown }).locks + } + }) + it('leaves the selected session queue to the mounted ChatBar drainer', async () => { const runtimeMap = { current: new Map([['stored-session-a', 'rt-session-a']]) } const submitText = vi.fn(async () => true) diff --git a/apps/desktop/src/app/session/hooks/use-background-queue-drain.ts b/apps/desktop/src/app/session/hooks/use-background-queue-drain.ts index 7c60904373..df1074ee80 100644 --- a/apps/desktop/src/app/session/hooks/use-background-queue-drain.ts +++ b/apps/desktop/src/app/session/hooks/use-background-queue-drain.ts @@ -6,12 +6,12 @@ import { resetBrowseState } from '@/store/composer-input-history' import { $parkedQueueSessions, $queuedPromptsBySession, - getQueuedPrompts, MAX_AUTO_DRAIN_ATTEMPTS, noteQueuedPromptDrainFailure, type QueuedPromptEntry, removeQueuedPrompt, - shouldAutoDrain + shouldAutoDrain, + withQueueDrainClaim } from '@/store/composer-queue' import { notify } from '@/store/notifications' import { @@ -164,36 +164,35 @@ export function useBackgroundQueueDrain({ scheduleRetry() } - void Promise.resolve() - .then(async () => { - const liveEntry = getQueuedPrompts(sessionKey).find(candidate => candidate.id === entry.id) - - if (!liveEntry) { - return true - } - - const runtimeSessionId = runtimeIdByStoredSessionIdRef.current.get(sessionKey) ?? null - - const accepted = await Promise.resolve( - submitTextRef.current(liveEntry.text, { - attachments: liveEntry.attachments, - fromQueue: true, - sessionId: runtimeSessionId, - storedSessionId: sessionKey - }) - ) - - if (accepted === false) { - return false - } - - drainFailuresRef.current.delete(liveEntry.id) - // Submit owns blob: previews after a successful drain handoff. - removeQueuedPrompt(sessionKey, liveEntry.id, { retainPreviewUrls: true }) - resetBrowseState(runtimeSessionId) + void withQueueDrainClaim(sessionKey, async queue => { + const liveEntry = queue.find(candidate => candidate.id === entry.id) + if (!liveEntry) { return true - }) + } + + const runtimeSessionId = runtimeIdByStoredSessionIdRef.current.get(sessionKey) ?? null + + const accepted = await Promise.resolve( + submitTextRef.current(liveEntry.text, { + attachments: liveEntry.attachments, + fromQueue: true, + sessionId: runtimeSessionId, + storedSessionId: sessionKey + }) + ) + + if (accepted === false) { + return false + } + + drainFailuresRef.current.delete(liveEntry.id) + // Submit owns blob: previews after a successful drain handoff. + removeQueuedPrompt(sessionKey, liveEntry.id, { retainPreviewUrls: true }) + resetBrowseState(runtimeSessionId) + + return true + }) .then(accepted => { if (!accepted) { onFail() diff --git a/apps/desktop/src/store/composer-queue.ts b/apps/desktop/src/store/composer-queue.ts index b5e83cd22e..515bc84f85 100644 --- a/apps/desktop/src/store/composer-queue.ts +++ b/apps/desktop/src/store/composer-queue.ts @@ -179,6 +179,22 @@ export const getQueuedPrompts = (key: string | null | undefined): QueuedPromptEn return sid ? queueFor(sid) : [] } +/** + * Run one drain of a session's queue while holding its cross-window claim, with + * `task` given that queue read fresh INSIDE the claim. Every idle window + * auto-drains the shared queue, so a renderer-local flag cannot stop two of + * them submitting the same entry — and the gateway runs the second copy as its + * own turn. Web Locks are arbitrated by the browser across windows and freed if + * the holder closes; a waiting window then finds an entry the holder sent + * already gone. Without Web Locks there is no other window to exclude. + */ +export const withQueueDrainClaim = (sid: string, task: (queue: QueuedPromptEntry[]) => Promise): Promise => { + const run = () => task(current()[sid] ?? []) + const locks = typeof navigator === 'undefined' ? undefined : navigator.locks + + return locks ? locks.request(`${STORAGE_KEY}.drain.${sid}`, run) : run() +} + export const enqueueQueuedPrompt = ( key: string | null | undefined, payload: { text: string; attachments: ComposerAttachment[]; displayText?: string; displayKind?: 'hidden' }