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
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ const encodeThreadLinkedPullRequest = Schema.encodeSync(
const encodeMessageContext = Schema.encodeEffect(
Schema.fromJsonString(OrchestrationMessageContext),
);
const encodeActivityPayload = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown));

it.effect("reads project shells without loading threads or resolving excluded projects", () => {
const resolved: string[] = [];
Expand Down Expand Up @@ -2902,6 +2903,97 @@ projectionSnapshotLayer("ProjectionSnapshotQuery windowed thread detail", (it) =
}),
);

it.effect("preserves agent lifecycle outside the turn window and activity cap", () =>
Effect.gen(function* () {
yield* seedFanOutThread();
const query = yield* ProjectionSnapshotQuery;
const sql = yield* SqlClient.SqlClient;
const expectedIds: string[] = [];

for (let index = 0; index < 17; index += 1) {
const taskId = `agent-${index}`;
const turnId = index < 8 ? "turn-1" : "turn-5";
const at = index < 8 ? "2026-03-01T00:00:30.000Z" : "2026-03-01T00:04:30.000Z";
for (const [offset, kind, payload] of [
[0, "task.started", { taskId, agentKind: "agent", title: taskId }],
[
1,
"task.progress",
{ taskId, agentKind: "agent", usageSnapshot: true, typedUsage: { totalTokens: 1000 } },
],
[2, "task.updated", { taskId, status: "idle" }],
[3, "tool.progress", { taskId, toolName: "Read" }],
] as const) {
const id = `${taskId}-${offset}`;
expectedIds.push(id);
yield* sql`
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
) VALUES (${id}, 'thread-w', ${turnId}, 'info', ${kind}, 'Agent activity',
${encodeActivityPayload(payload)}, ${index * 4 + offset}, ${at})
`;
}
}

const assertAgents = (activities: ReadonlyArray<{ id: string }>) => {
const ids = activities.map((activity) => activity.id);
assert.deepStrictEqual(
ids.filter((id) => id.startsWith("agent-")),
expectedIds,
);
assert.equal(ids.length, new Set(ids).size);
assert.ok(!ids.includes("old-background"));
};

yield* sql`
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
) VALUES
('malformed-task', 'thread-w', 'turn-1', 'info', 'task.started', 'Legacy task',
'{invalid', 68, '2026-03-01T00:00:30.000Z'),
('malformed-heartbeat', 'thread-w', 'turn-1', 'info', 'tool.progress', 'Legacy heartbeat',
'{invalid', 69, '2026-03-01T00:00:30.000Z')
`;

const window = Option.getOrThrow(
yield* query.getThreadDetailSnapshot(threadW, { turnLimit: 1 }),
);
assertAgents(window.thread.activities);
assert.ok(!window.thread.messages.some((message) => message.id === "user-msg-1"));
const older = Option.getOrThrow(
yield* query.getThreadDetailSnapshot(threadW, {
turnLimit: 1,
beforeCursor: window.page!.beforeCursor!,
}),
);
assertAgents(older.thread.activities);

yield* sql`
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
) VALUES ('old-background', 'thread-w', 'turn-1', 'info', 'task.started', 'Shell',
'{"taskId":"shell","agentKind":"background"}', 70, '2026-03-01T00:00:30.000Z')
`;
yield* sql`
WITH RECURSIVE rows(n) AS (
SELECT 1 UNION ALL SELECT n + 1 FROM rows WHERE n < 501
)
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
) SELECT 'noise-' || n, 'thread-w', 'turn-5', 'tool', 'tool.completed', 'Tool', '{}',
100 + n, '2026-03-01T00:04:40.000Z' FROM rows
`;
const detail = Option.getOrThrow(yield* query.getThreadDetailById(threadW));
const snapshot = Option.getOrThrow(
yield* query.getThreadDetailSnapshot(threadW, { turnLimit: 1 }),
);
for (const activities of [detail.activities, snapshot.thread.activities]) {
assertAgents(activities);
assert.equal(activities.length, 500 + expectedIds.length);
}
}),
);

it.effect("bounds activity hydration and preserves unresolved requests", () =>
Effect.gen(function* () {
yield* seedFanOutThread();
Expand Down
21 changes: 18 additions & 3 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1868,6 +1868,20 @@ pending_approval_requests AS (
AND json_extract(activity.payload_json, '$.requestId') IS NOT NULL
),
pinned_activity_ids AS (
SELECT activity_id
FROM projection_thread_activities
WHERE thread_id = ${threadId}
AND kind IN ('task.started', 'task.progress', 'task.updated', 'task.completed', 'tool.progress')
AND json_valid(payload_json)
AND json_extract(payload_json, '$.taskId') IN (
SELECT json_extract(payload_json, '$.taskId')
FROM projection_thread_activities
WHERE thread_id = ${threadId}
AND kind IN ('task.started', 'task.progress', 'task.updated', 'task.completed')
AND json_valid(payload_json)
AND json_extract(payload_json, '$.agentKind') = 'agent'
)
UNION ALL
SELECT activity_id
FROM pending_approval_activities
WHERE request_order = 1
Expand All @@ -1879,9 +1893,10 @@ pending_approval_requests AS (
)
`;

// Blocking request payloads must remain available even if they predate the
// recent activity window. Each CTE returns at most one unresolved row per
// request, so the merge below stays bounded by actionable work.
// Agent lifecycle and blocking requests outlive chat pagination. Keep every
// lifecycle row for known agents, including status rows without an agentKind
// stamp, so fresh clients fold the same roster and usage as live clients.
// Each request CTE still returns at most one unresolved row per request.
const listPinnedThreadActivityRowsByThread = SqlSchema.findAll({
Request: ThreadIdLookupInput,
Result: ProjectionThreadActivityDbRowSchema,
Expand Down