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 crate::operator_auth::OperatorAction;
11use chrono::{DateTime, SecondsFormat, Utc};
12use rusqlite::{params, Connection, OptionalExtension};
13use serde_json::Value;
14use std::path::{Path, PathBuf};
15use std::sync::Mutex;
16
17/// Persistent store wrapping a SQLite connection.
18pub struct PersistentStore {
19    conn: std::sync::Arc<Mutex<Connection>>,
20    integrity_key: Option<std::sync::Arc<[u8]>>,
21    #[cfg(test)]
22    terminal_projection_fault: std::sync::Arc<std::sync::atomic::AtomicBool>,
23    #[cfg(test)]
24    checkpoint_persistence_fault: std::sync::Arc<std::sync::atomic::AtomicBool>,
25    #[cfg(test)]
26    graph_delete_fault: std::sync::Arc<std::sync::atomic::AtomicBool>,
27}
28
29#[derive(Debug, Clone, PartialEq)]
30pub struct CheckpointRecord {
31    pub checkpoint_id: String,
32    pub run_id: String,
33    pub graph_id: String,
34    pub graph_version: String,
35    pub next_node_cursor: String,
36    pub state: Value,
37    pub state_digest: String,
38    pub budgets: Value,
39    pub budget_counters: Value,
40    pub dependency_summary: Value,
41    pub dependency_digest: String,
42    pub terminal_cursor: u64,
43    pub event_cursor: u64,
44    pub checkpoint_digest: String,
45    pub created_at: String,
46    pub consumed_at: Option<String>,
47}
48
49#[derive(Debug, Clone, PartialEq, Eq)]
50pub enum CheckpointError {
51    NotFound,
52    Consumed,
53    Integrity,
54    Persistence,
55    IntegrityKeyRequired,
56}
57
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub enum ApprovalError {
60    NotFound,
61    Conflict,
62    AlreadyDecided,
63    Expired,
64    DecisionNotAllowed,
65    Integrity,
66    Checkpoint(CheckpointError),
67    Persistence,
68    IntegrityKeyRequired,
69}
70
71impl ApprovalError {
72    pub fn code(&self) -> &'static str {
73        match self {
74            Self::NotFound => "APPROVAL_NOT_FOUND",
75            Self::Conflict => "APPROVAL_REQUEST_CONFLICT",
76            Self::AlreadyDecided => "APPROVAL_ALREADY_DECIDED",
77            Self::Expired => "APPROVAL_EXPIRED",
78            Self::DecisionNotAllowed => "APPROVAL_DECISION_NOT_ALLOWED",
79            Self::Integrity => "APPROVAL_INTEGRITY_FAILURE",
80            Self::Checkpoint(error) => error.code(),
81            Self::Persistence => "APPROVAL_PERSISTENCE_FAILURE",
82            Self::IntegrityKeyRequired => "INTEGRITY_KEY_REQUIRED",
83        }
84    }
85
86    pub fn message(&self) -> String {
87        match self {
88            Self::NotFound => "approval was not found".into(),
89            Self::Conflict => {
90                "a conflicting pending approval already exists for this checkpoint and audience"
91                    .into()
92            }
93            Self::AlreadyDecided => "approval has already been decided".into(),
94            Self::Expired => "approval has expired".into(),
95            Self::DecisionNotAllowed => "decision is not allowed by this approval".into(),
96            Self::Integrity => "approval integrity validation failed".into(),
97            Self::Checkpoint(error) => error.message().into(),
98            Self::Persistence => "approval persistence failed".into(),
99            Self::IntegrityKeyRequired => {
100                "integrity key is required for durable approval operations".into()
101            }
102        }
103    }
104}
105
106#[derive(Debug, Clone, PartialEq, Eq)]
107pub struct ApprovalRecord {
108    pub approval_id: String,
109    pub checkpoint_id: String,
110    pub run_id: String,
111    pub graph_id: String,
112    pub graph_version: String,
113    pub checkpoint_digest: String,
114    pub audience: String,
115    pub prompt_digest: String,
116    pub allowed_decisions: Vec<String>,
117    pub approval_digest: String,
118    pub status: String,
119    pub decision: Option<String>,
120    pub decided_by: Option<String>,
121    pub decided_at: Option<String>,
122    pub expires_at: String,
123    pub created_at: String,
124}
125
126#[derive(Debug, Clone)]
127pub struct ApprovedCheckpoint {
128    pub approval: ApprovalRecord,
129    pub checkpoint: CheckpointRecord,
130}
131
132impl CheckpointError {
133    pub fn code(&self) -> &'static str {
134        match self {
135            Self::NotFound => "CHECKPOINT_NOT_FOUND",
136            Self::Consumed => "CHECKPOINT_CONSUMED",
137            Self::Integrity => "CHECKPOINT_INTEGRITY_FAILURE",
138            Self::Persistence => "CHECKPOINT_PERSISTENCE_FAILURE",
139            Self::IntegrityKeyRequired => "INTEGRITY_KEY_REQUIRED",
140        }
141    }
142
143    pub fn message(&self) -> &'static str {
144        match self {
145            Self::NotFound => "checkpoint was not found",
146            Self::Consumed => "checkpoint has already been consumed",
147            Self::Integrity => "checkpoint integrity validation failed",
148            Self::Persistence => "checkpoint persistence failed; resumability was not advertised",
149            Self::IntegrityKeyRequired => {
150                "integrity key is required for durable checkpoint operations"
151            }
152        }
153    }
154}
155
156#[derive(Debug, Clone)]
157pub struct ExecutionContract {
158    pub graph_id: String,
159    pub graph_version: String,
160    pub input: Value,
161    pub budgets: Value,
162}
163
164type CheckpointParts = (
165    String,
166    String,
167    String,
168    Option<String>,
169    Option<String>,
170    Option<String>,
171    Option<String>,
172    Option<String>,
173    Option<String>,
174    Option<String>,
175    Option<i64>,
176    Option<i64>,
177    Option<String>,
178    Option<String>,
179    Option<String>,
180    Option<String>,
181);
182
183fn checkpoint_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<CheckpointParts> {
184    Ok((
185        row.get(0)?,
186        row.get(1)?,
187        row.get(2)?,
188        row.get(3)?,
189        row.get(4)?,
190        row.get(5)?,
191        row.get(6)?,
192        row.get(7)?,
193        row.get(8)?,
194        row.get(9)?,
195        row.get(10)?,
196        row.get(11)?,
197        row.get(12)?,
198        row.get(13)?,
199        row.get(14)?,
200        row.get(15)?,
201    ))
202}
203
204fn checkpoint_from_parts(parts: CheckpointParts) -> Option<CheckpointRecord> {
205    let (
206        run_id,
207        graph_id,
208        graph_version,
209        next_node_cursor,
210        state_json,
211        state_digest,
212        budgets_json,
213        counters_json,
214        dependency_json,
215        dependency_digest,
216        terminal_cursor,
217        event_cursor,
218        checkpoint_id,
219        checkpoint_digest,
220        created_at,
221        consumed_at,
222    ) = parts;
223    Some(CheckpointRecord {
224        checkpoint_id: checkpoint_id?,
225        run_id,
226        graph_id,
227        graph_version,
228        next_node_cursor: next_node_cursor?,
229        state: serde_json::from_str(&state_json?).ok()?,
230        state_digest: state_digest?,
231        budgets: serde_json::from_str(&budgets_json?).ok()?,
232        budget_counters: serde_json::from_str(&counters_json?).ok()?,
233        dependency_summary: serde_json::from_str(&dependency_json?).ok()?,
234        dependency_digest: dependency_digest?,
235        terminal_cursor: u64::try_from(terminal_cursor?).ok()?,
236        event_cursor: u64::try_from(event_cursor?).ok()?,
237        checkpoint_digest: checkpoint_digest?,
238        created_at: created_at?,
239        consumed_at,
240    })
241}
242
243fn checkpoint_digest(record: &CheckpointRecord, key: &[u8]) -> String {
244    hmac_sha256(
245        &serde_json::json!({
246            "checkpoint_id": record.checkpoint_id,
247            "run_id": record.run_id,
248            "graph_id": record.graph_id,
249            "graph_version": record.graph_version,
250            "next_node_cursor": record.next_node_cursor,
251            "state": record.state,
252            "state_digest": record.state_digest,
253            "budgets": record.budgets,
254            "budget_counters": record.budget_counters,
255            "dependency_summary": record.dependency_summary,
256            "dependency_digest": record.dependency_digest,
257            "terminal_cursor": record.terminal_cursor,
258            "event_cursor": record.event_cursor,
259            "created_at": record.created_at,
260        }),
261        key,
262    )
263}
264
265fn validate_checkpoint_record(
266    record: &CheckpointRecord,
267    key: &[u8],
268) -> Result<(), CheckpointError> {
269    if record.checkpoint_id != format!("checkpoint-{}-{}", record.run_id, record.next_node_cursor)
270        || record.state_digest != digest(&record.state)
271        || record.dependency_digest != digest(&record.dependency_summary)
272        || record.checkpoint_digest != checkpoint_digest(record, key)
273    {
274        return Err(CheckpointError::Integrity);
275    }
276    Ok(())
277}
278
279type ApprovalParts = (
280    String,
281    Option<String>,
282    String,
283    Option<String>,
284    Option<String>,
285    Option<String>,
286    String,
287    String,
288    Option<String>,
289    Option<String>,
290    String,
291    Option<String>,
292    Option<String>,
293    Option<String>,
294    String,
295    String,
296    String,
297);
298
299fn approval_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<ApprovalParts> {
300    Ok((
301        row.get(0)?,
302        row.get(1)?,
303        row.get(2)?,
304        row.get(3)?,
305        row.get(4)?,
306        row.get(5)?,
307        row.get(6)?,
308        row.get(7)?,
309        row.get(8)?,
310        row.get(9)?,
311        row.get(10)?,
312        row.get(11)?,
313        row.get(12)?,
314        row.get(13)?,
315        row.get(14)?,
316        row.get(15)?,
317        row.get(16)?,
318    ))
319}
320
321fn approval_from_parts(parts: ApprovalParts) -> Option<ApprovalRecord> {
322    let (
323        approval_id,
324        checkpoint_id,
325        run_id,
326        graph_id,
327        graph_version,
328        checkpoint_digest,
329        audience,
330        prompt_digest,
331        allowed_decisions,
332        approval_digest,
333        status,
334        decision,
335        decided_by,
336        decided_at,
337        expires_at,
338        created_at,
339        _legacy,
340    ) = parts;
341    let allowed_decisions = serde_json::from_str(&allowed_decisions?).ok()?;
342    Some(ApprovalRecord {
343        approval_id,
344        checkpoint_id: checkpoint_id?,
345        run_id,
346        graph_id: graph_id?,
347        graph_version: graph_version?,
348        checkpoint_digest: checkpoint_digest?,
349        audience,
350        prompt_digest,
351        allowed_decisions,
352        approval_digest: approval_digest?,
353        status,
354        decision,
355        decided_by,
356        decided_at,
357        expires_at,
358        created_at,
359    })
360}
361
362const APPROVAL_COLUMNS: &str = "approval_id, checkpoint_id, run_id, graph_id,
363    graph_version, checkpoint_digest, audience, prompt_digest,
364    allowed_decisions, approval_digest, status, decision, decided_by,
365    decided_at, expires_at, created_at, prompt";
366
367fn approval_digest(record: &ApprovalRecord, key: &[u8]) -> String {
368    hmac_sha256(
369        &serde_json::json!({
370            "approval_id": record.approval_id,
371            "checkpoint_id": record.checkpoint_id,
372            "run_id": record.run_id,
373            "graph_id": record.graph_id,
374            "graph_version": record.graph_version,
375            "checkpoint_digest": record.checkpoint_digest,
376            "audience": record.audience,
377            "prompt_digest": record.prompt_digest,
378            "allowed_decisions": record.allowed_decisions,
379            "status": record.status,
380            "decision": record.decision,
381            "decided_by": record.decided_by,
382            "decided_at": record.decided_at,
383            "expires_at": record.expires_at,
384            "created_at": record.created_at,
385        }),
386        key,
387    )
388}
389
390fn parse_approval_parts(parts: ApprovalParts, key: &[u8]) -> Result<ApprovalRecord, ApprovalError> {
391    let record = approval_from_parts(parts).ok_or(ApprovalError::Integrity)?;
392    if record.approval_digest != approval_digest(&record, key) {
393        return Err(ApprovalError::Integrity);
394    }
395    Ok(record)
396}
397
398fn uuid_like() -> String {
399    digest(&Value::String(
400        Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true),
401    ))
402    .trim_start_matches("sha256:")
403    .to_owned()
404}
405
406fn load_checkpoint_from_tx(
407    tx: &rusqlite::Transaction<'_>,
408    checkpoint_id: &str,
409    key: &[u8],
410) -> Result<CheckpointRecord, CheckpointError> {
411    let row = tx.query_row(
412        "SELECT run_id, graph_id, graph_version, next_cursor, state_json,
413                state_digest, budgets_json, budget_counters_json,
414                dependency_json, dependency_digest, terminal_cursor,
415                event_cursor, checkpoint_id, checkpoint_digest, created_at,
416                consumed_at
417         FROM checkpoints WHERE checkpoint_id = ?1",
418        params![checkpoint_id],
419        checkpoint_row,
420    );
421    let parts = match row {
422        Ok(parts) => parts,
423        Err(rusqlite::Error::QueryReturnedNoRows) => return Err(CheckpointError::NotFound),
424        Err(_) => return Err(CheckpointError::Persistence),
425    };
426    let record = checkpoint_from_parts(parts).ok_or(CheckpointError::Integrity)?;
427    validate_checkpoint_record(&record, key)?;
428    Ok(record)
429}
430
431#[derive(Debug, Clone, Copy, PartialEq, Eq)]
432pub enum GraphDeleteResult {
433    Deleted,
434    NotFound,
435    Referenced,
436    RetentionApprovalRequired,
437    ReferencedBySubgraph,
438}
439
440#[derive(Debug, Clone, PartialEq, Eq)]
441pub enum GraphRetentionError {
442    NotFound,
443    InvalidState,
444    InvalidTransition,
445    Referenced,
446    ReferencedBySubgraph,
447}
448
449#[derive(Debug, Clone, PartialEq, Eq)]
450pub struct GraphRetentionRecord {
451    pub graph_id: String,
452    pub state: String,
453    pub reason: String,
454    pub actor: String,
455    pub review_after: Option<String>,
456    pub updated_at: String,
457}
458
459#[derive(Debug, Clone, PartialEq, Eq)]
460pub struct GraphRetentionReport {
461    pub graph_id: String,
462    pub state: String,
463    pub reason: Option<String>,
464    pub actor: Option<String>,
465    pub review_after: Option<String>,
466    pub created_at: String,
467    pub updated_at: String,
468    pub version_count: u64,
469    pub execution_count: u64,
470    pub last_execution_at: Option<String>,
471    pub inbound_subgraph_refs: Vec<String>,
472    pub deletion_eligible: bool,
473    pub state_digest: String,
474    pub tombstoned: bool,
475}
476
477#[derive(Debug, Clone)]
478pub struct OperatorRetentionRequest {
479    pub request_digest: String,
480    pub action: OperatorAction,
481    pub graph_id: String,
482    pub expected_state_digest: String,
483    pub nonce: String,
484    pub operator_uid: u32,
485    pub daemon_instance_id: String,
486    pub issued_at: String,
487    pub expires_at: String,
488    pub state: Option<String>,
489    pub reason: Option<String>,
490    pub review_after: Option<String>,
491}
492
493#[derive(Debug, Clone, PartialEq, Eq)]
494pub enum OperatorRetentionResult {
495    Applied { receipt_id: String },
496    Replayed { receipt_id: String },
497}
498
499#[derive(Debug, Clone, PartialEq, Eq)]
500pub enum OperatorRetentionError {
501    NotFound,
502    StaleState,
503    NonceReplayed,
504    InvalidAction,
505    InvalidState,
506    InvalidTransition,
507    Referenced,
508    ReferencedBySubgraph,
509    Tombstoned,
510    Persistence,
511}
512
513impl Clone for PersistentStore {
514    fn clone(&self) -> Self {
515        Self {
516            conn: self.conn.clone(),
517            integrity_key: self.integrity_key.clone(),
518            #[cfg(test)]
519            terminal_projection_fault: self.terminal_projection_fault.clone(),
520            #[cfg(test)]
521            checkpoint_persistence_fault: self.checkpoint_persistence_fault.clone(),
522            #[cfg(test)]
523            graph_delete_fault: self.graph_delete_fault.clone(),
524        }
525    }
526}
527
528impl PersistentStore {
529    /// Open (or create) the SQLite database at `{data_dir}/agent-graph.db`.
530    pub fn open(data_dir: &Path) -> Result<Self, String> {
531        Self::open_with_integrity_key(data_dir, None)
532    }
533
534    /// Open (or create) the SQLite database with an explicit integrity key path.
535    /// If `integrity_key_path` is None, falls back to the
536    /// `AGENT_GRAPH_INTEGRITY_KEY_PATH` environment variable.
537    pub fn open_with_integrity_key(
538        data_dir: &Path,
539        integrity_key_path: Option<&Path>,
540    ) -> Result<Self, String> {
541        crate::fs_security::validate_data_store(data_dir, integrity_key_path)
542            .map_err(|e| format!("filesystem security check failed: {e}"))?;
543        std::fs::create_dir_all(data_dir).map_err(|e| format!("failed to create data dir: {e}"))?;
544        let db_path = data_dir.join("agent-graph.db");
545        let conn =
546            Connection::open(&db_path).map_err(|e| format!("failed to open database: {e}"))?;
547
548        conn.execute_batch(
549            "PRAGMA journal_mode = WAL;
550             PRAGMA foreign_keys = ON;
551             PRAGMA busy_timeout = 5000;",
552        )
553        .map_err(|e| format!("pragma error: {e}"))?;
554        crate::fs_security::validate_data_store(data_dir, integrity_key_path)
555            .map_err(|e| format!("filesystem security check failed after WAL init: {e}"))?;
556
557        let store = Self {
558            conn: std::sync::Arc::new(Mutex::new(conn)),
559            integrity_key: Self::load_integrity_key_from(integrity_key_path),
560            #[cfg(test)]
561            terminal_projection_fault: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(
562                false,
563            )),
564            #[cfg(test)]
565            checkpoint_persistence_fault: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(
566                false,
567            )),
568            #[cfg(test)]
569            graph_delete_fault: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
570        };
571        store.migrate()?;
572        Ok(store)
573    }
574
575    #[allow(dead_code)]
576    fn load_integrity_key() -> Option<std::sync::Arc<[u8]>> {
577        Self::load_integrity_key_from(None)
578    }
579
580    /// Load the integrity key from an explicit path, or fall back to the
581    /// `AGENT_GRAPH_INTEGRITY_KEY_PATH` environment variable.
582    fn load_integrity_key_from(explicit: Option<&Path>) -> Option<std::sync::Arc<[u8]>> {
583        let path = if let Some(p) = explicit {
584            std::ffi::OsStr::new(p).to_owned()
585        } else {
586            std::env::var_os("AGENT_GRAPH_INTEGRITY_KEY_PATH")?
587        };
588        let key = std::fs::read(path).ok()?;
589        (key.len() >= 32).then(|| std::sync::Arc::from(key))
590    }
591
592    fn require_integrity_key(&self) -> Result<&[u8], ()> {
593        self.integrity_key.as_deref().ok_or(())
594    }
595
596    pub fn has_integrity_key(&self) -> bool {
597        self.integrity_key.is_some()
598    }
599
600    fn migrate(&self) -> Result<(), String> {
601        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
602        conn.execute_batch(
603            "CREATE TABLE IF NOT EXISTS graphs (
604                name TEXT PRIMARY KEY,
605                spec_json TEXT NOT NULL,
606                spec_version TEXT NOT NULL DEFAULT '2',
607                topology_hash TEXT NOT NULL,
608                created_at TEXT NOT NULL DEFAULT (datetime('now')),
609                updated_at TEXT NOT NULL DEFAULT (datetime('now'))
610            );
611
612            CREATE TABLE IF NOT EXISTS graph_retention (
613                graph_name TEXT PRIMARY KEY,
614                state TEXT NOT NULL CHECK(state IN ('active', 'active_phantom_contaminated', 'pinned', 'archived', 'expired_pending_review', 'clear_approved', 'lineage_cleared', 'purged', 'delete_candidate', 'delete_approved')),
615                reason TEXT NOT NULL,
616                actor TEXT NOT NULL,
617                review_after TEXT,
618                created_at TEXT NOT NULL DEFAULT (datetime('now')),
619                updated_at TEXT NOT NULL DEFAULT (datetime('now'))
620            );
621
622            CREATE TABLE IF NOT EXISTS graph_tombstones (
623                graph_name TEXT PRIMARY KEY,
624                topology_hash TEXT NOT NULL,
625                reason TEXT NOT NULL,
626                actor TEXT NOT NULL,
627                deleted_at TEXT NOT NULL DEFAULT (datetime('now'))
628            );
629
630            CREATE TABLE IF NOT EXISTS operator_receipts (
631                receipt_id TEXT PRIMARY KEY,
632                request_digest TEXT NOT NULL,
633                action TEXT NOT NULL,
634                resource_kind TEXT NOT NULL,
635                resource_id TEXT NOT NULL,
636                state_digest TEXT NOT NULL,
637                operator_uid INTEGER NOT NULL,
638                daemon_instance_id TEXT NOT NULL,
639                nonce TEXT NOT NULL UNIQUE,
640                issued_at TEXT NOT NULL,
641                expires_at TEXT NOT NULL,
642                consumed_at TEXT NOT NULL DEFAULT (datetime('now'))
643            );
644            CREATE INDEX IF NOT EXISTS idx_operator_receipts_nonce ON operator_receipts(nonce);
645
646            CREATE TABLE IF NOT EXISTS graph_versions (
647                graph_name TEXT NOT NULL,
648                topology_hash TEXT NOT NULL,
649                spec_json TEXT NOT NULL,
650                created_at TEXT NOT NULL DEFAULT (datetime('now')),
651                PRIMARY KEY (graph_name, topology_hash)
652            );
653
654            CREATE TABLE IF NOT EXISTS executions (
655                run_id TEXT PRIMARY KEY,
656                graph_name TEXT NOT NULL,
657                graph_hash TEXT NOT NULL,
658                thread_id TEXT,
659                status TEXT NOT NULL,
660                input_json TEXT,
661                budgets_json TEXT,
662                final_state_json TEXT,
663                started_at TEXT NOT NULL,
664                finished_at TEXT,
665                total_nodes INTEGER DEFAULT 0,
666                failed_attempts INTEGER DEFAULT 0,
667                idempotency_key TEXT UNIQUE,
668                FOREIGN KEY (graph_name) REFERENCES graphs(name)
669            );
670
671            CREATE TABLE IF NOT EXISTS checkpoints (
672                run_id TEXT NOT NULL,
673                node_id TEXT NOT NULL,
674                attempt INTEGER NOT NULL,
675                input_json TEXT,
676                output_json TEXT,
677                status TEXT NOT NULL,
678                error TEXT,
679                recorded_at TEXT NOT NULL DEFAULT (datetime('now')),
680                checkpoint_id TEXT,
681                graph_id TEXT,
682                graph_version TEXT,
683                next_cursor TEXT,
684                state_json TEXT,
685                state_digest TEXT,
686                budgets_json TEXT,
687                budget_counters_json TEXT,
688                dependency_json TEXT,
689                dependency_digest TEXT,
690                terminal_cursor INTEGER,
691                event_cursor INTEGER,
692                checkpoint_digest TEXT,
693                created_at TEXT,
694                consumed_at TEXT,
695                PRIMARY KEY (run_id, node_id, attempt),
696                FOREIGN KEY (run_id) REFERENCES executions(run_id)
697            );
698
699            CREATE TABLE IF NOT EXISTS events (
700                run_id TEXT NOT NULL,
701                seq INTEGER NOT NULL,
702                event_type TEXT NOT NULL,
703                event_json TEXT NOT NULL,
704                emitted_at TEXT NOT NULL DEFAULT (datetime('now')),
705                PRIMARY KEY (run_id, seq),
706                FOREIGN KEY (run_id) REFERENCES executions(run_id)
707            );
708
709            CREATE TABLE IF NOT EXISTS terminal_receipts (
710                run_id TEXT PRIMARY KEY,
711                receipt_json TEXT NOT NULL,
712                bundle_json TEXT NOT NULL,
713                receipt_digest TEXT NOT NULL,
714                persisted_at TEXT NOT NULL DEFAULT (datetime('now')),
715                FOREIGN KEY (run_id) REFERENCES executions(run_id)
716            );
717
718            CREATE TABLE IF NOT EXISTS template_candidates (
719                template_id TEXT PRIMARY KEY, spec_digest TEXT NOT NULL,
720                graph_id TEXT NOT NULL, graph_version TEXT NOT NULL,
721                source_ref TEXT NOT NULL, state TEXT NOT NULL,
722                created_at TEXT NOT NULL DEFAULT (datetime('now')),
723                updated_at TEXT NOT NULL DEFAULT (datetime('now'))
724            );
725            CREATE TABLE IF NOT EXISTS template_outcome_links (
726                template_id TEXT NOT NULL, run_id TEXT NOT NULL,
727                terminal_receipt_id TEXT NOT NULL, receipt_digest TEXT NOT NULL,
728                disposition TEXT NOT NULL, evidence_digest TEXT NOT NULL,
729                recorded_at TEXT NOT NULL DEFAULT (datetime('now')),
730                PRIMARY KEY (template_id, run_id),
731                FOREIGN KEY (template_id) REFERENCES template_candidates(template_id),
732                FOREIGN KEY (run_id) REFERENCES executions(run_id),
733                FOREIGN KEY (run_id) REFERENCES terminal_receipts(run_id)
734            );
735            CREATE TABLE IF NOT EXISTS template_promotion_decisions (
736                template_id TEXT NOT NULL, from_state TEXT NOT NULL,
737                to_state TEXT NOT NULL, evidence_set_digest TEXT NOT NULL,
738                operator_receipt_id TEXT NOT NULL, decision_digest TEXT NOT NULL,
739                decided_at TEXT NOT NULL DEFAULT (datetime('now')),
740                FOREIGN KEY (template_id) REFERENCES template_candidates(template_id)
741            );
742
743            CREATE TABLE IF NOT EXISTS approval_requests (
744                approval_id TEXT PRIMARY KEY,
745                run_id TEXT NOT NULL,
746                node_id TEXT NOT NULL,
747                checkpoint_id TEXT,
748                graph_id TEXT,
749                graph_version TEXT,
750                checkpoint_digest TEXT,
751                audience TEXT NOT NULL,
752                prompt TEXT NOT NULL,
753                prompt_digest TEXT,
754                allowed_decisions TEXT NOT NULL,
755                approval_digest TEXT,
756                status TEXT NOT NULL DEFAULT 'pending',
757                decision TEXT,
758                decided_by TEXT,
759                decided_at TEXT,
760                expires_at TEXT NOT NULL,
761                created_at TEXT NOT NULL DEFAULT (datetime('now')),
762                FOREIGN KEY (run_id) REFERENCES executions(run_id)
763            );
764
765            CREATE TABLE IF NOT EXISTS idempotency_keys (
766                key TEXT PRIMARY KEY,
767                request_digest TEXT NOT NULL,
768                result_json TEXT NOT NULL,
769                valid INTEGER NOT NULL DEFAULT 1,
770                created_at TEXT NOT NULL DEFAULT (datetime('now'))
771            );
772
773            CREATE TABLE IF NOT EXISTS source_witnesses (
774                witness_id TEXT PRIMARY KEY,
775                locator TEXT NOT NULL,
776                content TEXT NOT NULL,
777                media_type TEXT NOT NULL,
778                authority_class TEXT NOT NULL,
779                retrieved_at TEXT NOT NULL,
780                digest TEXT NOT NULL UNIQUE,
781                created_at TEXT NOT NULL DEFAULT (datetime('now'))
782            );",
783        )
784        .map_err(|e| format!("migration error: {e}"))?;
785        let has_request_digest = conn
786            .prepare("PRAGMA table_info(idempotency_keys)")
787            .map_err(|e| format!("idempotency schema inspection error: {e}"))?
788            .query_map([], |row| row.get::<_, String>(1))
789            .map_err(|e| format!("idempotency schema inspection error: {e}"))?
790            .filter_map(Result::ok)
791            .any(|column| column == "request_digest");
792        if !has_request_digest {
793            conn.execute(
794                "ALTER TABLE idempotency_keys ADD COLUMN request_digest TEXT",
795                [],
796            )
797            .map_err(|e| format!("idempotency migration error: {e}"))?;
798        }
799        let has_idempotency_valid = conn
800            .prepare("PRAGMA table_info(idempotency_keys)")
801            .map_err(|e| format!("idempotency schema inspection error: {e}"))?
802            .query_map([], |row| row.get::<_, String>(1))
803            .map_err(|e| format!("idempotency schema inspection error: {e}"))?
804            .filter_map(Result::ok)
805            .any(|column| column == "valid");
806        if !has_idempotency_valid {
807            conn.execute(
808                "ALTER TABLE idempotency_keys ADD COLUMN valid INTEGER NOT NULL DEFAULT 1",
809                [],
810            )
811            .map_err(|e| format!("idempotency migration error: {e}"))?;
812        }
813        conn.execute(
814            "UPDATE idempotency_keys SET valid = 0 WHERE request_digest IS NULL",
815            [],
816        )
817        .map_err(|e| format!("idempotency quarantine error: {e}"))?;
818        for (table, column, definition) in [
819            ("executions", "budgets_json", "TEXT"),
820            ("checkpoints", "checkpoint_id", "TEXT"),
821            ("checkpoints", "graph_id", "TEXT"),
822            ("checkpoints", "graph_version", "TEXT"),
823            ("checkpoints", "next_cursor", "TEXT"),
824            ("checkpoints", "state_json", "TEXT"),
825            ("checkpoints", "state_digest", "TEXT"),
826            ("checkpoints", "budgets_json", "TEXT"),
827            ("checkpoints", "budget_counters_json", "TEXT"),
828            ("checkpoints", "dependency_json", "TEXT"),
829            ("checkpoints", "dependency_digest", "TEXT"),
830            ("checkpoints", "terminal_cursor", "INTEGER"),
831            ("checkpoints", "event_cursor", "INTEGER"),
832            ("checkpoints", "checkpoint_digest", "TEXT"),
833            ("checkpoints", "created_at", "TEXT"),
834            ("checkpoints", "consumed_at", "TEXT"),
835            ("approval_requests", "checkpoint_id", "TEXT"),
836            ("approval_requests", "graph_id", "TEXT"),
837            ("approval_requests", "graph_version", "TEXT"),
838            ("approval_requests", "checkpoint_digest", "TEXT"),
839            ("approval_requests", "prompt_digest", "TEXT"),
840            ("approval_requests", "approval_digest", "TEXT"),
841        ] {
842            let has_column = conn
843                .prepare(&format!("PRAGMA table_info({table})"))
844                .map_err(|e| format!("{table} schema inspection error: {e}"))?
845                .query_map([], |row| row.get::<_, String>(1))
846                .map_err(|e| format!("{table} schema inspection error: {e}"))?
847                .filter_map(Result::ok)
848                .any(|name| name == column);
849            if !has_column {
850                conn.execute(
851                    &format!("ALTER TABLE {table} ADD COLUMN {column} {definition}"),
852                    [],
853                )
854                .map_err(|e| format!("{table} migration error: {e}"))?;
855            }
856        }
857        conn.execute(
858            "CREATE UNIQUE INDEX IF NOT EXISTS checkpoints_checkpoint_id_idx
859             ON checkpoints(checkpoint_id) WHERE checkpoint_id IS NOT NULL",
860            [],
861        )
862        .map_err(|e| format!("checkpoint index migration error: {e}"))?;
863        Ok(())
864    }
865
866    // ── Local source witnesses ───────────────────────────────────────
867
868    pub fn capture_witness(&self, capture: WitnessCapture) -> Result<WitnessRecord, WitnessError> {
869        let key = self.require_integrity_key().map_err(|_| {
870            WitnessError::new(
871                "INTEGRITY_KEY_REQUIRED",
872                "an external integrity key is required for durable witness capture",
873            )
874        })?;
875        let expected = validate_witness_capture_with_key(capture, Some(key))
876            .map_err(|error| WitnessError::new(error.code, error.message))?;
877        let conn = self
878            .conn
879            .lock()
880            .map_err(|_| WitnessError::new("WITNESS_STORE_ERROR", "witness SQLite lock failed"))?;
881        conn.execute(
882            "INSERT OR IGNORE INTO source_witnesses
883             (witness_id, locator, content, media_type, authority_class, retrieved_at, digest)
884             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
885            params![
886                expected.witness_id,
887                expected.locator,
888                expected.content,
889                expected.media_type,
890                expected.authority_class,
891                expected.retrieved_at,
892                expected.digest
893            ],
894        )
895        .map_err(|_| WitnessError::new("WITNESS_STORE_ERROR", "witness SQLite write failed"))?;
896        drop(conn);
897        self.get_witness(&expected.witness_id)?.ok_or_else(|| {
898            WitnessError::new(
899                "WITNESS_STORE_ERROR",
900                "witness SQLite write did not produce a row",
901            )
902        })
903    }
904
905    pub fn get_witness(&self, witness_id: &str) -> Result<Option<WitnessRecord>, WitnessError> {
906        let key = self.require_integrity_key().map_err(|_| {
907            WitnessError::new(
908                "INTEGRITY_KEY_REQUIRED",
909                "an external integrity key is required for durable witness reads",
910            )
911        })?;
912        let conn = self
913            .conn
914            .lock()
915            .map_err(|_| WitnessError::new("WITNESS_STORE_ERROR", "witness SQLite lock failed"))?;
916        let row = conn.query_row(
917            "SELECT witness_id, locator, content, media_type, authority_class, retrieved_at, digest
918             FROM source_witnesses WHERE witness_id = ?1",
919            params![witness_id],
920            |row| {
921                Ok(WitnessRecord {
922                    witness_id: row.get(0)?,
923                    locator: row.get(1)?,
924                    content: row.get(2)?,
925                    media_type: row.get(3)?,
926                    authority_class: row.get(4)?,
927                    retrieved_at: row.get(5)?,
928                    digest: row.get(6)?,
929                })
930            },
931        );
932        match row {
933            Ok(record) => {
934                verify_witness_record_with_key(&record, Some(key))?;
935                Ok(Some(record))
936            }
937            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
938            Err(rusqlite::Error::FromSqlConversionFailure(..))
939            | Err(rusqlite::Error::InvalidColumnType(..))
940            | Err(rusqlite::Error::InvalidColumnName(_)) => Err(WitnessError::new(
941                "WITNESS_INTEGRITY_FAILURE",
942                "stored witness integrity validation failed",
943            )),
944            Err(_) => Err(WitnessError::new(
945                "WITNESS_STORE_ERROR",
946                "witness SQLite read failed",
947            )),
948        }
949    }
950
951    // ── Graphs ──────────────────────────────────────────────────────────
952
953    pub fn save_graph(
954        &self,
955        name: &str,
956        spec_json: &str,
957        topology_hash: &str,
958        overwrite: bool,
959    ) -> Result<(), String> {
960        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
961        let tombstoned: bool = conn
962            .query_row(
963                "SELECT EXISTS(SELECT 1 FROM graph_tombstones WHERE graph_name = ?1)",
964                params![name],
965                |row| row.get(0),
966            )
967            .map_err(|e| format!("graph tombstone check error: {e}"))?;
968        if tombstoned {
969            return Err(format!(
970                "graph '{name}' has a retention tombstone; choose a new graph ID"
971            ));
972        }
973        let current: Option<String> = conn
974            .query_row(
975                "SELECT topology_hash FROM graphs WHERE name = ?1",
976                params![name],
977                |row| row.get(0),
978            )
979            .ok();
980        if let Some(current) = current
981            .as_deref()
982            .filter(|current| *current != topology_hash)
983        {
984            let referenced: bool = conn
985                .query_row(
986                    "SELECT EXISTS(SELECT 1 FROM executions WHERE graph_name = ?1 AND graph_hash = ?2)",
987                    params![name, current],
988                    |row| row.get(0),
989                )
990                .map_err(|e| format!("graph version reference check error: {e}"))?;
991            if referenced && !overwrite {
992                return Err(format!("graph '{name}' current version {current} is referenced by a durable execution; explicit overwrite is required"));
993            }
994        }
995        conn.execute(
996            "INSERT OR IGNORE INTO graph_versions (graph_name, topology_hash, spec_json) VALUES (?1, ?2, ?3)",
997            params![name, topology_hash, spec_json],
998        )
999        .map_err(|e| format!("save graph version error: {e}"))?;
1000        conn.execute(
1001            "INSERT INTO graphs (name, spec_json, topology_hash)
1002             VALUES (?1, ?2, ?3)
1003             ON CONFLICT(name) DO UPDATE SET
1004                spec_json = excluded.spec_json,
1005                topology_hash = excluded.topology_hash,
1006                updated_at = datetime('now')",
1007            params![name, spec_json, topology_hash],
1008        )
1009        .map_err(|e| format!("save_graph error: {e}"))?;
1010        conn.execute(
1011            "INSERT INTO graph_retention (graph_name, state, reason, actor)
1012             VALUES (?1, 'active', 'created', 'system')
1013             ON CONFLICT(graph_name) DO NOTHING",
1014            params![name],
1015        )
1016        .map_err(|e| format!("save graph retention state error: {e}"))?;
1017        Ok(())
1018    }
1019
1020    pub fn load_graph(&self, name: &str) -> Result<Option<(String, String)>, String> {
1021        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1022        let mut stmt = conn
1023            .prepare("SELECT spec_json, topology_hash FROM graphs WHERE name = ?1")
1024            .map_err(|e| format!("load_graph error: {e}"))?;
1025        let result = stmt
1026            .query_row(params![name], |row| {
1027                Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
1028            })
1029            .ok();
1030        Ok(result)
1031    }
1032
1033    pub fn list_graphs(&self) -> Result<Vec<(String, String, String)>, String> {
1034        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1035        let mut stmt = conn
1036            .prepare("SELECT g.name, g.topology_hash, g.created_at FROM graphs AS g
1037                      WHERE NOT EXISTS (SELECT 1 FROM graph_tombstones AS t WHERE t.graph_name = g.name)
1038                      ORDER BY g.name")
1039            .map_err(|e| format!("list_graphs error: {e}"))?;
1040        let rows = stmt
1041            .query_map([], |row| {
1042                Ok((
1043                    row.get::<_, String>(0)?,
1044                    row.get::<_, String>(1)?,
1045                    row.get::<_, String>(2)?,
1046                ))
1047            })
1048            .map_err(|e| format!("list_graphs error: {e}"))?
1049            .filter_map(|r| r.ok())
1050            .collect();
1051        Ok(rows)
1052    }
1053
1054    pub fn graph_is_tombstoned(&self, graph_id: &str) -> Result<bool, String> {
1055        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1056        conn.query_row(
1057            "SELECT EXISTS(SELECT 1 FROM graph_tombstones WHERE graph_name = ?1)",
1058            params![graph_id],
1059            |row| row.get(0),
1060        )
1061        .map_err(|e| format!("graph tombstone query error: {e}"))
1062    }
1063
1064    fn inbound_subgraph_refs(conn: &Connection, name: &str) -> Result<Vec<String>, String> {
1065        let mut stmt = conn
1066            .prepare("SELECT name, spec_json FROM graphs WHERE name != ?1")
1067            .map_err(|e| format!("subgraph reference query error: {e}"))?;
1068        let rows = stmt
1069            .query_map(params![name], |row| {
1070                Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
1071            })
1072            .map_err(|e| format!("subgraph reference query error: {e}"))?;
1073        let mut refs = Vec::new();
1074        for row in rows {
1075            let (graph_name, spec_json) =
1076                row.map_err(|e| format!("subgraph reference row error: {e}"))?;
1077            let Ok(spec) = serde_json::from_str::<Value>(&spec_json) else {
1078                continue;
1079            };
1080            let references_target =
1081                spec.get("nodes")
1082                    .and_then(Value::as_array)
1083                    .is_some_and(|nodes| {
1084                        nodes.iter().any(|node| {
1085                            node.get("type").and_then(Value::as_str) == Some("subgraph")
1086                                && node
1087                                    .get("config")
1088                                    .and_then(|config| config.get("graph_name"))
1089                                    .and_then(Value::as_str)
1090                                    == Some(name)
1091                        })
1092                    });
1093            if references_target {
1094                refs.push(graph_name);
1095            }
1096        }
1097        refs.sort();
1098        Ok(refs)
1099    }
1100
1101    pub fn graph_retention_review(
1102        &self,
1103        graph_id: Option<&str>,
1104        state: Option<&str>,
1105        limit: usize,
1106    ) -> Result<Vec<GraphRetentionReport>, String> {
1107        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1108        let mut stmt = conn
1109            .prepare(
1110                "SELECT g.name, g.created_at, g.updated_at,
1111                        r.state, r.reason, r.actor, r.review_after,
1112                        COUNT(DISTINCT gv.topology_hash), COUNT(DISTINCT e.run_id),
1113                        MAX(COALESCE(e.finished_at, e.started_at))
1114                 FROM graphs AS g
1115                 LEFT JOIN graph_retention AS r ON r.graph_name = g.name
1116                 LEFT JOIN graph_versions AS gv ON gv.graph_name = g.name
1117                 LEFT JOIN executions AS e ON e.graph_name = g.name
1118                 WHERE (?1 IS NULL OR g.name = ?1)
1119                   AND (?2 IS NULL OR COALESCE(r.state, 'unclassified') = ?2)
1120                 GROUP BY g.name
1121                 ORDER BY g.updated_at ASC, g.name ASC
1122                 LIMIT ?3",
1123            )
1124            .map_err(|e| format!("graph retention review query error: {e}"))?;
1125        let rows = stmt
1126            .query_map(params![graph_id, state, limit as i64], |row| {
1127                Ok((
1128                    row.get::<_, String>(0)?,
1129                    row.get::<_, String>(1)?,
1130                    row.get::<_, String>(2)?,
1131                    row.get::<_, Option<String>>(3)?,
1132                    row.get::<_, Option<String>>(4)?,
1133                    row.get::<_, Option<String>>(5)?,
1134                    row.get::<_, Option<String>>(6)?,
1135                    row.get::<_, i64>(7)?,
1136                    row.get::<_, i64>(8)?,
1137                    row.get::<_, Option<String>>(9)?,
1138                ))
1139            })
1140            .map_err(|e| format!("graph retention review query error: {e}"))?;
1141        let mut reports = Vec::new();
1142        for row in rows {
1143            let (
1144                graph_id,
1145                created_at,
1146                updated_at,
1147                state,
1148                reason,
1149                actor,
1150                review_after,
1151                versions,
1152                executions,
1153                last_execution_at,
1154            ) = row.map_err(|e| format!("graph retention review row error: {e}"))?;
1155            let inbound_subgraph_refs = Self::inbound_subgraph_refs(&conn, &graph_id)?;
1156            let tombstoned: bool = conn
1157                .query_row(
1158                    "SELECT EXISTS(SELECT 1 FROM graph_tombstones WHERE graph_name = ?1)",
1159                    params![&graph_id],
1160                    |row| row.get(0),
1161                )
1162                .map_err(|e| format!("graph retention tombstone query error: {e}"))?;
1163            let state = state.unwrap_or_else(|| "unclassified".into());
1164            let state_digest = digest(&serde_json::json!({
1165                "graph_id": graph_id,
1166                "state": state,
1167                "reason": reason,
1168                "actor": actor,
1169                "review_after": review_after,
1170                "topology_hashes": versions,
1171                "execution_count": executions,
1172                "last_execution_at": last_execution_at,
1173                "inbound_subgraph_refs": inbound_subgraph_refs,
1174                "tombstoned": tombstoned,
1175            }));
1176            reports.push(GraphRetentionReport {
1177                graph_id,
1178                state,
1179                reason,
1180                actor,
1181                review_after,
1182                created_at,
1183                updated_at,
1184                version_count: versions as u64,
1185                execution_count: executions as u64,
1186                last_execution_at,
1187                deletion_eligible: executions == 0 && inbound_subgraph_refs.is_empty(),
1188                inbound_subgraph_refs,
1189                state_digest,
1190                tombstoned,
1191            });
1192        }
1193        Ok(reports)
1194    }
1195
1196    /// Apply one OS-authenticated retention mutation and its receipt atomically.
1197    /// The caller must have already validated peer credentials and the time window.
1198    pub fn apply_operator_retention(
1199        &self,
1200        request: &OperatorRetentionRequest,
1201    ) -> Result<OperatorRetentionResult, OperatorRetentionError> {
1202        let mut conn = self
1203            .conn
1204            .lock()
1205            .map_err(|_| OperatorRetentionError::Persistence)?;
1206
1207        let prior: Option<(String, String)> = conn
1208            .query_row(
1209                "SELECT request_digest, receipt_id FROM operator_receipts WHERE nonce = ?1",
1210                params![request.nonce],
1211                |row| Ok((row.get(0)?, row.get(1)?)),
1212            )
1213            .optional()
1214            .map_err(|_| OperatorRetentionError::Persistence)?;
1215        if let Some((digest, receipt_id)) = prior {
1216            if digest == request.request_digest {
1217                return Ok(OperatorRetentionResult::Replayed { receipt_id });
1218            }
1219            return Err(OperatorRetentionError::NonceReplayed);
1220        }
1221
1222        let snapshot: Option<(
1223            String,
1224            Option<String>,
1225            Option<String>,
1226            Option<String>,
1227            i64,
1228            i64,
1229            Option<String>,
1230        )> = conn
1231            .query_row(
1232                "SELECT COALESCE(r.state, 'unclassified'), r.reason, r.actor, r.review_after,
1233                        COUNT(DISTINCT gv.topology_hash), COUNT(DISTINCT e.run_id),
1234                        MAX(COALESCE(e.finished_at, e.started_at))
1235                 FROM graphs AS g
1236                 LEFT JOIN graph_retention AS r ON r.graph_name = g.name
1237                 LEFT JOIN graph_versions AS gv ON gv.graph_name = g.name
1238                 LEFT JOIN executions AS e ON e.graph_name = g.name
1239                 WHERE g.name = ?1
1240                 GROUP BY g.name",
1241                params![request.graph_id],
1242                |row| {
1243                    Ok((
1244                        row.get(0)?,
1245                        row.get(1)?,
1246                        row.get(2)?,
1247                        row.get(3)?,
1248                        row.get(4)?,
1249                        row.get(5)?,
1250                        row.get(6)?,
1251                    ))
1252                },
1253            )
1254            .optional()
1255            .map_err(|_| OperatorRetentionError::Persistence)?;
1256        let Some((
1257            current_state,
1258            current_reason,
1259            current_actor,
1260            current_review_after,
1261            versions,
1262            executions,
1263            last_execution_at,
1264        )) = snapshot
1265        else {
1266            return Err(OperatorRetentionError::NotFound);
1267        };
1268        let inbound = Self::inbound_subgraph_refs(&conn, &request.graph_id)
1269            .map_err(|_| OperatorRetentionError::Persistence)?;
1270        let tombstoned: bool = conn
1271            .query_row(
1272                "SELECT EXISTS(SELECT 1 FROM graph_tombstones WHERE graph_name = ?1)",
1273                params![request.graph_id],
1274                |row| row.get(0),
1275            )
1276            .map_err(|_| OperatorRetentionError::Persistence)?;
1277        if tombstoned {
1278            return Err(OperatorRetentionError::Tombstoned);
1279        }
1280        let current_digest = digest(&serde_json::json!({
1281            "graph_id": request.graph_id,
1282            "state": current_state,
1283            "reason": current_reason,
1284            "actor": current_actor,
1285            "review_after": current_review_after,
1286            "topology_hashes": versions,
1287            "execution_count": executions,
1288            "last_execution_at": last_execution_at.map_or(serde_json::Value::Null, serde_json::Value::String),
1289            "inbound_subgraph_refs": inbound,
1290            "tombstoned": tombstoned,
1291        }));
1292        if current_digest != request.expected_state_digest {
1293            return Err(OperatorRetentionError::StaleState);
1294        }
1295        if executions > 0 {
1296            return Err(OperatorRetentionError::Referenced);
1297        }
1298        if !inbound.is_empty() {
1299            return Err(OperatorRetentionError::ReferencedBySubgraph);
1300        }
1301
1302        let action = serde_json::to_string(&request.action)
1303            .map_err(|_| OperatorRetentionError::InvalidAction)?;
1304        let actor = format!("uid:{}", request.operator_uid);
1305        let tx = conn
1306            .transaction()
1307            .map_err(|_| OperatorRetentionError::Persistence)?;
1308        match request.action {
1309            OperatorAction::SetGraphRetention => {
1310                let state = request
1311                    .state
1312                    .as_deref()
1313                    .ok_or(OperatorRetentionError::InvalidState)?;
1314                if !matches!(state, "active" | "active_phantom_contaminated" | "pinned" | "archived" | "expired_pending_review" | "clear_approved" | "lineage_cleared" | "purged" | "delete_candidate") {
1315                    return Err(OperatorRetentionError::InvalidState);
1316                }
1317                tx.execute(
1318                    "INSERT INTO graph_retention (graph_name, state, reason, actor, review_after)
1319                     VALUES (?1, ?2, ?3, ?4, ?5)
1320                     ON CONFLICT(graph_name) DO UPDATE SET state=excluded.state,
1321                       reason=excluded.reason, actor=excluded.actor,
1322                       review_after=excluded.review_after, updated_at=datetime('now')",
1323                    params![
1324                        request.graph_id,
1325                        state,
1326                        request.reason.as_deref().unwrap_or("operator decision"),
1327                        actor,
1328                        request.review_after
1329                    ],
1330                )
1331                .map_err(|_| OperatorRetentionError::Persistence)?;
1332            }
1333            OperatorAction::ApproveGraphDeletion => {
1334                if current_state != "delete_candidate" {
1335                    return Err(OperatorRetentionError::InvalidTransition);
1336                }
1337                tx.execute(
1338                    "UPDATE graph_retention SET state='delete_approved', updated_at=datetime('now'), actor=?2, reason=?3 WHERE graph_name=?1",
1339                    params![request.graph_id, actor, request.reason.as_deref().unwrap_or("operator approval")],
1340                )
1341                .map_err(|_| OperatorRetentionError::Persistence)?;
1342            }
1343            OperatorAction::DeleteGraph => {
1344                if current_state != "delete_approved" {
1345                    return Err(OperatorRetentionError::InvalidTransition);
1346                }
1347                let topology_hash: String = tx
1348                    .query_row(
1349                        "SELECT topology_hash FROM graphs WHERE name=?1",
1350                        params![request.graph_id],
1351                        |row| row.get(0),
1352                    )
1353                    .map_err(|_| OperatorRetentionError::NotFound)?;
1354                tx.execute(
1355                    "INSERT INTO graph_tombstones (graph_name, topology_hash, reason, actor)
1356                     VALUES (?1, ?2, ?3, ?4)",
1357                    params![
1358                        request.graph_id,
1359                        topology_hash,
1360                        request.reason.as_deref().unwrap_or("operator deletion"),
1361                        actor
1362                    ],
1363                )
1364                .map_err(|_| OperatorRetentionError::Persistence)?;
1365            }
1366            OperatorAction::ClearExecutionLineage => {
1367                return Err(OperatorRetentionError::InvalidState);
1368                // TODO: implement per spec §4.1 — transactional DELETE of execution data
1369            }
1370            OperatorAction::PurgeGraph => {
1371                return Err(OperatorRetentionError::InvalidTransition);
1372                // TODO: implement per spec §4.2 — tombstone + graph deletion
1373            }
1374            OperatorAction::PromoteTemplate => {
1375                return Err(OperatorRetentionError::InvalidAction);
1376                // TODO: wire template promotion
1377            }
1378            OperatorAction::Migrate => {
1379                return Err(OperatorRetentionError::InvalidAction);
1380                // TODO: wire graph state migration
1381            }
1382            OperatorAction::Install => {
1383                // TODO: wire daemon installation — currently no-op for idempotency
1384            }
1385            _ => return Err(OperatorRetentionError::InvalidAction),
1386        }
1387
1388        let receipt_id = format!(
1389            "operator-{}",
1390            request.request_digest.trim_start_matches("sha256:")
1391        );
1392        tx.execute(
1393            "INSERT INTO operator_receipts
1394             (receipt_id, request_digest, action, resource_kind, resource_id, state_digest,
1395              operator_uid, daemon_instance_id, nonce, issued_at, expires_at)
1396             VALUES (?1, ?2, ?3, 'graph', ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
1397            params![
1398                receipt_id,
1399                request.request_digest,
1400                action,
1401                request.graph_id,
1402                request.expected_state_digest,
1403                request.operator_uid,
1404                request.daemon_instance_id,
1405                request.nonce,
1406                request.issued_at,
1407                request.expires_at
1408            ],
1409        )
1410        .map_err(|_| OperatorRetentionError::Persistence)?;
1411        tx.commit()
1412            .map_err(|_| OperatorRetentionError::Persistence)?;
1413        Ok(OperatorRetentionResult::Applied { receipt_id })
1414    }
1415
1416    pub fn set_graph_retention(
1417        &self,
1418        graph_id: &str,
1419        state: &str,
1420        reason: &str,
1421        actor: &str,
1422        review_after: Option<&str>,
1423    ) -> Result<GraphRetentionRecord, GraphRetentionError> {
1424        if !matches!(
1425            state,
1426            "active" | "pinned" | "archived" | "delete_candidate" | "delete_approved"
1427        ) {
1428            return Err(GraphRetentionError::InvalidState);
1429        }
1430        if reason.trim().is_empty() || actor.trim().is_empty() {
1431            return Err(GraphRetentionError::InvalidState);
1432        }
1433        let reports = self
1434            .graph_retention_review(Some(graph_id), None, 1)
1435            .map_err(|_| GraphRetentionError::NotFound)?;
1436        let Some(report) = reports.into_iter().next() else {
1437            return Err(GraphRetentionError::NotFound);
1438        };
1439        if matches!(state, "delete_candidate" | "delete_approved") {
1440            if report.execution_count > 0 {
1441                return Err(GraphRetentionError::Referenced);
1442            }
1443            if !report.inbound_subgraph_refs.is_empty() {
1444                return Err(GraphRetentionError::ReferencedBySubgraph);
1445            }
1446        }
1447        if state == "delete_approved" && report.state != "delete_candidate" {
1448            return Err(GraphRetentionError::InvalidTransition);
1449        }
1450        let conn = self
1451            .conn
1452            .lock()
1453            .map_err(|_| GraphRetentionError::NotFound)?;
1454        conn.execute(
1455            "INSERT INTO graph_retention (graph_name, state, reason, actor, review_after)
1456             VALUES (?1, ?2, ?3, ?4, ?5)
1457             ON CONFLICT(graph_name) DO UPDATE SET
1458                 state = excluded.state,
1459                 reason = excluded.reason,
1460                 actor = excluded.actor,
1461                 review_after = excluded.review_after,
1462                 updated_at = datetime('now')",
1463            params![graph_id, state, reason, actor, review_after],
1464        )
1465        .map_err(|_| GraphRetentionError::NotFound)?;
1466        let updated_at = conn
1467            .query_row(
1468                "SELECT updated_at FROM graph_retention WHERE graph_name = ?1",
1469                params![graph_id],
1470                |row| row.get(0),
1471            )
1472            .map_err(|_| GraphRetentionError::NotFound)?;
1473        Ok(GraphRetentionRecord {
1474            graph_id: graph_id.into(),
1475            state: state.into(),
1476            reason: reason.into(),
1477            actor: actor.into(),
1478            review_after: review_after.map(str::to_owned),
1479            updated_at,
1480        })
1481    }
1482
1483    pub fn graph_execution_allowed(&self, graph_id: &str) -> Result<bool, String> {
1484        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1485        let tombstoned: bool = conn
1486            .query_row(
1487                "SELECT EXISTS(SELECT 1 FROM graph_tombstones WHERE graph_name = ?1)",
1488                params![graph_id],
1489                |row| row.get(0),
1490            )
1491            .map_err(|e| format!("graph tombstone execution check error: {e}"))?;
1492        if tombstoned {
1493            return Ok(false);
1494        }
1495        let state: Option<String> = conn
1496            .query_row(
1497                "SELECT state FROM graph_retention WHERE graph_name = ?1",
1498                params![graph_id],
1499                |row| row.get(0),
1500            )
1501            .optional()
1502            .map_err(|e| format!("graph retention execution check error: {e}"))?;
1503        Ok(matches!(
1504            state.as_deref(),
1505            None | Some("active") | Some("pinned")
1506        ))
1507    }
1508
1509    pub fn delete_graph(&self, name: &str) -> Result<GraphDeleteResult, String> {
1510        let mut conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1511        let graph: Option<String> = conn
1512            .query_row(
1513                "SELECT topology_hash FROM graphs WHERE name = ?1",
1514                params![name],
1515                |row| row.get(0),
1516            )
1517            .optional()
1518            .map_err(|e| format!("delete_graph graph lookup error: {e}"))?;
1519        let Some(topology_hash) = graph else {
1520            return Ok(GraphDeleteResult::NotFound);
1521        };
1522        if !Self::inbound_subgraph_refs(&conn, name)?.is_empty() {
1523            return Ok(GraphDeleteResult::ReferencedBySubgraph);
1524        }
1525        let retention: Option<(String, String, String)> = conn
1526            .query_row(
1527                "SELECT state, reason, actor FROM graph_retention WHERE graph_name = ?1",
1528                params![name],
1529                |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
1530            )
1531            .optional()
1532            .map_err(|e| format!("delete_graph retention lookup error: {e}"))?;
1533        if !matches!(
1534            retention.as_ref().map(|record| record.0.as_str()),
1535            Some("delete_approved")
1536        ) {
1537            return Ok(GraphDeleteResult::RetentionApprovalRequired);
1538        }
1539        let (reason, actor) = retention
1540            .map(|(_, reason, actor)| (reason, actor))
1541            .expect("approved retention state exists");
1542        let tx = conn
1543            .transaction()
1544            .map_err(|e| format!("delete_graph transaction begin error: {e}"))?;
1545        let referenced: bool = tx
1546            .query_row(
1547                "SELECT EXISTS(SELECT 1 FROM executions WHERE graph_name = ?1)",
1548                params![name],
1549                |row| row.get(0),
1550            )
1551            .map_err(|e| format!("delete_graph reference check error: {e}"))?;
1552        if referenced {
1553            return Ok(GraphDeleteResult::Referenced);
1554        }
1555        tx.execute(
1556            "INSERT INTO graph_tombstones (graph_name, topology_hash, reason, actor)
1557             VALUES (?1, ?2, ?3, ?4)",
1558            params![name, topology_hash, reason, actor],
1559        )
1560        .map_err(|e| format!("delete_graph tombstone error: {e}"))?;
1561        let affected: i64 = tx
1562            .query_row(
1563                "SELECT COUNT(*) FROM graphs WHERE name = ?1",
1564                params![name],
1565                |row| row.get(0),
1566            )
1567            .map_err(|e| format!("delete_graph graph presence error: {e}"))?;
1568        if affected == 0 {
1569            return Ok(GraphDeleteResult::NotFound);
1570        }
1571        #[cfg(test)]
1572        if self
1573            .graph_delete_fault
1574            .swap(false, std::sync::atomic::Ordering::SeqCst)
1575        {
1576            return Err("injected graph deletion failure after graph row removal".into());
1577        }
1578        tx.commit()
1579            .map_err(|e| format!("delete_graph transaction commit error: {e}"))?;
1580        Ok(GraphDeleteResult::Deleted)
1581    }
1582
1583    #[cfg(test)]
1584    pub(crate) fn fail_graph_delete_after_graph_row(&self) {
1585        self.graph_delete_fault
1586            .store(true, std::sync::atomic::Ordering::SeqCst);
1587    }
1588
1589    #[cfg(test)]
1590    pub(crate) fn fail_terminal_projection_after_events(&self) {
1591        self.terminal_projection_fault
1592            .store(true, std::sync::atomic::Ordering::SeqCst);
1593    }
1594
1595    pub fn load_graph_version(
1596        &self,
1597        name: &str,
1598        topology_hash: &str,
1599    ) -> Result<Option<String>, String> {
1600        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1601        let result = conn
1602            .query_row(
1603                "SELECT spec_json FROM graph_versions WHERE graph_name = ?1 AND topology_hash = ?2",
1604                params![name, topology_hash],
1605                |row| row.get::<_, String>(0),
1606            )
1607            .ok();
1608        Ok(result)
1609    }
1610
1611    pub fn list_graph_versions(&self, name: &str) -> Result<Vec<String>, String> {
1612        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1613        let mut stmt = conn
1614            .prepare("SELECT topology_hash FROM graph_versions WHERE graph_name = ?1 ORDER BY created_at, topology_hash")
1615            .map_err(|e| format!("list graph versions error: {e}"))?;
1616        let versions = stmt
1617            .query_map(params![name], |row| row.get::<_, String>(0))
1618            .map_err(|e| format!("list graph versions error: {e}"))?
1619            .filter_map(Result::ok)
1620            .collect();
1621        Ok(versions)
1622    }
1623
1624    // ── Executions ──────────────────────────────────────────────────────
1625
1626    pub fn save_execution(
1627        &self,
1628        run_id: &str,
1629        graph_name: &str,
1630        graph_hash: &str,
1631        status: &str,
1632        input_json: &str,
1633    ) -> Result<(), String> {
1634        self.save_execution_with_budgets(run_id, graph_name, graph_hash, status, input_json, None)
1635    }
1636
1637    pub fn save_execution_with_budgets(
1638        &self,
1639        run_id: &str,
1640        graph_name: &str,
1641        graph_hash: &str,
1642        status: &str,
1643        input_json: &str,
1644        budgets_json: Option<&str>,
1645    ) -> Result<(), String> {
1646        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1647        conn.execute(
1648            "INSERT INTO executions (run_id, graph_name, graph_hash, status, input_json, budgets_json, started_at)
1649             VALUES (?1, ?2, ?3, ?4, ?5, ?6, datetime('now'))
1650             ON CONFLICT(run_id) DO UPDATE SET
1651                status = excluded.status,
1652                budgets_json = COALESCE(excluded.budgets_json, executions.budgets_json),
1653                finished_at = CASE WHEN excluded.status IN ('completed','failed','cancelled') THEN datetime('now') ELSE finished_at END",
1654            params![run_id, graph_name, graph_hash, status, input_json, budgets_json],
1655        )
1656        .map_err(|e| format!("save_execution error: {e}"))?;
1657        Ok(())
1658    }
1659
1660    pub fn update_execution_status(
1661        &self,
1662        run_id: &str,
1663        status: &str,
1664        final_state_json: Option<&str>,
1665        total_nodes: Option<usize>,
1666        failed_attempts: Option<usize>,
1667    ) -> Result<(), String> {
1668        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1669        let changed = conn.execute(
1670            "UPDATE executions SET
1671                status = ?2,
1672                final_state_json = COALESCE(?3, final_state_json),
1673                total_nodes = COALESCE(?4, total_nodes),
1674                failed_attempts = COALESCE(?5, failed_attempts),
1675                finished_at = CASE WHEN ?2 IN ('completed','failed','cancelled') THEN datetime('now') ELSE finished_at END
1676             WHERE run_id = ?1",
1677            params![
1678                run_id,
1679                status,
1680                final_state_json,
1681                total_nodes.map(|v| v as i64),
1682                failed_attempts.map(|v| v as i64)
1683            ],
1684        ).map_err(|e| format!("update_execution error: {e}"))?;
1685        if changed != 1 {
1686            return Err(format!(
1687                "update_execution error: run '{run_id}' was not found"
1688            ));
1689        }
1690        Ok(())
1691    }
1692
1693    pub fn persist_terminal_projection(
1694        &self,
1695        run_id: &str,
1696        status: &str,
1697        final_state_json: &str,
1698        total_nodes: usize,
1699        events: &[(u64, String, String)],
1700        receipt_json: &str,
1701        bundle_json: &str,
1702    ) -> Result<String, String> {
1703        let key = self
1704            .require_integrity_key()
1705            .map_err(|_| "INTEGRITY_KEY_REQUIRED".to_owned())?;
1706        let receipt: Value = serde_json::from_str(receipt_json)
1707            .map_err(|e| format!("terminal receipt JSON error: {e}"))?;
1708        let receipt_digest = hmac_sha256(&receipt, key);
1709        let mut conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1710        let tx = conn
1711            .transaction()
1712            .map_err(|e| format!("terminal transaction begin error: {e}"))?;
1713        let changed = tx.execute(
1714            "UPDATE executions SET status = ?2, final_state_json = ?3, total_nodes = ?4,
1715             finished_at = CASE WHEN ?2 IN ('completed','failed','cancelled') THEN datetime('now') ELSE finished_at END
1716             WHERE run_id = ?1",
1717            params![run_id, status, final_state_json, total_nodes as i64],
1718        ).map_err(|e| format!("terminal execution update error: {e}"))?;
1719        if changed != 1 {
1720            return Err(format!("terminal projection run '{run_id}' was not found"));
1721        }
1722        tx.execute("DELETE FROM events WHERE run_id = ?1", params![run_id])
1723            .map_err(|e| format!("terminal event reset error: {e}"))?;
1724        for (seq, event_type, event_json) in events {
1725            tx.execute(
1726                "INSERT INTO events (run_id, seq, event_type, event_json) VALUES (?1, ?2, ?3, ?4)",
1727                params![run_id, *seq, event_type, event_json],
1728            )
1729            .map_err(|e| format!("terminal event insert error: {e}"))?;
1730        }
1731        #[cfg(test)]
1732        if self
1733            .terminal_projection_fault
1734            .swap(false, std::sync::atomic::Ordering::SeqCst)
1735        {
1736            return Err("injected terminal projection failure after events".into());
1737        }
1738        tx.execute(
1739            "INSERT INTO terminal_receipts (run_id, receipt_json, bundle_json, receipt_digest) VALUES (?1, ?2, ?3, ?4)
1740             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')",
1741            params![run_id, receipt_json, bundle_json, receipt_digest],
1742        ).map_err(|e| format!("terminal receipt insert error: {e}"))?;
1743        tx.commit()
1744            .map_err(|e| format!("terminal transaction commit error: {e}"))?;
1745        Ok(receipt_digest)
1746    }
1747
1748    pub fn load_terminal_receipt(&self, run_id: &str) -> Result<Option<Value>, String> {
1749        let key = self
1750            .require_integrity_key()
1751            .map_err(|_| "INTEGRITY_KEY_REQUIRED".to_owned())?;
1752        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1753        let row: Result<(String, String), _> = conn.query_row(
1754            "SELECT receipt_json, receipt_digest FROM terminal_receipts WHERE run_id = ?1",
1755            params![run_id],
1756            |row| Ok((row.get(0)?, row.get(1)?)),
1757        );
1758        match row {
1759            Ok((receipt_json, receipt_digest)) => {
1760                let receipt: Value = serde_json::from_str(&receipt_json)
1761                    .map_err(|e| format!("stored receipt JSON error: {e}"))?;
1762                if hmac_sha256(&receipt, key) != receipt_digest {
1763                    return Err("RECEIPT_INTEGRITY_FAILURE".into());
1764                }
1765                Ok(Some(
1766                    serde_json::json!({"receipt":receipt,"receipt_digest":receipt_digest,"storage_class":"sqlite_terminal_receipt","replay_capability":"integrity_only"}),
1767                ))
1768            }
1769            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1770            Err(e) => Err(format!("load terminal receipt error: {e}")),
1771        }
1772    }
1773
1774    /// A server restart cannot resume an in-flight graph. Make that interruption
1775    /// explicit instead of leaving a permanently misleading `running` row.
1776    pub fn recover_incomplete_executions(&self) -> Result<(), String> {
1777        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1778        conn.execute(
1779            "UPDATE executions
1780             SET status = 'interrupted_non_resumable', finished_at = datetime('now')
1781             WHERE status IN ('accepted', 'running')",
1782            [],
1783        )
1784        .map_err(|e| format!("recover executions error: {e}"))?;
1785        Ok(())
1786    }
1787
1788    /// Return the terminal projection retained by SQLite. This is deliberately
1789    /// not a resumable checkpoint or replay artifact.
1790    pub fn load_execution(&self, run_id: &str) -> Result<Option<Value>, String> {
1791        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1792        let mut stmt = conn
1793            .prepare(
1794                "SELECT graph_name, graph_hash, status, final_state_json, started_at, finished_at
1795                 FROM executions WHERE run_id = ?1",
1796            )
1797            .map_err(|e| format!("load execution error: {e}"))?;
1798        let row = stmt.query_row(params![run_id], |row| {
1799            Ok((
1800                row.get::<_, String>(0)?,
1801                row.get::<_, String>(1)?,
1802                row.get::<_, String>(2)?,
1803                row.get::<_, Option<String>>(3)?,
1804                row.get::<_, String>(4)?,
1805                row.get::<_, Option<String>>(5)?,
1806            ))
1807        });
1808        match row {
1809            Ok((graph_id, graph_version, status, final_state_json, started_at, finished_at)) => {
1810                let final_state = final_state_json
1811                    .as_deref()
1812                    .map(serde_json::from_str)
1813                    .transpose()
1814                    .map_err(|e| format!("stored final state JSON error: {e}"))?
1815                    .unwrap_or(Value::Null);
1816                Ok(Some(serde_json::json!({
1817                    "run_id": run_id,
1818                    "graph_id": graph_id,
1819                    "graph_version": graph_version,
1820                    "status": status,
1821                    "success": status == "completed",
1822                    "final_state": final_state,
1823                    "started_at": started_at,
1824                    "finished_at": finished_at,
1825                    "storage_class": "sqlite_terminal_record",
1826                    "durable_resume": false,
1827                    "replay_capability": "integrity_only"
1828                })))
1829            }
1830            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1831            Err(e) => Err(format!("load execution error: {e}")),
1832        }
1833    }
1834
1835    pub fn load_execution_contract(
1836        &self,
1837        run_id: &str,
1838    ) -> Result<Option<ExecutionContract>, String> {
1839        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
1840        let result = conn.query_row(
1841            "SELECT graph_name, graph_hash, input_json, budgets_json
1842             FROM executions WHERE run_id = ?1",
1843            params![run_id],
1844            |row| {
1845                let input_json: Option<String> = row.get(2)?;
1846                let budgets_json: Option<String> = row.get(3)?;
1847                Ok(ExecutionContract {
1848                    graph_id: row.get(0)?,
1849                    graph_version: row.get(1)?,
1850                    input: input_json
1851                        .as_deref()
1852                        .map(serde_json::from_str)
1853                        .transpose()
1854                        .map_err(|error| {
1855                            rusqlite::Error::FromSqlConversionFailure(
1856                                2,
1857                                rusqlite::types::Type::Text,
1858                                Box::new(error),
1859                            )
1860                        })?
1861                        .unwrap_or(Value::Null),
1862                    budgets: budgets_json
1863                        .as_deref()
1864                        .map(serde_json::from_str)
1865                        .transpose()
1866                        .map_err(|error| {
1867                            rusqlite::Error::FromSqlConversionFailure(
1868                                3,
1869                                rusqlite::types::Type::Text,
1870                                Box::new(error),
1871                            )
1872                        })?
1873                        .unwrap_or(Value::Null),
1874                })
1875            },
1876        );
1877        match result {
1878            Ok(contract) => Ok(Some(contract)),
1879            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1880            Err(error) => Err(format!("load execution contract error: {error}")),
1881        }
1882    }
1883
1884    // ── Checkpoints ─────────────────────────────────────────────────────
1885
1886    pub fn create_resume_checkpoint(
1887        &self,
1888        run_id: &str,
1889        graph_id: &str,
1890        graph_version: &str,
1891        next_node_cursor: &str,
1892        state: &Value,
1893        budgets: &Value,
1894        budget_counters: &Value,
1895        dependency_summary: &Value,
1896        terminal_cursor: u64,
1897        event_cursor: u64,
1898    ) -> Result<CheckpointRecord, CheckpointError> {
1899        let key = self
1900            .require_integrity_key()
1901            .map_err(|_| CheckpointError::IntegrityKeyRequired)?;
1902        let state_json = serde_json::to_string(state).map_err(|_| CheckpointError::Persistence)?;
1903        let budgets_json =
1904            serde_json::to_string(budgets).map_err(|_| CheckpointError::Persistence)?;
1905        let counters_json =
1906            serde_json::to_string(budget_counters).map_err(|_| CheckpointError::Persistence)?;
1907        let dependency_json =
1908            serde_json::to_string(dependency_summary).map_err(|_| CheckpointError::Persistence)?;
1909        let state_digest = digest(state);
1910        let dependency_digest = digest(dependency_summary);
1911        let created_at = Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true);
1912        let checkpoint_id = format!("checkpoint-{run_id}-{next_node_cursor}");
1913        let mut record = CheckpointRecord {
1914            checkpoint_id,
1915            run_id: run_id.to_owned(),
1916            graph_id: graph_id.to_owned(),
1917            graph_version: graph_version.to_owned(),
1918            next_node_cursor: next_node_cursor.to_owned(),
1919            state: state.clone(),
1920            state_digest,
1921            budgets: budgets.clone(),
1922            budget_counters: budget_counters.clone(),
1923            dependency_summary: dependency_summary.clone(),
1924            dependency_digest,
1925            terminal_cursor,
1926            event_cursor,
1927            checkpoint_digest: String::new(),
1928            created_at,
1929            consumed_at: None,
1930        };
1931        record.checkpoint_digest = checkpoint_digest(&record, key);
1932
1933        let mut conn = self.conn.lock().map_err(|_| CheckpointError::Persistence)?;
1934        let tx = conn
1935            .transaction()
1936            .map_err(|_| CheckpointError::Persistence)?;
1937        #[cfg(test)]
1938        if self
1939            .checkpoint_persistence_fault
1940            .swap(false, std::sync::atomic::Ordering::SeqCst)
1941        {
1942            return Err(CheckpointError::Persistence);
1943        }
1944        tx.execute(
1945            "INSERT INTO checkpoints
1946             (run_id, node_id, attempt, input_json, status, checkpoint_id,
1947              graph_id, graph_version, next_cursor, state_json, state_digest,
1948              budgets_json, budget_counters_json, dependency_json,
1949              dependency_digest, terminal_cursor, event_cursor, checkpoint_digest,
1950              created_at, consumed_at)
1951             VALUES (?1, ?2, 0, ?3, 'available', ?4, ?5, ?6, ?7, ?3, ?8,
1952                     ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, NULL)",
1953            params![
1954                record.run_id,
1955                record.next_node_cursor,
1956                state_json,
1957                record.checkpoint_id,
1958                record.graph_id,
1959                record.graph_version,
1960                record.next_node_cursor,
1961                record.state_digest,
1962                budgets_json,
1963                counters_json,
1964                dependency_json,
1965                record.dependency_digest,
1966                record.terminal_cursor as i64,
1967                record.event_cursor as i64,
1968                record.checkpoint_digest,
1969                record.created_at,
1970            ],
1971        )
1972        .map_err(|_| CheckpointError::Persistence)?;
1973        tx.commit().map_err(|_| CheckpointError::Persistence)?;
1974        Ok(record)
1975    }
1976
1977    pub fn load_resume_checkpoint(
1978        &self,
1979        checkpoint_id: Option<&str>,
1980        run_id: Option<&str>,
1981    ) -> Result<Option<CheckpointRecord>, CheckpointError> {
1982        let key = self
1983            .require_integrity_key()
1984            .map_err(|_| CheckpointError::IntegrityKeyRequired)?;
1985        let conn = self.conn.lock().map_err(|_| CheckpointError::Persistence)?;
1986        let query = if checkpoint_id.is_some() {
1987            "SELECT run_id, graph_id, graph_version, next_cursor, state_json,
1988                    state_digest, budgets_json, budget_counters_json,
1989                    dependency_json, dependency_digest, terminal_cursor,
1990                    event_cursor, checkpoint_id, checkpoint_digest, created_at,
1991                    consumed_at
1992             FROM checkpoints WHERE checkpoint_id = ?1"
1993        } else {
1994            "SELECT run_id, graph_id, graph_version, next_cursor, state_json,
1995                    state_digest, budgets_json, budget_counters_json,
1996                    dependency_json, dependency_digest, terminal_cursor,
1997                    event_cursor, checkpoint_id, checkpoint_digest, created_at,
1998                    consumed_at
1999             FROM checkpoints WHERE run_id = ?1 AND checkpoint_id IS NOT NULL
2000             ORDER BY created_at DESC, checkpoint_id DESC LIMIT 1"
2001        };
2002        let selector = checkpoint_id.or(run_id).unwrap_or("");
2003        let row = conn.query_row(query, params![selector], checkpoint_row);
2004        let parts = match row {
2005            Ok(parts) => parts,
2006            Err(rusqlite::Error::QueryReturnedNoRows) => return Ok(None),
2007            Err(_) => return Err(CheckpointError::Persistence),
2008        };
2009        let record = checkpoint_from_parts(parts).ok_or(CheckpointError::Integrity)?;
2010        validate_checkpoint_record(&record, key)?;
2011        Ok(Some(record))
2012    }
2013
2014    pub fn consume_resume_checkpoint(
2015        &self,
2016        checkpoint_id: &str,
2017    ) -> Result<CheckpointRecord, CheckpointError> {
2018        let key = self
2019            .require_integrity_key()
2020            .map_err(|_| CheckpointError::IntegrityKeyRequired)?;
2021        let mut conn = self.conn.lock().map_err(|_| CheckpointError::Persistence)?;
2022        let tx = conn
2023            .transaction()
2024            .map_err(|_| CheckpointError::Persistence)?;
2025        let row = tx.query_row(
2026            "SELECT run_id, graph_id, graph_version, next_cursor, state_json,
2027                    state_digest, budgets_json, budget_counters_json,
2028                    dependency_json, dependency_digest, terminal_cursor,
2029                    event_cursor, checkpoint_id, checkpoint_digest, created_at,
2030                    consumed_at
2031             FROM checkpoints WHERE checkpoint_id = ?1",
2032            params![checkpoint_id],
2033            checkpoint_row,
2034        );
2035        let parts = match row {
2036            Ok(parts) => parts,
2037            Err(rusqlite::Error::QueryReturnedNoRows) => return Err(CheckpointError::NotFound),
2038            Err(_) => return Err(CheckpointError::Persistence),
2039        };
2040        let record = checkpoint_from_parts(parts).ok_or(CheckpointError::Integrity)?;
2041        validate_checkpoint_record(&record, key)?;
2042        if record.consumed_at.is_some() {
2043            return Err(CheckpointError::Consumed);
2044        }
2045        let consumed_at = Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true);
2046        let changed = tx
2047            .execute(
2048                "UPDATE checkpoints SET consumed_at = ?2, status = 'consumed'
2049                 WHERE checkpoint_id = ?1 AND consumed_at IS NULL",
2050                params![checkpoint_id, consumed_at],
2051            )
2052            .map_err(|_| CheckpointError::Persistence)?;
2053        if changed != 1 {
2054            return Err(CheckpointError::Consumed);
2055        }
2056        tx.commit().map_err(|_| CheckpointError::Persistence)?;
2057        let mut consumed = record;
2058        consumed.consumed_at = Some(consumed_at);
2059        Ok(consumed)
2060    }
2061
2062    // ── Durable checkpoint-bound approvals ───────────────────────────
2063
2064    pub fn create_checkpoint_approval(
2065        &self,
2066        checkpoint_id: &str,
2067        graph_id: &str,
2068        graph_version: &str,
2069        next_node_cursor: &str,
2070        expected_state: &Value,
2071        expected_budgets: &Value,
2072        expected_budget_counters: &Value,
2073        dependency_summary: &Value,
2074        audience: &str,
2075        prompt_digest: &str,
2076        allowed_decisions: &[String],
2077        expires_at: &str,
2078    ) -> Result<ApprovalRecord, ApprovalError> {
2079        let key = self
2080            .require_integrity_key()
2081            .map_err(|_| ApprovalError::IntegrityKeyRequired)?;
2082        let mut conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
2083        let tx = conn.transaction().map_err(|_| ApprovalError::Persistence)?;
2084        let checkpoint =
2085            load_checkpoint_from_tx(&tx, checkpoint_id, key).map_err(ApprovalError::Checkpoint)?;
2086        if checkpoint.consumed_at.is_some() {
2087            return Err(ApprovalError::Checkpoint(CheckpointError::Consumed));
2088        }
2089        if checkpoint.graph_id != graph_id
2090            || checkpoint.graph_version != graph_version
2091            || checkpoint.next_node_cursor != next_node_cursor
2092            || checkpoint.state != *expected_state
2093            || checkpoint.budgets != *expected_budgets
2094            || checkpoint.budget_counters != *expected_budget_counters
2095            || checkpoint.dependency_summary != *dependency_summary
2096            || checkpoint.terminal_cursor != 0
2097            || checkpoint.event_cursor != 0
2098        {
2099            return Err(ApprovalError::Checkpoint(CheckpointError::Integrity));
2100        }
2101
2102        let allowed_json =
2103            serde_json::to_string(allowed_decisions).map_err(|_| ApprovalError::Persistence)?;
2104        let existing = tx
2105            .query_row(
2106                &format!(
2107                    "SELECT {APPROVAL_COLUMNS} FROM approval_requests
2108                          WHERE checkpoint_id = ?1 AND audience = ?2 AND status = 'pending'"
2109                ),
2110                params![checkpoint_id, audience],
2111                approval_row,
2112            )
2113            .optional()
2114            .map_err(|_| ApprovalError::Persistence)?;
2115        if let Some(parts) = existing {
2116            let current = parse_approval_parts(parts, key)?;
2117            if current.graph_id == graph_id
2118                && current.graph_version == graph_version
2119                && current.checkpoint_digest == checkpoint.checkpoint_digest
2120                && current.prompt_digest == prompt_digest
2121                && current.allowed_decisions.as_slice() == allowed_decisions
2122                && current.expires_at == expires_at
2123            {
2124                tx.commit().map_err(|_| ApprovalError::Persistence)?;
2125                return Ok(current);
2126            }
2127            return Err(ApprovalError::Conflict);
2128        }
2129
2130        let created_at = Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true);
2131        let mut record = ApprovalRecord {
2132            approval_id: format!("approval-{}", uuid_like()),
2133            checkpoint_id: checkpoint.checkpoint_id.clone(),
2134            run_id: checkpoint.run_id.clone(),
2135            graph_id: graph_id.to_owned(),
2136            graph_version: graph_version.to_owned(),
2137            checkpoint_digest: checkpoint.checkpoint_digest.clone(),
2138            audience: audience.to_owned(),
2139            prompt_digest: prompt_digest.to_owned(),
2140            allowed_decisions: allowed_decisions.to_vec(),
2141            approval_digest: String::new(),
2142            status: "pending".into(),
2143            decision: None,
2144            decided_by: None,
2145            decided_at: None,
2146            expires_at: expires_at.to_owned(),
2147            created_at,
2148        };
2149        record.approval_digest = approval_digest(&record, key);
2150        tx.execute(
2151            "INSERT INTO approval_requests
2152             (approval_id, run_id, node_id, checkpoint_id, graph_id, graph_version,
2153              checkpoint_digest, audience, prompt, prompt_digest, allowed_decisions,
2154              approval_digest, status, expires_at, created_at)
2155             VALUES (?1, ?2, 'durable_checkpoint', ?3, ?4, ?5, ?6, ?7, '', ?8, ?9, ?10,
2156                     'pending', ?11, ?12)",
2157            params![
2158                record.approval_id,
2159                record.run_id,
2160                record.checkpoint_id,
2161                record.graph_id,
2162                record.graph_version,
2163                record.checkpoint_digest,
2164                record.audience,
2165                record.prompt_digest,
2166                allowed_json,
2167                record.approval_digest,
2168                record.expires_at,
2169                record.created_at,
2170            ],
2171        )
2172        .map_err(|_| ApprovalError::Persistence)?;
2173        tx.commit().map_err(|_| ApprovalError::Persistence)?;
2174        Ok(record)
2175    }
2176
2177    pub fn get_checkpoint_approval(
2178        &self,
2179        approval_id: &str,
2180    ) -> Result<Option<ApprovalRecord>, ApprovalError> {
2181        let key = self
2182            .require_integrity_key()
2183            .map_err(|_| ApprovalError::IntegrityKeyRequired)?;
2184        let conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
2185        let result = conn.query_row(
2186            &format!(
2187                "SELECT {APPROVAL_COLUMNS} FROM approval_requests
2188                      WHERE approval_id = ?1 AND checkpoint_id IS NOT NULL"
2189            ),
2190            params![approval_id],
2191            approval_row,
2192        );
2193        match result {
2194            Ok(parts) => Ok(Some(parse_approval_parts(parts, key)?)),
2195            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
2196            Err(_) => Err(ApprovalError::Persistence),
2197        }
2198    }
2199
2200    pub fn list_checkpoint_approvals(
2201        &self,
2202        run_id: Option<&str>,
2203        status: Option<&str>,
2204        limit: usize,
2205    ) -> Result<Vec<ApprovalRecord>, ApprovalError> {
2206        let key = self
2207            .require_integrity_key()
2208            .map_err(|_| ApprovalError::IntegrityKeyRequired)?;
2209        let conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
2210        let mut sql = format!(
2211            "SELECT {APPROVAL_COLUMNS} FROM approval_requests
2212                              WHERE checkpoint_id IS NOT NULL"
2213        );
2214        if run_id.is_some() {
2215            sql.push_str(" AND run_id = ?1");
2216        }
2217        if status.is_some() {
2218            sql.push_str(if run_id.is_some() {
2219                " AND status = ?2"
2220            } else {
2221                " AND status = ?1"
2222            });
2223        }
2224        sql.push_str(" ORDER BY created_at DESC, approval_id DESC LIMIT ?3");
2225        if run_id.is_none() && status.is_none() {
2226            sql = format!("SELECT {APPROVAL_COLUMNS} FROM approval_requests
2227                          WHERE checkpoint_id IS NOT NULL ORDER BY created_at DESC, approval_id DESC LIMIT ?1");
2228        } else if run_id.is_none() || status.is_none() {
2229            sql = sql.replace("LIMIT ?3", "LIMIT ?2");
2230        }
2231        let mut statement = conn.prepare(&sql).map_err(|_| ApprovalError::Persistence)?;
2232        let mut rows = if run_id.is_some() && status.is_some() {
2233            statement.query(params![run_id, status, limit.min(200) as i64])
2234        } else if run_id.is_some() {
2235            statement.query(params![run_id, limit.min(200) as i64])
2236        } else if status.is_some() {
2237            statement.query(params![status, limit.min(200) as i64])
2238        } else {
2239            statement.query(params![limit.min(200) as i64])
2240        }
2241        .map_err(|_| ApprovalError::Persistence)?;
2242        let mut approvals = Vec::new();
2243        while let Some(row) = rows.next().map_err(|_| ApprovalError::Persistence)? {
2244            approvals.push(parse_approval_parts(
2245                approval_row(row).map_err(|_| ApprovalError::Persistence)?,
2246                key,
2247            )?);
2248        }
2249        Ok(approvals)
2250    }
2251
2252    pub fn checkpoint_approval_status(
2253        &self,
2254        checkpoint_id: &str,
2255    ) -> Result<Option<String>, ApprovalError> {
2256        let conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
2257        conn.query_row(
2258            "SELECT status FROM approval_requests
2259             WHERE checkpoint_id = ?1 AND checkpoint_id IS NOT NULL
2260             ORDER BY created_at DESC LIMIT 1",
2261            params![checkpoint_id],
2262            |row| row.get(0),
2263        )
2264        .optional()
2265        .map_err(|_| ApprovalError::Persistence)
2266    }
2267
2268    pub fn decide_checkpoint_approval(
2269        &self,
2270        approval_id: &str,
2271        decision: &str,
2272        actor: &str,
2273        now: DateTime<Utc>,
2274    ) -> Result<ApprovedCheckpoint, ApprovalError> {
2275        let key = self
2276            .require_integrity_key()
2277            .map_err(|_| ApprovalError::IntegrityKeyRequired)?;
2278        let mut conn = self.conn.lock().map_err(|_| ApprovalError::Persistence)?;
2279        let tx = conn.transaction().map_err(|_| ApprovalError::Persistence)?;
2280        let parts = tx
2281            .query_row(
2282                &format!(
2283                    "SELECT {APPROVAL_COLUMNS} FROM approval_requests
2284                          WHERE approval_id = ?1 AND checkpoint_id IS NOT NULL"
2285                ),
2286                params![approval_id],
2287                approval_row,
2288            )
2289            .optional()
2290            .map_err(|_| ApprovalError::Persistence)?
2291            .ok_or(ApprovalError::NotFound)?;
2292        let mut approval = parse_approval_parts(parts, key)?;
2293        if approval.status == "expired" {
2294            return Err(ApprovalError::Expired);
2295        }
2296        if approval.status != "pending" {
2297            return Err(ApprovalError::AlreadyDecided);
2298        }
2299        let expired = DateTime::parse_from_rfc3339(&approval.expires_at)
2300            .map_err(|_| ApprovalError::Integrity)?
2301            .with_timezone(&Utc)
2302            <= now;
2303        if expired {
2304            let checkpoint = load_checkpoint_from_tx(&tx, &approval.checkpoint_id, key)
2305                .map_err(ApprovalError::Checkpoint)?;
2306            if checkpoint.run_id != approval.run_id
2307                || checkpoint.graph_id != approval.graph_id
2308                || checkpoint.graph_version != approval.graph_version
2309                || checkpoint.checkpoint_digest != approval.checkpoint_digest
2310            {
2311                return Err(ApprovalError::Integrity);
2312            }
2313            approval.status = "expired".into();
2314            approval.decision = None;
2315            approval.decided_by = None;
2316            approval.decided_at = Some(now.to_rfc3339_opts(SecondsFormat::Nanos, true));
2317            approval.approval_digest = approval_digest(&approval, key);
2318            let changed = tx
2319                .execute(
2320                    "UPDATE approval_requests SET status = 'expired', decision = NULL,
2321                     decided_by = NULL, decided_at = ?2, approval_digest = ?3
2322                     WHERE approval_id = ?1 AND status = 'pending'",
2323                    params![approval_id, approval.decided_at, approval.approval_digest],
2324                )
2325                .map_err(|_| ApprovalError::Persistence)?;
2326            if changed != 1 {
2327                return Err(ApprovalError::AlreadyDecided);
2328            }
2329            if checkpoint.consumed_at.is_none() {
2330                tx.execute(
2331                    "UPDATE checkpoints SET consumed_at = ?2, status = 'consumed'
2332                     WHERE checkpoint_id = ?1 AND consumed_at IS NULL",
2333                    params![
2334                        approval.checkpoint_id,
2335                        now.to_rfc3339_opts(SecondsFormat::Nanos, true)
2336                    ],
2337                )
2338                .map_err(|_| ApprovalError::Persistence)?;
2339            }
2340            tx.commit().map_err(|_| ApprovalError::Persistence)?;
2341            return Err(ApprovalError::Expired);
2342        }
2343
2344        if !approval
2345            .allowed_decisions
2346            .iter()
2347            .any(|allowed| allowed == decision)
2348        {
2349            return Err(ApprovalError::DecisionNotAllowed);
2350        }
2351
2352        let checkpoint = load_checkpoint_from_tx(&tx, &approval.checkpoint_id, key)
2353            .map_err(ApprovalError::Checkpoint)?;
2354        if checkpoint.run_id != approval.run_id
2355            || checkpoint.graph_id != approval.graph_id
2356            || checkpoint.graph_version != approval.graph_version
2357            || checkpoint.checkpoint_digest != approval.checkpoint_digest
2358        {
2359            return Err(ApprovalError::Integrity);
2360        }
2361        if checkpoint.consumed_at.is_some() {
2362            return Err(ApprovalError::Checkpoint(CheckpointError::Consumed));
2363        }
2364
2365        let decided_at = now.to_rfc3339_opts(SecondsFormat::Nanos, true);
2366        let status = if decision == "approve" {
2367            "approved"
2368        } else {
2369            "rejected"
2370        };
2371        approval.status = status.into();
2372        approval.decision = Some(decision.into());
2373        approval.decided_by = Some(actor.into());
2374        approval.decided_at = Some(decided_at.clone());
2375        approval.approval_digest = approval_digest(&approval, key);
2376        let changed = tx
2377            .execute(
2378                "UPDATE approval_requests SET status = ?2, decision = ?3,
2379                 decided_by = ?4, decided_at = ?5, approval_digest = ?6
2380                 WHERE approval_id = ?1 AND status = 'pending'",
2381                params![
2382                    approval_id,
2383                    status,
2384                    decision,
2385                    actor,
2386                    decided_at,
2387                    approval.approval_digest
2388                ],
2389            )
2390            .map_err(|_| ApprovalError::Persistence)?;
2391        if changed != 1 {
2392            return Err(ApprovalError::AlreadyDecided);
2393        }
2394        let consumed_at = now.to_rfc3339_opts(SecondsFormat::Nanos, true);
2395        let checkpoint_changed = tx
2396            .execute(
2397                "UPDATE checkpoints SET consumed_at = ?2, status = 'consumed'
2398                 WHERE checkpoint_id = ?1 AND consumed_at IS NULL",
2399                params![approval.checkpoint_id, consumed_at],
2400            )
2401            .map_err(|_| ApprovalError::Persistence)?;
2402        if checkpoint_changed != 1 {
2403            return Err(ApprovalError::Checkpoint(CheckpointError::Consumed));
2404        }
2405        tx.commit().map_err(|_| ApprovalError::Persistence)?;
2406        Ok(ApprovedCheckpoint {
2407            approval,
2408            checkpoint,
2409        })
2410    }
2411
2412    pub fn approval_receipt_value(approval: &ApprovalRecord) -> Value {
2413        let decision = approval.decision.clone().unwrap_or_default();
2414        let decided_by_digest = approval
2415            .decided_by
2416            .as_ref()
2417            .map(|actor| digest(&Value::String(actor.clone())));
2418        let decision_digest = digest(&serde_json::json!({
2419            "approval_digest": approval.approval_digest,
2420            "decision": decision,
2421            "decided_by_digest": decided_by_digest,
2422            "decided_at": approval.decided_at,
2423        }));
2424        serde_json::json!({
2425            "approval_id": approval.approval_id,
2426            "approval_digest": approval.approval_digest,
2427            "checkpoint_id": approval.checkpoint_id,
2428            "checkpoint_digest": approval.checkpoint_digest,
2429            "decision": approval.decision,
2430            "decided_by_digest": decided_by_digest,
2431            "decided_at": approval.decided_at,
2432            "allowed_decisions_digest": digest(&serde_json::json!(approval.allowed_decisions)),
2433            "audience_digest": digest(&Value::String(approval.audience.clone())),
2434            "prompt_digest": approval.prompt_digest,
2435            "decision_digest": decision_digest,
2436        })
2437    }
2438
2439    #[cfg(test)]
2440    pub(crate) fn fail_checkpoint_persistence(&self) {
2441        self.checkpoint_persistence_fault
2442            .store(true, std::sync::atomic::Ordering::SeqCst);
2443    }
2444
2445    pub fn save_checkpoint(
2446        &self,
2447        run_id: &str,
2448        node_id: &str,
2449        attempt: u32,
2450        input_json: &str,
2451        output_json: Option<&str>,
2452        status: &str,
2453        error: Option<&str>,
2454    ) -> Result<(), String> {
2455        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
2456        conn.execute(
2457            "INSERT INTO checkpoints (run_id, node_id, attempt, input_json, output_json, status, error)
2458             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
2459             ON CONFLICT(run_id, node_id, attempt) DO UPDATE SET
2460                output_json = COALESCE(excluded.output_json, output_json),
2461                status = excluded.status,
2462                error = COALESCE(excluded.error, error)",
2463            params![run_id, node_id, attempt, input_json, output_json, status, error],
2464        )
2465        .map_err(|e| format!("save_checkpoint error: {e}"))?;
2466        Ok(())
2467    }
2468
2469    // ── Events ──────────────────────────────────────────────────────────
2470
2471    pub fn save_event(
2472        &self,
2473        run_id: &str,
2474        seq: u64,
2475        event_type: &str,
2476        event_json: &str,
2477    ) -> Result<(), String> {
2478        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
2479        conn.execute(
2480            "INSERT OR IGNORE INTO events (run_id, seq, event_type, event_json)
2481             VALUES (?1, ?2, ?3, ?4)",
2482            params![run_id, seq, event_type, event_json],
2483        )
2484        .map_err(|e| format!("save_event error: {e}"))?;
2485        Ok(())
2486    }
2487
2488    pub fn load_events(
2489        &self,
2490        run_id: &str,
2491        cursor: u64,
2492        limit: usize,
2493    ) -> Result<Option<Value>, String> {
2494        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
2495        let mut stmt = conn
2496            .prepare("SELECT seq, event_json FROM events WHERE run_id = ?1 AND seq >= ?2 ORDER BY seq LIMIT ?3")
2497            .map_err(|e| format!("load events error: {e}"))?;
2498        let events: Vec<Value> = stmt
2499            .query_map(params![run_id, cursor, limit.min(200) as i64], |row| {
2500                let seq: u64 = row.get(0)?;
2501                let event_json: String = row.get(1)?;
2502                let event = serde_json::from_str::<Value>(&event_json)
2503                    .unwrap_or_else(|_| serde_json::json!({"receipt":"terminal event persisted with reduced fidelity"}));
2504                Ok(serde_json::json!({"cursor": seq, "event": event}))
2505            })
2506            .map_err(|e| format!("load events error: {e}"))?
2507            .filter_map(Result::ok)
2508            .collect();
2509        let exists: bool = conn
2510            .query_row(
2511                "SELECT EXISTS(SELECT 1 FROM events WHERE run_id = ?1)",
2512                params![run_id],
2513                |row| row.get(0),
2514            )
2515            .map_err(|e| format!("load events existence error: {e}"))?;
2516        if !exists {
2517            return Ok(None);
2518        }
2519        let first: u64 = conn
2520            .query_row(
2521                "SELECT MIN(seq) FROM events WHERE run_id = ?1",
2522                params![run_id],
2523                |row| row.get(0),
2524            )
2525            .map_err(|e| format!("load events first cursor error: {e}"))?;
2526        let next_cursor: u64 = conn
2527            .query_row(
2528                "SELECT COALESCE(MAX(seq) + 1, 0) FROM events WHERE run_id = ?1",
2529                params![run_id],
2530                |row| row.get(0),
2531            )
2532            .map_err(|e| format!("load events next cursor error: {e}"))?;
2533        Ok(Some(
2534            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}),
2535        ))
2536    }
2537
2538    // ── Idempotency ─────────────────────────────────────────────────────
2539
2540    /// Look up a cached idempotent response together with the canonical request
2541    /// digest it was bound to. A NULL digest is a pre-migration record and must
2542    /// never be replayed for a new request.
2543    pub fn check_idempotency(&self, key: &str) -> Result<Option<(Option<String>, Value)>, String> {
2544        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
2545        let mut stmt = conn
2546            .prepare("SELECT request_digest, result_json FROM idempotency_keys WHERE key = ?1 AND request_digest IS NOT NULL AND valid = 1")
2547            .map_err(|e| format!("idempotency error: {e}"))?;
2548        let result: Option<(Option<String>, String)> = stmt
2549            .query_row(params![key], |row| Ok((row.get(0)?, row.get(1)?)))
2550            .ok();
2551        match result {
2552            Some((request_digest, json_str)) => serde_json::from_str(&json_str)
2553                .map(|result| Some((request_digest, result)))
2554                .map_err(|e| format!("json parse error: {e}")),
2555            None => Ok(None),
2556        }
2557    }
2558
2559    pub fn save_idempotency(
2560        &self,
2561        key: &str,
2562        request_digest: &str,
2563        result_json: &str,
2564    ) -> Result<bool, String> {
2565        let conn = self.conn.lock().map_err(|e| format!("lock error: {e}"))?;
2566        let inserted = conn.execute(
2567            "INSERT OR IGNORE INTO idempotency_keys (key, request_digest, result_json) VALUES (?1, ?2, ?3)",
2568            params![key, request_digest, result_json],
2569        )
2570        .map_err(|e| format!("idempotency error: {e}"))?;
2571        Ok(inserted == 1)
2572    }
2573
2574    pub fn data_dir(&self) -> Option<PathBuf> {
2575        // Returns None since we don't store the path separately.
2576        // The store is identified by its existence, not a path reference.
2577        None
2578    }
2579}
2580
2581#[cfg(test)]
2582mod tests {
2583    use super::*;
2584
2585    fn configure_test_integrity_key() {
2586        let path = std::env::temp_dir().join("agent-graph-mcp-unit-integrity.key");
2587        std::fs::write(&path, [0x5au8; 32]).expect("test integrity key");
2588        std::env::set_var("AGENT_GRAPH_INTEGRITY_KEY_PATH", path);
2589    }
2590
2591    #[test]
2592    fn graph_delete_fault_rolls_back_graph_and_versions_together() {
2593        let temp = tempfile::tempdir().expect("graph database");
2594        let store = PersistentStore::open(temp.path()).expect("store");
2595        store
2596            .save_graph("atomic-delete", "{\"name\":\"atomic-delete\"}", "v1", false)
2597            .expect("graph");
2598        store
2599            .set_graph_retention(
2600                "atomic-delete",
2601                "delete_candidate",
2602                "test deletion rollback",
2603                "unit-test",
2604                None,
2605            )
2606            .expect("candidate");
2607        store
2608            .set_graph_retention(
2609                "atomic-delete",
2610                "delete_approved",
2611                "test deletion rollback approved",
2612                "unit-test",
2613                None,
2614            )
2615            .expect("approval");
2616        store.fail_graph_delete_after_graph_row();
2617        assert!(store.delete_graph("atomic-delete").is_err());
2618        assert!(store.load_graph("atomic-delete").unwrap().is_some());
2619        assert_eq!(
2620            store.list_graph_versions("atomic-delete").unwrap(),
2621            vec!["v1"]
2622        );
2623    }
2624
2625    #[test]
2626    fn checkpoint_persistence_fault_leaves_no_resumable_row() {
2627        configure_test_integrity_key();
2628        let temp = tempfile::tempdir().expect("checkpoint database");
2629        let store = PersistentStore::open(temp.path()).expect("store");
2630        store
2631            .save_graph("checkpoint-fault", "{}", "version", false)
2632            .expect("graph");
2633        store
2634            .save_execution_with_budgets(
2635                "run-checkpoint-fault",
2636                "checkpoint-fault",
2637                "version",
2638                "checkpointed",
2639                "{}",
2640                Some("null"),
2641            )
2642            .expect("execution");
2643        store.fail_checkpoint_persistence();
2644        let result = store.create_resume_checkpoint(
2645            "run-checkpoint-fault",
2646            "checkpoint-fault",
2647            "version",
2648            "entry",
2649            &serde_json::json!({"__input__":null}),
2650            &Value::Null,
2651            &serde_json::json!({"nodes":0,"llm_calls":0,"wall_clock_ms":0}),
2652            &serde_json::json!({"eligible":true}),
2653            0,
2654            0,
2655        );
2656        assert_eq!(result, Err(CheckpointError::Persistence));
2657        assert_eq!(
2658            store.load_resume_checkpoint(None, Some("run-checkpoint-fault")),
2659            Ok(None)
2660        );
2661    }
2662}