Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 47 additions & 0 deletions components/adapters/claude-code/src/context-cleaner/snapshot.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, Set<string>>();
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<string>();
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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<ReadonlyMap<string, string>> {
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,
Expand Down
116 changes: 100 additions & 16 deletions components/adapters/claude-code/src/gateway-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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;
Expand All @@ -84,6 +93,7 @@ type ClaudeCodeGatewayRuntimeDependencies = {
cloneRequestPayload?: typeof structuredClone;
resolveEstimator?: typeof resolveClaudeTaskStateEstimator;
persistTaskRegistry?: typeof persistSessionTaskRegistry;
readSnapshot?: typeof claudeContextRewriteBackend.readSnapshot;
saveSnapshot?: typeof saveLatestClaudeSnapshot;
};

Expand Down Expand Up @@ -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";
Expand All @@ -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 = [
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<ReturnType<typeof readLatestClaudeSnapshotRecord>>;
Expand All @@ -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.
Expand All @@ -559,7 +602,9 @@ export async function startClaudeCodeGatewayRuntime(params: {
| undefined;
let cleanerRegistry: Awaited<ReturnType<typeof loadSessionTaskRegistry>> | undefined;
try {
const baseCleanerSnapshot = await claudeContextRewriteBackend.readSnapshot({
const baseCleanerSnapshot = await (
params.dependencies?.readSnapshot ?? claudeContextRewriteBackend.readSnapshot
)({
sessionId,
request: {
sessionId,
Expand All @@ -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,
Expand Down Expand Up @@ -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 });
Expand All @@ -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();
}

Expand Down Expand Up @@ -849,20 +921,32 @@ 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<ReturnType<typeof prepareObservedBeforeCall<ClaudeReductionSummary>>>;
try {
prepared = await prepareObservedBeforeCall<ClaudeReductionSummary>({
envelope,
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: {
Expand Down
Loading
Loading