Skip to main content

stygian_graph/
mcp.rs

1//! MCP (Model Context Protocol) server for graph-based scraping.
2//!
3//! Exposes `stygian-graph` scraping and pipeline capabilities as an MCP server
4//! over stdin/stdout using the JSON-RPC 2.0 protocol. External tools (LLM
5//! agents, IDE plugins) can scrape URLs, query REST/GraphQL APIs, parse feeds,
6//! and execute full pipeline DAGs via the standardised MCP interface.
7//!
8//! ## Enabling
9//!
10//! ```toml
11//! [dependencies]
12//! stygian-graph = { version = "*", features = ["mcp"] }
13//! ```
14//!
15//! ## Running the server
16//!
17//! Add `stygian-graph` as a dependency with the `mcp` feature and call
18//! `McpGraphServer::run()` from your own binary:
19//!
20//! ```rust,no_run
21//! use stygian_graph::mcp::McpGraphServer;
22//!
23//! #[tokio::main]
24//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
25//!     McpGraphServer::run().await
26//! }
27//! ```
28//!
29//! ## Protocol
30//!
31//! Implements MCP 2026-07-28 over JSON-RPC 2.0 on stdin/stdout.
32//!
33//! | MCP Method | Description |
34//! | ----------- | ------------- |
35//! | `server/discover` | Advertise protocol versions, identity, capabilities |
36//! | `tools/list` | List available scraping and pipeline tools |
37//! | `tools/call` | Execute a scraping or pipeline tool |
38//!
39//! ## Migrating from MCP 2025-11-25
40//!
41//! The `initialize` / `notifications/initialized` handshake was removed. Clients
42//! must advertise their protocol version, identity, and capabilities in the
43//! `_meta` block of every request under the `io.modelcontextprotocol/*` keys.
44//! See the `extract_client_protocol_version` and `extract_meta` helpers for
45//! the reader-side extraction (used by future PRs in the [MCP-001] migration
46//! sequence; PR 1 only adds them as helpers — enforcement lands in PR 4
47//! alongside the aggregator).
48//!
49//! [MCP-001]: https://github.com/greysquirr3l/stygian/issues/95
50//!
51//! ## Tools
52//!
53//! | Tool | Key Parameters | Returns |
54//! | ------ | -------------- | ------- |
55//! | `scrape` | `url`, `timeout_secs?`, `proxy_url?`, `rotate_ua?` | `data`, `metadata` |
56//! | `scrape_rest` | `url`, `method?`, `auth?`, `query?`, `body?`, `headers?`, `pagination?`, `data_path?` | `data`, `metadata` |
57//! | `scrape_graphql` | `url`, `query`, `variables?`, `auth?`, `data_path?` | `data`, `metadata` |
58//! | `scrape_sitemap` | `url`, `max_depth?` | `data` (JSON array of entries), `metadata` |
59//! | `scrape_rss` | `url` | `data` (JSON array of items), `metadata` |
60//! | `pipeline_validate` | `toml` | `nodes`, `services`, `execution_order`, `valid` |
61//! | `pipeline_run` | `toml`, `timeout_secs?` | per-node `outputs`, `skipped`, `errors` |
62//! | `charon_*` | feature-gated HAR diagnostics and planning inputs | Charon report/policy JSON payloads |
63
64use std::collections::HashMap;
65use std::future::Future;
66#[cfg(feature = "acquisition-runner")]
67use std::sync::Arc;
68use std::time::Duration;
69
70#[cfg(feature = "charon")]
71use serde::de::DeserializeOwned;
72use serde_json::{Value, json};
73use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
74#[cfg(feature = "acquisition-runner")]
75use tokio::sync::OnceCell;
76use tracing::{debug, info, warn};
77
78#[cfg(feature = "acquisition-runner")]
79use stygian_browser::{
80    AcquisitionMode, AcquisitionRequest, AcquisitionRunner, BrowserConfig, BrowserPool,
81};
82#[cfg(feature = "acquisition-runner")]
83use stygian_charon::AcquisitionModeHint;
84#[cfg(feature = "charon")]
85use stygian_charon::{
86    AcquisitionPolicy, InvestigationBundle, InvestigationReport, RequirementsProfile,
87    RuntimePolicy, TargetClass, TransactionView, build_runtime_policy, classify_transaction,
88    infer_requirements_with_target_class, investigate_har, map_runtime_policy,
89};
90
91use crate::{
92    adapters::{
93        graphql::{GraphQlConfig, GraphQlService},
94        http::{HttpAdapter, HttpConfig},
95        rest_api::RestApiAdapter,
96        rss_feed::RssFeedAdapter,
97        sitemap::SitemapAdapter,
98    },
99    application::pipeline_parser::{NodeDecl, PipelineParser, ServiceDecl},
100    ports::{ScrapingService, ServiceInput},
101};
102
103// ─── Error response helpers ───────────────────────────────────────────────────
104
105fn error_response(id: &Value, code: i64, message: &str) -> Value {
106    json!({
107        "jsonrpc": "2.0",
108        "id": id,
109        "error": { "code": code, "message": message }
110    })
111}
112
113fn ok_response(id: &Value, result: Value) -> Value {
114    let mut map = serde_json::Map::new();
115    map.insert("jsonrpc".to_owned(), json!("2.0"));
116    map.insert("id".to_owned(), id.clone());
117    // MCP 2026-07-28 §8: every result carries a `resultType` field. `"complete"`
118    // for ordinary responses; `"input_required"` for MRTR interim responses.
119    // We only emit `"complete"` here — MRTR lands in a later migration PR.
120    let result = match result {
121        Value::Object(mut obj) => {
122            obj.insert("resultType".to_owned(), json!("complete"));
123            Value::Object(obj)
124        }
125        other => {
126            let mut obj = serde_json::Map::new();
127            obj.insert("resultType".to_owned(), json!("complete"));
128            obj.insert("value".to_owned(), other);
129            Value::Object(obj)
130        }
131    };
132    map.insert("result".to_owned(), result);
133    Value::Object(map)
134}
135
136#[cfg(feature = "charon")]
137fn json_content_response(id: &Value, payload: &Value) -> Value {
138    ok_response(
139        id,
140        json!({
141            "content": [{
142                "type": "text",
143                "text": serde_json::to_string(payload).unwrap_or_default()
144            }]
145        }),
146    )
147}
148
149#[cfg(feature = "charon")]
150fn decode_required_arg<T: DeserializeOwned>(args: &Value, key: &str) -> Result<T, String> {
151    let raw = args
152        .get(key)
153        .cloned()
154        .ok_or_else(|| format!("Missing required parameter: {key}"))?;
155    serde_json::from_value(raw).map_err(|e| format!("Invalid parameter '{key}': {e}"))
156}
157
158#[cfg(feature = "charon")]
159fn parse_target_class_json(value: Option<&Value>) -> Result<TargetClass, String> {
160    let Some(value) = value else {
161        return Ok(TargetClass::Unknown);
162    };
163    let Some(raw) = value.as_str() else {
164        return Err("target_class must be a string".to_string());
165    };
166
167    match raw.trim().to_ascii_lowercase().as_str() {
168        "api" => Ok(TargetClass::Api),
169        "content-site" | "content_site" | "contentsite" | "content" => Ok(TargetClass::ContentSite),
170        "high-security" | "high_security" | "highsecurity" => Ok(TargetClass::HighSecurity),
171        "unknown" => Ok(TargetClass::Unknown),
172        _ => Err(format!("Unknown target_class: {raw}")),
173    }
174}
175
176// ─── MCP 2026-07-28 _meta helpers ──────────────────────────────────────────────
177
178/// Read the `io.modelcontextprotocol/<key>` entry from a request's
179/// `params._meta` block. Returns `None` when the request omits `_meta` or
180/// the requested key is absent.
181///
182/// MCP 2026-07-28 §2: every request now carries its protocol version, client
183/// identity, and client capabilities under the `io.modelcontextprotocol/*`
184/// namespace within `params._meta`.
185#[allow(dead_code)] // Used by PR 4 (aggregator) and future per-request gates.
186fn extract_meta<'a>(req: &'a Value, key: &str) -> Option<&'a Value> {
187    let meta = req.get("params")?.get("_meta")?.as_object()?;
188    meta.get(&format!("io.modelcontextprotocol/{key}"))
189}
190
191/// Extract the client's advertised protocol version from a request's `_meta`.
192///
193/// Returns `None` when the field is absent. Spec mandates this be present on
194/// every 2026-07-28 request; enforcement is the aggregator's responsibility
195/// (lands in PR 4 of [MCP-001]).
196///
197/// [MCP-001]: https://github.com/greysquirr3l/stygian/issues/95
198#[allow(dead_code)] // Used by PR 4 (aggregator) and future per-request gates.
199fn extract_client_protocol_version(req: &Value) -> Option<String> {
200    extract_meta(req, "protocolVersion")
201        .and_then(Value::as_str)
202        .map(str::to_owned)
203}
204
205/// Compare a client-advertised protocol version against a list of versions the
206/// server supports. Returns `Ok(())` when the version is in the supported
207/// list, `Err(unsupported)` with the offending value otherwise.
208///
209/// MCP 2026-07-28 §2: version mismatch on any request returns
210/// `UnsupportedProtocolVersionError` (code `-32022`). PR 1 exposes the
211/// helper; PR 4 wires it into the aggregator's per-request gate.
212#[cfg(test)]
213fn is_supported_protocol_version(client: &str, supported: &[&str]) -> Result<(), String> {
214    if supported.contains(&client) {
215        Ok(())
216    } else {
217        Err(format!("Unsupported protocol version: {client}"))
218    }
219}
220
221// ─── Server ───────────────────────────────────────────────────────────────────
222
223/// MCP server exposing `stygian-graph` scraping and pipeline tools.
224///
225/// All tools are stateless — each invocation builds and runs the appropriate
226/// adapter directly without maintaining any server-side session state.
227///
228/// # Example
229///
230/// ```no_run
231/// use stygian_graph::mcp::McpGraphServer;
232///
233/// #[tokio::main]
234/// async fn main() -> Result<(), Box<dyn std::error::Error>> {
235///     McpGraphServer::run().await
236/// }
237/// ```
238pub struct McpGraphServer;
239
240impl McpGraphServer {
241    /// Create a new MCP graph server.
242    #[must_use]
243    pub const fn new() -> Self {
244        Self
245    }
246
247    /// Run the MCP server, reading JSON-RPC requests from stdin and writing
248    /// responses to stdout until EOF.
249    ///
250    /// # Errors
251    ///
252    /// Returns an `Err` if the underlying I/O fails unrecoverably.
253    pub async fn run() -> Result<(), Box<dyn std::error::Error>> {
254        info!("stygian-graph MCP server starting");
255
256        let stdin = tokio::io::stdin();
257        let mut reader = BufReader::new(stdin);
258        let mut stdout = tokio::io::stdout();
259        let mut line = String::new();
260
261        loop {
262            line.clear();
263            let bytes = reader.read_line(&mut line).await?;
264            if bytes == 0 {
265                break; // EOF
266            }
267
268            let trimmed = line.trim();
269            if trimmed.is_empty() {
270                continue;
271            }
272
273            debug!(request = trimmed, "received");
274
275            let response = match serde_json::from_str::<Value>(trimmed) {
276                Ok(req) => {
277                    let is_well_formed_notification = req.is_object()
278                        && req.get("jsonrpc").and_then(Value::as_str) == Some("2.0")
279                        && req.get("id").is_none()
280                        && req.get("method").and_then(Value::as_str).is_some();
281                    let response = Self::handle(&req).await;
282                    if is_well_formed_notification {
283                        continue;
284                    }
285                    response
286                }
287                Err(e) => json!({
288                    "jsonrpc": "2.0",
289                    "id": null,
290                    "error": { "code": -32700, "message": format!("Parse error: {e}") }
291                }),
292            };
293
294            let mut out = serde_json::to_string(&response)?;
295            out.push('\n');
296            stdout.write_all(out.as_bytes()).await?;
297            stdout.flush().await?;
298        }
299
300        info!("stygian-graph MCP server stopped");
301        Ok(())
302    }
303
304    /// Dispatch a single JSON-RPC request.
305    ///
306    /// Used by the `stygian-mcp` aggregator to route tool calls through this
307    /// server without running the full stdin/stdout loop.
308    ///
309    /// # Example
310    ///
311    /// ```
312    /// use stygian_graph::mcp::McpGraphServer;
313    /// use serde_json::json;
314    ///
315    /// # tokio_test::block_on(async {
316    /// let req = json!({"jsonrpc":"2.0","id":1,"method":"server/discover"});
317    /// let resp = McpGraphServer::handle_request(&req).await;
318    /// assert_eq!(
319    ///     resp.pointer("/result/protocolVersion").and_then(serde_json::Value::as_str),
320    ///     Some("2026-07-28")
321    /// );
322    /// # });
323    /// ```
324    pub async fn handle_request(req: &Value) -> Value {
325        Self::handle(req).await
326    }
327
328    async fn handle(req: &Value) -> Value {
329        let null = Value::Null;
330        let id = req.get("id").unwrap_or(&null);
331        let method = req.get("method").and_then(Value::as_str).unwrap_or("");
332
333        match method {
334            // MCP 2026-07-28 §3: `server/discover` replaces the `initialize`
335            // handshake. The handshake (`initialize` + `notifications/initialized`)
336            // and the unrelated `ping` RPC are removed.
337            "server/discover" => Self::handle_discover(id),
338            "tools/list" => Self::handle_tools_list(id),
339            "tools/call" => Self::handle_tools_call(id, req).await,
340            _ => error_response(id, -32601, &format!("Method not found: {method}")),
341        }
342    }
343
344    /// Advertise the server's identity, supported protocol versions, and
345    /// capabilities. Replaces the `initialize` handshake removed in
346    /// MCP 2026-07-28.
347    fn handle_discover(id: &Value) -> Value {
348        ok_response(
349            id,
350            json!({
351                "protocolVersion": "2026-07-28",
352                "supportedProtocolVersions": ["2026-07-28"],
353                "capabilities": {
354                    "tools":     { "listChanged": false },
355                    "resources": { "listChanged": false }
356                },
357                "serverInfo": {
358                    "name": "stygian-graph",
359                    "version": env!("CARGO_PKG_VERSION")
360                },
361                "extensions": []
362            }),
363        )
364    }
365
366    fn scraping_tool_defs() -> Vec<Value> {
367        vec![
368            json!({
369                "name": "scrape",
370                "description": "Fetch a URL with anti-bot UA rotation and retry logic. Returns raw HTML/JSON content and response metadata.",
371                "inputSchema": {
372                    "type": "object",
373                    "properties": {
374                        "url":          { "type": "string",  "description": "Target URL" },
375                        "timeout_secs": { "type": "integer", "description": "Request timeout in seconds (default: 30)" },
376                        "proxy_url":    { "type": "string",  "description": "HTTP/SOCKS5 proxy URL (e.g. socks5://user:pass@host:1080). Only pass this when the user has explicitly requested proxy use. Do NOT populate this field by default." },
377                        "rotate_ua":    { "type": "boolean", "description": "Rotate User-Agent on each request (default: true)" }
378                    },
379                    "required": ["url"]
380                }
381            }),
382            json!({
383                "name": "scrape_rest",
384                "description": "Call a REST/JSON API. Supports bearer/API-key auth, arbitrary HTTP methods, query parameters, request bodies, pagination, and response path extraction.",
385                "inputSchema": {
386                    "type": "object",
387                    "properties": {
388                        "url":        { "type": "string", "description": "API endpoint URL" },
389                        "method":     { "type": "string", "description": "HTTP method (GET, POST, PUT, PATCH, DELETE — default: GET)" },
390                        "auth":       {
391                            "type": "object",
392                            "description": "Authentication config",
393                            "properties": {
394                                "type":  { "type": "string", "description": "bearer | api_key | basic | header" },
395                                "token": { "type": "string", "description": "Token or credential value" },
396                                "header":{ "type": "string", "description": "Custom header name (for type=header)" }
397                            }
398                        },
399                        "query":      { "type": "object", "description": "URL query parameters as key-value pairs" },
400                        "body":       { "type": "object", "description": "Request body (JSON)" },
401                        "headers":    { "type": "object", "description": "Custom request headers" },
402                        "pagination": {
403                            "type": "object",
404                            "description": "Pagination config",
405                            "properties": {
406                                "strategy":  { "type": "string", "description": "link_header | offset | cursor" },
407                                "max_pages": { "type": "integer", "description": "Maximum pages to fetch (default: 1)" }
408                            }
409                        },
410                        "data_path":  { "type": "string", "description": "Dot-separated JSON path to extract (e.g. data.items)" }
411                    },
412                    "required": ["url"]
413                }
414            }),
415            json!({
416                "name": "scrape_graphql",
417                "description": "Execute a GraphQL query against any spec-compliant endpoint. Supports bearer/API-key auth, variables, and dot-path data extraction.",
418                "inputSchema": {
419                    "type": "object",
420                    "properties": {
421                        "url":       { "type": "string", "description": "GraphQL endpoint URL" },
422                        "query":     { "type": "string", "description": "GraphQL query or mutation string" },
423                        "variables": { "type": "object", "description": "Query variables (JSON object)" },
424                        "auth": {
425                            "type": "object",
426                            "description": "Auth config",
427                            "properties": {
428                                "kind":        { "type": "string", "description": "bearer | api_key | header | none" },
429                                "token":       { "type": "string", "description": "Auth token or key" },
430                                "header_name": { "type": "string", "description": "Custom header name (default: X-Api-Key)" }
431                            }
432                        },
433                        "data_path":     { "type": "string", "description": "Dot-separated path to extract from response (e.g. data.countries)" },
434                        "timeout_secs":  { "type": "integer", "description": "Request timeout in seconds (default: 30)" }
435                    },
436                    "required": ["url", "query"]
437                }
438            }),
439            json!({
440                "name": "scrape_sitemap",
441                "description": "Parse a sitemap.xml or sitemap index and return all discovered URLs with their priorities and change frequencies.",
442                "inputSchema": {
443                    "type": "object",
444                    "properties": {
445                        "url":       { "type": "string",  "description": "Sitemap URL (sitemap.xml or sitemap index)" },
446                        "max_depth": { "type": "integer", "description": "Maximum sitemap index recursion depth (default: 5)" }
447                    },
448                    "required": ["url"]
449                }
450            }),
451            json!({
452                "name": "scrape_rss",
453                "description": "Parse an RSS or Atom feed and return all entries as structured JSON.",
454                "inputSchema": {
455                    "type": "object",
456                    "properties": {
457                        "url": { "type": "string", "description": "RSS/Atom feed URL" }
458                    },
459                    "required": ["url"]
460                }
461            }),
462        ]
463    }
464
465    fn graph_tool_defs() -> Vec<Value> {
466        let mut tools = vec![
467            json!({
468                "name": "pipeline_validate",
469                "description": "Parse and validate a TOML pipeline definition without executing it. Returns the node list, service declarations, and computed execution order.",
470                "inputSchema": {
471                    "type": "object",
472                    "properties": {
473                        "toml": { "type": "string", "description": "TOML pipeline definition string" }
474                    },
475                    "required": ["toml"]
476                }
477            }),
478            json!({
479                "name": "pipeline_run",
480                "description": "Parse, validate, and execute a TOML pipeline DAG. HTTP, REST, GraphQL, sitemap, and RSS nodes are executed. AI nodes and browser nodes without opt-in acquisition config are recorded in the skipped list.",
481                "inputSchema": {
482                    "type": "object",
483                    "properties": {
484                        "toml":         { "type": "string",  "description": "TOML pipeline definition string" },
485                        "timeout_secs": { "type": "integer", "description": "Per-node timeout in seconds (default: 30)" }
486                    },
487                    "required": ["toml"]
488                }
489            }),
490            json!({
491                "name": "inspect",
492                "description": "Get a complete snapshot of a pipeline's graph structure including nodes, edges, execution waves, critical path, and connectivity metrics.",
493                "inputSchema": {
494                    "type": "object",
495                    "properties": {
496                        "toml": { "type": "string", "description": "TOML pipeline definition string" }
497                    },
498                    "required": ["toml"]
499                }
500            }),
501            json!({
502                "name": "node_info",
503                "description": "Get detailed information about a specific node in the pipeline graph, including its service type, depth, predecessors, and successors.",
504                "inputSchema": {
505                    "type": "object",
506                    "properties": {
507                        "toml":    { "type": "string", "description": "TOML pipeline definition string" },
508                        "node_id": { "type": "string", "description": "Node ID to inspect" }
509                    },
510                    "required": ["toml", "node_id"]
511                }
512            }),
513            json!({
514                "name": "impact",
515                "description": "Analyze what would be affected by changing a node. Returns all upstream dependencies and downstream dependents.",
516                "inputSchema": {
517                    "type": "object",
518                    "properties": {
519                        "toml":    { "type": "string", "description": "TOML pipeline definition string" },
520                        "node_id": { "type": "string", "description": "Node ID to analyze impact for" }
521                    },
522                    "required": ["toml", "node_id"]
523                }
524            }),
525            json!({
526                "name": "query_nodes",
527                "description": "Query nodes in the pipeline graph by various criteria: service type, root/leaf status, depth range, or ID pattern.",
528                "inputSchema": {
529                    "type": "object",
530                    "properties": {
531                        "toml":       { "type": "string",  "description": "TOML pipeline definition string" },
532                        "service":    { "type": "string",  "description": "Filter by service type (http, ai, browser, etc.)" },
533                        "id_pattern": { "type": "string",  "description": "Filter by node ID substring match" },
534                        "is_root":    { "type": "boolean", "description": "Only return root nodes (no predecessors)" },
535                        "is_leaf":    { "type": "boolean", "description": "Only return leaf nodes (no successors)" },
536                        "min_depth":  { "type": "integer", "description": "Minimum depth from root nodes" },
537                        "max_depth":  { "type": "integer", "description": "Maximum depth from root nodes" }
538                    },
539                    "required": ["toml"]
540                }
541            }),
542        ];
543
544        #[cfg(feature = "charon")]
545        tools.extend(Self::charon_tool_defs());
546
547        tools
548    }
549
550    #[cfg(feature = "charon")]
551    fn charon_tool_defs() -> Vec<Value> {
552        vec![
553            json!({
554                "name": "charon_classify_transaction",
555                "description": "Classify a single HTTP transaction for likely anti-bot provider signals.",
556                "inputSchema": {
557                    "type": "object",
558                    "properties": {
559                        "url": { "type": "string", "description": "Request URL" },
560                        "status": { "type": "integer", "description": "HTTP status code" },
561                        "response_headers": { "type": "object", "description": "Response headers as a string map" },
562                        "response_body_snippet": { "type": "string", "description": "Optional response body snippet" },
563                        "response_body_excerpt": { "type": "string", "description": "Alias for response_body_snippet" }
564                    },
565                    "required": ["url", "status"]
566                }
567            }),
568            json!({
569                "name": "charon_investigate_har",
570                "description": "Build a Charon investigation report from a HAR payload.",
571                "inputSchema": {
572                    "type": "object",
573                    "properties": {
574                        "har": { "type": "string", "description": "HAR JSON payload" },
575                        "target_class": { "type": "string", "description": "Optional target class: api | content-site | high-security | unknown" }
576                    },
577                    "required": ["har"]
578                }
579            }),
580            json!({
581                "name": "charon_infer_requirements",
582                "description": "Infer Charon operational requirements from an investigation report.",
583                "inputSchema": {
584                    "type": "object",
585                    "properties": {
586                        "report": { "type": "object", "description": "InvestigationReport JSON object" },
587                        "target_class": { "type": "string", "description": "Optional target class override: api | content-site | high-security | unknown" }
588                    },
589                    "required": ["report"]
590                }
591            }),
592            json!({
593                "name": "charon_build_runtime_policy",
594                "description": "Build a runtime policy from a Charon investigation report and inferred requirements profile.",
595                "inputSchema": {
596                    "type": "object",
597                    "properties": {
598                        "report": { "type": "object", "description": "InvestigationReport JSON object" },
599                        "requirements": { "type": "object", "description": "RequirementsProfile JSON object" }
600                    },
601                    "required": ["report", "requirements"]
602                }
603            }),
604            json!({
605                "name": "charon_map_runtime_policy",
606                "description": "Map a Charon runtime policy into acquisition hints for downstream runners.",
607                "inputSchema": {
608                    "type": "object",
609                    "properties": {
610                        "policy": { "type": "object", "description": "RuntimePolicy JSON object" }
611                    },
612                    "required": ["policy"]
613                }
614            }),
615            json!({
616                "name": "charon_analyze_and_plan",
617                "description": "Run end-to-end Charon HAR analysis, requirement inference, runtime policy planning, and acquisition mapping in one call.",
618                "inputSchema": {
619                    "type": "object",
620                    "properties": {
621                        "har": { "type": "string", "description": "HAR JSON payload" },
622                        "target_class": { "type": "string", "description": "Optional target class: api | content-site | high-security | unknown" }
623                    },
624                    "required": ["har"]
625                }
626            }),
627        ]
628    }
629
630    fn handle_tools_list(id: &Value) -> Value {
631        let mut tools = Self::scraping_tool_defs();
632        tools.extend(Self::graph_tool_defs());
633        ok_response(id, json!({ "tools": tools }))
634    }
635
636    async fn handle_tools_call(id: &Value, req: &Value) -> Value {
637        let null = Value::Null;
638        let params = req.get("params").unwrap_or(&null);
639        let name = params.get("name").and_then(Value::as_str).unwrap_or("");
640        let args = params.get("arguments").cloned().unwrap_or(Value::Null);
641
642        match name {
643            "scrape" => Self::tool_scrape(id, &args).await,
644            "scrape_rest" => Self::tool_scrape_rest(id, &args).await,
645            "scrape_graphql" => Self::tool_scrape_graphql(id, &args).await,
646            "scrape_sitemap" => Self::tool_scrape_sitemap(id, &args).await,
647            "scrape_rss" => Self::tool_scrape_rss(id, &args).await,
648            "pipeline_validate" => Self::tool_pipeline_validate(id, &args),
649            "pipeline_run" => Self::tool_pipeline_run(id, &args).await,
650            "inspect" => Self::tool_graph_inspect(id, &args),
651            "node_info" => Self::tool_graph_node_info(id, &args),
652            "impact" => Self::tool_graph_impact(id, &args),
653            "query_nodes" => Self::tool_graph_query(id, &args),
654            #[cfg(feature = "charon")]
655            "charon_classify_transaction" => Self::tool_charon_classify_transaction(id, &args),
656            #[cfg(feature = "charon")]
657            "charon_investigate_har" => Self::tool_charon_investigate_har(id, &args),
658            #[cfg(feature = "charon")]
659            "charon_infer_requirements" => Self::tool_charon_infer_requirements(id, &args),
660            #[cfg(feature = "charon")]
661            "charon_build_runtime_policy" => Self::tool_charon_build_runtime_policy(id, &args),
662            #[cfg(feature = "charon")]
663            "charon_map_runtime_policy" => Self::tool_charon_map_runtime_policy(id, &args),
664            #[cfg(feature = "charon")]
665            "charon_analyze_and_plan" => Self::tool_charon_analyze_and_plan(id, &args),
666            _ => error_response(id, -32602, &format!("Unknown tool: {name}")),
667        }
668    }
669
670    #[cfg(feature = "charon")]
671    fn tool_charon_classify_transaction(id: &Value, args: &Value) -> Value {
672        let Some(url) = args.get("url").and_then(Value::as_str) else {
673            return error_response(id, -32602, "Missing required parameter: url");
674        };
675        let Some(status_u64) = args.get("status").and_then(Value::as_u64) else {
676            return error_response(id, -32602, "Missing required parameter: status");
677        };
678        let Ok(status) = u16::try_from(status_u64) else {
679            return error_response(id, -32602, "status must fit in a 16-bit unsigned integer");
680        };
681
682        let response_headers = match args.get("response_headers") {
683            Some(value) if !value.is_null() => {
684                match serde_json::from_value::<std::collections::BTreeMap<String, String>>(
685                    value.clone(),
686                ) {
687                    Ok(headers) => headers,
688                    Err(e) => {
689                        return error_response(
690                            id,
691                            -32602,
692                            &format!("Invalid parameter 'response_headers': {e}"),
693                        );
694                    }
695                }
696            }
697            _ => std::collections::BTreeMap::new(),
698        };
699        let response_body_snippet = args
700            .get("response_body_snippet")
701            .or_else(|| args.get("response_body_excerpt"))
702            .and_then(Value::as_str)
703            .map(str::to_string);
704
705        let tx = TransactionView {
706            url: url.to_string(),
707            status,
708            response_headers,
709            response_body_snippet,
710        };
711        let detection = classify_transaction(&tx);
712        json_content_response(id, &json!({ "detection": detection }))
713    }
714
715    #[cfg(feature = "charon")]
716    fn tool_charon_investigate_har(id: &Value, args: &Value) -> Value {
717        let Some(har) = args.get("har").and_then(Value::as_str) else {
718            return error_response(id, -32602, "Missing required parameter: har");
719        };
720        let target_class = match parse_target_class_json(args.get("target_class")) {
721            Ok(target_class) => target_class,
722            Err(e) => return error_response(id, -32602, &e),
723        };
724
725        match investigate_har(har) {
726            Ok(mut report) => {
727                report.target_class = Some(target_class);
728                json_content_response(id, &json!({ "report": report }))
729            }
730            Err(e) => error_response(id, -32603, &format!("HAR investigation failed: {e}")),
731        }
732    }
733
734    #[cfg(feature = "charon")]
735    fn tool_charon_infer_requirements(id: &Value, args: &Value) -> Value {
736        let mut report: InvestigationReport = match decode_required_arg(args, "report") {
737            Ok(report) => report,
738            Err(e) => return error_response(id, -32602, &e),
739        };
740        let target_class = match parse_target_class_json(args.get("target_class")) {
741            Ok(TargetClass::Unknown) => report.target_class.unwrap_or(TargetClass::Unknown),
742            Ok(target_class) => target_class,
743            Err(e) => return error_response(id, -32602, &e),
744        };
745
746        report.target_class = Some(target_class);
747        let requirements = infer_requirements_with_target_class(&report, target_class);
748        json_content_response(id, &json!({ "requirements": requirements }))
749    }
750
751    #[cfg(feature = "charon")]
752    fn tool_charon_build_runtime_policy(id: &Value, args: &Value) -> Value {
753        let report: InvestigationReport = match decode_required_arg(args, "report") {
754            Ok(report) => report,
755            Err(e) => return error_response(id, -32602, &e),
756        };
757        let requirements: RequirementsProfile = match decode_required_arg(args, "requirements") {
758            Ok(requirements) => requirements,
759            Err(e) => return error_response(id, -32602, &e),
760        };
761
762        let policy = build_runtime_policy(&report, &requirements);
763        json_content_response(id, &json!({ "policy": policy }))
764    }
765
766    #[cfg(feature = "charon")]
767    fn tool_charon_map_runtime_policy(id: &Value, args: &Value) -> Value {
768        let policy: RuntimePolicy = match decode_required_arg(args, "policy") {
769            Ok(policy) => policy,
770            Err(e) => return error_response(id, -32602, &e),
771        };
772
773        let acquisition: AcquisitionPolicy = map_runtime_policy(&policy);
774        json_content_response(id, &json!({ "acquisition": acquisition }))
775    }
776
777    #[cfg(feature = "charon")]
778    fn tool_charon_analyze_and_plan(id: &Value, args: &Value) -> Value {
779        let Some(har) = args.get("har").and_then(Value::as_str) else {
780            return error_response(id, -32602, "Missing required parameter: har");
781        };
782        let target_class = match parse_target_class_json(args.get("target_class")) {
783            Ok(target_class) => target_class,
784            Err(e) => return error_response(id, -32602, &e),
785        };
786
787        match investigate_har(har) {
788            Ok(mut report) => {
789                report.target_class = Some(target_class);
790                let requirements = infer_requirements_with_target_class(&report, target_class);
791                let policy = build_runtime_policy(&report, &requirements);
792                let acquisition = map_runtime_policy(&policy);
793                let bundle = InvestigationBundle {
794                    report,
795                    requirements,
796                    policy,
797                };
798                json_content_response(id, &json!({ "bundle": bundle, "acquisition": acquisition }))
799            }
800            Err(e) => error_response(id, -32603, &format!("HAR investigation failed: {e}")),
801        }
802    }
803
804    // ── scrape ───────────────────────────────────────────────────────────────
805
806    async fn tool_scrape(id: &Value, args: &Value) -> Value {
807        let Some(url) = args.get("url").and_then(Value::as_str) else {
808            return error_response(id, -32602, "Missing required parameter: url");
809        };
810
811        let timeout_secs = args
812            .get("timeout_secs")
813            .and_then(Value::as_u64)
814            .unwrap_or(30);
815        let proxy_url = args
816            .get("proxy_url")
817            .and_then(Value::as_str)
818            .map(str::to_string);
819        let rotate_ua = args
820            .get("rotate_ua")
821            .and_then(Value::as_bool)
822            .unwrap_or(true);
823
824        let config = HttpConfig {
825            timeout: std::time::Duration::from_secs(timeout_secs),
826            proxy_url,
827            rotate_user_agent: rotate_ua,
828            // T112: opt in to catalogue-UA allow-list. The rotation
829            // pool only emits browser-class UAs, so this is a no-op
830            // in practice — but explicit is safer than implicit.
831            allow_plain_http: true,
832            ..HttpConfig::default()
833        };
834        let adapter = HttpAdapter::with_config(config);
835        let input = ServiceInput {
836            url: url.to_string(),
837            params: json!({}),
838        };
839
840        match adapter.execute(input).await {
841            Ok(output) => ok_response(
842                id,
843                json!({
844                    "content": [{
845                        "type": "text",
846                        "text": serde_json::to_string(&json!({
847                            "data": output.data,
848                            "metadata": output.metadata
849                        })).unwrap_or_default()
850                    }]
851                }),
852            ),
853            Err(e) => error_response(id, -32603, &format!("Scrape failed: {e}")),
854        }
855    }
856
857    // ── scrape_rest ──────────────────────────────────────────────────────────
858
859    async fn tool_scrape_rest(id: &Value, args: &Value) -> Value {
860        let Some(url) = args.get("url").and_then(Value::as_str) else {
861            return error_response(id, -32602, "Missing required parameter: url");
862        };
863
864        // Build params JSON from explicit fields only; extra keys in `args` are intentionally
865        // not forwarded — the REST adapter only reads the fields it recognises.
866        let mut map = serde_json::Map::new();
867        if let Some(method) = args.get("method").and_then(Value::as_str) {
868            map.insert("method".to_owned(), json!(method));
869        }
870        if let Some(auth) = args.get("auth").filter(|v| !v.is_null()) {
871            map.insert("auth".to_owned(), auth.clone());
872        }
873        if let Some(query) = args.get("query").filter(|v| !v.is_null()) {
874            map.insert("query".to_owned(), query.clone());
875        }
876        if let Some(body) = args.get("body").filter(|v| !v.is_null()) {
877            map.insert("body".to_owned(), body.clone());
878        }
879        if let Some(headers) = args.get("headers").filter(|v| !v.is_null()) {
880            map.insert("headers".to_owned(), headers.clone());
881        }
882        if let Some(pagination) = args.get("pagination").filter(|v| !v.is_null()) {
883            map.insert("pagination".to_owned(), pagination.clone());
884        }
885        if let Some(dp) = args.get("data_path").and_then(Value::as_str) {
886            map.insert("response".to_owned(), json!({ "data_path": dp }));
887        }
888        let params = Value::Object(map);
889
890        let adapter = RestApiAdapter::new();
891        let input = ServiceInput {
892            url: url.to_string(),
893            params,
894        };
895
896        match adapter.execute(input).await {
897            Ok(output) => ok_response(
898                id,
899                json!({
900                    "content": [{
901                        "type": "text",
902                        "text": serde_json::to_string(&json!({
903                            "data": output.data,
904                            "metadata": output.metadata
905                        })).unwrap_or_default()
906                    }]
907                }),
908            ),
909            Err(e) => error_response(id, -32603, &format!("REST scrape failed: {e}")),
910        }
911    }
912
913    // ── scrape_graphql ───────────────────────────────────────────────────────
914
915    async fn tool_scrape_graphql(id: &Value, args: &Value) -> Value {
916        let Some(url) = args.get("url").and_then(Value::as_str) else {
917            return error_response(id, -32602, "Missing required parameter: url");
918        };
919        let Some(query) = args.get("query").and_then(Value::as_str) else {
920            return error_response(id, -32602, "Missing required parameter: query");
921        };
922
923        let timeout_secs = args
924            .get("timeout_secs")
925            .and_then(Value::as_u64)
926            .unwrap_or(30);
927
928        let config = GraphQlConfig {
929            timeout_secs,
930            ..GraphQlConfig::default()
931        };
932        let service = GraphQlService::new(config, None);
933
934        let mut gql_map = serde_json::Map::new();
935        gql_map.insert("query".to_owned(), json!(query));
936        if let Some(variables) = args.get("variables").filter(|v| !v.is_null()) {
937            gql_map.insert("variables".to_owned(), variables.clone());
938        }
939        if let Some(auth) = args.get("auth").filter(|v| !v.is_null()) {
940            gql_map.insert("auth".to_owned(), auth.clone());
941        }
942        if let Some(dp) = args.get("data_path").and_then(Value::as_str) {
943            gql_map.insert("data_path".to_owned(), json!(dp));
944        }
945        let params = Value::Object(gql_map);
946
947        let input = ServiceInput {
948            url: url.to_string(),
949            params,
950        };
951
952        match service.execute(input).await {
953            Ok(output) => ok_response(
954                id,
955                json!({
956                    "content": [{
957                        "type": "text",
958                        "text": serde_json::to_string(&json!({
959                            "data": output.data,
960                            "metadata": output.metadata
961                        })).unwrap_or_default()
962                    }]
963                }),
964            ),
965            Err(e) => error_response(id, -32603, &format!("GraphQL scrape failed: {e}")),
966        }
967    }
968
969    // ── scrape_sitemap ───────────────────────────────────────────────────────
970
971    async fn tool_scrape_sitemap(id: &Value, args: &Value) -> Value {
972        let Some(url) = args.get("url").and_then(Value::as_str) else {
973            return error_response(id, -32602, "Missing required parameter: url");
974        };
975
976        let max_depth = args
977            .get("max_depth")
978            .and_then(Value::as_u64)
979            .map_or(5, |v| usize::try_from(v).unwrap_or(5));
980        let client = reqwest::Client::new();
981        let adapter = SitemapAdapter::new(client, max_depth);
982        let input = ServiceInput {
983            url: url.to_string(),
984            params: json!({}),
985        };
986
987        match adapter.execute(input).await {
988            Ok(output) => ok_response(
989                id,
990                json!({
991                    "content": [{
992                        "type": "text",
993                        "text": serde_json::to_string(&json!({
994                            "data": output.data,
995                            "metadata": output.metadata
996                        })).unwrap_or_default()
997                    }]
998                }),
999            ),
1000            Err(e) => error_response(id, -32603, &format!("Sitemap scrape failed: {e}")),
1001        }
1002    }
1003
1004    // ── scrape_rss ───────────────────────────────────────────────────────────
1005
1006    async fn tool_scrape_rss(id: &Value, args: &Value) -> Value {
1007        let Some(url) = args.get("url").and_then(Value::as_str) else {
1008            return error_response(id, -32602, "Missing required parameter: url");
1009        };
1010
1011        let client = reqwest::Client::new();
1012        let adapter = RssFeedAdapter::new(client);
1013        let input = ServiceInput {
1014            url: url.to_string(),
1015            params: json!({}),
1016        };
1017
1018        match adapter.execute(input).await {
1019            Ok(output) => ok_response(
1020                id,
1021                json!({
1022                    "content": [{
1023                        "type": "text",
1024                        "text": serde_json::to_string(&json!({
1025                            "data": output.data,
1026                            "metadata": output.metadata
1027                        })).unwrap_or_default()
1028                    }]
1029                }),
1030            ),
1031            Err(e) => error_response(id, -32603, &format!("RSS scrape failed: {e}")),
1032        }
1033    }
1034
1035    // ── pipeline_validate ────────────────────────────────────────────────────
1036
1037    fn tool_pipeline_validate(id: &Value, args: &Value) -> Value {
1038        let Some(toml) = args.get("toml").and_then(Value::as_str) else {
1039            return error_response(id, -32602, "Missing required parameter: toml");
1040        };
1041
1042        let def = match PipelineParser::from_str(toml) {
1043            Ok(d) => d,
1044            Err(e) => return error_response(id, -32603, &format!("Parse error: {e}")),
1045        };
1046
1047        if let Err(e) = def.validate() {
1048            return ok_response(
1049                id,
1050                json!({
1051                    "content": [{
1052                        "type": "text",
1053                        "text": serde_json::to_string(&json!({
1054                            "valid": false,
1055                            "error": e.to_string(),
1056                            "nodes": def.nodes.len(),
1057                            "services": def.services.len()
1058                        })).unwrap_or_default()
1059                    }]
1060                }),
1061            );
1062        }
1063
1064        let order = match def.topological_order() {
1065            Ok(o) => o,
1066            Err(e) => return error_response(id, -32603, &format!("Topology error: {e}")),
1067        };
1068
1069        let node_info: Vec<Value> = def
1070            .nodes
1071            .iter()
1072            .map(|n| {
1073                json!({
1074                    "name": n.name,
1075                    "service": n.service,
1076                    "url": n.url,
1077                    "depends_on": n.depends_on
1078                })
1079            })
1080            .collect();
1081
1082        let svc_info: Vec<Value> = def
1083            .services
1084            .iter()
1085            .map(|s| {
1086                json!({
1087                    "name": s.name,
1088                    "kind": s.kind,
1089                    "model": s.model
1090                })
1091            })
1092            .collect();
1093
1094        ok_response(
1095            id,
1096            json!({
1097                "content": [{
1098                    "type": "text",
1099                    "text": serde_json::to_string(&json!({
1100                        "valid": true,
1101                        "node_count": def.nodes.len(),
1102                        "service_count": def.services.len(),
1103                        "execution_order": order,
1104                        "nodes": node_info,
1105                        "services": svc_info
1106                    })).unwrap_or_default()
1107                }]
1108            }),
1109        )
1110    }
1111
1112    // ── pipeline_run ─────────────────────────────────────────────────────────
1113
1114    async fn tool_pipeline_run(id: &Value, args: &Value) -> Value {
1115        let Some(toml) = args.get("toml").and_then(Value::as_str) else {
1116            return error_response(id, -32602, "Missing required parameter: toml");
1117        };
1118
1119        let timeout_secs = args
1120            .get("timeout_secs")
1121            .and_then(Value::as_u64)
1122            .unwrap_or(30);
1123
1124        let def = match PipelineParser::from_str(toml) {
1125            Ok(d) => d,
1126            Err(e) => return error_response(id, -32603, &format!("Parse error: {e}")),
1127        };
1128
1129        if let Err(e) = def.validate() {
1130            return error_response(id, -32603, &format!("Validation error: {e}"));
1131        }
1132
1133        let order = match def.topological_order() {
1134            Ok(o) => o,
1135            Err(e) => return error_response(id, -32603, &format!("Topology error: {e}")),
1136        };
1137
1138        let svc_kinds: HashMap<String, ServiceDecl> = def
1139            .services
1140            .iter()
1141            .map(|s| (s.name.clone(), s.clone()))
1142            .collect();
1143
1144        let mut outputs: HashMap<String, Value> = HashMap::new();
1145        let mut skipped: Vec<String> = Vec::new();
1146        let mut errors: HashMap<String, String> = HashMap::new();
1147
1148        for node_name in &order {
1149            let Some(node) = def.nodes.iter().find(|n| n.name == *node_name) else {
1150                continue;
1151            };
1152
1153            let kind = svc_kinds
1154                .get(&node.service)
1155                .map_or(node.service.as_str(), |s| s.kind.as_str());
1156
1157            // Nodes without a URL are AI/transform nodes — skip.
1158            let Some(url) = node.url.as_deref() else {
1159                skipped.push(node_name.clone());
1160                continue;
1161            };
1162
1163            match execute_pipeline_node(kind, url, node_name, node, timeout_secs).await {
1164                Some(Ok(out)) => {
1165                    outputs.insert(node_name.clone(), out);
1166                }
1167                Some(Err(e)) => {
1168                    errors.insert(node_name.clone(), e);
1169                }
1170                None => {
1171                    skipped.push(node_name.clone());
1172                }
1173            }
1174        }
1175
1176        ok_response(
1177            id,
1178            json!({
1179                "content": [{
1180                    "type": "text",
1181                    "text": serde_json::to_string(&json!({
1182                        "execution_order": order,
1183                        "outputs": outputs,
1184                        "skipped": skipped,
1185                        "errors": errors
1186                    })).unwrap_or_default()
1187                }]
1188            }),
1189        )
1190    }
1191
1192    // ── Graph introspection tools ─────────────────────────────────────────────
1193
1194    fn tool_graph_inspect(id: &Value, args: &Value) -> Value {
1195        let Some(toml) = args.get("toml").and_then(Value::as_str) else {
1196            return error_response(id, -32602, "Missing required parameter: toml");
1197        };
1198
1199        let def = match PipelineParser::from_str(toml) {
1200            Ok(d) => d,
1201            Err(e) => return error_response(id, -32603, &format!("Parse error: {e}")),
1202        };
1203
1204        if let Err(e) = def.validate() {
1205            return error_response(id, -32603, &format!("Validation error: {e}"));
1206        }
1207
1208        // Build a Pipeline from the definition
1209        let mut pipeline = crate::domain::graph::Pipeline::new("pipeline");
1210        for node in &def.nodes {
1211            pipeline.add_node(crate::domain::graph::Node::with_metadata(
1212                &node.name,
1213                &node.service,
1214                serde_json::json!({
1215                    "url": node.url,
1216                    "params": toml_to_json(&toml::Value::Table(
1217                        node.params.iter()
1218                            .map(|(k, v)| (k.clone(), v.clone()))
1219                            .collect()
1220                    ))
1221                }),
1222                serde_json::Value::Null,
1223            ));
1224            for dep in &node.depends_on {
1225                pipeline.add_edge(crate::domain::graph::Edge::new(dep, &node.name));
1226            }
1227        }
1228
1229        let executor = match crate::domain::graph::DagExecutor::from_pipeline(&pipeline) {
1230            Ok(e) => e,
1231            Err(e) => return error_response(id, -32603, &format!("Graph build error: {e}")),
1232        };
1233
1234        let snapshot = executor.snapshot();
1235
1236        ok_response(
1237            id,
1238            json!({
1239                "content": [{
1240                    "type": "text",
1241                    "text": serde_json::to_string(&snapshot).unwrap_or_default()
1242                }]
1243            }),
1244        )
1245    }
1246
1247    fn tool_graph_node_info(id: &Value, args: &Value) -> Value {
1248        let Some(toml) = args.get("toml").and_then(Value::as_str) else {
1249            return error_response(id, -32602, "Missing required parameter: toml");
1250        };
1251        let Some(node_id) = args.get("node_id").and_then(Value::as_str) else {
1252            return error_response(id, -32602, "Missing required parameter: node_id");
1253        };
1254
1255        let def = match PipelineParser::from_str(toml) {
1256            Ok(d) => d,
1257            Err(e) => return error_response(id, -32603, &format!("Parse error: {e}")),
1258        };
1259
1260        if let Err(e) = def.validate() {
1261            return error_response(id, -32603, &format!("Validation error: {e}"));
1262        }
1263
1264        let mut pipeline = crate::domain::graph::Pipeline::new("pipeline");
1265        for node in &def.nodes {
1266            pipeline.add_node(crate::domain::graph::Node::with_metadata(
1267                &node.name,
1268                &node.service,
1269                serde_json::json!({
1270                    "url": node.url,
1271                    "params": toml_to_json(&toml::Value::Table(
1272                        node.params.iter()
1273                            .map(|(k, v)| (k.clone(), v.clone()))
1274                            .collect()
1275                    ))
1276                }),
1277                serde_json::Value::Null,
1278            ));
1279            for dep in &node.depends_on {
1280                pipeline.add_edge(crate::domain::graph::Edge::new(dep, &node.name));
1281            }
1282        }
1283
1284        let executor = match crate::domain::graph::DagExecutor::from_pipeline(&pipeline) {
1285            Ok(e) => e,
1286            Err(e) => return error_response(id, -32603, &format!("Graph build error: {e}")),
1287        };
1288
1289        executor.node_info(node_id).map_or_else(
1290            || error_response(id, -32602, &format!("Node not found: {node_id}")),
1291            |info| {
1292                ok_response(
1293                    id,
1294                    json!({
1295                        "content": [{
1296                            "type": "text",
1297                            "text": serde_json::to_string(&info).unwrap_or_default()
1298                        }]
1299                    }),
1300                )
1301            },
1302        )
1303    }
1304
1305    fn tool_graph_impact(id: &Value, args: &Value) -> Value {
1306        let Some(toml) = args.get("toml").and_then(Value::as_str) else {
1307            return error_response(id, -32602, "Missing required parameter: toml");
1308        };
1309        let Some(node_id) = args.get("node_id").and_then(Value::as_str) else {
1310            return error_response(id, -32602, "Missing required parameter: node_id");
1311        };
1312
1313        let def = match PipelineParser::from_str(toml) {
1314            Ok(d) => d,
1315            Err(e) => return error_response(id, -32603, &format!("Parse error: {e}")),
1316        };
1317
1318        if let Err(e) = def.validate() {
1319            return error_response(id, -32603, &format!("Validation error: {e}"));
1320        }
1321
1322        let mut pipeline = crate::domain::graph::Pipeline::new("pipeline");
1323        for node in &def.nodes {
1324            pipeline.add_node(crate::domain::graph::Node::with_metadata(
1325                &node.name,
1326                &node.service,
1327                serde_json::json!({
1328                    "url": node.url,
1329                    "params": toml_to_json(&toml::Value::Table(
1330                        node.params.iter()
1331                            .map(|(k, v)| (k.clone(), v.clone()))
1332                            .collect()
1333                    ))
1334                }),
1335                serde_json::Value::Null,
1336            ));
1337            for dep in &node.depends_on {
1338                pipeline.add_edge(crate::domain::graph::Edge::new(dep, &node.name));
1339            }
1340        }
1341
1342        let executor = match crate::domain::graph::DagExecutor::from_pipeline(&pipeline) {
1343            Ok(e) => e,
1344            Err(e) => return error_response(id, -32603, &format!("Graph build error: {e}")),
1345        };
1346
1347        let impact = executor.impact_analysis(node_id);
1348
1349        ok_response(
1350            id,
1351            json!({
1352                "content": [{
1353                    "type": "text",
1354                    "text": serde_json::to_string(&impact).unwrap_or_default()
1355                }]
1356            }),
1357        )
1358    }
1359
1360    fn tool_graph_query(id: &Value, args: &Value) -> Value {
1361        let Some(toml) = args.get("toml").and_then(Value::as_str) else {
1362            return error_response(id, -32602, "Missing required parameter: toml");
1363        };
1364
1365        let def = match PipelineParser::from_str(toml) {
1366            Ok(d) => d,
1367            Err(e) => return error_response(id, -32603, &format!("Parse error: {e}")),
1368        };
1369
1370        if let Err(e) = def.validate() {
1371            return error_response(id, -32603, &format!("Validation error: {e}"));
1372        }
1373
1374        let mut pipeline = crate::domain::graph::Pipeline::new("pipeline");
1375        for node in &def.nodes {
1376            pipeline.add_node(crate::domain::graph::Node::with_metadata(
1377                &node.name,
1378                &node.service,
1379                serde_json::json!({
1380                    "url": node.url,
1381                    "params": toml_to_json(&toml::Value::Table(
1382                        node.params.iter()
1383                            .map(|(k, v)| (k.clone(), v.clone()))
1384                            .collect()
1385                    ))
1386                }),
1387                serde_json::Value::Null,
1388            ));
1389            for dep in &node.depends_on {
1390                pipeline.add_edge(crate::domain::graph::Edge::new(dep, &node.name));
1391            }
1392        }
1393
1394        let executor = match crate::domain::graph::DagExecutor::from_pipeline(&pipeline) {
1395            Ok(e) => e,
1396            Err(e) => return error_response(id, -32603, &format!("Graph build error: {e}")),
1397        };
1398
1399        // Build the query from args
1400        let query = crate::domain::introspection::NodeQuery {
1401            service: args
1402                .get("service")
1403                .and_then(Value::as_str)
1404                .map(String::from),
1405            id: None,
1406            id_pattern: args
1407                .get("id_pattern")
1408                .and_then(Value::as_str)
1409                .map(String::from),
1410            is_root: args.get("is_root").and_then(Value::as_bool),
1411            is_leaf: args.get("is_leaf").and_then(Value::as_bool),
1412            min_depth: args
1413                .get("min_depth")
1414                .and_then(Value::as_u64)
1415                .map(|v| usize::try_from(v).unwrap_or(0)),
1416            max_depth: args
1417                .get("max_depth")
1418                .and_then(Value::as_u64)
1419                .map(|v| usize::try_from(v).unwrap_or(0)),
1420        };
1421
1422        let results = executor.query_nodes(&query);
1423
1424        ok_response(
1425            id,
1426            json!({
1427                "content": [{
1428                    "type": "text",
1429                    "text": serde_json::to_string(&results).unwrap_or_default()
1430                }]
1431            }),
1432        )
1433    }
1434}
1435
1436impl Default for McpGraphServer {
1437    fn default() -> Self {
1438        Self::new()
1439    }
1440}
1441
1442// ─── Helpers ──────────────────────────────────────────────────────────────────
1443
1444/// Build a [`GraphQlService`] and the corresponding [`ServiceInput`] from a pipeline node.
1445///
1446/// Extracts `query`, `variables`, `auth`, and `data_path` from the node's TOML params.
1447fn build_graphql_node_request(
1448    node: &NodeDecl,
1449    url: &str,
1450    timeout_secs: u64,
1451) -> (GraphQlService, ServiceInput) {
1452    let query = node
1453        .params
1454        .get("query")
1455        .and_then(|v| v.as_str())
1456        .unwrap_or("")
1457        .to_string();
1458    let config = GraphQlConfig {
1459        timeout_secs,
1460        ..GraphQlConfig::default()
1461    };
1462    let service = GraphQlService::new(config, None);
1463    let mut gql_map = serde_json::Map::new();
1464    gql_map.insert("query".to_owned(), json!(query));
1465    if let Some(variables) = node.params.get("variables") {
1466        gql_map.insert("variables".to_owned(), toml_to_json(variables));
1467    }
1468    if let Some(auth) = node.params.get("auth") {
1469        gql_map.insert("auth".to_owned(), toml_to_json(auth));
1470    }
1471    if let Some(dp) = node.params.get("data_path").and_then(|v| v.as_str()) {
1472        gql_map.insert("data_path".to_owned(), json!(dp));
1473    }
1474    (
1475        service,
1476        ServiceInput {
1477            url: url.to_string(),
1478            params: Value::Object(gql_map),
1479        },
1480    )
1481}
1482
1483#[derive(Debug, Clone, PartialEq)]
1484struct AcquisitionNodeConfig {
1485    mode: String,
1486    wait_for_selector: Option<String>,
1487    extraction_js: Option<String>,
1488    total_timeout: Option<Duration>,
1489    #[cfg(feature = "acquisition-runner")]
1490    target_class: Option<TargetClass>,
1491}
1492
1493#[cfg(feature = "acquisition-runner")]
1494fn parse_optional_target_class(value: &toml::Value) -> Result<TargetClass, String> {
1495    let raw = value
1496        .as_str()
1497        .ok_or_else(|| "acquisition.target_class must be a string".to_string())?;
1498    match raw {
1499        "api" => Ok(TargetClass::Api),
1500        "content-site" | "content_site" | "contentsite" | "content" => Ok(TargetClass::ContentSite),
1501        "high-security" | "high_security" | "highsecurity" | "high" => {
1502            Ok(TargetClass::HighSecurity)
1503        }
1504        "unknown" => Ok(TargetClass::Unknown),
1505        _ => Err(
1506            "acquisition.target_class must be one of: api, content-site, high-security, unknown"
1507                .to_string(),
1508        ),
1509    }
1510}
1511
1512#[cfg(feature = "acquisition-runner")]
1513fn parse_optional_positive_secs(value: &toml::Value) -> Result<Duration, String> {
1514    const MAX_ACQUISITION_TIMEOUT_SECS: u64 = 86_400;
1515    const MAX_ACQUISITION_TIMEOUT_SECS_F64: f64 = 86_400.0;
1516
1517    if let Some(seconds) = value.as_float() {
1518        if seconds.is_finite() && seconds > 0.0 && seconds <= MAX_ACQUISITION_TIMEOUT_SECS_F64 {
1519            return Ok(Duration::from_secs_f64(seconds));
1520        }
1521        return Err(format!(
1522            "acquisition.total_timeout_secs must be a positive finite number <= {MAX_ACQUISITION_TIMEOUT_SECS}"
1523        ));
1524    }
1525
1526    if let Some(seconds) = value.as_integer() {
1527        if seconds > 0 && seconds <= i64::try_from(MAX_ACQUISITION_TIMEOUT_SECS).unwrap_or(i64::MAX)
1528        {
1529            return Ok(Duration::from_secs(u64::try_from(seconds).map_err(
1530                |_| "acquisition.total_timeout_secs must fit into an unsigned integer".to_string(),
1531            )?));
1532        }
1533        return Err(format!(
1534            "acquisition.total_timeout_secs must be an integer in 1..={MAX_ACQUISITION_TIMEOUT_SECS}"
1535        ));
1536    }
1537
1538    Err("acquisition.total_timeout_secs must be a number".to_string())
1539}
1540
1541#[cfg(feature = "acquisition-runner")]
1542fn acquisition_config_from_node(node: &NodeDecl) -> Result<Option<AcquisitionNodeConfig>, String> {
1543    let Some(raw) = node.params.get("acquisition") else {
1544        return Ok(None);
1545    };
1546
1547    let table = raw
1548        .as_table()
1549        .ok_or_else(|| "acquisition must be a TOML table".to_string())?;
1550
1551    let enabled = table
1552        .get("enabled")
1553        .and_then(toml::Value::as_bool)
1554        .unwrap_or(true);
1555
1556    if !enabled {
1557        return Ok(None);
1558    }
1559
1560    let mode = table
1561        .get("mode")
1562        .and_then(toml::Value::as_str)
1563        .unwrap_or("resilient")
1564        .to_string();
1565
1566    let wait_for_selector = table
1567        .get("wait_for_selector")
1568        .or_else(|| table.get("selector_wait"))
1569        .and_then(toml::Value::as_str)
1570        .map(ToString::to_string);
1571
1572    let extraction_js = table
1573        .get("extraction_js")
1574        .and_then(toml::Value::as_str)
1575        .map(ToString::to_string);
1576
1577    let total_timeout = table
1578        .get("total_timeout_secs")
1579        .map(parse_optional_positive_secs)
1580        .transpose()?;
1581
1582    let target_class = table
1583        .get("target_class")
1584        .map(parse_optional_target_class)
1585        .transpose()?;
1586
1587    Ok(Some(AcquisitionNodeConfig {
1588        mode,
1589        wait_for_selector,
1590        extraction_js,
1591        total_timeout,
1592        target_class,
1593    }))
1594}
1595
1596#[cfg(feature = "acquisition-runner")]
1597fn parse_acquisition_mode(raw: &str) -> Result<AcquisitionMode, String> {
1598    match raw {
1599        "fast" => Ok(AcquisitionMode::Fast),
1600        "resilient" => Ok(AcquisitionMode::Resilient),
1601        "hostile" => Ok(AcquisitionMode::Hostile),
1602        "investigate" => Ok(AcquisitionMode::Investigate),
1603        other => Err(format!(
1604            "Invalid acquisition mode '{other}'. Use one of: fast, resilient, hostile, investigate"
1605        )),
1606    }
1607}
1608
1609#[cfg(feature = "acquisition-runner")]
1610const fn mode_rank(mode: AcquisitionMode) -> u8 {
1611    match mode {
1612        AcquisitionMode::Fast => 0,
1613        AcquisitionMode::Resilient => 1,
1614        AcquisitionMode::Hostile => 2,
1615        AcquisitionMode::Investigate => 3,
1616    }
1617}
1618
1619#[cfg(all(feature = "acquisition-runner", feature = "charon"))]
1620const fn mode_from_hint(hint: AcquisitionModeHint) -> AcquisitionMode {
1621    match hint {
1622        AcquisitionModeHint::Fast => AcquisitionMode::Fast,
1623        AcquisitionModeHint::Resilient => AcquisitionMode::Resilient,
1624        AcquisitionModeHint::Hostile => AcquisitionMode::Hostile,
1625        AcquisitionModeHint::Investigate => AcquisitionMode::Investigate,
1626    }
1627}
1628
1629#[cfg(feature = "acquisition-runner")]
1630fn build_status_only_har(url: &str, status: u16, body_excerpt: Option<&str>) -> String {
1631    let text = body_excerpt.unwrap_or_default();
1632    json!({
1633        "log": {
1634            "version": "1.2",
1635            "creator": {"name": "stygian-graph-acquisition-bridge", "version": "1.0"},
1636            "pages": [{
1637                "id": "page_1",
1638                "title": url,
1639                "startedDateTime": "2026-01-01T00:00:00.000Z",
1640                "pageTimings": {"onLoad": 0}
1641            }],
1642            "entries": [{
1643                "pageref": "page_1",
1644                "startedDateTime": "2026-01-01T00:00:00.000Z",
1645                "time": 0,
1646                "request": {
1647                    "method": "GET",
1648                    "url": url,
1649                    "httpVersion": "HTTP/2",
1650                    "headers": [],
1651                    "queryString": [],
1652                    "cookies": [],
1653                    "headersSize": -1,
1654                    "bodySize": 0
1655                },
1656                "response": {
1657                    "status": status,
1658                    "statusText": "bridge",
1659                    "httpVersion": "HTTP/2",
1660                    "headers": [],
1661                    "cookies": [],
1662                    "content": {"size": text.len(), "mimeType": "text/html", "text": text},
1663                    "redirectURL": "",
1664                    "headersSize": -1,
1665                    "bodySize": 0
1666                },
1667                "cache": {},
1668                "timings": {
1669                    "blocked": 0,
1670                    "dns": 0,
1671                    "connect": 0,
1672                    "send": 0,
1673                    "wait": 0,
1674                    "receive": 0,
1675                    "ssl": 0
1676                }
1677            }]
1678        }
1679    })
1680    .to_string()
1681}
1682
1683#[cfg(feature = "acquisition-runner")]
1684fn suggest_mode_from_slo(
1685    url: &str,
1686    status_code: Option<u16>,
1687    html_excerpt: Option<&str>,
1688    target_class: TargetClass,
1689) -> Option<AcquisitionMode> {
1690    let status = status_code.unwrap_or(200);
1691    let har = build_status_only_har(url, status, html_excerpt);
1692    let report = investigate_har(&har).ok()?;
1693    let requirements = infer_requirements_with_target_class(&report, target_class);
1694    let policy = build_runtime_policy(&report, &requirements);
1695    let mapped = map_runtime_policy(&policy);
1696    Some(mode_from_hint(mapped.mode))
1697}
1698
1699#[cfg(feature = "acquisition-runner")]
1700static ACQUISITION_BRIDGE_POOL: OnceCell<Arc<BrowserPool>> = OnceCell::const_new();
1701
1702#[cfg(feature = "acquisition-runner")]
1703async fn acquisition_bridge_pool() -> Result<Arc<BrowserPool>, String> {
1704    let pool = ACQUISITION_BRIDGE_POOL
1705        .get_or_try_init(|| async {
1706            BrowserPool::new(BrowserConfig::default())
1707                .await
1708                .map_err(|e| format!("acquisition bridge browser pool init failed: {e}"))
1709        })
1710        .await?;
1711    Ok(Arc::clone(pool))
1712}
1713
1714#[cfg(feature = "acquisition-runner")]
1715async fn run_acquisition_bridge(url: &str, cfg: &AcquisitionNodeConfig) -> Result<Value, String> {
1716    let configured_mode = parse_acquisition_mode(&cfg.mode)?;
1717    let pool = acquisition_bridge_pool().await?;
1718
1719    let runner = AcquisitionRunner::new(pool);
1720    let total_timeout = cfg
1721        .total_timeout
1722        .unwrap_or_else(|| AcquisitionRequest::default().total_timeout);
1723
1724    let mut result = runner
1725        .run(AcquisitionRequest {
1726            url: url.to_string(),
1727            mode: configured_mode,
1728            wait_for_selector: cfg.wait_for_selector.clone(),
1729            extraction_js: cfg.extraction_js.clone(),
1730            total_timeout,
1731            ..AcquisitionRequest::default()
1732        })
1733        .await;
1734
1735    let mut effective_mode = configured_mode;
1736    let mut slo_recommended_mode: Option<AcquisitionMode> = None;
1737    let mut slo_bridge_applied = false;
1738
1739    if let Some(target_class) = cfg.target_class
1740        && let Some(recommended_mode) = suggest_mode_from_slo(
1741            result.final_url.as_deref().unwrap_or(url),
1742            result.status_code,
1743            result.html_excerpt.as_deref(),
1744            target_class,
1745        )
1746    {
1747        slo_recommended_mode = Some(recommended_mode);
1748        if mode_rank(recommended_mode) > mode_rank(configured_mode) {
1749            let retried = runner
1750                .run(AcquisitionRequest {
1751                    url: url.to_string(),
1752                    mode: recommended_mode,
1753                    wait_for_selector: cfg.wait_for_selector.clone(),
1754                    extraction_js: cfg.extraction_js.clone(),
1755                    total_timeout,
1756                    ..AcquisitionRequest::default()
1757                })
1758                .await;
1759            if retried.success || !result.success {
1760                result = retried;
1761                effective_mode = recommended_mode;
1762                slo_bridge_applied = true;
1763            }
1764        }
1765    }
1766
1767    let strategy_used = serde_json::to_value(result.strategy_used).unwrap_or(Value::Null);
1768    let attempted = serde_json::to_value(&result.attempted).unwrap_or(Value::Array(Vec::new()));
1769    let failures = serde_json::to_value(&result.failures).unwrap_or(Value::Array(Vec::new()));
1770
1771    Ok(json!({
1772        "data": {
1773            "success": result.success,
1774            "strategy_used": strategy_used,
1775            "final_url": result.final_url,
1776            "status_code": result.status_code,
1777            "extracted": result.extracted,
1778            "html_excerpt": result.html_excerpt,
1779        },
1780        "metadata": {
1781            "acquisition_runner": true,
1782            "diagnostics": {
1783                "attempted": attempted,
1784                "timed_out": result.timed_out,
1785                "failure_count": result.failures.len(),
1786                "failures": failures,
1787                "configured_mode": format!("{configured_mode:?}"),
1788                "effective_mode": format!("{effective_mode:?}"),
1789                "slo_target_class": cfg.target_class.map(|tc| format!("{tc:?}")),
1790                "slo_recommended_mode": slo_recommended_mode.map(|mode| format!("{mode:?}")),
1791                "slo_bridge_applied": slo_bridge_applied,
1792            }
1793        }
1794    }))
1795}
1796
1797#[cfg(not(feature = "acquisition-runner"))]
1798#[allow(clippy::unused_async)]
1799async fn run_acquisition_bridge(_url: &str, _cfg: &AcquisitionNodeConfig) -> Result<Value, String> {
1800    Err(
1801        "acquisition bridge requested but stygian-graph was built without feature 'acquisition-runner'"
1802            .to_string(),
1803    )
1804}
1805
1806/// Execute a single pipeline node of a given service kind.
1807///
1808/// Returns `None` if the kind is not supported (node is skipped);
1809/// returns `Some(Ok(value))` on success or `Some(Err(message))` on failure.
1810async fn execute_pipeline_node(
1811    kind: &str,
1812    url: &str,
1813    node_name: &str,
1814    node: &NodeDecl,
1815    timeout_secs: u64,
1816) -> Option<Result<Value, String>> {
1817    execute_pipeline_node_with(
1818        kind,
1819        url,
1820        node_name,
1821        node,
1822        timeout_secs,
1823        |bridge_url, cfg| async move { run_acquisition_bridge(&bridge_url, &cfg).await },
1824    )
1825    .await
1826}
1827
1828async fn execute_pipeline_node_with<F, Fut>(
1829    kind: &str,
1830    url: &str,
1831    node_name: &str,
1832    node: &NodeDecl,
1833    timeout_secs: u64,
1834    run_acquisition: F,
1835) -> Option<Result<Value, String>>
1836where
1837    F: Fn(String, AcquisitionNodeConfig) -> Fut + Send + Sync,
1838    Fut: Future<Output = Result<Value, String>> + Send,
1839{
1840    match kind {
1841        "http" => {
1842            let config = HttpConfig {
1843                timeout: Duration::from_secs(timeout_secs),
1844                // T112: explicit opt-in — same rationale as the
1845                // scrape tool above. The rotation pool emits
1846                // browser-class UAs only.
1847                allow_plain_http: true,
1848                ..HttpConfig::default()
1849            };
1850            let adapter = HttpAdapter::with_config(config);
1851            let input = ServiceInput {
1852                url: url.to_string(),
1853                params: json!({}),
1854            };
1855            Some(
1856                adapter
1857                    .execute(input)
1858                    .await
1859                    .map(|out| json!({ "data": out.data, "metadata": out.metadata }))
1860                    .map_err(|e| e.to_string()),
1861            )
1862        }
1863        "rest" => {
1864            let params = build_rest_params_from_node(node);
1865            let adapter = RestApiAdapter::new();
1866            let input = ServiceInput {
1867                url: url.to_string(),
1868                params,
1869            };
1870            Some(
1871                adapter
1872                    .execute(input)
1873                    .await
1874                    .map(|out| json!({ "data": out.data, "metadata": out.metadata }))
1875                    .map_err(|e| e.to_string()),
1876            )
1877        }
1878        "graphql" => {
1879            let (service, input) = build_graphql_node_request(node, url, timeout_secs);
1880            Some(
1881                service
1882                    .execute(input)
1883                    .await
1884                    .map(|out| json!({ "data": out.data, "metadata": out.metadata }))
1885                    .map_err(|e| e.to_string()),
1886            )
1887        }
1888        "sitemap" => {
1889            let max_depth = node
1890                .params
1891                .get("max_depth")
1892                .and_then(toml::Value::as_integer)
1893                .map_or(5, |v| usize::try_from(v).unwrap_or(5));
1894            let client = reqwest::Client::new();
1895            let adapter = SitemapAdapter::new(client, max_depth);
1896            let input = ServiceInput {
1897                url: url.to_string(),
1898                params: json!({}),
1899            };
1900            Some(
1901                adapter
1902                    .execute(input)
1903                    .await
1904                    .map(|out| json!({ "data": out.data, "metadata": out.metadata }))
1905                    .map_err(|e| e.to_string()),
1906            )
1907        }
1908        "rss" => {
1909            let client = reqwest::Client::new();
1910            let adapter = RssFeedAdapter::new(client);
1911            let input = ServiceInput {
1912                url: url.to_string(),
1913                params: json!({}),
1914            };
1915            Some(
1916                adapter
1917                    .execute(input)
1918                    .await
1919                    .map(|out| json!({ "data": out.data, "metadata": out.metadata }))
1920                    .map_err(|e| e.to_string()),
1921            )
1922        }
1923        "browser" => execute_browser_pipeline_node(node, node_name, url, &run_acquisition).await,
1924        other => {
1925            warn!(
1926                kind = other,
1927                node = node_name,
1928                "skipping unsupported service kind in pipeline_run"
1929            );
1930            None
1931        }
1932    }
1933}
1934
1935#[cfg(feature = "acquisition-runner")]
1936async fn execute_browser_pipeline_node<F, Fut>(
1937    node: &NodeDecl,
1938    node_name: &str,
1939    url: &str,
1940    run_acquisition: &F,
1941) -> Option<Result<Value, String>>
1942where
1943    F: Fn(String, AcquisitionNodeConfig) -> Fut + Send + Sync,
1944    Fut: Future<Output = Result<Value, String>> + Send,
1945{
1946    let cfg = match acquisition_config_from_node(node) {
1947        Ok(Some(cfg)) => cfg,
1948        Ok(None) => return None,
1949        Err(err) => {
1950            return Some(Err(format!(
1951                "Invalid acquisition config for node '{node_name}': {err}"
1952            )));
1953        }
1954    };
1955
1956    Some(run_acquisition(url.to_string(), cfg).await)
1957}
1958
1959#[cfg(not(feature = "acquisition-runner"))]
1960#[allow(clippy::unused_async)]
1961async fn execute_browser_pipeline_node<F, Fut>(
1962    _node: &NodeDecl,
1963    _node_name: &str,
1964    _url: &str,
1965    _run_acquisition: &F,
1966) -> Option<Result<Value, String>>
1967where
1968    F: Fn(String, AcquisitionNodeConfig) -> Fut + Send + Sync,
1969    Fut: Future<Output = Result<Value, String>> + Send,
1970{
1971    None
1972}
1973
1974/// Convert a [`toml::Value`] to a [`serde_json::Value`] for adapter params.
1975fn toml_to_json(v: &toml::Value) -> Value {
1976    match v {
1977        toml::Value::String(s) => Value::String(s.clone()),
1978        toml::Value::Integer(i) => Value::Number((*i).into()),
1979        toml::Value::Float(f) => {
1980            serde_json::Number::from_f64(*f).map_or(Value::Null, Value::Number)
1981        }
1982        toml::Value::Boolean(b) => Value::Bool(*b),
1983        toml::Value::Array(arr) => Value::Array(arr.iter().map(toml_to_json).collect()),
1984        toml::Value::Table(tbl) => Value::Object(
1985            tbl.iter()
1986                .map(|(k, v)| (k.clone(), toml_to_json(v)))
1987                .collect(),
1988        ),
1989        toml::Value::Datetime(dt) => Value::String(dt.to_string()),
1990    }
1991}
1992
1993/// Build a `RestApiAdapter`-compatible `params` JSON value from a pipeline [`NodeDecl`].
1994///
1995/// Forwards all recognised REST parameters from the node's TOML declaration:
1996/// `method`, `auth`, `headers`, `query`, `body`, `pagination`, and `data_path`.
1997fn build_rest_params_from_node(node: &NodeDecl) -> Value {
1998    let mut map = serde_json::Map::new();
1999
2000    if let Some(method) = node.params.get("method").and_then(|v| v.as_str()) {
2001        map.insert("method".to_owned(), json!(method));
2002    }
2003    if let Some(auth) = node.params.get("auth") {
2004        map.insert("auth".to_owned(), toml_to_json(auth));
2005    }
2006    if let Some(headers) = node.params.get("headers") {
2007        map.insert("headers".to_owned(), toml_to_json(headers));
2008    }
2009    if let Some(query) = node.params.get("query") {
2010        map.insert("query".to_owned(), toml_to_json(query));
2011    }
2012    if let Some(body) = node.params.get("body") {
2013        map.insert("body".to_owned(), toml_to_json(body));
2014    }
2015    if let Some(pagination) = node.params.get("pagination") {
2016        map.insert("pagination".to_owned(), toml_to_json(pagination));
2017    }
2018    if let Some(dp) = node.params.get("data_path").and_then(|v| v.as_str()) {
2019        map.insert("response".to_owned(), json!({ "data_path": dp }));
2020    }
2021
2022    Value::Object(map)
2023}
2024
2025// ─── Tests ───────────────────────────────────────────────────────────────────
2026
2027#[cfg(test)]
2028#[allow(clippy::unwrap_used)]
2029mod tests {
2030    use super::*;
2031
2032    #[test]
2033    fn server_builds() {
2034        let _ = McpGraphServer::new();
2035    }
2036
2037    #[test]
2038    fn discover_response_advertises_protocol_version() {
2039        let id = json!(1);
2040        let resp = McpGraphServer::handle_discover(&id);
2041        assert_eq!(
2042            resp.pointer("/result/protocolVersion")
2043                .and_then(Value::as_str),
2044            Some("2026-07-28")
2045        );
2046        assert_eq!(
2047            resp.pointer("/result/supportedProtocolVersions")
2048                .and_then(Value::as_array)
2049                .and_then(|v| v.first())
2050                .and_then(Value::as_str),
2051            Some("2026-07-28")
2052        );
2053        assert_eq!(
2054            resp.pointer("/result/serverInfo/name")
2055                .and_then(Value::as_str),
2056            Some("stygian-graph")
2057        );
2058        assert_eq!(
2059            resp.pointer("/result/resultType").and_then(Value::as_str),
2060            Some("complete")
2061        );
2062        // MCP 2026-07-28 §8: extensions array present even when empty.
2063        assert!(
2064            resp.pointer("/result/extensions")
2065                .and_then(Value::as_array)
2066                .is_some()
2067        );
2068    }
2069
2070    #[test]
2071    fn initialize_method_is_no_longer_recognized() {
2072        // MCP 2026-07-28 removed the `initialize` handshake. The dispatcher
2073        // should return `Method not found` rather than silently re-creating
2074        // the legacy envelope.
2075        let resp = tokio_test::block_on(McpGraphServer::handle_request(&json!({
2076            "jsonrpc": "2.0",
2077            "id": 1,
2078            "method": "initialize",
2079            "params": {}
2080        })));
2081        assert_eq!(
2082            resp.pointer("/error/code").and_then(Value::as_i64),
2083            Some(-32601)
2084        );
2085    }
2086
2087    #[test]
2088    fn ping_method_is_no_longer_recognized() {
2089        // `ping` is removed in MCP 2026-07-28. Confirm the dispatcher returns
2090        // a clean `Method not found` error rather than a stale empty result.
2091        let resp = tokio_test::block_on(McpGraphServer::handle_request(&json!({
2092            "jsonrpc": "2.0",
2093            "id": 7,
2094            "method": "ping"
2095        })));
2096        assert_eq!(
2097            resp.pointer("/error/code").and_then(Value::as_i64),
2098            Some(-32601)
2099        );
2100    }
2101
2102    #[test]
2103    fn ok_response_threads_result_type_complete() {
2104        // MCP 2026-07-28 §8: every `ok_response` envelope carries a
2105        // `resultType: "complete"` field, even when the result is not an
2106        // object. The non-object branch is defensive — the spec only defines
2107        // object results, but a bug in a caller should not produce an invalid
2108        // envelope.
2109        let id = json!(42);
2110        let obj = McpGraphServer::handle_tools_list(&id);
2111        assert_eq!(
2112            obj.pointer("/result/resultType").and_then(Value::as_str),
2113            Some("complete")
2114        );
2115        assert!(obj.pointer("/result/tools").is_some());
2116
2117        let scalar = ok_response(&id, json!("plain-string"));
2118        assert_eq!(
2119            scalar.pointer("/result/resultType").and_then(Value::as_str),
2120            Some("complete")
2121        );
2122        assert_eq!(
2123            scalar.pointer("/result/value").and_then(Value::as_str),
2124            Some("plain-string")
2125        );
2126    }
2127
2128    #[test]
2129    fn extract_meta_reads_namespaced_keys() {
2130        // The `_meta` reader must look under `params._meta.io.modelcontextprotocol/*`.
2131        // `protocolVersion`, `clientInfo`, `clientCapabilities` are the three
2132        // carriers defined by MCP 2026-07-28 §2.
2133        let req = json!({
2134            "jsonrpc": "2.0",
2135            "id": 1,
2136            "method": "tools/list",
2137            "params": {
2138                "_meta": {
2139                    "io.modelcontextprotocol/protocolVersion": "2026-07-28",
2140                    "io.modelcontextprotocol/clientInfo": {
2141                        "name": "test-client",
2142                        "version": "0.0.1"
2143                    },
2144                    "io.modelcontextprotocol/clientCapabilities": {
2145                        "tools": {}
2146                    },
2147                    "unrelated": "ignored"
2148                }
2149            }
2150        });
2151        assert_eq!(
2152            extract_client_protocol_version(&req).as_deref(),
2153            Some("2026-07-28")
2154        );
2155        assert_eq!(
2156            extract_meta(&req, "clientInfo")
2157                .and_then(|v| v.get("name"))
2158                .and_then(Value::as_str),
2159            Some("test-client")
2160        );
2161        assert!(extract_meta(&req, "clientCapabilities").is_some());
2162        assert!(extract_meta(&req, "not-a-key").is_none());
2163
2164        // `_meta` absent → all helpers return None.
2165        let bare = json!({"jsonrpc": "2.0", "id": 1, "method": "tools/list"});
2166        assert!(extract_client_protocol_version(&bare).is_none());
2167        assert!(extract_meta(&bare, "protocolVersion").is_none());
2168    }
2169
2170    #[test]
2171    fn is_supported_protocol_version_accepts_listed_and_rejects_others() {
2172        // MCP 2026-07-28 §2: clients advertise their protocol version on every
2173        // request. Servers reject unknown versions with
2174        // `UnsupportedProtocolVersionError` (code `-32022`). The aggregator
2175        // (PR 4 of MCP-001) wires this check into the per-request gate; the
2176        // graph server itself stays permissive at this layer because it is
2177        // also called directly via `tools/list` / `tools/call` and the
2178        // dispatcher doesn't enforce `_meta` yet.
2179        assert!(is_supported_protocol_version("2026-07-28", &["2026-07-28"]).is_ok());
2180        assert!(is_supported_protocol_version("2026-07-28", &["2025-11-25"]).is_err());
2181        assert!(is_supported_protocol_version("2025-11-25", &["2026-07-28", "2025-11-25"]).is_ok());
2182    }
2183
2184    #[test]
2185    fn tools_list_contains_all_tools() {
2186        let id = json!(1);
2187        let resp = McpGraphServer::handle_tools_list(&id);
2188        let tools = resp
2189            .pointer("/result/tools")
2190            .and_then(Value::as_array)
2191            .unwrap();
2192        let names: Vec<&str> = tools
2193            .iter()
2194            .map(|t| t.get("name").and_then(Value::as_str).unwrap())
2195            .collect();
2196        assert!(names.contains(&"scrape"));
2197        assert!(names.contains(&"scrape_rest"));
2198        assert!(names.contains(&"scrape_graphql"));
2199        assert!(names.contains(&"scrape_sitemap"));
2200        assert!(names.contains(&"scrape_rss"));
2201        assert!(names.contains(&"pipeline_validate"));
2202        assert!(names.contains(&"pipeline_run"));
2203
2204        #[cfg(feature = "charon")]
2205        {
2206            assert!(names.contains(&"charon_classify_transaction"));
2207            assert!(names.contains(&"charon_investigate_har"));
2208            assert!(names.contains(&"charon_infer_requirements"));
2209            assert!(names.contains(&"charon_build_runtime_policy"));
2210            assert!(names.contains(&"charon_map_runtime_policy"));
2211            assert!(names.contains(&"charon_analyze_and_plan"));
2212        }
2213    }
2214
2215    #[cfg(feature = "charon")]
2216    #[test]
2217    fn charon_classify_transaction_returns_detection() {
2218        let id = json!(99);
2219        let args = json!({
2220            "url": "https://example.com/challenge",
2221            "status": 403,
2222            "response_headers": { "x-datadome": "1" },
2223            "response_body_snippet": "captcha-delivery.com"
2224        });
2225
2226        let resp = McpGraphServer::tool_charon_classify_transaction(&id, &args);
2227        let text = resp
2228            .pointer("/result/content/0/text")
2229            .and_then(Value::as_str)
2230            .unwrap_or_default();
2231        let payload: Value = serde_json::from_str(text).unwrap_or(Value::Null);
2232
2233        assert_eq!(
2234            payload
2235                .pointer("/detection/provider")
2236                .and_then(Value::as_str),
2237            Some("DataDome")
2238        );
2239    }
2240
2241    #[cfg(feature = "charon")]
2242    #[test]
2243    fn charon_analyze_and_plan_returns_policy_and_acquisition() {
2244        let id = json!(100);
2245        let args = json!({
2246            "har": json!({
2247                "log": {
2248                    "version": "1.2",
2249                    "creator": {"name": "test", "version": "1.0"},
2250                    "pages": [{
2251                        "id": "page_1",
2252                        "title": "https://example.com/challenge",
2253                        "startedDateTime": "2026-01-01T00:00:00.000Z",
2254                        "pageTimings": {"onLoad": 0}
2255                    }],
2256                    "entries": [{
2257                        "pageref": "page_1",
2258                        "startedDateTime": "2026-01-01T00:00:00.000Z",
2259                        "time": 0,
2260                        "request": {
2261                            "method": "GET",
2262                            "url": "https://example.com/challenge",
2263                            "httpVersion": "HTTP/2",
2264                            "headers": [],
2265                            "queryString": [],
2266                            "cookies": [],
2267                            "headersSize": -1,
2268                            "bodySize": 0
2269                        },
2270                        "response": {
2271                            "status": 403,
2272                            "statusText": "Forbidden",
2273                            "httpVersion": "HTTP/2",
2274                            "headers": [],
2275                            "cookies": [],
2276                            "content": {
2277                                "size": 0,
2278                                "mimeType": "text/html",
2279                                "text": "captcha-delivery.com"
2280                            },
2281                            "redirectURL": "",
2282                            "headersSize": -1,
2283                            "bodySize": 0
2284                        },
2285                        "cache": {},
2286                        "timings": {
2287                            "blocked": 0,
2288                            "dns": 0,
2289                            "connect": 0,
2290                            "send": 0,
2291                            "wait": 0,
2292                            "receive": 0,
2293                            "ssl": 0
2294                        }
2295                    }]
2296                }
2297            }).to_string(),
2298            "target_class": "api"
2299        });
2300
2301        let resp = McpGraphServer::tool_charon_analyze_and_plan(&id, &args);
2302        let text = resp
2303            .pointer("/result/content/0/text")
2304            .and_then(Value::as_str)
2305            .unwrap_or_default();
2306        let payload: Value = serde_json::from_str(text).unwrap_or(Value::Null);
2307
2308        assert!(payload.get("bundle").is_some());
2309        assert!(payload.pointer("/bundle/policy").is_some());
2310        assert!(payload.pointer("/acquisition/mode").is_some());
2311    }
2312
2313    #[test]
2314    fn pipeline_validate_rejects_bad_toml() {
2315        let id = json!(1);
2316        let args = json!({ "toml": "this is not valid toml [[[[" });
2317        let resp = McpGraphServer::tool_pipeline_validate(&id, &args);
2318        assert!(
2319            resp.get("error").is_some_and(Value::is_object)
2320                || resp
2321                    .pointer("/result/content/0/text")
2322                    .and_then(Value::as_str)
2323                    .unwrap_or("")
2324                    .contains("false")
2325        );
2326    }
2327
2328    #[test]
2329    fn pipeline_validate_accepts_valid_pipeline() {
2330        let id = json!(1);
2331        let toml = r#"
2332[[nodes]]
2333name = "fetch"
2334service = "http"
2335url = "https://example.com"
2336
2337[[nodes]]
2338name = "process"
2339service = "http"
2340url = "https://example.com/api"
2341depends_on = ["fetch"]
2342"#;
2343        let args = json!({ "toml": toml });
2344        let resp = McpGraphServer::tool_pipeline_validate(&id, &args);
2345        let text = resp
2346            .pointer("/result/content/0/text")
2347            .and_then(Value::as_str)
2348            .unwrap();
2349        let parsed: Value = serde_json::from_str(text).unwrap();
2350        assert_eq!(parsed.get("valid"), Some(&json!(true)));
2351        assert_eq!(parsed.get("node_count"), Some(&json!(2)));
2352    }
2353
2354    #[test]
2355    fn pipeline_validate_missing_toml_returns_error() {
2356        // Test is sync — just check param validation path
2357        let id = json!(1);
2358        let args = json!({});
2359        // We can't await in a non-async test context without a runtime,
2360        // so just verify the tool_pipeline_validate path for missing param.
2361        let resp = McpGraphServer::tool_pipeline_validate(&id, &args);
2362        assert!(resp.get("error").is_some_and(Value::is_object));
2363    }
2364
2365    #[tokio::test]
2366    async fn pipeline_browser_node_without_acquisition_is_skipped() {
2367        let node = NodeDecl {
2368            name: "render".to_string(),
2369            service: "browser".to_string(),
2370            depends_on: Vec::new(),
2371            url: Some("https://example.com".to_string()),
2372            params: HashMap::new(),
2373        };
2374
2375        let result = execute_pipeline_node_with(
2376            "browser",
2377            "https://example.com",
2378            "render",
2379            &node,
2380            30,
2381            |_url, _cfg| async { Ok(json!({"data": "should-not-run"})) },
2382        )
2383        .await;
2384
2385        assert!(result.is_none());
2386    }
2387
2388    #[cfg(feature = "acquisition-runner")]
2389    #[tokio::test]
2390    async fn pipeline_browser_node_with_acquisition_uses_bridge_path() {
2391        let mut acquisition = toml::map::Map::new();
2392        acquisition.insert("mode".to_string(), toml::Value::String("fast".to_string()));
2393        acquisition.insert(
2394            "wait_for_selector".to_string(),
2395            toml::Value::String("main".to_string()),
2396        );
2397
2398        let mut params = HashMap::new();
2399        params.insert("acquisition".to_string(), toml::Value::Table(acquisition));
2400
2401        let node = NodeDecl {
2402            name: "render".to_string(),
2403            service: "browser".to_string(),
2404            depends_on: Vec::new(),
2405            url: Some("https://example.com".to_string()),
2406            params,
2407        };
2408
2409        let result = execute_pipeline_node_with(
2410            "browser",
2411            "https://example.com",
2412            "render",
2413            &node,
2414            30,
2415            |url, cfg| async move {
2416                Ok(json!({
2417                    "data": {
2418                        "url": url,
2419                        "mode": cfg.mode,
2420                        "wait_for_selector": cfg.wait_for_selector,
2421                    },
2422                    "metadata": {"bridge": "mock"}
2423                }))
2424            },
2425        )
2426        .await;
2427
2428        let payload = match result {
2429            Some(Ok(payload)) => payload,
2430            other => {
2431                assert!(
2432                    matches!(other, Some(Ok(_))),
2433                    "browser acquisition should return Some(Ok(_))"
2434                );
2435                return;
2436            }
2437        };
2438
2439        assert_eq!(
2440            payload.pointer("/data/url").and_then(Value::as_str),
2441            Some("https://example.com")
2442        );
2443        assert_eq!(
2444            payload.pointer("/data/mode").and_then(Value::as_str),
2445            Some("fast")
2446        );
2447        assert_eq!(
2448            payload
2449                .pointer("/data/wait_for_selector")
2450                .and_then(Value::as_str),
2451            Some("main")
2452        );
2453    }
2454
2455    #[cfg(feature = "acquisition-runner")]
2456    #[test]
2457    fn acquisition_config_parses_target_class() {
2458        let mut acquisition = toml::map::Map::new();
2459        acquisition.insert(
2460            "mode".to_string(),
2461            toml::Value::String("resilient".to_string()),
2462        );
2463        acquisition.insert(
2464            "target_class".to_string(),
2465            toml::Value::String("content-site".to_string()),
2466        );
2467
2468        let mut params = HashMap::new();
2469        params.insert("acquisition".to_string(), toml::Value::Table(acquisition));
2470
2471        let node = NodeDecl {
2472            name: "render".to_string(),
2473            service: "browser".to_string(),
2474            depends_on: Vec::new(),
2475            url: Some("https://example.com".to_string()),
2476            params,
2477        };
2478
2479        let parsed = acquisition_config_from_node(&node);
2480        assert!(parsed.is_ok(), "target_class should parse");
2481        let Ok(Some(cfg)) = parsed else {
2482            return;
2483        };
2484        assert_eq!(cfg.target_class, Some(TargetClass::ContentSite));
2485    }
2486
2487    #[cfg(feature = "acquisition-runner")]
2488    #[test]
2489    fn slo_bridge_can_recommend_stronger_mode_for_blocked_status() {
2490        let recommended = suggest_mode_from_slo(
2491            "https://example.com/challenge",
2492            Some(403),
2493            Some("captcha-delivery.com"),
2494            TargetClass::Api,
2495        );
2496
2497        assert!(recommended.is_some(), "SLO bridge should return a mode");
2498        let Some(mode) = recommended else {
2499            return;
2500        };
2501        assert!(
2502            mode_rank(mode) >= mode_rank(AcquisitionMode::Resilient),
2503            "blocked scenarios should not downshift below resilient"
2504        );
2505    }
2506
2507    #[cfg(not(feature = "acquisition-runner"))]
2508    #[tokio::test]
2509    async fn pipeline_browser_node_with_acquisition_is_skipped_without_feature() {
2510        let mut acquisition = toml::map::Map::new();
2511        acquisition.insert("mode".to_string(), toml::Value::String("fast".to_string()));
2512
2513        let mut params = HashMap::new();
2514        params.insert("acquisition".to_string(), toml::Value::Table(acquisition));
2515
2516        let node = NodeDecl {
2517            name: "render".to_string(),
2518            service: "browser".to_string(),
2519            depends_on: Vec::new(),
2520            url: Some("https://example.com".to_string()),
2521            params,
2522        };
2523
2524        let result = execute_pipeline_node_with(
2525            "browser",
2526            "https://example.com",
2527            "render",
2528            &node,
2529            30,
2530            |_url, _cfg| async { Ok(json!({"data": "should-not-run"})) },
2531        )
2532        .await;
2533
2534        assert!(result.is_none());
2535    }
2536
2537    #[cfg(feature = "acquisition-runner")]
2538    #[tokio::test]
2539    async fn pipeline_browser_node_invalid_acquisition_timeout_returns_error() {
2540        let mut acquisition = toml::map::Map::new();
2541        acquisition.insert("mode".to_string(), toml::Value::String("fast".to_string()));
2542        acquisition.insert("total_timeout_secs".to_string(), toml::Value::Integer(0));
2543
2544        let mut params = HashMap::new();
2545        params.insert("acquisition".to_string(), toml::Value::Table(acquisition));
2546
2547        let node = NodeDecl {
2548            name: "render".to_string(),
2549            service: "browser".to_string(),
2550            depends_on: Vec::new(),
2551            url: Some("https://example.com".to_string()),
2552            params,
2553        };
2554
2555        let result = execute_pipeline_node_with(
2556            "browser",
2557            "https://example.com",
2558            "render",
2559            &node,
2560            30,
2561            |_url, _cfg| async { Ok(json!({"data": "unexpected"})) },
2562        )
2563        .await;
2564
2565        let err = match result {
2566            Some(Err(err)) => err,
2567            other => {
2568                assert!(
2569                    matches!(other, Some(Err(_))),
2570                    "invalid config should return Some(Err(_))"
2571                );
2572                return;
2573            }
2574        };
2575
2576        assert!(
2577            err.contains("total_timeout_secs") || err.contains("Invalid acquisition config"),
2578            "unexpected error: {err}"
2579        );
2580    }
2581}