Skip to content
Merged
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
63 changes: 63 additions & 0 deletions apps/desktop/src/main/__tests__/message-queue-ui-state.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -203,6 +203,69 @@ test('queue_update events drive the independent desktop queue projection', () =>
assert.equal(transientMessages.size, 0);
});

test('a rootless resubscription seed retires a stale queued card', () => {
// Switch away → the queue drains rootless → navigate back. The projector's
// rootless seed now carries the authoritative queue (apache/maka#5520
// review), so the card the client kept from before it left must go.
const controller = createAppShellSessionUiStateController();
const handlers = createAppShellSessionEventHandlers({
uiLocale: 'zh-CN',
activeIdRef: { current: 'session-1' },
liveTurnBySessionRef: controller.liveTurnBySessionRef,
refreshMessages: async () => true,
refreshSessions: async () => [],
setLiveTurnBySession: controller.setLiveTurnBySession,
setInteractionBySession: controller.setInteractionBySession,
setMessageQueueBySession: controller.setMessageQueueBySession,
removeTransientMessage: () => {},
showModelSetupToast() {},
toastApi: { error() {} },
});

handlers.handleEvent('session-1', {
type: 'queue_update',
id: 'queue-1',
turnId: 'turn-1',
ts: 1,
queueRevision: 3,
steering: ['adjust this run'],
followup: [],
steeringEntries: [
{
entryId: 'entry-steer',
messageId: 'message-steer',
content: { text: 'adjust this run' },
placement: 'current_turn' as const,
state: 'queued' as const,
},
],
followupEntries: [],
});
assert.ok(
controller.getState().messageQueueBySession['session-1'],
'the card is visible before the client leaves',
);

// The resubscription seed's authoritative empty queue: the drain landed
// while the Session was inactive, and the root Turn is gone.
handlers.handleEvent('session-1', {
type: 'queue_update',
id: 'host-queue:host-1:4',
turnId: '',
ts: 2,
queueRevision: 4,
steering: [],
followup: [],
steeringEntries: [],
followupEntries: [],
});
assert.equal(
controller.getState().messageQueueBySession['session-1'],
undefined,
'the stale card does not survive the resubscription',
);
});

test('steering delivery clears a promoted follow-up from the desktop queue', () => {
const controller = createAppShellSessionUiStateController();
const handlers = createAppShellSessionEventHandlers({
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,13 @@ test('projects root lifecycle without fabricating content events', async (t) =>
if (seed?.type === 'host_observation_seed') {
assert.deepEqual(seed.observerIds, ['execution-observer']);
assert.equal(seed.execution.rootTurn, null);
assert.deepEqual(seed.events, []);
// The rootless seed carries the authoritative queue once, so a stale
// queued card from an earlier observation retires on re-subscription
// (apache/maka#5520 review).
assert.deepEqual(
seed.events.map((event) => event.type),
['queue_update'],
);
}
const seededCount = messages.length;
events.push({
Expand Down
57 changes: 57 additions & 0 deletions packages/runtime-host/src/__tests__/session-projector.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -431,6 +431,63 @@ test('reseeds an empty queue after queued successors completed while disconnecte
assert.deepEqual(queue.followupEntries, []);
});

test('projects a queue drain that lands while no root Turn is live', () => {
// apache/maka#5520: a drain observed after the root Turn is gone must still
// reach the renderer, or a phantom queued card survives whose retract fails
// with not_found forever. Seeding stays silent for rootless snapshots — the
// Desktop observer pins an empty seed there — because a client that never
// observed the session has no stale card to clear.
const projector = new RuntimeHostSessionProjector(
snapshot({ queue: queue(2, [steeringEntry('queued')]) }),
createRuntimeHostSessionProjectionSeed([], snapshot()),
() => 10,
);

const drained = projector.accept({
kind: 'subscription.session_projection',
hostEpoch: 'host-1',
subscriptionId: 'subscription-1',
sequence: 1,
snapshot: snapshot({ projectionRevision: 2, rootTurn: null, queue: queue(3, []) }),
});
assert.deepEqual(
drained.events.map((event) => event.type),
['queue_update'],
);
const update = drained.events.find(
(event): event is Extract<SessionEvent, { type: 'queue_update' }> =>
event.type === 'queue_update',
);
assert.ok(update, 'the drained queue must be projected');
assert.deepEqual(update.steering, []);
assert.deepEqual(update.followup, []);
});

test('seeding a rootless snapshot conveys the authoritative queue', () => {
// A Desktop that navigates away unsubscribes; if the queue drains while the
// Session is inactive, the resubscribing client's stale queued card survives
// until a queue_update that the rootless seed never produced (apache/maka
// #5520 review). The rootless seed must carry the authoritative queue once.
const projector = new RuntimeHostSessionProjector(
snapshot({ rootTurn: null, queue: queue(3, []) }),
createRuntimeHostSessionProjectionSeed([], snapshot()),
() => 10,
);

const seeded = projector.seedActive(true);
assert.deepEqual(
seeded.map((event) => event.type),
['queue_update'],
);
const update = seeded.find(
(event): event is Extract<SessionEvent, { type: 'queue_update' }> =>
event.type === 'queue_update',
);
assert.ok(update);
assert.deepEqual(update.steering, []);
assert.deepEqual(update.followup, []);
});

test('reseeds the latest provider retry when the active Turn still carries one', () => {
const retry = {
phase: 'scheduled' as const,
Expand Down
29 changes: 26 additions & 3 deletions packages/runtime-host/src/adapter/session-projector.ts
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,20 @@ export class RuntimeHostSessionProjector {

seedActive(includeAssistantText: boolean): SessionEvent[] {
const root = this.#snapshot.rootTurn;
if (!root) return [];
if (!root) {
// A Session whose root Turn is gone still owns an authoritative queue:
// a client resubscribing after navigating away may hold a stale queued
// card that only this seed can retire, because no live drain will run
// while it is the active view (apache/maka#5520 review). Snapshots that
// carry no queue at all project as empty.
const queue = this.#snapshot.queue ?? {
hostEpoch: '',
queueRevision: 0,
steering: [],
followup: [],
};
return [projectQueueUpdate(this.#unplacedQueue(queue), '', this.#now())];
}
const events: SessionEvent[] = [];
const queueEvents =
this.#projectMessageAdmissions || queueHasEntries(this.#snapshot.queue)
Expand Down Expand Up @@ -403,8 +416,18 @@ export class RuntimeHostSessionProjector {
for (const interaction of newlyPendingInteractions(previousSnapshot, next)) {
events.push(...projectRuntimeHostInteractionRequest(interaction, this.#now()));
}
if (root && queueChanged(previousSnapshot.queue, next.queue)) {
events.push(projectQueueUpdate(this.#unplacedQueue(next.queue), root.turnId, this.#now()));
if (queueChanged(previousSnapshot.queue, next.queue)) {
// Project the authoritative queue even with no live Turn: a drain that
// lands after the root Turn is gone must still reach observers, or a
// queued card survives as a phantom whose retract fails with not_found
// (apache/maka#5520).
events.push(
projectQueueUpdate(
this.#unplacedQueue(next.queue),
root?.turnId ?? previousRoot?.turnId ?? '',
this.#now(),
),
);
}
if (startedTurn) this.#accumulators.clear();
// Emit the presentation-only compaction-started event when the root Turn
Expand Down