Skip to main content

eidos_kernel/
workflow.rs

1//! Deterministic task decomposition primitives.
2//!
3//! This module defines an internal, validated acyclic artifact for turning focused context into
4//! bounded inspection and verification hints. It does not execute anything and performs no IO.
5
6mod context;
7mod draft;
8mod run;
9mod schema;
10mod validate;
11
12pub use context::{
13    TaskContextError, TaskContextSlice, task_context_slice, task_context_slice_for_run,
14};
15pub use draft::draft_task_dag;
16pub use run::{TaskEventError, apply_task_event, initial_run, ready_tasks};
17pub use schema::{
18    RetryPolicy, TaskContextItem, TaskDag, TaskEvent, TaskNode, TaskOperator, TaskRisk, TaskRun,
19    TaskState,
20};
21pub use validate::{DagValidationError, validate_dag};
22
23#[cfg(test)]
24mod tests {
25    use super::*;
26
27    #[test]
28    fn validate_returns_stable_topological_order() {
29        let dag = TaskDag {
30            id: "task.test".to_string(),
31            title: "test".to_string(),
32            nodes: vec![
33                node("c", &["a", "b"]),
34                node("b", &["a"]),
35                node("a", &[]),
36                node("d", &["a"]),
37            ],
38        };
39
40        assert_eq!(validate_dag(&dag).unwrap(), ["a", "b", "c", "d"]);
41    }
42
43    #[test]
44    fn validate_rejects_duplicate_ids() {
45        let dag = TaskDag {
46            id: "task.dup".to_string(),
47            title: "dup".to_string(),
48            nodes: vec![node("a", &[]), node("a", &[])],
49        };
50
51        assert_eq!(
52            validate_dag(&dag),
53            Err(DagValidationError::DuplicateNodeId("a".to_string()))
54        );
55    }
56
57    #[test]
58    fn validate_rejects_missing_dependencies() {
59        let dag = TaskDag {
60            id: "task.missing".to_string(),
61            title: "missing".to_string(),
62            nodes: vec![node("a", &["missing"])],
63        };
64
65        assert_eq!(
66            validate_dag(&dag),
67            Err(DagValidationError::MissingDependency {
68                node_id: "a".to_string(),
69                dependency: "missing".to_string(),
70            })
71        );
72    }
73
74    #[test]
75    fn validate_rejects_cycles() {
76        let dag = TaskDag {
77            id: "task.cycle".to_string(),
78            title: "cycle".to_string(),
79            nodes: vec![node("a", &["b"]), node("b", &["a"])],
80        };
81
82        assert_eq!(
83            validate_dag(&dag),
84            Err(DagValidationError::Cycle {
85                nodes: vec!["a".to_string(), "b".to_string()],
86            })
87        );
88    }
89
90    #[test]
91    fn draft_task_dag_is_valid_and_receipt_backed() {
92        let dag = draft_task_dag(
93            "Fix retry policy",
94            &[
95                TaskContextItem {
96                    id: "fn.retry".to_string(),
97                    title: "retry".to_string(),
98                    kind: "function".to_string(),
99                    source: "src/retry.rs:10-30".to_string(),
100                },
101                TaskContextItem {
102                    id: "doc.retry-policy".to_string(),
103                    title: "Retry Policy".to_string(),
104                    kind: "doc".to_string(),
105                    source: "docs/retry.md".to_string(),
106                },
107            ],
108        );
109
110        assert_eq!(dag.id, "task.fix-retry-policy");
111        assert_eq!(
112            validate_dag(&dag).unwrap(),
113            [
114                "context",
115                "inspect-doc-retry-policy",
116                "inspect-fn-retry",
117                "decompose",
118                "change",
119                "verify",
120                "review",
121                "approval"
122            ]
123        );
124        let change = dag.nodes.iter().find(|node| node.id == "change").unwrap();
125        assert_eq!(change.operator, TaskOperator::Edit);
126        assert_eq!(change.risk, TaskRisk::High);
127        assert_eq!(change.tools, ["edit"]);
128        assert!(
129            change
130                .required_checks
131                .contains(&"compile_or_typecheck".to_string())
132        );
133        assert_eq!(change.receipts, ["docs/retry.md", "src/retry.rs:10-30"]);
134
135        let inspect = dag
136            .nodes
137            .iter()
138            .find(|node| node.id == "inspect-fn-retry")
139            .unwrap();
140        assert_eq!(inspect.depends_on, ["context"]);
141        assert_eq!(inspect.tools, ["read", "get_node"]);
142        assert_eq!(inspect.risk, TaskRisk::Low);
143    }
144
145    #[test]
146    fn draft_doc_only_task_uses_proposal_operator() {
147        let dag = draft_task_dag(
148            "Write workflow runbook",
149            &[TaskContextItem {
150                id: "doc.runbook".to_string(),
151                title: "Runbook".to_string(),
152                kind: "doc".to_string(),
153                source: "docs/runbook.md".to_string(),
154            }],
155        );
156
157        let change = dag.nodes.iter().find(|node| node.id == "change").unwrap();
158        assert_eq!(change.operator, TaskOperator::ProposeDoc);
159        assert_eq!(change.tools, ["propose", "sigil-diff"]);
160        assert!(
161            change
162                .required_checks
163                .contains(&"sigil_validation".to_string())
164        );
165    }
166
167    #[test]
168    fn draft_task_dag_batches_large_context_packs() {
169        let context = (0..30)
170            .map(|n| TaskContextItem {
171                id: format!("fn.item_{n}"),
172                title: format!("item {n}"),
173                kind: "function".to_string(),
174                source: format!("src/item.rs:{n}"),
175            })
176            .collect::<Vec<_>>();
177
178        let dag = draft_task_dag("Fix large context task", &context);
179
180        assert!(validate_dag(&dag).is_ok());
181        assert_eq!(dag.nodes.len(), 31);
182        let remaining = dag
183            .nodes
184            .iter()
185            .find(|node| node.id == "inspect-remaining-context")
186            .unwrap();
187        assert_eq!(remaining.inputs.len(), 6);
188        assert_eq!(remaining.receipts.len(), 6);
189        assert_eq!(remaining.depends_on, ["context"]);
190    }
191
192    #[test]
193    fn task_context_slice_resolves_task_inputs_to_context_items() {
194        let context = vec![
195            TaskContextItem {
196                id: "fn.retry".to_string(),
197                title: "retry".to_string(),
198                kind: "function".to_string(),
199                source: "src/retry.rs:10-30".to_string(),
200            },
201            TaskContextItem {
202                id: "doc.retry-policy".to_string(),
203                title: "Retry Policy".to_string(),
204                kind: "doc".to_string(),
205                source: "docs/retry.md".to_string(),
206            },
207        ];
208        let dag = draft_task_dag("Fix retry policy", &context);
209
210        let slice = task_context_slice(&dag, "inspect-fn-retry", &context).unwrap();
211
212        assert_eq!(slice.task_id, "inspect-fn-retry");
213        assert_eq!(slice.operator, TaskOperator::ReadContext);
214        assert_eq!(slice.context.len(), 1);
215        assert_eq!(slice.context[0].id, "fn.retry");
216        assert_eq!(slice.artifact_inputs, Vec::<String>::new());
217        assert_eq!(slice.tools, ["read", "get_node"]);
218        assert_eq!(slice.receipts, ["src/retry.rs:10-30"]);
219    }
220
221    #[test]
222    fn task_context_slice_keeps_artifact_inputs_explicit() {
223        let context = vec![TaskContextItem {
224            id: "doc.runbook".to_string(),
225            title: "Runbook".to_string(),
226            kind: "doc".to_string(),
227            source: "docs/runbook.md".to_string(),
228        }];
229        let dag = draft_task_dag("Write workflow runbook", &context);
230
231        let slice = task_context_slice(&dag, "change", &context).unwrap();
232
233        assert_eq!(slice.operator, TaskOperator::ProposeDoc);
234        assert!(slice.context.is_empty());
235        assert_eq!(slice.artifact_inputs, ["work_items"]);
236        assert_eq!(slice.receipts, ["docs/runbook.md"]);
237        assert_eq!(slice.required_checks, ["sigil_validation"]);
238    }
239
240    #[test]
241    fn task_context_slice_infers_dependency_outputs_when_inputs_are_omitted() {
242        let dag = TaskDag {
243            id: "task.template".to_string(),
244            title: "template task".to_string(),
245            nodes: vec![
246                TaskNode {
247                    id: "context".to_string(),
248                    title: "Gather context".to_string(),
249                    operator: TaskOperator::ReadContext,
250                    outputs: vec!["context_pack".to_string()],
251                    required_checks: vec!["receipt_coverage".to_string()],
252                    tools: vec!["focus".to_string()],
253                    risk: TaskRisk::Low,
254                    ..task_node_defaults()
255                },
256                TaskNode {
257                    id: "risk-review".to_string(),
258                    title: "Review risk".to_string(),
259                    operator: TaskOperator::Review,
260                    depends_on: vec!["context".to_string()],
261                    required_checks: vec!["risk_checked".to_string()],
262                    tools: vec!["neighbors".to_string(), "paths".to_string()],
263                    risk: TaskRisk::Medium,
264                    ..task_node_defaults()
265                },
266            ],
267        };
268
269        let slice = task_context_slice(&dag, "risk-review", &[]).unwrap();
270
271        assert_eq!(slice.inputs, Vec::<String>::new());
272        assert_eq!(slice.artifact_inputs, ["context_pack"]);
273        assert!(slice.context.is_empty());
274        assert_eq!(slice.tools, ["neighbors", "paths"]);
275        assert_eq!(slice.required_checks, ["risk_checked"]);
276    }
277
278    #[test]
279    fn explicit_task_inputs_override_dependency_output_inference() {
280        let context = vec![TaskContextItem {
281            id: "doc.runbook".to_string(),
282            title: "Runbook".to_string(),
283            kind: "doc".to_string(),
284            source: "docs/runbook.md".to_string(),
285        }];
286        let dag = TaskDag {
287            id: "task.explicit".to_string(),
288            title: "explicit task".to_string(),
289            nodes: vec![
290                TaskNode {
291                    id: "context".to_string(),
292                    title: "Gather context".to_string(),
293                    operator: TaskOperator::ReadContext,
294                    outputs: vec!["context_pack".to_string()],
295                    required_checks: vec!["receipt_coverage".to_string()],
296                    tools: vec!["focus".to_string()],
297                    risk: TaskRisk::Low,
298                    ..task_node_defaults()
299                },
300                TaskNode {
301                    id: "inspect-doc".to_string(),
302                    title: "Inspect doc".to_string(),
303                    operator: TaskOperator::ReadContext,
304                    depends_on: vec!["context".to_string()],
305                    inputs: vec!["doc.runbook".to_string()],
306                    required_checks: vec!["receipt_read".to_string()],
307                    tools: vec!["read".to_string()],
308                    risk: TaskRisk::Low,
309                    ..task_node_defaults()
310                },
311            ],
312        };
313
314        let slice = task_context_slice(&dag, "inspect-doc", &context).unwrap();
315
316        assert_eq!(slice.context.len(), 1);
317        assert_eq!(slice.context[0].id, "doc.runbook");
318        assert!(slice.artifact_inputs.is_empty());
319    }
320
321    #[test]
322    fn task_context_slice_rejects_unknown_task() {
323        let dag = TaskDag {
324            id: "task.unknown".to_string(),
325            title: "unknown".to_string(),
326            nodes: vec![node("a", &[])],
327        };
328
329        assert_eq!(
330            task_context_slice(&dag, "missing", &[]),
331            Err(TaskContextError::UnknownTask("missing".to_string()))
332        );
333    }
334
335    #[test]
336    fn task_context_slice_for_run_marks_ready_tasks() {
337        let context = vec![TaskContextItem {
338            id: "doc.runbook".to_string(),
339            title: "Runbook".to_string(),
340            kind: "doc".to_string(),
341            source: "docs/runbook.md".to_string(),
342        }];
343        let dag = draft_task_dag("Write workflow runbook", &context);
344        let run = initial_run(&dag).unwrap();
345        let run = apply_task_event(&dag, &run, event("context", TaskState::Running)).unwrap();
346        let run = apply_task_event(&dag, &run, event("context", TaskState::Passed)).unwrap();
347
348        let slice =
349            task_context_slice_for_run(&dag, &run, "inspect-doc-runbook", &context).unwrap();
350
351        assert_eq!(slice.run_state, Some(TaskState::Ready));
352        assert_eq!(slice.ready, Some(true));
353        assert!(slice.blocked_by.is_empty());
354        assert_eq!(slice.context[0].id, "doc.runbook");
355    }
356
357    #[test]
358    fn task_context_slice_for_run_reports_blocked_dependencies() {
359        let context = vec![TaskContextItem {
360            id: "doc.runbook".to_string(),
361            title: "Runbook".to_string(),
362            kind: "doc".to_string(),
363            source: "docs/runbook.md".to_string(),
364        }];
365        let dag = draft_task_dag("Write workflow runbook", &context);
366        let run = initial_run(&dag).unwrap();
367
368        let slice = task_context_slice_for_run(&dag, &run, "change", &context).unwrap();
369
370        assert_eq!(slice.run_state, Some(TaskState::Pending));
371        assert_eq!(slice.ready, Some(false));
372        assert_eq!(slice.blocked_by, ["decompose"]);
373        assert_eq!(slice.artifact_inputs, ["work_items"]);
374    }
375
376    #[test]
377    fn initial_run_marks_dependency_free_tasks_ready() {
378        let dag = TaskDag {
379            id: "task.ready".to_string(),
380            title: "ready".to_string(),
381            nodes: vec![node("a", &[]), node("b", &["a"]), node("c", &[])],
382        };
383
384        let run = initial_run(&dag).unwrap();
385
386        assert_eq!(ready_tasks(&dag, &run).unwrap(), ["a", "c"]);
387        assert_eq!(run.states["a"], TaskState::Ready);
388        assert_eq!(run.states["b"], TaskState::Pending);
389        assert_eq!(run.states["c"], TaskState::Ready);
390    }
391
392    #[test]
393    fn task_events_advance_ready_tasks_deterministically() {
394        let dag = TaskDag {
395            id: "task.run".to_string(),
396            title: "run".to_string(),
397            nodes: vec![node("a", &[]), node("b", &["a"])],
398        };
399        let run = initial_run(&dag).unwrap();
400
401        let run = apply_task_event(&dag, &run, event("a", TaskState::Running)).unwrap();
402        assert!(ready_tasks(&dag, &run).unwrap().is_empty());
403        assert_eq!(run.states["a"], TaskState::Running);
404
405        let run = apply_task_event(&dag, &run, event("a", TaskState::Passed)).unwrap();
406        assert_eq!(ready_tasks(&dag, &run).unwrap(), ["b"]);
407        assert_eq!(run.states["a"], TaskState::Passed);
408        assert_eq!(run.states["b"], TaskState::Ready);
409        assert_eq!(run.events.len(), 2);
410    }
411
412    #[test]
413    fn task_events_reject_running_before_dependencies_pass() {
414        let dag = TaskDag {
415            id: "task.blocked".to_string(),
416            title: "blocked".to_string(),
417            nodes: vec![node("a", &[]), node("b", &["a"])],
418        };
419        let run = initial_run(&dag).unwrap();
420
421        assert_eq!(
422            apply_task_event(&dag, &run, event("b", TaskState::Running)),
423            Err(TaskEventError::DependenciesNotPassed {
424                task_id: "b".to_string(),
425                blocked_by: vec!["a".to_string()],
426            })
427        );
428    }
429
430    #[test]
431    fn task_events_reject_running_from_non_ready_state() {
432        let dag = TaskDag {
433            id: "task.invalid-transition".to_string(),
434            title: "invalid transition".to_string(),
435            nodes: vec![node("a", &[])],
436        };
437        let mut run = initial_run(&dag).unwrap();
438        run.states.insert("a".to_string(), TaskState::Blocked);
439
440        assert_eq!(
441            apply_task_event(&dag, &run, event("a", TaskState::Running)),
442            Err(TaskEventError::InvalidTransition {
443                task_id: "a".to_string(),
444                from: TaskState::Blocked,
445                to: TaskState::Running,
446            })
447        );
448    }
449
450    #[test]
451    fn task_events_reject_terminal_task_changes() {
452        let dag = TaskDag {
453            id: "task.terminal".to_string(),
454            title: "terminal".to_string(),
455            nodes: vec![node("a", &[])],
456        };
457        let run = initial_run(&dag).unwrap();
458        let run = apply_task_event(&dag, &run, event("a", TaskState::Running)).unwrap();
459        let run = apply_task_event(&dag, &run, event("a", TaskState::Passed)).unwrap();
460
461        assert_eq!(
462            apply_task_event(&dag, &run, event("a", TaskState::Running)),
463            Err(TaskEventError::TerminalTask {
464                task_id: "a".to_string(),
465                state: TaskState::Passed,
466            })
467        );
468    }
469
470    fn node(id: &str, depends_on: &[&str]) -> TaskNode {
471        TaskNode {
472            id: id.to_string(),
473            title: id.to_string(),
474            operator: TaskOperator::Check,
475            state: TaskState::Pending,
476            depends_on: depends_on
477                .iter()
478                .map(std::string::ToString::to_string)
479                .collect(),
480            inputs: Vec::new(),
481            outputs: Vec::new(),
482            required_checks: Vec::new(),
483            tools: Vec::new(),
484            retry: RetryPolicy::default(),
485            risk: TaskRisk::Medium,
486            receipts: Vec::new(),
487        }
488    }
489
490    fn task_node_defaults() -> TaskNode {
491        node("placeholder", &[])
492    }
493
494    fn event(task_id: &str, to: TaskState) -> TaskEvent {
495        TaskEvent {
496            task_id: task_id.to_string(),
497            to,
498            evidence: Vec::new(),
499            note: None,
500        }
501    }
502}