diff --git a/Cargo.lock b/Cargo.lock index 098d71cbf..6a4a8069d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2659,6 +2659,7 @@ dependencies = [ "nom", "regex", "serde", + "serde_json", "strum 0.28.0", "thiserror 2.0.18", ] diff --git a/lib/crates/fabro-graphviz/Cargo.toml b/lib/crates/fabro-graphviz/Cargo.toml index 98e017097..3b33151d8 100644 --- a/lib/crates/fabro-graphviz/Cargo.toml +++ b/lib/crates/fabro-graphviz/Cargo.toml @@ -21,3 +21,6 @@ regex = { workspace = true } serde = { workspace = true } strum.workspace = true thiserror = { workspace = true } + +[dev-dependencies] +serde_json = { workspace = true } diff --git a/lib/crates/fabro-graphviz/src/fidelity.rs b/lib/crates/fabro-graphviz/src/fidelity.rs index 8b33c98e8..98605c135 100644 --- a/lib/crates/fabro-graphviz/src/fidelity.rs +++ b/lib/crates/fabro-graphviz/src/fidelity.rs @@ -1,9 +1,24 @@ +use serde::{Deserialize, Serialize}; use strum::{Display, EnumString, VariantArray}; /// Fidelity mode controlling how much prior context is provided to LLM /// sessions. -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Display, EnumString, VariantArray)] +#[derive( + Debug, + Clone, + Copy, + Default, + PartialEq, + Eq, + Hash, + Display, + EnumString, + VariantArray, + Serialize, + Deserialize, +)] #[strum(serialize_all = "lowercase")] +#[serde(rename_all = "lowercase")] pub enum Fidelity { /// Complete context, no summarization — sessions share a thread. Full, @@ -14,12 +29,15 @@ pub enum Fidelity { Compact, /// Brief textual summary (~600 token target). #[strum(serialize = "summary:low")] + #[serde(rename = "summary:low")] SummaryLow, /// Moderate textual summary (~1500 token target). #[strum(serialize = "summary:medium")] + #[serde(rename = "summary:medium")] SummaryMedium, /// Detailed per-stage Markdown report. #[strum(serialize = "summary:high")] + #[serde(rename = "summary:high")] SummaryHigh, } @@ -73,4 +91,14 @@ mod tests { fn fidelity_unknown_mode_errors() { assert!("bogus".parse::().is_err()); } + + #[test] + fn fidelity_serde_matches_strum_display() { + for mode in Fidelity::variants() { + let json = serde_json::to_value(mode).unwrap(); + assert_eq!(json, serde_json::Value::String(mode.to_string())); + let parsed: Fidelity = serde_json::from_value(json).unwrap(); + assert_eq!(parsed, *mode); + } + } } diff --git a/lib/crates/fabro-validate/src/rules/parallel_branch.rs b/lib/crates/fabro-validate/src/rules/parallel_branch.rs index 2f1cba2e7..d8a7d2fbd 100644 --- a/lib/crates/fabro-validate/src/rules/parallel_branch.rs +++ b/lib/crates/fabro-validate/src/rules/parallel_branch.rs @@ -35,22 +35,25 @@ impl<'a> ParallelBranches<'a> { .collect() } + /// True when every incoming edge of `node_id` comes from a parallel fork + /// (and there is at least one). Such a node only ever runs as a branch. pub(super) fn is_branch_only_node(&self, node_id: &str) -> bool { - self.branch_only_parents(node_id).is_some() + let incoming = self.graph.incoming_edges(node_id); + !incoming.is_empty() && incoming.iter().all(|edge| self.is_fork_edge(edge)) } + /// The sorted, deduplicated fork parents of a branch-only node, or `None` + /// when the node has a non-fork entry path (or no entry at all). pub(super) fn branch_only_parents(&self, node_id: &str) -> Option> { - let mut incoming = self + if !self.is_branch_only_node(node_id) { + return None; + } + let parents: BTreeSet<&str> = self .graph - .edges - .iter() - .filter(|edge| edge.to == node_id) - .peekable(); - incoming.peek()?; - - incoming - .map(|edge| self.is_fork_edge(edge).then(|| edge.from.clone())) - .collect::>>() - .map(|parents| parents.into_iter().collect()) + .incoming_edges(node_id) + .into_iter() + .map(|edge| edge.from.as_str()) + .collect(); + Some(parents.into_iter().map(String::from).collect()) } } diff --git a/lib/crates/fabro-validate/src/rules/parallel_branch_inert_attribute.rs b/lib/crates/fabro-validate/src/rules/parallel_branch_inert_attribute.rs index da36dbbea..bac461103 100644 --- a/lib/crates/fabro-validate/src/rules/parallel_branch_inert_attribute.rs +++ b/lib/crates/fabro-validate/src/rules/parallel_branch_inert_attribute.rs @@ -7,11 +7,15 @@ pub(super) fn rule() -> Box { Box::new(Rule) } -/// Attributes that parallel branch execution does not resolve. Per-branch -/// preambles now honor fidelity, while `thread_id` remains inert because -/// concurrent branches cannot share an LLM session. +/// Attributes that parallel branch execution does not resolve. Only +/// `thread_id` is inert on branches (concurrent branches cannot share an LLM +/// session); per-branch `fidelity` is honored via pre-rendered preambles. const BRANCH_IGNORED_ATTRS: &[&str] = &["thread_id"]; +const FULL_FIDELITY_MESSAGE: &str = "Parallel branches run at most at summary:high; full is degraded at runtime because branches cannot share a session"; + +const THREAD_ID_FIX: &str = "Remove 'thread_id': parallel branches inherit the thread resolved when the parallel node started"; + struct Rule; /// Renders one or more parallel-node ids as `'a'` or `'a', 'b'`. @@ -22,19 +26,6 @@ fn quoted_list(ids: &[String]) -> String { .join(", ") } -fn fix_message(attr: &str) -> String { - match attr { - "thread_id" => format!( - "Remove '{attr}': parallel branches inherit the thread resolved when the parallel node started" - ), - _ => format!("Remove '{attr}'"), - } -} - -fn full_fidelity_message() -> String { - "parallel branches run at most at summary:high; full is degraded at runtime because branches cannot share a session".to_string() -} - fn full_fidelity_fix(parallel_ids: &[String]) -> String { let parent = if parallel_ids.len() == 1 { format!("parallel node {}", quoted_list(parallel_ids)) @@ -46,6 +37,23 @@ fn full_fidelity_fix(parallel_ids: &[String]) -> String { ) } +fn full_fidelity_diagnostic( + rule_name: &str, + node_id: Option, + edge: Option<(String, String)>, + parallel_ids: &[String], +) -> Diagnostic { + Diagnostic { + rule: rule_name.to_string(), + severity: Severity::Warning, + message: FULL_FIDELITY_MESSAGE.to_string(), + node_id, + edge, + fix: Some(full_fidelity_fix(parallel_ids)), + ..Diagnostic::default() + } +} + impl LintRule for Rule { fn name(&self) -> &'static str { "parallel_branch_inert_attribute" @@ -64,15 +72,12 @@ impl LintRule for Rule { continue; } if edge.fidelity() == Some("full") { - diagnostics.push(Diagnostic { - rule: self.name().to_string(), - severity: Severity::Warning, - message: full_fidelity_message(), - node_id: None, - edge: Some((edge.from.clone(), edge.to.clone())), - fix: Some(full_fidelity_fix(std::slice::from_ref(&edge.from))), - ..Diagnostic::default() - }); + diagnostics.push(full_fidelity_diagnostic( + self.name(), + None, + Some((edge.from.clone(), edge.to.clone())), + std::slice::from_ref(&edge.from), + )); } for attr in BRANCH_IGNORED_ATTRS { if !edge.attrs.contains_key(*attr) { @@ -87,7 +92,7 @@ impl LintRule for Rule { ), node_id: None, edge: Some((edge.from.clone(), edge.to.clone())), - fix: Some(fix_message(attr)), + fix: Some(THREAD_ID_FIX.to_string()), ..Diagnostic::default() }); } @@ -103,15 +108,12 @@ impl LintRule for Rule { continue; }; if node.fidelity() == Some("full") { - diagnostics.push(Diagnostic { - rule: self.name().to_string(), - severity: Severity::Warning, - message: full_fidelity_message(), - node_id: Some(node.id.clone()), - edge: None, - fix: Some(full_fidelity_fix(&parents)), - ..Diagnostic::default() - }); + diagnostics.push(full_fidelity_diagnostic( + self.name(), + Some(node.id.clone()), + None, + &parents, + )); } for attr in BRANCH_IGNORED_ATTRS { if !node.attrs.contains_key(*attr) { @@ -127,7 +129,7 @@ impl LintRule for Rule { ), node_id: Some(node.id.clone()), edge: None, - fix: Some(fix_message(attr)), + fix: Some(THREAD_ID_FIX.to_string()), ..Diagnostic::default() }); } diff --git a/lib/crates/fabro-validate/src/rules/thread_id_requires_fidelity_full.rs b/lib/crates/fabro-validate/src/rules/thread_id_requires_fidelity_full.rs index a73335da5..157b2a032 100644 --- a/lib/crates/fabro-validate/src/rules/thread_id_requires_fidelity_full.rs +++ b/lib/crates/fabro-validate/src/rules/thread_id_requires_fidelity_full.rs @@ -23,6 +23,9 @@ impl LintRule for Rule { let graph_default_full = graph.default_fidelity() == Some("full"); let branches = ParallelBranches::new(graph); + // thread_id is inert on parallel branches, where + // parallel_branch_inert_attribute already says "remove thread_id" — + // advising fidelity="full" there would contradict it. for node in graph.nodes.values() { if node.thread_id().is_some() && !branches.is_branch_only_node(&node.id) diff --git a/lib/crates/fabro-workflow/src/artifact.rs b/lib/crates/fabro-workflow/src/artifact.rs index 975fee26d..bbf7f7ed8 100644 --- a/lib/crates/fabro-workflow/src/artifact.rs +++ b/lib/crates/fabro-workflow/src/artifact.rs @@ -78,12 +78,18 @@ pub fn format_artifact_reference(path: &str) -> String { pub fn durable_context_snapshot(context: &Context) -> HashMap { let mut snapshot = context.snapshot(); - snapshot.remove(context::keys::CURRENT_PREAMBLE); - snapshot.remove(context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES); + strip_transient_keys(&mut snapshot); normalize_durable_updates(&mut snapshot); snapshot } +/// Remove runtime-only keys that must never reach durable storage or events. +pub fn strip_transient_keys(values: &mut HashMap) { + for key in context::keys::TRANSIENT_CONTEXT_KEYS { + values.remove(*key); + } +} + pub fn normalize_durable_updates(updates: &mut HashMap) { for value in updates.values_mut() { normalize_durable_value(value); @@ -97,12 +103,7 @@ pub fn normalize_durable_outcomes(node_outcomes: &mut HashMap) } pub fn normalize_checkpoint_for_resume(checkpoint: &mut Checkpoint) { - checkpoint - .context_values - .remove(context::keys::CURRENT_PREAMBLE); - checkpoint - .context_values - .remove(context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES); + strip_transient_keys(&mut checkpoint.context_values); normalize_durable_updates(&mut checkpoint.context_values); normalize_durable_outcomes(&mut checkpoint.node_outcomes); } diff --git a/lib/crates/fabro-workflow/src/context.rs b/lib/crates/fabro-workflow/src/context.rs index b911e2210..a890c0e8f 100644 --- a/lib/crates/fabro-workflow/src/context.rs +++ b/lib/crates/fabro-workflow/src/context.rs @@ -25,6 +25,9 @@ pub mod keys { pub const INTERNAL_PARENT_PREAMBLE: &str = "internal.parent_preamble"; pub const INTERNAL_PARALLEL_GROUP_ID: &str = "internal.parallel_group_id"; pub const INTERNAL_PARALLEL_BRANCH_ID: &str = "internal.parallel_branch_id"; + /// Stash of pre-rendered per-branch preambles for a parallel node; see + /// [`super::ParallelBranchPreamble`] for the entry shape and the + /// producer/consumer contract. pub const INTERNAL_PARALLEL_BRANCH_PREAMBLES: &str = "internal.parallel_branch_preambles"; // --- current.* keys --- @@ -45,6 +48,12 @@ pub mod keys { pub const PARALLEL_FAN_IN_BEST_OUTCOME: &str = "parallel.fan_in.best_outcome"; pub const PARALLEL_FAN_IN_BEST_HEAD_SHA: &str = "parallel.fan_in.best_head_sha"; + /// Runtime-only keys stripped from durable context projections + /// (checkpoint snapshots and resume normalization). Add new transient + /// keys here so both strip sites stay in sync. + pub const TRANSIENT_CONTEXT_KEYS: &[&str] = + &[CURRENT_PREAMBLE, INTERNAL_PARALLEL_BRANCH_PREAMBLES]; + // --- Prefix constants (for filtering and dynamic keys) --- pub const GRAPH_PREFIX: &str = "graph."; pub const INTERNAL_PREFIX: &str = "internal."; @@ -136,9 +145,24 @@ pub mod keys { pub use fabro_core::Context; use fabro_graphviz::Fidelity; use fabro_types::{ParallelBranchId, StageId}; +use serde::{Deserialize, Serialize}; use crate::event::StageScope; +/// One entry of the [`keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES`] stash. +/// +/// The stash is a JSON array indexed by the parallel node's outgoing-edge +/// order (`Graph::outgoing_edges` preserves declaration order, so producer and +/// consumer align even with duplicate targets). `null` entries mean the branch +/// inherits the fork's preamble. `FidelityLifecycle::before_node` produces the +/// stash; `ParallelHandler::execute` consumes and clears it. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ParallelBranchPreamble { + pub fidelity: Fidelity, + pub preamble: String, +} + /// Domain-specific typed accessors for workflow context values. pub trait WorkflowContext { fn fidelity(&self) -> Fidelity; diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index ac4e4f3e1..32e4a8b10 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -4,14 +4,13 @@ use std::time::Instant; use async_trait::async_trait; use fabro_agent::{Sandbox, WorktreeOptions, WorktreeSandbox}; -use fabro_graphviz::Fidelity; use fabro_graphviz::graph::{AttrValue, Graph, Node}; use fabro_hooks::{HookContext, HookEvent}; use fabro_types::{ParallelBranchId, RunId, StageId}; use tokio::sync::Semaphore; use super::{EngineServices, Handler}; -use crate::context::{Context, WorkflowContext, keys}; +use crate::context::{Context, ParallelBranchPreamble, WorkflowContext, keys}; use crate::error::Error; use crate::event::{Event, RunNoticeCode, RunNoticeLevel, StageScope}; use crate::git::sanitize_ref_component; @@ -57,15 +56,15 @@ struct BranchResult { worktree_path: Option, } -struct BranchPreamble { - fidelity: Fidelity, - preamble: String, -} - +/// Parse the per-branch preamble stash produced by `FidelityLifecycle`. +/// +/// Outer `None` means the stash is absent, malformed, or has the wrong branch +/// count — every branch then inherits the fork context (legacy behavior). +/// Inner `None` means that single branch inherits. fn parse_branch_preambles( value: Option, branch_count: usize, -) -> Option>> { +) -> Option>> { let serde_json::Value::Array(entries) = value? else { return None; }; @@ -77,15 +76,7 @@ fn parse_branch_preambles( .into_iter() .map(|entry| match entry { serde_json::Value::Null => Some(None), - serde_json::Value::Object(entry) if entry.len() == 2 => { - let fidelity = entry.get("fidelity")?.as_str()?.parse().ok()?; - let preamble = entry.get("preamble")?.as_str()?; - Some(Some(BranchPreamble { - fidelity, - preamble: preamble.to_string(), - })) - } - _ => None, + entry => serde_json::from_value(entry).ok().map(Some), }) .collect() } @@ -258,6 +249,13 @@ impl Handler for ParallelHandler { context.get(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES), branches.len(), ); + // Clear the stash before forking so branch contexts never carry the + // outer array — a nested parallel branch target must not misread it as + // its own. The write-back diff also clears it on the run state. + context.set( + keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES, + serde_json::Value::Null, + ); let mut branch_setups: Vec = Vec::new(); for (branch_index, edge) in branches.iter().enumerate() { let target_id = edge.to.clone(); @@ -279,16 +277,15 @@ impl Handler for ParallelHandler { .and_then(|entries| entries.get(branch_index)) .and_then(Option::as_ref) { - branch_context.set(keys::CURRENT_PREAMBLE, serde_json::json!(&entry.preamble)); + branch_context.set( + keys::CURRENT_PREAMBLE, + serde_json::Value::String(entry.preamble.clone()), + ); branch_context.set( keys::INTERNAL_FIDELITY, - serde_json::json!(entry.fidelity.to_string()), + serde_json::Value::String(entry.fidelity.to_string()), ); } - branch_context.set( - keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES, - serde_json::Value::Null, - ); let (branch_sandbox, worktree_path): (Arc, Option) = if let ( Some(ref gs), @@ -350,10 +347,6 @@ impl Handler for ParallelHandler { worktree_path, }); } - context.set( - keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES, - serde_json::Value::Null, - ); // --- Fan out: concurrent execution --- let mut handles = Vec::new(); diff --git a/lib/crates/fabro-workflow/src/lifecycle/event.rs b/lib/crates/fabro-workflow/src/lifecycle/event.rs index 3965d7261..f547f5fe5 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/event.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/event.rs @@ -92,6 +92,10 @@ fn response_from_outcome(node_id: &str, outcome: &Outcome) -> Option { .and_then(|value| value.as_str().map(ToOwned::to_owned)) } +/// Context values for `StageCompleted` events. Unlike +/// `artifact::strip_transient_keys`, this keeps `CURRENT_PREAMBLE` — stage +/// events have always included the active preamble — and drops only the +/// parallel stash, which can embed every branch's rendered preamble. fn stage_context_values(workflow_context: &Context) -> Option> { let mut snapshot = workflow_context.snapshot(); snapshot.remove(context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES); diff --git a/lib/crates/fabro-workflow/src/lifecycle/fidelity.rs b/lib/crates/fabro-workflow/src/lifecycle/fidelity.rs index 3f0b24228..437550cc1 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/fidelity.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/fidelity.rs @@ -1,3 +1,4 @@ +use std::collections::HashMap; use std::path::PathBuf; use std::sync::{Arc, Mutex}; @@ -10,10 +11,10 @@ use fabro_core::state::ExecutionState; use fabro_graphviz::graph::types::{Edge as GvEdge, Graph as GvGraph, Node as GvNode}; use crate::artifact; -use crate::context::keys; +use crate::context::{Context, ParallelBranchPreamble, keys}; use crate::graph::{WorkflowGraph, WorkflowNode}; use crate::handler::llm::preamble::build_preamble; -use crate::outcome::BilledModelUsage; +use crate::outcome::{BilledModelUsage, Outcome}; use crate::runtime_store::RunStoreHandle; type WfRunState = ExecutionState>; @@ -61,6 +62,60 @@ impl FidelityLifecycle { "fidelity mutex should not be poisoned: no code panics while holding this lock", ) = flag; } + + /// Render the per-branch preamble stash for a parallel node, indexed by + /// outgoing-edge order (the same order `ParallelHandler` fans out in). + /// `Null` entries inherit the fork's preamble. + fn build_parallel_branch_preambles( + &self, + node_id: &str, + fork_fidelity: keys::Fidelity, + resolved_context: &Context, + resolved_outcomes: &HashMap, + completed_nodes: &[String], + ) -> Vec { + let mut rendered: HashMap = HashMap::new(); + self.graph + .outgoing_edges(node_id) + .into_iter() + .enumerate() + .map(|(branch_index, edge)| { + let Some(target_node) = self.graph.nodes.get(&edge.to) else { + return serde_json::Value::Null; + }; + let resolution = resolve_parallel_branch_fidelity(edge, target_node, fork_fidelity); + if resolution.requested == Some(keys::Fidelity::Full) { + tracing::warn!( + parallel_node = %node_id, + branch = %edge.to, + branch_index, + effective_fidelity = %keys::Fidelity::Full.degraded(), + "Parallel branch fidelity degraded from full" + ); + } + let Some(branch_fidelity) = resolution.effective else { + return serde_json::Value::Null; + }; + rendered + .entry(branch_fidelity) + .or_insert_with(|| { + let entry = ParallelBranchPreamble { + fidelity: branch_fidelity, + preamble: build_preamble( + branch_fidelity, + resolved_context, + &self.graph, + completed_nodes, + resolved_outcomes, + ), + }; + serde_json::to_value(entry) + .expect("ParallelBranchPreamble serialization cannot fail") + }) + .clone() + }) + .collect() + } } #[async_trait] @@ -143,47 +198,23 @@ impl RunLifecycle for FidelityLifecycle { .context .set(keys::CURRENT_PREAMBLE, serde_json::json!(preamble)); + // 5. Parallel nodes: pre-render per-branch preambles into the stash that + // ParallelHandler consumes at fan-out. if gv_node.handler_type() == Some("parallel") { - let mut branch_preambles = Vec::new(); - for (branch_index, edge) in self.graph.outgoing_edges(node.id()).iter().enumerate() { - let Some(target_node) = self.graph.nodes.get(&edge.to) else { - branch_preambles.push(serde_json::Value::Null); - continue; - }; - let resolution = resolve_parallel_branch_fidelity(edge, target_node, fidelity); - if resolution.requested() == Some(keys::Fidelity::Full) { - tracing::warn!( - parallel_node = %node.id(), - branch = %edge.to, - branch_index, - fidelity = %keys::Fidelity::Full, - effective_fidelity = %keys::Fidelity::SummaryHigh, - "Parallel branch fidelity degraded" - ); - } - let Some(branch_fidelity) = resolution.entry_fidelity() else { - branch_preambles.push(serde_json::Value::Null); - continue; - }; - let branch_preamble = build_preamble( - branch_fidelity, - &resolved_context, - &self.graph, - &state.completed_nodes, - &resolved_outcomes, - ); - branch_preambles.push(serde_json::json!({ - "fidelity": branch_fidelity.to_string(), - "preamble": branch_preamble, - })); - } + let branch_preambles = self.build_parallel_branch_preambles( + node.id(), + fidelity, + &resolved_context, + &resolved_outcomes, + &state.completed_nodes, + ); state.context.set( keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES, serde_json::Value::Array(branch_preambles), ); } - // 5. Thread ID resolution via resolve_thread_id: edge → node → graph default → + // 6. Thread ID resolution via resolve_thread_id: edge → node → graph default → // class → previous let thread_id = resolve_thread_id( incoming_edge_ref, @@ -192,13 +223,13 @@ impl RunLifecycle for FidelityLifecycle { state.previous_node_id.as_deref(), ); - // 6. Set thread.{tid}.current_node + // 7. Set thread.{tid}.current_node if let Some(ref tid) = thread_id { let key = keys::thread_current_node_key(tid); state.context.set(key, serde_json::json!(node.id())); } - // 7. Set INTERNAL_THREAD_ID (or null) + // 8. Set INTERNAL_THREAD_ID (or null) match thread_id { Some(tid) => { state @@ -212,7 +243,7 @@ impl RunLifecycle for FidelityLifecycle { } } - // 8. Set INTERNAL_NODE_VISIT_COUNT and CURRENT_NODE + // 9. Set INTERNAL_NODE_VISIT_COUNT and CURRENT_NODE let visits = state.node_visits.get(node.id()).copied().unwrap_or(1); state .context @@ -243,22 +274,14 @@ impl RunLifecycle for FidelityLifecycle { } } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, Copy)] struct ParallelBranchFidelityResolution { + /// The explicit fidelity requested on the edge or node, pre-degradation. requested: Option, + /// The fidelity to render an entry for; `None` inherits the fork preamble. effective: Option, } -impl ParallelBranchFidelityResolution { - fn requested(self) -> Option { - self.requested - } - - fn entry_fidelity(self) -> Option { - self.effective - } -} - /// Resolve explicit branch fidelity with edge-over-node precedence. /// /// Branches with no explicit fidelity inherit the parallel node's preamble. @@ -270,10 +293,7 @@ fn resolve_parallel_branch_fidelity( target_node: &GvNode, parallel_fidelity: keys::Fidelity, ) -> ParallelBranchFidelityResolution { - let requested = edge - .fidelity() - .and_then(|value| value.parse().ok()) - .or_else(|| target_node.fidelity().and_then(|value| value.parse().ok())); + let requested = explicit_fidelity(Some(edge), target_node).map(|(fidelity, _)| fidelity); let effective = requested .map(keys::Fidelity::degraded) .filter(|fidelity| *fidelity != parallel_fidelity); @@ -284,6 +304,23 @@ fn resolve_parallel_branch_fidelity( } } +/// Explicit fidelity from the incoming edge attribute, else the node +/// attribute, with the winning source labeled for logging. +fn explicit_fidelity( + incoming_edge: Option<&GvEdge>, + node: &GvNode, +) -> Option<(keys::Fidelity, &'static str)> { + incoming_edge + .and_then(|e| e.fidelity()) + .and_then(|s| s.parse().ok()) + .map(|f| (f, "edge")) + .or_else(|| { + node.fidelity() + .and_then(|s| s.parse().ok()) + .map(|f| (f, "node")) + }) +} + /// Resolve the context fidelity for a node, following the precedence: /// 1. Incoming edge `fidelity` attribute /// 2. Target node `fidelity` attribute @@ -294,13 +331,8 @@ fn resolve_fidelity( node: &GvNode, graph: &GvGraph, ) -> keys::Fidelity { - let (resolved, source) = if let Some(f) = incoming_edge - .and_then(|e| e.fidelity()) - .and_then(|s| s.parse().ok()) - { - (f, "edge") - } else if let Some(f) = node.fidelity().and_then(|s| s.parse().ok()) { - (f, "node") + let (resolved, source) = if let Some((f, source)) = explicit_fidelity(incoming_edge, node) { + (f, source) } else if let Some(f) = graph.default_fidelity().and_then(|s| s.parse().ok()) { (f, "graph") } else { @@ -362,7 +394,7 @@ mod tests { use crate::context::WorkflowContext; use crate::context::keys::Fidelity; - fn fidelity_attr(value: &str) -> AttrValue { + fn str_attr(value: &str) -> AttrValue { AttrValue::String(value.to_string()) } @@ -374,23 +406,23 @@ mod tests { let mut start = Node::new("start"); start .attrs - .insert("shape".to_string(), fidelity_attr("Mdiamond")); + .insert("shape".to_string(), str_attr("Mdiamond")); let mut fork = Node::new("fork"); fork.attrs - .insert("shape".to_string(), fidelity_attr("component")); + .insert("shape".to_string(), str_attr("component")); if let Some(fidelity) = fork_fidelity { fork.attrs - .insert("fidelity".to_string(), fidelity_attr(fidelity)); + .insert("fidelity".to_string(), str_attr(fidelity)); } let mut branch_a = Node::new("branch_a"); if let Some(fidelity) = branch_a_fidelity { branch_a .attrs - .insert("fidelity".to_string(), fidelity_attr(fidelity)); + .insert("fidelity".to_string(), str_attr(fidelity)); } let branch_b = Node::new("branch_b"); let mut work = Node::new("work"); - work.attrs.insert("shape".to_string(), fidelity_attr("box")); + work.attrs.insert("shape".to_string(), str_attr("box")); graph.nodes.insert(start.id.clone(), start); graph.nodes.insert(fork.id.clone(), fork); @@ -426,15 +458,15 @@ mod tests { fn parallel_branch_fidelity_edge_overrides_node() { let mut node = Node::new("branch"); node.attrs - .insert("fidelity".to_string(), fidelity_attr("compact")); + .insert("fidelity".to_string(), str_attr("compact")); let mut edge = Edge::new("fork", "branch"); edge.attrs - .insert("fidelity".to_string(), fidelity_attr("truncate")); + .insert("fidelity".to_string(), str_attr("truncate")); let resolved = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::SummaryHigh); - assert_eq!(resolved.requested(), Some(Fidelity::Truncate)); - assert_eq!(resolved.entry_fidelity(), Some(Fidelity::Truncate)); + assert_eq!(resolved.requested, Some(Fidelity::Truncate)); + assert_eq!(resolved.effective, Some(Fidelity::Truncate)); } #[test] @@ -444,47 +476,45 @@ mod tests { let resolution = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::Compact); - assert_eq!(resolution.requested(), None); - assert_eq!(resolution.entry_fidelity(), None); + assert_eq!(resolution.requested, None); + assert_eq!(resolution.effective, None); } #[test] fn parallel_branch_full_fidelity_degrades_to_summary_high() { let mut node = Node::new("branch"); - node.attrs - .insert("fidelity".to_string(), fidelity_attr("full")); + node.attrs.insert("fidelity".to_string(), str_attr("full")); let edge = Edge::new("fork", "branch"); let resolved = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::Compact); - assert_eq!(resolved.requested(), Some(Fidelity::Full)); - assert_eq!(resolved.entry_fidelity(), Some(Fidelity::SummaryHigh)); + assert_eq!(resolved.requested, Some(Fidelity::Full)); + assert_eq!(resolved.effective, Some(Fidelity::SummaryHigh)); } #[test] fn parallel_branch_fidelity_equal_to_fork_inherits() { let mut node = Node::new("branch"); node.attrs - .insert("fidelity".to_string(), fidelity_attr("summary:high")); + .insert("fidelity".to_string(), str_attr("summary:high")); let edge = Edge::new("fork", "branch"); let resolution = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::SummaryHigh); - assert_eq!(resolution.requested(), Some(Fidelity::SummaryHigh)); - assert_eq!(resolution.entry_fidelity(), None); + assert_eq!(resolution.requested, Some(Fidelity::SummaryHigh)); + assert_eq!(resolution.effective, None); } #[test] fn explicit_full_branch_equal_to_degraded_fork_inherits() { let mut node = Node::new("branch"); - node.attrs - .insert("fidelity".to_string(), fidelity_attr("full")); + node.attrs.insert("fidelity".to_string(), str_attr("full")); let edge = Edge::new("fork", "branch"); let resolution = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::SummaryHigh); - assert_eq!(resolution.requested(), Some(Fidelity::Full)); - assert_eq!(resolution.entry_fidelity(), None); + assert_eq!(resolution.requested, Some(Fidelity::Full)); + assert_eq!(resolution.effective, None); } #[test] @@ -494,8 +524,8 @@ mod tests { let resolution = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::Full); - assert_eq!(resolution.requested(), None); - assert_eq!(resolution.entry_fidelity(), None); + assert_eq!(resolution.requested, None); + assert_eq!(resolution.effective, None); } #[tokio::test]