21d152c70c
F6: defer the JSON -> YAML migration out of ConfigFile's constructor and
add an async path so the boot sequence stops blocking the event loop on
sync I/O. New ConfigFile.tryLoadAsync/loadAsync/loadOrDefaultAsync/
getMtimeMsAsync and static ConfigFile.warmup. The migration is now
idempotent (per-process cache) so relocate() does not re-run it.
ModelRegistry.create(authStorage, modelsPath?) is a new async factory
that runs the warmup before the sync constructor's bundled-model load.
Production call sites (main.ts, sdk.ts, task/executor.ts, commit
pipelines, SDK example) all switched. Sync new ModelRegistry(...)
constructor is still supported for tests.
F2: rewrite MemorySessionStorage's mirror as { chunks: string[]; byteLen;
mtimeMs } so writeLineSync appends a single chunk in O(1) instead of
read-modify-writing the entire file (which was O(N) per append, O(N^2)
per session). statSync now reports true UTF-8 byte length instead of
character count. readTextPrefix walks chunks until the byte budget is
exhausted instead of materialising the full mirror.
Also rolls in per-package CHANGELOG entries for F1-F8.
450 lines
12 KiB
TypeScript
450 lines
12 KiB
TypeScript
import * as fs from "node:fs";
|
|
import * as fsp from "node:fs/promises";
|
|
import * as path from "node:path";
|
|
import { isEnoent, peekFile, 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<void>;
|
|
/**
|
|
* 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<void>;
|
|
fsync(): Promise<void>;
|
|
close(): Promise<void>;
|
|
getError(): Error | undefined;
|
|
}
|
|
|
|
export interface SessionStorage {
|
|
ensureDirSync(dir: string): void;
|
|
existsSync(path: string): boolean;
|
|
writeTextSync(path: string, content: string): void;
|
|
readTextSync(path: string): string;
|
|
statSync(path: string): SessionStorageStat;
|
|
listFilesSync(dir: string, pattern: string): string[];
|
|
|
|
exists(path: string): Promise<boolean>;
|
|
readText(path: string): Promise<string>;
|
|
readTextPrefix(path: string, maxBytes: number): Promise<string>;
|
|
writeText(path: string, content: string): Promise<void>;
|
|
rename(path: string, nextPath: string): Promise<void>;
|
|
unlink(path: string): Promise<void>;
|
|
deleteSessionWithArtifacts(sessionPath: string): Promise<void>;
|
|
openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter;
|
|
}
|
|
|
|
// FinalizationRegistry to clean up leaked file descriptors
|
|
const writerRegistry = new FinalizationRegistry<number>(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<void> {
|
|
this.writeLineSync(line);
|
|
}
|
|
|
|
async flush(): Promise<void> {
|
|
if (this.#error) throw this.#error;
|
|
// OS buffers are flushed on fsync, nothing to do here
|
|
}
|
|
|
|
async fsync(): Promise<void> {
|
|
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<void> {
|
|
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);
|
|
}
|
|
|
|
readTextSync(fpath: string): string {
|
|
return fs.readFileSync(fpath, "utf-8");
|
|
}
|
|
|
|
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<boolean> {
|
|
try {
|
|
await fs.promises.access(path);
|
|
return true;
|
|
} catch (err) {
|
|
if (isEnoent(err)) return false;
|
|
throw err;
|
|
}
|
|
}
|
|
|
|
readText(path: string): Promise<string> {
|
|
return Bun.file(path).text();
|
|
}
|
|
|
|
async readTextPrefix(path: string, maxBytes: number): Promise<string> {
|
|
return peekFile(path, maxBytes, header => utf8Decoder.decode(header));
|
|
}
|
|
|
|
async writeText(path: string, content: string): Promise<void> {
|
|
await Bun.write(path, content, { createPath: true });
|
|
}
|
|
|
|
async rename(path: string, nextPath: string): Promise<void> {
|
|
try {
|
|
await fs.promises.rename(path, nextPath);
|
|
} catch (err) {
|
|
throw toError(err);
|
|
}
|
|
}
|
|
|
|
unlink(path: string): Promise<void> {
|
|
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<void> {
|
|
// 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) chunked append — see MemorySessionStorage.appendChunkSync.
|
|
this.#storage.appendChunkSync(this.#path, line);
|
|
} catch (err) {
|
|
throw this.#recordError(err);
|
|
}
|
|
}
|
|
|
|
async writeLine(line: string): Promise<void> {
|
|
this.writeLineSync(line);
|
|
}
|
|
|
|
async flush(): Promise<void> {
|
|
if (this.#error) throw this.#error;
|
|
}
|
|
|
|
async fsync(): Promise<void> {
|
|
// No-op for in-memory storage
|
|
if (this.#error) throw this.#error;
|
|
}
|
|
|
|
async close(): Promise<void> {
|
|
if (this.#closed) return;
|
|
this.#closed = true;
|
|
}
|
|
|
|
getError(): Error | undefined {
|
|
return this.#error;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Mirror entry stored per path. Chunks accumulate via O(1) `push` on
|
|
* `appendChunkSync`; readers materialise into a single string lazily.
|
|
* `byteLen` is kept in sync so `statSync` is O(1) (and returns true UTF-8
|
|
* bytes, not character count).
|
|
*/
|
|
interface MirrorEntry {
|
|
chunks: string[];
|
|
byteLen: number;
|
|
mtimeMs: number;
|
|
}
|
|
|
|
function materialiseMirror(entry: MirrorEntry): string {
|
|
if (entry.chunks.length === 0) return "";
|
|
if (entry.chunks.length === 1) return entry.chunks[0];
|
|
return entry.chunks.join("");
|
|
}
|
|
|
|
export class MemorySessionStorage implements SessionStorage {
|
|
#files = new Map<string, MirrorEntry>();
|
|
|
|
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, {
|
|
chunks: content.length === 0 ? [] : [content],
|
|
byteLen: Buffer.byteLength(content, "utf-8"),
|
|
mtimeMs: 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.
|
|
*/
|
|
appendChunkSync(path: string, chunk: string): void {
|
|
let entry = this.#files.get(path);
|
|
if (!entry) {
|
|
entry = { chunks: [], byteLen: 0, mtimeMs: Date.now() };
|
|
this.#files.set(path, entry);
|
|
}
|
|
entry.chunks.push(chunk);
|
|
entry.byteLen += Buffer.byteLength(chunk, "utf-8");
|
|
entry.mtimeMs = Date.now();
|
|
}
|
|
|
|
readTextSync(path: string): string {
|
|
const entry = this.#files.get(path);
|
|
if (!entry) throw new Error(`File not found: ${path}`);
|
|
return materialiseMirror(entry);
|
|
}
|
|
|
|
statSync(path: string): SessionStorageStat {
|
|
const entry = this.#files.get(path);
|
|
if (!entry) throw new Error(`File not found: ${path}`);
|
|
return {
|
|
size: entry.byteLen,
|
|
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<boolean> {
|
|
return Promise.resolve(this.existsSync(path));
|
|
}
|
|
|
|
readText(path: string): Promise<string> {
|
|
const entry = this.#files.get(path);
|
|
if (!entry) return Promise.reject(new Error(`File not found: ${path}`));
|
|
return Promise.resolve(materialiseMirror(entry));
|
|
}
|
|
|
|
readTextPrefix(path: string, maxBytes: number): Promise<string> {
|
|
const entry = this.#files.get(path);
|
|
if (!entry) return Promise.reject(new Error(`File not found: ${path}`));
|
|
if (entry.chunks.length === 0 || maxBytes <= 0) return Promise.resolve("");
|
|
|
|
// Walk chunks until the byte budget is exhausted. Avoids materialising
|
|
// the full mirror just to slice a prefix — bounded work for big files.
|
|
let accumulatedBytes = 0;
|
|
const out: string[] = [];
|
|
for (const chunk of entry.chunks) {
|
|
const chunkBytes = Buffer.byteLength(chunk, "utf-8");
|
|
if (accumulatedBytes + chunkBytes <= maxBytes) {
|
|
out.push(chunk);
|
|
accumulatedBytes += chunkBytes;
|
|
if (accumulatedBytes === maxBytes) break;
|
|
continue;
|
|
}
|
|
// Boundary chunk: slice in byte space and decode. Result MAY be
|
|
// shorter than the budget if a multi-byte codepoint straddles the
|
|
// boundary — matches `peekFile` semantics (partial decode at cap).
|
|
const remainingBytes = maxBytes - accumulatedBytes;
|
|
const utf8 = Buffer.from(chunk, "utf-8");
|
|
out.push(utf8Decoder.decode(utf8.subarray(0, remainingBytes)));
|
|
break;
|
|
}
|
|
return Promise.resolve(out.join(""));
|
|
}
|
|
|
|
writeText(path: string, content: string): Promise<void> {
|
|
this.writeTextSync(path, content);
|
|
return Promise.resolve();
|
|
}
|
|
|
|
rename(path: string, nextPath: string): Promise<void> {
|
|
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<void> {
|
|
this.#files.delete(path);
|
|
return Promise.resolve();
|
|
}
|
|
deleteSessionWithArtifacts(_sessionPath: string): Promise<void> {
|
|
return Promise.resolve();
|
|
}
|
|
|
|
openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter {
|
|
return new MemorySessionStorageWriter(this, path, options);
|
|
}
|
|
}
|