perf(server,web): cursor-paginate session history with tail-first loading (#202)

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Yaowei Zheng
2026-08-05 00:16:25 +08:00
committed by GitHub
parent fc7cb5b8ed
commit dd8ddbb69b
22 changed files with 2153 additions and 80 deletions
+44 -1
View File
@@ -661,15 +661,58 @@ export interface MessagesLiveTail {
fragments: OmniMessage[];
}
/**
* Pagination envelope of a windowed `GET /messages` (`tailLimit` / `before` requests
* only; the parameterless full read never carries it). A window is a run of whole
* message-bearing units — one unit = one Task in the Web reducer's sense, opened by a
* main-session user prompt — cut so that no pairing (tool_call/output), compaction span
* or steering group ever splits across windows.
*/
export interface MessagesPageInfo {
/**
* Cursor of this window's first unit (`<shardIndex>:<ordinal>`): pass it back as
* `before=` to fetch the previous window. Stable across requests and compaction —
* rotation opens a NEW shard and closed shards are immutable. Absent = this window
* reaches the very beginning of the transcript (no older history).
*/
before?: string;
/**
* Outline turns (the Web conversation outline's entry rule) opened BEFORE this
* window: the client offsets its global `第 N 轮` numbering by this, so a partial
* window never mis-numbers. 0 when the window starts at the beginning.
*/
earlierTurns: number;
/**
* Cumulative stats accrued before this window, seeded into the client's stats
* tracker so header chips and per-turn cumulative rows equal a full load:
* finished-Task elapsed, subagent token totals, and the last main-session
* session/context token readings.
*/
prior: {
subagentTokens: number;
elapsedMs: number;
sessionTokens: number;
contextTokens: number;
};
}
/** Message history: the full messages and events from concatenating all of this Session's Trace files in order (excludes partial_*). */
export interface MessagesResponse {
messages: OmniMessage[];
/**
* Present only while the Session is running/compacting: the in-progress stream tail
* (open streaming fragments + the channel cursor they cover), so a client joining
* mid-stream can render the currently streaming message. Omitted when idle.
* mid-stream can render the currently streaming message. Omitted when idle. On
* windowed requests it rides TAIL pages only — a `before` page is immutable history
* and never carries it.
*/
live?: MessagesLiveTail;
/**
* Present exactly on windowed requests (`tailLimit` / `before`): `messages` is then
* the requested window (subagent pointers inside it expanded as usual) rather than
* the full transcript. See MessagesPageInfo.
*/
page?: MessagesPageInfo;
}
// ---------------------------------------------------------------------------
+1
View File
@@ -30,6 +30,7 @@ export function openDatabase(dbPath: string): DatabaseSync {
ensureColumn(db, "sessions", "client", "TEXT");
ensureColumn(db, "sessions", "has_trace", "INTEGER NOT NULL DEFAULT 0");
ensureColumn(db, "auth_sessions", "via", "TEXT");
ensureColumn(db, "trace_files", "page_stats", "TEXT");
return db;
}
+44 -4
View File
@@ -62,17 +62,22 @@ export class TraceIndexRepo {
constructor(private readonly db: DatabaseSync) {}
upsertFile(row: TraceFileRow): void {
// page_stats is a derivation of the shard's CONTENT: any observed size change voids
// it (an unchanged size keeps the cached scan — reconcile passes re-upsert rows).
this.db
.prepare(
`INSERT INTO trace_files (project_id, agent_id, session_id, file_index, date, size_bytes)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT(project_id, agent_id, session_id, file_index)
DO UPDATE SET date = excluded.date, size_bytes = excluded.size_bytes`,
DO UPDATE SET date = excluded.date,
page_stats = CASE WHEN excluded.size_bytes = trace_files.size_bytes
THEN trace_files.page_stats ELSE NULL END,
size_bytes = excluded.size_bytes`,
)
.run(row.projectId, row.agentId, row.sessionId, row.fileIndex, row.date, row.sizeBytes);
}
/** Refreshes one shard's last-observed size (listings write back their page stats). */
/** Refreshes one shard's last-observed size (listings write back their page stats); a changed size voids the cached window scan. */
updateFileSize(
projectId: string,
agentId: string,
@@ -82,10 +87,45 @@ export class TraceIndexRepo {
): void {
this.db
.prepare(
`UPDATE trace_files SET size_bytes = ?
`UPDATE trace_files SET
page_stats = CASE WHEN size_bytes = ? THEN page_stats ELSE NULL END,
size_bytes = ?
WHERE project_id = ? AND agent_id = ? AND session_id = ? AND file_index = ?`,
)
.run(sizeBytes, projectId, agentId, sessionId, fileIndex);
.run(sizeBytes, sizeBytes, projectId, agentId, sessionId, fileIndex);
}
/** Cached message-window prefix record of one shard (see services/message-window.ts); null = none. */
getPageStats(
projectId: string,
agentId: string,
sessionId: string,
fileIndex: number,
): string | null {
const r = this.db
.prepare(
`SELECT page_stats FROM trace_files
WHERE project_id = ? AND agent_id = ? AND session_id = ? AND file_index = ?`,
)
.get(projectId, agentId, sessionId, fileIndex);
return r ? ((r.page_stats as string | null) ?? null) : null;
}
/** Stores a shard's message-window prefix record, guarded on the size it was computed for (a concurrent append voids the write). */
setPageStats(
projectId: string,
agentId: string,
sessionId: string,
fileIndex: number,
sizeBytes: number,
pageStats: string,
): void {
this.db
.prepare(
`UPDATE trace_files SET page_stats = ?
WHERE project_id = ? AND agent_id = ? AND session_id = ? AND file_index = ? AND size_bytes = ?`,
)
.run(pageStats, projectId, agentId, sessionId, fileIndex, sizeBytes);
}
deleteFile(projectId: string, agentId: string, sessionId: string, fileIndex: number): void {
+1
View File
@@ -128,6 +128,7 @@ CREATE TABLE IF NOT EXISTS trace_files ( -- DERIVED CACHE of the on-disk T
file_index INTEGER NOT NULL, -- shard index NNN of <session_id>_NNN.jsonl (name/path are reconstructed from the row, never stored: the data root may move)
date TEXT NOT NULL, -- date directory name (local yyyy-mm-dd)
size_bytes INTEGER NOT NULL, -- last observed size (listings stat-refresh the returned page and write back; an actively-appended shard may lag in between)
page_stats TEXT, -- cached message-window scan state at this shard's END (services/message-window.ts ShardPrefixRecord JSON); immutable shards only, NULL = not yet computed, nulled whenever size_bytes changes
PRIMARY KEY (project_id, agent_id, session_id, file_index)
);
CREATE INDEX IF NOT EXISTS idx_trace_files_agent_date ON trace_files(project_id, agent_id, date);
+78 -1
View File
@@ -17,6 +17,7 @@ import type {
FilesStatResponse,
GoalResponse,
MessagesLiveTail,
MessagesPageInfo,
MessagesResponse,
ServerEvent,
SessionCategory,
@@ -26,6 +27,8 @@ import type {
RetryNowResponse,
TaskCreateResponse,
} from "../../api/types.js";
import { decodeCursor } from "../../services/message-window.js";
import type { MessagesPageRequest } from "../../services/trace-service.js";
import { PREVIEW_TOKEN_TTL_MS, resolvePreviewTarget } from "../../services/preview-token.js";
import type { AppEnv } from "../../auth/middleware.js";
import type { SessionRow } from "../../db/repos/sessions.js";
@@ -72,6 +75,51 @@ export const APPROVAL_MODES: readonly ApprovalMode[] = [
/** The five valid per-turn thinking level names (TaskCreateRequest.thinkingLevel). */
const THINKING_LEVELS: readonly ThinkingLevelName[] = ["none", "low", "medium", "high", "xhigh"];
/** Unit-count bounds for windowed history reads (`tailLimit` / `limit`), and the `before` page's default. */
const MESSAGES_PAGE_LIMIT_MAX = 1000;
const MESSAGES_PAGE_LIMIT_DEFAULT = 200;
/** Parse one windowed-read unit-count param (positive integer, capped). */
function pageLimit(raw: string, name: string): number {
if (!/^\d{1,4}$/.test(raw)) throw badRequest(`${name} must be a positive integer.`);
const v = Number.parseInt(raw, 10);
if (v < 1 || v > MESSAGES_PAGE_LIMIT_MAX) {
throw badRequest(`${name} must be an integer between 1 and ${MESSAGES_PAGE_LIMIT_MAX}.`);
}
return v;
}
/**
* Parse GET /messages windowed-read params. No params → null: the legacy full-transcript
* read, byte-identical to the pre-pagination response (other consumers depend on it).
* `tailLimit=<n>` → the newest n units; `before=<cursor>[&limit=<n>]` → the n units
* preceding the cursor. The two forms are mutually exclusive, and `limit` belongs to
* `before` alone — mixing them is a caller bug worth a loud 400 rather than a guess.
*/
function messagesPageQuery(c: Context): MessagesPageRequest | null {
const rawTail = c.req.query("tailLimit");
const rawBefore = c.req.query("before");
const rawLimit = c.req.query("limit");
if (rawTail === undefined && rawBefore === undefined) {
if (rawLimit !== undefined) throw badRequest("limit requires before.");
return null;
}
if (rawTail !== undefined) {
if (rawBefore !== undefined) throw badRequest("tailLimit and before are mutually exclusive.");
if (rawLimit !== undefined) throw badRequest("limit only applies to before requests.");
return { kind: "tail", limit: pageLimit(rawTail, "tailLimit") };
}
const cursor = decodeCursor(rawBefore!);
if (cursor === null) {
throw badRequest("before must be a cursor of the form <shardIndex>:<ordinal>.");
}
return {
kind: "before",
cursor,
limit: rawLimit !== undefined ? pageLimit(rawLimit, "limit") : MESSAGES_PAGE_LIMIT_DEFAULT,
};
}
/** Accepted `category` query values of the list endpoint (SessionCategory, spelled out for validation). */
const SESSION_CATEGORIES: readonly SessionCategory[] = [
"active",
@@ -494,6 +542,7 @@ export function sessionsRoutes(deps: AppDeps): Hono<AppEnv> {
app.get("/:sessionId/messages", async (c) => {
const row = resolveSession(c);
const page = messagesPageQuery(c);
// Live tail (running/compacting sessions only): capture the channel cursor and the
// open-fragment snapshot together, synchronously — no await between the two, and both
// BEFORE the trace read starts. That ordering is what makes the client contract safe
@@ -504,13 +553,41 @@ export function sessionsRoutes(deps: AppDeps): Hono<AppEnv> {
// cursor — the client's overlap dedup against `messages` decides for them — so a
// complete message whose trace append is still in flight when the read starts is not
// lost either.
//
// Windowed reads keep the same contract on TAIL pages only: a tail page ends at the
// transcript's live edge, so the attachment's semantics are identical. A `before`
// page is immutable history — attaching in-flight fragments to it would seed them at
// the wrong position — so it never carries `live`.
let live: MessagesLiveTail | undefined;
if (deps.manager.statusOf(row.sessionId) !== "idle") {
if (page?.kind !== "before" && deps.manager.statusOf(row.sessionId) !== "idle") {
live = {
cursor: deps.channels.get(row.sessionId).lastEventId,
fragments: deps.manager.liveFragments(row.sessionId),
};
}
if (page !== null) {
const result = await deps.traceService.readMessagesPage(
row.projectId,
row.agentId,
row.sessionId,
page,
);
const info: MessagesPageInfo = {
...(result.before !== undefined ? { before: result.before } : {}),
earlierTurns: result.prior.turns,
prior: {
subagentTokens: result.prior.subagentTokens,
elapsedMs: result.prior.elapsedMs,
sessionTokens: result.prior.sessionTokens,
contextTokens: result.prior.contextTokens,
},
};
return c.json({
messages: result.messages,
...(live !== undefined ? { live } : {}),
page: info,
} satisfies MessagesResponse);
}
const messages = await deps.traceService.readMessages(
row.projectId,
row.agentId,
@@ -0,0 +1,430 @@
/**
* Message-window scanner: the pure logic behind cursor pagination of
* `GET /api/sessions/:id/messages` (TraceService.readMessagesPage).
*
* A window **unit** is one Task in the Web reducer's sense: it opens at a main-session
* user prompt (text or image) and runs until the next such prompt. Cutting only at these
* boundaries is what keeps every stream-model invariant intact inside a window — a
* tool_call is never separated from its tool_call_output (outputs land before the next
* user prompt), a compaction_begin/end span never splits (compaction-internal messages
* stay skippable), a steering chip keeps the images that ride behind it, and a
* mid-Task shard rotation (auto compaction with carry-over) stays glued to its Task.
*
* The scanner walks one shard's RAW messages (no subagent expansion — origin-carrying
* messages never appear in a shard) and mirrors, decision for decision, the Web reducer
* in web/src/lib/omni/stream-model.ts (pushMessage/startTask/finalizeOpenTask) plus the
* outline-entry rule in web/src/features/chat/outline-model.ts (buildOutline) — the four
* implementations of "what is one turn" (this file, stream-model, outline-model, and
* trace-service.analyze) must stay in step; tests on both sides pin the shared cases.
*
* Two distinct counts come out of one pass:
* - **unit boundaries** — the safe cut points described above (banner-only prompts and
* goal-round re-sends included: they start Tasks in the reducer);
* - **outline turns** — the subset that opens an entry in the Web's conversation
* outline (`第 N 轮` numbering). Machine-only prompts (handoff / model-switch
* blocks), goal rounds past 1, steering chips and compaction injections do NOT open
* entries, and consecutive user messages of one send merge into one entry — the
* count must match buildOutline exactly, or a paginated outline would mis-number.
*
* The pass also accumulates the **prior stats** a partial window needs seeded into the
* Web's stats tracker so header chips and per-turn cumulative rows keep telling the
* truth (see task-stats.ts `seedPriorStats`): elapsed time of finished Tasks, subagent
* token totals (via the caller-provided child expander), and the last main-session
* session/context token readings.
*
* Scan state is carried across shards (a Task can span a rotation). Cached per-shard
* prefix records (trace_files.page_stats) persist the carry so old shards are read at
* most once ever; bump CACHE_VERSION whenever any rule in this file changes.
*/
import { parseUserSteeringText } from "@prismshadow/penguin-core";
import {
parseGoalMessage,
parseHandoffMessage,
parseModelSwitchMessage,
} from "@prismshadow/penguin-core/markers";
import type { OmniMessage } from "@prismshadow/penguin-core";
/** Bump when any counting/boundary rule changes: cached page_stats records with an older version are recomputed. */
export const CACHE_VERSION = 1;
/** Cumulative totals at a point in the trace (all values are "before this point"). */
export interface WindowPriorStats {
/** Outline entries opened before this point (the Web outline's global numbering offset). */
turns: number;
/** Sum of subagent token_usage request totals (all descendant sessions) before this point. */
subagentTokens: number;
/** Sum of finished Tasks' elapsed time before this point (the Web's sessionElapsedMs basis). */
elapsedMs: number;
/** The last main-session token_usage `session.total` seen before this point (0 = none). */
sessionTokens: number;
/** The last main-session NON-compaction token_usage `request.total` before this point (context occupancy basis). */
contextTokens: number;
}
/** Open-Task bookkeeping carried across messages (mirrors StreamModel.task* fields). */
interface TaskCarry {
open: boolean;
/** Task first/latest message timestamps (ms; the degenerate-round elapsed fallback). */
firstTsMs: number;
lastTsMs: number;
/** Last non-compaction request_end (ms); the Task's true end when present. */
lastReqEndMs: number | null;
/** Whether this Task saw main usage / non-blank assistant text — decides whether a stats row would be emitted at its close (which breaks a user-item run). */
sawUsage: boolean;
sawReply: boolean;
/** Compaction usage held pending (committed to the Task by a later non-compaction request_end — mirrors task-stats.ts pendingCompaction*). */
pendingCompactionUsage: boolean;
}
/** Full scanner state, serializable into a shard's cached prefix record. */
export interface ScanState {
totals: WindowPriorStats;
task: TaskCarry;
/** Between a main-session compaction_begin and its compaction_end. */
compactionActive: boolean;
/** A steering chip is still collecting its trailing images (mirrors StreamModel.openSteering). */
steeringOpen: boolean;
/**
* Inside an unbroken, entry-carrying run of user items — buildOutline's `lastWasUser`,
* tracked at the item level: true after an entry-eligible user item, false after any
* non-user item (a stats row included) AND after banner/goal-round texts (which open no
* entry and break the run in buildOutline). Both the window-cut and the merge decision
* key off it, so a cut can never split what the outline merges.
*/
runOpen: boolean;
}
export function initialScanState(): ScanState {
return {
totals: { turns: 0, subagentTokens: 0, elapsedMs: 0, sessionTokens: 0, contextTokens: 0 },
task: {
open: false,
firstTsMs: 0,
lastTsMs: 0,
lastReqEndMs: null,
sawUsage: false,
sawReply: false,
pendingCompactionUsage: false,
},
compactionActive: false,
steeringOpen: false,
runOpen: false,
};
}
/** One safe cut point: `ordinal` indexes into the shard's parsed message array; `stats` is the cumulative prior AT the cut (the unit itself not included). */
export interface UnitBoundary {
ordinal: number;
stats: WindowPriorStats;
}
/**
* Aggregate a subagent pointer's child subtree without materializing its messages:
* the token sum feeds `subagentTokens`, the max timestamp advances the open Task's
* lastTs exactly as routeNested's touchTask would for every expanded child message.
* Null = child trace missing (the pointer event stays a plain event, contributing nothing).
*/
export interface ChildAggregate {
requestTokens: number;
maxTsMs: number | null;
}
function tsMs(timestamp: string): number | null {
const ms = Date.parse(timestamp);
return Number.isFinite(ms) ? ms : null;
}
/** Mirrors stream-model touchTask: advance the open Task's latest timestamp (compaction-internal messages excluded). */
function touchTask(state: ScanState, ms: number | null): void {
if (!state.task.open || state.compactionActive || ms === null) return;
if (ms > state.task.lastTsMs) state.task.lastTsMs = ms;
}
/** Mirrors stream-model finalizeOpenTask + task-stats endTask: settle the open Task's elapsed into the cumulative. */
function finalizeTask(state: ScanState): void {
const t = state.task;
if (!t.open) return;
t.open = false;
const endMs = t.lastReqEndMs ?? t.lastTsMs;
state.totals.elapsedMs += Math.max(0, endMs - t.firstTsMs);
t.pendingCompactionUsage = false;
}
/** Mirrors stream-model startTask: finalize the previous Task and open a new one at `ms`. */
function startTask(state: ScanState, ms: number | null): void {
finalizeTask(state);
const t = state.task;
t.open = true;
t.firstTsMs = ms ?? 0;
t.lastTsMs = t.firstTsMs;
t.lastReqEndMs = null;
t.sawUsage = false;
t.sawReply = false;
t.pendingCompactionUsage = false;
}
/** A non-user item entered the stream: the user-item run breaks (buildOutline's lastWasUser = false). */
function breakRuns(state: ScanState): void {
state.runOpen = false;
}
/**
* Handle one Task-starting user message (prompt text or non-steering image).
* `entryEligible` says whether the message would open an outline entry (buildOutline's
* rule); ineligible messages (banner blocks, goal rounds > 1) still start Tasks — and
* still cut when the run is broken — but never count and always break the run, exactly
* as buildOutline resets lastWasUser for them.
*/
function onTaskStart(
state: ScanState,
ms: number | null,
entryEligible: boolean,
onBoundary: (stats: WindowPriorStats) => void,
): void {
// A stats row emitted while closing the previous Task lands BEFORE this user item and
// breaks the item run (finalizeOpenTask inserts it ahead of the new prompt) — mirror
// that so two sends merged by the outline are never cut apart, while two sends
// separated by a stats row cut (and count) as two.
if (state.task.open && (state.task.sawUsage || state.task.sawReply)) breakRuns(state);
const boundary = !state.runOpen;
startTask(state, ms);
if (boundary) onBoundary({ ...state.totals });
if (entryEligible) {
if (boundary) state.totals.turns += 1;
state.runOpen = true;
} else {
state.runOpen = false;
}
}
/** buildOutline's eligibility for a user prompt TEXT: machine-only source blocks and goal rounds past 1 open no entry. */
function outlineEligibleText(text: string): boolean {
if (parseHandoffMessage(text) || parseModelSwitchMessage(text)) return false;
const goal = parseGoalMessage(text);
return !(goal !== null && goal.round > 1);
}
/**
* Scan one shard's raw messages, mutating `state` in place and reporting every unit
* boundary through `onBoundary` (with the cumulative priors AT the cut). `expandChild`
* resolves a subagent pointer's aggregate (may hit a per-request memo); pass null to
* skip child reads when the caller does not need subagent totals for this span.
*/
export async function scanMessages(
state: ScanState,
messages: readonly OmniMessage[],
onBoundary: (ordinal: number, stats: WindowPriorStats) => void,
expandChild: ((sessionId: string) => Promise<ChildAggregate | null>) | null,
fromOrdinal = 0,
toOrdinal = messages.length,
): Promise<void> {
for (let i = fromOrdinal; i < toOrdinal; i++) {
const msg = messages[i]!;
// Shards never contain origin-carrying messages (core's Writer filters them);
// defensively skip any that appear rather than mis-shaping the counts.
if (msg.origin !== undefined && msg.origin.length > 0) continue;
const ms = tsMs(msg.timestamp);
const p = msg.payload as Record<string, unknown> & { type?: string; role?: string };
// Mirrors pushMessage's first step: anything that is not a complete user image
// closes the steering-images window (session_meta and events included).
const isUserImage = msg.type === "model_msg" && p.type === "image_url";
if (!isUserImage) state.steeringOpen = false;
if (msg.type === "model_msg") {
// Compaction-internal model messages: never rendered, never counted (stream-model
// returns before any item/Task logic; only touchTask advances).
if (state.compactionActive) {
touchTask(state, ms);
continue;
}
if (p.type === "text" && p.role === "user" && typeof p.text === "string") {
const text = p.text;
// Compaction-summary injection: no item, no Task, no run change.
if (text.startsWith("[context_summary]") || text.startsWith("<context_summary>")) {
touchTask(state, ms);
continue;
}
// Mid-run steering: rendered as a user_steering item (not user_text) — it breaks
// the outline's user-item run but never starts a Task or cuts a window.
if (parseUserSteeringText(text) !== null) {
touchTask(state, ms);
breakRuns(state);
state.steeringOpen = true;
continue;
}
onTaskStart(state, ms, outlineEligibleText(text), (stats) => onBoundary(i, stats));
continue;
}
if (p.type === "image_url") {
// An image riding behind a steering chip folds into it: no item, no Task.
if (state.steeringOpen) {
touchTask(state, ms);
continue;
}
onTaskStart(state, ms, true, (stats) => onBoundary(i, stats));
continue;
}
if (p.type === "text" && p.role === "assistant" && typeof p.text === "string") {
touchTask(state, ms);
// Blank fidelity-only messages produce no item (stream-model discards them).
if (p.text.trim() !== "") {
state.task.sawReply = true;
breakRuns(state);
}
continue;
}
if (p.type === "thinking") {
touchTask(state, ms);
if (typeof p.thinking === "string" && p.thinking.trim() !== "") breakRuns(state);
continue;
}
if (p.type === "tool_call") {
touchTask(state, ms);
breakRuns(state);
continue;
}
// tool_call_output updates an existing card (no new item); inline_* render nothing.
touchTask(state, ms);
continue;
}
if (msg.type === "event_msg") {
touchTask(state, ms);
const t = p.type;
if (t === "compaction_begin") {
state.compactionActive = true;
breakRuns(state); // the compaction banner is an item
continue;
}
if (t === "compaction_end") {
// Closes the banner opened by the begin (no new item mid-window).
state.compactionActive = false;
continue;
}
if (t === "abort") {
breakRuns(state); // the interruption marker is an item
continue;
}
if (t === "request_end") {
if (state.compactionActive) continue; // compaction requests render nothing
if (ms !== null) state.task.lastReqEndMs = ms;
// Pending compaction usage commits at a later non-compaction request_end
// (compaction mid-Task) — from then on the Task WILL show a stats row.
if (state.task.pendingCompactionUsage) {
state.task.sawUsage = true;
state.task.pendingCompactionUsage = false;
}
const status = p.status;
// A retryable end renders a reconnect-hint item.
if (status === "failed" || status === "timeout" || status === "malformed") {
breakRuns(state);
}
continue;
}
if (t === "token_usage") {
const request = p.request as { total?: number } | undefined;
const session = p.session as { total?: number } | undefined;
// The session cumulative tracks the provider even during compaction
// (task-stats trackMainUsage does the same).
if (typeof session?.total === "number") state.totals.sessionTokens = session.total;
if (state.compactionActive) {
state.task.pendingCompactionUsage = true;
} else {
if (typeof request?.total === "number") state.totals.contextTokens = request.total;
state.task.sawUsage = true;
}
continue;
}
if (t === "subagent" && typeof p.session_id === "string" && expandChild !== null) {
// The pointer expands to the child's messages in the served transcript: their
// token_usage feeds the parent's subagent totals at every depth, and their
// timestamps advance the open Task exactly as routeNested's touchTask would.
const agg = await expandChild(p.session_id);
if (agg !== null) {
state.totals.subagentTokens += agg.requestTokens;
touchTask(state, agg.maxTsMs);
}
continue;
}
// approval_decision / request_begin / unexpanded pointers: no items, no run change.
continue;
}
// session_meta: no item (steering window already closed above).
}
}
/** Finalize the trailing open Task (end of the whole trace): the cumulative then covers every finished Task. */
export function finalizeScan(state: ScanState): void {
finalizeTask(state);
}
// ---------------------------------------------------------------------------
// Cursor encoding
// ---------------------------------------------------------------------------
/** A window cursor: shard file index + message ordinal within that shard's parsed array. */
export interface MessageCursor {
fileIndex: number;
ordinal: number;
}
/**
* `<fileIndex>:<ordinal>` — stable across requests and compaction: rotation always
* opens a NEW shard, closed shards are immutable, and the active shard is append-only,
* so a (shard, ordinal) pair never moves.
*/
export function encodeCursor(c: MessageCursor): string {
return `${c.fileIndex}:${c.ordinal}`;
}
/** Strict parse of a cursor string; null when malformed (callers turn that into a 400). */
export function decodeCursor(raw: string): MessageCursor | null {
const m = /^(\d{1,9}):(\d{1,9})$/.exec(raw);
if (!m) return null;
return { fileIndex: Number(m[1]), ordinal: Number(m[2]) };
}
// ---------------------------------------------------------------------------
// Cached per-shard prefix records
// ---------------------------------------------------------------------------
/**
* The persisted shape of trace_files.page_stats: the scan state at the END of a shard
* (cumulative from the very beginning of the session). Only immutable shards are cached
* — the newest shard still grows. `v` gates rule evolution: a record from an older
* CACHE_VERSION is recomputed as if absent.
*/
export interface ShardPrefixRecord {
v: number;
state: ScanState;
}
export function serializePrefix(state: ScanState): string {
return JSON.stringify({ v: CACHE_VERSION, state } satisfies ShardPrefixRecord);
}
/** Parse a cached record; null when absent, unparseable, or from another CACHE_VERSION. */
export function deserializePrefix(raw: string | null): ScanState | null {
if (raw === null) return null;
try {
const rec = JSON.parse(raw) as ShardPrefixRecord;
if (rec.v !== CACHE_VERSION || typeof rec.state !== "object" || rec.state === null) {
return null;
}
return rec.state;
} catch {
return null;
}
}
/** Deep-copy a scan state (cached records must not be mutated by a later scan). */
export function cloneScanState(state: ScanState): ScanState {
return {
totals: { ...state.totals },
task: { ...state.task },
compactionActive: state.compactionActive,
steeringOpen: state.steeringOpen,
runOpen: state.runOpen,
};
}
+301 -33
View File
@@ -40,6 +40,21 @@ import type { TraceFileRow, TraceSessionRow } from "../db/repos/trace-index.js";
import { HttpError } from "../http/errors.js";
import { formatLocalDate } from "../internal/dates.js";
import type { SessionSources } from "../runtime/session-sources.js";
import {
cloneScanState,
deserializePrefix,
encodeCursor,
finalizeScan,
initialScanState,
scanMessages,
serializePrefix,
} from "./message-window.js";
import type {
ChildAggregate,
MessageCursor,
ScanState,
WindowPriorStats,
} from "./message-window.js";
import { TraceIndexService, traceFilePath } from "./trace-index.js";
const TRACE_FILE_RE = /^(.+)_(\d{3})\.jsonl$/;
@@ -58,6 +73,34 @@ interface LocatedFile {
path: string;
date: string;
index: number;
/** Last observed size from the index row (the page-stats cache write guard). */
sizeBytes: number;
}
/** A windowed history read's request shape (see readMessagesPage). */
export type MessagesPageRequest =
{ kind: "tail"; limit: number } | { kind: "before"; cursor: MessageCursor; limit: number };
/** A windowed history read's result (route maps it onto MessagesResponse.page). */
export interface MessagesPageResult {
messages: OmniMessage[];
/** Cursor of the window's first unit; absent = the window reaches the beginning. */
before?: string;
/** Cumulative stats before the window (earlierTurns = prior.turns). */
prior: WindowPriorStats;
}
/**
* Per-request expansion context: `projectScanned` caps the miss-path force reconcile at
* one per request; `ancestry` guards cycles; `raw` memoizes each child session's
* concatenated shard messages so the window expansion and the subagent-token
* aggregation never read the same child twice within one request.
*/
interface ExpandCtx {
projectScanned: boolean;
ancestry: Set<string>;
depth: number;
raw: Map<string, OmniMessage[] | null>;
}
/**
@@ -92,6 +135,12 @@ export interface TraceServiceDeps {
sessions?: TraceSessionIndex;
/** Optional (narrow tests may omit): the shared in-process Session-origin registry (single source of truth for `source`). */
sources?: SessionSources;
/**
* Test observability: called with the path of every Trace shard this service reads
* from disk (windowed-read tests assert an old-window request never touches the
* newest shard and vice versa). Production wiring omits it.
*/
observeShardRead?: (path: string) => void;
}
/** One Session's classification result (see classify): its sidebar category + Workspace path ("" = unknown). */
@@ -127,9 +176,16 @@ export class TraceService {
path: traceFilePath(this.root, r),
date: r.date,
index: r.fileIndex,
sizeBytes: r.sizeBytes,
}));
}
/** All shard reads funnel through here (deps.observeShardRead is the windowed-read tests' proof of which files were touched). */
private async readShard(path: string): Promise<OmniMessage[]> {
this.deps.observeShardRead?.(path);
return readTraceTolerant(path);
}
/** Deletes all of this Session's Trace files (called when the Session is deleted); the index rows go with them. */
async deleteSessionTraces(projectId: string, agentId: string, sessionId: string): Promise<void> {
const files = await this.locateAll(projectId, agentId, sessionId);
@@ -160,11 +216,18 @@ export class TraceService {
agentId: string,
sessionId: string,
): Promise<OmniMessage[]> {
return this.readMessagesExpanded(projectId, agentId, sessionId, {
const ctx: ExpandCtx = {
projectScanned: false,
ancestry: new Set([sessionId]),
depth: 0,
});
raw: new Map(),
};
const files = await this.locateAll(projectId, agentId, sessionId);
const out: OmniMessage[] = [];
for (const file of files) {
out.push(...(await this.expandMessages(projectId, await this.readShard(file.path), ctx)));
}
return out;
}
/**
@@ -186,42 +249,247 @@ export class TraceService {
return this.deps.index.repo.findAgentBySession(projectId, childSid);
}
private async readMessagesExpanded(
/** A child session's concatenated raw shard messages, memoized per request; null = no Trace located. */
private async readChildRaw(
projectId: string,
childSid: string,
ctx: ExpandCtx,
): Promise<OmniMessage[] | null> {
const cached = ctx.raw.get(childSid);
if (cached !== undefined) return cached;
const childAgent = await this.resolveChildAgent(projectId, childSid, ctx);
let raw: OmniMessage[] | null = null;
if (childAgent) {
const files = await this.locateAll(projectId, childAgent, childSid);
if (files.length > 0) {
raw = [];
for (const file of files) raw.push(...(await this.readShard(file.path)));
}
}
ctx.raw.set(childSid, raw);
return raw;
}
/**
* Expand subagent pointers within one message span (the whole transcript on the full
* path, one window on the paged path — same recursive shape either way).
*/
private async expandMessages(
projectId: string,
messages: readonly OmniMessage[],
ctx: ExpandCtx,
): Promise<OmniMessage[]> {
const out: OmniMessage[] = [];
for (const msg of messages) {
// The depth cap guards against runaway recursion; ancestry guards against a
// cyclic pointer (a tampered Trace pointing to itself/an ancestor is not expanded).
const childSid = ctx.depth < MAX_SUBAGENT_DEPTH ? subagentPointer(msg) : null;
if (!childSid || ctx.ancestry.has(childSid)) {
out.push(msg);
continue;
}
const raw = await this.readChildRaw(projectId, childSid, ctx);
let nested: OmniMessage[] = [];
if (raw !== null) {
ctx.ancestry.add(childSid);
nested = await this.expandMessages(projectId, raw, { ...ctx, depth: ctx.depth + 1 });
ctx.ancestry.delete(childSid);
}
// Child Trace missing (deleted): keep the pointer event, since the sub-session's content can't be recovered.
if (nested.length === 0) {
out.push(msg);
continue;
}
for (const m of nested) out.push({ ...m, origin: [childSid, ...(m.origin ?? [])] });
}
return out;
}
/**
* A subagent pointer's subtree aggregate for the window scanner (message-window.ts):
* descendant token_usage request totals plus the subtree's max timestamp — the same
* contributions the expanded messages make to the Web's parent-level stats tracker.
* The recursion mirrors expandMessages' depth/ancestry rules exactly, so an
* unexpandable pointer contributes nothing on both paths.
*/
private async aggregateChild(
projectId: string,
childSid: string,
ctx: ExpandCtx,
): Promise<ChildAggregate | null> {
if (ctx.depth >= MAX_SUBAGENT_DEPTH || ctx.ancestry.has(childSid)) return null;
const raw = await this.readChildRaw(projectId, childSid, ctx);
if (raw === null || raw.length === 0) return null;
let requestTokens = 0;
let maxTsMs: number | null = null;
ctx.ancestry.add(childSid);
for (const msg of raw) {
const ms = Date.parse(msg.timestamp);
if (Number.isFinite(ms) && (maxTsMs === null || ms > maxTsMs)) maxTsMs = ms;
const p = msg.payload as { type?: string; request?: { total?: number } };
if (msg.type === "event_msg" && p.type === "token_usage") {
requestTokens += typeof p.request?.total === "number" ? p.request.total : 0;
continue;
}
const grandchild = subagentPointer(msg);
if (grandchild !== null) {
const agg = await this.aggregateChild(projectId, grandchild, {
...ctx,
depth: ctx.depth + 1,
});
if (agg !== null) {
requestTokens += agg.requestTokens;
if (agg.maxTsMs !== null && (maxTsMs === null || agg.maxTsMs > maxTsMs)) {
maxTsMs = agg.maxTsMs;
}
}
}
}
ctx.ancestry.delete(childSid);
return { requestTokens, maxTsMs };
}
/**
* End-of-shard scan states for files[0..upto] (cumulative from the transcript's
* beginning), served from the trace_files.page_stats cache. A missing/stale record
* costs one read of THAT shard — once ever: every shard here is immutable (rotation
* opens a new shard; the newest shard never appears in a prefix, because windows are
* suffixes and their first shard is always read anyway). The write-back is guarded on
* the indexed size, so an externally rewritten shard can only invalidate, never
* poison, the cache.
*/
private async prefixStates(
projectId: string,
agentId: string,
sessionId: string,
ctx: { projectScanned: boolean; ancestry: Set<string>; depth: number },
): Promise<OmniMessage[]> {
const files = await this.locateAll(projectId, agentId, sessionId);
const out: OmniMessage[] = [];
for (const file of files) {
for (const msg of await readTraceTolerant(file.path)) {
// The depth cap guards against runaway recursion; ancestry guards against a
// cyclic pointer (a tampered Trace pointing to itself/an ancestor is not expanded).
const childSid = ctx.depth < MAX_SUBAGENT_DEPTH ? subagentPointer(msg) : null;
if (!childSid || ctx.ancestry.has(childSid)) {
out.push(msg);
continue;
}
const childAgent = await this.resolveChildAgent(projectId, childSid, ctx);
let nested: OmniMessage[] = [];
if (childAgent) {
ctx.ancestry.add(childSid);
nested = await this.readMessagesExpanded(projectId, childAgent, childSid, {
...ctx,
depth: ctx.depth + 1,
});
ctx.ancestry.delete(childSid);
}
// Child Trace missing (deleted): keep the pointer event, since the sub-session's content can't be recovered.
if (nested.length === 0) {
out.push(msg);
continue;
}
for (const m of nested) out.push({ ...m, origin: [childSid, ...(m.origin ?? [])] });
files: LocatedFile[],
upto: number,
ctx: ExpandCtx,
): Promise<ScanState[]> {
const states: ScanState[] = [];
for (let j = 0; j <= upto; j++) {
const file = files[j]!;
const cached = deserializePrefix(
this.deps.index.repo.getPageStats(projectId, agentId, sessionId, file.index),
);
if (cached !== null) {
states.push(cached);
continue;
}
const state = cloneScanState(j === 0 ? initialScanState() : states[j - 1]!);
const messages = await this.readShard(file.path);
await scanMessages(
state,
messages,
() => {},
(sid) => this.aggregateChild(projectId, sid, ctx),
);
this.deps.index.repo.setPageStats(
projectId,
agentId,
sessionId,
file.index,
file.sizeBytes,
serializePrefix(state),
);
states.push(state);
}
return out;
return states;
}
/**
* Windowed history read (the `tailLimit` / `before` forms of GET /messages). The
* window is a run of whole units — cut points and unit semantics live in
* message-window.ts — assembled by reading ONLY the shards the window overlaps
* (plus, once ever per old shard, the prefix-cache backfill above). Subagent
* pointers are expanded exactly as the full path expands them, but only within the
* window: children referenced by older windows load when those windows do.
*/
async readMessagesPage(
projectId: string,
agentId: string,
sessionId: string,
req: MessagesPageRequest,
): Promise<MessagesPageResult> {
const files = await this.locateAll(projectId, agentId, sessionId);
const empty = (): MessagesPageResult => ({
messages: [],
prior: initialScanState().totals,
});
if (files.length === 0) return empty();
const ctx: ExpandCtx = {
projectScanned: false,
ancestry: new Set([sessionId]),
depth: 0,
raw: new Map(),
};
// The window's exclusive end: the newest shard's end (tail), or the cursor (before).
let endPos: number;
let endOrdinal: number | null = null; // null = to the shard's end
if (req.kind === "tail") {
endPos = files.length - 1;
} else {
endPos = files.findIndex((f) => f.index === req.cursor.fileIndex);
// Cursor shard no longer on disk (external deletion — locateAll already
// force-retried): the history it pointed into is gone; report end-of-history
// rather than a guess at what used to precede it.
if (endPos < 0) return empty();
endOrdinal = req.cursor.ordinal;
}
const prefixes = await this.prefixStates(projectId, agentId, sessionId, files, endPos - 1, ctx);
// Walk backward from the end, scanning whole shards (each from its cached carry-in)
// until the window has more units than requested or the beginning is reached.
const shardMessages = new Map<number, OmniMessage[]>();
let boundaries: Array<{ pos: number; ordinal: number; stats: WindowPriorStats }> = [];
let startPos = endPos + 1;
while (boundaries.length <= req.limit && startPos > 0) {
startPos -= 1;
const messages = await this.readShard(files[startPos]!.path);
shardMessages.set(startPos, messages);
const state = cloneScanState(startPos === 0 ? initialScanState() : prefixes[startPos - 1]!);
const shardBoundaries: Array<{ pos: number; ordinal: number; stats: WindowPriorStats }> = [];
const to = startPos === endPos && endOrdinal !== null ? endOrdinal : messages.length;
await scanMessages(
state,
messages,
(ordinal, stats) => shardBoundaries.push({ pos: startPos, ordinal, stats }),
(sid) => this.aggregateChild(projectId, sid, ctx),
0,
Math.min(to, messages.length),
);
boundaries = [...shardBoundaries, ...boundaries];
}
// Window start: the last `limit` units, or the very beginning (preamble included)
// when the whole remaining history fits — then there is no `before` cursor.
let start: { pos: number; ordinal: number };
let before: string | undefined;
let prior: WindowPriorStats;
if (boundaries.length > req.limit) {
const wb = boundaries[boundaries.length - req.limit]!;
start = { pos: wb.pos, ordinal: wb.ordinal };
before = encodeCursor({ fileIndex: files[wb.pos]!.index, ordinal: wb.ordinal });
prior = wb.stats;
} else {
start = { pos: 0, ordinal: 0 };
prior = initialScanState().totals;
}
const windowRaw: OmniMessage[] = [];
for (let pos = start.pos; pos <= endPos; pos++) {
const messages = shardMessages.get(pos) ?? (await this.readShard(files[pos]!.path));
const from = pos === start.pos ? start.ordinal : 0;
const to =
pos === endPos && endOrdinal !== null
? Math.min(endOrdinal, messages.length)
: messages.length;
for (let i = from; i < to; i++) windowRaw.push(messages[i]!);
}
const expanded = await this.expandMessages(projectId, windowRaw, ctx);
return { messages: expanded, ...(before !== undefined ? { before } : {}), prior };
}
/** List of Trace files (index / date / size / mtime). */