Skip to content
Closed
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 @@ -2421,6 +2421,71 @@ projectionSnapshotLayer("ProjectionSnapshotQuery windowed thread detail", (it) =
}),
);

it.effect("bounds activity hydration by serialized payload bytes", () =>
Effect.gen(function* () {
yield* seedFanOutThread();
const snapshotQuery = yield* ProjectionSnapshotQuery;
const sql = yield* SqlClient.SqlClient;

yield* sql`DELETE FROM projection_thread_activities`;
yield* sql`
WITH RECURSIVE activity_rows(sequence) AS (
SELECT 1
UNION ALL
SELECT sequence + 1 FROM activity_rows WHERE sequence < 31
)
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
)
SELECT
printf('activity-%04d', sequence),
'thread-w',
'turn-5',
'tool',
'tool.completed',
'ran tool',
json_object(
'output',
printf('%.*c', CASE WHEN sequence = 31 THEN 700000 ELSE 400000 END, 'x')
),
sequence,
'2026-03-01T00:04:00.000Z'
FROM activity_rows
`;

const assertBoundedActivities = (
activities: ReadonlyArray<{ readonly id: string; readonly payload: unknown }>,
) => {
assert.equal(activities.length, 20);
assert.equal(activities[0]?.id, asEventId("activity-0012"));
assert.equal(activities.at(-1)?.id, asEventId("activity-0031"));
assert.deepStrictEqual(activities.at(-1)?.payload, {
t3PayloadTruncated: true,
originalBytes: 700013,
});
const serializedPayloadBytes = activities.reduce(
(total, activity) => total + Buffer.byteLength(JSON.stringify(activity.payload)),
0,
);
assert.ok(serializedPayloadBytes <= 8 * 1024 * 1024);
};

const fullDetail = yield* snapshotQuery.getThreadDetailById(threadW);
assert.equal(fullDetail._tag, "Some");
if (fullDetail._tag === "Some") {
assertBoundedActivities(fullDetail.value.activities);
}

const windowedDetail = yield* snapshotQuery.getThreadDetailSnapshot(threadW, {
turnLimit: 2,
});
assert.equal(windowedDetail._tag, "Some");
if (windowedDetail._tag === "Some") {
assertBoundedActivities(windowedDetail.value.thread.activities);
}
}),
);

