Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions src/analyze/framework.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ import {
upsertAnalyzerDef,
upsertAnalyzerVersion,
} from "../db/analysis-queries.js";
import { prep } from "../db/prepared.js";
import { materializeProposalsFromNode, applyValidationFromNode } from "./proposal-materializer.js";
import { mapWithConcurrency } from "./concurrency.js";

Expand Down Expand Up @@ -731,8 +732,8 @@ function isDuplicateInputKey(err: unknown): boolean {
}

function db_loadMessages(db: Database.Database, sessionId: string): MessageRow[] {
return db
.prepare(
// Static SQL (two adjacent string literals) — stable text, so safe to cache.
return prep(db,
"SELECT id, session_id, parent_id, timestamp, role, content_text, content_thinking, tool_calls, tool_results " +
"FROM messages WHERE session_id = ? ORDER BY rowid ASC",
)
Expand Down
74 changes: 32 additions & 42 deletions src/db/analysis-queries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
*/

import type Database from "better-sqlite3";
import { prep } from "./prepared.js";
import type {
AnalysisEdgeRow,
AnalysisNodeRow,
Expand All @@ -25,7 +26,7 @@ import { EDGE_KINDS, REF_KINDS } from "../analyze/edge-kinds.js";
// ───────────────────────── analyzer registry ─────────────────────────

export function upsertAnalyzerDef(db: Database.Database, def: AnalyzerDef): void {
db.prepare(`
prep(db, `
INSERT INTO analyzer_defs (id, label, description, anchor_span, dependencies, created_at)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT(id) DO UPDATE SET
Expand All @@ -44,7 +45,7 @@ export function upsertAnalyzerDef(db: Database.Database, def: AnalyzerDef): void
}

export function upsertAnalyzerVersion(db: Database.Database, version: AnalyzerVersion): void {
db.prepare(`
prep(db, `
INSERT INTO analyzer_versions (analyzer_id, version_id, implementation_kind, code_ref, created_at)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(analyzer_id, version_id) DO NOTHING
Expand All @@ -58,7 +59,7 @@ export function upsertAnalyzerVersion(db: Database.Database, version: AnalyzerVe
}

export function registerPrompt(db: Database.Database, prompt: PromptVersion): void {
db.prepare(`
prep(db, `
INSERT INTO prompt_registry (hash, content, role, created_at)
VALUES (?, ?, ?, ?)
ON CONFLICT(hash) DO NOTHING
Expand All @@ -74,8 +75,7 @@ export function resolveConfig(
params: { analyzerId: string; configJson: Record<string, unknown>; label?: string },
): AnalyzerConfig {
const configHash = computeConfigHash(params.configJson);
const existing = db
.prepare("SELECT id, analyzer_id, config_hash, config_json, label FROM analyzer_configs WHERE config_hash = ?")
const existing = prep(db, "SELECT id, analyzer_id, config_hash, config_json, label FROM analyzer_configs WHERE config_hash = ?")
.get(configHash) as
| { id: string; analyzer_id: string; config_hash: string; config_json: string; label: string | null }
| undefined;
Expand All @@ -91,7 +91,7 @@ export function resolveConfig(
}

const id = uuidv7();
db.prepare(`
prep(db, `
INSERT INTO analyzer_configs (id, analyzer_id, config_hash, config_json, label, created_at)
VALUES (?, ?, ?, ?, ?, ?)
`).run(id, params.analyzerId, configHash, JSON.stringify(params.configJson), params.label ?? null, new Date().toISOString());
Expand Down Expand Up @@ -120,7 +120,7 @@ export function createRun(
modelSpec?: string;
},
): void {
db.prepare(`
prep(db, `
INSERT INTO analysis_runs
(id, analyzer_id, analyzer_version_id, config_id, session_id, mode, status, prompt_bundle_hash, model_spec, started_at)
VALUES (?, ?, ?, ?, ?, ?, 'ok', ?, ?, ?)
Expand Down Expand Up @@ -149,7 +149,7 @@ export function finishRun(
errorMessage?: string | null;
},
): void {
db.prepare(`
prep(db, `
UPDATE analysis_runs SET
status = ?, finished_at = ?, nodes_produced = ?, nodes_skipped = ?,
cost_usd = ?, tokens_used = ?, error_message = ?
Expand All @@ -167,7 +167,7 @@ export function finishRun(
}

export function getRun(db: Database.Database, runId: string): AnalysisRunRow | undefined {
return db.prepare("SELECT * FROM analysis_runs WHERE id = ?").get(runId) as AnalysisRunRow | undefined;
return prep(db, "SELECT * FROM analysis_runs WHERE id = ?").get(runId) as AnalysisRunRow | undefined;
}

// ───────────────────────── nodes ─────────────────────────
Expand All @@ -194,7 +194,7 @@ export function insertNode(
createdAt: string;
},
): void {
db.prepare(`
prep(db, `
INSERT INTO analysis_nodes
(id, session_id, analyzer_id, analyzer_version_id, config_id, run_id, node_kind,
content_json, source_set_hash, input_key, output_key, config_fingerprint, model_used, cost_usd, tokens_used, duration_ms, created_at)
Expand All @@ -221,7 +221,7 @@ export function insertNode(
}

export function getNode(db: Database.Database, id: string): AnalysisNodeRow | undefined {
return db.prepare("SELECT * FROM analysis_nodes WHERE id = ?").get(id) as AnalysisNodeRow | undefined;
return prep(db, "SELECT * FROM analysis_nodes WHERE id = ?").get(id) as AnalysisNodeRow | undefined;
}

/**
Expand All @@ -233,12 +233,12 @@ export function getNode(db: Database.Database, id: string): AnalysisNodeRow | un
*/
export function getNodeByOutputKey(db: Database.Database, outputKey: string): AnalysisNodeRow | undefined {
if (!outputKey) return undefined;
return db.prepare("SELECT * FROM analysis_nodes WHERE output_key = ? LIMIT 1").get(outputKey) as AnalysisNodeRow | undefined;
return prep(db, "SELECT * FROM analysis_nodes WHERE output_key = ? LIMIT 1").get(outputKey) as AnalysisNodeRow | undefined;
}

/** Idempotency lookup: a node produced by an exact recipe over an exact source set. */
export function findNodeByInputKey(db: Database.Database, inputKey: string): AnalysisNodeRow | undefined {
return db.prepare("SELECT * FROM analysis_nodes WHERE input_key = ?").get(inputKey) as AnalysisNodeRow | undefined;
return prep(db, "SELECT * FROM analysis_nodes WHERE input_key = ?").get(inputKey) as AnalysisNodeRow | undefined;
}

/**
Expand All @@ -251,35 +251,32 @@ export function findLatestNodeBySourceSet(
analyzerId: string,
sourceSetHash: string,
): AnalysisNodeRow | undefined {
return db
.prepare(
return prep(db,
"SELECT * FROM analysis_nodes WHERE analyzer_id = ? AND source_set_hash = ? AND node_kind != 'error' ORDER BY created_at DESC, rowid DESC LIMIT 1",
)
.get(analyzerId, sourceSetHash) as AnalysisNodeRow | undefined;
}

export function getSessionNodes(db: Database.Database, sessionId: string): AnalysisNodeRow[] {
return db.prepare("SELECT * FROM analysis_nodes WHERE session_id = ? ORDER BY created_at ASC, rowid ASC").all(sessionId) as AnalysisNodeRow[];
return prep(db, "SELECT * FROM analysis_nodes WHERE session_id = ? ORDER BY created_at ASC, rowid ASC").all(sessionId) as AnalysisNodeRow[];
}

/** Every analysis node, for integrity verification. */
export function getAllAnalysisNodes(db: Database.Database): AnalysisNodeRow[] {
return db.prepare("SELECT * FROM analysis_nodes ORDER BY created_at ASC, rowid ASC").all() as AnalysisNodeRow[];
return prep(db, "SELECT * FROM analysis_nodes ORDER BY created_at ASC, rowid ASC").all() as AnalysisNodeRow[];
}

/** A session's messages in stream order — for reconstructing turns verbatim. */
export function getSessionMessageRows(db: Database.Database, sessionId: string): MessageRow[] {
return db
.prepare(
return prep(db,
"SELECT id, session_id, parent_id, timestamp, role, content_text, content_thinking, tool_calls, tool_results " +
"FROM messages WHERE session_id = ? ORDER BY rowid ASC",
)
.all(sessionId) as MessageRow[];
}

export function getNodesByAnalyzer(db: Database.Database, analyzerId: string, sessionId: string): AnalysisNodeRow[] {
return db
.prepare("SELECT * FROM analysis_nodes WHERE analyzer_id = ? AND session_id = ? ORDER BY created_at ASC, rowid ASC")
return prep(db, "SELECT * FROM analysis_nodes WHERE analyzer_id = ? AND session_id = ? ORDER BY created_at ASC, rowid ASC")
.all(analyzerId, sessionId) as AnalysisNodeRow[];
}

Expand All @@ -295,8 +292,7 @@ export function getNodesByAnalyzer(db: Database.Database, analyzerId: string, se
* `source_set_hash`, newest first, errors excluded.
*/
export function getLatestNodesByAnalyzerAcrossSessions(db: Database.Database, analyzerId: string): AnalysisNodeRow[] {
return db
.prepare(
return prep(db,
`SELECT * FROM analysis_nodes n
WHERE n.analyzer_id = ?
AND n.node_kind != 'error'
Expand All @@ -319,36 +315,33 @@ export function insertEdge(
db: Database.Database,
edge: { fromNodeId: string; toRefKind: string; toRefId: string; edgeKind: string; ordinal: number },
): void {
db.prepare(`
prep(db, `
INSERT INTO analysis_edges (id, from_node_id, to_ref_kind, to_ref_id, edge_kind, ordinal)
VALUES (?, ?, ?, ?, ?, ?)
`).run(uuidv7(), edge.fromNodeId, edge.toRefKind, edge.toRefId, edge.edgeKind, edge.ordinal);
}

export function getEdgesFrom(db: Database.Database, nodeId: string): AnalysisEdgeRow[] {
return db.prepare("SELECT * FROM analysis_edges WHERE from_node_id = ? ORDER BY ordinal ASC").all(nodeId) as AnalysisEdgeRow[];
return prep(db, "SELECT * FROM analysis_edges WHERE from_node_id = ? ORDER BY ordinal ASC").all(nodeId) as AnalysisEdgeRow[];
}

export function getEdgesTo(db: Database.Database, toRefId: string, edgeKind?: string): AnalysisEdgeRow[] {
if (edgeKind) {
return db
.prepare("SELECT * FROM analysis_edges WHERE to_ref_id = ? AND edge_kind = ?")
return prep(db, "SELECT * FROM analysis_edges WHERE to_ref_id = ? AND edge_kind = ?")
.all(toRefId, edgeKind) as AnalysisEdgeRow[];
}
return db.prepare("SELECT * FROM analysis_edges WHERE to_ref_id = ?").all(toRefId) as AnalysisEdgeRow[];
return prep(db, "SELECT * FROM analysis_edges WHERE to_ref_id = ?").all(toRefId) as AnalysisEdgeRow[];
}

/** Message ids that a node anchors to (via `anchors` edges with message targets). */
export function getAnchoredMessageIds(db: Database.Database, nodeId: string): string[] {
const rows = db
.prepare("SELECT to_ref_id FROM analysis_edges WHERE from_node_id = ? AND edge_kind = ? AND to_ref_kind = ?")
const rows = prep(db, "SELECT to_ref_id FROM analysis_edges WHERE from_node_id = ? AND edge_kind = ? AND to_ref_kind = ?")
.all(nodeId, EDGE_KINDS.ANCHORS, REF_KINDS.MESSAGE) as Array<{ to_ref_id: string }>;
return rows.map((r) => r.to_ref_id);
}

export function getMessage(db: Database.Database, id: string): MessageRow | undefined {
return db
.prepare(
return prep(db,
"SELECT id, session_id, parent_id, timestamp, role, content_text, content_thinking, tool_calls, tool_results FROM messages WHERE id = ?",
)
.get(id) as MessageRow | undefined;
Expand All @@ -366,17 +359,15 @@ export function getNodeVersions(
analyzerId: string,
sourceSetHash: string,
): AnalysisNodeRow[] {
return db
.prepare(
return prep(db,
"SELECT * FROM analysis_nodes WHERE analyzer_id = ? AND source_set_hash = ? ORDER BY created_at ASC, rowid ASC",
)
.all(analyzerId, sourceSetHash) as AnalysisNodeRow[];
}

/** The node that `nodeId` revises (its immediate older-version predecessor), if any. */
export function getRevisedNode(db: Database.Database, nodeId: string): AnalysisNodeRow | undefined {
const edge = db
.prepare("SELECT to_ref_id FROM analysis_edges WHERE from_node_id = ? AND edge_kind = ? LIMIT 1")
const edge = prep(db, "SELECT to_ref_id FROM analysis_edges WHERE from_node_id = ? AND edge_kind = ? LIMIT 1")
.get(nodeId, EDGE_KINDS.REVISES) as { to_ref_id: string } | undefined;
if (!edge) return undefined;
// `revises` edges reference the predecessor's content-addressed output_key.
Expand All @@ -388,8 +379,7 @@ export function getRevisions(db: Database.Database, nodeId: string): AnalysisNod
// `revises` edges point at the predecessor's output_key, so match on that.
const node = getNode(db, nodeId);
if (!node || !node.output_key) return [];
const edges = db
.prepare("SELECT from_node_id FROM analysis_edges WHERE to_ref_id = ? AND edge_kind = ?")
const edges = prep(db, "SELECT from_node_id FROM analysis_edges WHERE to_ref_id = ? AND edge_kind = ?")
.all(node.output_key, EDGE_KINDS.REVISES) as Array<{ from_node_id: string }>;
const out: AnalysisNodeRow[] = [];
for (const e of edges) {
Expand All @@ -409,10 +399,10 @@ export interface AnalysisStats {
}

export function getAnalysisStats(db: Database.Database): AnalysisStats {
const nodes = (db.prepare("SELECT COUNT(*) AS c FROM analysis_nodes").get() as { c: number }).c;
const edges = (db.prepare("SELECT COUNT(*) AS c FROM analysis_edges").get() as { c: number }).c;
const runs = (db.prepare("SELECT COUNT(*) AS c FROM analysis_runs").get() as { c: number }).c;
const kindRows = db.prepare("SELECT node_kind, COUNT(*) AS c FROM analysis_nodes GROUP BY node_kind").all() as Array<{
const nodes = (prep(db, "SELECT COUNT(*) AS c FROM analysis_nodes").get() as { c: number }).c;
const edges = (prep(db, "SELECT COUNT(*) AS c FROM analysis_edges").get() as { c: number }).c;
const runs = (prep(db, "SELECT COUNT(*) AS c FROM analysis_runs").get() as { c: number }).c;
const kindRows = prep(db, "SELECT node_kind, COUNT(*) AS c FROM analysis_nodes GROUP BY node_kind").all() as Array<{
node_kind: string;
c: number;
}>;
Expand Down
36 changes: 36 additions & 0 deletions src/db/prepared.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
import type Database from "better-sqlite3";
import type { Statement } from "better-sqlite3";

/**
* Prepared-statement cache, keyed per `Database` connection.
*
* better-sqlite3 `Statement`s are bound to the connection they were prepared
* on, so the cache must be per-instance rather than module-global: a statement
* prepared on a closed database must never be handed to a later connection
* (tests create a fresh temp DB per case, so a module-level `Map<sql, Statement>`
* would leak across cases). We key a WeakMap by the `Database` object, and each
* connection lazily grows its own `Map<sql, Statement>` on first use.
*
* Population is lazy and every caller runs after `migrate()`, so the cache is
* only ever filled once the schema is final — a cached statement can never
* capture a pre-migration query plan. There is deliberately no
* `initializeStatementCache(db)` hook to call after migration: because the map
* fills on first use, there is simply nothing for a new connection path to
* forget, and no init call that could be omitted.
*/
const statementCache = new WeakMap<Database.Database, Map<string, Statement>>();

/** Return a cached, connection-bound prepared statement, preparing on miss. */
export function prep(db: Database.Database, sql: string): Statement {
let cache = statementCache.get(db);
if (!cache) {
cache = new Map();
statementCache.set(db, cache);
}
let stmt = cache.get(sql);
if (!stmt) {
stmt = db.prepare(sql);
cache.set(sql, stmt);
}
return stmt;
}
Loading