File size: 5,231 Bytes
f500658 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 | import { join } from "node:path";
import type { ModelsStore, ModelsStoreEntry, ModelsStoreOperationOptions } from "@earendil-works/pi-ai";
import { getAgentDir } from "../config.ts";
import { raceWithAbortSignal } from "../utils/abort.ts";
import { getFileRevision, normalizePath } from "../utils/paths.ts";
import { stripBom } from "../utils/text.ts";
import { type AuthStorageBackend, FileAuthStorageBackend } from "./auth-storage.ts";
type StoredModels = Record<string, ModelsStoreEntry>;
type ModelsFileReload = {
controller: AbortController;
promise: Promise<StoredModels>;
readers: number;
};
type ModelsFileReadState = {
data: StoredModels;
revision?: string;
reload?: ModelsFileReload;
};
// Optimize the common path without retaining an unbounded set of custom paths.
let sharedModelsFileReadState: { path: string; readState: ModelsFileReadState } | undefined;
export class InMemoryCodingAgentModelsStore implements ModelsStore {
private readonly entries = new Map<string, ModelsStoreEntry>();
async read(providerId: string, options?: ModelsStoreOperationOptions): Promise<ModelsStoreEntry | undefined> {
options?.signal?.throwIfAborted();
const entry = this.entries.get(providerId);
return entry ? structuredClone(entry) : undefined;
}
async write(providerId: string, entry: ModelsStoreEntry, options?: ModelsStoreOperationOptions): Promise<void> {
options?.signal?.throwIfAborted();
this.entries.set(providerId, structuredClone(entry));
}
async delete(providerId: string, options?: ModelsStoreOperationOptions): Promise<void> {
options?.signal?.throwIfAborted();
this.entries.delete(providerId);
}
}
/** Locked JSON-backed storage for dynamically refreshed provider catalogs. */
export class FileModelsStore implements ModelsStore {
private readonly storage: AuthStorageBackend;
private readonly path: string;
private readonly readState: ModelsFileReadState;
constructor(path: string = join(getAgentDir(), "models-store.json")) {
this.path = normalizePath(path);
this.storage = new FileAuthStorageBackend(this.path);
this.readState =
sharedModelsFileReadState?.path === this.path ? sharedModelsFileReadState.readState : { data: {} };
if (!sharedModelsFileReadState) {
sharedModelsFileReadState = { path: this.path, readState: this.readState };
}
}
private parse(content: string | undefined): StoredModels {
return content ? (JSON.parse(stripBom(content)) as StoredModels) : {};
}
private updateReadState(readState: ModelsFileReadState, data: StoredModels, revision?: string): void {
readState.data = data;
readState.revision = revision;
}
private reloadFromStorage(
readState: ModelsFileReadState,
options?: ModelsStoreOperationOptions,
): Promise<StoredModels> {
return this.storage.withLockAsync(async (content) => {
const data = this.parse(content);
this.updateReadState(readState, data, getFileRevision(this.path));
return { result: data };
}, options);
}
private async readLatest(
readState: ModelsFileReadState,
options?: ModelsStoreOperationOptions,
): Promise<StoredModels> {
options?.signal?.throwIfAborted();
const revision = getFileRevision(this.path);
if (revision !== undefined && revision === readState.revision) return readState.data;
if (!readState.reload) {
const controller = new AbortController();
const reload: ModelsFileReload = {
controller,
promise: this.reloadFromStorage(readState, { signal: controller.signal }),
readers: 0,
};
readState.reload = reload;
void reload.promise.then(
() => {
if (readState.reload === reload) readState.reload = undefined;
},
() => {
if (readState.reload === reload) readState.reload = undefined;
},
);
}
const reload = readState.reload;
reload.readers++;
try {
return await raceWithAbortSignal(reload.promise, options?.signal);
} finally {
reload.readers--;
if (reload.readers === 0 && readState.reload === reload) {
readState.reload = undefined;
reload.controller.abort();
}
}
}
async read(providerId: string, options?: ModelsStoreOperationOptions): Promise<ModelsStoreEntry | undefined> {
const entry = (await this.readLatest(this.readState, options))[providerId];
options?.signal?.throwIfAborted();
return entry ? structuredClone(entry) : undefined;
}
async write(providerId: string, entry: ModelsStoreEntry, options?: ModelsStoreOperationOptions): Promise<void> {
let latest: StoredModels | undefined;
await this.storage.withLockAsync(async (content) => {
const current = this.parse(content);
current[providerId] = structuredClone(entry);
latest = current;
return { result: undefined, next: JSON.stringify(current, null, 2) };
}, options);
if (latest) this.updateReadState(this.readState, latest);
}
async delete(providerId: string, options?: ModelsStoreOperationOptions): Promise<void> {
let latest: StoredModels | undefined;
await this.storage.withLockAsync(async (content) => {
const current = this.parse(content);
delete current[providerId];
latest = current;
return { result: undefined, next: JSON.stringify(current, null, 2) };
}, options);
if (latest) this.updateReadState(this.readState, latest);
}
}
|