Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 40 additions & 0 deletions src/core/starting-card-publication.ts
Original file line number Diff line number Diff line change
@@ -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<DaemonSession, Set<Promise<void>>>();

export function trackStartingCardPublication<T>(ds: DaemonSession, post: Promise<T>): Promise<T> {
const pending = posts.get(ds) ?? new Set<Promise<void>>();
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<void> | undefined {
if (!posts.get(ds)?.size) return undefined;
return (async () => {
let timer: ReturnType<typeof setTimeout> | undefined;
try {
await Promise.race([
(async () => {
while (posts.get(ds)?.size) await Promise.all([...posts.get(ds)!]);
})(),
new Promise<never>((_, reject) => {
timer = setTimeout(() => reject(new Error('starting card publication timed out')), 15_000);
timer.unref?.();
}),
]);
} finally {
if (timer) clearTimeout(timer);
}
})();
}
4 changes: 4 additions & 0 deletions src/core/worker-pool.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -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;
}
Expand Down
13 changes: 13 additions & 0 deletions src/im/lark/cot-message.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -574,6 +575,18 @@ async function pump(ds: DaemonSession, state: CotState): Promise<void> {
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
Expand Down
102 changes: 102 additions & 0 deletions test/cot-message.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>(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<void>((_, 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<void>(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<void>(resolve => { first = resolve; });
const b = new Promise<void>(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<void>((_, reject) => { fail = reject; }));
trackStartingCardPublication(ds, new Promise<void>(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<void>(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<void>(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<void>(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(); }
});
});
2 changes: 2 additions & 0 deletions test/recall-frozen-cards.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { pendingStartingCardPublication } from '../src/core/starting-card-publication.js';
/**
* Unit tests for recallFrozenCards (worker-pool.ts).
*
Expand Down Expand Up @@ -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');

Expand Down
Loading