Skip to main content

kiss_coding/workflows/
mod.rs

1//! Dynamic workflows: many child agents orchestrated by a script the model
2//! writes.
3//!
4//! A workflow is built on the same child sessions `spawn_agent` uses, so it is
5//! available only when subagents are on. The difference is where the plan
6//! lives: with `spawn_agent` the model decides what to start next on every
7//! turn, and every result lands in its context. With a workflow the plan is a
8//! script, and only the final answer comes back.
9
10mod prompt;
11mod runner;
12mod store;
13mod tool;
14
15pub use prompt::{authoring_prompt, workflow_trigger};
16pub use store::{SaveLocation, SavedWorkflow, discover, save};
17
18// Re-exported so the terminal renders workflow state without depending on
19// `kiss-workflow` directly.
20pub use kiss_workflow::{
21    AgentId, AgentOutcome, AgentSnapshot, AgentStatus, Journal, PhaseSnapshot, RunSnapshot,
22    RunStatus, Script,
23};
24
25use crate::session_runner::{AgentSession, WorkflowTurnStatus};
26use kiss_agent::DynTool;
27use kiss_workflow::{Limits, Workflow};
28use serde_json::Value;
29use std::sync::atomic::{AtomicU64, Ordering};
30use std::sync::{Arc, Mutex, Weak};
31use std::time::Duration;
32
33/// Identifies one run within a session.
34pub type RunId = u64;
35
36/// How many finished runs to keep, so `/workflows` can still show and save a
37/// run the user is only now getting back to.
38const MAX_RETAINED_RUNS: usize = 20;
39
40/// One run, with the script that produced it.
41pub struct RunRecord {
42    pub id: RunId,
43    pub name: String,
44    pub description: String,
45    /// The script text, kept so the user can read it and save it afterward.
46    pub source: Arc<str>,
47    workflow: Arc<Workflow>,
48}
49
50impl RunRecord {
51    pub fn snapshot(&self) -> Arc<RunSnapshot> {
52        self.workflow.snapshot()
53    }
54
55    pub fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
56        self.workflow.subscribe()
57    }
58
59    pub fn pause(&self) {
60        self.workflow.pause();
61    }
62
63    pub fn resume(&self) {
64        self.workflow.resume();
65    }
66
67    pub fn is_paused(&self) -> bool {
68        self.workflow.is_paused()
69    }
70
71    pub fn stop(&self) {
72        self.workflow.stop();
73    }
74
75    pub fn stop_agent(&self, agent: kiss_workflow::AgentId) {
76        self.workflow.stop_agent(agent);
77    }
78
79    pub fn restart_agent(&self, agent: kiss_workflow::AgentId) {
80        self.workflow.restart_agent(agent);
81    }
82
83    /// Results worth reusing if this run is started again.
84    pub fn journal(&self) -> Journal {
85        self.workflow.journal()
86    }
87}
88
89/// A run as listed in `/workflows`.
90#[derive(Debug, Clone)]
91pub struct RunSummary {
92    pub id: RunId,
93    pub name: String,
94    pub description: String,
95    pub status: RunStatus,
96    pub total_agents: usize,
97    pub finished_agents: usize,
98    pub tokens: u64,
99    pub elapsed: Duration,
100}
101
102/// What the user decided when asked to approve a run.
103#[derive(Debug, Clone, Copy, PartialEq, Eq)]
104pub enum ApprovalDecision {
105    Approve,
106    Cancel,
107}
108
109/// What the user is shown before a run starts.
110#[derive(Debug, Clone)]
111pub struct WorkflowPlan {
112    pub name: String,
113    pub description: String,
114    pub phases: Vec<String>,
115    /// `None` when the agent count depends on data the run has not fetched yet.
116    pub estimated_agents: Option<u32>,
117    pub source: Arc<str>,
118}
119
120impl WorkflowPlan {
121    pub fn from_script(script: &Script) -> WorkflowPlan {
122        WorkflowPlan {
123            name: script.meta().name.clone(),
124            description: script.meta().description.clone(),
125            phases: script.declared_phases().to_vec(),
126            estimated_agents: script.estimated_agents(),
127            source: Arc::from(script.source()),
128        }
129    }
130
131    /// The agent count as shown to the user.
132    ///
133    /// A count is only stated when the script fixes it. Otherwise this says so,
134    /// because a number in an approval prompt has to be one the user can rely
135    /// on.
136    pub fn agent_estimate(&self) -> String {
137        match self.estimated_agents {
138            Some(1) => "1 agent".into(),
139            Some(count) => format!("{count} agents"),
140            None => "an unbounded number of agents".into(),
141        }
142    }
143}
144
145/// Asks the user to approve a run. Interactive mode installs one. Other modes
146/// leave it unset and runs start without asking, since nothing can answer.
147pub type WorkflowApprover = Arc<
148    dyn Fn(WorkflowPlan) -> futures::future::BoxFuture<'static, ApprovalDecision> + Send + Sync,
149>;
150
151/// The workflow runs belonging to one session.
152pub struct WorkflowRuntime {
153    parent: Weak<AgentSession>,
154    runs: Mutex<Vec<Arc<RunRecord>>>,
155    next_id: AtomicU64,
156}
157
158impl WorkflowRuntime {
159    pub(crate) fn new(parent: Weak<AgentSession>) -> Arc<WorkflowRuntime> {
160        Arc::new(WorkflowRuntime {
161            parent,
162            runs: Mutex::new(Vec::new()),
163            next_id: AtomicU64::new(1),
164        })
165    }
166
167    /// The tool the model calls to run a workflow it has written.
168    pub(crate) fn tool(self: &Arc<Self>) -> DynTool {
169        Arc::new(tool::RunWorkflowTool::new(self.clone()))
170    }
171
172    pub(crate) fn emit_outcome(
173        &self,
174        run: Option<RunId>,
175        name: String,
176        status: WorkflowTurnStatus,
177    ) {
178        if let Some(parent) = self.parent.upgrade() {
179            parent.emit_workflow_outcome(run, name, status);
180        }
181    }
182
183    pub(crate) fn limits(&self) -> Limits {
184        Limits::default()
185    }
186
187    /// Ask the user to approve a run.
188    ///
189    /// This is a cost gate rather than a permission gate: one run can start
190    /// hundreds of child agents, and the user should see that before it is
191    /// spent. Modes with no way to answer, such as `-p`, start the run.
192    pub async fn approve(&self, plan: WorkflowPlan) -> ApprovalDecision {
193        let Some(parent) = self.parent.upgrade() else {
194            return ApprovalDecision::Cancel;
195        };
196        if !parent.settings().workflows.confirm {
197            return ApprovalDecision::Approve;
198        }
199        match parent.workflow_approver() {
200            Some(approver) => approver(plan).await,
201            None => ApprovalDecision::Approve,
202        }
203    }
204
205    /// Prepare a run and add it to the list, without starting it.
206    ///
207    /// Registering before running is what lets the progress view find the run
208    /// as soon as it is approved, rather than after its first agent answers.
209    pub fn prepare(
210        &self,
211        script: Script,
212        args: Value,
213        journal: Journal,
214    ) -> anyhow::Result<Arc<RunRecord>> {
215        let parent = self
216            .parent
217            .upgrade()
218            .ok_or_else(|| anyhow::anyhow!("the parent session has closed"))?;
219        let limits = self.limits();
220        let runner = Arc::new(runner::SessionAgentRunner::new(
221            self.parent.clone(),
222            limits.max_concurrency,
223        ));
224        let cwd = parent
225            .manager
226            .lock()
227            .map(|manager| manager.cwd().display().to_string())
228            .unwrap_or_default();
229
230        let id = self.next_id.fetch_add(1, Ordering::SeqCst);
231        let name = script.meta().name.clone();
232        let description = script.meta().description.clone();
233        let source: Arc<str> = Arc::from(script.source());
234        let workflow = Arc::new(Workflow::with_journal(
235            script, args, cwd, runner, limits, journal,
236        ));
237        let record = Arc::new(RunRecord {
238            id,
239            name,
240            description,
241            source,
242            workflow,
243        });
244
245        // Forward the workflow's watch channel into the normal session event
246        // stream. The terminal then sleeps until state changes instead of
247        // polling every frame.
248        let mut updates = record.subscribe();
249        let parent_for_updates = self.parent.clone();
250        let record_for_updates = record.clone();
251        tokio::spawn(async move {
252            while updates.changed().await.is_ok() {
253                let version = *updates.borrow_and_update();
254                let Some(parent) = parent_for_updates.upgrade() else {
255                    break;
256                };
257                parent.emit_workflow(id, version);
258                if record_for_updates.snapshot().status.is_finished() {
259                    break;
260                }
261            }
262        });
263
264        let mut runs = self
265            .runs
266            .lock()
267            .map_err(|_| anyhow::anyhow!("the workflow list is unavailable"))?;
268        runs.push(record.clone());
269        // Keep the newest runs, but never drop one that is still working.
270        while runs.len() > MAX_RETAINED_RUNS {
271            let Some(position) = runs
272                .iter()
273                .position(|run| run.snapshot().status.is_finished())
274            else {
275                break;
276            };
277            runs.remove(position);
278        }
279        Ok(record)
280    }
281
282    /// Run a prepared workflow to completion.
283    pub async fn run(&self, record: &Arc<RunRecord>) -> Result<Value, String> {
284        record
285            .workflow
286            .run()
287            .await
288            .map_err(|error| error.to_string())
289    }
290
291    pub fn get(&self, id: RunId) -> Option<Arc<RunRecord>> {
292        self.runs
293            .lock()
294            .ok()?
295            .iter()
296            .find(|run| run.id == id)
297            .cloned()
298    }
299
300    /// Every run in this session, oldest first.
301    pub fn summaries(&self) -> Vec<RunSummary> {
302        let Ok(runs) = self.runs.lock() else {
303            return Vec::new();
304        };
305        runs.iter()
306            .map(|run| {
307                let snapshot = run.snapshot();
308                RunSummary {
309                    id: run.id,
310                    name: run.name.clone(),
311                    description: run.description.clone(),
312                    status: snapshot.status,
313                    total_agents: snapshot.total_agents(),
314                    finished_agents: snapshot.finished_agents(),
315                    tokens: snapshot.tokens,
316                    elapsed: snapshot.elapsed,
317                }
318            })
319            .collect()
320    }
321
322    /// The run still working, if any.
323    pub fn active(&self) -> Option<Arc<RunRecord>> {
324        let runs = self.runs.lock().ok()?;
325        runs.iter()
326            .rev()
327            .find(|run| !run.snapshot().status.is_finished())
328            .cloned()
329    }
330
331    /// The most recent run, working or finished.
332    pub fn latest(&self) -> Option<Arc<RunRecord>> {
333        self.runs.lock().ok()?.last().cloned()
334    }
335
336    pub fn is_empty(&self) -> bool {
337        self.runs.lock().map(|runs| runs.is_empty()).unwrap_or(true)
338    }
339
340    /// Stop every run still working. Called when the session closes or when the
341    /// setting is turned off.
342    pub(crate) fn stop_all(&self) {
343        let Ok(runs) = self.runs.lock() else {
344            return;
345        };
346        for run in runs.iter() {
347            if !run.snapshot().status.is_finished() {
348                run.stop();
349            }
350        }
351    }
352}
353
354#[cfg(test)]
355mod tests {
356    use super::*;
357
358    #[test]
359    fn an_agent_estimate_never_guesses_a_number_it_does_not_know() {
360        let plan = |estimated_agents| WorkflowPlan {
361            name: "x".into(),
362            description: "y".into(),
363            phases: Vec::new(),
364            estimated_agents,
365            source: Arc::from(""),
366        };
367        assert_eq!(plan(Some(1)).agent_estimate(), "1 agent");
368        assert_eq!(plan(Some(12)).agent_estimate(), "12 agents");
369        assert_eq!(plan(None).agent_estimate(), "an unbounded number of agents");
370    }
371}