返回 CodeWhale
lib.rs
根目录 / crates / workflow / src / lib.rs
1 //! Typed Workflow IR and validation for CodeWhale.
2 //!
3 //! This crate deliberately stops at the Rust-owned IR boundary. Runtime tool
4 //! exposure, worktree application, replay, and model execution are layered on
5 //! top only after their cancellation and evidence semantics are proven.
6
7 mod elevation;
8 /// Setup-time Fleet composition: a suggestion schema with no runtime authority.
9 ///
10 /// Deliberately **not** re-exported from the crate root. The setup wizard uses
11 /// it by its explicit module path
12 /// (`codewhale_workflow::fleet_composition::…`), keeping the advisory boundary
13 /// visible instead of making the schema look like runtime Fleet authority.
14 pub mod fleet_composition;
15 pub mod fleet_exact;
16 pub mod fleet_preflight;
17 pub mod fleet_reasoning;
18 pub mod fleet_snapshot;
19 mod gates;
20 mod js_authoring;
21 mod model_policy;
22 mod named_fleet;
23 pub mod reasoning_router;
24 pub mod redaction;
25 mod role_resolve;
26
27 use std::collections::{BTreeMap, BTreeSet};
28 use std::path::Path;
29
30 use serde::{Deserialize, Serialize};
31 use thiserror::Error;
32
33 pub use elevation::{
34 DEFAULT_HIGH_BUDGET_THRESHOLD, ElevationOptions, PlanRiskHint, WorkflowPlanElevation,
35 assess_plan_risk_string, assess_workflow_elevation, is_shell_tool, is_write_tool,
36 };
37 pub use fleet_exact::{
38 EXACT_FLEET_SCHEMA_KIND, EXACT_FLEET_SCHEMA_REVISION, ExactFleet, ExactFleetError, ExactMember,
39 FrozenRoute, LEGACY_FLEET_SCHEMA_KIND, PermissionCeiling, ROLE_ALIASES, ROUTER_PUBLIC_ID,
40 ROUTER_PUBLIC_ROLE, ReasoningTier, RequestedReasoning, RouterMember, ShellCeiling,
41 canonical_member_key, canonical_role_key,
42 };
43 pub use fleet_preflight::{
44 CredentialReadiness, EndpointIdentity, PreflightError, PreflightedRoute, RoutePreflight,
45 };
46 pub use fleet_reasoning::{
47 EffectiveReasoning, EffectiveReasoningSource, FAITHFUL_WIRE_TIERS, FleetTaskReceipt,
48 ProviderEffectiveReasoning, ProviderReasoningControl, ROUTER_CALL_REASONING,
49 ROUTER_MAX_OUTPUT_TOKENS, ROUTER_REASONING_FIELD, ROUTER_SUMMARY_MAX_CHARS, ROUTING_SCOPE,
50 ReasoningCapability, ReasoningResolveError, ResolvedReasoning, RouterAvailability,
51 RouterCallDisclosure, RouterCallInput, RouterCallPlan, RouterDecision, RouterDecisionError,
52 RouterIdentity, RoutingDisclosure, RoutingPayload, TaskShape, bounded_routing_payload,
53 parse_router_decision, resolve_exact_member_reasoning, resolve_legacy_reasoning,
54 router_call_plan, router_system_prompt, router_user_message, transport_disclosure,
55 };
56 pub use fleet_snapshot::{
57 FleetSnapshot, FleetSnapshotLegacyRole, FleetSnapshotMember, FleetSnapshotRouter,
58 QualifiedFleetId, captured_legacy_inline_router, verify_snapshot_content_hash,
59 };
60 pub use gates::{
61 GateError, GateKind, GateOn, GateOnFail, GateOutcome, GateSpec, GateState, GateStatusLine,
62 HandoffArtifact, LaneGateBoard, stopship_gate_pipeline,
63 };
64 pub use js_authoring::{
65 JavascriptWorkflowError, JavascriptWorkflowResult, compile_javascript_workflow,
66 compile_typescript_workflow,
67 };
68 pub use model_policy::*;
69 pub use named_fleet::{
70 FleetDocument, FleetSchema, FleetSearchRoot, NamedFleet, NamedFleetError,
71 STOPSHIP_REQUIRED_ROLES, exact_schema_revision, load_named_fleet, load_named_fleet_file,
72 parse_named_fleet, split_qualified_fleet_name, validate_fleet_file_stem,
73 };
74 pub use reasoning_router::{
75 CapturedReasoningRouter, FleetRouterRef, LEGACY_INLINE_ROUTER_ORIGIN, QualifiedRouterId,
76 REASONING_ROUTER_DIR, REASONING_ROUTER_SCHEMA_KIND, REASONING_ROUTER_SERVICE_KIND,
77 ReasoningRouterError, ReasoningRouterProfile, RouterCallReasoning,
78 };
79 pub use redaction::{
80 REDACTION_ABSOLUTE_PATH, REDACTION_RELATIVE_PATH, REDACTION_SECRET, Redaction,
81 redact_for_disclosure,
82 };
83 pub use role_resolve::{
84 FleetRoleMap, FleetRoleResolveError, ResolvedWorkflowAgent, normalize_token,
85 resolve_workflow_agent, validate_role_token,
86 };
87
88 /// Default hard ceiling on total agents a Fleet-shaped Workflow plan may launch.
89 /// Matches the imperative VM lifetime cap (1_000 agents per run).
90 pub const DEFAULT_FLEET_WORKFLOW_MAX_AGENTS: usize = 1000;
91 pub const DEFAULT_FLEET_WORKFLOW_MAX_DEPTH: usize = 5;
92
93 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
94 pub struct WorkflowConfig {
95 pub goal: String,
96 #[serde(default = "default_max_concurrent")]
97 pub max_concurrent: u8,
98 #[serde(default)]
99 pub description: Option<String>,
100 #[serde(default)]
101 pub phases: Vec<Phase>,
102 }
103
104 impl WorkflowConfig {
105 pub fn validate(&self) -> Result<(), WorkflowValidationError> {
106 WorkflowPlan::from_config(self).map(|_| ())
107 }
108
109 pub fn compile(&self) -> Result<WorkflowPlan, WorkflowValidationError> {
110 WorkflowPlan::from_config(self)
111 }
112 }
113
114 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
115 pub struct WorkflowSpec {
116 #[serde(default)]
117 pub id: Option<String>,
118 pub goal: String,
119 #[serde(default)]
120 pub description: Option<String>,
121 #[serde(default)]
122 pub budget: BudgetSpec,
123 #[serde(default)]
124 pub permissions: PermissionSpec,
125 #[serde(default)]
126 pub model_policy: ModelPolicy,
127 #[serde(default)]
128 pub promotion_policy: PromotionPolicy,
129 /// Workflow-owned role gates and handoffs (#4179). Fleet supplies roles;
130 /// gate semantics stay attached to the Workflow definition.
131 #[serde(default, skip_serializing_if = "Vec::is_empty")]
132 pub gates: Vec<GateSpec>,
133 #[serde(default)]
134 pub nodes: Vec<WorkflowNode>,
135 }
136
137 impl WorkflowSpec {
138 pub fn validate_for_fleet(&self) -> Result<WorkflowFleetShape, WorkflowFleetLimitError> {
139 self.validate_for_fleet_with_limits(WorkflowFleetLimits::default())
140 }
141
142 pub fn validate_for_fleet_with_limits(
143 &self,
144 limits: WorkflowFleetLimits,
145 ) -> Result<WorkflowFleetShape, WorkflowFleetLimitError> {
146 validate_workflow_nodes(&self.nodes)
147 .map_err(|source| WorkflowFleetLimitError::InvalidWorkflow { source })?;
148 let shape = estimate_fleet_shape(&self.nodes)?;
149 if shape.total_agents > limits.max_total_agents {
150 return Err(WorkflowFleetLimitError::TooManyAgents {
151 total_agents: shape.total_agents,
152 max_total_agents: limits.max_total_agents,
153 });
154 }
155 if shape.max_depth > limits.max_depth {
156 return Err(WorkflowFleetLimitError::RecursionTooDeep {
157 depth: shape.max_depth,
158 max_depth: limits.max_depth,
159 });
160 }
161 Ok(shape)
162 }
163 }
164
165 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
166 #[serde(tag = "kind", content = "spec", rename_all = "snake_case")]
167 pub enum WorkflowNode {
168 BranchSet(BranchSpec),
169 Leaf(LeafSpec),
170 Sequence(SequenceSpec),
171 Reduce(ReduceSpec),
172 TeacherReview(TeacherReviewSpec),
173 LoopUntil(LoopUntilSpec),
174 Cond(CondSpec),
175 Expand(ExpandSpec),
176 }
177
178 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
179 pub struct BranchSpec {
180 pub id: String,
181 #[serde(default)]
182 pub description: Option<String>,
183 #[serde(default)]
184 pub parallel: bool,
185 #[serde(default)]
186 pub budget: BudgetSpec,
187 #[serde(default)]
188 pub permissions: PermissionSpec,
189 #[serde(default)]
190 pub model_policy: ModelPolicy,
191 #[serde(default)]
192 pub children: Vec<WorkflowNode>,
193 }
194
195 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
196 pub struct LeafSpec {
197 pub id: String,
198 pub prompt: String,
199 #[serde(default)]
200 pub agent_type: AgentType,
201 /// Fleet role this step should run as (e.g. `scout`, `implementer`).
202 /// Resolved via the fleet roster at dispatch time (#4177). Preferred
203 /// identity field for workflow steps; `profile` remains an explicit
204 /// AgentProfile override. Provider/model live on [`ModelPolicy`] as
205 /// optional overrides — they are **not** required identity fields.
206 #[serde(default, skip_serializing_if = "Option::is_none")]
207 pub role: Option<String>,
208 /// Named Fleet roster profile this agent should run as. Resolved against
209 /// the saved Fleet roster at dispatch time; unknown names fail validation
210 /// before any spawn. When set, role/model/loadout defaults come from the
211 /// roster member; explicit fields on this spec override the profile.
212 /// Precedence over `role` when both are set (#4177 / #4111).
213 #[serde(default, skip_serializing_if = "Option::is_none")]
214 pub profile: Option<String>,
215 #[serde(default)]
216 pub mode: TaskMode,
217 #[serde(default)]
218 pub isolation: IsolationMode,
219 #[serde(default)]
220 pub file_scope: Vec<String>,
221 /// Optional child working directory, repository-relative like
222 /// `task({cwd})`. Disambiguates multi-repo workspaces (#6232).
223 #[serde(default, skip_serializing_if = "Option::is_none")]
224 pub cwd: Option<String>,
225 #[serde(default)]
226 pub depends_on_results: Vec<String>,
227 #[serde(default)]
228 pub budget: BudgetSpec,
229 #[serde(default)]
230 pub permissions: PermissionSpec,
231 #[serde(default)]
232 pub model_policy: ModelPolicy,
233 }
234
235 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
236 pub struct SequenceSpec {
237 pub id: String,
238 #[serde(default)]
239 pub children: Vec<WorkflowNode>,
240 }
241
242 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
243 pub struct ReduceSpec {
244 pub id: String,
245 #[serde(default)]
246 pub inputs: Vec<String>,
247 pub prompt: String,
248 #[serde(default)]
249 pub model_policy: ModelPolicy,
250 }
251
252 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
253 pub struct TeacherReviewSpec {
254 pub id: String,
255 #[serde(default)]
256 pub candidates: Vec<String>,
257 #[serde(default)]
258 pub promotion_policy: PromotionPolicy,
259 }
260
261 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
262 pub struct LoopUntilSpec {
263 pub id: String,
264 pub condition: String,
265 #[serde(default)]
266 pub max_iterations: Option<u32>,
267 #[serde(default)]
268 pub children: Vec<WorkflowNode>,
269 }
270
271 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
272 pub struct CondSpec {
273 pub id: String,
274 pub condition: String,
275 #[serde(default)]
276 pub then_nodes: Vec<WorkflowNode>,
277 #[serde(default)]
278 pub else_nodes: Vec<WorkflowNode>,
279 }
280
281 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
282 pub struct ExpandSpec {
283 pub id: String,
284 pub source: String,
285 #[serde(default)]
286 pub max_children: Option<usize>,
287 #[serde(default)]
288 pub template: Option<Box<WorkflowNode>>,
289 }
290
291 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
292 pub struct BudgetSpec {
293 #[serde(default)]
294 pub max_steps: Option<u32>,
295 #[serde(default)]
296 pub timeout_secs: Option<u64>,
297 #[serde(default)]
298 pub max_parallel: Option<u8>,
299 #[serde(default)]
300 pub max_tokens: Option<u64>,
301 }
302
303 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
304 pub struct PermissionSpec {
305 #[serde(default)]
306 pub allow_write: bool,
307 #[serde(default)]
308 pub allow_network: bool,
309 /// Expose no tools to the child. This is distinct from an empty
310 /// `allowed_tools` list, which preserves the role's default tool surface.
311 #[serde(default)]
312 pub deny_all_tools: bool,
313 #[serde(default)]
314 pub allowed_tools: Vec<String>,
315 #[serde(default)]
316 pub file_scope: Vec<String>,
317 }
318
319 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
320 pub struct ModelPolicy {
321 #[serde(default)]
322 pub provider: Option<String>,
323 #[serde(default)]
324 pub model: Option<String>,
325 #[serde(default)]
326 pub fallback_models: Vec<String>,
327 }
328
329 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
330 pub struct PromotionPolicy {
331 #[serde(default)]
332 pub strategy: PromotionStrategy,
333 #[serde(default)]
334 pub require_teacher_review: bool,
335 #[serde(default)]
336 pub min_successful_branches: Option<u32>,
337 #[serde(default)]
338 pub promotion_gate: PromotionGate,
339 }
340
341 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
342 #[serde(rename_all = "snake_case")]
343 pub enum PromotionStrategy {
344 #[default]
345 All,
346 FirstSuccess,
347 BestScore,
348 TeacherSelected,
349 }
350
351 #[derive(Debug, Clone, PartialEq, Eq)]
352 pub struct WorkflowPlan {
353 goal: String,
354 max_concurrent: u8,
355 phases: Vec<PhasePlan>,
356 }
357
358 impl WorkflowPlan {
359 pub fn from_config(config: &WorkflowConfig) -> Result<Self, WorkflowValidationError> {
360 validate_non_empty("workflow goal", &config.goal)?;
361 if !(1..=20).contains(&config.max_concurrent) {
362 return Err(WorkflowValidationError::InvalidMaxConcurrent {
363 value: config.max_concurrent,
364 });
365 }
366 if config.phases.is_empty() {
367 return Err(WorkflowValidationError::EmptyWorkflow);
368 }
369
370 let mut phase_indices = BTreeMap::new();
371 let mut all_tasks = BTreeMap::new();
372 let mut task_phase = BTreeMap::new();
373
374 for (phase_index, phase) in config.phases.iter().enumerate() {
375 validate_non_empty("phase name", &phase.name)?;
376 if phase.tasks.is_empty() {
377 return Err(WorkflowValidationError::EmptyPhase {
378 phase: phase.name.clone(),
379 });
380 }
381 if phase_indices
382 .insert(phase.name.clone(), phase_index)
383 .is_some()
384 {
385 return Err(WorkflowValidationError::DuplicatePhase {
386 phase: phase.name.clone(),
387 });
388 }
389
390 for task in &phase.tasks {
391 validate_non_empty("task id", &task.id)?;
392 validate_non_empty("task prompt", &task.prompt)?;
393 if all_tasks.insert(task.id.clone(), task).is_some() {
394 return Err(WorkflowValidationError::DuplicateTask {
395 task: task.id.clone(),
396 });
397 }
398 task_phase.insert(task.id.clone(), phase.name.clone());
399 }
400 }
401
402 for phase in &config.phases {
403 for dependency in &phase.depends_on {
404 if dependency == &phase.name || !phase_indices.contains_key(dependency) {
405 return Err(WorkflowValidationError::InvalidPhaseDependency {
406 phase: phase.name.clone(),
407 dependency: dependency.clone(),
408 });
409 }
410 }
411 validate_parallel_write_scope(phase)?;
412 }
413
414 let ordered_phase_names = ordered_phases(config, &phase_indices)?;
415 let phase_order: BTreeMap<_, _> = ordered_phase_names
416 .iter()
417 .enumerate()
418 .map(|(index, phase)| (phase.clone(), index))
419 .collect();
420
421 for phase in &config.phases {
422 for task in &phase.tasks {
423 for dependency in &task.depends_on_results {
424 let Some(dependency_phase) = task_phase.get(dependency) else {
425 return Err(WorkflowValidationError::InvalidTaskResultDependency {
426 task: task.id.clone(),
427 dependency: dependency.clone(),
428 });
429 };
430 if phase_order[dependency_phase] >= phase_order[&phase.name] {
431 return Err(WorkflowValidationError::UnavailableTaskResultDependency {
432 task: task.id.clone(),
433 dependency: dependency.clone(),
434 dependency_phase: dependency_phase.clone(),
435 task_phase: phase.name.clone(),
436 });
437 }
438 }
439 }
440 }
441
442 let phases = ordered_phase_names
443 .iter()
444 .map(|phase_name| {
445 let phase = &config.phases[phase_indices[phase_name]];
446 PhasePlan {
447 name: phase.name.clone(),
448 parallel: phase.parallel,
449 on_failure: phase.on_failure,
450 tasks: phase.tasks.clone(),
451 }
452 })
453 .collect();
454
455 Ok(Self {
456 goal: config.goal.clone(),
457 max_concurrent: config.max_concurrent,
458 phases,
459 })
460 }
461
462 pub fn goal(&self) -> &str {
463 &self.goal
464 }
465
466 pub fn max_concurrent(&self) -> u8 {
467 self.max_concurrent
468 }
469
470 pub fn phases(&self) -> &[PhasePlan] {
471 &self.phases
472 }
473
474 pub fn phase_names(&self) -> impl Iterator<Item = &str> {
475 self.phases.iter().map(|phase| phase.name.as_str())
476 }
477 }
478
479 pub type WorkflowIr = WorkflowPlan;
480
481 #[derive(Debug, Clone, PartialEq, Eq)]
482 pub struct PhasePlan {
483 pub name: String,
484 pub parallel: bool,
485 pub on_failure: FailurePolicy,
486 pub tasks: Vec<Task>,
487 }
488
489 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
490 pub struct Phase {
491 pub name: String,
492 #[serde(default)]
493 pub description: Option<String>,
494 #[serde(default)]
495 pub depends_on: Vec<String>,
496 #[serde(default)]
497 pub parallel: bool,
498 #[serde(default)]
499 pub on_failure: FailurePolicy,
500 #[serde(default)]
501 pub tasks: Vec<Task>,
502 }
503
504 pub type WorkflowPhase = Phase;
505
506 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
507 #[serde(rename_all = "snake_case")]
508 pub enum FailurePolicy {
509 #[default]
510 SkipContinue,
511 Abort,
512 }
513
514 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
515 pub struct Task {
516 pub id: String,
517 pub prompt: String,
518 #[serde(default)]
519 pub agent_type: AgentType,
520 #[serde(default)]
521 pub mode: TaskMode,
522 #[serde(default)]
523 pub isolation: IsolationMode,
524 #[serde(default)]
525 pub file_scope: Vec<String>,
526 #[serde(default)]
527 pub depends_on_results: Vec<String>,
528 #[serde(default)]
529 pub max_steps: Option<u32>,
530 #[serde(default)]
531 pub timeout_secs: Option<u64>,
532 }
533
534 pub type WorkflowTask = Task;
535 pub type WorkflowRole = AgentType;
536
537 /// The workflow wire's agent type. The serialized spellings are the canonical
538 /// Codewhale role vocabulary (founder, 2026-09-02); the pre-rename JS
539 /// spellings stay accepted as aliases so checked-in and saved workflows keep
540 /// compiling.
541 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
542 #[serde(rename_all = "snake_case")]
543 pub enum AgentType {
544 #[default]
545 General,
546 #[serde(alias = "scout")]
547 Explore,
548 #[serde(rename = "planner", alias = "plan", alias = "awaiter")]
549 Plan,
550 #[serde(
551 alias = "review",
552 alias = "reviewer",
553 alias = "consultant",
554 alias = "oracle"
555 )]
556 Review,
557 #[serde(rename = "implement", alias = "implementer", alias = "builder")]
558 Implementer,
559 #[serde(rename = "test", alias = "verifier", alias = "verify")]
560 Verifier,
561 }
562
563 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
564 #[serde(rename_all = "snake_case")]
565 pub enum TaskMode {
566 #[default]
567 ReadOnly,
568 ReadWrite,
569 }
570
571 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
572 #[serde(rename_all = "snake_case")]
573 pub enum IsolationMode {
574 /// Runtime chooses isolation.
575 ///
576 /// Parallel write-capable children resolve to [`IsolationMode::Worktree`]
577 /// so concurrent writers do not collide in the parent checkout. Explicit
578 /// [`IsolationMode::Shared`] is the plan-level same-worktree override.
579 #[default]
580 Auto,
581 /// Share the parent checkout (same-worktree).
582 Shared,
583 /// Dedicated git worktree / branch.
584 Worktree,
585 }
586
587 impl IsolationMode {
588 /// Resolve [`IsolationMode::Auto`] for a leaf.
589 ///
590 /// When `parallel_write` is true (leaf is write-capable inside a parallel
591 /// branch), Auto becomes Worktree. Otherwise Auto becomes Shared.
592 #[must_use]
593 pub fn resolve(self, parallel_write: bool) -> Self {
594 match self {
595 Self::Auto if parallel_write => Self::Worktree,
596 Self::Auto => Self::Shared,
597 other => other,
598 }
599 }
600
601 /// Whether the resolved mode provisions a dedicated worktree.
602 #[must_use]
603 pub fn wants_worktree(self, parallel_write: bool) -> bool {
604 matches!(self.resolve(parallel_write), Self::Worktree)
605 }
606 }
607
608 /// A leaf is write-capable when it can mutate the workspace.
609 ///
610 /// Authority comes from the declared mode and permissions, not the agent's
611 /// role identity. In particular, an implementer may be deliberately confined
612 /// to a read-only verification task. Used by workflow lowering to decide the
613 /// default isolation for parallel children (#4120).
614 #[must_use]
615 pub fn leaf_is_write_capable(spec: &LeafSpec) -> bool {
616 spec.mode == TaskMode::ReadWrite || spec.permissions.allow_write
617 }
618
619 /// Effective worktree flag for a leaf given whether it is being lowered inside
620 /// a parallel branch.
621 #[must_use]
622 pub fn leaf_wants_worktree(spec: &LeafSpec, parallel: bool) -> bool {
623 let parallel_write = parallel && leaf_is_write_capable(spec);
624 spec.isolation.wants_worktree(parallel_write)
625 }
626
627 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
628 pub struct BranchResult {
629 pub branch_id: String,
630 pub task_id: String,
631 pub status: WorkflowRunStatus,
632 #[serde(default)]
633 pub usage: WorkflowUsage,
634 #[serde(default)]
635 pub memo_usage: WorkflowMemoUsage,
636 #[serde(default)]
637 pub artifacts: Vec<String>,
638 #[serde(default)]
639 pub notes: Option<String>,
640 }
641
642 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
643 pub struct LeafResult {
644 pub leaf_id: String,
645 pub task_id: String,
646 /// Fleet role the leaf was declared to run as, if any (#4177).
647 #[serde(default, skip_serializing_if = "Option::is_none")]
648 pub role: Option<String>,
649 /// Fleet roster profile the leaf was declared to run as, if any.
650 #[serde(default, skip_serializing_if = "Option::is_none")]
651 pub profile: Option<String>,
652 pub status: WorkflowRunStatus,
653 #[serde(default)]
654 pub usage: WorkflowUsage,
655 #[serde(default)]
656 pub memo_usage: WorkflowMemoUsage,
657 #[serde(default)]
658 pub output: Option<String>,
659 #[serde(default)]
660 pub artifacts: Vec<String>,
661 /// Post-hoc validation failure for the leaf's structured response.
662 #[serde(default, skip_serializing_if = "Option::is_none")]
663 pub schema_error: Option<String>,
664 }
665
666 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
667 pub struct WorkflowUsage {
668 #[serde(default, skip_serializing_if = "Option::is_none")]
669 pub input_tokens: Option<u64>,
670 #[serde(default, skip_serializing_if = "Option::is_none")]
671 pub output_tokens: Option<u64>,
672 #[serde(default, skip_serializing_if = "Option::is_none")]
673 pub cost_microusd: Option<u64>,
674 }
675
676 impl WorkflowUsage {
677 #[must_use]
678 pub fn total_tokens(self) -> Option<u64> {
679 self.input_tokens
680 .zip(self.output_tokens)
681 .map(|(input, output)| input.saturating_add(output))
682 }
683
684 /// Add two independently observed usage receipts. A field remains known
685 /// only when both contributors reported it; `Some(0)` is still observed.
686 pub(crate) fn add_assign(&mut self, other: Self) {
687 self.input_tokens = sum_reported(self.input_tokens, other.input_tokens);
688 self.output_tokens = sum_reported(self.output_tokens, other.output_tokens);
689 self.cost_microusd = sum_reported(self.cost_microusd, other.cost_microusd);
690 }
691 }
692
693 fn sum_reported(left: Option<u64>, right: Option<u64>) -> Option<u64> {
694 left.zip(right)
695 .map(|(left, right)| left.saturating_add(right))
696 }
697
698 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
699 pub struct WorkflowMemoUsage {
700 #[serde(default)]
701 pub armh_hits: u64,
702 #[serde(default)]
703 pub armh_misses: u64,
704 #[serde(default)]
705 pub armh_saved_estimated_tokens: u64,
706 #[serde(default)]
707 pub provider_prompt_cache_hits: u64,
708 #[serde(default)]
709 pub provider_prompt_cache_misses: u64,
710 }
711
712 impl WorkflowMemoUsage {
713 pub(crate) fn add_assign(&mut self, other: Self) {
714 self.armh_hits = self.armh_hits.saturating_add(other.armh_hits);
715 self.armh_misses = self.armh_misses.saturating_add(other.armh_misses);
716 self.armh_saved_estimated_tokens = self
717 .armh_saved_estimated_tokens
718 .saturating_add(other.armh_saved_estimated_tokens);
719 self.provider_prompt_cache_hits = self
720 .provider_prompt_cache_hits
721 .saturating_add(other.provider_prompt_cache_hits);
722 self.provider_prompt_cache_misses = self
723 .provider_prompt_cache_misses
724 .saturating_add(other.provider_prompt_cache_misses);
725 }
726 }
727
728 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
729 pub struct ControlNodeResult {
730 pub node_id: String,
731 pub kind: ControlNodeKind,
732 pub status: WorkflowRunStatus,
733 #[serde(default)]
734 pub selected_children: Vec<String>,
735 #[serde(default)]
736 pub summary: Option<String>,
737 }
738
739 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
740 #[serde(rename_all = "snake_case")]
741 pub enum WorkflowRunStatus {
742 #[default]
743 Pending,
744 Running,
745 Succeeded,
746 Failed,
747 Cancelled,
748 BudgetExceeded,
749 }
750
751 #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
752 #[serde(rename_all = "snake_case")]
753 pub enum ControlNodeKind {
754 BranchSet,
755 Leaf,
756 Sequence,
757 Reduce,
758 TeacherReview,
759 LoopUntil,
760 Cond,
761 Expand,
762 }
763
764 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
765 pub struct WorkflowExecution {
766 pub status: WorkflowRunStatus,
767 #[serde(default)]
768 pub usage: WorkflowUsage,
769 #[serde(default)]
770 pub memo_usage: WorkflowMemoUsage,
771 #[serde(default)]
772 pub leaf_results: Vec<LeafResult>,
773 #[serde(default)]
774 pub branch_results: Vec<BranchResult>,
775 #[serde(default)]
776 pub control_node_results: Vec<ControlNodeResult>,
777 }
778
779 impl Default for WorkflowExecution {
780 fn default() -> Self {
781 Self {
782 status: WorkflowRunStatus::Succeeded,
783 usage: WorkflowUsage::default(),
784 memo_usage: WorkflowMemoUsage::default(),
785 leaf_results: Vec::new(),
786 branch_results: Vec::new(),
787 control_node_results: Vec::new(),
788 }
789 }
790 }
791
792 impl WorkflowExecution {
793 pub fn mark_failed(&mut self) {
794 self.status = WorkflowRunStatus::Failed;
795 }
796
797 pub fn mark_cancelled(&mut self) {
798 self.status = WorkflowRunStatus::Cancelled;
799 }
800
801 pub fn mark_budget_exceeded(&mut self) {
802 self.status = WorkflowRunStatus::BudgetExceeded;
803 }
804
805 fn should_stop_mock_execution(&self) -> bool {
806 matches!(
807 self.status,
808 WorkflowRunStatus::Cancelled | WorkflowRunStatus::BudgetExceeded
809 )
810 }
811 }
812
813 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
814 pub struct MockLeafOutcome {
815 pub status: WorkflowRunStatus,
816 #[serde(default)]
817 pub usage: WorkflowUsage,
818 #[serde(default)]
819 pub memo_usage: WorkflowMemoUsage,
820 #[serde(default)]
821 pub output: Option<String>,
822 #[serde(default)]
823 pub artifacts: Vec<String>,
824 }
825
826 impl MockLeafOutcome {
827 pub fn succeeded(output: impl Into<String>) -> Self {
828 Self {
829 status: WorkflowRunStatus::Succeeded,
830 usage: WorkflowUsage::default(),
831 memo_usage: WorkflowMemoUsage::default(),
832 output: Some(output.into()),
833 artifacts: Vec::new(),
834 }
835 }
836
837 pub fn failed(output: impl Into<String>) -> Self {
838 Self {
839 status: WorkflowRunStatus::Failed,
840 usage: WorkflowUsage::default(),
841 memo_usage: WorkflowMemoUsage::default(),
842 output: Some(output.into()),
843 artifacts: Vec::new(),
844 }
845 }
846
847 pub fn with_usage(mut self, usage: WorkflowUsage) -> Self {
848 self.usage = usage;
849 self
850 }
851
852 pub fn with_memo_usage(mut self, memo_usage: WorkflowMemoUsage) -> Self {
853 self.memo_usage = memo_usage;
854 self
855 }
856 }
857
858 #[derive(Debug, Default, Clone)]
859 pub struct MockWorkflowExecutor {
860 leaf_outcomes: BTreeMap<String, MockLeafOutcome>,
861 predicate_results: BTreeMap<String, Vec<bool>>,
862 generated_nodes: BTreeMap<String, Vec<WorkflowNode>>,
863 cancelled: bool,
864 max_leaf_steps: Option<u32>,
865 leaf_steps_executed: u32,
866 max_leaf_tokens: Option<u64>,
867 leaf_tokens_used: u64,
868 }
869
870 impl MockWorkflowExecutor {
871 pub fn new() -> Self {
872 Self::default()
873 }
874
875 pub fn with_leaf_outcome(
876 mut self,
877 leaf_id: impl Into<String>,
878 outcome: MockLeafOutcome,
879 ) -> Self {
880 self.leaf_outcomes.insert(leaf_id.into(), outcome);
881 self
882 }
883
884 pub fn with_predicate_results(
885 mut self,
886 node_id: impl Into<String>,
887 results: Vec<bool>,
888 ) -> Self {
889 self.predicate_results.insert(node_id.into(), results);
890 self
891 }
892
893 pub fn with_generated_nodes(
894 mut self,
895 node_id: impl Into<String>,
896 nodes: Vec<WorkflowNode>,
897 ) -> Self {
898 self.generated_nodes.insert(node_id.into(), nodes);
899 self
900 }
901
902 pub fn with_cancelled(mut self) -> Self {
903 self.cancelled = true;
904 self
905 }
906
907 pub fn with_max_leaf_steps(mut self, max_leaf_steps: u32) -> Self {
908 self.max_leaf_steps = Some(max_leaf_steps);
909 self
910 }
911
912 pub fn with_max_leaf_tokens(mut self, max_leaf_tokens: u64) -> Self {
913 self.max_leaf_tokens = Some(max_leaf_tokens);
914 self
915 }
916
917 pub fn run(
918 &mut self,
919 spec: &WorkflowSpec,
920 ) -> Result<WorkflowExecution, WorkflowExecutionError> {
921 validate_workflow_nodes(&spec.nodes)?;
922 let mut execution = WorkflowExecution::default();
923 self.execute_nodes(&spec.nodes, &mut execution)?;
924 Ok(execution)
925 }
926
927 fn execute_nodes(
928 &mut self,
929 nodes: &[WorkflowNode],
930 execution: &mut WorkflowExecution,
931 ) -> Result<(), WorkflowExecutionError> {
932 for node in nodes {
933 if execution.should_stop_mock_execution() {
934 break;
935 }
936 self.execute_node(node, execution)?;
937 }
938 Ok(())
939 }
940
941 fn execute_node(
942 &mut self,
943 node: &WorkflowNode,
944 execution: &mut WorkflowExecution,
945 ) -> Result<(), WorkflowExecutionError> {
946 match node {
947 WorkflowNode::BranchSet(spec) => self.execute_branch_set(spec, execution),
948 WorkflowNode::Leaf(spec) => {
949 self.execute_leaf(spec, execution);
950 Ok(())
951 }
952 WorkflowNode::Sequence(spec) => {
953 self.execute_nodes(&spec.children, execution)?;
954 execution.control_node_results.push(ControlNodeResult {
955 node_id: spec.id.clone(),
956 kind: ControlNodeKind::Sequence,
957 status: execution.status,
958 selected_children: spec.children.iter().map(node_id).collect(),
959 summary: Some("sequence executed in declaration order".to_string()),
960 });
961 Ok(())
962 }
963 WorkflowNode::Reduce(spec) => {
964 execution.control_node_results.push(ControlNodeResult {
965 node_id: spec.id.clone(),
966 kind: ControlNodeKind::Reduce,
967 status: WorkflowRunStatus::Succeeded,
968 selected_children: spec.inputs.clone(),
969 summary: Some(spec.prompt.clone()),
970 });
971 Ok(())
972 }
973 WorkflowNode::TeacherReview(spec) => {
974 execution.control_node_results.push(ControlNodeResult {
975 node_id: spec.id.clone(),
976 kind: ControlNodeKind::TeacherReview,
977 status: WorkflowRunStatus::Succeeded,
978 selected_children: spec.candidates.clone(),
979 summary: Some(
980 "teacher review scaffold selected declared candidates".to_string(),
981 ),
982 });
983 Ok(())
984 }
985 WorkflowNode::LoopUntil(spec) => self.execute_loop_until(spec, execution),
986 WorkflowNode::Cond(spec) => self.execute_cond(spec, execution),
987 WorkflowNode::Expand(spec) => self.execute_expand(spec, execution),
988 }
989 }
990
991 fn execute_branch_set(
992 &mut self,
993 spec: &BranchSpec,
994 execution: &mut WorkflowExecution,
995 ) -> Result<(), WorkflowExecutionError> {
996 let before = execution.leaf_results.len();
997 self.execute_nodes(&spec.children, execution)?;
998 let status = aggregate_mock_status(&execution.leaf_results[before..]);
999 let mut usage = WorkflowUsage::default();
1000 let mut memo_usage = WorkflowMemoUsage::default();
1001 for (index, result) in execution.leaf_results[before..].iter().enumerate() {
1002 if index == 0 {
1003 usage = result.usage;
1004 } else {
1005 usage.add_assign(result.usage);
1006 }
1007 memo_usage.add_assign(result.memo_usage);
1008 }
1009 mark_execution_for_status(execution, status);
1010 execution.branch_results.push(BranchResult {
1011 branch_id: spec.id.clone(),
1012 task_id: spec.id.clone(),
1013 status,
1014 usage,
1015 memo_usage,
1016 artifacts: Vec::new(),
1017 notes: Some("mock branch set executed without runtime fanout".to_string()),
1018 });
1019 execution.control_node_results.push(ControlNodeResult {
1020 node_id: spec.id.clone(),
1021 kind: ControlNodeKind::BranchSet,
1022 status,
1023 selected_children: spec.children.iter().map(node_id).collect(),
1024 summary: Some("branch set scaffold executed children deterministically".to_string()),
1025 });
1026 Ok(())
1027 }
1028
1029 fn execute_leaf(&mut self, spec: &LeafSpec, execution: &mut WorkflowExecution) {
1030 let outcome = self.mock_leaf_outcome(spec);
1031 mark_execution_for_status(execution, outcome.status);
1032 if execution.leaf_results.is_empty() {
1033 execution.usage = outcome.usage;
1034 } else {
1035 execution.usage.add_assign(outcome.usage);
1036 }
1037 execution.memo_usage.add_assign(outcome.memo_usage);
1038 execution.leaf_results.push(LeafResult {
1039 leaf_id: spec.id.clone(),
1040 task_id: spec.id.clone(),
1041 role: spec.role.clone(),
1042 profile: spec.profile.clone(),
1043 status: outcome.status,
1044 usage: outcome.usage,
1045 memo_usage: outcome.memo_usage,
1046 output: outcome.output,
1047 artifacts: outcome.artifacts,
1048 schema_error: None,
1049 });
1050 }
1051
1052 fn execute_loop_until(
1053 &mut self,
1054 spec: &LoopUntilSpec,
1055 execution: &mut WorkflowExecution,
1056 ) -> Result<(), WorkflowExecutionError> {
1057 let max_iterations = spec.max_iterations.unwrap_or(1).max(1);
1058 let mut iterations = 0;
1059 let mut passed = false;
1060 while iterations < max_iterations {
1061 if execution.should_stop_mock_execution() {
1062 break;
1063 }
1064 iterations += 1;
1065 self.execute_nodes(&spec.children, execution)?;
1066 if execution.should_stop_mock_execution() {
1067 break;
1068 }
1069 if self.next_predicate_result(&spec.id) {
1070 passed = true;
1071 break;
1072 }
1073 }
1074 let status = if execution.should_stop_mock_execution() {
1075 execution.status
1076 } else if passed {
1077 WorkflowRunStatus::Succeeded
1078 } else {
1079 WorkflowRunStatus::Failed
1080 };
1081 mark_execution_for_status(execution, status);
1082 execution.control_node_results.push(ControlNodeResult {
1083 node_id: spec.id.clone(),
1084 kind: ControlNodeKind::LoopUntil,
1085 status,
1086 selected_children: spec.children.iter().map(node_id).collect(),
1087 summary: Some(format!("loop_until iterations={iterations}")),
1088 });
1089 Ok(())
1090 }
1091
1092 fn execute_cond(
1093 &mut self,
1094 spec: &CondSpec,
1095 execution: &mut WorkflowExecution,
1096 ) -> Result<(), WorkflowExecutionError> {
1097 let passed = self.next_predicate_result(&spec.id);
1098 let selected_nodes = if passed {
1099 &spec.then_nodes
1100 } else {
1101 &spec.else_nodes
1102 };
1103 self.execute_nodes(selected_nodes, execution)?;
1104 let status = if execution.should_stop_mock_execution() {
1105 execution.status
1106 } else {
1107 WorkflowRunStatus::Succeeded
1108 };
1109 execution.control_node_results.push(ControlNodeResult {
1110 node_id: spec.id.clone(),
1111 kind: ControlNodeKind::Cond,
1112 status,
1113 selected_children: selected_nodes.iter().map(node_id).collect(),
1114 summary: Some(format!("predicate_result={passed}")),
1115 });
1116 Ok(())
1117 }
1118
1119 fn execute_expand(
1120 &mut self,
1121 spec: &ExpandSpec,
1122 execution: &mut WorkflowExecution,
1123 ) -> Result<(), WorkflowExecutionError> {
1124 let mut nodes = self.generated_nodes.remove(&spec.id).unwrap_or_default();
1125 if let Some(max_children) = spec.max_children {
1126 nodes.truncate(max_children);
1127 }
1128 validate_workflow_node_shapes(&nodes)?;
1129 self.execute_nodes(&nodes, execution)?;
1130 let status = if execution.should_stop_mock_execution() {
1131 execution.status
1132 } else {
1133 WorkflowRunStatus::Succeeded
1134 };
1135 execution.control_node_results.push(ControlNodeResult {
1136 node_id: spec.id.clone(),
1137 kind: ControlNodeKind::Expand,
1138 status,
1139 selected_children: nodes.iter().map(node_id).collect(),
1140 summary: Some(format!("expanded_from={}", spec.source)),
1141 });
1142 Ok(())
1143 }
1144
1145 fn mock_leaf_outcome(&mut self, spec: &LeafSpec) -> MockLeafOutcome {
1146 if self.cancelled {
1147 return MockLeafOutcome {
1148 status: WorkflowRunStatus::Cancelled,
1149 usage: WorkflowUsage {
1150 input_tokens: Some(0),
1151 output_tokens: Some(0),
1152 cost_microusd: Some(0),
1153 },
1154 memo_usage: WorkflowMemoUsage::default(),
1155 output: Some("mock workflow cancelled before leaf execution".to_string()),
1156 artifacts: Vec::new(),
1157 };
1158 }
1159 if self
1160 .max_leaf_steps
1161 .is_some_and(|max| max > 0 && self.leaf_steps_executed >= max)
1162 {
1163 return MockLeafOutcome {
1164 status: WorkflowRunStatus::BudgetExceeded,
1165 // The leaf was rejected before execution, so zero usage is a
1166 // known observation rather than missing provider telemetry.
1167 usage: WorkflowUsage {
1168 input_tokens: Some(0),
1169 output_tokens: Some(0),
1170 cost_microusd: Some(0),
1171 },
1172 memo_usage: WorkflowMemoUsage::default(),
1173 output: Some("mock workflow leaf step budget exhausted".to_string()),
1174 artifacts: Vec::new(),
1175 };
1176 }
1177 if self
1178 .max_leaf_tokens
1179 .is_some_and(|max| self.leaf_tokens_used >= max)
1180 || spec.budget.max_tokens == Some(0)
1181 {
1182 return MockLeafOutcome {
1183 status: WorkflowRunStatus::BudgetExceeded,
1184 usage: WorkflowUsage {
1185 input_tokens: Some(0),
1186 output_tokens: Some(0),
1187 cost_microusd: Some(0),
1188 },
1189 memo_usage: WorkflowMemoUsage::default(),
1190 output: Some("mock workflow leaf token budget exhausted".to_string()),
1191 artifacts: Vec::new(),
1192 };
1193 }
1194 self.leaf_steps_executed = self.leaf_steps_executed.saturating_add(1);
1195 let outcome = self
1196 .leaf_outcomes
1197 .remove(&spec.id)
1198 .unwrap_or_else(|| MockLeafOutcome::succeeded(format!("mock leaf {}", spec.id)));
1199 let tokens = outcome.usage.total_tokens();
1200 if let (Some(per_leaf_token_cap), Some(tokens)) = (spec.budget.max_tokens, tokens)
1201 && tokens > per_leaf_token_cap
1202 {
1203 return MockLeafOutcome {
1204 status: WorkflowRunStatus::BudgetExceeded,
1205 usage: outcome.usage,
1206 memo_usage: outcome.memo_usage,
1207 output: Some(format!(
1208 "mock workflow leaf token budget exhausted ({tokens} > {per_leaf_token_cap})"
1209 )),
1210 artifacts: outcome.artifacts,
1211 };
1212 }
1213 if let Some(tokens) = tokens {
1214 self.leaf_tokens_used = self.leaf_tokens_used.saturating_add(tokens);
1215 }
1216 outcome
1217 }
1218
1219 fn next_predicate_result(&mut self, node_id: &str) -> bool {
1220 let Some(results) = self.predicate_results.get_mut(node_id) else {
1221 return false;
1222 };
1223 if results.is_empty() {
1224 return false;
1225 }
1226 results.remove(0)
1227 }
1228 }
1229
1230 fn aggregate_mock_status(results: &[LeafResult]) -> WorkflowRunStatus {
1231 if results
1232 .iter()
1233 .any(|result| result.status == WorkflowRunStatus::Cancelled)
1234 {
1235 WorkflowRunStatus::Cancelled
1236 } else if results
1237 .iter()
1238 .any(|result| result.status == WorkflowRunStatus::BudgetExceeded)
1239 {
1240 WorkflowRunStatus::BudgetExceeded
1241 } else if results
1242 .iter()
1243 .any(|result| result.status != WorkflowRunStatus::Succeeded)
1244 {
1245 WorkflowRunStatus::Failed
1246 } else {
1247 WorkflowRunStatus::Succeeded
1248 }
1249 }
1250
1251 fn mark_execution_for_status(execution: &mut WorkflowExecution, status: WorkflowRunStatus) {
1252 match status {
1253 WorkflowRunStatus::Succeeded | WorkflowRunStatus::Pending | WorkflowRunStatus::Running => {}
1254 WorkflowRunStatus::Failed => execution.mark_failed(),
1255 WorkflowRunStatus::Cancelled => execution.mark_cancelled(),
1256 WorkflowRunStatus::BudgetExceeded => execution.mark_budget_exceeded(),
1257 }
1258 }
1259
1260 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1261 pub struct BranchCandidate {
1262 pub branch_id: String,
1263 pub status: WorkflowRunStatus,
1264 pub score: u32,
1265 pub cost: u64,
1266 #[serde(default)]
1267 pub diversity_key: Option<String>,
1268 }
1269
1270 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1271 #[serde(rename_all = "snake_case")]
1272 pub enum TeacherCandidateKind {
1273 Note,
1274 WorkflowRecipe,
1275 SkillPatch,
1276 RegressionTest,
1277 CachePolicyPatch,
1278 BranchHeuristic,
1279 AuthoringPromptPatch,
1280 }
1281
1282 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
1283 #[serde(rename_all = "snake_case")]
1284 pub enum TeacherCandidateStatus {
1285 #[default]
1286 Proposed,
1287 Accepted,
1288 Rejected,
1289 Promoted,
1290 }
1291
1292 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1293 pub struct TeacherCandidate {
1294 pub candidate_id: String,
1295 pub kind: TeacherCandidateKind,
1296 #[serde(default)]
1297 pub status: TeacherCandidateStatus,
1298 pub source_node_id: String,
1299 #[serde(default)]
1300 pub source_branch_id: Option<String>,
1301 pub summary: String,
1302 #[serde(default)]
1303 pub evidence: Vec<String>,
1304 #[serde(default)]
1305 pub replay_results: Vec<StudentReplayResult>,
1306 }
1307
1308 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
1309 pub struct StudentReplayMetrics {
1310 #[serde(default)]
1311 pub score: i32,
1312 #[serde(default)]
1313 pub cost_microusd: u64,
1314 }
1315
1316 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1317 pub struct StudentReplayTestResult {
1318 pub name: String,
1319 pub passed: bool,
1320 }
1321
1322 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1323 pub struct StudentReplayResult {
1324 pub trace_id: String,
1325 pub candidate_id: String,
1326 pub baseline: StudentReplayMetrics,
1327 pub candidate: StudentReplayMetrics,
1328 #[serde(default)]
1329 pub required_tests: Vec<StudentReplayTestResult>,
1330 #[serde(default)]
1331 pub policy_violations: Vec<String>,
1332 #[serde(default)]
1333 pub stale: bool,
1334 #[serde(default)]
1335 pub notes: Option<String>,
1336 }
1337
1338 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1339 pub struct PromotionGate {
1340 #[serde(default = "default_min_replay_score_delta")]
1341 pub min_score_delta: i32,
1342 #[serde(default)]
1343 pub max_cost_delta_microusd: Option<i64>,
1344 #[serde(default = "default_true")]
1345 pub require_all_tests_pass: bool,
1346 #[serde(default = "default_true")]
1347 pub reject_policy_violations: bool,
1348 #[serde(default = "default_true")]
1349 pub reject_stale_replay: bool,
1350 }
1351
1352 impl Default for PromotionGate {
1353 fn default() -> Self {
1354 Self {
1355 min_score_delta: default_min_replay_score_delta(),
1356 max_cost_delta_microusd: None,
1357 require_all_tests_pass: true,
1358 reject_policy_violations: true,
1359 reject_stale_replay: true,
1360 }
1361 }
1362 }
1363
1364 impl PromotionGate {
1365 pub fn evaluate_candidate(&self, candidate: &TeacherCandidate) -> PromotionGateDecision {
1366 let Some(replay) = candidate.replay_results.last() else {
1367 return PromotionGateDecision {
1368 candidate_id: candidate.candidate_id.clone(),
1369 status: TeacherCandidateStatus::Rejected,
1370 score_delta: 0,
1371 cost_delta_microusd: 0,
1372 reasons: vec!["no student replay result recorded".to_string()],
1373 };
1374 };
1375 self.evaluate_replay(&candidate.candidate_id, replay)
1376 }
1377
1378 pub fn evaluate_replay(
1379 &self,
1380 candidate_id: &str,
1381 replay: &StudentReplayResult,
1382 ) -> PromotionGateDecision {
1383 let score_delta = replay.score_delta();
1384 let cost_delta_microusd = replay.cost_delta_microusd();
1385 let mut reasons = Vec::new();
1386
1387 if score_delta < self.min_score_delta {
1388 reasons.push(format!(
1389 "score delta {score_delta} is below required {}",
1390 self.min_score_delta
1391 ));
1392 }
1393 if let Some(max_cost_delta) = self.max_cost_delta_microusd
1394 && cost_delta_microusd > max_cost_delta
1395 {
1396 reasons.push(format!(
1397 "cost delta {cost_delta_microusd} exceeds allowed {max_cost_delta}"
1398 ));
1399 }
1400 if self.require_all_tests_pass {
1401 for test in replay.required_tests.iter().filter(|test| !test.passed) {
1402 reasons.push(format!("required test `{}` failed", test.name));
1403 }
1404 }
1405 if self.reject_policy_violations {
1406 for violation in &replay.policy_violations {
1407 reasons.push(format!("policy violation: {violation}"));
1408 }
1409 }
1410 if self.reject_stale_replay && replay.stale {
1411 reasons.push("student replay result is stale".to_string());
1412 }
1413
1414 let status = if reasons.is_empty() {
1415 TeacherCandidateStatus::Promoted
1416 } else {
1417 TeacherCandidateStatus::Rejected
1418 };
1419 PromotionGateDecision {
1420 candidate_id: candidate_id.to_string(),
1421 status,
1422 score_delta,
1423 cost_delta_microusd,
1424 reasons,
1425 }
1426 }
1427 }
1428
1429 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1430 pub struct PromotionGateDecision {
1431 pub candidate_id: String,
1432 pub status: TeacherCandidateStatus,
1433 pub score_delta: i32,
1434 pub cost_delta_microusd: i64,
1435 #[serde(default)]
1436 pub reasons: Vec<String>,
1437 }
1438
1439 impl PromotionGateDecision {
1440 pub fn promoted(&self) -> bool {
1441 self.status == TeacherCandidateStatus::Promoted
1442 }
1443 }
1444
1445 impl StudentReplayResult {
1446 pub fn score_delta(&self) -> i32 {
1447 self.candidate.score.saturating_sub(self.baseline.score)
1448 }
1449
1450 pub fn cost_delta_microusd(&self) -> i64 {
1451 signed_u64_delta(self.candidate.cost_microusd, self.baseline.cost_microusd)
1452 }
1453 }
1454
1455 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
1456 pub struct TeacherReviewReport {
1457 pub review_node_id: String,
1458 #[serde(default)]
1459 pub candidates: Vec<TeacherCandidate>,
1460 }
1461
1462 impl TeacherReviewReport {
1463 pub fn from_execution(review: &TeacherReviewSpec, execution: &WorkflowExecution) -> Self {
1464 let candidates = teacher_candidates_from_execution(review, execution);
1465 Self {
1466 review_node_id: review.id.clone(),
1467 candidates,
1468 }
1469 }
1470 }
1471
1472 pub fn teacher_candidates_from_execution(
1473 review: &TeacherReviewSpec,
1474 execution: &WorkflowExecution,
1475 ) -> Vec<TeacherCandidate> {
1476 let mut candidates = Vec::new();
1477 for source in &review.candidates {
1478 if let Some(branch) = execution
1479 .branch_results
1480 .iter()
1481 .find(|branch| branch.branch_id == *source || branch.task_id == *source)
1482 {
1483 candidates.push(teacher_candidate_from_branch(review, branch));
1484 continue;
1485 }
1486 if let Some(leaf) = execution
1487 .leaf_results
1488 .iter()
1489 .find(|leaf| leaf.leaf_id == *source || leaf.task_id == *source)
1490 {
1491 candidates.push(teacher_candidate_from_leaf(review, leaf));
1492 continue;
1493 }
1494 if let Some(control) = execution
1495 .control_node_results
1496 .iter()
1497 .find(|control| control.node_id == *source)
1498 {
1499 candidates.push(teacher_candidate_from_control(review, control));
1500 }
1501 }
1502 candidates
1503 }
1504
1505 fn teacher_candidate_from_branch(
1506 review: &TeacherReviewSpec,
1507 branch: &BranchResult,
1508 ) -> TeacherCandidate {
1509 let kind =
1510 if branch.memo_usage.armh_hits > 0 || branch.memo_usage.provider_prompt_cache_hits > 0 {
1511 TeacherCandidateKind::CachePolicyPatch
1512 } else if branch.status == WorkflowRunStatus::Succeeded {
1513 TeacherCandidateKind::WorkflowRecipe
1514 } else {
1515 TeacherCandidateKind::BranchHeuristic
1516 };
1517 let mut evidence = vec![format!("status={:?}", branch.status)];
1518 if branch.usage.total_tokens().is_some_and(|tokens| tokens > 0)
1519 || branch.usage.cost_microusd.is_some_and(|cost| cost > 0)
1520 {
1521 evidence.push(format!(
1522 "tokens={}, cost_microusd={}",
1523 branch
1524 .usage
1525 .total_tokens()
1526 .map_or_else(|| "unknown".to_string(), |value| value.to_string()),
1527 branch
1528 .usage
1529 .cost_microusd
1530 .map_or_else(|| "unknown".to_string(), |value| value.to_string())
1531 ));
1532 }
1533 if branch.memo_usage.armh_hits > 0 || branch.memo_usage.provider_prompt_cache_hits > 0 {
1534 evidence.push(format!(
1535 "armh_hits={}, provider_prompt_cache_hits={}",
1536 branch.memo_usage.armh_hits, branch.memo_usage.provider_prompt_cache_hits
1537 ));
1538 }
1539 if let Some(notes) = branch.notes.as_deref() {
1540 evidence.push(format!("notes={notes}"));
1541 }
1542 TeacherCandidate {
1543 candidate_id: format!("{}:{}", review.id, branch.branch_id),
1544 kind,
1545 status: TeacherCandidateStatus::Proposed,
1546 source_node_id: branch.task_id.clone(),
1547 source_branch_id: Some(branch.branch_id.clone()),
1548 summary: format!(
1549 "TeacherReview candidate from branch `{}` with {:?} status.",
1550 branch.branch_id, branch.status
1551 ),
1552 evidence,
1553 replay_results: Vec::new(),
1554 }
1555 }
1556
1557 fn teacher_candidate_from_leaf(review: &TeacherReviewSpec, leaf: &LeafResult) -> TeacherCandidate {
1558 let kind = if leaf.status == WorkflowRunStatus::Failed {
1559 TeacherCandidateKind::RegressionTest
1560 } else if leaf.memo_usage.armh_hits > 0 || leaf.memo_usage.provider_prompt_cache_hits > 0 {
1561 TeacherCandidateKind::CachePolicyPatch
1562 } else {
1563 TeacherCandidateKind::Note
1564 };
1565 let mut evidence = vec![format!("status={:?}", leaf.status)];
1566 if let Some(output) = leaf.output.as_deref() {
1567 evidence.push(format!("output={}", truncate_evidence(output)));
1568 }
1569 TeacherCandidate {
1570 candidate_id: format!("{}:{}", review.id, leaf.leaf_id),
1571 kind,
1572 status: TeacherCandidateStatus::Proposed,
1573 source_node_id: leaf.leaf_id.clone(),
1574 source_branch_id: None,
1575 summary: format!(
1576 "TeacherReview candidate from leaf `{}` with {:?} status.",
1577 leaf.leaf_id, leaf.status
1578 ),
1579 evidence,
1580 replay_results: Vec::new(),
1581 }
1582 }
1583
1584 fn teacher_candidate_from_control(
1585 review: &TeacherReviewSpec,
1586 control: &ControlNodeResult,
1587 ) -> TeacherCandidate {
1588 let mut evidence = vec![format!("status={:?}", control.status)];
1589 if !control.selected_children.is_empty() {
1590 evidence.push(format!(
1591 "selected_children={}",
1592 control.selected_children.join(",")
1593 ));
1594 }
1595 if let Some(summary) = control.summary.as_deref() {
1596 evidence.push(format!("summary={}", truncate_evidence(summary)));
1597 }
1598 TeacherCandidate {
1599 candidate_id: format!("{}:{}", review.id, control.node_id),
1600 kind: TeacherCandidateKind::AuthoringPromptPatch,
1601 status: TeacherCandidateStatus::Proposed,
1602 source_node_id: control.node_id.clone(),
1603 source_branch_id: None,
1604 summary: format!(
1605 "TeacherReview candidate from control node `{}` ({:?}).",
1606 control.node_id, control.kind
1607 ),
1608 evidence,
1609 replay_results: Vec::new(),
1610 }
1611 }
1612
1613 fn default_min_replay_score_delta() -> i32 {
1614 1
1615 }
1616
1617 fn default_true() -> bool {
1618 true
1619 }
1620
1621 fn signed_u64_delta(candidate: u64, baseline: u64) -> i64 {
1622 if candidate >= baseline {
1623 i64::try_from(candidate - baseline).unwrap_or(i64::MAX)
1624 } else {
1625 -i64::try_from(baseline - candidate).unwrap_or(i64::MAX)
1626 }
1627 }
1628
1629 fn truncate_evidence(value: &str) -> String {
1630 const MAX_EVIDENCE_CHARS: usize = 240;
1631 if value.chars().count() <= MAX_EVIDENCE_CHARS {
1632 return value.to_string();
1633 }
1634 let mut truncated = value
1635 .chars()
1636 .take(MAX_EVIDENCE_CHARS.saturating_sub(1))
1637 .collect::<String>();
1638 truncated.push_str("...");
1639 truncated
1640 }
1641
1642 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1643 pub struct BranchTournament {
1644 #[serde(default)]
1645 pub min_score: u32,
1646 #[serde(default)]
1647 pub ordering: TournamentOrdering,
1648 }
1649
1650 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
1651 #[serde(rename_all = "snake_case")]
1652 pub enum TournamentOrdering {
1653 /// Historical behavior: choose the cheapest passing branch, then score.
1654 #[default]
1655 CostThenScore,
1656 /// Quality-first behavior for explicitly configured evaluation workflows.
1657 ScoreThenCost,
1658 }
1659
1660 impl BranchTournament {
1661 pub fn select(&self, candidates: &[BranchCandidate]) -> Option<BranchCandidate> {
1662 let passing = || {
1663 candidates.iter().filter(|candidate| {
1664 candidate.status == WorkflowRunStatus::Succeeded
1665 && candidate.score >= self.min_score
1666 })
1667 };
1668 match self.ordering {
1669 TournamentOrdering::CostThenScore => passing()
1670 .min_by_key(|candidate| (candidate.cost, std::cmp::Reverse(candidate.score))),
1671 TournamentOrdering::ScoreThenCost => passing()
1672 .min_by_key(|candidate| (std::cmp::Reverse(candidate.score), candidate.cost)),
1673 }
1674 .cloned()
1675 }
1676 }
1677
1678 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1679 pub struct ParetoFrontier {
1680 #[serde(default = "default_frontier_limit")]
1681 pub max_items: usize,
1682 }
1683
1684 impl Default for ParetoFrontier {
1685 fn default() -> Self {
1686 Self {
1687 max_items: default_frontier_limit(),
1688 }
1689 }
1690 }
1691
1692 impl ParetoFrontier {
1693 pub fn select(&self, candidates: &[BranchCandidate]) -> Vec<BranchCandidate> {
1694 let mut frontier: Vec<_> = candidates
1695 .iter()
1696 .filter(|candidate| candidate.status == WorkflowRunStatus::Succeeded)
1697 .filter(|candidate| {
1698 !candidates.iter().any(|other| {
1699 other.status == WorkflowRunStatus::Succeeded
1700 && other.score >= candidate.score
1701 && other.cost <= candidate.cost
1702 && (other.score > candidate.score || other.cost < candidate.cost)
1703 })
1704 })
1705 .cloned()
1706 .collect();
1707 frontier.sort_by_key(|candidate| (std::cmp::Reverse(candidate.score), candidate.cost));
1708 frontier.truncate(self.max_items.max(1));
1709 frontier
1710 }
1711 }
1712
1713 #[derive(Debug, Clone, PartialEq, Eq, Error)]
1714 pub enum WorkflowExecutionError {
1715 #[error("{kind} node id must not be empty")]
1716 EmptyNodeId { kind: &'static str },
1717 #[error("leaf `{leaf}` prompt must not be empty")]
1718 EmptyLeafPrompt { leaf: String },
1719 #[error(
1720 "leaf `{leaf}` profile `{profile}` must be a non-empty token without whitespace, quotes, or `=`"
1721 )]
1722 InvalidLeafProfile { leaf: String, profile: String },
1723 #[error(
1724 "leaf `{leaf}` role `{role}` must be a non-empty token without whitespace, quotes, or `=`"
1725 )]
1726 InvalidLeafRole { leaf: String, role: String },
1727 #[error("duplicate workflow node `{node}`")]
1728 DuplicateNodeId { node: String },
1729 #[error("workflow node `{node}` has unknown {field} reference `{reference}`")]
1730 UnknownNodeReference {
1731 node: String,
1732 field: &'static str,
1733 reference: String,
1734 },
1735 }
1736
1737 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1738 pub struct WorkflowFleetLimits {
1739 pub max_total_agents: usize,
1740 pub max_depth: usize,
1741 }
1742
1743 impl Default for WorkflowFleetLimits {
1744 fn default() -> Self {
1745 Self {
1746 max_total_agents: DEFAULT_FLEET_WORKFLOW_MAX_AGENTS,
1747 max_depth: DEFAULT_FLEET_WORKFLOW_MAX_DEPTH,
1748 }
1749 }
1750 }
1751
1752 #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1753 pub struct WorkflowFleetShape {
1754 pub total_agents: usize,
1755 pub max_depth: usize,
1756 }
1757
1758 impl WorkflowFleetShape {
1759 fn add(self, other: Self) -> Self {
1760 Self {
1761 total_agents: self.total_agents.saturating_add(other.total_agents),
1762 max_depth: self.max_depth.max(other.max_depth),
1763 }
1764 }
1765
1766 fn repeat(self, times: usize) -> Self {
1767 Self {
1768 total_agents: self.total_agents.saturating_mul(times),
1769 max_depth: self.max_depth,
1770 }
1771 }
1772 }
1773
1774 #[derive(Debug, Clone, PartialEq, Eq, Error)]
1775 pub enum WorkflowFleetLimitError {
1776 #[error("workflow IR is invalid for Fleet: {source}")]
1777 InvalidWorkflow {
1778 #[from]
1779 source: WorkflowExecutionError,
1780 },
1781 #[error(
1782 "workflow would launch {total_agents} agents; Fleet Workflow limit is {max_total_agents}"
1783 )]
1784 TooManyAgents {
1785 total_agents: usize,
1786 max_total_agents: usize,
1787 },
1788 #[error("workflow reaches recursion depth {depth}; Fleet Workflow limit is {max_depth}")]
1789 RecursionTooDeep { depth: usize, max_depth: usize },
1790 #[error("expand node `{node}` must declare max_children before Fleet launch")]
1791 UnboundedExpand { node: String },
1792 #[error("expand node `{node}` must include a template before Fleet launch")]
1793 MissingExpandTemplate { node: String },
1794 #[error("loop_until node `{node}` must declare max_iterations before Fleet launch")]
1795 UnboundedLoop { node: String },
1796 }
1797
1798 fn estimate_fleet_shape(
1799 nodes: &[WorkflowNode],
1800 ) -> Result<WorkflowFleetShape, WorkflowFleetLimitError> {
1801 estimate_fleet_shape_at_depth(nodes, 1)
1802 }
1803
1804 fn estimate_fleet_shape_at_depth(
1805 nodes: &[WorkflowNode],
1806 depth: usize,
1807 ) -> Result<WorkflowFleetShape, WorkflowFleetLimitError> {
1808 nodes
1809 .iter()
1810 .try_fold(WorkflowFleetShape::default(), |shape, node| {
1811 Ok(shape.add(estimate_node_fleet_shape(node, depth)?))
1812 })
1813 }
1814
1815 fn estimate_node_fleet_shape(
1816 node: &WorkflowNode,
1817 depth: usize,
1818 ) -> Result<WorkflowFleetShape, WorkflowFleetLimitError> {
1819 match node {
1820 WorkflowNode::Leaf(_) => Ok(WorkflowFleetShape {
1821 total_agents: 1,
1822 max_depth: depth,
1823 }),
1824 WorkflowNode::BranchSet(spec) => estimate_fleet_shape_at_depth(&spec.children, depth + 1),
1825 WorkflowNode::Sequence(spec) => estimate_fleet_shape_at_depth(&spec.children, depth),
1826 WorkflowNode::Reduce(_) | WorkflowNode::TeacherReview(_) => Ok(WorkflowFleetShape {
1827 total_agents: 0,
1828 max_depth: 0,
1829 }),
1830 WorkflowNode::LoopUntil(spec) => {
1831 let iterations =
1832 spec.max_iterations
1833 .ok_or_else(|| WorkflowFleetLimitError::UnboundedLoop {
1834 node: spec.id.clone(),
1835 })? as usize;
1836 Ok(estimate_fleet_shape_at_depth(&spec.children, depth)?.repeat(iterations.max(1)))
1837 }
1838 WorkflowNode::Cond(spec) => Ok(estimate_fleet_shape_at_depth(&spec.then_nodes, depth)?
1839 .add(estimate_fleet_shape_at_depth(&spec.else_nodes, depth)?)),
1840 WorkflowNode::Expand(spec) => {
1841 let max_children =
1842 spec.max_children
1843 .ok_or_else(|| WorkflowFleetLimitError::UnboundedExpand {
1844 node: spec.id.clone(),
1845 })?;
1846 let template = spec.template.as_deref().ok_or_else(|| {
1847 WorkflowFleetLimitError::MissingExpandTemplate {
1848 node: spec.id.clone(),
1849 }
1850 })?;
1851 validate_workflow_node_shapes(std::slice::from_ref(template))
1852 .map_err(|source| WorkflowFleetLimitError::InvalidWorkflow { source })?;
1853 Ok(estimate_node_fleet_shape(template, depth)?.repeat(max_children))
1854 }
1855 }
1856 }
1857
1858 fn default_frontier_limit() -> usize {
1859 8
1860 }
1861
1862 fn node_id(node: &WorkflowNode) -> String {
1863 match node {
1864 WorkflowNode::BranchSet(spec) => spec.id.clone(),
1865 WorkflowNode::Leaf(spec) => spec.id.clone(),
1866 WorkflowNode::Sequence(spec) => spec.id.clone(),
1867 WorkflowNode::Reduce(spec) => spec.id.clone(),
1868 WorkflowNode::TeacherReview(spec) => spec.id.clone(),
1869 WorkflowNode::LoopUntil(spec) => spec.id.clone(),
1870 WorkflowNode::Cond(spec) => spec.id.clone(),
1871 WorkflowNode::Expand(spec) => spec.id.clone(),
1872 }
1873 }
1874
1875 pub(crate) fn validate_workflow_nodes(
1876 nodes: &[WorkflowNode],
1877 ) -> Result<(), WorkflowExecutionError> {
1878 let mut seen = BTreeSet::new();
1879 validate_workflow_nodes_inner(nodes, &mut seen)?;
1880 validate_workflow_references(nodes, &seen)
1881 }
1882
1883 pub(crate) fn validate_workflow_node_shapes(
1884 nodes: &[WorkflowNode],
1885 ) -> Result<(), WorkflowExecutionError> {
1886 let mut seen = BTreeSet::new();
1887 validate_workflow_nodes_inner(nodes, &mut seen)
1888 }
1889
1890 fn validate_workflow_nodes_inner(
1891 nodes: &[WorkflowNode],
1892 seen: &mut BTreeSet<String>,
1893 ) -> Result<(), WorkflowExecutionError> {
1894 for node in nodes {
1895 let id = node_id(node);
1896 let kind = control_kind_name(node);
1897 if id.trim().is_empty() {
1898 return Err(WorkflowExecutionError::EmptyNodeId { kind });
1899 }
1900 if !seen.insert(id.clone()) {
1901 return Err(WorkflowExecutionError::DuplicateNodeId { node: id });
1902 }
1903 match node {
1904 WorkflowNode::BranchSet(spec) => validate_workflow_nodes_inner(&spec.children, seen)?,
1905 WorkflowNode::Leaf(spec) => {
1906 if spec.prompt.trim().is_empty() {
1907 return Err(WorkflowExecutionError::EmptyLeafPrompt {
1908 leaf: spec.id.clone(),
1909 });
1910 }
1911 if let Some(role) = spec.role.as_deref() {
1912 validate_leaf_role(&spec.id, role)?;
1913 }
1914 if let Some(profile) = spec.profile.as_deref() {
1915 validate_leaf_profile(&spec.id, profile)?;
1916 }
1917 }
1918 WorkflowNode::Sequence(spec) => validate_workflow_nodes_inner(&spec.children, seen)?,
1919 WorkflowNode::LoopUntil(spec) => validate_workflow_nodes_inner(&spec.children, seen)?,
1920 WorkflowNode::Cond(spec) => {
1921 validate_workflow_nodes_inner(&spec.then_nodes, seen)?;
1922 validate_workflow_nodes_inner(&spec.else_nodes, seen)?;
1923 }
1924 WorkflowNode::Reduce(_) | WorkflowNode::TeacherReview(_) | WorkflowNode::Expand(_) => {}
1925 }
1926 }
1927 Ok(())
1928 }
1929
1930 fn validate_workflow_references(
1931 nodes: &[WorkflowNode],
1932 known_ids: &BTreeSet<String>,
1933 ) -> Result<(), WorkflowExecutionError> {
1934 for node in nodes {
1935 match node {
1936 WorkflowNode::BranchSet(spec) => {
1937 validate_workflow_references(&spec.children, known_ids)?;
1938 }
1939 WorkflowNode::Leaf(spec) => {
1940 validate_known_references(
1941 spec.id.as_str(),
1942 "depends_on_results",
1943 &spec.depends_on_results,
1944 known_ids,
1945 )?;
1946 }
1947 WorkflowNode::Sequence(spec) => {
1948 validate_workflow_references(&spec.children, known_ids)?;
1949 }
1950 WorkflowNode::Reduce(spec) => {
1951 validate_known_references(spec.id.as_str(), "inputs", &spec.inputs, known_ids)?;
1952 }
1953 WorkflowNode::TeacherReview(spec) => {
1954 validate_known_references(
1955 spec.id.as_str(),
1956 "candidates",
1957 &spec.candidates,
1958 known_ids,
1959 )?;
1960 }
1961 WorkflowNode::LoopUntil(spec) => {
1962 validate_workflow_references(&spec.children, known_ids)?;
1963 }
1964 WorkflowNode::Cond(spec) => {
1965 validate_workflow_references(&spec.then_nodes, known_ids)?;
1966 validate_workflow_references(&spec.else_nodes, known_ids)?;
1967 }
1968 WorkflowNode::Expand(_) => {}
1969 }
1970 }
1971 Ok(())
1972 }
1973
1974 // Token rule only. Roster membership is resolved by the dispatcher (tui crate)
1975 // at spawn time; this crate never sees the saved Fleet roster.
1976 fn validate_leaf_profile(leaf: &str, profile: &str) -> Result<(), WorkflowExecutionError> {
1977 let invalid = profile.is_empty()
1978 || profile
1979 .chars()
1980 .any(|ch| ch.is_whitespace() || matches!(ch, '"' | '\'' | '`' | '='));
1981 if invalid {
1982 return Err(WorkflowExecutionError::InvalidLeafProfile {
1983 leaf: leaf.to_string(),
1984 profile: profile.to_string(),
1985 });
1986 }
1987 Ok(())
1988 }
1989
1990 fn validate_leaf_role(leaf: &str, role: &str) -> Result<(), WorkflowExecutionError> {
1991 if validate_role_token(role).is_err() {
1992 return Err(WorkflowExecutionError::InvalidLeafRole {
1993 leaf: leaf.to_string(),
1994 role: role.to_string(),
1995 });
1996 }
1997 Ok(())
1998 }
1999
2000 fn validate_known_references(
2001 node: &str,
2002 field: &'static str,
2003 references: &[String],
2004 known_ids: &BTreeSet<String>,
2005 ) -> Result<(), WorkflowExecutionError> {
2006 for reference in references {
2007 if !known_ids.contains(reference) {
2008 return Err(WorkflowExecutionError::UnknownNodeReference {
2009 node: node.to_string(),
2010 field,
2011 reference: reference.clone(),
2012 });
2013 }
2014 }
2015 Ok(())
2016 }
2017
2018 fn control_kind_name(node: &WorkflowNode) -> &'static str {
2019 match node {
2020 WorkflowNode::BranchSet(_) => "branch_set",
2021 WorkflowNode::Leaf(_) => "leaf",
2022 WorkflowNode::Sequence(_) => "sequence",
2023 WorkflowNode::Reduce(_) => "reduce",
2024 WorkflowNode::TeacherReview(_) => "teacher_review",
2025 WorkflowNode::LoopUntil(_) => "loop_until",
2026 WorkflowNode::Cond(_) => "cond",
2027 WorkflowNode::Expand(_) => "expand",
2028 }
2029 }
2030
2031 #[derive(Debug, Clone, PartialEq, Eq, Error)]
2032 pub enum WorkflowValidationError {
2033 #[error("{field} must not be empty")]
2034 EmptyField { field: &'static str },
2035 #[error("workflow must contain at least one phase")]
2036 EmptyWorkflow,
2037 #[error("phase `{phase}` must contain at least one task")]
2038 EmptyPhase { phase: String },
2039 #[error("max_concurrent must be between 1 and 20, got {value}")]
2040 InvalidMaxConcurrent { value: u8 },
2041 #[error("duplicate workflow phase `{phase}`")]
2042 DuplicatePhase { phase: String },
2043 #[error("duplicate workflow task `{task}`")]
2044 DuplicateTask { task: String },
2045 #[error("phase `{phase}` has invalid dependency `{dependency}`")]
2046 InvalidPhaseDependency { phase: String, dependency: String },
2047 #[error("phase dependency cycle includes `{phase}`")]
2048 PhaseDependencyCycle { phase: String },
2049 #[error("task `{task}` has invalid result dependency `{dependency}`")]
2050 InvalidTaskResultDependency { task: String, dependency: String },
2051 #[error(
2052 "task `{task}` depends on result `{dependency}` from unavailable phase `{dependency_phase}` while running in `{task_phase}`"
2053 )]
2054 UnavailableTaskResultDependency {
2055 task: String,
2056 dependency: String,
2057 dependency_phase: String,
2058 task_phase: String,
2059 },
2060 #[error("parallel read-write task `{task}` must declare a file_scope")]
2061 MissingParallelWriteScope { task: String },
2062 #[error("parallel read-write tasks `{left}` and `{right}` have overlapping file scopes")]
2063 OverlappingParallelWriteScope { left: String, right: String },
2064 }
2065
2066 fn default_max_concurrent() -> u8 {
2067 4
2068 }
2069
2070 fn validate_non_empty(field: &'static str, value: &str) -> Result<(), WorkflowValidationError> {
2071 if value.trim().is_empty() {
2072 return Err(WorkflowValidationError::EmptyField { field });
2073 }
2074 Ok(())
2075 }
2076
2077 fn ordered_phases(
2078 config: &WorkflowConfig,
2079 phase_indices: &BTreeMap<String, usize>,
2080 ) -> Result<Vec<String>, WorkflowValidationError> {
2081 let mut visiting = BTreeSet::new();
2082 let mut visited = BTreeSet::new();
2083 let mut ordered = Vec::with_capacity(config.phases.len());
2084
2085 for phase in &config.phases {
2086 visit_phase(
2087 &phase.name,
2088 config,
2089 phase_indices,
2090 &mut visiting,
2091 &mut visited,
2092 &mut ordered,
2093 )?;
2094 }
2095
2096 Ok(ordered)
2097 }
2098
2099 fn visit_phase(
2100 phase_name: &str,
2101 config: &WorkflowConfig,
2102 phase_indices: &BTreeMap<String, usize>,
2103 visiting: &mut BTreeSet<String>,
2104 visited: &mut BTreeSet<String>,
2105 ordered: &mut Vec<String>,
2106 ) -> Result<(), WorkflowValidationError> {
2107 if visited.contains(phase_name) {
2108 return Ok(());
2109 }
2110 if !visiting.insert(phase_name.to_string()) {
2111 return Err(WorkflowValidationError::PhaseDependencyCycle {
2112 phase: phase_name.to_string(),
2113 });
2114 }
2115
2116 let phase = &config.phases[phase_indices[phase_name]];
2117 for dependency in &phase.depends_on {
2118 visit_phase(
2119 dependency,
2120 config,
2121 phase_indices,
2122 visiting,
2123 visited,
2124 ordered,
2125 )?;
2126 }
2127
2128 visiting.remove(phase_name);
2129 visited.insert(phase_name.to_string());
2130 ordered.push(phase_name.to_string());
2131 Ok(())
2132 }
2133
2134 fn validate_parallel_write_scope(phase: &Phase) -> Result<(), WorkflowValidationError> {
2135 if !phase.parallel {
2136 return Ok(());
2137 }
2138
2139 let write_tasks: Vec<_> = phase
2140 .tasks
2141 .iter()
2142 .filter(|task| task.mode == TaskMode::ReadWrite)
2143 .collect();
2144
2145 for task in &write_tasks {
2146 if task.file_scope.is_empty() {
2147 return Err(WorkflowValidationError::MissingParallelWriteScope {
2148 task: task.id.clone(),
2149 });
2150 }
2151 }
2152
2153 for (left_index, left) in write_tasks.iter().enumerate() {
2154 for right in write_tasks.iter().skip(left_index + 1) {
2155 if scopes_overlap(&left.file_scope, &right.file_scope) {
2156 return Err(WorkflowValidationError::OverlappingParallelWriteScope {
2157 left: left.id.clone(),
2158 right: right.id.clone(),
2159 });
2160 }
2161 }
2162 }
2163
2164 Ok(())
2165 }
2166
2167 pub fn scopes_overlap(left: &[String], right: &[String]) -> bool {
2168 left.iter().any(|left_scope| {
2169 right
2170 .iter()
2171 .any(|right_scope| scope_overlaps(left_scope, right_scope))
2172 })
2173 }
2174
2175 fn scope_overlaps(left: &str, right: &str) -> bool {
2176 let left = normalize_file_scope_root(left);
2177 let right = normalize_file_scope_root(right);
2178
2179 if left == right || left == "." || right == "." {
2180 return true;
2181 }
2182
2183 if left.contains('*') || right.contains('*') {
2184 return glob_prefix(&left) == glob_prefix(&right);
2185 }
2186
2187 let left_path = Path::new(&left);
2188 let right_path = Path::new(&right);
2189 left_path.starts_with(right_path) || right_path.starts_with(left_path)
2190 }
2191
2192 /// Normalize the suffix-glob spelling accepted by declarative workflow
2193 /// `file_scope` into the concrete directory root enforced at runtime.
2194 #[must_use]
2195 pub fn normalize_file_scope_root(scope: &str) -> String {
2196 let trimmed = scope.trim().trim_start_matches("./").trim_end_matches('/');
2197 trimmed
2198 .strip_suffix("/**")
2199 .or_else(|| trimmed.strip_suffix("/*"))
2200 .unwrap_or(trimmed)
2201 .to_string()
2202 }
2203
2204 fn glob_prefix(scope: &str) -> String {
2205 scope
2206 .split('*')
2207 .next()
2208 .unwrap_or(scope)
2209 .trim_end_matches('/')
2210 .to_string()
2211 }
2212
2213 #[cfg(test)]
2214 mod tests {
2215 use super::*;
2216
2217 fn task(id: &str) -> Task {
2218 Task {
2219 id: id.to_string(),
2220 prompt: format!("run {id}"),
2221 agent_type: AgentType::General,
2222 mode: TaskMode::ReadOnly,
2223 isolation: IsolationMode::Shared,
2224 file_scope: Vec::new(),
2225 depends_on_results: Vec::new(),
2226 max_steps: None,
2227 timeout_secs: None,
2228 }
2229 }
2230
2231 fn config(phases: Vec<Phase>) -> WorkflowConfig {
2232 WorkflowConfig {
2233 goal: "cache-change".to_string(),
2234 max_concurrent: 4,
2235 description: None,
2236 phases,
2237 }
2238 }
2239
2240 fn phase(name: &str, depends_on: &[&str], tasks: Vec<Task>) -> Phase {
2241 Phase {
2242 name: name.to_string(),
2243 description: None,
2244 depends_on: depends_on.iter().map(|value| value.to_string()).collect(),
2245 parallel: false,
2246 on_failure: FailurePolicy::SkipContinue,
2247 tasks,
2248 }
2249 }
2250
2251 fn leaf_node(id: &str) -> WorkflowNode {
2252 WorkflowNode::Leaf(LeafSpec {
2253 id: id.to_string(),
2254 prompt: format!("run {id}"),
2255 agent_type: AgentType::General,
2256 role: None,
2257 profile: None,
2258 mode: TaskMode::ReadOnly,
2259 isolation: IsolationMode::Shared,
2260 file_scope: Vec::new(),
2261 cwd: None,
2262 depends_on_results: Vec::new(),
2263 budget: BudgetSpec::default(),
2264 permissions: PermissionSpec::default(),
2265 model_policy: ModelPolicy::default(),
2266 })
2267 }
2268
2269 fn leaf_node_with_budget(id: &str, budget: BudgetSpec) -> WorkflowNode {
2270 WorkflowNode::Leaf(LeafSpec {
2271 id: id.to_string(),
2272 prompt: format!("run {id}"),
2273 agent_type: AgentType::General,
2274 role: None,
2275 profile: None,
2276 mode: TaskMode::ReadOnly,
2277 isolation: IsolationMode::Shared,
2278 file_scope: Vec::new(),
2279 cwd: None,
2280 depends_on_results: Vec::new(),
2281 budget,
2282 permissions: PermissionSpec::default(),
2283 model_policy: ModelPolicy::default(),
2284 })
2285 }
2286
2287 fn invalid_leaf_node(id: &str) -> WorkflowNode {
2288 WorkflowNode::Leaf(LeafSpec {
2289 id: id.to_string(),
2290 prompt: " ".to_string(),
2291 agent_type: AgentType::General,
2292 role: None,
2293 profile: None,
2294 mode: TaskMode::ReadOnly,
2295 isolation: IsolationMode::Shared,
2296 file_scope: Vec::new(),
2297 cwd: None,
2298 depends_on_results: Vec::new(),
2299 budget: BudgetSpec::default(),
2300 permissions: PermissionSpec::default(),
2301 model_policy: ModelPolicy::default(),
2302 })
2303 }
2304
2305 fn workflow_spec(nodes: Vec<WorkflowNode>) -> WorkflowSpec {
2306 WorkflowSpec {
2307 id: Some("mock-workflow".to_string()),
2308 goal: "prove mock executor control flow".to_string(),
2309 description: None,
2310 budget: BudgetSpec::default(),
2311 permissions: PermissionSpec::default(),
2312 model_policy: ModelPolicy::default(),
2313 promotion_policy: PromotionPolicy::default(),
2314 gates: Vec::new(),
2315 nodes,
2316 }
2317 }
2318
2319 fn control_result<'a>(
2320 execution: &'a WorkflowExecution,
2321 node_id: &str,
2322 ) -> &'a ControlNodeResult {
2323 execution
2324 .control_node_results
2325 .iter()
2326 .find(|result| result.node_id == node_id)
2327 .expect("control node result should exist")
2328 }
2329
2330 fn candidate(
2331 branch_id: &str,
2332 status: WorkflowRunStatus,
2333 score: u32,
2334 cost: u64,
2335 diversity_key: &str,
2336 ) -> BranchCandidate {
2337 BranchCandidate {
2338 branch_id: branch_id.to_string(),
2339 status,
2340 score,
2341 cost,
2342 diversity_key: Some(diversity_key.to_string()),
2343 }
2344 }
2345
2346 #[test]
2347 fn independent_phases_preserve_declaration_order() {
2348 let workflow = config(vec![
2349 phase("discover", &[], vec![task("scan")]),
2350 phase("report", &[], vec![task("summarize")]),
2351 ]);
2352
2353 let plan = workflow.compile().expect("workflow should compile");
2354
2355 assert_eq!(
2356 plan.phase_names().collect::<Vec<_>>(),
2357 vec!["discover", "report"]
2358 );
2359 }
2360
2361 #[test]
2362 fn dependencies_override_declaration_order_deterministically() {
2363 let workflow = config(vec![
2364 phase("review", &["implement"], vec![task("review-results")]),
2365 phase("discover", &[], vec![task("scan")]),
2366 phase("implement", &["discover"], vec![task("patch")]),
2367 phase("report", &["review"], vec![task("summarize")]),
2368 ]);
2369
2370 let plan = workflow.compile().expect("workflow should compile");
2371
2372 assert_eq!(
2373 plan.phase_names().collect::<Vec<_>>(),
2374 vec!["discover", "implement", "review", "report"]
2375 );
2376 }
2377
2378 #[test]
2379 fn rejects_empty_workflow() {
2380 let err = config(Vec::new())
2381 .validate()
2382 .expect_err("empty workflow should fail");
2383
2384 assert_eq!(err, WorkflowValidationError::EmptyWorkflow);
2385 }
2386
2387 #[test]
2388 fn rejects_empty_phase() {
2389 let err = config(vec![phase("empty", &[], Vec::new())])
2390 .validate()
2391 .expect_err("empty phase should fail");
2392
2393 assert_eq!(
2394 err,
2395 WorkflowValidationError::EmptyPhase {
2396 phase: "empty".to_string()
2397 }
2398 );
2399 }
2400
2401 #[test]
2402 fn rejects_invalid_max_concurrent() {
2403 let mut workflow = config(vec![phase("discover", &[], vec![task("scan")])]);
2404 workflow.max_concurrent = 0;
2405
2406 let err = workflow
2407 .validate()
2408 .expect_err("zero concurrency should fail");
2409
2410 assert_eq!(
2411 err,
2412 WorkflowValidationError::InvalidMaxConcurrent { value: 0 }
2413 );
2414 }
2415
2416 #[test]
2417 fn rejects_duplicate_phase_names() {
2418 let err = config(vec![
2419 phase("discover", &[], vec![task("scan")]),
2420 phase("discover", &[], vec![task("scan-again")]),
2421 ])
2422 .validate()
2423 .expect_err("duplicate phase should fail");
2424
2425 assert!(matches!(
2426 err,
2427 WorkflowValidationError::DuplicatePhase { .. }
2428 ));
2429 }
2430
2431 #[test]
2432 fn rejects_duplicate_task_ids() {
2433 let err = config(vec![
2434 phase("discover", &[], vec![task("scan")]),
2435 phase("report", &[], vec![task("scan")]),
2436 ])
2437 .validate()
2438 .expect_err("duplicate task should fail");
2439
2440 assert!(matches!(err, WorkflowValidationError::DuplicateTask { .. }));
2441 }
2442
2443 #[test]
2444 fn rejects_unknown_phase_dependency() {
2445 let err = config(vec![phase("report", &["missing"], vec![task("summarize")])])
2446 .validate()
2447 .expect_err("unknown dependency should fail");
2448
2449 assert!(matches!(
2450 err,
2451 WorkflowValidationError::InvalidPhaseDependency { .. }
2452 ));
2453 }
2454
2455 #[test]
2456 fn rejects_phase_dependency_cycles() {
2457 let workflow = config(vec![
2458 phase("a", &["b"], vec![task("a-task")]),
2459 phase("b", &["a"], vec![task("b-task")]),
2460 ]);
2461
2462 let err = workflow.validate().expect_err("cycle should fail");
2463
2464 assert!(matches!(
2465 err,
2466 WorkflowValidationError::PhaseDependencyCycle { .. }
2467 ));
2468 }
2469
2470 #[test]
2471 fn rejects_task_result_dependency_from_same_parallel_phase() {
2472 let mut first = task("first");
2473 first.depends_on_results.push("second".to_string());
2474 let mut parallel = phase("parallel", &[], vec![first, task("second")]);
2475 parallel.parallel = true;
2476
2477 let err = config(vec![parallel])
2478 .validate()
2479 .expect_err("same-phase result dependency should fail");
2480
2481 assert!(matches!(
2482 err,
2483 WorkflowValidationError::UnavailableTaskResultDependency { .. }
2484 ));
2485 }
2486
2487 #[test]
2488 fn rejects_task_result_dependency_from_later_phase() {
2489 let mut summarize = task("summarize");
2490 summarize.depends_on_results.push("scan".to_string());
2491 let workflow = config(vec![
2492 phase("report", &[], vec![summarize]),
2493 phase("discover", &[], vec![task("scan")]),
2494 ]);
2495
2496 let err = workflow
2497 .validate()
2498 .expect_err("later-phase result dependency should fail");
2499
2500 assert!(matches!(
2501 err,
2502 WorkflowValidationError::UnavailableTaskResultDependency { .. }
2503 ));
2504 }
2505
2506 #[test]
2507 fn allows_task_result_dependency_from_earlier_phase() {
2508 let upstream = phase("discover", &[], vec![task("scan")]);
2509 let mut summarize = task("summarize");
2510 summarize.depends_on_results.push("scan".to_string());
2511 let downstream = phase("report", &["discover"], vec![summarize]);
2512
2513 config(vec![upstream, downstream])
2514 .validate()
2515 .expect("earlier-phase result should be available");
2516 }
2517
2518 #[test]
2519 fn rejects_parallel_read_write_without_file_scope() {
2520 let mut write = task("write");
2521 write.mode = TaskMode::ReadWrite;
2522 let mut parallel = phase("parallel", &[], vec![write]);
2523 parallel.parallel = true;
2524
2525 let err = config(vec![parallel])
2526 .validate()
2527 .expect_err("write task needs a scope");
2528
2529 assert!(matches!(
2530 err,
2531 WorkflowValidationError::MissingParallelWriteScope { .. }
2532 ));
2533 }
2534
2535 #[test]
2536 fn detects_overlapping_parallel_write_scopes_with_path_boundaries() {
2537 let mut left = task("auth");
2538 left.mode = TaskMode::ReadWrite;
2539 left.file_scope = vec!["src/auth/**".to_string()];
2540 let mut right = task("auth-login");
2541 right.mode = TaskMode::ReadWrite;
2542 right.file_scope = vec!["src/auth/login.rs".to_string()];
2543 let mut parallel = phase("parallel", &[], vec![left, right]);
2544 parallel.parallel = true;
2545
2546 let err = config(vec![parallel])
2547 .validate()
2548 .expect_err("nested scopes should overlap");
2549
2550 assert!(matches!(
2551 err,
2552 WorkflowValidationError::OverlappingParallelWriteScope { .. }
2553 ));
2554 }
2555
2556 #[test]
2557 fn does_not_confuse_path_prefixes_for_overlapping_scopes() {
2558 let mut left = task("auth");
2559 left.mode = TaskMode::ReadWrite;
2560 left.file_scope = vec!["src/auth/**".to_string()];
2561 let mut right = task("auth-admin");
2562 right.mode = TaskMode::ReadWrite;
2563 right.file_scope = vec!["src/auth_admin/**".to_string()];
2564 let mut parallel = phase("parallel", &[], vec![left, right]);
2565 parallel.parallel = true;
2566
2567 config(vec![parallel])
2568 .validate()
2569 .expect("component boundary scopes should not overlap");
2570 }
2571
2572 #[test]
2573 fn json_roundtrip_keeps_snake_case_enum_names() {
2574 let mut task = task("patch");
2575 task.agent_type = AgentType::Implementer;
2576 task.mode = TaskMode::ReadWrite;
2577 task.isolation = IsolationMode::Worktree;
2578 task.file_scope = vec!["src/auth/**".to_string()];
2579 let mut parallel = phase("implement", &[], vec![task]);
2580 parallel.parallel = true;
2581 parallel.on_failure = FailurePolicy::Abort;
2582 let workflow = config(vec![parallel]);
2583
2584 let json = serde_json::to_string(&workflow).expect("serialize workflow");
2585
2586 assert!(json.contains("\"agent_type\":\"implement\""));
2587 assert!(json.contains("\"mode\":\"read_write\""));
2588 assert!(json.contains("\"isolation\":\"worktree\""));
2589 assert!(json.contains("\"on_failure\":\"abort\""));
2590
2591 let parsed: WorkflowConfig = serde_json::from_str(&json).expect("parse workflow");
2592 assert_eq!(parsed, workflow);
2593 }
2594
2595 #[test]
2596 fn json_accepts_pre_rename_agent_type_spellings_as_aliases() {
2597 for (legacy, expected) in [
2598 ("implementer", AgentType::Implementer),
2599 ("builder", AgentType::Implementer),
2600 ("verifier", AgentType::Verifier),
2601 ("verify", AgentType::Verifier),
2602 ("scout", AgentType::Explore),
2603 ("review", AgentType::Review),
2604 ] {
2605 let parsed: AgentType =
2606 serde_json::from_value(serde_json::json!(legacy)).expect("legacy alias must parse");
2607 assert_eq!(parsed, expected, "alias {legacy}");
2608 }
2609 // Canonical spellings parse and round-trip back as themselves.
2610 for (canonical, expected) in [
2611 ("implement", AgentType::Implementer),
2612 ("test", AgentType::Verifier),
2613 ("explore", AgentType::Explore),
2614 ("general", AgentType::General),
2615 ] {
2616 let parsed: AgentType =
2617 serde_json::from_value(serde_json::json!(canonical)).expect("canonical parses");
2618 assert_eq!(parsed, expected);
2619 assert_eq!(
2620 serde_json::to_value(parsed).expect("serialize"),
2621 serde_json::json!(canonical),
2622 "canonical spelling must serialize as itself"
2623 );
2624 }
2625 }
2626
2627 #[test]
2628 fn isolation_auto_defaults_parallel_write_to_worktree() {
2629 assert_eq!(IsolationMode::default(), IsolationMode::Auto);
2630 assert_eq!(
2631 IsolationMode::Auto.resolve(/* parallel_write */ true),
2632 IsolationMode::Worktree
2633 );
2634 assert_eq!(
2635 IsolationMode::Auto.resolve(/* parallel_write */ false),
2636 IsolationMode::Shared
2637 );
2638 // Explicit shared is the approved same-worktree override.
2639 assert_eq!(
2640 IsolationMode::Shared.resolve(/* parallel_write */ true),
2641 IsolationMode::Shared
2642 );
2643 assert!(IsolationMode::Worktree.wants_worktree(false));
2644 assert!(!IsolationMode::Shared.wants_worktree(true));
2645 }
2646
2647 #[test]
2648 fn leaf_write_capable_and_worktree_defaults() {
2649 let read_only = LeafSpec {
2650 id: "ro".to_string(),
2651 prompt: "inspect".to_string(),
2652 agent_type: AgentType::Explore,
2653 role: None,
2654 profile: None,
2655 mode: TaskMode::ReadOnly,
2656 isolation: IsolationMode::Auto,
2657 file_scope: Vec::new(),
2658 cwd: None,
2659 depends_on_results: Vec::new(),
2660 budget: BudgetSpec::default(),
2661 permissions: PermissionSpec::default(),
2662 model_policy: ModelPolicy::default(),
2663 };
2664 assert!(!leaf_is_write_capable(&read_only));
2665 assert!(!leaf_wants_worktree(&read_only, true));
2666
2667 let mut read_only_implementer = read_only.clone();
2668 read_only_implementer.id = "ro-implementer".to_string();
2669 read_only_implementer.agent_type = AgentType::Implementer;
2670 assert!(
2671 !leaf_is_write_capable(&read_only_implementer),
2672 "role identity must not grant write authority"
2673 );
2674 assert!(
2675 !leaf_wants_worktree(&read_only_implementer, true),
2676 "parallel read-only implementers stay shared under auto isolation"
2677 );
2678
2679 let mut write = read_only.clone();
2680 write.id = "rw".to_string();
2681 write.mode = TaskMode::ReadWrite;
2682 write.agent_type = AgentType::Implementer;
2683 assert!(leaf_is_write_capable(&write));
2684 // Parallel write-capable + Auto → worktree by default.
2685 assert!(leaf_wants_worktree(&write, true));
2686 // Sequential write-capable stays shared unless isolation is worktree.
2687 assert!(!leaf_wants_worktree(&write, false));
2688
2689 write.isolation = IsolationMode::Shared;
2690 assert!(
2691 !leaf_wants_worktree(&write, true),
2692 "explicit shared is the same-worktree override"
2693 );
2694
2695 write.isolation = IsolationMode::Worktree;
2696 assert!(leaf_wants_worktree(&write, true));
2697 assert!(leaf_wants_worktree(&write, false));
2698 }
2699
2700 #[test]
2701 fn workflow_ir_roundtrip() {
2702 let discover_leaf = LeafSpec {
2703 id: "scan-readme".to_string(),
2704 prompt: "Inspect README setup gaps".to_string(),
2705 agent_type: AgentType::Explore,
2706 role: None,
2707 profile: Some("scout".to_string()),
2708 mode: TaskMode::ReadOnly,
2709 isolation: IsolationMode::Shared,
2710 file_scope: vec!["README.md".to_string()],
2711 cwd: None,
2712 depends_on_results: Vec::new(),
2713 budget: BudgetSpec {
2714 max_steps: Some(8),
2715 timeout_secs: Some(300),
2716 max_parallel: None,
2717 max_tokens: None,
2718 },
2719 permissions: PermissionSpec::default(),
2720 model_policy: ModelPolicy {
2721 provider: Some("openai".to_string()),
2722 model: Some("gpt-5.4".to_string()),
2723 fallback_models: Vec::new(),
2724 },
2725 };
2726 let workflow = WorkflowSpec {
2727 id: Some("v090-readme-check".to_string()),
2728 goal: "tighten setup docs".to_string(),
2729 description: Some("metadata-only typed Workflow IR".to_string()),
2730 budget: BudgetSpec {
2731 max_steps: Some(30),
2732 timeout_secs: Some(1_800),
2733 max_parallel: Some(2),
2734 max_tokens: None,
2735 },
2736 permissions: PermissionSpec {
2737 allow_write: false,
2738 allow_network: false,
2739 deny_all_tools: false,
2740 allowed_tools: vec!["rg".to_string()],
2741 file_scope: vec!["README.md".to_string()],
2742 },
2743 model_policy: ModelPolicy {
2744 provider: Some("openai".to_string()),
2745 model: Some("gpt-5.4".to_string()),
2746 fallback_models: vec!["gpt-5.4-mini".to_string()],
2747 },
2748 promotion_policy: PromotionPolicy {
2749 strategy: PromotionStrategy::TeacherSelected,
2750 require_teacher_review: true,
2751 min_successful_branches: Some(1),
2752 promotion_gate: PromotionGate::default(),
2753 },
2754 gates: Vec::new(),
2755 nodes: vec![
2756 WorkflowNode::BranchSet(BranchSpec {
2757 id: "discover".to_string(),
2758 description: Some("parallel doc inspection".to_string()),
2759 parallel: true,
2760 budget: BudgetSpec {
2761 max_steps: Some(12),
2762 timeout_secs: Some(600),
2763 max_parallel: Some(2),
2764 max_tokens: None,
2765 },
2766 permissions: PermissionSpec::default(),
2767 model_policy: ModelPolicy::default(),
2768 children: vec![WorkflowNode::Leaf(discover_leaf)],
2769 }),
2770 WorkflowNode::Sequence(SequenceSpec {
2771 id: "review-and-reduce".to_string(),
2772 children: vec![
2773 WorkflowNode::TeacherReview(TeacherReviewSpec {
2774 id: "select-best".to_string(),
2775 candidates: vec!["scan-readme".to_string()],
2776 promotion_policy: PromotionPolicy {
2777 strategy: PromotionStrategy::BestScore,
2778 require_teacher_review: true,
2779 min_successful_branches: Some(1),
2780 promotion_gate: PromotionGate::default(),
2781 },
2782 }),
2783 WorkflowNode::Reduce(ReduceSpec {
2784 id: "summarize".to_string(),
2785 inputs: vec!["scan-readme".to_string()],
2786 prompt: "Summarize the smallest safe patch".to_string(),
2787 model_policy: ModelPolicy::default(),
2788 }),
2789 ],
2790 }),
2791 WorkflowNode::Cond(CondSpec {
2792 id: "maybe-expand".to_string(),
2793 condition: "summary identifies multiple independent gaps".to_string(),
2794 then_nodes: vec![WorkflowNode::Expand(ExpandSpec {
2795 id: "split-followups".to_string(),
2796 source: "summarize".to_string(),
2797 max_children: None,
2798 template: Some(Box::new(WorkflowNode::Leaf(LeafSpec {
2799 id: "followup-template".to_string(),
2800 prompt: "Patch one independent gap".to_string(),
2801 agent_type: AgentType::Implementer,
2802 role: None,
2803 profile: None,
2804 mode: TaskMode::ReadWrite,
2805 isolation: IsolationMode::Worktree,
2806 file_scope: vec!["README.md".to_string()],
2807 cwd: None,
2808 depends_on_results: Vec::new(),
2809 budget: BudgetSpec::default(),
2810 permissions: PermissionSpec {
2811 allow_write: true,
2812 allow_network: false,
2813 deny_all_tools: false,
2814 allowed_tools: Vec::new(),
2815 file_scope: vec!["README.md".to_string()],
2816 },
2817 model_policy: ModelPolicy::default(),
2818 }))),
2819 })],
2820 else_nodes: vec![WorkflowNode::LoopUntil(LoopUntilSpec {
2821 id: "verify-once".to_string(),
2822 condition: "local verification passes".to_string(),
2823 max_iterations: Some(1),
2824 children: Vec::new(),
2825 })],
2826 }),
2827 ],
2828 };
2829
2830 let json = serde_json::to_string_pretty(&workflow).expect("serialize workflow ir");
2831
2832 assert!(json.contains("\"kind\": \"branch_set\""));
2833 assert!(json.contains("\"strategy\": \"teacher_selected\""));
2834 assert!(json.contains("\"profile\": \"scout\""));
2835 let parsed: WorkflowSpec = serde_json::from_str(&json).expect("parse workflow ir");
2836 assert_eq!(parsed, workflow);
2837
2838 let minimal: WorkflowSpec = serde_json::from_str(r#"{"goal":"ship v0.9","nodes":[]}"#)
2839 .expect("parse minimal workflow ir");
2840 assert_eq!(minimal.budget, BudgetSpec::default());
2841 assert_eq!(minimal.permissions, PermissionSpec::default());
2842 assert_eq!(minimal.model_policy, ModelPolicy::default());
2843
2844 // Pre-profile leaf IR stays parseable and profile-less leaves omit the key.
2845 let legacy_leaf: LeafSpec = serde_json::from_str(r#"{"id":"scan","prompt":"scan safely"}"#)
2846 .expect("parse pre-profile leaf ir");
2847 assert_eq!(legacy_leaf.profile, None);
2848 let legacy_json = serde_json::to_string(&legacy_leaf).expect("serialize legacy leaf");
2849 assert!(!legacy_json.contains("profile"));
2850 }
2851
2852 #[test]
2853 fn fleet_validation_accepts_one_thousand_agents_and_variable_models() {
2854 let nodes = (0..DEFAULT_FLEET_WORKFLOW_MAX_AGENTS)
2855 .map(|index| {
2856 let mut leaf = match leaf_node(&format!("agent-{index}")) {
2857 WorkflowNode::Leaf(leaf) => leaf,
2858 _ => unreachable!("leaf helper returns a leaf"),
2859 };
2860 leaf.model_policy = if index == 0 {
2861 ModelPolicy {
2862 provider: Some("deepseek".to_string()),
2863 model: Some("deepseek-v4-pro".to_string()),
2864 fallback_models: Vec::new(),
2865 }
2866 } else {
2867 ModelPolicy {
2868 provider: Some("deepseek".to_string()),
2869 model: Some("deepseek-v4-flash".to_string()),
2870 fallback_models: Vec::new(),
2871 }
2872 };
2873 WorkflowNode::Leaf(leaf)
2874 })
2875 .collect();
2876 let workflow = workflow_spec(nodes);
2877
2878 let shape = workflow
2879 .validate_for_fleet()
2880 .expect("one thousand agents should fit the Fleet Workflow limit");
2881
2882 assert_eq!(shape.total_agents, DEFAULT_FLEET_WORKFLOW_MAX_AGENTS);
2883 assert_eq!(shape.max_depth, 1);
2884 }
2885
2886 #[test]
2887 fn fleet_validation_rejects_more_than_one_thousand_agents() {
2888 let nodes = (0..=DEFAULT_FLEET_WORKFLOW_MAX_AGENTS)
2889 .map(|index| leaf_node(&format!("agent-{index}")))
2890 .collect();
2891 let workflow = workflow_spec(nodes);
2892
2893 let err = workflow
2894 .validate_for_fleet()
2895 .expect_err("agent population should be bounded before Fleet launch");
2896
2897 assert_eq!(
2898 err,
2899 WorkflowFleetLimitError::TooManyAgents {
2900 total_agents: DEFAULT_FLEET_WORKFLOW_MAX_AGENTS + 1,
2901 max_total_agents: DEFAULT_FLEET_WORKFLOW_MAX_AGENTS,
2902 }
2903 );
2904 }
2905
2906 #[test]
2907 fn fleet_validation_rejects_depth_beyond_five() {
2908 let mut node = leaf_node("deep-leaf");
2909 for depth in (0..DEFAULT_FLEET_WORKFLOW_MAX_DEPTH).rev() {
2910 node = WorkflowNode::BranchSet(BranchSpec {
2911 id: format!("ring-{depth}"),
2912 description: None,
2913 parallel: true,
2914 budget: BudgetSpec::default(),
2915 permissions: PermissionSpec::default(),
2916 model_policy: ModelPolicy::default(),
2917 children: vec![node],
2918 });
2919 }
2920 let workflow = workflow_spec(vec![node]);
2921
2922 let err = workflow
2923 .validate_for_fleet()
2924 .expect_err("sixth agent ring should be rejected");
2925
2926 assert_eq!(
2927 err,
2928 WorkflowFleetLimitError::RecursionTooDeep {
2929 depth: DEFAULT_FLEET_WORKFLOW_MAX_DEPTH + 1,
2930 max_depth: DEFAULT_FLEET_WORKFLOW_MAX_DEPTH,
2931 }
2932 );
2933 }
2934
2935 #[test]
2936 fn fleet_validation_counts_loop_and_expand_fanout_conservatively() {
2937 let workflow = workflow_spec(vec![
2938 WorkflowNode::LoopUntil(LoopUntilSpec {
2939 id: "retry-ring".to_string(),
2940 condition: "verifier passes".to_string(),
2941 max_iterations: Some(3),
2942 children: vec![leaf_node("retry-worker")],
2943 }),
2944 WorkflowNode::Expand(ExpandSpec {
2945 id: "split".to_string(),
2946 source: "retry-ring".to_string(),
2947 max_children: Some(4),
2948 template: Some(Box::new(leaf_node("split-template"))),
2949 }),
2950 ]);
2951
2952 let shape = workflow
2953 .validate_for_fleet()
2954 .expect("bounded loop and expand should validate");
2955
2956 assert_eq!(shape.total_agents, 7);
2957 assert_eq!(shape.max_depth, 1);
2958 }
2959
2960 #[test]
2961 fn fleet_validation_rejects_unbounded_loop_or_expand_before_launch() {
2962 let workflow = workflow_spec(vec![
2963 WorkflowNode::LoopUntil(LoopUntilSpec {
2964 id: "retry-ring".to_string(),
2965 condition: "verifier passes".to_string(),
2966 max_iterations: None,
2967 children: vec![leaf_node("retry-worker")],
2968 }),
2969 WorkflowNode::Expand(ExpandSpec {
2970 id: "split".to_string(),
2971 source: "retry-ring".to_string(),
2972 max_children: Some(4),
2973 template: Some(Box::new(leaf_node("split-template"))),
2974 }),
2975 ]);
2976
2977 assert!(matches!(
2978 workflow.validate_for_fleet(),
2979 Err(WorkflowFleetLimitError::UnboundedLoop { node }) if node == "retry-ring"
2980 ));
2981
2982 let workflow = workflow_spec(vec![WorkflowNode::Expand(ExpandSpec {
2983 id: "split".to_string(),
2984 source: "retry-ring".to_string(),
2985 max_children: None,
2986 template: Some(Box::new(leaf_node("split-template"))),
2987 })]);
2988
2989 assert!(matches!(
2990 workflow.validate_for_fleet(),
2991 Err(WorkflowFleetLimitError::UnboundedExpand { node }) if node == "split"
2992 ));
2993 }
2994
2995 #[test]
2996 fn branch_result_serialization() {
2997 let result = BranchResult {
2998 branch_id: "discover".to_string(),
2999 task_id: "scan".to_string(),
3000 status: WorkflowRunStatus::Succeeded,
3001 usage: WorkflowUsage {
3002 input_tokens: Some(100),
3003 output_tokens: Some(25),
3004 cost_microusd: Some(42),
3005 },
3006 memo_usage: WorkflowMemoUsage::default(),
3007 artifacts: vec!["trace://branches/discover".to_string()],
3008 notes: Some("validated prompt surfaces".to_string()),
3009 };
3010
3011 let json = serde_json::to_string(&result).expect("serialize branch result");
3012
3013 assert!(json.contains("\"status\":\"succeeded\""));
3014 assert!(json.contains("\"cost_microusd\":42"));
3015 let parsed: BranchResult = serde_json::from_str(&json).expect("parse branch result");
3016 assert_eq!(parsed, result);
3017
3018 let minimal: BranchResult =
3019 serde_json::from_str(r#"{"branch_id":"discover","task_id":"scan","status":"pending"}"#)
3020 .expect("parse minimal branch result");
3021 assert_eq!(minimal.usage, WorkflowUsage::default());
3022 assert_eq!(minimal.memo_usage, WorkflowMemoUsage::default());
3023 assert!(minimal.artifacts.is_empty());
3024 assert_eq!(minimal.notes, None);
3025 }
3026
3027 #[test]
3028 fn workflow_usage_serialization_distinguishes_unknown_from_reported_zero() {
3029 let unknown = serde_json::to_value(WorkflowUsage::default()).expect("unknown usage JSON");
3030 assert_eq!(unknown, serde_json::json!({}));
3031
3032 let reported_zero = serde_json::to_value(WorkflowUsage {
3033 input_tokens: Some(0),
3034 output_tokens: Some(0),
3035 cost_microusd: Some(0),
3036 })
3037 .expect("reported zero usage JSON");
3038 assert_eq!(
3039 reported_zero,
3040 serde_json::json!({
3041 "input_tokens": 0,
3042 "output_tokens": 0,
3043 "cost_microusd": 0,
3044 })
3045 );
3046 }
3047
3048 #[test]
3049 fn leaf_result_serialization() {
3050 let result = LeafResult {
3051 leaf_id: "scan-readme".to_string(),
3052 task_id: "scan".to_string(),
3053 role: None,
3054 profile: Some("reviewer".to_string()),
3055 status: WorkflowRunStatus::Failed,
3056 usage: WorkflowUsage {
3057 input_tokens: Some(11),
3058 output_tokens: Some(7),
3059 cost_microusd: Some(3),
3060 },
3061 memo_usage: WorkflowMemoUsage {
3062 armh_hits: 1,
3063 armh_misses: 0,
3064 armh_saved_estimated_tokens: 128,
3065 provider_prompt_cache_hits: 2,
3066 provider_prompt_cache_misses: 1,
3067 },
3068 output: Some("README needs clearer setup steps".to_string()),
3069 artifacts: vec!["trace://leaves/scan-readme".to_string()],
3070 schema_error: None,
3071 };
3072
3073 let json = serde_json::to_string(&result).expect("serialize leaf result");
3074
3075 assert!(json.contains("\"status\":\"failed\""));
3076 assert!(json.contains("\"input_tokens\":11"));
3077 assert!(json.contains("\"armh_saved_estimated_tokens\":128"));
3078 assert!(json.contains("\"profile\":\"reviewer\""));
3079 let parsed: LeafResult = serde_json::from_str(&json).expect("parse leaf result");
3080 assert_eq!(parsed, result);
3081
3082 let minimal: LeafResult = serde_json::from_str(
3083 r#"{"leaf_id":"scan-readme","task_id":"scan","status":"pending"}"#,
3084 )
3085 .expect("parse minimal leaf result");
3086 assert_eq!(minimal.profile, None);
3087 assert_eq!(minimal.usage, WorkflowUsage::default());
3088 assert_eq!(minimal.memo_usage, WorkflowMemoUsage::default());
3089 assert_eq!(minimal.output, None);
3090 assert!(minimal.artifacts.is_empty());
3091 }
3092
3093 #[test]
3094 fn control_node_result_serialization() {
3095 let result = ControlNodeResult {
3096 node_id: "select-fix".to_string(),
3097 kind: ControlNodeKind::TeacherReview,
3098 status: WorkflowRunStatus::Running,
3099 selected_children: vec!["branch-a".to_string(), "branch-c".to_string()],
3100 summary: Some("teacher review is waiting on verifier evidence".to_string()),
3101 };
3102
3103 let json = serde_json::to_string(&result).expect("serialize control node result");
3104
3105 assert!(json.contains("\"kind\":\"teacher_review\""));
3106 assert!(json.contains("\"status\":\"running\""));
3107 let parsed: ControlNodeResult =
3108 serde_json::from_str(&json).expect("parse control node result");
3109 assert_eq!(parsed, result);
3110
3111 let minimal: ControlNodeResult = serde_json::from_str(
3112 r#"{"node_id":"select-fix","kind":"branch_set","status":"pending"}"#,
3113 )
3114 .expect("parse minimal control node result");
3115 assert!(minimal.selected_children.is_empty());
3116 assert_eq!(minimal.summary, None);
3117 }
3118
3119 #[test]
3120 fn run_mock_three_branch_workflow() {
3121 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3122 id: "discover".to_string(),
3123 description: None,
3124 parallel: true,
3125 budget: BudgetSpec::default(),
3126 permissions: PermissionSpec::default(),
3127 model_policy: ModelPolicy::default(),
3128 children: vec![
3129 leaf_node("scan-readme"),
3130 leaf_node("scan-config"),
3131 leaf_node("scan-tests"),
3132 ],
3133 })]);
3134
3135 let mut executor = MockWorkflowExecutor::new();
3136 let execution = executor.run(&workflow).expect("mock workflow should run");
3137
3138 assert_eq!(execution.status, WorkflowRunStatus::Succeeded);
3139 assert_eq!(
3140 execution
3141 .leaf_results
3142 .iter()
3143 .map(|result| result.leaf_id.as_str())
3144 .collect::<Vec<_>>(),
3145 vec!["scan-readme", "scan-config", "scan-tests"]
3146 );
3147 assert_eq!(execution.branch_results.len(), 1);
3148 assert_eq!(execution.branch_results[0].branch_id, "discover");
3149 assert_eq!(
3150 control_result(&execution, "discover").selected_children,
3151 vec!["scan-readme", "scan-config", "scan-tests"]
3152 );
3153 }
3154
3155 #[test]
3156 fn mock_executor_surfaces_leaf_profile() {
3157 let mut profiled_leaf = match leaf_node("review-change") {
3158 WorkflowNode::Leaf(leaf) => leaf,
3159 _ => unreachable!("leaf helper returns a leaf"),
3160 };
3161 profiled_leaf.profile = Some("reviewer".to_string());
3162 let workflow = workflow_spec(vec![
3163 WorkflowNode::Leaf(profiled_leaf),
3164 leaf_node("scan-readme"),
3165 ]);
3166
3167 let execution = MockWorkflowExecutor::new()
3168 .run(&workflow)
3169 .expect("mock workflow should run");
3170
3171 assert_eq!(execution.status, WorkflowRunStatus::Succeeded);
3172 assert_eq!(
3173 execution.leaf_results[0].profile.as_deref(),
3174 Some("reviewer")
3175 );
3176 assert_eq!(execution.leaf_results[1].profile, None);
3177 }
3178
3179 #[test]
3180 fn mock_executor_surfaces_leaf_role() {
3181 let mut role_leaf = match leaf_node("scout-issue") {
3182 WorkflowNode::Leaf(leaf) => leaf,
3183 _ => unreachable!("leaf helper returns a leaf"),
3184 };
3185 role_leaf.role = Some("scout".to_string());
3186 let workflow = workflow_spec(vec![WorkflowNode::Leaf(role_leaf)]);
3187
3188 let execution = MockWorkflowExecutor::new()
3189 .run(&workflow)
3190 .expect("mock workflow should run");
3191
3192 assert_eq!(execution.leaf_results[0].role.as_deref(), Some("scout"));
3193 }
3194
3195 #[test]
3196 fn leaf_role_token_rule_rejects_invalid_names() {
3197 for bad in ["", "has space", "role=scout"] {
3198 let mut leaf = match leaf_node("scan") {
3199 WorkflowNode::Leaf(leaf) => leaf,
3200 _ => unreachable!("leaf helper returns a leaf"),
3201 };
3202 leaf.role = Some(bad.to_string());
3203 let workflow = workflow_spec(vec![WorkflowNode::Leaf(leaf)]);
3204
3205 let err = MockWorkflowExecutor::new()
3206 .run(&workflow)
3207 .expect_err("invalid role token should fail validation");
3208
3209 assert!(
3210 matches!(&err, WorkflowExecutionError::InvalidLeafRole { role, .. } if role == bad),
3211 "role `{bad}` should be rejected, got {err:?}"
3212 );
3213 }
3214 }
3215
3216 #[test]
3217 fn leaf_role_roundtrips_without_required_provider_model() {
3218 let leaf = LeafSpec {
3219 id: "scout-1".to_string(),
3220 prompt: "Investigate #4090. Read-only.".to_string(),
3221 agent_type: AgentType::Explore,
3222 role: Some("scout".to_string()),
3223 profile: None,
3224 mode: TaskMode::ReadOnly,
3225 isolation: IsolationMode::Shared,
3226 file_scope: Vec::new(),
3227 cwd: None,
3228 depends_on_results: Vec::new(),
3229 budget: BudgetSpec::default(),
3230 permissions: PermissionSpec::default(),
3231 model_policy: ModelPolicy::default(),
3232 };
3233 let json = serde_json::to_string(&leaf).expect("serialize");
3234 assert!(json.contains("\"role\":\"scout\""));
3235 let parsed: LeafSpec = serde_json::from_str(&json).expect("parse");
3236 assert_eq!(parsed.role.as_deref(), Some("scout"));
3237 // Provider/model are optional overrides, not required identity fields.
3238 assert_eq!(parsed.model_policy.provider, None);
3239 assert_eq!(parsed.model_policy.model, None);
3240 assert_eq!(parsed.model_policy, ModelPolicy::default());
3241 }
3242
3243 #[test]
3244 fn leaf_profile_token_rule_rejects_invalid_names() {
3245 for bad in ["", "has space", "quote\"y", "role=reviewer", "back`tick"] {
3246 let mut leaf = match leaf_node("scan") {
3247 WorkflowNode::Leaf(leaf) => leaf,
3248 _ => unreachable!("leaf helper returns a leaf"),
3249 };
3250 leaf.profile = Some(bad.to_string());
3251 let workflow = workflow_spec(vec![WorkflowNode::Leaf(leaf)]);
3252
3253 let err = MockWorkflowExecutor::new()
3254 .run(&workflow)
3255 .expect_err("invalid profile token should fail validation");
3256
3257 assert!(
3258 matches!(&err, WorkflowExecutionError::InvalidLeafProfile { profile, .. } if profile == bad),
3259 "profile `{bad}` should be rejected, got {err:?}"
3260 );
3261 }
3262
3263 let mut leaf = match leaf_node("scan") {
3264 WorkflowNode::Leaf(leaf) => leaf,
3265 _ => unreachable!("leaf helper returns a leaf"),
3266 };
3267 leaf.profile = Some("reviewer".to_string());
3268 let workflow = workflow_spec(vec![WorkflowNode::Leaf(leaf)]);
3269 MockWorkflowExecutor::new()
3270 .run(&workflow)
3271 .expect("valid profile token should pass validation");
3272 }
3273
3274 #[test]
3275 fn mock_executor_aggregates_leaf_usage() {
3276 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3277 id: "discover".to_string(),
3278 description: None,
3279 parallel: true,
3280 budget: BudgetSpec::default(),
3281 permissions: PermissionSpec::default(),
3282 model_policy: ModelPolicy::default(),
3283 children: vec![leaf_node("scan-readme"), leaf_node("scan-tests")],
3284 })]);
3285
3286 let mut executor = MockWorkflowExecutor::new()
3287 .with_leaf_outcome(
3288 "scan-readme",
3289 MockLeafOutcome::succeeded("readme ok").with_usage(WorkflowUsage {
3290 input_tokens: Some(100),
3291 output_tokens: Some(25),
3292 cost_microusd: Some(500),
3293 }),
3294 )
3295 .with_leaf_outcome(
3296 "scan-tests",
3297 MockLeafOutcome::succeeded("tests ok").with_usage(WorkflowUsage {
3298 input_tokens: Some(50),
3299 output_tokens: Some(10),
3300 cost_microusd: Some(250),
3301 }),
3302 );
3303
3304 let execution = executor.run(&workflow).expect("mock workflow should run");
3305
3306 assert_eq!(
3307 execution.usage,
3308 WorkflowUsage {
3309 input_tokens: Some(150),
3310 output_tokens: Some(35),
3311 cost_microusd: Some(750),
3312 }
3313 );
3314 assert_eq!(execution.usage.total_tokens(), Some(185));
3315 assert_eq!(execution.branch_results[0].usage, execution.usage);
3316 assert_eq!(
3317 execution
3318 .leaf_results
3319 .iter()
3320 .map(|result| result.usage.cost_microusd)
3321 .collect::<Vec<_>>(),
3322 vec![Some(500), Some(250)]
3323 );
3324 }
3325
3326 #[test]
3327 fn mock_executor_aggregates_memo_usage() {
3328 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3329 id: "cache-branches".to_string(),
3330 description: None,
3331 parallel: true,
3332 budget: BudgetSpec::default(),
3333 permissions: PermissionSpec::default(),
3334 model_policy: ModelPolicy::default(),
3335 children: vec![leaf_node("rlm-hit"), leaf_node("rlm-miss")],
3336 })]);
3337
3338 let mut executor = MockWorkflowExecutor::new()
3339 .with_leaf_outcome(
3340 "rlm-hit",
3341 MockLeafOutcome::succeeded("memo hit").with_memo_usage(WorkflowMemoUsage {
3342 armh_hits: 1,
3343 armh_misses: 0,
3344 armh_saved_estimated_tokens: 4096,
3345 provider_prompt_cache_hits: 1,
3346 provider_prompt_cache_misses: 0,
3347 }),
3348 )
3349 .with_leaf_outcome(
3350 "rlm-miss",
3351 MockLeafOutcome::succeeded("memo miss").with_memo_usage(WorkflowMemoUsage {
3352 armh_hits: 0,
3353 armh_misses: 1,
3354 armh_saved_estimated_tokens: 0,
3355 provider_prompt_cache_hits: 0,
3356 provider_prompt_cache_misses: 1,
3357 }),
3358 );
3359
3360 let execution = executor.run(&workflow).expect("mock workflow should run");
3361
3362 assert_eq!(
3363 execution.memo_usage,
3364 WorkflowMemoUsage {
3365 armh_hits: 1,
3366 armh_misses: 1,
3367 armh_saved_estimated_tokens: 4096,
3368 provider_prompt_cache_hits: 1,
3369 provider_prompt_cache_misses: 1,
3370 }
3371 );
3372 assert_eq!(execution.branch_results[0].memo_usage, execution.memo_usage);
3373 assert_eq!(
3374 execution
3375 .leaf_results
3376 .iter()
3377 .map(|result| (result.memo_usage.armh_hits, result.memo_usage.armh_misses))
3378 .collect::<Vec<_>>(),
3379 vec![(1, 0), (0, 1)]
3380 );
3381 }
3382
3383 #[test]
3384 fn mock_executor_marks_cancelled_before_leaf() {
3385 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3386 id: "discover".to_string(),
3387 description: None,
3388 parallel: true,
3389 budget: BudgetSpec::default(),
3390 permissions: PermissionSpec::default(),
3391 model_policy: ModelPolicy::default(),
3392 children: vec![leaf_node("scan-readme"), leaf_node("scan-tests")],
3393 })]);
3394
3395 let mut executor = MockWorkflowExecutor::new().with_cancelled();
3396 let execution = executor.run(&workflow).expect("mock workflow should run");
3397
3398 assert_eq!(execution.status, WorkflowRunStatus::Cancelled);
3399 assert_eq!(execution.leaf_results.len(), 1);
3400 assert_eq!(
3401 execution.leaf_results[0].status,
3402 WorkflowRunStatus::Cancelled
3403 );
3404 assert_eq!(
3405 execution.branch_results[0].status,
3406 WorkflowRunStatus::Cancelled
3407 );
3408 assert_eq!(
3409 control_result(&execution, "discover").status,
3410 WorkflowRunStatus::Cancelled
3411 );
3412 }
3413
3414 #[test]
3415 fn mock_executor_stops_when_global_leaf_budget_is_exhausted() {
3416 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3417 id: "discover".to_string(),
3418 description: None,
3419 parallel: true,
3420 budget: BudgetSpec::default(),
3421 permissions: PermissionSpec::default(),
3422 model_policy: ModelPolicy::default(),
3423 children: vec![
3424 leaf_node("scan-readme"),
3425 leaf_node("scan-config"),
3426 leaf_node("scan-tests"),
3427 ],
3428 })]);
3429
3430 let mut executor = MockWorkflowExecutor::new().with_max_leaf_steps(1);
3431 let execution = executor.run(&workflow).expect("mock workflow should run");
3432
3433 assert_eq!(execution.status, WorkflowRunStatus::BudgetExceeded);
3434 assert_eq!(
3435 execution
3436 .leaf_results
3437 .iter()
3438 .map(|result| (result.leaf_id.as_str(), result.status))
3439 .collect::<Vec<_>>(),
3440 vec![
3441 ("scan-readme", WorkflowRunStatus::Succeeded),
3442 ("scan-config", WorkflowRunStatus::BudgetExceeded)
3443 ]
3444 );
3445 assert_eq!(
3446 execution.branch_results[0].status,
3447 WorkflowRunStatus::BudgetExceeded
3448 );
3449 }
3450
3451 #[test]
3452 fn mock_executor_treats_zero_step_budgets_as_unbounded() {
3453 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3454 id: "verify".to_string(),
3455 description: None,
3456 parallel: false,
3457 budget: BudgetSpec::default(),
3458 permissions: PermissionSpec::default(),
3459 model_policy: ModelPolicy::default(),
3460 children: vec![
3461 leaf_node_with_budget(
3462 "run-tests",
3463 BudgetSpec {
3464 max_steps: Some(0),
3465 timeout_secs: None,
3466 max_parallel: None,
3467 max_tokens: None,
3468 },
3469 ),
3470 leaf_node("summarize"),
3471 ],
3472 })]);
3473
3474 let mut executor = MockWorkflowExecutor::new().with_max_leaf_steps(0);
3475 let execution = executor.run(&workflow).expect("mock workflow should run");
3476
3477 assert_eq!(execution.status, WorkflowRunStatus::Succeeded);
3478 assert_eq!(execution.leaf_results.len(), 2);
3479 assert!(
3480 execution
3481 .leaf_results
3482 .iter()
3483 .all(|result| result.status == WorkflowRunStatus::Succeeded)
3484 );
3485 }
3486
3487 #[test]
3488 fn mock_executor_stops_when_global_token_budget_is_exhausted() {
3489 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3490 id: "discover".to_string(),
3491 description: None,
3492 parallel: true,
3493 budget: BudgetSpec::default(),
3494 permissions: PermissionSpec::default(),
3495 model_policy: ModelPolicy::default(),
3496 children: vec![
3497 leaf_node("scan-readme"),
3498 leaf_node("scan-config"),
3499 leaf_node("scan-tests"),
3500 ],
3501 })]);
3502
3503 // First leaf uses 600 tokens (300 in + 300 out); after the second leaf
3504 // (500 tokens) the running total is 1100, exceeding the 1000-token
3505 // global cap, so the third leaf hits the exhausted budget and halts the
3506 // run.
3507 let mut executor = MockWorkflowExecutor::new()
3508 .with_max_leaf_tokens(1000)
3509 .with_leaf_outcome(
3510 "scan-readme",
3511 MockLeafOutcome::succeeded("readme done").with_usage(WorkflowUsage {
3512 input_tokens: Some(300),
3513 output_tokens: Some(300),
3514 cost_microusd: Some(0),
3515 }),
3516 )
3517 .with_leaf_outcome(
3518 "scan-config",
3519 MockLeafOutcome::succeeded("config done").with_usage(WorkflowUsage {
3520 input_tokens: Some(250),
3521 output_tokens: Some(250),
3522 cost_microusd: Some(0),
3523 }),
3524 );
3525 let execution = executor.run(&workflow).expect("mock workflow should run");
3526
3527 assert_eq!(execution.status, WorkflowRunStatus::BudgetExceeded);
3528 // Leaves 1+2 consume 1100 tokens, exhausting the 1000-token global cap.
3529 // The third leaf is attempted, sees the budget already exceeded, and is
3530 // recorded as BudgetExceeded — the same boundary-leaf behaviour used by
3531 // step budgets (max_leaf_steps). The budget outcome carries no tokens,
3532 // so total usage stays at 1100.
3533 assert_eq!(execution.leaf_results.len(), 3);
3534 assert_eq!(
3535 execution.leaf_results[0].status,
3536 WorkflowRunStatus::Succeeded
3537 );
3538 assert_eq!(
3539 execution.leaf_results[1].status,
3540 WorkflowRunStatus::Succeeded
3541 );
3542 assert_eq!(
3543 execution.leaf_results[2].status,
3544 WorkflowRunStatus::BudgetExceeded
3545 );
3546 assert_eq!(execution.usage.total_tokens(), Some(1100));
3547 }
3548
3549 #[test]
3550 fn mock_executor_honors_zero_token_leaf_budget() {
3551 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3552 id: "verify".to_string(),
3553 description: None,
3554 parallel: false,
3555 budget: BudgetSpec::default(),
3556 permissions: PermissionSpec::default(),
3557 model_policy: ModelPolicy::default(),
3558 children: vec![
3559 leaf_node_with_budget(
3560 "run-tests",
3561 BudgetSpec {
3562 max_steps: None,
3563 timeout_secs: None,
3564 max_parallel: None,
3565 max_tokens: Some(0),
3566 },
3567 ),
3568 leaf_node("summarize"),
3569 ],
3570 })]);
3571
3572 let mut executor = MockWorkflowExecutor::new();
3573 let execution = executor.run(&workflow).expect("mock workflow should run");
3574
3575 assert_eq!(execution.status, WorkflowRunStatus::BudgetExceeded);
3576 assert_eq!(execution.leaf_results.len(), 1);
3577 assert_eq!(
3578 execution.leaf_results[0].status,
3579 WorkflowRunStatus::BudgetExceeded
3580 );
3581 assert!(
3582 execution.leaf_results[0]
3583 .output
3584 .as_deref()
3585 .unwrap_or_default()
3586 .contains("token budget exhausted")
3587 );
3588 }
3589
3590 #[test]
3591 fn mock_executor_honors_per_leaf_token_cap() {
3592 let workflow = workflow_spec(vec![WorkflowNode::BranchSet(BranchSpec {
3593 id: "review".to_string(),
3594 description: None,
3595 parallel: false,
3596 budget: BudgetSpec::default(),
3597 permissions: PermissionSpec::default(),
3598 model_policy: ModelPolicy::default(),
3599 children: vec![
3600 leaf_node_with_budget(
3601 "expensive-scan",
3602 BudgetSpec {
3603 max_steps: None,
3604 timeout_secs: None,
3605 max_parallel: None,
3606 max_tokens: Some(500),
3607 },
3608 ),
3609 leaf_node("summarize"),
3610 ],
3611 })]);
3612
3613 // The leaf outcome uses 800 tokens which exceeds the per-leaf cap of 500.
3614 let mut executor = MockWorkflowExecutor::new().with_leaf_outcome(
3615 "expensive-scan",
3616 MockLeafOutcome::succeeded("scan done").with_usage(WorkflowUsage {
3617 input_tokens: Some(500),
3618 output_tokens: Some(300),
3619 cost_microusd: Some(0),
3620 }),
3621 );
3622 let execution = executor.run(&workflow).expect("mock workflow should run");
3623
3624 assert_eq!(execution.status, WorkflowRunStatus::BudgetExceeded);
3625 assert_eq!(execution.leaf_results.len(), 1);
3626 assert_eq!(
3627 execution.leaf_results[0].status,
3628 WorkflowRunStatus::BudgetExceeded
3629 );
3630 assert!(
3631 execution.leaf_results[0]
3632 .output
3633 .as_deref()
3634 .unwrap_or_default()
3635 .contains("token budget exhausted")
3636 );
3637 }
3638
3639 #[test]
3640 fn budget_spec_serializes_max_tokens() {
3641 let budget = BudgetSpec {
3642 max_steps: Some(10),
3643 timeout_secs: Some(600),
3644 max_parallel: Some(4),
3645 max_tokens: Some(50_000),
3646 };
3647 let json = serde_json::to_string(&budget).expect("serialize budget");
3648 let parsed: BudgetSpec = serde_json::from_str(&json).expect("parse budget");
3649 assert_eq!(parsed, budget);
3650 assert!(json.contains("\"max_tokens\":50000"));
3651
3652 // Default (all None) round-trips without the field present.
3653 let default_json =
3654 serde_json::to_string(&BudgetSpec::default()).expect("serialize default");
3655 let parsed_default: BudgetSpec =
3656 serde_json::from_str(&default_json).expect("parse default budget");
3657 assert_eq!(parsed_default, BudgetSpec::default());
3658 assert!(parsed_default.max_tokens.is_none());
3659 }
3660
3661 #[test]
3662 fn loop_until_stops_on_pass() {
3663 let workflow = workflow_spec(vec![WorkflowNode::LoopUntil(LoopUntilSpec {
3664 id: "verify".to_string(),
3665 condition: "verification passed".to_string(),
3666 max_iterations: Some(5),
3667 children: vec![leaf_node("run-check")],
3668 })]);
3669
3670 let mut executor =
3671 MockWorkflowExecutor::new().with_predicate_results("verify", vec![false, false, true]);
3672 let execution = executor.run(&workflow).expect("loop should run");
3673
3674 assert_eq!(execution.status, WorkflowRunStatus::Succeeded);
3675 assert_eq!(execution.leaf_results.len(), 3);
3676 assert_eq!(
3677 control_result(&execution, "verify").summary.as_deref(),
3678 Some("loop_until iterations=3")
3679 );
3680 }
3681
3682 #[test]
3683 fn loop_until_honors_max_iters() {
3684 let workflow = workflow_spec(vec![WorkflowNode::LoopUntil(LoopUntilSpec {
3685 id: "verify".to_string(),
3686 condition: "verification passed".to_string(),
3687 max_iterations: Some(2),
3688 children: vec![leaf_node("run-check")],
3689 })]);
3690
3691 let mut executor =
3692 MockWorkflowExecutor::new().with_predicate_results("verify", vec![false, false, true]);
3693 let execution = executor.run(&workflow).expect("loop should run");
3694
3695 assert_eq!(execution.status, WorkflowRunStatus::Failed);
3696 assert_eq!(execution.leaf_results.len(), 2);
3697 assert_eq!(
3698 control_result(&execution, "verify").summary.as_deref(),
3699 Some("loop_until iterations=2")
3700 );
3701 }
3702
3703 #[test]
3704 fn cond_uses_logged_predicate_result() {
3705 let workflow = workflow_spec(vec![WorkflowNode::Cond(CondSpec {
3706 id: "should-fix".to_string(),
3707 condition: "finding requires a patch".to_string(),
3708 then_nodes: vec![leaf_node("patch")],
3709 else_nodes: vec![leaf_node("report-only")],
3710 })]);
3711
3712 let mut executor =
3713 MockWorkflowExecutor::new().with_predicate_results("should-fix", vec![true]);
3714 let execution = executor.run(&workflow).expect("cond should run");
3715
3716 assert_eq!(
3717 execution
3718 .leaf_results
3719 .iter()
3720 .map(|result| result.leaf_id.as_str())
3721 .collect::<Vec<_>>(),
3722 vec!["patch"]
3723 );
3724 assert_eq!(
3725 control_result(&execution, "should-fix").summary.as_deref(),
3726 Some("predicate_result=true")
3727 );
3728 }
3729
3730 #[test]
3731 fn expand_respects_max_children() {
3732 let workflow = workflow_spec(vec![WorkflowNode::Expand(ExpandSpec {
3733 id: "split".to_string(),
3734 source: "plan".to_string(),
3735 max_children: Some(2),
3736 template: None,
3737 })]);
3738
3739 let generated = vec![leaf_node("first"), leaf_node("second"), leaf_node("third")];
3740 let mut executor = MockWorkflowExecutor::new().with_generated_nodes("split", generated);
3741 let execution = executor.run(&workflow).expect("expand should run");
3742
3743 assert_eq!(
3744 execution
3745 .leaf_results
3746 .iter()
3747 .map(|result| result.leaf_id.as_str())
3748 .collect::<Vec<_>>(),
3749 vec!["first", "second"]
3750 );
3751 assert_eq!(
3752 control_result(&execution, "split").selected_children,
3753 vec!["first", "second"]
3754 );
3755 }
3756
3757 #[test]
3758 fn expand_generated_nodes_validate_before_run() {
3759 let workflow = workflow_spec(vec![WorkflowNode::Expand(ExpandSpec {
3760 id: "split".to_string(),
3761 source: "plan".to_string(),
3762 max_children: None,
3763 template: None,
3764 })]);
3765
3766 let mut executor = MockWorkflowExecutor::new()
3767 .with_generated_nodes("split", vec![invalid_leaf_node("bad")]);
3768 let err = executor
3769 .run(&workflow)
3770 .expect_err("invalid generated leaf should fail before execution");
3771
3772 assert_eq!(
3773 err,
3774 WorkflowExecutionError::EmptyLeafPrompt {
3775 leaf: "bad".to_string()
3776 }
3777 );
3778 }
3779
3780 #[test]
3781 fn workflow_spec_rejects_unknown_leaf_dependency() {
3782 let mut summarize = leaf_node("summarize");
3783 let WorkflowNode::Leaf(spec) = &mut summarize else {
3784 panic!("expected leaf");
3785 };
3786 spec.depends_on_results = vec!["missing-scan".to_string()];
3787 let workflow = workflow_spec(vec![summarize]);
3788
3789 let mut executor = MockWorkflowExecutor::new();
3790 let err = executor
3791 .run(&workflow)
3792 .expect_err("unknown leaf dependency should fail before execution");
3793
3794 assert_eq!(
3795 err,
3796 WorkflowExecutionError::UnknownNodeReference {
3797 node: "summarize".to_string(),
3798 field: "depends_on_results",
3799 reference: "missing-scan".to_string(),
3800 }
3801 );
3802 }
3803
3804 #[test]
3805 fn workflow_spec_rejects_unknown_reduce_input() {
3806 let workflow = workflow_spec(vec![
3807 leaf_node("scan"),
3808 WorkflowNode::Reduce(ReduceSpec {
3809 id: "summarize".to_string(),
3810 inputs: vec!["scan".to_string(), "missing-review".to_string()],
3811 prompt: "Summarize safe fixes".to_string(),
3812 model_policy: ModelPolicy::default(),
3813 }),
3814 ]);
3815
3816 let mut executor = MockWorkflowExecutor::new();
3817 let err = executor
3818 .run(&workflow)
3819 .expect_err("unknown reduce input should fail before execution");
3820
3821 assert_eq!(
3822 err,
3823 WorkflowExecutionError::UnknownNodeReference {
3824 node: "summarize".to_string(),
3825 field: "inputs",
3826 reference: "missing-review".to_string(),
3827 }
3828 );
3829 }
3830
3831 #[test]
3832 fn workflow_spec_rejects_unknown_teacher_candidate() {
3833 let workflow = workflow_spec(vec![
3834 leaf_node("candidate-a"),
3835 WorkflowNode::TeacherReview(TeacherReviewSpec {
3836 id: "teacher-review".to_string(),
3837 candidates: vec!["candidate-a".to_string(), "candidate-b".to_string()],
3838 promotion_policy: PromotionPolicy::default(),
3839 }),
3840 ]);
3841
3842 let mut executor = MockWorkflowExecutor::new();
3843 let err = executor
3844 .run(&workflow)
3845 .expect_err("unknown teacher candidate should fail before execution");
3846
3847 assert_eq!(
3848 err,
3849 WorkflowExecutionError::UnknownNodeReference {
3850 node: "teacher-review".to_string(),
3851 field: "candidates",
3852 reference: "candidate-b".to_string(),
3853 }
3854 );
3855 }
3856
3857 #[test]
3858 fn teacher_candidate_serialization() {
3859 let candidate = TeacherCandidate {
3860 candidate_id: "teacher-review:branch-a".to_string(),
3861 kind: TeacherCandidateKind::WorkflowRecipe,
3862 status: TeacherCandidateStatus::Proposed,
3863 source_node_id: "branch-a".to_string(),
3864 source_branch_id: Some("branch-a".to_string()),
3865 summary: "Winning branch found a reusable workflow recipe.".to_string(),
3866 evidence: vec![
3867 "status=Succeeded".to_string(),
3868 "tokens=42, cost_microusd=7".to_string(),
3869 ],
3870 replay_results: vec![StudentReplayResult {
3871 trace_id: "trace-a".to_string(),
3872 candidate_id: "teacher-review:branch-a".to_string(),
3873 baseline: StudentReplayMetrics {
3874 score: 70,
3875 cost_microusd: 10,
3876 },
3877 candidate: StudentReplayMetrics {
3878 score: 74,
3879 cost_microusd: 12,
3880 },
3881 required_tests: vec![StudentReplayTestResult {
3882 name: "cargo test -p codewhale-workflow".to_string(),
3883 passed: true,
3884 }],
3885 policy_violations: Vec::new(),
3886 stale: false,
3887 notes: Some("offline replay improved the constrained student".to_string()),
3888 }],
3889 };
3890
3891 let json = serde_json::to_string(&candidate).expect("serialize teacher candidate");
3892
3893 assert!(json.contains("\"kind\":\"workflow_recipe\""));
3894 assert!(json.contains("\"status\":\"proposed\""));
3895 assert!(json.contains("\"replay_results\""));
3896 let parsed: TeacherCandidate =
3897 serde_json::from_str(&json).expect("parse teacher candidate");
3898 assert_eq!(parsed, candidate);
3899 }
3900
3901 #[test]
3902 fn teacher_review_produces_candidate_from_trace() {
3903 let review = TeacherReviewSpec {
3904 id: "teacher-review".to_string(),
3905 candidates: vec!["winning-branch".to_string()],
3906 promotion_policy: PromotionPolicy::default(),
3907 };
3908 let execution = WorkflowExecution {
3909 branch_results: vec![BranchResult {
3910 branch_id: "winning-branch".to_string(),
3911 task_id: "winning-branch".to_string(),
3912 status: WorkflowRunStatus::Succeeded,
3913 usage: WorkflowUsage {
3914 input_tokens: Some(30),
3915 output_tokens: Some(12),
3916 cost_microusd: Some(7),
3917 },
3918 memo_usage: WorkflowMemoUsage::default(),
3919 artifacts: vec!["trace://branches/winning-branch".to_string()],
3920 notes: Some("branch produced a minimal verified patch".to_string()),
3921 }],
3922 ..WorkflowExecution::default()
3923 };
3924
3925 let report = TeacherReviewReport::from_execution(&review, &execution);
3926
3927 assert_eq!(report.review_node_id, "teacher-review");
3928 assert_eq!(report.candidates.len(), 1);
3929 assert_eq!(
3930 report.candidates[0].kind,
3931 TeacherCandidateKind::WorkflowRecipe
3932 );
3933 assert_eq!(
3934 report.candidates[0].status,
3935 TeacherCandidateStatus::Proposed
3936 );
3937 assert!(
3938 report.candidates[0]
3939 .evidence
3940 .iter()
3941 .any(|line| line.contains("tokens=42"))
3942 );
3943 }
3944
3945 #[test]
3946 fn failed_leaf_becomes_regression_test_candidate() {
3947 let review = TeacherReviewSpec {
3948 id: "teacher-review".to_string(),
3949 candidates: vec!["verify-failure".to_string()],
3950 promotion_policy: PromotionPolicy::default(),
3951 };
3952 let execution = WorkflowExecution {
3953 leaf_results: vec![LeafResult {
3954 leaf_id: "verify-failure".to_string(),
3955 task_id: "verify-failure".to_string(),
3956 role: None,
3957 profile: None,
3958 status: WorkflowRunStatus::Failed,
3959 usage: WorkflowUsage::default(),
3960 memo_usage: WorkflowMemoUsage::default(),
3961 output: Some("cargo test failed with a replay mismatch".to_string()),
3962 artifacts: Vec::new(),
3963 schema_error: None,
3964 }],
3965 ..WorkflowExecution::default()
3966 };
3967
3968 let candidates = teacher_candidates_from_execution(&review, &execution);
3969
3970 assert_eq!(candidates.len(), 1);
3971 assert_eq!(candidates[0].kind, TeacherCandidateKind::RegressionTest);
3972 assert_eq!(candidates[0].status, TeacherCandidateStatus::Proposed);
3973 assert!(
3974 candidates[0]
3975 .evidence
3976 .iter()
3977 .any(|line| { line.contains("cargo test failed with a replay mismatch") })
3978 );
3979 }
3980
3981 #[test]
3982 fn student_replay_promotes_only_on_delta() {
3983 let gate = PromotionGate {
3984 min_score_delta: 3,
3985 max_cost_delta_microusd: Some(25),
3986 ..PromotionGate::default()
3987 };
3988 let replay = StudentReplayResult {
3989 trace_id: "trace-a".to_string(),
3990 candidate_id: "teacher-review:branch-a".to_string(),
3991 baseline: StudentReplayMetrics {
3992 score: 80,
3993 cost_microusd: 100,
3994 },
3995 candidate: StudentReplayMetrics {
3996 score: 84,
3997 cost_microusd: 120,
3998 },
3999 required_tests: vec![StudentReplayTestResult {
4000 name: "workflow replay".to_string(),
4001 passed: true,
4002 }],
4003 policy_violations: Vec::new(),
4004 stale: false,
4005 notes: None,
4006 };
4007
4008 let promoted = gate.evaluate_replay("teacher-review:branch-a", &replay);
4009 assert!(promoted.promoted());
4010 assert_eq!(promoted.status, TeacherCandidateStatus::Promoted);
4011 assert_eq!(promoted.score_delta, 4);
4012
4013 let weak_replay = StudentReplayResult {
4014 candidate: StudentReplayMetrics {
4015 score: 82,
4016 cost_microusd: 120,
4017 },
4018 ..replay
4019 };
4020 let rejected = gate.evaluate_replay("teacher-review:branch-a", &weak_replay);
4021 assert!(!rejected.promoted());
4022 assert_eq!(rejected.status, TeacherCandidateStatus::Rejected);
4023 assert!(
4024 rejected
4025 .reasons
4026 .iter()
4027 .any(|reason| reason.contains("below required 3"))
4028 );
4029 }
4030
4031 #[test]
4032 fn promotion_gate_rejects_stale_policy_cost_and_failed_tests() {
4033 let gate = PromotionGate {
4034 min_score_delta: 1,
4035 max_cost_delta_microusd: Some(10),
4036 ..PromotionGate::default()
4037 };
4038 let replay = StudentReplayResult {
4039 trace_id: "trace-a".to_string(),
4040 candidate_id: "teacher-review:branch-a".to_string(),
4041 baseline: StudentReplayMetrics {
4042 score: 70,
4043 cost_microusd: 10,
4044 },
4045 candidate: StudentReplayMetrics {
4046 score: 90,
4047 cost_microusd: 30,
4048 },
4049 required_tests: vec![StudentReplayTestResult {
4050 name: "required regression".to_string(),
4051 passed: false,
4052 }],
4053 policy_violations: vec!["writes outside file scope".to_string()],
4054 stale: true,
4055 notes: None,
4056 };
4057
4058 let decision = gate.evaluate_replay("teacher-review:branch-a", &replay);
4059
4060 assert_eq!(decision.status, TeacherCandidateStatus::Rejected);
4061 assert!(
4062 decision
4063 .reasons
4064 .iter()
4065 .any(|reason| { reason.contains("cost delta 20 exceeds allowed 10") })
4066 );
4067 assert!(
4068 decision
4069 .reasons
4070 .iter()
4071 .any(|reason| { reason.contains("required test `required regression` failed") })
4072 );
4073 assert!(
4074 decision
4075 .reasons
4076 .iter()
4077 .any(|reason| { reason.contains("policy violation: writes outside file scope") })
4078 );
4079 assert!(
4080 decision
4081 .reasons
4082 .iter()
4083 .any(|reason| { reason.contains("student replay result is stale") })
4084 );
4085 }
4086
4087 #[test]
4088 fn promotion_gate_requires_recorded_replay_before_candidate_promotion() {
4089 let candidate = TeacherCandidate {
4090 candidate_id: "teacher-review:branch-a".to_string(),
4091 kind: TeacherCandidateKind::WorkflowRecipe,
4092 status: TeacherCandidateStatus::Proposed,
4093 source_node_id: "branch-a".to_string(),
4094 source_branch_id: Some("branch-a".to_string()),
4095 summary: "candidate waits for replay".to_string(),
4096 evidence: Vec::new(),
4097 replay_results: Vec::new(),
4098 };
4099
4100 let decision = PromotionGate::default().evaluate_candidate(&candidate);
4101
4102 assert_eq!(decision.status, TeacherCandidateStatus::Rejected);
4103 assert_eq!(
4104 decision.reasons,
4105 vec!["no student replay result recorded".to_string()]
4106 );
4107 }
4108
4109 #[test]
4110 fn tournament_selects_passing_minimal_branch() {
4111 let tournament = BranchTournament {
4112 min_score: 60,
4113 ordering: TournamentOrdering::CostThenScore,
4114 };
4115 let candidates = vec![
4116 candidate(
4117 "expensive-pass",
4118 WorkflowRunStatus::Succeeded,
4119 90,
4120 90,
4121 "quality",
4122 ),
4123 candidate("failed-cheap", WorkflowRunStatus::Failed, 100, 1, "broken"),
4124 candidate(
4125 "cheap-pass",
4126 WorkflowRunStatus::Succeeded,
4127 70,
4128 10,
4129 "minimal",
4130 ),
4131 candidate("too-low", WorkflowRunStatus::Succeeded, 40, 2, "weak"),
4132 ];
4133
4134 let selected = tournament
4135 .select(&candidates)
4136 .expect("one passing branch should be selected");
4137
4138 assert_eq!(selected.branch_id, "cheap-pass");
4139 }
4140
4141 #[test]
4142 fn tournament_can_select_score_before_cost_explicitly() {
4143 let tournament = BranchTournament {
4144 min_score: 60,
4145 ordering: TournamentOrdering::ScoreThenCost,
4146 };
4147 let candidates = vec![
4148 candidate("quality", WorkflowRunStatus::Succeeded, 95, 100, "quality"),
4149 candidate("minimal", WorkflowRunStatus::Succeeded, 70, 10, "minimal"),
4150 ];
4151
4152 let selected = tournament
4153 .select(&candidates)
4154 .expect("one passing branch should be selected");
4155
4156 assert_eq!(selected.branch_id, "quality");
4157 }
4158
4159 #[test]
4160 fn pareto_frontier_keeps_diverse_candidates() {
4161 let frontier = ParetoFrontier { max_items: 4 };
4162 let candidates = vec![
4163 candidate("quality", WorkflowRunStatus::Succeeded, 95, 100, "quality"),
4164 candidate("minimal", WorkflowRunStatus::Succeeded, 70, 10, "small"),
4165 candidate("dominated", WorkflowRunStatus::Succeeded, 60, 40, "middle"),
4166 candidate("failed", WorkflowRunStatus::Failed, 100, 1, "broken"),
4167 ];
4168
4169 let selected = frontier.select(&candidates);
4170
4171 assert_eq!(
4172 selected
4173 .iter()
4174 .map(|candidate| candidate.branch_id.as_str())
4175 .collect::<Vec<_>>(),
4176 vec!["quality", "minimal"]
4177 );
4178 assert_eq!(
4179 selected
4180 .iter()
4181 .filter_map(|candidate| candidate.diversity_key.as_deref())
4182 .collect::<Vec<_>>(),
4183 vec!["quality", "small"]
4184 );
4185 }
4186 }
4187
4187 lines RUST