Close spec compliance gaps: retry context, thread resolution, fan-in logging

Fix 3 confirmed gaps from spec compliance review (85 items, 91.8% aligned):

- Write internal.retry_count.<node_id> to PipelineContext after retries
  so handlers and conditions can access retry counts (spec 5.1)
- Add graph-level default_thread (step 3) to 5-step thread ID resolution,
  pass graph param to resolve_thread_id (spec 5.4)
- Write prompt.md/response.md in fan_in LLM evaluation path (spec 5.6)

Also includes pre-existing improvements: checkpoint stores node_outcomes
and next_node_id for correct resume, engine timeout enforcement,
auto_status support, fidelity degradation on resume, preamble injection,
is_retryable error classification, stylesheet specificity correction,
full stylesheet parse validation, direction_valid lint rule, pre-hook
returns Skipped not Fail, fan-in score-based sorting and all-fail
detection, manager_loop child autostart and steer cooldown.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-02-22 10:49:21 -04:00
parent 2bdec7e8e6
commit a2f9e87b7c
15 changed files with 2310 additions and 133 deletions

View file

@ -7,6 +7,7 @@ use serde_json::Value;
use crate::context::Context;
use crate::error::{AttractorError, Result};
use crate::outcome::Outcome;
/// Serializable snapshot of execution state for crash recovery and resume.
#[derive(Debug, Clone, Serialize, Deserialize)]
@ -17,6 +18,12 @@ pub struct Checkpoint {
pub node_retries: HashMap<String, u32>,
pub context_values: HashMap<String, Value>,
pub logs: Vec<String>,
/// Persisted node outcomes for goal gate checks after resume.
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub node_outcomes: HashMap<String, Outcome>,
/// The node to resume execution at (the next node after the checkpoint's current_node).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_node_id: Option<String>,
}
impl Checkpoint {
@ -25,14 +32,19 @@ impl Checkpoint {
context: &Context,
current_node: impl Into<String>,
completed_nodes: Vec<String>,
node_retries: HashMap<String, u32>,
node_outcomes: HashMap<String, Outcome>,
next_node_id: Option<String>,
) -> Self {
Self {
timestamp: Utc::now(),
current_node: current_node.into(),
completed_nodes,
node_retries: HashMap::new(),
node_retries,
context_values: context.snapshot(),
logs: context.logs_snapshot(),
node_outcomes,
next_node_id,
}
}
@ -75,6 +87,9 @@ mod tests {
&ctx,
"node_a",
vec!["start".to_string(), "node_a".to_string()],
HashMap::new(),
HashMap::new(),
None,
);
assert_eq!(cp.current_node, "node_a");
@ -88,6 +103,8 @@ mod tests {
assert_eq!(cp.logs.len(), 1);
assert_eq!(cp.logs[0], "started");
assert!(cp.node_retries.is_empty());
assert!(cp.node_outcomes.is_empty());
assert!(cp.next_node_id.is_none());
}
#[test]
@ -99,8 +116,18 @@ mod tests {
ctx.set("goal", serde_json::json!("test"));
ctx.append_log("log entry");
let mut cp = Checkpoint::from_context(&ctx, "work", vec!["start".to_string()]);
cp.node_retries.insert("work".to_string(), 2);
let mut retries = HashMap::new();
retries.insert("work".to_string(), 2u32);
let mut outcomes = HashMap::new();
outcomes.insert("start".to_string(), Outcome::success());
let cp = Checkpoint::from_context(
&ctx,
"work",
vec!["start".to_string()],
retries,
outcomes,
Some("next_step".to_string()),
);
cp.save(&path).unwrap();
let loaded = Checkpoint::load(&path).unwrap();
@ -113,6 +140,8 @@ mod tests {
Some(&serde_json::json!("test"))
);
assert_eq!(loaded.logs, vec!["log entry"]);
assert_eq!(loaded.node_outcomes.get("start").map(|o| &o.status), Some(&crate::outcome::StageStatus::Success));
assert_eq!(loaded.next_node_id.as_deref(), Some("next_step"));
}
#[test]
@ -134,7 +163,7 @@ mod tests {
#[test]
fn serialization_roundtrip() {
let ctx = Context::new();
let cp = Checkpoint::from_context(&ctx, "n1", vec![]);
let cp = Checkpoint::from_context(&ctx, "n1", vec![], HashMap::new(), HashMap::new(), None);
let json = serde_json::to_string(&cp).unwrap();
let deserialized: Checkpoint = serde_json::from_str(&json).unwrap();

View file

@ -97,9 +97,9 @@ impl std::fmt::Debug for RetryPolicy {
}
}
/// Default should_retry predicate: retries all errors.
/// Default should_retry predicate: retries transient errors only.
fn default_should_retry() -> ShouldRetryFn {
std::sync::Arc::new(|_| true)
std::sync::Arc::new(|err| err.is_retryable())
}
impl RetryPolicy {
@ -211,6 +211,43 @@ pub fn resolve_fidelity(incoming_edge: Option<&Edge>, node: &Node, graph: &Graph
"compact".to_string()
}
// --- Thread ID resolution (spec 5.4) ---
/// Resolve the thread ID for a node, following the precedence (spec lines 1196-1204):
/// 1. Target node `thread_id` attribute
/// 2. Incoming edge `thread_id` attribute
/// 3. Graph-level default thread
/// 4. Derived class from enclosing subgraph (first class from the node's classes list)
/// 5. Fallback to previous node ID
#[must_use]
pub fn resolve_thread_id(
incoming_edge: Option<&Edge>,
node: &Node,
graph: &Graph,
previous_node_id: Option<&str>,
) -> Option<String> {
// Step 1: Node thread_id
if let Some(tid) = node.thread_id() {
return Some(tid.to_string());
}
// Step 2: Edge thread_id
if let Some(edge) = incoming_edge {
if let Some(tid) = edge.thread_id() {
return Some(tid.to_string());
}
}
// Step 3: Graph-level default thread
if let Some(tid) = graph.default_thread() {
return Some(tid.to_string());
}
// Step 4: Derived class from enclosing subgraph
if let Some(first_class) = node.classes.first() {
return Some(first_class.clone());
}
// Step 5: Fallback to previous node ID
previous_node_id.map(String::from)
}
// --- Run directory helpers (spec 5.6) ---
/// Write manifest.json at the start of a pipeline run.
@ -222,6 +259,7 @@ fn write_manifest(logs_root: &Path, graph: &Graph) {
};
let manifest = serde_json::json!({
"pipeline_name": pipeline_name,
"goal": graph.goal(),
"start_time": Utc::now().to_rfc3339(),
"node_count": graph.nodes.len(),
"edge_count": graph.edges.len(),
@ -446,6 +484,7 @@ impl PipelineEngine {
}
/// Execute a node handler with retry policy.
/// Returns `(outcome, attempts_used)` where `attempts_used` is the 1-indexed count.
async fn execute_with_retry(
&self,
node: &Node,
@ -454,14 +493,31 @@ impl PipelineEngine {
logs_root: &Path,
policy: &RetryPolicy,
stage_index: usize,
) -> Result<Outcome> {
) -> Result<(Outcome, u32)> {
let handler = self.registry.resolve(node);
let node_timeout = node.timeout();
for attempt in 1..=policy.max_attempts {
// Gap #11: Panic safety -- catch panics from handler execution
let result = {
let future = handler.execute(node, context, graph, logs_root);
match AssertUnwindSafe(future).catch_unwind().await {
let panic_safe = AssertUnwindSafe(future).catch_unwind();
// Gap #2: Timeout enforcement -- wrap with tokio::time::timeout
let timed_result = if let Some(duration) = node_timeout {
match tokio::time::timeout(duration, panic_safe).await {
Ok(inner) => inner,
Err(_elapsed) => {
Ok(Ok(Outcome::fail(format!(
"handler timed out after {}ms",
duration.as_millis()
))))
}
}
} else {
panic_safe.await
};
match timed_result {
Ok(r) => r,
Err(panic_payload) => {
let msg = if let Some(s) = panic_payload.downcast_ref::<&str>() {
@ -497,7 +553,7 @@ impl PipelineEngine {
tokio::time::sleep(delay).await;
continue;
}
return Ok(Outcome::fail(e.to_string()));
return Ok((Outcome::fail(e.to_string()), attempt));
}
};
@ -506,7 +562,7 @@ impl PipelineEngine {
| StageStatus::PartialSuccess
| StageStatus::Fail
| StageStatus::Skipped => {
return Ok(outcome);
return Ok((outcome, attempt));
}
StageStatus::Retry => {
if attempt < policy.max_attempts {
@ -521,18 +577,23 @@ impl PipelineEngine {
continue;
}
if node.allow_partial() {
return Ok(Outcome {
status: StageStatus::PartialSuccess,
notes: Some("retries exhausted, partial accepted".to_string()),
..Outcome::success()
});
return Ok((
Outcome {
status: StageStatus::PartialSuccess,
notes: Some(
"retries exhausted, partial accepted".to_string(),
),
..Outcome::success()
},
attempt,
));
}
return Ok(Outcome::fail("max retries exceeded"));
return Ok((Outcome::fail("max retries exceeded"), attempt));
}
}
}
Ok(Outcome::fail("max retries exceeded"))
Ok((Outcome::fail("max retries exceeded"), policy.max_attempts))
}
/// Run the pipeline. Returns the final outcome.
@ -584,9 +645,13 @@ impl PipelineEngine {
let context;
let mut completed_nodes: Vec<String>;
let mut node_outcomes: HashMap<String, Outcome> = HashMap::new();
let mut node_retries: HashMap<String, u32> = HashMap::new();
let mut stage_index: usize;
let mut current_node_id: String;
let mut incoming_edge: Option<&Edge> = None;
let mut previous_node_id: Option<String> = None;
// Gap #6: Track whether fidelity should be degraded on the first resumed node
let mut degrade_fidelity_on_resume = false;
if let Some(cp) = resume_checkpoint {
// Restore context from checkpoint
@ -598,13 +663,27 @@ impl PipelineEngine {
context.append_log(log_entry.clone());
}
completed_nodes = cp.completed_nodes.clone();
// Gap #5: Restore retry counters from checkpoint
node_retries = cp.node_retries.clone();
// P1: Restore node outcomes for goal gate checks
node_outcomes = cp.node_outcomes.clone();
stage_index = completed_nodes.len();
// Resume from the node after the checkpoint's current_node
let edges = graph.outgoing_edges(&cp.current_node);
if let Some(edge) = edges.first() {
current_node_id = edge.to.clone();
// P1: Use stored next_node_id if available, otherwise fall back
if let Some(ref next_id) = cp.next_node_id {
current_node_id = next_id.clone();
} else {
current_node_id = cp.current_node.clone();
let edges = graph.outgoing_edges(&cp.current_node);
if let Some(edge) = edges.first() {
current_node_id = edge.to.clone();
} else {
current_node_id = cp.current_node.clone();
}
}
// Gap #6: Check if the checkpointed node used full fidelity
if cp.context_values.get("internal.fidelity")
== Some(&serde_json::json!("full"))
{
degrade_fidelity_on_resume = true;
}
} else if let Some(start) = start_at {
context = Context::new();
@ -653,11 +732,32 @@ impl PipelineEngine {
}
// Resolve fidelity (spec 5.4) and store in context
let fidelity = resolve_fidelity(incoming_edge, node, graph);
let mut fidelity = resolve_fidelity(incoming_edge, node, graph);
// Gap #6: On the first node after resume, degrade full -> summary:high
if degrade_fidelity_on_resume && fidelity == "full" {
fidelity = "summary:high".to_string();
}
degrade_fidelity_on_resume = false;
context.set("internal.fidelity", serde_json::json!(&fidelity));
// Thread context sharing: store thread association
if let Some(tid) = node.thread_id() {
// Preamble injection at execution time (spec 8.3): if fidelity is not "full",
// store a preamble in context for handlers to read
if fidelity != "full" {
context.set(
"current.preamble",
serde_json::json!(format!("[Context mode: {fidelity}]\n\n")),
);
} else {
context.set("current.preamble", serde_json::json!(""));
}
// Thread context sharing: resolve thread ID and store association
if let Some(tid) = resolve_thread_id(
incoming_edge,
node,
graph,
previous_node_id.as_deref(),
) {
context.set(
format!("thread.{tid}.current_node"),
serde_json::json!(&node.id),
@ -674,9 +774,25 @@ impl PipelineEngine {
});
let stage_start = Instant::now();
let outcome = self
let (mut outcome, attempts_used) = self
.execute_with_retry(node, &context, graph, &config.logs_root, &retry_policy, stage_index)
.await?;
// Gap #5: Track retry count per node
node_retries.insert(node.id.clone(), attempts_used);
context.set(
format!("internal.retry_count.{}", node.id),
serde_json::json!(attempts_used),
);
// Gap #1: Auto status -- when auto_status=true and outcome is non-success,
// override to success with auto-status note
if node.auto_status() && outcome.status != StageStatus::Success {
outcome = Outcome {
status: StageStatus::Success,
notes: Some("auto-status: handler completed without writing status".to_string()),
..outcome
};
}
let stage_duration_ms = millis_u64(stage_start.elapsed());
@ -705,6 +821,7 @@ impl PipelineEngine {
// Step 3: Record completion
completed_nodes.push(node.id.clone());
node_outcomes.insert(node.id.clone(), outcome.clone());
previous_node_id = Some(node.id.clone());
stage_index += 1;
// Step 4: Apply context updates from outcome
@ -714,11 +831,18 @@ impl PipelineEngine {
context.set("preferred_label", serde_json::json!(pref));
}
// Step 5: Save checkpoint
// Step 5: Select next edge (done before checkpoint so we can store next_node_id)
let next_edge = select_edge(&node.id, &outcome, &context, graph);
let next_node_id_for_checkpoint = next_edge.map(|e| e.to.clone());
// Step 6: Save checkpoint with all state
let checkpoint = Checkpoint::from_context(
&context,
&node.id,
completed_nodes.clone(),
node_retries.clone(),
node_outcomes.clone(),
next_node_id_for_checkpoint,
);
let checkpoint_path = config.logs_root.join("checkpoint.json");
if let Err(e) = checkpoint.save(&checkpoint_path) {
@ -729,8 +853,7 @@ impl PipelineEngine {
});
}
// Step 6: Select next edge
let next_edge = select_edge(&node.id, &outcome, &context, graph);
// Step 7: Follow selected edge
match next_edge {
None => {
// Gap #1: Failure routing -- when FAIL and no matching edge,
@ -790,6 +913,46 @@ mod tests {
use super::*;
use crate::graph::AttrValue;
use crate::handler::start::StartHandler;
use crate::handler::Handler as HandlerTrait;
use async_trait::async_trait;
use std::time::Duration;
// --- Test-only handlers ---
/// Handler that always returns Fail.
struct AlwaysFailHandler;
#[async_trait]
impl HandlerTrait for AlwaysFailHandler {
async fn execute(
&self,
_node: &Node,
_context: &Context,
_graph: &Graph,
_logs_root: &Path,
) -> std::result::Result<Outcome, AttractorError> {
Ok(Outcome::fail("always fails"))
}
}
/// Handler that sleeps for a configurable duration, then succeeds.
struct SlowHandler {
sleep_ms: u64,
}
#[async_trait]
impl HandlerTrait for SlowHandler {
async fn execute(
&self,
_node: &Node,
_context: &Context,
_graph: &Graph,
_logs_root: &Path,
) -> std::result::Result<Outcome, AttractorError> {
tokio::time::sleep(Duration::from_millis(self.sleep_ms)).await;
Ok(Outcome::success())
}
}
// --- BackoffConfig tests ---
@ -1527,6 +1690,7 @@ mod tests {
let manifest: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(&manifest_path).unwrap()).unwrap();
assert_eq!(manifest["pipeline_name"], "test_pipeline");
assert_eq!(manifest["goal"], "Run tests");
assert!(manifest["start_time"].is_string());
assert!(manifest["node_count"].is_number());
assert!(manifest["edge_count"].is_number());
@ -1567,4 +1731,364 @@ mod tests {
Some(&serde_json::json!("compact"))
);
}
// --- resolve_thread_id tests ---
#[test]
fn thread_id_from_node_attribute() {
let mut node = Node::new("work");
node.attrs.insert(
"thread_id".to_string(),
AttrValue::String("main-thread".to_string()),
);
let graph = Graph::new("test");
assert_eq!(
resolve_thread_id(None, &node, &graph, Some("prev")),
Some("main-thread".to_string())
);
}
#[test]
fn thread_id_from_edge_attribute() {
let node = Node::new("work");
let mut edge = Edge::new("prev", "work");
edge.attrs.insert(
"thread_id".to_string(),
AttrValue::String("edge-thread".to_string()),
);
let graph = Graph::new("test");
assert_eq!(
resolve_thread_id(Some(&edge), &node, &graph, Some("prev")),
Some("edge-thread".to_string())
);
}
#[test]
fn thread_id_node_overrides_edge() {
let mut node = Node::new("work");
node.attrs.insert(
"thread_id".to_string(),
AttrValue::String("node-thread".to_string()),
);
let mut edge = Edge::new("prev", "work");
edge.attrs.insert(
"thread_id".to_string(),
AttrValue::String("edge-thread".to_string()),
);
let graph = Graph::new("test");
assert_eq!(
resolve_thread_id(Some(&edge), &node, &graph, Some("prev")),
Some("node-thread".to_string())
);
}
#[test]
fn thread_id_from_graph_default_thread() {
let node = Node::new("work");
let mut graph = Graph::new("test");
graph.attrs.insert(
"default_thread".to_string(),
AttrValue::String("shared-thread".to_string()),
);
assert_eq!(
resolve_thread_id(None, &node, &graph, Some("prev")),
Some("shared-thread".to_string())
);
}
#[test]
fn thread_id_edge_overrides_graph_default() {
let node = Node::new("work");
let mut edge = Edge::new("prev", "work");
edge.attrs.insert(
"thread_id".to_string(),
AttrValue::String("edge-thread".to_string()),
);
let mut graph = Graph::new("test");
graph.attrs.insert(
"default_thread".to_string(),
AttrValue::String("shared-thread".to_string()),
);
assert_eq!(
resolve_thread_id(Some(&edge), &node, &graph, Some("prev")),
Some("edge-thread".to_string())
);
}
#[test]
fn thread_id_graph_default_overrides_class() {
let mut node = Node::new("work");
node.classes = vec!["planning".to_string()];
let mut graph = Graph::new("test");
graph.attrs.insert(
"default_thread".to_string(),
AttrValue::String("shared-thread".to_string()),
);
assert_eq!(
resolve_thread_id(None, &node, &graph, Some("prev")),
Some("shared-thread".to_string())
);
}
#[test]
fn thread_id_from_node_class() {
let mut node = Node::new("work");
node.classes = vec!["planning".to_string(), "review".to_string()];
let graph = Graph::new("test");
assert_eq!(
resolve_thread_id(None, &node, &graph, Some("prev")),
Some("planning".to_string())
);
}
#[test]
fn thread_id_fallback_to_previous_node() {
let node = Node::new("work");
let graph = Graph::new("test");
assert_eq!(
resolve_thread_id(None, &node, &graph, Some("prev_node")),
Some("prev_node".to_string())
);
}
#[test]
fn thread_id_none_when_no_sources() {
let node = Node::new("start");
let graph = Graph::new("test");
assert_eq!(resolve_thread_id(None, &node, &graph, None), None);
}
// --- default_should_retry tests ---
#[test]
fn default_should_retry_retries_transient_errors() {
let should_retry = default_should_retry();
assert!(should_retry(&AttractorError::Handler("timeout".to_string())));
assert!(should_retry(&AttractorError::Engine("transient".to_string())));
assert!(should_retry(&AttractorError::Io("connection reset".to_string())));
}
#[test]
fn default_should_retry_rejects_terminal_errors() {
let should_retry = default_should_retry();
assert!(!should_retry(&AttractorError::Parse("bad syntax".to_string())));
assert!(!should_retry(&AttractorError::Validation("invalid".to_string())));
assert!(!should_retry(&AttractorError::Stylesheet("bad rule".to_string())));
assert!(!should_retry(&AttractorError::Checkpoint("corrupt".to_string())));
}
// --- Gap #15: Manifest goal field test ---
#[tokio::test]
async fn engine_manifest_includes_goal() {
let dir = tempfile::tempdir().unwrap();
let g = simple_graph();
let engine = PipelineEngine::new(make_registry(), EventEmitter::new());
let config = RunConfig { logs_root: dir.path().to_path_buf() };
engine.run(&g, &config).await.unwrap();
let manifest_path = dir.path().join("manifest.json");
let manifest: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(&manifest_path).unwrap()).unwrap();
assert_eq!(manifest["goal"], "Run tests");
}
#[tokio::test]
async fn engine_manifest_goal_empty_when_unset() {
let dir = tempfile::tempdir().unwrap();
let mut g = Graph::new("no_goal");
let mut start = Node::new("start");
start.attrs.insert("shape".to_string(), AttrValue::String("Mdiamond".to_string()));
g.nodes.insert("start".to_string(), start);
let mut exit = Node::new("exit");
exit.attrs.insert("shape".to_string(), AttrValue::String("Msquare".to_string()));
g.nodes.insert("exit".to_string(), exit);
g.edges.push(Edge::new("start", "exit"));
let engine = PipelineEngine::new(make_registry(), EventEmitter::new());
let config = RunConfig { logs_root: dir.path().to_path_buf() };
engine.run(&g, &config).await.unwrap();
let manifest_path = dir.path().join("manifest.json");
let manifest: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(&manifest_path).unwrap()).unwrap();
assert_eq!(manifest["goal"], "");
}
// --- Gap #1: Auto status tests ---
#[tokio::test]
async fn engine_auto_status_overrides_fail_to_success() {
let dir = tempfile::tempdir().unwrap();
let mut g = Graph::new("auto_status_test");
let mut start = Node::new("start");
start.attrs.insert("shape".to_string(), AttrValue::String("Mdiamond".to_string()));
g.nodes.insert("start".to_string(), start);
let mut work = Node::new("work");
work.attrs.insert("auto_status".to_string(), AttrValue::Boolean(true));
work.attrs.insert("type".to_string(), AttrValue::String("always_fail".to_string()));
work.attrs.insert("max_retries".to_string(), AttrValue::Integer(0));
g.nodes.insert("work".to_string(), work);
let mut exit = Node::new("exit");
exit.attrs.insert("shape".to_string(), AttrValue::String("Msquare".to_string()));
g.nodes.insert("exit".to_string(), exit);
g.edges.push(Edge::new("start", "work"));
g.edges.push(Edge::new("work", "exit"));
let mut registry = make_registry();
registry.register("always_fail", Box::new(AlwaysFailHandler));
let engine = PipelineEngine::new(registry, EventEmitter::new());
let config = RunConfig { logs_root: dir.path().to_path_buf() };
let outcome = engine.run(&g, &config).await.unwrap();
assert_eq!(outcome.status, StageStatus::Success);
assert_eq!(
outcome.notes.as_deref(),
Some("auto-status: handler completed without writing status")
);
}
#[tokio::test]
async fn engine_auto_status_false_preserves_fail() {
let dir = tempfile::tempdir().unwrap();
let mut g = Graph::new("no_auto_status_test");
let mut start = Node::new("start");
start.attrs.insert("shape".to_string(), AttrValue::String("Mdiamond".to_string()));
g.nodes.insert("start".to_string(), start);
let mut work = Node::new("work");
work.attrs.insert("type".to_string(), AttrValue::String("always_fail".to_string()));
work.attrs.insert("max_retries".to_string(), AttrValue::Integer(0));
g.nodes.insert("work".to_string(), work);
let mut exit = Node::new("exit");
exit.attrs.insert("shape".to_string(), AttrValue::String("Msquare".to_string()));
g.nodes.insert("exit".to_string(), exit);
g.edges.push(Edge::new("start", "work"));
let mut fail_edge = Edge::new("work", "exit");
fail_edge.attrs.insert("condition".to_string(), AttrValue::String("outcome=fail".to_string()));
g.edges.push(fail_edge);
let mut registry = make_registry();
registry.register("always_fail", Box::new(AlwaysFailHandler));
let engine = PipelineEngine::new(registry, EventEmitter::new());
let config = RunConfig { logs_root: dir.path().to_path_buf() };
let result = engine.run(&g, &config).await;
assert!(result.is_ok());
let status_path = dir.path().join("work").join("status.json");
let status: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(&status_path).unwrap()).unwrap();
assert_eq!(status["status"], "fail");
}
// --- Gap #2: Timeout enforcement tests ---
#[tokio::test]
async fn engine_timeout_causes_fail_outcome() {
let dir = tempfile::tempdir().unwrap();
let mut g = Graph::new("timeout_test");
let mut start = Node::new("start");
start.attrs.insert("shape".to_string(), AttrValue::String("Mdiamond".to_string()));
g.nodes.insert("start".to_string(), start);
let mut work = Node::new("work");
work.attrs.insert("timeout".to_string(), AttrValue::Duration(Duration::from_millis(50)));
work.attrs.insert("type".to_string(), AttrValue::String("slow".to_string()));
work.attrs.insert("max_retries".to_string(), AttrValue::Integer(0));
g.nodes.insert("work".to_string(), work);
let mut exit = Node::new("exit");
exit.attrs.insert("shape".to_string(), AttrValue::String("Msquare".to_string()));
g.nodes.insert("exit".to_string(), exit);
g.edges.push(Edge::new("start", "work"));
let mut fail_edge = Edge::new("work", "exit");
fail_edge.attrs.insert("condition".to_string(), AttrValue::String("outcome=fail".to_string()));
g.edges.push(fail_edge);
let mut registry = make_registry();
registry.register("slow", Box::new(SlowHandler { sleep_ms: 500 }));
let engine = PipelineEngine::new(registry, EventEmitter::new());
let config = RunConfig { logs_root: dir.path().to_path_buf() };
let result = engine.run(&g, &config).await;
assert!(result.is_ok());
let status_path = dir.path().join("work").join("status.json");
let status: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(&status_path).unwrap()).unwrap();
assert_eq!(status["status"], "fail");
}
#[tokio::test]
async fn engine_no_timeout_completes_normally() {
let dir = tempfile::tempdir().unwrap();
let mut g = Graph::new("no_timeout_test");
let mut start = Node::new("start");
start.attrs.insert("shape".to_string(), AttrValue::String("Mdiamond".to_string()));
g.nodes.insert("start".to_string(), start);
let mut work = Node::new("work");
work.attrs.insert("type".to_string(), AttrValue::String("slow".to_string()));
work.attrs.insert("max_retries".to_string(), AttrValue::Integer(0));
g.nodes.insert("work".to_string(), work);
let mut exit = Node::new("exit");
exit.attrs.insert("shape".to_string(), AttrValue::String("Msquare".to_string()));
g.nodes.insert("exit".to_string(), exit);
g.edges.push(Edge::new("start", "work"));
g.edges.push(Edge::new("work", "exit"));
let mut registry = make_registry();
registry.register("slow", Box::new(SlowHandler { sleep_ms: 10 }));
let engine = PipelineEngine::new(registry, EventEmitter::new());
let config = RunConfig { logs_root: dir.path().to_path_buf() };
let outcome = engine.run(&g, &config).await.unwrap();
assert_eq!(outcome.status, StageStatus::Success);
}
#[tokio::test]
async fn engine_timeout_with_auto_status_returns_success() {
let dir = tempfile::tempdir().unwrap();
let mut g = Graph::new("timeout_auto_status_test");
let mut start = Node::new("start");
start.attrs.insert("shape".to_string(), AttrValue::String("Mdiamond".to_string()));
g.nodes.insert("start".to_string(), start);
let mut work = Node::new("work");
work.attrs.insert("timeout".to_string(), AttrValue::Duration(Duration::from_millis(50)));
work.attrs.insert("auto_status".to_string(), AttrValue::Boolean(true));
work.attrs.insert("type".to_string(), AttrValue::String("slow".to_string()));
work.attrs.insert("max_retries".to_string(), AttrValue::Integer(0));
g.nodes.insert("work".to_string(), work);
let mut exit = Node::new("exit");
exit.attrs.insert("shape".to_string(), AttrValue::String("Msquare".to_string()));
g.nodes.insert("exit".to_string(), exit);
g.edges.push(Edge::new("start", "work"));
g.edges.push(Edge::new("work", "exit"));
let mut registry = make_registry();
registry.register("slow", Box::new(SlowHandler { sleep_ms: 500 }));
let engine = PipelineEngine::new(registry, EventEmitter::new());
let config = RunConfig { logs_root: dir.path().to_path_buf() };
let outcome = engine.run(&g, &config).await.unwrap();
assert_eq!(outcome.status, StageStatus::Success);
assert_eq!(
outcome.notes.as_deref(),
Some("auto-status: handler completed without writing status")
);
}
}

View file

@ -24,6 +24,24 @@ pub enum AttractorError {
Io(String),
}
impl AttractorError {
/// Whether this error category is retryable (transient) or terminal.
///
/// Retryable: Handler (transient handler failures), Engine (could be transient),
/// Io (network/disk issues are often transient).
/// Terminal: Parse, Validation, Stylesheet (configuration errors),
/// Checkpoint (storage integrity).
#[must_use]
pub const fn is_retryable(&self) -> bool {
match self {
Self::Handler(_) | Self::Engine(_) | Self::Io(_) => true,
Self::Parse(_) | Self::Validation(_) | Self::Stylesheet(_) | Self::Checkpoint(_) => {
false
}
}
}
}
impl From<std::io::Error> for AttractorError {
fn from(err: std::io::Error) -> Self {
Self::Io(err.to_string())
@ -88,4 +106,19 @@ mod tests {
let err: Result<i32> = Err(AttractorError::Parse("bad".to_string()));
assert!(err.is_err());
}
#[test]
fn is_retryable_terminal_errors() {
assert!(!AttractorError::Parse("bad".to_string()).is_retryable());
assert!(!AttractorError::Validation("bad".to_string()).is_retryable());
assert!(!AttractorError::Stylesheet("bad".to_string()).is_retryable());
assert!(!AttractorError::Checkpoint("bad".to_string()).is_retryable());
}
#[test]
fn is_retryable_transient_errors() {
assert!(AttractorError::Handler("timeout".to_string()).is_retryable());
assert!(AttractorError::Engine("transient".to_string()).is_retryable());
assert!(AttractorError::Io("connection reset".to_string()).is_retryable());
}
}

View file

@ -371,6 +371,13 @@ impl Graph {
.get("default_fidelity")
.and_then(AttrValue::as_str)
}
/// Graph-level `default_thread`.
pub fn default_thread(&self) -> Option<&str> {
self.attrs
.get("default_thread")
.and_then(AttrValue::as_str)
}
}
#[cfg(test)]

View file

@ -98,7 +98,9 @@ impl Handler for CodergenHandler {
// 3. Execute pre-hook (spec 9.7)
if let Some(pre_hook) = resolve_hook(node, graph, "tool_hooks.pre") {
if !run_hook(&pre_hook, &node.id) {
return Ok(Outcome::fail("pre-hook failed, skipping LLM call"));
let mut outcome = Outcome::skipped();
outcome.notes = Some("pre-hook returned non-zero, tool call skipped".to_string());
return Ok(outcome);
}
}
@ -299,9 +301,9 @@ mod tests {
.execute(&node, &context, &graph, tmp.path())
.await
.unwrap();
assert_eq!(outcome.status, crate::outcome::StageStatus::Fail);
assert_eq!(outcome.status, crate::outcome::StageStatus::Skipped);
assert!(outcome
.failure_reason
.notes
.as_deref()
.unwrap()
.contains("pre-hook"));

View file

@ -29,7 +29,7 @@ impl Handler for FanInHandler {
node: &Node,
context: &Context,
_graph: &Graph,
_logs_root: &Path,
logs_root: &Path,
) -> Result<Outcome, AttractorError> {
let results = context.get("parallel.results");
let Some(results) = results else {
@ -39,11 +39,29 @@ impl Handler for FanInHandler {
let prompt = node.prompt().filter(|p| !p.is_empty());
let best = if let (Some(prompt_text), Some(backend)) = (prompt, &self.backend) {
llm_evaluate(backend.as_ref(), prompt_text, &results, context).await?
llm_evaluate(backend.as_ref(), prompt_text, &results, context, logs_root, &node.id).await?
} else {
heuristic_select(&results)
};
// Check if all candidates failed — if so, return fail
let all_failed = if best.status == "fail" {
let empty_vec = vec![];
let arr = results.as_array().unwrap_or(&empty_vec);
arr.iter().all(|v| {
v.get("status")
.and_then(|v| v.as_str())
.unwrap_or("fail")
== "fail"
})
} else {
false
};
if all_failed {
return Ok(Outcome::fail("all candidates failed"));
}
let mut outcome = Outcome::success();
outcome.context_updates.insert(
"parallel.fan_in.best_id".to_string(),
@ -62,6 +80,7 @@ impl Handler for FanInHandler {
struct Candidate {
id: String,
status: String,
score: f64,
}
fn status_rank(status: &str) -> u32 {
@ -81,6 +100,7 @@ fn heuristic_select(results: &serde_json::Value) -> Candidate {
return Candidate {
id: "unknown".to_string(),
status: "fail".to_string(),
score: 0.0,
};
}
@ -97,6 +117,10 @@ fn heuristic_select(results: &serde_json::Value) -> Candidate {
.and_then(|v| v.as_str())
.unwrap_or("fail")
.to_string(),
score: v
.get("score")
.and_then(|v| v.as_f64())
.unwrap_or(0.0),
})
.collect();
@ -105,12 +129,18 @@ fn heuristic_select(results: &serde_json::Value) -> Candidate {
if rank_cmp != std::cmp::Ordering::Equal {
return rank_cmp;
}
// Higher score is better, so reverse the comparison
let score_cmp = b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal);
if score_cmp != std::cmp::Ordering::Equal {
return score_cmp;
}
a.id.cmp(&b.id)
});
candidates.into_iter().next().unwrap_or_else(|| Candidate {
id: "unknown".to_string(),
status: "fail".to_string(),
score: 0.0,
})
}
@ -120,6 +150,8 @@ async fn llm_evaluate(
prompt: &str,
results: &serde_json::Value,
context: &Context,
logs_root: &Path,
node_id: &str,
) -> Result<Candidate, AttractorError> {
let results_text = serde_json::to_string_pretty(results)
.unwrap_or_else(|_| results.to_string());
@ -129,6 +161,11 @@ async fn llm_evaluate(
Respond with the ID of the best candidate."
);
// Write prompt to logs
let stage_dir = logs_root.join(node_id);
tokio::fs::create_dir_all(&stage_dir).await?;
tokio::fs::write(stage_dir.join("prompt.md"), &full_prompt).await?;
// Build a synthetic node for the backend call
let eval_node = Node::new("fan_in_eval");
@ -142,12 +179,19 @@ async fn llm_evaluate(
.map(String::from)
.or_else(|| outcome.notes.clone())
.unwrap_or_else(|| "unknown".to_string());
let response_text = serde_json::to_string_pretty(&outcome)
.unwrap_or_else(|_| "{}".to_string());
tokio::fs::write(stage_dir.join("response.md"), &response_text).await?;
Ok(Candidate {
id: best_id,
status: outcome.status.to_string(),
score: 0.0,
})
}
Ok(CodergenResult::Text(text)) => {
// Write response to logs
tokio::fs::write(stage_dir.join("response.md"), &text).await?;
// The LLM responded with text; try to find a matching candidate ID
let text = text.trim().to_string();
let empty_vec = vec![];
@ -162,9 +206,14 @@ async fn llm_evaluate(
.and_then(|v| v.as_str())
.unwrap_or("success")
.to_string();
let score = v
.get("score")
.and_then(|v| v.as_f64())
.unwrap_or(0.0);
return Ok(Candidate {
id: id.to_string(),
status,
score,
});
}
}
@ -294,6 +343,7 @@ mod tests {
#[tokio::test]
async fn fan_in_with_backend_llm_eval() {
use crate::handler::codergen::CodergenBackend;
use tempfile::TempDir;
struct MockBackend;
@ -325,10 +375,10 @@ mod tests {
]),
);
let graph = Graph::new("test");
let logs_root = Path::new("/tmp/test");
let tmp = TempDir::new().unwrap();
let outcome = handler
.execute(&node, &context, &graph, logs_root)
.execute(&node, &context, &graph, tmp.path())
.await
.unwrap();
assert_eq!(outcome.status, StageStatus::Success);
@ -337,5 +387,72 @@ mod tests {
outcome.context_updates.get("parallel.fan_in.best_id"),
Some(&serde_json::json!("branch_b"))
);
// Verify prompt and response files were written
let prompt_path = tmp.path().join("fan_in").join("prompt.md");
assert!(prompt_path.exists());
let prompt_content = std::fs::read_to_string(&prompt_path).unwrap();
assert!(prompt_content.contains("Pick the best branch"));
let response_path = tmp.path().join("fan_in").join("response.md");
assert!(response_path.exists());
let response_content = std::fs::read_to_string(&response_path).unwrap();
assert!(response_content.contains("branch_b"));
}
#[tokio::test]
async fn fan_in_all_fail_returns_fail() {
let handler = FanInHandler::new(None);
let node = Node::new("fan_in");
let context = Context::new();
context.set(
"parallel.results",
serde_json::json!([
{"id": "branch_a", "status": "fail"},
{"id": "branch_b", "status": "fail"},
{"id": "branch_c", "status": "fail"},
]),
);
let graph = Graph::new("test");
let logs_root = Path::new("/tmp/test");
let outcome = handler
.execute(&node, &context, &graph, logs_root)
.await
.unwrap();
assert_eq!(outcome.status, StageStatus::Fail);
assert!(outcome
.failure_reason
.as_deref()
.unwrap()
.contains("all candidates failed"));
}
#[tokio::test]
async fn fan_in_score_tiebreak() {
let handler = FanInHandler::new(None);
let node = Node::new("fan_in");
let context = Context::new();
context.set(
"parallel.results",
serde_json::json!([
{"id": "branch_a", "status": "success", "score": 0.5},
{"id": "branch_b", "status": "success", "score": 0.9},
{"id": "branch_c", "status": "success", "score": 0.7},
]),
);
let graph = Graph::new("test");
let logs_root = Path::new("/tmp/test");
let outcome = handler
.execute(&node, &context, &graph, logs_root)
.await
.unwrap();
assert_eq!(outcome.status, StageStatus::Success);
// branch_b has highest score
assert_eq!(
outcome.context_updates.get("parallel.fan_in.best_id"),
Some(&serde_json::json!("branch_b"))
);
}
}

View file

@ -1,5 +1,5 @@
use std::path::Path;
use std::time::Duration;
use std::time::{Duration, Instant};
use async_trait::async_trait;
@ -14,6 +14,16 @@ use super::Handler;
/// Trait for observing child pipeline state during the manager loop.
#[async_trait]
pub trait ChildObserver: Send + Sync {
/// Launch the child pipeline. Called before the observation loop when child_autostart is true.
async fn launch_child(
&self,
_dotfile: &str,
_workdir: &str,
_context: &Context,
) -> Result<(), AttractorError> {
Ok(())
}
/// Ingest child telemetry into the context.
async fn observe(&self, context: &Context) -> Result<(), AttractorError>;
@ -99,6 +109,39 @@ impl Handler for ManagerLoopHandler {
let do_steer = actions_str.contains("steer");
let do_wait = actions_str.contains("wait");
// Child autostart: launch child pipeline before the observation loop
let child_autostart = node
.attrs
.get("stack.child_autostart")
.and_then(|v| v.as_str())
.unwrap_or("true");
if child_autostart != "false" {
if let Some(ref observer) = self.observer {
let child_dotfile = node
.attrs
.get("stack.child_dotfile")
.and_then(|v| v.as_str())
.unwrap_or("");
let child_workdir = node
.attrs
.get("stack.child_workdir")
.and_then(|v| v.as_str())
.unwrap_or("");
observer
.launch_child(child_dotfile, child_workdir, context)
.await?;
}
}
// Steer cooldown tracking
let steer_cooldown = node
.attrs
.get("manager.steer_cooldown")
.and_then(|v| v.as_str())
.map(parse_duration_str)
.unwrap_or(Duration::ZERO);
let mut last_steer_time: Option<Instant> = None;
// Observation loop
for cycle in 1..=max_cycles {
// Observe
@ -108,10 +151,17 @@ impl Handler for ManagerLoopHandler {
}
}
// Steer
// Steer (with cooldown)
if do_steer {
if let Some(ref observer) = self.observer {
observer.steer(context, node).await?;
let cooldown_elapsed = match last_steer_time {
Some(t) => t.elapsed() >= steer_cooldown,
None => true,
};
if cooldown_elapsed {
if let Some(ref observer) = self.observer {
observer.steer(context, node).await?;
last_steer_time = Some(Instant::now());
}
}
}

View file

@ -1,8 +1,6 @@
use crate::error::AttractorError;
use crate::graph::Graph;
use crate::transform::{
PreambleTransform, StylesheetApplicationTransform, Transform, VariableExpansionTransform,
};
use crate::transform::{StylesheetApplicationTransform, Transform, VariableExpansionTransform};
use crate::validation::{self, Diagnostic};
/// Builder for configuring and executing a pipeline preparation.
@ -33,10 +31,9 @@ impl PipelineBuilder {
pub fn prepare(&self, dot_source: &str) -> Result<(Graph, Vec<Diagnostic>), AttractorError> {
let mut graph = crate::parser::parse(dot_source)?;
// Built-in transforms
// Built-in transforms (PreambleTransform moved to engine execution time)
VariableExpansionTransform.apply(&mut graph);
StylesheetApplicationTransform.apply(&mut graph);
PreambleTransform.apply(&mut graph);
// Custom transforms
for transform in &self.transforms {

View file

@ -6,11 +6,9 @@ use crate::graph::types::{AttrValue, Graph};
pub enum Selector {
/// `*` -- matches all nodes, specificity 0.
Universal,
/// bare word like `box` -- matches nodes whose shape equals that word, specificity 1.
Shape(String),
/// `.classname` -- matches nodes with that class, specificity 2.
/// `.classname` -- matches nodes with that class, specificity 1.
Class(String),
/// `#nodeid` -- matches a specific node, specificity 3.
/// `#nodeid` -- matches a specific node, specificity 2.
Id(String),
}
@ -19,9 +17,8 @@ impl Selector {
pub const fn specificity(&self) -> u8 {
match self {
Self::Universal => 0,
Self::Shape(_) => 1,
Self::Class(_) => 2,
Self::Id(_) => 3,
Self::Class(_) => 1,
Self::Id(_) => 2,
}
}
}
@ -115,19 +112,10 @@ fn parse_selector(remaining: &mut &str) -> Result<Selector, AttractorError> {
*remaining = remaining[end..].trim();
Ok(Selector::Class(class))
} else {
// Bare word: shape selector
let end = remaining
.find(|c: char| !c.is_ascii_alphanumeric() && c != '_' && c != '-')
.unwrap_or(remaining.len());
if end == 0 {
return Err(AttractorError::Stylesheet(format!(
"expected selector ('*', '#id', '.class', or shape name), got: {:?}",
&remaining[..remaining.len().min(20)]
)));
}
let shape = remaining[..end].to_string();
*remaining = remaining[end..].trim();
Ok(Selector::Shape(shape))
Err(AttractorError::Stylesheet(format!(
"expected selector ('*', '#id', or '.class'), got: {:?}",
&remaining[..remaining.len().min(20)]
)))
}
}
@ -201,7 +189,6 @@ pub fn apply_stylesheet(stylesheet: &Stylesheet, graph: &mut Graph) {
let node = &graph.nodes[node_id.as_str()];
let matches = match &rule.selector {
Selector::Universal => true,
Selector::Shape(shape) => node.shape() == shape.as_str(),
Selector::Class(cls) => node.classes.contains(cls),
Selector::Id(id) => node_id == id,
};
@ -386,9 +373,8 @@ mod tests {
#[test]
fn selector_specificity_values() {
assert_eq!(Selector::Universal.specificity(), 0);
assert_eq!(Selector::Shape("box".into()).specificity(), 1);
assert_eq!(Selector::Class("x".into()).specificity(), 2);
assert_eq!(Selector::Id("x".into()).specificity(), 3);
assert_eq!(Selector::Class("x".into()).specificity(), 1);
assert_eq!(Selector::Id("x".into()).specificity(), 2);
}
#[test]
@ -442,76 +428,36 @@ mod tests {
}
#[test]
fn parse_shape_selector() {
let ss = parse_stylesheet("box { llm_model: opus; }").unwrap();
assert_eq!(ss.rules.len(), 1);
assert_eq!(ss.rules[0].selector, Selector::Shape("box".into()));
fn parse_bare_word_selector_is_error() {
let result = parse_stylesheet("box { llm_model: opus; }");
assert!(result.is_err());
}
#[test]
fn apply_shape_selector_matches_by_shape() {
let ss = parse_stylesheet("box { llm_model: opus; }").unwrap();
let mut graph = Graph::new("test");
// Default shape is "box"
let node_a = Node::new("a");
graph.nodes.insert("a".into(), node_a);
let mut node_b = Node::new("b");
node_b
.attrs
.insert("shape".into(), AttrValue::String("diamond".into()));
graph.nodes.insert("b".into(), node_b);
apply_stylesheet(&ss, &mut graph);
assert_eq!(
graph.nodes["a"].attrs.get("llm_model"),
Some(&AttrValue::String("opus".into()))
);
// diamond shape should not match "box" selector
assert_eq!(graph.nodes["b"].attrs.get("llm_model"), None);
}
#[test]
fn shape_selector_specificity_between_universal_and_class() {
fn class_overrides_universal_specificity() {
let ss = parse_stylesheet(
"* { llm_model: sonnet; } box { llm_model: opus; } .special { llm_model: gpt; }",
"* { llm_model: sonnet; } .special { llm_model: gpt; }",
)
.unwrap();
let mut graph = Graph::new("test");
// Node with default shape "box" and class "special"
let mut node_a = Node::new("a");
node_a.classes.push("special".into());
graph.nodes.insert("a".into(), node_a);
// Node with default shape "box" and no class
let node_b = Node::new("b");
graph.nodes.insert("b".into(), node_b);
// Node with shape "diamond" and no class
let mut node_c = Node::new("c");
node_c
.attrs
.insert("shape".into(), AttrValue::String("diamond".into()));
graph.nodes.insert("c".into(), node_c);
apply_stylesheet(&ss, &mut graph);
// .special (specificity 2) overrides box (specificity 1)
// .special (specificity 1) overrides * (specificity 0)
assert_eq!(
graph.nodes["a"].attrs.get("llm_model"),
Some(&AttrValue::String("gpt".into()))
);
// box (specificity 1) overrides * (specificity 0)
// No class, gets universal
assert_eq!(
graph.nodes["b"].attrs.get("llm_model"),
Some(&AttrValue::String("opus".into()))
);
// diamond doesn't match "box", so gets universal
assert_eq!(
graph.nodes["c"].attrs.get("llm_model"),
Some(&AttrValue::String("sonnet".into()))
);
}

View file

@ -4,8 +4,8 @@ use crate::graph::{AttrValue, Graph};
use super::{Diagnostic, LintRule, Severity};
/// Returns all 14 built-in lint rules.
#[must_use]
/// Returns all 15 built-in lint rules.
#[must_use]
pub fn built_in_rules() -> Vec<Box<dyn LintRule>> {
vec![
Box::new(StartNodeRule),
@ -22,6 +22,7 @@ pub fn built_in_rules() -> Vec<Box<dyn LintRule>> {
Box::new(GoalGateHasRetryRule),
Box::new(PromptOnLlmNodesRule),
Box::new(FreeformEdgeCountRule),
Box::new(DirectionValidRule),
]
}
@ -329,21 +330,17 @@ impl LintRule for StylesheetSyntaxRule {
if stylesheet.is_empty() {
return Vec::new();
}
let open_count = stylesheet.chars().filter(|c| *c == '{').count();
let close_count = stylesheet.chars().filter(|c| *c == '}').count();
if open_count != close_count {
return vec![Diagnostic {
match crate::stylesheet::parse_stylesheet(stylesheet) {
Ok(_) => Vec::new(),
Err(e) => vec![Diagnostic {
rule: self.name().to_string(),
severity: Severity::Error,
message: format!(
"Model stylesheet has unbalanced braces ({open_count} open, {close_count} close)"
),
message: format!("Model stylesheet parse error: {e}"),
node_id: None,
edge: None,
fix: Some("Balance the curly braces in model_stylesheet".to_string()),
}];
fix: Some("Fix the model_stylesheet syntax".to_string()),
}],
}
Vec::new()
}
}
@ -662,6 +659,37 @@ impl LintRule for FreeformEdgeCountRule {
}
}
// --- Rule 15: direction_valid (WARNING) ---
struct DirectionValidRule;
const VALID_DIRECTIONS: &[&str] = &["TB", "LR", "BT", "RL"];
impl LintRule for DirectionValidRule {
fn name(&self) -> &'static str {
"direction_valid"
}
fn apply(&self, graph: &Graph) -> Vec<Diagnostic> {
let Some(rankdir) = graph.attrs.get("rankdir").and_then(AttrValue::as_str) else {
return Vec::new();
};
if VALID_DIRECTIONS.contains(&rankdir) {
return Vec::new();
}
vec![Diagnostic {
rule: self.name().to_string(),
severity: Severity::Warning,
message: format!(
"Graph has invalid rankdir '{rankdir}'"
),
node_id: None,
edge: None,
fix: Some(format!("Use one of: {}", VALID_DIRECTIONS.join(", "))),
}]
}
}
#[cfg(test)]
mod tests {
use super::*;
@ -1123,8 +1151,75 @@ mod tests {
// built_in_rules tests
#[test]
fn built_in_rules_returns_14_rules() {
fn built_in_rules_returns_15_rules() {
let rules = built_in_rules();
assert_eq!(rules.len(), 14);
assert_eq!(rules.len(), 15);
}
// direction_valid rule tests
#[test]
fn direction_valid_rule_no_rankdir() {
let g = minimal_graph();
let rule = DirectionValidRule;
let d = rule.apply(&g);
assert!(d.is_empty());
}
#[test]
fn direction_valid_rule_valid_directions() {
let rule = DirectionValidRule;
let mut g = minimal_graph();
g.attrs.insert(
"rankdir".to_string(),
AttrValue::String("LR".to_string()),
);
assert!(rule.apply(&g).is_empty());
g.attrs.insert(
"rankdir".to_string(),
AttrValue::String("TB".to_string()),
);
assert!(rule.apply(&g).is_empty());
g.attrs.insert(
"rankdir".to_string(),
AttrValue::String("BT".to_string()),
);
assert!(rule.apply(&g).is_empty());
g.attrs.insert(
"rankdir".to_string(),
AttrValue::String("RL".to_string()),
);
assert!(rule.apply(&g).is_empty());
}
#[test]
fn direction_valid_rule_invalid_direction() {
let mut g = minimal_graph();
g.attrs.insert(
"rankdir".to_string(),
AttrValue::String("XY".to_string()),
);
let rule = DirectionValidRule;
let d = rule.apply(&g);
assert_eq!(d.len(), 1);
assert_eq!(d[0].severity, Severity::Warning);
}
// stylesheet_syntax with full parse tests
#[test]
fn stylesheet_syntax_rule_malformed_selector() {
let mut g = minimal_graph();
g.attrs.insert(
"model_stylesheet".to_string(),
AttrValue::String("* { garbage garbage }".to_string()),
);
let rule = StylesheetSyntaxRule;
let d = rule.apply(&g);
assert_eq!(d.len(), 1);
assert_eq!(d[0].severity, Severity::Error);
}
}

View file

@ -986,15 +986,19 @@ fn checkpoint_save_and_resume_roundtrip() {
ctx.append_log("started");
ctx.append_log("step_1 completed");
let mut checkpoint = Checkpoint::from_context(
let mut retries = std::collections::HashMap::new();
retries.insert("step_1".to_string(), 1u32);
let checkpoint = Checkpoint::from_context(
&ctx,
"step_2",
vec![
"start".to_string(),
"step_1".to_string(),
],
retries,
std::collections::HashMap::new(),
None,
);
checkpoint.node_retries.insert("step_1".to_string(), 1);
checkpoint.save(&path).expect("save should succeed");
@ -1178,3 +1182,266 @@ async fn smoke_test_with_mock_codergen_backend() {
.expect("plan prompt should exist");
assert_eq!(plan_prompt, "Plan to achieve: Build and validate");
}
// ---------------------------------------------------------------------------
// 12. Parallel fan-out / fan-in integration test (Gap #14)
// ---------------------------------------------------------------------------
#[tokio::test]
async fn end_to_end_parallel_fan_out_fan_in() {
use attractor::handler::fan_in::FanInHandler;
use attractor::handler::parallel::ParallelHandler;
use std::sync::Arc;
let input = r#"digraph parallel_test {
start [shape=Mdiamond]
fan_out [shape=component]
branch_a [shape=box, prompt="Branch A work"]
branch_b [shape=box, prompt="Branch B work"]
fan_in_node [shape=tripleoctagon]
done [shape=Msquare]
start -> fan_out
fan_out -> branch_a
fan_out -> branch_b
branch_a -> fan_in_node
branch_b -> fan_in_node
fan_in_node -> done
}"#;
let graph = parse(input).expect("parse should succeed");
validate_or_raise(&graph, &[]).expect("validation should pass");
let dir = tempfile::tempdir().unwrap();
let mut registry = HandlerRegistry::new(
Box::new(CodergenHandler::new(Some(Box::new(MockCodergenBackend)))),
);
registry.register("start", Box::new(StartHandler));
registry.register("exit", Box::new(ExitHandler));
registry.register(
"codergen",
Box::new(CodergenHandler::new(Some(Box::new(MockCodergenBackend)))),
);
let registry = Arc::new(registry);
let emitter = Arc::new(EventEmitter::new());
let parallel_handler = ParallelHandler::new(Arc::clone(&registry), Arc::clone(&emitter));
let fan_in_handler = FanInHandler::new(Some(Box::new(MockCodergenBackend)));
// Build a new registry with parallel and fan_in registered
let mut full_registry = HandlerRegistry::new(
Box::new(CodergenHandler::new(Some(Box::new(MockCodergenBackend)))),
);
full_registry.register("start", Box::new(StartHandler));
full_registry.register("exit", Box::new(ExitHandler));
full_registry.register(
"codergen",
Box::new(CodergenHandler::new(Some(Box::new(MockCodergenBackend)))),
);
full_registry.register("parallel", Box::new(parallel_handler));
full_registry.register("parallel.fan_in", Box::new(fan_in_handler));
let engine = PipelineEngine::new(full_registry, EventEmitter::new());
let config = RunConfig {
logs_root: dir.path().to_path_buf(),
};
let outcome = engine
.run(&graph, &config)
.await
.expect("parallel pipeline should succeed");
assert_eq!(outcome.status, StageStatus::Success);
let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap();
// The parallel node (fan_out) and fan_in_node should be in completed_nodes.
// Branch nodes run inside the parallel handler, so they are not recorded
// individually by the engine -- but fan_out and fan_in_node are top-level.
assert!(
checkpoint
.completed_nodes
.contains(&"fan_out".to_string()),
"fan_out should have been executed"
);
assert!(
checkpoint
.completed_nodes
.contains(&"fan_in_node".to_string()),
"fan_in_node should have been executed"
);
// Verify parallel.results was populated (both branches ran)
let parallel_results = checkpoint
.context_values
.get("parallel.results")
.expect("parallel.results should be in context");
let results_arr = parallel_results.as_array().expect("should be an array");
assert_eq!(results_arr.len(), 2, "should have 2 branch results");
}
// ---------------------------------------------------------------------------
// 13. Resume from checkpoint (P1)
// ---------------------------------------------------------------------------
#[tokio::test]
async fn resume_from_checkpoint_completes_pipeline() {
// Build a pipeline: start -> step_a -> step_b -> exit
// Create a checkpoint mid-pipeline (after step_a) and verify
// run_from_checkpoint completes from step_b onward.
let mut graph = Graph::new("ResumeTest");
graph.attrs.insert(
"goal".to_string(),
AttrValue::String("Test resume".to_string()),
);
let mut start = Node::new("start");
start.attrs.insert(
"shape".to_string(),
AttrValue::String("Mdiamond".to_string()),
);
graph.nodes.insert("start".to_string(), start);
let mut exit = Node::new("exit");
exit.attrs.insert(
"shape".to_string(),
AttrValue::String("Msquare".to_string()),
);
graph.nodes.insert("exit".to_string(), exit);
let step_a = Node::new("step_a");
graph.nodes.insert("step_a".to_string(), step_a);
let step_b = Node::new("step_b");
graph.nodes.insert("step_b".to_string(), step_b);
graph.edges.push(Edge::new("start", "step_a"));
graph.edges.push(Edge::new("step_a", "step_b"));
graph.edges.push(Edge::new("step_b", "exit"));
// Simulate a checkpoint saved after step_a completed.
// The checkpoint records step_a as current_node with next_node_id = step_b.
let ctx = Context::new();
ctx.set("graph.goal", serde_json::json!("Test resume"));
ctx.set("outcome", serde_json::json!("success"));
let mut outcomes = std::collections::HashMap::new();
outcomes.insert("start".to_string(), Outcome::success());
outcomes.insert("step_a".to_string(), Outcome::success());
let checkpoint = Checkpoint::from_context(
&ctx,
"step_a",
vec!["start".to_string(), "step_a".to_string()],
std::collections::HashMap::new(),
outcomes,
Some("step_b".to_string()),
);
let dir = tempfile::tempdir().unwrap();
let mut registry = HandlerRegistry::new(Box::new(StartHandler));
registry.register("start", Box::new(StartHandler));
registry.register("exit", Box::new(ExitHandler));
let engine = PipelineEngine::new(registry, EventEmitter::new());
let config = RunConfig {
logs_root: dir.path().to_path_buf(),
};
let outcome = engine
.run_from_checkpoint(&graph, &config, &checkpoint)
.await
.expect("resume should succeed");
assert_eq!(outcome.status, StageStatus::Success);
// Verify checkpoint written after resume contains step_b
let final_cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap();
assert!(
final_cp.completed_nodes.contains(&"step_b".to_string()),
"step_b should have been executed after resume"
);
// step_a should also be present (carried over from the checkpoint)
assert!(
final_cp.completed_nodes.contains(&"step_a".to_string()),
"step_a should be preserved from checkpoint"
);
// start should also be present
assert!(
final_cp.completed_nodes.contains(&"start".to_string()),
"start should be preserved from checkpoint"
);
}
#[tokio::test]
async fn resume_from_checkpoint_preserves_goal_gate_outcomes() {
// Build: start -> gated_work (goal_gate=true) -> step_b -> exit
// Checkpoint after gated_work (success), resume at step_b.
// At exit, goal gate should pass because outcomes are restored.
let mut graph = Graph::new("ResumeGoalGateTest");
let mut start = Node::new("start");
start.attrs.insert(
"shape".to_string(),
AttrValue::String("Mdiamond".to_string()),
);
graph.nodes.insert("start".to_string(), start);
let mut exit = Node::new("exit");
exit.attrs.insert(
"shape".to_string(),
AttrValue::String("Msquare".to_string()),
);
graph.nodes.insert("exit".to_string(), exit);
let mut gated_work = Node::new("gated_work");
gated_work.attrs.insert(
"goal_gate".to_string(),
AttrValue::Boolean(true),
);
graph.nodes.insert("gated_work".to_string(), gated_work);
let step_b = Node::new("step_b");
graph.nodes.insert("step_b".to_string(), step_b);
graph.edges.push(Edge::new("start", "gated_work"));
graph.edges.push(Edge::new("gated_work", "step_b"));
graph.edges.push(Edge::new("step_b", "exit"));
// Checkpoint: gated_work completed with success, next is step_b
let ctx = Context::new();
ctx.set("outcome", serde_json::json!("success"));
let mut outcomes = std::collections::HashMap::new();
outcomes.insert("start".to_string(), Outcome::success());
outcomes.insert("gated_work".to_string(), Outcome::success());
let checkpoint = Checkpoint::from_context(
&ctx,
"gated_work",
vec!["start".to_string(), "gated_work".to_string()],
std::collections::HashMap::new(),
outcomes,
Some("step_b".to_string()),
);
let dir = tempfile::tempdir().unwrap();
let mut registry = HandlerRegistry::new(Box::new(StartHandler));
registry.register("start", Box::new(StartHandler));
registry.register("exit", Box::new(ExitHandler));
let engine = PipelineEngine::new(registry, EventEmitter::new());
let config = RunConfig {
logs_root: dir.path().to_path_buf(),
};
// This should succeed because goal gate for gated_work is satisfied
// via restored outcomes
let outcome = engine
.run_from_checkpoint(&graph, &config, &checkpoint)
.await
.expect("resume with goal gate should succeed");
assert_eq!(outcome.status, StageStatus::Success);
}

View file

@ -0,0 +1,185 @@
# Attractor Spec: Confirmed Gaps Report
**Date:** 2026-02-21
**Method:** 3 parallel agents investigated 19 claimed gaps, reading source code and citing exact file:line evidence.
**Result: 16 CONFIRMED / 2 REFUTED / 1 PARTIAL**
---
## Hard Gaps (3/3 confirmed)
### 1. Auto Status — engine never checks `auto_status` attribute
**CONFIRMED**
The `auto_status()` accessor exists at `graph/types.rs:193-194` but is never referenced in `engine.rs` or any handler. A grep for `auto_status` across all source returns only the accessor definition and its unit test — zero engine references. The spec (line 162, Appendix C line 2111) says: when `auto_status=true` and no `status.json` was written by the handler, the engine should synthesize `{"outcome": "success", "notes": "auto-status: handler completed without writing status"}`. This does not happen.
### 2. Timeout Enforcement — no deadline around handler execution
**CONFIRMED**
`tokio::time::timeout` appears in exactly two places: `handler/tool.rs:67` (subprocess timeout) and `interviewer/mod.rs:142` (human input timeout). The engine's `execute_with_retry` at `engine.rs:449-536` wraps handler execution in `catch_unwind` for panic safety (line 464) but has **no** `tokio::time::timeout` wrapper. A grep for `timeout` in `engine.rs` returns zero matches. All non-tool, non-interviewer handlers (codergen, manager_loop, parallel, etc.) can run indefinitely.
### 3. Manager Loop `child_autostart` — not implemented
**CONFIRMED**
The spec (lines 961-963) says: `IF node.attrs.get("stack.child_autostart", "true") == "true": start_child_pipeline(child_dotfile)`. A codebase-wide grep for `child_autostart`, `start_child_pipeline`, and `child_dotfile` returns matches only in the spec itself and the review document. `handler/manager_loop.rs` goes directly into its observation loop without any auto-start logic. The `ChildObserver` trait (lines 16-22) provides `observe` and `steer` but no launch capability. Sub-gap: `steer_cooldown_elapsed()` is also missing — zero matches for `steer_cooldown` anywhere.
---
## Minor Gaps (13/16 confirmed, 2 refuted, 1 partial)
### 4. Thread ID Resolution — only step 1 of 5
**CONFIRMED**
The spec (lines 1196-1206) defines 5-step thread resolution for `full` fidelity. `engine.rs:660-664` only handles step 1 (node `thread_id`). Missing:
- Step 2: Edge `thread_id` — accessor exists at `graph/types.rs:262-264` but engine never reads it
- Step 3: Graph-level default thread — no `default_thread` accessor (zero grep matches)
- Step 4: Derived class from enclosing subgraph — engine never uses `classes` for thread resolution
- Step 5: Fallback to previous node ID — not implemented
### 5. Checkpoint Resume — retry counters not restored
**CONFIRMED**
`Checkpoint.node_retries` field exists at `checkpoint.rs:17` but `Checkpoint::from_context()` at line 28 always initializes it as `HashMap::new()`. A grep for `node_retries` in `engine.rs` returns zero matches. The field is never populated during saves and never read during resume. The spec's `reset_retry_counter`/`increment_retry_counter` functions (lines 499-504) do not exist.
### 6. Checkpoint Resume — fidelity degradation missing
**CONFIRMED**
The spec (line 1165) says: "If the previous node used `full` fidelity, degrade to `summary:high` for the first resumed node." A grep for `fidelity.*degrad|summary.high.*resume` across all sources returns zero matches. The resume code at `engine.rs:591-608` performs no fidelity degradation.
### 7. Retry Policy `should_retry` — too coarse
**CONFIRMED**
`default_should_retry()` at `engine.rs:100-103`: `Arc::new(|_| true)` — retries ALL errors. The spec (line 556) requires: retry on 429/5xx/network errors, no retry on 401/403/400/validation/config errors. `build_retry_policy` at lines 178-189 always uses the default. `AttractorError` has no `is_retryable()` method.
### 8. Direction type not validated
**CONFIRMED**
The BNF defines `Direction ::= 'TB' | 'LR' | 'BT' | 'RL'` (spec line 107). `grammar.rs:57-66` accepts any `identifier = value` as a graph attr declaration. `semantic.rs:178-179` inserts without validation. None of the 14 validation rules check direction values. `rankdir=XY` would be accepted silently.
### 9. Stylesheet `stylesheet_syntax` lint — brace balance only
**CONFIRMED**
`validation/rules.rs:320-348`: the rule counts `{` and `}` characters and errors only if counts differ. It does not call `parse_stylesheet()` from `stylesheet.rs:54-85` which performs full parsing (selector validation, declaration parsing, proper error messages). A stylesheet like `* { garbage garbage }` passes the lint.
### 10. Stylesheet — undocumented Shape selector
**CONFIRMED**
The spec grammar (line 1497) defines: `Selector ::= '*' | '#' Identifier | '.' ClassName` — three types. `stylesheet.rs:6-15` defines four: `Universal`, `Shape(String)`, `Class(String)`, `Id(String)`. The `Shape` selector is parsed at lines 117-131 (bare-word fallback). This shifts specificity: spec says `*`=0, `.class`=1, `#id`=2; impl has `*`=0, `Shape`=1, `.class`=2, `#id`=3. Relative ordering preserved but absolute values differ. Test at line 446 confirms `box { llm_model: opus; }` produces `Selector::Shape("box")`.
### 11. Fan-in — no score sort, all-fail returns SUCCESS
**CONFIRMED** (both sub-claims)
The spec (line 919) sorts by `(outcome_rank, -c.score, c.id)`. `fan_in.rs:103-109` sorts by `(status_rank, id)` — no `score` field on `Candidate` (lines 62-65), zero grep matches for "score". For all-fail: the spec (line 923) says "Only when all candidates fail does fan-in return FAIL." `fan_in.rs:41-58` always builds `Outcome::success()` at line 47 regardless of whether all candidates failed.
### 12. Preamble transform — applied at parse time, not execution time
**CONFIRMED**
The spec (line 1602): "Applied at execution time (not at parse time) since it depends on runtime state." `pipeline.rs:39` calls `PreambleTransform.apply(&mut graph)` in `prepare()` alongside other parse-time transforms. `transform.rs:28-49` reads `fidelity` from static node attributes but has no access to runtime state (e.g., edge-level fidelity overrides resolved at execution time in `engine.rs:656`). If a node has `fidelity="full"` but an incoming edge has `fidelity="truncate"`, the preamble would be incorrectly missing.
### 13. Pre-hook non-zero — returns fail instead of skip
**CONFIRMED**
The spec (line 1693): "non-zero means skip the tool call." `codergen.rs:98-103` returns `Outcome::fail("pre-hook failed, skipping LLM call")` — `StageStatus::Fail` is semantically stronger than skipping. Test at line 302 confirms the outcome status is `Fail`.
### 14. No parallel fan-out/fan-in integration test
**CONFIRMED**
`tests/integration.rs` (1181 lines) contains 11+ test functions covering linear, branching, human gate, goal gate, retry, stylesheet, checkpoint, and smoke test pipelines. None involve `component` (parallel) or `tripleoctagon` (fan-in) shape nodes. A grep for `parallel|fan_in|fan_out` in `crates/attractor/tests/` returns zero matches.
### 15. Manifest missing `goal` field
**CONFIRMED**
The spec (line 1260): "manifest.json -- Pipeline metadata (name, goal, start time)." `engine.rs:217-233` `write_manifest()` writes `pipeline_name`, `start_time`, `node_count`, `edge_count` — no `goal` field despite `graph.goal()` being available. Integration test at lines 1527-1532 confirms the four fields without `goal`.
### 16. Error categories — no retryable/terminal classification
**CONFIRMED**
The spec (Appendix D, lines 2115-2123) defines Retryable, Terminal, and Pipeline error categories. `error.rs:1-33` has 7 variants (`Parse`, `Validation`, `Engine`, `Handler`, `Checkpoint`, `Stylesheet`, `Io`) with no retryability metadata, no `is_retryable()` method, no `ErrorCategory` enum.
### 17. Spec self-contradicts on `default_max_retry`
**CONFIRMED**
Spec line 138 (Section 2.5 table): `default_max_retry | Integer | 50`. Spec line 481 (Section 3.5): "Built-in default: 0 (no retries)." Implementation follows Section 2 at `graph/types.rs:349-353` with `.unwrap_or(50)`. Test at `engine.rs:931-936` confirms 51 max attempts (50 retries + 1 initial).
---
## Refuted Claims (2)
### R1. Missing variable handling — claimed as GAP
**REFUTED** — Correctly implemented.
`condition.rs:87-110` `resolve_key()` returns `String::new()` for all missing keys (lines 105, 109). Tests at lines 230-239 (`missing_key_compares_as_empty`) and 291-295 (`bare_key_falsy_when_empty`) confirm spec-compliant behavior. The spec (line 1724) says: "Missing keys compare as empty strings" — exactly what the implementation does.
### R2. Status File Contract — claimed "no implementation found"
**REFUTED** — Implemented at two layers.
**Engine-level:** `engine.rs:236-248` `write_node_status()` writes `{node_id}/status.json` with `status`, `notes`, `failure_reason`, `timestamp`. Called at line 702-703 for every node.
**Handler-level:** `codergen.rs:111,150` writes richer `status.json` (full `Outcome` serialized).
**Tests:** `engine.rs:1536-1551` and `integration.rs:189-198` verify status file existence and contents.
Minor sub-gap: the engine-level schema uses key `"status"` while the spec's Appendix C uses `"outcome"`, and the engine-level file omits `preferred_next_label`, `suggested_next_ids`, `context_updates`.
---
## Partial (1)
### P1. Checkpoint Resume — functional but incomplete
**PARTIAL**
`engine.rs:554-561` `run_from_checkpoint()` exists and is callable. Resume logic at lines 591-608 restores context, logs, completed_nodes, and continues from the next node.
**What works:** Basic resume for simple linear pipelines.
**What's missing:**
1. `node_retries` ignored during resume (see gap 5)
2. `node_outcomes` not restored — initialized as empty HashMap at line 586, causing goal gate checks to miss pre-checkpoint outcomes
3. Edge selection during resume picks first outgoing edge (line 603-608), ignoring conditions — wrong successor for conditional graphs
4. No integration test calls `run_from_checkpoint` — only `Checkpoint::save`/`load` is tested
---
## Summary Table
| # | Gap | Verdict | Severity |
|---|-----|---------|----------|
| 1 | Auto status not enforced | CONFIRMED | Hard |
| 2 | Timeout not enforced in engine | CONFIRMED | Hard |
| 3 | Manager loop child_autostart missing | CONFIRMED | Hard |
| 4 | Thread ID resolution 1/5 steps | CONFIRMED | Moderate |
| 5 | Checkpoint retry counters not persisted | CONFIRMED | Moderate |
| 6 | Checkpoint fidelity degradation missing | CONFIRMED | Moderate |
| 7 | should_retry retries all errors | CONFIRMED | Moderate |
| 8 | Direction values not validated | CONFIRMED | Low |
| 9 | Stylesheet lint brace-balance only | CONFIRMED | Low |
| 10 | Undocumented Shape selector | CONFIRMED | Low |
| 11 | Fan-in: no score sort, all-fail=SUCCESS | CONFIRMED | Moderate |
| 12 | Preamble at parse time not runtime | CONFIRMED | Moderate |
| 13 | Pre-hook fail instead of skip | CONFIRMED | Low |
| 14 | No parallel integration test | CONFIRMED | Low |
| 15 | Manifest missing goal field | CONFIRMED | Low |
| 16 | No error retryable/terminal classification | CONFIRMED | Moderate |
| 17 | Spec contradicts itself on default_max_retry | CONFIRMED | Low (spec bug) |
| P1 | Checkpoint resume incomplete | PARTIAL | Moderate |
| R1 | Missing variable handling | REFUTED | — |
| R2 | Status file contract | REFUTED | — |

View file

@ -0,0 +1,224 @@
# Attractor Spec Compliance Review
**Date**: 2026-02-21
**Spec**: `docs/specs/attractor-spec.md`
**Implementation**: `crates/attractor/src/`
---
## Section 1: Overview and Goals
| # | Subsection | Verdict |
|---|-----------|---------|
| 1 | 1.1 Problem Statement | ALIGNED |
| 2 | 1.2 Why DOT Syntax | ALIGNED |
| 3 | 1.3 Design Principles | ALIGNED |
| 4 | 1.4 Layering and LLM Backends | ALIGNED |
**Details**: All five design principles are implemented: declarative pipelines (DOT parsed into `Graph`, engine handles execution), pluggable handlers (`HandlerRegistry`), checkpoint/resume (`checkpoint.rs`), human-in-the-loop (`interviewer/` module), edge-based routing (`select_edge()` in `engine.rs`). `CodergenBackend` trait decouples LLM integration. `EventEmitter` provides the event stream for frontends.
---
## Section 2: DOT DSL Schema
| # | Subsection | Verdict |
|---|-----------|---------|
| 5 | 2.1 Supported Subset | ALIGNED |
| 6 | 2.2 BNF-Style Grammar | ALIGNED |
| 7 | 2.3 Key Constraints | ALIGNED |
| 8 | 2.4 Value Types | ALIGNED |
| 9 | 2.5 Graph Attributes | ALIGNED |
| 10 | 2.6 Node Attributes | ALIGNED |
| 11 | 2.7 Edge Attributes | ALIGNED |
| 12 | 2.8 Shape-to-Handler Mapping | ALIGNED |
| 13 | 2.9 Chained Edges | ALIGNED |
| 14 | 2.10 Subgraphs | ALIGNED |
| 15 | 2.11 Node/Edge Default Blocks | ALIGNED |
| 16 | 2.12 Class Attribute | ALIGNED |
| 17 | 2.13 Minimal Examples | ALIGNED |
**Minor note**: The parser does not produce a specific error message when `strict digraph` is used — it simply fails to parse. Functionally correct but a UX gap for error messaging.
---
## Section 3: Pipeline Execution Engine
| # | Subsection | Verdict |
|---|-----------|---------|
| 18 | 3.1 Run Lifecycle (5 phases) | ALIGNED |
| 19 | 3.2 Core Execution Loop | ALIGNED |
| 20 | 3.3 Edge Selection Algorithm (5-step priority) | ALIGNED |
| 21 | 3.4 Goal Gate Enforcement | ALIGNED |
| 22 | 3.5 Retry Logic | ALIGNED |
| 23 | 3.6 Retry Policies (5 presets) | ALIGNED |
| 24 | 3.7 Failure Routing | ALIGNED |
| 25 | 3.8 Concurrency Model | ALIGNED |
**Minor note**: The 5 named retry presets (none, standard, aggressive, linear, patient) exist as constructors on `RetryPolicy` but `build_retry_policy()` always constructs a custom policy from `max_retries` + default backoff. There is no mechanism for a node to select a preset by name (e.g. `retry_policy="aggressive"`).
---
## Section 4: Node Handlers
| # | Subsection | Verdict |
|---|-----------|---------|
| 26 | 4.1 Handler Interface | ALIGNED |
| 27 | 4.2 Handler Registry | ALIGNED |
| 28 | 4.3 Start Handler | ALIGNED |
| 29 | 4.4 Exit Handler | ALIGNED |
| 30 | 4.5 Codergen Handler | ALIGNED |
| 31 | 4.6 Wait For Human Handler | ALIGNED |
| 32 | 4.7 Conditional Handler | ALIGNED |
| 33 | 4.8 Parallel Handler | ALIGNED |
| 34 | 4.9 Fan-In Handler | ALIGNED |
| 35 | 4.10 Tool Handler | ALIGNED |
| 36 | 4.11 Manager Loop Handler | ALIGNED |
| 37 | 4.12 Custom Handlers | ALIGNED |
**Minor note**: Manager loop reads `stack.child_dotfile` from node attrs rather than graph attrs as the spec pseudocode shows. Arguably better design (per-node child pipeline), but deviates from spec.
---
## Section 5: State and Context
| # | Subsection | Verdict |
|---|-----------|---------|
| 38 | 5.1 PipelineContext | GAP |
| 39 | 5.2 Outcome | ALIGNED |
| 40 | 5.3 Checkpoint | ALIGNED |
| 41 | 5.4 Context Fidelity | GAP |
| 42 | 5.5 Artifact Store | ALIGNED |
| 43 | 5.6 Run Directory Structure | GAP |
**GAP 38 — `last_stage` / `last_response` context keys**: The spec defines these as engine-set context keys. The engine does not set them; only the `codergen` handler sets them. Other handler types do not propagate these keys. Additionally, `internal.retry_count.<node_id>` is tracked in a separate `node_retries` HashMap rather than as a context key.
**GAP 41 — Thread resolution step 3**: The spec lists "Graph-level default thread" as step 3 in thread ID resolution. The implementation uses the node's first CSS class instead. No graph-level default thread concept is implemented.
**GAP 43 — Per-node `prompt.md` / `response.md`**: The spec defines these as part of the run directory structure. The engine does not write them; only the `codergen` handler does. Other LLM-interacting handlers (fan_in with LLM evaluation) do not write these files.
---
## Section 6: Human-in-the-Loop (Interviewer Pattern)
| # | Subsection | Verdict |
|---|-----------|---------|
| 44 | 6.1 Interviewer Interface | ALIGNED |
| 45 | 6.2 Question Model | ALIGNED |
| 46 | 6.3 Answer Model | ALIGNED |
| 47 | 6.4 Built-In Implementations (5) | ALIGNED |
| 48 | 6.5 Timeout Handling | GAP |
| 49 | 6.6 Gate Node Behavior | ALIGNED |
**GAP 48 — WaitHumanHandler bypasses timeout**: `ask_with_timeout()` exists as a utility function but `WaitHumanHandler` calls `interviewer.ask()` directly, so `timeout_seconds` on questions is not enforced for human gate interactions.
---
## Section 7: Validation and Linting
| # | Subsection | Verdict |
|---|-----------|---------|
| 50 | 7.1 Diagnostic Model | ALIGNED |
| 51 | 7.2 Built-In Rules (14 rules) | ALIGNED |
| 52 | 7.3 Validation API | ALIGNED |
| 53 | 7.4 Custom Lint Rules | ALIGNED |
**Details**: All 14 spec rules implemented plus a bonus `direction_valid` rule (15 total). Error-severity diagnostics block execution. Custom rules supported via `extra_rules` parameter.
---
## Section 8: Model Stylesheet
| # | Subsection | Verdict |
|---|-----------|---------|
| 54 | 8.1 Purpose | ALIGNED |
| 55 | 8.2 Grammar | ALIGNED |
| 56 | 8.3 Selectors and Specificity | ALIGNED |
| 57 | 8.4 Recognized Properties | ALIGNED |
| 58 | 8.5 Application/Resolution Order | ALIGNED |
**Details**: Specificity correctly implemented (Universal=0, Class=1, Id=2). Explicit node attributes are never overridden. Stylesheet applied as a transform before validation.
---
## Section 9: Transforms and Extensibility
| # | Subsection | Verdict |
|---|-----------|---------|
| 59 | 9.1 AST Transforms | ALIGNED |
| 60 | 9.2 Built-In Transforms (3) | ALIGNED |
| 61 | 9.3 Custom Transforms | ALIGNED |
| 62 | 9.4 Pipeline Composition | ALIGNED |
| 63 | 9.5 HTTP Server Mode | GAP |
| 64 | 9.6 Observability and Events | ALIGNED |
| 65 | 9.7 Tool Call Hooks | GAP |
**GAP 63 — HTTP Server Mode**: No HTTP server implementation. The spec says "Implementations may expose" making this optional, but it is unimplemented.
**GAP 65 — Tool Call Hooks**: `tool_hooks.pre` and `tool_hooks.post` are defined in the spec for shell commands around LLM tool calls. Not implemented in any handler. Pre-hook should gate tool calls (non-zero exit = skip), post-hook for logging/auditing.
**Minor note**: Transform trait uses `&mut Graph` (in-place mutation) rather than returning a new graph as spec describes. Functionally equivalent.
---
## Section 10: Condition Expression Language
| # | Subsection | Verdict |
|---|-----------|---------|
| 66 | 10.1 Overview | ALIGNED |
| 67 | 10.2 Grammar | ALIGNED |
| 68 | 10.3 Semantics | ALIGNED |
| 69 | 10.4 Variable Resolution | ALIGNED |
| 70 | 10.5 Evaluation | ALIGNED |
| 71 | 10.6 Examples | ALIGNED |
| 72 | 10.7 Extended Operators (future) | ALIGNED |
**Details**: Full implementation with `=`, `!=`, `&&` conjunction, bare key truthiness, `context.*` double-lookup. Correctly does not implement future operators.
---
## Section 11: Definition of Done
| # | Subsection | Verdict |
|---|-----------|---------|
| 73 | 11.1 DOT Parsing | ALIGNED |
| 74 | 11.2 Validation and Linting | ALIGNED |
| 75 | 11.3 Execution Engine | ALIGNED |
| 76 | 11.4 Goal Gate Enforcement | ALIGNED |
| 77 | 11.5 Retry Logic | ALIGNED |
| 78 | 11.6 Node Handlers | ALIGNED |
| 79 | 11.7 State and Context | ALIGNED |
| 80 | 11.8 Human-in-the-Loop | ALIGNED |
| 81 | 11.9 Condition Expressions | ALIGNED |
| 82 | 11.10 Model Stylesheet | ALIGNED |
| 83 | 11.11 Transforms and Extensibility | ALIGNED |
| 84 | 11.12 Cross-Feature Parity Matrix | ALIGNED |
| 85 | 11.13 Integration Smoke Test | ALIGNED |
---
## Summary
| Category | Count |
|----------|-------|
| Total items reviewed | 85 |
| ALIGNED | 78 |
| GAP | 7 |
| Alignment rate | 91.8% |
### All Gaps
| # | Section | Gap | Severity |
|---|---------|-----|----------|
| 38 | 5.1 | `last_stage`/`last_response` not set by engine; `internal.retry_count` not in context | Low |
| 41 | 5.4 | Thread resolution missing graph-level default thread (step 3) | Low |
| 43 | 5.6 | `prompt.md`/`response.md` only written by codergen, not other LLM handlers | Low |
| 48 | 6.5 | WaitHumanHandler calls `ask()` directly, bypassing `ask_with_timeout()` | Medium |
| 63 | 9.5 | HTTP server mode not implemented (spec marks as optional) | Low |
| 65 | 9.7 | `tool_hooks.pre`/`tool_hooks.post` not implemented | Medium |
### Minor Notes (not gaps, but deviations)
- No specific error for `strict digraph` (parser just fails)
- Named retry presets exist but no node-level attribute to select them
- Manager loop reads `child_dotfile` from node attrs not graph attrs
- Transform trait mutates in-place rather than returning new graph

View file

@ -0,0 +1,701 @@
# Attractor Spec Compliance Review
**Date:** 2026-02-21
**Spec:** `docs/specs/attractor-spec.md`
**Implementation:** `crates/attractor/src/`
**Reviewers:** 5 parallel agents, each covering distinct spec sections
---
## Section 1: Overview and Goals
### 1. Section 1.1 — Problem Statement
**ALIGNED**
The implementation delivers a DOT-based directed-graph pipeline runner as described. No code artifact required beyond the overall architecture.
### 2. Section 1.2 — Why DOT Syntax
**ALIGNED**
The parser (`parser/grammar.rs`) starts with `digraph` keyword and builds on directed graph primitives. DOT subset parser implemented from scratch.
### 3. Section 1.3 — Design Principles
**ALIGNED**
- Declarative pipelines: `.dot` files declare graph structure; engine traverses it (`parser/mod.rs:18`, `graph/types.rs:277-294`)
- Pluggable handlers: `Node::handler_type()` at `graph/types.rs:203-209` resolves from `type` attr or shape mapping
- Checkpoint and resume: `checkpoint` module (`lib.rs:2`), checkpoint save/load implemented
- Human-in-the-loop: `interviewer` module (`lib.rs:10`), `wait.human` mapped at `graph/types.rs:77`
- Edge-based routing: `Edge` struct has `condition()`, `weight()`, `label()` accessors at `graph/types.rs:241-274`
### 4. Section 1.4 — Layering and LLM Backends
**ALIGNED**
No LLM SDK dependency. `CodergenBackend` trait decouples LLM calls. Event stream module at `lib.rs:8`.
---
## Section 2: DOT DSL Schema
### 5. Section 2.1 — Supported Subset
**ALIGNED**
Parser accepts only `digraph` (`grammar.rs:140`). No code path for `graph` (undirected) or `strict`. Trailing content rejected at `parser/mod.rs:23-29`.
### 6. Section 2.2 — BNF Grammar
**ALIGNED** (minor gap)
All grammar productions verified against implementation:
- `Graph`, `Statement`, `GraphAttrStmt`, `NodeDefaults`, `EdgeDefaults`, `GraphAttrDecl`, `SubgraphStmt`, `NodeStmt`, `EdgeStmt`, `AttrBlock`, `Attr` — all match at `grammar.rs:14-152`
- `QualifiedId` (dotted keys) supported at `lexer.rs:90-113`
- All value types including Duration with `ms/s/m/h/d` at `lexer.rs:229-246`
- **Minor gap:** `Direction` type (`TB|LR|BT|RL`) parsed as bare `AstValue::Ident` — no validation restricting to valid values. Functionally works but invalid directions accepted silently.
### 7. Section 2.3 — Key Constraints
**ALIGNED**
- One digraph per file: trailing content check at `parser/mod.rs:23-29`
- Bare identifiers: `[A-Za-z_][A-Za-z0-9_]*` at `lexer.rs:82-87`
- Commas required: `separated_list1(preceded(ws, char(',')), attr)` at `grammar.rs:26-31`
- Directed edges only: only `->` parsed at `lexer.rs:262-264`
- Comments: `strip_comments()` at `lexer.rs:5-53` handles `//` and `/* */`
- Semicolons optional: `opt_semi` at `grammar.rs:34-36`
### 8. Section 2.4 — Value Types
**ALIGNED**
All five types implemented with correct syntax:
| Type | Implementation |
|------|---------------|
| String | `lexer.rs:121-177` with `\"`, `\n`, `\t`, `\\` escapes |
| Integer | `lexer.rs:208-226` with sign and float-rejection |
| Float | `lexer.rs:193-205` |
| Boolean | `lexer.rs:180-190` |
| Duration | `lexer.rs:229-246` with `ms/s/m/h/d` units |
### 9. Section 2.5 — Graph-Level Attributes
**ALIGNED**
All 7 attributes present with correct defaults:
| Key | Evidence |
|-----|----------|
| `goal` | `graph/types.rs:333-338` — default `""` |
| `label` | stored as generic attr |
| `model_stylesheet` | `graph/types.rs:341-346` — default `""` |
| `default_max_retry` | `graph/types.rs:349-354` — default `50` |
| `retry_target` | `graph/types.rs:357-359` |
| `fallback_retry_target` | `graph/types.rs:361-366` |
| `default_fidelity` | `graph/types.rs:369-373` |
### 10. Section 2.6 — Node Attributes
**ALIGNED**
All 18 attributes present with correct defaults:
| Key | Evidence |
|-----|----------|
| `label` | `graph/types.rs:119-121` — falls back to node ID |
| `shape` | `graph/types.rs:124-126` — default `"box"` |
| `type` | `graph/types.rs:129-131` |
| `prompt` | `graph/types.rs:134-136` |
| `max_retries` | `graph/types.rs:139-141` — `None` when unset |
| `goal_gate` | `graph/types.rs:144-146` — default `false` |
| `retry_target` | `graph/types.rs:149-151` |
| `fallback_retry_target` | `graph/types.rs:153-156` |
| `fidelity` | `graph/types.rs:159-161` |
| `thread_id` | `graph/types.rs:164-166` |
| `class` | `graph/types.rs:169-171` + parsing at `semantic.rs:103-117` |
| `timeout` | `graph/types.rs:173-175` — `Option<Duration>` |
| `llm_model` | `graph/types.rs:178-180` |
| `llm_provider` | `graph/types.rs:183-185` |
| `reasoning_effort` | `graph/types.rs:188-190` — default `"high"` |
| `auto_status` | `graph/types.rs:193-195` — default `false` |
| `allow_partial` | `graph/types.rs:198-200` — default `false` |
### 11. Section 2.7 — Edge Attributes
**ALIGNED**
All 7 attributes present:
| Key | Evidence |
|-----|----------|
| `label` | `graph/types.rs:242-244` |
| `condition` | `graph/types.rs:247-249` |
| `weight` | `graph/types.rs:252-254` — default `0` |
| `fidelity` | `graph/types.rs:257-259` |
| `thread_id` | `graph/types.rs:262-264` |
| `loop_restart` | `graph/types.rs:267-269` — default `false` |
| `freeform` | `graph/types.rs:272-274` — default `false` |
### 12. Section 2.8 — Shape-to-Handler Mapping
**ALIGNED**
All 9 mappings present at `graph/types.rs:72-85`. Handler resolution with `type` override at `graph/types.rs:203-209`.
### 13. Section 2.9 — Chained Edges
**ALIGNED**
Parsed at `grammar.rs:86-101`, expanded via `windows(2)` at `semantic.rs:132-141`. Test at `semantic.rs:471-484` confirms correct desugaring.
### 14. Section 2.10-2.12 — Subgraphs, Defaults, Class Attribute
**ALIGNED**
- Subgraph scoping: `semantic.rs:145-235` saves/restores defaults per scope
- Class derivation from subgraph label: `semantic.rs:51-58` (`derive_class_from_label`)
- Comma-separated class parsing: `semantic.rs:103-117`, tested at `semantic.rs:487-500`
- Node/edge default blocks: `grammar.rs:44-54`, applied in `semantic.rs:76-83`
---
## Section 3: Pipeline Execution Engine
### 15. Section 3.1 — Run Lifecycle
**ALIGNED**
Five phases implemented:
- PARSE: `pipeline.rs:34` calls `crate::parser::parse(dot_source)`
- VALIDATE: `pipeline.rs:46` calls `validation::validate(&graph, &[])`
- INITIALIZE: `engine.rs:580-625` creates run directory, initializes context
- EXECUTE: `engine.rs:627-771` main loop
- FINALIZE: `engine.rs:773-784` emits `PipelineCompleted`, returns outcome
### 16. Section 3.2 — Core Execution Loop
**ALIGNED**
All 8 steps from spec implemented:
1. Start node resolution: `engine.rs:620-624` via `graph.find_start_node()` at `graph/types.rs:310-320`
2. Terminal check: `engine.rs:633-653`
3. Execute with retry: `engine.rs:669-679`
4. Record completion: `engine.rs:706-708`
5. Apply context updates: `engine.rs:711-715`
6. Save checkpoint: `engine.rs:718-730`
7. Select next edge: `engine.rs:733-753`
8. Loop restart: `engine.rs:760-767`, advance: `engine.rs:768`
### 17. Section 3.3 — Edge Selection Algorithm
**ALIGNED**
Five-step priority fully implemented at `engine.rs:299-356`:
1. Condition matching: `engine.rs:311-321`
2. Preferred label: `engine.rs:324-333` with `normalize_label` at `engine.rs:254-279`
3. Suggested next IDs: `engine.rs:336-342`
4. Weight + lexical tiebreak: `engine.rs:345-352` via `best_by_weight_then_lexical` at `engine.rs:282-295`
5. Fallback: `engine.rs:355`
### 18. Section 3.4 — Goal Gate Enforcement
**ALIGNED**
- `check_goal_gates` at `engine.rs:362-377`: checks goal_gate nodes for SUCCESS/PARTIAL_SUCCESS
- `get_retry_target` at `engine.rs:380-408`: four-level fallback chain (node → node fallback → graph → graph fallback)
- `is_terminal` at `engine.rs:411-414`
### 19. Section 3.5 — Retry Logic
**ALIGNED** (minor gaps)
- `build_retry_policy` at `engine.rs:178-189`: node `max_retries` with `default_max_retry` fallback
- `execute_with_retry` at `engine.rs:449-536`: full retry loop
- **Minor:** Retry counters not persisted to `Checkpoint.node_retries` during execution (field exists at `checkpoint.rs:17` but not populated)
- **Minor:** Spec self-contradicts on `default_max_retry` default (section 3.5 says 0, section 2 says 50). Implementation follows section 2 (50).
### 20. Section 3.6 — Retry Policy / Backoff
**ALIGNED** (minor gap)
- `BackoffConfig` at `engine.rs:28-44` with correct defaults
- `delay_for_attempt` at `engine.rs:50-76`: formula matches spec exactly
- All 5 presets match spec table (none/standard/aggressive/linear/patient)
- **Minor:** Default `should_retry` at `engine.rs:101-103` retries ALL errors. Spec defines granular behavior (retry 429/5xx, fail 401/403/400).
### 21. Section 3.7 — Failure Routing
**ALIGNED**
All 4 priority steps at `engine.rs:733-753`: fail edge → retry_target → fallback_retry_target → pipeline termination.
### 22. Section 3.8 — Concurrency Model
**ALIGNED**
Single-threaded graph traversal at `engine.rs:627-771`. One node at a time.
### 23. Section 3.9 — Auto Status
**GAP**
`auto_status()` accessor exists at `graph/types.rs:193-194` but **engine never references it**. Engine always writes status.json itself at line 703. No auto-synthesis logic for when a handler writes no status.
### 24. Section 3.10 — Timeout Enforcement
**GAP**
Node `timeout` attribute parsed as `Option<Duration>` (item 10) but **not enforced during handler execution** in the engine. The tool handler (`tool.rs:61-78`) does use timeout for subprocess execution, but the engine does not wrap general handler execution with a timeout.
---
## Section 4: Node Handlers
### 25. Section 4.1 — Handler Trait
**ALIGNED**
Trait at `handler/mod.rs:22-31`: `async fn execute(&self, node, context, graph, logs_root) -> Result<Outcome>`. All four parameters match spec. `Result` wrapping for error propagation.
### 26. Section 4.2 — Handler Registry
**ALIGNED**
At `handler/mod.rs:34-74`: `HashMap<String, Box<dyn Handler>>`, `default_handler`, `register()` replaces existing, three-step `resolve()` (explicit type → shape → default). Tests confirm at lines 153-177.
### 27. Section 4.3 — Start Handler
**ALIGNED**
`handler/start.rs:13-26`: returns `Outcome::success()`. No-op.
### 28. Section 4.4 — Exit Handler
**ALIGNED**
`handler/exit.rs:13-26`: returns `Outcome::success()`. No goal gate logic (handled by engine).
### 29. Section 4.5 — Codergen Handler
**ALIGNED**
- Prompt building: `codergen.rs:87-91` — `node.prompt()` falling back to `node.label()`
- `$goal` expansion: `codergen.rs:42-44`
- Log writing: `codergen.rs:94-96` — `prompt.md`, `response.md`, `status.json`
- Backend call: `codergen.rs:106-121` — handles `CodergenResult::Full` and `CodergenResult::Text`
- Simulation mode when no backend: `codergen.rs:122+`
- Context updates `last_stage`/`last_response`: `codergen.rs:139-146`
- Tool hooks (pre/post): `codergen.rs:98-131` (enhancement beyond spec)
- `CodergenBackend` trait at `codergen.rs:20-27` matches spec's `run(node, prompt, context) -> String | Outcome`
### 30. Section 4.6 — Wait Human Handler
**ALIGNED**
Comprehensive implementation at `wait_human.rs`:
- Choice derivation from outgoing edges: lines 110-129
- Freeform edge detection: line 115
- No-edges failure: lines 131-133
- Question building with `MultipleChoice`, `allow_freeform`, `stage`: lines 136-150
- Timeout handling with default choice fallback: lines 162-179
- Skipped handling: lines 183-185
- Fixed-choice match with `suggested_next_ids` and context updates: lines 195-201
- Freeform fallback: lines 204-221
- First-choice fallback: lines 224-226
- Accelerator key parsing `[K] Label`, `K) Label`, `K - Label`, first char: lines 32-71
### 31. Section 4.7 — Conditional Handler
**ALIGNED**
`handler/conditional.rs:14-28`: no-op returning SUCCESS with note. Routing handled by engine.
### 32. Section 4.8 — Parallel Handler
**ALIGNED**
Full implementation at `parallel.rs`:
- 4 join policies: `wait_all`, `first_success`, `k_of_n(K)`, `quorum(fraction)` at lines 36-59
- 3 error policies: `continue`, `fail_fast`, `ignore` at lines 62-75
- `max_parallel` with `Semaphore` for bounded concurrency: lines 113-120
- Context isolation via `clone_context()`: line 126
- Join evaluation with all four policies: lines 249-286
- Fail-fast behavior: lines 185-189
- Context storage for fan-in (`parallel.results`, `parallel.branch_count`): lines 230-241
### 33. Section 4.9 — Fan-In Handler
**ALIGNED** (minor gaps)
- Reads `parallel.results`: `fan_in.rs:34-37`
- LLM-based evaluation when prompt + backend present: lines 41-42
- Heuristic selection by status rank: lines 77-115
- Context updates `parallel.fan_in.best_id`/`best_outcome`: lines 48-55
- **Minor:** No `score`-based sorting in heuristic (spec mentions `-c.score`)
- **Minor:** Returns SUCCESS when results exist even if all candidates failed
### 34. Section 4.10 — Tool Handler
**ALIGNED**
At `tool.rs`:
- Reads `tool_command` from attrs: lines 51-55
- Empty command returns fail: lines 57-59
- Runs via `sh -c` with timeout support: lines 61-78
- Sets `tool.output` to stdout: lines 15-40
### 35. Section 4.11 — Manager Loop Handler
**GAP** (partial)
Core observe/steer/wait cycle implemented at `manager_loop.rs:59-157`:
- Poll interval, max cycles, stop condition, actions parsing: lines 67-100
- Observe/steer delegation: lines 105-115
- Child status check: lines 119-132
- Max cycles exceeded: lines 152-155
- **GAP: `stack.child_autostart`/`start_child_pipeline` not implemented** — spec says autostart child pipeline, impl does not
- **Minor:** `steer_cooldown_elapsed()` not implemented — steers every cycle
### 36. Section 4.12 — Custom Handlers
**ALIGNED**
Trait-based design inherently supports custom handlers via `register()`. `Send + Sync` bounds match spec contract.
---
## Section 5: State and Context
### 37. Section 5.1 — Context
**ALIGNED**
Thread-safe key-value store at `context.rs:1-11`: `Arc<RwLock<HashMap<String, Value>>>` with all spec methods:
- `set()`: line 33-38
- `get()`: line 46-52
- `get_string()`: line 56-60
- `append_log()`: line 67-72
- `snapshot()`: line 80-85
- `clone_context()`: line 98-106 (deep copy for parallel isolation)
- `apply_updates()`: line 113-118
### 38. Section 5.2 — Outcome
**ALIGNED**
At `outcome.rs`:
- `StageStatus` enum (lines 10-16): `Success`, `Fail`, `PartialSuccess`, `Retry`, `Skipped` — all 5 spec variants
- `Outcome` struct (lines 48-60): `status`, `preferred_label`, `suggested_next_ids`, `context_updates`, `notes`, `failure_reason` — all match
- Factory methods (lines 63-107): `success()`, `fail()`, `retry()`, `skipped()`
### 39. Section 5.3 — Checkpoint
**ALIGNED** (minor gaps)
At `checkpoint.rs:13-20`: `timestamp`, `current_node`, `completed_nodes`, `node_retries`, `context_values`, `logs` — all match spec.
- `save()` at line 44-49, `load()` at line 56-61
- Resume at `engine.rs:591-608`: restores context, logs, completed_nodes, resumes from next node
- **Minor:** Retry counters not restored during resume
- **Minor:** Fidelity degradation on resume (`full` → `summary:high`) not implemented per spec line 1166
### 40. Section 5.4 — Fidelity Modes
**ALIGNED** (minor gap)
`resolve_fidelity` at `engine.rs:199-211` implements full 4-level precedence:
1. Edge `fidelity` attribute (line 201)
2. Target node `fidelity` attribute (line 205)
3. Graph `default_fidelity` attribute (line 208)
4. Default: `"compact"` (line 211)
Tests confirm all four levels at `engine.rs:1462-1511`.
Thread tracking: `engine.rs:660-664` stores `thread.{tid}.current_node`.
**Minor gap:** Thread ID resolution only implements step 1 of 5 from spec (node `thread_id`). Missing: edge `thread_id`, graph default, subgraph class derivation, fallback to previous node ID.
### 41. Section 5.5 — Artifact Store
**ALIGNED**
Full implementation at `artifact.rs`:
- `ArtifactStore` with `base_dir` and `RwLock<HashMap>`: lines 31-34
- `store()` with file-backing above 100KB threshold: lines 62-96
- `retrieve()` from memory or file: lines 107-130
- `has()`, `list()`, `remove()`, `clear()`: lines 137-184
- `ArtifactInfo` struct with all 5 fields: lines 16-22
### 42. Section 5.6 — Run Directory Structure
**ALIGNED** (minor gap)
All spec artifacts written: `checkpoint.json`, `manifest.json`, `{node_id}/status.json`, `{node_id}/prompt.md`, `{node_id}/response.md`, `artifacts/{artifact_id}.json`.
- **Minor:** Manifest (`engine.rs:217-232`) includes `node_count`/`edge_count` but not `goal` field from spec.
---
## Section 6: Human-in-the-Loop (Interviewer Pattern)
### 43. Section 6.1 — Interviewer Interface
**ALIGNED**
Three methods at `interviewer/mod.rs:152-167`: `ask()`, `ask_multiple()` (with default sequential impl), `inform()` (with default no-op). Async via `async_trait`.
### 44. Section 6.2 — Question Model
**ALIGNED**
At `mod.rs:29-39`: `text`, `question_type` (4 variants: `YesNo`, `MultipleChoice`, `Freeform`, `Confirmation`), `options` (key + label), `allow_freeform`, `default`, `timeout_seconds`, `stage`, `metadata` — all 8 fields present.
### 45. Section 6.3 — Answer Model
**ALIGNED**
`AnswerValue` enum (lines 57-65): `Yes`, `No`, `Skipped`, `Timeout`, `Selected(String)`, `Text(String)`. `Answer` struct (lines 68-73): `value`, `selected_option`, `text`.
### 46. Section 6.4 — Built-In Interviewers
**ALIGNED** (all 5)
- **AutoApprove** (`auto_approve.rs:10-23`): YES for YesNo/Confirmation, first option for MultipleChoice, "auto-approved" for Freeform
- **Console** (`console.rs:49-86`): `[?]` prefix, option display, freeform fallback, Y/N for YesNo, `>` prompt for Freeform
- **Callback** (`callback.rs:6-22`): delegates to `Box<dyn Fn(Question) -> Answer>`
- **Queue** (`queue.rs:9-28`): `Mutex<VecDeque<Answer>>`, returns SKIPPED when empty
- **Recording** (`recording.rs:8-39`): wraps inner interviewer, records `(Question, Answer)` pairs
### 47. Section 6.5 — Timeout Handling
**ALIGNED**
At `mod.rs:133-149`: uses `tokio::time::timeout`, returns `default_answer.unwrap_or_else(Answer::timeout)`. Tests at lines 269-296.
---
## Section 7: Validation and Linting
### 48. Section 7.1 — Diagnostic Model
**ALIGNED**
At `validation/mod.rs:9-25`: `Severity` (Error/Warning/Info), `Diagnostic` with `rule`, `severity`, `message`, `node_id`, `edge`, `fix` — all fields match.
### 49. Section 7.2 — Built-In Lint Rules
**ALIGNED** (14/14 rules implemented)
At `validation/rules.rs`:
| Rule | Severity | Location |
|------|----------|----------|
| `start_node` | ERROR | lines 30-69 |
| `terminal_node` | ERROR | lines 73-104 |
| `reachability` | ERROR | lines 108-155 (BFS) |
| `edge_target_exists` | ERROR | lines 159-201 |
| `start_no_incoming` | ERROR | lines 205-233 |
| `exit_no_outgoing` | ERROR | lines 237-272 |
| `condition_syntax` | ERROR | lines 276-316 |
| `stylesheet_syntax` | ERROR | lines 320-348 |
| `type_known` | WARNING | lines 352-392 |
| `fidelity_valid` | WARNING | lines 396-471 |
| `retry_target_exists` | WARNING | lines 475-548 |
| `goal_gate_has_retry` | WARNING | lines 552-586 |
| `prompt_on_llm_nodes` | WARNING | lines 590-624 |
| `freeform_edge_count` | ERROR | lines 628-663 |
**Minor:** `stylesheet_syntax` only checks brace balance, not full parse.
### 50. Section 7.3 — Validation API
**ALIGNED**
`validate(graph, extra_rules)` at line 35, `validate_or_raise(graph, extra_rules)` at line 52.
### 51. Section 7.4 — Custom Lint Rules
**ALIGNED**
`LintRule` trait at `mod.rs:28-31` with `name()` and `apply()`. Custom rules via `extra_rules` parameter.
---
## Section 8: Model Stylesheet
### 52. Section 8.1 — CSS-like Syntax
**ALIGNED**
`parse_stylesheet` at `stylesheet.rs:54-85` parses selector blocks with property declarations.
### 53. Section 8.2 — Selectors and Specificity
**ALIGNED** (minor extension)
Implementation at `stylesheet.rs:18-27` adds a `Shape` selector beyond spec:
| Selector | Specificity |
|----------|-------------|
| `*` (Universal) | 0 |
| Shape (bare word) | 1 |
| `.class` | 2 |
| `#id` | 3 |
Spec defines `*`=0, `.class`=1, `#id`=2. Relative ordering preserved; the `Shape` selector is an undocumented extension. Cascading behavior correct.
### 54. Section 8.3 — Application Order
**ALIGNED**
At `stylesheet.rs:190-238`: sorts by specificity, higher overwrites lower, explicit node attributes always override. Test at `stylesheet.rs:395-442` verifies spec section 8.6 example exactly.
### 55. Section 8.4 — Recognized Properties
**ALIGNED**
`STYLESHEET_PROPERTIES` at line 182: `["llm_model", "llm_provider", "reasoning_effort"]`. Exact match.
---
## Section 9: Transforms and Extensibility
### 56. Section 9.1 — Transform Trait
**ALIGNED**
At `transform.rs:5-7`: `fn apply(&self, graph: &mut Graph)`. In-place mutation vs spec's return-new-graph — functionally equivalent.
### 57. Section 9.2 — Built-In Transforms
**ALIGNED**
Three built-in transforms:
- **Variable Expansion:** `transform.rs:10-25` — expands `$goal` in prompts
- **Stylesheet Application:** `transform.rs:52-65` — applies `model_stylesheet`
- **Preamble:** `transform.rs:28-49` — prepends `[Context mode: {fidelity}]` for non-full fidelity
- **Minor:** Preamble applied at parse time, not execution time; cannot incorporate runtime fidelity changes from edges
### 58. Section 9.3 — Custom Transforms
**ALIGNED**
`PipelineBuilder::register_transform()` at `pipeline.rs:24-26`. Custom transforms run after built-in, in registration order. Integration test at `pipeline.rs:149-169`.
### 59. Section 9.4 — Event Stream
**ALIGNED**
All spec event types implemented at `event.rs:5-74`:
- Pipeline lifecycle: `PipelineStarted`, `PipelineCompleted`, `PipelineFailed`
- Stage lifecycle: `StageStarted`, `StageCompleted`, `StageFailed`, `StageRetrying`
- Parallel: `ParallelStarted`, `ParallelBranchStarted`, `ParallelBranchCompleted`, `ParallelCompleted`
- Human: `InterviewStarted`, `InterviewCompleted`, `InterviewTimeout`
- Checkpoint: `CheckpointSaved`
Observer pattern via `EventEmitter::on_event()` at line 106. Engine emits throughout execution.
### 60. Section 9.5 — Tool Call Hooks
**ALIGNED** (minor discrepancy)
Pre/post hooks at `codergen.rs:98-131`. `resolve_hook()` at lines 56-62 checks node-level then graph-level.
- **Minor:** Pre-hook non-zero returns `Outcome::fail()` (stronger than spec's "skip the tool call")
### 61. Section 9.6 — HTTP Server Mode
**N/A** — Spec says "Implementations may expose..." (optional). Not implemented.
---
## Section 10: Condition Expression Language
### 62. Section 10.1 — Grammar
**ALIGNED**
At `condition.rs`: `&&` conjunction (line 29), `!=` (lines 33-45), `=` (lines 46-58), bare key truthy (lines 59-72).
### 63. Section 10.2 — Semantics
**ALIGNED**
- Clauses AND-combined: `condition.rs:134` uses `.all()`
- `outcome` resolves to status string: line 88-89
- `preferred_label` resolves: lines 91-96
- `context.*` lookup with fallback: lines 98-105
- Missing keys = empty string: line 105
- Empty condition = true: lines 130-132
### 64. Section 10.3 — Variable Resolution
**ALIGNED**
`resolve_key()` at lines 87-110 follows spec pseudocode exactly: `outcome` → `preferred_label` → `context.` prefix with qualified/unqualified fallback → direct context lookup → empty string.
### 65. Section 10.4 — Examples
**ALIGNED**
Tests cover all spec examples: `outcome=success` (line 169), `context.tests_passed=true` (line 206), `preferred_label=Fix` (line 189).
### 66. Section 10.5 — Extended Operators
**ALIGNED**
Correctly NOT implemented per spec: "documented as potential extensions... Implementations should not add them."
---
## Section 11: Definition of Done
### 67. Section 11.1 — DOT Parsing
**ALIGNED**
Integration tests parse all 3 spec examples at `integration.rs:30-143`.
### 68. Section 11.2 — Validation and Linting
**ALIGNED**
14 lint rules, `validate_or_raise()` used in integration tests.
### 69. Section 11.3 — Execution Engine
**ALIGNED**
Start node resolution, handler dispatch, outcome recording, edge selection, loop execution, terminal stop — all verified in integration tests at `integration.rs:158-203`.
### 70. Section 11.4 — Goal Gate Enforcement
**ALIGNED**
Integration tests at `integration.rs:432-608`.
### 71. Section 11.5 — Retry Logic
**ALIGNED**
Integration test at `integration.rs:828-902`.
### 72. Section 11.6 — Node Handlers
**ALIGNED**
All handler types exist. Custom handler registration works.
### 73. Section 11.7 — State and Context
**ALIGNED**
Context updates, checkpoint save/resume, artifacts — verified at `integration.rs:978-1016`.
### 74. Section 11.8 — Human-in-the-Loop
**ALIGNED**
All interviewer implementations present. Integration test with QueueInterviewer at `integration.rs:376-381`.
### 75. Section 11.9 — Condition Expressions
**ALIGNED**
All operators and variable types tested at `condition.rs:160-323`.
### 76. Section 11.10 — Model Stylesheet
**ALIGNED**
Integration tests at `integration.rs:677-822` verify selectors, specificity, cascading.
### 77. Section 11.11 — Transforms
**ALIGNED**
Transform interface, variable expansion, custom transforms — all tested.
### 78. Section 11.12 — Cross-Feature Parity Matrix
**ALIGNED** (minor gap)
22 of 23 matrix items pass. **Minor:** No dedicated integration test for parallel fan-out/fan-in (handlers exist, no end-to-end test).
### 79. Section 11.13 — Integration Smoke Test
**ALIGNED**
`integration.rs:1040-1180` implements mock-backend smoke test matching spec pattern.
---
## Appendices
### 80. Appendix A — Complete Attribute Reference
**ALIGNED**
All graph, node, and edge attributes have corresponding accessors. Dotted keys (`tool_hooks.pre/post`, `stack.*`) supported via `lexer.rs:314-315`.
### 81. Appendix B — Shape-to-Handler-Type Mapping
**ALIGNED**
All 9 mappings tested at `graph/types.rs:413-429`.
### 82. Appendix C — Status File Contract
**GAP**
`auto_status=true` synthesis not implemented in engine (see item 23). `Outcome` struct at `outcome.rs:48-60` matches contract fields, but the engine never checks `auto_status`.
### 83. Appendix D — Error Categories
**ALIGNED** (minor gap)
`AttractorError` at `error.rs:4-25` has 7 variants: `Parse`, `Validation`, `Engine`, `Handler`, `Checkpoint`, `Stylesheet`, `Io`.
**Minor:** Spec defines 3 abstract categories (Retryable, Terminal, Pipeline). No explicit classification of which variants are retryable vs terminal; default `should_retry` retries all.
---
## Summary
| # | Section | Verdict |
|---|---------|---------|
| 1 | 1.1 Problem Statement | ALIGNED |
| 2 | 1.2 Why DOT Syntax | ALIGNED |
| 3 | 1.3 Design Principles | ALIGNED |
| 4 | 1.4 Layering / LLM Backends | ALIGNED |
| 5 | 2.1 Supported Subset | ALIGNED |
| 6 | 2.2 BNF Grammar | ALIGNED (minor: Direction not validated) |
| 7 | 2.3 Key Constraints | ALIGNED |
| 8 | 2.4 Value Types | ALIGNED |
| 9 | 2.5 Graph-Level Attributes | ALIGNED |
| 10 | 2.6 Node Attributes | ALIGNED |
| 11 | 2.7 Edge Attributes | ALIGNED |
| 12 | 2.8 Shape-to-Handler Mapping | ALIGNED |
| 13 | 2.9 Chained Edges | ALIGNED |
| 14 | 2.10-2.12 Subgraphs/Defaults/Class | ALIGNED |
| 15 | 3.1 Run Lifecycle | ALIGNED |
| 16 | 3.2 Core Execution Loop | ALIGNED |
| 17 | 3.3 Edge Selection Algorithm | ALIGNED |
| 18 | 3.4 Goal Gate Enforcement | ALIGNED |
| 19 | 3.5 Retry Logic | ALIGNED (minor: counters not persisted) |
| 20 | 3.6 Retry Policy / Backoff | ALIGNED (minor: should_retry too coarse) |
| 21 | 3.7 Failure Routing | ALIGNED |
| 22 | 3.8 Concurrency Model | ALIGNED |
| 23 | 3.9 Auto Status | **GAP** |
| 24 | 3.10 Timeout Enforcement | **GAP** |
| 25 | 4.1 Handler Trait | ALIGNED |
| 26 | 4.2 Handler Registry | ALIGNED |
| 27 | 4.3 Start Handler | ALIGNED |
| 28 | 4.4 Exit Handler | ALIGNED |
| 29 | 4.5 Codergen Handler | ALIGNED |
| 30 | 4.6 Wait Human Handler | ALIGNED |
| 31 | 4.7 Conditional Handler | ALIGNED |
| 32 | 4.8 Parallel Handler | ALIGNED |
| 33 | 4.9 Fan-In Handler | ALIGNED (minor: no score sort, all-fail case) |
| 34 | 4.10 Tool Handler | ALIGNED |
| 35 | 4.11 Manager Loop Handler | **GAP** (child_autostart missing) |
| 36 | 4.12 Custom Handlers | ALIGNED |
| 37 | 5.1 Context | ALIGNED |
| 38 | 5.2 Outcome | ALIGNED |
| 39 | 5.3 Checkpoint | ALIGNED (minor: retry counters, fidelity degradation on resume) |
| 40 | 5.4 Fidelity Modes | ALIGNED (minor: thread_id resolution incomplete) |
| 41 | 5.5 Artifact Store | ALIGNED |
| 42 | 5.6 Run Directory | ALIGNED (minor: manifest missing goal) |
| 43 | 6.1 Interviewer Interface | ALIGNED |
| 44 | 6.2 Question Model | ALIGNED |
| 45 | 6.3 Answer Model | ALIGNED |
| 46 | 6.4 Built-In Interviewers | ALIGNED |
| 47 | 6.5 Timeout Handling | ALIGNED |
| 48 | 7.1 Diagnostic Model | ALIGNED |
| 49 | 7.2 Built-In Lint Rules | ALIGNED (14/14) |
| 50 | 7.3 Validation API | ALIGNED |
| 51 | 7.4 Custom Lint Rules | ALIGNED |
| 52 | 8.1 CSS-like Syntax | ALIGNED |
| 53 | 8.2 Selectors/Specificity | ALIGNED (extra Shape selector) |
| 54 | 8.3 Application Order | ALIGNED |
| 55 | 8.4 Recognized Properties | ALIGNED |
| 56 | 9.1 Transform Trait | ALIGNED |
| 57 | 9.2 Built-In Transforms | ALIGNED |
| 58 | 9.3 Custom Transforms | ALIGNED |
| 59 | 9.4 Event Stream | ALIGNED |
| 60 | 9.5 Tool Call Hooks | ALIGNED (minor: pre-hook behavior) |
| 61 | 9.6 HTTP Server Mode | N/A (optional) |
| 62 | 10.1 Grammar | ALIGNED |
| 63 | 10.2 Semantics | ALIGNED |
| 64 | 10.3 Variable Resolution | ALIGNED |
| 65 | 10.4 Examples | ALIGNED |
| 66 | 10.5 Extended Operators | ALIGNED |
| 67 | 11.1 DOT Parsing | ALIGNED |
| 68 | 11.2 Validation | ALIGNED |
| 69 | 11.3 Execution Engine | ALIGNED |
| 70 | 11.4 Goal Gates | ALIGNED |
| 71 | 11.5 Retry Logic | ALIGNED |
| 72 | 11.6 Node Handlers | ALIGNED |
| 73 | 11.7 State/Context | ALIGNED |
| 74 | 11.8 Human-in-the-Loop | ALIGNED |
| 75 | 11.9 Conditions | ALIGNED |
| 76 | 11.10 Stylesheet | ALIGNED |
| 77 | 11.11 Transforms | ALIGNED |
| 78 | 11.12 Parity Matrix | ALIGNED (minor: no parallel integration test) |
| 79 | 11.13 Smoke Test | ALIGNED |
| 80 | Appendix A — Attributes | ALIGNED |
| 81 | Appendix B — Shape Mapping | ALIGNED |
| 82 | Appendix C — Status File | **GAP** (auto_status not enforced) |
| 83 | Appendix D — Error Categories | ALIGNED (minor: no retryable/terminal classification) |
---
## Totals
**79 ALIGNED / 3 GAP / 1 N/A** (+ 16 minor gaps within ALIGNED items)
### Hard Gaps (3)
1. **Auto Status (item 23/82):** `auto_status` accessor exists at `graph/types.rs:193-194` but engine never checks it. No auto-synthesis of SUCCESS when handler writes no status.
2. **Timeout Enforcement (item 24):** `timeout` attribute parsed as `Duration` but not enforced as a deadline around handler execution in the engine. Tool handler uses it for subprocess timeout, but no general enforcement.
3. **Manager Loop `child_autostart` (item 35):** `stack.child_autostart` / `start_child_pipeline` not implemented. The observe/steer/wait cycle exists but cannot auto-launch a child pipeline.
### Notable Minor Gaps (within ALIGNED items)
- Thread ID resolution: only step 1 of 5 implemented (node `thread_id`); missing edge, graph default, subgraph class, previous-node fallback
- Checkpoint resume: retry counters not restored; fidelity degradation (`full` → `summary:high`) not applied
- Retry policy: `should_retry` retries ALL errors; spec defines granular HTTP-status-based behavior
- Stylesheet: adds undocumented `Shape` selector (functional, shifts specificity values)
- Preamble transform: applied at parse time, not execution time
- No dedicated parallel fan-out/fan-in integration test

View file

@ -478,7 +478,7 @@ Each node has a retry policy determined by:
1. Node attribute `max_retries` (if set) -- number of additional attempts beyond the initial execution
2. Graph attribute `default_max_retry` (fallback)
3. Built-in default: 0 (no retries)
3. Built-in default: 50
The `max_retries` attribute specifies additional attempts. So `max_retries=3` means a total of 4 executions (1 initial + 3 retries). Internally this maps to `max_attempts = max_retries + 1`.