use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::{Arc, Mutex, MutexGuard, OnceLock};
use std::time::{SystemTime, UNIX_EPOCH};
use async_trait::async_trait;
use codewhale_workflow::{
AgentType, BranchResult, BranchSpec, BudgetSpec, ControlNodeKind, ControlNodeResult,
FleetRoleMap, GateKind, GateOn, GateOutcome, GateSpec, GateState, GateStatusLine,
HandoffArtifact, LaneGateBoard, LeafResult, LeafSpec, ReduceSpec, SequenceSpec, TaskMode,
WorkflowExecution as IrWorkflowExecution, WorkflowMemoUsage, WorkflowNode,
WorkflowRunStatus as IrWorkflowRunStatus, WorkflowSpec, WorkflowUsage,
compile_javascript_workflow, compile_typescript_workflow, leaf_wants_worktree,
resolve_workflow_agent,
};
use codewhale_workflow_js::{
BudgetSnapshot, DriverError, ProgressEvent, SpawnedTask, TaskCompletion, TaskRequest,
WORKFLOW_MAX_CONCURRENT, WorkflowDriver, WorkflowRunCancel, WorkflowVm,
};
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc, oneshot};
use uuid::Uuid;
use crate::core::events::Event;
use crate::tools::spec::{
ApprovalRequirement, ToolCapability, ToolContext, ToolError, ToolResult, ToolSpec,
optional_bool, optional_str, optional_u64,
};
use crate::tools::subagent::{
SharedSubAgentManager, SubAgentCompletion, SubAgentManager, SubAgentResult, SubAgentRuntime,
SubAgentStatus, WorkflowTaskSpawnIdentity, WorkflowTaskSpawnMetadata, spawn_workflow_task,
};
use crate::tools::verifier::run_workflow_completion_gates;
use crate::tools::workflow_plan_approval::{
WorkflowPlanApprovalReceipt, analyze_workflow_plan_approval_with_config, analyze_workflow_spec,
workflow_approval_requirement_for,
};
use crate::utils::spawn_supervised;
use crate::work_graph::{
CancelOutcome, EvidenceKind, EvidenceRef, OperationIntent, OperationObservation,
OperationOwnerSnapshot, OwnerState, SharedWorkRuntime,
};
const WORKFLOW_HANDOFF_MAX_CHARS: usize = 4_000;
const WORKFLOW_RESULT_EVENTS_TAIL: usize = 50;
const WORKFLOW_RESULT_PROGRESS_TAIL: usize = 20;
const WORKFLOW_RESULT_DISPATCH_FAILURES_TAIL: usize = 12;
const WORKFLOW_RESULT_DISPATCH_FAILURE_FIELD_MAX_CHARS: usize = 320;
const WORKFLOW_RESULT_VALUE_MAX_CHARS: usize = 4_000;
const WORKFLOW_RESULT_LEAF_OUTPUT_MAX_CHARS: usize = 500;
const WORKFLOW_RESULT_MAX_CHARS: usize = 24_000;
const WORKFLOW_RUN_EVENTS_MAX_RETAINED: usize = 1_000;
#[derive(Clone)]
pub struct WorkflowTool {
manager: SharedSubAgentManager,
runtime: SubAgentRuntime,
approval_decision: &'static str,
}
impl WorkflowTool {
#[must_use]
pub fn new(manager: SharedSubAgentManager, runtime: SubAgentRuntime) -> Self {
Self {
manager,
runtime,
approval_decision: "approved",
}
}
#[must_use]
pub(crate) fn with_explicit_cli_approval(mut self) -> Self {
self.approval_decision = "approved_explicit_cli_command";
self
}
}
type SharedWorkflowRuns = Arc<Mutex<HashMap<String, WorkflowRunRecord>>>;
type SharedWorkflowControllers = Arc<Mutex<HashMap<String, Arc<WorkflowRunController>>>>;
type SharedWorkflowLifecycles = Arc<Mutex<HashMap<String, WorkflowWorkLifecycle>>>;
#[derive(Clone)]
struct WorkflowWorkLifecycle {
work: SharedWorkRuntime,
session_id: String,
external: String,
}
impl WorkflowWorkLifecycle {
fn register(
context: &ToolContext,
run_id: &str,
title: &str,
) -> Result<Option<Self>, ToolError> {
let Some(work) = context.runtime.work.clone() else {
return Ok(None);
};
let lifecycle = Self {
work,
session_id: context.state_namespace.clone(),
external: format!("workflow:{run_id}"),
};
lifecycle
.work
.register_operation(
&lifecycle.session_id,
OperationIntent::new(
lifecycle.external.clone(),
title,
true,
"workflow",
format!("workflow:{run_id}:start"),
),
)
.map_err(ToolError::execution_failed)?;
Ok(Some(lifecycle))
}
fn for_bound(context: &ToolContext, run_id: &str) -> Option<Self> {
let work = context.runtime.work.clone()?;
let external = format!("workflow:{run_id}");
work.has_operation_binding(Some(&context.state_namespace), &external)
.then(|| Self {
work,
session_id: context.state_namespace.clone(),
external,
})
}
fn reconcile_record(&self, record: &WorkflowRunRecord) -> Result<bool, String> {
let output = record.result.as_ref().and_then(|result| {
serde_json::to_vec(result).ok().and_then(|bytes| {
EvidenceRef::new(
EvidenceKind::Receipt {
owner: "workflow".to_string(),
},
format!("workflow:{}:result", record.run_id),
Some(u64::try_from(bytes.len()).unwrap_or(u64::MAX)),
false,
)
.ok()
})
});
let state = match record.status {
WorkflowRunStatus::Running => OwnerState::Running,
WorkflowRunStatus::Completed | WorkflowRunStatus::Degraded => OwnerState::Completed,
WorkflowRunStatus::Failed => OwnerState::Failed,
WorkflowRunStatus::Cancelled => OwnerState::Cancelled,
};
let mut snapshot = OperationOwnerSnapshot::new(
self.external.clone(),
state,
record.lifecycle_seq,
i64::try_from(record.completed_at_ms.unwrap_or(record.started_at_ms))
.unwrap_or(i64::MAX),
);
if let Some(output) = output {
snapshot = snapshot.with_output(output);
}
self.work.reconcile_operation(&self.session_id, snapshot)
}
fn reconcile_cancel(&self, outcome: CancelOutcome) -> Result<bool, String> {
self.work.reconcile_observation(
&self.session_id,
&self.external,
OperationObservation::CancelUpdate {
outcome,
at: i64::try_from(now_ms()).unwrap_or(i64::MAX),
},
)
}
fn reconcile_spawn_failure(&self) {
let _ = self.work.reconcile_operation(
&self.session_id,
OperationOwnerSnapshot::new(
self.external.clone(),
OwnerState::Failed,
1,
i64::try_from(now_ms()).unwrap_or(i64::MAX),
),
);
}
fn reconcile_missing(&self) {
let _ = self.work.reconcile_observation(
&self.session_id,
&self.external,
OperationObservation::OwnerMissing {
checked_at: i64::try_from(now_ms()).unwrap_or(i64::MAX),
},
);
}
}
struct WorkflowRunController {
driver: Arc<SubAgentWorkflowDriver>,
vm_cancel: WorkflowRunCancel,
run_handle: Mutex<Option<tokio::task::JoinHandle<()>>>,
}
impl WorkflowRunController {
fn new(driver: Arc<SubAgentWorkflowDriver>, vm_cancel: WorkflowRunCancel) -> Self {
Self {
driver,
vm_cancel,
run_handle: Mutex::new(None),
}
}
fn set_run_handle(&self, handle: tokio::task::JoinHandle<()>) {
if let Ok(mut guard) = self.run_handle.lock() {
*guard = Some(handle);
}
}
fn cancel(&self) {
self.vm_cancel.cancel();
self.driver.finalize_running_tasks_cancelled();
self.driver.force_cancel_all();
if let Ok(mut guard) = self.run_handle.lock()
&& let Some(handle) = guard.take()
{
handle.abort();
}
}
}
#[derive(Debug, Clone, Serialize)]
struct WorkflowRunSummary {
run_id: String,
status: WorkflowRunStatus,
lifecycle_seq: u64,
started_at_ms: u64,
completed_at_ms: Option<u64>,
source_path: Option<PathBuf>,
workflow_id: Option<String>,
workflow_goal: Option<String>,
token_budget: Option<u64>,
child_count: usize,
schema_error_count: usize,
dispatch_failure_count: usize,
progress_count: usize,
last_progress: Option<String>,
event_count: usize,
last_event_type: Option<String>,
leaf_count: usize,
branch_count: usize,
control_count: usize,
execution_status: Option<IrWorkflowRunStatus>,
gate_count: usize,
blocked_gate_count: usize,
gate_status: Vec<GateStatusLine>,
error: Option<String>,
usage: Option<WorkflowRunUsage>,
events_dropped: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct WorkflowSchemaError {
task_id: String,
message: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct WorkflowDispatchFailure {
at_ms: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
label: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
phase: Option<String>,
message: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct WorkflowUiEvent {
at_ms: u64,
#[serde(flatten)]
kind: WorkflowUiEventKind,
}
impl WorkflowUiEvent {
fn new(kind: WorkflowUiEventKind) -> Self {
Self {
at_ms: now_ms(),
kind,
}
}
fn at(at_ms: u64, kind: WorkflowUiEventKind) -> Self {
Self { at_ms, kind }
}
fn event_type(&self) -> &'static str {
self.kind.event_type()
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
enum WorkflowUiEventKind {
RunStarted {
workflow_id: Option<String>,
workflow_goal: Option<String>,
source_path: Option<PathBuf>,
token_budget: Option<u64>,
},
RunCompleted {
status: WorkflowRunStatus,
error: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
usage: Option<WorkflowRunUsage>,
},
RunCancelled {
reason: String,
},
PhaseStarted {
title: String,
},
TaskStarted(Box<WorkflowTaskStartedEvent>),
TaskCompleted {
task_id: String,
status: IrWorkflowRunStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
usage: Option<WorkflowTaskUsage>,
},
GateUpdated {
gate_id: String,
role: String,
gate: String,
state: String,
blocked_role: Option<String>,
blocked_reason: Option<String>,
},
HandoffPromoted {
artifact_id: String,
gate_id: String,
kind: String,
from_role: String,
to_role: String,
producer_task_id: String,
},
HandoffConsumed {
artifact_id: String,
kind: String,
from_role: String,
to_role: String,
consumer_task_id: String,
},
TaskSchemaValidationFailed {
task_id: String,
message: String,
},
TaskDispatchFailed {
#[serde(default, skip_serializing_if = "Option::is_none")]
label: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
phase: Option<String>,
message: String,
},
BudgetUpdated {
total: Option<u64>,
spent: u64,
remaining: Option<u64>,
},
Log {
message: String,
},
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
struct WorkflowTaskUsage {
#[serde(default, skip_serializing_if = "Option::is_none")]
input_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
output_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
total_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
cost_microusd: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
tool_calls: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
duration_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
result_ref: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
token_source: Option<WorkflowTokenSource>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
enum WorkflowTokenSource {
ProviderReported,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
struct WorkflowRunUsage {
#[serde(default, skip_serializing_if = "Option::is_none")]
input_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
output_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
total_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
cost_microusd: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
tool_calls: Option<u64>,
#[serde(default)]
tasks_reported: u64,
}
impl WorkflowRunUsage {
fn from_task(usage: &WorkflowTaskUsage) -> Self {
Self {
input_tokens: usage.input_tokens,
output_tokens: usage.output_tokens,
total_tokens: usage.total_tokens,
cost_microusd: usage.cost_microusd,
tool_calls: usage.tool_calls.map(u64::from),
tasks_reported: 1,
}
}
fn add_task(&mut self, usage: &WorkflowTaskUsage) {
self.input_tokens = sum_optional_usage(self.input_tokens, usage.input_tokens);
self.output_tokens = sum_optional_usage(self.output_tokens, usage.output_tokens);
self.total_tokens = sum_optional_usage(self.total_tokens, usage.total_tokens);
self.cost_microusd = sum_optional_usage(self.cost_microusd, usage.cost_microusd);
self.tool_calls = sum_optional_usage(self.tool_calls, usage.tool_calls.map(u64::from));
self.tasks_reported = self.tasks_reported.saturating_add(1);
}
}
fn sum_optional_usage(left: Option<u64>, right: Option<u64>) -> Option<u64> {
match (left, right) {
(Some(left), Some(right)) => Some(left.saturating_add(right)),
(Some(value), None) | (None, Some(value)) => Some(value),
(None, None) => None,
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct WorkflowTaskStartedEvent {
task_id: String,
label: Option<String>,
role: Option<String>,
profile: Option<String>,
model: Option<String>,
strength: Option<String>,
thinking: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
requested_reasoning: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
effective_reasoning: Option<String>,
resolved_role: Option<String>,
resolved_profile: Option<String>,
resolved_provider: String,
resolved_model: String,
route_source: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
child_route: Option<crate::tools::subagent::ChildRouteReceipt>,
worktree: bool,
workspace: Option<PathBuf>,
git_branch: Option<String>,
parent_task_id: Option<String>,
depth: u32,
workflow_run_id: Option<String>,
workflow_phase_id: Option<String>,
workflow_task_label: Option<String>,
workflow_child_index: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
fleet_receipt: Option<codewhale_workflow::FleetTaskReceipt>,
}
impl WorkflowUiEventKind {
fn event_type(&self) -> &'static str {
match self {
Self::RunStarted { .. } => "run_started",
Self::RunCompleted { .. } => "run_completed",
Self::RunCancelled { .. } => "run_cancelled",
Self::PhaseStarted { .. } => "phase_started",
Self::TaskStarted(_) => "task_started",
Self::TaskCompleted { .. } => "task_completed",
Self::GateUpdated { .. } => "gate_updated",
Self::HandoffPromoted { .. } => "handoff_promoted",
Self::HandoffConsumed { .. } => "handoff_consumed",
Self::TaskSchemaValidationFailed { .. } => "task_schema_validation_failed",
Self::TaskDispatchFailed { .. } => "task_dispatch_failed",
Self::BudgetUpdated { .. } => "budget_updated",
Self::Log { .. } => "log",
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct WorkflowRunRecord {
run_id: String,
status: WorkflowRunStatus,
#[serde(default)]
lifecycle_seq: u64,
started_at_ms: u64,
completed_at_ms: Option<u64>,
source_path: Option<PathBuf>,
workflow_id: Option<String>,
workflow_goal: Option<String>,
token_budget: Option<u64>,
child_ids: Vec<String>,
progress: Vec<String>,
#[serde(default)]
events: Vec<WorkflowUiEvent>,
schema_errors: Vec<WorkflowSchemaError>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
dispatch_failures: Vec<WorkflowDispatchFailure>,
result: Option<Value>,
execution: Option<IrWorkflowExecution>,
error: Option<String>,
#[serde(default)]
verify_on_complete: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
verification: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
plan_approval: Option<WorkflowPlanApprovalReceipt>,
#[serde(default)]
gate_status: Vec<GateStatusLine>,
#[serde(default, skip_serializing_if = "Option::is_none")]
usage: Option<WorkflowRunUsage>,
#[serde(default)]
events_total: u64,
#[serde(default)]
events_dropped: u64,
}
impl WorkflowRunRecord {
fn new(
run_id: String,
source_path: Option<PathBuf>,
token_budget: Option<u64>,
spec: Option<&WorkflowSpec>,
) -> Self {
let gate_status = spec
.map(|spec| initial_gate_status(&run_id, &spec.gates))
.unwrap_or_default();
Self {
run_id,
status: WorkflowRunStatus::Running,
lifecycle_seq: 1,
started_at_ms: now_ms(),
completed_at_ms: None,
source_path,
workflow_id: spec.and_then(|spec| spec.id.clone()),
workflow_goal: spec.map(|spec| spec.goal.clone()),
token_budget,
child_ids: Vec::new(),
progress: Vec::new(),
events: Vec::new(),
schema_errors: Vec::new(),
dispatch_failures: Vec::new(),
result: None,
execution: None,
error: None,
verify_on_complete: false,
verification: None,
plan_approval: None,
gate_status,
usage: None,
events_total: 0,
events_dropped: 0,
}
}
fn push_event(&mut self, event: WorkflowUiEvent) {
self.events_total = self.events_total.saturating_add(1);
self.events.push(event);
if self.events.len() > WORKFLOW_RUN_EVENTS_MAX_RETAINED {
let overflow = self.events.len() - WORKFLOW_RUN_EVENTS_MAX_RETAINED;
self.events.drain(..overflow);
self.events_dropped = self
.events_dropped
.saturating_add(u64::try_from(overflow).unwrap_or(u64::MAX));
}
}
fn summary(&self) -> WorkflowRunSummary {
WorkflowRunSummary {
run_id: self.run_id.clone(),
status: self.status,
lifecycle_seq: self.lifecycle_seq,
started_at_ms: self.started_at_ms,
completed_at_ms: self.completed_at_ms,
source_path: self.source_path.clone(),
workflow_id: self.workflow_id.clone(),
workflow_goal: self.workflow_goal.clone(),
token_budget: self.token_budget,
child_count: self.child_ids.len(),
schema_error_count: self.schema_errors.len(),
dispatch_failure_count: self.dispatch_failures.len(),
progress_count: self.progress.len(),
last_progress: self.progress.last().cloned(),
event_count: usize::try_from(self.events_total.max(self.events.len() as u64))
.unwrap_or(usize::MAX),
last_event_type: self
.events
.last()
.map(|event| event.event_type().to_string()),
leaf_count: self
.execution
.as_ref()
.map(|execution| execution.leaf_results.len())
.unwrap_or_default(),
branch_count: self
.execution
.as_ref()
.map(|execution| execution.branch_results.len())
.unwrap_or_default(),
control_count: self
.execution
.as_ref()
.map(|execution| execution.control_node_results.len())
.unwrap_or_default(),
execution_status: self.execution.as_ref().map(|execution| execution.status),
gate_count: self.gate_status.len(),
blocked_gate_count: self
.gate_status
.iter()
.filter(|line| line.blocked_reason.is_some())
.count(),
gate_status: self.gate_status.clone(),
error: self.error.clone(),
usage: self.usage.clone(),
events_dropped: self.events_dropped,
}
}
}
fn initial_gate_status(run_id: &str, gates: &[GateSpec]) -> Vec<GateStatusLine> {
if gates.is_empty() {
return Vec::new();
}
let mut board = LaneGateBoard::new(run_id);
board.install_gates(gates);
board.status_summary()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
enum WorkflowRunStatus {
Running,
Completed,
Degraded,
Failed,
Cancelled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum WorkflowAction {
Start,
Run,
Status,
Cancel,
}
fn parse_workflow_action(input: &Value) -> Result<WorkflowAction, ToolError> {
let Some(action) = optional_str(input, "action")? else {
return Ok(WorkflowAction::Start);
};
match action.trim().to_ascii_lowercase().as_str() {
"" | "start" | "spawn" => Ok(WorkflowAction::Start),
"run" | "wait" => Ok(WorkflowAction::Run),
"status" | "list" | "inspect" => Ok(WorkflowAction::Status),
"cancel" | "stop" | "abort" => Ok(WorkflowAction::Cancel),
other => Err(ToolError::invalid_input(format!(
"Invalid workflow action '{other}'. Use start, run, status, or cancel."
))),
}
}
#[async_trait]
impl ToolSpec for WorkflowTool {
fn name(&self) -> &'static str {
"workflow"
}
fn description(&self) -> &'static str {
concat!(
"Start, run, inspect, or cancel a Workflow. Workflows execute deterministic JS with args, phase/log progress, and task(...) calls that dispatch real sub-agents through Fleet/sub-agent scheduling. ",
"For parallel fan-out, pass an array of zero-argument thunks exactly like `await parallel([() => task({...}), () => task({...})])`; do not pass task promises as variadic arguments. ",
"Provide exactly one of script, source_path, or plan (structured planner JSON). ",
"Use action=start for detached orchestration and action=status with run_id to inspect progress. Use action=run when the model needs the final result before continuing. ",
"Start a workflow on your own only for broad or staged work (the session [workflow].automatic table, default on). An explicit /workflow invocation is authorization. Do not start a workflow for one-file edits or simple questions."
)
}
fn input_schema(&self) -> Value {
json!({
"type": "object",
"properties": {
"action": {
"type": "string",
"enum": ["start", "run", "status", "cancel"],
"description": "start (default) launches a Workflow in the background. run waits for completion. status lists runs or inspects run_id. cancel stops a run and its child agents."
},
"run_id": {
"type": "string",
"description": "Workflow run id for action=status or action=cancel."
},
"script": {
"type": "string",
"description": "Workflow JS source. The runtime provides args, task(...), parallel(thunks), pipeline(thunks), log(...), phase(...), and budget. Fan-out syntax: await parallel([() => task({...}), () => task({...})]). parallel() requires one array of zero-argument thunks, not variadic task promises."
},
"source_path": {
"type": "string",
"description": "Path to a .workflow.js script inside the workspace. Use instead of script for checked-in workflows."
},
"fleet": {
"type": "string",
"description": "Named Fleet to resolve task({ role }) declarations, loaded from $CODEWHALE_HOME/fleets/ or workspace fleets/. Accepts a qualified origin/name. A legacy roster maps roles to profiles. An exact Fleet (schema = \"exact\") is frozen at start: each member's provider, model, reasoning, and permission ceiling are fixed, and per-task model/thinking overrides are rejected."
},
"plan": {
"type": "object",
"description": "Structured planner plan JSON (#4124). Alternative to script/source_path. Accepts goal, risk, max_children, token_budget, phases[], and/or children[] (or IR nodes). risk must be exactly read_only, writes, or elevated. For a child, prefer role/profile without an explicit type; do not combine a role/profile with a conflicting type. Lowered to Workflow JS with parallel() partial-success semantics."
},
"args": {
"anyOf": [
{ "type": "null" },
{ "type": "boolean" },
{ "type": "integer" },
{ "type": "number" },
{ "type": "string" },
{ "type": "array" },
{
"type": "object",
"additionalProperties": {}
}
],
"description": "JSON value exposed to the script as args. Defaults to null."
},
"token_budget": {
"type": "integer",
"minimum": 1,
"description": "Optional shared Workflow admission hint. Usage is reconciled when children report completion; already-running parallel children can take aggregate spent past the hint, while later and descendant spawns are rejected once exhausted."
},
"wait": {
"type": "boolean",
"description": "For action=start, wait for completion instead of returning immediately."
},
"verify": {
"type": "boolean",
"default": false,
"description": "After a successful workflow completion, run quick workspace verifier gates (auto/quick profile)."
}
},
"required": [],
"additionalProperties": false
})
}
fn capabilities(&self) -> Vec<ToolCapability> {
vec![
ToolCapability::ExecutesCode,
ToolCapability::RequiresApproval,
]
}
fn approval_requirement(&self) -> ApprovalRequirement {
ApprovalRequirement::Required
}
fn approval_requirement_for(&self, input: &Value) -> ApprovalRequirement {
let config = workflow_config_for(&self.runtime);
workflow_approval_requirement_for(input, &config)
}
fn starts_detached_for(&self, input: &Value) -> bool {
matches!(parse_workflow_action(input), Ok(WorkflowAction::Start))
&& !optional_bool(input, "wait", false).unwrap_or(true)
}
fn supports_parallel_for(&self, input: &Value) -> bool {
matches!(parse_workflow_action(input), Ok(WorkflowAction::Status))
}
fn is_read_only_for(&self, input: &Value) -> bool {
matches!(parse_workflow_action(input), Ok(WorkflowAction::Status))
}
async fn execute(&self, input: Value, context: &ToolContext) -> Result<ToolResult, ToolError> {
let state = shared_workflow_state(&context.workspace);
attach_bound_workflow_lifecycles(context, &state)?;
let action = parse_workflow_action(&input)?;
codewhale_telemetry::session_counters().bump(codewhale_telemetry::Counter::WorkflowRun);
match action {
WorkflowAction::Start => {
let wait = optional_bool(&input, "wait", false)?;
start_workflow(
input,
context,
self.manager.clone(),
self.runtime.clone(),
state,
wait,
self.approval_decision,
)
.await
}
WorkflowAction::Run => {
start_workflow(
input,
context,
self.manager.clone(),
self.runtime.clone(),
state,
true,
self.approval_decision,
)
.await
}
WorkflowAction::Status => status_workflow(input, state),
WorkflowAction::Cancel => cancel_workflow(input, state).await,
}
}
}
fn attach_bound_workflow_lifecycles(
context: &ToolContext,
state: &Arc<WorkflowWorkspaceState>,
) -> Result<(), ToolError> {
let records = lock_mutex(&state.runs)?
.values()
.cloned()
.collect::<Vec<_>>();
for record in records {
if let Some(lifecycle) = WorkflowWorkLifecycle::for_bound(context, &record.run_id) {
state.attach_lifecycle(&record.run_id, lifecycle);
state.reconcile_snapshot(&record);
}
}
Ok(())
}
fn fail_workflow_start(state: &Arc<WorkflowWorkspaceState>, run_id: &str, message: String) {
let snapshot = state.runs.lock().ok().and_then(|mut runs| {
let record = runs.get_mut(run_id)?;
record.status = WorkflowRunStatus::Failed;
record.lifecycle_seq = record.lifecycle_seq.saturating_add(1);
record.completed_at_ms = Some(now_ms());
record.error = Some(message);
Some(record.clone())
});
let Some(snapshot) = snapshot else {
state.mark_owner_missing(run_id);
return;
};
if state.try_record_snapshot(&snapshot).is_ok() {
state.reconcile_snapshot(&snapshot);
} else {
state.mark_owner_missing(run_id);
}
}
fn fail_workflow_after_controller_registration(
state: &Arc<WorkflowWorkspaceState>,
run_id: &str,
controller: &Arc<WorkflowRunController>,
message: String,
) {
controller.cancel();
if let Ok(mut controllers) = state.controllers.lock() {
controllers.remove(run_id);
}
fail_workflow_start(state, run_id, message);
}
#[allow(clippy::too_many_arguments)]
async fn start_workflow(
input: Value,
context: &ToolContext,
manager: SharedSubAgentManager,
runtime: SubAgentRuntime,
state: Arc<WorkflowWorkspaceState>,
wait: bool,
approval_decision: &str,
) -> Result<ToolResult, ToolError> {
let source = workflow_source(&input, context)?;
let args = input.get("args").cloned().unwrap_or(Value::Null);
let token_budget = optional_u64(&input, "token_budget", 0)?;
let token_budget = (token_budget > 0).then_some(token_budget);
let verify_on_complete = optional_bool(&input, "verify", false)?;
let fleet = workflow_fleet_binding(&input, context, runtime.api_config.as_deref())?;
let run_id = format!("workflow_{}", &Uuid::new_v4().to_string()[..8]);
let gate_specs = source
.spec
.as_ref()
.map(|spec| spec.gates.clone())
.unwrap_or_default();
let workflow_cfg = workflow_config_for(&runtime);
let summary = source
.spec
.as_ref()
.map(|spec| analyze_workflow_spec(spec, token_budget, &workflow_cfg))
.unwrap_or_else(|| analyze_workflow_plan_approval_with_config(&input, &workflow_cfg));
let approval_decision = if summary.is_read_only_envelope() {
"auto_read_only"
} else {
approval_decision
};
let plan_approval = summary.to_receipt(approval_decision, now_ms());
let workflow_title = source
.spec
.as_ref()
.map(|spec| spec.goal.as_str())
.or_else(|| {
source
.path
.as_ref()
.and_then(|path| path.file_name()?.to_str())
})
.unwrap_or("Workflow run");
let lifecycle = WorkflowWorkLifecycle::register(context, &run_id, workflow_title)?;
{
let mut runs_guard = match lock_mutex(&state.runs) {
Ok(guard) => guard,
Err(err) => {
if let Some(lifecycle) = lifecycle.as_ref() {
lifecycle.reconcile_spawn_failure();
}
return Err(err);
}
};
let mut record = WorkflowRunRecord::new(
run_id.clone(),
source.path.clone(),
token_budget,
source.spec.as_ref(),
);
record.verify_on_complete = verify_on_complete;
record.plan_approval = Some(plan_approval.clone());
let started = WorkflowUiEvent::at(
record.started_at_ms,
WorkflowUiEventKind::RunStarted {
workflow_id: record.workflow_id.clone(),
workflow_goal: record.workflow_goal.clone(),
source_path: record.source_path.clone(),
token_budget: record.token_budget,
},
);
record.push_event(started.clone());
runs_guard.insert(run_id.clone(), record.clone());
if let Err(err) = state.try_record_snapshot(&record) {
runs_guard.remove(&run_id);
if let Some(lifecycle) = lifecycle.as_ref() {
lifecycle.reconcile_spawn_failure();
}
return Err(ToolError::execution_failed(format!(
"workflow journal snapshot failed before launch: {err}"
)));
}
if let Some(tx) = runtime.event_tx.as_ref()
&& let Ok(mut value) = serde_json::to_value(&started)
{
if let Some(obj) = value.as_object_mut() {
obj.insert("run_id".to_string(), json!(run_id));
}
let _ = tx.try_send(Event::WorkflowUi {
run_id: run_id.clone(),
event: value,
});
}
}
if let Some(lifecycle) = lifecycle {
state.attach_lifecycle(&run_id, lifecycle);
}
let mut runtime = runtime;
if let Some(operation) = fleet.exact() {
runtime.fleet_roster = operation.roster().clone();
}
let driver = SubAgentWorkflowDriver::new(
run_id.clone(),
manager,
runtime,
state.clone(),
token_budget,
fleet,
gate_specs,
);
let vm_cancel = WorkflowRunCancel::new();
let controller = Arc::new(WorkflowRunController::new(
driver.clone(),
vm_cancel.clone(),
));
if let Err(err) = lock_mutex(&state.controllers).map(|mut controllers_guard| {
controllers_guard.insert(run_id.clone(), controller.clone());
}) {
fail_workflow_start(&state, &run_id, err.to_string());
return Err(err);
}
let running_snapshot = {
let mut runs_guard = match lock_mutex(&state.runs) {
Ok(guard) => guard,
Err(err) => {
fail_workflow_after_controller_registration(
&state,
&run_id,
&controller,
err.to_string(),
);
return Err(err);
}
};
let Some(record) = runs_guard.get_mut(&run_id) else {
drop(runs_guard);
fail_workflow_after_controller_registration(
&state,
&run_id,
&controller,
"workflow owner record disappeared before launch".to_string(),
);
return Err(ToolError::execution_failed(
"workflow owner record disappeared before launch",
));
};
record.lifecycle_seq = record.lifecycle_seq.saturating_add(1);
record.clone()
};
if let Err(err) = state.try_record_snapshot(&running_snapshot) {
fail_workflow_after_controller_registration(
&state,
&run_id,
&controller,
format!("workflow journal failed while activating owner: {err}"),
);
return Err(ToolError::execution_failed(format!(
"workflow journal failed while activating owner: {err}"
)));
}
state.reconcile_snapshot(&running_snapshot);
let run = run_workflow_vm(
run_id.clone(),
source.source,
source.spec,
args,
driver,
state.clone(),
context.clone(),
vm_cancel,
);
if wait {
run.await;
} else {
let handle = spawn_supervised("workflow-run", std::panic::Location::caller(), run);
controller.set_run_handle(handle);
}
workflow_result_for(&run_id, state)
}
fn status_workflow(
input: Value,
state: Arc<WorkflowWorkspaceState>,
) -> Result<ToolResult, ToolError> {
if let Some(run_id) = optional_str(&input, "run_id")? {
return workflow_result_for(run_id, state);
}
let mut summaries = {
let runs_guard = lock_mutex(&state.runs)?;
runs_guard
.values()
.map(WorkflowRunRecord::summary)
.collect::<Vec<_>>()
};
summaries.sort_by_key(|record| record.started_at_ms);
ToolResult::json(&json!({
"action": "status",
"count": summaries.len(),
"runs": summaries,
}))
.map_err(|err| ToolError::execution_failed(err.to_string()))
}
async fn cancel_workflow(
input: Value,
state: Arc<WorkflowWorkspaceState>,
) -> Result<ToolResult, ToolError> {
let run_id =
optional_str(&input, "run_id")?.ok_or_else(|| ToolError::missing_field("run_id"))?;
cancel_workflow_run(run_id, state)
}
fn cancel_workflow_run(
run_id: &str,
state: Arc<WorkflowWorkspaceState>,
) -> Result<ToolResult, ToolError> {
let controller = {
let mut controllers_guard = lock_mutex(&state.controllers)?;
controllers_guard.remove(run_id)
};
let current_status = {
let runs_guard = lock_mutex(&state.runs)?;
let record = runs_guard.get(run_id).ok_or_else(|| {
ToolError::invalid_input(format!("Unknown workflow run_id '{run_id}'"))
})?;
record.status
};
if current_status != WorkflowRunStatus::Running {
state.reconcile_cancel(run_id, CancelOutcome::AlreadyFinished);
if let Ok(runs_guard) = state.runs.lock()
&& let Some(record) = runs_guard.get(run_id)
{
state.reconcile_snapshot(record);
}
return workflow_result_for(run_id, state);
}
let live = controller.is_some();
state.reconcile_cancel(
run_id,
if live {
CancelOutcome::Requested
} else {
CancelOutcome::Acknowledged
},
);
if let Some(controller) = controller.as_ref() {
controller.cancel();
}
let reason = if live {
"cancelled by workflow tool"
} else {
"cancelled; no live process to stop"
}
.to_string();
let cancelled_event = WorkflowUiEvent::new(WorkflowUiEventKind::RunCancelled {
reason: reason.clone(),
});
let snapshot = {
let mut runs_guard = lock_mutex(&state.runs)?;
let record = runs_guard.get_mut(run_id).ok_or_else(|| {
ToolError::invalid_input(format!("Unknown workflow run_id '{run_id}'"))
})?;
record.status = WorkflowRunStatus::Cancelled;
record.lifecycle_seq = record.lifecycle_seq.saturating_add(1);
record.completed_at_ms = Some(now_ms());
record.error = Some(reason);
record.push_event(cancelled_event.clone());
record.clone()
};
if let Err(err) = state.try_record_snapshot(&snapshot) {
state.mark_owner_missing(run_id);
return Err(ToolError::execution_failed(format!(
"workflow cancellation journal failed: {err}"
)));
}
state.reconcile_snapshot(&snapshot);
if let Some(controller) = controller {
controller.driver.emit_ui_event(&cancelled_event);
}
workflow_result_for(run_id, state)
}
fn workflow_fleet_name(input: &Value) -> Result<Option<String>, ToolError> {
let named = match optional_str(input, "fleet")? {
Some(name) => Some(name),
None => match input.get("args") {
Some(args) => optional_str(args, "fleet")?,
None => None,
},
};
Ok(named
.map(str::trim)
.filter(|name| !name.is_empty())
.map(str::to_string))
}
#[derive(Debug, Clone, Default)]
enum WorkflowFleetBinding {
#[default]
None,
Legacy {
name: String,
roles: FleetRoleMap,
},
Exact(Arc<crate::fleet::exact::ExactFleetWorkflow>),
}
impl WorkflowFleetBinding {
fn name(&self) -> Option<String> {
match self {
Self::None => None,
Self::Legacy { name, .. } => Some(name.clone()),
Self::Exact(operation) => Some(operation.snapshot().fleet().qualified()),
}
}
fn legacy_roles(&self) -> Option<&FleetRoleMap> {
match self {
Self::Legacy { roles, .. } => Some(roles),
Self::None | Self::Exact(_) => None,
}
}
fn exact(&self) -> Option<&Arc<crate::fleet::exact::ExactFleetWorkflow>> {
match self {
Self::Exact(operation) => Some(operation),
Self::None | Self::Legacy { .. } => None,
}
}
}
fn workflow_fleet_binding(
input: &Value,
context: &ToolContext,
api_config: Option<&crate::config::Config>,
) -> Result<WorkflowFleetBinding, ToolError> {
let Some(name) = workflow_fleet_name(input)? else {
return Ok(WorkflowFleetBinding::None);
};
let roots = crate::fleet::exact::fleet_search_roots(&context.workspace);
let (document, id) = crate::fleet::exact::load_fleet_document(&name, &context.workspace)
.map_err(|err| {
ToolError::invalid_input(format!(
"Failed to load workflow fleet '{name}' from {}: {err}",
roots
.iter()
.map(|root| format!("{}/{}", root.origin, root.root.display()))
.collect::<Vec<_>>()
.join(", ")
))
})?;
if let Some(legacy) = document.legacy() {
let roles = FleetRoleMap::from_pairs(
legacy
.roles
.iter()
.map(|(role, profile)| (role.as_str(), profile.as_str())),
)
.map_err(|err| ToolError::invalid_input(err.to_string()))?;
return Ok(WorkflowFleetBinding::Legacy { name, roles });
}
let operation = crate::fleet::exact::ExactFleetWorkflow::capture(
&document,
id,
chrono::Utc::now().to_rfc3339(),
api_config,
&roots,
)
.map_err(ToolError::invalid_input)?;
Ok(WorkflowFleetBinding::Exact(Arc::new(operation)))
}
fn apply_named_fleet_to_task_request(
fleet_roles: Option<&FleetRoleMap>,
request: &mut TaskRequest,
) -> Result<(), DriverError> {
let Some(fleet_roles) = fleet_roles else {
return Ok(());
};
let resolved = resolve_workflow_agent(
request.role.as_deref(),
request.profile.as_deref(),
fleet_roles,
true,
)
.map_err(|err| DriverError::Rejected(err.to_string()))?;
request.role = resolved.resolved_role;
request.profile = Some(resolved.resolved_profile);
Ok(())
}
fn bind_exact_fleet_task_request(
operation: &crate::fleet::exact::ExactFleetWorkflow,
session: codewhale_workflow::PermissionCeiling,
request: &mut TaskRequest,
) -> Result<crate::fleet::exact::ExactMemberBinding, DriverError> {
let fleet = operation.snapshot().fleet().qualified();
for (field, present) in [
("model", request.model.is_some()),
("model_strength", request.model_strength.is_some()),
("thinking", request.thinking.is_some()),
("subagent_type", request.subagent_type.is_some()),
("allowed_tools", request.allowed_tools.is_some()),
("write_authority", request.write_authority.is_some()),
] {
if present {
return Err(DriverError::Rejected(format!(
"fleet `{fleet}` is an exact fleet: task option `{field}` is not allowed. Every \
member's provider, model, reasoning, and permission ceiling are fixed by the \
saved Fleet — switch Fleets or edit the Fleet, do not override a member per \
task."
)));
}
}
let binding = operation
.bind_member(request.profile.as_deref(), request.role.as_deref(), session)
.map_err(|err| DriverError::Rejected(format!("fleet `{fleet}`: {err}")))?;
request.profile = Some(binding.member_id.clone());
request.role = Some(binding.member_role.clone());
request.subagent_type = None;
request.allowed_tools = binding.authority.allowed_tools.clone();
request.disallowed_tools = binding.authority.disallowed_tools.clone();
request.write_authority = Some(binding.authority.write_authority.to_string());
request.max_depth = Some(
request
.max_depth
.map_or(binding.authority.max_depth, |asked| {
asked.min(binding.authority.max_depth)
}),
);
validate_exact_write_scope(&fleet, &binding, request)?;
Ok(binding)
}
fn displayed_resolved_role(
fleet_receipt: Option<&codewhale_workflow::FleetTaskReceipt>,
metadata_role: Option<&str>,
request_role: Option<&str>,
) -> Option<String> {
fleet_receipt
.map(|receipt| receipt.member_role.clone())
.or_else(|| metadata_role.map(str::to_string))
.or_else(|| request_role.map(str::to_string))
}
fn orphaned_fleet_receipt_line(
receipt: &codewhale_workflow::FleetTaskReceipt,
error: &str,
) -> String {
format!(
"fleet route {} spawn_failed=true reason={}",
receipt.line(),
error.replace('\n', " ")
)
}
fn validate_exact_write_scope(
fleet: &str,
binding: &crate::fleet::exact::ExactMemberBinding,
request: &TaskRequest,
) -> Result<(), DriverError> {
let declares_scope = !request.write_roots.is_empty()
|| !request.exact_files.is_empty()
|| !request.coordination_contracts.is_empty();
if binding.authority.write_authority == "read_only" {
if declares_scope {
return Err(DriverError::Rejected(format!(
"fleet `{fleet}`: member `{}` is read-only under the clamped ceiling, so this \
task may not declare write_roots, exact_files, or coordination_contracts.",
binding.member_id
)));
}
return Ok(());
}
if !declares_scope {
return Err(DriverError::Rejected(format!(
"fleet `{fleet}`: member `{}` is write-capable, so this task must declare \
write_roots, exact_files, or coordination_contracts before it can start. An \
unbounded write claim is refused at the spawn boundary, and this task would spend a \
reasoning-router call on its way to that refusal.",
binding.member_id
)));
}
Ok(())
}
async fn route_admitted_exact_task(
operation: &crate::fleet::exact::ExactFleetWorkflow,
binding: &crate::fleet::exact::ExactMemberBinding,
request: &mut TaskRequest,
) -> Result<codewhale_workflow::FleetTaskReceipt, DriverError> {
let fleet = operation.snapshot().fleet().qualified();
let launch = operation
.route_admitted_task(binding, &request.description)
.await
.map_err(|err| DriverError::Rejected(format!("fleet `{fleet}`: {err}")))?;
request.thinking = Some(launch.thinking.clone());
apply_launch_authority(&fleet, &launch, request)?;
Ok(launch.receipt)
}
fn apply_launch_authority(
fleet: &str,
launch: &crate::fleet::exact::ExactMemberLaunch,
request: &mut TaskRequest,
) -> Result<(), DriverError> {
let authority = &launch.authority;
if request.profile.as_deref() != Some(launch.member_id.as_str())
|| request.role.as_deref() != Some(launch.member_role.as_str())
{
return Err(DriverError::Rejected(format!(
"fleet `{fleet}`: task identity drifted between admission and launch (request \
profile={:?} role={:?}, launch member `{}` role `{}`); the launch is refused rather \
than run under an envelope resolved for a different member.",
request.profile, request.role, launch.member_id, launch.member_role,
)));
}
request.allowed_tools = authority.allowed_tools.clone();
request.disallowed_tools = authority.disallowed_tools.clone();
request.write_authority = Some(authority.write_authority.to_string());
request.subagent_type = None;
request.max_depth = Some(
request
.max_depth
.map_or(authority.max_depth, |asked| asked.min(authority.max_depth)),
);
let expected = authority.fingerprint();
if launch.receipt.authority_fingerprint.as_deref() != Some(expected.as_str()) {
return Err(DriverError::Rejected(format!(
"fleet `{fleet}`: member `{}` produced a receipt whose authority fingerprint does not \
match the envelope being installed (receipt={:?} envelope={expected}). Failing \
closed.",
launch.member_id, launch.receipt.authority_fingerprint,
)));
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn run_workflow_vm(
run_id: String,
source: String,
spec: Option<WorkflowSpec>,
args: Value,
driver: Arc<SubAgentWorkflowDriver>,
state: Arc<WorkflowWorkspaceState>,
context: ToolContext,
vm_cancel: WorkflowRunCancel,
) {
let result = WorkflowVm::new()
.run_script_with_cancel(&source, args, driver.clone(), vm_cancel)
.await;
let mut status = WorkflowRunStatus::Completed;
let mut output = None;
let mut error = None;
match result {
Ok(value) => {
if let Some(gate_error) = driver.terminal_gate_failure() {
status = WorkflowRunStatus::Failed;
error = Some(gate_error);
} else {
output = Some(value);
}
}
Err(err) => {
status = WorkflowRunStatus::Failed;
error = Some(err.to_string());
}
}
let snapshot = {
let mut runs_guard = match state.runs.lock() {
Ok(guard) => guard,
Err(_) => {
state.mark_owner_missing(&run_id);
return;
}
};
let Some(record) = runs_guard.get_mut(&run_id) else {
state.mark_owner_missing(&run_id);
return;
};
if record.status != WorkflowRunStatus::Cancelled {
if status == WorkflowRunStatus::Completed {
let task_records = driver.task_records_snapshot();
let failed_children = task_records
.iter()
.filter(|task| task.status == IrWorkflowRunStatus::Failed)
.count();
let rejected = record.dispatch_failures.len();
if record.child_ids.is_empty() && rejected > 0 {
status = WorkflowRunStatus::Failed;
error = Some(format!(
"no child agents ran: all {rejected} task dispatch(es) were rejected; first: {}",
record.dispatch_failures[0].message
));
} else if failed_children > 0 || rejected > 0 {
status = WorkflowRunStatus::Degraded;
let mut parts = Vec::new();
if failed_children > 0 {
parts.push(format!(
"{failed_children} of {} task(s) failed",
task_records.len()
));
}
if rejected > 0 {
parts.push(format!("{rejected} dispatch(es) were rejected"));
}
error = Some(format!(
"completed with dropped slots: {}; the recorded result may be partial",
parts.join(" and ")
));
}
}
record.status = status;
record.result = output;
record.error = error.clone();
record.execution = spec.as_ref().map(|spec| {
execution_from_declarative_spec(spec, driver.task_records_snapshot(), status)
});
record.completed_at_ms = Some(now_ms());
}
record.clone()
};
let verify_on_complete = state
.runs
.lock()
.ok()
.and_then(|guard| guard.get(&run_id).map(|record| record.verify_on_complete))
.unwrap_or(false);
if status == WorkflowRunStatus::Completed && verify_on_complete {
match run_workflow_completion_gates(&context).await {
Ok(verification) => {
if let Ok(mut runs_guard) = state.runs.lock()
&& let Some(record) = runs_guard.get_mut(&run_id)
{
record.verification = Some(verification);
}
}
Err(err) => {
if let Ok(mut runs_guard) = state.runs.lock()
&& let Some(record) = runs_guard.get_mut(&run_id)
{
record.status = WorkflowRunStatus::Failed;
record.error = Some(format!("verification gates failed: {err}"));
}
}
}
}
let final_budget = driver.current_budget_snapshot();
let run_usage = run_usage_totals(&driver.task_records_snapshot());
let snapshot = state
.runs
.lock()
.ok()
.and_then(|mut guard| {
let record = guard.get_mut(&run_id)?;
if record.status != WorkflowRunStatus::Cancelled {
record.lifecycle_seq = record.lifecycle_seq.saturating_add(1);
if run_usage.is_some() {
record.usage = run_usage.clone();
}
let budget_event = WorkflowUiEvent::new(budget_event_kind(final_budget));
let completed = WorkflowUiEvent::new(WorkflowUiEventKind::RunCompleted {
status: record.status,
error: record.error.clone(),
usage: run_usage.clone(),
});
record.push_event(budget_event.clone());
record.push_event(completed.clone());
driver.emit_ui_event(&budget_event);
driver.emit_ui_event(&completed);
}
Some(record.clone())
})
.unwrap_or(snapshot);
if state.try_record_snapshot(&snapshot).is_ok() {
state.reconcile_snapshot(&snapshot);
} else {
state.mark_owner_missing(&run_id);
}
write_run_report_artifact(&context.workspace, &snapshot);
if let Ok(mut controllers_guard) = state.controllers.lock() {
controllers_guard.remove(&run_id);
}
}
fn write_run_report_artifact(workspace: &Path, record: &WorkflowRunRecord) {
if !matches!(
record.status,
WorkflowRunStatus::Completed
| WorkflowRunStatus::Degraded
| WorkflowRunStatus::Failed
| WorkflowRunStatus::Cancelled
) {
return;
}
let safe_id: String = record
.run_id
.chars()
.filter(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_'))
.collect();
if safe_id.is_empty() {
return;
}
let dir = workspace.join(".codewhale").join("reports");
if let Err(err) = std::fs::create_dir_all(&dir) {
crate::logging::warn(format!(
"workflow report dir {} not created: {err}",
dir.display()
));
return;
}
let path = dir.join(format!("{safe_id}.md"));
if let Err(err) = std::fs::write(&path, render_run_report(record)) {
crate::logging::warn(format!(
"workflow report {} not written: {err}",
path.display()
));
}
}
fn render_run_report(record: &WorkflowRunRecord) -> String {
let mut out = String::new();
out.push_str(&format!("# Workflow run {}\n\n", record.run_id));
out.push_str(&format!("- status: {:?}\n", record.status));
if let Some(goal) = record.workflow_goal.as_deref() {
out.push_str(&format!("- goal: {goal}\n"));
}
if let Some(source) = record.source_path.as_deref() {
out.push_str(&format!("- source: {}\n", source.display()));
}
out.push_str(&format!("- started_at_ms: {}\n", record.started_at_ms));
if let Some(completed) = record.completed_at_ms {
out.push_str(&format!("- completed_at_ms: {completed}\n"));
}
if let Some(budget) = record.token_budget {
out.push_str(&format!("- token_budget: {budget}\n"));
}
out.push_str(&format!("- child_agents: {}\n", record.child_ids.len()));
if let Some(error) = record.error.as_deref() {
out.push_str(&format!("- error: {error}\n"));
}
if !record.dispatch_failures.is_empty() {
out.push_str(&format!(
"\n## Dispatch failures ({})\n\n",
record.dispatch_failures.len()
));
for failure in &record.dispatch_failures {
let slot = failure
.label
.as_deref()
.or(failure.phase.as_deref())
.unwrap_or("task");
out.push_str(&format!("- {slot}: {}\n", failure.message));
}
}
if !record.gate_status.is_empty() {
out.push_str("\n## Gates\n\n");
for line in &record.gate_status {
out.push_str(&format!("- {line:?}\n"));
}
}
if !record.progress.is_empty() {
out.push_str("\n## Progress\n\n");
for line in &record.progress {
out.push_str(&format!("- {line}\n"));
}
}
if !record.schema_errors.is_empty() {
out.push_str(&format!(
"\n## Schema errors ({})\n\n",
record.schema_errors.len()
));
}
if let Some(result) = record.result.as_ref() {
out.push_str("\n## Result\n\n```json\n");
out.push_str(&serde_json::to_string_pretty(result).unwrap_or_else(|_| result.to_string()));
out.push_str("\n```\n");
}
if let Some(verification) = record.verification.as_ref() {
out.push_str("\n## Verification\n\n```json\n");
out.push_str(
&serde_json::to_string_pretty(verification)
.unwrap_or_else(|_| verification.to_string()),
);
out.push_str("\n```\n");
}
out
}
fn session_workflow_config_store()
-> &'static Mutex<HashMap<PathBuf, codewhale_config::WorkflowConfigToml>> {
static STORE: OnceLock<Mutex<HashMap<PathBuf, codewhale_config::WorkflowConfigToml>>> =
OnceLock::new();
STORE.get_or_init(|| Mutex::new(HashMap::new()))
}
fn workflow_session_key(workspace: &Path) -> PathBuf {
workspace
.canonicalize()
.unwrap_or_else(|_| workspace.to_path_buf())
}
pub(crate) fn set_session_workflow_config(
workspace: &Path,
config: codewhale_config::WorkflowConfigToml,
) {
session_workflow_config_store()
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.insert(workflow_session_key(workspace), config);
}
pub(crate) fn session_workflow_config(
workspace: &Path,
) -> Option<codewhale_config::WorkflowConfigToml> {
session_workflow_config_store()
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.get(&workflow_session_key(workspace))
.cloned()
}
fn workflow_config_for(runtime: &SubAgentRuntime) -> codewhale_config::WorkflowConfigToml {
session_workflow_config(&runtime.context.workspace).unwrap_or_else(|| {
runtime
.api_config
.as_deref()
.map(crate::config::Config::workflow_config)
.unwrap_or_default()
})
}
fn workflow_result_for(
run_id: &str,
state: Arc<WorkflowWorkspaceState>,
) -> Result<ToolResult, ToolError> {
let record = {
let runs_guard = lock_mutex(&state.runs)?;
runs_guard.get(run_id).cloned().ok_or_else(|| {
ToolError::invalid_input(format!("Unknown workflow run_id '{run_id}'"))
})?
};
let journal_path = state.journal_path().to_path_buf();
let (payload, bounds) = bounded_run_record_value(&record, &journal_path);
let mut result =
ToolResult::json(&payload).map_err(|err| ToolError::execution_failed(err.to_string()))?;
let summary = record.summary();
result.metadata = Some(json!({
"run_id": summary.run_id,
"status": summary.status,
"terminal": summary.status != WorkflowRunStatus::Running,
"child_count": summary.child_count,
"schema_error_count": summary.schema_error_count,
"dispatch_failure_count": summary.dispatch_failure_count,
"event_count": summary.event_count,
"events_returned": bounds.events_returned,
"events_omitted": bounds.events_omitted,
"dispatch_failures_returned": bounds.dispatch_failures_returned,
"dispatch_failures_omitted": bounds.dispatch_failures_omitted,
"events_dropped": summary.events_dropped,
"last_event_type": summary.last_event_type,
"leaf_count": summary.leaf_count,
"branch_count": summary.branch_count,
"control_count": summary.control_count,
"execution_status": summary.execution_status,
"gate_count": summary.gate_count,
"blocked_gate_count": summary.blocked_gate_count,
"gate_status": summary.gate_status,
"truncated": bounds.truncated(),
"payload_budget_chars": WORKFLOW_RESULT_MAX_CHARS,
"journal_path": journal_path.display().to_string(),
"plan_approval": record.plan_approval,
}));
Ok(result)
}
#[derive(Debug, Default)]
struct RunPayloadBounds {
events_returned: usize,
events_omitted: usize,
progress_omitted: usize,
dispatch_failures_returned: usize,
dispatch_failures_omitted: usize,
dispatch_failure_fields_truncated: usize,
result_truncated: bool,
leaf_outputs_truncated: usize,
}
impl RunPayloadBounds {
fn truncated(&self) -> bool {
self.events_omitted > 0
|| self.progress_omitted > 0
|| self.dispatch_failures_omitted > 0
|| self.dispatch_failure_fields_truncated > 0
|| self.result_truncated
|| self.leaf_outputs_truncated > 0
}
}
fn bounded_run_record_value(
record: &WorkflowRunRecord,
journal_path: &Path,
) -> (Value, RunPayloadBounds) {
let mut bounds = RunPayloadBounds::default();
let journal = journal_path.display().to_string();
let mut value = serde_json::to_value(record).unwrap_or_else(|_| json!({}));
let Some(obj) = value.as_object_mut() else {
return (value, bounds);
};
if let Some(events) = obj.get_mut("events").and_then(Value::as_array_mut) {
if events.len() > WORKFLOW_RESULT_EVENTS_TAIL {
let omitted = events.len() - WORKFLOW_RESULT_EVENTS_TAIL;
events.drain(..omitted);
bounds.events_omitted = omitted;
}
bounds.events_returned = events.len();
}
if bounds.events_omitted > 0 {
obj.insert(
"events_note".to_string(),
json!(format!(
"showing the newest {} of {} events; full stream: {journal}",
bounds.events_returned,
bounds.events_returned + bounds.events_omitted,
)),
);
}
if let Some(progress) = obj.get_mut("progress").and_then(Value::as_array_mut)
&& progress.len() > WORKFLOW_RESULT_PROGRESS_TAIL
{
let omitted = progress.len() - WORKFLOW_RESULT_PROGRESS_TAIL;
progress.drain(..omitted);
bounds.progress_omitted = omitted;
obj.insert(
"progress_note".to_string(),
json!(format!(
"showing the newest {WORKFLOW_RESULT_PROGRESS_TAIL} of {} progress lines; full log: {journal}",
omitted + WORKFLOW_RESULT_PROGRESS_TAIL,
)),
);
}
if let Some(failures) = obj
.get_mut("dispatch_failures")
.and_then(Value::as_array_mut)
{
if failures.len() > WORKFLOW_RESULT_DISPATCH_FAILURES_TAIL {
let omitted = failures.len() - WORKFLOW_RESULT_DISPATCH_FAILURES_TAIL;
failures.drain(..omitted);
bounds.dispatch_failures_omitted = omitted;
}
bounds.dispatch_failures_returned = failures.len();
for failure in failures {
let Some(fields) = failure.as_object_mut() else {
continue;
};
for key in ["label", "phase", "message"] {
let Some(slot) = fields.get_mut(key) else {
continue;
};
let Some(raw) = slot.as_str() else {
continue;
};
if raw.chars().count() > WORKFLOW_RESULT_DISPATCH_FAILURE_FIELD_MAX_CHARS {
*slot = Value::String(truncate_chars(
raw,
WORKFLOW_RESULT_DISPATCH_FAILURE_FIELD_MAX_CHARS,
));
bounds.dispatch_failure_fields_truncated += 1;
}
}
}
}
if bounds.dispatch_failures_omitted > 0 || bounds.dispatch_failure_fields_truncated > 0 {
obj.insert(
"dispatch_failures_note".to_string(),
json!(format!(
"showing {} of {} dispatch failures with bounded fields; full record: {journal}",
bounds.dispatch_failures_returned,
bounds.dispatch_failures_returned + bounds.dispatch_failures_omitted,
)),
);
}
for key in ["result", "verification"] {
let Some(raw) = obj.get(key).filter(|value| !value.is_null()) else {
continue;
};
let text = raw.to_string();
if text.chars().count() > WORKFLOW_RESULT_VALUE_MAX_CHARS {
obj.insert(
key.to_string(),
json!({
"truncated": true,
"omitted_chars": text.chars().count() - WORKFLOW_RESULT_VALUE_MAX_CHARS,
"preview": truncate_chars(&text, WORKFLOW_RESULT_VALUE_MAX_CHARS),
"full_detail": journal,
}),
);
bounds.result_truncated = true;
}
}
if let Some(leaves) = obj
.get_mut("execution")
.and_then(|execution| execution.get_mut("leaf_results"))
.and_then(Value::as_array_mut)
{
for leaf in leaves {
let too_long = leaf
.get("output")
.and_then(Value::as_str)
.is_some_and(|output| {
output.chars().count() > WORKFLOW_RESULT_LEAF_OUTPUT_MAX_CHARS
});
if too_long
&& let Some(slot) = leaf.get_mut("output")
&& let Some(output) = slot.as_str()
{
let clipped = format!(
"{} [leaf output truncated to {WORKFLOW_RESULT_LEAF_OUTPUT_MAX_CHARS} chars; full text: {journal}]",
truncate_chars(output, WORKFLOW_RESULT_LEAF_OUTPUT_MAX_CHARS),
);
*slot = Value::String(clipped);
bounds.leaf_outputs_truncated += 1;
}
}
}
(value, bounds)
}
fn truncate_chars(text: &str, max_chars: usize) -> String {
if let Some((idx, _)) = text.char_indices().nth(max_chars) {
if max_chars < 3 {
return text[..idx].to_string();
}
let truncate_at = text
.char_indices()
.nth(max_chars - 3)
.map(|(idx, _)| idx)
.unwrap_or(0);
format!("{}...", &text[..truncate_at])
} else {
text.to_string()
}
}
#[derive(Debug)]
struct WorkflowSource {
source: String,
path: Option<PathBuf>,
spec: Option<WorkflowSpec>,
}
fn workflow_source(input: &Value, context: &ToolContext) -> Result<WorkflowSource, ToolError> {
let script = match optional_str(input, "script")? {
Some(script) => Some(script),
None => optional_str(input, "source")?,
}
.map(str::to_string);
let source_path = match optional_str(input, "source_path")? {
Some(path) => Some(path),
None => optional_str(input, "path")?,
};
let plan = input.get("plan").filter(|value| !value.is_null());
let provided = [
script.as_ref().is_some_and(|s| !s.trim().is_empty()),
source_path.is_some(),
plan.is_some(),
]
.into_iter()
.filter(|present| *present)
.count();
if provided > 1 {
return Err(ToolError::invalid_input(
"Use exactly one of script, source_path, or plan",
));
}
match (script, source_path, plan) {
(Some(source), None, None) if !source.trim().is_empty() => {
workflow_source_from_raw(source, None)
}
(None, Some(path), None) => read_workflow_source_path(path, context),
(None, None, Some(plan_value)) => workflow_source_from_plan(plan_value),
_ => Err(ToolError::missing_field("script")),
}
}
fn workflow_source_from_plan(plan_value: &Value) -> Result<WorkflowSource, ToolError> {
let spec = structured_plan_to_workflow_spec(plan_value)?;
let lowered = lower_declarative_workflow_to_imperative_js(&spec)?;
Ok(WorkflowSource {
source: lowered,
path: None,
spec: Some(spec),
})
}
#[derive(Debug, Deserialize)]
struct StructuredWorkflowPlan {
goal: String,
#[serde(default)]
risk: Option<String>,
#[serde(default)]
max_children: Option<usize>,
#[serde(default)]
token_budget: Option<u64>,
#[serde(default)]
phases: Vec<StructuredPlanPhase>,
#[serde(default)]
children: Vec<StructuredPlanChild>,
#[serde(default)]
nodes: Option<Value>,
#[serde(default)]
gates: Vec<GateSpec>,
}
#[derive(Debug, Deserialize)]
struct StructuredPlanPhase {
#[serde(default)]
id: Option<String>,
#[serde(default)]
title: Option<String>,
#[serde(default)]
parallel: Option<bool>,
#[serde(default)]
children: Vec<StructuredPlanChild>,
}
#[derive(Debug, Deserialize)]
struct StructuredPlanChild {
#[serde(default)]
id: Option<String>,
#[serde(default)]
label: Option<String>,
#[serde(alias = "description")]
prompt: String,
#[serde(default, alias = "type", alias = "agent_type")]
agent_type: Option<String>,
#[serde(default)]
role: Option<String>,
#[serde(default)]
profile: Option<String>,
#[serde(default)]
mode: Option<String>,
#[serde(default)]
file_scope: Vec<String>,
}
fn structured_plan_to_workflow_spec(plan_value: &Value) -> Result<WorkflowSpec, ToolError> {
if !plan_value.is_object() {
return Err(ToolError::invalid_input(
"Workflow plan must be a JSON object with goal and phases/children (or nodes)",
));
}
let plan: StructuredWorkflowPlan =
serde_json::from_value(plan_value.clone()).map_err(|err| {
ToolError::invalid_input(format!("Invalid structured Workflow plan: {err}"))
})?;
let goal = plan.goal.trim();
if goal.is_empty() {
return Err(ToolError::invalid_input(
"Workflow plan goal must be a non-empty string",
));
}
if let Some(nodes) = plan.nodes.as_ref() {
if !nodes.is_array() {
return Err(ToolError::invalid_input(
"Workflow plan.nodes must be an array of workflow nodes",
));
}
let mut object = plan_value.clone();
if let Some(obj) = object.as_object_mut() {
obj.insert("goal".to_string(), Value::String(goal.to_string()));
if let Some(token_budget) = plan.token_budget {
let mut budget = obj.get("budget").cloned().unwrap_or_else(|| json!({}));
if let Some(budget_obj) = budget.as_object_mut() {
budget_obj.insert("max_tokens".to_string(), json!(token_budget));
}
obj.insert("budget".to_string(), budget);
}
}
let wrapped = format!("workflow({});", object);
return compile_javascript_workflow("<structured plan>", &wrapped).map_err(|err| {
ToolError::invalid_input(format!("Invalid structured Workflow plan nodes: {err}"))
});
}
let default_mode = plan_risk_to_mode(plan.risk.as_deref())?;
let mut nodes = Vec::new();
if !plan.phases.is_empty() {
for (phase_index, phase) in plan.phases.iter().enumerate() {
let phase_id = phase
.id
.as_deref()
.or(phase.title.as_deref())
.map(str::trim)
.filter(|id| !id.is_empty())
.map(str::to_string)
.unwrap_or_else(|| format!("phase-{}", phase_index + 1));
let children = plan_children_to_leaves(
&phase.children,
default_mode,
plan.token_budget,
&phase_id,
)?;
if children.is_empty() {
return Err(ToolError::invalid_input(format!(
"Workflow plan phase '{phase_id}' must declare at least one child"
)));
}
let parallel = phase.parallel.unwrap_or(children.len() > 1);
if parallel && children.len() > 1 {
nodes.push(WorkflowNode::BranchSet(BranchSpec {
id: phase_id,
description: phase.title.clone(),
parallel: true,
budget: BudgetSpec {
max_tokens: plan.token_budget,
..BudgetSpec::default()
},
permissions: Default::default(),
model_policy: Default::default(),
children: children.into_iter().map(WorkflowNode::Leaf).collect(),
}));
} else if children.len() == 1 {
nodes.push(WorkflowNode::Leaf(
children.into_iter().next().expect("one child"),
));
} else {
nodes.push(WorkflowNode::Sequence(SequenceSpec {
id: phase_id,
children: children.into_iter().map(WorkflowNode::Leaf).collect(),
}));
}
}
} else if !plan.children.is_empty() {
let children =
plan_children_to_leaves(&plan.children, default_mode, plan.token_budget, "plan")?;
if children.len() == 1 {
nodes.push(WorkflowNode::Leaf(
children.into_iter().next().expect("one child"),
));
} else {
nodes.push(WorkflowNode::BranchSet(BranchSpec {
id: "plan".to_string(),
description: Some(goal.to_string()),
parallel: true,
budget: BudgetSpec {
max_tokens: plan.token_budget,
..BudgetSpec::default()
},
permissions: Default::default(),
model_policy: Default::default(),
children: children.into_iter().map(WorkflowNode::Leaf).collect(),
}));
}
} else {
return Err(ToolError::invalid_input(
"Workflow plan must include phases, children, or nodes",
));
}
let mut total_children = 0usize;
count_plan_leaves(&nodes, &mut total_children);
if let Some(max_children) = plan.max_children
&& total_children > max_children
{
return Err(ToolError::invalid_input(format!(
"Workflow plan declares {total_children} children which exceeds max_children={max_children}"
)));
}
Ok(WorkflowSpec {
id: None,
goal: goal.to_string(),
description: plan.risk.clone(),
budget: BudgetSpec {
max_tokens: plan.token_budget,
..BudgetSpec::default()
},
permissions: Default::default(),
model_policy: Default::default(),
promotion_policy: Default::default(),
gates: plan.gates,
nodes,
})
}
fn plan_risk_to_mode(risk: Option<&str>) -> Result<TaskMode, ToolError> {
match risk.map(str::trim).filter(|s| !s.is_empty()) {
None | Some("read_only") | Some("readonly") | Some("low") | Some("safe") => {
Ok(TaskMode::ReadOnly)
}
Some("writes") | Some("write") | Some("read_write") | Some("readwrite")
| Some("medium") => Ok(TaskMode::ReadWrite),
Some("elevated") | Some("high") | Some("shell") | Some("network") => {
Ok(TaskMode::ReadWrite)
}
Some(other) => Err(ToolError::invalid_input(format!(
"Invalid plan risk '{other}'. Use read_only, writes, or elevated."
))),
}
}
fn plan_children_to_leaves(
children: &[StructuredPlanChild],
default_mode: TaskMode,
token_budget: Option<u64>,
phase_id: &str,
) -> Result<Vec<LeafSpec>, ToolError> {
if children.is_empty() {
return Ok(Vec::new());
}
let mut leaves = Vec::with_capacity(children.len());
for (index, child) in children.iter().enumerate() {
let prompt = child.prompt.trim();
if prompt.is_empty() {
return Err(ToolError::invalid_input(format!(
"Workflow plan child {} in phase '{phase_id}' must have a non-empty prompt",
index + 1
)));
}
let id = child
.id
.as_deref()
.or(child.label.as_deref())
.map(str::trim)
.filter(|id| !id.is_empty())
.map(str::to_string)
.unwrap_or_else(|| format!("{phase_id}-child-{}", index + 1));
let agent_type = parse_plan_agent_type(child.agent_type.as_deref())?;
let mode = match child.mode.as_deref().map(str::trim) {
None | Some("") => default_mode,
Some("read_only") | Some("readonly") => TaskMode::ReadOnly,
Some("read_write") | Some("readwrite") | Some("writes") | Some("write") => {
TaskMode::ReadWrite
}
Some(other) => {
return Err(ToolError::invalid_input(format!(
"Invalid plan child mode '{other}' on '{id}'. Use read_only or read_write."
)));
}
};
let role = child
.role
.as_deref()
.map(str::trim)
.filter(|r| !r.is_empty())
.map(|r| r.to_ascii_lowercase());
let profile = child
.profile
.as_deref()
.map(str::trim)
.filter(|p| !p.is_empty())
.map(|p| p.to_ascii_lowercase());
leaves.push(LeafSpec {
id,
prompt: prompt.to_string(),
agent_type,
role,
profile,
mode,
isolation: Default::default(),
file_scope: child.file_scope.clone(),
depends_on_results: Vec::new(),
budget: BudgetSpec {
max_tokens: token_budget,
..BudgetSpec::default()
},
permissions: Default::default(),
model_policy: Default::default(),
});
}
Ok(leaves)
}
fn parse_plan_agent_type(raw: Option<&str>) -> Result<AgentType, ToolError> {
let Some(kind) = raw.map(str::trim).filter(|s| !s.is_empty()) else {
return Ok(AgentType::General);
};
match kind.to_ascii_lowercase().as_str() {
"general" | "worker" | "delegate" => Ok(AgentType::General),
"explore" | "explorer" | "scout" => Ok(AgentType::Explore),
"plan" | "planner" | "awaiter" => Ok(AgentType::Plan),
"review" | "reviewer" | "consultant" | "oracle" | "advisor" => Ok(AgentType::Review),
"implementer" | "implement" | "builder" => Ok(AgentType::Implementer),
"verifier" | "verify" => Ok(AgentType::Verifier),
"custom" => Err(ToolError::invalid_input(
"Invalid sub-agent type 'custom' for a Workflow plan child: custom requires an \
explicit allowed_tools list, which plan children cannot declare. Use role/profile \
or another type.",
)),
_ => Err(ToolError::invalid_input(format!(
"Invalid sub-agent type '{kind}'. Use: worker, scout, planner, reviewer, builder, \
verifier (legacy aliases remain accepted: general, explore/explorer, plan/awaiter, \
review, implementer, consultant/oracle/advisor)."
))),
}
}
fn count_plan_leaves(nodes: &[WorkflowNode], total: &mut usize) {
for node in nodes {
match node {
WorkflowNode::Leaf(_) => *total += 1,
WorkflowNode::BranchSet(spec) => count_plan_leaves(&spec.children, total),
WorkflowNode::Sequence(spec) => count_plan_leaves(&spec.children, total),
WorkflowNode::Reduce(_)
| WorkflowNode::TeacherReview(_)
| WorkflowNode::LoopUntil(_)
| WorkflowNode::Cond(_)
| WorkflowNode::Expand(_) => {}
}
}
}
fn read_workflow_source_path(
path: &str,
context: &ToolContext,
) -> Result<WorkflowSource, ToolError> {
let raw = Path::new(path);
let joined = if raw.is_absolute() {
raw.to_path_buf()
} else {
context.workspace.join(raw)
};
let canonical = joined.canonicalize().map_err(|err| {
ToolError::invalid_input(format!(
"Failed to resolve workflow source_path '{path}': {err}"
))
})?;
if !context.trust_mode {
let workspace = context
.workspace
.canonicalize()
.unwrap_or_else(|_| context.workspace.clone());
let home_store = crate::config::effective_home_dir()
.map(|home| home.join(".codewhale").join("workflows"))
.and_then(|dir| dir.canonicalize().ok());
let inside_home_store = home_store
.as_deref()
.is_some_and(|dir| canonical.starts_with(dir));
if !canonical.starts_with(&workspace) && !inside_home_store {
return Err(ToolError::permission_denied(format!(
"workflow source_path must stay inside the workspace or ~/.codewhale/workflows: {}",
canonical.display()
)));
}
}
let source = std::fs::read_to_string(&canonical).map_err(|err| {
ToolError::execution_failed(format!(
"Failed to read workflow source_path '{}': {err}",
canonical.display()
))
})?;
workflow_source_from_raw(source, Some(canonical))
}
fn workflow_source_from_raw(
source: String,
path: Option<PathBuf>,
) -> Result<WorkflowSource, ToolError> {
let adapted = adapt_workflow_source(&source, path.as_deref())?;
Ok(WorkflowSource {
source: adapted.source,
path,
spec: adapted.spec,
})
}
struct AdaptedWorkflowSource {
source: String,
spec: Option<WorkflowSpec>,
}
fn adapt_workflow_source(
source: &str,
path: Option<&Path>,
) -> Result<AdaptedWorkflowSource, ToolError> {
if !looks_like_declarative_workflow(source) {
return Ok(AdaptedWorkflowSource {
source: source.to_string(),
spec: None,
});
}
let identifier = path
.map(|path| path.display().to_string())
.unwrap_or_else(|| "<inline workflow>".to_string());
let extension = path
.and_then(Path::extension)
.and_then(|extension| extension.to_str())
.unwrap_or_default();
let spec = if extension.eq_ignore_ascii_case("ts") {
compile_typescript_workflow(&identifier, source)
} else {
compile_javascript_workflow(&identifier, source)
}
.map_err(|err| {
ToolError::invalid_input(format!(
"Failed to compile declarative Workflow source '{identifier}': {err}"
))
})?;
let lowered = lower_declarative_workflow_to_imperative_js(&spec)?;
Ok(AdaptedWorkflowSource {
source: lowered,
spec: Some(spec),
})
}
fn looks_like_declarative_workflow(source: &str) -> bool {
source.lines().any(|line| {
let trimmed = line.trim_start();
trimmed.starts_with("workflow(") || trimmed.starts_with("export default workflow(")
})
}
fn lower_declarative_workflow_to_imperative_js(spec: &WorkflowSpec) -> Result<String, ToolError> {
let mut lowerer = DeclarativeWorkflowLowerer::default();
lowerer.line("\"use strict\";");
lowerer.line("const __results = {};");
lowerer.line(format!(
"phase({});",
js_string(&format!("workflow: {}", spec.goal))
));
for node in &spec.nodes {
lowerer.lower_node(node, None)?;
}
lowerer.line("return __results;");
Ok(lowerer.finish())
}
#[derive(Default)]
struct DeclarativeWorkflowLowerer {
source: String,
next_var: usize,
}
impl DeclarativeWorkflowLowerer {
fn finish(self) -> String {
self.source
}
fn line(&mut self, line: impl AsRef<str>) {
self.source.push_str(line.as_ref());
self.source.push('\n');
}
fn next_temp(&mut self, prefix: &str) -> String {
let value = format!("__{prefix}_{}", self.next_var);
self.next_var += 1;
value
}
fn lower_node(&mut self, node: &WorkflowNode, phase: Option<&str>) -> Result<(), ToolError> {
match node {
WorkflowNode::Leaf(spec) => self.lower_leaf(spec, phase, false),
WorkflowNode::BranchSet(spec) => self.lower_branch(spec),
WorkflowNode::Sequence(spec) => self.lower_sequence(spec),
WorkflowNode::Reduce(spec) => self.lower_reduce(spec),
WorkflowNode::TeacherReview(_) => Err(unsupported_declarative_node("teacher_review")),
WorkflowNode::LoopUntil(_) => Err(unsupported_declarative_node("loop_until")),
WorkflowNode::Cond(_) => Err(unsupported_declarative_node("cond")),
WorkflowNode::Expand(_) => Err(unsupported_declarative_node("expand")),
}
}
fn lower_leaf(
&mut self,
spec: &LeafSpec,
phase: Option<&str>,
parallel: bool,
) -> Result<(), ToolError> {
self.line(format!(
"__results[{}] = await task({});",
js_string(&spec.id),
leaf_task_options_expression(spec, phase, parallel)?
));
Ok(())
}
fn lower_branch(&mut self, spec: &BranchSpec) -> Result<(), ToolError> {
self.line(format!("phase({});", js_string(&spec.id)));
if spec.parallel {
let mut leaves = Vec::new();
for child in &spec.children {
let WorkflowNode::Leaf(leaf) = child else {
return Err(ToolError::invalid_input(format!(
"Declarative Workflow adapter only supports leaf children inside parallel branch '{}'",
spec.id
)));
};
leaves.push(leaf);
}
let temp = self.next_temp("parallel");
self.line(format!("const {temp} = await parallel(["));
for leaf in &leaves {
self.line(format!(
" () => task({}),",
leaf_task_options_expression(leaf, Some(&spec.id), true)?
));
}
self.line("]);");
for (index, leaf) in leaves.iter().enumerate() {
self.line(format!(
"__results[{}] = {temp}[{index}];",
js_string(&leaf.id)
));
}
return Ok(());
}
for child in &spec.children {
self.lower_node(child, Some(&spec.id))?;
}
Ok(())
}
fn lower_sequence(&mut self, spec: &SequenceSpec) -> Result<(), ToolError> {
self.line(format!("phase({});", js_string(&spec.id)));
for child in &spec.children {
self.lower_node(child, Some(&spec.id))?;
}
Ok(())
}
fn lower_reduce(&mut self, spec: &ReduceSpec) -> Result<(), ToolError> {
let inputs = result_inputs_expression(&spec.inputs);
self.line(format!(
"__results[{}] = await task({});",
js_string(&spec.id),
task_options_expression(
format!(
"{} + \"\\n\\nInputs:\\n\" + {inputs}",
js_string(&spec.prompt)
),
Some("plan"),
None,
None,
false,
None,
None,
None,
Some("read_only"),
&[],
&spec.id,
Some("reduce"),
None,
)
));
Ok(())
}
}
fn unsupported_declarative_node(kind: &str) -> ToolError {
ToolError::invalid_input(format!(
"Declarative Workflow adapter does not yet support {kind} nodes"
))
}
fn leaf_description(spec: &LeafSpec) -> String {
let mut description = spec.prompt.trim().to_string();
let mut metadata = Vec::new();
metadata.push(format!("Workflow leaf id: {}", spec.id));
metadata.push(format!("Mode: {}", task_mode_name(spec.mode)));
if !spec.file_scope.is_empty() {
metadata.push(format!("File scope: {}", spec.file_scope.join(", ")));
}
if !spec.depends_on_results.is_empty() {
metadata.push(format!(
"Depends on results: {}",
spec.depends_on_results.join(", ")
));
}
if spec.budget != BudgetSpec::default() {
let mut budget = Vec::new();
if let Some(max_steps) = spec.budget.max_steps {
budget.push(format!("max_steps={max_steps}"));
}
if let Some(timeout_secs) = spec.budget.timeout_secs {
budget.push(format!("timeout_secs={timeout_secs}"));
}
if let Some(max_parallel) = spec.budget.max_parallel {
budget.push(format!("max_parallel={max_parallel}"));
}
if let Some(max_tokens) = spec.budget.max_tokens {
budget.push(format!("max_tokens={max_tokens}"));
}
if !budget.is_empty() {
metadata.push(format!("Budget: {}", budget.join(", ")));
}
}
if !metadata.is_empty() {
description.push_str("\n\nWorkflow metadata:\n");
for item in metadata {
description.push_str("- ");
description.push_str(&item);
description.push('\n');
}
}
description
}
fn leaf_task_options_expression(
spec: &LeafSpec,
phase: Option<&str>,
parallel: bool,
) -> Result<String, ToolError> {
validate_leaf_runtime_contract(spec)?;
let worktree = leaf_wants_worktree(spec, parallel);
let write_authority = match spec.mode {
TaskMode::ReadOnly => "read_only",
TaskMode::ReadWrite if worktree => "worktree_write",
TaskMode::ReadWrite => "workspace_write",
};
let write_roots = if spec.mode == TaskMode::ReadWrite {
spec.file_scope
.iter()
.map(|scope| codewhale_workflow::normalize_file_scope_root(scope))
.collect::<Vec<_>>()
} else {
Vec::new()
};
Ok(task_options_expression(
leaf_description_expression(spec),
leaf_subagent_type(spec),
spec.role.as_deref(),
spec.profile.as_deref(),
worktree,
spec.budget.max_tokens,
spec.budget.max_steps,
spec.budget.timeout_secs,
Some(write_authority),
&write_roots,
&spec.id,
phase,
leaf_allowed_tools(spec)?,
))
}
fn validate_leaf_runtime_contract(spec: &LeafSpec) -> Result<(), ToolError> {
if spec.mode == TaskMode::ReadOnly && spec.permissions.allow_write {
return Err(ToolError::invalid_input(format!(
"Workflow leaf '{}' is read_only but requests allow_write permissions",
spec.id
)));
}
if spec.mode == TaskMode::ReadWrite && spec.file_scope.is_empty() {
return Err(ToolError::invalid_input(format!(
"Workflow leaf '{}' is read_write but declares no file_scope for its bounded write claim",
spec.id
)));
}
for scope in &spec.file_scope {
let normalized = codewhale_workflow::normalize_file_scope_root(scope);
if normalized.is_empty() || normalized.contains('*') {
return Err(ToolError::invalid_input(format!(
"Workflow leaf '{}' has unsupported file_scope '{}'; use a concrete path or a trailing /* or /** directory scope",
spec.id, scope
)));
}
}
if spec.mode == TaskMode::ReadWrite
&& matches!(
spec.agent_type,
AgentType::Explore | AgentType::Plan | AgentType::Review | AgentType::Verifier
)
{
return Err(ToolError::invalid_input(format!(
"Workflow leaf '{}' is read_write but uses read-only agent_type {}",
spec.id,
agent_type_name(spec.agent_type)
)));
}
if spec.mode == TaskMode::ReadOnly
&& spec
.permissions
.allowed_tools
.iter()
.any(|tool| is_write_or_shell_tool(tool))
{
return Err(ToolError::invalid_input(format!(
"Workflow leaf '{}' is read_only but requests write/shell allowed_tools",
spec.id
)));
}
if spec.permissions.deny_all_tools && !spec.permissions.allowed_tools.is_empty() {
return Err(ToolError::invalid_input(format!(
"Workflow leaf '{}' cannot combine deny_all_tools with allowed_tools",
spec.id
)));
}
Ok(())
}
fn leaf_description_expression(spec: &LeafSpec) -> String {
let description = js_string(&leaf_description(spec));
if spec.depends_on_results.is_empty() {
return description;
}
let inputs = result_inputs_expression(&spec.depends_on_results);
format!("{description} + \"\\n\\nInputs:\\n\" + {inputs}")
}
fn result_inputs_expression(inputs: &[String]) -> String {
let entries = inputs
.iter()
.map(|input| format!("[{}, __results[{}]]", js_string(input), js_string(input)))
.collect::<Vec<_>>()
.join(", ");
format!(
"[{entries}].map(([id, value]) => \"--- \" + id + \" ---\\n\" + String(value ?? \"\")).join(\"\\n\\n\")"
)
}
fn leaf_subagent_type(spec: &LeafSpec) -> Option<&'static str> {
if (spec.role.is_some() || spec.profile.is_some()) && spec.agent_type == AgentType::General {
return None;
}
if spec.mode == TaskMode::ReadOnly && spec.agent_type == AgentType::General {
return Some("review");
}
if spec.mode == TaskMode::ReadOnly
&& spec.agent_type == AgentType::Implementer
&& (spec.role.is_some() || spec.profile.is_some())
{
return None;
}
Some(agent_type_name(spec.agent_type))
}
fn leaf_allowed_tools(spec: &LeafSpec) -> Result<Option<Vec<String>>, ToolError> {
if spec.permissions.deny_all_tools {
return Ok(Some(Vec::new()));
}
if !spec.permissions.allowed_tools.is_empty() {
return Ok(Some(spec.permissions.allowed_tools.clone()));
}
if spec.mode != TaskMode::ReadOnly {
return Ok(None);
}
Ok(Some(
read_only_allowed_tools(spec.agent_type)
.iter()
.map(|tool| (*tool).to_string())
.collect(),
))
}
fn read_only_allowed_tools(agent_type: AgentType) -> &'static [&'static str] {
match agent_type {
AgentType::Verifier => &["File"],
_ => &["File"],
}
}
fn is_write_or_shell_tool(tool: &str) -> bool {
codewhale_workflow::is_write_tool(tool) || codewhale_workflow::is_shell_tool(tool)
}
#[allow(clippy::too_many_arguments)]
fn task_options_expression(
description_expr: String,
subagent_type: Option<&str>,
role: Option<&str>,
profile: Option<&str>,
worktree: bool,
token_budget: Option<u64>,
max_steps: Option<u32>,
wall_time_secs: Option<u64>,
write_authority: Option<&str>,
write_roots: &[String],
label: &str,
phase: Option<&str>,
allowed_tools: Option<Vec<String>>,
) -> String {
let mut fields = vec![format!("description: {description_expr}")];
if let Some(subagent_type) = subagent_type {
fields.push(format!("type: {}", js_string(subagent_type)));
}
fields.push(format!("label: {}", js_string(label)));
if let Some(phase) = phase {
fields.push(format!("phase: {}", js_string(phase)));
}
if let Some(role) = role {
fields.push(format!("role: {}", js_string(role)));
}
if let Some(profile) = profile {
fields.push(format!("profile: {}", js_string(profile)));
}
if worktree {
fields.push("worktree: true".to_string());
}
if let Some(token_budget) = token_budget {
fields.push(format!("tokenBudget: {token_budget}"));
}
if let Some(max_steps) = max_steps {
fields.push(format!("maxSteps: {max_steps}"));
}
if let Some(wall_time_secs) = wall_time_secs {
fields.push(format!("wallTimeSecs: {wall_time_secs}"));
}
if let Some(write_authority) = write_authority {
fields.push(format!("writeAuthority: {}", js_string(write_authority)));
}
if !write_roots.is_empty() {
fields.push(format!(
"writeRoots: {}",
serde_json::to_string(write_roots).expect("serializing write roots cannot fail")
));
}
if let Some(allowed_tools) = allowed_tools {
fields.push(format!(
"allowedTools: {}",
serde_json::to_string(&allowed_tools).expect("serializing tool names cannot fail")
));
}
format!("{{ {} }}", fields.join(", "))
}
fn js_string(value: &str) -> String {
serde_json::to_string(value).expect("serializing JS string cannot fail")
}
fn agent_type_name(agent_type: AgentType) -> &'static str {
match agent_type {
AgentType::General => "general",
AgentType::Explore => "explore",
AgentType::Plan => "plan",
AgentType::Review => "review",
AgentType::Implementer => "implementer",
AgentType::Verifier => "verifier",
}
}
fn task_mode_name(mode: TaskMode) -> &'static str {
match mode {
TaskMode::ReadOnly => "read_only",
TaskMode::ReadWrite => "read_write",
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ExplicitGateVerdict {
Approve,
Reject,
}
fn explicit_gate_verdict(output: Option<&str>) -> Option<ExplicitGateVerdict> {
let first_meaningful = output?
.lines()
.map(str::trim)
.find(|line| !line.is_empty())?;
if first_meaningful.eq_ignore_ascii_case("APPROVE")
|| first_meaningful.eq_ignore_ascii_case("PASS")
{
Some(ExplicitGateVerdict::Approve)
} else if first_meaningful.eq_ignore_ascii_case("BLOCK")
|| first_meaningful.eq_ignore_ascii_case("FAIL")
{
Some(ExplicitGateVerdict::Reject)
} else {
None
}
}
fn has_gate_artifact_body(output: Option<&str>) -> bool {
let Some(output) = output else {
return false;
};
let mut meaningful_lines = output
.lines()
.map(str::trim)
.filter(|line| !line.is_empty());
meaningful_lines.next();
meaningful_lines.next().is_some() && meaningful_lines.next().is_some()
}
fn gate_outcome_for_completed_role(
record: &RuntimeTaskRecord,
require_explicit_verdict: bool,
artifact_kind: Option<&str>,
) -> GateOutcome {
match record.status {
IrWorkflowRunStatus::Succeeded => match explicit_gate_verdict(record.output.as_deref()) {
Some(ExplicitGateVerdict::Reject) => GateOutcome::Fail {
reason: record
.output
.clone()
.unwrap_or_else(|| "child returned an explicit rejection verdict".into()),
},
Some(ExplicitGateVerdict::Approve)
if require_explicit_verdict
&& artifact_kind.is_some()
&& !has_gate_artifact_body(record.output.as_deref()) =>
{
GateOutcome::Fail {
reason: format!(
"task {} approved without the required {} artifact body",
record.agent_id,
artifact_kind.unwrap_or("gate")
),
}
}
Some(ExplicitGateVerdict::Approve) => GateOutcome::Pass,
None if require_explicit_verdict => GateOutcome::Fail {
reason: format!(
"task {} completed without the required first-line gate verdict; expected exactly APPROVE, PASS, BLOCK, or FAIL",
record.agent_id
),
},
None => GateOutcome::Pass,
},
_ => GateOutcome::Fail {
reason: record.output.clone().unwrap_or_else(|| {
format!("task {} ended as {:?}", record.agent_id, record.status)
}),
},
}
}
#[derive(Debug, Clone)]
struct RuntimeTaskRecord {
agent_id: String,
label: Option<String>,
role: Option<String>,
status: IrWorkflowRunStatus,
output: Option<String>,
schema_error: Option<String>,
usage: Option<WorkflowTaskUsage>,
}
struct SubAgentWorkflowDriver {
run_id: String,
manager: SharedSubAgentManager,
runtime: SubAgentRuntime,
state: Arc<WorkflowWorkspaceState>,
completion_tx: mpsc::UnboundedSender<SubAgentCompletion>,
completion_state: Arc<Mutex<CompletionState>>,
child_ids: Arc<Mutex<Vec<String>>>,
child_counter: AtomicU32,
current_phase: Mutex<Option<String>>,
task_records: Arc<Mutex<HashMap<String, RuntimeTaskRecord>>>,
total_budget: Option<u64>,
last_budget_event: Arc<Mutex<Option<BudgetSnapshot>>>,
gate_specs: Arc<Vec<GateSpec>>,
gate_board: Arc<Mutex<LaneGateBoard>>,
concurrent_gate: Arc<Semaphore>,
spawn_permits: Mutex<HashMap<String, OwnedSemaphorePermit>>,
fleet_name: Option<String>,
fleet: WorkflowFleetBinding,
}
impl SubAgentWorkflowDriver {
#[allow(clippy::too_many_arguments)]
fn new(
run_id: String,
manager: SharedSubAgentManager,
runtime: SubAgentRuntime,
state: Arc<WorkflowWorkspaceState>,
total_budget: Option<u64>,
fleet: WorkflowFleetBinding,
gate_specs: Vec<GateSpec>,
) -> Arc<Self> {
let fleet_name = fleet.name();
let (completion_tx, completion_rx) = mpsc::unbounded_channel();
let mut gate_board = LaneGateBoard::new(run_id.clone());
gate_board.install_gates(&gate_specs);
let driver = Arc::new(Self {
run_id,
manager,
runtime,
state,
completion_tx,
completion_state: Arc::new(Mutex::new(CompletionState::default())),
child_ids: Arc::new(Mutex::new(Vec::new())),
child_counter: AtomicU32::new(0),
current_phase: Mutex::new(None),
task_records: Arc::new(Mutex::new(HashMap::new())),
total_budget,
last_budget_event: Arc::new(Mutex::new(None)),
gate_specs: Arc::new(gate_specs),
gate_board: Arc::new(Mutex::new(gate_board)),
concurrent_gate: Arc::new(Semaphore::new(WORKFLOW_MAX_CONCURRENT.max(1))),
spawn_permits: Mutex::new(HashMap::new()),
fleet_name,
fleet,
});
spawn_completion_pump(driver.clone(), completion_rx);
driver
}
fn force_cancel_all(&self) {
let ids = self
.child_ids
.lock()
.map(|ids| ids.clone())
.unwrap_or_default();
if let Ok(mut permits) = self.spawn_permits.lock() {
permits.clear();
}
cancel_child_agents(self.manager.clone(), ids);
if let Ok(mut state) = self.completion_state.lock() {
for (_, waiter) in state.waiters.drain() {
let _ = waiter.send(TaskCompletion::Cancelled);
}
}
}
fn finalize_running_tasks_cancelled(&self) {
let ids = self
.child_ids
.lock()
.map(|ids| ids.clone())
.unwrap_or_default();
for id in &ids {
self.record_task_completion(id, &TaskCompletion::Cancelled, None);
}
}
fn record_child(&self, agent_id: &str) {
if let Ok(mut ids) = self.child_ids.lock()
&& !ids.iter().any(|id| id == agent_id)
{
ids.push(agent_id.to_string());
}
if let Ok(mut runs) = self.state.runs.lock()
&& let Some(record) = runs.get_mut(&self.run_id)
&& !record.child_ids.iter().any(|id| id == agent_id)
{
record.child_ids.push(agent_id.to_string());
}
}
fn current_budget_snapshot(&self) -> BudgetSnapshot {
let spent = self
.manager
.try_read()
.ok()
.map(|manager| manager.budget_spent_for_scope(&self.run_id))
.unwrap_or(0);
BudgetSnapshot {
total: self.total_budget,
spent,
}
}
fn terminal_gate_failure(&self) -> Option<String> {
let board = match self.gate_board.lock() {
Ok(board) => board,
Err(_) => {
return Some(
"workflow gate board was unavailable during terminal finalization".to_string(),
);
}
};
self.gate_specs.iter().find_map(|spec| {
let state = board.gates.get(&spec.id)?;
state.is_blocking().then(|| {
format!(
"workflow gate `{}` ended {}: {}",
spec.id,
state.as_str(),
gate_state_reason(state)
)
})
})
}
fn record_run_event(&self, event: WorkflowUiEvent) {
let recorded = if let Ok(mut runs) = self.state.runs.lock()
&& let Some(record) = runs.get_mut(&self.run_id)
{
record.push_event(event.clone());
true
} else {
false
};
if recorded {
self.state.record_event(&self.run_id, &event);
}
self.emit_ui_event(&event);
}
fn emit_ui_event(&self, event: &WorkflowUiEvent) {
let Some(tx) = self.runtime.event_tx.as_ref() else {
return;
};
let Ok(mut value) = serde_json::to_value(event) else {
return;
};
if let Some(obj) = value.as_object_mut() {
obj.insert("run_id".to_string(), json!(self.run_id));
}
let _ = tx.try_send(Event::WorkflowUi {
run_id: self.run_id.clone(),
event: value,
});
}
fn record_budget_snapshot(&self, snapshot: BudgetSnapshot) {
let changed = if let Ok(mut last) = self.last_budget_event.lock() {
if last.as_ref() == Some(&snapshot) {
false
} else {
*last = Some(snapshot);
true
}
} else {
false
};
let event = WorkflowUiEvent::new(budget_event_kind(snapshot));
if changed {
self.record_run_event(event);
} else {
self.emit_ui_event(&event);
}
}
fn prepare_request_for_gates(
&self,
request: &mut TaskRequest,
) -> Result<Vec<HandoffArtifact>, DriverError> {
let Some(role) = request.role.as_deref().filter(|role| !role.is_empty()) else {
return Ok(Vec::new());
};
if self.gate_specs.is_empty() {
return Ok(Vec::new());
}
let (blocked, handoffs) = {
let mut board = self
.gate_board
.lock()
.map_err(|_| DriverError::Rejected("workflow gate board lock poisoned".into()))?;
let blocked = board.role_is_blocked(&self.gate_specs, role).cloned();
let handoffs = if blocked.is_none() {
board.consume_handoffs_for(role, 4)
} else {
Vec::new()
};
(blocked, handoffs)
};
if let Some(state) = blocked {
return Err(DriverError::Rejected(format!(
"workflow gate blocks role `{role}`: {}",
gate_state_reason(&state)
)));
}
if !handoffs.is_empty() {
append_handoff_context(request, &handoffs);
}
Ok(handoffs)
}
fn update_gate_status(&self, status: Vec<GateStatusLine>) {
let snapshot = if let Ok(mut runs) = self.state.runs.lock()
&& let Some(record) = runs.get_mut(&self.run_id)
{
record.gate_status = status;
Some(record.clone())
} else {
None
};
if let Some(record) = snapshot {
self.state.record_snapshot(&record);
}
}
fn evaluate_gates_for_completed_role(&self, record: &RuntimeTaskRecord) {
let Some(role) = record.role.as_deref().filter(|role| !role.is_empty()) else {
return;
};
if self.gate_specs.is_empty() {
return;
}
let specs = self
.gate_specs
.iter()
.filter(|spec| spec.on == GateOn::RoleComplete && spec.role.eq_ignore_ascii_case(role))
.cloned()
.collect::<Vec<_>>();
if specs.is_empty() {
return;
}
let mut events = Vec::new();
let mut next_status = Vec::new();
if let Ok(mut board) = self.gate_board.lock() {
for spec in specs {
let outcome = gate_outcome_for_completed_role(
record,
spec.require_explicit_verdict,
spec.artifact_kind.as_deref(),
);
let mut state = match board.evaluate(&spec, outcome.clone()) {
Ok(state) => state,
Err(err) => {
let state = GateState::Blocked {
reason: err.to_string(),
};
board.gates.insert(spec.id.clone(), state.clone());
state
}
};
let mut promotion = None;
if matches!(state, GateState::Passed)
&& let (Some(kind), Some(to_role)) =
(spec.artifact_kind.as_deref(), spec.blocks_role.as_deref())
{
let artifact = HandoffArtifact {
id: format!("handoff_{}", Uuid::new_v4()),
lane_id: self.run_id.clone(),
from_role: spec.role.clone(),
to_role: to_role.to_string(),
kind: kind.to_string(),
payload: record.output.clone().unwrap_or_default(),
created_at: now_ms().to_string(),
};
match board.record_handoff(artifact.clone()) {
Ok(()) => {
promotion =
Some(WorkflowUiEvent::new(WorkflowUiEventKind::HandoffPromoted {
artifact_id: artifact.id,
gate_id: spec.id.clone(),
kind: artifact.kind,
from_role: artifact.from_role,
to_role: artifact.to_role,
producer_task_id: record.agent_id.clone(),
}));
}
Err(err) => {
state = GateState::Blocked {
reason: format!(
"gate passed but its handoff could not be recorded: {err}"
),
};
board.gates.insert(spec.id.clone(), state.clone());
}
}
}
events.push(WorkflowUiEvent::new(WorkflowUiEventKind::GateUpdated {
gate_id: spec.id.clone(),
role: spec.role.clone(),
gate: gate_kind_label(spec.gate).to_string(),
state: state.as_str().to_string(),
blocked_role: spec.blocks_role.clone(),
blocked_reason: state.blocked_reason().map(str::to_string),
}));
if let Some(event) = promotion {
events.push(event);
}
}
next_status = board.status_summary();
}
if !events.is_empty() || !next_status.is_empty() {
self.update_gate_status(next_status);
}
for event in events {
self.record_run_event(event);
}
}
fn record_task_started(
&self,
agent_id: &str,
request: &TaskRequest,
metadata: &WorkflowTaskSpawnMetadata,
result: &crate::tools::subagent::SubAgentResult,
fleet_receipt: Option<codewhale_workflow::FleetTaskReceipt>,
) {
let label = metadata
.workflow_task_label
.clone()
.or_else(|| request.label.clone());
self.record_run_event(WorkflowUiEvent::new(WorkflowUiEventKind::TaskStarted(
Box::new(WorkflowTaskStartedEvent {
task_id: agent_id.to_string(),
label,
role: request.role.clone(),
profile: request.profile.clone(),
model: request.model.clone(),
strength: request.model_strength.clone(),
thinking: request.thinking.clone(),
requested_reasoning: metadata
.requested_reasoning
.clone()
.or_else(|| request.thinking.clone()),
effective_reasoning: metadata.effective_reasoning.clone(),
resolved_role: displayed_resolved_role(
fleet_receipt.as_ref(),
metadata.resolved_role.as_deref(),
request.role.as_deref(),
),
resolved_profile: metadata
.resolved_profile
.clone()
.or_else(|| request.profile.clone()),
resolved_provider: metadata.resolved_provider.clone(),
resolved_model: metadata.resolved_model.clone(),
route_source: metadata.route_source.clone(),
child_route: Some(metadata.child_route.clone()),
worktree: request.worktree,
workspace: result.workspace.clone(),
git_branch: result.git_branch.clone(),
parent_task_id: metadata.parent_task_id.clone(),
depth: metadata.depth,
workflow_run_id: metadata.workflow_run_id.clone(),
workflow_phase_id: metadata.workflow_phase_id.clone(),
workflow_task_label: metadata.workflow_task_label.clone(),
workflow_child_index: metadata.workflow_child_index,
fleet_receipt: fleet_receipt.clone(),
}),
)));
if let Some(receipt) = fleet_receipt {
self.record_run_event(WorkflowUiEvent::new(WorkflowUiEventKind::Log {
message: format!("fleet route {}", receipt.line()),
}));
}
}
fn record_orphaned_fleet_receipt(
&self,
receipt: &codewhale_workflow::FleetTaskReceipt,
error: &str,
) {
self.record_run_event(WorkflowUiEvent::new(WorkflowUiEventKind::Log {
message: orphaned_fleet_receipt_line(receipt, error),
}));
}
fn record_task_request(&self, agent_id: &str, request: &TaskRequest) {
if let Ok(mut records) = self.task_records.lock() {
records.insert(
agent_id.to_string(),
RuntimeTaskRecord {
agent_id: agent_id.to_string(),
label: request.label.clone(),
role: request.role.clone(),
status: IrWorkflowRunStatus::Running,
output: None,
schema_error: None,
usage: None,
},
);
}
let pending_completion = self
.completion_state
.lock()
.ok()
.and_then(|state| state.pending.get(agent_id).cloned());
if let Some(completion) = pending_completion {
self.record_task_completion(agent_id, &completion.completion, completion.usage);
}
}
fn record_task_completion(
&self,
agent_id: &str,
completion: &TaskCompletion,
usage: Option<WorkflowTaskUsage>,
) {
let mut terminal_event = None;
let mut completed_record = None;
if let Ok(mut records) = self.task_records.lock()
&& let Some(record) = records.get_mut(agent_id)
{
let was_running = record.status == IrWorkflowRunStatus::Running;
let (status, output) = task_completion_status(completion);
record.status = status;
record.output = output;
if usage.is_some() {
record.usage = usage;
}
if was_running {
terminal_event = Some(WorkflowUiEvent::new(WorkflowUiEventKind::TaskCompleted {
task_id: agent_id.to_string(),
status,
usage: record.usage.clone(),
}));
completed_record = Some(record.clone());
}
}
if let Some(event) = terminal_event {
self.record_run_event(event);
}
if let Some(record) = completed_record.as_ref() {
self.evaluate_gates_for_completed_role(record);
}
}
fn record_dispatch_failure(
&self,
label: Option<String>,
phase: Option<String>,
message: String,
) {
let failure = WorkflowDispatchFailure {
at_ms: now_ms(),
label,
phase,
message,
};
let slot = failure
.label
.as_deref()
.or(failure.phase.as_deref())
.unwrap_or("task");
let progress_line = format!("dispatch failed for {slot}: {}", failure.message);
let ui_event = WorkflowUiEvent::new(WorkflowUiEventKind::TaskDispatchFailed {
label: failure.label.clone(),
phase: failure.phase.clone(),
message: failure.message.clone(),
});
if let Ok(mut runs) = self.state.runs.lock()
&& let Some(record) = runs.get_mut(&self.run_id)
{
record.progress.push(progress_line.clone());
record.push_event(ui_event.clone());
record.dispatch_failures.push(failure);
}
self.state.record_progress(&self.run_id, &progress_line);
self.state.record_event(&self.run_id, &ui_event);
self.emit_ui_event(&ui_event);
}
fn record_schema_validation_failure(&self, agent_id: &str, message: String) {
if let Ok(mut records) = self.task_records.lock()
&& let Some(record) = records.get_mut(agent_id)
{
record.status = IrWorkflowRunStatus::Failed;
record.schema_error = Some(message.clone());
record.output = Some(message);
}
}
fn task_records_snapshot(&self) -> Vec<RuntimeTaskRecord> {
self.task_records
.lock()
.map(|records| records.values().cloned().collect())
.unwrap_or_default()
}
fn add_waiter_or_complete(&self, agent_id: String, waiter: oneshot::Sender<TaskCompletion>) {
let mut state = self
.completion_state
.lock()
.unwrap_or_else(|poison| poison.into_inner());
if let Some(completion) = state.pending.remove(&agent_id) {
let _ = waiter.send(completion.completion);
} else {
state.waiters.insert(agent_id, waiter);
}
}
fn deliver_completion(
&self,
agent_id: String,
completion: TaskCompletion,
usage: Option<WorkflowTaskUsage>,
) {
self.record_task_completion(&agent_id, &completion, usage.clone());
if let Ok(mut permits) = self.spawn_permits.lock() {
permits.remove(&agent_id);
}
let mut state = self
.completion_state
.lock()
.unwrap_or_else(|poison| poison.into_inner());
if let Some(waiter) = state.waiters.remove(&agent_id) {
let _ = waiter.send(completion);
} else {
state
.pending
.insert(agent_id, PendingCompletion { completion, usage });
}
}
}
#[derive(Clone)]
struct PendingCompletion {
completion: TaskCompletion,
usage: Option<WorkflowTaskUsage>,
}
#[derive(Default)]
struct CompletionState {
waiters: HashMap<String, oneshot::Sender<TaskCompletion>>,
pending: HashMap<String, PendingCompletion>,
}
impl SubAgentWorkflowDriver {
async fn spawn_task_admitted(
&self,
mut request: TaskRequest,
) -> Result<SpawnedTask, DriverError> {
let exact_binding = if let Some(operation) = self.fleet.exact() {
if self.runtime.would_exceed_depth() {
return Err(DriverError::Rejected(format!(
"fleet `{}`: sub-agent depth limit reached (depth {}, max {}); this task \
cannot spawn a child, so it is refused before the reasoning router is asked \
anything.",
operation.snapshot().fleet().qualified(),
self.runtime.spawn_depth,
self.runtime.max_spawn_depth,
)));
}
Some(bind_exact_fleet_task_request(
operation,
crate::fleet::exact::session_permission_ceiling(&self.runtime),
&mut request,
)?)
} else {
apply_named_fleet_to_task_request(self.fleet.legacy_roles(), &mut request).map_err(
|err| {
if let Some(fleet) = self.fleet_name.as_deref() {
DriverError::Rejected(format!(
"fleet `{fleet}` role resolution failed: {err}"
))
} else {
err
}
},
)?;
None
};
let consumed_handoffs = self.prepare_request_for_gates(&mut request)?;
let permit = self
.concurrent_gate
.clone()
.acquire_owned()
.await
.map_err(|_| DriverError::Rejected("workflow concurrent admission closed".into()))?;
let fleet_receipt = match (self.fleet.exact(), exact_binding) {
(Some(operation), Some(binding)) => {
match route_admitted_exact_task(operation, &binding, &mut request).await {
Ok(receipt) => Some(receipt),
Err(err) => {
drop(permit);
return Err(err);
}
}
}
_ => None,
};
let runtime = self
.runtime
.clone()
.with_parent_completion_tx(self.completion_tx.clone());
let request_record = request.clone();
let workflow_child_index = self.child_counter.fetch_add(1, Ordering::SeqCst);
let workflow_phase_id = request
.phase
.as_ref()
.map(|phase| phase.trim())
.filter(|phase| !phase.is_empty())
.map(str::to_string)
.or_else(|| {
self.current_phase
.lock()
.ok()
.and_then(|phase| phase.clone())
});
let workflow_task_label = request
.label
.as_ref()
.map(|label| label.trim())
.filter(|label| !label.is_empty())
.map(str::to_string);
let identity = WorkflowTaskSpawnIdentity {
workflow_run_id: self.run_id.clone(),
workflow_phase_id,
workflow_task_label,
workflow_child_index,
fleet_authority_fingerprint: fleet_receipt
.as_ref()
.and_then(|receipt| receipt.authority_fingerprint.clone()),
};
let result =
match spawn_workflow_task(request, self.manager.clone(), runtime, identity).await {
Ok(result) => result,
Err(err) => {
drop(permit);
if let Some(receipt) = fleet_receipt {
self.record_orphaned_fleet_receipt(&receipt, &err.to_string());
}
return Err(DriverError::Rejected(err.to_string()));
}
};
let task_id = result.result.agent_id.clone();
if let Ok(mut permits) = self.spawn_permits.lock() {
permits.insert(task_id.clone(), permit);
}
self.record_child(&task_id);
self.record_task_started(
&task_id,
&request_record,
&result.metadata,
&result.result,
fleet_receipt,
);
for artifact in consumed_handoffs {
self.record_run_event(WorkflowUiEvent::new(WorkflowUiEventKind::HandoffConsumed {
artifact_id: artifact.id,
kind: artifact.kind,
from_role: artifact.from_role,
to_role: artifact.to_role,
consumer_task_id: task_id.clone(),
}));
}
self.record_task_request(&task_id, &request_record);
if let Some(limit) = self.total_budget {
let mut manager = self.manager.write().await;
manager.attach_shared_budget_scope(&task_id, &self.run_id, limit);
}
let (tx, rx) = oneshot::channel();
self.add_waiter_or_complete(task_id.clone(), tx);
Ok(SpawnedTask {
task_id,
completion: rx,
})
}
}
#[async_trait]
impl WorkflowDriver for SubAgentWorkflowDriver {
async fn spawn_task(&self, request: TaskRequest) -> Result<SpawnedTask, DriverError> {
let label = request
.label
.as_deref()
.map(str::trim)
.filter(|label| !label.is_empty())
.map(str::to_string);
let phase = request
.phase
.as_deref()
.map(str::trim)
.filter(|phase| !phase.is_empty())
.map(str::to_string)
.or_else(|| {
self.current_phase
.lock()
.ok()
.and_then(|phase| phase.clone())
});
let result = self.spawn_task_admitted(request).await;
if let Err(err) = &result {
self.record_dispatch_failure(label, phase, err.to_string());
}
result
}
fn cancel_all(&self) {
self.force_cancel_all();
}
fn budget(&self) -> BudgetSnapshot {
let snapshot = self.current_budget_snapshot();
self.record_budget_snapshot(snapshot);
snapshot
}
fn progress(&self, event: ProgressEvent) {
let mut schema_error = None;
let (message, ui_event) = match event {
ProgressEvent::TaskRejected {
label,
phase,
message,
} => {
self.record_dispatch_failure(label, phase, message);
return;
}
ProgressEvent::Log { message } => (
format!("log: {message}"),
WorkflowUiEvent::new(WorkflowUiEventKind::Log { message }),
),
ProgressEvent::Phase { title } => {
if let Ok(mut current) = self.current_phase.lock() {
*current = Some(title.clone());
}
(
format!("phase: {title}"),
WorkflowUiEvent::new(WorkflowUiEventKind::PhaseStarted { title }),
)
}
ProgressEvent::TaskSchemaValidationFailed { task_id, message } => {
self.record_schema_validation_failure(&task_id, message.clone());
schema_error = Some(WorkflowSchemaError {
task_id: task_id.clone(),
message: message.clone(),
});
(
format!("schema validation failed for {task_id}: {message}"),
WorkflowUiEvent::new(WorkflowUiEventKind::TaskSchemaValidationFailed {
task_id,
message,
}),
)
}
};
if let Ok(mut runs) = self.state.runs.lock()
&& let Some(record) = runs.get_mut(&self.run_id)
{
record.progress.push(message.clone());
record.push_event(ui_event.clone());
if let Some(schema_error) = schema_error {
record.schema_errors.push(schema_error);
}
}
self.state.record_progress(&self.run_id, &message);
self.state.record_event(&self.run_id, &ui_event);
self.emit_ui_event(&ui_event);
}
}
fn budget_event_kind(snapshot: BudgetSnapshot) -> WorkflowUiEventKind {
WorkflowUiEventKind::BudgetUpdated {
total: snapshot.total,
spent: snapshot.spent,
remaining: snapshot.remaining(),
}
}
fn gate_kind_label(kind: GateKind) -> &'static str {
match kind {
GateKind::Verify => "verify",
GateKind::Review => "review",
GateKind::Approve => "approve",
}
}
fn gate_state_reason(state: &GateState) -> String {
state
.blocked_reason()
.map(str::to_string)
.unwrap_or_else(|| state.as_str().to_string())
}
fn append_handoff_context(request: &mut TaskRequest, handoffs: &[HandoffArtifact]) {
request
.description
.push_str("\n\nWorkflow handoff artifacts available for this role:\n");
for artifact in handoffs {
request.description.push_str(&format!(
"- id: {} kind: {} from: {} to: {}\n payload: {}\n",
artifact.id,
artifact.kind,
artifact.from_role,
artifact.to_role,
compact_handoff_payload(&artifact.payload, WORKFLOW_HANDOFF_MAX_CHARS)
));
}
}
fn compact_handoff_payload(payload: &str, max_chars: usize) -> String {
let trimmed = payload.trim();
if trimmed.chars().count() <= max_chars {
return trimmed.to_string();
}
let mut out = trimmed.chars().take(max_chars).collect::<String>();
out.push_str("...");
out
}
fn task_completion_status(completion: &TaskCompletion) -> (IrWorkflowRunStatus, Option<String>) {
match completion {
TaskCompletion::Completed { text } => (IrWorkflowRunStatus::Succeeded, Some(text.clone())),
TaskCompletion::Failed { message } => (IrWorkflowRunStatus::Failed, Some(message.clone())),
TaskCompletion::Cancelled => (IrWorkflowRunStatus::Cancelled, None),
TaskCompletion::BudgetExhausted { message } => {
(IrWorkflowRunStatus::BudgetExceeded, Some(message.clone()))
}
}
}
fn run_usage_totals(records: &[RuntimeTaskRecord]) -> Option<WorkflowRunUsage> {
let mut usages = records.iter().filter_map(|record| record.usage.as_ref());
let mut totals = WorkflowRunUsage::from_task(usages.next()?);
for usage in usages {
totals.add_task(usage);
}
Some(totals)
}
fn workflow_usage_from_task(usage: &WorkflowTaskUsage) -> WorkflowUsage {
WorkflowUsage {
input_tokens: usage.input_tokens,
output_tokens: usage.output_tokens,
cost_microusd: usage.cost_microusd,
}
}
fn execution_from_declarative_spec(
spec: &WorkflowSpec,
records: Vec<RuntimeTaskRecord>,
terminal_status: WorkflowRunStatus,
) -> IrWorkflowExecution {
let by_label = records
.into_iter()
.filter_map(|record| record.label.clone().map(|label| (label, record)))
.collect::<HashMap<_, _>>();
let mut execution = IrWorkflowExecution::default();
for node in &spec.nodes {
push_execution_node(node, &by_label, &mut execution);
}
let mut leaf_usage = execution.leaf_results.iter().map(|leaf| leaf.usage);
execution.usage = leaf_usage
.next()
.map_or_else(WorkflowUsage::default, |first| {
leaf_usage.fold(first, |mut totals, usage| {
totals.input_tokens = sum_optional_usage(totals.input_tokens, usage.input_tokens);
totals.output_tokens =
sum_optional_usage(totals.output_tokens, usage.output_tokens);
totals.cost_microusd =
sum_optional_usage(totals.cost_microusd, usage.cost_microusd);
totals
})
});
match terminal_status {
WorkflowRunStatus::Completed | WorkflowRunStatus::Degraded => {}
WorkflowRunStatus::Failed => mark_ir_status(&mut execution, IrWorkflowRunStatus::Failed),
WorkflowRunStatus::Cancelled => {
mark_ir_status(&mut execution, IrWorkflowRunStatus::Cancelled);
}
WorkflowRunStatus::Running => {
execution.status = IrWorkflowRunStatus::Running;
}
}
execution
}
fn push_execution_node(
node: &WorkflowNode,
records: &HashMap<String, RuntimeTaskRecord>,
execution: &mut IrWorkflowExecution,
) {
match node {
WorkflowNode::Leaf(spec) => push_leaf_execution(spec, records, execution),
WorkflowNode::BranchSet(spec) => push_branch_execution(spec, records, execution),
WorkflowNode::Sequence(spec) => push_sequence_execution(spec, records, execution),
WorkflowNode::Reduce(spec) => push_control_execution(
spec.id.as_str(),
ControlNodeKind::Reduce,
records.get(&spec.id),
spec.inputs.clone(),
Some(spec.prompt.clone()),
execution,
),
WorkflowNode::TeacherReview(spec) => push_control_execution(
spec.id.as_str(),
ControlNodeKind::TeacherReview,
records.get(&spec.id),
spec.candidates.clone(),
Some("teacher review not lowered by the production adapter".to_string()),
execution,
),
WorkflowNode::LoopUntil(spec) => push_control_execution(
spec.id.as_str(),
ControlNodeKind::LoopUntil,
records.get(&spec.id),
spec.children.iter().map(declarative_node_id).collect(),
Some("loop_until not lowered by the production adapter".to_string()),
execution,
),
WorkflowNode::Cond(spec) => push_control_execution(
spec.id.as_str(),
ControlNodeKind::Cond,
records.get(&spec.id),
spec.then_nodes
.iter()
.chain(spec.else_nodes.iter())
.map(declarative_node_id)
.collect(),
Some("cond not lowered by the production adapter".to_string()),
execution,
),
WorkflowNode::Expand(spec) => push_control_execution(
spec.id.as_str(),
ControlNodeKind::Expand,
records.get(&spec.id),
Vec::new(),
Some(format!("expand not lowered from {}", spec.source)),
execution,
),
}
}
fn push_leaf_execution(
spec: &LeafSpec,
records: &HashMap<String, RuntimeTaskRecord>,
execution: &mut IrWorkflowExecution,
) {
let record = records.get(&spec.id);
let status = record
.map(|record| record.status)
.unwrap_or(IrWorkflowRunStatus::Pending);
mark_ir_status(execution, status);
execution.leaf_results.push(LeafResult {
leaf_id: spec.id.clone(),
task_id: record
.map(|record| record.agent_id.clone())
.unwrap_or_else(|| spec.id.clone()),
role: spec.role.clone(),
profile: spec.profile.clone(),
status,
usage: record
.and_then(|record| record.usage.as_ref())
.map(workflow_usage_from_task)
.unwrap_or_default(),
memo_usage: WorkflowMemoUsage::default(),
output: record.and_then(|record| record.output.clone()),
artifacts: Vec::new(),
schema_error: record.and_then(|record| record.schema_error.clone()),
});
}
fn push_branch_execution(
spec: &BranchSpec,
records: &HashMap<String, RuntimeTaskRecord>,
execution: &mut IrWorkflowExecution,
) {
let before = execution.leaf_results.len();
for child in &spec.children {
push_execution_node(child, records, execution);
}
let status = aggregate_ir_status(
execution.leaf_results[before..]
.iter()
.map(|result| result.status),
);
mark_ir_status(execution, status);
execution.branch_results.push(BranchResult {
branch_id: spec.id.clone(),
task_id: spec.id.clone(),
status,
usage: WorkflowUsage::default(),
memo_usage: WorkflowMemoUsage::default(),
artifacts: Vec::new(),
notes: Some("production driver branch receipt from child task outcomes".to_string()),
});
execution.control_node_results.push(ControlNodeResult {
node_id: spec.id.clone(),
kind: ControlNodeKind::BranchSet,
status,
selected_children: spec.children.iter().map(declarative_node_id).collect(),
summary: Some("branch set lowered into production child tasks".to_string()),
});
}
fn push_sequence_execution(
spec: &SequenceSpec,
records: &HashMap<String, RuntimeTaskRecord>,
execution: &mut IrWorkflowExecution,
) {
let before_leaf = execution.leaf_results.len();
let before_control = execution.control_node_results.len();
for child in &spec.children {
push_execution_node(child, records, execution);
}
let status = aggregate_ir_status(
execution.leaf_results[before_leaf..]
.iter()
.map(|result| result.status)
.chain(
execution.control_node_results[before_control..]
.iter()
.map(|result| result.status),
),
);
mark_ir_status(execution, status);
execution.control_node_results.push(ControlNodeResult {
node_id: spec.id.clone(),
kind: ControlNodeKind::Sequence,
status,
selected_children: spec.children.iter().map(declarative_node_id).collect(),
summary: Some("sequence lowered in declaration order".to_string()),
});
}
fn push_control_execution(
node_id: &str,
kind: ControlNodeKind,
record: Option<&RuntimeTaskRecord>,
selected_children: Vec<String>,
fallback_summary: Option<String>,
execution: &mut IrWorkflowExecution,
) {
let status = record
.map(|record| record.status)
.unwrap_or(IrWorkflowRunStatus::Pending);
mark_ir_status(execution, status);
execution.control_node_results.push(ControlNodeResult {
node_id: node_id.to_string(),
kind,
status,
selected_children,
summary: record
.and_then(|record| record.output.clone())
.or(fallback_summary),
});
}
fn aggregate_ir_status(
statuses: impl IntoIterator<Item = IrWorkflowRunStatus>,
) -> IrWorkflowRunStatus {
let mut saw_pending = false;
let mut saw_running = false;
for status in statuses {
match status {
IrWorkflowRunStatus::BudgetExceeded => return IrWorkflowRunStatus::BudgetExceeded,
IrWorkflowRunStatus::Cancelled => return IrWorkflowRunStatus::Cancelled,
IrWorkflowRunStatus::Failed | IrWorkflowRunStatus::ReplayDiverged => {
return IrWorkflowRunStatus::Failed;
}
IrWorkflowRunStatus::Running => saw_running = true,
IrWorkflowRunStatus::Pending => saw_pending = true,
IrWorkflowRunStatus::Succeeded => {}
}
}
if saw_running {
IrWorkflowRunStatus::Running
} else if saw_pending {
IrWorkflowRunStatus::Pending
} else {
IrWorkflowRunStatus::Succeeded
}
}
fn mark_ir_status(execution: &mut IrWorkflowExecution, status: IrWorkflowRunStatus) {
match status {
IrWorkflowRunStatus::Failed | IrWorkflowRunStatus::ReplayDiverged => {
execution.mark_failed()
}
IrWorkflowRunStatus::Cancelled => execution.mark_cancelled(),
IrWorkflowRunStatus::BudgetExceeded => execution.mark_budget_exceeded(),
IrWorkflowRunStatus::Running => {
if execution.status == IrWorkflowRunStatus::Succeeded {
execution.status = IrWorkflowRunStatus::Running;
}
}
IrWorkflowRunStatus::Pending => {
if execution.status == IrWorkflowRunStatus::Succeeded {
execution.status = IrWorkflowRunStatus::Pending;
}
}
IrWorkflowRunStatus::Succeeded => {}
}
}
fn declarative_node_id(node: &WorkflowNode) -> String {
match node {
WorkflowNode::BranchSet(spec) => spec.id.clone(),
WorkflowNode::Leaf(spec) => spec.id.clone(),
WorkflowNode::Sequence(spec) => spec.id.clone(),
WorkflowNode::Reduce(spec) => spec.id.clone(),
WorkflowNode::TeacherReview(spec) => spec.id.clone(),
WorkflowNode::LoopUntil(spec) => spec.id.clone(),
WorkflowNode::Cond(spec) => spec.id.clone(),
WorkflowNode::Expand(spec) => spec.id.clone(),
}
}
fn spawn_completion_pump(
driver: Arc<SubAgentWorkflowDriver>,
mut rx: mpsc::UnboundedReceiver<SubAgentCompletion>,
) {
spawn_supervised(
"workflow-completion-pump",
std::panic::Location::caller(),
async move {
while let Some(completion) = rx.recv().await {
let agent_id = completion.agent_id.clone();
let (task_completion, usage) =
completion_from_manager(driver.manager.clone(), &agent_id, completion.payload)
.await;
driver.deliver_completion(agent_id, task_completion, usage);
}
},
);
}
async fn completion_from_manager(
manager: SharedSubAgentManager,
agent_id: &str,
fallback_payload: String,
) -> (TaskCompletion, Option<WorkflowTaskUsage>) {
for _ in 0..50 {
let snapshot_and_usage = {
let manager = manager.read().await;
let snapshot = manager.get_result(agent_id).ok();
let usage = snapshot
.as_ref()
.filter(|snapshot| snapshot.status != SubAgentStatus::Running)
.map(|snapshot| task_usage_from_manager(&manager, agent_id, snapshot));
(snapshot, usage)
};
if let (Some(snapshot), usage) = snapshot_and_usage
&& snapshot.status != SubAgentStatus::Running
{
let completion = match snapshot.status {
SubAgentStatus::Completed => TaskCompletion::Completed {
text: snapshot.result.clone().unwrap_or(fallback_payload),
},
SubAgentStatus::Failed(ref message) => TaskCompletion::Failed {
message: message.clone(),
},
SubAgentStatus::Interrupted(ref message) => TaskCompletion::Failed {
message: message.clone(),
},
SubAgentStatus::Cancelled => TaskCompletion::Cancelled,
SubAgentStatus::BudgetExhausted => TaskCompletion::BudgetExhausted {
message: "sub-agent budget exhausted".to_string(),
},
SubAgentStatus::Running => unreachable!("guarded above"),
};
return (completion, usage);
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
(
TaskCompletion::Failed {
message: format!("sub-agent '{agent_id}' did not report a terminal status within 1s"),
},
None,
)
}
fn task_usage_from_manager(
manager: &SubAgentManager,
agent_id: &str,
snapshot: &SubAgentResult,
) -> WorkflowTaskUsage {
let record = manager.get_worker_record(agent_id);
let usage = record.as_ref().map(|record| &record.usage);
let result_ref = record.as_ref().and_then(|record| {
record
.artifacts
.iter()
.find(|artifact| artifact.kind == "transcript")
.or_else(|| record.artifacts.last())
.map(|artifact| artifact.target.clone())
});
let input_tokens = usage.and_then(|usage| usage.input_tokens);
let output_tokens = usage.and_then(|usage| usage.output_tokens);
let total_tokens = usage.and_then(|usage| usage.total_tokens);
let reported = provider_usage_was_reported(input_tokens, output_tokens, total_tokens);
WorkflowTaskUsage {
input_tokens: reported.then_some(input_tokens).flatten(),
output_tokens: reported.then_some(output_tokens).flatten(),
total_tokens: reported.then_some(total_tokens).flatten(),
cost_microusd: usage.and_then(|usage| usage.cost_microusd),
tool_calls: Some(snapshot.steps_taken),
duration_ms: Some(snapshot.duration_ms),
result_ref,
token_source: reported.then_some(WorkflowTokenSource::ProviderReported),
}
}
fn provider_usage_was_reported(
input_tokens: Option<u64>,
output_tokens: Option<u64>,
total_tokens: Option<u64>,
) -> bool {
[input_tokens, output_tokens, total_tokens]
.iter()
.any(Option::is_some)
}
fn cancel_child_agents(manager: SharedSubAgentManager, ids: Vec<String>) {
if ids.is_empty() {
return;
}
if let Ok(mut manager_guard) = manager.try_write() {
for id in ids {
let _ = manager_guard.cancel_agent(&id);
}
return;
}
if tokio::runtime::Handle::try_current().is_ok() {
spawn_supervised(
"workflow-cancel-children",
std::panic::Location::caller(),
async move {
let mut manager_guard = manager.write().await;
for id in ids {
let _ = manager_guard.cancel_agent(&id);
}
},
);
}
}
fn lock_mutex<T>(mutex: &Mutex<T>) -> Result<MutexGuard<'_, T>, ToolError> {
mutex
.lock()
.map_err(|_| ToolError::execution_failed("workflow state lock poisoned"))
}
fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.try_into()
.unwrap_or(u64::MAX)
}
mod journal {
use super::{
SharedWorkflowControllers, SharedWorkflowLifecycles, SharedWorkflowRuns, WorkflowRunRecord,
WorkflowRunStatus, WorkflowUiEvent, WorkflowWorkLifecycle,
};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::fs::OpenOptions;
use std::io::{BufRead, Write};
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, OnceLock};
use tracing::warn;
pub(super) const CODEWHALE_DIR: &str = ".codewhale";
pub(super) const WORKFLOW_RUNS_FILE: &str = "workflow-runs.jsonl";
pub(super) struct WorkflowWorkspaceState {
pub runs: SharedWorkflowRuns,
pub controllers: SharedWorkflowControllers,
lifecycles: SharedWorkflowLifecycles,
journal: WorkflowRunJournal,
}
impl WorkflowWorkspaceState {
pub fn open(workspace: &Path) -> Arc<Self> {
Self::open_inner(workspace, true)
}
pub fn open_preserving_running(workspace: &Path) -> Arc<Self> {
Self::open_inner(workspace, false)
}
fn open_inner(workspace: &Path, recover_orphans: bool) -> Arc<Self> {
let journal = WorkflowRunJournal::open(workspace);
let runs = Arc::new(Mutex::new(journal.hydrate_runs(recover_orphans)));
Arc::new(Self {
runs,
controllers: Arc::new(Mutex::new(HashMap::new())),
lifecycles: Arc::new(Mutex::new(HashMap::new())),
journal,
})
}
pub fn attach_lifecycle(&self, run_id: &str, lifecycle: WorkflowWorkLifecycle) {
self.lifecycles
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.entry(run_id.to_string())
.or_insert(lifecycle);
}
pub fn reconcile_snapshot(&self, record: &WorkflowRunRecord) {
let lifecycle = self
.lifecycles
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.get(&record.run_id)
.cloned();
if let Some(lifecycle) = lifecycle
&& let Err(err) = lifecycle.reconcile_record(record)
{
warn!(
run_id = record.run_id,
"workflow Work reconciliation failed: {err}"
);
}
}
pub fn reconcile_cancel(&self, run_id: &str, outcome: super::CancelOutcome) {
let lifecycle = self
.lifecycles
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.get(run_id)
.cloned();
if let Some(lifecycle) = lifecycle
&& let Err(err) = lifecycle.reconcile_cancel(outcome)
{
warn!(run_id, "workflow cancellation reconciliation failed: {err}");
}
}
pub fn mark_owner_missing(&self, run_id: &str) {
let lifecycle = self
.lifecycles
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.get(run_id)
.cloned();
if let Some(lifecycle) = lifecycle {
lifecycle.reconcile_missing();
}
}
pub fn try_record_snapshot(&self, record: &WorkflowRunRecord) -> Result<(), String> {
self.journal
.append_snapshot(record)
.map_err(|err| err.to_string())
}
pub fn record_snapshot(&self, record: &WorkflowRunRecord) {
if let Err(err) = self.try_record_snapshot(record) {
warn!("workflow journal snapshot failed: {err}");
}
}
pub fn record_progress(&self, run_id: &str, message: &str) {
if let Err(err) = self.journal.append_progress(run_id, message) {
warn!("workflow journal progress failed: {err}");
}
}
pub fn record_event(&self, run_id: &str, event: &WorkflowUiEvent) {
if let Err(err) = self.journal.append_event(run_id, event) {
warn!("workflow journal event failed: {err}");
}
}
pub fn journal_path(&self) -> &Path {
&self.journal.ledger_path
}
}
fn workspace_store() -> &'static Mutex<HashMap<PathBuf, Arc<WorkflowWorkspaceState>>> {
static STORE: OnceLock<Mutex<HashMap<PathBuf, Arc<WorkflowWorkspaceState>>>> =
OnceLock::new();
STORE.get_or_init(|| Mutex::new(HashMap::new()))
}
pub(super) fn shared_workflow_state(workspace: &Path) -> Arc<WorkflowWorkspaceState> {
let key = workspace
.canonicalize()
.unwrap_or_else(|_| workspace.to_path_buf());
let mut store = workspace_store()
.lock()
.unwrap_or_else(|poison| poison.into_inner());
store
.entry(key)
.or_insert_with(|| WorkflowWorkspaceState::open(workspace))
.clone()
}
pub(super) fn peek_shared_workflow_state(
workspace: &Path,
) -> Option<Arc<WorkflowWorkspaceState>> {
let key = workspace
.canonicalize()
.unwrap_or_else(|_| workspace.to_path_buf());
workspace_store()
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.get(&key)
.cloned()
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
enum WorkflowJournalRecord {
Snapshot {
run: Box<WorkflowRunRecord>,
},
Progress {
run_id: String,
message: String,
},
Event {
run_id: String,
event: Box<WorkflowUiEvent>,
},
}
#[derive(Debug)]
struct WorkflowRunJournal {
ledger_path: PathBuf,
}
impl WorkflowRunJournal {
fn open(workspace: &Path) -> Self {
let dir = workspace.join(CODEWHALE_DIR);
if let Err(err) = std::fs::create_dir_all(&dir) {
warn!(
"workflow journal dir create failed ({}): {err}",
dir.display()
);
}
let ledger_path = dir.join(WORKFLOW_RUNS_FILE);
if !ledger_path.exists()
&& let Err(err) = std::fs::write(&ledger_path, "")
{
warn!(
"workflow journal create failed ({}): {err}",
ledger_path.display()
);
}
Self { ledger_path }
}
fn hydrate_runs(&self, recover_orphans: bool) -> HashMap<String, WorkflowRunRecord> {
let file = match std::fs::File::open(&self.ledger_path) {
Ok(file) => file,
Err(_) => return HashMap::new(),
};
let mut runs = HashMap::new();
for line in std::io::BufReader::new(file).lines() {
let Ok(line) = line else { continue };
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
let record = match serde_json::from_str::<WorkflowJournalRecord>(trimmed) {
Ok(record) => record,
Err(err) => {
warn!("workflow journal skipped malformed line: {err}");
continue;
}
};
match record {
WorkflowJournalRecord::Snapshot { run } => {
runs.insert(run.run_id.clone(), *run);
}
WorkflowJournalRecord::Progress { run_id, message } => {
if let Some(run) = runs.get_mut(&run_id) {
run.progress.push(message);
}
}
WorkflowJournalRecord::Event { run_id, event } => {
if let Some(run) = runs.get_mut(&run_id) {
run.push_event(*event);
}
}
}
}
for run in runs.values_mut() {
run.events_total = run.events_total.max(run.events.len() as u64);
}
if recover_orphans {
let mut recovered = Vec::new();
for run in runs.values_mut() {
if run.status == WorkflowRunStatus::Running {
run.status = WorkflowRunStatus::Failed;
run.lifecycle_seq = run.lifecycle_seq.saturating_add(1);
run.completed_at_ms.get_or_insert_with(super::now_ms);
run.error = Some(
"process exited before the run completed (recovered on startup)"
.to_string(),
);
recovered.push(run.clone());
}
}
for run in recovered {
if let Err(err) = self.append_snapshot(&run) {
warn!(
run_id = run.run_id,
"workflow recovery snapshot append failed: {err}"
);
}
}
}
runs
}
fn append_record(&self, record: &WorkflowJournalRecord) -> std::io::Result<()> {
let mut line = serde_json::to_string(record)
.map_err(|err| std::io::Error::other(err.to_string()))?;
line.push('\n');
let mut file = OpenOptions::new()
.create(true)
.append(true)
.open(&self.ledger_path)?;
file.write_all(line.as_bytes())?;
file.flush()?;
Ok(())
}
fn append_snapshot(&self, record: &WorkflowRunRecord) -> std::io::Result<()> {
self.append_record(&WorkflowJournalRecord::Snapshot {
run: Box::new(record.clone()),
})
}
fn append_progress(&self, run_id: &str, message: &str) -> std::io::Result<()> {
self.append_record(&WorkflowJournalRecord::Progress {
run_id: run_id.to_string(),
message: message.to_string(),
})
}
fn append_event(&self, run_id: &str, event: &WorkflowUiEvent) -> std::io::Result<()> {
self.append_record(&WorkflowJournalRecord::Event {
run_id: run_id.to_string(),
event: Box::new(event.clone()),
})
}
}
#[cfg(test)]
mod tests {
use super::super::WorkflowUiEventKind;
use super::*;
fn sample_record(run_id: &str, status: WorkflowRunStatus) -> WorkflowRunRecord {
WorkflowRunRecord {
run_id: run_id.to_string(),
status,
lifecycle_seq: 1,
started_at_ms: 1,
completed_at_ms: None,
source_path: None,
workflow_id: Some("fixture".to_string()),
workflow_goal: Some("journal test".to_string()),
token_budget: None,
child_ids: Vec::new(),
progress: Vec::new(),
events: Vec::new(),
schema_errors: Vec::new(),
dispatch_failures: Vec::new(),
result: None,
execution: None,
error: None,
verify_on_complete: false,
verification: None,
plan_approval: None,
gate_status: Vec::new(),
usage: None,
events_total: 0,
events_dropped: 0,
}
}
#[test]
fn workflow_journal_hydrates_snapshots_and_progress() {
let tmp = tempfile::tempdir().expect("tempdir");
let state = WorkflowWorkspaceState::open(tmp.path());
let running = sample_record("workflow_abc", WorkflowRunStatus::Running);
state.record_snapshot(&running);
state.record_progress("workflow_abc", "phase: scan");
state.record_event(
"workflow_abc",
&WorkflowUiEvent::at(
5,
WorkflowUiEventKind::PhaseStarted {
title: "scan".to_string(),
},
),
);
let completed = WorkflowRunRecord {
status: WorkflowRunStatus::Completed,
completed_at_ms: Some(99),
progress: vec!["phase: scan".to_string()],
events: vec![WorkflowUiEvent::at(
5,
WorkflowUiEventKind::PhaseStarted {
title: "scan".to_string(),
},
)],
..sample_record("workflow_abc", WorkflowRunStatus::Completed)
};
state.record_snapshot(&completed);
state.record_event(
"workflow_abc",
&WorkflowUiEvent::at(
6,
WorkflowUiEventKind::HandoffPromoted {
artifact_id: "workflow_abc:scout-1:scout-gate:findings".to_string(),
gate_id: "scout-gate".to_string(),
kind: "findings".to_string(),
from_role: "scout".to_string(),
to_role: "implementer".to_string(),
producer_task_id: "scout-1".to_string(),
},
),
);
state.record_event(
"workflow_abc",
&WorkflowUiEvent::at(
7,
WorkflowUiEventKind::HandoffConsumed {
artifact_id: "workflow_abc:scout-1:scout-gate:findings".to_string(),
kind: "findings".to_string(),
from_role: "scout".to_string(),
to_role: "implementer".to_string(),
consumer_task_id: "implementer-1".to_string(),
},
),
);
let reloaded = WorkflowWorkspaceState::open(tmp.path());
let runs = reloaded
.runs
.lock()
.expect("runs lock")
.get("workflow_abc")
.cloned()
.expect("hydrated run");
assert_eq!(runs.status, WorkflowRunStatus::Completed);
assert_eq!(runs.progress, vec!["phase: scan"]);
assert_eq!(runs.events.len(), 3);
assert_eq!(runs.events[0].event_type(), "phase_started");
let promoted = serde_json::to_value(&runs.events[1]).expect("promoted receipt");
assert_eq!(promoted["type"], "handoff_promoted");
assert_eq!(
promoted["artifact_id"],
"workflow_abc:scout-1:scout-gate:findings"
);
assert_eq!(promoted["gate_id"], "scout-gate");
assert_eq!(promoted["producer_task_id"], "scout-1");
assert!(promoted.get("payload").is_none(), "{promoted}");
let consumed = serde_json::to_value(&runs.events[2]).expect("consumed receipt");
assert_eq!(consumed["type"], "handoff_consumed");
assert_eq!(consumed["artifact_id"], promoted["artifact_id"]);
assert_eq!(consumed["consumer_task_id"], "implementer-1");
assert!(consumed.get("payload").is_none(), "{consumed}");
assert_eq!(runs.completed_at_ms, Some(99));
reloaded.record_snapshot(&runs);
let reopened = WorkflowWorkspaceState::open(tmp.path());
let compacted = reopened
.runs
.lock()
.expect("runs lock")
.get("workflow_abc")
.cloned()
.expect("snapshot with handoff receipts");
assert_eq!(
compacted
.events
.iter()
.map(WorkflowUiEvent::event_type)
.collect::<Vec<_>>(),
vec!["phase_started", "handoff_promoted", "handoff_consumed"]
);
}
#[test]
fn workflow_journal_marks_orphaned_running_runs_failed() {
let tmp = tempfile::tempdir().expect("tempdir");
let state = WorkflowWorkspaceState::open(tmp.path());
state.record_snapshot(&sample_record(
"workflow_orphan",
WorkflowRunStatus::Running,
));
let reloaded = WorkflowWorkspaceState::open(tmp.path());
let run = reloaded
.runs
.lock()
.expect("runs lock")
.get("workflow_orphan")
.cloned()
.expect("hydrated run");
assert_eq!(run.status, WorkflowRunStatus::Failed);
assert_eq!(
run.lifecycle_seq, 2,
"restart recovery is a durable owner lifecycle transition"
);
assert!(
run.completed_at_ms.is_some(),
"restart recovery must terminalize the durable owner record"
);
assert!(
run.error
.as_deref()
.is_some_and(|error| error.contains("process exited")),
"expected orphan recovery error, got {:?}",
run.error
);
let reopened = WorkflowWorkspaceState::open(tmp.path());
let replayed = reopened
.runs
.lock()
.expect("runs lock")
.get("workflow_orphan")
.cloned()
.expect("durably recovered run");
assert_eq!(replayed.status, WorkflowRunStatus::Failed);
assert_eq!(
replayed.lifecycle_seq, 2,
"reopening must replay the recovery snapshot without another transition"
);
}
#[test]
fn host_cancel_hydrates_a_journal_without_live_process_state() {
let tmp = tempfile::tempdir().expect("tempdir");
let state = WorkflowWorkspaceState::open(tmp.path());
let record = sample_record("workflow_prior", WorkflowRunStatus::Running);
state.record_snapshot(&record);
drop(state);
assert!(
peek_shared_workflow_state(tmp.path()).is_none(),
"writing the journal must not insert process-wide live state"
);
let line = super::super::host_cancel_workflow(tmp.path(), "workflow_prior")
.expect("a journaled run must be visible to host cancel after restart");
assert_eq!(line.run_id, "workflow_prior");
assert_eq!(line.status, "cancelled");
assert!(
line.error
.as_deref()
.is_some_and(|error| error.contains("no live process")),
"controller-less cancel must leave an honest receipt, got {:?}",
line.error
);
let reopened = WorkflowWorkspaceState::open(tmp.path());
let replayed = reopened
.runs
.lock()
.expect("runs lock")
.get("workflow_prior")
.cloned()
.expect("cancelled journal line");
assert_eq!(replayed.status, WorkflowRunStatus::Cancelled);
}
}
}
use journal::{WorkflowWorkspaceState, peek_shared_workflow_state, shared_workflow_state};
pub(crate) fn structcopy_run_projection(workspace: &Path, run_id: &str) -> Option<Value> {
let state = peek_shared_workflow_state(workspace)?;
let runs = state
.runs
.lock()
.unwrap_or_else(|poison| poison.into_inner());
let record = runs.get(run_id)?;
let mut value = serde_json::to_value(record.summary()).ok()?;
if let Some(object) = value.as_object_mut() {
let source_file = record
.source_path
.as_ref()
.and_then(|path| path.file_name())
.and_then(|name| name.to_str())
.map(|name| Value::String(name.to_string()))
.unwrap_or(Value::Null);
object.remove("source_path");
object.insert("source_file".to_string(), source_file);
if record.execution.is_none() {
object.insert("leaf_count".to_string(), Value::Null);
object.insert("branch_count".to_string(), Value::Null);
object.insert("control_count".to_string(), Value::Null);
}
}
Some(value)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct HostWorkflowRunLine {
pub run_id: String,
pub status: &'static str,
pub label: String,
pub started_at_ms: u64,
pub completed_at_ms: Option<u64>,
pub child_count: usize,
pub last_progress: Option<String>,
pub error: Option<String>,
}
fn host_run_line(record: &WorkflowRunRecord) -> HostWorkflowRunLine {
let summary = record.summary();
let label = summary
.workflow_goal
.clone()
.or_else(|| summary.workflow_id.clone())
.or_else(|| {
summary
.source_path
.as_ref()
.and_then(|path| path.file_name())
.map(|name| name.to_string_lossy().into_owned())
})
.unwrap_or_else(|| "workflow".to_string());
HostWorkflowRunLine {
run_id: summary.run_id,
status: match summary.status {
WorkflowRunStatus::Running => "running",
WorkflowRunStatus::Completed => "completed",
WorkflowRunStatus::Degraded => "degraded",
WorkflowRunStatus::Failed => "failed",
WorkflowRunStatus::Cancelled => "cancelled",
},
label,
started_at_ms: summary.started_at_ms,
completed_at_ms: summary.completed_at_ms,
child_count: summary.child_count,
last_progress: summary.last_progress,
error: summary.error,
}
}
fn host_workflow_state(workspace: &Path) -> Option<Arc<WorkflowWorkspaceState>> {
if let Some(state) = peek_shared_workflow_state(workspace) {
return Some(state);
}
workflow_journal_exists(workspace).then(|| shared_workflow_state(workspace))
}
fn host_workflow_state_for_cancel(workspace: &Path) -> Option<Arc<WorkflowWorkspaceState>> {
if let Some(state) = peek_shared_workflow_state(workspace) {
return Some(state);
}
workflow_journal_exists(workspace)
.then(|| WorkflowWorkspaceState::open_preserving_running(workspace))
}
fn workflow_journal_exists(workspace: &Path) -> bool {
workspace
.join(journal::CODEWHALE_DIR)
.join(journal::WORKFLOW_RUNS_FILE)
.is_file()
}
pub(crate) fn host_workflow_runs(workspace: &Path) -> Vec<HostWorkflowRunLine> {
let Some(state) = host_workflow_state(workspace) else {
return Vec::new();
};
let runs = state
.runs
.lock()
.unwrap_or_else(|poison| poison.into_inner());
let mut lines: Vec<HostWorkflowRunLine> = runs.values().map(host_run_line).collect();
lines.sort_by_key(|line| line.started_at_ms);
lines
}
pub(crate) fn host_cancel_workflow(
workspace: &Path,
run_id: &str,
) -> Result<HostWorkflowRunLine, String> {
let Some(state) = host_workflow_state_for_cancel(workspace) else {
return Err(format!("Unknown workflow run '{run_id}'."));
};
match cancel_workflow_run(run_id, state.clone()) {
Ok(_) => {
let runs = state
.runs
.lock()
.unwrap_or_else(|poison| poison.into_inner());
runs.get(run_id)
.map(host_run_line)
.ok_or_else(|| format!("Unknown workflow run '{run_id}'."))
}
Err(err) => Err(err.to_string()),
}
}
#[cfg(test)]
pub(crate) fn structcopy_test_seed_run(workspace: &Path, run_id: &str) {
let state = shared_workflow_state(workspace);
let record = WorkflowRunRecord::new(run_id.to_string(), None, None, None);
state
.runs
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.insert(run_id.to_string(), record);
}
pub(crate) fn reconcile_persisted_workflow_bindings(
work: &SharedWorkRuntime,
session_id: &str,
workspace: &Path,
) -> Result<usize, String> {
let state = shared_workflow_state(workspace);
let records = state
.runs
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.values()
.cloned()
.collect::<Vec<_>>();
let candidates = work
.reconcilable_durable_bindings(Some(session_id))
.into_iter()
.filter(|external| external.starts_with("workflow:"))
.collect::<std::collections::HashSet<_>>();
let mut seen = std::collections::HashSet::new();
let mut changed = 0usize;
for record in records {
let external = format!("workflow:{}", record.run_id);
if !candidates.contains(&external) {
continue;
}
seen.insert(external.clone());
let lifecycle = WorkflowWorkLifecycle {
work: work.clone(),
session_id: session_id.to_string(),
external,
};
changed += usize::from(lifecycle.reconcile_record(&record)?);
}
for external in candidates.difference(&seen) {
changed += usize::from(work.reconcile_observation(
session_id,
external,
OperationObservation::OwnerMissing {
checked_at: i64::try_from(now_ms()).unwrap_or(i64::MAX),
},
)?);
}
Ok(changed)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::client::DeepSeekClient;
use crate::tools::ToolRegistryBuilder;
use crate::tools::subagent::{SubAgentRuntime, new_shared_subagent_manager};
use axum::{Json, Router, routing::post};
use codewhale_workflow::{IsolationMode, leaf_is_write_capable};
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn settled_runs_leave_a_report_artifact_under_codewhale_reports() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut record = WorkflowRunRecord::new("workflow_report_1".to_string(), None, None, None);
record.status = WorkflowRunStatus::Completed;
record.workflow_goal = Some("prove the report artifact".to_string());
record.progress.push("phase: scan".to_string());
record.result = Some(serde_json::json!({"confirmed": 2}));
write_run_report_artifact(tmp.path(), &record);
let path = tmp
.path()
.join(".codewhale")
.join("reports")
.join("workflow_report_1.md");
let body = std::fs::read_to_string(&path).expect("report written");
assert!(body.contains("# Workflow run workflow_report_1"), "{body}");
assert!(body.contains("status: Completed"), "{body}");
assert!(body.contains("prove the report artifact"), "{body}");
assert!(body.contains("phase: scan"), "{body}");
assert!(body.contains("\"confirmed\": 2"), "{body}");
}
#[test]
fn running_runs_write_no_report_artifact() {
let tmp = tempfile::tempdir().expect("tempdir");
let record = WorkflowRunRecord::new("workflow_report_2".to_string(), None, None, None);
write_run_report_artifact(tmp.path(), &record);
assert!(
!tmp.path().join(".codewhale").join("reports").exists(),
"running runs must not leave report files"
);
}
#[test]
fn source_path_accepts_the_home_workflow_store_and_rejects_elsewhere() {
let _lock = crate::test_support::lock_test_env();
let tmp = tempfile::tempdir().expect("tempdir");
let home = tmp.path().join("home");
let store = home.join(".codewhale").join("workflows");
std::fs::create_dir_all(&store).expect("store");
let _home_guard = crate::test_support::EnvVarGuard::set("HOME", &home);
let _userprofile_guard = crate::test_support::EnvVarGuard::set("USERPROFILE", &home);
let saved = store.join("triage.workflow.js");
std::fs::write(&saved, "phase('scan');\n").expect("write saved workflow");
let elsewhere = tmp.path().join("outside.workflow.js");
std::fs::write(&elsewhere, "phase('scan');\n").expect("write outside workflow");
let workspace = tmp.path().join("ws");
std::fs::create_dir_all(&workspace).expect("workspace");
let context = ToolContext::new(workspace);
let resolved = read_workflow_source_path(saved.to_str().expect("utf8 path"), &context)
.expect("home workflow store is a first-class source");
assert!(resolved.source.contains("phase('scan')"));
let err = read_workflow_source_path(elsewhere.to_str().expect("utf8 path"), &context)
.expect_err("arbitrary outside paths stay denied");
assert!(
err.to_string()
.contains("workspace or ~/.codewhale/workflows"),
"{err}"
);
}
#[test]
fn restored_workflow_binding_consumes_journal_recovery() {
let tmp = tempfile::tempdir().expect("tempdir");
let state = WorkflowWorkspaceState::open(tmp.path());
let record = WorkflowRunRecord::new("workflow_restore".to_string(), None, None, None);
state.record_snapshot(&record);
let work = crate::work_graph::new_shared_work_runtime(
crate::tools::todo::new_shared_todo_list(),
crate::tools::plan::new_shared_plan_state(),
);
work.register_operation(
"restored-workflow-session",
OperationIntent::new(
"workflow:workflow_restore",
"restored workflow",
true,
"workflow",
"restore-test",
),
)
.expect("register saved workflow binding");
work.reconcile_operation(
"restored-workflow-session",
OperationOwnerSnapshot::new("workflow:workflow_restore", OwnerState::Running, 1, 1),
)
.expect("saved running owner state");
work.register_operation(
"restored-workflow-session",
OperationIntent::new(
"workflow:workflow_absent",
"absent workflow",
true,
"workflow",
"absent-restore-test",
),
)
.expect("register absent workflow binding");
work.reconcile_operation(
"restored-workflow-session",
OperationOwnerSnapshot::new("workflow:workflow_absent", OwnerState::Running, 1, 1),
)
.expect("saved absent owner state");
assert_eq!(
reconcile_persisted_workflow_bindings(&work, "restored-workflow-session", tmp.path(),),
Ok(2)
);
let graph = work
.capture(Some("restored-workflow-session"))
.expect("capture restored workflow")
.expect("graph")
.graph;
let operation = graph
.nodes
.iter()
.find(|node| {
node.binding
.as_ref()
.is_some_and(|binding| binding.external == "workflow:workflow_restore")
})
.expect("workflow operation");
assert_eq!(operation.state, crate::work_graph::NodeState::Failed);
assert_eq!(
operation
.binding
.as_ref()
.and_then(|binding| binding.last_observation.as_ref())
.map(|observation| observation.seq),
Some(2),
"journal replay must advance the lost live owner before graph reconciliation"
);
let absent = graph
.nodes
.iter()
.find(|node| {
node.binding
.as_ref()
.is_some_and(|binding| binding.external == "workflow:workflow_absent")
})
.expect("absent workflow operation");
assert_eq!(absent.state, crate::work_graph::NodeState::Stale);
assert_eq!(
reconcile_persisted_workflow_bindings(&work, "restored-workflow-session", tmp.path(),),
Ok(0),
"rechecking an already stale missing owner must be idempotent"
);
}
#[tokio::test]
async fn cancellation_without_controller_marks_the_journal_cancelled() {
let tmp = tempfile::tempdir().expect("tempdir");
let state = WorkflowWorkspaceState::open(tmp.path());
let record =
WorkflowRunRecord::new("workflow_missing_controller".to_string(), None, None, None);
state
.runs
.lock()
.expect("runs lock")
.insert(record.run_id.clone(), record.clone());
state.record_snapshot(&record);
let work = crate::work_graph::new_shared_work_runtime(
crate::tools::todo::new_shared_todo_list(),
crate::tools::plan::new_shared_plan_state(),
);
work.register_operation(
"missing-controller-session",
OperationIntent::new(
"workflow:workflow_missing_controller",
"missing controller",
true,
"workflow",
"missing-controller-test",
),
)
.expect("register workflow");
work.reconcile_operation(
"missing-controller-session",
OperationOwnerSnapshot::new(
"workflow:workflow_missing_controller",
OwnerState::Running,
1,
1,
),
)
.expect("running workflow");
state.attach_lifecycle(
"workflow_missing_controller",
WorkflowWorkLifecycle {
work: work.clone(),
session_id: "missing-controller-session".to_string(),
external: "workflow:workflow_missing_controller".to_string(),
},
);
cancel_workflow(
json!({"run_id": "workflow_missing_controller"}),
state.clone(),
)
.await
.expect("controller-less cancel must still journal cancelled");
let record = state
.runs
.lock()
.expect("runs lock")
.get("workflow_missing_controller")
.cloned()
.expect("workflow owner");
assert_eq!(record.status, WorkflowRunStatus::Cancelled);
assert_eq!(record.lifecycle_seq, 2);
assert!(
record
.error
.as_deref()
.is_some_and(|error| error.contains("no live process")),
"expected an honest nothing-live receipt, got {:?}",
record.error
);
let operation = work
.capture(Some("missing-controller-session"))
.expect("capture")
.expect("graph")
.graph
.nodes
.into_iter()
.find(|node| node.kind == crate::work_graph::NodeKind::Operation)
.expect("workflow operation");
assert_eq!(operation.state, crate::work_graph::NodeState::Cancelled);
}
#[test]
fn handoff_compaction_preserves_release_sized_evidence() {
let payload = format!("APPROVE\n{}\nterminal: RunCompleted", "e".repeat(1_500));
assert_eq!(
compact_handoff_payload(&payload, WORKFLOW_HANDOFF_MAX_CHARS),
payload
);
}
#[test]
fn handoff_compaction_still_caps_oversized_artifacts() {
let payload = "e".repeat(WORKFLOW_HANDOFF_MAX_CHARS + 1);
let compacted = compact_handoff_payload(&payload, WORKFLOW_HANDOFF_MAX_CHARS);
assert_eq!(compacted.chars().count(), WORKFLOW_HANDOFF_MAX_CHARS + 3);
assert!(compacted.ends_with("..."));
}
#[test]
fn declarative_detection_matches_indented_and_nonleading_workflow_calls() {
assert!(looks_like_declarative_workflow("workflow({ tasks: [] })"));
assert!(looks_like_declarative_workflow(
"export default workflow({})"
));
assert!(looks_like_declarative_workflow(
"// build the run\n workflow({\n tasks: [],\n })"
));
assert!(!looks_like_declarative_workflow(
"return await parallel([() => task({ description: \"x\" })]);"
));
assert!(!looks_like_declarative_workflow("const x = myworkflow(1);"));
}
#[test]
fn workflow_action_defaults_to_start() {
assert_eq!(
parse_workflow_action(&json!({})).unwrap(),
WorkflowAction::Start
);
assert_eq!(
parse_workflow_action(&json!({"action": "run"})).unwrap(),
WorkflowAction::Run
);
}
#[test]
fn named_fleet_maps_workflow_role_to_profile_before_spawn() {
let fleet = FleetRoleMap::from_pairs([
("scout", "scout"),
("implementer", "builder"),
("reviewer", "reviewer"),
("verifier", "verifier"),
("release_lead", "manager"),
])
.expect("fleet");
let mut request = TaskRequest {
description: "fix it".to_string(),
subagent_type: None,
role: Some("implementer".to_string()),
profile: None,
model: None,
model_strength: None,
thinking: None,
cwd: None,
worktree: true,
write_authority: Some("worktree_write".to_string()),
write_roots: vec!["src".to_string()],
exact_files: Vec::new(),
coordination_contracts: vec!["test-contract".to_string()],
dependencies: Vec::new(),
acceptance: Vec::new(),
allowed_tools: None,
disallowed_tools: Vec::new(),
max_depth: None,
token_budget: None,
max_steps: None,
wall_time_secs: None,
response_schema: None,
label: Some("fix".to_string()),
phase: Some("implement".to_string()),
};
apply_named_fleet_to_task_request(Some(&fleet), &mut request).expect("resolve");
assert_eq!(request.role.as_deref(), Some("implementer"));
assert_eq!(request.profile.as_deref(), Some("builder"));
}
const EXACT_GLM_FLEET: &str = r#"
name = "glm-pair"
schema = "exact"
reasoning_router = "luna-low"
[[members]]
id = "implementer"
role = "builder"
provider = "zai"
model = "glm-5"
reasoning = "auto"
permissions = "read_write"
[[members]]
id = "auditor"
role = "reviewer"
provider = "zai"
model = "glm-5"
reasoning = "high"
permissions = "read_only"
"#;
fn exact_task_request(role: &str) -> TaskRequest {
TaskRequest {
description: "land the fix".to_string(),
subagent_type: None,
role: Some(role.to_string()),
profile: None,
model: None,
model_strength: None,
thinking: None,
cwd: None,
worktree: false,
write_authority: None,
write_roots: Vec::new(),
exact_files: Vec::new(),
coordination_contracts: Vec::new(),
dependencies: Vec::new(),
acceptance: Vec::new(),
allowed_tools: None,
disallowed_tools: Vec::new(),
max_depth: None,
token_budget: None,
max_steps: None,
wall_time_secs: None,
response_schema: None,
label: None,
phase: None,
}
}
fn exact_write_task_request(role: &str) -> TaskRequest {
TaskRequest {
write_roots: vec!["crates/tui".to_string()],
..exact_task_request(role)
}
}
fn exact_session() -> codewhale_workflow::PermissionCeiling {
codewhale_workflow::PermissionCeiling::preset("full").expect("preset")
}
fn exact_workflow_with(
text: &str,
router: Option<std::sync::Arc<crate::fleet::exact::StaticFleetRouter>>,
) -> crate::fleet::exact::ExactFleetWorkflow {
let document = codewhale_workflow::FleetDocument::parse(text).expect("exact fleet parses");
crate::fleet::exact::ExactFleetWorkflow::for_tests(
&document,
codewhale_workflow::QualifiedFleetId {
name: "glm-pair".to_string(),
origin: "workspace".to_string(),
},
router,
)
}
fn exact_workflow(text: &str) -> crate::fleet::exact::ExactFleetWorkflow {
exact_workflow_with(
text,
Some(crate::fleet::exact::StaticFleetRouter::new(
r#"{"reasoning":"max"}"#,
)),
)
}
#[tokio::test]
async fn exact_fleet_task_launch_resolves_the_member_route_and_ceiling() {
let operation = exact_workflow(EXACT_GLM_FLEET);
let mut request = exact_write_task_request("builder");
let binding = bind_exact_fleet_task_request(&operation, exact_session(), &mut request)
.expect("exact member resolves");
assert_eq!(request.profile.as_deref(), Some("implementer"));
assert_eq!(request.role.as_deref(), Some("builder"));
let member = operation.roster().get("implementer").expect("roster");
assert_eq!(member.profile.provider.as_deref(), Some("zai"));
assert_eq!(member.profile.model.as_deref(), Some("glm-5"));
assert_eq!(request.write_authority.as_deref(), Some("workspace_write"));
assert_eq!(request.max_depth, Some(0));
assert!(request.thinking.is_none(), "reasoning is not decided yet");
route_admitted_exact_task(&operation, &binding, &mut request)
.await
.expect("routing");
assert_eq!(request.thinking.as_deref(), Some("max"));
let mut auditor = exact_task_request("reviewer");
let auditor_binding =
bind_exact_fleet_task_request(&operation, exact_session(), &mut auditor)
.expect("auditor resolves");
route_admitted_exact_task(&operation, &auditor_binding, &mut auditor)
.await
.expect("routing");
assert_eq!(auditor.write_authority.as_deref(), Some("read_only"));
assert_eq!(auditor.thinking.as_deref(), Some("high"));
assert_eq!(auditor.role.as_deref(), Some("reviewer"));
assert_eq!(auditor.profile.as_deref(), Some("auditor"));
}
#[test]
fn workflow_tool_honors_the_session_workflow_config() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let mut runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx,
true,
None,
manager.clone(),
);
let input = json!({
"action": "start",
"plan": {
"goal": "write freely",
"risk": "writes",
"children": [{ "prompt": "edit", "type": "implementer" }]
}
});
let tool = WorkflowTool::new(Arc::clone(&manager), runtime.clone());
assert_eq!(
tool.approval_requirement_for(&input),
ApprovalRequirement::Required,
"product default: writes need approval"
);
let config = crate::config::Config {
workflow: Some(codewhale_config::WorkflowConfigToml {
require_approval_for_writes: false,
..Default::default()
}),
..Default::default()
};
runtime.api_config = Some(Arc::new(config));
let tool = WorkflowTool::new(Arc::clone(&manager), runtime.clone());
assert_eq!(
tool.approval_requirement_for(&input),
ApprovalRequirement::Auto,
"the session config must win"
);
let read_only = json!({
"action": "start",
"plan": {
"goal": "scout crates",
"risk": "read_only",
"children": [{ "prompt": "look", "type": "explore" }]
}
});
let config = crate::config::Config {
workflow: Some(codewhale_config::WorkflowConfigToml {
auto_start_read_only: false,
..Default::default()
}),
..Default::default()
};
runtime.api_config = Some(Arc::new(config));
let tool = WorkflowTool::new(manager, runtime);
assert_eq!(
tool.approval_requirement_for(&read_only),
ApprovalRequirement::Required,
"auto_start_read_only = false must still ask"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn a_role_keyed_gate_still_fires_when_the_member_id_differs() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_exact_gate".to_string();
let gates = vec![GateSpec {
id: "scout-findings".to_string(),
role: "scout".to_string(),
on: GateOn::RoleComplete,
gate: GateKind::Approve,
on_fail: codewhale_workflow::GateOnFail::Block,
blocks_role: Some("builder".to_string()),
max_retries: 0,
artifact_kind: Some("findings".to_string()),
require_explicit_verdict: false,
}];
state.runs.lock().expect("runs").insert(
run_id.clone(),
WorkflowRunRecord::new(run_id.clone(), None, None, None),
);
let driver = SubAgentWorkflowDriver::new(
run_id.clone(),
manager,
runtime,
state.clone(),
None,
WorkflowFleetBinding::None,
gates,
);
driver.evaluate_gates_for_completed_role(&RuntimeTaskRecord {
agent_id: "scout-agent".to_string(),
label: Some("scout".to_string()),
role: Some("scout".to_string()),
status: IrWorkflowRunStatus::Failed,
output: None,
schema_error: None,
usage: None,
});
let operation = exact_workflow(EXACT_GLM_FLEET);
let mut request = exact_write_task_request("builder");
bind_exact_fleet_task_request(&operation, exact_session(), &mut request).expect("bind");
assert_ne!(
request.role.as_deref(),
request.profile.as_deref(),
"this fleet's ids and roles differ, which is the whole point"
);
assert_eq!(request.role.as_deref(), Some("builder"));
let err = driver
.prepare_request_for_gates(&mut request)
.expect_err("a blocking gate on `builder` must still see `builder`");
assert!(err.to_string().contains("builder"), "{err}");
let mut by_id = exact_task_request("builder");
by_id.role = Some("implementer".to_string());
by_id.profile = None;
assert!(
driver.prepare_request_for_gates(&mut by_id).is_ok(),
"the profile id is not the semantic role the gate keys on"
);
}
#[test]
fn exact_fleet_rejects_a_conflicting_task_role_and_profile() {
let operation = exact_workflow(EXACT_GLM_FLEET);
let mut request = exact_task_request("reviewer");
request.profile = Some("implementer".to_string());
let err = bind_exact_fleet_task_request(&operation, exact_session(), &mut request)
.expect_err("conflicting identity");
let message = format!("{err:?}");
assert!(message.contains("different members"), "{message}");
}
#[test]
fn exact_fleet_rejects_task_level_route_overrides() {
let operation = exact_workflow(EXACT_GLM_FLEET);
for mutate in [
(|request: &mut TaskRequest| request.model = Some("glm-4".to_string()))
as fn(&mut TaskRequest),
|request: &mut TaskRequest| request.model_strength = Some("faster".to_string()),
|request: &mut TaskRequest| request.thinking = Some("off".to_string()),
] {
let mut request = exact_task_request("builder");
mutate(&mut request);
let err = bind_exact_fleet_task_request(&operation, exact_session(), &mut request)
.expect_err("an exact fleet member may not be re-routed per task");
let message = format!("{err:?}");
assert!(
message.contains("not allowed"),
"override must be rejected, not ignored: {message}"
);
}
}
#[test]
fn exact_fleet_rejects_task_level_posture_widening() {
let operation = exact_workflow(EXACT_GLM_FLEET);
for (field, mutate) in [
(
"subagent_type",
(|request: &mut TaskRequest| {
request.subagent_type = Some("general".to_string());
}) as fn(&mut TaskRequest),
),
("allowed_tools", |request: &mut TaskRequest| {
request.allowed_tools = Some(vec!["shell".to_string()]);
}),
("write_authority", |request: &mut TaskRequest| {
request.write_authority = Some("workspace_write".to_string());
}),
] {
let mut request = exact_task_request("reviewer");
mutate(&mut request);
let err = bind_exact_fleet_task_request(&operation, exact_session(), &mut request)
.expect_err("an exact ceiling must win over a task option");
let message = format!("{err:?}");
assert!(
message.contains(field) && message.contains("not allowed"),
"{field} must be rejected: {message}"
);
}
let mut clean = exact_task_request("reviewer");
bind_exact_fleet_task_request(&operation, exact_session(), &mut clean)
.expect("clean launch");
assert_eq!(clean.write_authority.as_deref(), Some("read_only"));
assert_eq!(clean.subagent_type, None);
}
#[test]
fn exact_fleet_ceilings_reach_the_spawn_request_as_a_tool_policy() {
let operation = exact_workflow(EXACT_GLM_FLEET);
let mut request = exact_task_request("reviewer");
bind_exact_fleet_task_request(&operation, exact_session(), &mut request).expect("bind");
assert_eq!(
request.allowed_tools, None,
"a tool-using member keeps full inheritance, narrowed by the deny list"
);
for denied in [
"web.run",
"web_search",
"fetch_url",
"wait_for_dev_server",
"mcp*",
] {
assert!(
request.disallowed_tools.iter().any(|name| name == denied),
"{denied} must be denied: {:?}",
request.disallowed_tools
);
}
assert!(
!request.disallowed_tools.iter().any(|name| name == "Web"),
"the Web family name must stay reachable so search/fetch survive: {:?}",
request.disallowed_tools
);
}
#[tokio::test]
async fn a_routing_receipt_survives_a_failed_spawn() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_orphaned_receipt".to_string();
state.runs.lock().expect("runs").insert(
run_id.clone(),
WorkflowRunRecord::new(run_id.clone(), None, None, None),
);
let driver = SubAgentWorkflowDriver::new(
run_id.clone(),
manager,
runtime,
state.clone(),
None,
WorkflowFleetBinding::None,
Vec::new(),
);
let operation = exact_workflow(EXACT_GLM_FLEET);
let mut request = exact_write_task_request("builder");
let binding =
bind_exact_fleet_task_request(&operation, exact_session(), &mut request).expect("bind");
let receipt = route_admitted_exact_task(&operation, &binding, &mut request)
.await
.expect("routing");
driver.record_orphaned_fleet_receipt(&receipt, "Sub-agent depth limit reached");
let events = state
.runs
.lock()
.expect("runs")
.get(&run_id)
.expect("run")
.events
.clone();
let logged = events
.iter()
.filter_map(|event| match &event.kind {
WorkflowUiEventKind::Log { message } => Some(message.clone()),
_ => None,
})
.find(|message| message.contains("spawn_failed=true"))
.expect("the receipt must outlive the failed spawn");
assert!(logged.contains("member=implementer"), "{logged}");
assert!(logged.contains("source=fleet_router"), "{logged}");
assert!(
logged.contains("reasoning_router:workspace/luna-low"),
"{logged}"
);
assert!(logged.contains("Sub-agent depth limit reached"), "{logged}");
assert!(!logged.contains("land the fix"), "{logged}");
}
#[tokio::test]
async fn a_started_event_shows_the_members_role_not_its_permission_posture() {
const AUDIT_FLEET: &str = r#"
name = "glm-pair"
schema = "exact"
[[members]]
id = "auditor"
role = "auditor"
provider = "zai"
model = "glm-5"
reasoning = "high"
permissions = "read_only"
"#;
let operation = exact_workflow_with(AUDIT_FLEET, None);
let mut request = exact_task_request("auditor");
let binding =
bind_exact_fleet_task_request(&operation, exact_session(), &mut request).expect("bind");
let receipt = route_admitted_exact_task(&operation, &binding, &mut request)
.await
.expect("routing");
let posture = operation
.roster()
.get("auditor")
.expect("roster entry")
.profile
.role
.name
.clone();
assert_eq!(posture, "scout");
assert_eq!(receipt.posture_role.as_deref(), Some("scout"));
assert_eq!(
displayed_resolved_role(Some(&receipt), Some(&posture), request.role.as_deref()),
Some("auditor".to_string()),
"the panel must not rename the operator's member to its posture"
);
assert_eq!(
displayed_resolved_role(None, Some("builder"), Some("reviewer")),
Some("builder".to_string())
);
assert_eq!(
displayed_resolved_role(None, None, Some("reviewer")),
Some("reviewer".to_string())
);
}
#[test]
fn a_predictably_invalid_write_scope_is_rejected_before_the_router_runs() {
let router = crate::fleet::exact::StaticFleetRouter::new(r#"{"reasoning":"max"}"#);
let operation = exact_workflow_with(EXACT_GLM_FLEET, Some(router.clone()));
let mut unbounded = exact_task_request("builder");
let err = bind_exact_fleet_task_request(&operation, exact_session(), &mut unbounded)
.expect_err("an unbounded write claim never reaches a spawn");
let message = format!("{err:?}");
assert!(message.contains("write_roots"), "{message}");
let mut scoped_read_only = exact_task_request("reviewer");
scoped_read_only.exact_files = vec!["crates/tui/src/main.rs".to_string()];
let err = bind_exact_fleet_task_request(&operation, exact_session(), &mut scoped_read_only)
.expect_err("a read-only member may not claim files");
assert!(format!("{err:?}").contains("read-only"), "{err:?}");
assert_eq!(
router.call_count(),
0,
"validation that the spawn will fail must precede the router call"
);
let mut bounded = exact_write_task_request("builder");
bind_exact_fleet_task_request(&operation, exact_session(), &mut bounded)
.expect("a bounded write claim is valid");
}
#[test]
fn a_read_only_session_narrows_a_write_capable_exact_member() {
let operation = exact_workflow(EXACT_GLM_FLEET);
let session = codewhale_workflow::PermissionCeiling {
write: false,
network_tool: false,
shell: codewhale_workflow::ShellCeiling::ReadOnly,
delegation_depth: 0,
tools: true,
};
let mut request = exact_task_request("builder");
bind_exact_fleet_task_request(&operation, session, &mut request).expect("bind");
assert_eq!(
request.write_authority.as_deref(),
Some("read_only"),
"a saved read_write member must not write inside a read-only session"
);
assert_eq!(request.max_depth, Some(0));
}
#[test]
fn a_rejected_or_unadmitted_task_spends_no_router_call() {
let router = crate::fleet::exact::StaticFleetRouter::new(r#"{"reasoning":"max"}"#);
let operation = exact_workflow_with(EXACT_GLM_FLEET, Some(router.clone()));
let mut overridden = exact_task_request("builder");
overridden.model = Some("glm-4".to_string());
assert!(
bind_exact_fleet_task_request(&operation, exact_session(), &mut overridden).is_err()
);
let mut unknown = exact_task_request("wizard");
assert!(bind_exact_fleet_task_request(&operation, exact_session(), &mut unknown).is_err());
let mut queued = exact_write_task_request("builder");
bind_exact_fleet_task_request(&operation, exact_session(), &mut queued).expect("bind");
assert_eq!(
router.call_count(),
0,
"binding must never contact the reasoning router"
);
}
#[tokio::test]
async fn an_exact_fleet_launch_produces_a_durable_routing_receipt() {
let operation = exact_workflow(EXACT_GLM_FLEET);
let mut request = exact_write_task_request("builder");
let binding =
bind_exact_fleet_task_request(&operation, exact_session(), &mut request).expect("bind");
let receipt = route_admitted_exact_task(&operation, &binding, &mut request)
.await
.expect("routing");
assert_eq!(receipt.fleet, "workspace/glm-pair");
assert_eq!(receipt.member_id, "implementer");
assert_eq!(receipt.member_role, "builder");
assert_eq!(receipt.provider, "zai");
assert_eq!(receipt.model, "glm-5");
assert_eq!(receipt.requested_reasoning, "auto");
assert_eq!(receipt.effective_reasoning, "max");
assert_eq!(receipt.selection_source, "fleet_router");
let router = receipt.router.as_ref().expect("router identity");
assert_eq!(router.service_kind, "reasoning_router");
assert_eq!(router.qualified(), "workspace/luna-low");
let call = router.call.as_ref().expect("call disclosure");
assert_eq!(call.requested, "low");
assert_eq!(call.effective, "low");
let event = WorkflowUiEvent::new(WorkflowUiEventKind::TaskStarted(Box::new(
WorkflowTaskStartedEvent {
task_id: "t1".to_string(),
label: None,
role: request.role.clone(),
profile: request.profile.clone(),
model: None,
strength: None,
thinking: request.thinking.clone(),
requested_reasoning: Some(receipt.requested_reasoning.clone()),
effective_reasoning: Some(receipt.effective_reasoning.clone()),
resolved_role: None,
resolved_profile: None,
resolved_provider: "zai".to_string(),
resolved_model: "glm-5".to_string(),
route_source: "fleet".to_string(),
child_route: None,
worktree: false,
workspace: None,
git_branch: None,
parent_task_id: None,
depth: 0,
workflow_run_id: None,
workflow_phase_id: None,
workflow_task_label: None,
workflow_child_index: None,
fleet_receipt: Some(receipt.clone()),
},
)));
let payload = serde_json::to_value(&event).expect("serialize");
let rendered = payload.to_string();
for expected in [
"fleet_receipt",
"\"selection_source\":\"fleet_router\"",
"\"provider_effective_reasoning\"",
"\"requested_reasoning\":\"auto\"",
"gpt-5.6-luna",
"\"service_kind\":\"reasoning_router\"",
] {
assert!(
rendered.contains(expected),
"{expected} missing: {rendered}"
);
}
for forbidden in ["/Users/", "/home/", ".toml", "api_key", "land the fix"] {
assert!(!rendered.contains(forbidden), "{forbidden} in {rendered}");
}
let line = receipt.line();
for expected in [
"requested=auto",
"effective=max",
"source=fleet_router",
"reasoning_router:workspace/luna-low",
"router_call_requested=low",
] {
assert!(line.contains(expected), "{expected} missing from {line}");
}
}
#[test]
fn exact_fleet_rejects_an_unknown_member() {
let operation = exact_workflow(EXACT_GLM_FLEET);
let mut request = exact_task_request("wizard");
let err = bind_exact_fleet_task_request(&operation, exact_session(), &mut request)
.expect_err("unknown member");
let message = format!("{err:?}");
assert!(message.contains("wizard"), "{message}");
assert!(message.contains("implementer"), "{message}");
let mut router_request = exact_task_request("luna-low");
assert!(
bind_exact_fleet_task_request(&operation, exact_session(), &mut router_request)
.is_err(),
"the reasoning router must not be launchable as a worker"
);
}
#[tokio::test]
async fn one_reasoning_router_profile_serves_two_fleets() {
let first = exact_workflow(EXACT_GLM_FLEET);
let second = {
let text = EXACT_GLM_FLEET.replace("name = \"glm-pair\"", "name = \"glm-solo\"");
let document =
codewhale_workflow::FleetDocument::parse(&text).expect("second fleet parses");
crate::fleet::exact::ExactFleetWorkflow::for_tests(
&document,
codewhale_workflow::QualifiedFleetId {
name: "glm-solo".to_string(),
origin: "workspace".to_string(),
},
Some(crate::fleet::exact::StaticFleetRouter::new(
r#"{"reasoning":"low"}"#,
)),
)
};
assert_eq!(
first.snapshot().router(),
second.snapshot().router(),
"both fleets reference the identical captured router service"
);
assert_ne!(
first.snapshot().fleet().qualified(),
second.snapshot().fleet().qualified()
);
for (operation, expected) in [(&first, "max"), (&second, "low")] {
let mut request = exact_write_task_request("builder");
let binding = bind_exact_fleet_task_request(operation, exact_session(), &mut request)
.expect("bind");
let receipt = route_admitted_exact_task(operation, &binding, &mut request)
.await
.expect("routing");
assert_eq!(receipt.effective_reasoning, expected);
assert_eq!(
receipt.router.as_ref().expect("router").qualified(),
"workspace/luna-low"
);
}
}
#[tokio::test]
async fn editing_the_fleet_file_after_start_does_not_move_a_running_route() {
let tmp = tempfile::tempdir().expect("tmp");
let fleets = tmp.path().join("fleets");
std::fs::create_dir_all(&fleets).expect("fleets dir");
let path = fleets.join("glm-pair.toml");
std::fs::write(&path, EXACT_GLM_FLEET).expect("write fleet");
let roots = vec![codewhale_workflow::FleetSearchRoot::new(
"workspace",
tmp.path(),
)];
let (document, id) =
codewhale_workflow::FleetDocument::load_by_name("glm-pair", &roots).expect("load");
let operation = crate::fleet::exact::ExactFleetWorkflow::for_tests(
&document,
id,
Some(crate::fleet::exact::StaticFleetRouter::new(
r#"{"reasoning":"max"}"#,
)),
);
let started_hash = operation.snapshot().content_hash().to_string();
std::fs::write(
&path,
EXACT_GLM_FLEET
.replace(
"model = \"glm-5\"\nreasoning = \"auto\"",
"model = \"glm-4\"\nreasoning = \"off\"",
)
.replace("permissions = \"read_write\"", "permissions = \"full\""),
)
.expect("rewrite fleet");
let mut request = exact_write_task_request("builder");
let binding = bind_exact_fleet_task_request(&operation, exact_session(), &mut request)
.expect("bind after the edit");
route_admitted_exact_task(&operation, &binding, &mut request)
.await
.expect("launch after the edit");
let member = operation.roster().get("implementer").expect("roster");
assert_eq!(
member.profile.model.as_deref(),
Some("glm-5"),
"the in-flight snapshot must keep the model it started with"
);
assert_eq!(request.thinking.as_deref(), Some("max"));
assert_eq!(
request.write_authority.as_deref(),
Some("workspace_write"),
"the edit must not widen the running ceiling to `full`"
);
assert_eq!(operation.snapshot().content_hash(), started_hash);
let (reloaded, reloaded_id) =
codewhale_workflow::FleetDocument::load_by_name("glm-pair", &roots).expect("reload");
let next = crate::fleet::exact::ExactFleetWorkflow::for_tests(
&reloaded,
reloaded_id,
Some(crate::fleet::exact::StaticFleetRouter::new(
r#"{"reasoning":"max"}"#,
)),
);
assert_eq!(
next.roster()
.get("implementer")
.expect("roster")
.profile
.model
.as_deref(),
Some("glm-4")
);
assert_ne!(next.snapshot().content_hash(), started_hash);
}
#[test]
fn named_fleet_rejects_unknown_workflow_role_before_spawn() {
let fleet = FleetRoleMap::from_pairs([("scout", "scout")]).expect("fleet");
let mut request = TaskRequest {
description: "fix it".to_string(),
subagent_type: None,
role: Some("wizard".to_string()),
profile: None,
model: None,
model_strength: None,
thinking: None,
cwd: None,
worktree: false,
write_authority: Some("read_only".to_string()),
write_roots: Vec::new(),
exact_files: Vec::new(),
coordination_contracts: Vec::new(),
dependencies: Vec::new(),
acceptance: Vec::new(),
allowed_tools: None,
disallowed_tools: Vec::new(),
max_depth: None,
token_budget: None,
max_steps: None,
wall_time_secs: None,
response_schema: None,
label: None,
phase: None,
};
let err = apply_named_fleet_to_task_request(Some(&fleet), &mut request)
.expect_err("unknown role should fail");
assert!(
err.to_string().contains("unknown fleet role `wizard`"),
"{err}"
);
}
#[test]
fn declarative_leaf_budget_reaches_task_runtime_options() {
let source = r#"
workflow({
"goal": "bound one child",
"nodes": [{
"agent": {
"id": "bounded",
"prompt": "Inspect bounded evidence.",
"budget": { "max_tokens": 5000, "max_steps": 4, "timeout_secs": 90 }
}
}]
});
"#;
let adapted = adapt_workflow_source(source, None).expect("lower bounded leaf");
assert!(
adapted.source.contains("tokenBudget: 5000"),
"{}",
adapted.source
);
assert!(adapted.source.contains("maxSteps: 4"), "{}", adapted.source);
assert!(
adapted.source.contains("wallTimeSecs: 90"),
"{}",
adapted.source
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn declarative_max_steps_zero_stops_before_provider_call() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("must not be called").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let tool = WorkflowTool::new(manager, runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": r#"
workflow({
"goal": "prove the child step cap reaches runtime",
"nodes": [{
"agent": {
"id": "zero-step",
"prompt": "Do not start a model turn.",
"budget": { "max_steps": 0, "timeout_secs": 90 }
}
}]
});
"#
}),
&ctx,
)
.await
.expect("failed workflow still returns its terminal receipt");
let payload: Value = serde_json::from_str(&result.content).expect("workflow JSON");
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"provider must not be called"
);
assert_eq!(payload["status"], "failed", "{payload}");
assert_eq!(
payload["execution"]["leaf_results"][0]["status"], "failed",
"{payload}"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn role_only_leaf_omits_type_and_resolves_through_named_fleet() {
let _retry_guard = workflow_test_retry_guard();
let _env_lock = crate::test_support::lock_test_env();
let tmp = tempfile::tempdir().expect("tempdir");
let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", tmp.path());
let fleet_dir = tmp.path().join("fleets");
std::fs::create_dir_all(&fleet_dir).expect("fleet dir");
std::fs::write(
fleet_dir.join("role-only-test.toml"),
r#"
name = "role-only-test"
[roles]
scout = "scout"
reviewer = "reviewer"
"#,
)
.expect("role-only fleet");
let source = r#"
export default workflow({
"goal": "resolve a role-only child",
"nodes": [
{
"agent": {
"id": "scout-source",
"prompt": "Inspect the source without editing.",
"role": "scout",
"mode": "read_only"
}
}
]
});
"#;
let adapted = adapt_workflow_source(source, None).expect("lower role-only workflow");
assert!(adapted.source.contains("role: \"scout\""));
assert!(
!adapted.source.contains("type:"),
"Fleet-addressed leaves must defer runtime type to the roster:\n{}",
adapted.source
);
let non_role = adapt_workflow_source(
r#"workflow({
"goal": "default non-role child",
"nodes": [{ "agent": { "id": "audit", "prompt": "Audit only." } }]
});"#,
None,
)
.expect("lower non-role workflow");
assert!(
non_role.source.contains("type: \"review\""),
"non-role read-only leaves retain the review default:\n{}",
non_role.source
);
let explicit_type_source = r#"
workflow({
"goal": "preserve an authored role type",
"nodes": [{
"agent": {
"id": "review-source",
"prompt": "Review the source without editing.",
"agent_type": "review",
"role": "reviewer",
"mode": "read_only"
}
}]
});
"#;
let explicit_type = adapt_workflow_source(explicit_type_source, None)
.expect("lower explicitly typed Fleet role");
assert!(
explicit_type.source.contains("type: \"review\""),
"an authored non-General type must remain a validated override:\n{}",
explicit_type.source
);
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("scout evidence").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": source,
"fleet": "role-only-test"
}),
&ctx,
)
.await
.expect("role-only workflow should resolve through its named Fleet");
let payload: Value = serde_json::from_str(&result.content).expect("workflow JSON");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(calls.load(Ordering::SeqCst), 1);
let started = payload["events"]
.as_array()
.expect("typed events")
.iter()
.find(|event| event["type"] == "task_started")
.expect("task_started receipt");
assert_eq!(started["role"], "scout");
assert_eq!(started["profile"], "scout");
assert_eq!(started["resolved_profile"], "scout");
let explicit_result = tool
.execute(
json!({
"action": "run",
"script": explicit_type_source,
"fleet": "role-only-test"
}),
&ctx,
)
.await
.expect("matching explicit role type should remain valid");
let explicit_payload: Value =
serde_json::from_str(&explicit_result.content).expect("workflow JSON");
assert_eq!(
explicit_payload["status"], "completed",
"{explicit_payload}"
);
assert_eq!(calls.load(Ordering::SeqCst), 2);
let conflicting_result = tool
.execute(
json!({
"action": "run",
"script": r#"workflow({
"goal": "reject a conflicting authored type",
"nodes": [{ "agent": {
"id": "bad-scout",
"prompt": "Review as a scout.",
"agent_type": "review",
"role": "scout",
"mode": "read_only"
} }]
});"#,
"fleet": "role-only-test"
}),
&ctx,
)
.await
.expect("conflicting type returns a terminal workflow record");
let conflicting_payload: Value =
serde_json::from_str(&conflicting_result.content).expect("workflow JSON");
assert_eq!(
conflicting_payload["status"], "failed",
"{conflicting_payload}"
);
assert!(
conflicting_payload["error"]
.as_str()
.is_some_and(|error| error
.contains("Fleet role conflicts with the explicit legacy agent type")),
"{conflicting_payload}"
);
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"conflicting explicit type must fail before the provider"
);
}
#[test]
fn parallel_write_children_default_to_worktree_isolation() {
let source = r#"
export default workflow({
"goal": "parallel write isolation default",
"nodes": [
{
"branch": {
"id": "implement",
"parallel": true,
"children": [
{
"agent": {
"id": "left",
"prompt": "Patch left lane",
"agent_type": "implementer",
"mode": "read_write",
"file_scope": ["src/left.rs"]
}
},
{
"agent": {
"id": "right",
"prompt": "Patch right lane",
"agent_type": "implementer",
"mode": "read_write",
"file_scope": ["src/right.rs"]
}
}
]
}
}
]
});
"#;
let adapted = adapt_workflow_source(source, None).expect("lower parallel write workflow");
let spec = adapted.spec.expect("declarative spec");
let WorkflowNode::BranchSet(branch) = &spec.nodes[0] else {
panic!("expected branch_set");
};
assert!(branch.parallel);
for child in &branch.children {
let WorkflowNode::Leaf(leaf) = child else {
panic!("expected leaf");
};
assert!(leaf_is_write_capable(leaf));
assert!(
leaf_wants_worktree(leaf, true),
"parallel write leaf {} should default to worktree",
leaf.id
);
assert_eq!(leaf.isolation, IsolationMode::Auto);
}
assert!(
adapted.source.contains("worktree: true"),
"lowered JS should request worktree isolation:\n{}",
adapted.source
);
assert_eq!(
adapted.source.matches("worktree: true").count(),
2,
"each parallel write child should get worktree: true:\n{}",
adapted.source
);
assert_eq!(
adapted
.source
.matches("writeAuthority: \"worktree_write\"")
.count(),
2,
"each isolated writer should carry enforced worktree authority:\n{}",
adapted.source
);
assert!(adapted.source.contains("writeRoots: [\"src/left.rs\"]"));
assert!(adapted.source.contains("writeRoots: [\"src/right.rs\"]"));
}
#[test]
fn parallel_write_same_worktree_requires_explicit_shared_isolation() {
let source = r#"
export default workflow({
"goal": "parallel write same-worktree override",
"nodes": [
{
"branch": {
"id": "implement",
"parallel": true,
"children": [
{
"agent": {
"id": "shared-writer",
"prompt": "Patch in the parent checkout",
"agent_type": "implementer",
"mode": "read_write",
"isolation": "shared",
"file_scope": ["src/shared.rs"]
}
},
{
"agent": {
"id": "isolated-writer",
"prompt": "Patch in a worktree",
"agent_type": "implementer",
"mode": "read_write",
"isolation": "worktree",
"file_scope": ["src/isolated.rs"]
}
}
]
}
}
]
});
"#;
let adapted =
adapt_workflow_source(source, None).expect("lower same-worktree override workflow");
let spec = adapted.spec.expect("declarative spec");
let WorkflowNode::BranchSet(branch) = &spec.nodes[0] else {
panic!("expected branch_set");
};
let leaves: Vec<&LeafSpec> = branch
.children
.iter()
.map(|child| match child {
WorkflowNode::Leaf(leaf) => leaf,
_ => panic!("expected leaf"),
})
.collect();
assert_eq!(leaves[0].isolation, IsolationMode::Shared);
assert!(
!leaf_wants_worktree(leaves[0], true),
"explicit shared should keep same-worktree"
);
assert_eq!(leaves[1].isolation, IsolationMode::Worktree);
assert!(leaf_wants_worktree(leaves[1], true));
assert_eq!(
adapted.source.matches("worktree: true").count(),
1,
"same-worktree override must not force worktree on shared leaf:\n{}",
adapted.source
);
assert!(
adapted.source.contains("shared-writer") && adapted.source.contains("isolated-writer"),
"both children should still be lowered:\n{}",
adapted.source
);
assert!(
adapted
.source
.contains("writeAuthority: \"workspace_write\"")
);
assert!(
adapted
.source
.contains("writeAuthority: \"worktree_write\"")
);
}
#[test]
fn parallel_read_only_children_do_not_default_to_worktree() {
let source = r#"
export default workflow({
"goal": "parallel read-only stays shared",
"nodes": [
{
"branch": {
"id": "audit",
"parallel": true,
"children": [
{
"agent": {
"id": "review-a",
"prompt": "Review A",
"agent_type": "review",
"mode": "read_only"
}
},
{
"agent": {
"id": "review-b",
"prompt": "Review B",
"agent_type": "verifier",
"mode": "read_only"
}
}
]
}
}
]
});
"#;
let adapted = adapt_workflow_source(source, None).expect("lower parallel read-only");
assert!(
!adapted.source.contains("worktree: true"),
"read-only parallel children should not get worktree isolation:\n{}",
adapted.source
);
assert_eq!(
adapted
.source
.matches("writeAuthority: \"read_only\"")
.count(),
2,
"read-only mode must reach the task authority contract:\n{}",
adapted.source
);
}
#[test]
fn sequential_write_children_do_not_default_to_worktree() {
let source = r#"
export default workflow({
"goal": "sequential write stays shared by default",
"nodes": [
{
"sequence": {
"id": "implement",
"children": [
{
"agent": {
"id": "writer",
"prompt": "Patch sequentially",
"agent_type": "implementer",
"mode": "read_write",
"file_scope": ["src/main.rs"]
}
}
]
}
}
]
});
"#;
let adapted = adapt_workflow_source(source, None).expect("lower sequential write");
assert!(
!adapted.source.contains("worktree: true"),
"sequential writes should not default to worktree:\n{}",
adapted.source
);
assert!(
adapted
.source
.contains("writeAuthority: \"workspace_write\"")
);
assert!(adapted.source.contains("writeRoots: [\"src/main.rs\"]"));
}
#[test]
fn write_scope_suffix_globs_lower_to_enforceable_roots() {
let source = r#"workflow({
"goal": "bounded auth patch",
"nodes": [{ "agent": {
"id": "writer",
"prompt": "Patch auth",
"agent_type": "implementer",
"mode": "read_write",
"file_scope": ["./src/auth/**"]
}}]
});"#;
let adapted = adapt_workflow_source(source, None).expect("lower trailing glob scope");
assert!(
adapted.source.contains("writeRoots: [\"src/auth\"]"),
"runtime claim must contain src/auth/login.rs:\n{}",
adapted.source
);
let unsupported = r#"workflow({
"goal": "reject ambiguous glob",
"nodes": [{ "agent": {
"id": "writer",
"prompt": "Patch auth",
"agent_type": "implementer",
"mode": "read_write",
"file_scope": ["src/*/auth"]
}}]
});"#;
let error = match adapt_workflow_source(unsupported, None) {
Ok(_) => panic!("internal globs cannot become literal runtime roots"),
Err(error) => error.to_string(),
};
assert!(error.contains("unsupported file_scope"), "{error}");
}
#[test]
fn write_leaves_require_scope_and_reduce_stays_read_only() {
let unscoped_writer = r#"workflow({
"goal": "reject an unbounded writer",
"nodes": [{ "agent": {
"id": "writer",
"prompt": "Patch it",
"agent_type": "implementer",
"mode": "read_write"
}}]
});"#;
let error = match adapt_workflow_source(unscoped_writer, None) {
Ok(_) => panic!("write-capable leaves must declare file_scope"),
Err(error) => error.to_string(),
};
assert!(error.contains("declares no file_scope"), "{error}");
let reduce = r#"workflow({
"goal": "reduce read-only evidence",
"nodes": [{ "reduce": {
"id": "summary",
"inputs": [],
"prompt": "Summarize the evidence"
}}]
});"#;
let adapted = adapt_workflow_source(reduce, None).expect("lower reduce");
assert!(adapted.source.contains("type: \"plan\""));
assert!(adapted.source.contains("writeAuthority: \"read_only\""));
}
#[test]
fn inline_script_and_source_path_are_mutually_exclusive() {
let ctx = ToolContext::new(".");
let err = workflow_source(
&json!({
"script": "return 1;",
"source_path": "workflow.js"
}),
&ctx,
)
.unwrap_err();
assert!(
err.to_string()
.contains("exactly one of script, source_path, or plan"),
"{err}"
);
}
#[test]
fn structured_plan_lowers_to_parallel_not_promise_all() {
let ctx = ToolContext::new(".");
let source = workflow_source(
&json!({
"plan": {
"goal": "audit two independent scopes",
"risk": "read_only",
"max_children": 8,
"token_budget": 120000,
"phases": [{
"id": "scout",
"title": "Scout",
"children": [
{
"id": "left",
"label": "left-lane",
"prompt": "Inspect crates/left",
"type": "explore"
},
{
"id": "right",
"prompt": "Inspect crates/right",
"type": "explore"
}
]
}]
}
}),
&ctx,
)
.expect("structured plan should lower");
assert!(
source.source.contains("await parallel(["),
"lowered JS must use parallel():\n{}",
source.source
);
assert!(
!source.source.contains("Promise.all"),
"lowered JS must not use raw Promise.all:\n{}",
source.source
);
assert!(
source.source.contains("() => task("),
"parallel slots should be thunks:\n{}",
source.source
);
let spec = source.spec.expect("plan should produce WorkflowSpec");
assert_eq!(spec.goal, "audit two independent scopes");
assert_eq!(spec.budget.max_tokens, Some(120000));
assert_eq!(spec.nodes.len(), 1);
let WorkflowNode::BranchSet(branch) = &spec.nodes[0] else {
panic!("expected parallel branch for multi-child phase");
};
assert!(branch.parallel);
assert_eq!(branch.children.len(), 2);
}
#[test]
fn structured_plan_validation_errors_are_typed() {
let ctx = ToolContext::new(".");
let missing_goal = workflow_source(
&json!({
"plan": {
"goal": " ",
"children": [{ "prompt": "do work" }]
}
}),
&ctx,
)
.unwrap_err();
assert!(missing_goal.to_string().contains("goal"), "{missing_goal}");
let over_limit = workflow_source(
&json!({
"plan": {
"goal": "too many children",
"max_children": 1,
"children": [
{ "id": "a", "prompt": "one" },
{ "id": "b", "prompt": "two" }
]
}
}),
&ctx,
)
.unwrap_err();
assert!(
over_limit.to_string().contains("max_children"),
"{over_limit}"
);
let bad_type = workflow_source(
&json!({
"plan": {
"goal": "bad type",
"children": [{ "prompt": "x", "type": "wizard" }]
}
}),
&ctx,
)
.unwrap_err();
assert!(
bad_type
.to_string()
.contains("Invalid sub-agent type 'wizard'"),
"{bad_type}"
);
let exclusive = workflow_source(
&json!({
"script": "return 1;",
"plan": { "goal": "x", "children": [{ "prompt": "y" }] }
}),
&ctx,
)
.unwrap_err();
assert!(
exclusive
.to_string()
.contains("exactly one of script, source_path, or plan"),
"{exclusive}"
);
}
#[test]
fn plan_child_type_shares_the_agent_tool_option_vocabulary() {
for (alias, expected) in [
("worker", AgentType::General),
("delegate", AgentType::General),
("scout", AgentType::Explore),
("Explorer", AgentType::Explore),
("planner", AgentType::Plan),
("awaiter", AgentType::Plan),
("reviewer", AgentType::Review),
("consultant", AgentType::Review),
("oracle", AgentType::Review),
("advisor", AgentType::Review),
("builder", AgentType::Implementer),
("verifier", AgentType::Verifier),
] {
assert_eq!(
parse_plan_agent_type(Some(alias))
.unwrap_or_else(|err| panic!("{alias} rejected: {err}")),
expected,
"{alias}"
);
}
let typo = parse_plan_agent_type(Some("wizard"))
.unwrap_err()
.to_string();
assert!(typo.contains("Invalid sub-agent type 'wizard'"), "{typo}");
assert!(
typo.contains("worker, scout, planner, reviewer, builder"),
"{typo}"
);
assert!(
typo.contains("consultant/oracle/advisor"),
"accepted advisory aliases must remain visible in the guidance: {typo}"
);
let custom = parse_plan_agent_type(Some("custom"))
.unwrap_err()
.to_string();
assert!(
custom.contains("Invalid sub-agent type 'custom'"),
"{custom}"
);
assert!(custom.contains("allowed_tools"), "{custom}");
}
#[test]
fn declarative_parallel_branch_uses_parallel_helper() {
let source = r#"
export default workflow({
"goal": "partial success fan-out",
"nodes": [
{
"branch": {
"id": "fan",
"parallel": true,
"children": [
{ "agent": { "id": "a", "prompt": "A", "agent_type": "explore", "mode": "read_only" } },
{ "agent": { "id": "b", "prompt": "B", "agent_type": "explore", "mode": "read_only" } }
]
}
}
]
});
"#;
let adapted = adapt_workflow_source(source, None).expect("lower declarative");
assert!(
adapted.source.contains("await parallel(["),
"declarative parallel must lower via parallel():\n{}",
adapted.source
);
assert!(
!adapted.source.contains("Promise.all"),
"must not emit raw Promise.all:\n{}",
adapted.source
);
}
#[test]
fn source_path_must_stay_inside_workspace_without_trust_mode() {
let workspace = tempfile::tempdir().expect("workspace tempdir");
let outside = tempfile::tempdir().expect("outside tempdir");
let outside_path = outside.path().join("outside.workflow.js");
std::fs::write(&outside_path, "return 1;").expect("outside workflow source");
let ctx = ToolContext::new(workspace.path().to_path_buf());
let err = workflow_source(
&json!({
"source_path": outside_path
}),
&ctx,
)
.expect_err("outside source_path should be denied");
assert!(
err.to_string().contains("must stay inside the workspace"),
"{err}"
);
}
#[test]
fn subagent_tool_surface_registers_workflow_and_agent() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let registry = ToolRegistryBuilder::new()
.with_subagent_tools(manager, runtime)
.build(ctx);
assert!(registry.contains("workflow"));
assert!(registry.contains("agent"));
assert!(registry.contains("agents/list"));
assert!(registry.contains("agents/message"));
assert!(registry.contains("agents/followup"));
assert!(registry.contains("agents/interrupt"));
assert!(registry.contains("agents/wait"));
assert!(
registry
.to_api_tools()
.iter()
.any(|tool| tool.name == "workflow")
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn workflow_run_dispatches_task_through_subagent_manager() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("child done").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let tool = WorkflowTool::new(manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": "phase('dispatch'); log('starting child'); const out = await task({ description: 'say done', allowedTools: [], label: 'inspect-child', model: 'deepseek-v4-flash', modelStrength: 'same', thinking: 'low' }); return { out };"
}),
&ctx,
)
.await
.expect("workflow run should complete");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(payload["result"]["out"], "child done");
assert_eq!(payload["child_ids"].as_array().unwrap().len(), 1);
assert_eq!(calls.load(Ordering::SeqCst), 1);
let child_id = payload["child_ids"][0].as_str().unwrap();
let events = payload["events"].as_array().expect("events array");
assert!(
events
.iter()
.any(|event| event["type"] == "phase_started" && event["title"] == "dispatch"),
"{events:#?}"
);
assert!(
events
.iter()
.any(|event| event["type"] == "log" && event["message"] == "starting child"),
"{events:#?}"
);
assert!(
events.iter().any(|event| event["type"] == "budget_updated"),
"{events:#?}"
);
let task_started = events
.iter()
.find(|event| event["type"] == "task_started")
.expect("task_started event");
assert_eq!(task_started["task_id"], child_id);
assert_eq!(task_started["label"], "inspect-child");
assert!(task_started["profile"].is_null());
assert_eq!(task_started["model"], "deepseek-v4-flash");
assert_eq!(task_started["strength"], "same");
assert_eq!(task_started["thinking"], "low");
assert_eq!(task_started["requested_reasoning"], "low");
assert!(
task_started["effective_reasoning"]
.as_str()
.is_some_and(|value| !value.is_empty()),
"{task_started}"
);
assert_eq!(task_started["resolved_provider"], "deepseek");
assert_eq!(task_started["resolved_model"], "deepseek-v4-flash");
assert_eq!(task_started["route_source"], "task.model");
assert_eq!(task_started["worktree"], false);
assert!(task_started["parent_task_id"].is_null());
assert_eq!(task_started["depth"], 1);
assert_eq!(
task_started["workflow_run_id"].as_str(),
payload["run_id"].as_str()
);
assert_eq!(task_started["workflow_phase_id"], "dispatch");
assert_eq!(task_started["workflow_task_label"], "inspect-child");
assert_eq!(task_started["workflow_child_index"], 0);
assert!(
events.iter().any(|event| event["type"] == "task_completed"
&& event["task_id"] == child_id
&& event["status"] == "succeeded"),
"{events:#?}"
);
let child = manager
.read()
.await
.get_result(child_id)
.expect("child result");
assert_eq!(child.status, SubAgentStatus::Completed);
assert_eq!(child.result.as_deref(), Some("child done"));
let reloaded = WorkflowWorkspaceState::open(tmp.path());
let persisted = reloaded
.runs
.lock()
.expect("reloaded workflow runs")
.get(payload["run_id"].as_str().expect("run id"))
.cloned()
.expect("persisted run");
let persisted_json = serde_json::to_value(&persisted).expect("persisted run JSON");
let persisted_started = persisted_json["events"]
.as_array()
.and_then(|events| events.iter().find(|event| event["type"] == "task_started"))
.expect("persisted task_started receipt");
assert_eq!(persisted_started["requested_reasoning"], "low");
assert_eq!(
persisted_started["effective_reasoning"],
task_started["effective_reasoning"]
);
assert_eq!(persisted_started["resolved_provider"], "deepseek");
assert_eq!(persisted_started["resolved_model"], "deepseek-v4-flash");
assert_eq!(persisted_started["route_source"], "task.model");
let mut panel =
crate::tui::widgets::workflow_panel::WorkflowPanel::from_run_json(&persisted_json)
.expect("journal should hydrate workflow panel");
let original_receipt = panel
.phases
.iter()
.flat_map(|phase| phase.rows.iter())
.find(|row| row.task_id == child_id)
.map(crate::tui::widgets::workflow_panel::row_receipt_text)
.expect("spawned child receipt");
assert!(original_receipt.contains("deepseek/deepseek-v4-flash"));
assert!(original_receipt.contains("reasoning low→"));
assert!(original_receipt.contains("via task.model"));
panel.apply_json_event(&json!({
"type": "task_started",
"at_ms": 9_000,
"task_id": "later-route",
"label": "later-route",
"resolved_role": "consultant",
"resolved_provider": "moonshot",
"resolved_model": "kimi-k3",
"requested_reasoning": "auto",
"effective_reasoning": "medium",
"route_source": "agent_profile.model",
"worktree": false,
}));
panel.apply_json_event(&json!({
"type": "task_completed",
"at_ms": 9_100,
"task_id": "later-route",
"status": "succeeded"
}));
let unchanged = panel
.phases
.iter()
.flat_map(|phase| phase.rows.iter())
.find(|row| row.task_id == child_id)
.map(crate::tui::widgets::workflow_panel::row_receipt_text)
.expect("original child after later route");
assert_eq!(unchanged, original_receipt);
let history =
crate::tui::widgets::workflow_panel::WorkflowPanel::from_run_json(&panel.to_run_json())
.expect("history receipt round trip");
let later_receipt = history
.phases
.iter()
.flat_map(|phase| phase.rows.iter())
.find(|row| row.task_id == "later-route")
.map(crate::tui::widgets::workflow_panel::row_receipt_text)
.expect("later route history receipt");
assert!(later_receipt.contains("moonshot/kimi-k3"));
assert!(later_receipt.contains("reasoning auto→medium"));
assert!(later_receipt.contains("tokens unknown"));
assert!(!later_receipt.contains("tokens 0"));
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn named_fleet_run_emits_role_resolved_receipt_and_rejects_unknown_before_provider() {
let _retry_guard = workflow_test_retry_guard();
let _env_lock = crate::test_support::lock_test_env();
let tmp = tempfile::tempdir().expect("tempdir");
let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", tmp.path());
std::fs::create_dir_all(tmp.path().join("fleets")).expect("fleets dir");
std::fs::write(
tmp.path().join("fleets/offline.toml"),
r#"
name = "offline"
[roles]
reviewer = "reviewer"
"#,
)
.expect("named fleet fixture");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("role-resolved child").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let tool = WorkflowTool::new(manager, runtime);
let completed = tool
.execute(
json!({
"action": "run",
"fleet": "offline",
"script": "return await task({ description: 'review it', type: 'review', role: 'reviewer', label: 'offline-review' });"
}),
&ctx,
)
.await
.expect("named fleet workflow");
let payload: Value = serde_json::from_str(&completed.content).expect("workflow JSON");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(payload["result"], "role-resolved child");
assert_eq!(calls.load(Ordering::SeqCst), 1);
let started = payload["events"]
.as_array()
.and_then(|events| events.iter().find(|event| event["type"] == "task_started"))
.expect("task_started receipt");
assert_eq!(started["role"], "reviewer");
assert_eq!(started["profile"], "reviewer");
assert_eq!(started["resolved_role"], "reviewer");
assert_eq!(started["resolved_profile"], "reviewer");
assert_eq!(started["resolved_provider"], "deepseek");
assert_eq!(started["resolved_model"], "deepseek-v4-flash");
assert_eq!(started["route_source"], "run.model");
assert!(
payload["events"]
.as_array()
.is_some_and(|events| events.iter().any(|event| event["type"] == "task_completed"))
);
let rejected = tool
.execute(
json!({
"action": "run",
"fleet": "offline",
"script": "return await task({ description: 'must not launch', type: 'review', role: 'wizard' });"
}),
&ctx,
)
.await
.expect("rejected workflow still returns its terminal record");
let rejected: Value = serde_json::from_str(&rejected.content).expect("rejected JSON");
assert_eq!(rejected["status"], "failed", "{rejected}");
assert!(
rejected["error"]
.as_str()
.is_some_and(|error| error.contains("unknown fleet role `wizard`")),
"{rejected}"
);
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"unknown role must fail before a second provider call"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn workflow_spawn_records_carry_child_index_and_phase_metadata() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
let (client, calls) = fake_chat_client("ok").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let tool = WorkflowTool::new(manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": "phase('alpha'); await task({ description: 'first', type: 'explore', allowedTools: [], label: 'one' }); phase('beta'); await task({ description: 'second', type: 'explore', allowedTools: [], label: 'two', phase: 'beta-explicit' }); return { ok: true };"
}),
&ctx,
)
.await
.expect("workflow run should complete");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(payload["child_ids"].as_array().unwrap().len(), 2);
assert_eq!(calls.load(Ordering::SeqCst), 2);
let mut started: Vec<&Value> = payload["events"]
.as_array()
.expect("events")
.iter()
.filter(|event| event["type"] == "task_started")
.collect();
started.sort_by_key(|event| event["workflow_child_index"].as_u64().unwrap_or(u64::MAX));
assert_eq!(started.len(), 2, "{started:#?}");
assert_eq!(started[0]["workflow_run_id"], payload["run_id"]);
assert_eq!(started[0]["workflow_phase_id"], "alpha");
assert_eq!(started[0]["workflow_task_label"], "one");
assert_eq!(started[0]["workflow_child_index"], 0);
assert_eq!(started[0]["label"], "one");
assert_eq!(started[1]["workflow_run_id"], payload["run_id"]);
assert_eq!(started[1]["workflow_phase_id"], "beta-explicit");
assert_eq!(started[1]["workflow_task_label"], "two");
assert_eq!(started[1]["workflow_child_index"], 1);
assert_eq!(started[1]["label"], "two");
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn declarative_parallel_spawn_failure_nulls_slot_and_continues() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("reduce-with-partial").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": r#"export default workflow({
"goal": "partial success fan-out",
"nodes": [
{
"branch": {
"id": "parallel",
"parallel": true,
"children": [
{
"agent": {
"id": "bad-profile",
"prompt": "This child should be rejected before model execution.",
"profile": "missing-profile"
}
}
]
}
},
{
"reduce": {
"id": "summary",
"inputs": ["bad-profile"],
"prompt": "Summarize whatever survived the parallel fan-out."
}
}
]
});"#
}),
&ctx,
)
.await
.expect("partial-success workflow still returns run record");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "degraded", "{payload}");
let degradation = payload["error"].as_str().expect("degradation surfaced");
assert!(
degradation.contains("result may be partial"),
"{degradation}"
);
assert!(
payload["result"]["bad-profile"].is_null(),
"failed parallel slot should be null: {}",
payload["result"]
);
assert_eq!(payload["result"]["summary"], "reduce-with-partial");
let progress = payload["progress"]
.as_array()
.expect("progress array")
.iter()
.filter_map(|v| v.as_str())
.collect::<Vec<_>>()
.join("\n");
assert!(
progress.contains("missing-profile") && progress.contains("dropped a failed slot"),
"breadcrumb should surface the spawn rejection:\n{progress}"
);
let failures = payload["dispatch_failures"]
.as_array()
.expect("dispatch failures surfaced on the run record");
assert_eq!(failures.len(), 1, "{payload}");
assert!(
failures[0]["message"]
.as_str()
.unwrap_or_default()
.contains("missing-profile"),
"{failures:?}"
);
assert_eq!(
result.metadata.as_ref().expect("metadata")["dispatch_failure_count"],
1
);
assert!(
calls.load(Ordering::SeqCst) >= 1,
"reduce should still run after a null parallel slot"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn parallel_slots_rejected_by_the_vm_cannot_report_plain_success() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("unused").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": "return await parallel([() => task({ description: 'bad a', type: 'explore', allowedTools: [], cwd: '/absolute/a' }), () => task({ description: 'bad b', type: 'explore', allowedTools: [], cwd: '/absolute/b' })]);"
}),
&ctx,
)
.await
.expect("vm-rejected fan-out still returns the run record");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "failed", "{payload}");
let error = payload["error"].as_str().expect("error surfaced");
assert!(
error.contains("all 2 task dispatch(es) were rejected"),
"error should name the total rejection: {error}"
);
let failures = payload["dispatch_failures"]
.as_array()
.expect("collected dispatch failures");
assert_eq!(failures.len(), 2, "{payload}");
for failure in failures {
let message = failure["message"].as_str().expect("failure message");
assert!(message.contains("bounded repo-relative paths"), "{message}");
}
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"no provider call should be spent on vm-rejected slots"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn partially_dropped_parallel_slots_degrade_the_run() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("child done").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": "return await parallel([() => task({ description: 'say done', type: 'explore', allowedTools: [] }), () => task({ description: 'bad slot', type: 'explore', allowedTools: [], cwd: '/absolute/path' })]);"
}),
&ctx,
)
.await
.expect("partially dropped fan-out still returns the run record");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "degraded", "{payload}");
let error = payload["error"].as_str().expect("degradation surfaced");
assert!(
error.contains("1 dispatch(es) were rejected"),
"error should count the dropped slot: {error}"
);
assert!(
error.contains("result may be partial"),
"error should warn the result is partial: {error}"
);
let results = payload["result"].as_array().expect("run kept its output");
assert_eq!(results.len(), 2, "{payload}");
assert!(results[1].is_null(), "the rejected slot stays null");
assert!(
calls.load(Ordering::SeqCst) >= 1,
"the healthy slot still ran"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn parallel_fan_out_with_every_dispatch_rejected_fails_the_run() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("unused").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": r#"export default workflow({
"goal": "total dispatch failure fan-out",
"nodes": [
{
"branch": {
"id": "parallel",
"parallel": true,
"children": [
{
"agent": {
"id": "bad-one",
"prompt": "Rejected before model execution.",
"profile": "missing-profile"
}
},
{
"agent": {
"id": "bad-two",
"prompt": "Also rejected before model execution.",
"profile": "missing-profile"
}
}
]
}
}
]
});"#
}),
&ctx,
)
.await
.expect("total dispatch failure still returns the run record");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "failed", "{payload}");
let error = payload["error"].as_str().expect("error surfaced");
assert!(
error.contains("all 2 task dispatch(es) were rejected"),
"error should name the total dispatch failure: {error}"
);
let failures = payload["dispatch_failures"]
.as_array()
.expect("collected dispatch failures");
assert_eq!(failures.len(), 2, "{payload}");
for failure in failures {
let message = failure["message"].as_str().expect("failure message");
assert!(message.contains("missing-profile"), "{message}");
}
assert_eq!(payload["child_ids"].as_array().unwrap().len(), 0);
assert_eq!(
result.metadata.as_ref().expect("metadata")["dispatch_failure_count"],
2
);
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"no provider call should be spent on an all-rejected fan-out"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn declarative_dependency_results_are_forwarded_to_downstream_prompt() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls, bodies) = fake_chat_client_capturing("upstream-output").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": r#"export default workflow({
"goal": "dependency forwarding",
"nodes": [
{
"agent": {
"id": "first",
"prompt": "Produce the upstream finding.",
"agent_type": "review"
}
},
{
"agent": {
"id": "second",
"prompt": "Use the upstream finding.",
"agent_type": "review",
"depends_on_results": ["first"]
}
}
]
});"#
}),
&ctx,
)
.await
.expect("dependency workflow should complete");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(payload["execution"]["status"], "succeeded");
assert_eq!(
payload["execution"]["leaf_results"][0]["output"],
"upstream-output"
);
assert_eq!(
payload["execution"]["leaf_results"][1]["output"],
"upstream-output"
);
assert_eq!(calls.load(Ordering::SeqCst), 2);
let bodies = bodies.lock().expect("captured bodies");
let second_body = bodies.get(1).expect("second provider call").to_string();
assert!(second_body.contains("--- first ---"), "{second_body}");
assert!(second_body.contains("upstream-output"), "{second_body}");
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn workflow_runtime_gates_promote_handoff_and_block_downstream_role() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_gate".to_string();
let gates = vec![GateSpec {
id: "scout-findings".to_string(),
role: "scout".to_string(),
on: GateOn::RoleComplete,
gate: GateKind::Approve,
on_fail: codewhale_workflow::GateOnFail::Block,
blocks_role: Some("implementer".to_string()),
max_retries: 0,
artifact_kind: Some("findings".to_string()),
require_explicit_verdict: false,
}];
let spec = WorkflowSpec {
id: Some("gate-fixture".to_string()),
goal: "gate fixture".to_string(),
description: None,
budget: BudgetSpec::default(),
permissions: Default::default(),
model_policy: Default::default(),
promotion_policy: Default::default(),
gates: gates.clone(),
nodes: Vec::new(),
};
state.runs.lock().expect("runs").insert(
run_id.clone(),
WorkflowRunRecord::new(run_id.clone(), None, None, Some(&spec)),
);
let driver = SubAgentWorkflowDriver::new(
run_id.clone(),
manager,
runtime,
state.clone(),
None,
WorkflowFleetBinding::None,
gates,
);
driver.evaluate_gates_for_completed_role(&RuntimeTaskRecord {
agent_id: "scout-agent".to_string(),
label: Some("scout".to_string()),
role: Some("scout".to_string()),
status: IrWorkflowRunStatus::Succeeded,
output: Some("findings: inspect tui exit path".to_string()),
schema_error: None,
usage: None,
});
let mut implementer = TaskRequest {
description: "Use the findings.".to_string(),
subagent_type: Some("implementer".to_string()),
role: Some("implementer".to_string()),
profile: None,
model: None,
model_strength: None,
thinking: None,
cwd: None,
worktree: false,
write_authority: Some("workspace_write".to_string()),
write_roots: vec!["src".to_string()],
exact_files: Vec::new(),
coordination_contracts: vec!["test-contract".to_string()],
dependencies: Vec::new(),
acceptance: Vec::new(),
allowed_tools: Some(Vec::new()),
disallowed_tools: Vec::new(),
max_depth: None,
token_budget: None,
max_steps: None,
wall_time_secs: None,
response_schema: None,
label: Some("fix".to_string()),
phase: None,
};
let handoffs = driver
.prepare_request_for_gates(&mut implementer)
.expect("passed gate should admit implementer");
assert_eq!(handoffs.len(), 1, "{handoffs:?}");
assert_eq!(handoffs[0].kind, "findings");
assert_eq!(handoffs[0].from_role, "scout");
assert_eq!(handoffs[0].to_role, "implementer");
assert!(
implementer
.description
.contains("Workflow handoff artifacts available"),
"{}",
implementer.description
);
assert!(
implementer.description.contains("inspect tui exit path"),
"{}",
implementer.description
);
driver.evaluate_gates_for_completed_role(&RuntimeTaskRecord {
agent_id: "scout-agent-2".to_string(),
label: Some("scout".to_string()),
role: Some("scout".to_string()),
status: IrWorkflowRunStatus::Failed,
output: Some("scout incomplete".to_string()),
schema_error: None,
usage: None,
});
let mut blocked = TaskRequest {
description: "Try after block.".to_string(),
role: Some("implementer".to_string()),
..implementer.clone()
};
let err = driver
.prepare_request_for_gates(&mut blocked)
.expect_err("blocked gate should reject downstream role");
assert!(err.to_string().contains("scout incomplete"), "{err}");
let run = state
.runs
.lock()
.expect("runs")
.get(&run_id)
.cloned()
.expect("run");
assert!(
run.gate_status
.iter()
.any(|line| line.gate_id == "scout-findings"
&& line.state == "blocked"
&& line.blocked_reason.as_deref() == Some("scout incomplete")),
"{:?}",
run.gate_status
);
assert!(
run.events
.iter()
.any(|event| event.event_type() == "gate_updated"),
"{:?}",
run.events
);
assert_eq!(
run.events
.iter()
.filter(|event| event.event_type() == "handoff_promoted")
.count(),
1,
"a later blocked gate must not publish another handoff: {:?}",
run.events
);
assert!(
run.events
.iter()
.all(|event| event.event_type() != "handoff_consumed"),
"request preparation alone must not claim consumption: {:?}",
run.events
);
}
#[tokio::test]
async fn workflow_handoff_is_delivered_to_exactly_one_task() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_handoff_once".to_string();
let gates = vec![GateSpec {
id: "scout-findings".to_string(),
role: "scout".to_string(),
on: GateOn::RoleComplete,
gate: GateKind::Approve,
on_fail: codewhale_workflow::GateOnFail::Block,
blocks_role: Some("implementer".to_string()),
max_retries: 0,
artifact_kind: Some("findings".to_string()),
require_explicit_verdict: false,
}];
let spec = WorkflowSpec {
id: Some("handoff-once-fixture".to_string()),
goal: "handoff delivered once".to_string(),
description: None,
budget: BudgetSpec::default(),
permissions: Default::default(),
model_policy: Default::default(),
promotion_policy: Default::default(),
gates: gates.clone(),
nodes: Vec::new(),
};
state.runs.lock().expect("runs").insert(
run_id.clone(),
WorkflowRunRecord::new(run_id.clone(), None, None, Some(&spec)),
);
let driver = SubAgentWorkflowDriver::new(
run_id.clone(),
manager,
runtime,
state.clone(),
None,
WorkflowFleetBinding::None,
gates,
);
driver.evaluate_gates_for_completed_role(&RuntimeTaskRecord {
agent_id: "scout-agent".to_string(),
label: Some("scout".to_string()),
role: Some("scout".to_string()),
status: IrWorkflowRunStatus::Succeeded,
output: Some("findings: exactly once".to_string()),
schema_error: None,
usage: None,
});
let implementer = TaskRequest {
description: "Use the findings.".to_string(),
subagent_type: Some("implementer".to_string()),
role: Some("implementer".to_string()),
profile: None,
model: None,
model_strength: None,
thinking: None,
cwd: None,
worktree: false,
write_authority: Some("workspace_write".to_string()),
write_roots: vec!["src".to_string()],
exact_files: Vec::new(),
coordination_contracts: Vec::new(),
dependencies: Vec::new(),
acceptance: Vec::new(),
allowed_tools: Some(Vec::new()),
disallowed_tools: Vec::new(),
max_depth: None,
token_budget: None,
max_steps: None,
wall_time_secs: None,
response_schema: None,
label: Some("fix".to_string()),
phase: None,
};
let mut first = implementer.clone();
let handoffs = driver
.prepare_request_for_gates(&mut first)
.expect("passed gate should admit first implementer");
assert_eq!(handoffs.len(), 1, "{handoffs:?}");
assert!(first.description.contains("exactly once"));
let mut second = TaskRequest {
description: "Second implementer task.".to_string(),
..implementer
};
let handoffs = driver
.prepare_request_for_gates(&mut second)
.expect("passed gate should admit second implementer");
assert!(
handoffs.is_empty(),
"handoff already consumed must not re-deliver: {handoffs:?}"
);
assert!(
!second
.description
.contains("Workflow handoff artifacts available"),
"{}",
second.description
);
}
#[tokio::test]
async fn workflow_gate_evaluation_error_persists_blocked_and_denies_target_role() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx,
true,
None,
manager.clone(),
);
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_malformed_gate".to_string();
let gates = vec![GateSpec {
id: String::new(),
role: "scout".to_string(),
on: GateOn::RoleComplete,
gate: GateKind::Approve,
on_fail: codewhale_workflow::GateOnFail::Block,
blocks_role: Some("implementer".to_string()),
max_retries: 0,
artifact_kind: Some("findings".to_string()),
require_explicit_verdict: false,
}];
let spec = WorkflowSpec {
id: Some("malformed-gate-fixture".to_string()),
goal: "malformed gate must fail closed".to_string(),
description: None,
budget: BudgetSpec::default(),
permissions: Default::default(),
model_policy: Default::default(),
promotion_policy: Default::default(),
gates: gates.clone(),
nodes: Vec::new(),
};
state.runs.lock().expect("runs").insert(
run_id.clone(),
WorkflowRunRecord::new(run_id.clone(), None, None, Some(&spec)),
);
let driver = SubAgentWorkflowDriver::new(
run_id.clone(),
manager,
runtime,
state.clone(),
None,
WorkflowFleetBinding::None,
gates,
);
driver.evaluate_gates_for_completed_role(&RuntimeTaskRecord {
agent_id: "scout-agent".to_string(),
label: Some("scout".to_string()),
role: Some("scout".to_string()),
status: IrWorkflowRunStatus::Succeeded,
output: Some("findings".to_string()),
schema_error: None,
usage: None,
});
let mut request = TaskRequest {
description: "Must not be admitted.".to_string(),
subagent_type: Some("implementer".to_string()),
role: Some("implementer".to_string()),
profile: None,
model: None,
model_strength: None,
thinking: None,
cwd: None,
worktree: false,
write_authority: Some("workspace_write".to_string()),
write_roots: vec!["src".to_string()],
exact_files: Vec::new(),
coordination_contracts: vec!["test-contract".to_string()],
dependencies: Vec::new(),
acceptance: Vec::new(),
allowed_tools: Some(Vec::new()),
disallowed_tools: Vec::new(),
max_depth: None,
token_budget: None,
max_steps: None,
wall_time_secs: None,
response_schema: None,
label: Some("blocked".to_string()),
phase: None,
};
let error = driver
.prepare_request_for_gates(&mut request)
.expect_err("malformed gate must deny its target role");
assert!(
error.to_string().contains("gate id must not be empty"),
"{error}"
);
let board = driver.gate_board.lock().expect("gate board");
assert!(matches!(
board.gates.get(""),
Some(GateState::Blocked { reason }) if reason.contains("gate id must not be empty")
));
assert!(board.artifacts.is_empty(), "{:?}", board.artifacts);
drop(board);
let run = state
.runs
.lock()
.expect("runs")
.get(&run_id)
.cloned()
.expect("run");
assert!(run.gate_status.iter().any(|line| {
line.gate_id.is_empty()
&& line.state == "blocked"
&& line
.blocked_reason
.as_deref()
.is_some_and(|reason| reason.contains("gate id must not be empty"))
}));
assert!(
run.events
.iter()
.all(|event| event.event_type() != "handoff_promoted"),
"{:?}",
run.events
);
}
#[tokio::test]
async fn workflow_handoff_record_error_changes_pass_to_blocked() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx,
true,
None,
manager.clone(),
);
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_handoff_record_error".to_string();
let gates = vec![GateSpec {
id: "scout-findings".to_string(),
role: "scout".to_string(),
on: GateOn::RoleComplete,
gate: GateKind::Approve,
on_fail: codewhale_workflow::GateOnFail::Block,
blocks_role: Some("implementer".to_string()),
max_retries: 0,
artifact_kind: Some("findings".to_string()),
require_explicit_verdict: false,
}];
let spec = WorkflowSpec {
id: Some("handoff-record-error-fixture".to_string()),
goal: "failed handoff recording must fail closed".to_string(),
description: None,
budget: BudgetSpec::default(),
permissions: Default::default(),
model_policy: Default::default(),
promotion_policy: Default::default(),
gates: gates.clone(),
nodes: Vec::new(),
};
state.runs.lock().expect("runs").insert(
run_id.clone(),
WorkflowRunRecord::new(run_id.clone(), None, None, Some(&spec)),
);
let driver = SubAgentWorkflowDriver::new(
run_id.clone(),
manager,
runtime,
state.clone(),
None,
WorkflowFleetBinding::None,
gates,
);
driver.gate_board.lock().expect("gate board").lane_id = "wrong-lane".to_string();
driver.evaluate_gates_for_completed_role(&RuntimeTaskRecord {
agent_id: "scout-agent".to_string(),
label: Some("scout".to_string()),
role: Some("scout".to_string()),
status: IrWorkflowRunStatus::Succeeded,
output: Some("findings".to_string()),
schema_error: None,
usage: None,
});
let mut request = TaskRequest {
description: "Must not be admitted.".to_string(),
subagent_type: Some("implementer".to_string()),
role: Some("implementer".to_string()),
profile: None,
model: None,
model_strength: None,
thinking: None,
cwd: None,
worktree: false,
write_authority: Some("workspace_write".to_string()),
write_roots: vec!["src".to_string()],
exact_files: Vec::new(),
coordination_contracts: vec!["test-contract".to_string()],
dependencies: Vec::new(),
acceptance: Vec::new(),
allowed_tools: Some(Vec::new()),
disallowed_tools: Vec::new(),
max_depth: None,
token_budget: None,
max_steps: None,
wall_time_secs: None,
response_schema: None,
label: Some("blocked".to_string()),
phase: None,
};
let error = driver
.prepare_request_for_gates(&mut request)
.expect_err("unrecorded handoff must deny its target role");
assert!(
error.to_string().contains("handoff could not be recorded"),
"{error}"
);
let board = driver.gate_board.lock().expect("gate board");
assert!(matches!(
board.gates.get("scout-findings"),
Some(GateState::Blocked { reason }) if reason.contains("does not match board lane")
));
assert!(board.artifacts.is_empty(), "{:?}", board.artifacts);
drop(board);
let run = state
.runs
.lock()
.expect("runs")
.get(&run_id)
.cloned()
.expect("run");
assert!(run.events.iter().any(|event| {
matches!(
&event.kind,
WorkflowUiEventKind::GateUpdated {
state,
blocked_reason: Some(reason),
..
} if state == "blocked" && reason.contains("handoff could not be recorded")
)
}));
assert!(
run.events
.iter()
.all(|event| event.event_type() != "handoff_promoted"),
"{:?}",
run.events
);
}
#[test]
fn explicit_gate_verdict_only_reads_first_standalone_token() {
assert_eq!(
explicit_gate_verdict(Some("\n APPROVE \nreview complete")),
Some(ExplicitGateVerdict::Approve)
);
assert_eq!(
explicit_gate_verdict(Some("PASS\nverification complete")),
Some(ExplicitGateVerdict::Approve)
);
assert_eq!(
explicit_gate_verdict(Some("BLOCK\nmissing receipt")),
Some(ExplicitGateVerdict::Reject)
);
assert_eq!(
explicit_gate_verdict(Some("\nFAIL\nmissing receipt")),
Some(ExplicitGateVerdict::Reject)
);
assert_eq!(
explicit_gate_verdict(Some("Review result: BLOCK")),
None,
"prose remains backward-compatible success output"
);
assert_eq!(
explicit_gate_verdict(Some("review notes\nBLOCK")),
None,
"later verdict words must not override the first meaningful line"
);
}
#[test]
fn required_explicit_gate_verdict_fails_closed_when_missing_or_malformed() {
let mut record = RuntimeTaskRecord {
agent_id: "reviewer-malformed".to_string(),
label: Some("reviewer".to_string()),
role: Some("reviewer".to_string()),
status: IrWorkflowRunStatus::Succeeded,
output: Some("Review result: BLOCK".to_string()),
schema_error: None,
usage: None,
};
match gate_outcome_for_completed_role(&record, true, None) {
GateOutcome::Fail { reason } => {
assert!(
reason.contains("required first-line gate verdict"),
"{reason}"
);
}
outcome => panic!("required malformed verdict must fail closed: {outcome:?}"),
}
assert_eq!(
gate_outcome_for_completed_role(&record, false, None),
GateOutcome::Pass,
"legacy gates retain pass-on-success behavior"
);
record.output = None;
assert!(matches!(
gate_outcome_for_completed_role(&record, true, None),
GateOutcome::Fail { .. }
));
}
#[test]
fn required_gate_artifact_rejects_bare_or_placeholder_approval() {
let mut record = RuntimeTaskRecord {
agent_id: "implementer-bare".to_string(),
label: Some("implementer".to_string()),
role: Some("implementer".to_string()),
status: IrWorkflowRunStatus::Succeeded,
output: Some("APPROVE".to_string()),
schema_error: None,
usage: None,
};
match gate_outcome_for_completed_role(&record, true, Some("verification_plan")) {
GateOutcome::Fail { reason } => {
assert!(
reason.contains("verification_plan artifact body"),
"{reason}"
);
}
outcome => panic!("bare approval must not promote an empty artifact: {outcome:?}"),
}
record.output = Some("APPROVE\nacceptance evidence".to_string());
match gate_outcome_for_completed_role(&record, true, Some("verification_plan")) {
GateOutcome::Fail { reason } => {
assert!(
reason.contains("verification_plan artifact body"),
"{reason}"
);
}
outcome => {
panic!("one placeholder line must not count as an artifact: {outcome:?}");
}
}
record.output = Some("APPROVE\nPLAN\n- verify the typed receipt".to_string());
assert_eq!(
gate_outcome_for_completed_role(&record, true, Some("verification_plan")),
GateOutcome::Pass
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn terminal_blocked_gate_fails_workflow_finalization() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) =
fake_chat_client("BLOCK\nFINAL RECEIPT\n- missing terminal evidence").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"script": r#"export default workflow({
"goal": "fail closed on the terminal release verdict",
"gates": [
{
"id": "terminal-release",
"role": "release_lead",
"on": "role_complete",
"gate": "approve",
"on_fail": "block",
"max_retries": 0,
"artifact_kind": "final_receipt",
"require_explicit_verdict": true
}
],
"nodes": [
{
"agent": {
"id": "release-receipt",
"prompt": "Return the terminal verdict and receipt.",
"agent_type": "general",
"role": "release_lead",
"mode": "read_only",
"permissions": { "deny_all_tools": true },
"budget": { "max_steps": 1 }
}
}
]
});"#
}),
&ctx,
)
.await
.expect("blocked terminal gate should return its failed run record");
let payload: Value = serde_json::from_str(&result.content).expect("workflow JSON");
assert_eq!(calls.load(Ordering::SeqCst), 1, "{payload}");
assert_eq!(payload["status"], "failed", "{payload}");
assert_eq!(payload["execution"]["status"], "failed", "{payload}");
assert!(
payload["error"]
.as_str()
.is_some_and(|error| error.contains("terminal-release")
&& error.contains("ended blocked")
&& error.contains("missing terminal evidence")),
"{payload}"
);
assert!(payload["gate_status"].as_array().is_some_and(|gates| {
gates
.iter()
.any(|gate| gate["gate_id"] == "terminal-release" && gate["state"] == "blocked")
}));
assert!(payload["events"].as_array().is_some_and(|events| {
events
.iter()
.any(|event| event["type"] == "run_completed" && event["status"] == "failed")
}));
}
#[tokio::test]
async fn workflow_runtime_gate_honors_explicit_reviewer_verdicts() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx,
true,
None,
manager.clone(),
);
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_explicit_verdict".to_string();
let gates = vec![GateSpec {
id: "review-findings".to_string(),
role: "reviewer".to_string(),
on: GateOn::RoleComplete,
gate: GateKind::Review,
on_fail: codewhale_workflow::GateOnFail::Block,
blocks_role: Some("verifier".to_string()),
max_retries: 0,
artifact_kind: Some("review_report".to_string()),
require_explicit_verdict: true,
}];
let spec = WorkflowSpec {
id: Some("explicit-verdict-fixture".to_string()),
goal: "honor reviewer verdict".to_string(),
description: None,
budget: BudgetSpec::default(),
permissions: Default::default(),
model_policy: Default::default(),
promotion_policy: Default::default(),
gates: gates.clone(),
nodes: Vec::new(),
};
state.runs.lock().expect("runs").insert(
run_id.clone(),
WorkflowRunRecord::new(run_id.clone(), None, None, Some(&spec)),
);
let driver = SubAgentWorkflowDriver::new(
run_id,
manager,
runtime,
state,
None,
WorkflowFleetBinding::None,
gates,
);
driver.evaluate_gates_for_completed_role(&RuntimeTaskRecord {
agent_id: "reviewer-block".to_string(),
label: Some("reviewer".to_string()),
role: Some("reviewer".to_string()),
status: IrWorkflowRunStatus::Succeeded,
output: Some("\nBLOCK\nmissing terminal receipt".to_string()),
schema_error: None,
usage: None,
});
let verifier_request = || TaskRequest {
description: "Verify the accepted review.".to_string(),
subagent_type: Some("verifier".to_string()),
role: Some("verifier".to_string()),
profile: None,
model: None,
model_strength: None,
thinking: None,
cwd: None,
worktree: false,
write_authority: Some("read_only".to_string()),
write_roots: Vec::new(),
exact_files: Vec::new(),
coordination_contracts: Vec::new(),
dependencies: Vec::new(),
acceptance: Vec::new(),
allowed_tools: Some(Vec::new()),
disallowed_tools: Vec::new(),
max_depth: None,
token_budget: None,
max_steps: None,
wall_time_secs: None,
response_schema: None,
label: Some("verify".to_string()),
phase: None,
};
let mut blocked_verifier = verifier_request();
let error = driver
.prepare_request_for_gates(&mut blocked_verifier)
.expect_err("successful reviewer BLOCK must not admit verifier");
assert!(error.to_string().contains("BLOCK"), "{error}");
{
let board = driver.gate_board.lock().expect("gate board");
assert!(
board.artifacts.is_empty(),
"rejected output must not produce a handoff: {:?}",
board.artifacts
);
assert!(matches!(
board.gates.get("review-findings"),
Some(GateState::Blocked { .. })
));
}
driver.evaluate_gates_for_completed_role(&RuntimeTaskRecord {
agent_id: "reviewer-approve".to_string(),
label: Some("reviewer".to_string()),
role: Some("reviewer".to_string()),
status: IrWorkflowRunStatus::Succeeded,
output: Some("APPROVE\nEVIDENCE REVIEW\n- all receipt owners confirmed".to_string()),
schema_error: None,
usage: None,
});
let mut admitted_verifier = verifier_request();
driver
.prepare_request_for_gates(&mut admitted_verifier)
.expect("explicit reviewer APPROVE should admit verifier");
assert!(
admitted_verifier
.description
.contains("all receipt owners confirmed"),
"{}",
admitted_verifier.description
);
let board = driver.gate_board.lock().expect("gate board");
assert!(
board.artifacts.is_empty(),
"the admitted verifier consumed the handoff; spent artifacts leave the board: {:?}",
board.artifacts
);
assert!(matches!(
board.gates.get("review-findings"),
Some(GateState::Passed)
));
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn workflow_status_lists_compact_typed_receipts() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, _calls) = fake_chat_client("status-output").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let run = tool
.execute(
json!({
"action": "run",
"script": r#"export default workflow({
"id": "status-fixture",
"goal": "status summary",
"nodes": [
{
"agent": {
"id": "inspect",
"prompt": "Inspect the code.",
"agent_type": "review"
}
}
]
});"#
}),
&ctx,
)
.await
.expect("workflow run");
let run_payload: Value = serde_json::from_str(&run.content).expect("run json");
let status = tool
.execute(json!({"action": "status"}), &ctx)
.await
.expect("workflow status");
let status_payload: Value = serde_json::from_str(&status.content).expect("status json");
let summary = &status_payload["runs"][0];
assert_eq!(status_payload["count"], 1);
assert_eq!(summary["run_id"], run_payload["run_id"]);
assert_eq!(summary["workflow_id"], "status-fixture");
assert_eq!(summary["workflow_goal"], "status summary");
assert_eq!(summary["status"], "completed");
assert_eq!(summary["execution_status"], "succeeded");
assert_eq!(summary["child_count"], 1);
assert_eq!(summary["leaf_count"], 1);
assert_eq!(summary["branch_count"], 0);
assert_eq!(summary["control_count"], 0);
assert!(summary["event_count"].as_u64().unwrap_or_default() >= 3);
assert_eq!(summary["last_event_type"], "run_completed");
assert!(summary.get("result").is_none());
assert!(summary.get("execution").is_none());
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn workflow_status_survives_tool_rebuild_via_journal() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, _calls) = fake_chat_client("journal-output").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let tool = WorkflowTool::new(manager.clone(), runtime.clone());
let run = tool
.execute(
json!({
"action": "run",
"script": "return { ok: true };"
}),
&ctx,
)
.await
.expect("workflow run");
let run_payload: Value = serde_json::from_str(&run.content).expect("run json");
let run_id = run_payload["run_id"].as_str().expect("run id");
let journal_path = tmp.path().join(".codewhale/workflow-runs.jsonl");
assert!(
journal_path.exists(),
"journal should be created under workspace"
);
let rebuilt = WorkflowTool::new(
manager.clone(),
SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
),
);
let status = rebuilt
.execute(json!({"action": "status", "run_id": run_id}), &ctx)
.await
.expect("workflow status after rebuild");
let status_payload: Value = serde_json::from_str(&status.content).expect("status json");
assert_eq!(status_payload["run_id"], run_id);
assert_eq!(status_payload["status"], "completed");
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn workflow_status_surfaces_schema_failure_instead_of_null_success() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, _calls) = fake_chat_client(r#"{"refuted":"yes"}"#).await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let tool = WorkflowTool::new(manager, runtime);
let run = tool
.execute(
json!({
"action": "run",
"script": r#"
return await parallel([
() => task({
description: "Return the schema fixture.",
responseSchema: {
type: "object",
properties: { refuted: { type: "boolean" } },
required: ["refuted"],
},
}),
]);
"#
}),
&ctx,
)
.await
.expect("workflow run returns a record");
let run_payload: Value = serde_json::from_str(&run.content).expect("run json");
assert_eq!(run_payload["status"], "failed");
assert!(run_payload["result"].is_null());
assert!(
run_payload["error"]
.as_str()
.unwrap()
.contains("responseSchema validation")
);
assert!(
run_payload["progress"]
.as_array()
.unwrap()
.iter()
.any(|message| message
.as_str()
.is_some_and(|message| message.contains("schema validation failed"))),
"schema validation error should be visible in the run receipt: {run_payload}"
);
assert!(
run_payload["events"]
.as_array()
.unwrap()
.iter()
.any(|event| event["type"] == "task_schema_validation_failed"
&& event["message"]
.as_str()
.is_some_and(|message| message.contains("responseSchema validation"))),
"schema validation event should be visible in the typed receipt: {run_payload}"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn declarative_issue_audit_fixture_runs_through_subagent_driver() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let workflow_dir = tmp.path().join("workflows");
std::fs::create_dir_all(&workflow_dir).expect("workflow dir");
let fixture = std::fs::read_to_string(
std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
.join("../../workflows/issue_audit.workflow.js"),
)
.expect("issue audit fixture");
std::fs::write(workflow_dir.join("issue_audit.workflow.js"), fixture)
.expect("write fixture into workspace");
let mut ctx = ToolContext::new(tmp.path().to_path_buf());
ctx.runtime.work = Some(crate::work_graph::new_shared_work_runtime(
crate::tools::todo::new_shared_todo_list(),
crate::tools::plan::new_shared_plan_state(),
));
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
let (client, calls) = fake_chat_client("audited").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"source_path": "workflows/issue_audit.workflow.js"
}),
&ctx,
)
.await
.expect("declarative workflow should complete");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(payload["result"]["code-audit"], "audited");
assert_eq!(payload["result"]["test-audit"], "audited");
assert_eq!(payload["result"]["docs-audit"], "audited");
assert_eq!(payload["result"]["synthesize-release-risk"], "audited");
assert_eq!(payload["execution"]["status"], "succeeded");
assert_eq!(
payload["execution"]["leaf_results"]
.as_array()
.expect("leaf results")
.len(),
3
);
assert_eq!(
payload["execution"]["branch_results"][0]["branch_id"],
"parallel-audit"
);
assert!(
payload["execution"]["control_node_results"]
.as_array()
.expect("control results")
.iter()
.any(|result| result["node_id"] == "synthesize-release-risk"
&& result["kind"] == "reduce"
&& result["status"] == "succeeded")
);
assert_eq!(payload["child_ids"].as_array().unwrap().len(), 4);
assert_eq!(calls.load(Ordering::SeqCst), 4);
assert!(
payload["progress"]
.as_array()
.unwrap()
.iter()
.any(|message| message == "phase: parallel-audit")
);
let work = ctx.runtime.work.as_ref().expect("work runtime");
let graph = work
.capture(Some(&ctx.state_namespace))
.expect("capture workflow work")
.expect("workflow graph")
.graph;
let workflow_external = format!(
"workflow:{}",
payload["run_id"].as_str().expect("workflow run id")
);
let workflow_operations = graph
.nodes
.iter()
.filter(|node| {
node.binding
.as_ref()
.is_some_and(|binding| binding.external == workflow_external)
})
.collect::<Vec<_>>();
assert_eq!(
workflow_operations.len(),
1,
"one accountable Workflow operation: {graph:#?}"
);
assert_eq!(
workflow_operations[0].state,
crate::work_graph::NodeState::Completed
);
for child_id in payload["child_ids"].as_array().expect("workflow child ids") {
let worker_external = format!(
"worker:{}",
child_id.as_str().expect("workflow child id string")
);
assert!(
graph.nodes.iter().any(|node| {
node.binding
.as_ref()
.is_some_and(|binding| binding.external == worker_external)
}),
"Workflow worker {worker_external} must remain inspectable in the same graph: {graph:#?}"
);
}
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn stopship_acceptance_fixture_emits_role_gate_and_terminal_receipts() {
let _retry_guard = workflow_test_retry_guard();
let _env_lock = crate::test_support::lock_test_env();
let tmp = tempfile::tempdir().expect("tempdir");
let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", tmp.path());
let workflow_dir = tmp.path().join("workflows");
let fleet_dir = tmp.path().join("fleets");
std::fs::create_dir_all(&workflow_dir).expect("workflow dir");
std::fs::create_dir_all(&fleet_dir).expect("fleet dir");
let repo_root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../..");
std::fs::copy(
repo_root.join("workflows/stopship.workflow.js"),
workflow_dir.join("stopship.workflow.js"),
)
.expect("copy stopship acceptance fixture");
std::fs::copy(
repo_root.join("fleets/stopship.toml"),
fleet_dir.join("stopship.toml"),
)
.expect("copy stopship fleet");
let source = std::fs::read_to_string(workflow_dir.join("stopship.workflow.js"))
.expect("read stopship acceptance fixture");
let compiled =
codewhale_workflow::compile_javascript_workflow("stopship.workflow.js", &source)
.expect("compile stopship acceptance fixture");
let codewhale_workflow::WorkflowNode::Sequence(sequence) = &compiled.nodes[0] else {
panic!("stopship fixture should be one ordered role chain");
};
for (index, node) in sequence.children.iter().enumerate() {
let codewhale_workflow::WorkflowNode::Leaf(leaf) = node else {
panic!("stopship role chain must contain only leaves");
};
let tools = leaf_allowed_tools(leaf).expect("lower stopship child tools");
if index == 0 {
assert!(tools.as_ref().is_some_and(|tools| !tools.is_empty()));
} else {
assert_eq!(
tools,
Some(Vec::<String>::new()),
"downstream handoff consumer {} must receive no tools",
leaf.id
);
}
}
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 8);
let responses = [
r#"APPROVE
SOURCE EVIDENCE
- crates/cli/src/lib.rs: load_named_fleet
- crates/workflow/src/role_resolve.rs: resolve_workflow_agent
- crates/cli/src/lib.rs: start_lane
- crates/tui/src/tools/workflow.rs: record_task_started
- crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::GateUpdated
- crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::RunCompleted -> terminal_completed_receipt
- crates/lane/src/runtime.rs: process_exit_receipt -> lane_reconciled"#,
r#"APPROVE
PLAN
- fleets/stopship.toml: name = "stopship" -> named Fleet loading
- crates/workflow/src/role_resolve.rs: resolve_workflow_agent -> role resolution
- crates/cli/src/lib.rs: start_lane -> tmux Lane launch
- crates/tui/src/tools/workflow.rs: record_task_started -> typed task_started
- crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::GateUpdated -> gate promotion
- crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::RunCompleted -> terminal_completed_receipt
- crates/lane/src/runtime.rs: process_exit_receipt -> tmux Lane reconciliation"#,
r#"APPROVE
EVIDENCE REVIEW
- fleets/stopship.toml: name = "stopship"
- crates/workflow/src/role_resolve.rs: resolve_workflow_agent
- crates/cli/src/lib.rs: start_lane
- crates/tui/src/tools/workflow.rs: record_task_started
- crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::GateUpdated
- crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::RunCompleted -> terminal_completed_receipt
- crates/lane/src/runtime.rs: process_exit_receipt -> lane_reconciled"#,
r#"APPROVE
EVIDENCE MATRIX
- fleet_load: fleets/stopship.toml: name = "stopship"
- role_resolution: crates/workflow/src/role_resolve.rs: resolve_workflow_agent
- lane_launch: crates/cli/src/lib.rs: start_lane
- task_started: crates/tui/src/tools/workflow.rs: record_task_started
- gate_updated: crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::GateUpdated
- run_completed: crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::RunCompleted -> terminal_completed_receipt
- lane_exit: crates/lane/src/runtime.rs: process_exit_receipt -> lane_reconciled"#,
r#"APPROVE
FINAL RECEIPT
- fleet_load: fleets/stopship.toml: name = "stopship"
- role_resolution: crates/workflow/src/role_resolve.rs: resolve_workflow_agent
- lane_launch: crates/cli/src/lib.rs: start_lane
- task_started: crates/tui/src/tools/workflow.rs: record_task_started
- gate_updated: crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::GateUpdated
- run_completed: crates/tui/src/tools/workflow.rs: WorkflowUiEventKind::RunCompleted -> terminal_completed_receipt
- lane_exit: crates/lane/src/runtime.rs: process_exit_receipt -> lane_reconciled"#,
];
let (client, calls) = fake_chat_client_responses(&responses).await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"source_path": "workflows/stopship.workflow.js",
"fleet": "stopship",
"token_budget": 60_000
}),
&ctx,
)
.await
.expect("stopship acceptance workflow returns a terminal record");
let payload: Value = serde_json::from_str(&result.content).expect("workflow JSON");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(payload["execution"]["status"], "succeeded", "{payload}");
assert_eq!(calls.load(Ordering::SeqCst), 5, "one child per Fleet role");
let approval = &payload["plan_approval"];
assert_eq!(approval["decision"], "auto_read_only", "{approval}");
assert_eq!(approval["token_budget"], 60_000, "{approval}");
assert_eq!(approval["writes"], false, "{approval}");
assert_eq!(approval["shell"], false, "{approval}");
assert_eq!(approval["network"], false, "{approval}");
assert_eq!(approval["high_budget"], false, "{approval}");
assert_eq!(approval["elevated"], false, "{approval}");
assert!(
approval["reasons"].as_array().is_some_and(Vec::is_empty),
"{approval}"
);
let events = payload["events"].as_array().expect("typed events");
let started = events
.iter()
.filter(|event| event["type"] == "task_started")
.collect::<Vec<_>>();
let expected_roles = [
("scout", "scout"),
("implementer", "builder"),
("reviewer", "reviewer"),
("verifier", "verifier"),
("release_lead", "manager"),
];
assert_eq!(started.len(), expected_roles.len(), "{started:#?}");
for (event, (role, profile)) in started.iter().zip(expected_roles) {
assert_eq!(event["role"], role);
assert_eq!(event["profile"], profile);
assert_eq!(event["resolved_profile"], profile);
assert_eq!(event["workflow_run_id"], payload["run_id"]);
}
let gates = events
.iter()
.filter(|event| event["type"] == "gate_updated")
.collect::<Vec<_>>();
assert_eq!(gates.len(), 5, "{gates:#?}");
assert!(gates.iter().all(|event| event["state"] == "passed"));
assert_eq!(gates[0]["role"], "scout");
assert_eq!(gates[0]["blocked_role"], "implementer");
assert_eq!(gates[3]["role"], "verifier");
assert_eq!(gates[3]["blocked_role"], "release_lead");
assert_eq!(gates[4]["role"], "release_lead");
assert!(gates[4]["blocked_role"].is_null());
let promoted = events
.iter()
.filter(|event| event["type"] == "handoff_promoted")
.collect::<Vec<_>>();
let consumed = events
.iter()
.filter(|event| event["type"] == "handoff_consumed")
.collect::<Vec<_>>();
let expected_handoffs = [
("scout", "implementer", "source_evidence"),
("implementer", "reviewer", "verification_plan"),
("reviewer", "verifier", "review_report"),
("verifier", "release_lead", "verification_report"),
];
assert_eq!(promoted.len(), expected_handoffs.len(), "{promoted:#?}");
assert_eq!(consumed.len(), expected_handoffs.len(), "{consumed:#?}");
let artifact_ids = promoted
.iter()
.map(|event| {
event["artifact_id"]
.as_str()
.filter(|id| id.starts_with("handoff_") && id.len() > "handoff_".len())
.expect("opaque non-empty handoff artifact id")
})
.collect::<std::collections::HashSet<_>>();
assert_eq!(
artifact_ids.len(),
promoted.len(),
"every promotion must have a unique artifact id: {promoted:#?}"
);
for (index, (from_role, to_role, kind)) in expected_handoffs.into_iter().enumerate() {
assert_eq!(promoted[index]["from_role"], from_role);
assert_eq!(promoted[index]["to_role"], to_role);
assert_eq!(promoted[index]["kind"], kind);
assert_eq!(promoted[index]["gate_id"], gates[index]["gate_id"]);
assert_eq!(
promoted[index]["producer_task_id"],
started[index]["task_id"]
);
assert!(
promoted[index].get("payload").is_none(),
"{:#?}",
promoted[index]
);
assert_eq!(
consumed[index]["artifact_id"],
promoted[index]["artifact_id"]
);
assert_eq!(consumed[index]["from_role"], from_role);
assert_eq!(consumed[index]["to_role"], to_role);
assert_eq!(consumed[index]["kind"], kind);
assert_eq!(
consumed[index]["consumer_task_id"],
started[index + 1]["task_id"]
);
assert!(
consumed[index].get("payload").is_none(),
"{:#?}",
consumed[index]
);
let producer_task_id = promoted[index]["producer_task_id"]
.as_str()
.expect("producer task id");
let consumer_task_id = consumed[index]["consumer_task_id"]
.as_str()
.expect("consumer task id");
let gate_id = promoted[index]["gate_id"].as_str().expect("gate id");
let artifact_id = promoted[index]["artifact_id"]
.as_str()
.expect("artifact id");
let task_completed_index = events
.iter()
.position(|event| {
event["type"] == "task_completed" && event["task_id"] == producer_task_id
})
.expect("producer completion receipt");
let gate_updated_index = events
.iter()
.position(|event| event["type"] == "gate_updated" && event["gate_id"] == gate_id)
.expect("gate update receipt");
let promoted_index = events
.iter()
.position(|event| {
event["type"] == "handoff_promoted" && event["artifact_id"] == artifact_id
})
.expect("handoff promotion receipt");
let consumer_started_index = events
.iter()
.position(|event| {
event["type"] == "task_started" && event["task_id"] == consumer_task_id
})
.expect("consumer start receipt");
let consumed_index = events
.iter()
.position(|event| {
event["type"] == "handoff_consumed" && event["artifact_id"] == artifact_id
})
.expect("handoff consumption receipt");
assert!(
task_completed_index < gate_updated_index
&& gate_updated_index < promoted_index
&& promoted_index < consumer_started_index
&& consumer_started_index < consumed_index,
"causal receipt order must be task_completed -> gate_updated -> handoff_promoted -> task_started -> handoff_consumed: {events:#?}"
);
}
let terminal_completed_receipt = events
.iter()
.any(|event| event["type"] == "run_completed" && event["status"] == "completed");
assert!(terminal_completed_receipt, "{events:#?}");
}
#[tokio::test]
async fn completion_from_manager_fails_closed_when_status_stays_running() {
let tmp = tempfile::tempdir().expect("tempdir");
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (completion, usage) =
completion_from_manager(manager, "missing_agent", "fallback".to_string()).await;
assert!(usage.is_none(), "fail-closed path carries no telemetry");
match completion {
TaskCompletion::Failed { message } => {
assert!(
message.contains("did not report a terminal status"),
"{message}"
);
}
other => panic!("expected timeout failure, got {other:?}"),
}
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn task_completed_and_run_completed_carry_usage_telemetry() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, calls) = fake_chat_client("telemetry-output").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager,
);
let tool = WorkflowTool::new(runtime.manager.clone(), runtime);
let run = tool
.execute(
json!({
"action": "run",
"script": r#"export default workflow({
"id": "telemetry-fixture",
"goal": "usage telemetry",
"nodes": [
{
"agent": {
"id": "inspect",
"prompt": "Inspect the code.",
"agent_type": "review"
}
}
]
});"#
}),
&ctx,
)
.await
.expect("workflow run");
let payload: Value = serde_json::from_str(&run.content).expect("run json");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(calls.load(Ordering::SeqCst), 1, "{payload}");
let events = payload["events"].as_array().expect("typed events");
let task_completed = events
.iter()
.find(|event| event["type"] == "task_completed")
.expect("task_completed event");
let usage = &task_completed["usage"];
let input = usage["input_tokens"].as_u64().expect("worker input tokens");
let output = usage["output_tokens"]
.as_u64()
.expect("worker output tokens");
let total = usage["total_tokens"].as_u64().expect("worker total tokens");
assert!(total >= input + output, "{task_completed}");
assert!(
usage["tool_calls"].as_u64().unwrap_or_default() >= 1,
"{task_completed}"
);
assert!(usage["duration_ms"].is_u64(), "{task_completed}");
assert_eq!(
usage["token_source"], "provider_reported",
"{task_completed}"
);
assert!(
usage["result_ref"]
.as_str()
.is_some_and(|target| target.starts_with("agent:")),
"{task_completed}"
);
let run_completed = events
.iter()
.find(|event| event["type"] == "run_completed")
.expect("run_completed event");
let run_usage = &run_completed["usage"];
assert_eq!(run_usage["total_tokens"], total, "{run_completed}");
assert_eq!(run_usage["input_tokens"], input, "{run_completed}");
assert_eq!(run_usage["output_tokens"], output, "{run_completed}");
assert_eq!(
run_usage["tool_calls"], usage["tool_calls"],
"{run_completed}"
);
assert_eq!(run_usage["tasks_reported"], 1, "{run_completed}");
assert_eq!(payload["usage"]["total_tokens"], total, "{payload}");
assert_eq!(
payload["execution"]["usage"]["input_tokens"], input,
"{payload}"
);
assert_eq!(
payload["execution"]["leaf_results"][0]["usage"]["input_tokens"], input,
"{payload}"
);
}
#[test]
fn provider_usage_presence_preserves_a_real_zero_receipt() {
assert!(!provider_usage_was_reported(None, None, None));
assert!(provider_usage_was_reported(Some(0), Some(0), Some(0)));
assert!(provider_usage_was_reported(None, Some(0), None));
}
#[test]
fn run_usage_totals_reconcile_task_telemetry() {
let task_usage = |total: u64, calls: u32| WorkflowTaskUsage {
input_tokens: Some(total / 2),
output_tokens: Some(total - total / 2),
total_tokens: Some(total),
cost_microusd: Some(total),
tool_calls: Some(calls),
duration_ms: Some(7),
result_ref: None,
token_source: Some(WorkflowTokenSource::ProviderReported),
};
let record = |agent_id: &str, usage: Option<WorkflowTaskUsage>| RuntimeTaskRecord {
agent_id: agent_id.to_string(),
label: None,
role: None,
status: IrWorkflowRunStatus::Succeeded,
output: None,
schema_error: None,
usage,
};
let records = vec![
record("a", Some(task_usage(100, 2))),
record("b", Some(task_usage(60, 1))),
record("c", None),
];
let totals = run_usage_totals(&records).expect("totals");
assert_eq!(totals.total_tokens, Some(160));
assert_eq!(totals.cost_microusd, Some(160));
assert_eq!(totals.input_tokens, Some(80));
assert_eq!(totals.output_tokens, Some(80));
assert_eq!(totals.tool_calls, Some(3));
assert_eq!(totals.tasks_reported, 2);
assert!(run_usage_totals(&[]).is_none());
assert!(run_usage_totals(&[record("d", None)]).is_none());
}
#[test]
fn run_and_ir_usage_keep_unknown_distinct_from_reported_zero() {
let record = |agent_id: &str, usage: WorkflowTaskUsage| RuntimeTaskRecord {
agent_id: agent_id.to_string(),
label: Some(agent_id.to_string()),
role: None,
status: IrWorkflowRunStatus::Succeeded,
output: None,
schema_error: None,
usage: Some(usage),
};
let unknown = WorkflowTaskUsage {
tool_calls: Some(1),
duration_ms: Some(4),
..WorkflowTaskUsage::default()
};
let reported_zero = WorkflowTaskUsage {
input_tokens: Some(0),
output_tokens: Some(0),
total_tokens: Some(0),
cost_microusd: Some(0),
tool_calls: Some(0),
duration_ms: Some(0),
token_source: Some(WorkflowTokenSource::ProviderReported),
..WorkflowTaskUsage::default()
};
let unknown_totals = run_usage_totals(&[record("unknown", unknown.clone())])
.expect("tool/duration receipt still creates run usage");
assert_eq!(unknown_totals.input_tokens, None);
assert_eq!(unknown_totals.output_tokens, None);
assert_eq!(unknown_totals.total_tokens, None);
assert_eq!(unknown_totals.tool_calls, Some(1));
let zero_totals = run_usage_totals(&[record("zero", reported_zero.clone())])
.expect("reported zero receipt");
assert_eq!(zero_totals.input_tokens, Some(0));
assert_eq!(zero_totals.output_tokens, Some(0));
assert_eq!(zero_totals.total_tokens, Some(0));
assert_eq!(zero_totals.cost_microusd, Some(0));
assert_eq!(zero_totals.tool_calls, Some(0));
let unknown_ir = workflow_usage_from_task(&unknown);
assert_eq!(unknown_ir.input_tokens, None);
assert_eq!(unknown_ir.output_tokens, None);
assert_eq!(unknown_ir.cost_microusd, None);
let zero_ir = workflow_usage_from_task(&reported_zero);
assert_eq!(zero_ir.input_tokens, Some(0));
assert_eq!(zero_ir.output_tokens, Some(0));
assert_eq!(zero_ir.cost_microusd, Some(0));
let mixed = run_usage_totals(&[record("zero", reported_zero), record("unknown", unknown)])
.expect("mixed receipts");
assert_eq!(
mixed.total_tokens,
Some(0),
"a missing contributor keeps the observed subtotal"
);
assert_eq!(mixed.input_tokens, Some(0));
assert_eq!(mixed.output_tokens, Some(0));
assert_eq!(mixed.cost_microusd, Some(0));
assert_eq!(mixed.tool_calls, Some(1));
}
#[test]
fn workflow_ui_event_usage_telemetry_serde_round_trip() {
let event = WorkflowUiEvent::at(
7,
WorkflowUiEventKind::TaskCompleted {
task_id: "child-1".to_string(),
status: IrWorkflowRunStatus::Succeeded,
usage: Some(WorkflowTaskUsage {
input_tokens: Some(128),
output_tokens: Some(32),
total_tokens: Some(160),
cost_microusd: Some(42),
tool_calls: Some(2),
duration_ms: Some(42),
result_ref: Some("agent:child-1".to_string()),
token_source: Some(WorkflowTokenSource::ProviderReported),
}),
},
);
let json = serde_json::to_value(&event).expect("serialize");
assert_eq!(json["usage"]["total_tokens"], 160);
let parsed: WorkflowUiEvent = serde_json::from_value(json).expect("deserialize round trip");
match parsed.kind {
WorkflowUiEventKind::TaskCompleted {
usage: Some(usage), ..
} => {
assert_eq!(usage.total_tokens, Some(160));
assert_eq!(usage.cost_microusd, Some(42));
assert_eq!(usage.tool_calls, Some(2));
assert_eq!(usage.duration_ms, Some(42));
}
other => panic!("expected task_completed with usage, got {other:?}"),
}
let legacy_task: WorkflowUiEvent = serde_json::from_str(
r#"{"at_ms":5,"type":"task_completed","task_id":"child-1","status":"succeeded"}"#,
)
.expect("legacy task_completed parses");
match legacy_task.kind {
WorkflowUiEventKind::TaskCompleted { usage: None, .. } => {}
other => panic!("expected legacy task_completed without usage, got {other:?}"),
}
let legacy_run: WorkflowUiEvent = serde_json::from_str(
r#"{"at_ms":6,"type":"run_completed","status":"completed","error":null}"#,
)
.expect("legacy run_completed parses");
match legacy_run.kind {
WorkflowUiEventKind::RunCompleted { usage: None, .. } => {}
other => panic!("expected legacy run_completed without usage, got {other:?}"),
}
let plain = serde_json::to_value(WorkflowUiEvent::at(
8,
WorkflowUiEventKind::RunCompleted {
status: WorkflowRunStatus::Completed,
error: None,
usage: None,
},
))
.expect("serialize plain run_completed");
assert!(plain.get("usage").is_none(), "{plain}");
}
#[test]
fn run_record_event_retention_is_bounded() {
let mut record = WorkflowRunRecord::new("workflow_tail".to_string(), None, None, None);
for index in 0..(WORKFLOW_RUN_EVENTS_MAX_RETAINED + 5) {
record.push_event(WorkflowUiEvent::at(
index as u64,
WorkflowUiEventKind::Log {
message: format!("log {index}"),
},
));
}
assert_eq!(record.events.len(), WORKFLOW_RUN_EVENTS_MAX_RETAINED);
assert_eq!(record.events_dropped, 5);
assert_eq!(
record.events_total,
(WORKFLOW_RUN_EVENTS_MAX_RETAINED + 5) as u64
);
let first = &record.events[0];
assert_eq!(first.at_ms, 5);
let summary = record.summary();
assert_eq!(summary.event_count, WORKFLOW_RUN_EVENTS_MAX_RETAINED + 5);
assert_eq!(summary.events_dropped, 5);
assert_eq!(summary.last_event_type.as_deref(), Some("log"));
}
#[test]
fn workflow_status_payload_bounds_oversized_run_records() {
let tmp = tempfile::tempdir().expect("tempdir");
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_big_run".to_string();
let mut record = WorkflowRunRecord::new(run_id.clone(), None, None, None);
record.status = WorkflowRunStatus::Completed;
for index in 0..200u64 {
record.push_event(WorkflowUiEvent::at(
index,
WorkflowUiEventKind::Log {
message: format!("fan-out event {index}"),
},
));
}
for index in 0..60 {
record.progress.push(format!("progress line {index}"));
}
for index in 0..40u64 {
record.dispatch_failures.push(WorkflowDispatchFailure {
at_ms: index,
label: Some(format!("rejected-{index}")),
phase: Some("fan-out".to_string()),
message: format!("dispatch rejected {index}: {}", "x".repeat(2_000)),
});
}
record.result = Some(json!({ "blob": "r".repeat(10_000) }));
record.execution = Some(IrWorkflowExecution {
status: IrWorkflowRunStatus::Succeeded,
usage: WorkflowUsage::default(),
memo_usage: WorkflowMemoUsage::default(),
leaf_results: vec![LeafResult {
leaf_id: "inspect".to_string(),
task_id: "agent_1".to_string(),
role: None,
profile: None,
status: IrWorkflowRunStatus::Succeeded,
usage: WorkflowUsage::default(),
memo_usage: WorkflowMemoUsage::default(),
output: Some("o".repeat(5_000)),
artifacts: Vec::new(),
schema_error: None,
}],
branch_results: Vec::new(),
control_node_results: Vec::new(),
});
state
.runs
.lock()
.expect("runs")
.insert(run_id.clone(), record);
let result = workflow_result_for(&run_id, state).expect("status result");
assert!(
result.content.len() < WORKFLOW_RESULT_MAX_CHARS,
"bounded payload must stay under {WORKFLOW_RESULT_MAX_CHARS} chars, got {}",
result.content.len()
);
let payload: Value = serde_json::from_str(&result.content).expect("payload json");
let events = payload["events"].as_array().expect("events");
assert_eq!(events.len(), WORKFLOW_RESULT_EVENTS_TAIL, "{payload}");
assert!(
payload["events_note"]
.as_str()
.is_some_and(|note| note.contains("workflow-runs.jsonl")),
"{payload}"
);
assert_eq!(events[0]["at_ms"], 150);
let progress = payload["progress"].as_array().expect("progress");
assert_eq!(progress.len(), WORKFLOW_RESULT_PROGRESS_TAIL, "{payload}");
assert!(payload.get("progress_note").is_some(), "{payload}");
let dispatch_failures = payload["dispatch_failures"]
.as_array()
.expect("dispatch failures");
assert_eq!(
dispatch_failures.len(),
WORKFLOW_RESULT_DISPATCH_FAILURES_TAIL,
"{payload}"
);
assert_eq!(dispatch_failures[0]["at_ms"], 28);
assert!(
dispatch_failures.iter().all(|failure| failure["message"]
.as_str()
.is_some_and(|message| message.chars().count()
<= WORKFLOW_RESULT_DISPATCH_FAILURE_FIELD_MAX_CHARS)),
"{dispatch_failures:?}"
);
assert!(
payload["dispatch_failures_note"]
.as_str()
.is_some_and(|note| note.contains("workflow-runs.jsonl")),
"{payload}"
);
assert_eq!(payload["result"]["truncated"], true, "{payload}");
assert!(
payload["result"]["preview"]
.as_str()
.is_some_and(|preview| preview.chars().count() <= WORKFLOW_RESULT_VALUE_MAX_CHARS),
"{payload}"
);
assert!(
payload["result"]["full_detail"]
.as_str()
.is_some_and(|path| path.contains("workflow-runs.jsonl")),
"{payload}"
);
let leaf_output = payload["execution"]["leaf_results"][0]["output"]
.as_str()
.expect("leaf output");
assert!(
leaf_output.contains("leaf output truncated"),
"{leaf_output}"
);
assert!(leaf_output.len() < 1_000, "{leaf_output}");
let metadata = result.metadata.expect("metadata");
assert_eq!(metadata["events_returned"], WORKFLOW_RESULT_EVENTS_TAIL);
assert_eq!(metadata["events_omitted"], 150);
assert_eq!(metadata["event_count"], 200);
assert_eq!(
metadata["dispatch_failures_returned"],
WORKFLOW_RESULT_DISPATCH_FAILURES_TAIL
);
assert_eq!(metadata["dispatch_failures_omitted"], 28);
assert_eq!(metadata["truncated"], true);
assert!(
metadata["journal_path"]
.as_str()
.is_some_and(|path| path.contains("workflow-runs.jsonl")),
"{metadata}"
);
}
#[test]
fn workflow_status_payload_keeps_small_records_intact() {
let tmp = tempfile::tempdir().expect("tempdir");
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_small_run".to_string();
let mut record = WorkflowRunRecord::new(run_id.clone(), None, None, None);
record.status = WorkflowRunStatus::Completed;
record.push_event(WorkflowUiEvent::at(
1,
WorkflowUiEventKind::PhaseStarted {
title: "scan".to_string(),
},
));
record.result = Some(json!({ "ok": true }));
state
.runs
.lock()
.expect("runs")
.insert(run_id.clone(), record);
let result = workflow_result_for(&run_id, state).expect("status result");
let payload: Value = serde_json::from_str(&result.content).expect("payload json");
assert_eq!(payload["events"].as_array().map(Vec::len), Some(1));
assert_eq!(payload["result"], json!({ "ok": true }));
assert!(payload.get("events_note").is_none(), "{payload}");
assert!(payload.get("progress_note").is_none(), "{payload}");
let metadata = result.metadata.expect("metadata");
assert_eq!(metadata["truncated"], false);
assert_eq!(metadata["events_omitted"], 0);
assert_eq!(metadata["events_returned"], 1);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn workflow_cancel_interrupts_vm_and_blocks_further_spawns() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
let (client, calls) = fake_chat_client("child done").await;
let (event_tx, mut event_rx) = tokio::sync::mpsc::channel(256);
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
Some(event_tx),
manager.clone(),
);
let tool = WorkflowTool::new(manager.clone(), runtime);
let started = tool
.execute(
json!({
"action": "start",
"script": r#"
let n = 0;
while (n < 20) {
await task({ description: `task ${n}`, type: 'explore', allowedTools: [] });
n++;
}
return n;
"#
}),
&ctx,
)
.await
.expect("workflow start");
let run_id = started
.metadata
.as_ref()
.and_then(|metadata| metadata.get("run_id"))
.and_then(Value::as_str)
.expect("run_id metadata");
tokio::time::timeout(std::time::Duration::from_secs(3), async {
while calls.load(Ordering::SeqCst) == 0 {
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
})
.await
.expect("workflow should spawn at least one child before cancel");
let calls_before_cancel = calls.load(Ordering::SeqCst);
assert!(calls_before_cancel >= 1);
let cancelled = tool
.execute(json!({"action": "cancel", "run_id": run_id}), &ctx)
.await
.expect("workflow cancel");
let cancelled_payload: Value =
serde_json::from_str(&cancelled.content).expect("cancel json");
assert_eq!(cancelled_payload["status"], "cancelled");
assert!(
cancelled_payload["events"]
.as_array()
.is_some_and(|events| events.iter().any(|event| event["type"] == "run_cancelled")),
"cancel receipt must include the authoritative terminal event: {cancelled_payload}"
);
let mut streamed_cancel = false;
while let Ok(event) = event_rx.try_recv() {
if let Event::WorkflowUi { event, .. } = event
&& event["type"] == "run_cancelled"
{
streamed_cancel = true;
}
}
assert!(
streamed_cancel,
"cancel must stream a terminal UI event after any racing completion"
);
let first_event_count = cancelled_payload["events"]
.as_array()
.expect("events")
.len();
let first_completed_at = cancelled_payload["completed_at_ms"].clone();
let cancelled_again = tool
.execute(json!({"action": "cancel", "run_id": run_id}), &ctx)
.await
.expect("second workflow cancel is a no-op");
let cancelled_again_payload: Value =
serde_json::from_str(&cancelled_again.content).expect("second cancel json");
assert_eq!(cancelled_again_payload["status"], "cancelled");
assert_eq!(
cancelled_again_payload["events"]
.as_array()
.expect("events")
.len(),
first_event_count,
"second cancel must not append a duplicate terminal event"
);
assert_eq!(
cancelled_again_payload["completed_at_ms"], first_completed_at,
"second cancel must preserve the original completion time"
);
tokio::time::sleep(std::time::Duration::from_millis(700)).await;
let calls_after_cancel = calls.load(Ordering::SeqCst);
assert!(
calls_after_cancel <= calls_before_cancel + 1,
"cancelled workflow kept spawning children: before={calls_before_cancel} after={calls_after_cancel}"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn workflow_budget_spent_delegates_to_manager_scope() {
let _retry_guard = workflow_test_retry_guard();
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (client, _calls) = fake_chat_client("budgeted").await;
let runtime = SubAgentRuntime::new(
client,
"deepseek-v4-flash".to_string(),
ctx.clone(),
true,
None,
manager.clone(),
);
let tool = WorkflowTool::new(manager.clone(), runtime);
let result = tool
.execute(
json!({
"action": "run",
"token_budget": 1000,
"script": r#"
await task({ description: 'budgeted work', type: 'explore', allowedTools: [] });
return { spent: budget.spent(), total: budget.total, remaining: budget.remaining() };
"#
}),
&ctx,
)
.await
.expect("budget workflow should complete");
let payload: Value = serde_json::from_str(&result.content).expect("json result");
assert_eq!(payload["status"], "completed", "{payload}");
assert_eq!(payload["result"]["spent"], 2);
assert_eq!(payload["result"]["total"], 1000);
assert_eq!(payload["result"]["remaining"], 998);
}
#[tokio::test]
async fn identical_budget_snapshots_emit_a_live_heartbeat_without_journal_duplication() {
let tmp = tempfile::tempdir().expect("tempdir");
let ctx = ToolContext::new(tmp.path().to_path_buf());
let manager = new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
let (event_tx, mut event_rx) = tokio::sync::mpsc::channel(8);
let runtime = SubAgentRuntime::new(
stub_client(),
"deepseek-v4-flash".to_string(),
ctx,
true,
Some(event_tx),
manager.clone(),
);
let state = WorkflowWorkspaceState::open(tmp.path());
let run_id = "workflow_budget_heartbeat".to_string();
state.runs.lock().expect("runs").insert(
run_id.clone(),
WorkflowRunRecord::new(run_id.clone(), None, None, None),
);
let driver = SubAgentWorkflowDriver::new(
run_id.clone(),
manager,
runtime,
state.clone(),
Some(1_000),
WorkflowFleetBinding::None,
Vec::new(),
);
let snapshot = BudgetSnapshot {
total: Some(1_000),
spent: 0,
};
driver.record_budget_snapshot(snapshot);
driver.record_budget_snapshot(snapshot);
let mut streamed = 0;
while let Ok(Event::WorkflowUi { event, .. }) = event_rx.try_recv() {
if event["type"] == "budget_updated" {
streamed += 1;
}
}
assert_eq!(
streamed, 2,
"unchanged budget still refreshes the live panel"
);
let recorded = state
.runs
.lock()
.expect("runs")
.get(&run_id)
.expect("run")
.events
.iter()
.filter(|event| event.event_type() == "budget_updated")
.count();
assert_eq!(recorded, 1, "heartbeat must not grow the durable journal");
}
fn stub_client() -> DeepSeekClient {
let _ = rustls::crypto::ring::default_provider().install_default();
let config = crate::config::Config {
api_key: Some("test-key".to_string()),
..crate::config::Config::default()
};
DeepSeekClient::new(&config).expect("stub client should construct")
}
async fn fake_chat_client(response_text: &str) -> (DeepSeekClient, Arc<AtomicUsize>) {
let (client, calls, _) = fake_chat_client_capturing(response_text).await;
(client, calls)
}
async fn fake_chat_client_responses(
response_texts: &[&str],
) -> (DeepSeekClient, Arc<AtomicUsize>) {
let (client, calls, _) = fake_chat_client_capturing_responses(response_texts).await;
(client, calls)
}
async fn fake_chat_client_capturing(
response_text: &str,
) -> (DeepSeekClient, Arc<AtomicUsize>, Arc<Mutex<Vec<Value>>>) {
fake_chat_client_capturing_responses(&[response_text]).await
}
async fn fake_chat_client_capturing_responses(
response_texts: &[&str],
) -> (DeepSeekClient, Arc<AtomicUsize>, Arc<Mutex<Vec<Value>>>) {
assert!(
!response_texts.is_empty(),
"fake chat client needs at least one response"
);
let calls = Arc::new(AtomicUsize::new(0));
let bodies = Arc::new(Mutex::new(Vec::new()));
let response_texts = Arc::new(
response_texts
.iter()
.map(|response| (*response).to_string())
.collect::<Vec<_>>(),
);
let app = Router::new().route(
"/{*path}",
post({
let calls = Arc::clone(&calls);
let bodies = Arc::clone(&bodies);
let response_texts = Arc::clone(&response_texts);
move |Json(body): Json<Value>| {
let calls = Arc::clone(&calls);
let bodies = Arc::clone(&bodies);
let response_texts = Arc::clone(&response_texts);
async move {
bodies.lock().expect("capture body").push(body);
let attempt = calls.fetch_add(1, Ordering::SeqCst) + 1;
let response_text = if response_texts.len() == 1 {
response_texts[0].clone()
} else {
response_texts
.get(attempt - 1)
.unwrap_or_else(|| {
panic!(
"fake chat server received call {attempt} but only {} responses were supplied",
response_texts.len()
)
})
.clone()
};
Json(json!({
"id": format!("chatcmpl-workflow-test-{attempt}"),
"model": "deepseek-v4-flash",
"choices": [{
"index": 0,
"message": {
"role": "assistant",
"content": response_text
},
"finish_reason": "stop"
}],
"usage": {
"prompt_tokens": 1,
"completion_tokens": 1,
"total_tokens": 2
}
}))
}
}
}),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind fake chat server");
let addr = listener.local_addr().expect("fake chat server addr");
tokio::spawn(async move {
let _ = axum::serve(listener, app).await;
});
let config = crate::config::Config {
api_key: Some("test-key".to_string()),
base_url: Some(format!("http://{addr}/v1")),
..crate::config::Config::default()
};
(
DeepSeekClient::new(&config).expect("fake chat client"),
calls,
bodies,
)
}
fn workflow_test_retry_guard() -> std::sync::MutexGuard<'static, ()> {
let guard = crate::retry_status::test_guard();
crate::retry_status::clear();
crate::retry_status::clear_rate_limit();
guard
}
}