Skip to main content

agent_graph_mcp/
server.rs

1//! MCP server handler using rmcp's #[tool_router] macro.
2//!
3//! Each #[tool] method becomes an MCP tool that Hermes can discover and call.
4//! The rmcp macro auto-generates JSON Schema from the parameter structs in tools.rs.
5
6use std::collections::HashMap;
7use std::sync::Mutex;
8use std::time::{Duration, Instant};
9
10use chrono::{SecondsFormat, Utc};
11use rmcp::{
12    handler::server::{router::tool::ToolRouter, wrapper::Parameters},
13    tool, tool_handler, tool_router, ErrorData, Json, ServerHandler,
14};
15use serde_json::Value;
16
17use std::path::PathBuf;
18
19use crate::evidence::{digest, validate_witness_capture, WitnessCapture, WitnessError};
20use crate::run_manager::{initial_state_for_input, RunBudgets, RunManager};
21use crate::spec::{ensure_size, parse_and_validate, GraphSpec, MAX_GRAPHS, MAX_INPUT_BYTES};
22use crate::store::{
23    ApprovalError, ApprovalRecord, CheckpointError, CheckpointRecord, GraphDeleteResult,
24    PersistentStore,
25};
26use crate::templates;
27use crate::tools::*;
28
29fn internal_error(message: impl Into<std::borrow::Cow<'static, str>>) -> ErrorData {
30    ErrorData::internal_error(message, None)
31}
32
33fn invalid_params(message: impl Into<std::borrow::Cow<'static, str>>) -> ErrorData {
34    ErrorData::invalid_params(message, None)
35}
36
37fn structured_output(value: Value) -> Json<StructuredOutput> {
38    Json(StructuredOutput {
39        ok: true,
40        status: None,
41        data: Some(value),
42        error: None,
43        error_code: None,
44        graph_id: None,
45        graph_version: None,
46        run_id: None,
47    })
48}
49
50fn error_output(message: impl Into<String>, code: impl Into<String>) -> Json<StructuredOutput> {
51    Json(StructuredOutput {
52        ok: false,
53        status: None,
54        data: None,
55        error: Some(message.into()),
56        error_code: Some(code.into()),
57        graph_id: None,
58        graph_version: None,
59        run_id: None,
60    })
61}
62
63fn structured_from_value(value: Value) -> Result<Json<StructuredOutput>, ErrorData> {
64    serde_json::from_value::<StructuredOutput>(value)
65        .map(Json)
66        .map_err(|e| internal_error(format!("cached idempotency decode: {e}")))
67}
68
69fn canonical_request_value(value: &Value) -> Value {
70    match value {
71        Value::String(raw) => serde_json::from_str(raw).unwrap_or_else(|_| value.clone()),
72        _ => value.clone(),
73    }
74}
75
76fn check_idempotency(
77    store: Option<&PersistentStore>,
78    key: Option<&str>,
79    request_digest: &str,
80) -> Result<Option<Json<StructuredOutput>>, ErrorData> {
81    let Some((store, key)) = store.zip(key) else {
82        return Ok(None);
83    };
84    let Some((stored_digest, cached)) = store.check_idempotency(key).map_err(internal_error)?
85    else {
86        return Ok(None);
87    };
88    if stored_digest.as_deref() == Some(request_digest) {
89        return structured_from_value(cached).map(Some);
90    }
91    Ok(Some(error_output(
92        "idempotency key is already bound to different request material",
93        "IDEMPOTENCY_CONFLICT",
94    )))
95}
96
97fn persist_idempotency(
98    store: &PersistentStore,
99    key: &str,
100    request_digest: &str,
101    output: &Json<StructuredOutput>,
102) -> Result<Option<Json<StructuredOutput>>, ErrorData> {
103    let result_json =
104        serde_json::to_string(&output.0).map_err(|e| internal_error(e.to_string()))?;
105    if store
106        .save_idempotency(key, request_digest, &result_json)
107        .map_err(internal_error)?
108    {
109        return Ok(None);
110    }
111    // Another request won the insert. Return its exact cached result so a
112    // concurrent same-key caller cannot observe a result that was not stored.
113    check_idempotency(Some(store), Some(key), request_digest)
114}
115
116fn output_with_meta(
117    data: Value,
118    graph_id: Option<&str>,
119    graph_version: Option<&str>,
120    run_id: Option<&str>,
121) -> Json<StructuredOutput> {
122    Json(StructuredOutput {
123        ok: true,
124        status: None,
125        data: Some(data),
126        error: None,
127        error_code: None,
128        graph_id: graph_id.map(String::from),
129        graph_version: graph_version.map(String::from),
130        run_id: run_id.map(String::from),
131    })
132}
133
134fn checkpoint_error_output(error: CheckpointError) -> Json<StructuredOutput> {
135    error_output(error.message(), error.code())
136}
137
138fn approval_error_output(error: ApprovalError) -> Json<StructuredOutput> {
139    error_output(error.message(), error.code())
140}
141
142fn approval_value(record: &ApprovalRecord) -> Value {
143    serde_json::json!({
144        "approval_id": record.approval_id,
145        "checkpoint_id": record.checkpoint_id,
146        "run_id": record.run_id,
147        "graph_id": record.graph_id,
148        "graph_version": record.graph_version,
149        "checkpoint_digest": record.checkpoint_digest,
150        "audience": record.audience,
151        "prompt_digest": record.prompt_digest,
152        "allowed_decisions": record.allowed_decisions,
153        "approval_digest": record.approval_digest,
154        "status": record.status,
155        "decision": record.decision,
156        "decided_by": record.decided_by,
157        "decided_at": record.decided_at,
158        "expires_at": record.expires_at,
159        "created_at": record.created_at,
160    })
161}
162
163fn checkpoint_value(record: &CheckpointRecord) -> Value {
164    serde_json::json!({
165        "checkpoint_id": record.checkpoint_id,
166        "run_id": record.run_id,
167        "graph_id": record.graph_id,
168        "graph_version": record.graph_version,
169        "next_node_cursor": record.next_node_cursor,
170        "state": record.state,
171        "state_digest": record.state_digest,
172        "budgets": record.budgets,
173        "budget_counters": record.budget_counters,
174        "dependency_summary": record.dependency_summary,
175        "dependency_digest": record.dependency_digest,
176        "terminal_cursor": record.terminal_cursor,
177        "event_cursor": record.event_cursor,
178        "checkpoint_digest": record.checkpoint_digest,
179        "created_at": record.created_at,
180        "consumed_at": record.consumed_at,
181        "status": if record.consumed_at.is_some() { "consumed" } else { "available" },
182        "resume_capability": "deterministic_local_resume",
183    })
184}
185
186#[derive(Clone)]
187struct RegisteredGraph {
188    spec: GraphSpec,
189    normalized: Value,
190    version: String,
191    warnings: Vec<String>,
192}
193
194pub struct AgentGraphServer {
195    tool_router: ToolRouter<Self>,
196    base_url: String,
197    default_model: String,
198    graphs: Mutex<HashMap<String, RegisteredGraph>>,
199    runs: Mutex<RunManager>,
200    store: Option<PersistentStore>,
201}
202
203impl AgentGraphServer {
204    fn graph_requires_witness_store(spec: &GraphSpec) -> bool {
205        spec.nodes.iter().any(|node| node.evidence_required)
206    }
207
208    fn witness_error_output(error: WitnessError) -> Json<StructuredOutput> {
209        error_output(error.message, error.code)
210    }
211
212    pub fn new(
213        base_url: String,
214        default_model: String,
215        data_dir: Option<PathBuf>,
216        integrity_key_path: Option<PathBuf>,
217    ) -> Result<Self, String> {
218        let store = match data_dir {
219            Some(ref dir) => Some(PersistentStore::open_with_integrity_key(
220                dir,
221                integrity_key_path.as_deref(),
222            )?),
223            None => None,
224        };
225        if let Some(ref store) = store {
226            store.recover_incomplete_executions()?;
227        }
228
229        let server = Self {
230            base_url,
231            default_model,
232            graphs: Mutex::new(HashMap::new()),
233            runs: Mutex::new(RunManager::default()),
234            store,
235            tool_router: Self::tool_router(),
236        };
237
238        // Restore persisted graphs on startup
239        if let Some(ref store) = server.store {
240            if let Ok(graphs) = store.list_graphs() {
241                for (name, hash, _created) in graphs {
242                    if let Ok(Some((spec_json, _))) = store.load_graph(&name) {
243                        if let Ok(spec) = serde_json::from_str::<GraphSpec>(&spec_json) {
244                            let normalized = serde_json::to_value(&spec).unwrap_or_default();
245                            server.graphs.lock().unwrap().insert(
246                                name,
247                                RegisteredGraph {
248                                    spec,
249                                    normalized,
250                                    version: hash,
251                                    warnings: Vec::new(),
252                                },
253                            );
254                        }
255                    }
256                }
257            }
258        }
259
260        Ok(server)
261    }
262
263    fn safe_provider_label(&self) -> String {
264        let url = &self.base_url;
265        let without_fragment = url.split(['?', '#']).next().unwrap_or(url);
266        if let Some((scheme, rest)) = without_fragment.split_once("://") {
267            let authority_and_path = rest.rsplit_once('@').map(|(_, safe)| safe).unwrap_or(rest);
268            format!("{scheme}://{authority_and_path}")
269        } else {
270            "server-configured".into()
271        }
272    }
273
274    fn persist_terminal(
275        store: Option<PersistentStore>,
276        record: crate::run_manager::RunRecord,
277    ) -> Result<(), String> {
278        let Some(store) = store else {
279            return Ok(());
280        };
281        let final_state = serde_json::to_string(&record.final_state)
282            .map_err(|e| format!("serialize terminal state error: {e}"))?;
283        // Persist one bounded terminal projection atomically. This is not replayable
284        // execution history and does not make the run resumable.
285        let events = record
286            .events
287            .iter()
288            .map(|entry| {
289                let seq = entry.get("cursor").and_then(Value::as_u64).unwrap_or(0);
290                let event = entry.get("event").cloned().unwrap_or_else(|| {
291                    serde_json::json!({"receipt": "terminal event persisted with reduced fidelity"})
292                });
293                let event_type = event
294                    .as_object()
295                    .and_then(|object| object.keys().next().cloned())
296                    .unwrap_or_else(|| "run_event".into());
297                Ok((seq, event_type, event.to_string()))
298            })
299            .collect::<Result<Vec<_>, String>>()?;
300        let mut durable_receipt = record.receipt.clone();
301        if let Some(object) = durable_receipt.as_object_mut() {
302            object.insert(
303                "persistence_status".into(),
304                Value::String("durable_terminal".into()),
305            );
306        }
307        let receipt = serde_json::to_string(&durable_receipt)
308            .map_err(|e| format!("serialize terminal receipt error: {e}"))?;
309        let durable_bundle = crate::evidence::bundle(
310            &record.run_id,
311            &record.graph_version,
312            &record.input,
313            &record.state,
314            &durable_receipt,
315        );
316        let bundle = serde_json::to_string(&durable_bundle)
317            .map_err(|e| format!("serialize terminal bundle error: {e}"))?;
318        store.persist_terminal_projection(
319            &record.run_id,
320            &record.status,
321            &final_state,
322            record.steps.len(),
323            &events,
324            &receipt,
325            &bundle,
326        )?;
327        Ok(())
328    }
329
330    fn persist_terminal_and_mark(
331        runs: crate::run_manager::RunManager,
332        store: Option<PersistentStore>,
333        record: crate::run_manager::RunRecord,
334    ) {
335        if store.is_none() {
336            runs.mark_persistence(&record.run_id, "volatile_no_store", None);
337            return;
338        }
339        match Self::persist_terminal(store, record.clone()) {
340            Ok(()) => runs.mark_persistence(&record.run_id, "durable_terminal", None),
341            Err(error) => {
342                tracing::error!(%error, "terminal run persistence failed; run remains volatile");
343                runs.mark_persistence(&record.run_id, "volatile_persistence_failed", Some(error));
344            }
345        }
346    }
347
348    fn stored_run(&self, run_id: &str) -> Result<Option<Value>, ErrorData> {
349        let Some(store) = &self.store else {
350            return Ok(None);
351        };
352        let Some(mut record) = store.load_execution(run_id).map_err(internal_error)? else {
353            return Ok(None);
354        };
355        if let Some(receipt) = store
356            .load_terminal_receipt(run_id)
357            .map_err(internal_error)?
358            .and_then(|value| value.get("receipt").cloned())
359        {
360            if let Some(object) = record.as_object_mut() {
361                for key in ["budgets", "budget_counters", "budget_exhausted"] {
362                    if let Some(value) = receipt.get(key) {
363                        object.insert(key.into(), value.clone());
364                    }
365                }
366                object.insert("receipt".into(), receipt);
367            }
368        }
369        Ok(Some(record))
370    }
371
372    fn resolve_graph(
373        &self,
374        graph_id: &str,
375        requested_version: Option<&str>,
376    ) -> Result<RegisteredGraph, ErrorData> {
377        let current = self
378            .graphs
379            .lock()
380            .map_err(|e| internal_error(e.to_string()))?
381            .get(graph_id)
382            .cloned()
383            .ok_or_else(|| invalid_params(format!("graph '{graph_id}' not found")))?;
384        let Some(requested_version) = requested_version else {
385            return Ok(current);
386        };
387        if requested_version == current.version {
388            return Ok(current);
389        }
390        let store = self.store.as_ref().ok_or_else(|| {
391            invalid_params("historical graph versions require SQLite persistence")
392        })?;
393        let serialized = store
394            .load_graph_version(graph_id, requested_version)
395            .map_err(internal_error)?
396            .ok_or_else(|| invalid_params("requested graph version was not found"))?;
397        let normalized: Value = serde_json::from_str(&serialized)
398            .map_err(|e| internal_error(format!("stored graph version JSON error: {e}")))?;
399        let spec = parse_and_validate(&normalized)
400            .map_err(|e| internal_error(format!("stored graph version validation error: {e}")))?;
401        let canonical = serde_json::to_value(&spec).map_err(|e| internal_error(e.to_string()))?;
402        let actual_version = digest(&canonical);
403        if actual_version != requested_version {
404            return Err(internal_error(
405                "stored graph version digest does not match its normalized specification",
406            ));
407        }
408        Ok(RegisteredGraph {
409            warnings: spec.warnings(),
410            spec,
411            normalized: canonical,
412            version: actual_version,
413        })
414    }
415
416    fn mermaid(spec: &GraphSpec) -> String {
417        let mut s = String::from("graph TD\n");
418        for edge in &spec.edges {
419            s.push_str(&format!("  {} --> {}\n", edge.from, edge.to));
420        }
421        s
422    }
423
424    fn delete_registered_graph(&self, graph_id: &str) -> Result<Json<StructuredOutput>, ErrorData> {
425        let exists = self
426            .graphs
427            .lock()
428            .map_err(|e| internal_error(e.to_string()))?
429            .contains_key(graph_id);
430        if !exists {
431            return Ok(error_output(
432                format!("graph '{graph_id}' not found"),
433                "GRAPH_NOT_FOUND",
434            ));
435        }
436
437        if let Some(store) = &self.store {
438            match store.delete_graph(graph_id).map_err(internal_error)? {
439                GraphDeleteResult::Deleted => {}
440                GraphDeleteResult::Referenced => {
441                    return Ok(error_output(
442                        format!("graph '{graph_id}' is referenced by a durable execution"),
443                        "GRAPH_REFERENCED",
444                    ));
445                }
446                GraphDeleteResult::NotFound => {
447                    return Ok(error_output(
448                        format!(
449                            "graph '{graph_id}' is present in memory but missing from durable storage"
450                        ),
451                        "GRAPH_PERSISTENCE_MISMATCH",
452                    ));
453                }
454            }
455        }
456
457        self.graphs
458            .lock()
459            .map_err(|e| internal_error(e.to_string()))?
460            .remove(graph_id);
461        Ok(output_with_meta(
462            serde_json::json!({"status": "deleted"}),
463            Some(graph_id),
464            None,
465            None,
466        ))
467    }
468}
469
470#[cfg(test)]
471mod tests {
472    use super::*;
473    use crate::run_manager::RunManager;
474    use crate::store::PersistentStore;
475
476    fn configure_test_integrity_key() {
477        let path = std::env::temp_dir().join("agent-graph-mcp-unit-integrity.key");
478        std::fs::write(&path, [0x5au8; 32]).expect("test integrity key");
479        std::env::set_var("AGENT_GRAPH_INTEGRITY_KEY_PATH", path);
480    }
481
482    #[test]
483    fn terminal_projection_failure_rolls_back_sqlite_and_marks_run_volatile() {
484        configure_test_integrity_key();
485        let temp = tempfile::tempdir().expect("temp graph database");
486        let store = PersistentStore::open(temp.path()).expect("store");
487        let spec: GraphSpec = serde_json::from_value(serde_json::json!({
488            "name":"fault-injection",
489            "entry":"x",
490            "nodes":[{"id":"x","type":"passthrough"}],
491            "edges":[{"from":"x","to":"END"}]
492        }))
493        .expect("graph spec");
494        let spec_json = serde_json::to_string(&spec).expect("spec JSON");
495        store
496            .save_graph("fault-injection", &spec_json, "version", false)
497            .expect("graph");
498
499        let runs = RunManager::default();
500        let run_id = runs
501            .allocate("fault-injection", "version", serde_json::json!({"x":1}))
502            .expect("run");
503        store
504            .save_execution(
505                &run_id,
506                "fault-injection",
507                "version",
508                "running",
509                "{\"x\":1}",
510            )
511            .expect("execution");
512        runs.execute(
513            &run_id,
514            spec,
515            "http://localhost".into(),
516            "test-model".into(),
517        )
518        .expect("execution completes");
519
520        store.fail_terminal_projection_after_events();
521        AgentGraphServer::persist_terminal_and_mark(
522            runs.clone(),
523            Some(store.clone()),
524            runs.get(&run_id).expect("terminal record"),
525        );
526
527        let public = runs.get(&run_id).expect("volatile record").public();
528        assert_eq!(public["persistence_status"], "volatile_persistence_failed");
529        assert_eq!(public["storage_class"], "volatile");
530
531        let reopened = PersistentStore::open(temp.path()).expect("fresh store");
532        assert_eq!(
533            reopened.load_execution(&run_id).unwrap().unwrap()["status"],
534            "running"
535        );
536        assert!(reopened.load_events(&run_id, 0, 100).unwrap().is_none());
537        assert!(reopened.load_terminal_receipt(&run_id).unwrap().is_none());
538    }
539
540    #[test]
541    fn capacity_is_reserved_before_direct_or_approved_checkpoint_consumption() {
542        configure_test_integrity_key();
543        let temp = tempfile::tempdir().expect("checkpoint database");
544        let server = AgentGraphServer::new(
545            "http://localhost".into(),
546            "test-model".into(),
547            Some(temp.path().to_owned()),
548            None,
549        )
550        .expect("server");
551        server
552            .graph_create(Parameters(GraphCreateParams {
553                spec: Some(serde_json::json!({
554                    "name":"capacity-resume", "entry":"first",
555                    "nodes":[
556                        {"id":"first","type":"passthrough"},
557                        {"id":"second","type":"state_transform","config":{"operations":[{"op":"set","path":"done","value":true}]}}
558                    ],
559                    "edges":[{"from":"first","to":"second"},{"from":"second","to":"END"}]
560                })),
561                action: None,
562                graph_id: None,
563                idempotency_key: None,
564                template: None,
565                overwrite: None,
566            }))
567            .expect("create graph");
568        let checkpoint = |server: &AgentGraphServer| {
569            server
570                .graph_run_start(Parameters(RunStartParams {
571                    graph_id: "capacity-resume".into(),
572                    input: None,
573                    graph_version: None,
574                    thread_id: None,
575                    idempotency_key: None,
576                    budgets: None,
577                    checkpoint: Some(true),
578                }))
579                .expect("checkpoint start")
580                .0
581                .data
582                .unwrap()["checkpoint_id"]
583                .as_str()
584                .unwrap()
585                .to_owned()
586        };
587        let direct_checkpoint = checkpoint(&server);
588        {
589            let runs = server.runs.lock().expect("runs");
590            for index in 0..8 {
591                let run_id = runs
592                    .allocate("capacity", "v1", serde_json::json!({"index":index}))
593                    .expect("slot record");
594                runs.admit_async(&run_id).expect("slot admission");
595            }
596        }
597        let direct = server
598            .graph_run_resume(Parameters(RunResumeParams {
599                checkpoint_id: Some(direct_checkpoint.clone()),
600                run_id: None,
601            }))
602            .expect("resume response");
603        assert_eq!(direct.0.error_code.as_deref(), Some("RUN_CAPACITY"));
604        let store = server.store.as_ref().expect("store");
605        assert!(store
606            .load_resume_checkpoint(Some(&direct_checkpoint), None)
607            .expect("checkpoint")
608            .expect("record")
609            .consumed_at
610            .is_none());
611
612        let approval_checkpoint = checkpoint(&server);
613        let approval = server
614            .graph_approval_request(Parameters(ApprovalRequestParams {
615                checkpoint_id: approval_checkpoint.clone(),
616                audience: "operator".into(),
617                prompt: "approve after capacity is available".into(),
618                allowed_decisions: vec!["approve".into()],
619                expiration: (Utc::now() + chrono::Duration::hours(1)).to_rfc3339(),
620            }))
621            .expect("approval request");
622        let approval_id = approval.0.data.unwrap()["approval_id"]
623            .as_str()
624            .unwrap()
625            .to_owned();
626        let decided = server
627            .graph_approval_decide(Parameters(ApprovalDecideParams {
628                approval_id: approval_id.clone(),
629                decision: "approve".into(),
630                claimed_actor_label: "operator".into(),
631            }))
632            .expect("approval response");
633        assert_eq!(
634            decided.0.error_code.as_deref(),
635            Some("AUTHENTICATED_OPERATOR_REQUIRED")
636        );
637        assert_eq!(
638            store
639                .get_checkpoint_approval(&approval_id)
640                .expect("approval")
641                .expect("approval row")
642                .status,
643            "pending"
644        );
645        assert!(store
646            .load_resume_checkpoint(Some(&approval_checkpoint), None)
647            .expect("checkpoint")
648            .expect("record")
649            .consumed_at
650            .is_none());
651    }
652}
653
654#[tool_router]
655impl AgentGraphServer {
656    // ── graph_create ──────────────────────────────────────────────────
657
658    #[tool(
659        description = "Create, validate, or delete a graph-orchestrated workflow from a JSON spec. Supports template instantiation and idempotency keys."
660    )]
661    fn graph_create(
662        &self,
663        Parameters(GraphCreateParams {
664            spec,
665            action,
666            graph_id,
667            idempotency_key,
668            template,
669            overwrite,
670        }): Parameters<GraphCreateParams>,
671    ) -> Result<Json<StructuredOutput>, ErrorData> {
672        let action = action.as_deref().unwrap_or("create");
673
674        let request_digest = digest(&serde_json::json!({
675            "operation": "graph_create",
676            "action": action,
677            "spec": spec.as_ref().map(canonical_request_value).unwrap_or(Value::Null),
678            "template": template.as_ref().map(canonical_request_value).unwrap_or(Value::Null),
679            "graph_id": graph_id,
680            "overwrite": overwrite.unwrap_or(false),
681        }));
682        if action != "delete" {
683            if let Some(cached) = check_idempotency(
684                self.store.as_ref(),
685                idempotency_key.as_deref(),
686                &request_digest,
687            )? {
688                return Ok(cached);
689            }
690        }
691
692        // ── delete ──
693        if action == "delete" {
694            let id = graph_id
695                .as_deref()
696                .ok_or_else(|| invalid_params("missing graph_id for delete action"))?;
697            return self.delete_registered_graph(id);
698        }
699
700        if action != "create" && action != "validate" {
701            return Ok(error_output(
702                format!("unsupported graph_create action '{action}'"),
703                "INVALID_ACTION",
704            ));
705        }
706
707        // ── create / validate ──
708        let raw = if let Some(ref tpl) = template {
709            let tpl_val = if let Value::String(s) = tpl {
710                serde_json::from_str(s).unwrap_or_else(|_| tpl.clone())
711            } else {
712                tpl.clone()
713            };
714            let tpl_id = tpl_val
715                .get("id")
716                .and_then(Value::as_str)
717                .ok_or_else(|| invalid_params("template.id required"))?;
718            let tpl_name = tpl_val
719                .get("name")
720                .and_then(Value::as_str)
721                .or_else(|| graph_id.as_deref())
722                .unwrap_or(tpl_id);
723            templates::instantiate(tpl_id, tpl_name)
724                .map_err(|e| internal_error(format!("template error: {e}")))?
725        } else {
726            let spec = spec
727                .clone()
728                .ok_or_else(|| invalid_params("missing spec for create/validate"))?;
729            if let Value::String(s) = spec {
730                serde_json::from_str(&s)
731                    .map_err(|e| invalid_params(format!("spec string parse error: {e}")))?
732            } else {
733                spec
734            }
735        };
736
737        let original_version = raw
738            .get("spec_version")
739            .and_then(Value::as_str)
740            .unwrap_or("1")
741            .to_owned();
742        let warnings_preview = serde_json::from_value::<GraphSpec>(raw.clone())
743            .ok()
744            .map(|s| s.warnings())
745            .unwrap_or_default();
746        let spec_parsed =
747            parse_and_validate(&raw).map_err(|e| invalid_params(format!("invalid spec: {e}")))?;
748        if let Some(node) = spec_parsed
749            .nodes
750            .iter()
751            .find(|node| crate::spec::GraphSpec::executable_node_type(&node.node_type).is_err())
752        {
753            return Ok(error_output(
754                format!("node '{}' declares an unsupported executable type", node.id),
755                "UNSUPPORTED_NODE_TYPE",
756            ));
757        }
758        let normalized =
759            serde_json::to_value(&spec_parsed).map_err(|e| internal_error(e.to_string()))?;
760        let version = digest(&normalized);
761        let warnings = if original_version == "1" {
762            warnings_preview
763        } else {
764            spec_parsed.warnings()
765        };
766
767        if action == "validate" {
768            let output = output_with_meta(
769                serde_json::json!({
770                    "graph_id": spec_parsed.name,
771                    "graph_version": version,
772                    "digest": version,
773                    "normalized_spec_version": "2",
774                    "warnings": warnings,
775                    "storage_class": "volatile",
776                    "status": "valid"
777                }),
778                Some(&spec_parsed.name),
779                Some(&version),
780                None,
781            );
782            if let Some(ref store) = self.store {
783                if let Some(idem) = idempotency_key {
784                    if let Some(cached) =
785                        persist_idempotency(store, &idem, &request_digest, &output)?
786                    {
787                        return Ok(cached);
788                    }
789                }
790            }
791            return Ok(output);
792        }
793
794        if Self::graph_requires_witness_store(&spec_parsed) && self.store.is_none() {
795            return Ok(error_output(
796                "evidence-required graphs require SQLite witness persistence",
797                "WITNESS_STORE_REQUIRED",
798            ));
799        }
800
801        // ── register ──
802        let mut graphs = self
803            .graphs
804            .lock()
805            .map_err(|e| internal_error(e.to_string()))?;
806        let name = spec_parsed.name.clone();
807        let overwrite = overwrite.unwrap_or(false);
808        if !overwrite && !graphs.contains_key(&name) && graphs.len() >= MAX_GRAPHS {
809            return Ok(error_output(
810                format!("graph limit ({MAX_GRAPHS}) reached"),
811                "LIMIT_EXCEEDED",
812            ));
813        }
814
815        let id = name.clone();
816        if let Some(ref store) = self.store {
817            let spec_str = serde_json::to_string(&normalized).unwrap_or_default();
818            if let Err(error) = store.save_graph(&id, &spec_str, &version, overwrite) {
819                return Ok(error_output(error, "GRAPH_VERSION_CONFLICT"));
820            }
821        }
822        graphs.insert(
823            id.clone(),
824            RegisteredGraph {
825                spec: spec_parsed,
826                normalized: normalized.clone(),
827                version: version.clone(),
828                warnings: warnings.clone(),
829            },
830        );
831        drop(graphs);
832
833        let output = output_with_meta(
834            serde_json::json!({
835                "graph_id": id,
836                "graph_version": version,
837                "digest": version,
838                "normalized_spec_version": "2",
839                "warnings": warnings,
840                "storage_class": "volatile",
841                "status": "created"
842            }),
843            Some(&id),
844            Some(&version),
845            None,
846        );
847
848        if let Some(ref store) = self.store {
849            if let Some(idem) = idempotency_key {
850                if let Some(cached) = persist_idempotency(store, &idem, &request_digest, &output)? {
851                    return Ok(cached);
852                }
853            }
854        }
855        Ok(output)
856    }
857
858    // ── graph_execute ─────────────────────────────────────────────────
859
860    #[tool(
861        description = "Execute a registered graph. Sync mode blocks until completion; async mode returns immediately with a run_id."
862    )]
863    fn graph_execute(
864        &self,
865        Parameters(GraphExecuteParams {
866            graph_id,
867            input,
868            graph_version,
869            thread_id,
870            mode,
871            idempotency_key,
872        }): Parameters<GraphExecuteParams>,
873    ) -> Result<Json<StructuredOutput>, ErrorData> {
874        let input = input.unwrap_or(Value::Null);
875        ensure_size(&input, MAX_INPUT_BYTES, "execution input").map_err(|e| invalid_params(e))?;
876
877        let graph = self.resolve_graph(&graph_id, graph_version.as_deref())?;
878
879        if Self::graph_requires_witness_store(&graph.spec) && self.store.is_none() {
880            return Ok(error_output(
881                "evidence-required graphs require SQLite witness persistence",
882                "WITNESS_STORE_REQUIRED",
883            ));
884        }
885
886        let request_digest = digest(&serde_json::json!({
887            "operation": "graph_execute",
888            "graph_id": graph_id,
889            "graph_spec": graph.normalized,
890            "graph_version": graph.version,
891            "input": input,
892            "mode": mode.clone().unwrap_or_else(|| "sync".into()),
893            "thread_id": thread_id,
894        }));
895
896        let runs = self
897            .runs
898            .lock()
899            .map_err(|e| internal_error(e.to_string()))?;
900        if let Some(idem) = idempotency_key.as_deref() {
901            if let Some(cached) =
902                check_idempotency(self.store.as_ref(), Some(idem), &request_digest)?
903            {
904                return Ok(cached);
905            }
906        }
907
908        let run_id = runs
909            .allocate(&graph_id, &graph.version, input.clone())
910            .map_err(|e| internal_error(e))?;
911
912        if let Err(e) = runs.admit_async(&run_id) {
913            runs.remove(&run_id);
914            return Ok(error_output(e, "RUN_CAPACITY"));
915        }
916        if let Some(ref store) = self.store {
917            let _ = store.save_execution(
918                &run_id,
919                &graph_id,
920                &graph.version,
921                "running",
922                &input.to_string(),
923            );
924        }
925
926        let is_async = mode.as_deref() == Some("async");
927        if is_async {
928            let terminal_store = self.store.clone();
929            let completion_runs = runs.clone();
930            runs.start_with_completion_with_store(
931                run_id.clone(),
932                graph.spec,
933                self.base_url.clone(),
934                self.default_model.clone(),
935                self.store.clone(),
936                move |record| {
937                    Self::persist_terminal_and_mark(completion_runs, terminal_store, record)
938                },
939            );
940            let output = output_with_meta(
941                serde_json::json!({
942                    "run_id": run_id,
943                    "status": "accepted",
944                    "thread_id": thread_id,
945                    "storage_class": "volatile",
946                    "cancellation": "provider_future_best_effort_drop; underlying_request_may_continue"
947                }),
948                Some(&graph_id),
949                Some(&graph.version),
950                Some(&run_id),
951            );
952            if let Some(ref store) = self.store {
953                if let Some(idem) = idempotency_key {
954                    if let Some(cached) =
955                        persist_idempotency(store, &idem, &request_digest, &output)?
956                    {
957                        return Ok(cached);
958                    }
959                }
960            }
961            return Ok(output);
962        }
963
964        let terminal_store = self.store.clone();
965        let completion_runs = runs.clone();
966        runs.start_with_completion_with_store(
967            run_id.clone(),
968            graph.spec,
969            self.base_url.clone(),
970            self.default_model.clone(),
971            self.store.clone(),
972            move |record| Self::persist_terminal_and_mark(completion_runs, terminal_store, record),
973        );
974
975        let deadline = Instant::now() + Duration::from_millis(300_000);
976        let output = loop {
977            let r = runs
978                .get(&run_id)
979                .ok_or_else(|| internal_error(format!("run '{run_id}' not found")))?;
980            if matches!(r.status.as_str(), "completed" | "failed" | "cancelled") {
981                break output_with_meta(
982                    r.public(),
983                    Some(&graph_id),
984                    Some(&graph.version),
985                    Some(&run_id),
986                );
987            }
988            if Instant::now() >= deadline {
989                let cancellation = runs.cancel(&run_id).unwrap_or_else(
990                    |_| serde_json::json!({"run_id": run_id, "status": "cancellation_requested"}),
991                );
992                break output_with_meta(
993                    serde_json::json!({
994                        "run_id": run_id,
995                        "status": r.status,
996                        "timed_out": true,
997                        "completion_unknown": true,
998                        "cancellation": "requested",
999                        "cancellation_result": cancellation,
1000                    }),
1001                    Some(&graph_id),
1002                    Some(&graph.version),
1003                    Some(&run_id),
1004                );
1005            }
1006            std::thread::sleep(Duration::from_millis(100));
1007        };
1008
1009        if let Some(ref store) = self.store {
1010            let status = output
1011                .0
1012                .data
1013                .as_ref()
1014                .and_then(|data| data.get("status").and_then(Value::as_str))
1015                .unwrap_or("failed");
1016            let final_state = output
1017                .0
1018                .data
1019                .as_ref()
1020                .and_then(|data| data.get("final_state").cloned())
1021                .map(|v| serde_json::to_string(&v).unwrap_or_default());
1022            let _ = store.save_execution(
1023                &run_id,
1024                &graph_id,
1025                &graph.version,
1026                status,
1027                &input.to_string(),
1028            );
1029            let _ =
1030                store.update_execution_status(&run_id, status, final_state.as_deref(), None, None);
1031        }
1032
1033        if let Some(ref store) = self.store {
1034            if let Some(idem) = idempotency_key {
1035                if let Some(cached) = persist_idempotency(store, &idem, &request_digest, &output)? {
1036                    return Ok(cached);
1037                }
1038            }
1039        }
1040        Ok(output)
1041    }
1042
1043    // ── Local source witness capture ─────────────────────────────────
1044
1045    #[tool(
1046        description = "Persist caller-supplied UTF-8 source content as a local witness receipt. The locator is metadata only; this tool never fetches or verifies it."
1047    )]
1048    fn graph_source_witness_capture(
1049        &self,
1050        Parameters(WitnessCaptureParams {
1051            locator,
1052            content,
1053            media_type,
1054            authority_class,
1055            retrieved_at,
1056        }): Parameters<WitnessCaptureParams>,
1057    ) -> Result<Json<StructuredOutput>, ErrorData> {
1058        let capture = WitnessCapture {
1059            locator,
1060            content,
1061            media_type,
1062            authority_class,
1063            retrieved_at: retrieved_at
1064                .unwrap_or_else(|| Utc::now().to_rfc3339_opts(SecondsFormat::Nanos, true)),
1065        };
1066        if let Err(error) = validate_witness_capture(capture.clone()) {
1067            return Ok(Self::witness_error_output(error));
1068        }
1069        let Some(store) = self.store.as_ref() else {
1070            return Ok(error_output(
1071                "SQLite persistence is required for source witness capture",
1072                "WITNESS_STORE_REQUIRED",
1073            ));
1074        };
1075        match store.capture_witness(capture) {
1076            Ok(record) => Ok(structured_output(serde_json::json!({
1077                "witness_id": record.witness_id,
1078                "digest": record.digest,
1079                "locator_digest": digest(&Value::String(record.locator)),
1080                "media_type": record.media_type,
1081                "authority_class": record.authority_class,
1082                "retrieved_at": record.retrieved_at,
1083                "content_bytes": record.content.len(),
1084                "storage_class": "sqlite_source_witness"
1085            }))),
1086            Err(error) => Ok(Self::witness_error_output(error)),
1087        }
1088    }
1089
1090    #[tool(
1091        description = "Read one exact local source witness ID, verifying its HMAC-SHA256 authentication tag before returning metadata and captured content."
1092    )]
1093    fn graph_source_witness_get(
1094        &self,
1095        Parameters(WitnessGetParams { witness_id }): Parameters<WitnessGetParams>,
1096    ) -> Result<Json<StructuredOutput>, ErrorData> {
1097        let Some(store) = self.store.as_ref() else {
1098            return Ok(error_output(
1099                "SQLite persistence is required for source witness reads",
1100                "WITNESS_STORE_REQUIRED",
1101            ));
1102        };
1103        match store.get_witness(&witness_id) {
1104            Ok(Some(record)) => {
1105                let locator_digest = digest(&Value::String(record.locator.clone()));
1106                Ok(structured_output(serde_json::json!({
1107                    "witness_id": record.witness_id,
1108                    "digest": record.digest,
1109                    "locator": record.locator,
1110                    "locator_digest": locator_digest,
1111                    "content": record.content,
1112                    "media_type": record.media_type,
1113                    "authority_class": record.authority_class,
1114                    "retrieved_at": record.retrieved_at,
1115                    "storage_class": "sqlite_source_witness"
1116                })))
1117            }
1118            Ok(None) => Ok(error_output(
1119                "source witness was not found",
1120                "WITNESS_NOT_FOUND",
1121            )),
1122            Err(error) => Ok(Self::witness_error_output(error)),
1123        }
1124    }
1125
1126    // ── graph_status ──────────────────────────────────────────────────
1127
1128    #[tool(
1129        description = "Query server state, graph details, run status, events, receipts, or templates."
1130    )]
1131    fn graph_status(
1132        &self,
1133        Parameters(GraphStatusParams {
1134            resource,
1135            graph_id,
1136            run_id,
1137            cursor,
1138            limit,
1139        }): Parameters<GraphStatusParams>,
1140    ) -> Result<Json<StructuredOutput>, ErrorData> {
1141        let resource = resource.as_deref();
1142
1143        // Server-level summary (no resource or resource="server")
1144        if resource.is_none() || resource == Some("server") {
1145            let graphs = self
1146                .graphs
1147                .lock()
1148                .map_err(|e| internal_error(e.to_string()))?;
1149            let graph_names: Vec<&String> = graphs.keys().collect();
1150            let runs = self
1151                .runs
1152                .lock()
1153                .map_err(|e| internal_error(e.to_string()))?;
1154            let run_ids = runs.list();
1155            let durable_integrity = self
1156                .store
1157                .as_ref()
1158                .is_some_and(PersistentStore::has_integrity_key);
1159
1160            return Ok(structured_output(serde_json::json!({
1161                "graphs": graph_names,
1162                "graph_count": graphs.len(),
1163                "execution_count": run_ids.len(),
1164                "retained_execution_count": run_ids.len(),
1165                "total_execution_count": run_ids.len(),
1166                "base_url": self.safe_provider_label(),
1167                "default_model": self.default_model,
1168                "storage_class": if self.store.is_none() {
1169                    "process_local"
1170                } else if durable_integrity {
1171                    "persisted_integrity_verified"
1172                } else {
1173                    "persisted_unverified"
1174                },
1175                "capabilities": {
1176                    "runtime": "agent_graph",
1177                    "async_start": true,
1178                    "cancellation": "provider_future_best_effort_drop; underlying_request_may_continue",
1179                    "durable_resume": if durable_integrity {
1180                        Value::String("deterministic_local_resume_only".into())
1181                    } else {
1182                        Value::Bool(false)
1183                    },
1184                    "terminal_persistence": if durable_integrity { "sqlite_projection_only" } else { "disabled_without_integrity_key" },
1185                    "checkpointing": if durable_integrity { "deterministic_local_pre_execution" } else { "unavailable" },
1186                    "events": if self.store.is_some() { "terminal_persisted_projection_with_sqlite_fallback" } else { "volatile_in_memory_only" },
1187                    "event_replay": "not_replayable_execution",
1188                    "restart_recovery": "interrupted_non_resumable",
1189                    "budgets": {
1190                        "max_wall_clock_ms": "enforced",
1191                        "max_nodes": "enforced_at_engine_superstep_boundary",
1192                        "max_llm_calls": "rejected_INVALID_BUDGETS_no_invocation_hook"
1193                    },
1194                    "state_write_conflicts": "rejected_without_explicit_reducer",
1195                    "evidence": "witness_bound_local_capture_only; locators_not_fetched; source_authority_not_verified",
1196                    "evidence_authority": "caller_supplied_unverified_or_local_primary_capture",
1197                    "hitl": if durable_integrity { "checkpoint_bound_durable_approval_only" } else { "unavailable" },
1198                    "replay": "integrity_only"
1199                },
1200                "limits": {"graphs": MAX_GRAPHS}
1201            })));
1202        }
1203
1204        match resource.unwrap() {
1205            "templates" => Ok(structured_output(templates::list())),
1206
1207            "graph" => {
1208                let id = graph_id
1209                    .as_deref()
1210                    .ok_or_else(|| invalid_params("missing graph_id"))?;
1211                let graphs = self
1212                    .graphs
1213                    .lock()
1214                    .map_err(|e| internal_error(e.to_string()))?;
1215                let g = graphs
1216                    .get(id)
1217                    .ok_or_else(|| invalid_params(format!("graph '{id}' not found")))?;
1218                Ok(output_with_meta(
1219                    serde_json::json!({
1220                        "graph_id": id,
1221                        "graph_version": g.version,
1222                        "normalized_spec": g.normalized,
1223                        "mermaid": Self::mermaid(&g.spec),
1224                        "warnings": g.warnings,
1225                        "storage_class": "volatile"
1226                    }),
1227                    Some(id),
1228                    Some(&g.version),
1229                    None,
1230                ))
1231            }
1232
1233            "run" => {
1234                let runs = self
1235                    .runs
1236                    .lock()
1237                    .map_err(|e| internal_error(e.to_string()))?;
1238                if run_id.is_none() {
1239                    // List all runs
1240                    return Ok(structured_output(serde_json::json!({
1241                        "runs": runs.list()
1242                    })));
1243                }
1244                let id = run_id.as_deref().unwrap();
1245                let r = runs
1246                    .get(id)
1247                    .ok_or_else(|| invalid_params(format!("run '{id}' not found")))?;
1248                Ok(structured_output(r.public()))
1249            }
1250
1251            "events" => {
1252                let id = run_id
1253                    .as_deref()
1254                    .ok_or_else(|| invalid_params("missing run_id for events"))?;
1255                let runs = self
1256                    .runs
1257                    .lock()
1258                    .map_err(|e| internal_error(e.to_string()))?;
1259                let cursor_val = cursor.unwrap_or(0);
1260                let limit_val = limit.unwrap_or(100) as usize;
1261                let result = runs
1262                    .events(self.store.as_ref(), id, cursor_val, limit_val)
1263                    .map_err(|e| invalid_params(e))?;
1264                Ok(output_with_meta(result, None, None, Some(id)))
1265            }
1266
1267            "receipt" => {
1268                let id = run_id
1269                    .as_deref()
1270                    .ok_or_else(|| invalid_params("missing run_id for receipt"))?;
1271                let runs = self
1272                    .runs
1273                    .lock()
1274                    .map_err(|e| internal_error(e.to_string()))?;
1275                let r = runs
1276                    .get(id)
1277                    .ok_or_else(|| invalid_params(format!("run '{id}' not found")))?;
1278                Ok(output_with_meta(r.receipt.clone(), None, None, Some(id)))
1279            }
1280
1281            "bundle" => {
1282                let id = run_id
1283                    .as_deref()
1284                    .ok_or_else(|| invalid_params("missing run_id for bundle"))?;
1285                let runs = self
1286                    .runs
1287                    .lock()
1288                    .map_err(|e| internal_error(e.to_string()))?;
1289                let r = runs
1290                    .get(id)
1291                    .ok_or_else(|| invalid_params(format!("run '{id}' not found")))?;
1292                Ok(output_with_meta(r.bundle.clone(), None, None, Some(id)))
1293            }
1294
1295            _ => Ok(error_output(
1296                format!("unknown status resource '{}'", resource.unwrap_or("")),
1297                "INVALID_RESOURCE",
1298            )),
1299        }
1300    }
1301
1302    // ── graph_list (NEW) ──────────────────────────────────────────────
1303
1304    #[tool(
1305        description = "List all registered graphs with metadata (name, node count, edge count, version)."
1306    )]
1307    fn graph_list(
1308        &self,
1309        Parameters(GraphListParams { query, limit }): Parameters<GraphListParams>,
1310    ) -> Result<Json<StructuredOutput>, ErrorData> {
1311        let graphs = self
1312            .graphs
1313            .lock()
1314            .map_err(|e| internal_error(e.to_string()))?;
1315
1316        let mut entries: Vec<Value> = graphs
1317            .iter()
1318            .filter(|(name, _)| {
1319                query
1320                    .as_ref()
1321                    .map(|q| name.contains(q.as_str()))
1322                    .unwrap_or(true)
1323            })
1324            .take(limit.unwrap_or(50) as usize)
1325            .map(|(name, g)| {
1326                let version_history = self
1327                    .store
1328                    .as_ref()
1329                    .and_then(|store| store.list_graph_versions(name).ok())
1330                    .unwrap_or_else(|| vec![g.version.clone()]);
1331                serde_json::json!({
1332                    "name": name,
1333                    "version": g.version,
1334                    "current_version": g.version,
1335                    "version_history": version_history,
1336                    "historical_specs": self.store.is_some(),
1337                    "node_count": g.spec.nodes.len(),
1338                    "edge_count": g.spec.edges.len(),
1339                    "entry": g.spec.entry,
1340                    "warnings": g.warnings,
1341                })
1342            })
1343            .collect();
1344
1345        entries.sort_by(|a, b| {
1346            a.get("name")
1347                .and_then(Value::as_str)
1348                .cmp(&b.get("name").and_then(Value::as_str))
1349        });
1350
1351        Ok(structured_output(serde_json::json!({
1352            "graphs": entries,
1353            "count": entries.len(),
1354        })))
1355    }
1356
1357    // ── graph_delete (NEW) ────────────────────────────────────────────
1358
1359    #[allow(dead_code)]
1360    fn graph_delete(
1361        &self,
1362        Parameters(GraphDeleteParams { graph_id }): Parameters<GraphDeleteParams>,
1363    ) -> Result<Json<StructuredOutput>, ErrorData> {
1364        self.delete_registered_graph(&graph_id)
1365    }
1366
1367    // ── graph_inspect (NEW) ───────────────────────────────────────────
1368
1369    #[tool(
1370        description = "Get a graph's full topology: nodes, edges, Mermaid diagram, and topology hash."
1371    )]
1372    fn graph_inspect(
1373        &self,
1374        Parameters(GraphInspectParams { graph_id }): Parameters<GraphInspectParams>,
1375    ) -> Result<Json<StructuredOutput>, ErrorData> {
1376        let graphs = self
1377            .graphs
1378            .lock()
1379            .map_err(|e| internal_error(e.to_string()))?;
1380        let g = graphs
1381            .get(&graph_id)
1382            .ok_or_else(|| invalid_params(format!("graph '{graph_id}' not found")))?;
1383
1384        let nodes: Vec<Value> = g
1385            .spec
1386            .nodes
1387            .iter()
1388            .map(|n| {
1389                serde_json::json!({
1390                    "id": n.id,
1391                    "type": n.node_type,
1392                    "config": n.config,
1393                })
1394            })
1395            .collect();
1396
1397        let edges: Vec<Value> = g
1398            .spec
1399            .edges
1400            .iter()
1401            .map(|e| {
1402                serde_json::json!({
1403                    "from": e.from,
1404                    "to": e.to,
1405                })
1406            })
1407            .collect();
1408
1409        Ok(output_with_meta(
1410            serde_json::json!({
1411                "name": graph_id,
1412                "version": g.version,
1413                "current_version": g.version,
1414                "version_history": self.store.as_ref().and_then(|store| store.list_graph_versions(&graph_id).ok()).unwrap_or_else(|| vec![g.version.clone()]),
1415                "historical_specs": self.store.is_some(),
1416                "entry": g.spec.entry,
1417                "max_iterations": g.spec.max_iterations,
1418                "max_parallelism": g.spec.max_parallelism,
1419                "nodes": nodes,
1420                "node_count": nodes.len(),
1421                "edges": edges,
1422                "edge_count": edges.len(),
1423                "mermaid": Self::mermaid(&g.spec),
1424                "topology_hash": g.version,
1425                "reducers": g.spec.reducers,
1426                "warnings": g.warnings,
1427            }),
1428            Some(&graph_id),
1429            Some(&g.version),
1430            None,
1431        ))
1432    }
1433
1434    // ── Approval lifecycle ────────────────────────────────────────────
1435
1436    fn validate_resume_checkpoint(
1437        &self,
1438        store: &PersistentStore,
1439        checkpoint: &CheckpointRecord,
1440    ) -> Result<
1441        (
1442            crate::store::ExecutionContract,
1443            RegisteredGraph,
1444            Option<RunBudgets>,
1445        ),
1446        (String, String),
1447    > {
1448        let contract = store
1449            .load_execution_contract(&checkpoint.run_id)
1450            .map_err(|error| (error, "CHECKPOINT_PERSISTENCE_FAILURE".into()))?
1451            .ok_or_else(|| {
1452                (
1453                    "checkpoint execution contract was not found".into(),
1454                    "CHECKPOINT_INTEGRITY_FAILURE".into(),
1455                )
1456            })?;
1457        if contract.graph_id != checkpoint.graph_id
1458            || contract.graph_version != checkpoint.graph_version
1459            || checkpoint.terminal_cursor != 0
1460            || checkpoint.event_cursor != 0
1461        {
1462            return Err((
1463                "checkpoint integrity validation failed".into(),
1464                "CHECKPOINT_INTEGRITY_FAILURE".into(),
1465            ));
1466        }
1467        let graph = self
1468            .resolve_graph(&checkpoint.graph_id, Some(&checkpoint.graph_version))
1469            .map_err(|_| {
1470                (
1471                    "checkpoint graph version is unavailable".into(),
1472                    "CHECKPOINT_INTEGRITY_FAILURE".into(),
1473                )
1474            })?;
1475        let eligibility = graph.spec.resume_eligibility().map_err(|_| {
1476            (
1477                "checkpoint graph is no longer in the deterministic local resume subset".into(),
1478                "RESUME_INELIGIBLE".into(),
1479            )
1480        })?;
1481        if graph.version != checkpoint.graph_version
1482            || checkpoint.next_node_cursor != eligibility.next_node_cursor
1483            || checkpoint.dependency_summary != eligibility.dependency_summary
1484            || checkpoint.dependency_digest != digest(&eligibility.dependency_summary)
1485            || checkpoint.state != initial_state_for_input(&contract.input)
1486            || checkpoint.budgets != contract.budgets
1487            || checkpoint.budget_counters
1488                != serde_json::json!({"nodes":0,"llm_calls":0,"wall_clock_ms":0})
1489        {
1490            return Err((
1491                "checkpoint integrity validation failed".into(),
1492                "CHECKPOINT_INTEGRITY_FAILURE".into(),
1493            ));
1494        }
1495        let budgets = RunBudgets::parse(Some(&checkpoint.budgets)).map_err(|_| {
1496            (
1497                "checkpoint budgets failed validation".into(),
1498                "CHECKPOINT_INTEGRITY_FAILURE".into(),
1499            )
1500        })?;
1501        Ok((contract, graph, budgets))
1502    }
1503
1504    fn launch_resumed(
1505        &self,
1506        checkpoint: CheckpointRecord,
1507        contract: crate::store::ExecutionContract,
1508        graph: RegisteredGraph,
1509        budgets: Option<RunBudgets>,
1510        approval: Option<Value>,
1511    ) -> Result<Json<StructuredOutput>, ErrorData> {
1512        let runs = self
1513            .runs
1514            .lock()
1515            .map_err(|e| internal_error(e.to_string()))?;
1516        if runs.get(&checkpoint.run_id).is_some() {
1517            let _ = runs.remove(&checkpoint.run_id);
1518        }
1519        let run_id = match runs.allocate_resumed(
1520            &checkpoint.run_id,
1521            &checkpoint.graph_id,
1522            &checkpoint.graph_version,
1523            contract.input,
1524            checkpoint.state.clone(),
1525            budgets,
1526            &checkpoint.checkpoint_id,
1527            &checkpoint.checkpoint_digest,
1528            approval.clone(),
1529        ) {
1530            Ok(run_id) => run_id,
1531            Err(error) => {
1532                runs.release_async_slot();
1533                return Ok(error_output(error, "RUN_CAPACITY"));
1534            }
1535        };
1536        if let Err(error) = runs.admit_reserved_async(&run_id) {
1537            runs.remove(&run_id);
1538            runs.release_async_slot();
1539            return Ok(error_output(error, "RUN_CAPACITY"));
1540        }
1541        self.store
1542            .as_ref()
1543            .expect("resumed launch requires SQLite")
1544            .update_execution_status(&run_id, "running", None, None, None)
1545            .map_err(internal_error)?;
1546        let terminal_store = self.store.clone();
1547        let completion_runs = runs.clone();
1548        runs.start_resumed_with_completion(
1549            run_id.clone(),
1550            graph.spec,
1551            self.base_url.clone(),
1552            self.default_model.clone(),
1553            self.store.clone(),
1554            move |record| Self::persist_terminal_and_mark(completion_runs, terminal_store, record),
1555        );
1556        Ok(output_with_meta(
1557            serde_json::json!({
1558                "run_id": run_id,
1559                "status": "running",
1560                "checkpoint": checkpoint_value(&checkpoint),
1561                "resume_capability": "deterministic_local_resume",
1562                "approval": approval,
1563            }),
1564            Some(&checkpoint.graph_id),
1565            Some(&checkpoint.graph_version),
1566            Some(&run_id),
1567        ))
1568    }
1569
1570    #[tool(
1571        description = "Create a durable approval request bound to one unconsumed deterministic-local checkpoint."
1572    )]
1573    fn graph_approval_request(
1574        &self,
1575        Parameters(ApprovalRequestParams {
1576            checkpoint_id,
1577            audience,
1578            prompt,
1579            allowed_decisions,
1580            expiration,
1581        }): Parameters<ApprovalRequestParams>,
1582    ) -> Result<Json<StructuredOutput>, ErrorData> {
1583        let Some(store) = self.store.as_ref() else {
1584            return Ok(error_output(
1585                "SQLite persistence is required for durable approvals",
1586                "APPROVAL_STORE_REQUIRED",
1587            ));
1588        };
1589        if audience.trim().is_empty() || audience.len() > 256 {
1590            return Ok(error_output(
1591                "audience must be non-empty and at most 256 bytes",
1592                "INVALID_PARAMS",
1593            ));
1594        }
1595        if allowed_decisions.is_empty()
1596            || allowed_decisions
1597                .iter()
1598                .any(|decision| !matches!(decision.as_str(), "approve" | "reject"))
1599        {
1600            return Ok(error_output(
1601                "allowed_decisions must be a non-empty subset of approve and reject",
1602                "INVALID_PARAMS",
1603            ));
1604        }
1605        if chrono::DateTime::parse_from_rfc3339(&expiration).is_err() {
1606            return Ok(error_output("expiration must be RFC3339", "INVALID_PARAMS"));
1607        }
1608        if prompt.len() > 16 * 1024 {
1609            return Ok(error_output(
1610                "prompt exceeds the bounded approval prompt size",
1611                "INVALID_PARAMS",
1612            ));
1613        }
1614        let checkpoint = match store.load_resume_checkpoint(Some(&checkpoint_id), None) {
1615            Ok(Some(checkpoint)) => checkpoint,
1616            Ok(None) => return Ok(checkpoint_error_output(CheckpointError::NotFound)),
1617            Err(error) => return Ok(checkpoint_error_output(error)),
1618        };
1619        if checkpoint.consumed_at.is_some() {
1620            return Ok(checkpoint_error_output(CheckpointError::Consumed));
1621        }
1622        if let Err((message, code)) = self.validate_resume_checkpoint(store, &checkpoint) {
1623            return Ok(error_output(message, code));
1624        }
1625        let prompt_digest = digest(&Value::String(prompt));
1626        let approval = match store.create_checkpoint_approval(
1627            &checkpoint.checkpoint_id,
1628            &checkpoint.graph_id,
1629            &checkpoint.graph_version,
1630            &checkpoint.next_node_cursor,
1631            &checkpoint.state,
1632            &checkpoint.budgets,
1633            &checkpoint.budget_counters,
1634            &checkpoint.dependency_summary,
1635            &audience,
1636            &prompt_digest,
1637            &allowed_decisions,
1638            &expiration,
1639        ) {
1640            Ok(approval) => approval,
1641            Err(error) => return Ok(approval_error_output(error)),
1642        };
1643        Ok(output_with_meta(
1644            approval_value(&approval),
1645            Some(&approval.graph_id),
1646            Some(&approval.graph_version),
1647            Some(&approval.run_id),
1648        ))
1649    }
1650
1651    #[tool(
1652        description = "Read durable checkpoint-bound approval metadata from SQLite without raw prompt or checkpoint state."
1653    )]
1654    fn graph_approval_list(
1655        &self,
1656        Parameters(ApprovalListParams {
1657            run_id,
1658            status,
1659            limit,
1660        }): Parameters<ApprovalListParams>,
1661    ) -> Result<Json<StructuredOutput>, ErrorData> {
1662        let Some(store) = self.store.as_ref() else {
1663            return Ok(error_output(
1664                "SQLite persistence is required for durable approvals",
1665                "APPROVAL_STORE_REQUIRED",
1666            ));
1667        };
1668        let approvals = store
1669            .list_checkpoint_approvals(
1670                run_id.as_deref(),
1671                status.as_deref(),
1672                limit.unwrap_or(50) as usize,
1673            )
1674            .map_err(|error| internal_error(error.message()))?;
1675        Ok(structured_output(serde_json::json!({
1676            "approvals": approvals.iter().map(approval_value).collect::<Vec<_>>(),
1677            "count": approvals.len(),
1678            "storage_class": "sqlite_durable_approval_metadata",
1679        })))
1680    }
1681
1682    #[tool(
1683        description = "Read one durable checkpoint-bound approval's metadata from SQLite without raw prompt or checkpoint state."
1684    )]
1685    fn graph_approval_get(
1686        &self,
1687        Parameters(ApprovalGetParams { approval_id }): Parameters<ApprovalGetParams>,
1688    ) -> Result<Json<StructuredOutput>, ErrorData> {
1689        let Some(store) = self.store.as_ref() else {
1690            return Ok(error_output(
1691                "SQLite persistence is required for durable approvals",
1692                "APPROVAL_STORE_REQUIRED",
1693            ));
1694        };
1695        match store
1696            .get_checkpoint_approval(&approval_id)
1697            .map_err(|error| internal_error(error.message()))?
1698        {
1699            Some(approval) => Ok(output_with_meta(
1700                approval_value(&approval),
1701                Some(&approval.graph_id),
1702                Some(&approval.graph_version),
1703                Some(&approval.run_id),
1704            )),
1705            None => Ok(approval_error_output(ApprovalError::NotFound)),
1706        }
1707    }
1708
1709    #[allow(dead_code)]
1710    fn graph_approval_decide(
1711        &self,
1712        Parameters(ApprovalDecideParams {
1713            approval_id: _,
1714            decision: _,
1715            claimed_actor_label: _,
1716        }): Parameters<ApprovalDecideParams>,
1717    ) -> Result<Json<StructuredOutput>, ErrorData> {
1718        return Ok(error_output(
1719            "approval decisions require authenticated operator transport",
1720            "AUTHENTICATED_OPERATOR_REQUIRED",
1721        ));
1722    }
1723
1724    // ── Async run lifecycle ───────────────────────────────────────────
1725
1726    #[tool(
1727        description = "Start an async graph run. Returns run_id immediately; use graph_run_wait to block on completion. Optional budgets accept only positive integer max_wall_clock_ms or max_nodes fields; max_llm_calls is rejected until a real invocation hook exists."
1728    )]
1729    fn graph_run_start(
1730        &self,
1731        Parameters(RunStartParams {
1732            graph_id,
1733            input,
1734            graph_version,
1735            thread_id,
1736            idempotency_key,
1737            budgets,
1738            checkpoint,
1739        }): Parameters<RunStartParams>,
1740    ) -> Result<Json<StructuredOutput>, ErrorData> {
1741        let requested_budgets = match RunBudgets::parse(budgets.as_ref()) {
1742            Ok(budgets) => budgets,
1743            Err(error) => return Ok(error_output(error, "INVALID_BUDGETS")),
1744        };
1745        let input = input.unwrap_or(Value::Null);
1746        let checkpoint_requested = checkpoint.unwrap_or(false);
1747        ensure_size(&input, MAX_INPUT_BYTES, "execution input").map_err(|e| invalid_params(e))?;
1748
1749        let RegisteredGraph {
1750            spec,
1751            normalized,
1752            version,
1753            ..
1754        } = self.resolve_graph(&graph_id, graph_version.as_deref())?;
1755
1756        if Self::graph_requires_witness_store(&spec) && self.store.is_none() {
1757            return Ok(error_output(
1758                "evidence-required graphs require SQLite witness persistence",
1759                "WITNESS_STORE_REQUIRED",
1760            ));
1761        }
1762
1763        let eligibility = if checkpoint_requested {
1764            match spec.resume_eligibility() {
1765                Ok(eligibility) => Some(eligibility),
1766                Err(reason) => return Ok(error_output(reason, "RESUME_INELIGIBLE")),
1767            }
1768        } else {
1769            None
1770        };
1771
1772        let request_digest = digest(&serde_json::json!({
1773            "operation": "graph_run_start",
1774            "graph_id": graph_id,
1775            "graph_spec": normalized,
1776            "graph_version": version,
1777            "input": input,
1778            "thread_id": thread_id,
1779            "budgets": requested_budgets
1780                .as_ref()
1781                .map(RunBudgets::requested_value)
1782                .unwrap_or(Value::Null),
1783            "checkpoint": checkpoint_requested,
1784        }));
1785
1786        let runs = self
1787            .runs
1788            .lock()
1789            .map_err(|e| internal_error(e.to_string()))?;
1790        if let Some(idem) = idempotency_key.as_deref() {
1791            if let Some(cached) =
1792                check_idempotency(self.store.as_ref(), Some(idem), &request_digest)?
1793            {
1794                return Ok(cached);
1795            }
1796        }
1797
1798        if checkpoint_requested {
1799            let Some(store) = self.store.as_ref() else {
1800                return Ok(error_output(
1801                    "SQLite persistence is required for deterministic checkpoints",
1802                    "CHECKPOINT_STORE_REQUIRED",
1803                ));
1804            };
1805            let eligibility = eligibility.expect("checkpoint eligibility");
1806            let state = initial_state_for_input(&input);
1807            let budgets_value = requested_budgets
1808                .as_ref()
1809                .map(RunBudgets::requested_value)
1810                .unwrap_or(Value::Null);
1811            let counters = serde_json::json!({"nodes":0,"llm_calls":0,"wall_clock_ms":0});
1812            let run_id = runs
1813                .allocate_with_budgets(
1814                    &graph_id,
1815                    &version,
1816                    input.clone(),
1817                    requested_budgets.clone(),
1818                )
1819                .map_err(|e| internal_error(e))?;
1820            if let Err(error) = store.save_execution_with_budgets(
1821                &run_id,
1822                &graph_id,
1823                &version,
1824                "checkpointed",
1825                &input.to_string(),
1826                Some(&budgets_value.to_string()),
1827            ) {
1828                runs.remove(&run_id);
1829                return Ok(error_output(error, "CHECKPOINT_PERSISTENCE_FAILURE"));
1830            }
1831            let checkpoint_record = match store.create_resume_checkpoint(
1832                &run_id,
1833                &graph_id,
1834                &version,
1835                &eligibility.next_node_cursor,
1836                &state,
1837                &budgets_value,
1838                &counters,
1839                &eligibility.dependency_summary,
1840                0,
1841                0,
1842            ) {
1843                Ok(record) => record,
1844                Err(error) => {
1845                    let _ = store.update_execution_status(&run_id, "failed", None, None, None);
1846                    runs.remove(&run_id);
1847                    return Ok(checkpoint_error_output(error));
1848                }
1849            };
1850            runs.mark_checkpointed(
1851                &run_id,
1852                &checkpoint_record.checkpoint_id,
1853                &checkpoint_record.checkpoint_digest,
1854            )
1855            .map_err(internal_error)?;
1856            let output = output_with_meta(
1857                serde_json::json!({
1858                    "run_id": run_id,
1859                    "status": "checkpointed",
1860                    "thread_id": thread_id,
1861                    "checkpoint_id": checkpoint_record.checkpoint_id,
1862                    "checkpoint_digest": checkpoint_record.checkpoint_digest,
1863                    "checkpoint": checkpoint_value(&checkpoint_record),
1864                    "resume_capability": "deterministic_local_resume",
1865                }),
1866                Some(&graph_id),
1867                Some(&version),
1868                Some(&run_id),
1869            );
1870            if let Some(idem) = idempotency_key {
1871                if let Some(cached) = persist_idempotency(store, &idem, &request_digest, &output)? {
1872                    return Ok(cached);
1873                }
1874            }
1875            return Ok(output);
1876        }
1877
1878        let run_id = runs
1879            .allocate_with_budgets(&graph_id, &version, input.clone(), requested_budgets)
1880            .map_err(|e| internal_error(e))?;
1881        if let Err(e) = runs.admit_async(&run_id) {
1882            runs.remove(&run_id);
1883            return Ok(error_output(e, "RUN_CAPACITY"));
1884        }
1885
1886        if let Some(ref store) = self.store {
1887            let _ =
1888                store.save_execution(&run_id, &graph_id, &version, "running", &input.to_string());
1889        }
1890
1891        let terminal_store = self.store.clone();
1892        let completion_runs = runs.clone();
1893        runs.start_with_completion_with_store(
1894            run_id.clone(),
1895            spec,
1896            self.base_url.clone(),
1897            self.default_model.clone(),
1898            self.store.clone(),
1899            move |record| Self::persist_terminal_and_mark(completion_runs, terminal_store, record),
1900        );
1901
1902        let output = output_with_meta(
1903            serde_json::json!({
1904                "run_id": run_id,
1905                "status": "running",
1906                "thread_id": thread_id,
1907            }),
1908            Some(&graph_id),
1909            Some(&version),
1910            Some(&run_id),
1911        );
1912        if let Some(ref store) = self.store {
1913            if let Some(idem) = idempotency_key {
1914                if let Some(cached) = persist_idempotency(store, &idem, &request_digest, &output)? {
1915                    return Ok(cached);
1916                }
1917            }
1918        }
1919
1920        Ok(output)
1921    }
1922
1923    #[tool(
1924        description = "Read one durable deterministic-local checkpoint, including its integrity-bound state and resume metadata."
1925    )]
1926    fn graph_run_checkpoint(
1927        &self,
1928        Parameters(RunCheckpointParams {
1929            run_id,
1930            checkpoint_id,
1931        }): Parameters<RunCheckpointParams>,
1932    ) -> Result<Json<StructuredOutput>, ErrorData> {
1933        let Some(store) = self.store.as_ref() else {
1934            return Ok(error_output(
1935                "SQLite persistence is required for checkpoint reads",
1936                "CHECKPOINT_STORE_REQUIRED",
1937            ));
1938        };
1939        if run_id.is_none() && checkpoint_id.is_none() {
1940            return Ok(error_output(
1941                "run_id or checkpoint_id is required for checkpoint reads",
1942                "INVALID_PARAMS",
1943            ));
1944        }
1945        match store.load_resume_checkpoint(checkpoint_id.as_deref(), run_id.as_deref()) {
1946            Ok(Some(record))
1947                if run_id
1948                    .as_deref()
1949                    .is_none_or(|run_id| record.run_id == run_id) =>
1950            {
1951                Ok(output_with_meta(
1952                    checkpoint_value(&record),
1953                    Some(&record.graph_id),
1954                    Some(&record.graph_version),
1955                    Some(&record.run_id),
1956                ))
1957            }
1958            Ok(Some(_)) => Ok(checkpoint_error_output(CheckpointError::Integrity)),
1959            Ok(None) => Ok(checkpoint_error_output(CheckpointError::NotFound)),
1960            Err(error) => Ok(checkpoint_error_output(error)),
1961        }
1962    }
1963
1964    #[tool(
1965        description = "Consume one deterministic-local checkpoint atomically and resume its pinned run exactly once."
1966    )]
1967    fn graph_run_resume(
1968        &self,
1969        Parameters(RunResumeParams {
1970            checkpoint_id,
1971            run_id,
1972        }): Parameters<RunResumeParams>,
1973    ) -> Result<Json<StructuredOutput>, ErrorData> {
1974        let Some(store) = self.store.as_ref() else {
1975            return Ok(error_output(
1976                "SQLite persistence is required for deterministic resume",
1977                "CHECKPOINT_STORE_REQUIRED",
1978            ));
1979        };
1980        if checkpoint_id.is_none() && run_id.is_none() {
1981            return Ok(error_output(
1982                "checkpoint_id or run_id is required for resume",
1983                "INVALID_PARAMS",
1984            ));
1985        }
1986        let checkpoint =
1987            match store.load_resume_checkpoint(checkpoint_id.as_deref(), run_id.as_deref()) {
1988                Ok(Some(record)) => record,
1989                Ok(None) => return Ok(checkpoint_error_output(CheckpointError::NotFound)),
1990                Err(error) => return Ok(checkpoint_error_output(error)),
1991            };
1992        if store
1993            .checkpoint_approval_status(&checkpoint.checkpoint_id)
1994            .map_err(|error| internal_error(error.message()))?
1995            .as_deref()
1996            == Some("pending")
1997        {
1998            return Ok(error_output(
1999                "checkpoint resume is pending its durable approval decision",
2000                "APPROVAL_PENDING",
2001            ));
2002        }
2003        if checkpoint.consumed_at.is_some() {
2004            return Ok(checkpoint_error_output(CheckpointError::Consumed));
2005        }
2006        if run_id
2007            .as_deref()
2008            .is_some_and(|run_id| run_id != checkpoint.run_id)
2009        {
2010            return Ok(checkpoint_error_output(CheckpointError::Integrity));
2011        }
2012        let Some(contract) = store
2013            .load_execution_contract(&checkpoint.run_id)
2014            .map_err(internal_error)?
2015        else {
2016            return Ok(checkpoint_error_output(CheckpointError::Integrity));
2017        };
2018        if contract.graph_id != checkpoint.graph_id
2019            || contract.graph_version != checkpoint.graph_version
2020            || checkpoint.terminal_cursor != 0
2021            || checkpoint.event_cursor != 0
2022        {
2023            return Ok(checkpoint_error_output(CheckpointError::Integrity));
2024        }
2025        let graph = match self.resolve_graph(&checkpoint.graph_id, Some(&checkpoint.graph_version))
2026        {
2027            Ok(graph) => graph,
2028            Err(_) => return Ok(checkpoint_error_output(CheckpointError::Integrity)),
2029        };
2030        if graph.version != checkpoint.graph_version {
2031            return Ok(checkpoint_error_output(CheckpointError::Integrity));
2032        }
2033        let eligibility = match graph.spec.resume_eligibility() {
2034            Ok(eligibility) => eligibility,
2035            Err(_) => {
2036                return Ok(error_output(
2037                    "checkpoint graph is no longer in the deterministic local resume subset",
2038                    "RESUME_INELIGIBLE",
2039                ))
2040            }
2041        };
2042        if checkpoint.next_node_cursor != eligibility.next_node_cursor
2043            || checkpoint.dependency_summary != eligibility.dependency_summary
2044            || checkpoint.dependency_digest != digest(&eligibility.dependency_summary)
2045            || checkpoint.state != initial_state_for_input(&contract.input)
2046            || checkpoint.budgets != contract.budgets
2047            || checkpoint.budget_counters
2048                != serde_json::json!({"nodes":0,"llm_calls":0,"wall_clock_ms":0})
2049        {
2050            return Ok(checkpoint_error_output(CheckpointError::Integrity));
2051        }
2052        let budgets = match RunBudgets::parse(Some(&checkpoint.budgets)) {
2053            Ok(budgets) => budgets,
2054            Err(_) => return Ok(checkpoint_error_output(CheckpointError::Integrity)),
2055        };
2056        let reserved_runs = self
2057            .runs
2058            .lock()
2059            .map_err(|e| internal_error(e.to_string()))?;
2060        if let Err(error) = reserved_runs.reserve_async_slot() {
2061            return Ok(error_output(error, "RUN_CAPACITY"));
2062        }
2063        drop(reserved_runs);
2064        let consumed = match store.consume_resume_checkpoint(&checkpoint.checkpoint_id) {
2065            Ok(record) => record,
2066            Err(error) => {
2067                if let Ok(runs) = self.runs.lock() {
2068                    runs.release_async_slot();
2069                }
2070                return Ok(checkpoint_error_output(error));
2071            }
2072        };
2073        self.launch_resumed(consumed, contract, graph, budgets, None)
2074    }
2075
2076    #[tool(description = "Wait for an async run to complete, with optional timeout.")]
2077    fn graph_run_wait(
2078        &self,
2079        Parameters(RunWaitParams { run_id, timeout_ms }): Parameters<RunWaitParams>,
2080    ) -> Result<Json<StructuredOutput>, ErrorData> {
2081        let timeout = Duration::from_millis(timeout_ms.unwrap_or(300_000));
2082        let deadline = Instant::now() + timeout;
2083        loop {
2084            let r = {
2085                let runs = self
2086                    .runs
2087                    .lock()
2088                    .map_err(|e| internal_error(e.to_string()))?;
2089                runs.get(&run_id)
2090                    .ok_or_else(|| invalid_params(format!("run '{run_id}' not found")))?
2091            };
2092            if matches!(r.status.as_str(), "completed" | "failed" | "cancelled") {
2093                let persist = Self::persist_terminal(self.store.clone(), r.clone());
2094                if let Ok(runs) = self.runs.lock() {
2095                    if self.store.is_none() {
2096                        runs.mark_persistence(&run_id, "volatile_no_store", None);
2097                    } else {
2098                        match persist {
2099                            Ok(()) => runs.mark_persistence(&run_id, "durable_terminal", None),
2100                            Err(error) => runs.mark_persistence(
2101                                &run_id,
2102                                "volatile_persistence_failed",
2103                                Some(error),
2104                            ),
2105                        }
2106                    }
2107                }
2108                let public = self
2109                    .runs
2110                    .lock()
2111                    .ok()
2112                    .and_then(|runs| runs.get(&run_id).map(|record| record.public()))
2113                    .unwrap_or_else(|| r.public());
2114                return Ok(output_with_meta(public, None, None, Some(&run_id)));
2115            }
2116            if Instant::now() >= deadline {
2117                return Ok(output_with_meta(
2118                    serde_json::json!({
2119                        "run_id": run_id,
2120                        "status": r.status,
2121                        "timed_out": true,
2122                    }),
2123                    None,
2124                    None,
2125                    Some(&run_id),
2126                ));
2127            }
2128            std::thread::sleep(Duration::from_millis(100));
2129        }
2130    }
2131
2132    #[tool(description = "Cancel a running execution.")]
2133    fn graph_run_cancel(
2134        &self,
2135        Parameters(RunCancelParams { run_id, reason: _ }): Parameters<RunCancelParams>,
2136    ) -> Result<Json<StructuredOutput>, ErrorData> {
2137        let runs = self
2138            .runs
2139            .lock()
2140            .map_err(|e| internal_error(e.to_string()))?;
2141        match runs.cancel(&run_id) {
2142            Ok(_) => {}
2143            Err(error) if error == "RUN_NOT_CANCELLABLE" => {
2144                return Ok(error_output(
2145                    "terminal or checkpointed runs cannot be cancelled",
2146                    "RUN_NOT_CANCELLABLE",
2147                ));
2148            }
2149            Err(error) if error == "run not found" => {
2150                drop(runs);
2151                if self
2152                    .store
2153                    .as_ref()
2154                    .and_then(|store| store.load_execution(&run_id).ok().flatten())
2155                    .is_some_and(|stored| {
2156                        matches!(
2157                            stored.get("status").and_then(Value::as_str),
2158                            Some("completed" | "failed" | "cancelled" | "checkpointed")
2159                        )
2160                    })
2161                {
2162                    return Ok(error_output(
2163                        "terminal or checkpointed runs cannot be cancelled",
2164                        "RUN_NOT_CANCELLABLE",
2165                    ));
2166                }
2167                return Err(invalid_params(error));
2168            }
2169            Err(error) => return Err(invalid_params(error)),
2170        }
2171        Ok(output_with_meta(
2172            serde_json::json!({
2173                "run_id": run_id,
2174                "status": "cancellation_requested",
2175                "cancellation_effect": "best_effort_drop_provider_future",
2176                "provider_request_may_still_be_in_flight": true,
2177                "effective_at": "provider_completion_or_cancellation_observation"
2178            }),
2179            None,
2180            None,
2181            Some(&run_id),
2182        ))
2183    }
2184
2185    #[tool(description = "Get current run status, budget usage, and pending approvals.")]
2186    fn graph_run_get(
2187        &self,
2188        Parameters(RunGetParams { run_id }): Parameters<RunGetParams>,
2189    ) -> Result<Json<StructuredOutput>, ErrorData> {
2190        let runs = self
2191            .runs
2192            .lock()
2193            .map_err(|e| internal_error(e.to_string()))?;
2194        if let Some(r) = runs.get(&run_id) {
2195            return Ok(output_with_meta(r.public(), None, None, Some(&run_id)));
2196        }
2197        drop(runs);
2198        if let Some(record) = self.stored_run(&run_id)? {
2199            return Ok(output_with_meta(record, None, None, Some(&run_id)));
2200        }
2201        Err(invalid_params(format!("run '{run_id}' not found")))
2202    }
2203
2204    #[tool(
2205        description = "Read the in-memory state projection from a live run; use graph_run_checkpoint for a durable checkpoint state."
2206    )]
2207    fn graph_run_state(
2208        &self,
2209        Parameters(RunStateParams {
2210            run_id,
2211            checkpoint_id: _,
2212            json_pointer,
2213        }): Parameters<RunStateParams>,
2214    ) -> Result<Json<StructuredOutput>, ErrorData> {
2215        let runs = self
2216            .runs
2217            .lock()
2218            .map_err(|e| internal_error(e.to_string()))?;
2219        let r = runs
2220            .get(&run_id)
2221            .ok_or_else(|| invalid_params(format!("run '{run_id}' not found")))?;
2222        let state = if let Some(pointer) = json_pointer.as_deref() {
2223            if pointer.is_empty() {
2224                r.state.clone()
2225            } else {
2226                r.state.pointer(pointer).cloned().unwrap_or(Value::Null)
2227            }
2228        } else {
2229            r.state.clone()
2230        };
2231        Ok(output_with_meta(
2232            serde_json::json!({
2233                "state": state,
2234                "run_id": run_id,
2235                "status": r.status,
2236            }),
2237            None,
2238            None,
2239            Some(&run_id),
2240        ))
2241    }
2242
2243    #[tool(
2244        description = "Read bounded events. With SQLite, terminal emitted events remain available as a persisted projection after restart; this is not replayable execution or resume support."
2245    )]
2246    fn graph_run_events(
2247        &self,
2248        Parameters(RunEventsParams {
2249            run_id,
2250            cursor,
2251            limit,
2252        }): Parameters<RunEventsParams>,
2253    ) -> Result<Json<StructuredOutput>, ErrorData> {
2254        let runs = self
2255            .runs
2256            .lock()
2257            .map_err(|e| internal_error(e.to_string()))?;
2258        let result = runs
2259            .events(
2260                self.store.as_ref(),
2261                &run_id,
2262                cursor.unwrap_or(0),
2263                limit.unwrap_or(100) as usize,
2264            )
2265            .map_err(|e| invalid_params(e))?;
2266        Ok(output_with_meta(result, None, None, Some(&run_id)))
2267    }
2268
2269    #[tool(description = "Fetch the canonical execution receipt for a run.")]
2270    fn graph_run_receipt(
2271        &self,
2272        Parameters(RunReceiptParams { run_id }): Parameters<RunReceiptParams>,
2273    ) -> Result<Json<StructuredOutput>, ErrorData> {
2274        if let Some(r) = self
2275            .runs
2276            .lock()
2277            .map_err(|e| internal_error(e.to_string()))?
2278            .get(&run_id)
2279        {
2280            return Ok(output_with_meta(
2281                r.receipt.clone(),
2282                None,
2283                None,
2284                Some(&run_id),
2285            ));
2286        }
2287        if let Some(store) = &self.store {
2288            match store.load_terminal_receipt(&run_id) {
2289                Ok(Some(receipt)) => {
2290                    return Ok(output_with_meta(receipt, None, None, Some(&run_id)));
2291                }
2292                Ok(None) => {}
2293                Err(error) if error == "RECEIPT_INTEGRITY_FAILURE" => {
2294                    return Ok(error_output(
2295                        "terminal receipt integrity validation failed",
2296                        "RECEIPT_INTEGRITY_FAILURE",
2297                    ));
2298                }
2299                Err(error) if error == "INTEGRITY_KEY_REQUIRED" => {
2300                    return Ok(error_output(
2301                        "an external integrity key is required for terminal receipt reads",
2302                        "INTEGRITY_KEY_REQUIRED",
2303                    ));
2304                }
2305                Err(error) => return Err(internal_error(error)),
2306            }
2307        }
2308        Ok(error_output(
2309            format!("run '{run_id}' not found"),
2310            "RUN_NOT_FOUND",
2311        ))
2312    }
2313
2314    // ── Policy + render ───────────────────────────────────────────────
2315
2316    #[tool(description = "Preflight a graph against policy before execution.")]
2317    fn graph_policy_check(
2318        &self,
2319        Parameters(PolicyCheckParams { graph_id, input: _ }): Parameters<PolicyCheckParams>,
2320    ) -> Result<Json<StructuredOutput>, ErrorData> {
2321        let graphs = self
2322            .graphs
2323            .lock()
2324            .map_err(|e| internal_error(e.to_string()))?;
2325        let g = graphs
2326            .get(&graph_id)
2327            .ok_or_else(|| invalid_params(format!("graph '{graph_id}' not found")))?;
2328
2329        let node_count = g.spec.nodes.len();
2330        let edge_count = g.spec.edges.len();
2331        let issues: Vec<String> = Vec::new();
2332
2333        Ok(structured_output(serde_json::json!({
2334            "graph_id": graph_id,
2335            "passed": issues.is_empty(),
2336            "issues": issues,
2337            "stats": {
2338                "node_count": node_count,
2339                "edge_count": edge_count,
2340                "max_iterations": g.spec.max_iterations,
2341                "max_parallelism": g.spec.max_parallelism,
2342            },
2343            "capabilities": {
2344                "models": [self.default_model.clone()],
2345                "tools": [],
2346            }
2347        })))
2348    }
2349
2350    #[tool(description = "Render a graph as Mermaid diagram or JSON topology.")]
2351    fn graph_render(
2352        &self,
2353        Parameters(RenderParams { graph_id, format }): Parameters<RenderParams>,
2354    ) -> Result<Json<StructuredOutput>, ErrorData> {
2355        let graphs = self
2356            .graphs
2357            .lock()
2358            .map_err(|e| internal_error(e.to_string()))?;
2359        let g = graphs
2360            .get(&graph_id)
2361            .ok_or_else(|| invalid_params(format!("graph '{graph_id}' not found")))?;
2362        let fmt = format.as_deref().unwrap_or("mermaid");
2363
2364        match fmt {
2365            "json" => Ok(output_with_meta(
2366                serde_json::json!({
2367                    "name": graph_id,
2368                    "nodes": g.spec.nodes.iter().map(|n| serde_json::json!({
2369                        "id": n.id, "type": n.node_type
2370                    })).collect::<Vec<_>>(),
2371                    "edges": g.spec.edges.iter().map(|e| serde_json::json!({
2372                        "from": e.from, "to": e.to
2373                    })).collect::<Vec<_>>(),
2374                }),
2375                Some(&graph_id),
2376                Some(&g.version),
2377                None,
2378            )),
2379            _ => Ok(output_with_meta(
2380                serde_json::json!({
2381                    "mermaid": Self::mermaid(&g.spec),
2382                    "name": graph_id,
2383                }),
2384                Some(&graph_id),
2385                Some(&g.version),
2386                None,
2387            )),
2388        }
2389    }
2390
2391    // ── Templates ─────────────────────────────────────────────────────
2392
2393    #[tool(description = "List available built-in graph templates.")]
2394    fn graph_template_list(
2395        &self,
2396        Parameters(TemplateListParams { query: _ }): Parameters<TemplateListParams>,
2397    ) -> Result<Json<StructuredOutput>, ErrorData> {
2398        Ok(structured_output(templates::list()))
2399    }
2400
2401    #[tool(
2402        description = "Instantiate a template into a graph spec that can be passed to graph_create."
2403    )]
2404    fn graph_template_instantiate(
2405        &self,
2406        Parameters(TemplateInstantiateParams { template_id, name }): Parameters<
2407            TemplateInstantiateParams,
2408        >,
2409    ) -> Result<Json<StructuredOutput>, ErrorData> {
2410        match templates::instantiate(&template_id, &name) {
2411            Ok(spec) => Ok(structured_output(serde_json::json!({
2412                "template_id": template_id,
2413                "name": name,
2414                "spec": spec,
2415            }))),
2416            Err(e) => Ok(error_output(e, "GRAPH_INVALID")),
2417        }
2418    }
2419    #[tool(description = "Read-only list of template promotion candidates.")]
2420    fn graph_template_candidates(
2421        &self,
2422        Parameters(TemplateCandidatesParams { state: _ }): Parameters<TemplateCandidatesParams>,
2423    ) -> Result<Json<StructuredOutput>, ErrorData> {
2424        Ok(structured_output(serde_json::json!({ "candidates": [] })))
2425    }
2426
2427    #[tool(description = "Read-only list of recorded outcomes for a template.")]
2428    fn graph_template_outcomes(
2429        &self,
2430        Parameters(TemplateOutcomesParams { template_id }): Parameters<TemplateOutcomesParams>,
2431    ) -> Result<Json<StructuredOutput>, ErrorData> {
2432        Ok(structured_output(serde_json::json!({
2433            "template_id": template_id,
2434            "outcomes": [],
2435        })))
2436    }
2437}
2438
2439#[tool_handler(
2440    router = self.tool_router,
2441    name = "agent-graph-mcp",
2442    version = "0.2.0",
2443    instructions = "Graph orchestration for bounded multi-step LLM workflows with parallel fan-out, conditional routing, state transforms, joins, cooperative cancellation, and optional enforced max_wall_clock_ms/max_nodes run budgets. max_llm_calls is rejected with INVALID_BUDGETS because no real invocation hook exists in this runtime path. Parallel unordered state writes require an explicit reducer. Cancellation can drop the local provider future on request, best effort; an underlying provider request may continue. Optional SQLite stores terminal projections plus explicit pre-execution checkpoints. Durable checkpoints, approvals, terminal receipts, and source witnesses require an external key file named by AGENT_GRAPH_INTEGRITY_KEY_PATH; without it their operations fail closed with INTEGRITY_KEY_REQUIRED. Deterministic local resume is limited to linear passthrough/state_transform chains and is never generic replay; uncheckpointed or ineligible runs remain interrupted_non_resumable after restart. SQLite-backed approvals can decide only an immutable deterministic-local checkpoint and resume that checkpoint; HumanApproval nodes and arbitrary external actions remain unsupported. Source witnesses are caller-supplied local captures: locators are never fetched, HMAC-authenticated witness integrity and bounded evidence spans are checked against SQLite, and source authority is not independently verified. Receipts provide integrity_only except a successfully resumed deterministic-local path, which reports deterministic_local_resume. Define graphs with graph_create, execute with graph_execute or graph_run_start, checkpoint with checkpoint:true, inspect with graph_run_get/wait/cancel/state/events/receipt/checkpoint, request or decide checkpoint approvals with graph_approval_request/decide, and resume with graph_run_resume."
2444)]
2445impl ServerHandler for AgentGraphServer {}