diff --git a/src/core/starting-card-publication.ts b/src/core/starting-card-publication.ts new file mode 100644 index 0000000000..b339642efc --- /dev/null +++ b/src/core/starting-card-publication.ts @@ -0,0 +1,40 @@ +import type { DaemonSession } from './types.js'; + +// Cosmetic ordering only. A worker never waits for this promise. Include a +// predecessor POST: its completion may synchronously start the next turn's +// pending card, and thinking must not jump ahead of that follow-up POST. +const posts = new WeakMap>>(); + +export function trackStartingCardPublication(ds: DaemonSession, post: Promise): Promise { + const pending = posts.get(ds) ?? new Set>(); + posts.set(ds, pending); + const settled = post.then(() => {}, () => {}); + pending.add(settled); + void settled.then(() => { + pending.delete(settled); + if (!pending.size && posts.get(ds) === pending) posts.delete(ds); + }); + return post; +} + +/** No extra microtask for sessions with no starting card. Disabled cards and + * non-Lark sessions retain their existing independent thinking policy. */ +export function pendingStartingCardPublication(ds: DaemonSession): Promise | undefined { + if (!posts.get(ds)?.size) return undefined; + return (async () => { + let timer: ReturnType | undefined; + try { + await Promise.race([ + (async () => { + while (posts.get(ds)?.size) await Promise.all([...posts.get(ds)!]); + })(), + new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error('starting card publication timed out')), 15_000); + timer.unref?.(); + }), + ]); + } finally { + if (timer) clearTimeout(timer); + } + })(); +} diff --git a/src/core/worker-pool.ts b/src/core/worker-pool.ts index f9abd4777a..f539e59ba9 100644 --- a/src/core/worker-pool.ts +++ b/src/core/worker-pool.ts @@ -1,3 +1,4 @@ +import { trackStartingCardPublication } from './starting-card-publication.js'; import { handoffCardClosed, handoffCardBlocksStreaming, applyHandoffCardEvent, type HandoffCardEvent } from './handoff-card-lifecycle.js'; import { commitTriggerStreamingCard, discardTriggerStreamingCard, hasPendingTriggerStreamingCard } from './trigger-streaming-card.js'; import { sessionPromptInjection } from './prompt-injection.js'; @@ -3914,6 +3915,9 @@ export async function postTurnStartingCard( // durable reply card. Start both before awaiting either to preserve the // terminal card's synchronous generation/sentinel fence during slow POSTs. const statusPost = postTurnStartingStatusCard(ds, sessionReply, turnId); + // Track each POST separately: one rejection must not release the other card. + trackStartingCardPublication(ds, replyPost); + trackStartingCardPublication(ds, statusPost); const [replyPosted, statusPosted] = await Promise.all([replyPost, statusPost]); return replyPosted || statusPosted; } diff --git a/src/im/lark/cot-message.ts b/src/im/lark/cot-message.ts index 417ccf7ac9..f81ebc6921 100644 --- a/src/im/lark/cot-message.ts +++ b/src/im/lark/cot-message.ts @@ -46,6 +46,7 @@ import { getBot, getBotClient } from '../../bot-registry.js'; import { boundSubjectForTitle, subjectFromArgsString, type ToolSubject } from '../../services/cot-subject.js'; import { fallbackTurnId, frozenReplyContextForTurn } from '../../core/reply-target.js'; import { isSilentScheduledTurn } from '../../core/silent-schedule-turns.js'; +import { pendingStartingCardPublication } from '../../core/starting-card-publication.js'; import { config } from '../../config.js'; import { logger } from '../../utils/logger.js'; import { localeForBot, t } from '../../i18n/index.js'; @@ -574,6 +575,18 @@ async function pump(ds: DaemonSession, state: CotState): Promise { try { while (!state.disabled) { if (!state.cotId) { + const startingCard = pendingStartingCardPublication(ds); + if (startingCard) { + await startingCard; + // A newer turn/stop can arrive during the POST. Never resurrect its + // predecessor's not-yet-visible bubble below the current work card. + if (state.disabled || state.settled || states.get(ds) !== state + || state.finishStatus === 'interrupted' + || (ds.currentTurnId && ds.currentTurnId !== state.turnId)) { + state.disabled = true; + break; + } + } await apiCreate(ds, state); // Record the orphan marker the moment the bubble exists — before the // prologue append. If the prologue fails (or the daemon restarts diff --git a/test/cot-message.test.ts b/test/cot-message.test.ts index 88c84c3c44..9dfc1691e4 100644 --- a/test/cot-message.test.ts +++ b/test/cot-message.test.ts @@ -958,3 +958,105 @@ describe('superseded turn (type-ahead: next turn starts before the previous one expect(complete![0].params ?? complete![0].data).toMatchObject({ reason: 'error' }); }); }); + +describe('starting work card and thinking publication order', () => { + it('waits for a delayed work card before creating the bubble and preserves buffered output', async () => { + const { trackStartingCardPublication } = await import('../src/core/starting-card-publication.js'); + const ds = makeDs(); + let finish!: () => void; + trackStartingCardPublication(ds, new Promise(resolve => { finish = resolve; })); + handleCotThinkingUpdate(ds, upd([think('first')])); + handleCotThinkingUpdate(ds, upd([think('first'), say('second')])); + finalizeCotMessage(ds, 'om_turn1', 'completed'); + await flush(); + expect(request).not.toHaveBeenCalled(); + finish(); await flush(); await flush(); + expect(request.mock.calls.filter(([req]) => req.url === '/open-apis/im/v1/message_cot' && req.method === 'POST')).toHaveLength(1); + const events = pushedEvents(); + expect(events.some(e => e.content.delta === 'second')).toBe(true); + expect(events.some(e => e.type === 'RUN_FINISHED')).toBe(true); + }); + it('does not let a failed card POST block thinking or another session', async () => { + const { trackStartingCardPublication } = await import('../src/core/starting-card-publication.js'); + const ds = makeDs(), other = makeDs({session:{sessionId:'other'}}); + let reject!: (e: Error) => void; + const post = new Promise((_, r) => { reject = r; }); + trackStartingCardPublication(ds, post).catch(() => {}); + handleCotThinkingUpdate(ds, upd([think('pending')])); + handleCotThinkingUpdate(other, upd([think('independent')], 'om_other')); + await flush(); expect(request.mock.calls.filter(([r]) => r.method === 'POST')).toHaveLength(1); + reject(new Error('card failed')); await flush(); await flush(); + expect(request.mock.calls.filter(([r]) => r.method === 'POST')).toHaveLength(2); + }); + it('drops a superseded not-yet-visible bubble instead of placing it below the successor card', async () => { + const { trackStartingCardPublication } = await import('../src/core/starting-card-publication.js'); + const ds = makeDs(); let finish!: () => void; + trackStartingCardPublication(ds, new Promise(resolve => { finish = resolve; })); + handleCotThinkingUpdate(ds, upd([think('old')], 'om_old')); + handleCotThinkingUpdate(ds, upd([think('new')], 'om_new')); + finish(); await flush(); await flush(); + expect(request.mock.calls.filter(([r]) => r.method === 'POST')).toHaveLength(1); + expect(pushedEvents().some(e => e.content.delta === 'old')).toBe(false); + expect(pushedEvents().some(e => e.content.delta === 'new')).toBe(true); + }); + it('follows a pending successor card started while the predecessor POST settles', async () => { + const { trackStartingCardPublication } = await import('../src/core/starting-card-publication.js'); + const ds = makeDs(); let first!: () => void, second!: () => void; + const a = new Promise(resolve => { first = resolve; }); + const b = new Promise(resolve => { second = resolve; }); + trackStartingCardPublication(ds, a.then(() => { trackStartingCardPublication(ds, b); })); + handleCotThinkingUpdate(ds, upd([think('new')])); + first(); await flush(); expect(request).not.toHaveBeenCalled(); + second(); await flush(); await flush(); + expect(request.mock.calls.filter(([r]) => r.method === 'POST')).toHaveLength(1); + }); + it('keeps waiting for another card after one publication rejects', async () => { + const { trackStartingCardPublication } = await import('../src/core/starting-card-publication.js'); + const ds = makeDs(); + let fail!: (error: Error) => void; + let finish!: () => void; + const first = trackStartingCardPublication(ds, new Promise((_, reject) => { fail = reject; })); + trackStartingCardPublication(ds, new Promise(resolve => { finish = resolve; })); + handleCotThinkingUpdate(ds, upd([think('buffered')])); + fail(new Error('first card failed')); + await expect(first).rejects.toThrow('first card failed'); + await flush(); + expect(request).not.toHaveBeenCalled(); + finish(); await flush(); await flush(); + expect(request.mock.calls.filter(([req]) => req.method === 'POST')).toHaveLength(1); + expect(pushedEvents().some(event => event.content.delta === 'buffered')).toBe(true); + }); + it('does not publish a stopped turn after its pending card finishes', async () => { + const { trackStartingCardPublication } = await import('../src/core/starting-card-publication.js'); + const ds = makeDs(); + let finish!: () => void; + trackStartingCardPublication(ds, new Promise(resolve => { finish = resolve; })); + handleCotThinkingUpdate(ds, upd([think('cancelled')])); + abortCotMessage(ds); + finish(); await flush(); await flush(); + expect(request).not.toHaveBeenCalled(); + }); + it('drops a predecessor when the current turn changes before its next thinking update', async () => { + const { trackStartingCardPublication } = await import('../src/core/starting-card-publication.js'); + const ds = makeDs({ currentTurnId: 'om_turn1' }); + let finish!: () => void; + trackStartingCardPublication(ds, new Promise(resolve => { finish = resolve; })); + handleCotThinkingUpdate(ds, upd([think('old')])); + ds.currentTurnId = 'om_turn2'; + finish(); await flush(); await flush(); + expect(request).not.toHaveBeenCalled(); + }); + it('bounds a stuck card without blocking turn settlement or publishing a detached bubble', async () => { + const { trackStartingCardPublication } = await import('../src/core/starting-card-publication.js'); + vi.useFakeTimers(); + try { + const ds = makeDs(); let finish!: () => void; + trackStartingCardPublication(ds, new Promise(resolve => { finish = resolve; })); + handleCotThinkingUpdate(ds, upd([think('buffered')])); + await vi.advanceTimersByTimeAsync(15_000); + expect(request).not.toHaveBeenCalled(); + finish(); await flush(); + expect(request).not.toHaveBeenCalled(); + } finally { vi.useRealTimers(); } + }); +}); diff --git a/test/recall-frozen-cards.test.ts b/test/recall-frozen-cards.test.ts index 590f65ed91..d1e797b2ee 100644 --- a/test/recall-frozen-cards.test.ts +++ b/test/recall-frozen-cards.test.ts @@ -1,3 +1,4 @@ +import { pendingStartingCardPublication } from '../src/core/starting-card-publication.js'; /** * Unit tests for recallFrozenCards (worker-pool.ts). * @@ -660,6 +661,7 @@ describe('postTurnStartingCard', () => { activate(ds); const post = postTurnStartingCard(ds, sessionReply, 'om_turn_1'); + expect(pendingStartingCardPublication(ds)).toBeDefined(); expect(sessionReply).toHaveBeenCalledTimes(1); expect(buildStreamingCardMock.mock.calls[0]?.[5]).toBe('starting');