Skip to content

Commit ba1d031

Browse files
committed
Fix evidence archive blob keys for the context store
Context store blob keys cannot contain slashes, so archive occurrences failed to persist. Hash overflow blobs after scrub and admit string send the same way as inbound messages.
1 parent 6db955d commit ba1d031

5 files changed

Lines changed: 384 additions & 20 deletions

File tree

‎src/plugins/result-truncation-plugin.ts‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -259,10 +259,11 @@ export function resultTruncationPlugin(options: ResultTruncationPluginOptions =
259259
const truncated = await truncateToolResultContent(content, MAX_RESULT_CHARS, spill);
260260
if (truncated === content) return result;
261261
if (archive !== undefined && spill !== undefined && content.length > MAX_RESULT_CHARS) {
262+
const spilled = scrubSecretShapedContent(materializeToolResultContent(content).text);
262263
await archive.recordExistingBlobReference({
263264
kind: "overflow_blob",
264265
blobKey: spillBlobKey(call.id),
265-
contentHash: hashAuthorizedBytes(new TextEncoder().encode(content)),
266+
contentHash: hashAuthorizedBytes(new TextEncoder().encode(spilled)),
266267
callId: call.id,
267268
provenance: "result-truncation:full",
268269
});
@@ -275,6 +276,16 @@ export function resultTruncationPlugin(options: ResultTruncationPluginOptions =
275276
const compact = JSON.stringify(record);
276277
if (compact.length <= MAX_RESULT_CHARS) return result;
277278
const truncated = await truncateToolResultRecord(record, MAX_RESULT_CHARS, spill);
279+
if (archive !== undefined && spill !== undefined) {
280+
const spilled = scrubSecretShapedContent(materializeToolResultRecord(record).text);
281+
await archive.recordExistingBlobReference({
282+
kind: "overflow_blob",
283+
blobKey: spillBlobKey(call.id),
284+
contentHash: hashAuthorizedBytes(new TextEncoder().encode(spilled)),
285+
callId: call.id,
286+
provenance: "result-truncation:full",
287+
});
288+
}
278289
return { ...result, content: truncated };
279290
}
280291

‎src/session/assemble-runtime.ts‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@ import {
5555
createCompactionArchive,
5656
createPrimaryDeliveryAdmission,
5757
wrapAuthorizeWithEvidenceArchive,
58+
wrapCompactorWithCompletenessGate,
5859
type CompactionArchive,
5960
} from "./compaction-archive.js";
6061
import path from "node:path";
@@ -500,7 +501,10 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent {
500501
defaultId: `${ID_PREFIX}/chat`,
501502
}),
502503
compactors: {
503-
"pruning-compactor": wiring.getCompactor(),
504+
"pruning-compactor":
505+
primaryArchive === undefined
506+
? wiring.getCompactor()
507+
: wrapCompactorWithCompletenessGate(wiring.getCompactor(), primaryArchive),
504508
},
505509
});
506510
const admittedAgent =

‎src/session/attachment-store.ts‎

Lines changed: 3 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ import {
1212
formatAgedImageMarker,
1313
parseAgedImageMarker,
1414
} from "./attachment-uri.js";
15-
import { hashAuthorizedBytes, type CompactionArchive } from "./compaction-archive.js";
15+
import { type CompactionArchive } from "./compaction-archive.js";
1616

