use crate::domain::Provider;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DSLWorkflow {
pub name: String,
pub version: String,
#[serde(
default = "default_dsl_version",
skip_serializing_if = "is_default_dsl_version"
)]
pub dsl_version: String,
#[serde(default, skip_serializing_if = "is_default_provider")]
pub provider: Provider,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cwd: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub create_cwd: Option<bool>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub secrets: HashMap<String, SecretSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub inputs: HashMap<String, InputSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub outputs: HashMap<String, OutputSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub agents: HashMap<String, AgentSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub tasks: HashMap<String, TaskSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub workflows: HashMap<String, WorkflowSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tools: Option<ToolsConfig>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub communication: Option<CommunicationConfig>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub mcp_servers: HashMap<String, McpServerSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub subflows: HashMap<String, SubflowSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub imports: HashMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub notifications: Option<NotificationDefaults>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub limits: Option<LimitsConfig>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentSpec {
pub description: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provider: Option<Provider>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub system_prompt: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cwd: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub create_cwd: Option<bool>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub inputs: HashMap<String, InputSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub outputs: HashMap<String, OutputSpec>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub tools: Vec<String>,
#[serde(default, skip_serializing_if = "is_default_permissions")]
pub permissions: PermissionsSpec,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_turns: Option<u32>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct PermissionsSpec {
#[serde(
default = "default_permission_mode",
skip_serializing_if = "is_default_permission_mode"
)]
pub mode: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub allowed_directories: Vec<String>,
}
impl Default for PermissionsSpec {
fn default() -> Self {
Self {
mode: "default".to_string(),
allowed_directories: Vec::new(),
}
}
}
fn default_permission_mode() -> String {
"default".to_string()
}
fn default_dsl_version() -> String {
"1.0.0".to_string()
}
fn is_default_dsl_version(version: &str) -> bool {
version == "1.0.0"
}
fn is_default_provider(provider: &Provider) -> bool {
*provider == Provider::Claude
}
fn is_default_permission_mode(mode: &str) -> bool {
mode == "default"
}
fn is_default_permissions(perms: &PermissionsSpec) -> bool {
perms == &PermissionsSpec::default()
}
fn is_zero(value: &u32) -> bool {
*value == 0
}
fn is_zero_u32(value: &u32) -> bool {
*value == 0
}
fn is_false(value: &bool) -> bool {
!*value
}
fn is_default_retry_delay(value: &u64) -> bool {
*value == 1
}
fn is_sequential_mode(mode: &ExecutionMode) -> bool {
*mode == ExecutionMode::Sequential
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(untagged)]
pub enum ConditionSpec {
Single(Condition),
And {
and: Vec<ConditionSpec>,
},
Or {
or: Vec<ConditionSpec>,
},
Not {
not: Box<ConditionSpec>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum Condition {
TaskStatus {
task: String,
status: TaskStatusCondition,
},
StateEquals {
key: String,
value: serde_json::Value,
},
StateExists { key: String },
Always,
Never,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum TaskStatusCondition {
Completed,
Failed,
Running,
Pending,
Skipped,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct TaskSpec {
pub description: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub subflow: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub uses: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub embed: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub overrides: Option<serde_yaml::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub script: Option<ScriptSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub command: Option<CommandSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub http: Option<HttpSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub mcp_tool: Option<McpToolSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub uses_workflow: Option<String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub inputs: HashMap<String, serde_json::Value>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub outputs: HashMap<String, OutputSpec>,
#[serde(default, skip_serializing_if = "is_zero")]
pub priority: u32,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub subtasks: Vec<HashMap<String, TaskSpec>>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub depends_on: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub parallel_with: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub on_complete: Option<ActionSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub on_error: Option<ErrorHandlingSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub condition: Option<ConditionSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub definition_of_done: Option<DefinitionOfDone>,
#[serde(default, skip_serializing_if = "Option::is_none", rename = "loop")]
pub loop_spec: Option<LoopSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub loop_control: Option<LoopControl>,
#[serde(default, skip_serializing_if = "is_false")]
pub inject_context: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub limits: Option<LimitsConfig>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<ContextConfig>,
}
impl TaskSpec {
pub fn parse_workflow_reference(reference: &str) -> Option<(&str, &str)> {
let parts: Vec<&str> = reference.split(':').collect();
if parts.len() == 2 && !parts[0].is_empty() && !parts[1].is_empty() {
Some((parts[0], parts[1]))
} else {
None
}
}
pub fn has_execution_type(&self) -> bool {
self.agent.is_some()
|| self.subflow.is_some()
|| self.uses.is_some()
|| self.embed.is_some()
|| self.script.is_some()
|| self.command.is_some()
|| self.http.is_some()
|| self.mcp_tool.is_some()
|| self.uses_workflow.is_some()
}
pub fn execution_type_count(&self) -> usize {
let mut count = 0;
if self.agent.is_some() {
count += 1;
}
if self.subflow.is_some() {
count += 1;
}
if self.uses.is_some() {
count += 1;
}
if self.embed.is_some() {
count += 1;
}
if self.script.is_some() {
count += 1;
}
if self.command.is_some() {
count += 1;
}
if self.http.is_some() {
count += 1;
}
if self.mcp_tool.is_some() {
count += 1;
}
if self.uses_workflow.is_some() {
count += 1;
}
count
}
pub fn inherit_from_parent(&mut self, parent: &TaskSpec) {
if self.agent.is_none() && parent.agent.is_some() && !self.has_execution_type() {
self.agent = parent.agent.clone();
}
if parent.inject_context && !self.inject_context {
self.inject_context = true;
}
if self.on_error.is_none() && parent.on_error.is_some() {
self.on_error = parent.on_error.clone();
}
if self.priority == 0 && parent.priority != 0 {
self.priority = parent.priority;
}
if self.loop_control.is_none() && parent.loop_control.is_some() {
self.loop_control = parent.loop_control.clone();
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DefinitionOfDone {
pub criteria: Vec<DoneCriterion>,
#[serde(
default = "default_dod_retries",
skip_serializing_if = "is_default_dod_retries"
)]
pub max_retries: u32,
#[serde(default = "default_fail_on_unmet", skip_serializing_if = "is_true")]
pub fail_on_unmet: bool,
#[serde(default, skip_serializing_if = "is_false")]
pub auto_elevate_permissions: bool,
}
fn default_dod_retries() -> u32 {
3
}
fn is_default_dod_retries(value: &u32) -> bool {
*value == 3
}
fn default_fail_on_unmet() -> bool {
true
}
fn is_true(value: &bool) -> bool {
*value
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum DoneCriterion {
FileExists { path: String, description: String },
FileContains {
path: String,
pattern: String,
description: String,
},
FileNotContains {
path: String,
pattern: String,
description: String,
},
CommandSucceeds {
command: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
args: Vec<String>,
description: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
working_dir: Option<String>,
},
OutputMatches {
source: OutputSource,
pattern: String,
description: String,
},
DirectoryExists { path: String, description: String },
TestsPassed {
command: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
args: Vec<String>,
description: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OutputSource {
File { path: String },
TaskOutput,
}
#[derive(Debug, Clone)]
pub struct CriterionResult {
pub met: bool,
pub description: String,
pub details: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ActionSpec {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub notify: Option<NotificationSpec>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ErrorHandlingSpec {
#[serde(default, skip_serializing_if = "is_zero_u32")]
pub retry: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub fallback_agent: Option<String>,
#[serde(
default = "default_retry_delay",
skip_serializing_if = "is_default_retry_delay"
)]
pub retry_delay_secs: u64,
#[serde(default, skip_serializing_if = "is_false")]
pub exponential_backoff: bool,
}
fn default_retry_delay() -> u64 {
1
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkflowSpec {
pub description: String,
pub steps: Vec<StageSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub hooks: Option<HooksSpec>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct WorkflowImport {
pub namespace: String,
pub group_reference: String,
}
impl WorkflowImport {
pub fn new(namespace: String, group_reference: String) -> Self {
Self {
namespace,
group_reference,
}
}
pub fn from_entry(namespace: &str, group_reference: &str) -> Self {
Self {
namespace: namespace.to_string(),
group_reference: group_reference.to_string(),
}
}
pub fn validate_namespace(namespace: &str) -> bool {
!namespace.is_empty()
&& !namespace.starts_with(|c: char| c.is_numeric())
&& namespace
.chars()
.all(|c| c.is_alphanumeric() || c == '-' || c == '_')
}
pub fn parse_group_reference(reference: &str) -> Option<(String, String)> {
let parts: Vec<&str> = reference.split('@').collect();
if parts.len() == 2 && !parts[0].is_empty() && !parts[1].is_empty() {
Some((parts[0].to_string(), parts[1].to_string()))
} else {
None
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StageSpec {
pub stage: String,
pub agents: Vec<String>,
pub tasks: Vec<HashMap<String, TaskSpec>>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub depends_on: Vec<String>,
#[serde(default, skip_serializing_if = "is_sequential_mode")]
pub mode: ExecutionMode,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
#[derive(Default)]
pub enum ExecutionMode {
#[default]
Sequential,
Parallel,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct HooksSpec {
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub pre_workflow: Vec<HookCommand>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub post_workflow: Vec<HookCommand>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub on_stage_complete: Vec<HookCommand>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub on_error: Vec<HookCommand>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(untagged)]
pub enum HookCommand {
Command(String),
CommandSpec {
command: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
description: Option<String>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct ToolsConfig {
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub allowed: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub disallowed: Vec<String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub constraints: HashMap<String, ToolConstraints>,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct ToolConstraints {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout: Option<u64>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub allowed_commands: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_file_size: Option<u64>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub allowed_extensions: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub rate_limit: Option<u32>,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct CommunicationConfig {
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub channels: HashMap<String, ChannelSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub message_types: HashMap<String, MessageTypeSpec>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ChannelSpec {
pub description: String,
pub participants: Vec<String>,
pub message_format: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MessageTypeSpec {
pub schema: serde_json::Value,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct McpServerSpec {
#[serde(rename = "type")]
pub server_type: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub command: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub args: Vec<String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub env: HashMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub url: Option<String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub headers: HashMap<String, String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum LoopSpec {
ForEach {
collection: CollectionSource,
iterator: String,
#[serde(default, skip_serializing_if = "is_false")]
parallel: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
max_parallel: Option<usize>,
},
While {
condition: Box<ConditionSpec>,
max_iterations: usize,
#[serde(default, skip_serializing_if = "Option::is_none")]
iteration_variable: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
delay_between_secs: Option<u64>,
},
RepeatUntil {
condition: Box<ConditionSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
min_iterations: Option<usize>,
max_iterations: usize,
#[serde(default, skip_serializing_if = "Option::is_none")]
iteration_variable: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
delay_between_secs: Option<u64>,
},
Repeat {
count: usize,
#[serde(default, skip_serializing_if = "Option::is_none")]
iterator: Option<String>,
#[serde(default, skip_serializing_if = "is_false")]
parallel: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
max_parallel: Option<usize>,
},
}
impl LoopSpec {
pub fn max_parallel(&self) -> Option<usize> {
match self {
LoopSpec::ForEach { max_parallel, .. } => *max_parallel,
LoopSpec::Repeat { max_parallel, .. } => *max_parallel,
_ => None,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "source", rename_all = "snake_case")]
pub enum CollectionSource {
State {
key: String,
},
File {
path: String,
format: FileFormat,
},
Range {
start: i64,
end: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
step: Option<i64>,
},
Inline {
items: Vec<serde_json::Value>,
},
Http {
url: String,
#[serde(default = "default_http_method")]
method: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
headers: Option<HashMap<String, String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
body: Option<String>,
#[serde(default = "default_response_format")]
format: FileFormat,
#[serde(default, skip_serializing_if = "Option::is_none")]
json_path: Option<String>,
},
}
fn default_http_method() -> String {
"GET".to_string()
}
fn default_response_format() -> FileFormat {
FileFormat::Json
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FileFormat {
Json,
JsonLines,
Csv,
Lines,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LoopControl {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub break_condition: Option<ConditionSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub continue_condition: Option<ConditionSpec>,
#[serde(default, skip_serializing_if = "is_false")]
pub collect_results: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub result_key: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout_secs: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub checkpoint_interval: Option<usize>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SubflowSpec {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source: Option<SubflowSource>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub agents: HashMap<String, AgentSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub tasks: HashMap<String, TaskSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub inputs: HashMap<String, InputSpec>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub outputs: HashMap<String, OutputSpec>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "lowercase")]
pub enum SubflowSource {
File {
path: String,
},
Git {
url: String,
path: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
reference: Option<String>,
},
Http {
url: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
checksum: Option<String>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct InputSpec {
#[serde(rename = "type")]
pub param_type: String,
#[serde(default, skip_serializing_if = "is_false")]
pub required: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub default: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct OutputSpec {
pub source: OutputDataSource,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum OutputDataSource {
File {
path: String,
},
State {
key: String,
},
TaskOutput {
task: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SecretSpec {
pub source: SecretSource,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum SecretSource {
Env {
var: String,
},
File {
path: String,
},
Value {
value: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ScriptSpec {
pub language: ScriptLanguage,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub content: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub file: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub working_dir: Option<String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub env: HashMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout_secs: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum ScriptLanguage {
Python,
JavaScript,
Bash,
Ruby,
Perl,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CommandSpec {
pub executable: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub args: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub working_dir: Option<String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub env: HashMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout_secs: Option<u64>,
#[serde(default = "default_true")]
pub capture_stdout: bool,
#[serde(default = "default_true")]
pub capture_stderr: bool,
}
fn default_true() -> bool {
true
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HttpSpec {
pub method: HttpMethod,
pub url: String,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub headers: HashMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub body: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth: Option<HttpAuth>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout_secs: Option<u64>,
#[serde(default = "default_true")]
pub follow_redirects: bool,
#[serde(default = "default_true")]
pub verify_tls: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "UPPERCASE")]
pub enum HttpMethod {
Get,
Post,
Put,
Patch,
Delete,
Head,
Options,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum HttpAuth {
Bearer {
token: String,
},
Basic {
username: String,
password: String,
},
ApiKey {
header: String,
key: String,
},
Custom {
headers: HashMap<String, String>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct McpToolSpec {
pub server: String,
pub tool: String,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub parameters: HashMap<String, serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout_secs: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(untagged)]
pub enum NotificationSpec {
Simple(String),
Structured {
message: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
channels: Vec<NotificationChannel>,
#[serde(default, skip_serializing_if = "Option::is_none")]
title: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
priority: Option<NotificationPriority>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
metadata: HashMap<String, String>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum NotificationChannel {
Console {
#[serde(default = "default_true")]
colored: bool,
#[serde(default = "default_true")]
timestamp: bool,
},
Email {
to: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
cc: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
bcc: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
subject: Option<String>,
smtp: SmtpConfig,
},
Slack {
credential: String,
channel: String,
#[serde(default = "default_slack_method")]
method: SlackMethod,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
attachments: Vec<SlackAttachment>,
},
Discord {
webhook_url: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
username: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
avatar_url: Option<String>,
#[serde(default, skip_serializing_if = "is_false")]
tts: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
embed: Option<DiscordEmbed>,
},
Teams {
webhook_url: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
theme_color: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
facts: Vec<TeamsFact>,
},
Telegram {
bot_token: String,
chat_id: String,
#[serde(default = "default_parse_mode")]
parse_mode: TelegramParseMode,
#[serde(default, skip_serializing_if = "is_false")]
disable_preview: bool,
#[serde(default, skip_serializing_if = "is_false")]
silent: bool,
},
PagerDuty {
integration_key: String,
action: PagerDutyAction,
#[serde(default = "default_pagerduty_severity")]
severity: PagerDutySeverity,
#[serde(default, skip_serializing_if = "Option::is_none")]
dedup_key: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
custom_details: Option<serde_json::Value>,
},
Webhook {
url: String,
#[serde(default = "default_webhook_method")]
method: HttpMethod,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
headers: HashMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
auth: Option<HttpAuth>,
#[serde(default, skip_serializing_if = "Option::is_none")]
body_template: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
timeout_secs: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
retry: Option<RetryConfig>,
},
File {
path: String,
#[serde(default = "default_true")]
append: bool,
#[serde(default = "default_true")]
timestamp: bool,
#[serde(default = "default_file_format")]
format: FileNotificationFormat,
},
Ntfy {
#[serde(default = "default_ntfy_server")]
server: String,
topic: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
title: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
priority: Option<u8>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
tags: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
click_url: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
attach_url: Option<String>,
#[serde(default, skip_serializing_if = "is_false")]
markdown: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
auth_token: Option<String>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct NotificationDefaults {
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub default_channels: Vec<NotificationChannel>,
#[serde(default, skip_serializing_if = "is_false")]
pub notify_on_completion: bool,
#[serde(default = "default_true")]
pub notify_on_failure: bool,
#[serde(default, skip_serializing_if = "is_false")]
pub notify_on_start: bool,
#[serde(default = "default_true")]
pub notify_on_workflow_completion: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum NotificationPriority {
Low,
Normal,
High,
Critical,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SmtpConfig {
pub host: String,
pub port: u16,
pub username: String,
pub password: String,
pub from: String,
#[serde(default = "default_true")]
pub use_tls: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum SlackMethod {
Webhook,
Bot,
}
fn default_slack_method() -> SlackMethod {
SlackMethod::Webhook
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SlackAttachment {
pub text: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub color: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub fields: Vec<SlackField>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SlackField {
pub title: String,
pub value: String,
#[serde(default, skip_serializing_if = "is_false")]
pub short: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DiscordEmbed {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub color: Option<u32>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub fields: Vec<DiscordField>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub footer: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timestamp: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DiscordField {
pub name: String,
pub value: String,
#[serde(default, skip_serializing_if = "is_false")]
pub inline: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TeamsFact {
pub name: String,
pub value: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "PascalCase")]
pub enum TelegramParseMode {
Markdown,
Html,
None,
}
fn default_parse_mode() -> TelegramParseMode {
TelegramParseMode::Markdown
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum PagerDutyAction {
Trigger,
Acknowledge,
Resolve,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum PagerDutySeverity {
Critical,
Error,
Warning,
Info,
}
fn default_pagerduty_severity() -> PagerDutySeverity {
PagerDutySeverity::Error
}
fn default_webhook_method() -> HttpMethod {
HttpMethod::Post
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RetryConfig {
#[serde(default = "default_retry_attempts")]
pub max_attempts: u32,
#[serde(default = "default_retry_delay")]
pub delay_secs: u64,
#[serde(default, skip_serializing_if = "is_false")]
pub exponential_backoff: bool,
}
fn default_retry_attempts() -> u32 {
3
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum FileNotificationFormat {
Text,
Json,
JsonLines,
}
fn default_file_format() -> FileNotificationFormat {
FileNotificationFormat::Text
}
fn default_ntfy_server() -> String {
"https://ntfy.sh".to_string()
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LimitsConfig {
#[serde(default = "default_max_stdout_bytes")]
pub max_stdout_bytes: usize,
#[serde(default = "default_max_stderr_bytes")]
pub max_stderr_bytes: usize,
#[serde(default = "default_max_combined_bytes")]
pub max_combined_bytes: usize,
#[serde(default = "default_truncation_strategy")]
pub truncation_strategy: TruncationStrategy,
#[serde(default = "default_max_context_bytes")]
pub max_context_bytes: usize,
#[serde(default = "default_max_context_tasks")]
pub max_context_tasks: usize,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub external_storage_threshold: Option<usize>,
#[serde(default = "default_external_storage_dir")]
pub external_storage_dir: String,
#[serde(default = "default_true")]
pub compress_external: bool,
#[serde(default = "default_cleanup_strategy")]
pub cleanup_strategy: CleanupStrategy,
}
impl Default for LimitsConfig {
fn default() -> Self {
Self {
max_stdout_bytes: 1_048_576, max_stderr_bytes: 262_144, max_combined_bytes: 1_572_864, truncation_strategy: TruncationStrategy::Tail,
max_context_bytes: 102_400, max_context_tasks: 10,
external_storage_threshold: Some(5_242_880), external_storage_dir: ".workflow_state/task_outputs".to_string(),
compress_external: true,
cleanup_strategy: CleanupStrategy::MostRecent { keep_count: 20 },
}
}
}
fn default_max_stdout_bytes() -> usize {
1_048_576 }
fn default_max_stderr_bytes() -> usize {
262_144 }
fn default_max_combined_bytes() -> usize {
1_572_864 }
fn default_max_context_bytes() -> usize {
102_400 }
fn default_max_context_tasks() -> usize {
10
}
fn default_external_storage_dir() -> String {
".workflow_state/task_outputs".to_string()
}
fn default_truncation_strategy() -> TruncationStrategy {
TruncationStrategy::Tail
}
fn default_cleanup_strategy() -> CleanupStrategy {
CleanupStrategy::MostRecent { keep_count: 20 }
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum TruncationStrategy {
Head,
Tail,
Both,
Summary,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum CleanupStrategy {
MostRecent { keep_count: usize },
HighestRelevance { keep_count: usize },
Lru { keep_count: usize },
DirectDependencies,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ContextConfig {
#[serde(default = "default_context_mode")]
pub mode: ContextMode,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub include_tasks: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub exclude_tasks: Vec<String>,
#[serde(default = "default_min_relevance")]
pub min_relevance: f64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_bytes: Option<usize>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_tasks: Option<usize>,
}
fn default_context_mode() -> ContextMode {
ContextMode::Automatic
}
fn default_min_relevance() -> f64 {
0.5
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum ContextMode {
Automatic,
Manual,
None,
}
#[cfg(test)]
#[allow(clippy::field_reassign_with_default)]
mod tests {
use super::*;
#[test]
fn test_execution_mode_default() {
let mode: ExecutionMode = Default::default();
assert_eq!(mode, ExecutionMode::Sequential);
}
#[test]
fn test_permissions_spec_default() {
let perms: PermissionsSpec = Default::default();
assert_eq!(perms.mode, "default");
assert!(perms.allowed_directories.is_empty());
}
#[test]
fn test_task_spec_default() {
let task: TaskSpec = Default::default();
assert!(task.description.is_empty());
assert_eq!(task.priority, 0);
assert!(task.subtasks.is_empty());
}
#[test]
fn test_inherit_from_parent_agent() {
let mut parent = TaskSpec::default();
parent.agent = Some("parent_agent".to_string());
parent.description = "Parent task".to_string();
let mut child = TaskSpec::default();
child.description = "Child task".to_string();
child.inherit_from_parent(&parent);
assert_eq!(child.agent, Some("parent_agent".to_string()));
assert_eq!(child.description, "Child task");
}
#[test]
fn test_inherit_from_parent_agent_not_overridden() {
let mut parent = TaskSpec::default();
parent.agent = Some("parent_agent".to_string());
let mut child = TaskSpec::default();
child.agent = Some("child_agent".to_string());
child.inherit_from_parent(&parent);
assert_eq!(child.agent, Some("child_agent".to_string()));
}
#[test]
fn test_inherit_from_parent_inject_context() {
let mut parent = TaskSpec::default();
parent.inject_context = true;
let mut child = TaskSpec::default();
child.inject_context = false;
child.inherit_from_parent(&parent);
assert!(child.inject_context);
}
#[test]
fn test_inherit_from_parent_priority() {
let mut parent = TaskSpec::default();
parent.priority = 5;
let mut child = TaskSpec::default();
child.priority = 0;
child.inherit_from_parent(&parent);
assert_eq!(child.priority, 5);
}
#[test]
fn test_inherit_from_parent_priority_not_overridden() {
let mut parent = TaskSpec::default();
parent.priority = 5;
let mut child = TaskSpec::default();
child.priority = 10;
child.inherit_from_parent(&parent);
assert_eq!(child.priority, 10);
}
#[test]
fn test_inherit_from_parent_error_handling() {
let mut parent = TaskSpec::default();
parent.on_error = Some(ErrorHandlingSpec {
retry: 3,
retry_delay_secs: 5,
exponential_backoff: true,
fallback_agent: Some("fallback".to_string()),
});
let mut child = TaskSpec::default();
child.inherit_from_parent(&parent);
assert!(child.on_error.is_some());
let error_handling = child.on_error.as_ref().unwrap();
assert_eq!(error_handling.retry, 3);
assert_eq!(error_handling.retry_delay_secs, 5);
assert!(error_handling.exponential_backoff);
}
#[test]
fn test_inherit_from_parent_multiple_attributes() {
let mut parent = TaskSpec::default();
parent.agent = Some("shared_agent".to_string());
parent.inject_context = true;
parent.priority = 5;
parent.on_error = Some(ErrorHandlingSpec {
retry: 2,
retry_delay_secs: 3,
exponential_backoff: false,
fallback_agent: None,
});
let mut child = TaskSpec::default();
child.description = "Child task".to_string();
child.inherit_from_parent(&parent);
assert_eq!(child.agent, Some("shared_agent".to_string()));
assert!(child.inject_context);
assert_eq!(child.priority, 5);
assert!(child.on_error.is_some());
assert_eq!(child.on_error.as_ref().unwrap().retry, 2);
assert_eq!(child.description, "Child task");
}
#[test]
fn test_inherit_from_parent_loop_control() {
let mut parent = TaskSpec::default();
parent.loop_control = Some(LoopControl {
break_condition: None,
continue_condition: None,
collect_results: true,
result_key: Some("results".to_string()),
timeout_secs: Some(300),
checkpoint_interval: None,
});
let mut child = TaskSpec::default();
child.inherit_from_parent(&parent);
assert!(child.loop_control.is_some());
let control = child.loop_control.as_ref().unwrap();
assert!(control.collect_results);
assert_eq!(control.result_key, Some("results".to_string()));
assert_eq!(control.timeout_secs, Some(300));
}
#[test]
fn test_notification_spec_simple() {
let yaml = r#""Task completed successfully""#;
let spec: NotificationSpec = serde_yaml::from_str(yaml).unwrap();
match spec {
NotificationSpec::Simple(msg) => assert_eq!(msg, "Task completed successfully"),
_ => panic!("Expected Simple variant"),
}
}
#[test]
fn test_notification_spec_structured() {
let yaml = r#"
message: "Deployment finished"
title: "Production Update"
priority: high
channels:
- type: console
colored: true
timestamp: true
"#;
let spec: NotificationSpec = serde_yaml::from_str(yaml).unwrap();
match spec {
NotificationSpec::Structured {
message,
title,
priority,
channels,
..
} => {
assert_eq!(message, "Deployment finished");
assert_eq!(title, Some("Production Update".to_string()));
assert_eq!(priority, Some(NotificationPriority::High));
assert_eq!(channels.len(), 1);
}
_ => panic!("Expected Structured variant"),
}
}
#[test]
fn test_notification_channel_console() {
let yaml = r#"
type: console
colored: true
timestamp: false
"#;
let channel: NotificationChannel = serde_yaml::from_str(yaml).unwrap();
match channel {
NotificationChannel::Console { colored, timestamp } => {
assert!(colored);
assert!(!timestamp);
}
_ => panic!("Expected Console variant"),
}
}
#[test]
fn test_notification_channel_slack() {
let yaml = r##"
type: slack
credential: "${secret.slack_webhook}"
channel: "#general"
method: webhook
"##;
let channel: NotificationChannel = serde_yaml::from_str(yaml).unwrap();
match channel {
NotificationChannel::Slack {
credential,
channel,
method,
..
} => {
assert_eq!(credential, "${secret.slack_webhook}");
assert_eq!(channel, "#general");
assert_eq!(method, SlackMethod::Webhook);
}
_ => panic!("Expected Slack variant"),
}
}
#[test]
fn test_notification_channel_ntfy() {
let yaml = r#"
type: ntfy
server: "https://ntfy.sh"
topic: "workflow-updates"
priority: 4
tags: ["tada", "rocket"]
markdown: true
"#;
let channel: NotificationChannel = serde_yaml::from_str(yaml).unwrap();
match channel {
NotificationChannel::Ntfy {
server,
topic,
priority,
tags,
markdown,
..
} => {
assert_eq!(server, "https://ntfy.sh");
assert_eq!(topic, "workflow-updates");
assert_eq!(priority, Some(4));
assert_eq!(tags.len(), 2);
assert!(markdown);
}
_ => panic!("Expected Ntfy variant"),
}
}
#[test]
fn test_notification_defaults() {
let defaults = NotificationDefaults {
notify_on_completion: true,
notify_on_failure: true,
notify_on_start: false,
notify_on_workflow_completion: true,
default_channels: vec![],
};
assert!(defaults.notify_on_completion);
assert!(defaults.notify_on_failure);
assert!(!defaults.notify_on_start);
assert!(defaults.notify_on_workflow_completion);
}
#[test]
fn test_action_spec_with_notification() {
let yaml = r#"
notify: "Build completed"
"#;
let action: ActionSpec = serde_yaml::from_str(yaml).unwrap();
assert!(action.notify.is_some());
match action.notify.unwrap() {
NotificationSpec::Simple(msg) => assert_eq!(msg, "Build completed"),
_ => panic!("Expected Simple notification"),
}
}
#[test]
fn test_workflow_with_notification_defaults() {
let yaml = r#"
name: "Test Workflow"
version: "1.0.0"
notifications:
notify_on_completion: true
notify_on_failure: true
default_channels:
- type: console
colored: true
timestamp: true
"#;
let workflow: DSLWorkflow = serde_yaml::from_str(yaml).unwrap();
assert!(workflow.notifications.is_some());
let notif = workflow.notifications.unwrap();
assert!(notif.notify_on_completion);
assert!(notif.notify_on_failure);
assert_eq!(notif.default_channels.len(), 1);
}
}