diff --git a/components/adapters/claude-code/src/context-cleaner/snapshot.ts b/components/adapters/claude-code/src/context-cleaner/snapshot.ts index 7674628a..86710e53 100644 --- a/components/adapters/claude-code/src/context-cleaner/snapshot.ts +++ b/components/adapters/claude-code/src/context-cleaner/snapshot.ts @@ -3,6 +3,53 @@ import type { ContextItemRef, ModelContextSnapshot } from "@lightrsi/host-adapte import { buildToolResultSegments } from "../eviction.js"; +export function scheduledCleanerAttributionUnavailable(params: { + selectedTaskIds: readonly string[]; + approvalSnapshot: ModelContextSnapshot; + currentSnapshot: ModelContextSnapshot; +}): boolean { + const selectedTaskIds = new Set(params.selectedTaskIds); + const expectedTaskIdsByCallId = new Map>(); + for (const item of params.approvalSnapshot.items) { + if (!item.callId) continue; + const expectedTaskIds = (item.taskIds ?? []).filter((taskId) => selectedTaskIds.has(taskId)); + if (expectedTaskIds.length === 0) continue; + const accumulated = expectedTaskIdsByCallId.get(item.callId) ?? new Set(); + for (const taskId of expectedTaskIds) accumulated.add(taskId); + expectedTaskIdsByCallId.set(item.callId, accumulated); + } + + let unavailableAttribution = false; + for (const [callId, expectedTaskIds] of expectedTaskIdsByCallId) { + const approvalPair = params.approvalSnapshot.items.filter((item) => item.callId === callId + && (item.kind === "tool_call" || item.kind === "tool_result")); + const currentPair = params.currentSnapshot.items.filter((item) => item.callId === callId + && (item.kind === "tool_call" || item.kind === "tool_result")); + const approvalCalls = approvalPair.filter((item) => item.kind === "tool_call"); + const approvalResults = approvalPair.filter((item) => item.kind === "tool_result"); + const currentCalls = currentPair.filter((item) => item.kind === "tool_call"); + const currentResults = currentPair.filter((item) => item.kind === "tool_result"); + // Missing, partial, or ambiguous pairs are genuine context drift and stay + // on the existing stale-validation path. That deterministic drift wins + // over unavailable attribution on any other approved pair. + if (approvalCalls.length !== 1 + || approvalResults.length !== 1 + || currentCalls.length !== 1 + || currentResults.length !== 1) return false; + if (approvalCalls[0]!.fingerprint !== currentCalls[0]!.fingerprint + || approvalResults[0]!.fingerprint !== currentResults[0]!.fingerprint) return false; + const retainsApprovedTaskProof = [...expectedTaskIds].some((taskId) => ( + currentPair.every((item) => item.taskIds?.includes(taskId)) + )); + if (retainsApprovedTaskProof) continue; + // A different non-empty task attribution is also deterministic drift. + // Defer only when the complete pair has lost attribution evidence entirely. + if (currentPair.every((item) => (item.taskIds?.length ?? 0) > 0)) return false; + unavailableAttribution = true; + } + return unavailableAttribution; +} + function provenTaskIds( registry: SessionTaskRegistry, segmentId: string, diff --git a/components/adapters/claude-code/src/context-rewrite/semantic-pipeline.ts b/components/adapters/claude-code/src/context-rewrite/semantic-pipeline.ts index 6c2f8664..d5ad1057 100644 --- a/components/adapters/claude-code/src/context-rewrite/semantic-pipeline.ts +++ b/components/adapters/claude-code/src/context-rewrite/semantic-pipeline.ts @@ -77,6 +77,25 @@ export function buildUniqueToolCallTurnMap( return result; } +/** + * Recover historical tool-call ownership without claiming or persisting a new + * semantic turn. Manual Cleaner schedules intentionally pause the lifecycle + * planner, but the Cleaner still needs this read-only map to re-attribute the + * approved task scope on the next request. + */ +export async function loadPersistedToolCallTurnMap(params: { + stateDir: string; + sessionId: string; +}): Promise> { + const allSeqs = await listRawSemanticTurnSeqs(params.stateDir, params.sessionId); + const loaded = await Promise.all( + allSeqs.map((seq) => loadRawSemanticTurnRecord(params.stateDir, params.sessionId, seq)), + ); + return buildUniqueToolCallTurnMap( + loaded.filter((record): record is RawSemanticTurnRecord => record !== null), + ); +} + /** * Prepare the (lastProcessed, now] semantic-delta materials for one Claude * request: claim the per-session turn seq, persist this turn's raw record, diff --git a/components/adapters/claude-code/src/gateway-runtime.ts b/components/adapters/claude-code/src/gateway-runtime.ts index 1aef6cec..6ee8d72d 100644 --- a/components/adapters/claude-code/src/gateway-runtime.ts +++ b/components/adapters/claude-code/src/gateway-runtime.ts @@ -44,7 +44,10 @@ import { } from "./context-rewrite/snapshot-store.js"; import { appendOverlayHistory } from "./context-rewrite/overlay-history.js"; import { resolveClaudeTaskStateEstimator } from "./context-rewrite/estimator-config.js"; -import { prepareSemanticDelta } from "./context-rewrite/semantic-pipeline.js"; +import { + loadPersistedToolCallTurnMap, + prepareSemanticDelta, +} from "./context-rewrite/semantic-pipeline.js"; import { buildSegmentToStableIdMap } from "./context-rewrite/segment-stable-id-map.js"; import { buildContextMutationPlan, @@ -67,13 +70,19 @@ import { buildAnthropicGatewayModelList, mapClaudeVisibleModelToUpstreamModel } import { resolveLatestClaudeCodeSessionId } from "./session-state.js"; import { lookupRealSessionId, recordSessionMapping } from "./context-rewrite/session-map.js"; import { initializeClaudeCodeTokenPilotPreset } from "./preset.js"; -import { attributeClaudeSnapshotTasks } from "./context-cleaner/snapshot.js"; +import { + attributeClaudeSnapshotTasks, + scheduledCleanerAttributionUnavailable, +} from "./context-cleaner/snapshot.js"; import { abandonClaudeCleanerOverlay, finalizeClaudeCleanerOverlay, prepareClaudeCleanerOverlay, } from "./context-cleaner/runtime.js"; -import { readClaudeCleanerSchedule } from "./context-cleaner/scheduler.js"; +import { + readClaudeCleanerSchedule, + type ClaudeCleanerScheduledRecord, +} from "./context-cleaner/scheduler.js"; export type ClaudeCodeGatewayRuntime = { baseUrl: string; @@ -84,6 +93,7 @@ type ClaudeCodeGatewayRuntimeDependencies = { cloneRequestPayload?: typeof structuredClone; resolveEstimator?: typeof resolveClaudeTaskStateEstimator; persistTaskRegistry?: typeof persistSessionTaskRegistry; + readSnapshot?: typeof claudeContextRewriteBackend.readSnapshot; saveSnapshot?: typeof saveLatestClaudeSnapshot; }; @@ -413,6 +423,9 @@ export async function startClaudeCodeGatewayRuntime(params: { { outcome: "prepared" } > | undefined; let manualCleanerSuppressesAutomaticEviction = cleanerApplyControlRequest; + let manualCleanerSchedulePending = false; + let manualCleanerExecutionDeferred = false; + let manualCleanerSchedule: ClaudeCleanerScheduledRecord | undefined; if (cleanerApplyControlRequest) { lifecyclePlannerStatus = "deferred"; @@ -429,6 +442,10 @@ export async function startClaudeCodeGatewayRuntime(params: { sessionId, }); if (manualSchedule.outcome === "ready" || manualSchedule.outcome === "bypassed") { + manualCleanerSchedulePending = manualSchedule.outcome === "ready"; + manualCleanerSchedule = manualSchedule.outcome === "ready" + ? manualSchedule.record + : undefined; manualCleanerSuppressesAutomaticEviction = true; lifecyclePlannerStatus = "deferred"; lifecyclePlannerReasonCodes = [ @@ -478,7 +495,9 @@ export async function startClaudeCodeGatewayRuntime(params: { .digest("hex") .slice(0, 32); const { bindings: plannerBindings } = buildToolResultSegments(plannerMessages); - const plannerSnapshot = await claudeContextRewriteBackend.readSnapshot({ + const plannerSnapshot = await ( + params.dependencies?.readSnapshot ?? claudeContextRewriteBackend.readSnapshot + )({ sessionId, request: { sessionId, @@ -542,6 +561,26 @@ export async function startClaudeCodeGatewayRuntime(params: { } } + // A pending manual schedule skips the lifecycle planner so task state + // cannot advance before the approved rewrite. Restore the already- + // persisted tool-call map separately so current-scope validation can + // still prove that the relocated items belong to the selected task. + if (manualCleanerSuppressesAutomaticEviction && !semanticTurnByToolCallId) { + try { + semanticTurnByToolCallId = await loadPersistedToolCallTurnMap({ + stateDir: config.stateDir, + sessionId, + }); + } catch (error) { + manualCleanerExecutionDeferred = manualCleanerSchedulePending; + logger.warn( + `context cleaner historical task attribution failed (${manualCleanerExecutionDeferred + ? "scheduled clean deferred" + : "ignored"}): ${String(error)}`, + ); + } + } + // A scheduled manual clean must validate against the canonical snapshot // from approval time. Read that base before this request replaces it. let previousCleanerSnapshot: Awaited>; @@ -550,6 +589,10 @@ export async function startClaudeCodeGatewayRuntime(params: { } catch (error) { logger.warn(`context cleaner base snapshot read failed (ignored): ${String(error)}`); } + if (manualCleanerSchedulePending && !previousCleanerSnapshot) { + manualCleanerExecutionDeferred = true; + logger.warn("context cleaner approval snapshot unavailable (scheduled clean deferred)"); + } // Build the current canonical snapshot after lifecycle update, then let a // pending manual Cleaner plan preempt automatic eviction for this request. @@ -559,7 +602,9 @@ export async function startClaudeCodeGatewayRuntime(params: { | undefined; let cleanerRegistry: Awaited> | undefined; try { - const baseCleanerSnapshot = await claudeContextRewriteBackend.readSnapshot({ + const baseCleanerSnapshot = await ( + params.dependencies?.readSnapshot ?? claudeContextRewriteBackend.readSnapshot + )({ sessionId, request: { sessionId, @@ -576,16 +621,36 @@ export async function startClaudeCodeGatewayRuntime(params: { registry: cleanerRegistry, turnAbsIdByToolCallId: semanticTurnByToolCallId, }); + if (manualCleanerSchedule + && previousCleanerSnapshot + && scheduledCleanerAttributionUnavailable({ + selectedTaskIds: manualCleanerSchedule.selectedTaskIds, + approvalSnapshot: previousCleanerSnapshot.snapshot, + currentSnapshot: cleanerSnapshot, + })) { + manualCleanerExecutionDeferred = true; + logger.warn("context cleaner task attribution unavailable (scheduled clean deferred)"); + } } catch (error) { - // Task attribution is optional. Keep the canonical snapshot even when - // registry recovery fails so Cleaner can still inspect unassigned context. - logger.warn(`context cleaner task attribution failed (ignored): ${String(error)}`); + // Ordinary requests may keep an unassigned snapshot, but a pending + // manual clean must retain its approval-time snapshot and retry. + manualCleanerExecutionDeferred = manualCleanerSchedulePending; + logger.warn( + `context cleaner task attribution failed (${manualCleanerExecutionDeferred + ? "scheduled clean deferred" + : "ignored"}): ${String(error)}`, + ); } } catch (error) { - logger.warn(`context cleaner snapshot preparation failed (ignored): ${String(error)}`); + manualCleanerExecutionDeferred = manualCleanerSchedulePending; + logger.warn( + `context cleaner snapshot preparation failed (${manualCleanerExecutionDeferred + ? "scheduled clean deferred" + : "ignored"}): ${String(error)}`, + ); } - if (cleanerSnapshot && cleanerRegistry) { + if (cleanerSnapshot && cleanerRegistry && !manualCleanerExecutionDeferred) { const manual = await prepareClaudeCleanerOverlay({ stateDir: config.stateDir, sessionId, @@ -613,7 +678,12 @@ export async function startClaudeCodeGatewayRuntime(params: { throw error; } } else if (manual.outcome === "reserved") { - logger.warn(`context cleaner manual overlay deferred (ignored): ${manual.reasonCodes.join(",")}`); + manualCleanerExecutionDeferred = manualCleanerSchedulePending; + logger.warn( + `context cleaner manual overlay deferred (${manualCleanerExecutionDeferred + ? "scheduled clean deferred" + : "ignored"}): ${manual.reasonCodes.join(",")}`, + ); } } else if (!manualCleanerSuppressesAutomaticEviction) { const schedule = await readClaudeCleanerSchedule({ stateDir: config.stateDir, sessionId }); @@ -636,9 +706,11 @@ export async function startClaudeCodeGatewayRuntime(params: { } }; - // A prepared manual overlay keeps the approval-time snapshot as its retry - // anchor until the upstream accepts and the applied receipt commits. - if (!manualCleanerOverlay && !cleanerApplyControlRequest) { + // A prepared manual overlay or any deferred prerequisite keeps the + // approval-time snapshot as its retry anchor until execution can finish. + if (!manualCleanerOverlay + && !cleanerApplyControlRequest + && !manualCleanerExecutionDeferred) { await persistCleanerSnapshot(); } @@ -849,6 +921,18 @@ export async function startClaudeCodeGatewayRuntime(params: { const authorization = typeof req.headers.authorization === "string" ? req.headers.authorization : undefined; const model = envelope.model; const workspaceHint = extractWorkspaceHint(envelope); + // A deferred manual clean owns this context boundary. Disable the other + // context transformations so its frozen request can be retried unchanged. + const observedPreparationConfig = manualCleanerExecutionDeferred + ? { + ...config, + modules: { + ...config.modules, + stabilizer: false, + reduction: false, + }, + } + : config; let prepared: Awaited>>; try { prepared = await prepareObservedBeforeCall({ @@ -856,13 +940,13 @@ export async function startClaudeCodeGatewayRuntime(params: { codec, config: { mode: "normal" }, prepareStablePrefix(nextEnvelope) { - return prepareClaudeStablePrefix(nextEnvelope, config); + return prepareClaudeStablePrefix(nextEnvelope, observedPreparationConfig); }, async applyBeforeCallReduction({ envelope: nextEnvelope, codec: nextCodec }) { return reduceClaudeRequestEnvelope({ envelope: nextEnvelope, codec: nextCodec, - config, + config: observedPreparationConfig, }); }, observability: { diff --git a/components/adapters/claude-code/tests/context-cleaner-gateway.test.ts b/components/adapters/claude-code/tests/context-cleaner-gateway.test.ts index 29f79515..f52d0e03 100644 --- a/components/adapters/claude-code/tests/context-cleaner-gateway.test.ts +++ b/components/adapters/claude-code/tests/context-cleaner-gateway.test.ts @@ -1,5 +1,5 @@ import assert from "node:assert/strict"; -import { mkdtemp, rm } from "node:fs/promises"; +import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { createServer } from "node:net"; import { tmpdir } from "node:os"; import { join } from "node:path"; @@ -14,19 +14,29 @@ import { type ContextCleanPlan, } from "@lightrsi/cleaner"; import type { HostGatewayForwarder } from "@lightrsi/host-adapter"; -import type { SessionTaskRegistry } from "@lightrsi/history"; +import { + persistRawSemanticTurnRecord, + persistSessionTaskRegistry, + rawSemanticTurnRecordPath, + sessionTaskRegistryPath, + type SessionTaskRegistry, +} from "@lightrsi/history"; import { attributeClaudeSnapshotTasks } from "../src/context-cleaner/snapshot.js"; -import { scheduleClaudeCleanerPlan } from "../src/context-cleaner/scheduler.js"; +import { + acquireClaudeCleanerScheduleLock, + scheduleClaudeCleanerPlan, +} from "../src/context-cleaner/scheduler.js"; import { normalizeTokenPilotClaudeCodeConfig } from "../src/config.js"; import { startClaudeCodeGatewayRuntime } from "../src/gateway-runtime.js"; import { createConsoleLogger } from "../src/logger.js"; +import { buildRawSemanticTurnRecord } from "../src/context-rewrite/semantic-mapping.js"; +import { claudeContextRewriteBackend } from "../src/context-rewrite/backend.js"; import { buildClaudeContextSnapshot } from "../src/context-rewrite/snapshot.js"; import { readLatestClaudeSnapshotRecord, saveLatestClaudeSnapshot, } from "../src/context-rewrite/snapshot-store.js"; -import { persistSessionTaskRegistry } from "@lightrsi/history"; const SESSION = "claude-cleaner-gateway-session"; const PLAN = "claude-cleaner-gateway-plan"; @@ -48,6 +58,12 @@ async function reserveUnusedPort(): Promise { }); } +function withoutCodecDefaults(value: unknown): unknown { + return JSON.parse(JSON.stringify(value), (key, entry) => ( + key === "is_error" && entry === false ? undefined : entry + )); +} + function registry(): SessionTaskRegistry { return { sessionId: SESSION, @@ -169,7 +185,7 @@ test("slash clean apply control request preserves the approval-time snapshot", a } }); -test("scheduled Claude clean retries after upstream rejection and commits only an accepted overlay", async () => { +test("scheduled Claude clean defers unavailable prerequisites and commits only an accepted overlay", async () => { const root = await mkdtemp(join(tmpdir(), "lightrsi-claude-cleaner-gateway-")); const stateDir = join(root, "state"); const proxyPort = await reserveUnusedPort(); @@ -242,11 +258,13 @@ test("scheduled Claude clean retries after upstream rejection and commits only a }; const forwarded: Array> = []; + let rejectNext = false; const forwarder: HostGatewayForwarder = { async requestRaw() { throw new Error("requestRaw not used"); }, async request(params) { forwarded.push(params.payload as Record); - if (forwarded.length === 1) { + if (rejectNext) { + rejectNext = false; return { status: 503, headers: { "content-type": "application/json" }, @@ -268,11 +286,23 @@ test("scheduled Claude clean retries after upstream rejection and commits only a async requestStream() { throw new Error("stream not used"); }, }; let resolverCalls = 0; + let failSnapshotRead = false; const runtime = await startClaudeCodeGatewayRuntime({ config: normalizeTokenPilotClaudeCodeConfig({ stateDir, proxyPort, - modules: { stabilizer: false, reduction: false, eviction: true }, + modules: { stabilizer: true, reduction: true, eviction: true }, + reduction: { + triggerMinChars: 256, + maxToolChars: 300, + passes: { + readStateCompaction: false, + toolPayloadTrim: true, + htmlSlimming: false, + execOutputTruncation: true, + agentsStartupOptimization: false, + }, + }, eviction: { enabled: true, minBlockChars: 1 }, taskStateEstimator: { enabled: true, batchTurns: 1 }, }), @@ -283,11 +313,28 @@ test("scheduled Claude clean retries after upstream rejection and commits only a resolverCalls += 1; return undefined; }, + async readSnapshot( + params: Parameters[0], + ) { + if (failSnapshotRead) throw new Error("simulated current snapshot failure"); + return claudeContextRewriteBackend.readSnapshot(params); + }, }, }); try { - await persistSessionTaskRegistry(stateDir, registry(), { expectedVersion: 0 }); + const runtimeRegistry = registry(); + runtimeRegistry.blockToTaskIds = {}; + runtimeRegistry.turnToTaskIds = { + [`${SESSION}:t1`]: ["task-completed"], + }; + await persistSessionTaskRegistry(stateDir, runtimeRegistry, { expectedVersion: 0 }); + const rawTurn = buildRawSemanticTurnRecord({ + sessionId: SESSION, + turnSeq: 1, + messages: historicalMessages.slice(0, 2), + }); + await persistRawSemanticTurnRecord(stateDir, rawTurn); assert.deepEqual(await saveLatestClaudeSnapshot(stateDir, SESSION, baseSnapshot), { saved: true }); assert.equal((await saveContextCleanPlan({ stateDir, plan })).outcome, "stored"); const pending: Omit = { @@ -325,6 +372,129 @@ test("scheduled Claude clean retries after upstream rejection and commits only a ], max_tokens: 128, }); + const originalMessages = (JSON.parse(requestBody) as { messages: Array> }).messages; + + await rm(join(stateDir, "claude-context", "sessions"), { recursive: true, force: true }); + const missingApprovalSnapshot = await fetch(`${runtime.baseUrl}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", "x-session-id": SESSION }, + body: requestBody, + }); + assert.equal(missingApprovalSnapshot.status, 200); + assert.equal((await readContextCleanReceipt({ stateDir, planId: PLAN })).value?.status, "scheduled"); + assert.equal(forwarded.length, 1); + assert.deepEqual(withoutCodecDefaults(forwarded[0]!.messages), originalMessages); + + assert.deepEqual(await saveLatestClaudeSnapshot(stateDir, SESSION, baseSnapshot), { saved: true }); + await writeFile(rawSemanticTurnRecordPath(stateDir, SESSION, 1), "{not-json", "utf8"); + const deferred = await fetch(`${runtime.baseUrl}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", "x-session-id": SESSION }, + body: requestBody, + }); + assert.equal(deferred.status, 200); + assert.equal((await readContextCleanReceipt({ stateDir, planId: PLAN })).value?.status, "scheduled"); + assert.equal( + (await readLatestClaudeSnapshotRecord(stateDir, SESSION))?.snapshot.revision, + REVISION, + ); + assert.equal(forwarded.length, 2); + const unchangedMessages = forwarded[1]!.messages as Array>; + assert.deepEqual(withoutCodecDefaults(unchangedMessages), originalMessages); + + await rm(rawSemanticTurnRecordPath(stateDir, SESSION, 1), { force: true }); + const missingRawTurnDeferred = await fetch(`${runtime.baseUrl}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", "x-session-id": SESSION }, + body: requestBody, + }); + assert.equal(missingRawTurnDeferred.status, 200); + assert.equal((await readContextCleanReceipt({ stateDir, planId: PLAN })).value?.status, "scheduled"); + assert.equal( + (await readLatestClaudeSnapshotRecord(stateDir, SESSION))?.snapshot.revision, + REVISION, + ); + assert.equal(forwarded.length, 3); + assert.deepEqual(withoutCodecDefaults(forwarded[2]!.messages), originalMessages); + + await persistRawSemanticTurnRecord(stateDir, rawTurn); + await writeFile(sessionTaskRegistryPath(stateDir, SESSION), "{not-json", "utf8"); + const registryDeferred = await fetch(`${runtime.baseUrl}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", "x-session-id": SESSION }, + body: requestBody, + }); + assert.equal(registryDeferred.status, 200); + assert.equal((await readContextCleanReceipt({ stateDir, planId: PLAN })).value?.status, "scheduled"); + assert.equal( + (await readLatestClaudeSnapshotRecord(stateDir, SESSION))?.snapshot.revision, + REVISION, + ); + assert.equal(forwarded.length, 4); + assert.deepEqual(withoutCodecDefaults(forwarded[3]!.messages), originalMessages); + + await rm(sessionTaskRegistryPath(stateDir, SESSION), { force: true }); + const missingRegistryDeferred = await fetch(`${runtime.baseUrl}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", "x-session-id": SESSION }, + body: requestBody, + }); + assert.equal(missingRegistryDeferred.status, 200); + assert.equal((await readContextCleanReceipt({ stateDir, planId: PLAN })).value?.status, "scheduled"); + assert.equal( + (await readLatestClaudeSnapshotRecord(stateDir, SESSION))?.snapshot.revision, + REVISION, + ); + assert.equal(forwarded.length, 5); + assert.deepEqual(withoutCodecDefaults(forwarded[4]!.messages), originalMessages); + + await persistSessionTaskRegistry(stateDir, runtimeRegistry); + failSnapshotRead = true; + const snapshotDeferred = await fetch(`${runtime.baseUrl}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", "x-session-id": SESSION }, + body: requestBody, + }); + failSnapshotRead = false; + assert.equal(snapshotDeferred.status, 200); + assert.equal((await readContextCleanReceipt({ stateDir, planId: PLAN })).value?.status, "scheduled"); + assert.equal( + (await readLatestClaudeSnapshotRecord(stateDir, SESSION))?.snapshot.revision, + REVISION, + ); + assert.equal(forwarded.length, 6); + assert.deepEqual(withoutCodecDefaults(forwarded[5]!.messages), originalMessages); + + const heldCleanerLock = await acquireClaudeCleanerScheduleLock({ stateDir, sessionId: SESSION }); + assert.ok(heldCleanerLock); + let lockDeferred: Response; + try { + lockDeferred = await fetch(`${runtime.baseUrl}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", "x-session-id": SESSION }, + body: requestBody, + }); + } finally { + await heldCleanerLock.release(); + } + assert.equal(lockDeferred.status, 200); + assert.equal((await readContextCleanReceipt({ stateDir, planId: PLAN })).value?.status, "scheduled"); + assert.equal( + (await readLatestClaudeSnapshotRecord(stateDir, SESSION))?.snapshot.revision, + REVISION, + ); + assert.equal(forwarded.length, 7); + assert.deepEqual(withoutCodecDefaults(forwarded[6]!.messages), originalMessages); + const deferredTrace = (await readFile(join(stateDir, "event-trace.jsonl"), "utf8")) + .trim() + .split("\n") + .map((line) => JSON.parse(line) as Record) + .filter((entry) => entry.stage === "gateway_before_call") + .at(-1); + assert.equal(deferredTrace?.stablePrefixApplied, false); + assert.equal(deferredTrace?.reductionApplied, false); + + rejectNext = true; const rejected = await fetch(`${runtime.baseUrl}/v1/messages`, { method: "POST", headers: { "content-type": "application/json", "x-session-id": SESSION }, @@ -341,22 +511,24 @@ test("scheduled Claude clean retries after upstream rejection and commits only a assert.equal(response.status, 200); assert.equal(resolverCalls, 0); - assert.equal(forwarded.length, 2); - const messages = forwarded[1]!.messages as Array>; - assert.deepEqual((messages[0]!.content as Array>)[0], { - type: "tool_use", - id: "toolu_cleaner_gateway", - name: "Read", - input: {}, - }); - assert.match( - String((messages[1]!.content as Array>)[0]!.content), - /^\[(Tool payload trimmed|evicted: earlier tool result)/, - ); - assert.equal( - (messages.at(-1)!.content as Array>)[0]!.text, - "KEEP_CURRENT_REQUEST", - ); + assert.equal(forwarded.length, 9); + for (const forwardedIndex of [7, 8]) { + const messages = forwarded[forwardedIndex]!.messages as Array>; + assert.deepEqual((messages[0]!.content as Array>)[0], { + type: "tool_use", + id: "toolu_cleaner_gateway", + name: "Read", + input: {}, + }); + assert.match( + String((messages[1]!.content as Array>)[0]!.content), + /^\[(Tool payload trimmed|evicted: earlier tool result)/, + ); + assert.equal( + (messages.at(-1)!.content as Array>)[0]!.text, + "KEEP_CURRENT_REQUEST", + ); + } const receipt = await readContextCleanReceipt({ stateDir, planId: PLAN }); assert.equal(receipt.value?.status, "applied"); assert.deepEqual(receipt.value?.evidence?.itemIds, approvedItems.map((item) => item.stableId)); diff --git a/components/adapters/claude-code/tests/context-cleaner-snapshot.test.ts b/components/adapters/claude-code/tests/context-cleaner-snapshot.test.ts index 48cf5b5a..d82781d6 100644 --- a/components/adapters/claude-code/tests/context-cleaner-snapshot.test.ts +++ b/components/adapters/claude-code/tests/context-cleaner-snapshot.test.ts @@ -3,7 +3,10 @@ import test from "node:test"; import type { SessionTaskRegistry } from "@lightrsi/history"; -import { attributeClaudeSnapshotTasks } from "../src/context-cleaner/snapshot.js"; +import { + attributeClaudeSnapshotTasks, + scheduledCleanerAttributionUnavailable, +} from "../src/context-cleaner/snapshot.js"; import { buildClaudeContextSnapshot } from "../src/context-rewrite/snapshot.js"; const SESSION = "claude-task-attribution"; @@ -159,3 +162,121 @@ test("rejects attribution from a registry for another session", () => { }); assert.ok(attributed.items.every((item) => item.taskIds === undefined)); }); + +test("treats explicit task reattribution as stale context rather than unavailable evidence", () => { + const inbound = messages(); + const baseSnapshot = buildClaudeContextSnapshot({ + sessionId: SESSION, + revision: "revision-reattributed", + messages: inbound as any, + }); + const approvalSnapshot = attributeClaudeSnapshotTasks({ + snapshot: baseSnapshot, + messages: inbound, + registry: registry({ "anthropic-tool-result:toolu_read": ["task-read"] }), + }); + const currentSnapshot = attributeClaudeSnapshotTasks({ + snapshot: baseSnapshot, + messages: inbound, + registry: registry( + { "anthropic-tool-result:toolu_read": ["task-other"] }, + ["task-read", "task-other"], + ), + }); + + assert.equal(scheduledCleanerAttributionUnavailable({ + selectedTaskIds: ["task-read"], + approvalSnapshot, + currentSnapshot, + }), false); +}); + +test("lets deterministic pair drift override unavailable attribution on another pair", () => { + const approvalMessages = [ + ...messages().slice(0, -1), + { + role: "assistant", + content: [{ type: "tool_use", id: "toolu_other", name: "Read", input: {} }], + }, + { + role: "user", + content: [{ type: "tool_result", tool_use_id: "toolu_other", content: "other" }], + }, + { role: "user", content: [{ type: "text", text: "current request" }] }, + ]; + const approvalSnapshot = attributeClaudeSnapshotTasks({ + snapshot: buildClaudeContextSnapshot({ + sessionId: SESSION, + revision: "revision-multi-pair", + messages: approvalMessages as any, + }), + messages: approvalMessages, + registry: registry({ + "anthropic-tool-result:toolu_read": ["task-read"], + "anthropic-tool-result:toolu_other": ["task-read"], + }), + }); + const currentMessages = approvalMessages.filter((message) => { + const block = Array.isArray(message.content) + ? message.content[0] as Record | undefined + : undefined; + return block?.id !== "toolu_read" && block?.tool_use_id !== "toolu_read"; + }); + const currentSnapshot = buildClaudeContextSnapshot({ + sessionId: SESSION, + revision: "revision-multi-pair-current", + messages: currentMessages as any, + }); + + assert.equal(scheduledCleanerAttributionUnavailable({ + selectedTaskIds: ["task-read"], + approvalSnapshot, + currentSnapshot, + }), false); +}); + +test("lets fingerprint drift override unavailable attribution on another pair", () => { + const approvalMessages = [ + ...messages().slice(0, -1), + { + role: "assistant", + content: [{ type: "tool_use", id: "toolu_other", name: "Read", input: {} }], + }, + { + role: "user", + content: [{ type: "tool_result", tool_use_id: "toolu_other", content: "other" }], + }, + { role: "user", content: [{ type: "text", text: "current request" }] }, + ]; + const approvalSnapshot = attributeClaudeSnapshotTasks({ + snapshot: buildClaudeContextSnapshot({ + sessionId: SESSION, + revision: "revision-fingerprint-drift", + messages: approvalMessages as any, + }), + messages: approvalMessages, + registry: registry({ + "anthropic-tool-result:toolu_read": ["task-read"], + "anthropic-tool-result:toolu_other": ["task-read"], + }), + }); + const currentMessages = structuredClone(approvalMessages) as Array<{ + content: Array>; + }>; + const changedResult = currentMessages + .flatMap((message) => message.content) + .find((block) => block.tool_use_id === "toolu_other"); + assert.ok(changedResult); + changedResult.content = "changed other result"; + const currentSnapshot = buildClaudeContextSnapshot({ + sessionId: SESSION, + revision: "revision-fingerprint-drift-current", + messages: currentMessages as any, + }); + + assert.equal(scheduledCleanerAttributionUnavailable({ + selectedTaskIds: ["task-read"], + approvalSnapshot, + currentSnapshot, + }), false); +}); diff --git a/components/adapters/claude-code/tests/context-rewrite-semantic-pipeline.test.ts b/components/adapters/claude-code/tests/context-rewrite-semantic-pipeline.test.ts index eabcbe34..b0991918 100644 --- a/components/adapters/claude-code/tests/context-rewrite-semantic-pipeline.test.ts +++ b/components/adapters/claude-code/tests/context-rewrite-semantic-pipeline.test.ts @@ -6,6 +6,7 @@ import { join } from "node:path"; import { loadSessionTaskRegistry, loadRawSemanticTurnRecord, + persistRawSemanticTurnRecord, persistSessionTaskRegistry, type SessionTaskRegistry, type DeltaView, @@ -13,8 +14,10 @@ import { import type { TaskStateEstimator } from "@lightrsi/eviction"; import { buildUniqueToolCallTurnMap, + loadPersistedToolCallTurnMap, runSemanticPipeline, } from "../src/context-rewrite/semantic-pipeline.js"; +import { buildRawSemanticTurnRecord } from "../src/context-rewrite/semantic-mapping.js"; import { updateRegistryFromDelta as realUpdateRegistryFromDelta } from "../src/context-rewrite/task-registry-update.js"; async function tempStateDir(): Promise { @@ -290,6 +293,31 @@ test("tool call attribution fails closed when a call id appears in multiple turn ); }); +test("loads the persisted tool-call turn map without claiming a new turn", async () => { + const stateDir = await tempStateDir(); + const sessionId = "sess-persisted-map"; + const record = buildRawSemanticTurnRecord({ + sessionId, + turnSeq: 3, + messages: [ + { + role: "assistant", + content: [{ type: "tool_use", id: "toolu_persisted", name: "Read", input: {} }], + }, + { + role: "user", + content: [{ type: "tool_result", tool_use_id: "toolu_persisted", content: "done" }], + }, + ], + }); + await persistRawSemanticTurnRecord(stateDir, record); + + const result = await loadPersistedToolCallTurnMap({ stateDir, sessionId }); + + assert.equal(result.get("toolu_persisted"), `${sessionId}:t3`); + assert.equal(await loadRawSemanticTurnRecord(stateDir, sessionId, 4), null); +}); + test("internal Claude metadata requests do not advance or invoke the semantic estimator", async () => { const stateDir = await tempStateDir(); const sessionId = "sess-internal-request";