use crate::adapters::primary::PeriplonSDKClient;
use crate::dsl::hooks::{ErrorRecovery, HooksExecutor};
use crate::dsl::loop_context::{substitute_task_variables, LoopContext};
use crate::dsl::message_bus::MessageBus;
use crate::dsl::notifications::{NotificationContext, NotificationManager};
use crate::dsl::schema::{AgentSpec, CollectionSource, DSLWorkflow, FileFormat, LoopSpec};
use crate::dsl::state::{StatePersistence, WorkflowState};
use crate::dsl::task_graph::{TaskGraph, TaskStatus};
use crate::error::{Error, Result};
use crate::options::AgentOptions;
use futures::StreamExt;
use std::collections::HashMap;
use std::sync::Arc;
use std::time::Instant;
use tokio::sync::{Mutex, Semaphore};
struct ExecutionContext<'a> {
workflow_inputs: &'a HashMap<String, serde_json::Value>,
agents: &'a Arc<Mutex<HashMap<String, PeriplonSDKClient>>>,
task_graph: &'a Arc<Mutex<TaskGraph>>,
state: &'a Arc<Mutex<Option<WorkflowState>>>,
workflow_name: &'a Arc<String>,
json_output: bool,
}
pub struct DSLExecutor {
workflow: DSLWorkflow,
agents: HashMap<String, PeriplonSDKClient>,
task_graph: TaskGraph,
message_bus: Arc<MessageBus>,
state: Option<WorkflowState>,
state_persistence: Option<StatePersistence>,
resolved_inputs: HashMap<String, serde_json::Value>,
notification_manager: Arc<NotificationManager>,
workflow_start_time: Option<Instant>,
json_output: bool,
}
impl DSLExecutor {
pub fn new(workflow: DSLWorkflow) -> Result<Self> {
let resolved_inputs = Self::resolve_workflow_inputs(&workflow);
let notification_manager = Arc::new(NotificationManager::new());
Ok(DSLExecutor {
workflow,
agents: HashMap::new(),
task_graph: TaskGraph::new(),
message_bus: Arc::new(MessageBus::new()),
state: None,
state_persistence: None,
resolved_inputs,
notification_manager,
workflow_start_time: None,
json_output: false,
})
}
pub fn task_graph(&self) -> &TaskGraph {
&self.task_graph
}
fn resolve_workflow_inputs(workflow: &DSLWorkflow) -> HashMap<String, serde_json::Value> {
let mut resolved = HashMap::new();
for (key, input_spec) in &workflow.inputs {
if let Some(default_value) = &input_spec.default {
resolved.insert(key.clone(), default_value.clone());
}
}
resolved
}
fn substitute_variables(
text: &str,
workflow_inputs: &HashMap<String, serde_json::Value>,
task_inputs: &HashMap<String, serde_json::Value>,
) -> String {
let mut result = text.to_string();
for (key, value) in workflow_inputs {
let placeholder = format!("${{workflow.{}}}", key);
let value_str = match value {
serde_json::Value::String(s) => s.clone(),
serde_json::Value::Number(n) => n.to_string(),
serde_json::Value::Bool(b) => b.to_string(),
_ => serde_json::to_string(value).unwrap_or_default(),
};
result = result.replace(&placeholder, &value_str);
}
for (key, value) in task_inputs {
let placeholder = format!("${{task.{}}}", key);
let value_str = match value {
serde_json::Value::String(s) => s.clone(),
serde_json::Value::Number(n) => n.to_string(),
serde_json::Value::Bool(b) => b.to_string(),
_ => serde_json::to_string(value).unwrap_or_default(),
};
result = result.replace(&placeholder, &value_str);
}
result
}
fn create_notification_context(
&self,
task_id: &str,
task_status: &str,
duration: Option<std::time::Duration>,
error_msg: Option<&str>,
) -> NotificationContext {
let mut context = NotificationContext::new();
for (key, value) in &self.resolved_inputs {
let value_str = match value {
serde_json::Value::String(s) => s.clone(),
serde_json::Value::Number(n) => n.to_string(),
serde_json::Value::Bool(b) => b.to_string(),
_ => serde_json::to_string(value).unwrap_or_default(),
};
context = context.with_workflow_var(key, value_str);
}
context = context
.with_metadata("task_id", task_id)
.with_metadata("task_status", task_status)
.with_metadata("workflow_name", &self.workflow.name);
if let Some(dur) = duration {
context = context.with_metadata("duration_secs", dur.as_secs().to_string());
context = context.with_metadata("duration_human", format!("{:.2}s", dur.as_secs_f64()));
}
if let Some(err) = error_msg {
context = context.with_metadata("error", err);
}
for secret_name in self.workflow.secrets.keys() {
context = context.with_secret(secret_name, format!("${{secret.{}}}", secret_name));
}
context
}
pub fn enable_state_persistence(&mut self, state_dir: Option<&str>) -> Result<()> {
let persistence = if let Some(dir) = state_dir {
StatePersistence::new(dir)?
} else {
StatePersistence::default()
};
self.state_persistence = Some(persistence);
if let Some(ref persistence) = self.state_persistence {
if persistence.has_state(&self.workflow.name) {
println!(
"Found existing state for workflow '{}' - will resume if possible",
self.workflow.name
);
}
}
Ok(())
}
pub fn set_json_output(&mut self, json: bool) {
self.json_output = json;
}
pub fn try_resume(&mut self) -> Result<bool> {
if let Some(ref persistence) = self.state_persistence {
if persistence.has_state(&self.workflow.name) {
let saved_state = persistence.load_state(&self.workflow.name)?;
if saved_state.can_resume() {
println!(
"Resuming workflow '{}' from checkpoint (progress: {:.1}%)",
self.workflow.name,
saved_state.get_progress() * 100.0
);
self.state = Some(saved_state);
return Ok(true);
} else {
println!(
"Cannot resume workflow '{}' - status: {:?}",
self.workflow.name, saved_state.status
);
}
}
}
Ok(false)
}
fn checkpoint_state(&mut self) -> Result<()> {
if let (Some(ref mut state), Some(ref persistence)) =
(&mut self.state, &self.state_persistence)
{
persistence.save_state(state)?;
}
Ok(())
}
pub fn message_bus(&self) -> Arc<MessageBus> {
self.message_bus.clone()
}
pub fn get_state(&self) -> Option<&WorkflowState> {
self.state.as_ref()
}
pub async fn initialize(&mut self) -> Result<()> {
for agent_name in self.workflow.agents.keys() {
self.message_bus.register_agent(agent_name.clone()).await?;
}
if let Some(comm_config) = &self.workflow.communication {
for (channel_name, channel_spec) in &comm_config.channels {
self.message_bus
.create_channel(
channel_name.clone(),
channel_spec.description.clone(),
channel_spec.participants.clone(),
channel_spec.message_format.clone(),
)
.await?;
}
}
let mut var_context = crate::dsl::variables::VariableContext::new();
for (key, value) in &self.resolved_inputs {
var_context.insert(&crate::dsl::variables::Scope::Workflow, key, value.clone());
}
for (name, spec) in &self.workflow.agents {
let options = self.agent_spec_to_options(
spec,
self.workflow.cwd.as_deref(),
self.workflow.create_cwd,
&var_context,
)?;
let mut client = PeriplonSDKClient::new(options);
client.connect(None).await?;
self.agents.insert(name.clone(), client);
}
let tasks: Vec<(String, crate::dsl::schema::TaskSpec)> = self
.workflow
.tasks
.iter()
.map(|(name, spec)| (name.clone(), spec.clone()))
.collect();
for (name, spec) in tasks {
self.add_hierarchical_task(&name, &spec, None)?;
}
println!(
"Initialized {} agents and {} channels",
self.message_bus.agent_count().await,
self.message_bus.channel_count().await
);
if self.state_persistence.is_some() && self.state.is_none() {
let mut state =
WorkflowState::new(self.workflow.name.clone(), self.workflow.version.clone());
for task_id in self.task_graph.topological_sort()? {
state.update_task_status(&task_id, TaskStatus::Pending);
}
self.state = Some(state);
println!("Initialized workflow state tracking");
}
Ok(())
}
fn add_hierarchical_task(
&mut self,
task_name: &str,
task_spec: &crate::dsl::schema::TaskSpec,
parent_name: Option<&str>,
) -> Result<()> {
let mut task_spec_flat = task_spec.clone();
let has_subtasks = !task_spec.subtasks.is_empty();
let sibling_names: std::collections::HashSet<String> = if let Some(parent) = parent_name {
if let Some(parent_task) = self.workflow.tasks.get(parent) {
parent_task
.subtasks
.iter()
.flat_map(|map| map.keys())
.map(|s| s.to_string())
.collect()
} else {
std::collections::HashSet::new()
}
} else {
std::collections::HashSet::new()
};
if let Some(parent) = parent_name {
task_spec_flat.depends_on = task_spec_flat
.depends_on
.iter()
.map(|dep| {
if sibling_names.contains(dep) {
format!("{}.{}", parent, dep)
} else {
dep.clone()
}
})
.collect();
task_spec_flat.parallel_with = task_spec_flat
.parallel_with
.iter()
.map(|task| {
if sibling_names.contains(task) {
format!("{}.{}", parent, task)
} else {
task.clone()
}
})
.collect();
}
let has_loop = task_spec_flat.loop_spec.is_some();
if !has_loop {
task_spec_flat.subtasks.clear();
}
let is_executable = !has_subtasks || has_loop;
if let Some(parent) = parent_name {
let parent_is_executable = self
.workflow
.tasks
.get(parent)
.map(|parent_spec| parent_spec.subtasks.is_empty())
.unwrap_or(true);
if parent_is_executable && !task_spec_flat.depends_on.contains(&parent.to_string()) {
task_spec_flat.depends_on.push(parent.to_string());
}
}
let mut resolved_dependencies = Vec::new();
for dep in &task_spec_flat.depends_on {
if let Some(dep_task) = self.workflow.tasks.get(dep) {
let dep_has_loop = dep_task.loop_spec.is_some();
let is_non_executable_parent = !dep_task.subtasks.is_empty() && !dep_has_loop;
if is_non_executable_parent {
for subtask_map in &dep_task.subtasks {
for subtask_name in subtask_map.keys() {
resolved_dependencies.push(format!("{}.{}", dep, subtask_name));
}
}
} else {
resolved_dependencies.push(dep.clone());
}
} else {
resolved_dependencies.push(dep.clone());
}
}
task_spec_flat.depends_on = resolved_dependencies;
if is_executable {
self.task_graph
.add_task(task_name.to_string(), task_spec_flat);
}
if has_subtasks && !has_loop {
for subtask_map in &task_spec.subtasks {
for (subtask_name, subtask_spec) in subtask_map {
let full_subtask_name = format!("{}.{}", task_name, subtask_name);
let mut inherited_subtask_spec = subtask_spec.clone();
inherited_subtask_spec.inherit_from_parent(task_spec);
self.add_hierarchical_task(
&full_subtask_name,
&inherited_subtask_spec,
Some(task_name),
)?;
}
}
}
Ok(())
}
pub async fn execute(&mut self) -> Result<()> {
if let Some(workflows) = self.workflow.workflows.values().next() {
if let Some(hooks) = &workflows.hooks {
if !hooks.pre_workflow.is_empty() {
HooksExecutor::execute_pre_workflow(&hooks.pre_workflow, &self.workflow.name)
.await?;
}
}
}
let execution_result = self.execute_tasks().await;
if let Some(workflows) = self.workflow.workflows.values().next() {
if let Some(hooks) = &workflows.hooks {
if !hooks.post_workflow.is_empty() {
let _ = HooksExecutor::execute_post_workflow(
&hooks.post_workflow,
&self.workflow.name,
)
.await;
}
}
}
if let Err(ref e) = execution_result {
if let Some(workflows) = self.workflow.workflows.values().next() {
if let Some(hooks) = &workflows.hooks {
if !hooks.on_error.is_empty() {
let _ = HooksExecutor::execute_error(
&hooks.on_error,
&self.workflow.name,
&e.to_string(),
)
.await;
}
}
}
if let Some(ref mut state) = self.state {
state.mark_failed();
let _ = self.checkpoint_state();
}
} else {
if let Some(ref mut state) = self.state {
state.mark_completed();
let _ = self.checkpoint_state();
}
}
execution_result
}
async fn execute_tasks(&mut self) -> Result<()> {
self.workflow_start_time = Some(Instant::now());
let order = self.task_graph.topological_sort()?;
println!("Executing workflow: {}", self.workflow.name);
println!("Task execution order: {:?}", order);
if let Some(notif_config) = &self.workflow.notifications {
if notif_config.notify_on_start && !notif_config.default_channels.is_empty() {
let context = self.create_notification_context("workflow", "started", None, None);
let spec = crate::dsl::schema::NotificationSpec::Structured {
message: format!("Workflow '{}' started", self.workflow.name),
title: Some(format!("Workflow Started: {}", self.workflow.name)),
priority: Some(crate::dsl::schema::NotificationPriority::Normal),
channels: notif_config.default_channels.clone(),
metadata: HashMap::new(),
};
if let Err(e) = self.notification_manager.send(&spec, &context).await {
eprintln!("Warning: Failed to send workflow start notification: {}", e);
}
}
}
let agents = Arc::new(Mutex::new(std::mem::take(&mut self.agents)));
let task_graph = Arc::new(Mutex::new(std::mem::take(&mut self.task_graph)));
let state = Arc::new(Mutex::new(self.state.take()));
let workflow_inputs = Arc::new(self.resolved_inputs.clone());
let state_persistence = Arc::new(self.state_persistence.clone());
let workflow_name = Arc::new(self.workflow.name.clone());
let mut processed = std::collections::HashSet::new();
for task_id in order {
if processed.contains(&task_id) {
continue;
}
if let Some(ref workflow_state) = *state.lock().await {
if workflow_state.get_task_status(&task_id) == Some(TaskStatus::Completed) {
println!("Skipping already completed task: {}", task_id);
processed.insert(task_id.clone());
continue;
}
}
let parallel_tasks = {
let graph = task_graph.lock().await;
graph.get_parallel_tasks(&task_id)
};
if parallel_tasks.is_empty() {
self.execute_task_parallel(
&task_id,
agents.clone(),
task_graph.clone(),
state.clone(),
workflow_inputs.clone(),
state_persistence.clone(),
)
.await?;
processed.insert(task_id.clone());
} else {
println!(
"Executing tasks in parallel: {} with {:?}",
task_id, parallel_tasks
);
let mut handles = vec![];
{
let task_id = task_id.clone();
let agents = agents.clone();
let graph = task_graph.clone();
let workflow_state = state.clone();
let inputs = workflow_inputs.clone();
let persistence = state_persistence.clone();
let wf_name = workflow_name.clone();
let json_out = self.json_output;
let handle = tokio::spawn(async move {
execute_task_static(
task_id,
agents,
graph,
workflow_state,
inputs,
persistence,
wf_name,
json_out,
)
.await
});
handles.push(handle);
}
for parallel_id in ¶llel_tasks {
let status = {
let graph = task_graph.lock().await;
graph.get_task_status(parallel_id)
};
if status == Some(TaskStatus::Pending) {
let parallel_id = parallel_id.clone();
let agents = agents.clone();
let graph = task_graph.clone();
let workflow_state = state.clone();
let inputs = workflow_inputs.clone();
let persistence = state_persistence.clone();
let wf_name = workflow_name.clone();
let json_out = self.json_output;
let handle = tokio::spawn(async move {
execute_task_static(
parallel_id,
agents,
graph,
workflow_state,
inputs,
persistence,
wf_name,
json_out,
)
.await
});
handles.push(handle);
}
}
for handle in handles {
match handle.await {
Ok(Ok(())) => {}
Ok(Err(e)) => return Err(e),
Err(e) => return Err(Error::InvalidInput(format!("Task panicked: {}", e))),
}
}
processed.insert(task_id.clone());
for parallel_id in parallel_tasks {
processed.insert(parallel_id);
}
}
}
self.agents = Arc::try_unwrap(agents)
.map_err(|_| Error::InvalidInput("Failed to unwrap agents".to_string()))?
.into_inner();
self.task_graph = Arc::try_unwrap(task_graph)
.map_err(|_| Error::InvalidInput("Failed to unwrap task graph".to_string()))?
.into_inner();
self.state = Arc::try_unwrap(state)
.map_err(|_| Error::InvalidInput("Failed to unwrap state".to_string()))?
.into_inner();
if self.state.is_some() {
let _ = self.checkpoint_state();
}
println!("Workflow execution completed");
if let Some(notif_config) = &self.workflow.notifications {
if notif_config.notify_on_workflow_completion
&& !notif_config.default_channels.is_empty()
{
let duration = self.workflow_start_time.map(|start| start.elapsed());
let context =
self.create_notification_context("workflow", "completed", duration, None);
let spec = crate::dsl::schema::NotificationSpec::Structured {
message: format!("Workflow '{}' completed successfully", self.workflow.name),
title: Some(format!("Workflow Completed: {}", self.workflow.name)),
priority: Some(crate::dsl::schema::NotificationPriority::Normal),
channels: notif_config.default_channels.clone(),
metadata: HashMap::new(),
};
if let Err(e) = self.notification_manager.send(&spec, &context).await {
eprintln!(
"Warning: Failed to send workflow completion notification: {}",
e
);
}
}
}
Ok(())
}
async fn execute_task_parallel(
&self,
task_id: &str,
agents: Arc<Mutex<HashMap<String, PeriplonSDKClient>>>,
task_graph: Arc<Mutex<TaskGraph>>,
state: Arc<Mutex<Option<WorkflowState>>>,
workflow_inputs: Arc<HashMap<String, serde_json::Value>>,
state_persistence: Arc<Option<StatePersistence>>,
) -> Result<()> {
let workflow_name = Arc::new(self.workflow.name.clone());
execute_task_static(
task_id.to_string(),
agents,
task_graph,
state,
workflow_inputs,
state_persistence,
workflow_name,
self.json_output,
)
.await
}
fn agent_spec_to_options(
&self,
spec: &AgentSpec,
workflow_cwd: Option<&str>,
workflow_create_cwd: Option<bool>,
var_context: &crate::dsl::variables::VariableContext,
) -> Result<AgentOptions> {
let add_dirs = spec
.permissions
.allowed_directories
.iter()
.map(std::path::PathBuf::from)
.collect();
let interpolated_workflow_cwd = workflow_cwd.map(|cwd| {
var_context
.interpolate(cwd)
.unwrap_or_else(|_| cwd.to_string())
});
let interpolated_agent_cwd = spec
.cwd
.as_ref()
.map(|cwd| var_context.interpolate(cwd).unwrap_or_else(|_| cwd.clone()));
let cwd = interpolated_agent_cwd
.as_deref()
.or(interpolated_workflow_cwd.as_deref())
.map(std::path::PathBuf::from);
let create_cwd = spec.create_cwd.or(workflow_create_cwd).unwrap_or(false);
let options = AgentOptions {
allowed_tools: spec.tools.clone(),
model: spec.model.clone(),
max_turns: spec.max_turns,
permission_mode: Some(spec.permissions.mode.clone()),
add_dirs,
cwd,
create_cwd,
..Default::default()
};
Ok(options)
}
pub async fn shutdown(&mut self) -> Result<()> {
println!("Shutting down executor...");
for (name, agent) in &mut self.agents {
println!("Disconnecting agent: {}", name);
agent.disconnect().await?;
}
println!("Executor shutdown complete");
Ok(())
}
pub fn get_workflow_info(&self) -> (&str, &str) {
(&self.workflow.name, &self.workflow.version)
}
pub fn get_task_count(&self) -> usize {
self.task_graph.task_count()
}
pub fn is_complete(&self) -> bool {
self.task_graph.is_complete()
}
}
#[allow(clippy::too_many_arguments)]
async fn execute_task_static(
task_id: String,
agents: Arc<Mutex<HashMap<String, PeriplonSDKClient>>>,
task_graph: Arc<Mutex<TaskGraph>>,
state: Arc<Mutex<Option<WorkflowState>>>,
workflow_inputs: Arc<HashMap<String, serde_json::Value>>,
state_persistence: Arc<Option<StatePersistence>>,
workflow_name: Arc<String>,
json_output: bool,
) -> Result<()> {
let (spec, recovery_strategy) = {
let graph = task_graph.lock().await;
let task_node = graph
.get_task(&task_id)
.ok_or_else(|| Error::InvalidInput(format!("Task '{}' not found", task_id)))?;
let strategy = ErrorRecovery::get_strategy_from_spec(task_node.spec.on_error.as_ref());
(task_node.spec.clone(), strategy)
};
println!("Executing task: {} - {}", task_id, spec.description);
if let Some(ref loop_spec) = spec.loop_spec {
println!(
"Task '{}' has loop specification - delegating to loop executor",
task_id
);
let ctx = ExecutionContext {
workflow_inputs: &workflow_inputs,
agents: &agents,
task_graph: &task_graph,
state: &state,
workflow_name: &workflow_name,
json_output,
};
return execute_task_with_loop(&task_id, &spec, loop_spec, &ctx).await;
}
if let Some(ref condition) = spec.condition {
let condition_met = {
let graph = task_graph.lock().await;
let workflow_state = state.lock().await;
let state_ref = workflow_state.as_ref();
evaluate_condition(condition, &graph, state_ref)
};
if !condition_met {
println!("Task '{}' condition not met - skipping", task_id);
{
let mut graph = task_graph.lock().await;
graph.update_task_status(&task_id, TaskStatus::Skipped)?;
}
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.update_task_status(&task_id, TaskStatus::Skipped);
}
return Ok(());
}
println!(
"Task '{}' condition met - proceeding with execution",
task_id
);
}
let mut error_attempt = 0;
let mut dod_attempt = 0;
let mut last_dod_feedback: Option<String> = None;
loop {
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.record_task_attempt(&task_id);
}
{
let mut graph = task_graph.lock().await;
graph.update_task_status(&task_id, TaskStatus::Running)?;
}
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.update_task_status(&task_id, TaskStatus::Running);
}
let mut var_context = crate::dsl::variables::VariableContext::new();
for (key, value) in workflow_inputs.iter() {
var_context.insert(&crate::dsl::variables::Scope::Workflow, key, value.clone());
}
let interpolated_description = var_context
.interpolate(&spec.description)
.unwrap_or_else(|_| spec.description.clone());
let task_description = if let Some(ref feedback) = last_dod_feedback {
format!("{}\n\n{}", interpolated_description, feedback)
} else {
interpolated_description
};
match execute_task_attempt(
&task_id,
&task_description,
&spec,
&workflow_inputs,
&agents,
error_attempt,
&state,
&workflow_name,
json_output,
)
.await
{
Ok(task_output) => {
if let Some(ref dod) = spec.definition_of_done {
println!("Checking definition of done for task: {}", task_id);
let dod_results =
check_definition_of_done(dod, task_output.as_deref(), &var_context).await;
let all_met = dod_results.iter().all(|r| r.met);
if !all_met {
let mut unmet_feedback = format_unmet_criteria(&dod_results);
unmet_feedback = enhance_feedback_with_permission_hints(
unmet_feedback,
task_output.as_deref().unwrap_or(""),
&dod_results,
dod.auto_elevate_permissions,
);
println!("Definition of done not met for task '{}':", task_id);
println!("{}", unmet_feedback);
dod_attempt += 1;
if dod_attempt <= dod.max_retries {
println!(
"Retrying task '{}' (DoD attempt {}/{})",
task_id, dod_attempt, dod.max_retries
);
if dod.auto_elevate_permissions
&& detect_permission_issue(
task_output.as_deref().unwrap_or(""),
&dod_results,
)
{
println!(
" 🔓 Auto-elevating permissions to 'bypassPermissions' for retry..."
);
if let Some(ref agent_id) = spec.agent {
let mut agents_guard = agents.lock().await;
if let Some(agent) = agents_guard.get_mut(agent_id) {
if let Err(e) =
agent.set_permission_mode("bypassPermissions").await
{
eprintln!(
"Warning: Failed to elevate permissions: {}",
e
);
} else {
println!(" ✓ Permissions elevated successfully");
}
}
}
}
last_dod_feedback = Some(unmet_feedback);
continue;
} else {
println!(
"Task '{}' exhausted all DoD retries ({}/{})",
task_id, dod_attempt, dod.max_retries
);
if dod.fail_on_unmet {
{
let mut graph = task_graph.lock().await;
graph.update_task_status(&task_id, TaskStatus::Failed)?;
}
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.update_task_status(&task_id, TaskStatus::Failed);
workflow_state.record_task_error(
&task_id,
&format!(
"Definition of done not met after {} retries",
dod.max_retries
),
);
}
return Err(Error::InvalidInput(format!(
"Task '{}' failed: definition of done not met after {} retries",
task_id, dod.max_retries
)));
}
}
} else {
println!("✓ Definition of done met for task: {}", task_id);
}
}
{
let mut graph = task_graph.lock().await;
graph.update_task_status(&task_id, TaskStatus::Completed)?;
}
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.update_task_status(&task_id, TaskStatus::Completed);
if let Some(ref output) = task_output {
workflow_state.record_task_result(&task_id, output);
}
}
if let (Some(ref workflow_state), Some(ref persistence)) =
(&*state.lock().await, &*state_persistence)
{
if let Err(e) = persistence.save_state(workflow_state) {
eprintln!(
"Warning: Failed to checkpoint state after task '{}': {}",
task_id, e
);
}
}
println!("Task completed: {}", task_id);
if let Some(on_complete) = &spec.on_complete {
if let Some(notify_spec) = &on_complete.notify {
let message = match notify_spec {
crate::dsl::NotificationSpec::Simple(msg) => msg.clone(),
crate::dsl::NotificationSpec::Structured { message, .. } => {
message.clone()
}
};
println!("Notification: {}", message);
}
}
return Ok(());
}
Err(e) => {
println!(
"Task '{}' failed (attempt {}): {}",
task_id,
error_attempt + 1,
e
);
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.record_task_error(&task_id, &e.to_string());
}
error_attempt += 1;
if !ErrorRecovery::should_retry(&recovery_strategy, error_attempt) {
{
let mut graph = task_graph.lock().await;
graph.update_task_status(&task_id, TaskStatus::Failed)?;
}
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.update_task_status(&task_id, TaskStatus::Failed);
}
if let Some(fallback_agent) =
ErrorRecovery::get_fallback_agent(&recovery_strategy)
{
println!("Attempting fallback with agent: {}", fallback_agent);
{
let mut graph = task_graph.lock().await;
graph.update_task_status(&task_id, TaskStatus::Running)?;
}
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.update_task_status(&task_id, TaskStatus::Running);
}
match execute_task_with_agent(
&task_id,
&spec,
&agents,
fallback_agent,
error_attempt,
json_output,
)
.await
{
Ok(()) => {
println!("Task completed with fallback agent: {}", task_id);
{
let mut graph = task_graph.lock().await;
graph.update_task_status(&task_id, TaskStatus::Completed)?;
}
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state
.update_task_status(&task_id, TaskStatus::Completed);
}
return Ok(());
}
Err(fallback_err) => {
println!("Fallback agent also failed: {}", fallback_err);
{
let mut graph = task_graph.lock().await;
graph.update_task_status(&task_id, TaskStatus::Failed)?;
}
if let Some(ref mut workflow_state) = *state.lock().await {
workflow_state.update_task_status(&task_id, TaskStatus::Failed);
workflow_state.record_task_error(
&task_id,
&format!(
"Primary and fallback failed: {} / {}",
e, fallback_err
),
);
}
return Err(fallback_err);
}
}
}
return Err(e);
}
let retry_delay = calculate_retry_delay(&recovery_strategy, error_attempt);
println!(
"Retrying task '{}' (attempt {}) in {}s...",
task_id,
error_attempt + 1,
retry_delay
);
tokio::time::sleep(tokio::time::Duration::from_secs(retry_delay)).await;
}
}
}
}
async fn execute_script_task(
_task_id: &str,
script_spec: &crate::dsl::schema::ScriptSpec,
workflow_inputs: &HashMap<String, serde_json::Value>,
task_inputs: &HashMap<String, serde_json::Value>,
attempt: u32,
) -> Result<Option<String>> {
use crate::dsl::schema::ScriptLanguage;
use tokio::process::Command;
let (interpreter, args): (&str, Vec<&str>) = match script_spec.language {
ScriptLanguage::Python => ("python3", vec!["-c"]),
ScriptLanguage::JavaScript => ("node", vec!["-e"]),
ScriptLanguage::Bash => ("bash", vec!["-c"]),
ScriptLanguage::Ruby => ("ruby", vec!["-e"]),
ScriptLanguage::Perl => ("perl", vec!["-e"]),
};
let raw_script_content = if let Some(content) = &script_spec.content {
content.clone()
} else if let Some(file_path) = &script_spec.file {
let interpolated_file_path =
DSLExecutor::substitute_variables(file_path, workflow_inputs, task_inputs);
tokio::fs::read_to_string(&interpolated_file_path)
.await
.map_err(|e| {
Error::InvalidInput(format!(
"Failed to read script file '{}': {}",
interpolated_file_path, e
))
})?
} else {
return Err(Error::InvalidInput(
"Script must specify either 'content' or 'file'".to_string(),
));
};
let script_content =
DSLExecutor::substitute_variables(&raw_script_content, workflow_inputs, task_inputs);
let mut cmd = Command::new(interpreter);
cmd.args(&args);
cmd.arg(&script_content);
if let Some(working_dir) = &script_spec.working_dir {
let interpolated_working_dir =
DSLExecutor::substitute_variables(working_dir, workflow_inputs, task_inputs);
cmd.current_dir(interpolated_working_dir);
}
for (key, value) in &script_spec.env {
let interpolated_value =
DSLExecutor::substitute_variables(value, workflow_inputs, task_inputs);
cmd.env(key, interpolated_value);
}
cmd.stdout(std::process::Stdio::piped());
cmd.stderr(std::process::Stdio::piped());
if attempt > 0 {
println!(
" [Retry {}] Executing {} script...",
attempt,
format!("{:?}", script_spec.language).to_lowercase()
);
} else {
println!(
" Executing {} script...",
format!("{:?}", script_spec.language).to_lowercase()
);
}
let output = if let Some(timeout_secs) = script_spec.timeout_secs {
tokio::time::timeout(std::time::Duration::from_secs(timeout_secs), cmd.output())
.await
.map_err(|_| {
Error::InvalidInput(format!(
"Script execution timed out after {} seconds",
timeout_secs
))
})?
.map_err(|e| Error::InvalidInput(format!("Failed to execute script: {}", e)))?
} else {
cmd.output()
.await
.map_err(|e| Error::InvalidInput(format!("Failed to execute script: {}", e)))?
};
let stdout = String::from_utf8_lossy(&output.stdout).to_string();
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
if !stdout.is_empty() {
print!("{}", stdout);
}
if !stderr.is_empty() {
eprint!("{}", stderr);
}
if !output.status.success() {
return Err(Error::InvalidInput(format!(
"Script failed with exit code: {:?}",
output.status.code()
)));
}
let combined_output = format!("{}{}", stdout, stderr);
Ok(if combined_output.is_empty() {
None
} else {
Some(combined_output)
})
}
async fn execute_command_task(
_task_id: &str,
command_spec: &crate::dsl::schema::CommandSpec,
workflow_inputs: &HashMap<String, serde_json::Value>,
task_inputs: &HashMap<String, serde_json::Value>,
attempt: u32,
) -> Result<Option<String>> {
use tokio::process::Command;
let executable =
DSLExecutor::substitute_variables(&command_spec.executable, workflow_inputs, task_inputs);
let args: Vec<String> = command_spec
.args
.iter()
.map(|arg| DSLExecutor::substitute_variables(arg, workflow_inputs, task_inputs))
.collect();
let mut cmd = Command::new(&executable);
cmd.args(&args);
if let Some(working_dir) = &command_spec.working_dir {
let working_dir =
DSLExecutor::substitute_variables(working_dir, workflow_inputs, task_inputs);
cmd.current_dir(working_dir);
}
for (key, value) in &command_spec.env {
let value = DSLExecutor::substitute_variables(value, workflow_inputs, task_inputs);
cmd.env(key, value);
}
if command_spec.capture_stdout {
cmd.stdout(std::process::Stdio::piped());
}
if command_spec.capture_stderr {
cmd.stderr(std::process::Stdio::piped());
}
if attempt > 0 {
println!(
" [Retry {}] Executing command: {} {}",
attempt,
executable,
args.join(" ")
);
} else {
println!(" Executing command: {} {}", executable, args.join(" "));
}
let output = if let Some(timeout_secs) = command_spec.timeout_secs {
tokio::time::timeout(std::time::Duration::from_secs(timeout_secs), cmd.output())
.await
.map_err(|_| {
Error::InvalidInput(format!(
"Command execution timed out after {} seconds",
timeout_secs
))
})?
.map_err(|e| Error::InvalidInput(format!("Failed to execute command: {}", e)))?
} else {
cmd.output()
.await
.map_err(|e| Error::InvalidInput(format!("Failed to execute command: {}", e)))?
};
let stdout = if command_spec.capture_stdout {
String::from_utf8_lossy(&output.stdout).to_string()
} else {
String::new()
};
let stderr = if command_spec.capture_stderr {
String::from_utf8_lossy(&output.stderr).to_string()
} else {
String::new()
};
if !stdout.is_empty() {
print!("{}", stdout);
}
if !stderr.is_empty() {
eprint!("{}", stderr);
}
if !output.status.success() {
return Err(Error::InvalidInput(format!(
"Command '{}' failed with exit code: {:?}",
executable,
output.status.code()
)));
}
let combined_output = format!("{}{}", stdout, stderr);
Ok(if combined_output.is_empty() {
None
} else {
Some(combined_output)
})
}
#[allow(clippy::too_many_arguments)]
async fn execute_task_attempt(
_task_id: &str,
task_description: &str,
_spec: &crate::dsl::schema::TaskSpec,
workflow_inputs: &HashMap<String, serde_json::Value>,
agents: &Arc<Mutex<HashMap<String, PeriplonSDKClient>>>,
attempt: u32,
workflow_state: &Arc<Mutex<Option<crate::dsl::state::WorkflowState>>>,
workflow_name: &str,
json_output: bool,
) -> Result<Option<String>> {
if let Some(script_spec) = &_spec.script {
return execute_script_task(
_task_id,
script_spec,
workflow_inputs,
&_spec.inputs,
attempt,
)
.await;
}
if let Some(command_spec) = &_spec.command {
return execute_command_task(
_task_id,
command_spec,
workflow_inputs,
&_spec.inputs,
attempt,
)
.await;
}
let agent_name = _spec.agent.as_ref().ok_or_else(|| {
Error::InvalidInput(format!("No agent specified for task '{}'", _task_id))
})?;
let workflow_context = if _spec.inject_context {
if let Some(ref state) = *workflow_state.lock().await {
let has_workflow_context =
!state.get_completed_tasks().is_empty() || !state.get_failed_tasks().is_empty();
let has_output_file = _spec.output.is_some();
if has_workflow_context || has_output_file {
state.build_context_summary(workflow_name, None, _spec.output.as_deref())
} else {
String::new()
}
} else {
String::new()
}
} else {
String::new()
};
let enhanced_description = if _spec.inject_context && !workflow_context.is_empty() {
format!("{}\n{}", workflow_context, task_description)
} else {
task_description.to_string()
};
let mut agents_guard = agents.lock().await;
let agent = agents_guard
.get_mut(agent_name)
.ok_or_else(|| Error::InvalidInput(format!("Agent '{}' not found", agent_name)))?;
agent.query(&enhanced_description).await?;
let stream = agent.receive_response()?;
futures::pin_mut!(stream);
let mut output = String::new();
while let Some(msg) = stream.next().await {
output.push_str(&format!("{:?}\n", msg));
if attempt > 0 {
let prefix_str = format!("Retry {}", attempt);
let formatted =
crate::dsl::message_formatter::format_message(&msg, json_output, Some(&prefix_str));
println!("{}", formatted);
} else {
let formatted = crate::dsl::message_formatter::format_message(&msg, json_output, None);
println!("{}", formatted);
}
}
Ok(if output.is_empty() {
None
} else {
Some(output)
})
}
async fn execute_task_with_agent(
_task_id: &str,
spec: &crate::dsl::schema::TaskSpec,
agents: &Arc<Mutex<HashMap<String, PeriplonSDKClient>>>,
agent_name: &str,
attempt: u32,
json_output: bool,
) -> Result<()> {
let mut agents_guard = agents.lock().await;
let agent = agents_guard
.get_mut(agent_name)
.ok_or_else(|| Error::InvalidInput(format!("Agent '{}' not found", agent_name)))?;
agent.query(&spec.description).await?;
let stream = agent.receive_response()?;
futures::pin_mut!(stream);
while let Some(msg) = stream.next().await {
if attempt > 0 {
let prefix_str = "Fallback".to_string();
let formatted =
crate::dsl::message_formatter::format_message(&msg, json_output, Some(&prefix_str));
println!("{}", formatted);
} else {
let formatted = crate::dsl::message_formatter::format_message(&msg, json_output, None);
println!("{}", formatted);
}
}
Ok(())
}
fn calculate_retry_delay(
recovery_strategy: &crate::dsl::hooks::RecoveryStrategy,
attempt: u32,
) -> u64 {
use crate::dsl::hooks::RecoveryStrategy;
match recovery_strategy {
RecoveryStrategy::Retry { config, .. } => {
let base_delay = config.as_ref().map(|c| c.retry_delay_secs).unwrap_or(1);
let use_backoff = config
.as_ref()
.map(|c| c.exponential_backoff)
.unwrap_or(false);
if use_backoff {
let exponent = attempt.saturating_sub(1);
let delay = base_delay * 2u64.pow(exponent);
delay.min(60)
} else {
base_delay
}
}
_ => 1, }
}
fn evaluate_condition(
condition: &crate::dsl::schema::ConditionSpec,
task_graph: &TaskGraph,
workflow_state: Option<&WorkflowState>,
) -> bool {
use crate::dsl::schema::ConditionSpec;
match condition {
ConditionSpec::Single(cond) => evaluate_single_condition(cond, task_graph, workflow_state),
ConditionSpec::And { and } => and
.iter()
.all(|c| evaluate_condition(c, task_graph, workflow_state)),
ConditionSpec::Or { or } => or
.iter()
.any(|c| evaluate_condition(c, task_graph, workflow_state)),
ConditionSpec::Not { not } => !evaluate_condition(not, task_graph, workflow_state),
}
}
fn evaluate_single_condition(
condition: &crate::dsl::schema::Condition,
task_graph: &TaskGraph,
workflow_state: Option<&WorkflowState>,
) -> bool {
use crate::dsl::schema::{Condition, TaskStatusCondition};
match condition {
Condition::TaskStatus { task, status } => {
let task_status = task_graph.get_task_status(task);
matches!(
(task_status, status),
(Some(TaskStatus::Completed), TaskStatusCondition::Completed)
| (Some(TaskStatus::Failed), TaskStatusCondition::Failed)
| (Some(TaskStatus::Running), TaskStatusCondition::Running)
| (Some(TaskStatus::Pending), TaskStatusCondition::Pending)
| (Some(TaskStatus::Skipped), TaskStatusCondition::Skipped)
)
}
Condition::StateEquals { key, value } => {
if let Some(state) = workflow_state {
state.get_metadata(key) == Some(value)
} else {
false
}
}
Condition::StateExists { key } => {
if let Some(state) = workflow_state {
state.get_metadata(key).is_some()
} else {
false
}
}
Condition::Always => true,
Condition::Never => false,
}
}
async fn check_definition_of_done(
dod: &crate::dsl::schema::DefinitionOfDone,
_task_output: Option<&str>,
var_context: &crate::dsl::variables::VariableContext,
) -> Vec<crate::dsl::schema::CriterionResult> {
let mut results = Vec::new();
for criterion in &dod.criteria {
let result = check_criterion(criterion, _task_output, var_context).await;
results.push(result);
}
results
}
async fn check_criterion(
criterion: &crate::dsl::schema::DoneCriterion,
_task_output: Option<&str>,
var_context: &crate::dsl::variables::VariableContext,
) -> crate::dsl::schema::CriterionResult {
use crate::dsl::schema::{CriterionResult, DoneCriterion, OutputSource};
match criterion {
DoneCriterion::FileExists { path, description } => {
let interpolated_path = var_context
.interpolate(path)
.unwrap_or_else(|_| path.clone());
let exists = std::path::Path::new(&interpolated_path).exists();
CriterionResult {
met: exists,
description: description.clone(),
details: if exists {
format!("File '{}' exists", interpolated_path)
} else {
format!("File '{}' does not exist", interpolated_path)
},
}
}
DoneCriterion::FileContains {
path,
pattern,
description,
} => {
let interpolated_path = var_context
.interpolate(path)
.unwrap_or_else(|_| path.clone());
let interpolated_pattern = var_context
.interpolate(pattern)
.unwrap_or_else(|_| pattern.clone());
match std::fs::read_to_string(&interpolated_path) {
Ok(content) => {
let contains = match regex::Regex::new(&interpolated_pattern) {
Ok(re) => re.is_match(&content),
Err(_) => content.contains(&interpolated_pattern), };
CriterionResult {
met: contains,
description: description.clone(),
details: if contains {
format!(
"File '{}' contains pattern '{}'",
interpolated_path, interpolated_pattern
)
} else {
format!(
"File '{}' does not contain pattern '{}'",
interpolated_path, interpolated_pattern
)
},
}
}
Err(e) => CriterionResult {
met: false,
description: description.clone(),
details: format!("Failed to read file '{}': {}", interpolated_path, e),
},
}
}
DoneCriterion::FileNotContains {
path,
pattern,
description,
} => {
let interpolated_path = var_context
.interpolate(path)
.unwrap_or_else(|_| path.clone());
let interpolated_pattern = var_context
.interpolate(pattern)
.unwrap_or_else(|_| pattern.clone());
match std::fs::read_to_string(&interpolated_path) {
Ok(content) => {
let contains = match regex::Regex::new(&interpolated_pattern) {
Ok(re) => re.is_match(&content),
Err(_) => content.contains(&interpolated_pattern), };
let not_contains = !contains;
CriterionResult {
met: not_contains,
description: description.clone(),
details: if not_contains {
format!(
"File '{}' does not contain pattern '{}'",
interpolated_path, interpolated_pattern
)
} else {
format!(
"File '{}' contains pattern '{}' (should not)",
interpolated_path, interpolated_pattern
)
},
}
}
Err(e) => CriterionResult {
met: false,
description: description.clone(),
details: format!("Failed to read file '{}': {}", interpolated_path, e),
},
}
}
DoneCriterion::CommandSucceeds {
command,
args,
description,
working_dir,
} => {
let interpolated_command = var_context
.interpolate(command)
.unwrap_or_else(|_| command.clone());
let interpolated_args: Vec<String> = args
.iter()
.map(|arg| var_context.interpolate(arg).unwrap_or_else(|_| arg.clone()))
.collect();
let interpolated_working_dir = working_dir
.as_ref()
.map(|dir| var_context.interpolate(dir).unwrap_or_else(|_| dir.clone()));
let mut cmd = tokio::process::Command::new(&interpolated_command);
cmd.args(&interpolated_args);
if let Some(ref dir) = interpolated_working_dir {
cmd.current_dir(dir);
}
match cmd.output().await {
Ok(output) => {
let success = output.status.success();
let details = if success {
format!("Command '{}' succeeded", interpolated_command)
} else {
let stdout = String::from_utf8_lossy(&output.stdout);
let stderr = String::from_utf8_lossy(&output.stderr);
let mut error_msg = format!(
"Command '{}' failed with exit code: {}\n\n",
interpolated_command,
output.status.code().unwrap_or(-1)
);
if !stderr.is_empty() {
error_msg.push_str("STDERR:\n");
error_msg.push_str("```\n");
error_msg.push_str(&stderr);
error_msg.push_str("```\n\n");
}
if !stdout.is_empty() {
error_msg.push_str("STDOUT:\n");
error_msg.push_str("```\n");
error_msg.push_str(&stdout);
error_msg.push_str("```\n");
}
error_msg
};
CriterionResult {
met: success,
description: description.clone(),
details,
}
}
Err(e) => CriterionResult {
met: false,
description: description.clone(),
details: format!(
"Failed to execute command '{}': {}",
interpolated_command, e
),
},
}
}
DoneCriterion::OutputMatches {
source,
pattern,
description,
} => {
let interpolated_pattern = var_context
.interpolate(pattern)
.unwrap_or_else(|_| pattern.clone());
let content = match source {
OutputSource::File { path } => {
let interpolated_path = var_context
.interpolate(path)
.unwrap_or_else(|_| path.clone());
std::fs::read_to_string(&interpolated_path).ok()
}
OutputSource::TaskOutput => _task_output.map(|s| s.to_string()),
};
match content {
Some(text) => {
let matches = text.contains(&interpolated_pattern);
CriterionResult {
met: matches,
description: description.clone(),
details: if matches {
format!("Output matches pattern '{}'", interpolated_pattern)
} else {
format!("Output does not match pattern '{}'", interpolated_pattern)
},
}
}
None => CriterionResult {
met: false,
description: description.clone(),
details: "Failed to read output".to_string(),
},
}
}
DoneCriterion::DirectoryExists { path, description } => {
let interpolated_path = var_context
.interpolate(path)
.unwrap_or_else(|_| path.clone());
let exists = std::path::Path::new(&interpolated_path).is_dir();
CriterionResult {
met: exists,
description: description.clone(),
details: if exists {
format!("Directory '{}' exists", interpolated_path)
} else {
format!("Directory '{}' does not exist", interpolated_path)
},
}
}
DoneCriterion::TestsPassed {
command,
args,
description,
} => {
let interpolated_command = var_context
.interpolate(command)
.unwrap_or_else(|_| command.clone());
let interpolated_args: Vec<String> = args
.iter()
.map(|arg| var_context.interpolate(arg).unwrap_or_else(|_| arg.clone()))
.collect();
let mut cmd = tokio::process::Command::new(&interpolated_command);
cmd.args(&interpolated_args);
match cmd.output().await {
Ok(output) => {
let success = output.status.success();
let details = if success {
format!("Tests passed ({})", interpolated_command)
} else {
let stdout = String::from_utf8_lossy(&output.stdout);
let stderr = String::from_utf8_lossy(&output.stderr);
let mut error_msg = format!(
"Tests failed ({}) with exit code: {}\n\n",
interpolated_command,
output.status.code().unwrap_or(-1)
);
if !stderr.is_empty() {
error_msg.push_str("STDERR:\n");
error_msg.push_str("```\n");
error_msg.push_str(&stderr);
error_msg.push_str("```\n\n");
}
if !stdout.is_empty() {
error_msg.push_str("STDOUT:\n");
error_msg.push_str("```\n");
error_msg.push_str(&stdout);
error_msg.push_str("```\n");
}
error_msg
};
CriterionResult {
met: success,
description: description.clone(),
details,
}
}
Err(e) => CriterionResult {
met: false,
description: description.clone(),
details: format!("Failed to run tests '{}': {}", interpolated_command, e),
},
}
}
}
}
fn format_unmet_criteria(results: &[crate::dsl::schema::CriterionResult]) -> String {
let unmet: Vec<_> = results.iter().filter(|r| !r.met).collect();
if unmet.is_empty() {
return String::new();
}
let mut feedback = String::from("\n\n=== DEFINITION OF DONE - UNMET CRITERIA ===\n\n");
feedback.push_str("The following criteria were not met:\n\n");
for (i, result) in unmet.iter().enumerate() {
feedback.push_str(&format!(
"{}. {}\n Status: ✗ FAILED\n Details: {}\n\n",
i + 1,
result.description,
result.details
));
}
feedback.push_str("Please address these issues and retry the task.\n");
feedback
}
fn detect_permission_issue(output: &str, results: &[crate::dsl::schema::CriterionResult]) -> bool {
let permission_keywords = [
"permission",
"permissions",
"write access",
"read access",
"file write",
"cannot create",
"cannot write",
"access denied",
"forbidden",
];
let output_lower = output.to_lowercase();
let has_permission_mention = permission_keywords
.iter()
.any(|keyword| output_lower.contains(keyword));
let has_file_failures = results.iter().any(|r| {
!r.met
&& (r.description.to_lowercase().contains("file")
|| r.details.to_lowercase().contains("does not exist")
|| r.details.to_lowercase().contains("not found"))
});
has_permission_mention || has_file_failures
}
fn enhance_feedback_with_permission_hints(
mut feedback: String,
output: &str,
results: &[crate::dsl::schema::CriterionResult],
auto_elevate: bool,
) -> String {
if detect_permission_issue(output, results) {
feedback.push_str("\n⚠️ PERMISSION HINT:\n");
feedback.push_str("The failure appears to be related to file access or permissions.\n");
if auto_elevate {
feedback.push_str(
"Auto-elevation is enabled - 'bypassPermissions' mode will be granted on retry.\n",
);
feedback
.push_str("All tool operations will be automatically approved for this retry.\n");
} else {
feedback.push_str("Consider:\n");
feedback.push_str("1. Ensuring required files exist before checking\n");
feedback.push_str("2. Creating necessary files if they don't exist\n");
feedback.push_str("3. Requesting write permissions if needed\n\n");
feedback.push_str(
"TIP: Add 'auto_elevate_permissions: true' to the definition_of_done config\n",
);
feedback.push_str(" to automatically grant enhanced permissions on retry.\n");
}
}
feedback
}
fn extract_json_path(value: &serde_json::Value, path: &str) -> Result<serde_json::Value> {
let mut current = value.clone();
for segment in path.split('.') {
if segment.is_empty() {
continue;
}
if let Some(bracket_pos) = segment.find('[') {
let key = &segment[..bracket_pos];
let index_str = &segment[bracket_pos + 1..segment.len() - 1];
let index: usize = index_str.parse().map_err(|_| {
Error::InvalidInput(format!("Invalid array index in JSON path: {}", segment))
})?;
if !key.is_empty() {
current = current
.get(key)
.ok_or_else(|| {
Error::InvalidInput(format!("JSON path key '{}' not found", key))
})?
.clone();
}
current = current
.get(index)
.ok_or_else(|| {
Error::InvalidInput(format!("JSON path index {} out of bounds", index))
})?
.clone();
} else {
current = current
.get(segment)
.ok_or_else(|| {
Error::InvalidInput(format!("JSON path key '{}' not found", segment))
})?
.clone();
}
}
Ok(current)
}
async fn resolve_collection(
collection: &CollectionSource,
state: Option<&WorkflowState>,
) -> Result<Vec<serde_json::Value>> {
match collection {
CollectionSource::State { key } => {
let state = state.ok_or_else(|| {
Error::InvalidInput(
"Workflow state not available for state-based collection".to_string(),
)
})?;
let value = state
.get_metadata(key)
.ok_or_else(|| Error::InvalidInput(format!("State key '{}' not found", key)))?;
match value {
serde_json::Value::Array(arr) => Ok(arr.clone()),
_ => Err(Error::InvalidInput(format!(
"State key '{}' is not an array",
key
))),
}
}
CollectionSource::File { path, format } => {
let content = std::fs::read_to_string(path).map_err(|e| {
Error::InvalidInput(format!("Failed to read collection file '{}': {}", path, e))
})?;
match format {
FileFormat::Json => {
let value: serde_json::Value = serde_json::from_str(&content).map_err(|e| {
Error::InvalidInput(format!("Failed to parse JSON file '{}': {}", path, e))
})?;
match value {
serde_json::Value::Array(arr) => Ok(arr),
_ => Err(Error::InvalidInput(format!(
"JSON file '{}' does not contain an array",
path
))),
}
}
FileFormat::JsonLines => {
let items: Result<Vec<serde_json::Value>> = content
.lines()
.enumerate()
.map(|(i, line)| {
serde_json::from_str(line).map_err(|e| {
Error::InvalidInput(format!(
"Failed to parse JSON line {} in '{}': {}",
i + 1,
path,
e
))
})
})
.collect();
items
}
FileFormat::Csv => {
let items: Vec<serde_json::Value> = content
.lines()
.skip(1) .map(|line| {
let fields: Vec<&str> = line.split(',').collect();
serde_json::Value::Array(
fields
.iter()
.map(|f| serde_json::Value::String(f.to_string()))
.collect(),
)
})
.collect();
Ok(items)
}
FileFormat::Lines => {
let items: Vec<serde_json::Value> = content
.lines()
.map(|line| serde_json::Value::String(line.to_string()))
.collect();
Ok(items)
}
}
}
CollectionSource::Range { start, end, step } => {
let step_size = step.unwrap_or(1);
let mut items = Vec::new();
let mut current = *start;
while current < *end {
items.push(serde_json::Value::Number(current.into()));
current += step_size;
}
Ok(items)
}
CollectionSource::Inline { items } => {
Ok(items.clone())
}
CollectionSource::Http {
url,
method,
headers,
body,
format,
json_path,
} => {
let client = reqwest::Client::new();
let mut request = match method.to_uppercase().as_str() {
"GET" => client.get(url),
"POST" => client.post(url),
"PUT" => client.put(url),
"DELETE" => client.delete(url),
"PATCH" => client.patch(url),
_ => {
return Err(Error::InvalidInput(format!(
"Unsupported HTTP method: {}",
method
)))
}
};
if let Some(headers_map) = headers {
for (key, value) in headers_map {
request = request.header(key, value);
}
}
if let Some(body_content) = body {
request = request.body(body_content.clone());
}
let response = request.send().await.map_err(|e| {
Error::InvalidInput(format!("HTTP request to '{}' failed: {}", url, e))
})?;
if !response.status().is_success() {
return Err(Error::InvalidInput(format!(
"HTTP request to '{}' returned status: {}",
url,
response.status()
)));
}
let response_text = response.text().await.map_err(|e| {
Error::InvalidInput(format!("Failed to read response from '{}': {}", url, e))
})?;
let mut value: serde_json::Value = match format {
FileFormat::Json => serde_json::from_str(&response_text).map_err(|e| {
Error::InvalidInput(format!(
"Failed to parse JSON response from '{}': {}",
url, e
))
})?,
FileFormat::JsonLines => {
let items: Result<Vec<serde_json::Value>> = response_text
.lines()
.enumerate()
.map(|(i, line)| {
serde_json::from_str(line).map_err(|e| {
Error::InvalidInput(format!(
"Failed to parse JSON line {} from '{}': {}",
i + 1,
url,
e
))
})
})
.collect();
serde_json::Value::Array(items?)
}
FileFormat::Lines => {
let items: Vec<serde_json::Value> = response_text
.lines()
.map(|line| serde_json::Value::String(line.to_string()))
.collect();
serde_json::Value::Array(items)
}
FileFormat::Csv => {
let items: Vec<serde_json::Value> = response_text
.lines()
.skip(1) .map(|line| {
let fields: Vec<&str> = line.split(',').collect();
serde_json::Value::Array(
fields
.iter()
.map(|f| serde_json::Value::String(f.to_string()))
.collect(),
)
})
.collect();
serde_json::Value::Array(items)
}
};
if let Some(path) = json_path {
value = extract_json_path(&value, path)?;
}
match value {
serde_json::Value::Array(arr) => Ok(arr),
_ => Err(Error::InvalidInput(format!(
"HTTP response from '{}' does not contain an array",
url
))),
}
}
}
}
async fn execute_task_with_loop(
task_id: &str,
spec: &crate::dsl::schema::TaskSpec,
loop_spec: &LoopSpec,
ctx: &ExecutionContext<'_>,
) -> Result<()> {
match loop_spec {
LoopSpec::ForEach {
collection,
iterator,
parallel,
..
} => {
let items = {
let state_guard = ctx.state.lock().await;
resolve_collection(collection, state_guard.as_ref()).await?
};
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.init_loop(task_id, Some(items.len()));
}
}
println!(
"Executing ForEach loop for task '{}': {} items",
task_id,
items.len()
);
if *parallel {
let max_parallel = loop_spec.max_parallel().unwrap_or(items.len().min(10)); let workflow_inputs_arc = Arc::new(ctx.workflow_inputs.clone());
execute_foreach_parallel(
task_id,
spec,
&items,
iterator,
max_parallel,
workflow_inputs_arc,
ctx,
)
.await?;
} else {
execute_foreach_sequential(task_id, spec, &items, iterator, ctx).await?;
}
println!("ForEach loop completed for task '{}'", task_id);
Ok(())
}
LoopSpec::Repeat {
count,
iterator,
parallel,
..
} => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.init_loop(task_id, Some(*count));
}
}
println!(
"Executing Repeat loop for task '{}': {} iterations",
task_id, count
);
if *parallel {
let max_parallel = loop_spec.max_parallel().unwrap_or((*count).min(10)); let workflow_inputs_arc = Arc::new(ctx.workflow_inputs.clone());
execute_repeat_parallel(
task_id,
spec,
*count,
iterator.as_deref(),
max_parallel,
workflow_inputs_arc,
ctx,
)
.await?;
} else {
execute_repeat_sequential(task_id, spec, *count, iterator.as_deref(), ctx).await?;
}
println!("Repeat loop completed for task '{}'", task_id);
Ok(())
}
LoopSpec::While {
condition,
max_iterations,
iteration_variable,
delay_between_secs,
} => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.init_loop(task_id, None);
}
}
println!(
"Executing While loop for task '{}': max {} iterations",
task_id, max_iterations
);
execute_while_sequential(
task_id,
spec,
condition,
*max_iterations,
iteration_variable.as_deref(),
*delay_between_secs,
ctx,
)
.await?;
println!("While loop completed for task '{}'", task_id);
Ok(())
}
LoopSpec::RepeatUntil {
condition,
min_iterations,
max_iterations,
iteration_variable,
delay_between_secs,
} => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.init_loop(task_id, None);
}
}
println!(
"Executing RepeatUntil loop for task '{}': min {}, max {} iterations",
task_id,
min_iterations.unwrap_or(1),
max_iterations
);
execute_repeat_until_sequential(
task_id,
spec,
condition,
*min_iterations,
*max_iterations,
iteration_variable.as_deref(),
*delay_between_secs,
ctx,
)
.await?;
println!("RepeatUntil loop completed for task '{}'", task_id);
Ok(())
}
}
}
#[allow(clippy::too_many_arguments)]
async fn execute_subtasks_in_loop_iteration(
parent_task_id: &str,
substituted_parent: &crate::dsl::schema::TaskSpec,
workflow_inputs: &HashMap<String, serde_json::Value>,
agents: &Arc<Mutex<HashMap<String, PeriplonSDKClient>>>,
task_graph: &Arc<Mutex<TaskGraph>>,
state: &Arc<Mutex<Option<WorkflowState>>>,
workflow_name: &Arc<String>,
json_output: bool,
) -> Result<Option<String>> {
for subtask_map in &substituted_parent.subtasks {
for (subtask_name, subtask_spec) in subtask_map {
let full_subtask_id = format!("{}.{}", parent_task_id, subtask_name);
println!(
" Executing subtask: {} - {}",
full_subtask_id, subtask_spec.description
);
if !subtask_spec.depends_on.is_empty() {
let deps_met = {
let graph = task_graph.lock().await;
subtask_spec.depends_on.iter().all(|dep| {
let dep_full_id = if dep.contains('.') {
dep.clone()
} else {
format!("{}.{}", parent_task_id, dep)
};
if let Some(dep_task) = graph.get_task(&dep_full_id) {
dep_task.status == TaskStatus::Completed
} else {
true
}
})
};
if !deps_met {
println!(
" Skipping subtask {} - dependencies not met",
subtask_name
);
continue;
}
}
if let Some(ref condition) = subtask_spec.condition {
let condition_met = {
let graph = task_graph.lock().await;
let workflow_state = state.lock().await;
evaluate_condition(condition, &graph, workflow_state.as_ref())
};
if !condition_met {
println!(
" Skipping subtask {} - condition not met",
subtask_name
);
continue;
}
}
let result = execute_task_attempt(
&full_subtask_id,
&subtask_spec.description,
subtask_spec,
workflow_inputs,
agents,
0,
state,
workflow_name,
json_output,
)
.await;
match result {
Ok(_) => {
let mut graph = task_graph.lock().await;
let _ = graph.update_task_status(&full_subtask_id, TaskStatus::Completed);
println!(" ✓ Subtask {} completed", subtask_name);
}
Err(e) => {
let mut graph = task_graph.lock().await;
let _ = graph.update_task_status(&full_subtask_id, TaskStatus::Failed);
println!(" ✗ Subtask {} failed: {}", subtask_name, e);
return Err(e);
}
}
}
}
Ok(None)
}
async fn execute_foreach_sequential(
task_id: &str,
spec: &crate::dsl::schema::TaskSpec,
items: &[serde_json::Value],
iterator: &str,
ctx: &ExecutionContext<'_>,
) -> Result<()> {
let timeout_duration = spec
.loop_control
.as_ref()
.and_then(|lc| lc.timeout_secs)
.map(tokio::time::Duration::from_secs);
let loop_future = async {
let checkpoint_interval = spec
.loop_control
.as_ref()
.and_then(|lc| lc.checkpoint_interval)
.filter(|&interval| interval > 0);
for (iteration, item) in items.iter().enumerate() {
let already_completed = {
let state_guard = ctx.state.lock().await;
if let Some(ref workflow_state) = *state_guard {
workflow_state.is_iteration_completed(task_id, iteration)
} else {
false
}
};
if already_completed {
println!(
" Iteration {}: Already completed, skipping (resume)",
iteration + 1
);
continue;
}
if let Some(ref loop_control) = spec.loop_control {
if let Some(ref continue_cond) = loop_control.continue_condition {
let should_continue = {
let task_graph_guard = ctx.task_graph.lock().await;
let state_guard = ctx.state.lock().await;
evaluate_condition(continue_cond, &task_graph_guard, state_guard.as_ref())
};
if should_continue {
println!(
" Iteration {}: Skipping due to continue condition",
iteration + 1
);
continue;
}
}
}
println!(
" Iteration {}/{}: Processing item: {:?}",
iteration + 1,
items.len(),
item
);
let mut context = LoopContext::new(iteration);
context.set_variable(iterator.to_string(), item.clone());
let substituted_task = substitute_task_variables(spec, &context);
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Running,
Some(item.clone()),
);
}
}
let result = if !substituted_task.subtasks.is_empty() {
println!(
" Task has {} subtasks - executing within loop iteration",
substituted_task.subtasks.len()
);
execute_subtasks_in_loop_iteration(
task_id,
&substituted_task,
ctx.workflow_inputs,
ctx.agents,
ctx.task_graph,
ctx.state,
ctx.workflow_name,
ctx.json_output,
)
.await
} else {
execute_task_attempt(
&format!("{}[{}]", task_id, iteration),
&substituted_task.description,
&substituted_task,
ctx.workflow_inputs,
ctx.agents,
0,
ctx.state,
ctx.workflow_name,
ctx.json_output,
)
.await
};
match result {
Ok(output) => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Completed,
Some(item.clone()),
);
if spec
.loop_control
.as_ref()
.is_some_and(|lc| lc.collect_results)
{
if let Some(output_str) = output {
workflow_state.store_loop_result(
task_id,
serde_json::Value::String(output_str),
);
}
}
}
}
println!(" Iteration {} completed successfully", iteration + 1);
if let Some(interval) = checkpoint_interval {
if (iteration + 1) % interval == 0 {
let state_guard = ctx.state.lock().await;
if let Some(ref workflow_state) = *state_guard {
if let Err(e) = workflow_state.save_checkpoint() {
println!(" Warning: Failed to save checkpoint: {}", e);
} else {
println!(" Checkpoint saved at iteration {}", iteration + 1);
}
}
}
}
}
Err(e) => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Failed,
Some(item.clone()),
);
}
}
let break_on_error = spec
.loop_control
.as_ref()
.and_then(|lc| lc.break_condition.as_ref())
.is_some();
if break_on_error {
return Err(e);
}
println!(" Iteration {} failed: {}", iteration + 1, e);
}
}
if let Some(ref loop_control) = spec.loop_control {
if let Some(ref break_cond) = loop_control.break_condition {
let should_break = {
let task_graph_guard = ctx.task_graph.lock().await;
let state_guard = ctx.state.lock().await;
evaluate_condition(break_cond, &task_graph_guard, state_guard.as_ref())
};
if should_break {
println!(
" Breaking loop due to break condition at iteration {}",
iteration + 1
);
break;
}
}
}
}
Ok(())
};
if let Some(timeout) = timeout_duration {
match tokio::time::timeout(timeout, loop_future).await {
Ok(result) => result,
Err(_) => {
println!(" Loop timed out after {:?}", timeout);
Err(Error::InvalidInput(format!(
"Loop '{}' timed out after {} seconds",
task_id,
timeout.as_secs()
)))
}
}
} else {
loop_future.await
}
}
async fn execute_repeat_sequential(
task_id: &str,
spec: &crate::dsl::schema::TaskSpec,
count: usize,
iterator: Option<&str>,
ctx: &ExecutionContext<'_>,
) -> Result<()> {
for iteration in 0..count {
if let Some(ref loop_control) = spec.loop_control {
if let Some(ref continue_cond) = loop_control.continue_condition {
let should_continue = {
let task_graph_guard = ctx.task_graph.lock().await;
let state_guard = ctx.state.lock().await;
evaluate_condition(continue_cond, &task_graph_guard, state_guard.as_ref())
};
if should_continue {
println!(
" Iteration {}: Skipping due to continue condition",
iteration + 1
);
continue;
}
}
}
println!(" Iteration {}/{}", iteration + 1, count);
let mut context = LoopContext::new(iteration);
if let Some(iter_name) = iterator {
context.set_variable(
iter_name.to_string(),
serde_json::Value::Number(iteration.into()),
);
}
let substituted_task = substitute_task_variables(spec, &context);
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Running,
Some(serde_json::Value::Number(iteration.into())),
);
}
}
match execute_task_attempt(
&format!("{}[{}]", task_id, iteration),
&substituted_task.description,
&substituted_task,
ctx.workflow_inputs,
ctx.agents,
0,
ctx.state,
ctx.workflow_name,
ctx.json_output,
)
.await
{
Ok(output) => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Completed,
Some(serde_json::Value::Number(iteration.into())),
);
if spec
.loop_control
.as_ref()
.is_some_and(|lc| lc.collect_results)
{
if let Some(output_str) = output {
workflow_state.store_loop_result(
task_id,
serde_json::Value::String(output_str),
);
}
}
}
}
println!(" Iteration {} completed successfully", iteration + 1);
}
Err(e) => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Failed,
Some(serde_json::Value::Number(iteration.into())),
);
}
}
let break_on_error = spec
.loop_control
.as_ref()
.and_then(|lc| lc.break_condition.as_ref())
.is_some();
if break_on_error {
return Err(e);
}
println!(" Iteration {} failed: {}", iteration + 1, e);
}
}
if let Some(ref loop_control) = spec.loop_control {
if let Some(ref break_cond) = loop_control.break_condition {
let should_break = {
let task_graph_guard = ctx.task_graph.lock().await;
let state_guard = ctx.state.lock().await;
evaluate_condition(break_cond, &task_graph_guard, state_guard.as_ref())
};
if should_break {
println!(
" Breaking loop due to break condition at iteration {}",
iteration + 1
);
break;
}
}
}
}
Ok(())
}
async fn execute_while_sequential(
task_id: &str,
spec: &crate::dsl::schema::TaskSpec,
condition: &crate::dsl::schema::ConditionSpec,
max_iterations: usize,
iteration_variable: Option<&str>,
delay_between_secs: Option<u64>,
ctx: &ExecutionContext<'_>,
) -> Result<()> {
let mut iteration = 0;
loop {
if iteration >= max_iterations {
println!(
" While loop reached max iterations limit: {}",
max_iterations
);
break;
}
let condition_met = {
let state_guard = ctx.state.lock().await;
let task_graph = TaskGraph::new(); evaluate_condition(condition, &task_graph, state_guard.as_ref())
};
if !condition_met {
println!(
" While loop condition became false at iteration {}",
iteration
);
break;
}
println!(
" Iteration {}: Condition is true, executing...",
iteration + 1
);
let mut context = LoopContext::new(iteration);
if let Some(iter_name) = iteration_variable {
context.set_variable(
iter_name.to_string(),
serde_json::Value::Number(iteration.into()),
);
}
let substituted_task = substitute_task_variables(spec, &context);
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Running,
Some(serde_json::Value::Number(iteration.into())),
);
}
}
match execute_task_attempt(
&format!("{}[{}]", task_id, iteration),
&substituted_task.description,
&substituted_task,
ctx.workflow_inputs,
ctx.agents,
0,
ctx.state,
ctx.workflow_name,
ctx.json_output,
)
.await
{
Ok(output) => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Completed,
Some(serde_json::Value::Number(iteration.into())),
);
if spec
.loop_control
.as_ref()
.is_some_and(|lc| lc.collect_results)
{
if let Some(output_str) = output {
workflow_state.store_loop_result(
task_id,
serde_json::Value::String(output_str),
);
}
}
}
}
println!(" Iteration {} completed successfully", iteration + 1);
}
Err(e) => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Failed,
Some(serde_json::Value::Number(iteration.into())),
);
}
}
let break_on_error = spec
.loop_control
.as_ref()
.and_then(|lc| lc.break_condition.as_ref())
.is_some();
if break_on_error {
return Err(e);
}
println!(" Iteration {} failed: {}", iteration + 1, e);
}
}
iteration += 1;
if let Some(delay_secs) = delay_between_secs {
if delay_secs > 0 {
println!(" Waiting {} seconds before next iteration...", delay_secs);
tokio::time::sleep(tokio::time::Duration::from_secs(delay_secs)).await;
}
}
}
println!(" While loop completed after {} iterations", iteration);
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn execute_repeat_until_sequential(
task_id: &str,
spec: &crate::dsl::schema::TaskSpec,
condition: &crate::dsl::schema::ConditionSpec,
min_iterations: Option<usize>,
max_iterations: usize,
iteration_variable: Option<&str>,
delay_between_secs: Option<u64>,
ctx: &ExecutionContext<'_>,
) -> Result<()> {
let min = min_iterations.unwrap_or(1);
let mut iteration = 0;
loop {
if iteration >= max_iterations {
println!(
" RepeatUntil loop reached max iterations limit: {}",
max_iterations
);
break;
}
println!(" Iteration {}", iteration + 1);
let mut context = LoopContext::new(iteration);
if let Some(iter_name) = iteration_variable {
context.set_variable(
iter_name.to_string(),
serde_json::Value::Number(iteration.into()),
);
}
let substituted_task = substitute_task_variables(spec, &context);
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Running,
Some(serde_json::Value::Number(iteration.into())),
);
}
}
match execute_task_attempt(
&format!("{}[{}]", task_id, iteration),
&substituted_task.description,
&substituted_task,
ctx.workflow_inputs,
ctx.agents,
0,
ctx.state,
ctx.workflow_name,
ctx.json_output,
)
.await
{
Ok(output) => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Completed,
Some(serde_json::Value::Number(iteration.into())),
);
if spec
.loop_control
.as_ref()
.is_some_and(|lc| lc.collect_results)
{
if let Some(output_str) = output {
workflow_state.store_loop_result(
task_id,
serde_json::Value::String(output_str),
);
}
}
}
}
println!(" Iteration {} completed successfully", iteration + 1);
}
Err(e) => {
{
let mut state_guard = ctx.state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
task_id,
iteration,
TaskStatus::Failed,
Some(serde_json::Value::Number(iteration.into())),
);
}
}
let break_on_error = spec
.loop_control
.as_ref()
.and_then(|lc| lc.break_condition.as_ref())
.is_some();
if break_on_error {
return Err(e);
}
println!(" Iteration {} failed: {}", iteration + 1, e);
}
}
iteration += 1;
let condition_met = {
let state_guard = ctx.state.lock().await;
let task_graph = TaskGraph::new(); evaluate_condition(condition, &task_graph, state_guard.as_ref())
};
if condition_met && iteration >= min {
println!(
" RepeatUntil condition became true at iteration {}",
iteration
);
break;
}
if let Some(delay_secs) = delay_between_secs {
if delay_secs > 0 {
println!(" Waiting {} seconds before next iteration...", delay_secs);
tokio::time::sleep(tokio::time::Duration::from_secs(delay_secs)).await;
}
}
}
println!(
" RepeatUntil loop completed after {} iterations",
iteration
);
Ok(())
}
async fn execute_foreach_parallel(
task_id: &str,
spec: &crate::dsl::schema::TaskSpec,
items: &[serde_json::Value],
iterator: &str,
max_parallel: usize,
workflow_inputs: Arc<HashMap<String, serde_json::Value>>,
ctx: &ExecutionContext<'_>,
) -> Result<()> {
use tokio::task::JoinSet;
println!(
" Executing ForEach in parallel: {} items, max {} concurrent",
items.len(),
max_parallel
);
let items_owned: Vec<_> = items.to_vec();
let total_items = items_owned.len();
let semaphore = Arc::new(Semaphore::new(max_parallel));
let mut join_set = JoinSet::new();
for (iteration, item) in items_owned.into_iter().enumerate() {
let task_id = task_id.to_string();
let spec = spec.clone();
let iterator = iterator.to_string();
let agents = ctx.agents.clone();
let state = ctx.state.clone();
let semaphore = semaphore.clone();
let workflow_inputs_clone = workflow_inputs.clone();
let workflow_name_clone = ctx.workflow_name.clone();
let json_output = ctx.json_output;
join_set.spawn(async move {
let _permit = semaphore.acquire().await.unwrap();
println!(
" [Parallel] Iteration {}/{}: Processing item: {:?}",
iteration + 1,
total_items,
item
);
let mut context = LoopContext::new(iteration);
context.set_variable(iterator, item.clone());
let substituted_task = substitute_task_variables(&spec, &context);
{
let mut state_guard = state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
&task_id,
iteration,
TaskStatus::Running,
Some(item.clone()),
);
}
}
let result = execute_task_attempt(
&format!("{}[{}]", task_id, iteration),
&substituted_task.description,
&substituted_task,
&workflow_inputs_clone,
&agents,
0,
&state,
&workflow_name_clone,
json_output,
)
.await;
match result {
Ok(output) => {
let mut state_guard = state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
&task_id,
iteration,
TaskStatus::Completed,
Some(item.clone()),
);
if spec
.loop_control
.as_ref()
.is_some_and(|lc| lc.collect_results)
{
if let Some(output_str) = output {
workflow_state.store_loop_result(
&task_id,
serde_json::Value::String(output_str),
);
}
}
}
println!(
" [Parallel] Iteration {} completed successfully",
iteration + 1
);
Ok(())
}
Err(e) => {
let mut state_guard = state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
&task_id,
iteration,
TaskStatus::Failed,
Some(item),
);
}
println!(" [Parallel] Iteration {} failed: {}", iteration + 1, e);
Err(e)
}
}
});
}
let mut errors = Vec::new();
while let Some(result) = join_set.join_next().await {
match result {
Ok(Ok(())) => {
}
Ok(Err(e)) => {
errors.push(e);
}
Err(join_err) => {
errors.push(Error::InvalidInput(format!("Task panicked: {}", join_err)));
}
}
}
if let Some(error) = errors.into_iter().next() {
return Err(error);
}
Ok(())
}
async fn execute_repeat_parallel(
task_id: &str,
spec: &crate::dsl::schema::TaskSpec,
count: usize,
iterator: Option<&str>,
max_parallel: usize,
workflow_inputs: Arc<HashMap<String, serde_json::Value>>,
ctx: &ExecutionContext<'_>,
) -> Result<()> {
use tokio::task::JoinSet;
println!(
" Executing Repeat in parallel: {} iterations, max {} concurrent",
count, max_parallel
);
let semaphore = Arc::new(Semaphore::new(max_parallel));
let mut join_set = JoinSet::new();
for iteration in 0..count {
let task_id = task_id.to_string();
let spec = spec.clone();
let iterator = iterator.map(|s| s.to_string());
let agents = ctx.agents.clone();
let state = ctx.state.clone();
let semaphore = semaphore.clone();
let workflow_inputs_clone = workflow_inputs.clone();
let workflow_name_clone = ctx.workflow_name.clone();
let json_output = ctx.json_output;
join_set.spawn(async move {
let _permit = semaphore.acquire().await.unwrap();
println!(" [Parallel] Iteration {}/{}", iteration + 1, count);
let mut context = LoopContext::new(iteration);
if let Some(iter_name) = iterator {
context.set_variable(iter_name, serde_json::Value::Number(iteration.into()));
}
let substituted_task = substitute_task_variables(&spec, &context);
{
let mut state_guard = state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
&task_id,
iteration,
TaskStatus::Running,
Some(serde_json::Value::Number(iteration.into())),
);
}
}
let result = execute_task_attempt(
&format!("{}[{}]", task_id, iteration),
&substituted_task.description,
&substituted_task,
&workflow_inputs_clone,
&agents,
0,
&state,
&workflow_name_clone,
json_output,
)
.await;
match result {
Ok(output) => {
let mut state_guard = state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
&task_id,
iteration,
TaskStatus::Completed,
Some(serde_json::Value::Number(iteration.into())),
);
if spec
.loop_control
.as_ref()
.is_some_and(|lc| lc.collect_results)
{
if let Some(output_str) = output {
workflow_state.store_loop_result(
&task_id,
serde_json::Value::String(output_str),
);
}
}
}
println!(
" [Parallel] Iteration {} completed successfully",
iteration + 1
);
Ok(())
}
Err(e) => {
let mut state_guard = state.lock().await;
if let Some(ref mut workflow_state) = *state_guard {
workflow_state.update_loop_iteration(
&task_id,
iteration,
TaskStatus::Failed,
Some(serde_json::Value::Number(iteration.into())),
);
}
println!(" [Parallel] Iteration {} failed: {}", iteration + 1, e);
Err(e)
}
}
});
}
let mut errors = Vec::new();
while let Some(result) = join_set.join_next().await {
match result {
Ok(Ok(())) => {
}
Ok(Err(e)) => {
errors.push(e);
}
Err(join_err) => {
errors.push(Error::InvalidInput(format!("Task panicked: {}", join_err)));
}
}
}
if let Some(error) = errors.into_iter().next() {
return Err(error);
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::domain::Provider;
use crate::dsl::schema::PermissionsSpec;
#[test]
fn test_executor_creation() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test Workflow".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: None,
create_cwd: None,
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let (name, version) = executor.get_workflow_info();
assert_eq!(name, "Test Workflow");
assert_eq!(version, "1.0.0");
}
#[test]
fn test_agent_spec_to_options() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: None,
create_cwd: None,
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let agent_spec = AgentSpec {
provider: None,
description: "Test agent".to_string(),
model: Some("claude-sonnet-4-5".to_string()),
system_prompt: None,
cwd: None,
create_cwd: None,
inputs: HashMap::new(),
outputs: HashMap::new(),
tools: vec!["Read".to_string(), "Write".to_string()],
permissions: PermissionsSpec {
mode: "acceptEdits".to_string(),
allowed_directories: vec!["./src".to_string(), "/tmp/test".to_string()],
},
max_turns: Some(10),
};
let var_context = crate::dsl::variables::VariableContext::new();
let options = executor
.agent_spec_to_options(&agent_spec, None, None, &var_context)
.unwrap();
assert_eq!(options.allowed_tools.len(), 2);
assert_eq!(options.model, Some("claude-sonnet-4-5".to_string()));
assert_eq!(options.max_turns, Some(10));
assert_eq!(options.permission_mode, Some("acceptEdits".to_string()));
assert_eq!(options.add_dirs.len(), 2);
assert_eq!(options.add_dirs[0], std::path::PathBuf::from("./src"));
assert_eq!(options.add_dirs[1], std::path::PathBuf::from("/tmp/test"));
assert_eq!(options.cwd, None); assert!(!options.create_cwd); }
#[test]
fn test_cwd_workflow_level_only() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: Some("/workflow/dir".to_string()),
create_cwd: None,
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let agent_spec = AgentSpec {
provider: None,
description: "Test agent".to_string(),
model: None,
system_prompt: None,
cwd: None,
create_cwd: None,
inputs: HashMap::new(),
outputs: HashMap::new(),
tools: vec![],
permissions: PermissionsSpec::default(),
max_turns: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let options = executor
.agent_spec_to_options(&agent_spec, Some("/workflow/dir"), None, &var_context)
.unwrap();
assert_eq!(options.cwd, Some(std::path::PathBuf::from("/workflow/dir")));
}
#[test]
fn test_cwd_agent_level_only() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: None,
create_cwd: None,
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let agent_spec = AgentSpec {
provider: None,
description: "Test agent".to_string(),
model: None,
system_prompt: None,
cwd: Some("/agent/dir".to_string()),
create_cwd: None,
inputs: HashMap::new(),
outputs: HashMap::new(),
tools: vec![],
permissions: PermissionsSpec::default(),
max_turns: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let options = executor
.agent_spec_to_options(&agent_spec, None, None, &var_context)
.unwrap();
assert_eq!(options.cwd, Some(std::path::PathBuf::from("/agent/dir")));
}
#[test]
fn test_cwd_agent_overrides_workflow() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: Some("/workflow/dir".to_string()),
create_cwd: None,
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let agent_spec = AgentSpec {
provider: None,
description: "Test agent".to_string(),
model: None,
system_prompt: None,
cwd: Some("/agent/dir".to_string()),
create_cwd: None,
inputs: HashMap::new(),
outputs: HashMap::new(),
tools: vec![],
permissions: PermissionsSpec::default(),
max_turns: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let options = executor
.agent_spec_to_options(&agent_spec, Some("/workflow/dir"), None, &var_context)
.unwrap();
assert_eq!(options.cwd, Some(std::path::PathBuf::from("/agent/dir")));
}
#[test]
fn test_cwd_backwards_compatibility() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: None,
create_cwd: None,
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let agent_spec = AgentSpec {
provider: None,
description: "Test agent".to_string(),
model: None,
system_prompt: None,
cwd: None,
create_cwd: None,
inputs: HashMap::new(),
outputs: HashMap::new(),
tools: vec![],
permissions: PermissionsSpec::default(),
max_turns: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let options = executor
.agent_spec_to_options(&agent_spec, None, None, &var_context)
.unwrap();
assert_eq!(options.cwd, None);
}
#[test]
fn test_create_cwd_workflow_level() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: Some("/tmp/test".to_string()),
create_cwd: Some(true),
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let agent_spec = AgentSpec {
provider: None,
description: "Test agent".to_string(),
model: None,
system_prompt: None,
cwd: None,
create_cwd: None,
inputs: HashMap::new(),
outputs: HashMap::new(),
tools: vec![],
permissions: PermissionsSpec::default(),
max_turns: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let options = executor
.agent_spec_to_options(&agent_spec, Some("/tmp/test"), Some(true), &var_context)
.unwrap();
assert!(options.create_cwd);
}
#[test]
fn test_create_cwd_agent_overrides_workflow() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: Some("/tmp/test".to_string()),
create_cwd: Some(true),
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let agent_spec = AgentSpec {
provider: None,
description: "Test agent".to_string(),
model: None,
system_prompt: None,
cwd: Some("/agent/dir".to_string()),
create_cwd: Some(false), inputs: HashMap::new(),
outputs: HashMap::new(),
tools: vec![],
permissions: PermissionsSpec::default(),
max_turns: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let options = executor
.agent_spec_to_options(&agent_spec, Some("/tmp/test"), Some(true), &var_context)
.unwrap();
assert!(!options.create_cwd);
}
#[test]
fn test_create_cwd_defaults_to_false() {
let workflow = DSLWorkflow {
provider: Provider::Claude,
model: None,
name: "Test".to_string(),
version: "1.0.0".to_string(),
dsl_version: "1.0.0".to_string(),
cwd: Some("/tmp/test".to_string()),
create_cwd: None,
secrets: HashMap::new(),
inputs: HashMap::new(),
outputs: HashMap::new(),
agents: HashMap::new(),
tasks: HashMap::new(),
workflows: HashMap::new(),
tools: None,
communication: None,
mcp_servers: HashMap::new(),
subflows: HashMap::new(),
imports: HashMap::new(),
notifications: None,
limits: None,
};
let executor = DSLExecutor::new(workflow).unwrap();
let agent_spec = AgentSpec {
provider: None,
description: "Test agent".to_string(),
model: None,
system_prompt: None,
cwd: None,
create_cwd: None,
inputs: HashMap::new(),
outputs: HashMap::new(),
tools: vec![],
permissions: PermissionsSpec::default(),
max_turns: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let options = executor
.agent_spec_to_options(&agent_spec, Some("/tmp/test"), None, &var_context)
.unwrap();
assert!(!options.create_cwd);
}
#[test]
fn test_condition_always() {
use crate::dsl::schema::{Condition, ConditionSpec};
let task_graph = TaskGraph::new();
let condition = ConditionSpec::Single(Condition::Always);
assert!(evaluate_condition(&condition, &task_graph, None));
}
#[test]
fn test_condition_never() {
use crate::dsl::schema::{Condition, ConditionSpec};
let task_graph = TaskGraph::new();
let condition = ConditionSpec::Single(Condition::Never);
assert!(!evaluate_condition(&condition, &task_graph, None));
}
#[test]
fn test_condition_task_status_completed() {
use crate::dsl::schema::{Condition, ConditionSpec, TaskSpec, TaskStatusCondition};
let mut task_graph = TaskGraph::new();
task_graph.add_task("task1".to_string(), TaskSpec::default());
task_graph
.update_task_status("task1", TaskStatus::Completed)
.unwrap();
let condition = ConditionSpec::Single(Condition::TaskStatus {
task: "task1".to_string(),
status: TaskStatusCondition::Completed,
});
assert!(evaluate_condition(&condition, &task_graph, None));
}
#[test]
fn test_condition_task_status_failed() {
use crate::dsl::schema::{Condition, ConditionSpec, TaskSpec, TaskStatusCondition};
let mut task_graph = TaskGraph::new();
task_graph.add_task("task1".to_string(), TaskSpec::default());
task_graph
.update_task_status("task1", TaskStatus::Failed)
.unwrap();
let condition = ConditionSpec::Single(Condition::TaskStatus {
task: "task1".to_string(),
status: TaskStatusCondition::Failed,
});
assert!(evaluate_condition(&condition, &task_graph, None));
}
#[test]
fn test_condition_task_status_skipped() {
use crate::dsl::schema::{Condition, ConditionSpec, TaskSpec, TaskStatusCondition};
let mut task_graph = TaskGraph::new();
task_graph.add_task("task1".to_string(), TaskSpec::default());
task_graph
.update_task_status("task1", TaskStatus::Skipped)
.unwrap();
let condition = ConditionSpec::Single(Condition::TaskStatus {
task: "task1".to_string(),
status: TaskStatusCondition::Skipped,
});
assert!(evaluate_condition(&condition, &task_graph, None));
}
#[test]
fn test_condition_state_equals() {
use crate::dsl::schema::{Condition, ConditionSpec};
let task_graph = TaskGraph::new();
let mut workflow_state = WorkflowState::new("test".to_string(), "1.0.0".to_string());
workflow_state.add_metadata("environment".to_string(), serde_json::json!("production"));
let condition = ConditionSpec::Single(Condition::StateEquals {
key: "environment".to_string(),
value: serde_json::json!("production"),
});
assert!(evaluate_condition(
&condition,
&task_graph,
Some(&workflow_state)
));
}
#[test]
fn test_condition_state_exists() {
use crate::dsl::schema::{Condition, ConditionSpec};
let task_graph = TaskGraph::new();
let mut workflow_state = WorkflowState::new("test".to_string(), "1.0.0".to_string());
workflow_state.add_metadata("key1".to_string(), serde_json::json!("value1"));
let condition = ConditionSpec::Single(Condition::StateExists {
key: "key1".to_string(),
});
assert!(evaluate_condition(
&condition,
&task_graph,
Some(&workflow_state)
));
let condition_not_exists = ConditionSpec::Single(Condition::StateExists {
key: "key2".to_string(),
});
assert!(!evaluate_condition(
&condition_not_exists,
&task_graph,
Some(&workflow_state)
));
}
#[test]
fn test_condition_and() {
use crate::dsl::schema::{Condition, ConditionSpec};
let task_graph = TaskGraph::new();
let condition = ConditionSpec::And {
and: vec![
ConditionSpec::Single(Condition::Always),
ConditionSpec::Single(Condition::Always),
],
};
assert!(evaluate_condition(&condition, &task_graph, None));
let condition_with_never = ConditionSpec::And {
and: vec![
ConditionSpec::Single(Condition::Always),
ConditionSpec::Single(Condition::Never),
],
};
assert!(!evaluate_condition(
&condition_with_never,
&task_graph,
None
));
}
#[test]
fn test_condition_or() {
use crate::dsl::schema::{Condition, ConditionSpec};
let task_graph = TaskGraph::new();
let condition = ConditionSpec::Or {
or: vec![
ConditionSpec::Single(Condition::Always),
ConditionSpec::Single(Condition::Never),
],
};
assert!(evaluate_condition(&condition, &task_graph, None));
let condition_all_never = ConditionSpec::Or {
or: vec![
ConditionSpec::Single(Condition::Never),
ConditionSpec::Single(Condition::Never),
],
};
assert!(!evaluate_condition(&condition_all_never, &task_graph, None));
}
#[test]
fn test_condition_not() {
use crate::dsl::schema::{Condition, ConditionSpec};
let task_graph = TaskGraph::new();
let condition = ConditionSpec::Not {
not: Box::new(ConditionSpec::Single(Condition::Never)),
};
assert!(evaluate_condition(&condition, &task_graph, None));
let condition_not_always = ConditionSpec::Not {
not: Box::new(ConditionSpec::Single(Condition::Always)),
};
assert!(!evaluate_condition(
&condition_not_always,
&task_graph,
None
));
}
#[test]
fn test_condition_complex() {
use crate::dsl::schema::{Condition, ConditionSpec, TaskSpec, TaskStatusCondition};
let mut task_graph = TaskGraph::new();
task_graph.add_task("task1".to_string(), TaskSpec::default());
task_graph
.update_task_status("task1", TaskStatus::Completed)
.unwrap();
let mut workflow_state = WorkflowState::new("test".to_string(), "1.0.0".to_string());
workflow_state.add_metadata("env".to_string(), serde_json::json!("prod"));
let condition = ConditionSpec::And {
and: vec![
ConditionSpec::Single(Condition::TaskStatus {
task: "task1".to_string(),
status: TaskStatusCondition::Completed,
}),
ConditionSpec::Single(Condition::StateEquals {
key: "env".to_string(),
value: serde_json::json!("prod"),
}),
],
};
assert!(evaluate_condition(
&condition,
&task_graph,
Some(&workflow_state)
));
}
#[tokio::test]
async fn test_dod_file_exists_met() {
use crate::dsl::schema::DoneCriterion;
let temp_file = std::env::temp_dir().join("test_dod_file_exists.txt");
std::fs::write(&temp_file, "test content").unwrap();
let criterion = DoneCriterion::FileExists {
path: temp_file.to_string_lossy().to_string(),
description: "Test file must exist".to_string(),
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(result.met);
assert_eq!(result.description, "Test file must exist");
std::fs::remove_file(temp_file).ok();
}
#[tokio::test]
async fn test_dod_file_exists_not_met() {
use crate::dsl::schema::DoneCriterion;
let nonexistent_file = std::env::temp_dir().join("nonexistent_file_xyz.txt");
let criterion = DoneCriterion::FileExists {
path: nonexistent_file.to_string_lossy().to_string(),
description: "Nonexistent file".to_string(),
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(!result.met);
assert!(result.details.contains("does not exist"));
}
#[tokio::test]
async fn test_dod_file_contains_met() {
use crate::dsl::schema::DoneCriterion;
let temp_file = std::env::temp_dir().join("test_dod_file_contains.txt");
std::fs::write(&temp_file, "Hello World\nTest Content").unwrap();
let criterion = DoneCriterion::FileContains {
path: temp_file.to_string_lossy().to_string(),
pattern: "Hello World".to_string(),
description: "File must contain greeting".to_string(),
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(result.met);
std::fs::remove_file(temp_file).ok();
}
#[tokio::test]
async fn test_dod_file_contains_regex_pattern() {
use crate::dsl::schema::DoneCriterion;
let temp_file = std::env::temp_dir().join("test_dod_file_contains_regex.txt");
std::fs::write(
&temp_file,
"group_list\ngroup_install\ngroup_update\ngroup_validate",
)
.unwrap();
let criterion = DoneCriterion::FileContains {
path: temp_file.to_string_lossy().to_string(),
pattern: "group.*list".to_string(),
description: "File must contain group_list".to_string(),
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(
result.met,
"Regex pattern 'group.*list' should match 'group_list'"
);
let criterion2 = DoneCriterion::FileContains {
path: temp_file.to_string_lossy().to_string(),
pattern: "group.*(update|validate)".to_string(),
description: "File must contain group_update or group_validate".to_string(),
};
let result2 = check_criterion(&criterion2, None, &var_context).await;
assert!(
result2.met,
"Regex pattern 'group.*(update|validate)' should match"
);
std::fs::remove_file(temp_file).ok();
}
#[tokio::test]
async fn test_dod_file_not_contains_met() {
use crate::dsl::schema::DoneCriterion;
let temp_file = std::env::temp_dir().join("test_dod_file_not_contains.txt");
std::fs::write(&temp_file, "Clean code").unwrap();
let criterion = DoneCriterion::FileNotContains {
path: temp_file.to_string_lossy().to_string(),
pattern: "TODO".to_string(),
description: "No TODOs allowed".to_string(),
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(result.met);
std::fs::remove_file(temp_file).ok();
}
#[tokio::test]
async fn test_dod_file_not_contains_failed() {
use crate::dsl::schema::DoneCriterion;
let temp_file = std::env::temp_dir().join("test_dod_file_not_contains_fail.txt");
std::fs::write(&temp_file, "Code with TODO: fix this").unwrap();
let criterion = DoneCriterion::FileNotContains {
path: temp_file.to_string_lossy().to_string(),
pattern: "TODO".to_string(),
description: "No TODOs allowed".to_string(),
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(!result.met);
assert!(result.details.contains("should not"));
std::fs::remove_file(temp_file).ok();
}
#[tokio::test]
async fn test_dod_directory_exists() {
use crate::dsl::schema::DoneCriterion;
let temp_dir = std::env::temp_dir().join("test_dod_directory");
std::fs::create_dir_all(&temp_dir).unwrap();
let criterion = DoneCriterion::DirectoryExists {
path: temp_dir.to_string_lossy().to_string(),
description: "Output directory must exist".to_string(),
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(result.met);
std::fs::remove_dir(temp_dir).ok();
}
#[tokio::test]
async fn test_dod_command_succeeds() {
use crate::dsl::schema::DoneCriterion;
let criterion = DoneCriterion::CommandSucceeds {
command: "echo".to_string(),
args: vec!["test".to_string()],
description: "Echo command should succeed".to_string(),
working_dir: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(result.met);
}
#[tokio::test]
async fn test_dod_command_fails() {
use crate::dsl::schema::DoneCriterion;
let criterion = DoneCriterion::CommandSucceeds {
command: "false".to_string(),
args: vec![],
description: "False command should fail".to_string(),
working_dir: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(!result.met);
assert!(result.details.contains("failed"));
}
#[tokio::test]
async fn test_dod_command_captures_output() {
use crate::dsl::schema::DoneCriterion;
let criterion = DoneCriterion::CommandSucceeds {
command: "sh".to_string(),
args: vec![
"-c".to_string(),
"echo 'This is stdout output'; echo 'This is stderr output' >&2; exit 1"
.to_string(),
],
description: "Command that outputs to stdout/stderr and fails".to_string(),
working_dir: None,
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, None, &var_context).await;
assert!(!result.met);
assert!(result.details.contains("STDOUT"));
assert!(result.details.contains("This is stdout output"));
assert!(result.details.contains("STDERR"));
assert!(result.details.contains("This is stderr output"));
}
#[tokio::test]
async fn test_dod_output_matches() {
use crate::dsl::schema::{DoneCriterion, OutputSource};
let task_output = "Build succeeded\nAll tests passed";
let criterion = DoneCriterion::OutputMatches {
source: OutputSource::TaskOutput,
pattern: "tests passed".to_string(),
description: "Output must indicate tests passed".to_string(),
};
let var_context = crate::dsl::variables::VariableContext::new();
let result = check_criterion(&criterion, Some(task_output), &var_context).await;
assert!(result.met);
}
#[test]
fn test_format_unmet_criteria() {
use crate::dsl::schema::CriterionResult;
let results = vec![
CriterionResult {
met: true,
description: "File exists".to_string(),
details: "File found".to_string(),
},
CriterionResult {
met: false,
description: "Tests must pass".to_string(),
details: "Tests failed with 3 errors".to_string(),
},
CriterionResult {
met: false,
description: "No TODO comments".to_string(),
details: "Found TODO in line 42".to_string(),
},
];
let feedback = format_unmet_criteria(&results);
assert!(feedback.contains("DEFINITION OF DONE"));
assert!(feedback.contains("Tests must pass"));
assert!(feedback.contains("No TODO comments"));
assert!(feedback.contains("✗ FAILED"));
assert!(!feedback.contains("File exists")); }
#[tokio::test]
async fn test_check_definition_of_done_all_met() {
use crate::dsl::schema::{DefinitionOfDone, DoneCriterion};
let temp_file = std::env::temp_dir().join("test_dod_all_met.txt");
std::fs::write(&temp_file, "Complete").unwrap();
let temp_file_str = temp_file.to_string_lossy().to_string();
let dod = DefinitionOfDone {
criteria: vec![
DoneCriterion::FileExists {
path: temp_file_str.clone(),
description: "File must exist".to_string(),
},
DoneCriterion::FileContains {
path: temp_file_str,
pattern: "Complete".to_string(),
description: "File must be complete".to_string(),
},
],
max_retries: 3,
fail_on_unmet: true,
auto_elevate_permissions: false,
};
let var_context = crate::dsl::variables::VariableContext::new();
let results = check_definition_of_done(&dod, None, &var_context).await;
assert_eq!(results.len(), 2);
assert!(results.iter().all(|r| r.met));
std::fs::remove_file(temp_file).ok();
}
#[tokio::test]
async fn test_check_definition_of_done_some_unmet() {
use crate::dsl::schema::{DefinitionOfDone, DoneCriterion};
let existing_file = std::env::temp_dir().join("existing_file.txt");
let nonexistent_file = std::env::temp_dir().join("nonexistent_xyz_abc.txt");
let dod = DefinitionOfDone {
criteria: vec![
DoneCriterion::FileExists {
path: existing_file.to_string_lossy().to_string(),
description: "File must exist".to_string(),
},
DoneCriterion::FileExists {
path: nonexistent_file.to_string_lossy().to_string(),
description: "Another file must exist".to_string(),
},
],
max_retries: 2,
fail_on_unmet: true,
auto_elevate_permissions: false,
};
std::fs::write(&existing_file, "test").unwrap();
let var_context = crate::dsl::variables::VariableContext::new();
let results = check_definition_of_done(&dod, None, &var_context).await;
assert_eq!(results.len(), 2);
assert!(results[0].met);
assert!(!results[1].met);
std::fs::remove_file(existing_file).ok();
}
}