Skip to main content

agent_graph_mcp/
store.rs

1//! SQLite-backed persistent storage for graphs, executions, checkpoints, and events.
2//!
3//! Enabled when `--data-dir` is passed to the server binary. Without it,
4//! everything is in-memory and lost on restart.
5
6use crate::evidence::{
7    digest, hmac_sha256, validate_witness_capture_with_key, verify_witness_record_with_key,
8    WitnessCapture, WitnessError, WitnessRecord,
9};
10use chrono::{DateTime, SecondsFormat, Utc};
11use rusqlite::{params, Connection, OptionalExtension};
12use serde_json::Value;
13use std::path::{Path, PathBuf};
14use std::sync::Mutex;
15
16/// Persistent store wrapping a SQLite connection.
17pub struct PersistentStore {
18    conn: std::sync::Arc<Mutex<Connection>>,
19    integrity_key: Option<std::sync::Arc<[u8]>>,
20    #[cfg(test)]
21    terminal_projection_fault: std::sync::Arc<std::sync::atomic::AtomicBool>,
22    #[cfg(test)]
23    checkpoint_persistence_fault: std::sync::Arc<std::sync::atomic::AtomicBool>,
24    #[cfg(test)]
25    graph_delete_fault: std::sync::Arc<std::sync::atomic::AtomicBool>,
26}
27
28#[derive(Debug, Clone, PartialEq)]
29pub struct CheckpointRecord {
30    pub checkpoint_id: String,
31    pub run_id: String,
32    pub graph_id: String,
33    pub graph_version: String,
34    pub next_node_cursor: String,
35    pub state: Value,
36    pub state_digest: String,
37    pub budgets: Value,
38    pub budget_counters: Value,
39    pub dependency_summary: Value,
40    pub dependency_digest: String,
41    pub terminal_cursor: u64,
42    pub event_cursor: u64,
43    pub checkpoint_digest: String,
44    pub created_at: String,
45    pub consumed_at: Option<String>,
46}
47
48#[derive(Debug, Clone, PartialEq, Eq)]
49pub enum CheckpointError {
50    NotFound,
51    Consumed,
52    Integrity,
53    Persistence,
54    IntegrityKeyRequired,
55}
56
57#[derive(Debug, Clone, PartialEq, Eq)]
58pub enum ApprovalError {
59    NotFound,
60    Conflict,
61    AlreadyDecided,
62    Expired,
63    DecisionNotAllowed,
64    Integrity,
65    Checkpoint(CheckpointError),
66    Persistence,
67    IntegrityKeyRequired,
68}
69
70impl ApprovalError {
71    pub fn code(&self) -> &'static str {
72        match self {
73            Self::NotFound => "APPROVAL_NOT_FOUND",
74            Self::Conflict => "APPROVAL_REQUEST_CONFLICT",
75            Self::AlreadyDecided => "APPROVAL_ALREADY_DECIDED",
76            Self::Expired => "APPROVAL_EXPIRED",
77            Self::DecisionNotAllowed => "APPROVAL_DECISION_NOT_ALLOWED",
78            Self::Integrity => "APPROVAL_INTEGRITY_FAILURE",
79            Self::Checkpoint(error) => error.code(),
80            Self::Persistence => "APPROVAL_PERSISTENCE_FAILURE",
81            Self::IntegrityKeyRequired => "INTEGRITY_KEY_REQUIRED",
82        }
83    }
84
85    pub fn message(&self) -> String {
86        match self {
87            Self::NotFound => "approval was not found".into(),
88            Self::Conflict => {
89                "a conflicting pending approval already exists for this checkpoint and audience"
90                    .into()
91            }
92            Self::AlreadyDecided => "approval has already been decided".into(),
93            Self::Expired => "approval has expired".into(),
94            Self::DecisionNotAllowed => "decision is not allowed by this approval".into(),
95            Self::Integrity => "approval integrity validation failed".into(),
96            Self::Checkpoint(error) => error.message().into(),
97            Self::Persistence => "approval persistence failed".into(),
98            Self::IntegrityKeyRequired => {
99                "integrity key is required for durable approval operations".into()
100            }
101        }
102    }
103}
104
105#[derive(Debug, Clone, PartialEq, Eq)]
106pub struct ApprovalRecord {
107    pub approval_id: String,
108    pub checkpoint_id: String,
109    pub run_id: String,
110    pub graph_id: String,
111    pub graph_version: String,
112    pub checkpoint_digest: String,
113    pub audience: String,
114    pub prompt_digest: String,
115    pub allowed_decisions: Vec<String>,
116    pub approval_digest: String,
117    pub status: String,
118    pub decision: Option<String>,
119    pub decided_by: Option<String>,
120    pub decided_at: Option<String>,
121    pub expires_at: String,
122    pub created_at: String,
123}
124
125#[derive(Debug, Clone)]
126pub struct ApprovedCheckpoint {
127    pub approval: ApprovalRecord,
128    pub checkpoint: CheckpointRecord,
129}
130
131impl CheckpointError {
132    pub fn code(&self) -> &'static str {
133        match self {
134            Self::NotFound => "CHECKPOINT_NOT_FOUND",
135            Self::Consumed => "CHECKPOINT_CONSUMED",
136            Self::Integrity => "CHECKPOINT_INTEGRITY_FAILURE",
137            Self::Persistence => "CHECKPOINT_PERSISTENCE_FAILURE",
138            Self::IntegrityKeyRequired => "INTEGRITY_KEY_REQUIRED",
139        }
140    }
141
142    pub fn message(&self) -> &'static str {
143        match self {
144            Self::NotFound => "checkpoint was not found",
145            Self::Consumed => "checkpoint has already been consumed",
146            Self::Integrity => "checkpoint integrity validation failed",
147            Self::Persistence => "checkpoint persistence failed; resumability was not advertised",
148            Self::IntegrityKeyRequired => {
149                "integrity key is required for durable checkpoint operations"
150            }
151        }
152    }
153}
154
155#[derive(Debug, Clone)]
156pub struct ExecutionContract {
157    pub graph_id: String,
158    pub graph_version: String,
159    pub input: Value,
160    pub budgets: Value,
161}
162
163type CheckpointParts = (
164    String,
165    String,
166    String,
167    Option<String>,
168    Option<String>,
169    Option<String>,
170    Option<String>,
171    Option<String>,
172    Option<String>,
173    Option<String>,
174    Option<i64>,
175    Option<i64>,
176    Option<String>,
177    Option<String>,
178    Option<String>,
179    Option<String>,
180);
181
182fn checkpoint_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<CheckpointParts> {
183    Ok((
184        row.get(0)?,
185        row.get(1)?,
186        row.get(2)?,
187        row.get(3)?,
188        row.get(4)?,
189        row.get(5)?,
190        row.get(6)?,
191        row.get(7)?,
192        row.get(8)?,
193        row.get(9)?,
194        row.get(10)?,
195        row.get(11)?,
196        row.get(12)?,
197        row.get(13)?,
198        row.get(14)?,
199        row.get(15)?,
200    ))
201}
202
203fn checkpoint_from_parts(parts: CheckpointParts) -> Option<CheckpointRecord> {
204    let (
205        run_id,
206        graph_id,
207        graph_version,
208        next_node_cursor,
209        state_json,
210        state_digest,
211        budgets_json,
212        counters_json,
213        dependency_json,
214        dependency_digest,
215        terminal_cursor,
216        event_cursor,
217        checkpoint_id,
218        checkpoint_digest,
219        created_at,
220        consumed_at,
221    ) = parts;
222    Some(CheckpointRecord {
223        checkpoint_id: checkpoint_id?,
224        run_id,
225        graph_id,
226        graph_version,
227        next_node_cursor: next_node_cursor?,
228        state: serde_json::from_str(&state_json?).ok()?,
229        state_digest: state_digest?,
230        budgets: serde_json::from_str(&budgets_json?).ok()?,
231        budget_counters: serde_json::from_str(&counters_json?).ok()?,
232        dependency_summary: serde_json::from_str(&dependency_json?).ok()?,
233        dependency_digest: dependency_digest?,
234        terminal_cursor: u64::try_from(terminal_cursor?).ok()?,
235        event_cursor: u64::try_from(event_cursor?).ok()?,
236        checkpoint_digest: checkpoint_digest?,
237        created_at: created_at?,
238        consumed_at,
239    })
240}
241
242fn checkpoint_digest(record: &CheckpointRecord, key: &[u8]) -> String {
243    hmac_sha256(
244        &serde_json::json!({
245            "checkpoint_id": record.checkpoint_id,
246            "run_id": record.run_id,
247            "graph_id": record.graph_id,
248            "graph_version": record.graph_version,
249            "next_node_cursor": record.next_node_cursor,
250            "state": record.state,
251            "state_digest": record.state_digest,
252            "budgets": record.budgets,
253            "budget_counters": record.budget_counters,
254            "dependency_summary": record.dependency_summary,
255            "dependency_digest": record.dependency_digest,
256            "terminal_cursor": record.terminal_cursor,
257            "event_cursor": record.event_cursor,
258            "created_at": record.created_at,
259        }),
260        key,
261    )
262}
263
264fn validate_checkpoint_record(
265    record: &CheckpointRecord,
266    key: &[u8],
267) -> Result<(), CheckpointError> {
268    if record.checkpoint_id != format!("checkpoint-{}-{}", record.run_id, record.next_node_cursor)
269        || record.state_digest != digest(&record.state)
270        || record.dependency_digest != digest(&record.dependency_summary)
271        || record.checkpoint_digest != checkpoint_digest(record, key)
272    {
273        return Err(CheckpointError::Integrity);
274    }
275    Ok(())
276}
277
278type ApprovalParts = (
279    String,
280    Option<String>,
281    String,
282    Option<String>,
283    Option<String>,
284    Option<String>,
285    String,
286    String,
287    Option<String>,
288    Option<String>,
289    String,
290    Option<String>,
291    Option<String>,
292    Option<String>,
293    String,
294    String,
295    String,
296);
297
298fn approval_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<ApprovalParts> {
299    Ok((
300        row.get(0)?,
301        row.get(1)?,
302        row.get(2)?,
303        row.get(3)?,
304        row.get(4)?,
305        row.get(5)?,
306        row.get(6)?,
307        row.get(7)?,
308        row.get(8)?,
309        row.get(9)?,
310        row.get(10)?,
311        row.get(11)?,
312        row.get(12)?,
313        row.get(13)?,
314        row.get(14)?,
315        row.get(15)?,
316        row.get(16)?,
317    ))
318}
319
320fn approval_from_parts(parts: ApprovalParts) -> Option<ApprovalRecord> {
321    let (
322        approval_id,
323        checkpoint_id,
324        run_id,
325        graph_id,
326        graph_version,
327        checkpoint_digest,
328        audience,
329        prompt_digest,
330        allowed_decisions,
331        approval_digest,
332        status,
333        decision,
334        decided_by,
335        decided_at,
336        expires_at,
337        created_at,
338        _legacy,
339    ) = parts;
340    let allowed_decisions = serde_json::from_str(&allowed_decisions?).ok()?;
341    Some(ApprovalRecord {
342        approval_id,
343        checkpoint_id: checkpoint_id?,
344        run_id,
345        graph_id: graph_id?,
346        graph_version: graph_version?,
347        checkpoint_digest: checkpoint_digest?,
348        audience,
349        prompt_digest,
350        allowed_decisions,
351        approval_digest: approval_digest?,
352        status,
353        decision,
354        decided_by,
355        decided_at,
356        expires_at,
357        created_at,
358    })
359}
360
361const APPROVAL_COLUMNS: &str = "approval_id, checkpoint_id, run_id, graph_id,
362    graph_version, checkpoint_digest, audience, prompt_digest,
363    allowed_decisions, approval_digest, status, decision, decided_by,
364    decided_at, expires_at, created_at, prompt";
365
366fn approval_digest(record: &ApprovalRecord, key: &[u8]) -> String {
367    hmac_sha256(
368        &serde_json::json!({
369            "approval_id": record.approval_id,
370            "checkpoint_id": record.checkpoint_id,
371            "run_id": record.run_id,
372            "graph_id": record.graph_id,
373            "graph_version": record.graph_version,
374            "checkpoint_digest": record.checkpoint_digest,
375            "audience": record.audience,
376            "prompt_digest": record.prompt_digest,
377            "allowed_decisions": record.allowed_decisions,
378            "status": record.status,
379            "decision": record.decision,
380            "decided_by": record.decided_by,
381            "decided_at": record.decided_at,
382            "expires_at": record.expires_at,
383            "created_at": record.created_at,
384        }),
385        key,
386    )
387}
388
389fn parse_approval_parts(parts: ApprovalParts, key: &[u8]) -> Result<ApprovalRecord, ApprovalError> {
390    let record = approval_from_parts(parts).ok_or(ApprovalError::Integrity)?;
391    if record.approval_digest != approval_digest(&record, key) {
392        return Err(ApprovalError::Integrity);
393    }
394    Ok(record)
395}
396
397fn uuid_like() -> String {
398    digest(&Value::String(
399        Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true),
400    ))
401    .trim_start_matches("sha256:")
402    .to_owned()
403}
404
405fn load_checkpoint_from_tx(
406    tx: &rusqlite::Transaction<'_>,
407    checkpoint_id: &str,
408    key: &[u8],
409) -> Result<CheckpointRecord, CheckpointError> {
410    let row = tx.query_row(
411        "SELECT run_id, graph_id, graph_version, next_cursor, state_json,
412                state_digest, budgets_json, budget_counters_json,
413                dependency_json, dependency_digest, terminal_cursor,
414                event_cursor, checkpoint_id, checkpoint_digest, created_at,
415                consumed_at
416         FROM checkpoints WHERE checkpoint_id = ?1",
417        params![checkpoint_id],
418        checkpoint_row,
419    );
420    let parts = match row {
421        Ok(parts) => parts,
422        Err(rusqlite::Error::QueryReturnedNoRows) => return Err(CheckpointError::NotFound),
423        Err(_) => return Err(CheckpointError::Persistence),
424    };
425    let record = checkpoint_from_parts(parts).ok_or(CheckpointError::Integrity)?;
426    validate_checkpoint_record(&record, key)?;
427    Ok(record)
428}
429
430#[derive(Debug, Clone, Copy, PartialEq, Eq)]
431pub enum GraphDeleteResult {
432    Deleted,
433    NotFound,
434    Referenced,
435}
436
437impl Clone for PersistentStore {
438    fn clone(&self) -> Self {
439        Self {
440            conn: self.conn.clone(),
441            integrity_key: self.integrity_key.clone(),
442            #[cfg(test)]
443            terminal_projection_fault: self.terminal_projection_fault.clone(),
444            #[cfg(test)]
445            checkpoint_persistence_fault: self.checkpoint_persistence_fault.clone(),
446            #[cfg(test)]
447            graph_delete_fault: self.graph_delete_fault.clone(),
448        }
449    }
450}
451
452impl PersistentStore {
453    /// Open (or create) the SQLite database at `{data_dir}/agent-graph.db`.
454    pub fn open(data_dir: &Path) -> Result<Self, String> {
455        Self::open_with_integrity_key(data_dir, None)
456    }
457
458    /// Open (or create) the SQLite database with an explicit integrity key path.
459    /// If `integrity_key_path` is None, falls back to the
460    /// `AGENT_GRAPH_INTEGRITY_KEY_PATH` environment variable.
461    pub fn open_with_integrity_key(
462        data_dir: &Path,
463        integrity_key_path: Option<&Path>,
464    ) -> Result<Self, String> {
465        crate::fs_security::validate_data_store(data_dir, integrity_key_path)
466            .map_err(|e| format!("filesystem security check failed: {e}"))?;
467        std::fs::create_dir_all(data_dir).map_err(|e| format!("failed to create data dir: {e}"))?;
468        let db_path = data_dir.join("agent-graph.db");
469        let conn =
470            Connection::open(&db_path).map_err(|e| format!("failed to open database: {e}"))?;
471
472        conn.execute_batch(
473            "PRAGMA journal_mode = WAL;
474             PRAGMA foreign_keys = ON;
475             PRAGMA busy_timeout = 5000;",
476        )
477        .map_err(|e| format!("pragma error: {e}"))?;
478        crate::fs_security::validate_data_store(data_dir, integrity_key_path)
479            .map_err(|e| format!("filesystem security check failed after WAL init: {e}"))?;
480
481        let store = Self {
482            conn: std::sync::Arc::new(Mutex::new(conn)),
483            integrity_key: Self::load_integrity_key_from(integrity_key_path),
484            #[cfg(test)]
485            terminal_projection_fault: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(
486                false,
487            )),
488            #[cfg(test)]
489            checkpoint_persistence_fault: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(
490                false,
491            )),
492            #[cfg(test)]
493            graph_delete_fault: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
494        };
495        store.migrate()?;
496        Ok(store)
497    }
498
499    #[allow(dead_code)]
500    fn load_integrity_key() -> Option<std::sync::Arc<[u8]>> {
501        Self::load_integrity_key_from(None)
502    }
503
504    /// Load the integrity key from an explicit path, or fall back to the
505    /// `AGENT_GRAPH_INTEGRITY_KEY_PATH` environment variable.
506    fn load_integrity_key_from(explicit: Option<&Path>) -> Option<std::sync::Arc<[u8]>> {
507        let path = if let Some(p) = explicit {
508            std::ffi::OsStr::new(p).to_owned()
509        } else {
510            std::env::var_os("AGENT_GRAPH_INTEGRITY_KEY_PATH")?
511        };
512        let key = std::fs::read(path).ok()?;
513        (key.len() >= 32).then(|| std::sync::Arc::from(key))
514    }
515
516    fn require_integrity_key(&self) -> Result<&[u8], ()> {
517        self.integrity_key.as_deref().ok_or(())
518    }
519
520    pub fn has_integrity_key(&self) -> bool {
521        self.integrity_key.is_some()
522    }
523
524    fn migrate(&self) -> Result<(), String> {
525        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
526        conn.execute_batch(
527            "CREATE TABLE IF NOT EXISTS graphs (
528                name TEXT PRIMARY KEY,
529                spec_json TEXT NOT NULL,
530                spec_version TEXT NOT NULL DEFAULT '2',
531                topology_hash TEXT NOT NULL,
532                created_at TEXT NOT NULL DEFAULT (datetime('now')),
533                updated_at TEXT NOT NULL DEFAULT (datetime('now'))
534            );
535
536            CREATE TABLE IF NOT EXISTS graph_versions (
537                graph_name TEXT NOT NULL,
538                topology_hash TEXT NOT NULL,
539                spec_json TEXT NOT NULL,
540                created_at TEXT NOT NULL DEFAULT (datetime('now')),
541                PRIMARY KEY (graph_name, topology_hash)
542            );
543
544            CREATE TABLE IF NOT EXISTS executions (
545                run_id TEXT PRIMARY KEY,
546                graph_name TEXT NOT NULL,
547                graph_hash TEXT NOT NULL,
548                thread_id TEXT,
549                status TEXT NOT NULL,
550                input_json TEXT,
551                budgets_json TEXT,
552                final_state_json TEXT,
553                started_at TEXT NOT NULL,
554                finished_at TEXT,
555                total_nodes INTEGER DEFAULT 0,
556                failed_attempts INTEGER DEFAULT 0,
557                idempotency_key TEXT UNIQUE,
558                FOREIGN KEY (graph_name) REFERENCES graphs(name)
559            );
560
561            CREATE TABLE IF NOT EXISTS checkpoints (
562                run_id TEXT NOT NULL,
563                node_id TEXT NOT NULL,
564                attempt INTEGER NOT NULL,
565                input_json TEXT,
566                output_json TEXT,
567                status TEXT NOT NULL,
568                error TEXT,
569                recorded_at TEXT NOT NULL DEFAULT (datetime('now')),
570                checkpoint_id TEXT,
571                graph_id TEXT,
572                graph_version TEXT,
573                next_cursor TEXT,
574                state_json TEXT,
575                state_digest TEXT,
576                budgets_json TEXT,
577                budget_counters_json TEXT,
578                dependency_json TEXT,
579                dependency_digest TEXT,
580                terminal_cursor INTEGER,
581                event_cursor INTEGER,
582                checkpoint_digest TEXT,
583                created_at TEXT,
584                consumed_at TEXT,
585                PRIMARY KEY (run_id, node_id, attempt),
586                FOREIGN KEY (run_id) REFERENCES executions(run_id)
587            );
588
589            CREATE TABLE IF NOT EXISTS events (
590                run_id TEXT NOT NULL,
591                seq INTEGER NOT NULL,
592                event_type TEXT NOT NULL,
593                event_json TEXT NOT NULL,
594                emitted_at TEXT NOT NULL DEFAULT (datetime('now')),
595                PRIMARY KEY (run_id, seq),
596                FOREIGN KEY (run_id) REFERENCES executions(run_id)
597            );
598
599            CREATE TABLE IF NOT EXISTS terminal_receipts (
600                run_id TEXT PRIMARY KEY,
601                receipt_json TEXT NOT NULL,
602                bundle_json TEXT NOT NULL,
603                receipt_digest TEXT NOT NULL,
604                persisted_at TEXT NOT NULL DEFAULT (datetime('now')),
605                FOREIGN KEY (run_id) REFERENCES executions(run_id)
606            );
607
608            CREATE TABLE IF NOT EXISTS template_candidates (
609                template_id TEXT PRIMARY KEY, spec_digest TEXT NOT NULL,
610                graph_id TEXT NOT NULL, graph_version TEXT NOT NULL,
611                source_ref TEXT NOT NULL, state TEXT NOT NULL,
612                created_at TEXT NOT NULL DEFAULT (datetime('now')),
613                updated_at TEXT NOT NULL DEFAULT (datetime('now'))
614            );
615            CREATE TABLE IF NOT EXISTS template_outcome_links (
616                template_id TEXT NOT NULL, run_id TEXT NOT NULL,
617                terminal_receipt_id TEXT NOT NULL, receipt_digest TEXT NOT NULL,
618                disposition TEXT NOT NULL, evidence_digest TEXT NOT NULL,
619                recorded_at TEXT NOT NULL DEFAULT (datetime('now')),
620                PRIMARY KEY (template_id, run_id),
621                FOREIGN KEY (template_id) REFERENCES template_candidates(template_id),
622                FOREIGN KEY (run_id) REFERENCES executions(run_id),
623                FOREIGN KEY (run_id) REFERENCES terminal_receipts(run_id)
624            );
625            CREATE TABLE IF NOT EXISTS template_promotion_decisions (
626                template_id TEXT NOT NULL, from_state TEXT NOT NULL,
627                to_state TEXT NOT NULL, evidence_set_digest TEXT NOT NULL,
628                operator_receipt_id TEXT NOT NULL, decision_digest TEXT NOT NULL,
629                decided_at TEXT NOT NULL DEFAULT (datetime('now')),
630                FOREIGN KEY (template_id) REFERENCES template_candidates(template_id)
631            );
632
633            CREATE TABLE IF NOT EXISTS approval_requests (
634                approval_id TEXT PRIMARY KEY,
635                run_id TEXT NOT NULL,
636                node_id TEXT NOT NULL,
637                checkpoint_id TEXT,
638                graph_id TEXT,
639                graph_version TEXT,
640                checkpoint_digest TEXT,
641                audience TEXT NOT NULL,
642                prompt TEXT NOT NULL,
643                prompt_digest TEXT,
644                allowed_decisions TEXT NOT NULL,
645                approval_digest TEXT,
646                status TEXT NOT NULL DEFAULT 'pending',
647                decision TEXT,
648                decided_by TEXT,
649                decided_at TEXT,
650                expires_at TEXT NOT NULL,
651                created_at TEXT NOT NULL DEFAULT (datetime('now')),
652                FOREIGN KEY (run_id) REFERENCES executions(run_id)
653            );
654
655            CREATE TABLE IF NOT EXISTS idempotency_keys (
656                key TEXT PRIMARY KEY,
657                request_digest TEXT NOT NULL,
658                result_json TEXT NOT NULL,
659                valid INTEGER NOT NULL DEFAULT 1,
660                created_at TEXT NOT NULL DEFAULT (datetime('now'))
661            );
662
663            CREATE TABLE IF NOT EXISTS source_witnesses (
664                witness_id TEXT PRIMARY KEY,
665                locator TEXT NOT NULL,
666                content TEXT NOT NULL,
667                media_type TEXT NOT NULL,
668                authority_class TEXT NOT NULL,
669                retrieved_at TEXT NOT NULL,
670                digest TEXT NOT NULL UNIQUE,
671                created_at TEXT NOT NULL DEFAULT (datetime('now'))
672            );",
673        )
674        .map_err(|e| format!("migration error: {e}"))?;
675        let has_request_digest = conn
676            .prepare("PRAGMA table_info(idempotency_keys)")
677            .map_err(|e| format!("idempotency schema inspection error: {e}"))?
678            .query_map([], |row| row.get::<_, String>(1))
679            .map_err(|e| format!("idempotency schema inspection error: {e}"))?
680            .filter_map(Result::ok)
681            .any(|column| column == "request_digest");
682        if !has_request_digest {
683            conn.execute(
684                "ALTER TABLE idempotency_keys ADD COLUMN request_digest TEXT",
685                [],
686            )
687            .map_err(|e| format!("idempotency migration error: {e}"))?;
688        }
689        let has_idempotency_valid = conn
690            .prepare("PRAGMA table_info(idempotency_keys)")
691            .map_err(|e| format!("idempotency schema inspection error: {e}"))?
692            .query_map([], |row| row.get::<_, String>(1))
693            .map_err(|e| format!("idempotency schema inspection error: {e}"))?
694            .filter_map(Result::ok)
695            .any(|column| column == "valid");
696        if !has_idempotency_valid {
697            conn.execute(
698                "ALTER TABLE idempotency_keys ADD COLUMN valid INTEGER NOT NULL DEFAULT 1",
699                [],
700            )
701            .map_err(|e| format!("idempotency migration error: {e}"))?;
702        }
703        conn.execute(
704            "UPDATE idempotency_keys SET valid = 0 WHERE request_digest IS NULL",
705            [],
706        )
707        .map_err(|e| format!("idempotency quarantine error: {e}"))?;
708        for (table, column, definition) in [
709            ("executions", "budgets_json", "TEXT"),
710            ("checkpoints", "checkpoint_id", "TEXT"),
711            ("checkpoints", "graph_id", "TEXT"),
712            ("checkpoints", "graph_version", "TEXT"),
713            ("checkpoints", "next_cursor", "TEXT"),
714            ("checkpoints", "state_json", "TEXT"),
715            ("checkpoints", "state_digest", "TEXT"),
716            ("checkpoints", "budgets_json", "TEXT"),
717            ("checkpoints", "budget_counters_json", "TEXT"),
718            ("checkpoints", "dependency_json", "TEXT"),
719            ("checkpoints", "dependency_digest", "TEXT"),
720            ("checkpoints", "terminal_cursor", "INTEGER"),
721            ("checkpoints", "event_cursor", "INTEGER"),
722            ("checkpoints", "checkpoint_digest", "TEXT"),
723            ("checkpoints", "created_at", "TEXT"),
724            ("checkpoints", "consumed_at", "TEXT"),
725            ("approval_requests", "checkpoint_id", "TEXT"),
726            ("approval_requests", "graph_id", "TEXT"),
727            ("approval_requests", "graph_version", "TEXT"),
728            ("approval_requests", "checkpoint_digest", "TEXT"),
729            ("approval_requests", "prompt_digest", "TEXT"),
730            ("approval_requests", "approval_digest", "TEXT"),
731        ] {
732            let has_column = conn
733                .prepare(&format!("PRAGMA table_info({table})"))
734                .map_err(|e| format!("{table} schema inspection error: {e}"))?
735                .query_map([], |row| row.get::<_, String>(1))
736                .map_err(|e| format!("{table} schema inspection error: {e}"))?
737                .filter_map(Result::ok)
738                .any(|name| name == column);
739            if !has_column {
740                conn.execute(
741                    &format!("ALTER TABLE {table} ADD COLUMN {column} {definition}"),
742                    [],
743                )
744                .map_err(|e| format!("{table} migration error: {e}"))?;
745            }
746        }
747        conn.execute(
748            "CREATE UNIQUE INDEX IF NOT EXISTS checkpoints_checkpoint_id_idx
749             ON checkpoints(checkpoint_id) WHERE checkpoint_id IS NOT NULL",
750            [],
751        )
752        .map_err(|e| format!("checkpoint index migration error: {e}"))?;
753        Ok(())
754    }
755
756    // ── Local source witnesses ───────────────────────────────────────
757
758    pub fn capture_witness(&self, capture: WitnessCapture) -> Result<WitnessRecord, WitnessError> {
759        let key = self.require_integrity_key().map_err(|_| {
760            WitnessError::new(
761                "INTEGRITY_KEY_REQUIRED",
762                "an external integrity key is required for durable witness capture",
763            )
764        })?;
765        let expected = validate_witness_capture_with_key(capture, Some(key))
766            .map_err(|error| WitnessError::new(error.code, error.message))?;
767        let conn = self
768            .conn
769            .lock()
770            .map_err(|_| WitnessError::new("WITNESS_STORE_ERROR", "witness SQLite lock failed"))?;
771        conn.execute(
772            "INSERT OR IGNORE INTO source_witnesses
773             (witness_id, locator, content, media_type, authority_class, retrieved_at, digest)
774             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
775            params![
776                expected.witness_id,
777                expected.locator,
778                expected.content,
779                expected.media_type,
780                expected.authority_class,
781                expected.retrieved_at,
782                expected.digest
783            ],
784        )
785        .map_err(|_| WitnessError::new("WITNESS_STORE_ERROR", "witness SQLite write failed"))?;
786        drop(conn);
787        self.get_witness(&expected.witness_id)?.ok_or_else(|| {
788            WitnessError::new(
789                "WITNESS_STORE_ERROR",
790                "witness SQLite write did not produce a row",
791            )
792        })
793    }
794
795    pub fn get_witness(&self, witness_id: &str) -> Result<Option<WitnessRecord>, WitnessError> {
796        let key = self.require_integrity_key().map_err(|_| {
797            WitnessError::new(
798                "INTEGRITY_KEY_REQUIRED",
799                "an external integrity key is required for durable witness reads",
800            )
801        })?;
802        let conn = self
803            .conn
804            .lock()
805            .map_err(|_| WitnessError::new("WITNESS_STORE_ERROR", "witness SQLite lock failed"))?;
806        let row = conn.query_row(
807            "SELECT witness_id, locator, content, media_type, authority_class, retrieved_at, digest
808             FROM source_witnesses WHERE witness_id = ?1",
809            params![witness_id],
810            |row| {
811                Ok(WitnessRecord {
812                    witness_id: row.get(0)?,
813                    locator: row.get(1)?,
814                    content: row.get(2)?,
815                    media_type: row.get(3)?,
816                    authority_class: row.get(4)?,
817                    retrieved_at: row.get(5)?,
818                    digest: row.get(6)?,
819                })
820            },
821        );
822        match row {
823            Ok(record) => {
824                verify_witness_record_with_key(&record, Some(key))?;
825                Ok(Some(record))
826            }
827            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
828            Err(rusqlite::Error::FromSqlConversionFailure(..))
829            | Err(rusqlite::Error::InvalidColumnType(..))
830            | Err(rusqlite::Error::InvalidColumnName(_)) => Err(WitnessError::new(
831                "WITNESS_INTEGRITY_FAILURE",
832                "stored witness integrity validation failed",
833            )),
834            Err(_) => Err(WitnessError::new(
835                "WITNESS_STORE_ERROR",
836                "witness SQLite read failed",
837            )),
838        }
839    }
840
841    // ── Graphs ──────────────────────────────────────────────────────────
842
843    pub fn save_graph(
844        &self,
845        name: &str,
846        spec_json: &str,
847        topology_hash: &str,
848        overwrite: bool,
849    ) -> Result<(), String> {
850        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
851        let current: Option<String> = conn
852            .query_row(
853                "SELECT topology_hash FROM graphs WHERE name = ?1",
854                params![name],
855                |row| row.get(0),
856            )
857            .ok();
858        if let Some(current) = current
859            .as_deref()
860            .filter(|current| *current != topology_hash)
861        {
862            let referenced: bool = conn
863                .query_row(
864                    "SELECT EXISTS(SELECT 1 FROM executions WHERE graph_name = ?1 AND graph_hash = ?2)",
865                    params![name, current],
866                    |row| row.get(0),
867                )
868                .map_err(|e| format!("graph version reference check error: {e}"))?;
869            if referenced && !overwrite {
870                return Err(format!("graph '{name}' current version {current} is referenced by a durable execution; explicit overwrite is required"));
871            }
872        }
873        conn.execute(
874            "INSERT OR IGNORE INTO graph_versions (graph_name, topology_hash, spec_json) VALUES (?1, ?2, ?3)",
875            params![name, topology_hash, spec_json],
876        )
877        .map_err(|e| format!("save graph version error: {e}"))?;
878        conn.execute(
879            "INSERT INTO graphs (name, spec_json, topology_hash)
880             VALUES (?1, ?2, ?3)
881             ON CONFLICT(name) DO UPDATE SET
882                spec_json = excluded.spec_json,
883                topology_hash = excluded.topology_hash,
884                updated_at = datetime('now')",
885            params![name, spec_json, topology_hash],
886        )
887        .map_err(|e| format!("save_graph error: {e}"))?;
888        Ok(())
889    }
890
891    pub fn load_graph(&self, name: &str) -> Result<Option<(String, String)>, String> {
892        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
893        let mut stmt = conn
894            .prepare("SELECT spec_json, topology_hash FROM graphs WHERE name = ?1")
895            .map_err(|e| format!("load_graph error: {e}"))?;
896        let result = stmt
897            .query_row(params![name], |row| {
898                Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
899            })
900            .ok();
901        Ok(result)
902    }
903
904    pub fn list_graphs(&self) -> Result<Vec<(String, String, String)>, String> {
905        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
906        let mut stmt = conn
907            .prepare("SELECT name, topology_hash, created_at FROM graphs ORDER BY name")
908            .map_err(|e| format!("list_graphs error: {e}"))?;
909        let rows = stmt
910            .query_map([], |row| {
911                Ok((
912                    row.get::<_, String>(0)?,
913                    row.get::<_, String>(1)?,
914                    row.get::<_, String>(2)?,
915                ))
916            })
917            .map_err(|e| format!("list_graphs error: {e}"))?
918            .filter_map(|r| r.ok())
919            .collect();
920        Ok(rows)
921    }
922
923    pub fn delete_graph(&self, name: &str) -> Result<GraphDeleteResult, String> {
924        let mut conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
925        let tx = conn
926            .transaction()
927            .map_err(|e| format!("delete_graph transaction begin error: {e}"))?;
928        let referenced: bool = tx
929            .query_row(
930                "SELECT EXISTS(SELECT 1 FROM executions WHERE graph_name = ?1)",
931                params![name],
932                |row| row.get(0),
933            )
934            .map_err(|e| format!("delete_graph reference check error: {e}"))?;
935        if referenced {
936            return Ok(GraphDeleteResult::Referenced);
937        }
938        let affected = tx
939            .execute("DELETE FROM graphs WHERE name = ?1", params![name])
940            .map_err(|e| format!("delete_graph error: {e}"))?;
941        if affected == 0 {
942            return Ok(GraphDeleteResult::NotFound);
943        }
944        #[cfg(test)]
945        if self
946            .graph_delete_fault
947            .swap(false, std::sync::atomic::Ordering::SeqCst)
948        {
949            return Err("injected graph deletion failure after graph row removal".into());
950        }
951        tx.execute(
952            "DELETE FROM graph_versions WHERE graph_name = ?1",
953            params![name],
954        )
955        .map_err(|e| format!("delete_graph versions error: {e}"))?;
956        tx.commit()
957            .map_err(|e| format!("delete_graph transaction commit error: {e}"))?;
958        Ok(GraphDeleteResult::Deleted)
959    }
960
961    #[cfg(test)]
962    pub(crate) fn fail_graph_delete_after_graph_row(&self) {
963        self.graph_delete_fault
964            .store(true, std::sync::atomic::Ordering::SeqCst);
965    }
966
967    #[cfg(test)]
968    pub(crate) fn fail_terminal_projection_after_events(&self) {
969        self.terminal_projection_fault
970            .store(true, std::sync::atomic::Ordering::SeqCst);
971    }
972
973    pub fn load_graph_version(
974        &self,
975        name: &str,
976        topology_hash: &str,
977    ) -> Result<Option<String>, String> {
978        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
979        let result = conn
980            .query_row(
981                "SELECT spec_json FROM graph_versions WHERE graph_name = ?1 AND topology_hash = ?2",
982                params![name, topology_hash],
983                |row| row.get::<_, String>(0),
984            )
985            .ok();
986        Ok(result)
987    }
988
989    pub fn list_graph_versions(&self, name: &str) -> Result<Vec<String>, String> {
990        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
991        let mut stmt = conn
992            .prepare("SELECT topology_hash FROM graph_versions WHERE graph_name = ?1 ORDER BY created_at, topology_hash")
993            .map_err(|e| format!("list graph versions error: {e}"))?;
994        let versions = stmt
995            .query_map(params![name], |row| row.get::<_, String>(0))
996            .map_err(|e| format!("list graph versions error: {e}"))?
997            .filter_map(Result::ok)
998            .collect();
999        Ok(versions)
1000    }
1001
1002    // ── Executions ──────────────────────────────────────────────────────
1003
1004    pub fn save_execution(
1005        &self,
1006        run_id: &str,
1007        graph_name: &str,
1008        graph_hash: &str,
1009        status: &str,
1010        input_json: &str,
1011    ) -> Result<(), String> {
1012        self.save_execution_with_budgets(run_id, graph_name, graph_hash, status, input_json, None)
1013    }
1014
1015    pub fn save_execution_with_budgets(
1016        &self,
1017        run_id: &str,
1018        graph_name: &str,
1019        graph_hash: &str,
1020        status: &str,
1021        input_json: &str,
1022        budgets_json: Option<&str>,
1023    ) -> Result<(), String> {
1024        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1025        conn.execute(
1026            "INSERT INTO executions (run_id, graph_name, graph_hash, status, input_json, budgets_json, started_at)
1027             VALUES (?1, ?2, ?3, ?4, ?5, ?6, datetime('now'))
1028             ON CONFLICT(run_id) DO UPDATE SET
1029                status = excluded.status,
1030                budgets_json = COALESCE(excluded.budgets_json, executions.budgets_json),
1031                finished_at = CASE WHEN excluded.status IN ('completed','failed','cancelled') THEN datetime('now') ELSE finished_at END",
1032            params![run_id, graph_name, graph_hash, status, input_json, budgets_json],
1033        )
1034        .map_err(|e| format!("save_execution error: {e}"))?;
1035        Ok(())
1036    }
1037
1038    pub fn update_execution_status(
1039        &self,
1040        run_id: &str,
1041        status: &str,
1042        final_state_json: Option<&str>,
1043        total_nodes: Option<usize>,
1044        failed_attempts: Option<usize>,
1045    ) -> Result<(), String> {
1046        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1047        let changed = conn.execute(
1048            "UPDATE executions SET
1049                status = ?2,
1050                final_state_json = COALESCE(?3, final_state_json),
1051                total_nodes = COALESCE(?4, total_nodes),
1052                failed_attempts = COALESCE(?5, failed_attempts),
1053                finished_at = CASE WHEN ?2 IN ('completed','failed','cancelled') THEN datetime('now') ELSE finished_at END
1054             WHERE run_id = ?1",
1055            params![
1056                run_id,
1057                status,
1058                final_state_json,
1059                total_nodes.map(|v| v as i64),
1060                failed_attempts.map(|v| v as i64)
1061            ],
1062        ).map_err(|e| format!("update_execution error: {e}"))?;
1063        if changed != 1 {
1064            return Err(format!(
1065                "update_execution error: run '{run_id}' was not found"
1066            ));
1067        }
1068        Ok(())
1069    }
1070
1071    pub fn persist_terminal_projection(
1072        &self,
1073        run_id: &str,
1074        status: &str,
1075        final_state_json: &str,
1076        total_nodes: usize,
1077        events: &[(u64, String, String)],
1078        receipt_json: &str,
1079        bundle_json: &str,
1080    ) -> Result<String, String> {
1081        let key = self
1082            .require_integrity_key()
1083            .map_err(|_| "INTEGRITY_KEY_REQUIRED".to_owned())?;
1084        let receipt: Value = serde_json::from_str(receipt_json)
1085            .map_err(|e| format!("terminal receipt JSON error: {e}"))?;
1086        let receipt_digest = hmac_sha256(&receipt, key);
1087        let mut conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1088        let tx = conn
1089            .transaction()
1090            .map_err(|e| format!("terminal transaction begin error: {e}"))?;
1091        let changed = tx.execute(
1092            "UPDATE executions SET status = ?2, final_state_json = ?3, total_nodes = ?4,
1093             finished_at = CASE WHEN ?2 IN ('completed','failed','cancelled') THEN datetime('now') ELSE finished_at END
1094             WHERE run_id = ?1",
1095            params![run_id, status, final_state_json, total_nodes as i64],
1096        ).map_err(|e| format!("terminal execution update error: {e}"))?;
1097        if changed != 1 {
1098            return Err(format!("terminal projection run '{run_id}' was not found"));
1099        }
1100        tx.execute("DELETE FROM events WHERE run_id = ?1", params![run_id])
1101            .map_err(|e| format!("terminal event reset error: {e}"))?;
1102        for (seq, event_type, event_json) in events {
1103            tx.execute(
1104                "INSERT INTO events (run_id, seq, event_type, event_json) VALUES (?1, ?2, ?3, ?4)",
1105                params![run_id, *seq, event_type, event_json],
1106            )
1107            .map_err(|e| format!("terminal event insert error: {e}"))?;
1108        }
1109        #[cfg(test)]
1110        if self
1111            .terminal_projection_fault
1112            .swap(false, std::sync::atomic::Ordering::SeqCst)
1113        {
1114            return Err("injected terminal projection failure after events".into());
1115        }
1116        tx.execute(
1117            "INSERT INTO terminal_receipts (run_id, receipt_json, bundle_json, receipt_digest) VALUES (?1, ?2, ?3, ?4)
1118             ON CONFLICT(run_id) DO UPDATE SET receipt_json=excluded.receipt_json, bundle_json=excluded.bundle_json, receipt_digest=excluded.receipt_digest, persisted_at=datetime('now')",
1119            params![run_id, receipt_json, bundle_json, receipt_digest],
1120        ).map_err(|e| format!("terminal receipt insert error: {e}"))?;
1121        tx.commit()
1122            .map_err(|e| format!("terminal transaction commit error: {e}"))?;
1123        Ok(receipt_digest)
1124    }
1125
1126    pub fn load_terminal_receipt(&self, run_id: &str) -> Result<Option<Value>, String> {
1127        let key = self
1128            .require_integrity_key()
1129            .map_err(|_| "INTEGRITY_KEY_REQUIRED".to_owned())?;
1130        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1131        let row: Result<(String, String), _> = conn.query_row(
1132            "SELECT receipt_json, receipt_digest FROM terminal_receipts WHERE run_id = ?1",
1133            params![run_id],
1134            |row| Ok((row.get(0)?, row.get(1)?)),
1135        );
1136        match row {
1137            Ok((receipt_json, receipt_digest)) => {
1138                let receipt: Value = serde_json::from_str(&receipt_json)
1139                    .map_err(|e| format!("stored receipt JSON error: {e}"))?;
1140                if hmac_sha256(&receipt, key) != receipt_digest {
1141                    return Err("RECEIPT_INTEGRITY_FAILURE".into());
1142                }
1143                Ok(Some(
1144                    serde_json::json!({"receipt":receipt,"receipt_digest":receipt_digest,"storage_class":"sqlite_terminal_receipt","replay_capability":"integrity_only"}),
1145                ))
1146            }
1147            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1148            Err(e) => Err(format!("load terminal receipt error: {e}")),
1149        }
1150    }
1151
1152    /// A server restart cannot resume an in-flight graph. Make that interruption
1153    /// explicit instead of leaving a permanently misleading `running` row.
1154    pub fn recover_incomplete_executions(&self) -> Result<(), String> {
1155        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1156        conn.execute(
1157            "UPDATE executions
1158             SET status = 'interrupted_non_resumable', finished_at = datetime('now')
1159             WHERE status IN ('accepted', 'running')",
1160            [],
1161        )
1162        .map_err(|e| format!("recover executions error: {e}"))?;
1163        Ok(())
1164    }
1165
1166    /// Return the terminal projection retained by SQLite. This is deliberately
1167    /// not a resumable checkpoint or replay artifact.
1168    pub fn load_execution(&self, run_id: &str) -> Result<Option<Value>, String> {
1169        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1170        let mut stmt = conn
1171            .prepare(
1172                "SELECT graph_name, graph_hash, status, final_state_json, started_at, finished_at
1173                 FROM executions WHERE run_id = ?1",
1174            )
1175            .map_err(|e| format!("load execution error: {e}"))?;
1176        let row = stmt.query_row(params![run_id], |row| {
1177            Ok((
1178                row.get::<_, String>(0)?,
1179                row.get::<_, String>(1)?,
1180                row.get::<_, String>(2)?,
1181                row.get::<_, Option<String>>(3)?,
1182                row.get::<_, String>(4)?,
1183                row.get::<_, Option<String>>(5)?,
1184            ))
1185        });
1186        match row {
1187            Ok((graph_id, graph_version, status, final_state_json, started_at, finished_at)) => {
1188                let final_state = final_state_json
1189                    .as_deref()
1190                    .map(serde_json::from_str)
1191                    .transpose()
1192                    .map_err(|e| format!("stored final state JSON error: {e}"))?
1193                    .unwrap_or(Value::Null);
1194                Ok(Some(serde_json::json!({
1195                    "run_id": run_id,
1196                    "graph_id": graph_id,
1197                    "graph_version": graph_version,
1198                    "status": status,
1199                    "success": status == "completed",
1200                    "final_state": final_state,
1201                    "started_at": started_at,
1202                    "finished_at": finished_at,
1203                    "storage_class": "sqlite_terminal_record",
1204                    "durable_resume": false,
1205                    "replay_capability": "integrity_only"
1206                })))
1207            }
1208            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1209            Err(e) => Err(format!("load execution error: {e}")),
1210        }
1211    }
1212
1213    pub fn load_execution_contract(
1214        &self,
1215        run_id: &str,
1216    ) -> Result<Option<ExecutionContract>, String> {
1217        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1218        let result = conn.query_row(
1219            "SELECT graph_name, graph_hash, input_json, budgets_json
1220             FROM executions WHERE run_id = ?1",
1221            params![run_id],
1222            |row| {
1223                let input_json: Option<String> = row.get(2)?;
1224                let budgets_json: Option<String> = row.get(3)?;
1225                Ok(ExecutionContract {
1226                    graph_id: row.get(0)?,
1227                    graph_version: row.get(1)?,
1228                    input: input_json
1229                        .as_deref()
1230                        .map(serde_json::from_str)
1231                        .transpose()
1232                        .map_err(|error| {
1233                            rusqlite::Error::FromSqlConversionFailure(
1234                                2,
1235                                rusqlite::types::Type::Text,
1236                                Box::new(error),
1237                            )
1238                        })?
1239                        .unwrap_or(Value::Null),
1240                    budgets: budgets_json
1241                        .as_deref()
1242                        .map(serde_json::from_str)
1243                        .transpose()
1244                        .map_err(|error| {
1245                            rusqlite::Error::FromSqlConversionFailure(
1246                                3,
1247                                rusqlite::types::Type::Text,
1248                                Box::new(error),
1249                            )
1250                        })?
1251                        .unwrap_or(Value::Null),
1252                })
1253            },
1254        );
1255        match result {
1256            Ok(contract) => Ok(Some(contract)),
1257            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1258            Err(error) => Err(format!("load execution contract error: {error}")),
1259        }
1260    }
1261
1262    // ── Checkpoints ─────────────────────────────────────────────────────
1263
1264    pub fn create_resume_checkpoint(
1265        &self,
1266        run_id: &str,
1267        graph_id: &str,
1268        graph_version: &str,
1269        next_node_cursor: &str,
1270        state: &Value,
1271        budgets: &Value,
1272        budget_counters: &Value,
1273        dependency_summary: &Value,
1274        terminal_cursor: u64,
1275        event_cursor: u64,
1276    ) -> Result<CheckpointRecord, CheckpointError> {
1277        let key = self
1278            .require_integrity_key()
1279            .map_err(|_| CheckpointError::IntegrityKeyRequired)?;
1280        let state_json = serde_json::to_string(state).map_err(|_| CheckpointError::Persistence)?;
1281        let budgets_json =
1282            serde_json::to_string(budgets).map_err(|_| CheckpointError::Persistence)?;
1283        let counters_json =
1284            serde_json::to_string(budget_counters).map_err(|_| CheckpointError::Persistence)?;
1285        let dependency_json =
1286            serde_json::to_string(dependency_summary).map_err(|_| CheckpointError::Persistence)?;
1287        let state_digest = digest(state);
1288        let dependency_digest = digest(dependency_summary);
1289        let created_at = Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true);
1290        let checkpoint_id = format!("checkpoint-{run_id}-{next_node_cursor}");
1291        let mut record = CheckpointRecord {
1292            checkpoint_id,
1293            run_id: run_id.to_owned(),
1294            graph_id: graph_id.to_owned(),
1295            graph_version: graph_version.to_owned(),
1296            next_node_cursor: next_node_cursor.to_owned(),
1297            state: state.clone(),
1298            state_digest,
1299            budgets: budgets.clone(),
1300            budget_counters: budget_counters.clone(),
1301            dependency_summary: dependency_summary.clone(),
1302            dependency_digest,
1303            terminal_cursor,
1304            event_cursor,
1305            checkpoint_digest: String::new(),
1306            created_at,
1307            consumed_at: None,
1308        };
1309        record.checkpoint_digest = checkpoint_digest(&record, key);
1310
1311        let mut conn = self.conn.lock().map_err(|_| CheckpointError::Persistence)?;
1312        let tx = conn
1313            .transaction()
1314            .map_err(|_| CheckpointError::Persistence)?;
1315        #[cfg(test)]
1316        if self
1317            .checkpoint_persistence_fault
1318            .swap(false, std::sync::atomic::Ordering::SeqCst)
1319        {
1320            return Err(CheckpointError::Persistence);
1321        }
1322        tx.execute(
1323            "INSERT INTO checkpoints
1324             (run_id, node_id, attempt, input_json, status, checkpoint_id,
1325              graph_id, graph_version, next_cursor, state_json, state_digest,
1326              budgets_json, budget_counters_json, dependency_json,
1327              dependency_digest, terminal_cursor, event_cursor, checkpoint_digest,
1328              created_at, consumed_at)
1329             VALUES (?1, ?2, 0, ?3, 'available', ?4, ?5, ?6, ?7, ?3, ?8,
1330                     ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, NULL)",
1331            params![
1332                record.run_id,
1333                record.next_node_cursor,
1334                state_json,
1335                record.checkpoint_id,
1336                record.graph_id,
1337                record.graph_version,
1338                record.next_node_cursor,
1339                record.state_digest,
1340                budgets_json,
1341                counters_json,
1342                dependency_json,
1343                record.dependency_digest,
1344                record.terminal_cursor as i64,
1345                record.event_cursor as i64,
1346                record.checkpoint_digest,
1347                record.created_at,
1348            ],
1349        )
1350        .map_err(|_| CheckpointError::Persistence)?;
1351        tx.commit().map_err(|_| CheckpointError::Persistence)?;
1352        Ok(record)
1353    }
1354
1355    pub fn load_resume_checkpoint(
1356        &self,
1357        checkpoint_id: Option<&str>,
1358        run_id: Option<&str>,
1359    ) -> Result<Option<CheckpointRecord>, CheckpointError> {
1360        let key = self
1361            .require_integrity_key()
1362            .map_err(|_| CheckpointError::IntegrityKeyRequired)?;
1363        let conn = self.conn.lock().map_err(|_| CheckpointError::Persistence)?;
1364        let query = if checkpoint_id.is_some() {
1365            "SELECT run_id, graph_id, graph_version, next_cursor, state_json,
1366                    state_digest, budgets_json, budget_counters_json,
1367                    dependency_json, dependency_digest, terminal_cursor,
1368                    event_cursor, checkpoint_id, checkpoint_digest, created_at,
1369                    consumed_at
1370             FROM checkpoints WHERE checkpoint_id = ?1"
1371        } else {
1372            "SELECT run_id, graph_id, graph_version, next_cursor, state_json,
1373                    state_digest, budgets_json, budget_counters_json,
1374                    dependency_json, dependency_digest, terminal_cursor,
1375                    event_cursor, checkpoint_id, checkpoint_digest, created_at,
1376                    consumed_at
1377             FROM checkpoints WHERE run_id = ?1 AND checkpoint_id IS NOT NULL
1378             ORDER BY created_at DESC, checkpoint_id DESC LIMIT 1"
1379        };
1380        let selector = checkpoint_id.or(run_id).unwrap_or("");
1381        let row = conn.query_row(query, params![selector], checkpoint_row);
1382        let parts = match row {
1383            Ok(parts) => parts,
1384            Err(rusqlite::Error::QueryReturnedNoRows) => return Ok(None),
1385            Err(_) => return Err(CheckpointError::Persistence),
1386        };
1387        let record = checkpoint_from_parts(parts).ok_or(CheckpointError::Integrity)?;
1388        validate_checkpoint_record(&record, key)?;
1389        Ok(Some(record))
1390    }
1391
1392    pub fn consume_resume_checkpoint(
1393        &self,
1394        checkpoint_id: &str,
1395    ) -> Result<CheckpointRecord, CheckpointError> {
1396        let key = self
1397            .require_integrity_key()
1398            .map_err(|_| CheckpointError::IntegrityKeyRequired)?;
1399        let mut conn = self.conn.lock().map_err(|_| CheckpointError::Persistence)?;
1400        let tx = conn
1401            .transaction()
1402            .map_err(|_| CheckpointError::Persistence)?;
1403        let row = tx.query_row(
1404            "SELECT run_id, graph_id, graph_version, next_cursor, state_json,
1405                    state_digest, budgets_json, budget_counters_json,
1406                    dependency_json, dependency_digest, terminal_cursor,
1407                    event_cursor, checkpoint_id, checkpoint_digest, created_at,
1408                    consumed_at
1409             FROM checkpoints WHERE checkpoint_id = ?1",
1410            params![checkpoint_id],
1411            checkpoint_row,
1412        );
1413        let parts = match row {
1414            Ok(parts) => parts,
1415            Err(rusqlite::Error::QueryReturnedNoRows) => return Err(CheckpointError::NotFound),
1416            Err(_) => return Err(CheckpointError::Persistence),
1417        };
1418        let record = checkpoint_from_parts(parts).ok_or(CheckpointError::Integrity)?;
1419        validate_checkpoint_record(&record, key)?;
1420        if record.consumed_at.is_some() {
1421            return Err(CheckpointError::Consumed);
1422        }
1423        let consumed_at = Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true);
1424        let changed = tx
1425            .execute(
1426                "UPDATE checkpoints SET consumed_at = ?2, status = 'consumed'
1427                 WHERE checkpoint_id = ?1 AND consumed_at IS NULL",
1428                params![checkpoint_id, consumed_at],
1429            )
1430            .map_err(|_| CheckpointError::Persistence)?;
1431        if changed != 1 {
1432            return Err(CheckpointError::Consumed);
1433        }
1434        tx.commit().map_err(|_| CheckpointError::Persistence)?;
1435        let mut consumed = record;
1436        consumed.consumed_at = Some(consumed_at);
1437        Ok(consumed)
1438    }
1439
1440    // ── Durable checkpoint-bound approvals ───────────────────────────
1441
1442    pub fn create_checkpoint_approval(
1443        &self,
1444        checkpoint_id: &str,
1445        graph_id: &str,
1446        graph_version: &str,
1447        next_node_cursor: &str,
1448        expected_state: &Value,
1449        expected_budgets: &Value,
1450        expected_budget_counters: &Value,
1451        dependency_summary: &Value,
1452        audience: &str,
1453        prompt_digest: &str,
1454        allowed_decisions: &[String],
1455        expires_at: &str,
1456    ) -> Result<ApprovalRecord, ApprovalError> {
1457        let key = self
1458            .require_integrity_key()
1459            .map_err(|_| ApprovalError::IntegrityKeyRequired)?;
1460        let mut conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
1461        let tx = conn.transaction().map_err(|_| ApprovalError::Persistence)?;
1462        let checkpoint =
1463            load_checkpoint_from_tx(&tx, checkpoint_id, key).map_err(ApprovalError::Checkpoint)?;
1464        if checkpoint.consumed_at.is_some() {
1465            return Err(ApprovalError::Checkpoint(CheckpointError::Consumed));
1466        }
1467        if checkpoint.graph_id != graph_id
1468            || checkpoint.graph_version != graph_version
1469            || checkpoint.next_node_cursor != next_node_cursor
1470            || checkpoint.state != *expected_state
1471            || checkpoint.budgets != *expected_budgets
1472            || checkpoint.budget_counters != *expected_budget_counters
1473            || checkpoint.dependency_summary != *dependency_summary
1474            || checkpoint.terminal_cursor != 0
1475            || checkpoint.event_cursor != 0
1476        {
1477            return Err(ApprovalError::Checkpoint(CheckpointError::Integrity));
1478        }
1479
1480        let allowed_json =
1481            serde_json::to_string(allowed_decisions).map_err(|_| ApprovalError::Persistence)?;
1482        let existing = tx
1483            .query_row(
1484                &format!(
1485                    "SELECT {APPROVAL_COLUMNS} FROM approval_requests
1486                          WHERE checkpoint_id = ?1 AND audience = ?2 AND status = 'pending'"
1487                ),
1488                params![checkpoint_id, audience],
1489                approval_row,
1490            )
1491            .optional()
1492            .map_err(|_| ApprovalError::Persistence)?;
1493        if let Some(parts) = existing {
1494            let current = parse_approval_parts(parts, key)?;
1495            if current.graph_id == graph_id
1496                && current.graph_version == graph_version
1497                && current.checkpoint_digest == checkpoint.checkpoint_digest
1498                && current.prompt_digest == prompt_digest
1499                && current.allowed_decisions.as_slice() == allowed_decisions
1500                && current.expires_at == expires_at
1501            {
1502                tx.commit().map_err(|_| ApprovalError::Persistence)?;
1503                return Ok(current);
1504            }
1505            return Err(ApprovalError::Conflict);
1506        }
1507
1508        let created_at = Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true);
1509        let mut record = ApprovalRecord {
1510            approval_id: format!("approval-{}", uuid_like()),
1511            checkpoint_id: checkpoint.checkpoint_id.clone(),
1512            run_id: checkpoint.run_id.clone(),
1513            graph_id: graph_id.to_owned(),
1514            graph_version: graph_version.to_owned(),
1515            checkpoint_digest: checkpoint.checkpoint_digest.clone(),
1516            audience: audience.to_owned(),
1517            prompt_digest: prompt_digest.to_owned(),
1518            allowed_decisions: allowed_decisions.to_vec(),
1519            approval_digest: String::new(),
1520            status: "pending".into(),
1521            decision: None,
1522            decided_by: None,
1523            decided_at: None,
1524            expires_at: expires_at.to_owned(),
1525            created_at,
1526        };
1527        record.approval_digest = approval_digest(&record, key);
1528        tx.execute(
1529            "INSERT INTO approval_requests
1530             (approval_id, run_id, node_id, checkpoint_id, graph_id, graph_version,
1531              checkpoint_digest, audience, prompt, prompt_digest, allowed_decisions,
1532              approval_digest, status, expires_at, created_at)
1533             VALUES (?1, ?2, 'durable_checkpoint', ?3, ?4, ?5, ?6, ?7, '', ?8, ?9, ?10,
1534                     'pending', ?11, ?12)",
1535            params![
1536                record.approval_id,
1537                record.run_id,
1538                record.checkpoint_id,
1539                record.graph_id,
1540                record.graph_version,
1541                record.checkpoint_digest,
1542                record.audience,
1543                record.prompt_digest,
1544                allowed_json,
1545                record.approval_digest,
1546                record.expires_at,
1547                record.created_at,
1548            ],
1549        )
1550        .map_err(|_| ApprovalError::Persistence)?;
1551        tx.commit().map_err(|_| ApprovalError::Persistence)?;
1552        Ok(record)
1553    }
1554
1555    pub fn get_checkpoint_approval(
1556        &self,
1557        approval_id: &str,
1558    ) -> Result<Option<ApprovalRecord>, ApprovalError> {
1559        let key = self
1560            .require_integrity_key()
1561            .map_err(|_| ApprovalError::IntegrityKeyRequired)?;
1562        let conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
1563        let result = conn.query_row(
1564            &format!(
1565                "SELECT {APPROVAL_COLUMNS} FROM approval_requests
1566                      WHERE approval_id = ?1 AND checkpoint_id IS NOT NULL"
1567            ),
1568            params![approval_id],
1569            approval_row,
1570        );
1571        match result {
1572            Ok(parts) => Ok(Some(parse_approval_parts(parts, key)?)),
1573            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1574            Err(_) => Err(ApprovalError::Persistence),
1575        }
1576    }
1577
1578    pub fn list_checkpoint_approvals(
1579        &self,
1580        run_id: Option<&str>,
1581        status: Option<&str>,
1582        limit: usize,
1583    ) -> Result<Vec<ApprovalRecord>, ApprovalError> {
1584        let key = self
1585            .require_integrity_key()
1586            .map_err(|_| ApprovalError::IntegrityKeyRequired)?;
1587        let conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
1588        let mut sql = format!(
1589            "SELECT {APPROVAL_COLUMNS} FROM approval_requests
1590                              WHERE checkpoint_id IS NOT NULL"
1591        );
1592        if run_id.is_some() {
1593            sql.push_str(" AND run_id = ?1");
1594        }
1595        if status.is_some() {
1596            sql.push_str(if run_id.is_some() {
1597                " AND status = ?2"
1598            } else {
1599                " AND status = ?1"
1600            });
1601        }
1602        sql.push_str(" ORDER BY created_at DESC, approval_id DESC LIMIT ?3");
1603        if run_id.is_none() && status.is_none() {
1604            sql = format!("SELECT {APPROVAL_COLUMNS} FROM approval_requests
1605                          WHERE checkpoint_id IS NOT NULL ORDER BY created_at DESC, approval_id DESC LIMIT ?1");
1606        } else if run_id.is_none() || status.is_none() {
1607            sql = sql.replace("LIMIT ?3", "LIMIT ?2");
1608        }
1609        let mut statement = conn.prepare(&sql).map_err(|_| ApprovalError::Persistence)?;
1610        let mut rows = if run_id.is_some() && status.is_some() {
1611            statement.query(params![run_id, status, limit.min(200) as i64])
1612        } else if run_id.is_some() {
1613            statement.query(params![run_id, limit.min(200) as i64])
1614        } else if status.is_some() {
1615            statement.query(params![status, limit.min(200) as i64])
1616        } else {
1617            statement.query(params![limit.min(200) as i64])
1618        }
1619        .map_err(|_| ApprovalError::Persistence)?;
1620        let mut approvals = Vec::new();
1621        while let Some(row) = rows.next().map_err(|_| ApprovalError::Persistence)? {
1622            approvals.push(parse_approval_parts(
1623                approval_row(row).map_err(|_| ApprovalError::Persistence)?,
1624                key,
1625            )?);
1626        }
1627        Ok(approvals)
1628    }
1629
1630    pub fn checkpoint_approval_status(
1631        &self,
1632        checkpoint_id: &str,
1633    ) -> Result<Option<String>, ApprovalError> {
1634        let conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
1635        conn.query_row(
1636            "SELECT status FROM approval_requests
1637             WHERE checkpoint_id = ?1 AND checkpoint_id IS NOT NULL
1638             ORDER BY created_at DESC LIMIT 1",
1639            params![checkpoint_id],
1640            |row| row.get(0),
1641        )
1642        .optional()
1643        .map_err(|_| ApprovalError::Persistence)
1644    }
1645
1646    pub fn decide_checkpoint_approval(
1647        &self,
1648        approval_id: &str,
1649        decision: &str,
1650        actor: &str,
1651        now: DateTime<Utc>,
1652    ) -> Result<ApprovedCheckpoint, ApprovalError> {
1653        let key = self
1654            .require_integrity_key()
1655            .map_err(|_| ApprovalError::IntegrityKeyRequired)?;
1656        let mut conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
1657        let tx = conn.transaction().map_err(|_| ApprovalError::Persistence)?;
1658        let parts = tx
1659            .query_row(
1660                &format!(
1661                    "SELECT {APPROVAL_COLUMNS} FROM approval_requests
1662                          WHERE approval_id = ?1 AND checkpoint_id IS NOT NULL"
1663                ),
1664                params![approval_id],
1665                approval_row,
1666            )
1667            .optional()
1668            .map_err(|_| ApprovalError::Persistence)?
1669            .ok_or(ApprovalError::NotFound)?;
1670        let mut approval = parse_approval_parts(parts, key)?;
1671        if approval.status == "expired" {
1672            return Err(ApprovalError::Expired);
1673        }
1674        if approval.status != "pending" {
1675            return Err(ApprovalError::AlreadyDecided);
1676        }
1677        let expired = DateTime::parse_from_rfc3339(&approval.expires_at)
1678            .map_err(|_| ApprovalError::Integrity)?
1679            .with_timezone(&Utc)
1680            <= now;
1681        if expired {
1682            let checkpoint = load_checkpoint_from_tx(&tx, &approval.checkpoint_id, key)
1683                .map_err(ApprovalError::Checkpoint)?;
1684            if checkpoint.run_id != approval.run_id
1685                || checkpoint.graph_id != approval.graph_id
1686                || checkpoint.graph_version != approval.graph_version
1687                || checkpoint.checkpoint_digest != approval.checkpoint_digest
1688            {
1689                return Err(ApprovalError::Integrity);
1690            }
1691            approval.status = "expired".into();
1692            approval.decision = None;
1693            approval.decided_by = None;
1694            approval.decided_at = Some(now.to_rfc3339_opts(SecondsFormat::Nanos, true));
1695            approval.approval_digest = approval_digest(&approval, key);
1696            let changed = tx
1697                .execute(
1698                    "UPDATE approval_requests SET status = 'expired', decision = NULL,
1699                     decided_by = NULL, decided_at = ?2, approval_digest = ?3
1700                     WHERE approval_id = ?1 AND status = 'pending'",
1701                    params![approval_id, approval.decided_at, approval.approval_digest],
1702                )
1703                .map_err(|_| ApprovalError::Persistence)?;
1704            if changed != 1 {
1705                return Err(ApprovalError::AlreadyDecided);
1706            }
1707            if checkpoint.consumed_at.is_none() {
1708                tx.execute(
1709                    "UPDATE checkpoints SET consumed_at = ?2, status = 'consumed'
1710                     WHERE checkpoint_id = ?1 AND consumed_at IS NULL",
1711                    params![
1712                        approval.checkpoint_id,
1713                        now.to_rfc3339_opts(SecondsFormat::Nanos, true)
1714                    ],
1715                )
1716                .map_err(|_| ApprovalError::Persistence)?;
1717            }
1718            tx.commit().map_err(|_| ApprovalError::Persistence)?;
1719            return Err(ApprovalError::Expired);
1720        }
1721
1722        if !approval
1723            .allowed_decisions
1724            .iter()
1725            .any(|allowed| allowed == decision)
1726        {
1727            return Err(ApprovalError::DecisionNotAllowed);
1728        }
1729
1730        let checkpoint = load_checkpoint_from_tx(&tx, &approval.checkpoint_id, key)
1731            .map_err(ApprovalError::Checkpoint)?;
1732        if checkpoint.run_id != approval.run_id
1733            || checkpoint.graph_id != approval.graph_id
1734            || checkpoint.graph_version != approval.graph_version
1735            || checkpoint.checkpoint_digest != approval.checkpoint_digest
1736        {
1737            return Err(ApprovalError::Integrity);
1738        }
1739        if checkpoint.consumed_at.is_some() {
1740            return Err(ApprovalError::Checkpoint(CheckpointError::Consumed));
1741        }
1742
1743        let decided_at = now.to_rfc3339_opts(SecondsFormat::Nanos, true);
1744        let status = if decision == "approve" {
1745            "approved"
1746        } else {
1747            "rejected"
1748        };
1749        approval.status = status.into();
1750        approval.decision = Some(decision.into());
1751        approval.decided_by = Some(actor.into());
1752        approval.decided_at = Some(decided_at.clone());
1753        approval.approval_digest = approval_digest(&approval, key);
1754        let changed = tx
1755            .execute(
1756                "UPDATE approval_requests SET status = ?2, decision = ?3,
1757                 decided_by = ?4, decided_at = ?5, approval_digest = ?6
1758                 WHERE approval_id = ?1 AND status = 'pending'",
1759                params![
1760                    approval_id,
1761                    status,
1762                    decision,
1763                    actor,
1764                    decided_at,
1765                    approval.approval_digest
1766                ],
1767            )
1768            .map_err(|_| ApprovalError::Persistence)?;
1769        if changed != 1 {
1770            return Err(ApprovalError::AlreadyDecided);
1771        }
1772        let consumed_at = now.to_rfc3339_opts(SecondsFormat::Nanos, true);
1773        let checkpoint_changed = tx
1774            .execute(
1775                "UPDATE checkpoints SET consumed_at = ?2, status = 'consumed'
1776                 WHERE checkpoint_id = ?1 AND consumed_at IS NULL",
1777                params![approval.checkpoint_id, consumed_at],
1778            )
1779            .map_err(|_| ApprovalError::Persistence)?;
1780        if checkpoint_changed != 1 {
1781            return Err(ApprovalError::Checkpoint(CheckpointError::Consumed));
1782        }
1783        tx.commit().map_err(|_| ApprovalError::Persistence)?;
1784        Ok(ApprovedCheckpoint {
1785            approval,
1786            checkpoint,
1787        })
1788    }
1789
1790    pub fn approval_receipt_value(approval: &ApprovalRecord) -> Value {
1791        let decision = approval.decision.clone().unwrap_or_default();
1792        let decided_by_digest = approval
1793            .decided_by
1794            .as_ref()
1795            .map(|actor| digest(&Value::String(actor.clone())));
1796        let decision_digest = digest(&serde_json::json!({
1797            "approval_digest": approval.approval_digest,
1798            "decision": decision,
1799            "decided_by_digest": decided_by_digest,
1800            "decided_at": approval.decided_at,
1801        }));
1802        serde_json::json!({
1803            "approval_id": approval.approval_id,
1804            "approval_digest": approval.approval_digest,
1805            "checkpoint_id": approval.checkpoint_id,
1806            "checkpoint_digest": approval.checkpoint_digest,
1807            "decision": approval.decision,
1808            "decided_by_digest": decided_by_digest,
1809            "decided_at": approval.decided_at,
1810            "allowed_decisions_digest": digest(&serde_json::json!(approval.allowed_decisions)),
1811            "audience_digest": digest(&Value::String(approval.audience.clone())),
1812            "prompt_digest": approval.prompt_digest,
1813            "decision_digest": decision_digest,
1814        })
1815    }
1816
1817    #[cfg(test)]
1818    pub(crate) fn fail_checkpoint_persistence(&self) {
1819        self.checkpoint_persistence_fault
1820            .store(true, std::sync::atomic::Ordering::SeqCst);
1821    }
1822
1823    pub fn save_checkpoint(
1824        &self,
1825        run_id: &str,
1826        node_id: &str,
1827        attempt: u32,
1828        input_json: &str,
1829        output_json: Option<&str>,
1830        status: &str,
1831        error: Option<&str>,
1832    ) -> Result<(), String> {
1833        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1834        conn.execute(
1835            "INSERT INTO checkpoints (run_id, node_id, attempt, input_json, output_json, status, error)
1836             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
1837             ON CONFLICT(run_id, node_id, attempt) DO UPDATE SET
1838                output_json = COALESCE(excluded.output_json, output_json),
1839                status = excluded.status,
1840                error = COALESCE(excluded.error, error)",
1841            params![run_id, node_id, attempt, input_json, output_json, status, error],
1842        )
1843        .map_err(|e| format!("save_checkpoint error: {e}"))?;
1844        Ok(())
1845    }
1846
1847    // ── Events ──────────────────────────────────────────────────────────
1848
1849    pub fn save_event(
1850        &self,
1851        run_id: &str,
1852        seq: u64,
1853        event_type: &str,
1854        event_json: &str,
1855    ) -> Result<(), String> {
1856        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1857        conn.execute(
1858            "INSERT OR IGNORE INTO events (run_id, seq, event_type, event_json)
1859             VALUES (?1, ?2, ?3, ?4)",
1860            params![run_id, seq, event_type, event_json],
1861        )
1862        .map_err(|e| format!("save_event error: {e}"))?;
1863        Ok(())
1864    }
1865
1866    pub fn load_events(
1867        &self,
1868        run_id: &str,
1869        cursor: u64,
1870        limit: usize,
1871    ) -> Result<Option<Value>, String> {
1872        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1873        let mut stmt = conn
1874            .prepare("SELECT seq, event_json FROM events WHERE run_id = ?1 AND seq >= ?2 ORDER BY seq LIMIT ?3")
1875            .map_err(|e| format!("load events error: {e}"))?;
1876        let events: Vec<Value> = stmt
1877            .query_map(params![run_id, cursor, limit.min(200) as i64], |row| {
1878                let seq: u64 = row.get(0)?;
1879                let event_json: String = row.get(1)?;
1880                let event = serde_json::from_str::<Value>(&event_json)
1881                    .unwrap_or_else(|_| serde_json::json!({"receipt":"terminal event persisted with reduced fidelity"}));
1882                Ok(serde_json::json!({"cursor": seq, "event": event}))
1883            })
1884            .map_err(|e| format!("load events error: {e}"))?
1885            .filter_map(Result::ok)
1886            .collect();
1887        let exists: bool = conn
1888            .query_row(
1889                "SELECT EXISTS(SELECT 1 FROM events WHERE run_id = ?1)",
1890                params![run_id],
1891                |row| row.get(0),
1892            )
1893            .map_err(|e| format!("load events existence error: {e}"))?;
1894        if !exists {
1895            return Ok(None);
1896        }
1897        let first: u64 = conn
1898            .query_row(
1899                "SELECT MIN(seq) FROM events WHERE run_id = ?1",
1900                params![run_id],
1901                |row| row.get(0),
1902            )
1903            .map_err(|e| format!("load events first cursor error: {e}"))?;
1904        let next_cursor: u64 = conn
1905            .query_row(
1906                "SELECT COALESCE(MAX(seq) + 1, 0) FROM events WHERE run_id = ?1",
1907                params![run_id],
1908                |row| row.get(0),
1909            )
1910            .map_err(|e| format!("load events next cursor error: {e}"))?;
1911        Ok(Some(
1912            serde_json::json!({"run_id":run_id,"events":events,"next_cursor":next_cursor,"gap":cursor<first,"truncated":false,"dropped":first,"projection":"terminal_persisted_projection","replay_capability":"integrity_only","replayable_execution":false,"resume_supported":false}),
1913        ))
1914    }
1915
1916    // ── Idempotency ─────────────────────────────────────────────────────
1917
1918    /// Look up a cached idempotent response together with the canonical request
1919    /// digest it was bound to. A NULL digest is a pre-migration record and must
1920    /// never be replayed for a new request.
1921    pub fn check_idempotency(&self, key: &str) -> Result<Option<(Option<String>, Value)>, String> {
1922        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1923        let mut stmt = conn
1924            .prepare("SELECT request_digest, result_json FROM idempotency_keys WHERE key = ?1 AND request_digest IS NOT NULL AND valid = 1")
1925            .map_err(|e| format!("idempotency error: {e}"))?;
1926        let result: Option<(Option<String>, String)> = stmt
1927            .query_row(params![key], |row| Ok((row.get(0)?, row.get(1)?)))
1928            .ok();
1929        match result {
1930            Some((request_digest, json_str)) => serde_json::from_str(&json_str)
1931                .map(|result| Some((request_digest, result)))
1932                .map_err(|e| format!("json parse error: {e}")),
1933            None => Ok(None),
1934        }
1935    }
1936
1937    pub fn save_idempotency(
1938        &self,
1939        key: &str,
1940        request_digest: &str,
1941        result_json: &str,
1942    ) -> Result<bool, String> {
1943        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1944        let inserted = conn.execute(
1945            "INSERT OR IGNORE INTO idempotency_keys (key, request_digest, result_json) VALUES (?1, ?2, ?3)",
1946            params![key, request_digest, result_json],
1947        )
1948        .map_err(|e| format!("idempotency error: {e}"))?;
1949        Ok(inserted == 1)
1950    }
1951
1952    pub fn data_dir(&self) -> Option<PathBuf> {
1953        // Returns None since we don't store the path separately.
1954        // The store is identified by its existence, not a path reference.
1955        None
1956    }
1957}
1958
1959#[cfg(test)]
1960mod tests {
1961    use super::*;
1962
1963    fn configure_test_integrity_key() {
1964        let path = std::env::temp_dir().join("agent-graph-mcp-unit-integrity.key");
1965        std::fs::write(&path, [0x5au8; 32]).expect("test integrity key");
1966        std::env::set_var("AGENT_GRAPH_INTEGRITY_KEY_PATH", path);
1967    }
1968
1969    #[test]
1970    fn graph_delete_fault_rolls_back_graph_and_versions_together() {
1971        let temp = tempfile::tempdir().expect("graph database");
1972        let store = PersistentStore::open(temp.path()).expect("store");
1973        store
1974            .save_graph("atomic-delete", "{\"name\":\"atomic-delete\"}", "v1", false)
1975            .expect("graph");
1976        store.fail_graph_delete_after_graph_row();
1977        assert!(store.delete_graph("atomic-delete").is_err());
1978        assert!(store.load_graph("atomic-delete").unwrap().is_some());
1979        assert_eq!(
1980            store.list_graph_versions("atomic-delete").unwrap(),
1981            vec!["v1"]
1982        );
1983    }
1984
1985    #[test]
1986    fn checkpoint_persistence_fault_leaves_no_resumable_row() {
1987        configure_test_integrity_key();
1988        let temp = tempfile::tempdir().expect("checkpoint database");
1989        let store = PersistentStore::open(temp.path()).expect("store");
1990        store
1991            .save_graph("checkpoint-fault", "{}", "version", false)
1992            .expect("graph");
1993        store
1994            .save_execution_with_budgets(
1995                "run-checkpoint-fault",
1996                "checkpoint-fault",
1997                "version",
1998                "checkpointed",
1999                "{}",
2000                Some("null"),
2001            )
2002            .expect("execution");
2003        store.fail_checkpoint_persistence();
2004        let result = store.create_resume_checkpoint(
2005            "run-checkpoint-fault",
2006            "checkpoint-fault",
2007            "version",
2008            "entry",
2009            &serde_json::json!({"__input__":null}),
2010            &Value::Null,
2011            &serde_json::json!({"nodes":0,"llm_calls":0,"wall_clock_ms":0}),
2012            &serde_json::json!({"eligible":true}),
2013            0,
2014            0,
2015        );
2016        assert_eq!(result, Err(CheckpointError::Persistence));
2017        assert_eq!(
2018            store.load_resume_checkpoint(None, Some("run-checkpoint-fault")),
2019            Ok(None)
2020        );
2021    }
2022}