From 3cd02dffc6115127d1b933acbb3e606a19a55d05 Mon Sep 17 00:00:00 2001 From: _Kerman Date: Wed, 27 May 2026 19:47:56 +0800 Subject: [PATCH 01/16] feat: offload large base64 media payloads to external blob files Wire.jsonl grows very large when conversations contain images, audio, or video embedded as data: URIs. This change extracts base64 payloads exceeding a threshold (default 4KB) into per-agent blob files stored alongside wire.jsonl, replacing the inline data with a lightweight blobref URL. - New blobref.ts handles offload / rehydrate / missing-blob downgrade. - Blobref format mirrors data URI: blobref:;. - FileSystemAgentRecordPersistence transparently offloads on write. - KosongLLM rehydrates blobrefs back to data URIs before calling the LLM provider, so the rest of the stack is unaware of blobrefs. - Export (zip) naturally includes blob directories via recursive file collection, no extra changes needed. Key design decisions: - Offload scans only known record types (turn.prompt, turn.steer, context.append_message, context.append_loop_event tool.result) to avoid accidentally touching unrelated strings. - Within a ContentPart, media URLs are found by structural pattern ({ xxxUrl: { url } }) rather than hard-coding field names, so new media types are supported automatically. - Blob files are named by SHA-256 of the base64 payload, providing natural deduplication. - Missing blobs degrade to a text placeholder so the LLM provider does not receive an invalid media block. --- .changeset/wire-blob-offloading.md | 6 + packages/agent-core/src/agent/index.ts | 2 + .../agent-core/src/agent/records/blobref.ts | 234 ++++++++++++++++++ .../agent-core/src/agent/records/index.ts | 7 + .../src/agent/records/persistence.ts | 14 ++ .../agent-core/src/agent/turn/kosong-llm.ts | 16 +- .../test/agent/records/blobref.test.ts | 187 ++++++++++++++ .../test/agent/records/persistence.test.ts | 32 ++- 8 files changed, 496 insertions(+), 2 deletions(-) create mode 100644 .changeset/wire-blob-offloading.md create mode 100644 packages/agent-core/src/agent/records/blobref.ts create mode 100644 packages/agent-core/test/agent/records/blobref.test.ts diff --git a/.changeset/wire-blob-offloading.md b/.changeset/wire-blob-offloading.md new file mode 100644 index 0000000000..55e90ed88e --- /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. diff --git a/packages/agent-core/src/agent/index.ts b/packages/agent-core/src/agent/index.ts index cdb1def9af..c5f89b7f99 100644 --- a/packages/agent-core/src/agent/index.ts +++ b/packages/agent-core/src/agent/index.ts @@ -140,6 +140,7 @@ export class Agent { onError: (error) => { this.emitRecordsWriteError(error); }, + blobsDir: join(config.homedir, 'blobs'), }) : undefined), ); @@ -194,6 +195,7 @@ export class Agent { capability: this.config.modelCapabilities, generate: this.generate, completionBudgetConfig, + blobsDir: this.homedir ? join(this.homedir, 'blobs') : 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..5051b81c3e --- /dev/null +++ b/packages/agent-core/src/agent/records/blobref.ts @@ -0,0 +1,234 @@ +import { createHash } from 'node:crypto'; +import { mkdir, open, readFile } from 'node:fs/promises'; +import { join } from 'pathe'; +import type { ContentPart, Message } from '@moonshot-ai/kosong'; +import type { AgentRecord } from './types'; + +export interface BlobOffloadOptions { + readonly blobsDir: string; + readonly threshold?: number; +} + +const DEFAULT_THRESHOLD = 4096; +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 async function writeBlob( + blobsDir: string, + mimeType: string, + base64Payload: string, +): Promise { + await mkdir(blobsDir, { recursive: true, mode: 0o700 }); + const hash = createHash('sha256').update(base64Payload, 'utf8').digest('hex'); + const blobPath = join(blobsDir, hash); + try { + const fh = await open(blobPath, 'wx'); + try { + await fh.writeFile(base64Payload, 'utf8'); + await fh.sync(); + } finally { + await fh.close(); + } + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code !== 'EEXIST') throw error; + } + return `${BLOBREF_PROTOCOL}${mimeType};${hash}`; +} + +export async function readBlob(blobsDir: string, hash: string): Promise { + try { + return await readFile(join(blobsDir, hash), 'utf8'); + } catch { + return undefined; + } +} + +export async function offloadRecordBlobs( + record: AgentRecord, + options: BlobOffloadOptions, +): Promise { + const threshold = options.threshold ?? DEFAULT_THRESHOLD; + + switch (record.type) { + case 'turn.prompt': + case 'turn.steer': + await offloadContentParts(record.input, options.blobsDir, threshold); + break; + case 'context.append_message': + await offloadContentParts(record.message.content, options.blobsDir, threshold); + break; + case 'context.append_loop_event': { + const event = record.event; + if (event.type === 'tool.result' && typeof event.result.output !== 'string') { + await offloadContentParts(event.result.output, options.blobsDir, threshold); + } + break; + } + default: + break; + } + + return record; +} + +export async function rehydrateRecordBlobs( + record: AgentRecord, + blobsDir: string, +): Promise { + switch (record.type) { + case 'turn.prompt': + case 'turn.steer': + await rehydrateContentParts(record.input, blobsDir); + break; + case 'context.append_message': + await rehydrateContentParts(record.message.content, blobsDir); + break; + case 'context.append_loop_event': { + const event = record.event; + if (event.type === 'tool.result' && typeof event.result.output !== 'string') { + await rehydrateContentParts(event.result.output, blobsDir); + } + break; + } + default: + break; + } + return record; +} + +export async function rehydrateMessagesBlobs( + messages: readonly Message[], + blobsDir: string, +): Promise { + return Promise.all( + messages.map(async (msg): Promise => { + if (msg.content.length === 0) return msg; + const clone = structuredClone(msg) as Message; + for (const part of clone.content) { + await rehydrateContentPart(part, blobsDir); + } + return { + role: clone.role, + name: clone.name, + content: clone.content.map((part) => downgradeMissingMedia(part)), + toolCalls: clone.toolCalls, + toolCallId: clone.toolCallId, + partial: clone.partial, + }; + }), + ); +} + +async function offloadContentParts( + parts: readonly ContentPart[], + blobsDir: string, + threshold: number, +): Promise { + for (const part of parts) { + await offloadContentPart(part, blobsDir, threshold); + } +} + +async function rehydrateContentParts(parts: readonly ContentPart[], blobsDir: string): Promise { + for (const part of parts) { + await rehydrateContentPart(part, blobsDir); + } +} + +async function offloadContentPart( + part: ContentPart, + blobsDir: string, + threshold: number, +): 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') continue; + + const newUrl = await maybeOffloadString(url, blobsDir, threshold); + if (newUrl !== url) { + mediaObj.url = newUrl; + } + } +} + +async function rehydrateContentPart(part: ContentPart, blobsDir: string): 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 rehydrateBlobRefUrl(url, blobsDir); + mediaObj.url = newUrl ?? MISSING_MEDIA_PLACEHOLDER; + } +} + +function downgradeMissingMedia(part: ContentPart): ContentPart { + const record = part as unknown as Record; + for (const value of Object.values(record)) { + const mediaObj = asMediaContainer(value); + if (mediaObj === undefined) continue; + if (mediaObj.url === MISSING_MEDIA_PLACEHOLDER) { + return { type: 'text', text: MISSING_MEDIA_PLACEHOLDER }; + } + } + return part; +} + +async function maybeOffloadString( + value: string, + blobsDir: string, + threshold: number, +): 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 < threshold) { + return value; + } + return writeBlob(blobsDir, mimeType, payload); +} + +async function rehydrateBlobRefUrl(url: string, blobsDir: 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 readBlob(blobsDir, hash); + if (payload === undefined) { + return undefined; + } + return `data:${mimeType};base64,${payload}`; +} + +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..69d82ec8b7 100644 --- a/packages/agent-core/src/agent/records/index.ts +++ b/packages/agent-core/src/agent/records/index.ts @@ -16,6 +16,13 @@ export { InMemoryAgentRecordPersistence, } from './persistence'; export type { FileSystemAgentRecordPersistenceOptions } from './persistence'; +export { + isBlobRef, + offloadRecordBlobs, + rehydrateRecordBlobs, + rehydrateMessagesBlobs, +} from './blobref'; +export type { BlobOffloadOptions } 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 diff --git a/packages/agent-core/src/agent/records/persistence.ts b/packages/agent-core/src/agent/records/persistence.ts index f2fc4a93b2..48797460d9 100644 --- a/packages/agent-core/src/agent/records/persistence.ts +++ b/packages/agent-core/src/agent/records/persistence.ts @@ -3,10 +3,13 @@ import { mkdir, open } from 'node:fs/promises'; import { dirname } from 'pathe'; import { syncDir } from '../../utils/fs'; +import { offloadRecordBlobs } from './blobref'; import { type AgentRecord, type AgentRecordPersistence } from './types'; export interface FileSystemAgentRecordPersistenceOptions { readonly onError?: ((error: unknown) => void) | undefined; + readonly blobsDir?: string | undefined; + readonly blobThreshold?: number | undefined; } export interface InMemoryAgentRecordPersistenceOptions { @@ -169,6 +172,17 @@ export class FileSystemAgentRecordPersistence implements AgentRecordPersistence const batch = this.pendingRecords.splice(0); this.shouldClear = false; + if (this.options.blobsDir !== undefined) { + await Promise.all( + batch.map((record) => + offloadRecordBlobs(record, { + blobsDir: this.options.blobsDir!, + threshold: this.options.blobThreshold, + }), + ), + ); + } + const content = batch.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..a467769c2f 100644 --- a/packages/agent-core/src/agent/turn/kosong-llm.ts +++ b/packages/agent-core/src/agent/turn/kosong-llm.ts @@ -34,6 +34,7 @@ import type { LLMRequestLogContext, LLMStreamTiming, } from '../../loop'; +import { rehydrateMessagesBlobs } from '../records/blobref'; import { applyCompletionBudget, type CompletionBudgetConfig, @@ -63,6 +64,12 @@ export interface KosongLLMConfig { * final cap is applied to each request. */ readonly completionBudgetConfig?: CompletionBudgetConfig | undefined; + /** + * Directory holding offloaded base64 blobs. When set, media URLs in + * messages are rehydrated from `blobref://` back to `data:` URIs before + * being sent to the LLM. + */ + readonly blobsDir?: string | undefined; } export class KosongLLM implements LLM { @@ -73,6 +80,7 @@ export class KosongLLM implements LLM { private readonly provider: ChatProvider; private readonly generate: GenerateFn; private readonly completionBudgetConfig: CompletionBudgetConfig | undefined; + private readonly blobsDir: string | undefined; constructor(config: KosongLLMConfig) { this.provider = config.provider; @@ -81,9 +89,15 @@ export class KosongLLM implements LLM { this.capability = config.capability; this.generate = config.generate ?? kosongGenerate; this.completionBudgetConfig = config.completionBudgetConfig; + this.blobsDir = config.blobsDir; } async chat(params: LLMChatParams): Promise { + let messages = params.messages; + if (this.blobsDir !== undefined) { + messages = await rehydrateMessagesBlobs(params.messages, this.blobsDir); + } + let requestStartedAt = Date.now(); let firstChunkAt: number | undefined; let streamEndedAt: number | undefined; @@ -114,7 +128,7 @@ export class KosongLLM implements LLM { effectiveProvider, this.systemPrompt, [...params.tools], - [...params.messages], + 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..5bf02454ed --- /dev/null +++ b/packages/agent-core/test/agent/records/blobref.test.ts @@ -0,0 +1,187 @@ +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 { + isBlobRef, + offloadRecordBlobs, + rehydrateRecordBlobs, + rehydrateMessagesBlobs, +} 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(() => {}); + } +}); + +async function makeBlobsDir(): Promise { + const dir = join(tmpdir(), `blobref-test-${randomBytes(6).toString('hex')}`); + await mkdir(dir, { recursive: true }); + cleanups.push(dir); + return dir; +} + +describe('blobref', () => { + it('offloads large data URIs and replaces with blobref', async () => { + const blobsDir = await makeBlobsDir(); + 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' }, + }; + + await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + + const url = (record.input as unknown as [{ imageUrl: { url: string } }])[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]!), 'utf8')).toBe(payload); + }); + + it('skips small data URIs below threshold', async () => { + const blobsDir = await makeBlobsDir(); + 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' }, + }; + + await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + + const url = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; + expect(url).toBe(dataUri); + const files = await readdir(blobsDir).catch(() => []); + expect(files).toHaveLength(0); + }); + + it('skips existing blobrefs during offload', async () => { + const blobsDir = await makeBlobsDir(); + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: 'blobref:image/png;abc' } }], + origin: { kind: 'user' }, + }; + + await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + + const url = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; + expect(url).toBe('blobref:image/png;abc'); + }); + + it('rehydrates blobrefs back to data URIs', async () => { + const blobsDir = await makeBlobsDir(); + 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' }, + }; + + await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + await rehydrateRecordBlobs(record, blobsDir); + + const url = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; + expect(url).toBe(dataUri); + }); + + it('replaces missing blobs with placeholder text', async () => { + const blobsDir = await makeBlobsDir(); + const record: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: 'blobref:image/png;deadbeef' } }], + origin: { kind: 'user' }, + }; + + await rehydrateRecordBlobs(record, blobsDir); + + 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 blobsDir = await makeBlobsDir(); + 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 offloadRecordBlobs(record1, { blobsDir, threshold: 4096 }); + await offloadRecordBlobs(record2, { blobsDir, threshold: 4096 }); + + const files = await readdir(blobsDir); + expect(files).toHaveLength(1); + }); + + it('rehydrates messages with media parts only', async () => { + const blobsDir = await makeBlobsDir(); + const payload = 'D'.repeat(5000); + const dataUri = `data:audio/wav;base64,${payload}`; + + const record: AgentRecord = { + type: 'turn.prompt', + input: [ + { type: 'text', text: 'hello' }, + { type: 'audio_url', audioUrl: { url: dataUri } }, + ], + origin: { kind: 'user' }, + }; + + await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + + const messages = [ + { + role: 'user' as const, + content: [...record.input], + toolCalls: [], + }, + ]; + + const hydrated = await rehydrateMessagesBlobs(messages, blobsDir); + const firstMsg = hydrated[0]!; + expect(firstMsg.content[0]).toEqual({ type: 'text', text: 'hello' }); + expect((firstMsg.content[1]! as { audioUrl: { url: string } }).audioUrl.url).toBe(dataUri); + }); + + it('degrades missing blob in messages to text placeholder', async () => { + const blobsDir = await makeBlobsDir(); + const messages = [ + { + role: 'user' as const, + content: [{ type: 'image_url' as const, imageUrl: { url: 'blobref:image/png;missing' } }], + toolCalls: [], + }, + ]; + + const hydrated = await rehydrateMessagesBlobs(messages, blobsDir); + const firstMsg = hydrated[0]!; + expect(firstMsg.content[0]).toEqual({ type: 'text', text: '[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..b862f5cede 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'; @@ -217,6 +217,36 @@ 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, { 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]!), 'utf8')).toBe(payload); + }); }); describe('InMemoryAgentRecordPersistence', () => { From 8eec7f0204374fc0203053ad591719b1062f70fd Mon Sep 17 00:00:00 2001 From: _Kerman Date: Wed, 27 May 2026 20:23:47 +0800 Subject: [PATCH 02/16] fix --- packages/agent-core/src/agent/index.ts | 9 +- .../agent-core/src/agent/records/blobref.ts | 216 ++++++++---------- .../agent-core/src/agent/records/index.ts | 9 +- .../src/agent/records/persistence.ts | 14 +- .../agent-core/src/agent/turn/kosong-llm.ts | 17 +- .../test/agent/records/blobref.test.ts | 55 ++--- .../test/agent/records/persistence.test.ts | 5 +- 7 files changed, 148 insertions(+), 177 deletions(-) diff --git a/packages/agent-core/src/agent/index.ts b/packages/agent-core/src/agent/index.ts index c5f89b7f99..25a3dd6943 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,7 +145,7 @@ export class Agent { onError: (error) => { this.emitRecordsWriteError(error); }, - blobsDir: join(config.homedir, 'blobs'), + blobStore: this.blobStore, }) : undefined), ); @@ -195,7 +200,7 @@ export class Agent { capability: this.config.modelCapabilities, generate: this.generate, completionBudgetConfig, - blobsDir: this.homedir ? join(this.homedir, 'blobs') : undefined, + blobStore: this.blobStore, }); } diff --git a/packages/agent-core/src/agent/records/blobref.ts b/packages/agent-core/src/agent/records/blobref.ts index 5051b81c3e..2a6190f6ec 100644 --- a/packages/agent-core/src/agent/records/blobref.ts +++ b/packages/agent-core/src/agent/records/blobref.ts @@ -4,11 +4,6 @@ import { join } from 'pathe'; import type { ContentPart, Message } from '@moonshot-ai/kosong'; import type { AgentRecord } from './types'; -export interface BlobOffloadOptions { - readonly blobsDir: string; - readonly threshold?: number; -} - const DEFAULT_THRESHOLD = 4096; const BLOBREF_PROTOCOL = 'blobref:'; const DATA_URI_HEADER_RE = /^data:([^;]+);base64,/; @@ -18,126 +13,92 @@ export function isBlobRef(url: string): boolean { return url.startsWith(BLOBREF_PROTOCOL); } -export async function writeBlob( - blobsDir: string, - mimeType: string, - base64Payload: string, -): Promise { - await mkdir(blobsDir, { recursive: true, mode: 0o700 }); - const hash = createHash('sha256').update(base64Payload, 'utf8').digest('hex'); - const blobPath = join(blobsDir, hash); - try { - const fh = await open(blobPath, 'wx'); - try { - await fh.writeFile(base64Payload, 'utf8'); - await fh.sync(); - } finally { - await fh.close(); - } - } catch (error) { - const code = (error as NodeJS.ErrnoException).code; - if (code !== 'EEXIST') throw error; - } - return `${BLOBREF_PROTOCOL}${mimeType};${hash}`; -} - -export async function readBlob(blobsDir: string, hash: string): Promise { - try { - return await readFile(join(blobsDir, hash), 'utf8'); - } catch { - return undefined; - } +export interface BlobStoreOptions { + readonly blobsDir: string; + readonly threshold?: number; } -export async function offloadRecordBlobs( - record: AgentRecord, - options: BlobOffloadOptions, -): Promise { - const threshold = options.threshold ?? DEFAULT_THRESHOLD; - - switch (record.type) { - case 'turn.prompt': - case 'turn.steer': - await offloadContentParts(record.input, options.blobsDir, threshold); - break; - case 'context.append_message': - await offloadContentParts(record.message.content, options.blobsDir, threshold); - break; - case 'context.append_loop_event': { - const event = record.event; - if (event.type === 'tool.result' && typeof event.result.output !== 'string') { - await offloadContentParts(event.result.output, options.blobsDir, threshold); +export class BlobStore { + private readonly blobsDir: string; + private readonly threshold: number; + + constructor(options: BlobStoreOptions) { + this.blobsDir = options.blobsDir; + this.threshold = options.threshold ?? DEFAULT_THRESHOLD; + } + + async offload(record: AgentRecord): Promise { + switch (record.type) { + case 'turn.prompt': + case 'turn.steer': + for (const part of record.input) { + await offloadContentPart(part, this.blobsDir, this.threshold); + } + break; + case 'context.append_message': + for (const part of record.message.content) { + await offloadContentPart(part, this.blobsDir, this.threshold); + } + break; + case 'context.append_loop_event': { + const event = record.event; + if (event.type === 'tool.result' && typeof event.result.output !== 'string') { + for (const part of event.result.output) { + await offloadContentPart(part, this.blobsDir, this.threshold); + } + } + break; } - break; + default: + break; } - default: - break; } - return record; -} - -export async function rehydrateRecordBlobs( - record: AgentRecord, - blobsDir: string, -): Promise { - switch (record.type) { - case 'turn.prompt': - case 'turn.steer': - await rehydrateContentParts(record.input, blobsDir); - break; - case 'context.append_message': - await rehydrateContentParts(record.message.content, blobsDir); - break; - case 'context.append_loop_event': { - const event = record.event; - if (event.type === 'tool.result' && typeof event.result.output !== 'string') { - await rehydrateContentParts(event.result.output, blobsDir); + async rehydrate(record: AgentRecord): Promise { + switch (record.type) { + case 'turn.prompt': + case 'turn.steer': + for (const part of record.input) { + await rehydrateContentPart(part, this.blobsDir); + } + break; + case 'context.append_message': + for (const part of record.message.content) { + await rehydrateContentPart(part, this.blobsDir); + } + break; + case 'context.append_loop_event': { + const event = record.event; + if (event.type === 'tool.result' && typeof event.result.output !== 'string') { + for (const part of event.result.output) { + await rehydrateContentPart(part, this.blobsDir); + } + } + break; } - break; + default: + break; } - default: - break; } - return record; -} - -export async function rehydrateMessagesBlobs( - messages: readonly Message[], - blobsDir: string, -): Promise { - return Promise.all( - messages.map(async (msg): Promise => { - if (msg.content.length === 0) return msg; - const clone = structuredClone(msg) as Message; - for (const part of clone.content) { - await rehydrateContentPart(part, blobsDir); - } - return { - role: clone.role, - name: clone.name, - content: clone.content.map((part) => downgradeMissingMedia(part)), - toolCalls: clone.toolCalls, - toolCallId: clone.toolCallId, - partial: clone.partial, - }; - }), - ); -} -async function offloadContentParts( - parts: readonly ContentPart[], - blobsDir: string, - threshold: number, -): Promise { - for (const part of parts) { - await offloadContentPart(part, blobsDir, threshold); - } -} - -async function rehydrateContentParts(parts: readonly ContentPart[], blobsDir: string): Promise { - for (const part of parts) { - await rehydrateContentPart(part, blobsDir); + async rehydrateMessages(messages: readonly Message[]): Promise { + return Promise.all( + messages.map(async (msg): Promise => { + if (msg.content.length === 0) return msg; + const clone = structuredClone(msg) as Message; + for (const part of clone.content) { + await rehydrateContentPart(part, this.blobsDir); + } + return { + role: clone.role, + name: clone.name, + content: clone.content.map((part) => downgradeMissingMedia(part)), + toolCalls: clone.toolCalls, + toolCallId: clone.toolCallId, + partial: clone.partial, + }; + }), + ); } } @@ -218,13 +179,36 @@ async function rehydrateBlobRefUrl(url: string, blobsDir: string): Promise undefined); if (payload === undefined) { return undefined; } return `data:${mimeType};base64,${payload}`; } +async function writeBlob( + blobsDir: string, + mimeType: string, + base64Payload: string, +): Promise { + await mkdir(blobsDir, { recursive: true, mode: 0o700 }); + const hash = createHash('sha256').update(base64Payload, 'utf8').digest('hex'); + const blobPath = join(blobsDir, hash); + try { + const fh = await open(blobPath, 'wx'); + try { + await fh.writeFile(base64Payload, 'utf8'); + await fh.sync(); + } finally { + await fh.close(); + } + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code !== 'EEXIST') throw error; + } + return `${BLOBREF_PROTOCOL}${mimeType};${hash}`; +} + function asMediaContainer(value: unknown): { url: unknown } | undefined { if (value === null || typeof value !== 'object' || Array.isArray(value)) { return undefined; diff --git a/packages/agent-core/src/agent/records/index.ts b/packages/agent-core/src/agent/records/index.ts index 69d82ec8b7..a860acaef4 100644 --- a/packages/agent-core/src/agent/records/index.ts +++ b/packages/agent-core/src/agent/records/index.ts @@ -16,13 +16,8 @@ export { InMemoryAgentRecordPersistence, } from './persistence'; export type { FileSystemAgentRecordPersistenceOptions } from './persistence'; -export { - isBlobRef, - offloadRecordBlobs, - rehydrateRecordBlobs, - rehydrateMessagesBlobs, -} from './blobref'; -export type { BlobOffloadOptions } from './blobref'; +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 diff --git a/packages/agent-core/src/agent/records/persistence.ts b/packages/agent-core/src/agent/records/persistence.ts index 48797460d9..45ab603fff 100644 --- a/packages/agent-core/src/agent/records/persistence.ts +++ b/packages/agent-core/src/agent/records/persistence.ts @@ -3,13 +3,12 @@ import { mkdir, open } from 'node:fs/promises'; import { dirname } from 'pathe'; import { syncDir } from '../../utils/fs'; -import { offloadRecordBlobs } from './blobref'; +import type { BlobStore } from './blobref'; import { type AgentRecord, type AgentRecordPersistence } from './types'; export interface FileSystemAgentRecordPersistenceOptions { readonly onError?: ((error: unknown) => void) | undefined; - readonly blobsDir?: string | undefined; - readonly blobThreshold?: number | undefined; + readonly blobStore?: BlobStore | undefined; } export interface InMemoryAgentRecordPersistenceOptions { @@ -172,14 +171,9 @@ export class FileSystemAgentRecordPersistence implements AgentRecordPersistence const batch = this.pendingRecords.splice(0); this.shouldClear = false; - if (this.options.blobsDir !== undefined) { + if (this.options.blobStore !== undefined) { await Promise.all( - batch.map((record) => - offloadRecordBlobs(record, { - blobsDir: this.options.blobsDir!, - threshold: this.options.blobThreshold, - }), - ), + batch.map((record) => this.options.blobStore!.offload(record)), ); } diff --git a/packages/agent-core/src/agent/turn/kosong-llm.ts b/packages/agent-core/src/agent/turn/kosong-llm.ts index a467769c2f..eb9cc80321 100644 --- a/packages/agent-core/src/agent/turn/kosong-llm.ts +++ b/packages/agent-core/src/agent/turn/kosong-llm.ts @@ -34,7 +34,7 @@ import type { LLMRequestLogContext, LLMStreamTiming, } from '../../loop'; -import { rehydrateMessagesBlobs } from '../records/blobref'; +import type { BlobStore } from '../records/blobref'; import { applyCompletionBudget, type CompletionBudgetConfig, @@ -64,12 +64,7 @@ export interface KosongLLMConfig { * final cap is applied to each request. */ readonly completionBudgetConfig?: CompletionBudgetConfig | undefined; - /** - * Directory holding offloaded base64 blobs. When set, media URLs in - * messages are rehydrated from `blobref://` back to `data:` URIs before - * being sent to the LLM. - */ - readonly blobsDir?: string | undefined; + readonly blobStore?: BlobStore | undefined; } export class KosongLLM implements LLM { @@ -80,7 +75,7 @@ export class KosongLLM implements LLM { private readonly provider: ChatProvider; private readonly generate: GenerateFn; private readonly completionBudgetConfig: CompletionBudgetConfig | undefined; - private readonly blobsDir: string | undefined; + private readonly blobStore: BlobStore | undefined; constructor(config: KosongLLMConfig) { this.provider = config.provider; @@ -89,13 +84,13 @@ export class KosongLLM implements LLM { this.capability = config.capability; this.generate = config.generate ?? kosongGenerate; this.completionBudgetConfig = config.completionBudgetConfig; - this.blobsDir = config.blobsDir; + this.blobStore = config.blobStore; } async chat(params: LLMChatParams): Promise { let messages = params.messages; - if (this.blobsDir !== undefined) { - messages = await rehydrateMessagesBlobs(params.messages, this.blobsDir); + if (this.blobStore !== undefined) { + messages = await this.blobStore.rehydrateMessages(params.messages); } let requestStartedAt = Date.now(); diff --git a/packages/agent-core/test/agent/records/blobref.test.ts b/packages/agent-core/test/agent/records/blobref.test.ts index 5bf02454ed..9cc8557b82 100644 --- a/packages/agent-core/test/agent/records/blobref.test.ts +++ b/packages/agent-core/test/agent/records/blobref.test.ts @@ -5,12 +5,7 @@ import { join } from 'pathe'; import { afterEach, describe, expect, it } from 'vitest'; -import { - isBlobRef, - offloadRecordBlobs, - rehydrateRecordBlobs, - rehydrateMessagesBlobs, -} from '../../../src/agent/records/blobref'; +import { BlobStore, isBlobRef } from '../../../src/agent/records/blobref'; import type { AgentRecord } from '../../../src/agent/records'; const cleanups: string[] = []; @@ -21,16 +16,16 @@ afterEach(async () => { } }); -async function makeBlobsDir(): Promise { - const dir = join(tmpdir(), `blobref-test-${randomBytes(6).toString('hex')}`); - await mkdir(dir, { recursive: true }); - cleanups.push(dir); - return dir; +async function makeStore(): 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: 4096 }), blobsDir }; } describe('blobref', () => { it('offloads large data URIs and replaces with blobref', async () => { - const blobsDir = await makeBlobsDir(); + const { store, blobsDir } = await makeStore(); const payload = 'A'.repeat(5000); const dataUri = `data:image/png;base64,${payload}`; @@ -40,7 +35,7 @@ describe('blobref', () => { origin: { kind: 'user' }, }; - await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + await store.offload(record); const url = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; expect(isBlobRef(url)).toBe(true); @@ -53,7 +48,7 @@ describe('blobref', () => { }); it('skips small data URIs below threshold', async () => { - const blobsDir = await makeBlobsDir(); + const { store, blobsDir } = await makeStore(); const payload = 'short'; const dataUri = `data:image/png;base64,${payload}`; @@ -63,7 +58,7 @@ describe('blobref', () => { origin: { kind: 'user' }, }; - await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + await store.offload(record); const url = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; expect(url).toBe(dataUri); @@ -72,21 +67,21 @@ describe('blobref', () => { }); it('skips existing blobrefs during offload', async () => { - const blobsDir = await makeBlobsDir(); + const { store } = await makeStore(); const record: AgentRecord = { type: 'turn.prompt', input: [{ type: 'image_url', imageUrl: { url: 'blobref:image/png;abc' } }], origin: { kind: 'user' }, }; - await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + await store.offload(record); const url = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; expect(url).toBe('blobref:image/png;abc'); }); it('rehydrates blobrefs back to data URIs', async () => { - const blobsDir = await makeBlobsDir(); + const { store } = await makeStore(); const payload = 'B'.repeat(5000); const dataUri = `data:image/jpeg;base64,${payload}`; @@ -96,29 +91,29 @@ describe('blobref', () => { origin: { kind: 'user' }, }; - await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); - await rehydrateRecordBlobs(record, blobsDir); + await store.offload(record); + await store.rehydrate(record); const url = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; expect(url).toBe(dataUri); }); it('replaces missing blobs with placeholder text', async () => { - const blobsDir = await makeBlobsDir(); + const { store } = await makeStore(); const record: AgentRecord = { type: 'turn.prompt', input: [{ type: 'image_url', imageUrl: { url: 'blobref:image/png;deadbeef' } }], origin: { kind: 'user' }, }; - await rehydrateRecordBlobs(record, blobsDir); + 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 blobsDir = await makeBlobsDir(); + const { store, blobsDir } = await makeStore(); const payload = 'C'.repeat(5000); const dataUri = `data:image/png;base64,${payload}`; @@ -133,15 +128,15 @@ describe('blobref', () => { origin: { kind: 'user' }, }; - await offloadRecordBlobs(record1, { blobsDir, threshold: 4096 }); - await offloadRecordBlobs(record2, { blobsDir, threshold: 4096 }); + await store.offload(record1); + await store.offload(record2); const files = await readdir(blobsDir); expect(files).toHaveLength(1); }); it('rehydrates messages with media parts only', async () => { - const blobsDir = await makeBlobsDir(); + const { store } = await makeStore(); const payload = 'D'.repeat(5000); const dataUri = `data:audio/wav;base64,${payload}`; @@ -154,7 +149,7 @@ describe('blobref', () => { origin: { kind: 'user' }, }; - await offloadRecordBlobs(record, { blobsDir, threshold: 4096 }); + await store.offload(record); const messages = [ { @@ -164,14 +159,14 @@ describe('blobref', () => { }, ]; - const hydrated = await rehydrateMessagesBlobs(messages, blobsDir); + const hydrated = await store.rehydrateMessages(messages); const firstMsg = hydrated[0]!; expect(firstMsg.content[0]).toEqual({ type: 'text', text: 'hello' }); expect((firstMsg.content[1]! as { audioUrl: { url: string } }).audioUrl.url).toBe(dataUri); }); it('degrades missing blob in messages to text placeholder', async () => { - const blobsDir = await makeBlobsDir(); + const { store } = await makeStore(); const messages = [ { role: 'user' as const, @@ -180,7 +175,7 @@ describe('blobref', () => { }, ]; - const hydrated = await rehydrateMessagesBlobs(messages, blobsDir); + const hydrated = await store.rehydrateMessages(messages); const firstMsg = hydrated[0]!; expect(firstMsg.content[0]).toEqual({ type: 'text', text: '[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 b862f5cede..8337eb9176 100644 --- a/packages/agent-core/test/agent/records/persistence.test.ts +++ b/packages/agent-core/test/agent/records/persistence.test.ts @@ -7,6 +7,7 @@ import { afterEach, describe, expect, it } from 'vitest'; import { AGENT_WIRE_PROTOCOL_VERSION, + BlobStore, FileSystemAgentRecordPersistence, InMemoryAgentRecordPersistence, type AgentRecord, @@ -225,7 +226,9 @@ describe('FileSystemAgentRecordPersistence', () => { const wirePath = join(dir, 'wire.jsonl'); const blobsDir = join(dir, 'blobs'); - const persistence = new FileSystemAgentRecordPersistence(wirePath, { blobsDir }); + const persistence = new FileSystemAgentRecordPersistence(wirePath, { + blobStore: new BlobStore({ blobsDir }), + }); const payload = 'X'.repeat(5000); const dataUri = `data:image/png;base64,${payload}`; From 547b9dcbea98ef65144da6a382302f109e3963ad Mon Sep 17 00:00:00 2001 From: _Kerman Date: Wed, 27 May 2026 20:39:16 +0800 Subject: [PATCH 03/16] add comment --- packages/agent-core/src/agent/records/blobref.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/packages/agent-core/src/agent/records/blobref.ts b/packages/agent-core/src/agent/records/blobref.ts index 2a6190f6ec..034dd6290c 100644 --- a/packages/agent-core/src/agent/records/blobref.ts +++ b/packages/agent-core/src/agent/records/blobref.ts @@ -204,6 +204,7 @@ async function writeBlob( } } catch (error) { const code = (error as NodeJS.ErrnoException).code; + // EEXIST means the identical payload was already written; deduplication. if (code !== 'EEXIST') throw error; } return `${BLOBREF_PROTOCOL}${mimeType};${hash}`; From cfa00dc5ad3fac7046570e98f93d1efec4b2e8fd Mon Sep 17 00:00:00 2001 From: _Kerman Date: Wed, 27 May 2026 21:27:07 +0800 Subject: [PATCH 04/16] fix --- .../agent-core/src/agent/records/blobref.ts | 188 +++++++++--------- .../test/agent/records/blobref.test.ts | 53 +++++ 2 files changed, 147 insertions(+), 94 deletions(-) diff --git a/packages/agent-core/src/agent/records/blobref.ts b/packages/agent-core/src/agent/records/blobref.ts index 034dd6290c..8b568d990e 100644 --- a/packages/agent-core/src/agent/records/blobref.ts +++ b/packages/agent-core/src/agent/records/blobref.ts @@ -21,6 +21,7 @@ export interface BlobStoreOptions { export class BlobStore { private readonly blobsDir: string; private readonly threshold: number; + private readonly cache = new Map(); constructor(options: BlobStoreOptions) { this.blobsDir = options.blobsDir; @@ -32,19 +33,19 @@ export class BlobStore { case 'turn.prompt': case 'turn.steer': for (const part of record.input) { - await offloadContentPart(part, this.blobsDir, this.threshold); + await this.offloadContentPart(part); } break; case 'context.append_message': for (const part of record.message.content) { - await offloadContentPart(part, this.blobsDir, this.threshold); + await this.offloadContentPart(part); } break; case 'context.append_loop_event': { const event = record.event; if (event.type === 'tool.result' && typeof event.result.output !== 'string') { for (const part of event.result.output) { - await offloadContentPart(part, this.blobsDir, this.threshold); + await this.offloadContentPart(part); } } break; @@ -59,19 +60,19 @@ export class BlobStore { case 'turn.prompt': case 'turn.steer': for (const part of record.input) { - await rehydrateContentPart(part, this.blobsDir); + await this.rehydrateContentPart(part); } break; case 'context.append_message': for (const part of record.message.content) { - await rehydrateContentPart(part, this.blobsDir); + await this.rehydrateContentPart(part); } break; case 'context.append_loop_event': { const event = record.event; if (event.type === 'tool.result' && typeof event.result.output !== 'string') { for (const part of event.result.output) { - await rehydrateContentPart(part, this.blobsDir); + await this.rehydrateContentPart(part); } } break; @@ -87,7 +88,7 @@ export class BlobStore { if (msg.content.length === 0) return msg; const clone = structuredClone(msg) as Message; for (const part of clone.content) { - await rehydrateContentPart(part, this.blobsDir); + await this.rehydrateContentPart(part); } return { role: clone.role, @@ -100,114 +101,113 @@ export class BlobStore { }), ); } -} -async function offloadContentPart( - part: ContentPart, - blobsDir: string, - threshold: number, -): Promise { - const record = part as unknown as Record; - for (const value of Object.values(record)) { - const mediaObj = asMediaContainer(value); - if (mediaObj === undefined) continue; + private async offloadContentPart(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') continue; + const url = mediaObj.url; + if (typeof url !== 'string') continue; - const newUrl = await maybeOffloadString(url, blobsDir, threshold); - if (newUrl !== url) { - mediaObj.url = newUrl; + const newUrl = await this.maybeOffloadString(url); + if (newUrl !== url) { + mediaObj.url = newUrl; + } } } -} -async function rehydrateContentPart(part: ContentPart, blobsDir: string): Promise { - const record = part as unknown as Record; - for (const value of Object.values(record)) { - const mediaObj = asMediaContainer(value); - if (mediaObj === undefined) continue; + 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 url = mediaObj.url; + if (typeof url !== 'string' || !isBlobRef(url)) continue; - const newUrl = await rehydrateBlobRefUrl(url, blobsDir); - mediaObj.url = newUrl ?? MISSING_MEDIA_PLACEHOLDER; + const newUrl = await this.rehydrateBlobRefUrl(url); + mediaObj.url = newUrl ?? MISSING_MEDIA_PLACEHOLDER; + } } -} -function downgradeMissingMedia(part: ContentPart): ContentPart { - const record = part as unknown as Record; - for (const value of Object.values(record)) { - const mediaObj = asMediaContainer(value); - if (mediaObj === undefined) continue; - if (mediaObj.url === MISSING_MEDIA_PLACEHOLDER) { - return { type: 'text', text: 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}`; } - return part; -} -async function maybeOffloadString( - value: string, - blobsDir: string, - threshold: number, -): 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 < threshold) { - return value; + private async readBlob(hash: string): Promise { + const cached = this.cache.get(hash); + if (cached !== undefined) return cached; + const payload = await readFile(join(this.blobsDir, hash), 'utf8').catch(() => undefined); + if (payload !== undefined) { + this.cache.set(hash, payload); + } + return payload; } - return writeBlob(blobsDir, mimeType, payload); -} -async function rehydrateBlobRefUrl(url: string, blobsDir: 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; + 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); } - const payload = await readFile(join(blobsDir, hash), 'utf8').catch(() => undefined); - if (payload === undefined) { - return undefined; + + 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); + try { + const fh = await open(blobPath, 'wx'); + try { + await fh.writeFile(base64Payload, 'utf8'); + 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.cache.set(hash, base64Payload); + return `${BLOBREF_PROTOCOL}${mimeType};${hash}`; } - return `data:${mimeType};base64,${payload}`; } -async function writeBlob( - blobsDir: string, - mimeType: string, - base64Payload: string, -): Promise { - await mkdir(blobsDir, { recursive: true, mode: 0o700 }); - const hash = createHash('sha256').update(base64Payload, 'utf8').digest('hex'); - const blobPath = join(blobsDir, hash); - try { - const fh = await open(blobPath, 'wx'); - try { - await fh.writeFile(base64Payload, 'utf8'); - await fh.sync(); - } finally { - await fh.close(); +function downgradeMissingMedia(part: ContentPart): ContentPart { + const record = part as unknown as Record; + for (const value of Object.values(record)) { + const mediaObj = asMediaContainer(value); + if (mediaObj === undefined) continue; + if (mediaObj.url === MISSING_MEDIA_PLACEHOLDER) { + return { type: 'text', text: MISSING_MEDIA_PLACEHOLDER }; } - } catch (error) { - const code = (error as NodeJS.ErrnoException).code; - // EEXIST means the identical payload was already written; deduplication. - if (code !== 'EEXIST') throw error; } - return `${BLOBREF_PROTOCOL}${mimeType};${hash}`; + return part; } function asMediaContainer(value: unknown): { url: unknown } | undefined { diff --git a/packages/agent-core/test/agent/records/blobref.test.ts b/packages/agent-core/test/agent/records/blobref.test.ts index 9cc8557b82..7ae2bec4e9 100644 --- a/packages/agent-core/test/agent/records/blobref.test.ts +++ b/packages/agent-core/test/agent/records/blobref.test.ts @@ -179,4 +179,57 @@ describe('blobref', () => { const firstMsg = hydrated[0]!; expect(firstMsg.content[0]).toEqual({ type: 'text', text: '[media missing]' }); }); + + 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' }, + }; + + 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(record); + const url = (record.input as unknown as [{ imageUrl: { url: string } }])[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' }, + }; + + await store.offload(record); + await store.rehydrate(record); + + 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. + const blobrefUrl = (record.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; + const record2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: blobrefUrl } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(record2); + const url2 = (record2.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; + expect(url2).toBe(dataUri); + }); }); From b4d1dfc6075201aefa42b57c1a055b360c78a5ae Mon Sep 17 00:00:00 2001 From: _Kerman Date: Wed, 27 May 2026 21:37:15 +0800 Subject: [PATCH 05/16] lru cache --- .changeset/wire-blob-read-cache.md | 6 ++ .../agent-core/src/agent/records/blobref.ts | 44 +++++++++- .../test/agent/records/blobref.test.ts | 83 ++++++++++++++++++- 3 files changed, 128 insertions(+), 5 deletions(-) create mode 100644 .changeset/wire-blob-read-cache.md diff --git a/.changeset/wire-blob-read-cache.md b/.changeset/wire-blob-read-cache.md new file mode 100644 index 0000000000..058197d754 --- /dev/null +++ b/.changeset/wire-blob-read-cache.md @@ -0,0 +1,6 @@ +--- +"@moonshot-ai/agent-core": patch +"@moonshot-ai/kimi-code": patch +--- + +Add an in-memory read-through cache to `BlobStore` so repeated rehydration avoids redundant disk reads. diff --git a/packages/agent-core/src/agent/records/blobref.ts b/packages/agent-core/src/agent/records/blobref.ts index 8b568d990e..c0ead47969 100644 --- a/packages/agent-core/src/agent/records/blobref.ts +++ b/packages/agent-core/src/agent/records/blobref.ts @@ -5,6 +5,7 @@ import type { ContentPart, Message } from '@moonshot-ai/kosong'; import type { AgentRecord } from './types'; const DEFAULT_THRESHOLD = 4096; +const DEFAULT_MAX_CACHE_SIZE = 500 * 1024 * 1024; const BLOBREF_PROTOCOL = 'blobref:'; const DATA_URI_HEADER_RE = /^data:([^;]+);base64,/; const MISSING_MEDIA_PLACEHOLDER = '[media missing]'; @@ -16,16 +17,21 @@ export function isBlobRef(url: string): boolean { 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 { @@ -152,10 +158,15 @@ export class BlobStore { private async readBlob(hash: string): Promise { const cached = this.cache.get(hash); - if (cached !== undefined) return cached; + 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), 'utf8').catch(() => undefined); if (payload !== undefined) { - this.cache.set(hash, payload); + this.setCache(hash, payload); } return payload; } @@ -193,9 +204,36 @@ export class BlobStore { // EEXIST means the identical payload was already written; deduplication. if (code !== 'EEXIST') throw error; } - this.cache.set(hash, base64Payload); + this.setCache(hash, base64Payload); return `${BLOBREF_PROTOCOL}${mimeType};${hash}`; } + + private setCache(hash: string, payload: string): void { + const size = Buffer.byteLength(payload, 'utf8'); + 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 { + 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 downgradeMissingMedia(part: ContentPart): ContentPart { diff --git a/packages/agent-core/test/agent/records/blobref.test.ts b/packages/agent-core/test/agent/records/blobref.test.ts index 7ae2bec4e9..31d1f0f7e9 100644 --- a/packages/agent-core/test/agent/records/blobref.test.ts +++ b/packages/agent-core/test/agent/records/blobref.test.ts @@ -16,11 +16,18 @@ afterEach(async () => { } }); -async function makeStore(): Promise<{ store: BlobStore; blobsDir: string }> { +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: 4096 }), blobsDir }; + return { + store: new BlobStore({ + blobsDir, + threshold: options?.threshold ?? 4096, + maxCacheSize: options?.maxCacheSize, + }), + blobsDir, + }; } describe('blobref', () => { @@ -232,4 +239,76 @@ describe('blobref', () => { const url2 = (record2.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url; expect(url2).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); // 4 bytes + const payloadB = 'B'.repeat(4); // 4 bytes + const payloadC = 'C'.repeat(4); // 4 bytes + + 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' }, + }; + + await store.offload(recordA); + 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: (recordA.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordA_touch); + + // Adding C should evict B (the least-recently-used), not A. + 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: (recordA.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordA2); + expect((recordA2.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url).toBe(`data:image/png;base64,${payloadA}`); + + // B should have been evicted. + const recordB2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: (recordB.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordB2); + expect((recordB2.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url).toBe('[media missing]'); + + // C should still be cached. + const recordC2: AgentRecord = { + type: 'turn.prompt', + input: [{ type: 'image_url', imageUrl: { url: (recordC.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url } }], + origin: { kind: 'user' }, + }; + await store.rehydrate(recordC2); + expect((recordC2.input as unknown as [{ imageUrl: { url: string } }])[0].imageUrl.url).toBe(`data:image/png;base64,${payloadC}`); + }); }); From e232cec914719d41cb3969b4dc3ddfe6390be7c7 Mon Sep 17 00:00:00 2001 From: _Kerman Date: Wed, 27 May 2026 21:38:08 +0800 Subject: [PATCH 06/16] 50M cache --- packages/agent-core/src/agent/records/blobref.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/agent-core/src/agent/records/blobref.ts b/packages/agent-core/src/agent/records/blobref.ts index c0ead47969..3db610953e 100644 --- a/packages/agent-core/src/agent/records/blobref.ts +++ b/packages/agent-core/src/agent/records/blobref.ts @@ -5,7 +5,7 @@ import type { ContentPart, Message } from '@moonshot-ai/kosong'; import type { AgentRecord } from './types'; const DEFAULT_THRESHOLD = 4096; -const DEFAULT_MAX_CACHE_SIZE = 500 * 1024 * 1024; +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]'; From 113e2502c21ce8e378c5034c1b72f4e03ee352f3 Mon Sep 17 00:00:00 2001 From: _Kerman Date: Wed, 27 May 2026 22:23:21 +0800 Subject: [PATCH 07/16] fix --- .../agent-core/src/agent/records/blobref.ts | 18 +++++++++--------- packages/agent-core/src/agent/records/index.ts | 5 +++++ 2 files changed, 14 insertions(+), 9 deletions(-) diff --git a/packages/agent-core/src/agent/records/blobref.ts b/packages/agent-core/src/agent/records/blobref.ts index 3db610953e..f1e7f8138f 100644 --- a/packages/agent-core/src/agent/records/blobref.ts +++ b/packages/agent-core/src/agent/records/blobref.ts @@ -65,21 +65,15 @@ export class BlobStore { switch (record.type) { case 'turn.prompt': case 'turn.steer': - for (const part of record.input) { - await this.rehydrateContentPart(part); - } + await this.rehydrateParts(record.input); break; case 'context.append_message': - for (const part of record.message.content) { - await this.rehydrateContentPart(part); - } + 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') { - for (const part of event.result.output) { - await this.rehydrateContentPart(part); - } + await this.rehydrateParts(event.result.output); } break; } @@ -88,6 +82,12 @@ export class BlobStore { } } + async rehydrateParts(parts: readonly ContentPart[]): Promise { + for (const part of parts) { + await this.rehydrateContentPart(part); + } + } + async rehydrateMessages(messages: readonly Message[]): Promise { return Promise.all( messages.map(async (msg): Promise => { diff --git a/packages/agent-core/src/agent/records/index.ts b/packages/agent-core/src/agent/records/index.ts index a860acaef4..06ec2de66b 100644 --- a/packages/agent-core/src/agent/records/index.ts +++ b/packages/agent-core/src/agent/records/index.ts @@ -177,6 +177,11 @@ export class AgentRecords { replayedRecords.push(migratedRecord); this.restore(migratedRecord); } + if (this.agent.blobStore !== undefined) { + for (const msg of this.agent.context.history) { + await this.agent.blobStore.rehydrateParts(msg.content); + } + } if (shouldRewrite) { this.persistence.rewrite(replayedRecords); await this.persistence.flush(); From 7639023b2a1a0ce60838d49e541f15f64e75a18c Mon Sep 17 00:00:00 2001 From: _Kerman Date: Wed, 27 May 2026 22:29:34 +0800 Subject: [PATCH 08/16] bump wire --- .../src/agent/records/migration/index.ts | 9 +++++++-- .../src/agent/records/migration/v1.3.ts | 18 ++++++++++++++++++ packages/agent-core/test/agent/turn.test.ts | 2 +- 3 files changed, 26 insertions(+), 3 deletions(-) create mode 100644 packages/agent-core/src/agent/records/migration/v1.3.ts 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/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": "