import * as fs from "node:fs"; import * as fsp from "node:fs/promises"; import * as path from "node:path"; import { isEnoent, peekFileEnds, toError } from "@oh-my-pi/pi-utils"; const utf8Decoder = new TextDecoder("utf-8"); export interface SessionStorageStat { size: number; mtimeMs: number; mtime: Date; } export interface SessionStorageWriter { writeLine(line: string): Promise; /** * Synchronously append a single line. Returns once the bytes are handed to the kernel * (page cache), so the data survives a non-graceful process death (OOM, SIGKILL, etc.) * even though it has not yet been fsynced to the underlying disk. * * `line` MUST already include the trailing newline. Throws synchronously on I/O error. */ writeLineSync(line: string): void; flush(): Promise; fsync(): Promise; close(): Promise; getError(): Error | undefined; } export interface SessionStorage { ensureDirSync(dir: string): void; existsSync(path: string): boolean; writeTextSync(path: string, content: string): void; statSync(path: string): SessionStorageStat; listFilesSync(dir: string, pattern: string): string[]; exists(path: string): Promise; readText(path: string): Promise; /** Read the requested UTF-8 byte windows from the head and tail of the file. */ readTextSlices(path: string, prefixBytes: number, suffixBytes: number): Promise<[string, string]>; writeText(path: string, content: string): Promise; rename(path: string, nextPath: string): Promise; unlink(path: string): Promise; deleteSessionWithArtifacts(sessionPath: string): Promise; openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter; } // FinalizationRegistry to clean up leaked file descriptors const writerRegistry = new FinalizationRegistry(fd => { try { fs.closeSync(fd); } catch { // Ignore - fd may already be closed or invalid } }); class FileSessionStorageWriter implements SessionStorageWriter { #fd: number; #closed = false; #error: Error | undefined; #onError: ((err: Error) => void) | undefined; constructor(fpath: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }) { this.#onError = options?.onError; const flags = options?.flags ?? "a"; // Ensure parent directory exists const dir = path.dirname(fpath); if (!fs.existsSync(dir)) { fs.mkdirSync(dir, { recursive: true }); } // Open file once, keep fd for lifetime this.#fd = fs.openSync(fpath, flags === "w" ? "w" : "a"); // Register for cleanup if abandoned without close() writerRegistry.register(this, this.#fd, this); } #recordError(err: unknown): Error { const error = toError(err); if (!this.#error) this.#error = error; this.#onError?.(error); return error; } writeLineSync(line: string): void { if (this.#closed) throw new Error("Writer closed"); if (this.#error) throw this.#error; try { const buf = Buffer.from(line, "utf-8"); let offset = 0; while (offset < buf.length) { const written = fs.writeSync(this.#fd, buf, offset, buf.length - offset); if (written === 0) { throw new Error("Short write"); } offset += written; } } catch (err) { throw this.#recordError(err); } } async writeLine(line: string): Promise { this.writeLineSync(line); } async flush(): Promise { if (this.#error) throw this.#error; // OS buffers are flushed on fsync, nothing to do here } async fsync(): Promise { if (this.#closed) throw new Error("Writer closed"); if (this.#error) throw this.#error; try { fs.fsyncSync(this.#fd); } catch (err) { throw this.#recordError(err); } } async close(): Promise { if (this.#closed) return; this.#closed = true; // Unregister from finalization - we're closing properly writerRegistry.unregister(this); try { fs.closeSync(this.#fd); } catch { // Ignore close errors } } getError(): Error | undefined { return this.#error; } } export class FileSessionStorage implements SessionStorage { ensureDirSync(dir: string): void { if (!fs.existsSync(dir)) { fs.mkdirSync(dir, { recursive: true }); } } existsSync(path: string): boolean { return fs.existsSync(path); } writeTextSync(fpath: string, content: string): void { this.ensureDirSync(path.dirname(fpath)); fs.writeFileSync(fpath, content); } statSync(path: string): SessionStorageStat { const stats = fs.statSync(path); return { size: stats.size, mtimeMs: stats.mtimeMs, mtime: stats.mtime }; } listFilesSync(dir: string, pattern: string): string[] { try { return Array.from(new Bun.Glob(pattern).scanSync(dir)).map(name => path.join(dir, name)); } catch { return []; } } async exists(path: string): Promise { try { await fs.promises.access(path); return true; } catch (err) { if (isEnoent(err)) return false; throw err; } } readText(path: string): Promise { return Bun.file(path).text(); } async readTextSlices(path: string, prefixBytes: number, suffixBytes: number): Promise<[string, string]> { return peekFileEnds(path, prefixBytes, suffixBytes, (head, tail) => [ utf8Decoder.decode(head), utf8Decoder.decode(tail), ]); } async writeText(path: string, content: string): Promise { await Bun.write(path, content, { createPath: true }); } async rename(path: string, nextPath: string): Promise { try { await fs.promises.rename(path, nextPath); } catch (err) { throw toError(err); } } unlink(path: string): Promise { return fs.promises.unlink(path); } openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter { return new FileSessionStorageWriter(path, options); } /** * Delete a session file and its artifacts directory. * Artifacts are stored in a sibling directory with the same name minus .jsonl extension. */ async deleteSessionWithArtifacts(sessionPath: string): Promise { // Delete the session file itself await this.unlink(sessionPath); // Compute artifacts directory: /path/to/session.jsonl -> /path/to/session const artifactsDir = sessionPath.slice(0, -6); // Delete artifacts directory if it exists. Missing directories are fine, but // surface real cleanup failures because the session file is already gone. try { await fsp.rm(artifactsDir, { recursive: true, force: true }); } catch (err) { const error = toError(err); throw new Error( `Session file deleted but failed to remove artifacts directory ${artifactsDir}: ${error.message}`, { cause: error, }, ); } } } function matchesPattern(name: string, pattern: string): boolean { if (pattern === "*") return true; if (pattern.startsWith("*.")) { return name.endsWith(pattern.slice(1)); } return name === pattern; } class MemorySessionStorageWriter implements SessionStorageWriter { #storage: MemorySessionStorage; #path: string; #closed = false; #error: Error | undefined; #onError: ((err: Error) => void) | undefined; constructor( storage: MemorySessionStorage, path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }, ) { this.#storage = storage; this.#path = path; this.#onError = options?.onError; if ((options?.flags ?? "a") === "w") { this.#storage.writeTextSync(path, ""); } } #recordError(err: unknown): Error { const error = toError(err); if (!this.#error) this.#error = error; this.#onError?.(error); return error; } writeLineSync(line: string): void { if (this.#closed) throw new Error("Writer closed"); if (this.#error) throw this.#error; try { // O(1) append — push onto the path's indexed in-memory entry. this.#storage.appendSync(this.#path, line); } catch (err) { throw this.#recordError(err); } } async writeLine(line: string): Promise { this.writeLineSync(line); } async flush(): Promise { if (this.#error) throw this.#error; } async fsync(): Promise { // No-op for in-memory storage if (this.#error) throw this.#error; } async close(): Promise { if (this.#closed) return; this.#closed = true; } getError(): Error | undefined { return this.#error; } } interface MemoryFileEntry { chunks: string[]; cumulativeBytes: number[]; size: number; mtimeMs: number; } function createMemoryFileEntry(content: string, mtimeMs: number): MemoryFileEntry { const size = Buffer.byteLength(content, "utf-8"); return { chunks: size === 0 ? [] : [content], cumulativeBytes: size === 0 ? [] : [size], size, mtimeMs, }; } function appendMemoryChunk(entry: MemoryFileEntry, chunk: string): void { const chunkSize = Buffer.byteLength(chunk, "utf-8"); if (chunkSize === 0) return; entry.size += chunkSize; entry.chunks.push(chunk); entry.cumulativeBytes.push(entry.size); } function normalizeByteLimit(maxBytes: number, size: number): number { if (!(maxBytes > 0) || size === 0) return 0; return Math.min(Math.trunc(maxBytes), size); } function lowerBound(values: readonly number[], target: number): number { let lo = 0; let hi = values.length; while (lo < hi) { const mid = (lo + hi) >>> 1; if (values[mid] < target) { lo = mid + 1; } else { hi = mid; } } return lo; } function upperBound(values: readonly number[], target: number): number { let lo = 0; let hi = values.length; while (lo < hi) { const mid = (lo + hi) >>> 1; if (values[mid] <= target) { lo = mid + 1; } else { hi = mid; } } return lo; } function joinChunkRange(chunks: readonly string[], start: number, end: number): string { const count = end - start; if (count <= 0) return ""; if (count === 1) return chunks[start] ?? ""; let content = ""; for (let i = start; i < end; i++) { content += chunks[i]; } return content; } function decodeChunkByteRange(chunk: string, startByte: number, endByte: number, chunkSize: number): string { if (startByte >= endByte) return ""; if (startByte === 0 && endByte === chunkSize) return chunk; if (chunk.length === chunkSize) return chunk.slice(startByte, endByte); const bytes = Buffer.from(chunk, "utf-8"); return utf8Decoder.decode(bytes.subarray(startByte, endByte)); } function materializeMemoryEntry(entry: MemoryFileEntry): string { const { chunks } = entry; if (chunks.length === 0) return ""; if (chunks.length === 1) return chunks[0]; const content = chunks.join(""); entry.chunks = [content]; entry.cumulativeBytes = [entry.size]; return content; } function sliceChunksHead(entry: MemoryFileEntry, maxBytes: number): string { const limit = normalizeByteLimit(maxBytes, entry.size); if (limit === 0) return ""; if (limit >= entry.size) return materializeMemoryEntry(entry); const boundaryIndex = lowerBound(entry.cumulativeBytes, limit); const chunkStart = boundaryIndex === 0 ? 0 : entry.cumulativeBytes[boundaryIndex - 1]; const chunkEnd = entry.cumulativeBytes[boundaryIndex]; if (chunkEnd === limit) return joinChunkRange(entry.chunks, 0, boundaryIndex + 1); const chunk = entry.chunks[boundaryIndex]; const chunkPrefix = decodeChunkByteRange(chunk, 0, limit - chunkStart, chunkEnd - chunkStart); return joinChunkRange(entry.chunks, 0, boundaryIndex) + chunkPrefix; } function sliceChunksTail(entry: MemoryFileEntry, maxBytes: number): string { const limit = normalizeByteLimit(maxBytes, entry.size); if (limit === 0) return ""; if (limit >= entry.size) return materializeMemoryEntry(entry); const startByte = entry.size - limit; const boundaryIndex = upperBound(entry.cumulativeBytes, startByte); const chunkStart = boundaryIndex === 0 ? 0 : entry.cumulativeBytes[boundaryIndex - 1]; const chunkEnd = entry.cumulativeBytes[boundaryIndex]; const chunkOffset = startByte - chunkStart; if (chunkOffset === 0) return joinChunkRange(entry.chunks, boundaryIndex, entry.chunks.length); const chunk = entry.chunks[boundaryIndex]; const chunkSuffix = decodeChunkByteRange(chunk, chunkOffset, chunkEnd - chunkStart, chunkEnd - chunkStart); return chunkSuffix + joinChunkRange(entry.chunks, boundaryIndex + 1, entry.chunks.length); } export class MemorySessionStorage implements SessionStorage { // Each path keeps appended string chunks plus cumulative UTF-8 byte offsets. // Full reads materialize the chunks into one string chunk, so repeated reads // do not re-join stale history. Later appends still stay O(1) by pushing // after that materialized chunk. Prefix/suffix reads binary-search byte // offsets and join only the requested window. #files = new Map(); #requireEntry(path: string): MemoryFileEntry { const entry = this.#files.get(path); if (!entry) throw new Error(`File not found: ${path}`); return entry; } ensureDirSync(_dir: string): void { // No-op for in-memory storage. } existsSync(path: string): boolean { return this.#files.has(path); } writeTextSync(path: string, content: string): void { this.#files.set(path, createMemoryFileEntry(content, Date.now())); } /** * Internal O(1) append used by {@link MemorySessionStorageWriter}. Lazily * creates the entry. External callers should go through `openWriter()` * rather than touching the mirror directly. */ appendSync(path: string, chunk: string): void { const mtimeMs = Date.now(); let entry = this.#files.get(path); if (!entry) { entry = createMemoryFileEntry("", mtimeMs); this.#files.set(path, entry); } appendMemoryChunk(entry, chunk); entry.mtimeMs = mtimeMs; } statSync(path: string): SessionStorageStat { const entry = this.#requireEntry(path); return { size: entry.size, mtimeMs: entry.mtimeMs, mtime: new Date(entry.mtimeMs), }; } listFilesSync(dir: string, pattern: string): string[] { const prefix = dir.endsWith("/") ? dir : `${dir}/`; const files: string[] = []; for (const path of this.#files.keys()) { if (!path.startsWith(prefix)) continue; const name = path.slice(prefix.length); if (name.includes("/") || name.includes("\\")) continue; if (!matchesPattern(name, pattern)) continue; files.push(path); } return files; } exists(path: string): Promise { return Promise.resolve(this.existsSync(path)); } readText(path: string): Promise { const entry = this.#files.get(path); if (!entry) return Promise.reject(new Error(`File not found: ${path}`)); return Promise.resolve(materializeMemoryEntry(entry)); } readTextSlices(path: string, prefixBytes: number, suffixBytes: number): Promise<[string, string]> { const entry = this.#files.get(path); if (!entry) return Promise.reject(new Error(`File not found: ${path}`)); return Promise.resolve([sliceChunksHead(entry, prefixBytes), sliceChunksTail(entry, suffixBytes)]); } writeText(path: string, content: string): Promise { this.writeTextSync(path, content); return Promise.resolve(); } rename(path: string, nextPath: string): Promise { const entry = this.#files.get(path); if (!entry) return Promise.reject(new Error(`File not found: ${path}`)); this.#files.set(nextPath, entry); this.#files.delete(path); return Promise.resolve(); } unlink(path: string): Promise { this.#files.delete(path); return Promise.resolve(); } deleteSessionWithArtifacts(_sessionPath: string): Promise { return Promise.resolve(); } openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter { return new MemorySessionStorageWriter(this, path, options); } }