1717
export interface AgeImageResult {
1818
turn: ConversationTurn;
@@ -71,14 +71,8 @@ export async function ageImageBlocks(
7171
bytes,
7272
contentType: block.source.mimeType,
7373
});
74-
if (options.archive !== undefined) {
75-
await options.archive.recordExistingBlobReference({
76-
kind: "attachment",
77-
blobKey: id,
78-
contentHash: hashAuthorizedBytes(bytes),
79-
provenance: "attachment-age:base64",
80-
});
81-
}
74+
// recordExistingBlobReference before persistBlobs marks a false gap.
75+
// Callers that already wrote the blob may record after this returns.
8276
content.push({
8377
type: "text",
8478
text: formatAgedImageMarker({ uri, mimeType: block.source.mimeType }),

‎src/session/compaction-archive.test.ts‎

Lines changed: 233 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import {
1515
hashAuthorizedBytes,
1616
isControlOrEmptyInbound,
1717
} from "./compaction-archive.js";
18+
import { createOptimizedContextStore } from "./optimized-context-store.js";
1819

1920
function tempDir(): string {
2021
return fs.mkdtempSync(path.join(os.tmpdir(), "compaction-archive-"));
@@ -122,8 +123,12 @@ describe("primary message admission", () => {
122123
deliver(message: InboundMessage) {
123124
delivered.push(message);
124125
},
125-
async send(message: InboundMessage) {
126-
delivered.push(message);
126+
async send(content: string | InboundMessage) {
127+
if (typeof content === "string") {
128+
delivered.push(inbound({ content }));
129+
return { ok: true as const };
130+
}
131+
delivered.push(content);
127132
return { ok: true as const };
128133
},
129134
};
@@ -155,8 +160,12 @@ describe("primary message admission", () => {
155160
deliver(message: InboundMessage) {
156161
delivered.push(message);
157162
},
158-
async send(message: InboundMessage) {
159-
delivered.push(message);
163+
async send(content: string | InboundMessage) {
164+
if (typeof content === "string") {
165+
delivered.push(inbound({ content }));
166+
return { ok: true as const };
167+
}
168+
delivered.push(content);
160169
return { ok: true as const };
161170
},
162171
};
@@ -173,6 +182,43 @@ describe("primary message admission", () => {
173182
const archived = await archive.readAuthorizedPayload(occurrences[0]!.occurrenceId);
174183
expect(archived).toBe(admitted);
175184
});
185+
186+
test("send(string) admits and archives like InboundMessage", async () => {
187+
const dir = tempDir();
188+
const blobs = new Map<string, Uint8Array>();
189+
const archive = createCompactionArchive({
190+
sessionId: "sess-send-string",
191+
contextDir: dir,
192+
writeBlob: async (key, bytes) => {
193+
blobs.set(key, bytes);
194+
},
195+
readBlob: async (key) => {
196+
const bytes = blobs.get(key);
197+
if (bytes === undefined) throw new Error(`missing blob ${key}`);
198+
return bytes;
199+
},
200+
});
201+
const sent: string[] = [];
202+
const agent = {
203+
deliver(_message: InboundMessage) {
204+
/* unused */
205+
},
206+
async send(content: string | InboundMessage) {
207+
if (typeof content !== "string") throw new Error("expected string send");
208+
sent.push(content);
209+
return { ok: true as const };
210+
},
211+
};
212+
const wrapped = createPrimaryDeliveryAdmission(agent, archive);
213+
const secret = `sk-${"a".repeat(24)}`;
214+
await wrapped.send(`constraint ${secret}`);
215+
expect(sent).toHaveLength(1);
216+
expect(sent[0]).toContain(CREDENTIAL_REDACTION);
217+
expect(sent[0]).not.toContain(secret);
218+
const occurrences = await archive.listOccurrences();
219+
expect(occurrences).toHaveLength(1);
220+
expect(await archive.readAuthorizedPayload(occurrences[0]!.occurrenceId)).toBe(sent[0]!);
221+
});
176222
});
177223

178224
describe("compaction archive storage", () => {
@@ -451,4 +497,187 @@ describe("compaction archive storage", () => {
451497
const cert = await archive.certifyRange([occ.occurrenceId]);
452498
expect(cert.status).toBe("complete");
453499
});
500+
501+
test("recordAuthorizedPayload writes store-legal keys through createOptimizedContextStore", async () => {
502+
const dir = tempDir();
503+
const store = await createOptimizedContextStore(dir);
504+
const archive = createCompactionArchive({
505+
sessionId: "sess-store-keys",
506+
contextDir: dir,
507+
writeBlob: (key, bytes, contentType) => store.writeBlob(key, bytes, contentType),
508+
readBlob: (key) => store.readBlob(key),
509+
});
510+
const occ = await archive.recordAuthorizedPayload({
511+
kind: "user_message",
512+
payload: "hello-store",
513+
});
514+
expect(occ.blobKey.includes("/")).toBe(false);
515+
expect(occ.blobKey.includes("..")).toBe(false);
516+
expect(await archive.readAuthorizedPayload(occ.occurrenceId)).toBe("hello-store");
517+
});
518+
});
519+
520+
describe("wrapCompactorWithCompletenessGate", () => {
521+
const ctx = { trigger: "test" } as unknown as import("@intx/types/runtime").StrategyContext;
522+
523+
function truncating(name: string): import("@intx/types/runtime").Compactor {
524+
return {
525+
name,
526+
version: "1",
527+
async apply(turns) {
528+
return {
529+
output: turns.slice(-1),
530+
blobs: [
531+
{
532+
key: "stats",
533+
bytes: new TextEncoder().encode("{}"),
534+
contentType: "application/json",
535+
},
536+
],
537+
record: {
538+
strategy: name,
539+
version: "1",
540+
parameters: {},
541+
reason: "compact",
542+
decisions: { dropped: turns.length - 1 },
543+
},
544+
};
545+
},
546+
};
547+
}
548+
549+
function memoryArchive() {
550+
const dir = tempDir();
551+
const blobs = new Map<string, Uint8Array>();
552+
const archive = createCompactionArchive({
553+
sessionId: "sess-gate",
554+
contextDir: dir,
555+
writeBlob: async (key, bytes) => {
556+
blobs.set(key, bytes);
557+
},
558+
readBlob: async (key) => {
559+
const bytes = blobs.get(key);
560+
if (bytes === undefined) throw new Error(`missing ${key}`);
561+
return bytes;
562+
},
563+
});
564+
return { archive, blobs };
565+
}
566+
567+
test("incomplete archive returns identity history and drops stats blobs", async () => {
568+
const { wrapCompactorWithCompletenessGate } = await import("./compaction-archive.js");
569+
const { archive } = memoryArchive();
570+
const inner = truncating("pruning-compactor");
571+
const wrapped = wrapCompactorWithCompletenessGate(inner, archive);
572+
const turns: import("@intx/types/runtime").ConversationTurn[] = [
573+
{
574+
role: "user",
575+
content: [{ type: "text", text: "secret-fact" }],
576+
timestamp: 1,
577+
},
578+
{
579+
role: "assistant",
580+
content: [
581+
{ type: "tool_call", id: "call-drop", name: "read_file", arguments: { path: "a.ts" } },
582+
],
583+
timestamp: 2,
584+
},
585+
{
586+
role: "user",
587+
content: [
588+
{ type: "tool_result", callId: "call-drop", content: [{ type: "text", text: "ok" }] },
589+
],
590+
timestamp: 3,
591+
},
592+
];
593+
594+
const result = await wrapped.apply(turns, ctx);
595+
expect(result.output).toBe(turns);
596+
expect(result.blobs).toBeUndefined();
597+
expect(result.record.reason).toBe("incomplete-evidence-archive");
598+
});
599+
600+
test("complete archive covering dropped callIds allows the rewrite", async () => {
601+
const { wrapCompactorWithCompletenessGate } = await import("./compaction-archive.js");
602+
const { archive } = memoryArchive();
603+
await archive.recordAuthorizedPayload({
604+
kind: "tool_args",
605+
payload: { name: "read_file", arguments: { path: "a.ts" } },
606+
callId: "call-drop",
607+
});
608+
await archive.recordAuthorizedPayload({
609+
kind: "tool_result",
610+
payload: "ok",
611+
callId: "call-drop",
612+
});
613+
await archive.recordAuthorizedPayload({
614+
kind: "user_message",
615+
payload: "secret-fact",
616+
});
617+
618+
const inner = truncating("pruning-compactor");
619+
const wrapped = wrapCompactorWithCompletenessGate(inner, archive);
620+
const turns: import("@intx/types/runtime").ConversationTurn[] = [
621+
{
622+
role: "user",
623+
content: [{ type: "text", text: "secret-fact" }],
624+
timestamp: 1,
625+
},
626+
{
627+
role: "assistant",
628+
content: [
629+
{ type: "tool_call", id: "call-drop", name: "read_file", arguments: { path: "a.ts" } },
630+
],
631+
timestamp: 2,
632+
},
633+
{
634+
role: "user",
635+
content: [
636+
{ type: "tool_result", callId: "call-drop", content: [{ type: "text", text: "ok" }] },
637+
],
638+
timestamp: 3,
639+
},
640+
];
641+
642+
const result = await wrapped.apply(turns, ctx);
643+
expect(result.output).toHaveLength(1);
644+
expect(result.blobs?.some((b) => b.key === "stats")).toBe(true);
645+
expect(result.record.reason).toBe("compact");
646+
});
647+
648+
test("explicit gap records are not required for completeness", async () => {
649+
const { wrapCompactorWithCompletenessGate } = await import("./compaction-archive.js");
650+
const { archive } = memoryArchive();
651+
await archive.importHistoricalEvidence([
652+
{ kind: "user_message", available: false, callId: "historical-gap" },
653+
]);
654+
await archive.recordAuthorizedPayload({
655+
kind: "tool_result",
656+
payload: "ok",
657+
callId: "call-drop",
658+
});
659+
await archive.recordAuthorizedPayload({
660+
kind: "user_message",
661+
payload: "secret-fact",
662+
});
663+
664+
const wrapped = wrapCompactorWithCompletenessGate(truncating("pruning-compactor"), archive);
665+
const turns: import("@intx/types/runtime").ConversationTurn[] = [
666+
{
667+
role: "user",
668+
content: [{ type: "text", text: "secret-fact" }],
669+
timestamp: 1,
670+
},
671+
{
672+
role: "user",
673+
content: [
674+
{ type: "tool_result", callId: "call-drop", content: [{ type: "text", text: "ok" }] },
675+
],
676+
timestamp: 2,
677+
},
678+
];
679+
const result = await wrapped.apply(turns, ctx);
680+
expect(result.output).toHaveLength(1);
681+
expect(result.record.reason).toBe("compact");
682+
});
454683
});

0 commit comments

Comments
 (0)