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