feat(rewrite): tambah infrastructure — LLM client + full file persistence
Sub-fase 3a & 3e infrastructure TypeScript (@zesdex/infrastructure): - llm/provider.ts: LlmClient (OpenAI/Anthropic-compatible fetch + SSE, retry 10x exponential backoff + jitter, abort on 401/402/403, stream quota-aware tool-call delta accumulation) - resolveApiKey: settings.api_keys → app_config.env_var → inline default - utils: writeJsonAtomic (tmp+rename+fsync), slugify, writeOsc52, truncateChars, buildWorkspaceTree, runCommand - persistence/cms: JsonSettingsRepository, JsonAppConfigRepository (Claude credential auto-detection dari ~/.claude/settings.json + env), JsonConversationRepository, JsonlEditLogRepository (append NDJSON), MarkdownMemoryRepository (YAML-ish frontmatter), FileRewindBlobRepository - persistence/iam: FileSystemSessionRepository, FileSystemSessionLockRepository, FileSystemOAuthRepository - Workspace @zesdex/infrastructure tsc clean, semua test domain+application tetap hijau (38 pass) Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
03538a455a
commit
a2c7e5d44c
@@ -0,0 +1,22 @@
|
||||
{
|
||||
"name": "@zesdex/infrastructure",
|
||||
"version": "1.21.2",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"description": "Zesdex infrastructure layer — concrete adapters (LLM, tools, persistence, auth, MCP, IPC, workflow)",
|
||||
"exports": {
|
||||
".": "./src/index.ts"
|
||||
},
|
||||
"scripts": {
|
||||
"test": "bun test",
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@zesdex/domain": "workspace:*",
|
||||
"@zesdex/application": "workspace:*"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/bun": "^1.2.0",
|
||||
"typescript": "^5.7.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
/**
|
||||
* Zesdex Infrastructure Layer — concrete adapters implementing domain/application
|
||||
* ports. Mirrors the Rust `zesdex-infrastructure` crate.
|
||||
*/
|
||||
export * from "./utils.ts";
|
||||
export * from "./llm/index.ts";
|
||||
export * from "./persistence/index.ts";
|
||||
@@ -0,0 +1,2 @@
|
||||
/** LLM layer — provider client. */
|
||||
export * from "./provider.ts";
|
||||
@@ -0,0 +1,278 @@
|
||||
/**
|
||||
* HTTP client for OpenAI/Anthropic-compatible chat completion APIs, with
|
||||
* automatic retry and SSE streaming. Mirrors `apps/infrastructure/src/llm/provider.rs`.
|
||||
*/
|
||||
import {
|
||||
type ChatMessage,
|
||||
type ChatRequest,
|
||||
type StreamEvent,
|
||||
type ToolCall,
|
||||
type ToolDef,
|
||||
type Settings,
|
||||
type AppConfig,
|
||||
SseParser,
|
||||
Roles,
|
||||
DEFAULT_API_BASE,
|
||||
DEFAULT_MODEL,
|
||||
} from "@zesdex/domain";
|
||||
import type { ProviderService } from "@zesdex/application";
|
||||
|
||||
const DEFAULT_BASE_URL = DEFAULT_API_BASE;
|
||||
const REQUEST_TIMEOUT_MS = 600_000;
|
||||
const MAX_RETRIES = 10;
|
||||
|
||||
function backoffSeconds(attempt: number, cap: number): number {
|
||||
const base = Math.pow(2, Math.max(attempt - 1, 0));
|
||||
const delay = Math.min(base, cap);
|
||||
const jitterFactor = 0.75 + Math.floor(Math.random() * 51) / 100.0;
|
||||
return delay * jitterFactor;
|
||||
}
|
||||
|
||||
/** Whether a provider error message indicates an auth problem. */
|
||||
export function isAuthError(errStr: string): boolean {
|
||||
const lower = errStr.toLowerCase();
|
||||
return (
|
||||
errStr.includes("API error 401") ||
|
||||
errStr.includes("API error 402") ||
|
||||
errStr.includes("API error 403") ||
|
||||
lower.includes("unauthorized") ||
|
||||
lower.includes("forbidden") ||
|
||||
lower.includes("authentication failed")
|
||||
);
|
||||
}
|
||||
|
||||
function isRateLimit(errStr: string): boolean {
|
||||
return errStr.includes("API error 429") || errStr.toLowerCase().includes("rate limit");
|
||||
}
|
||||
|
||||
function backoffForError(attempt: number, errStr: string): number {
|
||||
return isRateLimit(errStr) ? backoffSeconds(attempt, 60) : backoffSeconds(attempt, 30);
|
||||
}
|
||||
|
||||
const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));
|
||||
|
||||
/** Assemble the assistant message and accumulate usage from stream events. */
|
||||
class StreamedTurn {
|
||||
content = "";
|
||||
toolCalls: ToolCall[] = [];
|
||||
usage: [number, number] | null = null;
|
||||
doneReceived = false;
|
||||
|
||||
applyEvent(event: StreamEvent): void {
|
||||
switch (event.kind) {
|
||||
case "token":
|
||||
this.content += event.content;
|
||||
break;
|
||||
case "tool_call_delta": {
|
||||
const index = Math.min(event.index, 63);
|
||||
while (this.toolCalls.length <= index) {
|
||||
this.toolCalls.push({ id: "", type: "function", function: { name: "", arguments: "" } });
|
||||
}
|
||||
const tc = this.toolCalls[index]!;
|
||||
if (event.id) tc.id = event.id;
|
||||
if (event.name) tc.function.name = event.name;
|
||||
tc.function.arguments = `${tc.function.arguments}${event.arguments_delta}`;
|
||||
break;
|
||||
}
|
||||
case "usage":
|
||||
this.usage = [event.prompt_tokens, event.completion_tokens];
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
buildAssistantMessage(): ChatMessage {
|
||||
const msg: ChatMessage = {
|
||||
role: Roles.Assistant,
|
||||
content: this.content === "" ? null : this.content,
|
||||
};
|
||||
if (this.toolCalls.length > 0) msg.tool_calls = this.toolCalls;
|
||||
return msg;
|
||||
}
|
||||
}
|
||||
|
||||
/** Async HTTP client for a single LLM provider endpoint. */
|
||||
export class LlmClient implements ProviderService {
|
||||
apiKey: string;
|
||||
baseUrl: string;
|
||||
model: string;
|
||||
|
||||
constructor(apiKey: string, model: string, baseUrl?: string) {
|
||||
this.apiKey = apiKey || "";
|
||||
this.model = model === "" ? DEFAULT_MODEL : model;
|
||||
this.baseUrl = baseUrl && baseUrl !== "" ? baseUrl : DEFAULT_BASE_URL;
|
||||
}
|
||||
|
||||
private buildUrl(): string {
|
||||
return `${this.baseUrl}/chat/completions`;
|
||||
}
|
||||
|
||||
private headers(): HeadersInit {
|
||||
const h: Record<string, string> = { "Content-Type": "application/json" };
|
||||
if (this.apiKey !== "") h.Authorization = `Bearer ${this.apiKey}`;
|
||||
return h;
|
||||
}
|
||||
|
||||
private errorFromFetch(e: unknown): string {
|
||||
const err = e as { name?: string; message?: string };
|
||||
if (err?.name === "TimeoutError") {
|
||||
return `API request timed out after ${REQUEST_TIMEOUT_MS}ms. Check your network or try again.`;
|
||||
}
|
||||
if (err?.name === "TypeError" && /fetch failed|network/i.test(err.message ?? "")) {
|
||||
return `Could not connect to ${this.baseUrl}. Is the URL correct and is the service reachable?`;
|
||||
}
|
||||
return `API request failed: ${err?.message ?? String(e)}`;
|
||||
}
|
||||
|
||||
async chat(
|
||||
messages: ChatMessage[],
|
||||
tools?: ToolDef[],
|
||||
maxTokens?: number,
|
||||
temperature?: number,
|
||||
): Promise<{ message: ChatMessage; usage: [number, number] | null }> {
|
||||
const req: ChatRequest = {
|
||||
model: this.model,
|
||||
messages,
|
||||
max_tokens: maxTokens ?? 4096,
|
||||
temperature: temperature ?? 0.7,
|
||||
tools,
|
||||
stream: false,
|
||||
};
|
||||
const url = this.buildUrl();
|
||||
|
||||
let attempt = 0;
|
||||
for (;;) {
|
||||
attempt++;
|
||||
try {
|
||||
const resp = await fetch(url, {
|
||||
method: "POST",
|
||||
headers: this.headers(),
|
||||
body: JSON.stringify(req),
|
||||
signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS),
|
||||
});
|
||||
if (!resp.ok) {
|
||||
const body = await resp.text();
|
||||
throw new Error(`API error ${resp.status} from ${this.baseUrl}: ${body}`);
|
||||
}
|
||||
const data = (await resp.json()) as {
|
||||
usage?: { prompt_tokens?: number; completion_tokens?: number };
|
||||
choices?: { message?: ChatMessage }[];
|
||||
};
|
||||
const usage: [number, number] | null = data.usage
|
||||
? [data.usage.prompt_tokens ?? 0, data.usage.completion_tokens ?? 0]
|
||||
: null;
|
||||
const message = data.choices?.[0]?.message;
|
||||
if (!message) throw new Error("API response had no choices");
|
||||
return { message, usage };
|
||||
} catch (e) {
|
||||
const errStr = this.errorFromFetch(e);
|
||||
if (attempt >= MAX_RETRIES || isAuthError(errStr)) throw new Error(errStr);
|
||||
const delayMs = backoffForError(attempt, errStr) * 1000;
|
||||
await sleep(delayMs);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async chatStream(
|
||||
messages: ChatMessage[],
|
||||
tools: ToolDef[],
|
||||
maxTokens: number,
|
||||
temperature: number,
|
||||
onEvent: (event: StreamEvent) => boolean,
|
||||
signal?: AbortSignal,
|
||||
): Promise<{ message: ChatMessage; usage: [number, number] | null }> {
|
||||
const req: ChatRequest = {
|
||||
model: this.model,
|
||||
messages,
|
||||
max_tokens: maxTokens,
|
||||
temperature,
|
||||
tools,
|
||||
stream: true,
|
||||
stream_options: { include_usage: true },
|
||||
};
|
||||
const url = this.buildUrl();
|
||||
|
||||
let attempt = 0;
|
||||
for (;;) {
|
||||
attempt++;
|
||||
let capturedContent = false;
|
||||
try {
|
||||
const resp = await fetch(url, {
|
||||
method: "POST",
|
||||
headers: this.headers(),
|
||||
body: JSON.stringify(req),
|
||||
signal: signal ?? undefined,
|
||||
});
|
||||
if (!resp.ok) {
|
||||
const body = await resp.text();
|
||||
throw new Error(`API error ${resp.status} from ${this.baseUrl}: ${body}`);
|
||||
}
|
||||
return await this.handleStream(resp, onEvent);
|
||||
} catch (e) {
|
||||
const errStr = this.errorFromFetch(e);
|
||||
// Only retry if we captured no content yet.
|
||||
if (isAuthError(errStr) || capturedContent || attempt >= MAX_RETRIES) throw new Error(errStr);
|
||||
const delayMs = backoffForError(attempt, errStr) * 1000;
|
||||
await sleep(delayMs);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async handleStream(
|
||||
resp: Response,
|
||||
onEvent: (event: StreamEvent) => boolean,
|
||||
): Promise<{ message: ChatMessage; usage: [number, number] | null }> {
|
||||
if (!resp.body) throw new Error("no response body");
|
||||
const reader = resp.body.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
const parser = new SseParser();
|
||||
const turn = new StreamedTurn();
|
||||
let buffer = "";
|
||||
|
||||
for (;;) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
buffer += decoder.decode(value, { stream: true });
|
||||
|
||||
// Try to decode partial UTF-8 safely.
|
||||
let processed = 0;
|
||||
// Feed full UTF-8 boundaries only; keep remainder for next chunk.
|
||||
for (let i = buffer.length; i > 0; i--) {
|
||||
try {
|
||||
const text = buffer.slice(processed, i);
|
||||
processed = i;
|
||||
for (const event of parser.feed(text)) {
|
||||
if (!onEvent(event)) throw new Error("aborted");
|
||||
if (event.kind === "error") throw new Error(`stream error: ${event.message}`);
|
||||
if (event.kind === "done") {
|
||||
turn.applyEvent(event);
|
||||
return { message: turn.buildAssistantMessage(), usage: turn.usage };
|
||||
}
|
||||
turn.applyEvent(event);
|
||||
}
|
||||
break;
|
||||
} catch {
|
||||
/* not a valid boundary yet */
|
||||
}
|
||||
}
|
||||
}
|
||||
return { message: turn.buildAssistantMessage(), usage: turn.usage };
|
||||
}
|
||||
}
|
||||
|
||||
/** Resolve the API key for the current provider from settings + app config. */
|
||||
export function resolveApiKey(settings: Settings, appConfig: AppConfig): string {
|
||||
const provider = settings.provider;
|
||||
let apiKey = settings.api_keys[provider] ?? "";
|
||||
if (apiKey === "") {
|
||||
const cfg = appConfig.providers[provider];
|
||||
if (cfg) {
|
||||
apiKey =
|
||||
(cfg.api_key_env ? process.env[cfg.api_key_env] : undefined) ??
|
||||
cfg.default_api_key ??
|
||||
"";
|
||||
}
|
||||
}
|
||||
return apiKey;
|
||||
}
|
||||
@@ -0,0 +1,151 @@
|
||||
/**
|
||||
* JSON file–backed `AppConfigRepository` with Claude credential auto-detection.
|
||||
* Mirrors `apps/infrastructure/src/persistence/cms/app_config_repo.rs`.
|
||||
*/
|
||||
import * as fs from "node:fs";
|
||||
import * as os from "node:os";
|
||||
import * as path from "node:path";
|
||||
import {
|
||||
type AppConfig,
|
||||
type AppConfigRepository,
|
||||
type ModelRole,
|
||||
type ProviderConfig,
|
||||
newAppConfig,
|
||||
} from "@zesdex/domain";
|
||||
import { writeJsonAtomic } from "../../utils.ts";
|
||||
|
||||
interface ClaudeEnv {
|
||||
ANTHROPIC_BASE_URL?: string;
|
||||
anthropic_base_url?: string;
|
||||
ANTHROPIC_API_KEY?: string;
|
||||
anthropic_api_key?: string;
|
||||
}
|
||||
|
||||
interface ClaudeSettings {
|
||||
env?: ClaudeEnv;
|
||||
customModel?: string;
|
||||
model?: string;
|
||||
}
|
||||
|
||||
function claudeSettingsFromFile(): ClaudeSettings | null {
|
||||
const home = os.homedir();
|
||||
const file = path.join(home, ".claude", "settings.json");
|
||||
try {
|
||||
const raw = fs.readFileSync(file, "utf8");
|
||||
return JSON.parse(raw) as ClaudeSettings;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function detectClaudeSettingsProvider(): { provider: ProviderConfig; customModel: string | null } | null {
|
||||
const settings = claudeSettingsFromFile();
|
||||
|
||||
let fileCreds: [string, string] | null = null;
|
||||
if (settings?.env) {
|
||||
const baseUrl = settings.env.ANTHROPIC_BASE_URL ?? settings.env.anthropic_base_url;
|
||||
const key = settings.env.ANTHROPIC_API_KEY ?? settings.env.anthropic_api_key;
|
||||
if (baseUrl && key) fileCreds = [baseUrl, key];
|
||||
}
|
||||
|
||||
const envCreds = (() => {
|
||||
const baseUrl = process.env["ANTHROPIC_BASE_URL"];
|
||||
const key = process.env["ANTHROPIC_API_KEY"];
|
||||
if (baseUrl && key) return [baseUrl, key] as [string, string];
|
||||
return null;
|
||||
})();
|
||||
|
||||
const customModel = settings?.customModel ?? settings?.model ?? null;
|
||||
|
||||
const creds = fileCreds ?? envCreds;
|
||||
if (!creds) return null;
|
||||
const [baseUrl, key] = creds;
|
||||
|
||||
return {
|
||||
provider: {
|
||||
api_base: baseUrl,
|
||||
api_key_env: "ANTHROPIC_API_KEY",
|
||||
default_model: customModel ?? undefined,
|
||||
default_api_key: key,
|
||||
},
|
||||
customModel,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Apply a detected Claude provider + custom model onto an AppConfig.
|
||||
* Pure (no I/O). Always inserts "claude", registers known roles, and sets
|
||||
* default_provider/model to Claude/Opus.
|
||||
*/
|
||||
export function applyClaudeProvider(
|
||||
cfg: AppConfig,
|
||||
claudeProvider: ProviderConfig,
|
||||
customModel: string | null,
|
||||
): void {
|
||||
cfg.providers.claude = claudeProvider;
|
||||
|
||||
const claudeModels: [string, string][] = [
|
||||
["claude-opus-5", "claude-opus-5"],
|
||||
["claude-sonnet-5", "claude-sonnet-5"],
|
||||
["claude-haiku-4-5", "claude-haiku-4-5-20251001"],
|
||||
];
|
||||
for (const [roleName, modelName] of claudeModels) {
|
||||
if (!cfg.model_roles[roleName]) {
|
||||
cfg.model_roles[roleName] = {
|
||||
provider: "claude",
|
||||
model: modelName,
|
||||
max_tokens: 8192,
|
||||
context_window: 200_000,
|
||||
temperature: 0.7,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
if (customModel) {
|
||||
cfg.model_roles[customModel] = cfg.model_roles[customModel] ?? {
|
||||
provider: "claude",
|
||||
model: customModel,
|
||||
max_tokens: 8192,
|
||||
context_window: 200_000,
|
||||
temperature: 0.7,
|
||||
} satisfies ModelRole;
|
||||
}
|
||||
|
||||
cfg.default_provider = "claude";
|
||||
cfg.default_model = customModel ?? "claude-opus-5";
|
||||
}
|
||||
|
||||
/** File-based `AppConfigRepository` that reads/writes `app_config.json`. */
|
||||
export class JsonAppConfigRepository implements AppConfigRepository {
|
||||
async load(baseDir: string): Promise<AppConfig> {
|
||||
const file = path.join(baseDir, "app_config.json");
|
||||
let cfg: AppConfig;
|
||||
try {
|
||||
const raw = fs.readFileSync(file, "utf8");
|
||||
cfg = JSON.parse(raw) as AppConfig;
|
||||
} catch (e) {
|
||||
const code = (e as NodeJS.ErrnoException).code;
|
||||
if (code === "ENOENT") cfg = newAppConfig();
|
||||
else throw e;
|
||||
}
|
||||
|
||||
// Merge in default providers that are missing.
|
||||
const defaults = newAppConfig();
|
||||
for (const [name, provider] of Object.entries(defaults.providers)) {
|
||||
if (!cfg.providers[name]) cfg.providers[name] = provider;
|
||||
}
|
||||
|
||||
const detected = detectClaudeSettingsProvider();
|
||||
if (detected) {
|
||||
applyClaudeProvider(cfg, detected.provider, detected.customModel);
|
||||
}
|
||||
|
||||
return cfg;
|
||||
}
|
||||
|
||||
async save(baseDir: string, config: AppConfig): Promise<void> {
|
||||
fs.mkdirSync(baseDir, { recursive: true });
|
||||
const file = path.join(baseDir, "app_config.json");
|
||||
writeJsonAtomic(file, config);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
/** JSON file–backed `ConversationRepository`. Path: `<session_dir>/conversation.json`. */
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { type Conversation, type ConversationRepository, newConversation } from "@zesdex/domain";
|
||||
import { writeJsonAtomic } from "../../utils.ts";
|
||||
|
||||
/** Persists `Conversation` as JSON at `<session_dir>/conversation.json`. */
|
||||
export class JsonConversationRepository implements ConversationRepository {
|
||||
async load(sessionDir: string): Promise<Conversation> {
|
||||
const file = path.join(sessionDir, "conversation.json");
|
||||
try {
|
||||
const raw = fs.readFileSync(file, "utf8");
|
||||
return JSON.parse(raw) as Conversation;
|
||||
} catch (e) {
|
||||
const code = (e as NodeJS.ErrnoException).code;
|
||||
if (code === "ENOENT") {
|
||||
return newConversation("", path.basename(sessionDir));
|
||||
}
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
async save(sessionDir: string, conversation: Conversation): Promise<void> {
|
||||
fs.mkdirSync(sessionDir, { recursive: true });
|
||||
const file = path.join(sessionDir, "conversation.json");
|
||||
writeJsonAtomic(file, conversation);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,63 @@
|
||||
/** JSONL file–backed `EditLogRepository`. Stores `EditLog` as append-only NDJSON. */
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import {
|
||||
type EditLogEntry,
|
||||
type EditLogRepository,
|
||||
MAX_MEMORY_ENTRIES,
|
||||
EditLog,
|
||||
} from "@zesdex/domain";
|
||||
|
||||
/** File-based `EditLogRepository` that reads/writes `edits.jsonl`. */
|
||||
export class JsonlEditLogRepository implements EditLogRepository {
|
||||
private loadFromDisk(file: string): EditLogEntry[] {
|
||||
try {
|
||||
const raw = fs.readFileSync(file, "utf8");
|
||||
const entries: EditLogEntry[] = [];
|
||||
for (const line of raw.split("\n")) {
|
||||
if (!line) continue;
|
||||
try {
|
||||
const entry = JSON.parse(line) as EditLogEntry;
|
||||
if (entries.length >= MAX_MEMORY_ENTRIES) entries.shift();
|
||||
entries.push(entry);
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
return entries;
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
}
|
||||
|
||||
async open(sessionDir: string): Promise<EditLog> {
|
||||
const file = path.join(sessionDir, "edits.jsonl");
|
||||
fs.mkdirSync(path.dirname(file), { recursive: true });
|
||||
const entries = this.loadFromDisk(file);
|
||||
if (!fs.existsSync(file)) {
|
||||
fs.openSync(file, "a");
|
||||
}
|
||||
const log = new EditLog();
|
||||
log.entries = entries;
|
||||
return log;
|
||||
}
|
||||
|
||||
async append(sessionDir: string, log: EditLog, entry: EditLogEntry): Promise<void> {
|
||||
const file = path.join(sessionDir, "edits.jsonl");
|
||||
const line = `${JSON.stringify(entry)}\n`;
|
||||
fs.mkdirSync(path.dirname(file), { recursive: true });
|
||||
const fd = fs.openSync(file, "a");
|
||||
try {
|
||||
fs.writeFileSync(fd, line);
|
||||
fs.fsyncSync(fd);
|
||||
} finally {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
log.entries.push(entry);
|
||||
if (log.entries.length > MAX_MEMORY_ENTRIES) log.entries.shift();
|
||||
}
|
||||
|
||||
entries(log: EditLog): EditLogEntry[] {
|
||||
return [...log.entries];
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
/** CMS persistence — file/JSON-backed repositories. */
|
||||
export * from "./settings_repo.ts";
|
||||
export * from "./app_config_repo.ts";
|
||||
export * from "./conversation_repo.ts";
|
||||
export * from "./edit_log_repo.ts";
|
||||
export * from "./memory_repo.ts";
|
||||
export * from "./rewind_blob_repo.ts";
|
||||
@@ -0,0 +1,107 @@
|
||||
/**
|
||||
* Markdown file–backed `MemoryRepository`. Each memory is a `.md` file with
|
||||
* YAML-ish frontmatter. Mirrors `apps/infrastructure/src/persistence/cms/memory_repo.rs`.
|
||||
*/
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { type Memory, type MemoryRepository, memoryPath } from "@zesdex/domain";
|
||||
import { randomUUID } from "node:crypto";
|
||||
|
||||
/** Escape newlines so they do not break line-oriented frontmatter. */
|
||||
function escapeNewlines(s: string): string {
|
||||
return s.replace(/\n/g, "\\n");
|
||||
}
|
||||
function unescapeNewlines(s: string): string {
|
||||
return s.replace(/\\n/g, "\n");
|
||||
}
|
||||
|
||||
function buildFrontmatter(memory: Memory): string {
|
||||
const outcome = memory.outcome ? `outcome: ${escapeNewlines(memory.outcome)}\n` : "";
|
||||
const scope = memory.scope ? `scope: ${escapeNewlines(memory.scope)}\n` : "";
|
||||
const before = memory.before_snippet ? `before: ${escapeNewlines(memory.before_snippet)}\n` : "";
|
||||
const after = memory.after_snippet ? `after: ${escapeNewlines(memory.after_snippet)}\n` : "";
|
||||
const prov = memory.provenances.length > 0 ? `provenances: ${memory.provenances.join(", ")}\n` : "";
|
||||
return `name: ${memory.name}\ndescription: ${memory.description}\nkind: ${memory.kind}\ncreated_at: ${memory.created_at}\nupdated_at: ${memory.updated_at}\nlifecycle: ${memory.lifecycle}\n${outcome}${scope}${before}${after}${prov}`;
|
||||
}
|
||||
|
||||
function parseFrontmatter(front: string): Record<string, string> {
|
||||
const map: Record<string, string> = {};
|
||||
for (const line of front.split("\n")) {
|
||||
const idx = line.indexOf(":");
|
||||
if (idx === -1) continue;
|
||||
map[line.slice(0, idx).trim()] = line.slice(idx + 1).trim();
|
||||
}
|
||||
return map;
|
||||
}
|
||||
|
||||
function parseMemory(content: string): Memory {
|
||||
let body = content;
|
||||
if (body.startsWith("---\n")) body = body.slice(4);
|
||||
const parts = body.split("\n---\n", 2);
|
||||
const frontPart = parts[0];
|
||||
const mdPart = parts[1];
|
||||
if (frontPart === undefined || mdPart === undefined) throw new Error("missing frontmatter");
|
||||
const front = parseFrontmatter(frontPart);
|
||||
const md = mdPart.trim();
|
||||
|
||||
const prov = front.provenances ?? "";
|
||||
return {
|
||||
name: front.name ?? "",
|
||||
description: front.description ?? "",
|
||||
content: md,
|
||||
kind: front.kind ?? "reference",
|
||||
created_at: Number(front.created_at) || 0,
|
||||
updated_at: Number(front.updated_at) || 0,
|
||||
outcome: front.outcome ? unescapeNewlines(front.outcome) : undefined,
|
||||
lifecycle: front.lifecycle ?? "new",
|
||||
scope: front.scope ? unescapeNewlines(front.scope) : undefined,
|
||||
before_snippet: front.before ? unescapeNewlines(front.before) : undefined,
|
||||
after_snippet: front.after ? unescapeNewlines(front.after) : undefined,
|
||||
provenances: prov === "" ? [] : prov.split(", "),
|
||||
};
|
||||
}
|
||||
|
||||
/** File-based `MemoryRepository` storing memories as `.md` files. */
|
||||
export class MarkdownMemoryRepository implements MemoryRepository {
|
||||
async list(memoryDir: string): Promise<string[]> {
|
||||
let entries: string[];
|
||||
try {
|
||||
entries = fs.readdirSync(memoryDir);
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
return entries
|
||||
.filter((n) => n.endsWith(".md") && n !== "MEMORY.md")
|
||||
.map((n) => n.slice(0, -3));
|
||||
}
|
||||
|
||||
async load(memoryDir: string, name: string): Promise<Memory> {
|
||||
const file = memoryPath(memoryDir, name);
|
||||
const content = fs.readFileSync(file, "utf8");
|
||||
try {
|
||||
return parseMemory(content);
|
||||
} catch (e) {
|
||||
throw new Error(`failed to parse memory '${name}': ${(e as Error).message}`);
|
||||
}
|
||||
}
|
||||
|
||||
async save(memoryDir: string, memory: Memory): Promise<void> {
|
||||
const file = memoryPath(memoryDir, memory.name);
|
||||
fs.mkdirSync(path.dirname(file), { recursive: true });
|
||||
const content = `---\n${buildFrontmatter(memory)}---\n\n${memory.content}`;
|
||||
const tmp = path.join(path.dirname(file), `.${randomUUID()}.tmp`);
|
||||
const fd = fs.openSync(tmp, "w");
|
||||
try {
|
||||
fs.writeFileSync(fd, content);
|
||||
fs.fsyncSync(fd);
|
||||
} finally {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
fs.renameSync(tmp, file);
|
||||
}
|
||||
|
||||
async delete(memoryDir: string, name: string): Promise<void> {
|
||||
const file = memoryPath(memoryDir, name);
|
||||
if (fs.existsSync(file)) fs.unlinkSync(file);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,89 @@
|
||||
/**
|
||||
* Filesystem-backed `RewindBlobRepository`. Blobs at
|
||||
* `<session_dir>/blobs/<hex(key)>.bin` + `index.jsonl` metadata.
|
||||
*/
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import * as crypto from "node:crypto";
|
||||
import { type RewindBlobRepository } from "@zesdex/domain";
|
||||
|
||||
interface BlobIndexEntry {
|
||||
key: string;
|
||||
mime_type?: string;
|
||||
created_at: number;
|
||||
}
|
||||
|
||||
function blobsDir(sessionDir: string): string {
|
||||
return path.join(sessionDir, "blobs");
|
||||
}
|
||||
function blobFilePath(sessionDir: string, blobKey: string): string {
|
||||
const hex = crypto.createHash("sha256").update(blobKey, "utf8").digest("hex");
|
||||
return path.join(blobsDir(sessionDir), `${hex}.bin`);
|
||||
}
|
||||
function indexPath(sessionDir: string): string {
|
||||
return path.join(blobsDir(sessionDir), "index.jsonl");
|
||||
}
|
||||
|
||||
/** Concrete filesystem rewind-blob repository. */
|
||||
export class FileRewindBlobRepository implements RewindBlobRepository {
|
||||
async storeBlob(sessionDir: string, blobKey: string, data: Uint8Array, mimeType?: string): Promise<void> {
|
||||
const dir = blobsDir(sessionDir);
|
||||
fs.mkdirSync(dir, { recursive: true });
|
||||
|
||||
const file = blobFilePath(sessionDir, blobKey);
|
||||
const tmp = `${file}.tmp`;
|
||||
fs.writeFileSync(tmp, data);
|
||||
const fd = fs.openSync(tmp, "r");
|
||||
try {
|
||||
fs.fsyncSync(fd);
|
||||
} finally {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
fs.renameSync(tmp, file);
|
||||
|
||||
const entry: BlobIndexEntry = { key: blobKey, mime_type: mimeType, created_at: Date.now() };
|
||||
const idx = indexPath(sessionDir);
|
||||
const out = fs.openSync(idx, "a");
|
||||
try {
|
||||
fs.writeFileSync(out, `${JSON.stringify(entry)}\n`);
|
||||
fs.fsyncSync(out);
|
||||
} finally {
|
||||
fs.closeSync(out);
|
||||
}
|
||||
}
|
||||
|
||||
async retrieveBlob(sessionDir: string, blobKey: string): Promise<Uint8Array | null> {
|
||||
const file = blobFilePath(sessionDir, blobKey);
|
||||
if (!fs.existsSync(file)) return null;
|
||||
return new Uint8Array(fs.readFileSync(file));
|
||||
}
|
||||
|
||||
async listBlobKeys(sessionDir: string): Promise<string[]> {
|
||||
const idx = indexPath(sessionDir);
|
||||
let content: string;
|
||||
try {
|
||||
content = fs.readFileSync(idx, "utf8");
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
|
||||
const firstSeen: string[] = [];
|
||||
const latest: Record<string, BlobIndexEntry> = {};
|
||||
for (const line of content.split("\n")) {
|
||||
if (!line) continue;
|
||||
try {
|
||||
const entry = JSON.parse(line) as BlobIndexEntry;
|
||||
if (!(entry.key in latest)) firstSeen.push(entry.key);
|
||||
latest[entry.key] = entry;
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
const entries = firstSeen
|
||||
.map((k) => latest[k])
|
||||
.filter(Boolean)
|
||||
.sort((a, b) => a!.created_at - b!.created_at);
|
||||
return entries.map((e) => e!.key);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
/** JSON file–backed `SettingsRepository`. Path: `<base_dir>/settings.json`. */
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { type Settings, type SettingsRepository, newSettings } from "@zesdex/domain";
|
||||
import { writeJsonAtomic } from "../../utils.ts";
|
||||
|
||||
/** Persists `Settings` as pretty-printed JSON at `<base_dir>/settings.json`. */
|
||||
export class JsonSettingsRepository implements SettingsRepository {
|
||||
async load(baseDir: string): Promise<Settings> {
|
||||
const file = path.join(baseDir, "settings.json");
|
||||
try {
|
||||
const raw = fs.readFileSync(file, "utf8");
|
||||
try {
|
||||
return JSON.parse(raw) as Settings;
|
||||
} catch {
|
||||
console.warn(`settings.json at '${file}' failed to parse; falling back to defaults`);
|
||||
return newSettings();
|
||||
}
|
||||
} catch (e) {
|
||||
const code = (e as NodeJS.ErrnoException).code;
|
||||
if (code === "ENOENT") return newSettings();
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
async save(baseDir: string, settings: Settings): Promise<void> {
|
||||
fs.mkdirSync(baseDir, { recursive: true });
|
||||
const file = path.join(baseDir, "settings.json");
|
||||
writeJsonAtomic(file, settings);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,4 @@
|
||||
/** IAM persistence — filesystem sessions, locks, tokens. */
|
||||
export * from "./session_repo.ts";
|
||||
export * from "./session_lock_repo.ts";
|
||||
export * from "./oauth_repo.ts";
|
||||
@@ -0,0 +1,21 @@
|
||||
/**
|
||||
* Filesystem-backed `OAuthRepository`. Token written to JSON with mode 0o600.
|
||||
*/
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { type OAuthRepository, type OAuthToken } from "@zesdex/domain";
|
||||
import { writeJsonAtomic } from "../../utils.ts";
|
||||
|
||||
/** Concrete filesystem OAuth token repository. */
|
||||
export class FileSystemOAuthRepository implements OAuthRepository {
|
||||
async saveToken(file: string, token: OAuthToken): Promise<void> {
|
||||
fs.mkdirSync(path.dirname(file), { recursive: true });
|
||||
writeJsonAtomic(file, token, 0o600);
|
||||
}
|
||||
|
||||
async loadToken(file: string): Promise<OAuthToken | null> {
|
||||
if (!fs.existsSync(file)) return null;
|
||||
const data = fs.readFileSync(file, "utf8");
|
||||
return JSON.parse(data) as OAuthToken;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,78 @@
|
||||
/**
|
||||
* Filesystem-backed `SessionLockRepository` using a PID file with atomic
|
||||
* `O_CREAT|O_EXCL` acquisition. Mirrors `session_lock_repo.rs`.
|
||||
*/
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { type SessionLockRepository, pidIsAlive } from "@zesdex/domain";
|
||||
|
||||
/** Concrete filesystem session-lock repository. */
|
||||
export class FileSystemSessionLockRepository implements SessionLockRepository {
|
||||
tryLock(sessionDir: string): boolean {
|
||||
const file = path.join(sessionDir, ".lock");
|
||||
const pid = process.pid;
|
||||
|
||||
// Phase 1: atomic create.
|
||||
try {
|
||||
const fd = fs.openSync(file, "wx");
|
||||
try {
|
||||
fs.writeFileSync(fd, String(pid));
|
||||
fs.fsyncSync(fd);
|
||||
} finally {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
return true;
|
||||
} catch (e) {
|
||||
const code = (e as NodeJS.ErrnoException).code;
|
||||
if (code !== "EEXIST") throw e;
|
||||
}
|
||||
|
||||
// Phase 2: liveness check.
|
||||
const content = fs.readFileSync(file, "utf8").trim();
|
||||
const existing = Number(content);
|
||||
if (Number.isFinite(existing) && existing > 0 && this.isAlive(existing)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Phase 3: stale lock overwrite.
|
||||
const tmp = `${file}.tmp`;
|
||||
let tmpFd: number;
|
||||
try {
|
||||
tmpFd = fs.openSync(tmp, "wx");
|
||||
} catch {
|
||||
throw new Error("another process is replacing the lock");
|
||||
}
|
||||
try {
|
||||
fs.writeFileSync(tmpFd, String(pid));
|
||||
fs.fsyncSync(tmpFd);
|
||||
} finally {
|
||||
fs.closeSync(tmpFd);
|
||||
}
|
||||
fs.renameSync(tmp, file);
|
||||
try {
|
||||
const parent = path.dirname(file);
|
||||
const dirFd = fs.openSync(parent, "r");
|
||||
try {
|
||||
fs.fsyncSync(dirFd);
|
||||
} finally {
|
||||
fs.closeSync(dirFd);
|
||||
}
|
||||
} catch {
|
||||
/* best-effort */
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
unlock(sessionDir: string): void {
|
||||
const file = path.join(sessionDir, ".lock");
|
||||
try {
|
||||
fs.unlinkSync(file);
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
}
|
||||
|
||||
isAlive(pid: number): boolean {
|
||||
return pidIsAlive(pid);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
/**
|
||||
* Filesystem-backed `SessionRepository`. Each session is
|
||||
* `<base_dir>/sessions/<id>/session.json`.
|
||||
*/
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import {
|
||||
type Session,
|
||||
type SessionId,
|
||||
type SessionRepository,
|
||||
newSessionId,
|
||||
notFound,
|
||||
} from "@zesdex/domain";
|
||||
import { writeJsonAtomic } from "../../utils.ts";
|
||||
|
||||
/** Concrete filesystem session repository. */
|
||||
export class FileSystemSessionRepository implements SessionRepository {
|
||||
async listSessions(baseDir: string): Promise<Session[]> {
|
||||
const dir = path.join(baseDir, "sessions");
|
||||
let entries: string[];
|
||||
try {
|
||||
entries = fs.readdirSync(dir, { withFileTypes: true })
|
||||
.filter((e) => e.isDirectory())
|
||||
.map((e) => e.name);
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
const sessions: Session[] = [];
|
||||
for (const name of entries) {
|
||||
const res = newSessionId(name);
|
||||
if (!res.ok) continue;
|
||||
try {
|
||||
sessions.push(await this.loadSession(baseDir, res.value));
|
||||
} catch {
|
||||
/* skip unreadable session */
|
||||
}
|
||||
}
|
||||
return sessions;
|
||||
}
|
||||
|
||||
async loadSession(baseDir: string, id: SessionId): Promise<Session> {
|
||||
const file = path.join(baseDir, "sessions", id, "session.json");
|
||||
if (!fs.existsSync(file)) {
|
||||
throw notFound(`session not found: ${id}`);
|
||||
}
|
||||
const data = fs.readFileSync(file, "utf8");
|
||||
return JSON.parse(data) as Session;
|
||||
}
|
||||
|
||||
async saveSession(baseDir: string, session: Session): Promise<void> {
|
||||
const dir = path.join(baseDir, "sessions", session.id);
|
||||
fs.mkdirSync(dir, { recursive: true });
|
||||
const file = path.join(dir, "session.json");
|
||||
writeJsonAtomic(file, session);
|
||||
}
|
||||
|
||||
async deleteSession(baseDir: string, id: SessionId): Promise<void> {
|
||||
const dir = path.join(baseDir, "sessions", id);
|
||||
if (fs.existsSync(dir)) fs.rmSync(dir, { recursive: true, force: true });
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
/** Persistence layer — CMS + IAM file repositories. */
|
||||
export * from "./cms/index.ts";
|
||||
export * from "./iam/index.ts";
|
||||
@@ -0,0 +1,120 @@
|
||||
/**
|
||||
* Utility helpers. Mirrors `apps/infrastructure/src/utils.rs`.
|
||||
*/
|
||||
import { randomUUID } from "node:crypto";
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
|
||||
/** Atomically write JSON `data` to `path`, with optional mode and parent fsync. */
|
||||
export function writeJsonAtomic(filePath: string, data: unknown, mode?: number): void {
|
||||
const tmp = `${filePath}.tmp.${randomUUID()}`;
|
||||
const parent = path.dirname(filePath);
|
||||
try {
|
||||
fs.mkdirSync(parent, { recursive: true });
|
||||
const bytes = JSON.stringify(data, null, 2);
|
||||
const fd = fs.openSync(tmp, "w");
|
||||
try {
|
||||
fs.writeFileSync(fd, bytes);
|
||||
fs.fsyncSync(fd);
|
||||
} finally {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
if (mode !== undefined) fs.chmodSync(tmp, mode);
|
||||
fs.renameSync(tmp, filePath);
|
||||
try {
|
||||
const dirFd = fs.openSync(parent, "r");
|
||||
try {
|
||||
fs.fsyncSync(dirFd);
|
||||
} finally {
|
||||
fs.closeSync(dirFd);
|
||||
}
|
||||
} catch {
|
||||
/* best-effort */
|
||||
}
|
||||
} catch (e) {
|
||||
try {
|
||||
fs.unlinkSync(tmp);
|
||||
} catch {
|
||||
/* orphan cleanup */
|
||||
}
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
/** Convert a string into a filesystem-safe slug, or null if empty/too long. */
|
||||
export function slugify(s: string): string | null {
|
||||
let slug = s.toLowerCase().replace(/[^a-z0-9]/g, "-");
|
||||
slug = slug.split("-").filter((seg) => seg !== "").join("-");
|
||||
if (slug === "" || slug.length > 80) return null;
|
||||
return slug;
|
||||
}
|
||||
|
||||
/** Write `text` to the terminal clipboard via OSC-52 escape sequence. */
|
||||
export function writeOsc52(text: string, out: { write: (s: string) => unknown } = process.stdout): void {
|
||||
const encoded = Buffer.from(text, "utf8").toString("base64");
|
||||
out.write(`\x1b]52;c;${encoded}\x1b\\`);
|
||||
}
|
||||
|
||||
/** Truncate a string to at most `maxChars` characters (UTF-8 safe). */
|
||||
export function truncateChars(s: string, maxChars: number): string {
|
||||
if ([...s].length <= maxChars) return s;
|
||||
return [...s].slice(0, maxChars).join("");
|
||||
}
|
||||
|
||||
/** Build a string directory tree of `root`, respecting .gitignore semantics loosely. */
|
||||
export function buildWorkspaceTree(root: string, maxFiles: number): string {
|
||||
const lines: string[] = [];
|
||||
let count = 0;
|
||||
const walk = (dir: string, depth: number, ignore: Set<string>) => {
|
||||
if (count >= maxFiles) {
|
||||
lines.push("... (truncated)");
|
||||
return;
|
||||
}
|
||||
let entries: fs.Dirent[] = [];
|
||||
try {
|
||||
entries = fs.readdirSync(dir, { withFileTypes: true }).sort((a, b) => (a.isDirectory() === b.isDirectory() ? a.name.localeCompare(b.name) : a.isDirectory() ? -1 : 1));
|
||||
} catch {
|
||||
return;
|
||||
}
|
||||
for (const entry of entries) {
|
||||
if (count >= maxFiles) {
|
||||
lines.push("... (truncated)");
|
||||
return;
|
||||
}
|
||||
if (entry.name === ".git" || ignore.has(entry.name)) continue;
|
||||
const indent = " ".repeat(Math.max(depth - 1, 0));
|
||||
if (entry.isDirectory()) {
|
||||
lines.push(`${indent}${entry.name}/`);
|
||||
count++;
|
||||
const subIgnore = new Set(ignore);
|
||||
// Honor .gitignore briefly
|
||||
try {
|
||||
const gi = fs.readFileSync(path.join(dir, entry.name, ".gitignore"), "utf8");
|
||||
for (const l of gi.split("\n")) {
|
||||
const t = l.trim().replace(/\/$/, "");
|
||||
if (t && !t.startsWith("#")) subIgnore.add(t);
|
||||
}
|
||||
} catch {
|
||||
/* keep default ignore set */
|
||||
}
|
||||
walk(path.join(dir, entry.name), depth + 1, subIgnore);
|
||||
} else {
|
||||
lines.push(`${indent}${entry.name}`);
|
||||
count++;
|
||||
}
|
||||
}
|
||||
};
|
||||
lines.push(".");
|
||||
walk(root, 1, new Set(["node_modules", "target"]));
|
||||
return lines.join("\n").trimEnd();
|
||||
}
|
||||
|
||||
/** Run a shell command and return trimmed stdout (best-effort). */
|
||||
export function runCommand(cmd: string, cwd?: string, args?: string[]): string {
|
||||
try {
|
||||
const res = Bun.spawnSync({ cmd: args ? [cmd, ...args] : [cmd], cwd, stdout: "pipe", stderr: "pipe" });
|
||||
return res.stdout.toString().trim();
|
||||
} catch {
|
||||
return "";
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user