diff --git a/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json b/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json index 00a1846395c5..041a4faaac1e 100644 --- a/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json @@ -1 +1 @@ -{"contentHash":"7c5c92898e62f52b09c51a1dd55bcb758eb8a40a983cf9c079c913d95dba3b58","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"} +{"contentHash":"301f657c77e4a26a00896ada2ae1134e2f45cd2c102c03967d5910424520c8b7","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/agent-harness.json b/docs/.generated/plugin-sdk-api-baseline/agent-harness.json index eb842b6f71af..c482cd15f0b6 100644 --- a/docs/.generated/plugin-sdk-api-baseline/agent-harness.json +++ b/docs/.generated/plugin-sdk-api-baseline/agent-harness.json @@ -1 +1 @@ -{"contentHash":"0069a9c1609324d148c10ae79fed31195bbd8d866a864d546ef8df3a9eb6f2eb","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"} +{"contentHash":"f65b4dc2a33a84047a3f43f6f0db7f5d4779a892eb91b2e085be3666774ad760","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-core.json b/docs/.generated/plugin-sdk-api-baseline/channel-core.json index ea79a9a77bcd..dd79cbb1ae1d 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-core.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-core.json @@ -1 +1 @@ -{"contentHash":"7a28969db8490d4b3c3f6f8c05db564ec89508030523a95a8a4248f8d9072edf","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"} +{"contentHash":"92d1c4e69384ca9fb89d6d07077090b96b02f9390d4a51e95a35ea40f530934c","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json b/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json index df49669131e1..87a4b7ac32c5 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json @@ -1 +1 @@ -{"contentHash":"532feed49e6fd72b81b4ce744d7284519158999cebb31db751174470fb63a3cc","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"} +{"contentHash":"342b52192d1694e76eb09912576417ba0cf493d5f47f02134ed88ad4f5dbb786","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-message.json b/docs/.generated/plugin-sdk-api-baseline/channel-message.json index 935ab373e9c6..01f2d15589de 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-message.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-message.json @@ -1 +1 @@ -{"contentHash":"68796f883fb54466351cb0a077f3fd1f71ad6d6d1df53248f0d382e77a5f486f","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"} +{"contentHash":"5f20abd3059f4a159e1db5caa62d185bf02ccd08551a9b41cd1ba670418dc89d","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json b/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json index 9519f5d7802f..68fa4f298195 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json @@ -1 +1 @@ -{"contentHash":"69a92cf59ee4dd7c66431a80aed3f32e0b4e830789e4cab7b4173ec2d502bcfc","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"} +{"contentHash":"60da79b335ff21d12842d56fe5dcb29d24acd686ac9a07b0c54ccb56737ebaef","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json b/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json index 7804a4fc9a80..9d8c61631828 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json @@ -1 +1 @@ -{"contentHash":"22ee5f9631d28579af7419f2e1263ecfa903f557b3b3e9962769e3a1fb61489d","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"} +{"contentHash":"23b1677325bf4ba2b6dca7209f902dcbafc9e57c019d65c9819788f07ec02294","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"} diff --git a/docs/.generated/plugin-sdk-api-baseline/core.json b/docs/.generated/plugin-sdk-api-baseline/core.json index 03b4e2c561c9..b101f4bd1840 100644 --- a/docs/.generated/plugin-sdk-api-baseline/core.json +++ b/docs/.generated/plugin-sdk-api-baseline/core.json @@ -1 +1 @@ -{"contentHash":"432dc3f2d69b09ae1f8c64bf5c1b4f43d2a748ab79d90862ff2a9aadaa4af43f","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"} +{"contentHash":"838d86a164f74fe9d31abb32f31ec34580e2814d6e5611593e2616aa91d7ee8f","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"} diff --git a/docs/.generated/plugin-sdk-api-baseline/discord.json b/docs/.generated/plugin-sdk-api-baseline/discord.json index c4eb214fdeb3..9c2535d02c47 100644 --- a/docs/.generated/plugin-sdk-api-baseline/discord.json +++ b/docs/.generated/plugin-sdk-api-baseline/discord.json @@ -1 +1 @@ -{"contentHash":"1f02b35bfc99cc8bdeea415dc430d832a74d9b1ff7e52243fea224a7a006425b","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"} +{"contentHash":"9eacfda79c3ddaf6f79d7c53f40368cde1c0e81c5c52f9d13ea09e55fe9f23b3","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"} diff --git a/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json b/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json index 31051a3d2432..7ec40c1c652f 100644 --- a/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json +++ b/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json @@ -1 +1 @@ -{"contentHash":"6bff1e7ef600d59cfbe7adb7a67ed9ffe8173e50f89ba618f6f190e90fc74da9","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"} +{"contentHash":"df01b92e7a05975696443bcf565aa07053451fc4b8693303a758ba59ec063191","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"} diff --git a/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json b/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json index f1e14495da41..bb1e172cc7a9 100644 --- a/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json @@ -1 +1 @@ -{"contentHash":"8e07f0c6667ff6dfcafdf2dc98be90b0cc6aa8d8449bd1e35227976c63c4f8df","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"} +{"contentHash":"8b7b1cc3ec236e1de88076c2866c1ef8485a1b1135c90a0a671a305650b67031","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json b/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json index 793be537e40e..e983b5a62d04 100644 --- a/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json +++ b/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json @@ -1 +1 @@ -{"contentHash":"9be8b6ce4a747459accac5b8b4a64023f464d437d1117e3148b1199461dbb402","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"} +{"contentHash":"576a0385d3fe222dd1be03d8eaf965583bd29e55377073a63b7734e6052bfa1d","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"} diff --git a/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json b/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json index 1e30f507e7b0..c82ed0fb3a5b 100644 --- a/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json @@ -1 +1 @@ -{"contentHash":"c2e4b880b08c352ba57f781a0c629d15066c5842a287dac3727b1dd5d8561811","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"} +{"contentHash":"8bc3fff3377a2ebbe8ef572c4023085e81ad4a8a408d7b922cde80a2c890ccca","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json b/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json index 973f26435acd..84b74bca68cd 100644 --- a/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json @@ -1 +1 @@ -{"contentHash":"e4df60d58a485eed9e2b1509dded473da589102ad14a4ae7ba5cedb2250dc33c","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"} +{"contentHash":"4d4dea12b2e6f3ff2b6b55ec2ea97a6fd4a9c4c0e550b5ac6c14763ace7cabb3","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json b/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json index d590c1404cc3..f36e7ce69810 100644 --- a/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json +++ b/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json @@ -1 +1 @@ -{"contentHash":"06fad1eff295b30e25b2aa52c6259e67f2e74c15a5bd7f15e9e5e6c3d9be4934","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"} +{"contentHash":"cd826485d2bc10f25aaa441a3b298a09addaa6cfe84a1d8d3f9797985c48fcc4","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"} diff --git a/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json b/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json index 01d28b4c5232..339ddea2f0b2 100644 --- a/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json +++ b/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json @@ -1 +1 @@ -{"contentHash":"6e792826511277b2940a7e788cffdaa647dff60767073ec0a0278d3d6319535d","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"} +{"contentHash":"1f8dadb2465061b3ddf87bb0445ea639a3fcf10b7622df3a4947079d329b2ee6","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"} diff --git a/src/config/sessions/session-accessor.sqlite-active-events.test.ts b/src/config/sessions/session-accessor.sqlite-active-events.test.ts index caa0ac975a1d..2205f60c899d 100644 --- a/src/config/sessions/session-accessor.sqlite-active-events.test.ts +++ b/src/config/sessions/session-accessor.sqlite-active-events.test.ts @@ -13,6 +13,7 @@ import { readRecentSessionTranscriptMessageEvents, readSessionTranscriptActiveLeafEvents, readSessionTranscriptActiveStats, + readSessionTranscriptBoundedContextMessageTailPage, readSessionTranscriptBoundedMessageTailPage, readSessionTranscriptMessageAnchorPage, readSessionTranscriptMessageEventById, @@ -872,5 +873,40 @@ describe("SQLite active transcript event projection", () => { await reconciliation; expect(order).toEqual(["event-loop-responsive", "live-write", "reconciled"]); expect(readSessionTranscriptMessageEventCount(scope)).toBe(100_001); + + await appendTranscriptEvent(scope, { + type: "compaction", + id: "large-context-boundary", + parentId: "m100001", + timestamp: "2026-08-11T00:00:00.000Z", + summary: "bounded large-history summary", + firstKeptEntryId: "m1", + tokensBefore: 100_000, + }); + await persistSessionTranscriptTurn(scope, { + messages: Array.from({ length: 40 }, (_, index) => ({ + eventId: `post-boundary-${index}`, + parentId: index === 0 ? "large-context-boundary" : `post-boundary-${index - 1}`, + message: { + role: index % 2 === 0 ? ("user" as const) : ("assistant" as const), + content: `post-boundary ${index}`, + }, + })), + touchSessionEntry: false, + }); + + const contextPage = readSessionTranscriptBoundedContextMessageTailPage(scope, { + maxBytes: 1024 * 1024, + maxMessages: 40, + maxScannedMessages: 4096, + }); + expect(contextPage).toMatchObject({ + authoritative: true, + contextSummary: { text: "bounded large-history summary" }, + empty: false, + }); + expect(contextPage.events).toHaveLength(40); + expect(contextPage.events.at(0)?.event).toMatchObject({ id: "post-boundary-0" }); + expect(contextPage.events.at(-1)?.event).toMatchObject({ id: "post-boundary-39" }); }, 60_000); }); diff --git a/src/config/sessions/session-accessor.sqlite-active-events.ts b/src/config/sessions/session-accessor.sqlite-active-events.ts index 49e2036f3f3d..a9d418b02a6e 100644 --- a/src/config/sessions/session-accessor.sqlite-active-events.ts +++ b/src/config/sessions/session-accessor.sqlite-active-events.ts @@ -14,10 +14,8 @@ import type { TranscriptEvent, } from "./session-accessor.sqlite-contract.js"; import { + readBoundedContextMessageTail, readVisibleMessageRange, - resolveContextBoundarySummary, - resolveContextMessagePositionRange, - resolveContextMessagePositions, resolveVisibleMessagePositionRange, resolveVisibleMessagePositions, } from "./session-accessor.sqlite-reset-window.js"; @@ -63,8 +61,11 @@ export type SessionTranscriptBoundedMessageTailPage = SessionTranscriptMessageEv serializedBytes: number; }; -type SessionTranscriptBoundedContextMessageTailPage = SessionTranscriptBoundedMessageTailPage & { +type SessionTranscriptBoundedContextMessageTailPage = { + authoritative: boolean; contextSummary?: { text: string; ts: number }; + empty: boolean; + events: SessionTranscriptMessageEvent[]; }; function parseMessageEventRow(row: { @@ -562,19 +563,10 @@ export function readSessionTranscriptBoundedMessageTailPage( */ export function readSessionTranscriptBoundedContextMessageTailPage( scope: SessionTranscriptReadScope, - options: { maxBytes: number; maxMessages: number; offset: number }, + options: { maxBytes: number; maxMessages: number; maxScannedMessages: number }, ): SessionTranscriptBoundedContextMessageTailPage { - return readBoundedMessageTailPage( - scope, - options, - { - resolvePositionRange: resolveContextMessagePositionRange, - resolvePositions: resolveContextMessagePositions, - }, - (projection) => { - const contextSummary = resolveContextBoundarySummary(projection); - return contextSummary ? { contextSummary } : {}; - }, + return withCurrentProjectionSnapshot(scope, (projection) => + readBoundedContextMessageTail(projection, options), ); } diff --git a/src/config/sessions/session-accessor.sqlite-reset-window.ts b/src/config/sessions/session-accessor.sqlite-reset-window.ts index ff31383e249f..f6fc49503036 100644 --- a/src/config/sessions/session-accessor.sqlite-reset-window.ts +++ b/src/config/sessions/session-accessor.sqlite-reset-window.ts @@ -1,5 +1,6 @@ // Reset and model-context boundaries project logical message windows without // rewriting raw cursor positions. +import { sql } from "kysely"; import { executeSqliteQuerySync, executeSqliteQueryTakeFirstSync, @@ -42,24 +43,26 @@ type ContextBoundarySummary = { ts: number; }; -type BoundaryMessageWindow = { - contextSummary?: ContextBoundarySummary; +type ResetMessageWindow = { generation: string | undefined; indexedSeq: number; keptMessagePositions: number[]; postBoundaryMessagePosition: number; }; -type BoundaryMessageWindowCacheEntry = { +type ResetMessageWindowCacheEntry = { generation: string | undefined; indexedSeq: number; - window: BoundaryMessageWindow | null; + window: ResetMessageWindow | null; }; -type BoundaryWindowMode = "reset-only" | "context"; +type SessionTranscriptContextWindow = { + contextSummary?: ContextBoundarySummary; + scanStartActivePosition: number; +}; -const boundaryMessageWindowCache = new Map(); -const MAX_BOUNDARY_MESSAGE_WINDOW_CACHE = 128; +const resetMessageWindowCache = new Map(); +const MAX_MESSAGE_WINDOW_CACHE = 64; function getResetWindowKysely(database: OpenClawAgentDatabase) { return getNodeSqliteKysely(database.db); @@ -105,20 +108,8 @@ function readMessageRange( ).rows.map(parseMessageEventRow); } -function parseTranscriptEventType(eventJson: string): string | undefined { - try { - const parsed = JSON.parse(eventJson) as { type?: unknown }; - return typeof parsed.type === "string" ? parsed.type : undefined; - } catch { - return undefined; - } -} - -function boundaryMessageWindowCacheKey( - projection: ResetWindowProjection, - mode: BoundaryWindowMode, -): string { - return `${projection.database.path}\0${projection.resolved.sessionId}\0${mode}`; +function messageWindowCacheKey(projection: ResetWindowProjection): string { + return `${projection.database.path}\0${projection.resolved.sessionId}`; } function readTranscriptGeneration(projection: ResetWindowProjection): string | undefined { @@ -131,52 +122,82 @@ function readTranscriptGeneration(projection: ResetWindowProjection): string | u )?.generation; } -function cacheBoundaryMessageWindow(key: string, entry: BoundaryMessageWindowCacheEntry): void { - boundaryMessageWindowCache.delete(key); - boundaryMessageWindowCache.set(key, entry); - pruneMapToMaxSize(boundaryMessageWindowCache, MAX_BOUNDARY_MESSAGE_WINDOW_CACHE); -} - -function findLatestBoundaryMessageWindow( +function readLatestActiveBoundaryByType( projection: ResetWindowProjection, - generation: string | undefined, - mode: BoundaryWindowMode, -): BoundaryMessageWindow | null { + eventType: "compaction" | "reset", +) { const db = getResetWindowKysely(projection.database); - const nonMessageRows = executeSqliteQuerySync( + return executeSqliteQueryTakeFirstSync( projection.database.db, db .selectFrom("session_transcript_active_events as active") + .innerJoin("transcript_event_identities as identity", (join) => + join + .onRef("identity.session_id", "=", "active.session_id") + .onRef("identity.seq", "=", "active.event_seq"), + ) .innerJoin("transcript_events as event", (join) => join .onRef("event.session_id", "=", "active.session_id") .onRef("event.seq", "=", "active.event_seq"), ) - .select(["active.active_position", "event.event_json"]) + .select(["active.active_position", "identity.event_type", "identity.seq", "event.event_json"]) .where("active.session_id", "=", projection.resolved.sessionId) - .where("active.message_position", "is", null) - .orderBy("active.active_position", "desc"), - ).rows; - const latestBoundaryRow = nonMessageRows.find((row) => { - const type = parseTranscriptEventType(row.event_json); - return type === "reset" || type === "compaction"; - }); - const boundaryType = latestBoundaryRow - ? parseTranscriptEventType(latestBoundaryRow.event_json) + .where("identity.event_type", "=", eventType) + .orderBy("identity.seq", "desc") + .limit(1), + ); +} + +function readLatestActiveBoundary(projection: ResetWindowProjection) { + const reset = readLatestActiveBoundaryByType(projection, "reset"); + const compaction = readLatestActiveBoundaryByType(projection, "compaction"); + if (!reset) { + return compaction; + } + if (!compaction) { + return reset; + } + return reset.seq > compaction.seq ? reset : compaction; +} + +function readFirstKeptActivePosition( + projection: ResetWindowProjection, + firstKeptEntryId: unknown, + boundaryActivePosition: number, +): number | undefined { + if (typeof firstKeptEntryId !== "string") { + return undefined; + } + const db = getResetWindowKysely(projection.database); + const firstKept = executeSqliteQueryTakeFirstSync( + projection.database.db, + db + .selectFrom("transcript_event_identities as identity") + .innerJoin("session_transcript_active_events as active", (join) => + join + .onRef("active.session_id", "=", "identity.session_id") + .onRef("active.event_seq", "=", "identity.seq"), + ) + .select("active.active_position") + .where("identity.session_id", "=", projection.resolved.sessionId) + .where("identity.event_id", "=", firstKeptEntryId), + ); + return firstKept && firstKept.active_position < boundaryActivePosition + ? firstKept.active_position : undefined; - if ( - !latestBoundaryRow || - (mode === "reset-only" && boundaryType !== "reset") || - (boundaryType !== "reset" && boundaryType !== "compaction") - ) { +} + +function findLatestResetMessageWindow( + projection: ResetWindowProjection, + generation: string | undefined, +): ResetMessageWindow | null { + const db = getResetWindowKysely(projection.database); + const latestBoundaryRow = readLatestActiveBoundary(projection); + if (!latestBoundaryRow || latestBoundaryRow.event_type !== "reset") { return null; } - const boundaryRow = latestBoundaryRow; - const boundary = JSON.parse(boundaryRow.event_json) as { - firstKeptEntryId?: unknown; - summary?: unknown; - timestamp?: unknown; - }; + const boundary = JSON.parse(latestBoundaryRow.event_json) as { firstKeptEntryId?: unknown }; const postBoundaryMessagePosition = executeSqliteQueryTakeFirstSync( projection.database.db, @@ -184,61 +205,73 @@ function findLatestBoundaryMessageWindow( .selectFrom("session_transcript_active_events") .select("message_position") .where("session_id", "=", projection.resolved.sessionId) - .where("active_position", ">", boundaryRow.active_position) + .where("active_position", ">", latestBoundaryRow.active_position) .where("message_position", "is not", null) .orderBy("active_position", "asc") .limit(1), )?.message_position ?? projection.state.activeMessageCount; let keptMessagePositions: number[] = []; - if (typeof boundary.firstKeptEntryId === "string") { - const firstKept = executeSqliteQueryTakeFirstSync( + const firstKeptActivePosition = readFirstKeptActivePosition( + projection, + boundary.firstKeptEntryId, + latestBoundaryRow.active_position, + ); + if (firstKeptActivePosition !== undefined) { + keptMessagePositions = executeSqliteQuerySync( projection.database.db, db - .selectFrom("transcript_event_identities as identity") - .innerJoin("session_transcript_active_events as active", (join) => + .selectFrom("session_transcript_active_events as active") + .innerJoin("transcript_events as event", (join) => join - .onRef("active.session_id", "=", "identity.session_id") - .onRef("active.event_seq", "=", "identity.seq"), + .onRef("event.session_id", "=", "active.session_id") + .onRef("event.seq", "=", "active.event_seq"), ) - .select("active.active_position") - .where("identity.session_id", "=", projection.resolved.sessionId) - .where("identity.event_id", "=", boundary.firstKeptEntryId), - ); - if (firstKept && firstKept.active_position < boundaryRow.active_position) { - keptMessagePositions = executeSqliteQuerySync( - projection.database.db, - db - .selectFrom("session_transcript_active_events as active") - .innerJoin("transcript_events as event", (join) => - join - .onRef("event.session_id", "=", "active.session_id") - .onRef("event.seq", "=", "active.event_seq"), - ) - .select(["active.message_position", "event.event_json"]) - .where("active.session_id", "=", projection.resolved.sessionId) - .where("active.active_position", ">=", firstKept.active_position) - .where("active.active_position", "<", boundaryRow.active_position) - .where("active.message_position", "is not", null) - .orderBy("active.active_position", "asc"), - ).rows.flatMap((row) => { - if (row.message_position === null) { - return []; - } - try { - const role = (JSON.parse(row.event_json) as { message?: { role?: unknown } }).message - ?.role; - if (boundaryType === "compaction") { - return [row.message_position]; - } - return role === "user" || role === "assistant" ? [row.message_position] : []; - } catch { - return []; - } - }); - } + .select(["active.message_position", "event.event_json"]) + .where("active.session_id", "=", projection.resolved.sessionId) + .where("active.active_position", ">=", firstKeptActivePosition) + .where("active.active_position", "<", latestBoundaryRow.active_position) + .where("active.message_position", "is not", null) + .orderBy("active.active_position", "asc"), + ).rows.flatMap((row) => { + if (row.message_position === null) { + return []; + } + try { + const role = (JSON.parse(row.event_json) as { message?: { role?: unknown } }).message?.role; + return role === "user" || role === "assistant" ? [row.message_position] : []; + } catch { + return []; + } + }); } return { - ...(boundaryType === "compaction" && typeof boundary.summary === "string" + generation, + indexedSeq: projection.state.indexedSeq, + keptMessagePositions, + postBoundaryMessagePosition, + }; +} + +function findContextMessageWindow( + projection: ResetWindowProjection, +): SessionTranscriptContextWindow | null { + const latestBoundaryRow = readLatestActiveBoundary(projection); + if (!latestBoundaryRow) { + return null; + } + const boundary = JSON.parse(latestBoundaryRow.event_json) as { + firstKeptEntryId?: unknown; + summary?: unknown; + timestamp?: unknown; + }; + const retainedStartActivePosition = readFirstKeptActivePosition( + projection, + boundary.firstKeptEntryId, + latestBoundaryRow.active_position, + ); + return { + scanStartActivePosition: retainedStartActivePosition ?? latestBoundaryRow.active_position + 1, + ...(latestBoundaryRow.event_type === "compaction" && typeof boundary.summary === "string" ? { contextSummary: { text: boundary.summary, @@ -251,60 +284,39 @@ function findLatestBoundaryMessageWindow( }, } : {}), - generation, - indexedSeq: projection.state.indexedSeq, - keptMessagePositions, - postBoundaryMessagePosition, }; } -function resolveBoundaryMessageWindow( - projection: ResetWindowProjection, - mode: BoundaryWindowMode, -): BoundaryMessageWindow | null { - const key = boundaryMessageWindowCacheKey(projection, mode); - const cached = boundaryMessageWindowCache.get(key); +function resolveResetMessageWindow(projection: ResetWindowProjection): ResetMessageWindow | null { + const key = messageWindowCacheKey(projection); + const cached = resetMessageWindowCache.get(key); const generation = readTranscriptGeneration(projection); if (cached) { if (cached.generation === generation && cached.indexedSeq === projection.state.indexedSeq) { return cached.window; } } - const window = findLatestBoundaryMessageWindow(projection, generation, mode); - cacheBoundaryMessageWindow(key, { + const window = findLatestResetMessageWindow(projection, generation); + resetMessageWindowCache.delete(key); + resetMessageWindowCache.set(key, { generation, indexedSeq: projection.state.indexedSeq, window, }); + pruneMapToMaxSize(resetMessageWindowCache, MAX_MESSAGE_WINDOW_CACHE); return window; } +function resolveContextMessageWindow( + projection: ResetWindowProjection, +): SessionTranscriptContextWindow | null { + return findContextMessageWindow(projection); +} + export function resolveVisibleMessagePositions( projection: ResetWindowProjection, ): VisibleMessagePositions { - return toVisibleMessagePositions( - projection, - resolveBoundaryMessageWindow(projection, "reset-only"), - ); -} - -/** Mirrors model replay's latest reset/compaction boundary without materializing history. */ -export function resolveContextMessagePositions( - projection: ResetWindowProjection, -): VisibleMessagePositions { - return toVisibleMessagePositions(projection, resolveBoundaryMessageWindow(projection, "context")); -} - -export function resolveContextBoundarySummary( - projection: ResetWindowProjection, -): ContextBoundarySummary | undefined { - return resolveBoundaryMessageWindow(projection, "context")?.contextSummary; -} - -function toVisibleMessagePositions( - projection: ResetWindowProjection, - window: BoundaryMessageWindow | null, -): VisibleMessagePositions { + const window = resolveResetMessageWindow(projection); if (!window) { return { kept: [], postStart: 0, total: projection.state.activeMessageCount }; } @@ -345,16 +357,13 @@ export function readVisibleMessageRange( return [...keptEvents, ...postEvents]; } -function resolveMessagePositionRange( +/** Maps a logical transcript-visible range to materialized message positions. */ +export function resolveVisibleMessagePositionRange( projection: ResetWindowProjection, start: number, endExclusive: number, - resolvePositions: (projection: ResetWindowProjection) => VisibleMessagePositions, ): number[] { - if (endExclusive <= start) { - return []; - } - const visible = resolvePositions(projection); + const visible = resolveVisibleMessagePositions(projection); const boundedStart = Math.min(Math.max(0, start), visible.total); const boundedEnd = Math.min(Math.max(boundedStart, endExclusive), visible.total); const keptEnd = Math.min(boundedEnd, visible.kept.length); @@ -367,30 +376,83 @@ function resolveMessagePositionRange( return positions; } -/** Maps a logical transcript-visible range to materialized message positions. */ -export function resolveVisibleMessagePositionRange( +/** Reads one authoritative bounded model-context tail from the active semantic window. */ +export function readBoundedContextMessageTail( projection: ResetWindowProjection, - start: number, - endExclusive: number, -): number[] { - return resolveMessagePositionRange( - projection, - start, - endExclusive, - resolveVisibleMessagePositions, - ); -} - -/** Maps a logical model-context range to materialized message positions. */ -export function resolveContextMessagePositionRange( - projection: ResetWindowProjection, - start: number, - endExclusive: number, -): number[] { - return resolveMessagePositionRange( - projection, - start, - endExclusive, - resolveContextMessagePositions, - ); + options: { maxBytes: number; maxMessages: number; maxScannedMessages: number }, +) { + const maxMessages = Math.max(0, Math.floor(options.maxMessages)); + const maxScannedMessages = Math.max(0, Math.floor(options.maxScannedMessages)); + const maxBytes = Math.max(0, Math.floor(options.maxBytes)); + const contextWindow = resolveContextMessageWindow(projection); + const db = getResetWindowKysely(projection.database); + const metadata = executeSqliteQuerySync( + projection.database.db, + db + .selectFrom("session_transcript_active_events as active") + .innerJoin("transcript_events as event", (join) => + join + .onRef("event.session_id", "=", "active.session_id") + .onRef("event.seq", "=", "active.event_seq"), + ) + .select([ + "active.message_position", + /* kysely-allow-raw: inspect role and bytes without materializing message payloads. */ + sql`CASE WHEN json_valid(event.event_json) + THEN json_extract(event.event_json, '$.message.role') ELSE NULL END`.as("message_role"), + sql`LENGTH(CAST(event.event_json AS BLOB)) + 1`.as("serialized_bytes"), + ]) + .where("active.session_id", "=", projection.resolved.sessionId) + .where("active.message_position", "is not", null) + .$if(contextWindow !== null, (query) => + query.where("active.active_position", ">=", contextWindow?.scanStartActivePosition ?? 0), + ) + .orderBy("active.active_position", "desc") + .limit(maxScannedMessages + 1), + ).rows; + const selectedPositions: number[] = []; + let serializedBytes = 0; + let blockedByBytes = false; + for (const row of metadata.slice(0, maxScannedMessages)) { + if ( + row.message_position === null || + (row.message_role !== "assistant" && row.message_role !== "user") + ) { + continue; + } + if (selectedPositions.length >= maxMessages) { + break; + } + if (serializedBytes + row.serialized_bytes > maxBytes) { + blockedByBytes = true; + break; + } + selectedPositions.push(row.message_position); + serializedBytes += row.serialized_bytes; + } + const events = + selectedPositions.length === 0 + ? [] + : executeSqliteQuerySync( + projection.database.db, + db + .selectFrom("session_transcript_active_events as active") + .innerJoin("transcript_events as event", (join) => + join + .onRef("event.session_id", "=", "active.session_id") + .onRef("event.seq", "=", "active.event_seq"), + ) + .select(["active.message_position", "event.event_json"]) + .where("active.session_id", "=", projection.resolved.sessionId) + .where("active.message_position", "in", selectedPositions) + .orderBy("active.message_position", "asc"), + ).rows.map(parseMessageEventRow); + return { + authoritative: + !blockedByBytes && + (selectedPositions.length >= maxMessages || metadata.length <= maxScannedMessages), + ...(contextWindow?.contextSummary ? { contextSummary: contextWindow.contextSummary } : {}), + empty: metadata.length === 0 && !contextWindow?.contextSummary, + events, + }; } diff --git a/src/gateway/session-companion-ask.ts b/src/gateway/session-companion-ask.ts index f14a4efab465..fd4d340d9389 100644 --- a/src/gateway/session-companion-ask.ts +++ b/src/gateway/session-companion-ask.ts @@ -68,6 +68,17 @@ type SessionCompanionAskRuntimeParams = SessionCompanionAskDeps & { isDisposed: () => boolean; }; +type SessionCompanionCancellationKind = + | "backing-session-revoked" + | "disposed" + | "explicit-reset" + | "timeout"; + +type SessionCompanionActiveAsk = { + cancellation?: SessionCompanionCancellationKind; + controller: AbortController; +}; + type SessionCompanionAskErrorReason = | "busy" | "context-unavailable" @@ -367,7 +378,7 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt const run = params.run ?? defaultRun; const setTimeoutFn = params.setTimeoutFn ?? setTimeout; const clearTimeoutFn = params.clearTimeoutFn ?? clearTimeout; - const controllers = new Map(); + const activeAsks = new Map(); const preparations = new Map>(); const admissions: Array<{ connId: string; admittedAt: number }> = []; @@ -404,10 +415,7 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt const preparation = (async () => { const result = await contextReader.read({ agentId, sessionKey, signal }); if (signal.aborted || params.isDisposed()) { - throw contextError( - "context-unavailable", - "The selected session changed before its history was ready.", - ); + throw new Error("session companion preparation was cancelled"); } if (result.kind === "missing") { throw contextError("session-missing", "The selected session is no longer available."); @@ -456,7 +464,7 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt throw new SessionCompanionAskError("unavailable", "Session companion is unavailable."); } const existing = params.threads.get(sessionKey); - if (existing?.busy || controllers.has(sessionKey)) { + if (existing?.busy || activeAsks.has(sessionKey)) { throw new SessionCompanionAskError( "busy", "The session companion is answering another question.", @@ -482,7 +490,7 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt ) : 0; if ( - controllers.size >= MAX_CONCURRENT_ASKS || + activeAsks.size >= MAX_CONCURRENT_ASKS || globalRetryAfterMs > 0 || connectionRetryAfterMs > 0 ) { @@ -490,7 +498,7 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt "rate-limited", "The session companion has reached its question limit. Try again shortly.", Math.max( - controllers.size >= MAX_CONCURRENT_ASKS ? ASK_TIMEOUT_MS : 0, + activeAsks.size >= MAX_CONCURRENT_ASKS ? ASK_TIMEOUT_MS : 0, globalRetryAfterMs, connectionRetryAfterMs, ), @@ -499,8 +507,16 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt admissions.push({ connId: request.connId, admittedAt }); const controller = new AbortController(); - controllers.set(sessionKey, controller); - const timeout = setTimeoutFn(() => controller.abort(), ASK_TIMEOUT_MS); + const activeAsk: SessionCompanionActiveAsk = { controller }; + activeAsks.set(sessionKey, activeAsk); + const abort = (cancellation: SessionCompanionCancellationKind) => { + if (activeAsks.get(sessionKey) !== activeAsk || activeAsk.cancellation) { + return; + } + activeAsk.cancellation = cancellation; + controller.abort(); + }; + const timeout = setTimeoutFn(() => abort("timeout"), ASK_TIMEOUT_MS); const aborted = new Promise((_resolve, reject) => { controller.signal.addEventListener( "abort", @@ -508,8 +524,16 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt { once: true }, ); }); + let ownedAgentId: string | undefined; + let ownedThread: SessionCompanionThread | undefined; + const discardOwnedThread = () => { + if (ownedThread && params.threads.get(sessionKey) === ownedThread) { + params.threads.delete(sessionKey); + } + }; try { const thread = await prepareThread(sessionKey, controller.signal); + ownedThread = thread; notifySessionCompanionPrepared({ connId: request.connId, empty: thread.context.empty, @@ -524,6 +548,7 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt thread.busy = true; thread.lastUsedAt = admittedAt; const { agentId, cfg } = resolveTarget(sessionKey); + ownedAgentId = agentId; if (currentSessionId(sessionKey, agentId) !== thread.context.sessionId) { params.threads.delete(sessionKey); throw contextError( @@ -565,13 +590,25 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt }), aborted, ]); + if (activeAsk.cancellation === "backing-session-revoked") { + discardOwnedThread(); + throw contextError( + "context-unavailable", + "The selected session changed before the companion could answer.", + ); + } + if (activeAsk.cancellation || params.isDisposed()) { + throw new Error("session companion ask was cancelled"); + } if ( - controller.signal.aborted || - params.isDisposed() || params.threads.get(sessionKey) !== thread || currentSessionId(sessionKey, agentId) !== thread.context.sessionId ) { - throw new Error("session companion ask is no longer active"); + discardOwnedThread(); + throw contextError( + "context-unavailable", + "The selected session changed before the companion could answer.", + ); } const answer = sanitizeAnswer(rawAnswer); if (!answer) { @@ -588,33 +625,68 @@ export function createSessionCompanionAskRuntime(params: SessionCompanionAskRunt if (error instanceof SessionCompanionAskError) { throw error; } + if (activeAsk.cancellation === "backing-session-revoked") { + discardOwnedThread(); + throw contextError( + "context-unavailable", + "The selected session changed before the companion could answer.", + ); + } + if (!activeAsk.cancellation && ownedThread) { + if ( + params.threads.get(sessionKey) !== ownedThread || + (ownedAgentId && + currentSessionId(sessionKey, ownedAgentId) !== ownedThread.context.sessionId) + ) { + discardOwnedThread(); + throw contextError( + "context-unavailable", + "The selected session changed before the companion could answer.", + ); + } + } companionLog.warn("session companion ask failed", { sessionKey, error }); throw new SessionCompanionAskError( "unavailable", - "The session companion could not answer right now.", + activeAsk.cancellation === "timeout" + ? "The session companion timed out." + : activeAsk.cancellation === "explicit-reset" + ? "The session companion request was cancelled." + : "The session companion could not answer right now.", ); } finally { clearTimeoutFn(timeout); - if (controllers.get(sessionKey) === controller) { - controllers.delete(sessionKey); + if (activeAsks.get(sessionKey) === activeAsk) { + activeAsks.delete(sessionKey); } - const thread = params.threads.get(sessionKey); - if (thread) { - thread.busy = false; + if (ownedThread && params.threads.get(sessionKey) === ownedThread) { + ownedThread.busy = false; } } }; return { ask, - cancel(sessionKey: string) { - controllers.get(sessionKey)?.abort(); + cancel( + sessionKey: string, + cancellation: Extract< + SessionCompanionCancellationKind, + "backing-session-revoked" | "explicit-reset" + >, + ) { + const activeAsk = activeAsks.get(sessionKey); + if (!activeAsk || activeAsk.cancellation) { + return; + } + activeAsk.cancellation = cancellation; + activeAsk.controller.abort(); }, dispose() { - for (const controller of controllers.values()) { - controller.abort(); + for (const activeAsk of activeAsks.values()) { + activeAsk.cancellation ??= "disposed"; + activeAsk.controller.abort(); } - controllers.clear(); + activeAsks.clear(); preparations.clear(); admissions.length = 0; }, diff --git a/src/gateway/session-companion-context.test.ts b/src/gateway/session-companion-context.test.ts index 72afcb9325ef..c33211cff047 100644 --- a/src/gateway/session-companion-context.test.ts +++ b/src/gateway/session-companion-context.test.ts @@ -110,6 +110,72 @@ describe("session companion context", () => { ); }); + it.each([ + { expectedKind: "ready", unsupportedCount: 4095 }, + { expectedKind: "unavailable", unsupportedCount: 4096 }, + ] as const)( + "fails closed only when unsupported roles exceed the bounded scan ($unsupportedCount)", + async ({ expectedKind, unsupportedCount }) => { + const scope = createScope(`companion-context-unsupported-${unsupportedCount}`); + await upsertSessionEntryCore(scope, { sessionId: scope.sessionId, updatedAt: 1 }); + const messages = [ + { + eventId: "visible-user", + parentId: null, + message: { role: "user" as const, content: "authoritative question", timestamp: 1 }, + }, + ...Array.from({ length: unsupportedCount }, (_, index) => ({ + eventId: `custom-${index}`, + parentId: index === 0 ? "visible-user" : `custom-${index - 1}`, + message: { + role: "custom" as const, + customType: "test-context", + content: `unsupported ${index}`, + timestamp: index + 2, + }, + })), + ]; + await persistSessionTranscriptTurn(scope, { messages, touchSessionEntry: true }); + + const result = await defaultSessionCompanionContextReader.read(scope); + + expect(result.kind).toBe(expectedKind); + if (result.kind === "ready") { + expect(result.context.messages.map((message) => message.text)).toEqual([ + "authoritative question", + ]); + } + }, + ); + + it("returns unavailable instead of backfilling past an oversized latest user message", async () => { + const scope = createScope("companion-context-oversized-latest"); + await upsertSessionEntryCore(scope, { sessionId: scope.sessionId, updatedAt: 1 }); + await persistSessionTranscriptTurn(scope, { + messages: [ + { + eventId: "older", + parentId: null, + message: { role: "user" as const, content: "stale older question", timestamp: 1 }, + }, + { + eventId: "oversized", + parentId: "older", + message: { + role: "user" as const, + content: "x".repeat(1024 * 1024), + timestamp: 2, + }, + }, + ], + touchSessionEntry: true, + }); + + await expect(defaultSessionCompanionContextReader.read(scope)).resolves.toEqual({ + kind: "unavailable", + }); + }); + it("preserves the latest compaction summary and retained context without resurrecting history", async () => { const scope = createScope("companion-context-compaction"); await upsertSessionEntryCore(scope, { sessionId: scope.sessionId, updatedAt: 1 }); diff --git a/src/gateway/session-companion-context.ts b/src/gateway/session-companion-context.ts index 1f3947974911..6eaaeb2a8ed3 100644 --- a/src/gateway/session-companion-context.ts +++ b/src/gateway/session-companion-context.ts @@ -17,7 +17,6 @@ import { loadSessionEntryReadOnly } from "./session-utils.js"; const CONTEXT_MAX_MESSAGES = 40; const CONTEXT_MAX_BYTES = 24 * 1024; const CONTEXT_MESSAGE_MAX_CHARS = 4000; -const CONTEXT_READ_PAGE_MESSAGES = CONTEXT_MAX_MESSAGES * 4; const CONTEXT_READ_MAX_SCANNED_MESSAGES = 4096; const CONTEXT_READ_MAX_BYTES = 1024 * 1024; @@ -142,48 +141,22 @@ async function readSessionCompanionContext(params: { sessionKey: params.sessionKey, storePath: loaded.storePath, }; - const messages: unknown[] = []; - let activeLeafEntryId: string | null | undefined; - let contextSummary: { text: string; ts: number } | undefined; - let offset = 0; - let serializedBytes = 0; - let totalMessages = Number.POSITIVE_INFINITY; - while ( - stripToolMessages(messages).length < CONTEXT_MAX_MESSAGES && - offset < totalMessages && - offset < CONTEXT_READ_MAX_SCANNED_MESSAGES && - serializedBytes < CONTEXT_READ_MAX_BYTES - ) { - if (params.signal?.aborted) { - return { kind: "unavailable" }; - } - const page = readSessionTranscriptBoundedContextMessageTailPage(scope, { - maxBytes: CONTEXT_READ_MAX_BYTES - serializedBytes, - maxMessages: Math.min( - CONTEXT_READ_PAGE_MESSAGES, - CONTEXT_READ_MAX_SCANNED_MESSAGES - offset, - ), - offset, - }); - if (activeLeafEntryId === undefined) { - activeLeafEntryId = page.activeLeafEntryId; - contextSummary = page.contextSummary; - } else if (page.activeLeafEntryId !== activeLeafEntryId) { - return { kind: "unavailable" }; - } - totalMessages = page.totalMessages; - if (page.scannedMessages === 0) { - break; - } - messages.unshift(...readPageMessages(page.events)); - offset += page.scannedMessages; - serializedBytes += page.serializedBytes; + if (params.signal?.aborted) { + return { kind: "unavailable" }; } - const selected = sanitizeContextMessages(messages, contextSummary); + const page = readSessionTranscriptBoundedContextMessageTailPage(scope, { + maxBytes: CONTEXT_READ_MAX_BYTES, + maxMessages: CONTEXT_MAX_MESSAGES, + maxScannedMessages: CONTEXT_READ_MAX_SCANNED_MESSAGES, + }); + if (!page.authoritative || params.signal?.aborted) { + return { kind: "unavailable" }; + } + const selected = sanitizeContextMessages(readPageMessages(page.events), page.contextSummary); return { kind: "ready", context: { - empty: totalMessages === 0 && !contextSummary, + empty: page.empty, messages: selected, sessionId, }, diff --git a/src/gateway/session-companion.test.ts b/src/gateway/session-companion.test.ts index 3b0d26240d5b..7b054e845eb2 100644 --- a/src/gateway/session-companion.test.ts +++ b/src/gateway/session-companion.test.ts @@ -262,10 +262,19 @@ describe("session companion asks", () => { it("discards an answer when the backing session identity changes", async () => { vi.useFakeTimers(); let sessionId = "session-1"; + let runCount = 0; const pending = deferred(); const harness = createHarness({ currentSessionId: () => sessionId, - run: async () => await pending.promise, + readContext: async () => ({ + kind: "ready", + context: { + empty: false, + messages: [{ role: "user", text: `question for ${sessionId}`, ts: 1 }], + sessionId, + }, + }), + run: async () => (runCount++ === 0 ? await pending.promise : "fresh answer"), }); const active = harness.service.ask({ sessionKey: "agent:main:main", @@ -277,9 +286,18 @@ describe("session companion asks", () => { pending.resolve("stale answer"); await expect(active).rejects.toMatchObject({ - reason: "unavailable", + reason: "context-unavailable", } satisfies Partial); expect(harness.service.state("agent:main:main")).toEqual({ exchanges: [] }); + + await expect( + harness.service.ask({ + sessionKey: "agent:main:main", + question: "Which session now?", + connId: "conn-1", + }), + ).resolves.toMatchObject({ answer: "fresh answer" }); + expect(harness.readContext).toHaveBeenCalledTimes(2); harness.service.dispose(); }); @@ -454,6 +472,69 @@ describe("session companion asks", () => { harness.service.dispose(); }); + it("makes a committed backing-session reset retryable and ignores the late model result", async () => { + vi.useFakeTimers(); + const pending = deferred(); + const harness = createHarness({ run: async () => await pending.promise }); + const active = harness.service.ask({ + sessionKey: "agent:main:main", + question: "Still the same backing session?", + connId: "conn-1", + }); + await vi.waitFor(() => expect(harness.run).toHaveBeenCalledOnce()); + + notifyGatewaySessionReset("agent:main:main"); + pending.resolve("stale answer"); + + await expect(active).rejects.toMatchObject({ + reason: "context-unavailable", + } satisfies Partial); + expect(harness.service.state("agent:main:main")).toEqual({ exchanges: [] }); + harness.service.dispose(); + }); + + it("disposal cancels an active ask without committing its late model result", async () => { + vi.useFakeTimers(); + const pending = deferred(); + const harness = createHarness({ run: async () => await pending.promise }); + const active = harness.service.ask({ + sessionKey: "agent:main:main", + question: "Will this survive shutdown?", + connId: "conn-1", + }); + await vi.waitFor(() => expect(harness.run).toHaveBeenCalledOnce()); + + harness.service.dispose(); + pending.resolve("late answer"); + + await expect(active).rejects.toMatchObject({ + reason: "unavailable", + } satisfies Partial); + expect(harness.service.state("agent:main:main")).toEqual({ exchanges: [] }); + }); + + it("keeps provider failures terminal after one model call", async () => { + vi.useFakeTimers(); + const harness = createHarness({ + run: async () => { + throw new Error("provider unavailable"); + }, + }); + + await expect( + harness.service.ask({ + sessionKey: "agent:main:main", + question: "Can the provider answer?", + connId: "conn-1", + }), + ).rejects.toMatchObject({ + reason: "unavailable", + } satisfies Partial); + expect(harness.run).toHaveBeenCalledOnce(); + expect(harness.service.state("agent:main:main")).toEqual({ exchanges: [] }); + harness.service.dispose(); + }); + it("clears a thread when the committed gateway reset path notifies", async () => { vi.useFakeTimers(); const harness = createHarness(); diff --git a/src/gateway/session-companion.ts b/src/gateway/session-companion.ts index 71115df01833..58bd998293a5 100644 --- a/src/gateway/session-companion.ts +++ b/src/gateway/session-companion.ts @@ -41,12 +41,15 @@ export function createSessionCompanion(deps: SessionCompanionDeps): SessionCompa isDisposed: () => disposed, }); - const reset = (sessionKey: string) => { + const reset = ( + sessionKey: string, + cancellation: "backing-session-revoked" | "explicit-reset", + ) => { const key = sessionKey.trim(); if (!key) { return; } - askRuntime.cancel(key); + askRuntime.cancel(key, cancellation); threads.delete(key); }; @@ -54,13 +57,15 @@ export function createSessionCompanion(deps: SessionCompanionDeps): SessionCompa const cutoff = now() - SESSION_COMPANION_IDLE_TTL_MS; for (const [sessionKey, thread] of threads) { if (!thread.busy && thread.lastUsedAt <= cutoff) { - reset(sessionKey); + reset(sessionKey, "explicit-reset"); } } }; const sweepTimer = setIntervalFn(sweep, SESSION_COMPANION_SWEEP_INTERVAL_MS); sweepTimer.unref?.(); - const unsubscribeReset = onGatewaySessionReset(reset); + const unsubscribeReset = onGatewaySessionReset((sessionKey) => + reset(sessionKey, "backing-session-revoked"), + ); return { ask: askRuntime.ask, @@ -75,7 +80,9 @@ export function createSessionCompanion(deps: SessionCompanionDeps): SessionCompa exchanges: thread.exchanges.map(({ question, answer, ts }) => ({ question, answer, ts })), }; }, - reset, + reset(sessionKey) { + reset(sessionKey, "explicit-reset"); + }, dispose() { if (disposed) { return; diff --git a/ui/src/pages/chat/chat-session-companion.ts b/ui/src/pages/chat/chat-session-companion.ts index 214bc9c1d7d3..2e2cd3c79282 100644 --- a/ui/src/pages/chat/chat-session-companion.ts +++ b/ui/src/pages/chat/chat-session-companion.ts @@ -181,6 +181,7 @@ export class ChatSessionCompanionThreads { if (!isCurrent()) { throw Object.assign(new Error("stale companion answer"), { details: { reason: "context-unavailable" }, + retryable: true, }); } thread.exchanges = [ @@ -192,6 +193,11 @@ export class ChatSessionCompanionThreads { return; } thread.failedQuestion = normalized; + if (!isCurrent()) { + thread.hint = "history-unavailable"; + thread.retryable = true; + return; + } const reason = errorDetailReason(error); thread.hint = errorDetailCode(error) === COMPANION_BUSY_DETAIL_CODE diff --git a/ui/src/pages/chat/chat-session-rail.test.ts b/ui/src/pages/chat/chat-session-rail.test.ts index f1659bc43eef..87b84e638f21 100644 --- a/ui/src/pages/chat/chat-session-rail.test.ts +++ b/ui/src/pages/chat/chat-session-rail.test.ts @@ -430,9 +430,67 @@ describe("ChatSessionCompanionThreads", () => { failedQuestion: "Which connection owns this?", hint: "history-unavailable", pendingQuestion: null, + retryable: true, }); }); + it("makes a rejected stale connection settlement retryable", async () => { + let current = true; + let rejectAnswer!: (error: Error) => void; + const threads = new ChatSessionCompanionThreads(); + const pending = threads.submit( + "one", + "Which connection rejected this?", + (_sessionKey, _question, prepared) => { + prepared(); + return new Promise((_resolve, reject) => { + rejectAnswer = reject; + }); + }, + () => current, + ); + await vi.waitFor(() => expect(threads.view("one").phase).toBe("answering")); + current = false; + rejectAnswer(new Error("old socket closed")); + await pending; + + expect(threads.view("one")).toMatchObject({ + exchanges: [], + failedQuestion: "Which connection rejected this?", + hint: "history-unavailable", + retryable: true, + }); + }); + + it.each(["resolve", "reject"] as const)( + "does not resurrect a reset pending request after late $outcome", + async (outcome) => { + let resolveAnswer!: (value: { answer: string; ts: number }) => void; + let rejectAnswer!: (error: Error) => void; + const threads = new ChatSessionCompanionThreads(); + const pending = threads.submit("one", "Will reset keep this?", () => { + return new Promise((resolve, reject) => { + resolveAnswer = resolve; + rejectAnswer = reject; + }); + }); + + await threads.reset("one", async () => ({ ok: true })); + if (outcome === "resolve") { + resolveAnswer({ answer: "late answer", ts: 5 }); + } else { + rejectAnswer(new Error("late error")); + } + await pending; + + expect(threads.view("one")).toMatchObject({ + exchanges: [], + failedQuestion: null, + pendingQuestion: null, + }); + }, + ); + it("clears local state only after the reset RPC succeeds", async () => { const threads = new ChatSessionCompanionThreads(); await threads.hydrate("one", async () => ({