{"version":3,"file":"storage.d.ts","sourceRoot":"","sources":["../../../../src/harness/session/jsonl/storage.ts"],"names":[],"mappings":"AAEA,OAAO,KAAK,EAAE,OAAO,EAAE,MAAM,kBAAkB,CAAC;AAGhD,OAAO,EAAE,KAAK,uBAAuB,EAAE,KAAK,kBAAkB,EAAsB,MAAM,YAAY,CAAC;AAWvG,OAAO,KAAK,EACX,YAAY,EACZ,KAAK,EACL,SAAS,EACT,cAAc,EACd,YAAY,EACZ,OAAO,EACP,iBAAiB,EACjB,QAAQ,EACR,SAAS,EACT,KAAK,EACL,MAAM,aAAa,CAAC;AACrB,OAAO,KAAK,EAAE,WAAW,EAAE,eAAe,EAAE,WAAW,EAAE,KAAK,EAAE,SAAS,EAAE,MAAM,cAAc,CAAC;AAGhG,OAAO,EAAyB,KAAK,kBAAkB,EAAE,KAAK,mBAAmB,EAAE,MAAM,YAAY,CAAC;AA4FtG,iEAAiE;AACjE,qBAAa,YAAa,YAAW,OAAO;IAC3C,OAAO,CAAC,QAAQ,CAAC,UAAU,CAAa;IACxC,OAAO,CAAC,QAAQ,CAAC,IAAI,CAAS;IAC9B,OAAO,CAAC,QAAQ,CAAC,GAAG,CAAe;IACnC,QAAQ,CAAC,MAAM,EAAE,kBAAkB,CAAC;IACpC,OAAO,CAAC,OAAO,CAAe;IAC9B,OAAO,CAAC,QAAQ,CAAC,YAAY,CAA8B;IAC3D,OAAO,CAAC,WAAW,CAAoC;IACvD,OAAO,CAAC,KAAK,CAAyC;IACtD,OAAO,CAAC,YAAY,CAA4B;IAEhD,OAAO,eAMN;IAED,OAAa,MAAM,CAClB,OAAO,EAAE,mBAAmB,EAC5B,MAAM,EAAE,kBAAkB,EAC1B,aAAa,EAAE,KAAK,EAAE,EACtB,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,YAAY,CAAC,CAOvB;IAED,mEAAmE;IACnE,OAAa,sBAAsB,CAClC,OAAO,EAAE,mBAAmB,EAC5B,MAAM,EAAE,kBAAkB,EAC1B,QAAQ,EAAE,uBAAuB,EACjC,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,YAAY,CAAC,CAavB;IAED,OAAa,IAAI,CAAC,OAAO,EAAE,mBAAmB,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,YAAY,CAAC,CAiCvF;mBAEoB,YAAY;IAoBjC,OAAO,CAAC,eAAe;IAKjB,MAAM,CAAC,MAAM,EAAE,KAAK,EAAE,EAAE,OAAO,EAAE,OAAO,GAAG,OAAO,CAAC,YAAY,CAAC,CAQrE;YAEa,WAAW;YAgBX,mBAAmB;IAuCjC,UAAU,CAAC,GAAG,EAAE,MAAM,EAAE,EAAE,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,GAAG,CAAC,MAAM,EAAE,KAAK,CAAC,CAAC,CAGxE;IAED,QAAQ,CAAC,CAAC,EAAE,OAAO,EAAE,KAAK,CAAC,CAAC,CAAC,EAAE,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,WAAW,CAAC,CAAC,CAAC,GAAG,SAAS,CAAC,CAGrF;IAED,UAAU,CAAC,CAAC,EAAE,MAAM,EAAE,KAAK,CAAC,CAAC,CAAC,EAAE,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,WAAW,CAAC,CAAC,CAAC,EAAE,CAAC,CAG5E;IAEK,QAAQ,CAAC,CAAC,EACf,OAAO,EAAE,SAAS,CAAC,CAAC,CAAC,EACrB,OAAO,EAAE,eAAe,GAAG,SAAS,EACpC,QAAQ,EAAE,OAAO,GACf,OAAO,CAAC,WAAW,CAAC,CAAC,CAAC,EAAE,CAAC,CAG3B;IAEK,UAAU,CAAC,KAAK,EAAE,iBAAiB,EAAE,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,KAAK,EAAE,CAAC,CAG9E;IAEK,mBAAmB,CAAC,KAAK,EAAE,iBAAiB,EAAE,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,cAAc,EAAE,CAAC,CAGhG;IAED,WAAW,CAAC,KAAK,EAAE,SAAS,EAAE,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,KAAK,EAAE,CAAC,CAGjE;IAED,SAAS,CAAC,KAAK,EAAE,SAAS,EAAE,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,QAAQ,EAAE,CAAC,CAGlE;IAED,QAAQ,CAAC,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,YAAY,CAAC,CAGjD;IAED,OAAO,CAAC,iBAAiB;IAIzB,mFAAmF;IACnF,iBAAiB,CAAC,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,kBAAkB,CAAC,CAQhE;IAED,KAAK,CAAC,QAAQ,EAAE,OAAO,GAAG,OAAO,CAAC,IAAI,CAAC,CAOtC;CACD","sourcesContent":["import type { Usage } from \"@earendil-works/pi-ai\";\nimport { uuidv7 } from \"@earendil-works/pi-ai/utils/uuid\";\nimport type { Context } from \"../../context.ts\";\nimport type { FileError, FileSystem, Result } from \"../../types.ts\";\nimport { insertUsage } from \"../commit.ts\";\nimport { type ForkDestinationSnapshot, type ForkSourceSnapshot, forkSnapshotWrites } from \"../fork.ts\";\nimport {\n\ttype CommittedEntryWrite,\n\ttype CommittedListAppendWrite,\n\ttype CommittedListDeleteWrite,\n\ttype CommittedUsageWrite,\n\ttype CommittedValueDeleteWrite,\n\ttype CommittedValueSetWrite,\n\ttype CommittedWrite,\n\tInMemoryStorageState,\n} from \"../in-memory-storage-state.ts\";\nimport type {\n\tCommitResult,\n\tEntry,\n\tEntryScan,\n\tEntryStructure,\n\tSessionStats,\n\tStorage,\n\tStorageBranchScan,\n\tUsageRow,\n\tUsageScan,\n\tWrite,\n} from \"../types.ts\";\nimport type { ListElement, ListReadOptions, StoredValue, Value, ValueList } from \"../values.ts\";\nimport { type LegacyV3SessionHeader, parseJsonlSessionHeader } from \"./codec.ts\";\nimport { normalizeLegacyV3Header, normalizeLegacyV3Records } from \"./legacy-v3.ts\";\nimport { JSONL_STORAGE_VERSION, type JsonlStorageHeader, type JsonlStorageOptions } from \"./types.ts\";\n\nfunction fileValue<T>(result: Result<T, FileError>, action: string): T {\n\tif (!result.ok) throw new Error(`${action}: ${result.error.message}`, { cause: result.error });\n\treturn result.value;\n}\n\nfunction isRecord(value: unknown): value is Record<string, unknown> {\n\treturn typeof value === \"object\" && value !== null && !Array.isArray(value);\n}\n\nfunction requireSafeInteger(value: unknown, field: string, minimum: number): void {\n\tif (!Number.isSafeInteger(value) || (value as number) < minimum) throw new Error(`Invalid JSONL ${field}`);\n}\n\nfunction parseCommittedWrite(value: unknown): CommittedWrite {\n\tif (!isRecord(value)) throw new Error(\"Invalid JSONL transaction write\");\n\trequireSafeInteger(value.seq, \"write seq\", 1);\n\tswitch (value.kind) {\n\t\tcase \"entry\":\n\t\t\trequireSafeInteger(value.timestamp, \"entry timestamp\", 0);\n\t\t\treturn value as unknown as CommittedEntryWrite;\n\t\tcase \"usage\":\n\t\t\treturn value as unknown as CommittedUsageWrite;\n\t\tcase \"value\":\n\t\t\tif (value.op === \"set\") return value as unknown as CommittedValueSetWrite;\n\t\t\tif (value.op === \"delete\") return value as unknown as CommittedValueDeleteWrite;\n\t\t\tthrow new Error(`Invalid JSONL value operation: ${String(value.op)}`);\n\t\tcase \"list\":\n\t\t\tif (value.op === \"append\") return value as unknown as CommittedListAppendWrite;\n\t\t\tif (value.op === \"delete\") return value as unknown as CommittedListDeleteWrite;\n\t\t\tthrow new Error(`Invalid JSONL list operation: ${String(value.op)}`);\n\t\tdefault:\n\t\t\tthrow new Error(`Invalid JSONL write kind: ${String(value.kind)}`);\n\t}\n}\n\nfunction parseTransaction(line: string): CommittedWrite[] {\n\tlet value: unknown;\n\ttry {\n\t\tvalue = JSON.parse(line);\n\t} catch (error) {\n\t\tthrow new Error(\"Invalid JSONL transaction: not valid JSON\", { cause: error });\n\t}\n\treturn (Array.isArray(value) ? value : [value]).map(parseCommittedWrite);\n}\n\nfunction serializeTransaction(writes: CommittedWrite[]): string {\n\treturn JSON.stringify(writes.length === 1 ? writes[0] : writes);\n}\n\nfunction serializeStorage(header: JsonlStorageHeader, transactions: CommittedWrite[][]): string {\n\treturn `${[JSON.stringify(header), ...transactions.map(serializeTransaction)].join(\"\\n\")}\\n`;\n}\n\nfunction splitCompleteLines(content: string): { lines: string[]; torn: boolean } {\n\tif (content.endsWith(\"\\n\")) return { lines: content.slice(0, -1).split(\"\\n\"), torn: false };\n\tconst lastNewline = content.lastIndexOf(\"\\n\");\n\tif (lastNewline === -1) return { lines: [], torn: true };\n\treturn { lines: content.slice(0, lastNewline).split(\"\\n\"), torn: true };\n}\n\nasync function publishFileAtomically(\n\tfileSystem: FileSystem,\n\tdestinationPath: string,\n\tcontent: string,\n\tcontext: Context,\n): Promise<void> {\n\tconst tempPath = `${destinationPath}.tmp`;\n\ttry {\n\t\tfileValue(\n\t\t\tawait fileSystem.writeFile(tempPath, content, context),\n\t\t\t`Failed to stage JSONL storage ${destinationPath}`,\n\t\t);\n\t\tfileValue(\n\t\t\tawait fileSystem.renameFile(tempPath, destinationPath, context),\n\t\t\t`Failed to publish JSONL storage ${destinationPath}`,\n\t\t);\n\t} catch (error) {\n\t\tawait fileSystem.remove(tempPath, { force: true }, context);\n\t\tthrow error;\n\t}\n}\n\ntype LegacyV3Backing = {\n\tkind: \"v3\";\n\timportedUsage: Usage;\n\tbaselineWrites: readonly CommittedWrite[];\n};\n\ntype JsonlBacking = { kind: \"v4\" } | LegacyV3Backing;\n\n/** JSONL storage backed by an injected filesystem capability. */\nexport class JsonlStorage implements Storage {\n\tprivate readonly fileSystem: FileSystem;\n\tprivate readonly path: string;\n\tprivate readonly now: () => number;\n\treadonly header: JsonlStorageHeader;\n\tprivate backing: JsonlBacking;\n\tprivate readonly storageState = new InMemoryStorageState();\n\tprivate commitQueue: Promise<void> = Promise.resolve();\n\tprivate state: \"open\" | \"closing\" | \"closed\" = \"open\";\n\tprivate closePromise: Promise<void> | undefined;\n\n\tprivate constructor(options: JsonlStorageOptions, header: JsonlStorageHeader, backing: JsonlBacking) {\n\t\tthis.fileSystem = options.fileSystem;\n\t\tthis.path = options.path;\n\t\tthis.now = options.now ?? Date.now;\n\t\tthis.header = header;\n\t\tthis.backing = backing;\n\t}\n\n\tstatic async create(\n\t\toptions: JsonlStorageOptions,\n\t\theader: JsonlStorageHeader,\n\t\tinitialWrites: Write[],\n\t\tcontext: Context,\n\t): Promise<JsonlStorage> {\n\t\tconst storage = new JsonlStorage(options, header, { kind: \"v4\" });\n\t\tconst prepared = storage.storageState.prepareCommit(initialWrites, storage.now());\n\t\tconst transactions = prepared.writes.length === 0 ? [] : [prepared.writes];\n\t\tawait publishFileAtomically(options.fileSystem, options.path, serializeStorage(header, transactions), context);\n\t\tstorage.storageState.applyValidated(prepared.writes);\n\t\treturn storage;\n\t}\n\n\t/** Atomically create storage from a complete prepared snapshot. */\n\tstatic async createFromForkSnapshot(\n\t\toptions: JsonlStorageOptions,\n\t\theader: JsonlStorageHeader,\n\t\tsnapshot: ForkDestinationSnapshot,\n\t\tcontext: Context,\n\t): Promise<JsonlStorage> {\n\t\tconst writes = forkSnapshotWrites(snapshot);\n\t\tconst snapshotHeader = { ...header, nextSeq: snapshot.nextSeq };\n\t\tawait publishFileAtomically(\n\t\t\toptions.fileSystem,\n\t\t\toptions.path,\n\t\t\tserializeStorage(\n\t\t\t\tsnapshotHeader,\n\t\t\t\twrites.map((write) => [write]),\n\t\t\t),\n\t\t\tcontext,\n\t\t);\n\t\treturn JsonlStorage.open(options, context);\n\t}\n\n\tstatic async open(options: JsonlStorageOptions, context: Context): Promise<JsonlStorage> {\n\t\tconst content = fileValue(\n\t\t\tawait options.fileSystem.readTextFile(options.path, context),\n\t\t\t`Failed to read JSONL storage ${options.path}`,\n\t\t);\n\t\tconst { lines, torn } = splitCompleteLines(content);\n\t\tif (lines[0] === undefined || lines[0] === \"\") {\n\t\t\tthrow new Error(`Invalid JSONL storage ${options.path}: missing header`);\n\t\t}\n\t\tconst parsedHeader = parseJsonlSessionHeader(lines[0]);\n\t\tif (!parsedHeader.ok) {\n\t\t\tthrow new Error(`Invalid JSONL storage ${options.path}: invalid header`, { cause: parsedHeader.error });\n\t\t}\n\t\tif (parsedHeader.value.format === \"v3-legacy\") {\n\t\t\treturn JsonlStorage.openLegacyV3(options, parsedHeader.value.header, lines.slice(1), context);\n\t\t}\n\n\t\tconst header = parsedHeader.value.header;\n\t\tif (header.storageVersion !== JSONL_STORAGE_VERSION) {\n\t\t\tthrow new Error(`Session ${header.id} uses unsupported storage version ${header.storageVersion}`);\n\t\t}\n\t\tconst storage = new JsonlStorage(options, header, { kind: \"v4\" });\n\t\tfor (let index = 1; index < lines.length; index++) {\n\t\t\tconst line = lines[index]!;\n\t\t\ttry {\n\t\t\t\tstorage.replayCommitted(parseTransaction(line));\n\t\t\t} catch (error) {\n\t\t\t\tthrow new Error(`Invalid JSONL storage ${options.path}: line ${index + 1}`, { cause: error });\n\t\t\t}\n\t\t}\n\t\tif (header.nextSeq !== undefined) storage.storageState.advanceNextSeq(header.nextSeq);\n\t\tif (torn) await publishFileAtomically(options.fileSystem, options.path, `${lines.join(\"\\n\")}\\n`, context);\n\t\treturn storage;\n\t}\n\n\tprivate static async openLegacyV3(\n\t\toptions: JsonlStorageOptions,\n\t\theader: LegacyV3SessionHeader,\n\t\trecordLines: readonly string[],\n\t\tcontext: Context,\n\t): Promise<JsonlStorage> {\n\t\tconst { writes, importedUsage, nextSeq } = normalizeLegacyV3Records(recordLines);\n\t\tconst targetHeader = {\n\t\t\t...(await normalizeLegacyV3Header(options.fileSystem, header, context)),\n\t\t\tnextSeq,\n\t\t};\n\t\tconst storage = new JsonlStorage(options, targetHeader, {\n\t\t\tkind: \"v3\",\n\t\t\timportedUsage,\n\t\t\tbaselineWrites: writes,\n\t\t});\n\t\tstorage.replayCommitted(writes);\n\t\treturn storage;\n\t}\n\n\tprivate replayCommitted(writes: readonly CommittedWrite[]): void {\n\t\tthis.storageState.validateCommitted(writes);\n\t\tthis.storageState.applyValidated(writes);\n\t}\n\n\tasync commit(writes: Write[], context: Context): Promise<CommitResult> {\n\t\tif (this.state !== \"open\") throw new Error(\"JsonlStorage is closed\");\n\t\tconst result = this.commitQueue.then(() => this.applyCommit(writes, context));\n\t\tthis.commitQueue = result.then(\n\t\t\t() => undefined,\n\t\t\t() => undefined,\n\t\t);\n\t\treturn result;\n\t}\n\n\tprivate async applyCommit(writes: Write[], context: Context): Promise<CommitResult> {\n\t\tif (this.backing.kind === \"v3\" && writes.length !== 0) {\n\t\t\treturn this.upgradeLegacyV3ToV4(this.backing, writes, context);\n\t\t}\n\t\tconst prepared = this.storageState.prepareCommit(writes, this.now());\n\t\tif (prepared.writes.length !== 0) {\n\t\t\tfileValue(\n\t\t\t\tawait this.fileSystem.appendFile(this.path, `${serializeTransaction(prepared.writes)}\\n`, context),\n\t\t\t\t`Failed to append JSONL storage ${this.path}`,\n\t\t\t);\n\t\t}\n\t\tconst stats = this.storageState.applyValidated(prepared.writes);\n\t\treturn { ...prepared.result, stats: this.withImportedUsage(stats) };\n\t}\n\n\t/** Atomically upgrade legacy v3 backing and preserve the first caller write as a v4 transaction. */\n\tprivate async upgradeLegacyV3ToV4(\n\t\tbacking: LegacyV3Backing,\n\t\tcallerWrites: Write[],\n\t\tcontext: Context,\n\t): Promise<CommitResult> {\n\t\tconst timestamp = this.now();\n\t\tconst prepared = this.storageState.prepareCommit(\n\t\t\t[\n\t\t\t\tinsertUsage({\n\t\t\t\t\tid: uuidv7(timestamp),\n\t\t\t\t\tusage: backing.importedUsage,\n\t\t\t\t\tadjustment: true,\n\t\t\t\t\tdetails: { source: \"v3-import\" },\n\t\t\t\t}),\n\t\t\t\t...callerWrites,\n\t\t\t],\n\t\t\ttimestamp,\n\t\t);\n\n\t\tconst nextSeq = prepared.result.firstSeq + prepared.writes.length;\n\t\tconst upgradedHeader = { ...this.header, nextSeq };\n\t\tawait publishFileAtomically(\n\t\t\tthis.fileSystem,\n\t\t\tthis.path,\n\t\t\tserializeStorage(upgradedHeader, [...backing.baselineWrites.map((write) => [write]), prepared.writes]),\n\t\t\tcontext,\n\t\t);\n\n\t\tconst stats = this.storageState.applyValidated(prepared.writes);\n\t\tthis.backing = { kind: \"v4\" };\n\t\t// The first sequence belongs to the internal usage adjustment; return only caller-write sequences.\n\t\treturn {\n\t\t\t...prepared.result,\n\t\t\tfirstSeq: prepared.result.firstSeq + 1,\n\t\t\tseqs: prepared.result.seqs.slice(1),\n\t\t\tstats,\n\t\t};\n\t}\n\n\tgetEntries(ids: string[], _context: Context): Promise<Map<string, Entry>> {\n\t\tif (this.state !== \"open\") return Promise.reject(new Error(\"JsonlStorage is closed\"));\n\t\treturn Promise.resolve(this.storageState.getEntries(ids));\n\t}\n\n\tgetValue<T>(address: Value<T>, _context: Context): Promise<StoredValue<T> | undefined> {\n\t\tif (this.state !== \"open\") return Promise.reject(new Error(\"JsonlStorage is closed\"));\n\t\treturn Promise.resolve(this.storageState.getValue(address));\n\t}\n\n\tscanValues<T>(prefix: Value<T>, _context: Context): Promise<StoredValue<T>[]> {\n\t\tif (this.state !== \"open\") return Promise.reject(new Error(\"JsonlStorage is closed\"));\n\t\treturn Promise.resolve(this.storageState.scanValues(prefix));\n\t}\n\n\tasync readList<T>(\n\t\taddress: ValueList<T>,\n\t\toptions: ListReadOptions | undefined,\n\t\t_context: Context,\n\t): Promise<ListElement<T>[]> {\n\t\tif (this.state !== \"open\") throw new Error(\"JsonlStorage is closed\");\n\t\treturn this.storageState.readList(address, options);\n\t}\n\n\tasync scanBranch(query: StorageBranchScan, _context: Context): Promise<Entry[]> {\n\t\tif (this.state !== \"open\") throw new Error(\"JsonlStorage is closed\");\n\t\treturn this.storageState.scanBranch(query);\n\t}\n\n\tasync scanBranchStructure(query: StorageBranchScan, _context: Context): Promise<EntryStructure[]> {\n\t\tif (this.state !== \"open\") throw new Error(\"JsonlStorage is closed\");\n\t\treturn this.storageState.scanBranchStructure(query);\n\t}\n\n\tscanEntries(query: EntryScan, _context: Context): Promise<Entry[]> {\n\t\tif (this.state !== \"open\") return Promise.reject(new Error(\"JsonlStorage is closed\"));\n\t\treturn Promise.resolve(this.storageState.scanEntries(query));\n\t}\n\n\tscanUsage(query: UsageScan, _context: Context): Promise<UsageRow[]> {\n\t\tif (this.state !== \"open\") return Promise.reject(new Error(\"JsonlStorage is closed\"));\n\t\treturn Promise.resolve(this.storageState.scanUsage(query));\n\t}\n\n\tgetStats(_context: Context): Promise<SessionStats> {\n\t\tif (this.state !== \"open\") return Promise.reject(new Error(\"JsonlStorage is closed\"));\n\t\treturn Promise.resolve(this.withImportedUsage(this.storageState.getStats()));\n\t}\n\n\tprivate withImportedUsage(stats: SessionStats): SessionStats {\n\t\treturn this.backing.kind === \"v4\" ? stats : { ...stats, usage: this.backing.importedUsage };\n\t}\n\n\t/** Capture the state needed to fork at one serialized boundary between commits. */\n\tcaptureForkSource(_context: Context): Promise<ForkSourceSnapshot> {\n\t\tif (this.state !== \"open\") return Promise.reject(new Error(\"JsonlStorage is closed\"));\n\t\tconst result = this.commitQueue.then(() => this.storageState.snapshotEntriesAndValues());\n\t\tthis.commitQueue = result.then(\n\t\t\t() => undefined,\n\t\t\t() => undefined,\n\t\t);\n\t\treturn result;\n\t}\n\n\tclose(_context: Context): Promise<void> {\n\t\tif (this.closePromise !== undefined) return this.closePromise;\n\t\tthis.state = \"closing\";\n\t\tthis.closePromise = this.commitQueue.then(() => {\n\t\t\tthis.state = \"closed\";\n\t\t});\n\t\treturn this.closePromise;\n\t}\n}\n"]}