Skip to main content

roder_core/
dynamic_workflows.rs

1use async_trait::async_trait;
2use roder_api::catalog::REASONING_NONE;
3use roder_api::dynamic_workflows::{
4    WorkflowCostEstimate, WorkflowRunLimits, WorkflowRunStatus, WorkflowScript,
5    WorkflowScriptSource, WorkflowScriptSourceKind,
6};
7use roder_api::inference::{InstructionBundle, ReasoningConfig, RuntimeProfile};
8use roder_dynamic_workflows::{
9    WorkflowDefinition, WorkflowRuntimeOptions, parse_workflow_definition, workflow_script_hash,
10};
11use time::OffsetDateTime;
12
13use crate::speed_policy::reasoning_for_supported_effort;
14
15#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
16pub enum DynamicWorkflowEffortProfile {
17    #[default]
18    Standard,
19    Ultracode,
20}
21
22impl DynamicWorkflowEffortProfile {
23    pub fn as_str(self) -> &'static str {
24        match self {
25            Self::Standard => "standard",
26            Self::Ultracode => "ultracode",
27        }
28    }
29
30    pub fn auto_workflows_enabled(self) -> bool {
31        matches!(self, Self::Ultracode)
32    }
33}
34
35#[derive(Debug, Clone, PartialEq)]
36pub struct RuntimeDynamicWorkflowConfig {
37    pub enabled: bool,
38    pub trigger_word_enabled: bool,
39    pub auto_with_ultracode: bool,
40    pub effort_profile: DynamicWorkflowEffortProfile,
41    pub limits: WorkflowRunLimits,
42}
43
44impl Default for RuntimeDynamicWorkflowConfig {
45    fn default() -> Self {
46        Self {
47            enabled: true,
48            trigger_word_enabled: true,
49            auto_with_ultracode: true,
50            effort_profile: DynamicWorkflowEffortProfile::Standard,
51            limits: WorkflowRunLimits::default(),
52        }
53    }
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub enum WorkflowTriggerKind {
58    TriggerWord,
59    UltracodeAuto,
60}
61
62#[derive(Debug, Clone, Copy, PartialEq, Eq)]
63pub enum WorkflowTriggerSuppression {
64    Disabled,
65    NoTrigger,
66    SlashCommand,
67    ApprovalReply,
68    TinyCommand,
69    PureChatQuestion,
70}
71
72#[derive(Debug, Clone, Copy, PartialEq, Eq)]
73pub enum WorkflowTriggerDecision {
74    Plan(WorkflowTriggerKind),
75    Ignore(WorkflowTriggerSuppression),
76}
77
78impl WorkflowTriggerDecision {
79    pub fn should_plan(self) -> bool {
80        matches!(self, Self::Plan(_))
81    }
82}
83
84impl RuntimeDynamicWorkflowConfig {
85    pub fn classify_trigger(&self, input: &str) -> WorkflowTriggerDecision {
86        classify_workflow_trigger(input, self)
87    }
88}
89
90pub fn classify_workflow_trigger(
91    input: &str,
92    config: &RuntimeDynamicWorkflowConfig,
93) -> WorkflowTriggerDecision {
94    if !config.enabled {
95        return WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::Disabled);
96    }
97
98    let trimmed = input.trim();
99    if trimmed.is_empty() {
100        return WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::TinyCommand);
101    }
102    if trimmed.starts_with('/') {
103        return WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::SlashCommand);
104    }
105    if is_approval_reply(trimmed) {
106        return WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::ApprovalReply);
107    }
108    if is_pure_chat_question(trimmed) {
109        return WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::PureChatQuestion);
110    }
111
112    let words = word_count(trimmed);
113    if words <= 2 {
114        return WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::TinyCommand);
115    }
116
117    if config.trigger_word_enabled && contains_word(trimmed, "workflow") {
118        return WorkflowTriggerDecision::Plan(WorkflowTriggerKind::TriggerWord);
119    }
120
121    if config.auto_with_ultracode
122        && config.effort_profile.auto_workflows_enabled()
123        && is_substantive_task(trimmed, words)
124    {
125        return WorkflowTriggerDecision::Plan(WorkflowTriggerKind::UltracodeAuto);
126    }
127
128    WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::NoTrigger)
129}
130
131#[derive(Debug, Clone, PartialEq, Eq)]
132pub struct UltracodeReasoningDecision {
133    pub desired_reasoning: String,
134    pub applied_reasoning: Option<String>,
135    pub supported: bool,
136    pub reasoning: ReasoningConfig,
137}
138
139pub fn ultracode_reasoning_decision(
140    model: &str,
141    desired_reasoning: &str,
142    fallback: ReasoningConfig,
143) -> UltracodeReasoningDecision {
144    let (reasoning, supported) = reasoning_for_supported_effort(model, desired_reasoning, fallback);
145    UltracodeReasoningDecision {
146        desired_reasoning: desired_reasoning.to_string(),
147        applied_reasoning: supported.then(|| desired_reasoning.to_string()),
148        supported,
149        reasoning,
150    }
151}
152
153pub fn ultracode_reasoning_level_for_model(
154    model: &str,
155    desired_reasoning: &str,
156    fallback_level: &str,
157) -> String {
158    let fallback = reasoning_config_from_level(fallback_level);
159    ultracode_reasoning_decision(model, desired_reasoning, fallback)
160        .reasoning
161        .level
162        .unwrap_or_else(|| REASONING_NONE.to_string())
163}
164
165#[derive(Debug, Clone, PartialEq, Eq)]
166pub struct WorkflowModelPolicy {
167    pub provider: String,
168    pub model: String,
169    pub reasoning: ReasoningConfig,
170    pub runtime_profile: RuntimeProfile,
171    pub effort_profile: DynamicWorkflowEffortProfile,
172}
173
174#[derive(Debug, Clone, PartialEq)]
175pub struct WorkflowPlannerRequest {
176    pub prompt: String,
177    pub workspace: Option<String>,
178    pub source_path: Option<String>,
179    pub provider: String,
180    pub model: String,
181    pub reasoning: ReasoningConfig,
182    pub runtime_profile: RuntimeProfile,
183    pub effort_profile: DynamicWorkflowEffortProfile,
184    pub trigger: WorkflowTriggerKind,
185    pub instructions: InstructionBundle,
186    pub limits: WorkflowRunLimits,
187}
188
189#[derive(Debug, Clone, PartialEq)]
190pub struct WorkflowPlanDraft {
191    pub status: WorkflowRunStatus,
192    pub script: WorkflowScript,
193    pub definition: WorkflowDefinition,
194    pub phase_names: Vec<String>,
195    pub capability_scope: Vec<String>,
196    pub model_policy: WorkflowModelPolicy,
197    pub cost_estimate: WorkflowCostEstimate,
198}
199
200#[async_trait]
201pub trait WorkflowScriptDraftSource: Send + Sync {
202    async fn draft_workflow_script(
203        &self,
204        request: &WorkflowPlannerRequest,
205    ) -> anyhow::Result<String>;
206}
207
208#[derive(Debug, Clone)]
209pub struct WorkflowPlanner<P> {
210    draft_source: P,
211    options: WorkflowRuntimeOptions,
212}
213
214impl<P> WorkflowPlanner<P> {
215    pub fn new(draft_source: P) -> Self {
216        Self {
217            draft_source,
218            options: WorkflowRuntimeOptions::default(),
219        }
220    }
221
222    pub fn with_options(mut self, options: WorkflowRuntimeOptions) -> Self {
223        self.options = options;
224        self
225    }
226}
227
228impl<P> WorkflowPlanner<P>
229where
230    P: WorkflowScriptDraftSource,
231{
232    pub async fn plan(&self, request: WorkflowPlannerRequest) -> anyhow::Result<WorkflowPlanDraft> {
233        let source = self.draft_source.draft_workflow_script(&request).await?;
234        let mut options = self.options.clone();
235        options.limits = request.limits.clone();
236        let definition = parse_workflow_definition(&source, &options)
237            .map_err(|err| anyhow::anyhow!("generated workflow script failed validation: {err}"))?;
238        let hash = workflow_script_hash(&source);
239        let now = OffsetDateTime::now_utc();
240        let script = WorkflowScript {
241            script_id: generated_script_id(&hash),
242            name: definition.name.clone(),
243            description: definition.description.clone(),
244            source: WorkflowScriptSource {
245                kind: WorkflowScriptSourceKind::Generated,
246                path: request.source_path.clone(),
247                command_name: None,
248                extension_id: None,
249            },
250            hash,
251            host_api_version: definition.host_api_version,
252            arguments_schema: definition.arguments_schema.clone(),
253            body: Some(source.clone()),
254            limits: definition.limits.clone(),
255            created_at: now,
256            updated_at: now,
257        };
258
259        Ok(WorkflowPlanDraft {
260            status: WorkflowRunStatus::AwaitingApproval,
261            phase_names: phase_names(&definition),
262            capability_scope: capability_scope(&source),
263            model_policy: WorkflowModelPolicy {
264                provider: request.provider,
265                model: request.model,
266                reasoning: request.reasoning,
267                runtime_profile: request.runtime_profile,
268                effort_profile: request.effort_profile,
269            },
270            cost_estimate: cost_estimate(&definition, &source),
271            script,
272            definition,
273        })
274    }
275}
276
277pub fn workflow_planner_instructions(prompt: &str) -> String {
278    format!(
279        "Draft a Roder workflow script for this task. Return only JavaScript that calls workflow.define with phase names, limits, and an async handler. The script may only coordinate ctx.agents, ctx.checkpoint, and ctx.report host APIs.\n\nTask:\n{prompt}"
280    )
281}
282
283fn is_approval_reply(input: &str) -> bool {
284    matches!(
285        normalize_sentence(input).as_str(),
286        "y" | "yes"
287            | "ok"
288            | "okay"
289            | "approve"
290            | "approved"
291            | "run it"
292            | "go ahead"
293            | "continue"
294            | "deny"
295            | "denied"
296            | "no"
297    )
298}
299
300fn is_pure_chat_question(input: &str) -> bool {
301    let normalized = normalize_sentence(input);
302    if !normalized.ends_with('?') {
303        return false;
304    }
305    normalized.starts_with("what ")
306        || normalized.starts_with("what's ")
307        || normalized.starts_with("how ")
308        || normalized.starts_with("why ")
309        || normalized.starts_with("can you explain ")
310        || normalized.starts_with("explain ")
311        || normalized.starts_with("tell me about ")
312}
313
314fn is_substantive_task(input: &str, words: usize) -> bool {
315    if words < 8 {
316        return false;
317    }
318    let normalized = normalize_sentence(input);
319    [
320        "audit",
321        "build",
322        "create",
323        "implement",
324        "investigate",
325        "migrate",
326        "plan",
327        "research",
328        "review",
329        "verify",
330    ]
331    .iter()
332    .any(|verb| contains_word(&normalized, verb))
333}
334
335fn contains_word(input: &str, needle: &str) -> bool {
336    input
337        .split(|ch: char| !ch.is_ascii_alphanumeric())
338        .any(|word| word.eq_ignore_ascii_case(needle))
339}
340
341fn word_count(input: &str) -> usize {
342    input
343        .split(|ch: char| !ch.is_ascii_alphanumeric())
344        .filter(|word| !word.is_empty())
345        .count()
346}
347
348fn normalize_sentence(input: &str) -> String {
349    input
350        .split_whitespace()
351        .collect::<Vec<_>>()
352        .join(" ")
353        .to_ascii_lowercase()
354}
355
356fn reasoning_config_from_level(level: &str) -> ReasoningConfig {
357    match level {
358        "" | REASONING_NONE => ReasoningConfig::default(),
359        level => ReasoningConfig {
360            enabled: true,
361            level: Some(level.to_string()),
362        },
363    }
364}
365
366fn generated_script_id(hash: &str) -> String {
367    let short_hash = hash.get(..12).unwrap_or(hash);
368    format!("generated-{short_hash}")
369}
370
371fn phase_names(definition: &WorkflowDefinition) -> Vec<String> {
372    if definition.phases.is_empty() {
373        vec![definition.name.clone()]
374    } else {
375        definition.phases.clone()
376    }
377}
378
379fn capability_scope(source: &str) -> Vec<String> {
380    let mut capabilities = Vec::new();
381    if source.contains("ctx.agents") {
382        capabilities.push("childAgents".to_string());
383    }
384    if source.contains("ctx.checkpoint") {
385        capabilities.push("checkpoints".to_string());
386    }
387    if source.contains("ctx.report") {
388        capabilities.push("reports".to_string());
389    }
390    capabilities
391}
392
393fn cost_estimate(definition: &WorkflowDefinition, source: &str) -> WorkflowCostEstimate {
394    let uses_agents = source.contains("ctx.agents");
395    let phase_count = definition.phases.len().max(1) as u32;
396    let min_child_agents = if uses_agents { phase_count } else { 0 };
397    let max_child_agents = if uses_agents {
398        definition.limits.max_agents_per_run
399    } else {
400        0
401    };
402    let estimated_prompt_tokens = Some((source.len() as u64 / 4).max(1));
403    let warning = (max_child_agents >= 100).then(|| {
404        format!(
405            "This workflow may launch up to {max_child_agents} child agents; review cost and runtime before approving."
406        )
407    });
408
409    WorkflowCostEstimate {
410        min_child_agents,
411        max_child_agents,
412        estimated_prompt_tokens,
413        estimated_completion_tokens: None,
414        warning,
415    }
416}
417
418#[cfg(test)]
419mod tests {
420    use super::*;
421    use roder_api::catalog::REASONING_XHIGH;
422
423    #[test]
424    fn trigger_word_plans_but_chat_questions_and_commands_are_ignored() {
425        let config = RuntimeDynamicWorkflowConfig::default();
426
427        assert_eq!(
428            classify_workflow_trigger(
429                "Use a workflow to audit all crates for auth regressions",
430                &config
431            ),
432            WorkflowTriggerDecision::Plan(WorkflowTriggerKind::TriggerWord)
433        );
434        assert_eq!(
435            classify_workflow_trigger("What is a workflow?", &config),
436            WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::PureChatQuestion)
437        );
438        assert_eq!(
439            classify_workflow_trigger("/workflows", &config),
440            WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::SlashCommand)
441        );
442    }
443
444    #[test]
445    fn ultracode_auto_plans_only_for_substantive_tasks() {
446        let config = RuntimeDynamicWorkflowConfig {
447            effort_profile: DynamicWorkflowEffortProfile::Ultracode,
448            ..RuntimeDynamicWorkflowConfig::default()
449        };
450
451        assert_eq!(
452            classify_workflow_trigger(
453                "Audit the workspace for security issues and verify every finding with tests",
454                &config
455            ),
456            WorkflowTriggerDecision::Plan(WorkflowTriggerKind::UltracodeAuto)
457        );
458        assert_eq!(
459            classify_workflow_trigger("yes", &config),
460            WorkflowTriggerDecision::Ignore(WorkflowTriggerSuppression::ApprovalReply)
461        );
462    }
463
464    #[test]
465    fn ultracode_reasoning_degrades_on_models_without_xhigh() {
466        let fallback = ReasoningConfig::default();
467        let decision = ultracode_reasoning_decision("mock", REASONING_XHIGH, fallback.clone());
468
469        assert_eq!(decision.desired_reasoning, REASONING_XHIGH);
470        assert_eq!(decision.applied_reasoning, None);
471        assert!(!decision.supported);
472        assert_eq!(decision.reasoning, fallback);
473    }
474}