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
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,9 @@ function validSessionState(value: unknown, sessionId: string): value is {
&& (value.disclosedReadPaths === undefined
|| (Array.isArray(value.disclosedReadPaths)
&& value.disclosedReadPaths.every((path) => typeof path === "string")))
&& (value.disclosedReadOwners === undefined
|| (isRecord(value.disclosedReadOwners)
&& Object.values(value.disclosedReadOwners).every((owner) => owner === null || typeof owner === "string")))
&& optionalCount(value.lastToolInputChars)
&& optionalCount(value.lastToolOutputChars)
&& optionalCount(value.requestChars)
Expand Down
4 changes: 4 additions & 0 deletions components/adapters/claude-code/src/gateway-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -232,6 +232,7 @@ async function recordClaudeGatewayTurn(params: {
responseId?: string;
previousResponseId?: string;
disclosedReadPaths?: string[];
disclosedReadOwners?: Record<string, string | null>;
requestChars: number;
responseChars: number;
assistantChars: number;
Expand All @@ -249,6 +250,7 @@ async function recordClaudeGatewayTurn(params: {
latestModel: params.model,
workspaceHint: params.workspaceHint,
disclosedReadPaths: params.disclosedReadPaths,
disclosedReadOwners: params.disclosedReadOwners,
requestChars: params.requestChars,
responseChars: params.responseChars,
assistantChars: params.assistantChars,
Expand Down Expand Up @@ -1066,6 +1068,7 @@ export async function startClaudeCodeGatewayRuntime(params: {
responseId,
previousResponseId,
disclosedReadPaths: reductionSummary?.disclosedReadPaths,
disclosedReadOwners: reductionSummary?.disclosedReadOwners,
requestChars: body.length,
responseChars: rawStreamText.length,
assistantChars: snapshot.assistantText.length,
Expand Down Expand Up @@ -1158,6 +1161,7 @@ export async function startClaudeCodeGatewayRuntime(params: {
responseId,
previousResponseId,
disclosedReadPaths: reductionSummary?.disclosedReadPaths,
disclosedReadOwners: reductionSummary?.disclosedReadOwners,
requestChars: body.length,
responseChars: upstreamResp.text.length,
assistantChars,
Expand Down
108 changes: 106 additions & 2 deletions components/adapters/claude-code/src/reduction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ type ClaudeSegmentBinding = {
blockIndex?: number;
field: "content" | "text";
toolName?: string;
toolUseId?: string;
};

type ClaudeReductionInstruction = {
Expand Down Expand Up @@ -80,9 +81,13 @@ export type ClaudeReductionSummary = {
diagnostics: ClaudeReductionDiagnostics;
visualSegments?: ClaudeReductionVisualSegment[];
disclosedReadPaths?: string[];
/** Normalized read path -> tool_use_id of the read that first disclosed it (null: not attributable). */
disclosedReadOwners?: DisclosedReadOwners;
skippedReason?: string;
};

export type DisclosedReadOwners = Record<string, string | null>;

function normalizeDisclosedReadPaths(value: unknown): string[] | undefined {
if (!Array.isArray(value)) return undefined;
const next = new Set<string>();
Expand All @@ -94,6 +99,83 @@ function normalizeDisclosedReadPaths(value: unknown): string[] | undefined {
return next.size > 0 ? [...next] : undefined;
}

export function normalizeDisclosedReadOwners(value: unknown): DisclosedReadOwners | undefined {
if (!value || typeof value !== "object" || Array.isArray(value)) return undefined;
const next: DisclosedReadOwners = {};
for (const [path, owner] of Object.entries(value as Record<string, unknown>)) {
const normalized = path.trim().toLowerCase();
if (!normalized) continue;
next[normalized] = typeof owner === "string" && owner ? owner : null;
}
return Object.keys(next).length > 0 ? next : undefined;
}

/**
* Paths to hand the pass as "already disclosed". Claude Code resends the whole
* history every request, so a disclosing read that is still in it is re-detected by
* the pass itself. Carrying its path too would make the pass treat that same read as
* a repeat and send it untrimmed, which breaks the prompt cache. Only paths whose
* disclosing read has left the history (compaction, a new session branch) are carried.
*/
export function carriedDisclosedReadPaths(
owners: DisclosedReadOwners | undefined,
presentToolUseIds: ReadonlySet<string>,
): string[] | undefined {
if (!owners) return undefined;
const carried = Object.entries(owners)
.filter(([, owner]) => owner === null || !presentToolUseIds.has(owner))
.map(([path]) => path);
return carried.length > 0 ? carried : undefined;
}

/**
* Rebuild owners for the reported path set; new paths belong to a matching read
* actually trimmed in this request.
*/
export function recordDisclosedReadOwners(
owners: DisclosedReadOwners | undefined,
reportedPaths: unknown,
segments: readonly ContextSegment[],
bindings: readonly ClaudeSegmentBinding[],
trimmedSegmentIds: ReadonlySet<string>,
): DisclosedReadOwners | undefined {
const paths = normalizeDisclosedReadPaths(reportedPaths);
if (!paths) return owners;
const next: DisclosedReadOwners = {};
const bindingBySegment = new Map(bindings.map((binding) => [binding.segmentId, binding]));
for (const path of paths) {
if (owners && Object.hasOwn(owners, path)) {
next[path] = owners[path] ?? null;
continue;
}
const owner = segments.find((segment) => {
if (!trimmedSegmentIds.has(segment.id)) return false;
const binding = bindingBySegment.get(segment.id);
const toolName = binding?.toolName?.trim().toLowerCase();
if (toolName !== "read" && toolName !== "file_read") return false;
const segmentPath = asRecord(segment.metadata).path;
return typeof segmentPath === "string" && segmentPath.trim().toLowerCase() === path;
});
next[path] = (owner && bindingBySegment.get(owner.id)?.toolUseId) || null;
}
return Object.keys(next).length > 0 ? next : undefined;
}

function presentToolResultIds(payload: any): Set<string> {
const ids = new Set<string>();
for (const message of Array.isArray(payload?.messages) ? payload.messages : []) {
const content = asRecord(message).content;
if (!Array.isArray(content)) continue;
for (const block of content) {
const entry = asRecord(block);
if (String(entry.type ?? "").toLowerCase() === "tool_result" && typeof entry.tool_use_id === "string" && entry.tool_use_id) {
ids.add(entry.tool_use_id);
}
}
}
return ids;
}

function asRecord(value: unknown): Record<string, unknown> {
return value && typeof value === "object" && !Array.isArray(value)
? value as Record<string, unknown>
Expand Down Expand Up @@ -312,6 +394,7 @@ function buildTurnContext(
blockIndex,
field: typeof entry.text === "string" ? "text" : "content",
toolName: typeof entry.name === "string" ? entry.name : hint?.toolName,
...(toolUseId ? { toolUseId } : {}),
});
});
});
Expand Down Expand Up @@ -552,8 +635,11 @@ export async function applyBeforeCallReductionToClaudePayload(params: {
}

const snapshot = await loadClaudeCodeSessionSnapshot(config.stateDir, sessionId);
// Snapshots written before owners were tracked only have `disclosedReadPaths`; those
// are ignored rather than carried, because carrying them reproduces the cache miss.
const disclosedReadOwners = normalizeDisclosedReadOwners(snapshot?.disclosedReadOwners);
const built = buildTurnContext(payload, sessionId, {
disclosedReadPaths: normalizeDisclosedReadPaths(snapshot?.disclosedReadPaths),
disclosedReadPaths: carriedDisclosedReadPaths(disclosedReadOwners, presentToolResultIds(payload)),
});
const analyzerInstructions = buildAnalyzerReductionInstructions(built.turnCtx.segments, config);
const fallbackInstructions = buildFallbackReductionInstructions(built.turnCtx.segments, config);
Expand Down Expand Up @@ -585,9 +671,13 @@ export async function applyBeforeCallReductionToClaudePayload(params: {
const { turnCtx: reducedCtx, report } = await runReductionBeforeCall({ turnCtx, passes });
const passEffects = summarizePassEffects(report);
const changedSegmentIds = new Set<string>();
const trimmedSegmentIds = new Set<string>();
for (const entry of report) {
if (!entry.changed) continue;
for (const id of entry.touchedSegmentIds ?? []) changedSegmentIds.add(id);
for (const id of entry.touchedSegmentIds ?? []) {
changedSegmentIds.add(id);
if (entry.id === "tool_payload_trim") trimmedSegmentIds.add(id);
}
}

if (changedSegmentIds.size === 0) {
Expand All @@ -601,6 +691,13 @@ export async function applyBeforeCallReductionToClaudePayload(params: {
passEffects,
diagnostics: built.diagnostics,
disclosedReadPaths: normalizeDisclosedReadPaths(reducedCtx.metadata?.disclosedReadPaths),
disclosedReadOwners: recordDisclosedReadOwners(
disclosedReadOwners,
reducedCtx.metadata?.disclosedReadPaths,
built.turnCtx.segments,
bindings,
trimmedSegmentIds,
),
skippedReason: "pipeline_no_effect",
};
}
Expand Down Expand Up @@ -659,6 +756,13 @@ export async function applyBeforeCallReductionToClaudePayload(params: {
diagnostics: built.diagnostics,
visualSegments,
disclosedReadPaths: normalizeDisclosedReadPaths(reducedCtx.metadata?.disclosedReadPaths),
disclosedReadOwners: recordDisclosedReadOwners(
disclosedReadOwners,
reducedCtx.metadata?.disclosedReadPaths,
built.turnCtx.segments,
bindings,
trimmedSegmentIds,
),
};
}

Expand Down
3 changes: 3 additions & 0 deletions components/adapters/claude-code/src/session-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ export type ClaudeCodeSessionSnapshot = {
latestModel?: string;
workspaceHint?: string;
disclosedReadPaths?: string[];
/** Normalized read path -> tool_use_id of the read that first disclosed it (null: not attributable). */
disclosedReadOwners?: Record<string, string | null>;
lastHookEvent?: string;
/** pid of the `claude` process that owns this session, when it could be resolved. */
hostPid?: number;
Expand Down Expand Up @@ -69,6 +71,7 @@ export async function upsertClaudeCodeSessionSnapshot(
latestModel: patch.latestModel ?? current?.latestModel,
workspaceHint: patch.workspaceHint ?? current?.workspaceHint,
disclosedReadPaths: patch.disclosedReadPaths ?? current?.disclosedReadPaths,
disclosedReadOwners: patch.disclosedReadOwners ?? current?.disclosedReadOwners,
lastHookEvent: patch.lastHookEvent ?? current?.lastHookEvent,
hostPid: patch.hostPid ?? current?.hostPid,
lastToolName: patch.lastToolName ?? current?.lastToolName,
Expand Down
121 changes: 121 additions & 0 deletions components/adapters/claude-code/tests/gateway-runtime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -963,6 +963,127 @@ export function saveConfig(file: string, text: string) {
}
});

test("gateway runtime keeps a read trimmed while that same read is still in the history", async () => {
const dir = await mkdtemp(join(tmpdir(), "lightrsi-claude-gateway-same-read-"));
const proxyPort = await reserveUnusedPort();
const codePayload = `
export function loadConfig(file: string) {
return file.trim();
}

export function saveConfig(file: string, text: string) {
return text + file;
}
`.repeat(30);
const seenPayloads: Array<Record<string, unknown>> = [];
const forwarder: HostGatewayForwarder = {
requestRaw: async () => { throw new Error("requestRaw not used in test"); },
async request(params) {
seenPayloads.push(params.payload as Record<string, unknown>);
return {
status: 200,
headers: {
"content-type": "application/json",
},
text: JSON.stringify({
id: "msg_disclosed_1",
type: "message",
role: "assistant",
content: [{ type: "text", text: "done" }],
usage: { input_tokens: 20, output_tokens: 5 },
stop_reason: "end_turn",
}),
};
},
async requestStream() {
throw new Error("stream path should not be used in this test");
},
};

const runtime = await startClaudeCodeGatewayRuntime({
config: normalizeTokenPilotClaudeCodeConfig({
stateDir: join(dir, "state"),
proxyPort,
reduction: {
triggerMinChars: 256,
maxToolChars: 300,
passes: {
readStateCompaction: false,
toolPayloadTrim: true,
htmlSlimming: false,
execOutputTruncation: false,
agentsStartupOptimization: false,
},
},
}),
logger: createConsoleLogger(false),
forwarder,
});

try {
for (let turn = 0; turn < 2; turn += 1) {
const response = await fetch(`${runtime.baseUrl}/v1/messages`, {
method: "POST",
headers: {
"content-type": "application/json",
"x-session-id": "sess-same-read-1",
},
body: JSON.stringify({
model: "claude-sonnet-4-6",
stream: false,
system: "Your working directory is: /repo/demo",
messages: [
{
role: "assistant",
content: [
{
type: "tool_use",
id: "toolu_read_0",
name: "Read",
input: { path: "/repo/src/config.ts" },
},
],
},
{
role: "user",
content: [
{ type: "text", text: "summarize this" },
{ type: "tool_result", tool_use_id: "toolu_read_0", content: codePayload },
],
},
...(turn === 0 ? [] : [
{ role: "assistant", content: [{ type: "tool_use", id: "toolu_write_1", name: "Write", input: { path: "/repo/out.txt", content: "x" } }] },
{ role: "user", content: [{ type: "tool_result", tool_use_id: "toolu_write_1", content: "wrote 1 byte" }] },
]),
],
max_tokens: 256,
}),
});
assert.equal(response.status, 200);
}

assert.equal(seenPayloads.length, 2);
const firstMessages = seenPayloads[0]?.messages as Array<Record<string, unknown>>;
const secondMessages = seenPayloads[1]?.messages as Array<Record<string, unknown>>;
const firstToolResult = ((firstMessages?.[1]?.content as Array<Record<string, unknown>>)?.[1] ?? {}) as Record<string, unknown>;
const secondToolResult = ((secondMessages?.[1]?.content as Array<Record<string, unknown>>)?.[1] ?? {}) as Record<string, unknown>;

const firstText = String(firstToolResult.content ?? firstToolResult.text ?? "");
const secondText = String(secondToolResult.content ?? secondToolResult.text ?? "");
assert.match(firstText, /\[code outlined lines=/);
assert.match(secondText, /\[code outlined lines=/, "the same read must not be treated as a repeat read");
assert.equal(secondText, firstText);

const snapshot = JSON.parse(
await readFile(join(dir, "state", "session-state", "sessions", "sess-same-read-1.json"), "utf8"),
) as { disclosedReadOwners?: Record<string, string | null> };
assert.deepEqual(snapshot.disclosedReadOwners, { "/repo/src/config.ts": "toolu_read_0" });
} finally {
await runtime.close();
await rm(dir, { recursive: true, force: true });
}
});

test("gateway runtime does not record ux-effects when reduced request fails upstream", async () => {
const dir = await mkdtemp(join(tmpdir(), "lightrsi-claude-gateway-failed-"));
const proxyPort = await reserveUnusedPort();
Expand Down
Loading
Loading