Repository navigation
fix(server): bound thread activity payload hydration #8991
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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)), | ||
|
|
@@ -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, | ||
|
|
@@ -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), | ||
| ${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, | ||
|
|
@@ -1342,7 +1369,14 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { | |
| activity.tone, | ||
| activity.kind, | ||
| activity.summary, | ||
| activity.payload_json AS "payload", | ||
| CASE | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🟠 High Oversized pinned 🤖 Copy this AI Prompt to have your agent fix this: |
||
| 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 | ||
|
|
@@ -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, | ||
|
|
@@ -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, | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🟠 High
Layers/ProjectionSnapshotQuery.ts:1046length(payload_json)counts Unicode code points for SQLiteTEXT, 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 aslength(CAST(payload_json AS BLOB))consistently for the cutoff, truncation marker, and cumulative budget.🤖 Copy this AI Prompt to have your agent fix this: