1use std::collections::{BTreeMap, BTreeSet};
2
3use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5
6use crate::event::{Event, FlowNodeStatus, FlowStatus, TurnId};
7use crate::permission::{PermissionGroupId, PermissionRequestId};
8use crate::permission_audit::{PermissionGroupAudit, PermissionRequestAudit};
9
10#[cfg(test)]
11#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
12struct PerfCounters {
13 workflow_node_visits: u64,
14 permission_table_visits: u64,
15}
16
17#[cfg(test)]
18thread_local! {
19 static PERF_COUNTERS: std::cell::Cell<PerfCounters> = const {
20 std::cell::Cell::new(PerfCounters {
21 workflow_node_visits: 0,
22 permission_table_visits: 0,
23 })
24 };
25}
26
27#[cfg(test)]
28fn reset_perf_counters() {
29 PERF_COUNTERS.with(|counters| counters.set(PerfCounters::default()));
30}
31
32#[cfg(test)]
33fn perf_counters() -> PerfCounters {
34 PERF_COUNTERS.with(std::cell::Cell::get)
35}
36
37#[cfg(test)]
38fn update_perf_counters(update: impl FnOnce(&mut PerfCounters)) {
39 PERF_COUNTERS.with(|counters| {
40 let mut value = counters.get();
41 update(&mut value);
42 counters.set(value);
43 });
44}
45
46#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
47pub struct WorkflowGraph {
48 pub turn_id: TurnId,
49 pub root: Vec<WorkflowNode>,
50 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
51 pub permission_requests: BTreeMap<WorkflowPermissionIdentity, WorkflowPermissionRequest>,
52 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
53 pub permission_groups: BTreeMap<PermissionGroupId, PermissionGroupAudit>,
54 #[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
55 pub resolved_permission_groups: BTreeSet<PermissionGroupId>,
56}
57
58#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)]
59#[serde(tag = "kind", rename_all = "snake_case")]
60pub enum WorkflowPermissionIdentity {
61 Canonical {
62 request_id: PermissionRequestId,
63 },
64 Legacy {
65 seq: u64,
66 run_id: String,
67 tool_use_id: String,
68 },
69}
70
71#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
72pub struct WorkflowPermissionRequest {
73 pub payload: PermissionRequestAudit,
74 pub state: WorkflowPermissionState,
75}
76
77#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
78#[serde(rename_all = "snake_case")]
79pub enum WorkflowPermissionState {
80 Pending,
81 Approved,
82 Denied,
83 Cancelled,
84 Interrupted,
85 Unrestricted,
86}
87
88impl WorkflowPermissionState {
89 pub fn is_pending(self) -> bool {
90 matches!(self, Self::Pending)
91 }
92}
93
94#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
95pub struct WorkflowNode {
96 pub id: String,
97 pub kind: WorkflowNodeKind,
98 pub label: String,
99 pub status: NodeStatus,
100 pub started_at: Option<DateTime<Utc>>,
101 pub ended_at: Option<DateTime<Utc>>,
102 pub output_preview: Option<String>,
103 pub children: Vec<WorkflowNode>,
104 pub parallelism: Parallelism,
105 #[serde(default, skip_serializing_if = "Option::is_none")]
106 pub approval: Option<ApprovalState>,
107 #[serde(default, skip_serializing_if = "Option::is_none")]
108 pub llm_stats: Option<LlmStats>,
109}
110
111#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
112pub struct LlmStats {
113 pub model: String,
114 #[serde(default)]
115 pub provider: String,
116 #[serde(default)]
117 pub context_call_purpose: crate::context_plan::ContextCallPurpose,
118 #[serde(default)]
119 pub context_call_scope: crate::context_plan::ContextCallScope,
120 pub input_tokens: u64,
121 pub output_tokens: u64,
122 pub cache_read: u64,
123 pub cache_write: u64,
124 pub ttft_ms: u64,
125 pub tokens_per_second: f64,
126 pub wallclock_ms: u64,
127}
128
129#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
130#[serde(tag = "type", rename_all = "snake_case")]
131pub enum WorkflowNodeKind {
132 Flow {
133 run_id: String,
134 flow_name: String,
135 },
136 Stmt {
137 node_kind: crate::nodegraph::NodeKind,
138 },
139 ToolCall {
140 tool_use_id: String,
141 tool: String,
142 args_preview: String,
143 #[serde(default, skip_serializing_if = "Option::is_none")]
144 call_intent: Option<crate::message::ToolCallIntent>,
145 result_preview: Option<String>,
146 },
147 Subflow {
148 run_id: String,
149 flow_name: String,
150 },
151 FanoutBranch {
152 branch_index: usize,
153 },
154}
155
156#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
157#[serde(rename_all = "snake_case")]
158pub enum NodeStatus {
159 Pending,
160 Running,
161 Ok,
162 Err,
163 Cancelled,
164}
165
166#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
167#[serde(rename_all = "snake_case", tag = "kind")]
168pub enum ApprovalState {
169 Pending {
170 level: String,
171 #[serde(default, skip_serializing_if = "Option::is_none")]
172 preview: Option<String>,
173 },
174 Approved,
175 Denied {
176 reason: String,
177 },
178}
179
180#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
181#[serde(rename_all = "snake_case")]
182pub enum Parallelism {
183 Serial,
184 Parallel,
185}
186
187impl WorkflowGraph {
188 pub fn new(turn_id: TurnId) -> Self {
189 Self {
190 turn_id,
191 root: Vec::new(),
192 permission_requests: BTreeMap::new(),
193 permission_groups: BTreeMap::new(),
194 resolved_permission_groups: BTreeSet::new(),
195 }
196 }
197
198 pub fn apply_event(&mut self, event: &Event) {
199 match event {
200 Event::FlowStart {
201 run_id,
202 flow_name,
203 parent_run_id,
204 parent_node_id,
205 ..
206 } => {
207 let run_id_str = run_id.0.to_string();
208 let node = WorkflowNode {
209 id: run_id_str.clone(),
210 kind: if parent_run_id.is_some() {
211 WorkflowNodeKind::Subflow {
212 run_id: run_id_str,
213 flow_name: flow_name.clone(),
214 }
215 } else {
216 WorkflowNodeKind::Flow {
217 run_id: run_id_str,
218 flow_name: flow_name.clone(),
219 }
220 },
221 label: flow_name.clone(),
222 status: NodeStatus::Running,
223 started_at: Some(chrono::Utc::now()),
224 ended_at: None,
225 output_preview: None,
226 children: Vec::new(),
227 parallelism: Parallelism::Serial,
228 approval: None,
229 llm_stats: None,
230 };
231 match (parent_run_id.as_ref(), parent_node_id.as_deref()) {
232 (Some(prid), Some(pid)) => {
233 let scoped = scope_id(&prid.0.to_string(), pid);
234 if let Some(parent) = find_node_mut(&mut self.root, &scoped) {
235 parent.children.push(node);
236 }
237 }
238 _ => self.root.push(node),
239 }
240 }
241 Event::FlowEnd { run_id, status, .. } => {
242 let id = run_id.0.to_string();
243 if let Some(n) = find_node_mut(&mut self.root, &id) {
244 let new_status = match status {
245 FlowStatus::Ok => NodeStatus::Ok,
246 FlowStatus::Errored { .. } => NodeStatus::Err,
247 FlowStatus::Cancelled => NodeStatus::Cancelled,
248 };
249 n.status = new_status;
250 n.ended_at = Some(chrono::Utc::now());
251 for child in n.children.iter_mut() {
252 if matches!(child.status, NodeStatus::Running | NodeStatus::Pending) {
253 child.status = new_status;
254 child.ended_at = Some(chrono::Utc::now());
255 }
256 }
257 }
258 }
259 Event::FlowNodeStart {
260 run_id,
261 node_id,
262 kind: nk,
263 label,
264 parent_node_id,
265 ..
266 } => {
267 let rid = run_id.0.to_string();
268 let scoped_id = scope_id(&rid, node_id);
269 let parent_id = parent_node_id
270 .as_deref()
271 .map(|p| scope_id(&rid, p))
272 .unwrap_or_else(|| rid.clone());
273 let kind = if let Some(idx) = parse_branch_index(node_id) {
274 WorkflowNodeKind::FanoutBranch { branch_index: idx }
275 } else {
276 WorkflowNodeKind::Stmt {
277 node_kind: nk.clone(),
278 }
279 };
280 let node = WorkflowNode {
281 id: scoped_id,
282 kind,
283 label: label.clone(),
284 status: NodeStatus::Running,
285 started_at: Some(chrono::Utc::now()),
286 ended_at: None,
287 output_preview: None,
288 children: Vec::new(),
289 parallelism: Parallelism::Serial,
290 approval: None,
291 llm_stats: None,
292 };
293 if let Some(parent) = find_node_mut(&mut self.root, &parent_id) {
294 if matches!(node.kind, WorkflowNodeKind::FanoutBranch { .. }) {
295 parent.parallelism = Parallelism::Parallel;
296 }
297 parent.children.push(node);
298 }
299 }
300 Event::FlowNodeEnd {
301 run_id,
302 node_id,
303 status,
304 output_preview,
305 ..
306 } => {
307 let scoped = scope_id(&run_id.0.to_string(), node_id);
308 if let Some(n) = find_node_mut(&mut self.root, &scoped) {
309 let new_status = match status {
310 FlowNodeStatus::Ok => NodeStatus::Ok,
311 FlowNodeStatus::Err => NodeStatus::Err,
312 FlowNodeStatus::Cancelled => NodeStatus::Cancelled,
313 };
314 n.status = new_status;
315 n.ended_at = Some(chrono::Utc::now());
316 if let Some(p) = output_preview {
317 n.output_preview = Some(p.clone());
318 }
319 for child in n.children.iter_mut() {
320 if matches!(child.status, NodeStatus::Running | NodeStatus::Pending) {
321 child.status = new_status;
322 child.ended_at = Some(chrono::Utc::now());
323 }
324 }
325 }
326 }
327 Event::LlmCall {
328 run_id,
329 node_id,
330 model,
331 provider,
332 context_call_purpose,
333 context_call_identity,
334 usage,
335 wallclock_ms,
336 ttft_ms,
337 tokens_per_second,
338 ..
339 } => {
340 if let (Some(rid), Some(nid)) = (run_id, node_id) {
341 let scoped = scope_id(&rid.0.to_string(), nid);
342 if let Some(n) = find_node_mut(&mut self.root, &scoped) {
343 n.llm_stats = Some(LlmStats {
344 model: model.clone(),
345 provider: provider.clone(),
346 context_call_purpose: context_call_purpose.unwrap_or_default(),
347 context_call_scope: context_call_identity
348 .as_ref()
349 .map(|identity| identity.scope)
350 .unwrap_or_else(|| {
351 if run_id.is_some() {
352 crate::context_plan::ContextCallScope::Root
353 } else {
354 crate::context_plan::ContextCallScope::Detached
355 }
356 }),
357 input_tokens: usage.input,
358 output_tokens: usage.output,
359 cache_read: usage.cached_input,
360 cache_write: usage.cache_write,
361 ttft_ms: ttft_ms.unwrap_or(0),
362 tokens_per_second: tokens_per_second.unwrap_or(0.0),
363 wallclock_ms: *wallclock_ms,
364 });
365 }
366 }
367 }
368 Event::ToolNode {
369 run_id,
370 parent_node_id,
371 tool_use_id,
372 tool_name,
373 args_preview,
374 call_intent,
375 ..
376 } => {
377 let rid = run_id.0.to_string();
378 let scoped_parent = scope_id(&rid, parent_node_id);
379 let id = tool_node_id(&rid, tool_use_id);
380 if find_node(&self.root, &id).is_some() {
381 return;
382 }
383 let node = WorkflowNode {
384 id,
385 kind: WorkflowNodeKind::ToolCall {
386 tool_use_id: tool_use_id.clone(),
387 tool: tool_name.clone(),
388 args_preview: args_preview.clone(),
389 call_intent: call_intent.clone(),
390 result_preview: None,
391 },
392 label: tool_name.clone(),
393 status: NodeStatus::Running,
394 started_at: Some(chrono::Utc::now()),
395 ended_at: None,
396 output_preview: None,
397 children: Vec::new(),
398 parallelism: Parallelism::Serial,
399 approval: None,
400 llm_stats: None,
401 };
402 if let Some(parent) = find_node_mut(&mut self.root, &scoped_parent) {
403 parent.children.push(node);
404 self.refresh_permission_tool_approvals();
405 }
406 }
407 Event::AssistantMsg { .. } => {}
408 Event::ToolResultMsg {
409 flow_run_id,
410 message,
411 ..
412 } => {
413 let flow_id = flow_run_id.as_ref().map(|r| r.0.to_string());
414 for part in &message.parts {
415 if let crate::message::MessagePart::ToolResult {
416 tool_use_id,
417 content,
418 is_error,
419 } = part
420 {
421 let node = match flow_id.as_deref() {
422 Some(rid) => {
423 let id = tool_node_id(rid, tool_use_id);
424 find_node_mut(&mut self.root, &id)
425 }
426 None => find_tool_node_by_tool_use_id(&mut self.root, tool_use_id),
427 };
428 if let Some(n) = node {
429 n.status = if *is_error {
430 NodeStatus::Err
431 } else {
432 NodeStatus::Ok
433 };
434 n.ended_at = Some(chrono::Utc::now());
435 let preview: String = content.chars().take(300).collect();
436 n.output_preview = Some(preview.clone());
437 if let WorkflowNodeKind::ToolCall { result_preview, .. } = &mut n.kind {
438 *result_preview = Some(preview);
439 }
440 }
441 }
442 }
443 }
444 Event::ToolPendingApproval {
445 run_id,
446 tool_use_id,
447 level,
448 preview,
449 ..
450 } => {
451 let rid = run_id.0.to_string();
452 let id = tool_node_id(&rid, tool_use_id);
453 if let Some(n) = find_node_mut(&mut self.root, &id) {
454 n.approval = Some(ApprovalState::Pending {
455 level: level.clone(),
456 preview: preview.clone(),
457 });
458 }
459 }
460 Event::ToolApproved {
461 run_id,
462 tool_use_id,
463 ..
464 } => {
465 let rid = run_id.0.to_string();
466 let id = tool_node_id(&rid, tool_use_id);
467 if let Some(n) = find_node_mut(&mut self.root, &id) {
468 n.approval = Some(ApprovalState::Approved);
469 }
470 }
471 Event::ToolDenied {
472 run_id,
473 tool_use_id,
474 reason,
475 ..
476 } => {
477 let rid = run_id.0.to_string();
478 let id = tool_node_id(&rid, tool_use_id);
479 if let Some(n) = find_node_mut(&mut self.root, &id) {
480 n.approval = Some(ApprovalState::Denied {
481 reason: reason.clone(),
482 });
483 }
484 }
485 Event::PermissionRequestCreated { payload }
486 | Event::PermissionRequestTargeted { payload }
487 | Event::PermissionRequestDeferred { payload } => {
488 self.apply_permission_request(payload, WorkflowPermissionState::Pending);
489 }
490 Event::PermissionRequestApproved { payload } => {
491 self.apply_permission_request(payload, WorkflowPermissionState::Approved);
492 }
493 Event::PermissionRequestDenied { payload } => {
494 self.apply_permission_request(payload, WorkflowPermissionState::Denied);
495 }
496 Event::PermissionRequestCancelled { payload } => {
497 self.apply_permission_request(payload, WorkflowPermissionState::Cancelled);
498 }
499 Event::UnrestrictedExecution { payload } => {
500 self.apply_permission_request(payload, WorkflowPermissionState::Unrestricted);
501 }
502 Event::PermissionGroupCreated { payload }
503 | Event::PermissionGroupUpdated { payload } => {
504 self.apply_permission_group(payload, false);
505 }
506 Event::PermissionGroupResolved { payload } => {
507 self.apply_permission_group(payload, true);
508 }
509 _ => {}
510 }
511 }
512
513 pub fn find_node(&self, id: &str) -> Option<&WorkflowNode> {
514 find_node(&self.root, id)
515 }
516
517 pub fn find_node_mut(&mut self, id: &str) -> Option<&mut WorkflowNode> {
518 find_node_mut(&mut self.root, id)
519 }
520
521 pub fn descendant_pending_permissions(&self, flow_node_id: &str) -> usize {
522 let Some(node) = find_node(&self.root, flow_node_id) else {
523 return 0;
524 };
525 let mut run_ids = Vec::new();
526 collect_flow_run_ids(node, &mut run_ids);
527 self.permission_requests
528 .values()
529 .filter(|request| {
530 request.state.is_pending()
531 && run_ids.contains(&request.payload.requesting_run_id.0.to_string())
532 })
533 .count()
534 }
535
536 pub fn interrupt_pending_permissions(&mut self) {
537 for request in self.permission_requests.values_mut() {
538 if request.state.is_pending() {
539 request.state = WorkflowPermissionState::Interrupted;
540 request.payload.reason = Some("interrupted at end of persisted history".into());
541 }
542 }
543 self.refresh_permission_tool_approvals();
544 }
545
546 pub fn apply_permission_request(
547 &mut self,
548 payload: &PermissionRequestAudit,
549 state: WorkflowPermissionState,
550 ) {
551 let Some(request_id) = payload.request_id.clone() else {
552 return;
553 };
554 self.apply_permission_request_with_identity(
555 WorkflowPermissionIdentity::Canonical { request_id },
556 payload,
557 state,
558 );
559 }
560
561 pub fn apply_permission_request_with_identity(
562 &mut self,
563 identity: WorkflowPermissionIdentity,
564 payload: &PermissionRequestAudit,
565 state: WorkflowPermissionState,
566 ) {
567 self.apply_permission_requests([(identity, payload.clone(), state)]);
568 }
569
570 pub fn apply_permission_requests(
571 &mut self,
572 requests: impl IntoIterator<
573 Item = (
574 WorkflowPermissionIdentity,
575 PermissionRequestAudit,
576 WorkflowPermissionState,
577 ),
578 >,
579 ) {
580 for (identity, payload, state) in requests {
581 self.permission_requests
582 .insert(identity, WorkflowPermissionRequest { payload, state });
583 }
584 self.refresh_permission_tool_approvals();
585 }
586
587 pub fn apply_permission_group(&mut self, payload: &PermissionGroupAudit, resolved: bool) {
588 if resolved {
589 self.permission_groups.remove(&payload.group_id);
590 self.resolved_permission_groups
591 .insert(payload.group_id.clone());
592 } else if !self.resolved_permission_groups.contains(&payload.group_id) {
593 self.permission_groups
594 .insert(payload.group_id.clone(), payload.clone());
595 }
596 }
597
598 fn refresh_permission_tool_approvals(&mut self) {
599 if self.permission_requests.is_empty() {
600 return;
601 }
602 let mut exact: BTreeMap<(String, String), &WorkflowPermissionRequest> = BTreeMap::new();
603 for request in self.permission_requests.values() {
604 #[cfg(test)]
605 update_perf_counters(|counters| {
606 counters.permission_table_visits =
607 counters.permission_table_visits.saturating_add(1);
608 });
609 let key = (
610 request.payload.requesting_run_id.0.to_string(),
611 request.payload.tool_use_id.clone(),
612 );
613 exact
614 .entry(key)
615 .and_modify(|current| {
616 if (!current.state.is_pending() && request.state.is_pending())
617 || (current.state.is_pending() == request.state.is_pending()
618 && request.payload.at > current.payload.at)
619 {
620 *current = request;
621 }
622 })
623 .or_insert(request);
624 }
625 let approvals = exact
626 .into_iter()
627 .map(|((run_id, tool_use_id), request)| {
628 let approval = match request.state {
629 WorkflowPermissionState::Pending => ApprovalState::Pending {
630 level: format!("{:?}", request.payload.tier).to_lowercase(),
631 preview: None,
632 },
633 WorkflowPermissionState::Approved | WorkflowPermissionState::Unrestricted => {
634 ApprovalState::Approved
635 }
636 WorkflowPermissionState::Denied
637 | WorkflowPermissionState::Cancelled
638 | WorkflowPermissionState::Interrupted => ApprovalState::Denied {
639 reason: request
640 .payload
641 .reason
642 .clone()
643 .unwrap_or_else(|| "permission denied".into()),
644 },
645 };
646 (tool_node_id(&run_id, &tool_use_id), approval)
647 })
648 .collect();
649 apply_permission_approvals(&mut self.root, &approvals);
650 }
651
652 pub fn apply_stream_frame(&mut self, frame: &crate::stream::StreamFrame) {
653 self.apply_stream_frame_at(frame, None);
654 }
655
656 pub fn apply_stream_frame_at(
657 &mut self,
658 frame: &crate::stream::StreamFrame,
659 override_ts: Option<chrono::DateTime<chrono::Utc>>,
660 ) {
661 use crate::stream::StreamFrame;
662 let now = override_ts.unwrap_or_else(Utc::now);
663 match frame {
664 StreamFrame::FlowGraph { run_id, graph } => {
665 if self.find_node(run_id).is_none() {
666 self.root.push(WorkflowNode {
667 id: run_id.clone(),
668 kind: WorkflowNodeKind::Flow {
669 run_id: run_id.clone(),
670 flow_name: graph.flow_name.clone(),
671 },
672 label: graph.flow_name.clone(),
673 status: NodeStatus::Running,
674 started_at: Some(now),
675 ended_at: None,
676 output_preview: None,
677 children: Vec::new(),
678 parallelism: Parallelism::Serial,
679 approval: None,
680 llm_stats: None,
681 });
682 }
683 }
684 StreamFrame::FlowStart {
685 run_id,
686 flow_name,
687 parent_run_id,
688 parent_node_id,
689 } => {
690 if self.find_node(run_id).is_some() {
691 return;
692 }
693 let kind = if parent_run_id.is_some() {
694 WorkflowNodeKind::Subflow {
695 run_id: run_id.clone(),
696 flow_name: flow_name.clone(),
697 }
698 } else {
699 WorkflowNodeKind::Flow {
700 run_id: run_id.clone(),
701 flow_name: flow_name.clone(),
702 }
703 };
704 let node = WorkflowNode {
705 id: run_id.clone(),
706 kind,
707 label: flow_name.clone(),
708 status: NodeStatus::Running,
709 started_at: Some(now),
710 ended_at: None,
711 output_preview: None,
712 children: Vec::new(),
713 parallelism: Parallelism::Serial,
714 approval: None,
715 llm_stats: None,
716 };
717 match (parent_run_id.as_deref(), parent_node_id.as_deref()) {
718 (Some(prid), Some(pid)) => {
719 let scoped = scope_id(prid, pid);
720 if let Some(parent) = find_node_mut(&mut self.root, &scoped) {
721 parent.children.push(node);
722 } else {
723 self.root.push(node);
724 }
725 }
726 _ => self.root.push(node),
727 }
728 }
729 StreamFrame::FlowNodeStart {
730 run_id,
731 node_id,
732 kind: nk,
733 label,
734 parent_node_id,
735 } => {
736 let scoped_id = scope_id(run_id, node_id);
737 let parent_id = parent_node_id
738 .as_deref()
739 .map(|p| scope_id(run_id, p))
740 .unwrap_or_else(|| run_id.clone());
741 let kind = if let Some(idx) = parse_branch_index(node_id) {
742 WorkflowNodeKind::FanoutBranch { branch_index: idx }
743 } else {
744 WorkflowNodeKind::Stmt {
745 node_kind: nk.clone(),
746 }
747 };
748 let node = WorkflowNode {
749 id: scoped_id,
750 kind,
751 label: label.clone(),
752 status: NodeStatus::Running,
753 started_at: Some(now),
754 ended_at: None,
755 output_preview: None,
756 children: Vec::new(),
757 parallelism: Parallelism::Serial,
758 approval: None,
759 llm_stats: None,
760 };
761 if let Some(parent) = find_node_mut(&mut self.root, &parent_id) {
762 if matches!(node.kind, WorkflowNodeKind::FanoutBranch { .. }) {
763 parent.parallelism = Parallelism::Parallel;
764 }
765 parent.children.push(node);
766 }
767 }
768 StreamFrame::FlowNodeEnd {
769 run_id,
770 node_id,
771 status,
772 output_preview,
773 ..
774 } => {
775 let scoped = scope_id(run_id, node_id);
776 if let Some(n) = find_node_mut(&mut self.root, &scoped) {
777 let new_status = match status {
778 FlowNodeStatus::Ok => NodeStatus::Ok,
779 FlowNodeStatus::Err => NodeStatus::Err,
780 FlowNodeStatus::Cancelled => NodeStatus::Cancelled,
781 };
782 n.status = new_status;
783 n.ended_at = Some(now);
784 if let Some(p) = output_preview {
785 n.output_preview = Some(p.clone());
786 }
787 for child in n.children.iter_mut() {
788 if matches!(child.status, NodeStatus::Running | NodeStatus::Pending) {
789 child.status = new_status;
790 child.ended_at = Some(now);
791 }
792 }
793 }
794 }
795 StreamFrame::LlmCallStats {
796 model,
797 provider,
798 context_call_purpose,
799 context_call_scope,
800 input_tokens,
801 output_tokens,
802 cache_read,
803 cache_write,
804 ttft_ms,
805 tokens_per_second,
806 wallclock_ms,
807 run_id,
808 node_id,
809 } => {
810 if let (Some(rid), Some(nid)) = (run_id.as_deref(), node_id.as_deref()) {
811 let scoped = scope_id(rid, nid);
812 if let Some(n) = find_node_mut(&mut self.root, &scoped) {
813 n.llm_stats = Some(LlmStats {
814 model: model.clone(),
815 provider: provider.clone(),
816 context_call_purpose: *context_call_purpose,
817 context_call_scope: *context_call_scope,
818 input_tokens: *input_tokens,
819 output_tokens: *output_tokens,
820 cache_read: *cache_read,
821 cache_write: *cache_write,
822 ttft_ms: *ttft_ms,
823 tokens_per_second: *tokens_per_second,
824 wallclock_ms: *wallclock_ms,
825 });
826 }
827 }
828 }
829 StreamFrame::ToolNode {
830 run_id,
831 parent_node_id,
832 tool_use_id,
833 tool,
834 args_preview,
835 call_intent,
836 ..
837 } => {
838 let scoped_parent = scope_id(run_id, parent_node_id);
839 let id = tool_node_id(run_id, tool_use_id);
840 if find_node(&self.root, &id).is_some() {
841 return;
842 }
843 let node = WorkflowNode {
844 id,
845 kind: WorkflowNodeKind::ToolCall {
846 tool_use_id: tool_use_id.clone(),
847 tool: tool.clone(),
848 args_preview: args_preview.clone(),
849 call_intent: call_intent.clone(),
850 result_preview: None,
851 },
852 label: tool.clone(),
853 status: NodeStatus::Running,
854 started_at: Some(now),
855 ended_at: None,
856 output_preview: None,
857 children: Vec::new(),
858 parallelism: Parallelism::Serial,
859 approval: None,
860 llm_stats: None,
861 };
862 if let Some(parent) = find_node_mut(&mut self.root, &scoped_parent) {
863 parent.children.push(node);
864 self.refresh_permission_tool_approvals();
865 }
866 }
867 StreamFrame::ToolUseDone {
868 id, ok, preview, ..
869 } => {
870 if let Some(n) = find_tool_node_by_tool_use_id(&mut self.root, id) {
871 n.status = if *ok { NodeStatus::Ok } else { NodeStatus::Err };
872 n.ended_at = Some(now);
873 n.output_preview = Some(preview.clone());
874 }
875 }
876 StreamFrame::FlowDone {
877 run_id,
878 ok,
879 cancelled,
880 ..
881 } => {
882 if let Some(n) = find_node_mut(&mut self.root, run_id) {
883 let status = if *cancelled {
884 NodeStatus::Cancelled
885 } else if *ok {
886 NodeStatus::Ok
887 } else {
888 NodeStatus::Err
889 };
890 cascade_terminate(n, status, now);
891 }
892 }
893 StreamFrame::AssistantMsg {
894 flow_run_id,
895 message,
896 } => {
897 let Some(rid_str) = flow_run_id else { return };
898 let Ok(uuid) = uuid::Uuid::parse_str(rid_str) else {
899 return;
900 };
901 self.apply_event(&Event::AssistantMsg {
902 turn_id: crate::event::TurnId::now(),
903 flow_run_id: Some(crate::event::FlowRunId(uuid)),
904 message: message.clone(),
905 });
906 }
907 StreamFrame::ToolResultMsg {
908 flow_run_id,
909 message,
910 } => {
911 let scoped_run_id = flow_run_id
912 .as_deref()
913 .and_then(|rid| uuid::Uuid::parse_str(rid).ok())
914 .map(crate::event::FlowRunId);
915 self.apply_event(&Event::ToolResultMsg {
916 turn_id: crate::event::TurnId::now(),
917 flow_run_id: scoped_run_id,
918 message: message.clone(),
919 });
920 }
921 StreamFrame::ToolPendingApproval {
922 run_id,
923 tool_use_id,
924 level,
925 preview,
926 ..
927 } => {
928 let id = tool_node_id(run_id, tool_use_id);
929 if let Some(n) = find_node_mut(&mut self.root, &id) {
930 n.approval = Some(ApprovalState::Pending {
931 level: level.clone(),
932 preview: preview.clone(),
933 });
934 }
935 }
936 StreamFrame::ToolApproved {
937 run_id,
938 tool_use_id,
939 ..
940 } => {
941 let id = tool_node_id(run_id, tool_use_id);
942 if let Some(n) = find_node_mut(&mut self.root, &id) {
943 n.approval = Some(ApprovalState::Approved);
944 }
945 }
946 StreamFrame::ToolDenied {
947 run_id,
948 tool_use_id,
949 reason,
950 } => {
951 let id = tool_node_id(run_id, tool_use_id);
952 if let Some(n) = find_node_mut(&mut self.root, &id) {
953 n.approval = Some(ApprovalState::Denied {
954 reason: reason.clone(),
955 });
956 }
957 }
958 StreamFrame::PermissionRequestCreated { payload, .. }
959 | StreamFrame::PermissionRequestTargeted { payload, .. }
960 | StreamFrame::PermissionRequestDeferred { payload, .. } => {
961 self.apply_permission_request(payload, WorkflowPermissionState::Pending);
962 }
963 StreamFrame::PermissionRequestApproved { payload, .. } => {
964 self.apply_permission_request(payload, WorkflowPermissionState::Approved);
965 }
966 StreamFrame::PermissionRequestDenied { payload, .. } => {
967 self.apply_permission_request(payload, WorkflowPermissionState::Denied);
968 }
969 StreamFrame::PermissionRequestCancelled { payload, .. } => {
970 self.apply_permission_request(payload, WorkflowPermissionState::Cancelled);
971 }
972 StreamFrame::UnrestrictedExecution { payload, .. } => {
973 self.apply_permission_request(payload, WorkflowPermissionState::Unrestricted);
974 }
975 StreamFrame::PermissionGroupCreated { payload, .. }
976 | StreamFrame::PermissionGroupUpdated { payload, .. } => {
977 self.apply_permission_group(payload, false);
978 }
979 StreamFrame::PermissionGroupResolved { payload, .. } => {
980 self.apply_permission_group(payload, true);
981 }
982 _ => {}
983 }
984 }
985}
986
987fn cascade_terminate(n: &mut WorkflowNode, status: NodeStatus, now: DateTime<Utc>) {
988 if matches!(n.status, NodeStatus::Running | NodeStatus::Pending) {
989 n.status = status;
990 n.ended_at = Some(now);
991 }
992 for child in n.children.iter_mut() {
993 cascade_terminate(child, status, now);
994 }
995}
996
997fn apply_permission_approvals(
998 nodes: &mut [WorkflowNode],
999 approvals: &BTreeMap<String, ApprovalState>,
1000) {
1001 for node in nodes {
1002 #[cfg(test)]
1003 update_perf_counters(|counters| {
1004 counters.workflow_node_visits = counters.workflow_node_visits.saturating_add(1);
1005 });
1006 if let Some(approval) = approvals.get(&node.id) {
1007 node.approval = Some(approval.clone());
1008 }
1009 apply_permission_approvals(&mut node.children, approvals);
1010 }
1011}
1012
1013fn find_node<'a>(nodes: &'a [WorkflowNode], id: &str) -> Option<&'a WorkflowNode> {
1014 for n in nodes {
1015 #[cfg(test)]
1016 update_perf_counters(|counters| {
1017 counters.workflow_node_visits = counters.workflow_node_visits.saturating_add(1);
1018 });
1019 if n.id == id {
1020 return Some(n);
1021 }
1022 if let Some(hit) = find_node(&n.children, id) {
1023 return Some(hit);
1024 }
1025 }
1026 None
1027}
1028
1029fn find_node_mut<'a>(nodes: &'a mut [WorkflowNode], id: &str) -> Option<&'a mut WorkflowNode> {
1030 for n in nodes.iter_mut() {
1031 #[cfg(test)]
1032 update_perf_counters(|counters| {
1033 counters.workflow_node_visits = counters.workflow_node_visits.saturating_add(1);
1034 });
1035 if n.id == id {
1036 return Some(n);
1037 }
1038 if let Some(hit) = find_node_mut(&mut n.children, id) {
1039 return Some(hit);
1040 }
1041 }
1042 None
1043}
1044
1045fn scope_id(run_id: &str, node_id: &str) -> String {
1046 format!("{run_id}::{node_id}")
1047}
1048
1049fn tool_node_id(run_id: &str, tool_use_id: &str) -> String {
1050 format!("tool:{run_id}:{tool_use_id}")
1051}
1052
1053fn collect_flow_run_ids(node: &WorkflowNode, out: &mut Vec<String>) {
1054 match &node.kind {
1055 WorkflowNodeKind::Flow { run_id, .. } | WorkflowNodeKind::Subflow { run_id, .. } => {
1056 out.push(run_id.clone());
1057 }
1058 _ => {}
1059 }
1060 for child in &node.children {
1061 collect_flow_run_ids(child, out);
1062 }
1063}
1064
1065pub fn permission_preview(payload: &PermissionRequestAudit) -> Option<String> {
1066 let mut lines = vec![format!("intent: {} ({:?})", payload.tool, payload.tier)];
1067 if let Some(call_intent) = &payload.call_intent {
1068 lines.push(format!("purpose: {}", call_intent.as_str()));
1069 }
1070 let provenance = &payload.provenance;
1071 if let Some(path) = &provenance.path {
1072 lines.push(format!("path: {path}"));
1073 }
1074 if let Some(cwd) = &provenance.cwd {
1075 lines.push(format!("cwd: {cwd}"));
1076 }
1077 if let Some(workspace) = &provenance.workspace_root {
1078 lines.push(format!("workspace: {workspace}"));
1079 }
1080 if provenance.network {
1081 lines.push("network: true".into());
1082 }
1083 if !provenance.risks.is_empty() {
1084 lines.push(format!(
1085 "risks: {}",
1086 provenance
1087 .risks
1088 .iter()
1089 .cloned()
1090 .collect::<Vec<_>>()
1091 .join(", ")
1092 ));
1093 }
1094 if !provenance.targets.is_empty() {
1095 lines.push(format!("targets: {}", provenance.targets.join(", ")));
1096 }
1097 lines.push(format!("policy: {}", payload.policy.snapshot_id));
1098 lines.push(format!("rule: {}", payload.policy.rule_id));
1099 Some(lines.join("\n"))
1100}
1101
1102fn find_tool_node_by_tool_use_id<'a>(
1103 nodes: &'a mut [WorkflowNode],
1104 tool_use_id: &str,
1105) -> Option<&'a mut WorkflowNode> {
1106 for n in nodes.iter_mut() {
1107 if let WorkflowNodeKind::ToolCall {
1108 tool_use_id: tid, ..
1109 } = &n.kind
1110 && tid == tool_use_id
1111 {
1112 return Some(n);
1113 }
1114 if let Some(hit) = find_tool_node_by_tool_use_id(&mut n.children, tool_use_id) {
1115 return Some(hit);
1116 }
1117 }
1118 None
1119}
1120
1121fn parse_branch_index(node_id: &str) -> Option<usize> {
1122 let start = node_id.rfind(".branch[")?;
1123 let rest = &node_id[start + ".branch[".len()..];
1124 let end = rest.find(']')?;
1125 rest[..end].parse().ok()
1126}
1127
1128pub fn rebuild_workflow_tree(events: &[crate::event::Event]) -> WorkflowGraph {
1133 let mut g = WorkflowGraph::new(crate::event::TurnId::now());
1134 for ev in events {
1135 g.apply_event(ev);
1136 }
1137 g
1138}
1139
1140pub fn rebuild_messages_for_run(
1144 events: &[crate::event::Event],
1145 run_id: &crate::event::FlowRunId,
1146) -> Vec<crate::message::Message> {
1147 events
1148 .iter()
1149 .filter_map(|ev| match ev {
1150 crate::event::Event::AssistantMsg {
1151 flow_run_id: Some(rid),
1152 message,
1153 ..
1154 } if rid == run_id => Some(message.clone()),
1155 crate::event::Event::ToolResultMsg {
1156 flow_run_id: Some(rid),
1157 message,
1158 ..
1159 } if rid == run_id => Some(message.clone()),
1160 _ => None,
1161 })
1162 .collect()
1163}
1164
1165#[cfg(test)]
1166mod tests {
1167 use super::*;
1168 use crate::event::{FlowRunId, FlowStatus};
1169 use crate::nodegraph::NodeKind;
1170 use std::collections::BTreeSet;
1171
1172 fn flow_start(run_id: FlowRunId, name: &str) -> Event {
1173 Event::FlowStart {
1174 run_id,
1175 flow_name: name.into(),
1176 parent_run_id: None,
1177 parent_node_id: None,
1178 spawned: false,
1179 }
1180 }
1181
1182 fn subflow_start(child: FlowRunId, parent: FlowRunId, parent_node: &str, name: &str) -> Event {
1183 Event::FlowStart {
1184 run_id: child,
1185 flow_name: name.into(),
1186 parent_run_id: Some(parent),
1187 parent_node_id: Some(parent_node.into()),
1188 spawned: false,
1189 }
1190 }
1191
1192 fn stmt_start(run_id: FlowRunId, node_id: &str, parent: Option<&str>) -> Event {
1193 Event::FlowNodeStart {
1194 run_id,
1195 node_id: node_id.into(),
1196 kind: NodeKind::UserConfirm,
1197 label: node_id.into(),
1198 parent_node_id: parent.map(String::from),
1199 }
1200 }
1201
1202 fn stmt_end(run_id: FlowRunId, node_id: &str, status: FlowNodeStatus) -> Event {
1203 Event::FlowNodeEnd {
1204 run_id,
1205 node_id: node_id.into(),
1206 status,
1207 output_preview: None,
1208 }
1209 }
1210
1211 fn request_id() -> PermissionRequestId {
1212 PermissionRequestId(uuid::Uuid::now_v7())
1213 }
1214
1215 fn permission_payload(
1216 request_id: PermissionRequestId,
1217 requesting_run_id: FlowRunId,
1218 root_run_id: FlowRunId,
1219 tool_use_id: &str,
1220 at: chrono::DateTime<chrono::Utc>,
1221 ) -> PermissionRequestAudit {
1222 PermissionRequestAudit {
1223 request_id: Some(request_id),
1224 revision: 1,
1225 session_id: "session".into(),
1226 requesting_run_id,
1227 parent_run_id: None,
1228 root_run_id,
1229 tool_use_id: tool_use_id.into(),
1230 tool: "fs.read".into(),
1231 call_intent: None,
1232 tier: crate::tool::Tier::Two,
1233 execution_boundary: Default::default(),
1234 provenance: crate::permission_audit::PermissionProvenanceSummary {
1235 cwd: None,
1236 path: None,
1237 path_origin: None,
1238 workspace_id: None,
1239 workspace_root: None,
1240 repository_root: None,
1241 network: false,
1242 risks: BTreeSet::new(),
1243 targets: Vec::new(),
1244 },
1245 target: crate::permission_audit::PermissionAuditTarget::User,
1246 group_ids: Vec::new(),
1247 policy: crate::permission_audit::PermissionPolicyReference {
1248 snapshot_id: "snapshot".into(),
1249 rule_id: "rule".into(),
1250 },
1251 escalation_path: Vec::new(),
1252 decision_id: None,
1253 actor: None,
1254 scope: None,
1255 reason: None,
1256 at,
1257 }
1258 }
1259
1260 fn add_tool(graph: &mut WorkflowGraph, run_id: &FlowRunId, tool_use_id: &str) {
1261 graph.apply_event(&flow_start(run_id.clone(), "agent_loop"));
1262 graph.apply_event(&stmt_start(run_id.clone(), "dispatch_all", None));
1263 graph.apply_event(&Event::ToolNode {
1264 run_id: run_id.clone(),
1265 parent_node_id: "dispatch_all".into(),
1266 tool_use_id: tool_use_id.into(),
1267 tool_name: "fs.read".into(),
1268 args_preview: String::new(),
1269 call_intent: None,
1270 });
1271 }
1272
1273 #[test]
1274 #[ignore = "captures the pre-index permission refresh traversal"]
1275 fn baseline_permission_transition_visits_the_full_table_and_tree() {
1276 const ENTRIES: usize = 10_000;
1277 let run_id = FlowRunId::now();
1278 let at = chrono::Utc::now();
1279 let mut graph = WorkflowGraph::new(TurnId::now());
1280 graph.root.push(WorkflowNode {
1281 id: run_id.to_string(),
1282 kind: WorkflowNodeKind::Flow {
1283 run_id: run_id.to_string(),
1284 flow_name: "baseline".into(),
1285 },
1286 label: "baseline".into(),
1287 status: NodeStatus::Running,
1288 started_at: Some(at),
1289 ended_at: None,
1290 output_preview: None,
1291 children: (0..ENTRIES)
1292 .map(|idx| WorkflowNode {
1293 id: tool_node_id(&run_id.to_string(), &format!("tool-{idx}")),
1294 kind: WorkflowNodeKind::ToolCall {
1295 tool_use_id: format!("tool-{idx}"),
1296 tool: "fs.read".into(),
1297 args_preview: String::new(),
1298 call_intent: None,
1299 result_preview: None,
1300 },
1301 label: "fs.read".into(),
1302 status: NodeStatus::Running,
1303 started_at: Some(at),
1304 ended_at: None,
1305 output_preview: None,
1306 children: Vec::new(),
1307 parallelism: Parallelism::Serial,
1308 approval: None,
1309 llm_stats: None,
1310 })
1311 .collect(),
1312 parallelism: Parallelism::Serial,
1313 approval: None,
1314 llm_stats: None,
1315 });
1316 for idx in 0..ENTRIES {
1317 let id = request_id();
1318 graph.permission_requests.insert(
1319 WorkflowPermissionIdentity::Canonical {
1320 request_id: id.clone(),
1321 },
1322 WorkflowPermissionRequest {
1323 payload: permission_payload(
1324 id,
1325 run_id.clone(),
1326 run_id.clone(),
1327 &format!("tool-{idx}"),
1328 at,
1329 ),
1330 state: WorkflowPermissionState::Pending,
1331 },
1332 );
1333 }
1334
1335 reset_perf_counters();
1336 let new_id = request_id();
1337 graph.apply_permission_request_with_identity(
1338 WorkflowPermissionIdentity::Canonical {
1339 request_id: new_id.clone(),
1340 },
1341 &permission_payload(
1342 new_id,
1343 run_id.clone(),
1344 run_id,
1345 "tool-0",
1346 at + chrono::Duration::seconds(1),
1347 ),
1348 WorkflowPermissionState::Approved,
1349 );
1350 let counters = perf_counters();
1351 assert_eq!(counters.permission_table_visits, (ENTRIES + 1) as u64);
1352 assert_eq!(counters.workflow_node_visits, (ENTRIES + 1) as u64);
1353 }
1354
1355 #[test]
1356 fn top_level_flow_becomes_root_child() {
1357 let mut g = WorkflowGraph::new(TurnId::now());
1358 let rid = FlowRunId::now();
1359 g.apply_event(&flow_start(rid.clone(), "main"));
1360 assert_eq!(g.root.len(), 1);
1361 let flow = &g.root[0];
1362 assert!(matches!(flow.kind, WorkflowNodeKind::Flow { .. }));
1363 assert_eq!(flow.status, NodeStatus::Running);
1364 assert_eq!(flow.id, rid.0.to_string());
1365 }
1366
1367 #[test]
1368 fn subflow_attaches_under_parent_node() {
1369 let mut g = WorkflowGraph::new(TurnId::now());
1370 let parent_flow = FlowRunId::now();
1371 let child_flow = FlowRunId::now();
1372 g.apply_event(&flow_start(parent_flow.clone(), "outer"));
1373 g.apply_event(&stmt_start(parent_flow.clone(), "stmt_0", None));
1374 g.apply_event(&subflow_start(
1375 child_flow.clone(),
1376 parent_flow.clone(),
1377 "stmt_0",
1378 "inner",
1379 ));
1380 let scoped = scope_id(&parent_flow.0.to_string(), "stmt_0");
1381 let stmt = g.find_node(&scoped).unwrap();
1382 assert_eq!(stmt.children.len(), 1);
1383 assert!(matches!(
1384 stmt.children[0].kind,
1385 WorkflowNodeKind::Subflow { .. }
1386 ));
1387 assert_eq!(stmt.children[0].id, child_flow.0.to_string());
1388 }
1389
1390 #[test]
1391 fn tool_node_attaches_and_flow_end_marks_status() {
1392 let mut g = WorkflowGraph::new(TurnId::now());
1393 let rid = FlowRunId::now();
1394 g.apply_event(&flow_start(rid.clone(), "main"));
1395 g.apply_event(&stmt_start(rid.clone(), "stmt_0", None));
1396 g.apply_event(&Event::ToolNode {
1397 run_id: rid.clone(),
1398 parent_node_id: "stmt_0".into(),
1399 tool_use_id: "tu_1".into(),
1400 tool_name: "fs.read".into(),
1401 args_preview: "{\"path\":\"a\"}".into(),
1402 call_intent: crate::message::ToolCallIntent::new("Inspect source"),
1403 });
1404 g.apply_event(&stmt_end(rid.clone(), "stmt_0", FlowNodeStatus::Ok));
1405 let tool = g
1406 .find_node(&tool_node_id(&rid.0.to_string(), "tu_1"))
1407 .expect("tool node");
1408 assert!(matches!(
1409 &tool.kind,
1410 WorkflowNodeKind::ToolCall { call_intent: Some(intent), .. }
1411 if intent.as_str() == "Inspect source"
1412 ));
1413 g.apply_event(&Event::FlowEnd {
1414 run_id: rid.clone(),
1415 flow_name: "main".into(),
1416 status: FlowStatus::Ok,
1417 });
1418 let scoped = scope_id(&rid.0.to_string(), "stmt_0");
1419 let stmt = g.find_node(&scoped).unwrap();
1420 assert_eq!(stmt.status, NodeStatus::Ok);
1421 assert_eq!(stmt.children.len(), 1);
1422 let tool = &stmt.children[0];
1423 assert_eq!(tool.id, tool_node_id(&rid.0.to_string(), "tu_1"));
1424 assert!(matches!(tool.kind, WorkflowNodeKind::ToolCall { .. }));
1425 assert_eq!(g.root[0].status, NodeStatus::Ok);
1426 }
1427
1428 #[test]
1429 fn assistant_tool_use_waits_for_scoped_tool_node_and_result() {
1430 use crate::message::{Message, MessageOrigin, MessagePart, MessageRole};
1431 use crate::stream::StreamFrame;
1432
1433 let mut graph = WorkflowGraph::new(TurnId::now());
1434 let run_id = FlowRunId::now();
1435 let run = run_id.0.to_string();
1436 graph.apply_event(&flow_start(run_id.clone(), "agent_loop"));
1437 graph.apply_event(&stmt_start(run_id.clone(), "llm", None));
1438 graph.apply_stream_frame(&StreamFrame::AssistantMsg {
1439 flow_run_id: Some(run.clone()),
1440 message: Message {
1441 role: MessageRole::Assistant,
1442 parts: vec![MessagePart::ToolUse {
1443 id: "tu_1".into(),
1444 name: "fs.read".into(),
1445 input: serde_json::json!({"path": "a.rs"}),
1446 intent: None,
1447 }],
1448 turn_id: TurnId::now(),
1449 origin: MessageOrigin::User,
1450 },
1451 });
1452 assert!(graph.find_node(&tool_node_id(&run, "tu_1")).is_none());
1453
1454 graph.apply_event(&stmt_start(run_id.clone(), "dispatch_all", None));
1455 graph.apply_stream_frame(&StreamFrame::ToolNode {
1456 run_id: run.clone(),
1457 parent_node_id: "dispatch_all".into(),
1458 tool_use_id: "tu_1".into(),
1459 tool: "fs.read".into(),
1460 args_preview: "{path: a.rs}".into(),
1461 call_intent: None,
1462 });
1463 graph.apply_stream_frame(&StreamFrame::ToolResultMsg {
1464 flow_run_id: Some(run.clone()),
1465 message: Message {
1466 role: MessageRole::Tool,
1467 parts: vec![MessagePart::ToolResult {
1468 tool_use_id: "tu_1".into(),
1469 content: "contents".into(),
1470 is_error: false,
1471 }],
1472 turn_id: TurnId::now(),
1473 origin: MessageOrigin::User,
1474 },
1475 });
1476
1477 let dispatch = graph
1478 .find_node(&scope_id(&run, "dispatch_all"))
1479 .expect("dispatch statement");
1480 assert_eq!(dispatch.children.len(), 1);
1481 assert_eq!(dispatch.children[0].status, NodeStatus::Ok);
1482 assert_eq!(
1483 dispatch.children[0].output_preview.as_deref(),
1484 Some("contents")
1485 );
1486 assert!(
1487 graph.root[0]
1488 .children
1489 .iter()
1490 .all(|node| !matches!(node.kind, WorkflowNodeKind::ToolCall { .. }))
1491 );
1492 }
1493
1494 #[test]
1495 fn llm_stream_stats_preserve_route_metadata() {
1496 use crate::stream::StreamFrame;
1497
1498 let mut graph = WorkflowGraph::new(TurnId::now());
1499 let run_id = FlowRunId::now();
1500 let run = run_id.0.to_string();
1501 graph.apply_event(&flow_start(run_id.clone(), "agent_loop"));
1502 graph.apply_event(&stmt_start(run_id, "llm", None));
1503 graph.apply_stream_frame(&StreamFrame::LlmCallStats {
1504 model: "helper-model".into(),
1505 provider: "helper-provider".into(),
1506 context_call_purpose: crate::context_plan::ContextCallPurpose::Extraction,
1507 context_call_scope: crate::context_plan::ContextCallScope::Child,
1508 input_tokens: 100,
1509 output_tokens: 10,
1510 cache_read: 20,
1511 cache_write: 30,
1512 ttft_ms: 40,
1513 tokens_per_second: 50.0,
1514 wallclock_ms: 60,
1515 run_id: Some(run.clone()),
1516 node_id: Some("llm".into()),
1517 });
1518
1519 let stats = graph
1520 .find_node(&scope_id(&run, "llm"))
1521 .and_then(|node| node.llm_stats.as_ref())
1522 .expect("llm stats");
1523 assert_eq!(stats.provider, "helper-provider");
1524 assert_eq!(
1525 stats.context_call_purpose,
1526 crate::context_plan::ContextCallPurpose::Extraction
1527 );
1528 assert_eq!(
1529 stats.context_call_scope,
1530 crate::context_plan::ContextCallScope::Child
1531 );
1532 }
1533
1534 #[test]
1535 fn scoped_tool_result_updates_only_matching_run() {
1536 use crate::message::{Message, MessageOrigin, MessagePart, MessageRole};
1537 use crate::stream::StreamFrame;
1538
1539 let mut graph = WorkflowGraph::new(TurnId::now());
1540 let run_a = FlowRunId::now();
1541 let run_b = FlowRunId::now();
1542 for run_id in [&run_a, &run_b] {
1543 graph.apply_event(&flow_start(run_id.clone(), "agent_loop"));
1544 graph.apply_event(&stmt_start(run_id.clone(), "dispatch_all", None));
1545 graph.apply_stream_frame(&StreamFrame::ToolNode {
1546 run_id: run_id.0.to_string(),
1547 parent_node_id: "dispatch_all".into(),
1548 tool_use_id: "same_id".into(),
1549 tool: "fs.read".into(),
1550 args_preview: String::new(),
1551 call_intent: None,
1552 });
1553 }
1554 graph.apply_stream_frame(&StreamFrame::ToolResultMsg {
1555 flow_run_id: Some(run_a.0.to_string()),
1556 message: Message {
1557 role: MessageRole::Tool,
1558 parts: vec![MessagePart::ToolResult {
1559 tool_use_id: "same_id".into(),
1560 content: "done".into(),
1561 is_error: false,
1562 }],
1563 turn_id: TurnId::now(),
1564 origin: MessageOrigin::User,
1565 },
1566 });
1567
1568 assert_eq!(
1569 graph
1570 .find_node(&tool_node_id(&run_a.0.to_string(), "same_id"))
1571 .unwrap()
1572 .status,
1573 NodeStatus::Ok
1574 );
1575 assert_eq!(
1576 graph
1577 .find_node(&tool_node_id(&run_b.0.to_string(), "same_id"))
1578 .unwrap()
1579 .status,
1580 NodeStatus::Running
1581 );
1582 }
1583
1584 #[test]
1585 fn permission_tool_approval_correlates_by_exact_run_and_tool() {
1586 let mut graph = WorkflowGraph::new(TurnId::now());
1587 let run_a = FlowRunId::now();
1588 let run_b = FlowRunId::now();
1589 add_tool(&mut graph, &run_a, "shared");
1590 add_tool(&mut graph, &run_b, "shared");
1591 let payload = permission_payload(
1592 request_id(),
1593 run_a.clone(),
1594 run_a.clone(),
1595 "shared",
1596 Utc::now(),
1597 );
1598 graph.apply_event(&Event::PermissionRequestCreated { payload });
1599
1600 assert!(matches!(
1601 graph
1602 .find_node(&tool_node_id(&run_a.0.to_string(), "shared"))
1603 .unwrap()
1604 .approval,
1605 Some(ApprovalState::Pending { .. })
1606 ));
1607 assert_eq!(
1608 graph
1609 .find_node(&tool_node_id(&run_b.0.to_string(), "shared"))
1610 .unwrap()
1611 .approval,
1612 None
1613 );
1614 }
1615
1616 #[test]
1617 fn permission_arriving_before_tool_refreshes_when_tool_is_added() {
1618 let mut graph = WorkflowGraph::new(TurnId::now());
1619 let run_id = FlowRunId::now();
1620 let payload = permission_payload(
1621 request_id(),
1622 run_id.clone(),
1623 run_id.clone(),
1624 "late-tool",
1625 Utc::now(),
1626 );
1627 graph.apply_event(&Event::PermissionRequestCreated { payload });
1628 add_tool(&mut graph, &run_id, "late-tool");
1629 assert!(matches!(
1630 graph
1631 .find_node(&tool_node_id(&run_id.0.to_string(), "late-tool"))
1632 .unwrap()
1633 .approval,
1634 Some(ApprovalState::Pending { .. })
1635 ));
1636 }
1637
1638 #[test]
1639 fn batched_permission_projection_matches_incremental_updates() {
1640 let run_id = FlowRunId::now();
1641 let request_id = request_id();
1642 let pending = permission_payload(
1643 request_id.clone(),
1644 run_id.clone(),
1645 run_id.clone(),
1646 "batched-tool",
1647 Utc::now(),
1648 );
1649 let mut approved = pending.clone();
1650 approved.at += chrono::Duration::seconds(1);
1651 approved.reason = Some("approved".into());
1652
1653 let mut incremental = WorkflowGraph::new(TurnId::now());
1654 add_tool(&mut incremental, &run_id, "batched-tool");
1655 incremental.apply_permission_request_with_identity(
1656 WorkflowPermissionIdentity::Canonical {
1657 request_id: request_id.clone(),
1658 },
1659 &pending,
1660 WorkflowPermissionState::Pending,
1661 );
1662 incremental.apply_permission_request_with_identity(
1663 WorkflowPermissionIdentity::Canonical {
1664 request_id: request_id.clone(),
1665 },
1666 &approved,
1667 WorkflowPermissionState::Approved,
1668 );
1669
1670 let mut batched = WorkflowGraph::new(TurnId::now());
1671 add_tool(&mut batched, &run_id, "batched-tool");
1672 batched.apply_permission_requests([
1673 (
1674 WorkflowPermissionIdentity::Canonical {
1675 request_id: request_id.clone(),
1676 },
1677 pending,
1678 WorkflowPermissionState::Pending,
1679 ),
1680 (
1681 WorkflowPermissionIdentity::Canonical { request_id },
1682 approved,
1683 WorkflowPermissionState::Approved,
1684 ),
1685 ]);
1686
1687 assert_eq!(batched.permission_requests, incremental.permission_requests);
1688 assert_eq!(
1689 batched
1690 .find_node(&tool_node_id(&run_id.0.to_string(), "batched-tool"))
1691 .unwrap()
1692 .approval,
1693 incremental
1694 .find_node(&tool_node_id(&run_id.0.to_string(), "batched-tool"))
1695 .unwrap()
1696 .approval
1697 );
1698 }
1699
1700 #[test]
1701 fn canonical_and_legacy_permission_lifecycles_count_and_interrupt_independently() {
1702 let mut graph = WorkflowGraph::new(TurnId::now());
1703 let root = FlowRunId::now();
1704 let child = FlowRunId::now();
1705 graph.apply_event(&flow_start(root.clone(), "root"));
1706 graph.apply_event(&stmt_start(root.clone(), "spawn", None));
1707 graph.apply_event(&subflow_start(
1708 child.clone(),
1709 root.clone(),
1710 "spawn",
1711 "child",
1712 ));
1713 let now = Utc::now();
1714 let canonical_id = request_id();
1715 let canonical = permission_payload(
1716 canonical_id.clone(),
1717 child.clone(),
1718 root.clone(),
1719 "canonical",
1720 now,
1721 );
1722 let legacy = permission_payload(
1723 request_id(),
1724 child.clone(),
1725 root.clone(),
1726 "legacy",
1727 now + chrono::Duration::seconds(1),
1728 );
1729 graph.apply_permission_request_with_identity(
1730 WorkflowPermissionIdentity::Canonical {
1731 request_id: canonical_id.clone(),
1732 },
1733 &canonical,
1734 WorkflowPermissionState::Pending,
1735 );
1736 graph.apply_permission_request_with_identity(
1737 WorkflowPermissionIdentity::Legacy {
1738 seq: 7,
1739 run_id: child.0.to_string(),
1740 tool_use_id: "legacy".into(),
1741 },
1742 &legacy,
1743 WorkflowPermissionState::Pending,
1744 );
1745 assert_eq!(graph.descendant_pending_permissions(&root.0.to_string()), 2);
1746
1747 let mut approved = canonical.clone();
1748 approved.actor = Some(crate::permission_audit::PermissionProjectionActor::Flow {
1749 session_id: "session".into(),
1750 run_id: root.clone(),
1751 });
1752 approved.reason = Some("approved by parent".into());
1753 graph.apply_event(&Event::PermissionRequestApproved { payload: approved });
1754 assert_eq!(graph.descendant_pending_permissions(&root.0.to_string()), 1);
1755 assert_eq!(
1756 graph
1757 .permission_requests
1758 .get(&WorkflowPermissionIdentity::Canonical {
1759 request_id: canonical_id,
1760 })
1761 .unwrap()
1762 .state,
1763 WorkflowPermissionState::Approved
1764 );
1765
1766 graph.interrupt_pending_permissions();
1767 assert_eq!(graph.descendant_pending_permissions(&root.0.to_string()), 0);
1768 assert!(graph.permission_requests.values().any(|request| {
1769 request.state == WorkflowPermissionState::Interrupted
1770 && request.payload.tool_use_id == "legacy"
1771 }));
1772 }
1773
1774 #[test]
1775 fn permission_event_and_stream_reducers_converge_on_canonical_state() {
1776 let run_id = FlowRunId::now();
1777 let request_id = request_id();
1778 let created = permission_payload(
1779 request_id,
1780 run_id.clone(),
1781 run_id.clone(),
1782 "tool",
1783 Utc::now(),
1784 );
1785 let mut approved = created.clone();
1786 approved.actor = Some(crate::permission_audit::PermissionProjectionActor::User {
1787 session_id: "session".into(),
1788 principal_id: Some("operator".into()),
1789 });
1790 approved.reason = Some("accepted".into());
1791 approved.at += chrono::Duration::seconds(1);
1792
1793 let mut from_events = WorkflowGraph::new(TurnId::now());
1794 from_events.apply_event(&Event::PermissionRequestCreated {
1795 payload: created.clone(),
1796 });
1797 from_events.apply_event(&Event::PermissionRequestApproved {
1798 payload: approved.clone(),
1799 });
1800 let mut from_stream = WorkflowGraph::new(TurnId::now());
1801 from_stream.apply_stream_frame(&crate::stream::StreamFrame::PermissionRequestCreated {
1802 run_id: run_id.0.to_string(),
1803 payload: created,
1804 });
1805 from_stream.apply_stream_frame(&crate::stream::StreamFrame::PermissionRequestApproved {
1806 run_id: run_id.0.to_string(),
1807 payload: approved,
1808 });
1809
1810 assert_eq!(
1811 from_events.permission_requests,
1812 from_stream.permission_requests
1813 );
1814 let request = from_events.permission_requests.values().next().unwrap();
1815 assert_eq!(request.state, WorkflowPermissionState::Approved);
1816 assert_eq!(request.payload.reason.as_deref(), Some("accepted"));
1817 assert!(matches!(
1818 request.payload.actor,
1819 Some(crate::permission_audit::PermissionProjectionActor::User { .. })
1820 ));
1821 }
1822
1823 #[test]
1824 fn s8_canonical_jsonl_replay_converges_with_event_and_stream_reducers() {
1825 use crate::permission::PermissionGroupId;
1826 use crate::permission_audit::{PermissionGroupAudit, PermissionProjectionActor};
1827 use crate::projection::message_window::{TranscriptEntry, replay_transcript_from};
1828 use crate::stream::StreamFrame;
1829
1830 let root = FlowRunId::now();
1831 let child = FlowRunId::now();
1832 let now = Utc::now();
1833 let approved_id = request_id();
1834 let unrestricted_id = request_id();
1835 let interrupted_id = request_id();
1836 let mut approved = permission_payload(
1837 approved_id.clone(),
1838 child.clone(),
1839 root.clone(),
1840 "approved",
1841 now,
1842 );
1843 approved.parent_run_id = Some(root.clone());
1844 let mut approved_final = approved.clone();
1845 approved_final.actor = Some(PermissionProjectionActor::User {
1846 session_id: "session".into(),
1847 principal_id: Some("operator".into()),
1848 });
1849 approved_final.reason = Some("accepted".into());
1850 approved_final.at += chrono::Duration::seconds(2);
1851 let mut unrestricted = permission_payload(
1852 unrestricted_id,
1853 child.clone(),
1854 root.clone(),
1855 "unrestricted",
1856 now,
1857 );
1858 unrestricted.parent_run_id = Some(root.clone());
1859 unrestricted.actor = Some(PermissionProjectionActor::Policy {
1860 policy_version: "snapshot".into(),
1861 rule_id: "unrestricted".into(),
1862 });
1863 let mut interrupted = permission_payload(
1864 interrupted_id,
1865 child.clone(),
1866 root.clone(),
1867 "interrupted",
1868 now,
1869 );
1870 interrupted.parent_run_id = Some(root.clone());
1871 let group = PermissionGroupAudit {
1872 group_id: PermissionGroupId(uuid::Uuid::now_v7()),
1873 owner: crate::permission_audit::PermissionGroupAuditOwner::Flow {
1874 run_id: root.clone(),
1875 },
1876 label: "acceptance".into(),
1877 request_ids: vec![approved_id.clone()],
1878 revision: 1,
1879 at: now,
1880 };
1881
1882 let mut created_json = serde_json::to_value(&approved).unwrap();
1883 created_json.as_object_mut().unwrap().remove("provenance");
1884 let lines = [
1885 serde_json::json!({"type":"permission_request_created","seq":1,"payload":created_json}).to_string(),
1886 "{\"type\":\"permission_request_targeted\",\"payload\":".into(),
1887 serde_json::json!({"type":"permission_request_targeted","seq":2,"payload":approved}).to_string(),
1888 serde_json::json!({"type":"tool_pending_approval","seq":3,"run_id":child.0.to_string(),"tool_use_id":"approved","tool_name":"fs.read","args_preview":"{}","level":"approve"}).to_string(),
1889 serde_json::json!({"type":"permission_group_created","seq":4,"payload":group}).to_string(),
1890 serde_json::json!({"type":"permission_request_approved","seq":5,"payload":approved_final}).to_string(),
1891 serde_json::json!({"type":"tool_approved","seq":6,"run_id":child.0.to_string(),"tool_use_id":"approved","decided_by":"broker"}).to_string(),
1892 serde_json::json!({"type":"permission_group_resolved","seq":7,"payload":group}).to_string(),
1893 serde_json::json!({"type":"unrestricted_execution","seq":8,"payload":unrestricted}).to_string(),
1894 serde_json::json!({"type":"tool_approved","seq":9,"run_id":child.0.to_string(),"tool_use_id":"unrestricted","decided_by":"unrestricted"}).to_string(),
1895 serde_json::json!({"type":"permission_request_created","seq":10,"payload":interrupted}).to_string(),
1896 serde_json::json!({"type":"tool_pending_approval","seq":11,"run_id":child.0.to_string(),"tool_use_id":"interrupted","tool_name":"fs.read","args_preview":"{}","level":"approve"}).to_string(),
1897 ];
1898 let dir = tempfile::tempdir().unwrap();
1899 let path = dir.path().join("events.jsonl");
1900 std::fs::write(&path, lines.join("\n")).unwrap();
1901 let transcript = replay_transcript_from(&path).unwrap();
1902
1903 let mut from_history = WorkflowGraph::new(TurnId::now());
1904 let mut from_events = WorkflowGraph::new(TurnId::now());
1905 let mut from_stream = WorkflowGraph::new(TurnId::now());
1906 for graph in [&mut from_history, &mut from_events, &mut from_stream] {
1907 graph.apply_event(&flow_start(root.clone(), "root"));
1908 graph.apply_event(&stmt_start(root.clone(), "spawn", None));
1909 graph.apply_event(&subflow_start(
1910 child.clone(),
1911 root.clone(),
1912 "spawn",
1913 "child",
1914 ));
1915 graph.apply_event(&stmt_start(child.clone(), "dispatch_all", None));
1916 for tool in ["approved", "unrestricted", "interrupted"] {
1917 graph.apply_event(&Event::ToolNode {
1918 run_id: child.clone(),
1919 parent_node_id: "dispatch_all".into(),
1920 tool_use_id: tool.into(),
1921 tool_name: "fs.read".into(),
1922 args_preview: String::new(),
1923 call_intent: None,
1924 });
1925 }
1926 }
1927 for entry in &transcript {
1928 match entry {
1929 TranscriptEntry::PermissionRequest {
1930 identity,
1931 payload,
1932 state,
1933 } => {
1934 from_history.apply_permission_request_with_identity(
1935 identity.clone(),
1936 payload,
1937 *state,
1938 );
1939 }
1940 TranscriptEntry::PermissionGroup { payload, resolved } => {
1941 from_history.apply_permission_group(payload, *resolved);
1942 }
1943 _ => {}
1944 }
1945 }
1946 let events = [
1947 Event::PermissionRequestCreated {
1948 payload: approved.clone(),
1949 },
1950 Event::PermissionRequestTargeted {
1951 payload: approved.clone(),
1952 },
1953 Event::PermissionGroupCreated {
1954 payload: group.clone(),
1955 },
1956 Event::PermissionRequestApproved {
1957 payload: approved_final.clone(),
1958 },
1959 Event::PermissionGroupResolved {
1960 payload: group.clone(),
1961 },
1962 Event::UnrestrictedExecution {
1963 payload: unrestricted.clone(),
1964 },
1965 Event::PermissionRequestCreated {
1966 payload: interrupted.clone(),
1967 },
1968 ];
1969 for event in &events {
1970 from_events.apply_event(event);
1971 }
1972 from_events.interrupt_pending_permissions();
1973 let run_id = child.0.to_string();
1974 let frames = [
1975 StreamFrame::PermissionRequestCreated {
1976 run_id: run_id.clone(),
1977 payload: approved.clone(),
1978 },
1979 StreamFrame::PermissionRequestTargeted {
1980 run_id: run_id.clone(),
1981 payload: approved,
1982 },
1983 StreamFrame::PermissionGroupCreated {
1984 run_id: root.0.to_string(),
1985 payload: group.clone(),
1986 },
1987 StreamFrame::PermissionRequestApproved {
1988 run_id: run_id.clone(),
1989 payload: approved_final,
1990 },
1991 StreamFrame::PermissionGroupResolved {
1992 run_id: root.0.to_string(),
1993 payload: group,
1994 },
1995 StreamFrame::UnrestrictedExecution {
1996 run_id: run_id.clone(),
1997 payload: unrestricted,
1998 },
1999 StreamFrame::PermissionRequestCreated {
2000 run_id,
2001 payload: interrupted,
2002 },
2003 ];
2004 for frame in &frames {
2005 from_stream.apply_stream_frame(frame);
2006 }
2007 from_stream.interrupt_pending_permissions();
2008
2009 assert_eq!(
2010 from_history.permission_requests,
2011 from_events.permission_requests
2012 );
2013 assert_eq!(
2014 from_history.permission_requests,
2015 from_stream.permission_requests
2016 );
2017 assert_eq!(
2018 from_history.permission_groups,
2019 from_events.permission_groups
2020 );
2021 assert_eq!(
2022 from_history.permission_groups,
2023 from_stream.permission_groups
2024 );
2025 fn without_timestamps(mut nodes: Vec<WorkflowNode>) -> Vec<WorkflowNode> {
2026 fn clear(nodes: &mut [WorkflowNode]) {
2027 for node in nodes {
2028 node.started_at = None;
2029 node.ended_at = None;
2030 clear(&mut node.children);
2031 }
2032 }
2033 clear(&mut nodes);
2034 nodes
2035 }
2036 assert_eq!(
2037 without_timestamps(from_history.root.clone()),
2038 without_timestamps(from_events.root.clone())
2039 );
2040 assert_eq!(
2041 without_timestamps(from_history.root.clone()),
2042 without_timestamps(from_stream.root.clone())
2043 );
2044 assert_eq!(
2045 from_history.descendant_pending_permissions(&root.0.to_string()),
2046 0
2047 );
2048 assert_eq!(from_history.permission_requests.len(), 3);
2049 assert!(
2050 from_history
2051 .permission_requests
2052 .keys()
2053 .all(|identity| matches!(identity, WorkflowPermissionIdentity::Canonical { .. }))
2054 );
2055 assert!(from_history.permission_requests.contains_key(
2056 &WorkflowPermissionIdentity::Canonical {
2057 request_id: approved_id
2058 }
2059 ));
2060 assert!(from_history.permission_requests.values().any(|request| {
2061 request.state == WorkflowPermissionState::Interrupted
2062 && request.payload.reason.as_deref()
2063 == Some("interrupted at end of persisted history")
2064 }));
2065 }
2066
2067 #[test]
2068 fn s8_resolved_group_tombstones_prevent_live_stream_and_replay_resurrection() {
2069 use crate::permission::PermissionGroupId;
2070 use crate::permission_audit::PermissionGroupAudit;
2071 use crate::projection::message_window::{TranscriptEntry, replay_transcript_from};
2072 use crate::stream::StreamFrame;
2073
2074 let owner = FlowRunId::now();
2075 let groups = [
2076 PermissionGroupAudit {
2077 group_id: PermissionGroupId(uuid::Uuid::now_v7()),
2078 owner: crate::permission_audit::PermissionGroupAuditOwner::Flow {
2079 run_id: owner.clone(),
2080 },
2081 label: "resolved before ungroup".into(),
2082 request_ids: vec![request_id()],
2083 revision: 1,
2084 at: Utc::now(),
2085 },
2086 PermissionGroupAudit {
2087 group_id: PermissionGroupId(uuid::Uuid::now_v7()),
2088 owner: crate::permission_audit::PermissionGroupAuditOwner::Flow {
2089 run_id: owner.clone(),
2090 },
2091 label: "terminal-only".into(),
2092 request_ids: Vec::new(),
2093 revision: 1,
2094 at: Utc::now(),
2095 },
2096 ];
2097 let mut from_events = WorkflowGraph::new(TurnId::now());
2098 let mut from_stream = WorkflowGraph::new(TurnId::now());
2099 let mut lines = Vec::new();
2100 let mut seq = 1_u64;
2101 for group in &groups {
2102 let mut updated = group.clone();
2103 updated.request_ids.clear();
2104 updated.revision += 1;
2105
2106 for event in [
2107 Event::PermissionGroupCreated {
2108 payload: group.clone(),
2109 },
2110 Event::PermissionGroupResolved {
2111 payload: group.clone(),
2112 },
2113 Event::PermissionGroupUpdated {
2114 payload: updated.clone(),
2115 },
2116 ] {
2117 from_events.apply_event(&event);
2118 }
2119 for frame in [
2120 StreamFrame::PermissionGroupCreated {
2121 run_id: owner.0.to_string(),
2122 payload: group.clone(),
2123 },
2124 StreamFrame::PermissionGroupResolved {
2125 run_id: owner.0.to_string(),
2126 payload: group.clone(),
2127 },
2128 StreamFrame::PermissionGroupUpdated {
2129 run_id: owner.0.to_string(),
2130 payload: updated.clone(),
2131 },
2132 ] {
2133 from_stream.apply_stream_frame(&frame);
2134 }
2135 for (kind, payload) in [
2136 ("permission_group_created", group),
2137 ("permission_group_resolved", group),
2138 ("permission_group_updated", &updated),
2139 ] {
2140 lines
2141 .push(serde_json::json!({"type":kind,"seq":seq,"payload":payload}).to_string());
2142 seq += 1;
2143 }
2144 }
2145
2146 let dir = tempfile::tempdir().unwrap();
2147 let path = dir.path().join("events.jsonl");
2148 std::fs::write(&path, lines.join("\n")).unwrap();
2149 let mut from_history = WorkflowGraph::new(TurnId::now());
2150 for entry in replay_transcript_from(&path).unwrap() {
2151 if let TranscriptEntry::PermissionGroup { payload, resolved } = entry {
2152 from_history.apply_permission_group(&payload, resolved);
2153 }
2154 }
2155
2156 for graph in [&from_events, &from_stream, &from_history] {
2157 assert!(graph.permission_groups.is_empty());
2158 assert_eq!(graph.resolved_permission_groups.len(), groups.len());
2159 assert!(
2160 groups
2161 .iter()
2162 .all(|group| graph.resolved_permission_groups.contains(&group.group_id))
2163 );
2164 }
2165 let restored: WorkflowGraph =
2166 serde_json::from_value(serde_json::to_value(&from_events).unwrap()).unwrap();
2167 assert_eq!(
2168 restored.resolved_permission_groups,
2169 from_events.resolved_permission_groups
2170 );
2171 }
2172
2173 #[test]
2174 fn fanout_branch_marks_parent_parallel() {
2175 let mut g = WorkflowGraph::new(TurnId::now());
2176 let rid = FlowRunId::now();
2177 g.apply_event(&flow_start(rid.clone(), "main"));
2178 g.apply_event(&stmt_start(rid.clone(), "stmt_1", None));
2179 g.apply_event(&stmt_start(rid.clone(), "stmt_1.branch[0]", Some("stmt_1")));
2180 g.apply_event(&stmt_start(rid.clone(), "stmt_1.branch[1]", Some("stmt_1")));
2181 let scoped = scope_id(&rid.0.to_string(), "stmt_1");
2182 let parent = g.find_node(&scoped).unwrap();
2183 assert_eq!(parent.parallelism, Parallelism::Parallel);
2184 assert_eq!(parent.children.len(), 2);
2185 assert!(matches!(
2186 parent.children[0].kind,
2187 WorkflowNodeKind::FanoutBranch { branch_index: 0 }
2188 ));
2189 assert!(matches!(
2190 parent.children[1].kind,
2191 WorkflowNodeKind::FanoutBranch { branch_index: 1 }
2192 ));
2193 }
2194
2195 #[test]
2196 fn out_of_order_events_silently_dropped() {
2197 let mut g = WorkflowGraph::new(TurnId::now());
2198 g.apply_event(&stmt_start(FlowRunId::now(), "stmt_0", Some("missing")));
2199 g.apply_event(&Event::ToolNode {
2200 run_id: FlowRunId::now(),
2201 parent_node_id: "missing".into(),
2202 tool_use_id: "tu".into(),
2203 tool_name: "t".into(),
2204 args_preview: "{}".into(),
2205 call_intent: None,
2206 });
2207 assert!(g.root.is_empty());
2208 }
2209}
2210
2211#[test]
2212fn rebuild_workflow_tree_restores_root_and_subflow() {
2213 use crate::event::{Event, FlowRunId};
2214 let root = FlowRunId::now();
2215 let child = FlowRunId::now();
2216 let events = vec![
2217 Event::FlowStart {
2218 run_id: root.clone(),
2219 flow_name: "agent".into(),
2220 parent_run_id: None,
2221 parent_node_id: None,
2222 spawned: false,
2223 },
2224 Event::FlowNodeStart {
2225 run_id: root.clone(),
2226 node_id: "stmt_0".into(),
2227 kind: crate::nodegraph::NodeKind::Llm { model: None },
2228 label: "llm".into(),
2229 parent_node_id: None,
2230 },
2231 Event::FlowStart {
2232 run_id: child.clone(),
2233 flow_name: "subagent".into(),
2234 parent_run_id: Some(root.clone()),
2235 parent_node_id: Some("stmt_0".into()),
2236 spawned: false,
2237 },
2238 Event::FlowEnd {
2239 run_id: child.clone(),
2240 flow_name: "subagent".into(),
2241 status: crate::event::FlowStatus::Ok,
2242 },
2243 Event::FlowEnd {
2244 run_id: root.clone(),
2245 flow_name: "agent".into(),
2246 status: crate::event::FlowStatus::Ok,
2247 },
2248 ];
2249 let g = rebuild_workflow_tree(&events);
2250 assert_eq!(g.root.len(), 1, "one root flow");
2251 assert_eq!(g.root[0].id, root.0.to_string());
2252 assert!(!g.root[0].children.is_empty(), "subflow nested under root");
2254}
2255
2256#[test]
2257fn rebuild_messages_for_run_filters_by_run_id() {
2258 use crate::event::{Event, FlowRunId, TurnId};
2259 use crate::message::{Message, MessageOrigin, MessagePart, MessageRole};
2260 let root = FlowRunId::now();
2261 let child = FlowRunId::now();
2262 let tid = TurnId::now();
2263 let mk = |role, text, rid: Option<FlowRunId>| Event::AssistantMsg {
2264 turn_id: tid.clone(),
2265 flow_run_id: rid,
2266 message: Message {
2267 role,
2268 parts: vec![MessagePart::Text { text }],
2269 turn_id: tid.clone(),
2270 origin: MessageOrigin::User,
2271 },
2272 };
2273 let events = vec![
2274 mk(MessageRole::Assistant, "root reply".into(), None),
2275 mk(
2276 MessageRole::Assistant,
2277 "child reply".into(),
2278 Some(child.clone()),
2279 ),
2280 mk(MessageRole::Assistant, "root again".into(), None),
2281 ];
2282 let root_msgs = rebuild_messages_for_run(&events, &root);
2285 assert_eq!(root_msgs.len(), 0);
2286 let child_msgs = rebuild_messages_for_run(&events, &child);
2287 assert_eq!(child_msgs.len(), 1);
2288 assert_eq!(child_msgs[0].text_concat(), "child reply");
2289}