1use 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
103fn 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 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#[allow(dead_code)] fn 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#[allow(dead_code)] fn 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#[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
221pub struct McpGraphServer;
239
240impl McpGraphServer {
241 #[must_use]
243 pub const fn new() -> Self {
244 Self
245 }
246
247 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; }
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 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 "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 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 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 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 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 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 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 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 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 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 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 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 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 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 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
1442fn 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
1806async 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 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
1974fn 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
1993fn 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#[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 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 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 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 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 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 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 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 let id = json!(1);
2358 let args = json!({});
2359 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}