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
30 changes: 30 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
);
}),
);

Expand Down
8 changes: 4 additions & 4 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2666,7 +2666,7 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
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
Expand All @@ -2675,7 +2675,7 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
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 (
Expand Down Expand Up @@ -3155,7 +3155,7 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
localWindow !== undefined &&
localWindow.rowLimit > 0 &&
projection.turnItems.length >= localWindow.rowLimit &&
!projection.turnItems.some(isThreadHistoryTurnStart)
!projection.turnItems.some(isThreadHistoryUserTurn)
) {
return withLocalVisibleTurnItems(projection);
}
Expand Down Expand Up @@ -6048,7 +6048,7 @@ export const layerMemory: Layer.Layer<ProjectionStoreV2> = Layer.effect(
);
const anchorLimit = (options.userTurnLimit ?? 0) + 2;
const start =
turnAnchors.length > 0
anchors.length > 0
? anchors.length < anchorLimit
? rawStart
: anchors.at(-anchorLimit)!
Expand Down
30 changes: 30 additions & 0 deletions apps/server/src/orchestration-v2/threadHistoryPaging.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/orchestration-v2/threadHistoryPaging.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Loading