import { closeSync, existsSync, openSync, readFileSync, rmSync, statSync, writeFileSync } from "node:fs"; import { dirname } from "node:path"; import { mkdirSync } from "node:fs"; import { atomicWriteText } from "./file-editor.ts"; import { IMPORT_KINDS, type ArchivedEntry, type BlockRecord, type ConfigFormat, type ImportKind, type McpcState } from "./types.ts"; const STATE_VERSION = 1; const STALE_LOCK_MS = 30_000; function isRecord(value: unknown): value is Record { return value !== null && typeof value === "object" && !Array.isArray(value); } function isFormat(value: unknown): value is ConfigFormat { return value === "jsonc" || value === "toml"; } function isImportKind(value: unknown): value is ImportKind { return typeof value === "string" && (IMPORT_KINDS as readonly string[]).includes(value); } function parseArchivedEntry(value: unknown, recordName: string): ArchivedEntry { if (!isRecord(value)) throw new Error(`Invalid archived entry for ${recordName}`); if (typeof value.path !== "string" || !isFormat(value.format)) { throw new Error(`Invalid archived path or format for ${recordName}`); } if (!Array.isArray(value.keyPath) || value.keyPath.length === 0 || value.keyPath.some((key) => typeof key !== "string")) { throw new Error(`Invalid archived key path for ${recordName}`); } if (value.source !== "config" && value.source !== "import") { throw new Error(`Invalid archived source for ${recordName}`); } return { path: value.path, format: value.format, keyPath: [...value.keyPath] as string[], value: structuredClone(value.value), source: value.source, ...(typeof value.sourceId === "string" ? { sourceId: value.sourceId } : {}), ...(isImportKind(value.importKind) ? { importKind: value.importKind } : {}), }; } function parseBlockRecord(value: unknown): BlockRecord { if (!isRecord(value) || typeof value.name !== "string" || !value.name.trim()) { throw new Error("Invalid MCPC block record"); } if (typeof value.blockedAt !== "string" || typeof value.cwd !== "string" || !Array.isArray(value.entries)) { throw new Error(`Invalid MCPC block record for ${value.name}`); } const entries = value.entries.map((entry) => parseArchivedEntry(entry, value.name as string)); if (entries.length === 0) throw new Error(`MCPC block record has no archived entries: ${value.name}`); return { name: value.name, blockedAt: value.blockedAt, cwd: value.cwd, entries, }; } export function emptyState(): McpcState { return { version: STATE_VERSION, records: [] }; } export function loadState(statePath: string): McpcState { if (!existsSync(statePath)) return emptyState(); let parsed: unknown; try { parsed = JSON.parse(readFileSync(statePath, "utf8")); } catch (error) { const message = error instanceof Error ? error.message : String(error); throw new Error(`Failed to read MCPC state ${statePath}: ${message}`, { cause: error }); } if (!isRecord(parsed) || parsed.version !== STATE_VERSION || !Array.isArray(parsed.records)) { throw new Error(`Unsupported or invalid MCPC state file: ${statePath}`); } const records = parsed.records.map(parseBlockRecord); const names = new Set(); for (const record of records) { if (names.has(record.name)) throw new Error(`Duplicate MCPC block record: ${record.name}`); names.add(record.name); } return { version: STATE_VERSION, records }; } export function saveState(statePath: string, state: McpcState): void { atomicWriteText(statePath, `${JSON.stringify(state, null, 2)}\n`, 0o600); } function removeStaleLock(lockPath: string): void { if (!existsSync(lockPath)) return; try { if (Date.now() - statSync(lockPath).mtimeMs > STALE_LOCK_MS) rmSync(lockPath, { force: true }); } catch { // The exclusive open below is the final lock check. } } export function withStateLock(statePath: string, operation: () => T): T { const lockPath = `${statePath}.lock`; mkdirSync(dirname(lockPath), { recursive: true }); removeStaleLock(lockPath); let descriptor: number | undefined; try { descriptor = openSync(lockPath, "wx", 0o600); writeFileSync(descriptor, `${process.pid}\n${new Date().toISOString()}\n`, "utf8"); } catch (error) { if (descriptor !== undefined) { try { closeSync(descriptor); rmSync(lockPath, { force: true }); } catch { // Preserve the lock acquisition error. } } throw new Error(`Another MCPC operation is already in progress (${lockPath})`, { cause: error }); } try { return operation(); } finally { try { closeSync(descriptor); } finally { rmSync(lockPath, { force: true }); } } }