import { randomUUID } from "node:crypto"; import * as fs from "node:fs"; import * as path from "node:path"; import { workerLevelConfigSignature } from "./config-core.ts"; import { collectQualityUpgrades, isLearningEvidence, learningRetentionMs, unexpiredEvidence } from "./learning-core.ts"; import type { ClusterConfig, ClusterState, LearningEvidence } from "./types.ts"; async function readEvidenceFile(filePath: string): Promise { let content: string; try { content = await fs.promises.readFile(filePath, "utf8"); } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; return undefined; } try { const value: unknown = JSON.parse(content); return isLearningEvidence(value) ? value : undefined; } catch { return undefined; } } export async function loadLearningEvidenceAt( evidenceDirectory: string, config: ClusterConfig, now = Date.now(), ): Promise { let entries: fs.Dirent[]; try { entries = await fs.promises.readdir(evidenceDirectory, { withFileTypes: true }); } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return []; throw new Error(`无法读取全局学习证据目录 ${evidenceDirectory}:${error instanceof Error ? error.message : String(error)}`); } const evidence = await Promise.all( entries .filter((entry) => entry.isFile() && entry.name.endsWith(".json")) .map(async (entry) => ({ entry, evidence: await readEvidenceFile(path.join(evidenceDirectory, entry.name)) })), ); const retentionMs = learningRetentionMs(config); const validEvidence: LearningEvidence[] = []; await Promise.all(evidence.map(async ({ entry, evidence: item }) => { if (!item) return; if (item.completedAt < now - retentionMs) { await fs.promises.rm(path.join(evidenceDirectory, entry.name), { force: true }).catch(() => undefined); return; } if (item.completedAt > now || item.configSignature !== workerLevelConfigSignature(config)) return; validEvidence.push(item); })); return unexpiredEvidence(validEvidence, now, retentionMs).sort((left, right) => left.completedAt - right.completedAt || left.projectKey.localeCompare(right.projectKey) || left.runId.localeCompare(right.runId) || left.taskFingerprint.localeCompare(right.taskFingerprint), ); } async function writeEvidenceAt(evidenceDirectory: string, evidence: LearningEvidence): Promise { await fs.promises.mkdir(evidenceDirectory, { recursive: true, mode: 0o700 }); const temporaryPath = path.join(evidenceDirectory, `.evidence-${process.pid}-${randomUUID()}.tmp`); const evidencePath = path.join(evidenceDirectory, `${randomUUID()}.json`); try { await fs.promises.writeFile(temporaryPath, `${JSON.stringify(evidence, null, 2)}\n`, { encoding: "utf8", mode: 0o600 }); await fs.promises.rename(temporaryPath, evidencePath); } finally { await fs.promises.rm(temporaryPath, { force: true }).catch(() => undefined); } } export async function appendLearningEvidenceAt( evidenceDirectory: string, state: ClusterState, config: ClusterConfig, now = Date.now(), ): Promise { if (!config.learning.enabled) return []; const additions = collectQualityUpgrades(state, config, now); await Promise.all(additions.map((evidence) => writeEvidenceAt(evidenceDirectory, evidence))); return additions; } export async function clearLearningEvidenceAt(evidenceDirectory: string): Promise { let entries: fs.Dirent[]; try { entries = await fs.promises.readdir(evidenceDirectory, { withFileTypes: true }); } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return false; throw new Error(`无法检查全局学习证据目录 ${evidenceDirectory}:${error instanceof Error ? error.message : String(error)}`); } if (entries.length === 0) return false; try { await fs.promises.rm(evidenceDirectory, { recursive: true, force: true }); return true; } catch (error) { throw new Error(`无法清除全局学习证据 ${evidenceDirectory}:${error instanceof Error ? error.message : String(error)}`); } }