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 <brooklyn.bb.nicholson@gmail.com>
This commit is contained in:
committed by
brooklyn!
parent
7975f40088
commit
ebdde89ea1
@@ -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<boolean> => {
|
||||
async (pickEntry: (entries: QueuedPromptEntry[]) => QueuedPromptEntry | undefined): Promise<boolean | null> => {
|
||||
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()
|
||||
}
|
||||
})
|
||||
|
||||
@@ -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<string, Promise<unknown>>()
|
||||
|
||||
const request = (name: string, callback: () => Promise<unknown>) => {
|
||||
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<boolean>(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(<Harness runtimeMap={runtimeMap} selectedStoredSessionId="stored-session-b" submitText={submitText} />)
|
||||
render(<Harness runtimeMap={runtimeMap} selectedStoredSessionId="stored-session-c" submitText={submitText} />)
|
||||
|
||||
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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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 = <T>(sid: string, task: (queue: QueuedPromptEntry[]) => Promise<T>): Promise<T> => {
|
||||
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' }
|
||||
|
||||
Reference in New Issue
Block a user