From c2d88faef4cd994fd18d91d8119e6409694a6734 Mon Sep 17 00:00:00 2001 From: Bil0000 <62337003+Bil0000@users.noreply.github.com> Date: Mon, 28 Sep 2026 17:05:22 +0200 Subject: [PATCH] fix(orchestration-v2): page agent-only transcripts by activity --- .../orchestration-v2/ProjectionStore.test.ts | 30 +++++++++++++++++++ .../src/orchestration-v2/ProjectionStore.ts | 8 ++--- .../threadHistoryPaging.test.ts | 30 +++++++++++++++++++ .../orchestration-v2/threadHistoryPaging.ts | 2 +- 4 files changed, 65 insertions(+), 5 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index f3b6348c5d31..25a5e9e3d0b8 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -1377,6 +1377,36 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { (row) => row.sourceItemId === interruptResultId, ), ); + + yield* sql` + UPDATE orchestration_v2_projection_turn_items + SET payload_json = json_set(payload_json, '$.createdBy', 'agent') + WHERE thread_id = ${threadId} AND type = 'user_message' + `; + const agentPromptId = TurnItemId.make("turn-item:bounded-sql-history:interrupt-filler:1281"); + yield* sql` + UPDATE orchestration_v2_projection_turn_items + SET type = 'user_message', + payload_json = json_set(payload_json, + '$.type', 'user_message', '$.inputIntent', 'turn_start', + '$.createdBy', 'agent', '$.creationSource', 'provider', + '$.messageId', 'message:bounded-sql-history:agent-prompt', + '$.text', 'Continue the child task', '$.attachments', json('[]')) + WHERE turn_item_id = ${agentPromptId} + `; + const agentWindow = yield* projectionStore.getThreadSnapshotWindow(threadId, { + rowLimit: sqlPageLimit, + userTurnLimit: THREAD_HISTORY_PAGE_POLICY.maxUserTurns, + }); + assert.lengthOf( + agentWindow.projection.visibleTurnItems.filter( + (row) => row.item.type !== "run_interrupt_request", + ), + sqlPageLimit, + ); + assert.isTrue( + agentWindow.projection.visibleTurnItems.some((row) => row.sourceItemId === agentPromptId), + ); }), ); diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 6718240d9e7a..f12f56572264 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -2666,7 +2666,7 @@ export const layer: Layer.Layer = WHEN (SELECT COUNT(*) FROM turn_anchors) >= ${THREAD_HISTORY_MAX_RAW_TURNS + 2} THEN (SELECT MIN(ordinal) FROM turn_anchors) ELSE 0 - END AS ordinal, (SELECT COUNT(*) FROM turn_anchors) AS anchors + END AS ordinal FROM user_anchors ), selected AS ( SELECT payload_json, ordinal, turn_item_id, run_id, type @@ -2675,7 +2675,7 @@ export const layer: Layer.Layer = ORDER BY ordinal DESC, turn_item_id DESC LIMIT CASE WHEN ${window.rowLimit} = 0 THEN 0 - WHEN (SELECT anchors FROM boundary) > 0 THEN -1 + WHEN (SELECT COUNT(*) FROM user_anchors) > 0 THEN -1 ELSE ${window.rowLimit} END ), retained AS ( @@ -3155,7 +3155,7 @@ export const layer: Layer.Layer = localWindow !== undefined && localWindow.rowLimit > 0 && projection.turnItems.length >= localWindow.rowLimit && - !projection.turnItems.some(isThreadHistoryTurnStart) + !projection.turnItems.some(isThreadHistoryUserTurn) ) { return withLocalVisibleTurnItems(projection); } @@ -6048,7 +6048,7 @@ export const layerMemory: Layer.Layer = Layer.effect( ); const anchorLimit = (options.userTurnLimit ?? 0) + 2; const start = - turnAnchors.length > 0 + anchors.length > 0 ? anchors.length < anchorLimit ? rawStart : anchors.at(-anchorLimit)! diff --git a/apps/server/src/orchestration-v2/threadHistoryPaging.test.ts b/apps/server/src/orchestration-v2/threadHistoryPaging.test.ts index cccc826478b7..9c1115ca0df7 100644 --- a/apps/server/src/orchestration-v2/threadHistoryPaging.test.ts +++ b/apps/server/src/orchestration-v2/threadHistoryPaging.test.ts @@ -209,6 +209,36 @@ describe("threadHistoryPaging", () => { ); }); + it("pages agent-only child transcripts instead of dropping their earlier activity", () => { + const commandRows = Array.from({ length: 90 }, (_, index) => makeRow(index + 1)); + const first = makeRow(0); + const prompt = { + ...first, + item: { + ...first.item, + type: "user_message" as const, + createdBy: "agent" as const, + creationSource: "provider" as const, + inputIntent: "turn_start" as const, + messageId: MessageId.make("child-prompt"), + text: "Inspect this project", + attachments: [], + }, + } as OrchestrationV2ProjectedTurnItem; + const items = [prompt, ...commandRows]; + const recent = selectRecentTimelineWindow({ items, snapshotSequence: 1 }); + + expect(recent.items).toHaveLength(THREAD_HISTORY_PAGE_POLICY.maxItems); + expect(recent.hasMoreHistory).toBe(true); + const older = selectHistoryPageFromCursor({ + items, + cursor: recent.nextCursor!, + snapshotSequence: 1, + }); + expect(older.items[0]?.item.type).toBe("user_message"); + expect([...older.items, ...recent.items]).toHaveLength(items.length); + }); + it("encodes opaque cursors with stable source identity", () => { const cursor = encodeThreadHistoryCursor({ snapshotSequence: 9, diff --git a/apps/server/src/orchestration-v2/threadHistoryPaging.ts b/apps/server/src/orchestration-v2/threadHistoryPaging.ts index 98ae11afdad0..a2e6c87e20df 100644 --- a/apps/server/src/orchestration-v2/threadHistoryPaging.ts +++ b/apps/server/src/orchestration-v2/threadHistoryPaging.ts @@ -185,7 +185,7 @@ function selectOlderTimelinePage(input: { let encodedBytes = 0; let userTurns = 0; let rawTurns = 0; - const turnLimit = input.items.slice(0, end).some((row) => isThreadHistoryTurnStart(row.item)) + const turnLimit = input.items.slice(0, end).some((row) => isThreadHistoryUserTurn(row.item)) ? policy.maxUserTurns : undefined; for (let index = end - 1; index >= 0; index -= 1) {