import { DatabaseSync } from "node:sqlite"; import { chmodSync } from "node:fs"; import type { BaseEvent } from "../types.js"; const AGENTS_SCHEMA = "CREATE TABLE agents (node_id TEXT NOT NULL, task_id TEXT NOT NULL, parent_node_id TEXT, role TEXT, lens TEXT, model TEXT, thinking_level TEXT, topology TEXT, spawned_at TEXT, completed_at TEXT, status TEXT, cost_credits REAL, facts_total INTEGER, facts_used INTEGER, unique_facts_used INTEGER, duplicate_facts INTEGER, output_tokens INTEGER, PRIMARY KEY (task_id, node_id));"; export const PROJECTED_EXPERIMENT_EVENTS = new Set(["experiment.started", "experiment.stopped"]); export const SCHEMA = ` CREATE TABLE IF NOT EXISTS tasks (task_id TEXT PRIMARY KEY, session_id TEXT NOT NULL, started_at TEXT NOT NULL, completed_at TEXT, profile TEXT NOT NULL, privacy_class TEXT NOT NULL, request_export TEXT, request_hash TEXT NOT NULL, task_type TEXT, breadth TEXT, coupling TEXT, uncertainty TEXT, risk TEXT, verifiability TEXT, topology TEXT, policy TEXT, config_version TEXT NOT NULL, success INTEGER, verified_success INTEGER, human_feedback TEXT, rework_count INTEGER DEFAULT 0); CREATE TABLE IF NOT EXISTS model_calls (call_id TEXT PRIMARY KEY, task_id TEXT NOT NULL, node_id TEXT NOT NULL, model TEXT NOT NULL, provider TEXT NOT NULL, thinking_level TEXT, role TEXT, input_tokens INTEGER, cached_input_tokens INTEGER, output_tokens INTEGER, credits REAL, latency_ms INTEGER, success INTEGER, error_class TEXT); ${AGENTS_SCHEMA.replace("CREATE TABLE agents", "CREATE TABLE IF NOT EXISTS agents")} CREATE TABLE IF NOT EXISTS verifications (verification_id TEXT PRIMARY KEY, task_id TEXT NOT NULL, node_id TEXT, command_hash TEXT, fingerprint TEXT, passed INTEGER, duration_ms INTEGER, attempt INTEGER, error_class TEXT); CREATE TABLE IF NOT EXISTS route_decisions (id TEXT PRIMARY KEY, task_id TEXT NOT NULL, proposed_topology TEXT, final_topology TEXT, reason_codes TEXT, fanout_proposed INTEGER, fanout_final INTEGER, model_proposed TEXT, model_final TEXT, overridden INTEGER, override_reason TEXT); CREATE TABLE IF NOT EXISTS config_versions (version TEXT PRIMARY KEY, created_at TEXT NOT NULL, config_json TEXT NOT NULL, config_hash TEXT NOT NULL, parent_version TEXT, change_reason TEXT); CREATE TABLE IF NOT EXISTS experiments (experiment_id TEXT PRIMARY KEY, name TEXT NOT NULL, status TEXT NOT NULL, champion_version TEXT NOT NULL, challenger_version TEXT NOT NULL, allocation REAL NOT NULL, started_at TEXT NOT NULL, stopped_at TEXT, stop_reason TEXT); `; const value = (event: BaseEvent, key: string): string | number | null => { const raw = event[key]; if (typeof raw === "string" || typeof raw === "number") return raw; if (typeof raw === "boolean") return Number(raw); return null; }; const stringArray = (event: BaseEvent, key: string): string[] => Array.isArray(event[key]) ? event[key].filter((entry): entry is string => typeof entry === "string") : []; const jsonObject = (event: BaseEvent, key: string): string | null => { const raw = event[key]; return raw && typeof raw === "object" && !Array.isArray(raw) ? JSON.stringify(raw) : null; }; export class SqliteProjector { readonly db: DatabaseSync; constructor(path: string) { this.db = new DatabaseSync(path); chmodSync(path, 0o600); this.db.exec(SCHEMA); const primaryKey = (this.db.prepare("PRAGMA table_info(agents)").all() as Array<{ name: string; pk: number }>).filter((column) => column.pk > 0).sort((left, right) => left.pk - right.pk).map((column) => column.name).join(","); if (primaryKey !== "task_id,node_id") this.db.exec(`BEGIN; ALTER TABLE agents RENAME TO agents_legacy; ${AGENTS_SCHEMA} INSERT OR REPLACE INTO agents SELECT node_id, task_id, parent_node_id, role, lens, model, thinking_level, topology, spawned_at, completed_at, status, cost_credits, facts_total, facts_used, unique_facts_used, duplicate_facts, output_tokens FROM agents_legacy; DROP TABLE agents_legacy; COMMIT;`); } project(event: BaseEvent): void { const globalEvent = event.eventType === "config.changed" || event.eventType === "budget.updated" || event.eventType === "experiment.started" || event.eventType === "experiment.evaluated" || event.eventType === "experiment.stopped" || event.eventType === "recommendation.accepted" || event.eventType === "recommendation.rejected" || event.eventType === "export.created"; if (!event.eventType.startsWith("session.") && !globalEvent) this.db.prepare("INSERT OR IGNORE INTO tasks (task_id, session_id, started_at, profile, privacy_class, request_hash, config_version) VALUES (?, ?, ?, ?, ?, ?, ?)").run(event.taskId, event.sessionId, event.timestamp, event.profile, value(event, "privacyClass") ?? "restricted", value(event, "requestHash") ?? event.taskId, event.configVersion); if (event.eventType === "request.received") this.db.prepare("UPDATE tasks SET privacy_class=COALESCE(?, privacy_class), request_hash=COALESCE(?, request_hash), config_version=? WHERE task_id=?").run(value(event, "privacyClass"), value(event, "requestHash"), event.configVersion, event.taskId); if (event.eventType === "request.sanitized") this.db.prepare("UPDATE tasks SET request_export=? WHERE task_id=?").run(value(event, "requestExport") ?? null, event.taskId); if (event.eventType === "request.classified") this.db.prepare("UPDATE tasks SET task_type=?, breadth=?, coupling=?, uncertainty=?, risk=?, verifiability=?, topology=?, policy=? WHERE task_id=?").run(value(event, "intent") ?? null, value(event, "breadth") ?? null, value(event, "coupling") ?? null, value(event, "uncertainty") ?? null, value(event, "risk") ?? null, value(event, "verifiability") ?? null, value(event, "topology") ?? null, value(event, "policy") ?? null, event.taskId); if (event.eventType === "result.completed") this.db.prepare("UPDATE tasks SET completed_at=?, success=?, verified_success=? WHERE task_id=?").run(event.timestamp, Number(Boolean(value(event, "success"))), Number(Boolean(value(event, "verified"))), event.taskId); if (event.eventType === "feedback.explicit" || event.eventType === "feedback.implicit" || event.eventType === "user.feedback") { const kind = value(event, "kind"); const implicit = event.eventType === "feedback.implicit" || value(event, "explicit") === 0; this.db.prepare("UPDATE tasks SET human_feedback=COALESCE(?, human_feedback), rework_count=rework_count+? WHERE task_id=?").run(kind ?? (implicit ? "implicit" : null), Number(kind === "fixed" || implicit), event.taskId); } if (event.eventType === "model.call.completed" || event.eventType === "model.call.failed") this.db.prepare("INSERT OR REPLACE INTO model_calls VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)").run(event.eventId, event.taskId, value(event, "nodeId") ?? "root", value(event, "model") ?? "unknown", value(event, "provider") ?? "unknown", value(event, "thinkingLevel") ?? null, value(event, "role") ?? null, value(event, "inputTokens") ?? 0, value(event, "cachedInputTokens") ?? 0, value(event, "outputTokens") ?? 0, value(event, "creditsEstimated") ?? null, value(event, "latencyMs") ?? null, Number(event.eventType === "model.call.completed"), value(event, "errorClass") ?? null); if (event.eventType.startsWith("agent.")) this.projectAgent(event); if (event.eventType === "attribution.created") this.db.prepare("UPDATE agents SET facts_total=?, facts_used=?, unique_facts_used=?, duplicate_facts=? WHERE task_id=? AND node_id=?").run(value(event, "factsTotal") ?? null, value(event, "factsUsed") ?? null, value(event, "uniqueFactsUsed") ?? null, value(event, "duplicateFacts") ?? null, event.taskId, value(event, "nodeId") ?? "root"); if (event.eventType === "verification.completed") this.db.prepare("INSERT OR REPLACE INTO verifications VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)").run(event.eventId, event.taskId, value(event, "nodeId") ?? null, value(event, "commandHash") ?? null, value(event, "fingerprint") ?? null, Number(Boolean(value(event, "passed"))), value(event, "durationMs") ?? null, value(event, "attempt") ?? null, value(event, "errorClass") ?? null); if (event.eventType === "route.validated") this.db.prepare("INSERT OR REPLACE INTO route_decisions VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)").run(event.eventId, event.taskId, value(event, "proposedTopology") ?? null, value(event, "finalTopology") ?? null, JSON.stringify(stringArray(event, "reasonCodes")), value(event, "fanoutProposed") ?? null, value(event, "fanoutFinal") ?? null, value(event, "modelProposed") ?? null, value(event, "modelFinal") ?? null, Number(Boolean(value(event, "overridden"))), value(event, "overrideReason") ?? null); if (event.eventType === "root.selection.completed") { const actual = value(event, "actualModel"); const planned = value(event, "plannedModel"); if (typeof actual === "string") this.db.prepare("UPDATE route_decisions SET model_final=?, overridden=?, override_reason=? WHERE task_id=? AND final_topology='direct'").run(actual, Number(actual !== planned), actual === planned ? null : value(event, "outcome"), event.taskId); } if (event.eventType === "config.changed") this.projectConfig(event); if (PROJECTED_EXPERIMENT_EVENTS.has(event.eventType)) this.projectExperiment(event); } private projectAgent(event: BaseEvent): void { const node = String(value(event, "nodeId") ?? event.eventId); if (event.eventType === "agent.spawn.proposed") this.db.prepare("INSERT OR REPLACE INTO agents (node_id, task_id, parent_node_id, role, lens, model, thinking_level, topology, spawned_at, status) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)").run(node, event.taskId, value(event, "parentNodeId") ?? null, value(event, "role") ?? null, value(event, "lens") ?? null, value(event, "model") ?? null, value(event, "thinkingLevel") ?? null, value(event, "topology") ?? null, event.timestamp, "proposed"); if (event.eventType === "agent.spawned") this.db.prepare("INSERT INTO agents (node_id, task_id, parent_node_id, role, lens, model, thinking_level, topology, spawned_at, status) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(task_id, node_id) DO UPDATE SET spawned_at=excluded.spawned_at, status='running', parent_node_id=COALESCE(agents.parent_node_id, excluded.parent_node_id), role=COALESCE(agents.role, excluded.role), lens=COALESCE(agents.lens, excluded.lens), model=COALESCE(agents.model, excluded.model), thinking_level=COALESCE(agents.thinking_level, excluded.thinking_level), topology=COALESCE(agents.topology, excluded.topology)").run(node, event.taskId, value(event, "parentNodeId") ?? null, value(event, "role") ?? null, value(event, "lens") ?? null, value(event, "model") ?? null, value(event, "thinkingLevel") ?? null, value(event, "topology") ?? null, event.timestamp, "running"); if (["agent.completed", "agent.failed", "agent.cancelled"].includes(event.eventType)) this.db.prepare("UPDATE agents SET completed_at=?, status=?, cost_credits=?, facts_total=?, output_tokens=? WHERE task_id=? AND node_id=?").run(event.timestamp, event.eventType.slice("agent.".length), value(event, "costCredits") ?? null, value(event, "factsTotal") ?? null, value(event, "outputTokens") ?? null, event.taskId, node); } private projectConfig(event: BaseEvent): void { const version = value(event, "version") ?? event.configVersion; const configJson = jsonObject(event, "configJson"); const configHash = value(event, "configHash"); if (typeof version !== "string" || !configJson || typeof configHash !== "string") return; this.db.prepare("INSERT OR REPLACE INTO config_versions VALUES (?, ?, ?, ?, ?, ?)").run(version, value(event, "createdAt") ?? event.timestamp, configJson, configHash, value(event, "parentVersion"), value(event, "changeReason")); } private projectExperiment(event: BaseEvent): void { const id = value(event, "experimentId"); if (typeof id !== "string") return; const fields = [value(event, "name"), value(event, "status"), value(event, "championVersion"), value(event, "challengerVersion"), value(event, "allocation"), value(event, "startedAt")] as const; if (fields.some((field) => field === null)) { if (event.eventType === "experiment.stopped") this.db.prepare("UPDATE experiments SET status='stopped', stopped_at=?, stop_reason=? WHERE experiment_id=?").run(value(event, "stoppedAt"), value(event, "stopReason"), id); return; } this.db.prepare("INSERT INTO experiments VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(experiment_id) DO UPDATE SET name=excluded.name, status=CASE WHEN experiments.status='stopped' OR excluded.status='stopped' THEN 'stopped' ELSE excluded.status END, champion_version=excluded.champion_version, challenger_version=excluded.challenger_version, allocation=excluded.allocation, started_at=excluded.started_at, stopped_at=COALESCE(experiments.stopped_at, excluded.stopped_at), stop_reason=COALESCE(experiments.stop_reason, excluded.stop_reason)").run(id, ...fields, value(event, "stoppedAt"), value(event, "stopReason")); } rebuild(events: readonly BaseEvent[]): void { this.db.exec("DELETE FROM tasks; DELETE FROM model_calls; DELETE FROM agents; DELETE FROM verifications; DELETE FROM route_decisions; DELETE FROM config_versions; DELETE FROM experiments;"); for (const event of events) this.project(event); } close(): void { this.db.close(); } }