1mod prompt;
11mod runner;
12mod store;
13mod tool;
14
15pub use prompt::{authoring_prompt, workflow_trigger};
16pub use store::{SaveLocation, SavedWorkflow, discover, save};
17
18pub 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
33pub type RunId = u64;
35
36const MAX_RETAINED_RUNS: usize = 20;
39
40pub struct RunRecord {
42 pub id: RunId,
43 pub name: String,
44 pub description: String,
45 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 pub fn journal(&self) -> Journal {
85 self.workflow.journal()
86 }
87}
88
89#[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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
104pub enum ApprovalDecision {
105 Approve,
106 Cancel,
107}
108
109#[derive(Debug, Clone)]
111pub struct WorkflowPlan {
112 pub name: String,
113 pub description: String,
114 pub phases: Vec<String>,
115 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 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
145pub type WorkflowApprover = Arc<
148 dyn Fn(WorkflowPlan) -> futures::future::BoxFuture<'static, ApprovalDecision> + Send + Sync,
149>;
150
151pub 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 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 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 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 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 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 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 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 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 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 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}