agent-graph-mcp 0.2.0

MCP server exposing agent-graph + llm-pipeline for graph-orchestrated LLM workflows — typed rmcp tools, parallel fan-out, interrupt/resume, durable persistence
Documentation
use agent_graph_mcp::{cli, proxy, transport, AgentGraphServer};
use rmcp::ServiceExt;
use std::io::{self, BufRead, Write};

fn main() {
    let args: Vec<String> = std::env::args().skip(1).collect();
    if args.first().map(String::as_str) == Some("--direct") {
        let direct: Vec<String> = args.into_iter().skip(1).collect();
        run_direct(&direct);
        return;
    }
    let cfg = match cli::parse_proxy_args(&args) {
        Ok(c) => c,
        Err(e) => {
            eprintln!("agent-graph-mcp: {}", e.message);
            std::process::exit(e.exit_code);
        }
    };
    let mut socket = match proxy::connect_timeout(&cfg.socket, cfg.timeout_ms) {
        Ok(s) => s,
        Err(_) => {
            eprintln!("DAEMON_UNAVAILABLE");
            std::process::exit(69);
        }
    };
    let stdin = io::stdin();
    let mut out = io::stdout();
    for line in stdin.lock().lines() {
        let line = match line {
            Ok(v) => v,
            Err(_) => break,
        };
        if transport::write_frame(&mut socket, line.as_bytes()).is_err() {
            break;
        }
        // JSON-RPC notifications have no "id" field and expect no response.
        // Skip read_frame for notifications to avoid blocking.
        let is_notification = line.contains("\"method\"") && !line.contains("\"id\"");
        if !is_notification {
            let response = match transport::read_frame(&mut socket) {
                Ok(v) => v,
                Err(_) => break,
            };
            let _ = out.write_all(&response);
            let _ = out.write_all(b"\n");
            let _ = out.flush();
        }
    }
}

fn run_direct(args: &[String]) {
    eprintln!("agent-graph-mcp: --direct is deprecated; use agent-graph-mcpd plus the proxy");
    let config = match cli::parse_args(args) {
        Ok(c) => c,
        Err(e) => {
            if e.exit_code != 0 {
                eprintln!("agent-graph-mcp: {}", e.message);
            }
            std::process::exit(e.exit_code);
        }
    };
    let key = config.integrity_key_path.clone().or_else(|| {
        std::env::var("AGENT_GRAPH_INTEGRITY_KEY_PATH")
            .ok()
            .map(std::path::PathBuf::from)
    });
    let rt = tokio::runtime::Runtime::new().expect("tokio runtime");
    rt.block_on(async move {
        let server =
            AgentGraphServer::new(config.base_url, config.default_model, config.data_dir, key)
                .map_err(anyhow::Error::msg)
                .unwrap();
        let service = server.serve(rmcp::transport::stdio()).await.unwrap();
        service.waiting().await.unwrap();
    });
}