Skip to main content

stygian_graph/domain/
pipeline.rs

1//! Pipeline types with typestate pattern
2//!
3//! The typestate pattern ensures pipelines can only transition through valid states:
4//! Unvalidated → Validated → Executing → Complete
5//!
6//! # Example
7//!
8//! ```
9//! use stygian_graph::domain::pipeline::PipelineUnvalidated;
10//! use serde_json::json;
11//!
12//! # fn example() -> Result<(), Box<dyn std::error::Error>> {
13//! let unvalidated = PipelineUnvalidated::new(json!({"nodes": []}));
14//! let validated = unvalidated.validate()?;
15//! let executing = validated.execute();
16//! let complete = executing.complete(json!({"status": "success"}));
17//! # Ok(())
18//! # }
19//! ```
20
21use serde::{Deserialize, Serialize};
22
23use super::error::{GraphError, StygianError};
24use super::policy::RobotsPolicy;
25use crate::ports::robots_policy::{
26    PolicyOutcome, RobotsPolicyGuard, apply_policy, validate_guard_pair,
27};
28
29/// Pipeline in unvalidated state
30///
31/// Initial state after loading configuration from a file or API.
32/// Must be validated before execution.
33#[derive(Debug, Clone, Serialize, Deserialize)]
34pub struct PipelineUnvalidated {
35    /// Pipeline configuration (unvalidated)
36    pub config: serde_json::Value,
37}
38
39/// Pipeline in validated state
40///
41/// Configuration has been validated and is ready for execution.
42#[derive(Debug, Clone)]
43pub struct PipelineValidated {
44    /// Validated configuration
45    pub config: serde_json::Value,
46}
47
48/// Pipeline in executing state
49///
50/// Pipeline is actively being executed. Contains runtime context.
51#[derive(Debug)]
52pub struct PipelineExecuting {
53    /// Execution context and state
54    pub context: serde_json::Value,
55}
56
57/// Pipeline in completed state
58///
59/// Pipeline execution has finished. Contains final results.
60#[derive(Debug)]
61pub struct PipelineComplete {
62    /// Execution results
63    pub results: serde_json::Value,
64}
65
66impl PipelineUnvalidated {
67    /// Create a new unvalidated pipeline from raw configuration
68    ///
69    /// # Example
70    ///
71    /// ```
72    /// use stygian_graph::domain::pipeline::PipelineUnvalidated;
73    /// use serde_json::json;
74    ///
75    /// let pipeline = PipelineUnvalidated::new(json!({
76    ///     "nodes": [{"id": "fetch", "service": "http"}],
77    ///     "edges": []
78    /// }));
79    /// ```
80    #[must_use]
81    pub const fn new(config: serde_json::Value) -> Self {
82        Self { config }
83    }
84
85    /// Validate the pipeline configuration
86    ///
87    /// Transitions from `Unvalidated` to `Validated` state.
88    ///
89    /// # Panics
90    ///
91    /// Panics if the validated DAG contains an edge whose source node is missing
92    /// from the adjacency map. This is guarded by the cycle check above and is
93    /// unreachable in well-formed input.
94    ///
95    /// # Errors
96    ///
97    /// Returns `GraphError::InvalidPipeline` if validation fails.
98    ///
99    /// # Example
100    ///
101    /// ```
102    /// use stygian_graph::domain::pipeline::PipelineUnvalidated;
103    /// use serde_json::json;
104    ///
105    /// # fn example() -> Result<(), Box<dyn std::error::Error>> {
106    /// let pipeline = PipelineUnvalidated::new(json!({"nodes": []}));
107    /// let validated = pipeline.validate()?;
108    /// # Ok(())
109    /// # }
110    /// ```
111    #[allow(clippy::too_many_lines, clippy::unwrap_used, clippy::indexing_slicing)]
112    pub fn validate(self) -> Result<PipelineValidated, StygianError> {
113        use std::collections::{HashMap, HashSet, VecDeque};
114
115        // Extract nodes and edges from config
116        let nodes = self
117            .config
118            .get("nodes")
119            .and_then(|n| n.as_array())
120            .ok_or_else(|| {
121                GraphError::InvalidPipeline("Pipeline must contain a 'nodes' array".to_string())
122            })?;
123
124        let empty_edges = vec![];
125        let edges = self
126            .config
127            .get("edges")
128            .and_then(|e| e.as_array())
129            .unwrap_or(&empty_edges);
130
131        // Rule 1: At least one node
132        if nodes.is_empty() {
133            return Err(GraphError::InvalidPipeline(
134                "Pipeline must contain at least one node".to_string(),
135            )
136            .into());
137        }
138
139        // Build node map and validate individual nodes
140        let mut node_map: HashMap<String, usize> = HashMap::new();
141        let valid_services = [
142            "http",
143            "http_escalating",
144            "browser",
145            "ai_claude",
146            "ai_openai",
147            "ai_gemini",
148            "ai_github",
149            "ai_ollama",
150            "javascript",
151            "graphql",
152            "storage",
153        ];
154
155        for (idx, node) in nodes.iter().enumerate() {
156            let node_obj = node.as_object().ok_or_else(|| {
157                GraphError::InvalidPipeline(format!("Node at index {idx}: must be an object"))
158            })?;
159
160            // Rule 2 & 3: Validate node ID
161            let node_id = node_obj.get("id").and_then(|v| v.as_str()).ok_or_else(|| {
162                GraphError::InvalidPipeline(format!(
163                    "Node at index {idx}: 'id' field is required and must be a string"
164                ))
165            })?;
166
167            if node_id.is_empty() {
168                return Err(GraphError::InvalidPipeline(format!(
169                    "Node at index {idx}: id cannot be empty"
170                ))
171                .into());
172            }
173
174            // Check for duplicate node IDs
175            if node_map.insert(node_id.to_string(), idx).is_some() {
176                return Err(
177                    GraphError::InvalidPipeline(format!("Duplicate node id: '{node_id}'")).into(),
178                );
179            }
180
181            // Rule 4: Validate service type
182            let service = node_obj
183                .get("service")
184                .and_then(|v| v.as_str())
185                .ok_or_else(|| {
186                    GraphError::InvalidPipeline(format!(
187                        "Node '{node_id}': 'service' field is required and must be a string"
188                    ))
189                })?;
190
191            if !valid_services.contains(&service) {
192                return Err(GraphError::InvalidPipeline(format!(
193                    "Node '{node_id}': service type '{service}' is not recognized"
194                ))
195                .into());
196            }
197        }
198
199        // Rule 5 & 6: Validate edges
200        let mut adjacency: HashMap<String, Vec<String>> = HashMap::new();
201        let mut in_degree: HashMap<String, usize> = HashMap::new();
202
203        // Initialize in_degree for all nodes
204        for node in nodes {
205            if let Some(id) = node.get("id").and_then(|v| v.as_str()) {
206                in_degree.insert(id.to_string(), 0);
207                adjacency.insert(id.to_string(), Vec::new());
208            }
209        }
210
211        for (edge_idx, edge) in edges.iter().enumerate() {
212            let edge_obj = edge.as_object().ok_or_else(|| {
213                GraphError::InvalidPipeline(format!("Edge at index {edge_idx}: must be an object"))
214            })?;
215
216            let from = edge_obj
217                .get("from")
218                .and_then(|v| v.as_str())
219                .ok_or_else(|| {
220                    GraphError::InvalidPipeline(format!(
221                        "Edge at index {edge_idx}: 'from' field is required and must be a string"
222                    ))
223                })?;
224
225            let to = edge_obj.get("to").and_then(|v| v.as_str()).ok_or_else(|| {
226                GraphError::InvalidPipeline(format!(
227                    "Edge at index {edge_idx}: 'to' field is required and must be a string"
228                ))
229            })?;
230
231            // Source node must exist
232            if !node_map.contains_key(from) {
233                return Err(GraphError::InvalidPipeline(format!(
234                    "Edge {from} -> {to}: source node '{from}' not found"
235                ))
236                .into());
237            }
238
239            // Target node must exist
240            if !node_map.contains_key(to) {
241                return Err(GraphError::InvalidPipeline(format!(
242                    "Edge {from} -> {to}: target node '{to}' not found"
243                ))
244                .into());
245            }
246
247            // Source and target cannot be the same
248            if from == to {
249                return Err(GraphError::InvalidPipeline(format!(
250                    "Self-loop detected at node '{from}'"
251                ))
252                .into());
253            }
254
255            // Build adjacency list and track in-degrees
256            adjacency.get_mut(from).unwrap().push(to.to_string());
257            *in_degree.get_mut(to).unwrap() += 1;
258        }
259
260        // Rule 7: Detect cycles using Kahn's algorithm (topological sort)
261        let mut in_degree_copy = in_degree.clone();
262        let mut queue: VecDeque<String> = VecDeque::new();
263
264        // Add all nodes with no incoming edges (entry points)
265        let entry_points: Vec<String> = in_degree_copy
266            .iter()
267            .filter(|(_, degree)| **degree == 0)
268            .map(|(node_id, _)| node_id.clone())
269            .collect();
270        for node_id in entry_points {
271            queue.push_back(node_id);
272        }
273
274        let mut sorted_count = 0;
275        while let Some(node_id) = queue.pop_front() {
276            sorted_count += 1;
277
278            // For each neighbor of this node
279            if let Some(neighbors) = adjacency.get(&node_id) {
280                let neighbors_copy = neighbors.clone();
281                for neighbor in neighbors_copy {
282                    *in_degree_copy.get_mut(&neighbor).unwrap() -= 1;
283                    if in_degree_copy[&neighbor] == 0 {
284                        queue.push_back(neighbor);
285                    }
286                }
287            }
288        }
289
290        // If we didn't sort all nodes, there's a cycle
291        if sorted_count != node_map.len() {
292            return Err(GraphError::InvalidPipeline(
293                "Cycle detected in pipeline graph".to_string(),
294            )
295            .into());
296        }
297
298        // Rule 8: Check for unreachable nodes (isolated components)
299        // All nodes must form a single connected DAG with one or more entry points
300        // Only start reachability from the FIRST entry point to ensure all nodes are connected
301        let mut visited: HashSet<String> = HashSet::new();
302        let mut to_visit: VecDeque<String> = VecDeque::new();
303
304        // Find first entry point (node with in_degree == 0)
305        let mut entry_points = Vec::new();
306        for (node_id, degree) in &in_degree {
307            if *degree == 0 {
308                entry_points.push(node_id.clone());
309            }
310        }
311
312        if entry_points.is_empty() {
313            // Should not happen if cycle check passed, but be safe
314            return Err(GraphError::InvalidPipeline(
315                "No entry points found (all nodes have incoming edges)".to_string(),
316            )
317            .into());
318        }
319
320        // Start BFS from ONLY the first entry point to ensure single connected component
321        to_visit.push_back(entry_points[0].clone());
322
323        // BFS from first entry point
324        while let Some(node_id) = to_visit.pop_front() {
325            if visited.insert(node_id.clone()) {
326                // Explore outgoing edges
327                if let Some(neighbors) = adjacency.get(&node_id) {
328                    for neighbor in neighbors {
329                        to_visit.push_back(neighbor.clone());
330                    }
331                }
332
333                // Also explore reverse adjacency (incoming edges) to handle branching
334                for (source, targets) in &adjacency {
335                    if targets.contains(&node_id) && !visited.contains(source) {
336                        to_visit.push_back(source.clone());
337                    }
338                }
339            }
340        }
341
342        // Check for unreachable nodes
343        let all_node_ids: HashSet<String> = node_map.keys().cloned().collect();
344        let unreachable: Vec<_> = all_node_ids.difference(&visited).collect();
345
346        if !unreachable.is_empty() {
347            let unreachable_str = unreachable
348                .iter()
349                .map(|s| s.as_str())
350                .collect::<Vec<_>>()
351                .join("', '");
352            return Err(GraphError::InvalidPipeline(format!(
353                "Unreachable nodes found: '{unreachable_str}' (ensure all nodes are connected in a single DAG)"
354            ))
355            .into());
356        }
357
358        Ok(PipelineValidated {
359            config: self.config,
360        })
361    }
362
363    /// Compute the effective [`RobotsPolicy`] for this pipeline.
364    ///
365    /// Looks for `config["robots_policy"]` and
366    /// `config["pipeline"]["robots_policy"]`; falls back to
367    /// [`RobotsPolicy::Obey`] when neither is set.
368    ///
369    /// # Errors
370    ///
371    /// Returns [`StygianError::Config`] when a value is present but
372    /// cannot be parsed as a `RobotsPolicy`.
373    pub fn effective_robots_policy(&self) -> Result<RobotsPolicy, StygianError> {
374        if let Some(raw) = self.config.get("robots_policy").and_then(|v| v.as_str()) {
375            return raw.parse::<RobotsPolicy>();
376        }
377        if let Some(raw) = self
378            .config
379            .get("pipeline")
380            .and_then(|p| p.get("robots_policy"))
381            .and_then(|v| v.as_str())
382        {
383            return raw.parse::<RobotsPolicy>();
384        }
385        Ok(RobotsPolicy::Obey)
386    }
387
388    /// Apply the pipeline's [`RobotsPolicy`] to every URL in the
389    /// pipeline against the supplied [`RobotsPolicyGuard`].
390    ///
391    /// This is the **content** check that complements the structural
392    /// `validate()` above — it refuses to build a pipeline whose URLs
393    /// would be forbidden at run time. Both recon and production share
394    /// this method, so the policy cannot drift between spec-build and
395    /// execution.
396    ///
397    /// Behaviour:
398    ///
399    /// - [`RobotsPolicy::Obey`] — every URL that the guard marks
400    ///   `Forbid` or `Unknown` is refused; the pipeline fails with
401    ///   [`GraphError::InvalidPipeline`] before any production
402    ///   traffic.
403    /// - [`RobotsPolicy::IgnoreWithAudit`] — every URL is permitted;
404    ///   the [`PolicyOutcome::FetchWithAudit`] entries are surfaced
405    ///   through the returned `Vec<RobotsAuditEvent>` so the operator
406    ///   can see which URLs were ignored.
407    /// - [`RobotsPolicy::IgnoreSilently`] — every URL is permitted
408    ///   unconditionally; the returned audit vec is always empty.
409    ///
410    /// # Errors
411    ///
412    /// - [`GraphError::InvalidPipeline`] if `RobotsPolicy::Obey` is in
413    ///   effect and any URL is refused.
414    /// - The wrapped guard error if the guard itself fails to decide.
415    pub async fn check_robots_policy(
416        &self,
417        guard: &dyn RobotsPolicyGuard,
418    ) -> Result<Vec<RobotsAuditEvent>, StygianError> {
419        let policy = self.effective_robots_policy()?;
420        validate_guard_pair(policy, guard)?;
421
422        let urls = collect_node_urls(&self.config);
423        let mut audit = Vec::with_capacity(urls.len());
424
425        for (node_id, url) in &urls {
426            let decision = guard.decide(url).await?;
427            let outcome = apply_policy(policy, decision);
428            match outcome {
429                PolicyOutcome::Fetch => {}
430                PolicyOutcome::FetchWithAudit { reason } => {
431                    audit.push(RobotsAuditEvent {
432                        node_id: node_id.clone(),
433                        url: url.clone(),
434                        reason,
435                    });
436                }
437                PolicyOutcome::Refuse { reason } => {
438                    return Err(GraphError::InvalidPipeline(format!(
439                        "robots policy '{policy}' refuses node '{node_id}' \
440                         (url: {url}): {reason}"
441                    ))
442                    .into());
443                }
444            }
445        }
446
447        Ok(audit)
448    }
449}
450
451impl PipelineValidated {
452    /// Begin executing the validated pipeline
453    ///
454    /// Transitions from `Validated` to `Executing` state.
455    ///
456    /// # Example
457    ///
458    /// ```
459    /// use stygian_graph::domain::pipeline::PipelineUnvalidated;
460    /// use serde_json::json;
461    ///
462    /// # fn example() -> Result<(), Box<dyn std::error::Error>> {
463    /// let pipeline = PipelineUnvalidated::new(json!({"nodes": []}))
464    ///     .validate()?;
465    /// let executing = pipeline.execute();
466    /// # Ok(())
467    /// # }
468    /// ```
469    #[must_use]
470    pub fn execute(self) -> PipelineExecuting {
471        PipelineExecuting {
472            context: self.config,
473        }
474    }
475}
476
477impl PipelineExecuting {
478    /// Mark the pipeline as complete with results
479    ///
480    /// Transitions from `Executing` to `Complete` state.
481    ///
482    /// # Example
483    ///
484    /// ```
485    /// use stygian_graph::domain::pipeline::PipelineUnvalidated;
486    /// use serde_json::json;
487    ///
488    /// # fn example() -> Result<(), Box<dyn std::error::Error>> {
489    /// let pipeline = PipelineUnvalidated::new(json!({"nodes": []}))
490    ///     .validate()?
491    ///     .execute();
492    ///
493    /// let complete = pipeline.complete(json!({"status": "success"}));
494    /// # Ok(())
495    /// # }
496    /// ```
497    #[must_use]
498    pub fn complete(self, results: serde_json::Value) -> PipelineComplete {
499        PipelineComplete { results }
500    }
501
502    /// Abort execution with an error
503    ///
504    /// Transitions from `Executing` to `Complete` state with error details.
505    ///
506    /// # Example
507    ///
508    /// ```
509    /// use stygian_graph::domain::pipeline::PipelineUnvalidated;
510    /// use serde_json::json;
511    ///
512    /// # fn example() -> Result<(), Box<dyn std::error::Error>> {
513    /// let pipeline = PipelineUnvalidated::new(json!({"nodes": []}))
514    ///     .validate()?
515    ///     .execute();
516    ///
517    /// let complete = pipeline.abort("Network timeout");
518    /// # Ok(())
519    /// # }
520    /// ```
521    #[must_use]
522    pub fn abort(self, error: &str) -> PipelineComplete {
523        PipelineComplete {
524            results: serde_json::json!({
525                "status": "error",
526                "error": error
527            }),
528        }
529    }
530}
531
532impl PipelineComplete {
533    /// Check if the pipeline completed successfully
534    ///
535    /// # Example
536    ///
537    /// ```
538    /// use stygian_graph::domain::pipeline::PipelineUnvalidated;
539    /// use serde_json::json;
540    ///
541    /// # fn example() -> Result<(), Box<dyn std::error::Error>> {
542    /// let pipeline = PipelineUnvalidated::new(json!({"nodes": []}))
543    ///     .validate()?
544    ///     .execute()
545    ///     .complete(json!({"status": "success"}));
546    ///
547    /// assert!(pipeline.is_success());
548    /// # Ok(())
549    /// # }
550    /// ```
551    #[must_use]
552    pub fn is_success(&self) -> bool {
553        self.results
554            .get("status")
555            .and_then(|s| s.as_str())
556            .is_some_and(|s| s == "success")
557    }
558
559    /// Get the execution results
560    #[must_use]
561    pub const fn results(&self) -> &serde_json::Value {
562        &self.results
563    }
564}
565
566/// One URL the pipeline was allowed to fetch despite a guard
567/// `Forbid` verdict, surfaced under [`RobotsPolicy::IgnoreWithAudit`].
568#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
569pub struct RobotsAuditEvent {
570    /// Node id this URL was attached to.
571    pub node_id: String,
572    /// The URL that was fetched.
573    pub url: String,
574    /// Reason recorded by the guard.
575    pub reason: String,
576}
577
578/// Collect every `(node_id, url)` pair the pipeline declares.
579///
580/// Walks the `nodes` array; for each node, looks at the `url` field
581/// or `params.url` (the canonical shapes used by the example TOML
582/// configs and the MCP server). Nodes without a URL are skipped —
583/// non-fetching nodes (e.g. `ai_claude`, `storage`) shouldn't trigger
584/// a robots check.
585fn collect_node_urls(config: &serde_json::Value) -> Vec<(String, String)> {
586    let mut out = Vec::new();
587    let Some(nodes) = config.get("nodes").and_then(|n| n.as_array()) else {
588        return out;
589    };
590    for node in nodes {
591        let Some(node_obj) = node.as_object() else {
592            continue;
593        };
594        let Some(id) = node_obj.get("id").and_then(|v| v.as_str()) else {
595            continue;
596        };
597
598        // Shape 1: top-level `url` on the node.
599        if let Some(url) = node_obj.get("url").and_then(|v| v.as_str()) {
600            out.push((id.to_string(), url.to_string()));
601            continue;
602        }
603
604        // Shape 2: `params.url` (the example TOML convention).
605        if let Some(url) = node_obj
606            .get("params")
607            .and_then(|p| p.get("url"))
608            .and_then(|v| v.as_str())
609        {
610            out.push((id.to_string(), url.to_string()));
611        }
612    }
613    out
614}
615
616#[cfg(test)]
617#[allow(clippy::unwrap_used, clippy::indexing_slicing, clippy::panic)]
618mod tests {
619    use super::*;
620    use serde_json::json;
621
622    #[test]
623    fn validate_empty_nodes_array() {
624        let pipe = PipelineUnvalidated::new(json!({"nodes": [], "edges": []}));
625        let result = pipe.validate();
626        assert!(result.is_err());
627        assert!(
628            result
629                .unwrap_err()
630                .to_string()
631                .contains("at least one node")
632        );
633    }
634
635    #[test]
636    fn validate_missing_nodes_field() {
637        let pipe = PipelineUnvalidated::new(json!({"edges": []}));
638        let result = pipe.validate();
639        assert!(result.is_err());
640    }
641
642    #[test]
643    fn validate_missing_node_id() {
644        let pipe = PipelineUnvalidated::new(json!({
645            "nodes": [{"service": "http"}],
646            "edges": []
647        }));
648        let result = pipe.validate();
649        assert!(result.is_err());
650        assert!(
651            result
652                .unwrap_err()
653                .to_string()
654                .contains("'id' field is required")
655        );
656    }
657
658    #[test]
659    fn validate_empty_node_id() {
660        let pipe = PipelineUnvalidated::new(json!({
661            "nodes": [{"id": "", "service": "http"}],
662            "edges": []
663        }));
664        let result = pipe.validate();
665        assert!(result.is_err());
666        assert!(
667            result
668                .unwrap_err()
669                .to_string()
670                .contains("id cannot be empty")
671        );
672    }
673
674    #[test]
675    fn validate_duplicate_node_ids() {
676        let pipe = PipelineUnvalidated::new(json!({
677            "nodes": [
678                {"id": "fetch", "service": "http"},
679                {"id": "fetch", "service": "browser"}
680            ],
681            "edges": []
682        }));
683        let result = pipe.validate();
684        assert!(result.is_err());
685        assert!(
686            result
687                .unwrap_err()
688                .to_string()
689                .contains("Duplicate node id")
690        );
691    }
692
693    #[test]
694    fn validate_invalid_service_type() {
695        let pipe = PipelineUnvalidated::new(json!({
696            "nodes": [{"id": "fetch", "service": "invalid_service"}],
697            "edges": []
698        }));
699        let result = pipe.validate();
700        assert!(result.is_err());
701        assert!(result.unwrap_err().to_string().contains("not recognized"));
702    }
703
704    #[test]
705    fn validate_edge_nonexistent_source() {
706        let pipe = PipelineUnvalidated::new(json!({
707            "nodes": [{"id": "extract", "service": "ai_claude"}],
708            "edges": [{"from": "fetch", "to": "extract"}]
709        }));
710        let result = pipe.validate();
711        assert!(result.is_err());
712        assert!(
713            result
714                .unwrap_err()
715                .to_string()
716                .contains("source node 'fetch' not found")
717        );
718    }
719
720    #[test]
721    fn validate_edge_nonexistent_target() {
722        let pipe = PipelineUnvalidated::new(json!({
723            "nodes": [{"id": "fetch", "service": "http"}],
724            "edges": [{"from": "fetch", "to": "extract"}]
725        }));
726        let result = pipe.validate();
727        assert!(result.is_err());
728        assert!(
729            result
730                .unwrap_err()
731                .to_string()
732                .contains("target node 'extract' not found")
733        );
734    }
735
736    #[test]
737    fn validate_self_loop() {
738        let pipe = PipelineUnvalidated::new(json!({
739            "nodes": [{"id": "node1", "service": "http"}],
740            "edges": [{"from": "node1", "to": "node1"}]
741        }));
742        let result = pipe.validate();
743        assert!(result.is_err());
744        assert!(result.unwrap_err().to_string().contains("Self-loop"));
745    }
746
747    #[test]
748    fn validate_cycle_detection() {
749        let pipe = PipelineUnvalidated::new(json!({
750            "nodes": [
751                {"id": "a", "service": "http"},
752                {"id": "b", "service": "ai_claude"},
753                {"id": "c", "service": "browser"}
754            ],
755            "edges": [
756                {"from": "a", "to": "b"},
757                {"from": "b", "to": "c"},
758                {"from": "c", "to": "a"}
759            ]
760        }));
761        let result = pipe.validate();
762        assert!(result.is_err());
763        assert!(result.unwrap_err().to_string().contains("Cycle"));
764    }
765
766    #[test]
767    fn validate_unreachable_nodes() {
768        let pipe = PipelineUnvalidated::new(json!({
769            "nodes": [
770                {"id": "a", "service": "http"},
771                {"id": "orphan", "service": "browser"}
772            ],
773            "edges": []
774        }));
775        let result = pipe.validate();
776        assert!(result.is_err());
777        assert!(result.unwrap_err().to_string().contains("Unreachable"));
778    }
779
780    #[test]
781    fn validate_valid_single_node() {
782        let pipe = PipelineUnvalidated::new(json!({
783            "nodes": [{"id": "fetch", "service": "http"}],
784            "edges": []
785        }));
786        assert!(pipe.validate().is_ok());
787    }
788
789    #[test]
790    fn validate_valid_linear_pipeline() {
791        let pipe = PipelineUnvalidated::new(json!({
792            "nodes": [
793                {"id": "fetch", "service": "http"},
794                {"id": "extract", "service": "ai_claude"},
795                {"id": "store", "service": "storage"}
796            ],
797            "edges": [
798                {"from": "fetch", "to": "extract"},
799                {"from": "extract", "to": "store"}
800            ]
801        }));
802        assert!(pipe.validate().is_ok());
803    }
804
805    #[test]
806    fn validate_valid_dag_branching() {
807        let pipe = PipelineUnvalidated::new(json!({
808            "nodes": [
809                {"id": "fetch", "service": "http"},
810                {"id": "extract_ai", "service": "ai_claude"},
811                {"id": "extract_browser", "service": "browser"},
812                {"id": "merge", "service": "storage"}
813            ],
814            "edges": [
815                {"from": "fetch", "to": "extract_ai"},
816                {"from": "fetch", "to": "extract_browser"},
817                {"from": "extract_ai", "to": "merge"},
818                {"from": "extract_browser", "to": "merge"}
819            ]
820        }));
821        assert!(pipe.validate().is_ok());
822    }
823
824    // ── T111 robots-policy integration tests ──────────────────────────
825
826    use crate::domain::policy::RobotsDecision;
827
828    /// Programmable guard for the tests — decides URLs from a
829    /// pre-populated allow/forbid table.
830    struct ScriptedGuard {
831        name: &'static str,
832        /// URL → decision (None means Unknown).
833        table: std::collections::HashMap<String, crate::domain::policy::RobotsDecision>,
834    }
835
836    #[async_trait::async_trait]
837    impl crate::ports::robots_policy::RobotsPolicyGuard for ScriptedGuard {
838        fn name(&self) -> &'static str {
839            self.name
840        }
841        async fn decide(&self, url: &str) -> Result<RobotsDecision, StygianError> {
842            Ok(self
843                .table
844                .get(url)
845                .cloned()
846                .unwrap_or(RobotsDecision::Unknown))
847        }
848    }
849
850    fn guard_allowing_all() -> std::sync::Arc<ScriptedGuard> {
851        std::sync::Arc::new(ScriptedGuard {
852            name: "test-allow",
853            table: std::collections::HashMap::new(),
854        })
855    }
856
857    fn guard_forbidding(url: &str) -> std::sync::Arc<ScriptedGuard> {
858        let mut table = std::collections::HashMap::new();
859        table.insert(
860            url.to_string(),
861            RobotsDecision::Forbid {
862                reason: "Disallow: /private".to_string(),
863            },
864        );
865        std::sync::Arc::new(ScriptedGuard {
866            name: "test-forbid",
867            table,
868        })
869    }
870
871    #[test]
872    fn effective_robots_policy_defaults_to_obey() {
873        let pipe = PipelineUnvalidated::new(json!({
874            "nodes": [{"id": "fetch", "service": "http", "url": "https://example.com"}]
875        }));
876        assert_eq!(
877            pipe.effective_robots_policy().unwrap(),
878            crate::domain::policy::RobotsPolicy::Obey
879        );
880    }
881
882    #[test]
883    fn effective_robots_policy_reads_top_level_key() {
884        let pipe = PipelineUnvalidated::new(json!({
885            "robots_policy": "ignore_with_audit",
886            "nodes": [{"id": "fetch", "service": "http", "url": "https://example.com"}]
887        }));
888        assert_eq!(
889            pipe.effective_robots_policy().unwrap(),
890            crate::domain::policy::RobotsPolicy::IgnoreWithAudit
891        );
892    }
893
894    #[test]
895    fn effective_robots_policy_reads_nested_pipeline_key() {
896        let pipe = PipelineUnvalidated::new(json!({
897            "pipeline": {"robots_policy": "ignore_silently"},
898            "nodes": [{"id": "fetch", "service": "http", "url": "https://example.com"}]
899        }));
900        assert_eq!(
901            pipe.effective_robots_policy().unwrap(),
902            crate::domain::policy::RobotsPolicy::IgnoreSilently
903        );
904    }
905
906    #[test]
907    fn effective_robots_policy_rejects_unknown_variant() {
908        let pipe = PipelineUnvalidated::new(json!({
909            "robots_policy": "always_obey",
910            "nodes": [{"id": "fetch", "service": "http", "url": "https://example.com"}]
911        }));
912        let err = pipe.effective_robots_policy().unwrap_err();
913        assert!(format!("{err}").contains("always_obey"));
914    }
915
916    #[test]
917    fn collect_node_urls_picks_both_shapes() {
918        let cfg = json!({
919            "nodes": [
920                {"id": "a", "service": "http", "url": "https://a.example"},
921                {"id": "b", "service": "browser", "params": {"url": "https://b.example"}},
922                {"id": "c", "service": "ai_claude"},          // no URL — skip
923                {"id": "d", "service": "http", "url": 42}     // wrong type — skip
924            ]
925        });
926        let urls = collect_node_urls(&cfg);
927        let pairs: std::collections::HashSet<(String, String)> = urls.into_iter().collect();
928        assert!(pairs.contains(&("a".to_string(), "https://a.example".to_string())));
929        assert!(pairs.contains(&("b".to_string(), "https://b.example".to_string())));
930        assert_eq!(pairs.len(), 2);
931    }
932
933    #[tokio::test]
934    async fn obey_with_forbidden_url_refuses_to_validate() {
935        let pipe = PipelineUnvalidated::new(json!({
936            "robots_policy": "obey",
937            "nodes": [{"id": "fetch", "service": "http", "url": "https://private.example/x"}]
938        }));
939        let guard = guard_forbidding("https://private.example/x");
940        let err = pipe.check_robots_policy(&*guard).await.unwrap_err();
941        let msg = format!("{err}");
942        assert!(msg.contains("refuses node 'fetch'"), "{msg}");
943        assert!(msg.contains("Disallow: /private"), "{msg}");
944    }
945
946    #[tokio::test]
947    async fn obey_with_unknown_decision_refuses() {
948        // Guard returns `Unknown` (empty table) and the policy is
949        // `Obey` — refuse. This is the conservative behaviour: an
950        // unknown answer must not let a forbidden URL slip through
951        // just because the guard has no data.
952        let pipe = PipelineUnvalidated::new(json!({
953            "robots_policy": "obey",
954            "nodes": [{"id": "fetch", "service": "http", "url": "https://public.example/"}]
955        }));
956        let guard = guard_allowing_all();
957        let err = pipe.check_robots_policy(&*guard).await.unwrap_err();
958        let msg = format!("{err}");
959        assert!(msg.contains("Unknown under Obey"), "{msg}");
960    }
961
962    #[tokio::test]
963    async fn obey_with_explicit_allow_passes() {
964        // Guard explicitly returns `Allow` for the URL — Obey
965        // permits. The audit vec is empty (no policy violation to
966        // record).
967        let url = "https://public.example/";
968        let mut table = std::collections::HashMap::new();
969        table.insert(
970            url.to_string(),
971            RobotsDecision::Allow {
972                reason: "no rule matched".to_string(),
973            },
974        );
975        let guard = std::sync::Arc::new(ScriptedGuard {
976            name: "test-allow-explicit",
977            table,
978        });
979        let pipe = PipelineUnvalidated::new(json!({
980            "robots_policy": "obey",
981            "nodes": [{"id": "fetch", "service": "http", "url": url}]
982        }));
983        let audit = pipe.check_robots_policy(&*guard).await.unwrap();
984        assert!(audit.is_empty());
985    }
986
987    #[tokio::test]
988    async fn ignore_with_audit_emits_audit_event_for_forbidden_url() {
989        let url = "https://private.example/x";
990        let pipe = PipelineUnvalidated::new(json!({
991            "robots_policy": "ignore_with_audit",
992            "nodes": [{"id": "fetch", "service": "http", "url": url}]
993        }));
994        let guard = guard_forbidding(url);
995        let audit = pipe.check_robots_policy(&*guard).await.unwrap();
996        assert_eq!(audit.len(), 1);
997        assert_eq!(audit[0].node_id, "fetch");
998        assert_eq!(audit[0].url, url);
999        assert!(audit[0].reason.contains("Disallow"));
1000    }
1001
1002    #[tokio::test]
1003    async fn ignore_silently_passes_forbidden_url_without_audit() {
1004        let url = "https://private.example/x";
1005        let pipe = PipelineUnvalidated::new(json!({
1006            "robots_policy": "ignore_silently",
1007            "nodes": [{"id": "fetch", "service": "http", "url": url}]
1008        }));
1009        let guard = guard_forbidding(url);
1010        let audit = pipe.check_robots_policy(&*guard).await.unwrap();
1011        assert!(audit.is_empty());
1012    }
1013
1014    #[tokio::test]
1015    async fn obey_with_permissive_guard_is_rejected_at_check_time() {
1016        // Default-permissive guard + Obey policy = contradiction
1017        // (pipeline claims to obey but guard has no data). Must fail
1018        // loudly.
1019        let pipe = PipelineUnvalidated::new(json!({
1020            "robots_policy": "obey",
1021            "nodes": [{"id": "fetch", "service": "http", "url": "https://example.com"}]
1022        }));
1023        let guard = crate::ports::robots_policy::permissive_guard();
1024        let err = pipe.check_robots_policy(&*guard).await.unwrap_err();
1025        let msg = format!("{err}");
1026        assert!(msg.contains("Obey"), "{msg}");
1027        assert!(msg.contains("permissive"), "{msg}");
1028    }
1029}