1use 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 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 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 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 #[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 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 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 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 #[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 #[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 #[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 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 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 #[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 #[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 #[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 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 #[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 #[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 #[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 {}