//! MCP tools for event-based pipeline triggers: //! `schedule_event_trigger`, `list_event_triggers`, `cancel_event_trigger`. use crate::http::context::AppContext; use crate::service::event_triggers::store::{parse_action, parse_mode, parse_predicate}; use serde_json::{Value, json}; /// Register a new event trigger that fires when a `TransitionFired` event matches the predicate. pub(crate) fn tool_schedule_event_trigger( args: &Value, ctx: &AppContext, ) -> Result { let predicate = parse_predicate(args)?; let action = parse_action(args)?; let mode = parse_mode(args); let trigger = ctx.event_trigger_store.add(predicate, action, mode)?; serde_json::to_string_pretty(&json!({ "id": trigger.id, "mode": format!("{:?}", trigger.mode).to_lowercase(), "created_at": trigger.created_at.to_rfc3339(), "message": format!("Trigger {} registered.", trigger.id), })) .map_err(|e| format!("Serialization error: {e}")) } /// List all currently registered event triggers. pub(crate) fn tool_list_event_triggers(ctx: &AppContext) -> Result { let triggers = ctx.event_trigger_store.list(); let items: Vec = triggers .iter() .map(|t| { json!({ "id": t.id, "mode": format!("{:?}", t.mode).to_lowercase(), "created_at": t.created_at.to_rfc3339(), "predicate": { "story_id": t.predicate.story_id, "from_stage": t.predicate.from_stage, "to_stage": t.predicate.to_stage, "event_kind": t.predicate.event_kind, }, "action": serde_json::to_value(&t.action).unwrap_or(json!(null)), }) }) .collect(); serde_json::to_string_pretty(&json!({ "triggers": items, "count": items.len() })) .map_err(|e| format!("Serialization error: {e}")) } /// Cancel (remove) a registered event trigger by its ID. pub(crate) fn tool_cancel_event_trigger(args: &Value, ctx: &AppContext) -> Result { let id = args .get("id") .and_then(|v| v.as_str()) .ok_or("Missing required argument: id")?; if ctx.event_trigger_store.cancel(id) { Ok(format!("Trigger {id} cancelled.")) } else { Err(format!("No trigger found with id '{id}'.")) } } #[cfg(test)] mod tests { use super::*; use crate::http::test_helpers::test_ctx; #[test] fn schedule_and_list() { let tmp = tempfile::tempdir().unwrap(); let ctx = test_ctx(tmp.path()); let result = tool_schedule_event_trigger( &json!({ "predicate": { "to_stage": "Done" }, "action": { "type": "mcp", "method": "get_pipeline_status", "args": {} }, "mode": "once" }), &ctx, ) .unwrap(); let parsed: Value = serde_json::from_str(&result).unwrap(); let id = parsed["id"].as_str().unwrap(); assert!(!id.is_empty()); let list_result = tool_list_event_triggers(&ctx).unwrap(); let list: Value = serde_json::from_str(&list_result).unwrap(); assert_eq!(list["count"], 1); } #[test] fn cancel_existing_trigger() { let tmp = tempfile::tempdir().unwrap(); let ctx = test_ctx(tmp.path()); let result = tool_schedule_event_trigger( &json!({ "predicate": {}, "action": { "type": "prompt", "text": "investigate" }, "mode": "persistent" }), &ctx, ) .unwrap(); let parsed: Value = serde_json::from_str(&result).unwrap(); let id = parsed["id"].as_str().unwrap().to_string(); let cancel = tool_cancel_event_trigger(&json!({ "id": id }), &ctx).unwrap(); assert!(cancel.contains("cancelled")); let list: Value = serde_json::from_str(&tool_list_event_triggers(&ctx).unwrap()).unwrap(); assert_eq!(list["count"], 0); } #[test] fn cancel_missing_trigger_errors() { let tmp = tempfile::tempdir().unwrap(); let ctx = test_ctx(tmp.path()); let result = tool_cancel_event_trigger(&json!({ "id": "nonexistent-id" }), &ctx); assert!(result.is_err()); assert!(result.unwrap_err().contains("No trigger found")); } #[test] fn schedule_missing_predicate_errors() { let tmp = tempfile::tempdir().unwrap(); let ctx = test_ctx(tmp.path()); let result = tool_schedule_event_trigger( &json!({ "action": { "type": "mcp", "method": "get_version", "args": {} } }), &ctx, ); assert!(result.is_err()); assert!(result.unwrap_err().contains("predicate")); } #[test] fn schedule_missing_action_errors() { let tmp = tempfile::tempdir().unwrap(); let ctx = test_ctx(tmp.path()); let result = tool_schedule_event_trigger(&json!({ "predicate": { "to_stage": "Done" } }), &ctx); assert!(result.is_err()); assert!(result.unwrap_err().contains("action")); } }