diff --git a/.changeset/wire-blob-offloading.md b/.changeset/wire-blob-offloading.md new file mode 100644 index 0000000000..5a86a80650 --- /dev/null +++ b/.changeset/wire-blob-offloading.md @@ -0,0 +1,6 @@ +--- +"@moonshot-ai/agent-core": patch +"@moonshot-ai/kimi-code": patch +--- + +Offload large base64 media payloads from wire.jsonl into external blob files to reduce wire size and memory pressure during session replay. Includes an in-memory read-through cache on `BlobStore` so repeated rehydration avoids redundant disk reads. diff --git a/apps/vis/server/src/app.ts b/apps/vis/server/src/app.ts index e2670b1277..e8ec4e8c02 100644 --- a/apps/vis/server/src/app.ts +++ b/apps/vis/server/src/app.ts @@ -4,6 +4,7 @@ import { join, resolve } from 'node:path'; import { Hono } from 'hono'; +import { blobsRoute } from './routes/blobs'; import { contextRoute } from './routes/context'; import { sessionDetailRoute } from './routes/session-detail'; import { sessionsRoute } from './routes/sessions'; @@ -87,6 +88,7 @@ export async function createApp(options: CreateAppOptions = {}): Promise { api.route('/sessions', sessionDetailRoute()); api.route('/sessions', wireRoute()); api.route('/sessions', subagentsRoute()); + api.route('/sessions', blobsRoute()); // Mount contextRoute last because it currently uses a catch-all stub // (Phase C scope) that would otherwise shadow more specific routes // registered below it. diff --git a/apps/vis/server/src/lib/blob-resolver.ts b/apps/vis/server/src/lib/blob-resolver.ts new file mode 100644 index 0000000000..dbb5ba6e31 --- /dev/null +++ b/apps/vis/server/src/lib/blob-resolver.ts @@ -0,0 +1,97 @@ +import type { ContentPart, WireEntry } from './agent-record-types'; + +const BLOBREF_PROTOCOL = 'blobref:'; + +function isBlobRef(url: string): boolean { + return url.startsWith(BLOBREF_PROTOCOL); +} + +/** Convert a `blobref:;` URL into a vis-server blob route. + * Non-blobref URLs are returned unchanged. */ +export function resolveBlobRefUrl( + url: string, + sessionId: string, + agentId: string, + baseUrl: string = '', +): string { + if (!isBlobRef(url)) return url; + const rest = url.slice(BLOBREF_PROTOCOL.length); + const semiIdx = rest.indexOf(';'); + if (semiIdx === -1) return url; + const mimeType = rest.slice(0, semiIdx); + const hash = rest.slice(semiIdx + 1); + if (hash.length === 0) return url; + const path = `/api/sessions/${encodeURIComponent(sessionId)}/blobs/${encodeURIComponent(hash)}?agent=${encodeURIComponent(agentId)}&mime=${encodeURIComponent(mimeType)}`; + return baseUrl ? `${baseUrl}${path}` : path; +} + +/** Walk every record in a wire and replace blobref URLs with vis-server + * blob routes so the UI can render them. Only mutates `entry.data`; + * `entry.raw` is left untouched. */ +export function rehydrateWireEntries( + entries: readonly WireEntry[], + sessionId: string, + agentId: string, + baseUrl: string = '', +): void { + for (const entry of entries) { + rehydrateRecord(entry.data as Record, sessionId, agentId, baseUrl); + } +} + +function rehydrateRecord( + record: Record, + sessionId: string, + agentId: string, + baseUrl: string, +): void { + const type = record['type']; + if (type === 'turn.prompt' || type === 'turn.steer') { + rehydrateParts(record['input'] as unknown as ContentPart[], sessionId, agentId, baseUrl); + return; + } + if (type === 'context.append_message') { + const message = record['message'] as { content: ContentPart[] }; + rehydrateParts(message.content, sessionId, agentId, baseUrl); + return; + } + if (type === 'context.append_loop_event') { + const event = record['event'] as Record; + if (event['type'] === 'tool.result') { + const result = event['result'] as Record; + if (typeof result['output'] !== 'string') { + rehydrateParts(result['output'] as ContentPart[], sessionId, agentId, baseUrl); + } + } else if (event['type'] === 'content.part') { + rehydrateParts([event['part'] as ContentPart], sessionId, agentId, baseUrl); + } + return; + } +} + +function rehydrateParts( + parts: ContentPart[], + sessionId: string, + agentId: string, + baseUrl: string, +): void { + for (const part of parts) { + switch (part.type) { + case 'image_url': + part.imageUrl.url = resolveBlobRefUrl(part.imageUrl.url, sessionId, agentId, baseUrl); + break; + case 'audio_url': + part.audioUrl.url = resolveBlobRefUrl(part.audioUrl.url, sessionId, agentId, baseUrl); + break; + case 'video_url': + part.videoUrl.url = resolveBlobRefUrl(part.videoUrl.url, sessionId, agentId, baseUrl); + break; + default: + break; + } + } +} + +export function isSafeBlobHash(hash: string): boolean { + return /^[a-f0-9]{64}$/.test(hash); +} diff --git a/apps/vis/server/src/lib/context-projector.ts b/apps/vis/server/src/lib/context-projector.ts index bdee45c97e..81569f9e3c 100644 --- a/apps/vis/server/src/lib/context-projector.ts +++ b/apps/vis/server/src/lib/context-projector.ts @@ -170,7 +170,7 @@ export function projectContext(entries: ReadonlyArray): ContextProjec }]; break; case 'usage.record': { - const scope = rec.usageScope ?? 'session'; + const scope = (rec.usageScope ?? 'session') as 'session' | 'turn'; addUsage(usage.byScope[scope], rec.usage); if (!usage.byModel[rec.model]) usage.byModel[rec.model] = { ...ZERO }; addUsage(usage.byModel[rec.model]!, rec.usage); diff --git a/apps/vis/server/src/lib/wire-reader.ts b/apps/vis/server/src/lib/wire-reader.ts index baad8ca293..b077e1773e 100644 --- a/apps/vis/server/src/lib/wire-reader.ts +++ b/apps/vis/server/src/lib/wire-reader.ts @@ -52,28 +52,28 @@ export async function readAgentWire(path: string): Promise { let parsed: unknown; try { parsed = JSON.parse(line); - } catch (err) { - warnings.push(`line ${lineNo}: invalid JSON (${(err as Error).message})`); + } catch (error) { + warnings.push(`line ${lineNo}: invalid JSON (${(error as Error).message})`); continue; } - if (!isObject(parsed) || typeof parsed.type !== 'string') { + if (!isObject(parsed) || typeof parsed['type'] !== 'string') { warnings.push(`line ${lineNo}: missing 'type' field`); continue; } if (metadata === null) { - if (parsed.type !== 'metadata') { + if (parsed['type'] !== 'metadata') { throw new Error(`Wire file missing metadata header at line ${lineNo}`); } const pv = parsed['protocol_version']; const ca = parsed['created_at']; if (typeof pv !== 'string' || typeof ca !== 'number') { - throw new Error(`Wire metadata malformed at line ${lineNo}`); + throw new TypeError(`Wire metadata malformed at line ${lineNo}`); } try { migrations = resolveWireMigrations(pv); - } catch (err) { + } catch (error) { warnings.push( - `unrecognised protocol_version "${pv}" — parsing as best-effort (${(err as Error).message})`, + `unrecognised protocol_version "${pv}" — parsing as best-effort (${(error as Error).message})`, ); migrations = bestEffortMigrations(); } @@ -85,18 +85,18 @@ export async function readAgentWire(path: string): Promise { try { migrated = migrations.length === 0 - ? raw + ? (structuredClone(raw) as Record) : (migrateWireRecord( raw as Record & { type: string }, migrations, ) as Record); - } catch (err) { + } catch (error) { // A single record that won't migrate is not fatal — keep the raw // payload so the UI can still render whatever fields it understands. warnings.push( - `line ${lineNo}: migration failed (${(err as Error).message}); using raw record`, + `line ${lineNo}: migration failed (${(error as Error).message}); using raw record`, ); - migrated = raw; + migrated = structuredClone(raw) as Record; } records.push({ lineNo, data: migrated as AgentRecord, raw }); } diff --git a/apps/vis/server/src/routes/blobs.ts b/apps/vis/server/src/routes/blobs.ts new file mode 100644 index 0000000000..3c39a215f4 --- /dev/null +++ b/apps/vis/server/src/routes/blobs.ts @@ -0,0 +1,45 @@ +import { Hono } from 'hono'; +import { join } from 'node:path'; +import { readFile } from 'node:fs/promises'; + +import { KIMI_CODE_HOME } from '../config'; +import { isSafeAgentId, readSessionDetail } from '../lib/session-store'; +import { isSafeBlobHash } from '../lib/blob-resolver'; + +export function blobsRoute(home: string = KIMI_CODE_HOME): Hono { + const r = new Hono(); + r.get('/:id/blobs/:hash', async (c) => { + const id = c.req.param('id'); + const agentId = c.req.query('agent') ?? 'main'; + const hash = c.req.param('hash'); + if (!isSafeAgentId(agentId)) { + return c.json({ error: 'invalid agent id', code: 'BAD_REQUEST' }, 400); + } + if (!isSafeBlobHash(hash)) { + return c.json({ error: 'invalid blob hash', code: 'BAD_REQUEST' }, 400); + } + const detail = await readSessionDetail(home, id); + if (!detail) { + return c.json({ error: 'session not found', code: 'NOT_FOUND' }, 404); + } + const agent = detail.agents.find((a) => a.agentId === agentId); + if (!agent) { + return c.json( + { error: `agent "${agentId}" not found`, code: 'NOT_FOUND' }, + 404, + ); + } + const blobPath = join(agent.homedir, 'blobs', hash); + let content: Buffer; + try { + content = await readFile(blobPath); + } catch { + return c.json({ error: 'blob not found', code: 'NOT_FOUND' }, 404); + } + const mimeType = c.req.query('mime') ?? 'application/octet-stream'; + return new Response(content, { + headers: { 'content-type': mimeType }, + }); + }); + return r; +} diff --git a/apps/vis/server/src/routes/context.ts b/apps/vis/server/src/routes/context.ts index df28eb91c6..1cd26473e4 100644 --- a/apps/vis/server/src/routes/context.ts +++ b/apps/vis/server/src/routes/context.ts @@ -3,6 +3,7 @@ import { join } from 'node:path'; import { KIMI_CODE_HOME } from '../config'; import { isSafeAgentId, readSessionDetail } from '../lib/session-store'; +import { rehydrateWireEntries } from '../lib/blob-resolver'; import { readAgentWire } from '../lib/wire-reader'; import { projectContext } from '../lib/context-projector'; @@ -26,6 +27,8 @@ export function contextRoute(): Hono { const wire = await readAgentWire( join(detail.sessionDir, 'agents', agentId, 'wire.jsonl'), ); + const baseUrl = new URL(c.req.url).origin; + rehydrateWireEntries(wire.records, id, agentId, baseUrl); const proj = projectContext(wire.records); return c.json({ sessionId: id, diff --git a/apps/vis/server/src/routes/wire.ts b/apps/vis/server/src/routes/wire.ts index a35fd4dc84..c2e46aab84 100644 --- a/apps/vis/server/src/routes/wire.ts +++ b/apps/vis/server/src/routes/wire.ts @@ -3,6 +3,7 @@ import { join } from 'node:path'; import { KIMI_CODE_HOME } from '../config'; import { isSafeAgentId, readSessionDetail } from '../lib/session-store'; +import { rehydrateWireEntries } from '../lib/blob-resolver'; import { readAgentWire } from '../lib/wire-reader'; export function wireRoute(): Hono { @@ -28,6 +29,8 @@ export function wireRoute(): Hono { const result = await readAgentWire( join(detail.sessionDir, 'agents', agentId, 'wire.jsonl'), ); + const baseUrl = new URL(c.req.url).origin; + rehydrateWireEntries(result.records, id, agentId, baseUrl); return c.json({ sessionId: id, agentId, diff --git a/apps/vis/server/test/lib/blob-resolver.test.ts b/apps/vis/server/test/lib/blob-resolver.test.ts new file mode 100644 index 0000000000..bf29441c14 --- /dev/null +++ b/apps/vis/server/test/lib/blob-resolver.test.ts @@ -0,0 +1,140 @@ +import { describe, it, expect } from 'vitest'; +import { resolveBlobRefUrl, isSafeBlobHash, rehydrateWireEntries } from '../../src/lib/blob-resolver'; + +describe('blob-resolver', () => { + describe('resolveBlobRefUrl', () => { + it('converts a well-formed blobref into a relative route by default', () => { + const url = resolveBlobRefUrl( + 'blobref:image/png;abc123def456', + 'sess-1', + 'main', + ); + expect(url).toBe( + '/api/sessions/sess-1/blobs/abc123def456?agent=main&mime=image%2Fpng', + ); + }); + + it('returns an absolute URL when baseUrl is provided', () => { + const url = resolveBlobRefUrl( + 'blobref:image/png;abc123def456', + 'sess-1', + 'main', + 'http://localhost:3001', + ); + expect(url).toBe( + 'http://localhost:3001/api/sessions/sess-1/blobs/abc123def456?agent=main&mime=image%2Fpng', + ); + }); + + it('returns non-blobref URLs unchanged', () => { + expect(resolveBlobRefUrl('https://example.com/x.png', 's', 'a')).toBe( + 'https://example.com/x.png', + ); + expect(resolveBlobRefUrl('data:image/png;base64,abc', 's', 'a')).toBe( + 'data:image/png;base64,abc', + ); + }); + + it('returns malformed blobrefs unchanged', () => { + expect(resolveBlobRefUrl('blobref:nosemicolon', 's', 'a')).toBe( + 'blobref:nosemicolon', + ); + expect(resolveBlobRefUrl('blobref:image/png;', 's', 'a')).toBe( + 'blobref:image/png;', + ); + }); + }); + + describe('isSafeBlobHash', () => { + it('accepts a 64-char hex string', () => { + expect(isSafeBlobHash('a'.repeat(64))).toBe(true); + expect(isSafeBlobHash('0'.repeat(64))).toBe(true); + expect(isSafeBlobHash('f'.repeat(64))).toBe(true); + }); + + it('rejects non-hex, wrong length, and path-traversal strings', () => { + expect(isSafeBlobHash('')).toBe(false); + expect(isSafeBlobHash('a'.repeat(63))).toBe(false); + expect(isSafeBlobHash('a'.repeat(65))).toBe(false); + expect(isSafeBlobHash('x'.repeat(64))).toBe(false); + expect(isSafeBlobHash('../etc/passwd')).toBe(false); + expect(isSafeBlobHash('a'.repeat(64) + '\n')).toBe(false); + }); + }); + + describe('rehydrateWireEntries', () => { + it('mutates entry.data but leaves entry.raw untouched', () => { + const data: Record = { + type: 'turn.prompt', + input: [ + { + type: 'image_url', + imageUrl: { url: 'blobref:image/png;hashA' }, + }, + ], + }; + const raw = JSON.parse(JSON.stringify(data)); + const entries = [{ lineNo: 1, data: data as any, raw }]; + + rehydrateWireEntries(entries, 'sess-1', 'main'); + + expect((entries[0]!.data as any).input[0].imageUrl.url).toBe( + '/api/sessions/sess-1/blobs/hashA?agent=main&mime=image%2Fpng', + ); + expect((entries[0]!.raw as any).input[0].imageUrl.url).toBe( + 'blobref:image/png;hashA', + ); + }); + + it('handles audio_url, video_url, and nested tool result parts', () => { + const data: Record = { + type: 'context.append_loop_event', + event: { + type: 'tool.result', + result: { + output: [ + { + type: 'audio_url', + audioUrl: { url: 'blobref:audio/wav;hashB' }, + }, + ], + }, + }, + }; + const entries = [{ lineNo: 1, data: data as any, raw: {} }]; + + rehydrateWireEntries(entries, 'sess-2', 'sub-1'); + + const parts = (entries[0]!.data as any).event.result.output; + expect(parts[0].audioUrl.url).toBe( + '/api/sessions/sess-2/blobs/hashB?agent=sub-1&mime=audio%2Fwav', + ); + }); + + it('resolves blobrefs to absolute URLs when baseUrl is provided', () => { + const data: Record = { + type: 'turn.prompt', + input: [ + { + type: 'image_url', + imageUrl: { url: 'blobref:image/png;hashC' }, + }, + ], + }; + const entries = [{ lineNo: 1, data: data as any, raw: {} }]; + + rehydrateWireEntries(entries, 'sess-3', 'main', 'http://localhost:3001'); + + expect((entries[0]!.data as any).input[0].imageUrl.url).toBe( + 'http://localhost:3001/api/sessions/sess-3/blobs/hashC?agent=main&mime=image%2Fpng', + ); + }); + + it('ignores records without media URLs', () => { + const data = { type: 'config.update', cwd: '/tmp' }; + const entries = [{ lineNo: 1, data: data as any, raw: {} }]; + rehydrateWireEntries(entries, 's', 'a'); + expect(entries[0]!.data).toEqual(data); + }); + }); +}); diff --git a/apps/vis/server/test/lib/wire-reader.test.ts b/apps/vis/server/test/lib/wire-reader.test.ts index 7efd1d49b2..6f67d62855 100644 --- a/apps/vis/server/test/lib/wire-reader.test.ts +++ b/apps/vis/server/test/lib/wire-reader.test.ts @@ -78,7 +78,7 @@ describe('wire-reader', () => { // written" view in the detail panel relies on. const rawMsg = (entry.raw as { message: { toolCalls: Array<{ function: unknown; name?: unknown }> } }).message; expect(rawMsg.toolCalls[0]).toHaveProperty('function'); - expect(rawMsg.toolCalls[0].function).toEqual({ name: 'Read', arguments: '{"path":"/x"}' }); + expect(rawMsg.toolCalls[0]!.function).toEqual({ name: 'Read', arguments: '{"path":"/x"}' }); expect(rawMsg.toolCalls[0]).not.toHaveProperty('name'); } finally { await rm(dir, { recursive: true, force: true }); @@ -91,12 +91,12 @@ describe('wire-reader', () => { const path = join(sessionDir, 'agents', 'main', 'wire.jsonl'); const { writeFile, readFile } = await import('node:fs/promises'); const lines = (await readFile(path, 'utf8')).split('\n'); - lines[0] = '{"type":"metadata","protocol_version":"2.2","created_at":1}'; + lines[0] = '{"type":"metadata","protocol_version":"0.9","created_at":1}'; await writeFile(path, lines.join('\n')); const result = await readAgentWire(path); - expect(result.metadata.protocolVersion).toBe('2.2'); + expect(result.metadata.protocolVersion).toBe('0.9'); expect(result.records.length).toBeGreaterThan(0); - expect(result.warnings.some((w) => /unrecognised protocol_version.*2\.2/i.test(w))).toBe( + expect(result.warnings.some((w) => /unrecognised protocol_version.*0\.9/i.test(w))).toBe( true, ); }); diff --git a/apps/vis/server/test/routes/blobs.test.ts b/apps/vis/server/test/routes/blobs.test.ts new file mode 100644 index 0000000000..cd1e25a922 --- /dev/null +++ b/apps/vis/server/test/routes/blobs.test.ts @@ -0,0 +1,93 @@ +import { describe, it, expect, afterEach } from 'vitest'; +import { writeFile, mkdir } from 'node:fs/promises'; +import { join } from 'node:path'; +import { buildSessionFixture } from '../fixtures/build'; +import { blobsRoute } from '../../src/routes/blobs'; + +describe('blobs route', () => { + let cleanup: (() => Promise) | null = null; + afterEach(async () => { if (cleanup) await cleanup(); cleanup = null; }); + + it('serves a blob with the requested content-type', async () => { + const { home, sessionDir, cleanup: c } = await buildSessionFixture('sample-main'); + cleanup = c; + const blobDir = join(sessionDir, 'agents', 'main', 'blobs'); + await mkdir(blobDir, { recursive: true }); + const hash = 'a'.repeat(64); + await writeFile(join(blobDir, hash), Buffer.from('binary-content')); + + const app = blobsRoute(home); + const res = await app.request( + `/session_fixture/blobs/${hash}?agent=main&mime=image/png`, + ); + expect(res.status).toBe(200); + expect(res.headers.get('content-type')).toBe('image/png'); + const body = await res.text(); + expect(body).toBe('binary-content'); + }); + + it('defaults mime to application/octet-stream', async () => { + const { home, sessionDir, cleanup: c } = await buildSessionFixture('sample-main'); + cleanup = c; + const blobDir = join(sessionDir, 'agents', 'main', 'blobs'); + await mkdir(blobDir, { recursive: true }); + const hash = 'b'.repeat(64); + await writeFile(join(blobDir, hash), Buffer.from('x')); + + const app = blobsRoute(home); + const res = await app.request( + `/session_fixture/blobs/${hash}?agent=main`, + ); + expect(res.status).toBe(200); + expect(res.headers.get('content-type')).toBe('application/octet-stream'); + }); + + it('returns 404 for missing session', async () => { + const app = blobsRoute(); + const res = await app.request( + `/no-such-session/blobs/${'c'.repeat(64)}?agent=main&mime=image/png`, + ); + expect(res.status).toBe(404); + expect(await res.json()).toMatchObject({ code: 'NOT_FOUND' }); + }); + + it('returns 404 for missing agent', async () => { + const { home, cleanup: c } = await buildSessionFixture('sample-main'); + cleanup = c; + const app = blobsRoute(home); + const res = await app.request( + `/session_fixture/blobs/${'d'.repeat(64)}?agent=no-such-agent&mime=image/png`, + ); + expect(res.status).toBe(404); + expect(await res.json()).toMatchObject({ code: 'NOT_FOUND' }); + }); + + it('returns 404 for missing blob file', async () => { + const { home, cleanup: c } = await buildSessionFixture('sample-main'); + cleanup = c; + const app = blobsRoute(home); + const res = await app.request( + `/session_fixture/blobs/${'e'.repeat(64)}?agent=main&mime=image/png`, + ); + expect(res.status).toBe(404); + expect(await res.json()).toMatchObject({ code: 'NOT_FOUND' }); + }); + + it('returns 400 for invalid agent id', async () => { + const app = blobsRoute(); + const res = await app.request( + `/session_fixture/blobs/${'f'.repeat(64)}?agent=../escape&mime=image/png`, + ); + expect(res.status).toBe(400); + expect(await res.json()).toMatchObject({ code: 'BAD_REQUEST' }); + }); + + it('returns 400 for invalid blob hash', async () => { + const app = blobsRoute(); + const res = await app.request( + `/session_fixture/blobs/not-a-hash?agent=main&mime=image/png`, + ); + expect(res.status).toBe(400); + expect(await res.json()).toMatchObject({ code: 'BAD_REQUEST' }); + }); +}); diff --git a/apps/vis/web/src/components/shared/ImagePreview.tsx b/apps/vis/web/src/components/shared/ImagePreview.tsx index dd97010179..9ec523bdf1 100644 --- a/apps/vis/web/src/components/shared/ImagePreview.tsx +++ b/apps/vis/web/src/components/shared/ImagePreview.tsx @@ -16,7 +16,8 @@ interface ImagePreviewProps { export function ImagePreview({ url, label = 'image_url' }: ImagePreviewProps) { const [open, setOpen] = useState(false); const [failed, setFailed] = useState(false); - const supported = url.startsWith('data:image/') || /^https?:\/\//.test(url); + const supported = + url.startsWith('data:image/') || /^https?:\/\//.test(url) || url.startsWith('/'); const sizeLabel = url.startsWith('data:image/') ? `${url.length.toLocaleString()} chars` : new URL(url, window.location.href).hostname; diff --git a/packages/agent-core/src/agent/index.ts b/packages/agent-core/src/agent/index.ts index cdb1def9af..786206ee94 100644 --- a/packages/agent-core/src/agent/index.ts +++ b/packages/agent-core/src/agent/index.ts @@ -40,6 +40,7 @@ import { PermissionManager, type PermissionManagerOptions } from './permission'; import { PlanMode } from './plan'; import { AgentRecords, + BlobStore, FileSystemAgentRecordPersistence, type AgentRecord, type AgentRecordPersistence, @@ -96,6 +97,7 @@ export class Agent { readonly hooks: HookEngine | undefined; readonly type: AgentType; + readonly blobStore: BlobStore | undefined; readonly records: AgentRecords; readonly fullCompaction: FullCompaction; readonly context: ContextMemory; @@ -132,6 +134,9 @@ export class Agent { this.rpc = config.rpc; this.telemetry = config.telemetry ?? noopTelemetryClient; + this.blobStore = config.homedir + ? new BlobStore({ blobsDir: join(config.homedir, 'blobs') }) + : undefined; this.records = new AgentRecords( this, config.persistence ?? @@ -140,6 +145,7 @@ export class Agent { onError: (error) => { this.emitRecordsWriteError(error); }, + blobStore: this.blobStore, }) : undefined), ); diff --git a/packages/agent-core/src/agent/records/blobref.ts b/packages/agent-core/src/agent/records/blobref.ts new file mode 100644 index 0000000000..997b695f1e --- /dev/null +++ b/packages/agent-core/src/agent/records/blobref.ts @@ -0,0 +1,249 @@ +import { createHash } from 'node:crypto'; +import { mkdir, open, readFile } from 'node:fs/promises'; +import { join } from 'pathe'; +import type { ContentPart } from '@moonshot-ai/kosong'; +import type { AgentRecord } from './types'; + +const DEFAULT_THRESHOLD = 4096; +const DEFAULT_MAX_CACHE_SIZE = 50 * 1024 * 1024; +const BLOBREF_PROTOCOL = 'blobref:'; +const DATA_URI_HEADER_RE = /^data:([^;]+);base64,/; +const MISSING_MEDIA_PLACEHOLDER = '[media missing]'; + +export function isBlobRef(url: string): boolean { + return url.startsWith(BLOBREF_PROTOCOL); +} + +export interface BlobStoreOptions { + readonly blobsDir: string; + readonly threshold?: number; + readonly maxCacheSize?: number; +} + +export class BlobStore { + private readonly blobsDir: string; + private readonly threshold: number; + private readonly maxCacheSize: number; + private readonly cache = new Map(); + private readonly cacheSizes = new Map(); + private currentCacheSize = 0; + + constructor(options: BlobStoreOptions) { + this.blobsDir = options.blobsDir; + this.threshold = options.threshold ?? DEFAULT_THRESHOLD; + this.maxCacheSize = options.maxCacheSize ?? DEFAULT_MAX_CACHE_SIZE; + } + + async offload(record: AgentRecord): Promise { + switch (record.type) { + case 'turn.prompt': + case 'turn.steer': { + const input = await this.offloadParts(record.input); + return input === record.input ? record : { ...record, input }; + } + case 'context.append_message': { + const content = await this.offloadParts(record.message.content); + return content === record.message.content + ? record + : { ...record, message: { ...record.message, content } }; + } + case 'context.append_loop_event': { + const event = record.event; + if (event.type !== 'tool.result' || typeof event.result.output === 'string') { + return record; + } + const output = await this.offloadParts(event.result.output); + if (output === event.result.output) return record; + return { + ...record, + event: { + ...event, + result: { ...event.result, output }, + }, + }; + } + default: + return record; + } + } + + private async offloadParts(parts: readonly ContentPart[]): Promise { + let changed = false; + const out: ContentPart[] = []; + for (const part of parts) { + const next = await this.offloadContentPart(part); + if (next !== part) changed = true; + out.push(next); + } + return changed ? out : (parts as ContentPart[]); + } + + async rehydrate(record: AgentRecord): Promise { + switch (record.type) { + case 'turn.prompt': + case 'turn.steer': + await this.rehydrateParts(record.input); + break; + case 'context.append_message': + await this.rehydrateParts(record.message.content); + break; + case 'context.append_loop_event': { + const event = record.event; + if (event.type === 'tool.result' && typeof event.result.output !== 'string') { + await this.rehydrateParts(event.result.output); + } + break; + } + default: + break; + } + } + + async rehydrateParts(parts: readonly ContentPart[]): Promise { + for (const part of parts) { + await this.rehydrateContentPart(part); + } + } + + private async offloadContentPart(part: ContentPart): Promise { + let updated: Record | undefined; + for (const [key, value] of Object.entries(part)) { + const mediaObj = asMediaContainer(value); + if (mediaObj === undefined) continue; + + const url = mediaObj.url; + if (typeof url !== 'string') continue; + + const newUrl = await this.maybeOffloadString(url); + if (newUrl === url) continue; + + if (updated === undefined) updated = { ...part }; + updated[key] = { ...(value as object), url: newUrl }; + } + return updated === undefined ? part : (updated as unknown as ContentPart); + } + + private async rehydrateContentPart(part: ContentPart): Promise { + const record = part as unknown as Record; + for (const value of Object.values(record)) { + const mediaObj = asMediaContainer(value); + if (mediaObj === undefined) continue; + + const url = mediaObj.url; + if (typeof url !== 'string' || !isBlobRef(url)) continue; + + const newUrl = await this.rehydrateBlobRefUrl(url); + mediaObj.url = newUrl ?? MISSING_MEDIA_PLACEHOLDER; + } + } + + private async rehydrateBlobRefUrl(url: string): Promise { + const rest = url.slice(BLOBREF_PROTOCOL.length); + const semiIdx = rest.indexOf(';'); + if (semiIdx === -1) { + return undefined; + } + const mimeType = rest.slice(0, semiIdx); + const hash = rest.slice(semiIdx + 1); + if (hash.length === 0) { + return undefined; + } + const payload = await this.readBlob(hash); + if (payload === undefined) { + return undefined; + } + return `data:${mimeType};base64,${payload.toString('base64')}`; + } + + private async readBlob(hash: string): Promise { + const cached = this.cache.get(hash); + if (cached !== undefined) { + // Move the entry to the end so it lives longer than less-recently-used items. + this.cache.delete(hash); + this.cache.set(hash, cached); + return cached; + } + const payload = await readFile(join(this.blobsDir, hash)).catch(() => undefined); + if (payload !== undefined) { + this.setCache(hash, payload); + } + return payload; + } + + private async maybeOffloadString(value: string): Promise { + if (value.startsWith(BLOBREF_PROTOCOL)) { + return value; + } + const match = DATA_URI_HEADER_RE.exec(value); + if (match === null) { + return value; + } + const mimeType = match[1]!; + const payload = value.slice(match[0].length); + if (payload.length < this.threshold) { + return value; + } + return this.writeBlob(mimeType, payload); + } + + private async writeBlob(mimeType: string, base64Payload: string): Promise { + await mkdir(this.blobsDir, { recursive: true, mode: 0o700 }); + const hash = createHash('sha256').update(base64Payload, 'utf8').digest('hex'); + const blobPath = join(this.blobsDir, hash); + const binary = Buffer.from(base64Payload, 'base64'); + try { + const fh = await open(blobPath, 'wx'); + try { + await fh.writeFile(binary); + await fh.sync(); + } finally { + await fh.close(); + } + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + // EEXIST means the identical payload was already written; deduplication. + if (code !== 'EEXIST') throw error; + } + this.setCache(hash, binary); + return `${BLOBREF_PROTOCOL}${mimeType};${hash}`; + } + + private setCache(hash: string, payload: Buffer): void { + const size = payload.byteLength; + const alreadyCached = this.cache.has(hash); + if (alreadyCached) { + const oldSize = this.cacheSizes.get(hash) ?? 0; + this.currentCacheSize += size - oldSize; + // Re-insert to update LRU position. + this.cache.delete(hash); + } else { + if (size > this.maxCacheSize) { + // Skip caching a single blob that exceeds the entire cap. + return; + } + while (this.currentCacheSize + size > this.maxCacheSize && this.cache.size > 0) { + this.evictLRU(); + } + this.currentCacheSize += size; + } + this.cache.set(hash, payload); + this.cacheSizes.set(hash, size); + } + + private evictLRU(): void { + const lru = this.cache.keys().next().value; + if (lru === undefined) return; + const size = this.cacheSizes.get(lru) ?? 0; + this.currentCacheSize -= size; + this.cache.delete(lru); + this.cacheSizes.delete(lru); + } +} + +function asMediaContainer(value: unknown): { url: unknown } | undefined { + if (value === null || typeof value !== 'object' || Array.isArray(value)) { + return undefined; + } + const obj = value as Record; + return 'url' in obj ? (obj as { url: unknown }) : undefined; +} diff --git a/packages/agent-core/src/agent/records/index.ts b/packages/agent-core/src/agent/records/index.ts index 2e4ede744e..e21dbaa4a5 100644 --- a/packages/agent-core/src/agent/records/index.ts +++ b/packages/agent-core/src/agent/records/index.ts @@ -16,6 +16,8 @@ export { InMemoryAgentRecordPersistence, } from './persistence'; export type { FileSystemAgentRecordPersistenceOptions } from './persistence'; +export { BlobStore, isBlobRef } from './blobref'; +export type { BlobStoreOptions } from './blobref'; // Contract: restore MUST NOT emit UI events, call the LLM, execute tools, or // touch the filesystem in a way that triggers external side effects. Each case @@ -179,6 +181,11 @@ export class AgentRecords { this.persistence.rewrite(replayedRecords); await this.persistence.flush(); } + if (this.agent.blobStore !== undefined) { + for (const msg of this.agent.context.history) { + await this.agent.blobStore.rehydrateParts(msg.content); + } + } return { warning }; } diff --git a/packages/agent-core/src/agent/records/migration/index.ts b/packages/agent-core/src/agent/records/migration/index.ts index d87a214ebb..6bd443dcb7 100644 --- a/packages/agent-core/src/agent/records/migration/index.ts +++ b/packages/agent-core/src/agent/records/migration/index.ts @@ -1,8 +1,9 @@ import { migrateV1_0ToV1_1 } from './v1.1'; import { migrateV1_1ToV1_2 } from './v1.2'; +import { migrateV1_2ToV1_3 } from './v1.3'; // Wire protocol versions currently support only the `number.number` format. -export const AGENT_WIRE_PROTOCOL_VERSION = '1.2'; +export const AGENT_WIRE_PROTOCOL_VERSION = '1.3'; export interface WireMigrationRecord { readonly type: string; @@ -15,7 +16,11 @@ export interface WireMigration { migrateRecord(record: WireMigrationRecord): WireMigrationRecord; } -const MIGRATIONS: readonly WireMigration[] = [migrateV1_0ToV1_1, migrateV1_1ToV1_2]; +const MIGRATIONS: readonly WireMigration[] = [ + migrateV1_0ToV1_1, + migrateV1_1ToV1_2, + migrateV1_2ToV1_3, +]; export function isNewerWireVersion(readVersion: string): boolean { return compareWireVersions(readVersion, AGENT_WIRE_PROTOCOL_VERSION) > 0; diff --git a/packages/agent-core/src/agent/records/migration/v1.3.ts b/packages/agent-core/src/agent/records/migration/v1.3.ts new file mode 100644 index 0000000000..670b06cd2e --- /dev/null +++ b/packages/agent-core/src/agent/records/migration/v1.3.ts @@ -0,0 +1,18 @@ +import type { WireMigration, WireMigrationRecord } from './index'; + +/** + * v1.2 -> v1.3 is a bump-only migration. + * + * v1.3 introduces blobref offloading for large base64 media payloads. + * Records written by v1.3+ may contain `blobref:;` URLs in + * message content instead of inline `data:` URIs. Wire records are still + * valid JSON and do not require transformation; the blobref format is + * transparently handled at read/write time by BlobStore. + */ +export const migrateV1_2ToV1_3: WireMigration = { + sourceVersion: '1.2', + targetVersion: '1.3', + migrateRecord(record: WireMigrationRecord): WireMigrationRecord { + return record; + }, +}; diff --git a/packages/agent-core/src/agent/records/persistence.ts b/packages/agent-core/src/agent/records/persistence.ts index f2fc4a93b2..0e51942faf 100644 --- a/packages/agent-core/src/agent/records/persistence.ts +++ b/packages/agent-core/src/agent/records/persistence.ts @@ -3,10 +3,12 @@ import { mkdir, open } from 'node:fs/promises'; import { dirname } from 'pathe'; import { syncDir } from '../../utils/fs'; +import type { BlobStore } from './blobref'; import { type AgentRecord, type AgentRecordPersistence } from './types'; export interface FileSystemAgentRecordPersistenceOptions { readonly onError?: ((error: unknown) => void) | undefined; + readonly blobStore?: BlobStore | undefined; } export interface InMemoryAgentRecordPersistenceOptions { @@ -169,7 +171,13 @@ export class FileSystemAgentRecordPersistence implements AgentRecordPersistence const batch = this.pendingRecords.splice(0); this.shouldClear = false; - const content = batch.map((e) => JSON.stringify(e) + '\n').join(''); + const writable = this.options.blobStore !== undefined + ? await Promise.all( + batch.map((record) => this.options.blobStore!.offload(record)), + ) + : batch; + + const content = writable.map((e) => JSON.stringify(e) + '\n').join(''); const directory = dirname(this.filePath); await mkdir(directory, { recursive: true }); diff --git a/packages/agent-core/src/agent/turn/kosong-llm.ts b/packages/agent-core/src/agent/turn/kosong-llm.ts index f8e254e071..b33614e43e 100644 --- a/packages/agent-core/src/agent/turn/kosong-llm.ts +++ b/packages/agent-core/src/agent/turn/kosong-llm.ts @@ -114,7 +114,7 @@ export class KosongLLM implements LLM { effectiveProvider, this.systemPrompt, [...params.tools], - [...params.messages], + params.messages, callbacks, generateOptions(params, { onRequestStart: markRequestStart, diff --git a/packages/agent-core/test/agent/records/blobref.test.ts b/packages/agent-core/test/agent/records/blobref.test.ts new file mode 100644 index 0000000000..b59176e348 --- /dev/null +++ b/packages/agent-core/test/agent/records/blobref.test.ts @@ -0,0 +1,398 @@ +import { randomBytes } from 'node:crypto'; +import { mkdir, readdir, readFile, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'pathe'; + +import { afterEach, describe, expect, it } from 'vitest'; + +import { BlobStore, isBlobRef } from '../../../src/agent/records/blobref'; +import type { AgentRecord } from '../../../src/agent/records'; + +const cleanups: string[] = []; + +afterEach(async () => { + for (const dir of cleanups.splice(0)) { + await rm(dir, { recursive: true, force: true }).catch(() => {}); + } +}); + +function firstImageUrl(record: AgentRecord): string { + return (record as unknown as { input: [{ imageUrl: { url: string } }] }).input[0].imageUrl.url; +} + +async function makeStore(options?: { maxCacheSize?: number; threshold?: number }): Promise<{ store: BlobStore; blobsDir: string }> { + const blobsDir = join(tmpdir(), `blobref-test-${randomBytes(6).toString('hex')}`); + await mkdir(blobsDir, { recursive: true }); + cleanups.push(blobsDir); + return { + store: new BlobStore({ + blobsDir, + threshold: options?.threshold ?? 4096, + maxCacheSize: options?.maxCacheSize, + }), + blobsDir, + }; +} + +describe('blobref', () => { + it('offloads large data URIs and replaces with blobref', async () => { + const { store, blobsDir } = await makeStore(); + const payload = 'A'.repeat(5000); + const dataUri = `data:image/png;base64,${payload}`; + + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: dataUri } }], + origin: { kind: 'user' }, + }; + + const offloaded = await store.offload(record); + + const url = (offloaded as unknown as { input: [{ imageUrl: { url: string } }] }).input[0] + .imageUrl.url; + expect(isBlobRef(url)).toBe(true); + expect(url.startsWith('blobref:')).toBe(true); + expect(url.startsWith('blobref:image/png;')).toBe(true); + + const files = await readdir(blobsDir); + expect(files).toHaveLength(1); + expect((await readFile(join(blobsDir, files[0]!))).toString('base64')).toBe(payload); + }); + + it('does not mutate the input record or its content parts', async () => { + const { store } = await makeStore(); + const payload = 'M'.repeat(5000); + const dataUri = `data:image/png;base64,${payload}`; + const innerImageUrl = { url: dataUri }; + const part = { type: 'image_url', imageUrl: innerImageUrl } as const; + const record: AgentRecord = { + type: 'turn.prompt', + input: [part as unknown as { type: 'image_url'; imageUrl: { url: string } }], + origin: { kind: 'user' }, + } as unknown as AgentRecord; + + const offloaded = await store.offload(record); + + // The original record/parts must remain untouched. + expect( + (record as unknown as { input: unknown[] }).input[0], + ).toBe(part); + expect(part.imageUrl).toBe(innerImageUrl); + expect(innerImageUrl.url).toBe(dataUri); + + // The returned record carries the blobref URL. + expect(offloaded).not.toBe(record); + const returnedUrl = ( + offloaded as unknown as { input: [{ imageUrl: { url: string } }] } + ).input[0].imageUrl.url; + expect(returnedUrl.startsWith('blobref:image/png;')).toBe(true); + }); + + it('offloads tool.result media parts in context.append_loop_event records', async () => { + const { store, blobsDir } = await makeStore(); + const payload = 'X'.repeat(5000); + const dataUri = `data:image/png;base64,${payload}`; + const innerImageUrl = { url: dataUri }; + const part = { type: 'image_url', imageUrl: innerImageUrl } as const; + const record: AgentRecord = { + type: 'context.append_loop_event', + event: { + type: 'tool.result', + parentUuid: 'p', + toolCallId: 'tc', + result: { isError: false, output: [part] }, + }, + } as unknown as AgentRecord; + + const offloaded = await store.offload(record); + + // Input record/part untouched — same path that the agent's in-memory + // history shares with this record reference. + expect(innerImageUrl.url).toBe(dataUri); + expect(part.imageUrl).toBe(innerImageUrl); + + // Returned record has blobref URL on a fresh imageUrl object. + const returned = offloaded as unknown as { + event: { result: { output: [{ imageUrl: { url: string } }] } }; + }; + expect(returned.event.result.output[0].imageUrl).not.toBe(innerImageUrl); + expect(returned.event.result.output[0].imageUrl.url.startsWith('blobref:image/png;')).toBe(true); + + const files = await readdir(blobsDir); + expect(files).toHaveLength(1); + }); + + it('returns the same record reference when nothing needs offloading', async () => { + const { store } = await makeStore(); + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'text', text: 'just text' }], + origin: { kind: 'user' }, + } as unknown as AgentRecord; + + const offloaded = await store.offload(record); + expect(offloaded).toBe(record); + }); + + it('skips small data URIs below threshold', async () => { + const { store, blobsDir } = await makeStore(); + const payload = 'short'; + const dataUri = `data:image/png;base64,${payload}`; + + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: dataUri } }], + origin: { kind: 'user' }, + }; + + const offloaded = await store.offload(record); + + // Below threshold: nothing happens; original record is returned as-is. + expect(offloaded).toBe(record); + const files = await readdir(blobsDir).catch(() => []); + expect(files).toHaveLength(0); + }); + + it('skips existing blobrefs during offload', async () => { + const { store } = await makeStore(); + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: 'blobref:image/png;abc' } }], + origin: { kind: 'user' }, + }; + + const offloaded = await store.offload(record); + + // Already a blobref: nothing to do; the same reference is returned. + expect(offloaded).toBe(record); + }); + + it('rehydrates blobrefs back to data URIs', async () => { + const { store } = await makeStore(); + const payload = 'B'.repeat(5000); + const dataUri = `data:image/jpeg;base64,${payload}`; + + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: dataUri } }], + origin: { kind: 'user' }, + }; + + const offloaded = await store.offload(record); + await store.rehydrate(offloaded); + + const url = (offloaded as unknown as { input: [{ imageUrl: { url: string } }] }).input[0] + .imageUrl.url; + expect(url).toBe(dataUri); + }); + + it('replaces missing blobs with placeholder text', async () => { + const { store } = await makeStore(); + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: 'blobref:image/png;deadbeef' } }], + origin: { kind: 'user' }, + }; + + await store.rehydrate(record); + + const url = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; + expect(url).toBe('[media missing]'); + }); + + it('deduplicates identical payloads by hash', async () => { + const { store, blobsDir } = await makeStore(); + const payload = 'C'.repeat(5000); + const dataUri = `data:image/png;base64,${payload}`; + + const record1: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: dataUri } }], + origin: { kind: 'user' }, + }; + const record2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: dataUri } }], + origin: { kind: 'user' }, + }; + + await store.offload(record1); + await store.offload(record2); + + const files = await readdir(blobsDir); + expect(files).toHaveLength(1); + }); + + it('rehydrates from write-through cache after blob file is deleted', async () => { + const { store, blobsDir } = await makeStore(); + const payload = 'E'.repeat(5000); + const dataUri = `data:image/png;base64,${payload}`; + + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: dataUri } }], + origin: { kind: 'user' }, + }; + + const offloaded = await store.offload(record); + const files = await readdir(blobsDir); + expect(files).toHaveLength(1); + await rm(join(blobsDir, files[0]!)); + + // Should still rehydrate because offload populated the cache. + await store.rehydrate(offloaded); + const url = (offloaded as unknown as { input: [{ imageUrl: { url: string } }] }).input[0] + .imageUrl.url; + expect(url).toBe(dataUri); + }); + + it('rehydrates from read cache after first disk read', async () => { + const { store, blobsDir } = await makeStore(); + const payload = 'F'.repeat(5000); + const dataUri = `data:image/png;base64,${payload}`; + + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: dataUri } }], + origin: { kind: 'user' }, + }; + + const offloaded = await store.offload(record); + await store.rehydrate(offloaded); + + const files = await readdir(blobsDir); + expect(files).toHaveLength(1); + await rm(join(blobsDir, files[0]!)); + + // Second rehydrate (of a fresh record pointing to the same blobref) + // should still succeed because the first rehydrate populated the read cache. + // After the rehydrate above, `offloaded` carries the data URI again — pull + // the blobref back out by re-offloading or re-using the original offload run. + const offloadedFresh = await store.offload(record); + const record2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: firstImageUrl(offloadedFresh) } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(record2); + expect(firstImageUrl(record2)).toBe(dataUri); + }); + + it('evicts least-recently-used entries when cache size limit is exceeded', async () => { + const limit = 8; // bytes + const { store, blobsDir } = await makeStore({ maxCacheSize: limit, threshold: 1 }); + + const payloadA = 'A'.repeat(4); // 3 bytes after base64 decode + const payloadB = 'B'.repeat(4); // 3 bytes after base64 decode + const payloadC = 'C'.repeat(4); // 3 bytes after base64 decode + + const recordA: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: `data:image/png;base64,${payloadA}` } }], + origin: { kind: 'user' }, + }; + const recordB: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: `data:image/png;base64,${payloadB}` } }], + origin: { kind: 'user' }, + }; + const recordC: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: `data:image/png;base64,${payloadC}` } }], + origin: { kind: 'user' }, + }; + + const offloadedA = await store.offload(recordA); + const offloadedB = await store.offload(recordB); + + // Touch A so it becomes more recent than B. + const recordA_touch: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: firstImageUrl(offloadedA) } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordA_touch); + + // Adding C should evict B (the least-recently-used), not A. + const offloadedC = await store.offload(recordC); + + // Delete all files so only cache can satisfy rehydration. + const files = await readdir(blobsDir); + for (const f of files) { + await rm(join(blobsDir, f)); + } + + // A should still be cached because it was touched after B. + const recordA2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: firstImageUrl(offloadedA) } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordA2); + expect(firstImageUrl(recordA2)).toBe(`data:image/png;base64,${payloadA}`); + + // B should have been evicted. + const recordB2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: firstImageUrl(offloadedB) } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordB2); + expect(firstImageUrl(recordB2)).toBe('[media missing]'); + + // C should still be cached. + const recordC2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: firstImageUrl(offloadedC) } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordC2); + expect(firstImageUrl(recordC2)).toBe(`data:image/png;base64,${payloadC}`); + }); + + it('skips caching a blob larger than the entire cache cap', async () => { + const limit = 8; // bytes + const { store, blobsDir } = await makeStore({ maxCacheSize: limit, threshold: 1 }); + + const small = 'S'.repeat(4); + const large = 'L'.repeat(16); + + const recordSmall: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: `data:image/png;base64,${small}` } }], + origin: { kind: 'user' }, + }; + const recordLarge: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: `data:image/png;base64,${large}` } }], + origin: { kind: 'user' }, + }; + + const offloadedSmall = await store.offload(recordSmall); + const offloadedLarge = await store.offload(recordLarge); + + // Delete all files so only cache can satisfy rehydration. + const files = await readdir(blobsDir); + for (const f of files) { + await rm(join(blobsDir, f)); + } + + // The small blob is still cached. + const recordSmall2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: firstImageUrl(offloadedSmall) } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordSmall2); + expect(firstImageUrl(recordSmall2)).toBe(`data:image/png;base64,${small}`); + + // The large blob was never cached, so rehydration fails. + const recordLarge2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: firstImageUrl(offloadedLarge) } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordLarge2); + expect(firstImageUrl(recordLarge2)).toBe('[media missing]'); + }); +}); diff --git a/packages/agent-core/test/agent/records/persistence.test.ts b/packages/agent-core/test/agent/records/persistence.test.ts index 701515e003..e3b5bb699d 100644 --- a/packages/agent-core/test/agent/records/persistence.test.ts +++ b/packages/agent-core/test/agent/records/persistence.test.ts @@ -1,5 +1,5 @@ import { randomBytes } from 'node:crypto'; -import { mkdir, readFile, rm } from 'node:fs/promises'; +import { mkdir, readFile, readdir, rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'pathe'; @@ -7,6 +7,7 @@ import { afterEach, describe, expect, it } from 'vitest'; import { AGENT_WIRE_PROTOCOL_VERSION, + BlobStore, FileSystemAgentRecordPersistence, InMemoryAgentRecordPersistence, type AgentRecord, @@ -217,6 +218,38 @@ describe('FileSystemAgentRecordPersistence', () => { }).toThrow(); await expect(persistence.flush()).rejects.toBeInstanceOf(Error); }); + + it('offloads large data URIs to blobsDir during append', async () => { + const dir = join(tmpdir(), `wire-blob-test-${randomBytes(6).toString('hex')}`); + await mkdir(dir, { recursive: true }); + cleanups.push(dir); + + const wirePath = join(dir, 'wire.jsonl'); + const blobsDir = join(dir, 'blobs'); + const persistence = new FileSystemAgentRecordPersistence(wirePath, { + blobStore: new BlobStore({ blobsDir }), + }); + + const payload = 'X'.repeat(5000); + const dataUri = `data:image/png;base64,${payload}`; + + persistence.append({ + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: dataUri } }], + origin: { kind: 'user' }, + }); + await persistence.close(); + + const lines = await readLines(wirePath); + expect(lines).toHaveLength(1); + const record = JSON.parse(lines[0]!) as unknown as Record; + const url = ((record['input'] as unknown[])[0] as { imageUrl: { url: string } }).imageUrl.url; + expect(url.startsWith('blobref:')).toBe(true); + + const blobFiles = await readdir(blobsDir); + expect(blobFiles).toHaveLength(1); + expect((await readFile(join(blobsDir, blobFiles[0]!))).toString('base64')).toBe(payload); + }); }); describe('InMemoryAgentRecordPersistence', () => { diff --git a/packages/agent-core/test/agent/turn.test.ts b/packages/agent-core/test/agent/turn.test.ts index c761140cf8..42d3cad6a1 100644 --- a/packages/agent-core/test/agent/turn.test.ts +++ b/packages/agent-core/test/agent/turn.test.ts @@ -255,7 +255,7 @@ describe('Agent turn flow', () => { await ctx.rpc.prompt({ input: [{ type: 'text', text: 'Hello without login' }] }); expect(await ctx.untilTurnEnd()).toMatchInlineSnapshot(` - [wire] metadata { "protocol_version": "1.2", "created_at": "