it.effect("a thread with no turns returns its content unwindowed on the first page", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
Expand Down
132 changes: 96 additions & 36 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -70,10 +70,11 @@ import {
const decodeReadModel = Schema.decodeUnknownEffect(OrchestrationReadModel);
const decodeShellSnapshot = Schema.decodeUnknownEffect(OrchestrationShellSnapshot);
const decodeThread = Schema.decodeUnknownEffect(OrchestrationThread);
// Keep detail reads consistent with the in-memory projector's retained
// activity window. Applying the limit in SQL avoids decoding an unbounded
// payload_json set before the projector can enforce that invariant.
// Apply both row and serialized-payload limits in SQL so oversized tool output
// cannot exhaust the server heap before JSON decoding can enforce a boundary.
const THREAD_DETAIL_ACTIVITY_LIMIT = 500;
const THREAD_DETAIL_ACTIVITY_PAYLOAD_LIMIT_BYTES = 512 * 1024;
const THREAD_DETAIL_ACTIVITY_PAYLOAD_BUDGET_BYTES = 8 * 1024 * 1024;
const ProjectionProjectDbRowSchema = ProjectionProject.mapFields(
Struct.assign({
defaultModelSelection: Schema.NullOr(Schema.fromJsonString(ModelSelection)),
Expand Down Expand Up @@ -1019,17 +1020,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
Result: ProjectionThreadActivityDbRowSchema,
execute: ({ threadId }) =>
sql`
SELECT
activity_id AS "activityId",
thread_id AS "threadId",
turn_id AS "turnId",
tone,
kind,
summary,
payload_json AS "payload",
sequence,
created_at AS "createdAt"
FROM (
WITH recent_activities AS (
SELECT
activity_id,
thread_id,
Expand All @@ -1042,12 +1033,48 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
created_at
FROM projection_thread_activities
WHERE thread_id = ${threadId}
ORDER BY
sequence DESC,
created_at DESC,
activity_id DESC
ORDER BY sequence DESC, created_at DESC, activity_id DESC
LIMIT ${THREAD_DETAIL_ACTIVITY_LIMIT}
) AS recent_activities
),
ordered_activities AS (
SELECT
*,
ROW_NUMBER() OVER (
ORDER BY sequence DESC, created_at DESC, activity_id DESC
) AS recent_order,
MIN(
length(payload_json),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 High Layers/ProjectionSnapshotQuery.ts:1046

length(payload_json) counts Unicode code points for SQLite TEXT, not UTF-8 bytes, so multi-byte JSON payloads bypass the 512 KiB per-payload cutoff and the 8 MiB cumulative budget; this can hydrate tens of MiB. Use a byte-length expression such as length(CAST(payload_json AS BLOB)) consistently for the cutoff, truncation marker, and cumulative budget.

🤖 Copy this AI Prompt to have your agent fix this:
In file @apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts around line 1046:

`length(payload_json)` counts Unicode code points for SQLite `TEXT`, not UTF-8 bytes, so multi-byte JSON payloads bypass the 512 KiB per-payload cutoff and the 8 MiB cumulative budget; this can hydrate tens of MiB. Use a byte-length expression such as `length(CAST(payload_json AS BLOB))` consistently for the cutoff, truncation marker, and cumulative budget.

${THREAD_DETAIL_ACTIVITY_PAYLOAD_LIMIT_BYTES}
) AS hydrated_payload_bytes
FROM recent_activities
),
budgeted_activities AS (
SELECT
*,
SUM(hydrated_payload_bytes) OVER (
ORDER BY recent_order ASC
) AS cumulative_payload_bytes
FROM ordered_activities
)
SELECT
activity_id AS "activityId",
thread_id AS "threadId",
turn_id AS "turnId",
tone,
kind,
summary,
CASE
WHEN length(payload_json) > ${THREAD_DETAIL_ACTIVITY_PAYLOAD_LIMIT_BYTES}
THEN json_object(
't3PayloadTruncated', json('true'),
'originalBytes', length(payload_json)
)
ELSE payload_json
END AS "payload",
sequence,
created_at AS "createdAt"
FROM budgeted_activities
WHERE cumulative_payload_bytes <= ${THREAD_DETAIL_ACTIVITY_PAYLOAD_BUDGET_BYTES}
ORDER BY
sequence ASC,
created_at ASC,
Expand Down Expand Up @@ -1342,7 +1369,14 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
activity.tone,
activity.kind,
activity.summary,
activity.payload_json AS "payload",
CASE

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 High Layers/ProjectionSnapshotQuery.ts:1372

Oversized pinned approval.requested and user-input.requested activities are returned without requestId (and without questions for user input), so the client cannot derive an actionable request and the provider turn remains blocked. The truncation branch replaces the entire payload with a generic marker; keep pinned unresolved-request payloads intact (or truncate while preserving the required request fields).

🤖 Copy this AI Prompt to have your agent fix this:
In file @apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts around line 1372:

Oversized pinned `approval.requested` and `user-input.requested` activities are returned without `requestId` (and without `questions` for user input), so the client cannot derive an actionable request and the provider turn remains blocked. The truncation branch replaces the entire payload with a generic marker; keep pinned unresolved-request payloads intact (or truncate while preserving the required request fields).

WHEN length(activity.payload_json) > ${THREAD_DETAIL_ACTIVITY_PAYLOAD_LIMIT_BYTES}
THEN json_object(
't3PayloadTruncated', json('true'),
'originalBytes', length(activity.payload_json)
)
ELSE activity.payload_json
END AS "payload",
activity.sequence,
activity.created_at AS "createdAt"
FROM pinned_activity_ids AS pinned
Expand All @@ -1357,17 +1391,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
Result: ProjectionThreadActivityDbRowSchema,
execute: ({ threadId, minAnchorAt, minTurnKey, beforeAnchorAt, beforeTurnKey }) =>
sql`
SELECT
activity_id AS "activityId",
thread_id AS "threadId",
turn_id AS "turnId",
tone,
kind,
summary,
payload_json AS "payload",
sequence,
created_at AS "createdAt"
FROM (
WITH recent_activities AS (
SELECT
activity_id,
thread_id,
Expand Down Expand Up @@ -1406,12 +1430,48 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
AND created_at < ${beforeAnchorAt}
)
)
ORDER BY
sequence DESC,
created_at DESC,
activity_id DESC
ORDER BY sequence DESC, created_at DESC, activity_id DESC
LIMIT ${THREAD_DETAIL_ACTIVITY_LIMIT}
) AS recent_activities
),
ordered_activities AS (
SELECT
*,
ROW_NUMBER() OVER (
ORDER BY sequence DESC, created_at DESC, activity_id DESC
) AS recent_order,
MIN(
length(payload_json),
${THREAD_DETAIL_ACTIVITY_PAYLOAD_LIMIT_BYTES}
) AS hydrated_payload_bytes
FROM recent_activities
),
budgeted_activities AS (
SELECT
*,
SUM(hydrated_payload_bytes) OVER (
ORDER BY recent_order ASC
) AS cumulative_payload_bytes
FROM ordered_activities
)
SELECT
activity_id AS "activityId",
thread_id AS "threadId",
turn_id AS "turnId",
tone,
kind,
summary,
CASE
WHEN length(payload_json) > ${THREAD_DETAIL_ACTIVITY_PAYLOAD_LIMIT_BYTES}
THEN json_object(
't3PayloadTruncated', json('true'),
'originalBytes', length(payload_json)
)
ELSE payload_json
END AS "payload",
sequence,
created_at AS "createdAt"
FROM budgeted_activities
WHERE cumulative_payload_bytes <= ${THREAD_DETAIL_ACTIVITY_PAYLOAD_BUDGET_BYTES}
ORDER BY
sequence ASC,
created_at ASC,
Expand Down
Loading