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