1use 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
17pub 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 pub fn open(data_dir: &Path) -> Result<Self, String> {
531 Self::open_with_integrity_key(data_dir, None)
532 }
533
534 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 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 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 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 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 }
1370 OperatorAction::PurgeGraph => {
1371 return Err(OperatorRetentionError::InvalidTransition);
1372 }
1374 OperatorAction::PromoteTemplate => {
1375 return Err(OperatorRetentionError::InvalidAction);
1376 }
1378 OperatorAction::Migrate => {
1379 return Err(OperatorRetentionError::InvalidAction);
1380 }
1382 OperatorAction::Install => {
1383 }
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 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 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 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 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 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 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 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 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}