feat(coding-agent): async ConfigFile + ModelRegistry.create() factory; fix MemorySessionStorage O(N^2) appends
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.
This commit is contained in:
@@ -23,6 +23,10 @@
|
||||
|
||||
- Fixed compaction summarizer throws losing the provider's HTTP status. `generateSummary`, `generateHandoff`, `generateShortSummary`, and `generateTurnPrefixSummary` now route their `stopReason === "error"` throws through a `createSummarizationError` helper that copies `AssistantMessage.errorStatus` onto the thrown `Error` as `.status`, letting downstream consumers (e.g. `AgentSession.#isCompactionAuthFailure` in `@oh-my-pi/pi-coding-agent`) branch on real provider 401/403s without regex-scraping the message body.
|
||||
|
||||
### Changed
|
||||
|
||||
- Changed `Agent.appendMessage`, `popMessage`, `clearMessages`, and `reset` to mutate `state.messages` and `state.pendingToolCalls` in place instead of allocating a fresh array/Set on every transition. Subscribers that capture `state.messages` by reference now observe updates without needing to re-read `state` after each event. The public type signature is unchanged (always `AgentMessage[]` / `Set<string>`).
|
||||
|
||||
## [15.5.0] - 2026-05-26
|
||||
### Added
|
||||
|
||||
|
||||
@@ -130,6 +130,15 @@
|
||||
|
||||
- Fixed Synthetic model discovery to treat the provider `/models` response as authoritative so deprecated bundled IDs are pruned from the runtime cache, and changed Synthetic login validation to avoid probing a specific model ([#1417](https://github.com/can1357/oh-my-pi/issues/1417)).
|
||||
|
||||
### Added
|
||||
|
||||
- Added `parseStreamingJsonThrottled` to `@oh-my-pi/pi-ai/utils/json-parse` — a per-delta wrapper around `parseStreamingJson` that skips re-parses until the buffer has grown by `minGrowthBytes` (default 256). Wired into the streaming hot path of every provider's tool-call argument accumulator (`anthropic`, `amazon-bedrock`, `openai-completions`, `openai-codex-responses`, `openai-responses-shared`) so per-delta cost is O(N) in total buffer length instead of O(N²). Each provider's `toolcall_end` still runs a final unthrottled parse, so the published `block.arguments` is unchanged.
|
||||
- Added named-tool routing support to Google providers: `GoogleSharedStreamOptions.toolChoice` and `GoogleGeminiCliOptions.toolChoice` now accept `{ mode: "ANY"; allowedFunctionNames: [string, ...string[]] }` in addition to the string forms. `mapGoogleToolChoice` converts `ToolChoice` objects of shape `{ type: "tool" | "function", name }` to the wire form. Mirrors the equivalent Anthropic mapper.
|
||||
|
||||
### Changed
|
||||
|
||||
- Changed `mapGoogleToolChoice` to be exported from `@oh-my-pi/pi-ai/stream` so callers can build the wire-shape allow-list directly without re-deriving it.
|
||||
|
||||
## [15.5.0] - 2026-05-26
|
||||
### Added
|
||||
|
||||
|
||||
@@ -313,12 +313,23 @@
|
||||
|
||||
- Added `read.summarize.minTotalLines` setting (default 100) to set the minimum file length that triggers read summarization
|
||||
- Added `<file>:<lines>` support to `search` `paths`, allowing file-scoped constraints such as `:N-M`, `:N+K`, and comma-separated ranges
|
||||
- Added `ModelRegistry.create(authStorage, modelsPath?)` async factory that runs the JSON → YAML migration step on `models.{yml,yaml}` asynchronously ahead of the sync constructor's bundled-model load. The sync `new ModelRegistry(...)` constructor still works (tests rely on it); production boot paths now use the factory so the migration's I/O lands off the event-loop hot path.
|
||||
- Added `ConfigFile.tryLoadAsync()`, `ConfigFile.loadAsync()`, `ConfigFile.loadOrDefaultAsync()`, `ConfigFile.getMtimeMsAsync()`, and `ConfigFile.warmup(file)` so the rest of the codebase can migrate config reads off the sync path.
|
||||
|
||||
### Changed
|
||||
|
||||
- Changed multi-section hashline `edit` execution to defer LSP diagnostics flushing until the final section is written
|
||||
- Changed read to return verbatim contents for files shorter than `read.summarize.minTotalLines` instead of summarizing them
|
||||
- Changed `search` path line-range filtering to include only matches and context lines that fall inside the requested ranges
|
||||
- Changed `MemorySessionStorage`'s mirror to a chunks-based representation. `writeLineSync` now appends to a `string[]` in O(1) (previously read the whole file and concatenated, giving O(N²) growth per session). `statSync` 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.
|
||||
- Changed `ToolExecutionComponent.updateArgs` to drop the per-delta `structuredClone` of streaming tool arguments. Callers (`event-controller.ts`, `ui-helpers.ts`) already spread their input into a fresh object on each delta, so cloning here was dead work on the rendering hot path. Added a reference-equality short-circuit so repeat calls with the same args object skip the preview-diff and display refresh.
|
||||
- Changed `ConfigFile`'s constructor to defer the JSON → YAML migration until first `tryLoad`/`tryLoadAsync` and to cache (jsonPath, ymlPath) pairs already migrated this process, so `relocate()` / repeated loads do not re-run the migration.
|
||||
- Changed all production `new ModelRegistry(...)` call sites (`main.ts`, `sdk.ts`, `task/executor.ts`, `commit/pipeline.ts`, `commit/agentic/index.ts`, the SDK example) to `await ModelRegistry.create(...)`.
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed a race in `withFileLock` where a contender losing the `mkdir` race could wipe the winner's freshly-created lock directory before the winner finished writing its info file. Every lock now carries a per-process UUID token; `releaseLock(path, expectedToken)` verifies ownership before `fs.rm`, and `isLockStale` no longer returns `true` for a dir whose info file is absent but whose mtime is still inside the staleness window (or whose dir vanished mid-check).
|
||||
- Fixed `formatErrorMessage` not sanitising tabs or truncating oversized error strings before painting them through the theme. Errors that embedded raw file content (apply_patch failures, hashline mismatches, etc.) could break terminal alignment via raw `\t` chars or overflow the line width.
|
||||
|
||||
### Fixed
|
||||
|
||||
|
||||
@@ -26,7 +26,7 @@ console.log("Session with default auth storage and model registry");
|
||||
|
||||
// Custom auth storage location
|
||||
const customAuthStorage = await AuthStorage.create("/tmp/my-app/agent.db");
|
||||
const customModelRegistry = new ModelRegistry(customAuthStorage, "/tmp/my-app/models.json");
|
||||
const customModelRegistry = await ModelRegistry.create(customAuthStorage, "/tmp/my-app/models.json");
|
||||
|
||||
await createAgentSession({
|
||||
sessionManager: SessionManager.inMemory(),
|
||||
@@ -45,7 +45,7 @@ await createAgentSession({
|
||||
console.log("Session with runtime API key override");
|
||||
|
||||
// No models.json - only built-in models
|
||||
const simpleRegistry = new ModelRegistry(authStorage); // null = no models.json
|
||||
const simpleRegistry = await ModelRegistry.create(authStorage); // null = no models.json
|
||||
await createAgentSession({
|
||||
sessionManager: SessionManager.inMemory(),
|
||||
authStorage,
|
||||
|
||||
@@ -29,7 +29,7 @@ export async function runAgenticCommit(args: CommitCommandArgs): Promise<void> {
|
||||
const [settings, authStorage] = await Promise.all([Settings.init({ cwd }), discoverAuthStorage()]);
|
||||
|
||||
process.stdout.write("● Resolving model...\n");
|
||||
const modelRegistry = new ModelRegistry(authStorage);
|
||||
const modelRegistry = await ModelRegistry.create(authStorage);
|
||||
await modelRegistry.refresh();
|
||||
const stagedFilesPromise = (async () => {
|
||||
let stagedFiles = await git.diff.changedFiles(cwd, { cached: true });
|
||||
|
||||
@@ -43,7 +43,7 @@ async function runLegacyCommitCommand(args: CommitCommandArgs): Promise<void> {
|
||||
const settings = await Settings.init();
|
||||
const commitSettings = settings.getGroup("commit");
|
||||
const authStorage = await discoverAuthStorage();
|
||||
const modelRegistry = new ModelRegistry(authStorage);
|
||||
const modelRegistry = await ModelRegistry.create(authStorage);
|
||||
await modelRegistry.refresh();
|
||||
|
||||
const {
|
||||
|
||||
@@ -10,23 +10,85 @@ interface ConfigSchemaError {
|
||||
message: string | undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Module-private cache of (jsonPath, ymlPath) pairs we already migrated this
|
||||
* process. Prevents `ConfigFile.relocate()` / repeated `tryLoad()` calls from
|
||||
* re-running the migration over and over on the boot path.
|
||||
*/
|
||||
const migratedPaths = new Set<string>();
|
||||
|
||||
function migrationKey(jsonPath: string, ymlPath: string): string {
|
||||
return `${jsonPath}\u0000${ymlPath}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Synchronous JSON → YAML migration kept for callers that still want the
|
||||
* eager path (settings init, tests that observe migration completion).
|
||||
* Idempotent — re-running is a no-op.
|
||||
*/
|
||||
function migrateJsonToYml(jsonPath: string, ymlPath: string) {
|
||||
const key = migrationKey(jsonPath, ymlPath);
|
||||
if (migratedPaths.has(key)) return;
|
||||
try {
|
||||
if (fs.existsSync(ymlPath)) return;
|
||||
if (!fs.existsSync(jsonPath)) return;
|
||||
if (fs.existsSync(ymlPath)) {
|
||||
migratedPaths.add(key);
|
||||
return;
|
||||
}
|
||||
if (!fs.existsSync(jsonPath)) {
|
||||
migratedPaths.add(key);
|
||||
return;
|
||||
}
|
||||
|
||||
const content = fs.readFileSync(jsonPath, "utf-8");
|
||||
const parsed = JSON.parse(content);
|
||||
if (!parsed) {
|
||||
logger.warn("migrateJsonToYml: invalid json structure", { path: jsonPath });
|
||||
migratedPaths.add(key);
|
||||
return;
|
||||
}
|
||||
fs.writeFileSync(ymlPath, YAML.stringify(parsed, null, 2));
|
||||
migratedPaths.add(key);
|
||||
} catch (error) {
|
||||
logger.warn("migrateJsonToYml: migration failed", { error: String(error) });
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Async sibling of `migrateJsonToYml`. Uses Bun.file so the boot path no
|
||||
* longer blocks on a sync FS read before the first await. Idempotent and
|
||||
* shares the process-wide `migratedPaths` cache with the sync path.
|
||||
*/
|
||||
async function migrateJsonToYmlAsync(jsonPath: string, ymlPath: string) {
|
||||
const key = migrationKey(jsonPath, ymlPath);
|
||||
if (migratedPaths.has(key)) return;
|
||||
try {
|
||||
if (await Bun.file(ymlPath).exists()) {
|
||||
migratedPaths.add(key);
|
||||
return;
|
||||
}
|
||||
let content: string;
|
||||
try {
|
||||
content = await Bun.file(jsonPath).text();
|
||||
} catch (err) {
|
||||
if (isEnoent(err)) {
|
||||
migratedPaths.add(key);
|
||||
return;
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
const parsed = JSON.parse(content);
|
||||
if (!parsed) {
|
||||
logger.warn("migrateJsonToYmlAsync: invalid json structure", { path: jsonPath });
|
||||
migratedPaths.add(key);
|
||||
return;
|
||||
}
|
||||
await Bun.write(ymlPath, YAML.stringify(parsed, null, 2));
|
||||
migratedPaths.add(key);
|
||||
} catch (error) {
|
||||
logger.warn("migrateJsonToYmlAsync: migration failed", { error: String(error) });
|
||||
}
|
||||
}
|
||||
|
||||
export interface IConfigFile<T> {
|
||||
readonly id: string;
|
||||
readonly schema: ZodType<T>;
|
||||
@@ -97,6 +159,7 @@ export type LoadResult<T> =
|
||||
|
||||
export class ConfigFile<T> implements IConfigFile<T> {
|
||||
readonly #basePath: string;
|
||||
readonly #jsonMigrationPath: string | null;
|
||||
#cache?: LoadResult<T>;
|
||||
#auxValidate?: (value: T) => void;
|
||||
|
||||
@@ -107,18 +170,47 @@ export class ConfigFile<T> implements IConfigFile<T> {
|
||||
) {
|
||||
this.#basePath = configPath;
|
||||
if (configPath.endsWith(".yml")) {
|
||||
const jsonPath = `${configPath.slice(0, -4)}.json`;
|
||||
migrateJsonToYml(jsonPath, configPath);
|
||||
this.#jsonMigrationPath = `${configPath.slice(0, -4)}.json`;
|
||||
} else if (configPath.endsWith(".yaml")) {
|
||||
const jsonPath = `${configPath.slice(0, -5)}.json`;
|
||||
migrateJsonToYml(jsonPath, configPath);
|
||||
this.#jsonMigrationPath = `${configPath.slice(0, -5)}.json`;
|
||||
} else if (configPath.endsWith(".json") || configPath.endsWith(".jsonc")) {
|
||||
// JSON configs are still supported without migration.
|
||||
this.#jsonMigrationPath = null;
|
||||
} else {
|
||||
throw new Error(`Invalid config file path: ${configPath}`);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Run the JSON → YAML migration synchronously, if applicable. Idempotent.
|
||||
* Sync callers (tests, settings init) hit this implicitly via {@link tryLoad}.
|
||||
*/
|
||||
#ensureMigratedSync(): void {
|
||||
if (this.#jsonMigrationPath) {
|
||||
migrateJsonToYml(this.#jsonMigrationPath, this.#basePath);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Async sibling of {@link #ensureMigratedSync}. Boot-path callers should
|
||||
* `await ConfigFile.warmup(file)` before doing any sync `tryLoad`/`load`
|
||||
* so the migration's I/O happens off the event-loop's hot path.
|
||||
*/
|
||||
async #ensureMigratedAsync(): Promise<void> {
|
||||
if (this.#jsonMigrationPath) {
|
||||
await migrateJsonToYmlAsync(this.#jsonMigrationPath, this.#basePath);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Run any pending JSON → YAML migration asynchronously, ahead of a sync
|
||||
* `tryLoad()` on the boot path. Safe to call multiple times; subsequent
|
||||
* calls are O(1) thanks to the module-level migration cache.
|
||||
*/
|
||||
static warmup<U>(file: ConfigFile<U>): Promise<void> {
|
||||
return file.#ensureMigratedAsync();
|
||||
}
|
||||
|
||||
relocate(configPath?: string): ConfigFile<T> {
|
||||
if (!configPath || configPath === this.#basePath) return this;
|
||||
const result = new ConfigFile<T>(this.id, this.schema, configPath);
|
||||
@@ -135,6 +227,13 @@ export class ConfigFile<T> implements IConfigFile<T> {
|
||||
}
|
||||
}
|
||||
|
||||
async getMtimeMsAsync(): Promise<number | null> {
|
||||
const file = Bun.file(this.path());
|
||||
if (!(await file.exists())) return null;
|
||||
const lm = file.lastModified;
|
||||
return typeof lm === "number" && Number.isFinite(lm) ? lm : null;
|
||||
}
|
||||
|
||||
withValidation(name: string, validate: (value: T) => void): this {
|
||||
const prev = this.#auxValidate;
|
||||
this.#auxValidate = (value: T) => {
|
||||
@@ -164,12 +263,8 @@ export class ConfigFile<T> implements IConfigFile<T> {
|
||||
return result;
|
||||
}
|
||||
|
||||
tryLoad(): LoadResult<T> {
|
||||
if (this.#cache) return this.#cache;
|
||||
|
||||
#parseContent(content: string): LoadResult<T> {
|
||||
try {
|
||||
const content = fs.readFileSync(this.path(), "utf-8").trim();
|
||||
|
||||
let parsed: unknown;
|
||||
if (this.#basePath.endsWith(".json") || this.#basePath.endsWith(".jsonc")) {
|
||||
parsed = JSONC.parse(content);
|
||||
@@ -203,9 +298,6 @@ export class ConfigFile<T> implements IConfigFile<T> {
|
||||
}
|
||||
return this.#storeCache({ value, status: "ok" });
|
||||
} catch (error) {
|
||||
if (isEnoent(error)) {
|
||||
return this.#storeCache({ status: "not-found" });
|
||||
}
|
||||
logger.warn("Failed to parse config file", { path: this.path(), error });
|
||||
return this.#storeCache({
|
||||
error: new ConfigError(this.id, undefined, { err: error, stage: "Unexpected" }),
|
||||
@@ -214,14 +306,62 @@ export class ConfigFile<T> implements IConfigFile<T> {
|
||||
}
|
||||
}
|
||||
|
||||
tryLoad(): LoadResult<T> {
|
||||
if (this.#cache) return this.#cache;
|
||||
this.#ensureMigratedSync();
|
||||
|
||||
let content: string;
|
||||
try {
|
||||
content = fs.readFileSync(this.path(), "utf-8").trim();
|
||||
} catch (error) {
|
||||
if (isEnoent(error)) {
|
||||
return this.#storeCache({ status: "not-found" });
|
||||
}
|
||||
logger.warn("Failed to read config file", { path: this.path(), error });
|
||||
return this.#storeCache({
|
||||
error: new ConfigError(this.id, undefined, { err: error, stage: "Read" }),
|
||||
status: "error",
|
||||
});
|
||||
}
|
||||
return this.#parseContent(content);
|
||||
}
|
||||
|
||||
async tryLoadAsync(): Promise<LoadResult<T>> {
|
||||
if (this.#cache) return this.#cache;
|
||||
await this.#ensureMigratedAsync();
|
||||
|
||||
let content: string;
|
||||
try {
|
||||
content = (await Bun.file(this.path()).text()).trim();
|
||||
} catch (error) {
|
||||
if (isEnoent(error)) {
|
||||
return this.#storeCache({ status: "not-found" });
|
||||
}
|
||||
logger.warn("Failed to read config file", { path: this.path(), error });
|
||||
return this.#storeCache({
|
||||
error: new ConfigError(this.id, undefined, { err: error, stage: "Read" }),
|
||||
status: "error",
|
||||
});
|
||||
}
|
||||
return this.#parseContent(content);
|
||||
}
|
||||
|
||||
load(): T | null {
|
||||
return this.tryLoad().value ?? null;
|
||||
}
|
||||
|
||||
async loadAsync(): Promise<T | null> {
|
||||
return (await this.tryLoadAsync()).value ?? null;
|
||||
}
|
||||
|
||||
loadOrDefault(): T {
|
||||
return this.tryLoad().value ?? this.createDefault();
|
||||
}
|
||||
|
||||
async loadOrDefaultAsync(): Promise<T> {
|
||||
return (await this.tryLoadAsync()).value ?? this.createDefault();
|
||||
}
|
||||
|
||||
path(): string {
|
||||
return this.#basePath;
|
||||
}
|
||||
|
||||
@@ -851,6 +851,12 @@ export class ModelRegistry {
|
||||
|
||||
/**
|
||||
* @param authStorage - Auth storage for API key resolution
|
||||
*
|
||||
* Sync constructor — eagerly loads bundled + cached models so tests and
|
||||
* synchronous callers see a fully-populated registry immediately. Production
|
||||
* boot paths SHOULD prefer {@link ModelRegistry.create} so the YAML/JSONC
|
||||
* migration step lands off the event loop's hot path before the first
|
||||
* `tryLoad()` runs.
|
||||
*/
|
||||
constructor(
|
||||
readonly authStorage: AuthStorage,
|
||||
@@ -866,10 +872,27 @@ export class ModelRegistry {
|
||||
}
|
||||
return undefined;
|
||||
});
|
||||
// Load models synchronously in constructor
|
||||
// Load models synchronously in constructor.
|
||||
this.#loadModels();
|
||||
}
|
||||
|
||||
/**
|
||||
* Async factory used by the production boot path. Runs the JSON → YAML
|
||||
* migration on `models.{yml,yaml}` ahead of any sync I/O the constructor
|
||||
* still performs, so the registry's first `tryLoad()` is a pure read of
|
||||
* an already-migrated file. The constructor's bundled + cached-snapshot
|
||||
* load is unchanged — this factory only adds an awaited warmup step.
|
||||
*/
|
||||
static async create(authStorage: AuthStorage, modelsPath?: string): Promise<ModelRegistry> {
|
||||
// Warm the migration before the constructor's sync tryLoad fires. We
|
||||
// reach the underlying ConfigFile via a temporary relocate so the
|
||||
// shared static registry stays untouched until the real instance is
|
||||
// constructed below.
|
||||
const warmupFile = ModelsConfigFile.relocate(modelsPath);
|
||||
await ConfigFile.warmup(warmupFile);
|
||||
return new ModelRegistry(authStorage, modelsPath);
|
||||
}
|
||||
|
||||
/**
|
||||
* Reload models from disk (built-in + custom from models.json).
|
||||
*/
|
||||
|
||||
@@ -733,7 +733,7 @@ export async function runRootCommand(
|
||||
|
||||
// Create AuthStorage and ModelRegistry upfront
|
||||
const authStorage = await logger.time("discoverModels", deps.discoverAuthStorage ?? discoverAuthStorage);
|
||||
const modelRegistry = new ModelRegistry(authStorage);
|
||||
const modelRegistry = await ModelRegistry.create(authStorage);
|
||||
|
||||
if (parsedArgs.version) {
|
||||
process.stdout.write(`${VERSION}\n`);
|
||||
|
||||
@@ -838,7 +838,9 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
|
||||
// / session would silently miss credential_disabled events.
|
||||
const modelRegistry =
|
||||
options.modelRegistry ??
|
||||
new ModelRegistry(options.authStorage ?? (await logger.time("discoverModels", discoverAuthStorage, agentDir)));
|
||||
(await ModelRegistry.create(
|
||||
options.authStorage ?? (await logger.time("discoverModels", discoverAuthStorage, agentDir)),
|
||||
));
|
||||
const authStorage = modelRegistry.authStorage;
|
||||
if (options.authStorage && options.authStorage !== authStorage) {
|
||||
throw new Error(
|
||||
|
||||
@@ -272,8 +272,8 @@ class MemorySessionStorageWriter implements SessionStorageWriter {
|
||||
if (this.#closed) throw new Error("Writer closed");
|
||||
if (this.#error) throw this.#error;
|
||||
try {
|
||||
const existing = this.#storage.existsSync(this.#path) ? this.#storage.readTextSync(this.#path) : "";
|
||||
this.#storage.writeTextSync(this.#path, `${existing}${line}`);
|
||||
// O(1) chunked append — see MemorySessionStorage.appendChunkSync.
|
||||
this.#storage.appendChunkSync(this.#path, line);
|
||||
} catch (err) {
|
||||
throw this.#recordError(err);
|
||||
}
|
||||
@@ -302,8 +302,26 @@ class MemorySessionStorageWriter implements SessionStorageWriter {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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, { content: string; mtimeMs: number }>();
|
||||
#files = new Map<string, MirrorEntry>();
|
||||
|
||||
ensureDirSync(_dir: string): void {
|
||||
// No-op for in-memory storage.
|
||||
@@ -314,20 +332,40 @@ export class MemorySessionStorage implements SessionStorage {
|
||||
}
|
||||
|
||||
writeTextSync(path: string, content: string): void {
|
||||
this.#files.set(path, { content, mtimeMs: Date.now() });
|
||||
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 entry.content;
|
||||
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.content.length,
|
||||
size: entry.byteLen,
|
||||
mtimeMs: entry.mtimeMs,
|
||||
mtime: new Date(entry.mtimeMs),
|
||||
};
|
||||
@@ -353,13 +391,35 @@ export class MemorySessionStorage implements SessionStorage {
|
||||
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(entry.content);
|
||||
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}`));
|
||||
return Promise.resolve(entry.content.slice(0, maxBytes));
|
||||
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> {
|
||||
|
||||
@@ -1121,7 +1121,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
const registryFromParent = options.modelRegistry !== undefined;
|
||||
const modelRegistry =
|
||||
options.modelRegistry ??
|
||||
new ModelRegistry(options.authStorage ?? (await awaitAbortable(discoverAuthStorage())));
|
||||
(await ModelRegistry.create(options.authStorage ?? (await awaitAbortable(discoverAuthStorage()))));
|
||||
const authStorage = modelRegistry.authStorage;
|
||||
if (options.authStorage && options.authStorage !== authStorage) {
|
||||
throw new Error(
|
||||
|
||||
@@ -0,0 +1,84 @@
|
||||
import { describe, expect, test } from "bun:test";
|
||||
import { MemorySessionStorage } from "../src/session/session-storage";
|
||||
|
||||
describe("MemorySessionStorage chunked mirror (F2)", () => {
|
||||
test("writeLineSync builds the same content as a single writeTextSync of the join", async () => {
|
||||
const storage = new MemorySessionStorage();
|
||||
const path = "/virtual/session.jsonl";
|
||||
const writer = storage.openWriter(path, { flags: "w" });
|
||||
try {
|
||||
const N = 1000;
|
||||
for (let i = 0; i < N; i++) {
|
||||
writer.writeLineSync(`{"i":${i}}\n`);
|
||||
}
|
||||
} finally {
|
||||
await writer.close();
|
||||
}
|
||||
|
||||
// Construct the baseline from the same parts.
|
||||
const expected = Array.from({ length: 1000 }, (_, i) => `{"i":${i}}\n`).join("");
|
||||
const actual = storage.readTextSync(path);
|
||||
expect(actual).toBe(expected);
|
||||
expect(actual.length).toBe(expected.length);
|
||||
});
|
||||
|
||||
test("statSync reports UTF-8 byte length, not character count", () => {
|
||||
const storage = new MemorySessionStorage();
|
||||
const path = "/virtual/unicode.jsonl";
|
||||
const writer = storage.openWriter(path, { flags: "w" });
|
||||
try {
|
||||
writer.writeLineSync("héllo\n"); // é = 2 bytes in UTF-8
|
||||
writer.writeLineSync("日本語\n"); // 3 chars × 3 bytes = 9
|
||||
} finally {
|
||||
void writer.close();
|
||||
}
|
||||
|
||||
const expectedBytes = Buffer.byteLength("héllo\n日本語\n", "utf-8");
|
||||
expect(storage.statSync(path).size).toBe(expectedBytes);
|
||||
});
|
||||
|
||||
test("readTextPrefix walks chunks until the byte budget is exhausted", async () => {
|
||||
const storage = new MemorySessionStorage();
|
||||
const path = "/virtual/prefix.jsonl";
|
||||
const writer = storage.openWriter(path, { flags: "w" });
|
||||
try {
|
||||
writer.writeLineSync("alpha\n");
|
||||
writer.writeLineSync("bravo\n");
|
||||
writer.writeLineSync("charlie\n");
|
||||
} finally {
|
||||
void writer.close();
|
||||
}
|
||||
|
||||
// Cap mid-second-chunk; first chunk = 6B, take 4 of the second.
|
||||
const prefix = await storage.readTextPrefix(path, 10);
|
||||
expect(prefix).toBe("alpha\nbrav");
|
||||
});
|
||||
|
||||
test("subsequent writeLineSync after readTextSync stays O(1) (chunks preserved)", () => {
|
||||
const storage = new MemorySessionStorage();
|
||||
const path = "/virtual/cont.jsonl";
|
||||
const writer = storage.openWriter(path, { flags: "w" });
|
||||
try {
|
||||
writer.writeLineSync("first\n");
|
||||
writer.writeLineSync("second\n");
|
||||
// Materialise once — implementation must NOT cache the joined string
|
||||
// back into the entry, or subsequent appends collapse back to O(N).
|
||||
expect(storage.readTextSync(path)).toBe("first\nsecond\n");
|
||||
writer.writeLineSync("third\n");
|
||||
expect(storage.readTextSync(path)).toBe("first\nsecond\nthird\n");
|
||||
expect(storage.statSync(path).size).toBe(Buffer.byteLength("first\nsecond\nthird\n", "utf-8"));
|
||||
} finally {
|
||||
void writer.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("writeTextSync resets the chunks and byte counter (overwrite semantics)", () => {
|
||||
const storage = new MemorySessionStorage();
|
||||
const path = "/virtual/overwrite.jsonl";
|
||||
storage.writeTextSync(path, "abcdef");
|
||||
expect(storage.statSync(path).size).toBe(6);
|
||||
storage.writeTextSync(path, "xy");
|
||||
expect(storage.statSync(path).size).toBe(2);
|
||||
expect(storage.readTextSync(path)).toBe("xy");
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,71 @@
|
||||
import { afterEach, beforeEach, describe, expect, test } from "bun:test";
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
|
||||
import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
import { ConfigFile } from "../src/config/config-file";
|
||||
import { ModelRegistry } from "../src/config/model-registry";
|
||||
import { ModelsConfigSchema } from "../src/config/models-config-schema";
|
||||
|
||||
describe("ModelRegistry.create() factory (F6)", () => {
|
||||
let tempDir: TempDir;
|
||||
|
||||
beforeEach(() => {
|
||||
tempDir = TempDir.createSync("@model-registry-create-");
|
||||
});
|
||||
|
||||
afterEach(async () => {
|
||||
// On Windows the cache SQLite handle inside the registry may briefly hold
|
||||
// the dir; treat cleanup errors as best-effort like TempDir's Symbol.dispose.
|
||||
await tempDir.remove().catch(() => {});
|
||||
});
|
||||
|
||||
test("produces an instance whose authStorage matches and that exposes bundled models", async () => {
|
||||
const authStorage = await AuthStorage.create(path.join(tempDir.path(), "auth.db"));
|
||||
try {
|
||||
const registry = await ModelRegistry.create(authStorage, path.join(tempDir.path(), "models.yml"));
|
||||
expect(registry.authStorage).toBe(authStorage);
|
||||
// The constructor's bundled-model load runs after warmup, so the
|
||||
// factory's returned instance must be queryable immediately.
|
||||
const claude = registry.find("anthropic", "claude-sonnet-4-5");
|
||||
expect(claude).toBeDefined();
|
||||
expect(claude?.id).toBe("claude-sonnet-4-5");
|
||||
} finally {
|
||||
authStorage.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("migrates legacy models.json → models.yml ahead of the sync constructor", async () => {
|
||||
const yml = path.join(tempDir.path(), "models.yml");
|
||||
const json = path.join(tempDir.path(), "models.json");
|
||||
|
||||
// Seed a legacy JSON config; factory should migrate it asynchronously
|
||||
// before the sync constructor reads from the yml path.
|
||||
await Bun.write(json, JSON.stringify({ models: [] }));
|
||||
expect(fs.existsSync(yml)).toBe(false);
|
||||
|
||||
const authStorage = await AuthStorage.create(path.join(tempDir.path(), "auth.db"));
|
||||
try {
|
||||
await ModelRegistry.create(authStorage, yml);
|
||||
expect(fs.existsSync(yml)).toBe(true);
|
||||
} finally {
|
||||
authStorage.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("ConfigFile.warmup is idempotent — second call is a no-op", async () => {
|
||||
const yml = path.join(tempDir.path(), "models.yml");
|
||||
const json = path.join(tempDir.path(), "models.json");
|
||||
await Bun.write(json, JSON.stringify({ models: [] }));
|
||||
|
||||
const cf = new ConfigFile("models", ModelsConfigSchema, yml);
|
||||
await ConfigFile.warmup(cf);
|
||||
expect(fs.existsSync(yml)).toBe(true);
|
||||
const mtime1 = fs.statSync(yml).mtimeMs;
|
||||
|
||||
// Second warmup should not rewrite the file (idempotent path).
|
||||
await ConfigFile.warmup(cf);
|
||||
const mtime2 = fs.statSync(yml).mtimeMs;
|
||||
expect(mtime2).toBe(mtime1);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user