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 self.apply_permission_inner(payload, WorkflowPermissionState::Pending)
504 }
505 Event::PermissionRequestApproved { payload } => {
506 self.apply_permission_inner(payload, WorkflowPermissionState::Approved)
507 }
508 Event::PermissionRequestDenied { payload } => {
509 self.apply_permission_inner(payload, WorkflowPermissionState::Denied)
510 }
511 Event::PermissionRequestCancelled { payload } => {
512 self.apply_permission_inner(payload, WorkflowPermissionState::Cancelled)
513 }
514 Event::UnrestrictedExecution { payload } => {
515 self.apply_permission_inner(payload, WorkflowPermissionState::Unrestricted)
516 }
517 Event::PermissionGroupCreated { payload }
518 | Event::PermissionGroupUpdated { payload } => {
519 self.apply_permission_group_inner(payload, false)
520 }
521 Event::PermissionGroupResolved { payload } => {
522 self.apply_permission_group_inner(payload, true)
523 }
524 _ => PendingDelta::default(),
525 }
526 }
527
528 fn apply_stream_frame_inner(
529 &mut self,
530 frame: &StreamFrame,
531 override_ts: Option<DateTime<Utc>>,
532 ) -> PendingDelta {
533 let now = override_ts.unwrap_or_else(Utc::now);
534 match frame {
535 StreamFrame::FlowGraph { run_id, graph } => self.insert_flow(
536 run_id.clone(),
537 graph.flow_name.clone(),
538 None,
539 None,
540 now,
541 true,
542 ),
543 StreamFrame::FlowStart {
544 run_id,
545 flow_name,
546 parent_run_id,
547 parent_node_id,
548 } => self.insert_flow(
549 run_id.clone(),
550 flow_name.clone(),
551 parent_run_id.clone(),
552 parent_node_id.clone(),
553 now,
554 true,
555 ),
556 StreamFrame::FlowNodeStart {
557 run_id,
558 node_id,
559 kind,
560 label,
561 parent_node_id,
562 } => {
563 self.insert_flow_node(run_id, node_id, kind, label, parent_node_id.as_deref(), now)
564 }
565 StreamFrame::FlowNodeEnd {
566 run_id,
567 node_id,
568 status,
569 output_preview,
570 ..
571 } => {
572 let status = match status {
573 FlowNodeStatus::Ok => NodeStatus::Ok,
574 FlowNodeStatus::Err => NodeStatus::Err,
575 FlowNodeStatus::Cancelled => NodeStatus::Cancelled,
576 };
577 self.finish_node(
578 &scope_id(run_id, node_id),
579 status,
580 output_preview.as_deref(),
581 now,
582 false,
583 )
584 }
585 StreamFrame::LlmCallStats {
586 model,
587 provider,
588 context_call_purpose,
589 context_call_scope,
590 input_tokens,
591 output_tokens,
592 cache_read,
593 cache_write,
594 ttft_ms,
595 tokens_per_second,
596 wallclock_ms,
597 run_id,
598 node_id,
599 } => {
600 let Some((run_id, node_id)) = run_id.as_deref().zip(node_id.as_deref()) else {
601 return PendingDelta::default();
602 };
603 self.set_llm_stats(
604 &scope_id(run_id, node_id),
605 LlmStats {
606 model: model.clone(),
607 provider: provider.clone(),
608 context_call_purpose: *context_call_purpose,
609 context_call_scope: *context_call_scope,
610 input_tokens: *input_tokens,
611 output_tokens: *output_tokens,
612 cache_read: *cache_read,
613 cache_write: *cache_write,
614 ttft_ms: *ttft_ms,
615 tokens_per_second: *tokens_per_second,
616 wallclock_ms: *wallclock_ms,
617 },
618 )
619 }
620 StreamFrame::ToolNode {
621 run_id,
622 parent_node_id,
623 tool_use_id,
624 tool,
625 args_preview,
626 call_intent,
627 ..
628 } => self.insert_tool(
629 run_id,
630 parent_node_id,
631 tool_use_id,
632 tool,
633 args_preview,
634 call_intent.clone(),
635 now,
636 ),
637 StreamFrame::ToolUseDone {
638 id, ok, preview, ..
639 } => self.finish_tool(None, id, *ok, preview, now, false),
640 StreamFrame::FlowDone {
641 run_id,
642 ok,
643 cancelled,
644 ..
645 } => {
646 let status = if *cancelled {
647 NodeStatus::Cancelled
648 } else if *ok {
649 NodeStatus::Ok
650 } else {
651 NodeStatus::Err
652 };
653 self.finish_node(run_id, status, None, now, true)
654 }
655 StreamFrame::ToolResultMsg {
656 flow_run_id,
657 message,
658 } => self.apply_tool_results(flow_run_id.as_deref(), message, now),
659 StreamFrame::ToolPendingApproval {
660 run_id,
661 tool_use_id,
662 level,
663 preview,
664 ..
665 } => self.set_tool_approval(
666 run_id,
667 tool_use_id,
668 ApprovalState::Pending {
669 level: level.clone(),
670 preview: preview.clone(),
671 },
672 ),
673 StreamFrame::ToolApproved {
674 run_id,
675 tool_use_id,
676 ..
677 } => self.set_tool_approval(run_id, tool_use_id, ApprovalState::Approved),
678 StreamFrame::ToolDenied {
679 run_id,
680 tool_use_id,
681 reason,
682 } => self.set_tool_approval(
683 run_id,
684 tool_use_id,
685 ApprovalState::Denied {
686 reason: reason.clone(),
687 },
688 ),
689 StreamFrame::PermissionRequestCreated { payload, .. }
690 | StreamFrame::PermissionRequestTargeted { payload, .. }
691 | StreamFrame::PermissionRequestDeferred { payload, .. } => {
692 self.apply_permission_inner(payload, WorkflowPermissionState::Pending)
693 }
694 StreamFrame::PermissionRequestApproved { payload, .. } => {
695 self.apply_permission_inner(payload, WorkflowPermissionState::Approved)
696 }
697 StreamFrame::PermissionRequestDenied { payload, .. } => {
698 self.apply_permission_inner(payload, WorkflowPermissionState::Denied)
699 }
700 StreamFrame::PermissionRequestCancelled { payload, .. } => {
701 self.apply_permission_inner(payload, WorkflowPermissionState::Cancelled)
702 }
703 StreamFrame::UnrestrictedExecution { payload, .. } => {
704 self.apply_permission_inner(payload, WorkflowPermissionState::Unrestricted)
705 }
706 StreamFrame::PermissionGroupCreated { payload, .. }
707 | StreamFrame::PermissionGroupUpdated { payload, .. } => {
708 self.apply_permission_group_inner(payload, false)
709 }
710 StreamFrame::PermissionGroupResolved { payload, .. } => {
711 self.apply_permission_group_inner(payload, true)
712 }
713 _ => PendingDelta::default(),
714 }
715 }
716
717 fn insert_flow(
718 &mut self,
719 run_id: String,
720 flow_name: String,
721 parent_run_id: Option<String>,
722 parent_node_id: Option<String>,
723 now: DateTime<Utc>,
724 fallback_to_root: bool,
725 ) -> PendingDelta {
726 if self.index.node_paths.contains_key(&run_id) {
727 return PendingDelta::default();
728 }
729 let kind = if parent_run_id.is_some() {
730 WorkflowNodeKind::Subflow {
731 run_id: run_id.clone(),
732 flow_name: flow_name.clone(),
733 }
734 } else {
735 WorkflowNodeKind::Flow {
736 run_id: run_id.clone(),
737 flow_name: flow_name.clone(),
738 }
739 };
740 let node = WorkflowNode {
741 id: run_id.clone(),
742 kind,
743 label: flow_name,
744 status: NodeStatus::Running,
745 started_at: Some(now),
746 ended_at: None,
747 output_preview: None,
748 children: Vec::new(),
749 parallelism: Parallelism::Serial,
750 approval: None,
751 llm_stats: None,
752 };
753 let parent_id = parent_run_id
754 .as_deref()
755 .zip(parent_node_id.as_deref())
756 .map(|(parent_run_id, parent_node_id)| scope_id(parent_run_id, parent_node_id));
757 if let Some(parent_id) = parent_id.as_deref()
758 && self.append_child(parent_id, node.clone(), None)
759 {
760 let mut delta = PendingDelta::default();
761 delta.mark_structure(Some(parent_id), &run_id);
762 return delta;
763 }
764 if parent_id.is_none() || fallback_to_root {
765 self.append_root(node, None);
766 let mut delta = PendingDelta::default();
767 delta.mark_structure(None, &run_id);
768 return delta;
769 }
770 PendingDelta::default()
771 }
772
773 fn insert_flow_node(
774 &mut self,
775 run_id: &str,
776 node_id: &str,
777 node_kind: &crate::nodegraph::NodeKind,
778 label: &str,
779 parent_node_id: Option<&str>,
780 now: DateTime<Utc>,
781 ) -> PendingDelta {
782 let id = scope_id(run_id, node_id);
783 if self.index.node_paths.contains_key(&id) {
784 return PendingDelta::default();
785 }
786 let parent_id = parent_node_id
787 .map(|parent| scope_id(run_id, parent))
788 .unwrap_or_else(|| run_id.to_string());
789 let kind = parse_branch_index(node_id).map_or_else(
790 || WorkflowNodeKind::Stmt {
791 node_kind: node_kind.clone(),
792 },
793 |branch_index| WorkflowNodeKind::FanoutBranch { branch_index },
794 );
795 let parallel = matches!(kind, WorkflowNodeKind::FanoutBranch { .. });
796 let node = WorkflowNode {
797 id: id.clone(),
798 kind,
799 label: label.to_string(),
800 status: NodeStatus::Running,
801 started_at: Some(now),
802 ended_at: None,
803 output_preview: None,
804 children: Vec::new(),
805 parallelism: Parallelism::Serial,
806 approval: None,
807 llm_stats: None,
808 };
809 if !self.append_child(&parent_id, node, None) {
810 return PendingDelta::default();
811 }
812 if parallel {
813 self.mutate_node(&parent_id, |parent| {
814 if parent.parallelism == Parallelism::Parallel {
815 false
816 } else {
817 parent.parallelism = Parallelism::Parallel;
818 true
819 }
820 });
821 }
822 let mut delta = PendingDelta::default();
823 delta.mark_structure(Some(&parent_id), &id);
824 delta
825 }
826
827 #[allow(clippy::too_many_arguments)]
828 fn insert_tool(
829 &mut self,
830 run_id: &str,
831 parent_node_id: &str,
832 tool_use_id: &str,
833 tool: &str,
834 args_preview: &str,
835 call_intent: Option<crate::message::ToolCallIntent>,
836 now: DateTime<Utc>,
837 ) -> PendingDelta {
838 let id = tool_node_id(run_id, tool_use_id);
839 if self.index.node_paths.contains_key(&id) {
840 return PendingDelta::default();
841 }
842 let parent_id = scope_id(run_id, parent_node_id);
843 let tool_key = (run_id.to_string(), tool_use_id.to_string());
844 let node = WorkflowNode {
845 id: id.clone(),
846 kind: WorkflowNodeKind::ToolCall {
847 tool_use_id: tool_use_id.to_string(),
848 tool: tool.to_string(),
849 args_preview: args_preview.to_string(),
850 call_intent,
851 result_preview: None,
852 },
853 label: tool.to_string(),
854 status: NodeStatus::Running,
855 started_at: Some(now),
856 ended_at: None,
857 output_preview: None,
858 children: Vec::new(),
859 parallelism: Parallelism::Serial,
860 approval: self.permissions.approval_for_tool(&self.graph, &tool_key),
861 llm_stats: None,
862 };
863 if !self.append_child(&parent_id, node, Some((run_id, tool_use_id))) {
864 return PendingDelta::default();
865 }
866 let mut delta = PendingDelta::default();
867 delta.mark_structure(Some(&parent_id), &id);
868 delta
869 }
870
871 fn finish_node(
872 &mut self,
873 id: &str,
874 status: NodeStatus,
875 output_preview: Option<&str>,
876 now: DateTime<Utc>,
877 recursive: bool,
878 ) -> PendingDelta {
879 let Some(path) = self.index.node_paths.get(id).cloned() else {
880 return PendingDelta::default();
881 };
882 let Some(node) = node_at_path_mut(&mut self.graph.root, &path) else {
883 return PendingDelta::default();
884 };
885 let changed = {
886 let before = (node.status, node.ended_at, node.output_preview.clone());
887 if recursive {
888 cascade_terminate(node, status, now);
889 } else {
890 node.status = status;
891 node.ended_at = Some(now);
892 for child in &mut node.children {
893 if matches!(child.status, NodeStatus::Running | NodeStatus::Pending) {
894 child.status = status;
895 child.ended_at = Some(now);
896 }
897 }
898 }
899 if let Some(output_preview) = output_preview {
900 node.output_preview = Some(output_preview.to_string());
901 }
902 before != (node.status, node.ended_at, node.output_preview.clone())
903 };
904 let summary_changed = if recursive {
905 self.summary.sync_subtree(node, path.len() == 1)
906 } else {
907 let mut changed = self.summary.sync_node(node, path.len() == 1);
908 for child in &node.children {
909 changed |= self.summary.sync_node(child, false);
910 }
911 changed
912 };
913 let mut delta = PendingDelta::default();
914 if changed || summary_changed {
915 delta.mark(id);
916 }
917 delta
918 }
919
920 fn set_llm_stats(&mut self, id: &str, stats: LlmStats) -> PendingDelta {
921 let changed = self.mutate_node(id, |node| {
922 if node.llm_stats.as_ref() == Some(&stats) {
923 false
924 } else {
925 node.llm_stats = Some(stats);
926 true
927 }
928 });
929 let mut delta = PendingDelta::default();
930 if changed {
931 delta.mark(id);
932 }
933 delta
934 }
935
936 fn apply_tool_results(
937 &mut self,
938 run_id: Option<&str>,
939 message: &crate::message::Message,
940 now: DateTime<Utc>,
941 ) -> PendingDelta {
942 let mut delta = PendingDelta::default();
943 for part in &message.parts {
944 let crate::message::MessagePart::ToolResult {
945 tool_use_id,
946 content,
947 is_error,
948 } = part
949 else {
950 continue;
951 };
952 delta.merge(self.finish_tool(run_id, tool_use_id, !*is_error, content, now, true));
953 }
954 delta
955 }
956
957 fn finish_tool(
958 &mut self,
959 run_id: Option<&str>,
960 tool_use_id: &str,
961 ok: bool,
962 preview: &str,
963 now: DateTime<Utc>,
964 set_kind_preview: bool,
965 ) -> PendingDelta {
966 let path = match run_id {
967 Some(run_id) => self
968 .index
969 .tool_paths
970 .get(&(run_id.to_string(), tool_use_id.to_string())),
971 None => self.index.first_tool_paths.get(tool_use_id),
972 }
973 .cloned();
974 let Some(path) = path else {
975 return PendingDelta::default();
976 };
977 let Some(node) = node_at_path_mut(&mut self.graph.root, &path) else {
978 return PendingDelta::default();
979 };
980 let content = if set_kind_preview {
981 preview.chars().take(300).collect::<String>()
982 } else {
983 preview.to_string()
984 };
985 let status = if ok { NodeStatus::Ok } else { NodeStatus::Err };
986 let changed = node.status != status
987 || node.ended_at != Some(now)
988 || node.output_preview.as_deref() != Some(content.as_str());
989 node.status = status;
990 node.ended_at = Some(now);
991 node.output_preview = Some(content.clone());
992 if set_kind_preview
993 && let WorkflowNodeKind::ToolCall { result_preview, .. } = &mut node.kind
994 {
995 *result_preview = Some(content);
996 }
997 let node_id = node.id.clone();
998 let summary_changed = self.summary.sync_node(node, path.len() == 1);
999 let mut delta = PendingDelta::default();
1000 if changed || summary_changed {
1001 delta.mark(node_id);
1002 }
1003 delta
1004 }
1005
1006 fn set_tool_approval(
1007 &mut self,
1008 run_id: &str,
1009 tool_use_id: &str,
1010 approval: ApprovalState,
1011 ) -> PendingDelta {
1012 let id = tool_node_id(run_id, tool_use_id);
1013 let changed = self.mutate_node(&id, |node| {
1014 if node.approval.as_ref() == Some(&approval) {
1015 false
1016 } else {
1017 node.approval = Some(approval);
1018 true
1019 }
1020 });
1021 let mut delta = PendingDelta::default();
1022 if changed {
1023 delta.mark(id);
1024 }
1025 delta
1026 }
1027
1028 fn apply_permission_inner(
1029 &mut self,
1030 payload: &PermissionRequestAudit,
1031 state: WorkflowPermissionState,
1032 ) -> PendingDelta {
1033 let Some(request_id) = payload.request_id.clone() else {
1034 return PendingDelta::default();
1035 };
1036 let identity = WorkflowPermissionIdentity::Canonical { request_id };
1037 self.apply_permission_transition(identity, payload.clone(), state)
1038 }
1039
1040 fn apply_permission_group_inner(
1041 &mut self,
1042 payload: &PermissionGroupAudit,
1043 resolved: bool,
1044 ) -> PendingDelta {
1045 self.apply_permission_group_update(payload, resolved)
1046 }
1047
1048 fn apply_permission_transition(
1049 &mut self,
1050 identity: WorkflowPermissionIdentity,
1051 payload: PermissionRequestAudit,
1052 state: WorkflowPermissionState,
1053 ) -> PendingDelta {
1054 let update = self
1055 .permissions
1056 .apply_request(&mut self.graph, identity, payload, state);
1057 if !update.changed {
1058 return PendingDelta::default();
1059 }
1060 let mut delta = PendingDelta {
1061 projection_changed: true,
1062 layout_changed: update.group_progress_changed,
1063 ..PendingDelta::default()
1064 };
1065 for (run_id, tool_use_id) in update.winner_tools {
1066 let node_id = tool_node_id(&run_id, &tool_use_id);
1067 if !self.index.node_paths.contains_key(&node_id) {
1068 continue;
1069 }
1070 let approval = self
1071 .permissions
1072 .approval_for_tool(&self.graph, &(run_id, tool_use_id));
1073 self.mutate_node(&node_id, |node| {
1074 if node.approval == approval {
1075 false
1076 } else {
1077 node.approval = approval;
1078 true
1079 }
1080 });
1081 delta.mark(node_id);
1082 }
1083 delta
1084 }
1085
1086 fn apply_permission_group_update(
1087 &mut self,
1088 payload: &PermissionGroupAudit,
1089 resolved: bool,
1090 ) -> PendingDelta {
1091 if !self
1092 .permissions
1093 .apply_group(&mut self.graph, payload, resolved)
1094 {
1095 return PendingDelta::default();
1096 }
1097 PendingDelta {
1098 layout_changed: true,
1099 projection_changed: true,
1100 ..PendingDelta::default()
1101 }
1102 }
1103
1104 fn append_root(&mut self, node: WorkflowNode, tool: Option<(&str, &str)>) {
1105 let path = vec![self.graph.root.len()];
1106 let id = node.id.clone();
1107 self.summary.insert_node(&node, &path, true);
1108 self.graph.root.push(node);
1109 self.index.node_paths.insert(id, path.clone());
1110 if let Some((run_id, tool_use_id)) = tool {
1111 self.index_tool(run_id, tool_use_id, path);
1112 }
1113 }
1114
1115 fn append_child(
1116 &mut self,
1117 parent_id: &str,
1118 node: WorkflowNode,
1119 tool: Option<(&str, &str)>,
1120 ) -> bool {
1121 let Some(mut path) = self.index.node_paths.get(parent_id).cloned() else {
1122 return false;
1123 };
1124 let Some(parent) = node_at_path_mut(&mut self.graph.root, &path) else {
1125 return false;
1126 };
1127 path.push(parent.children.len());
1128 let id = node.id.clone();
1129 self.summary.remove_leaf(parent_id);
1130 self.summary.insert_node(&node, &path, false);
1131 parent.children.push(node);
1132 self.index.node_paths.insert(id, path.clone());
1133 if let Some((run_id, tool_use_id)) = tool {
1134 self.index_tool(run_id, tool_use_id, path);
1135 }
1136 true
1137 }
1138
1139 fn index_tool(&mut self, run_id: &str, tool_use_id: &str, path: NodePath) {
1140 self.index.tool_keys_by_node.insert(
1141 tool_node_id(run_id, tool_use_id),
1142 (run_id.to_string(), tool_use_id.to_string()),
1143 );
1144 self.index
1145 .tool_paths
1146 .entry((run_id.to_string(), tool_use_id.to_string()))
1147 .and_modify(|current| {
1148 if path < *current {
1149 *current = path.clone();
1150 }
1151 })
1152 .or_insert_with(|| path.clone());
1153 self.index
1154 .first_tool_paths
1155 .entry(tool_use_id.to_string())
1156 .and_modify(|current| {
1157 if path < *current {
1158 *current = path.clone();
1159 }
1160 })
1161 .or_insert(path);
1162 }
1163
1164 fn mutate_node(&mut self, id: &str, mutation: impl FnOnce(&mut WorkflowNode) -> bool) -> bool {
1165 let Some(path) = self.index.node_paths.get(id).cloned() else {
1166 return false;
1167 };
1168 let Some(node) = node_at_path_mut(&mut self.graph.root, &path) else {
1169 return false;
1170 };
1171 let changed = mutation(node);
1172 let summary_changed = self.summary.sync_node(node, path.len() == 1);
1173 changed || summary_changed
1174 }
1175
1176 fn rebuild_index(&mut self) {
1177 self.index = WorkflowIndex::default();
1178 let mut path = Vec::new();
1179 index_nodes(&self.graph.root, &mut path, None, &mut self.index);
1180 self.permissions = PermissionProjection::rebuild(&self.graph);
1181 self.summary = WorkflowSummary::rebuild(&self.graph);
1182 }
1183
1184 fn commit(&mut self, pending: PendingDelta) -> ProjectionDelta {
1185 if pending.changed() {
1186 self.revision = self.revision.wrapping_add(1);
1187 }
1188 ProjectionDelta {
1189 revision: self.revision,
1190 dirty_nodes: pending.dirty_nodes,
1191 structural_changed: pending.structural_changed,
1192 layout_changed: pending.layout_changed,
1193 }
1194 }
1195}
1196
1197fn index_nodes(
1198 nodes: &[WorkflowNode],
1199 path: &mut NodePath,
1200 inherited_run_id: Option<&str>,
1201 index: &mut WorkflowIndex,
1202) {
1203 for (child_index, node) in nodes.iter().enumerate() {
1204 path.push(child_index);
1205 index
1206 .node_paths
1207 .entry(node.id.clone())
1208 .or_insert_with(|| path.clone());
1209 let run_id = match &node.kind {
1210 WorkflowNodeKind::Flow { run_id, .. } | WorkflowNodeKind::Subflow { run_id, .. } => {
1211 Some(run_id.as_str())
1212 }
1213 _ => inherited_run_id,
1214 };
1215 if let WorkflowNodeKind::ToolCall { tool_use_id, .. } = &node.kind
1216 && let Some(run_id) = run_id
1217 {
1218 index
1219 .tool_keys_by_node
1220 .insert(node.id.clone(), (run_id.to_string(), tool_use_id.clone()));
1221 index
1222 .tool_paths
1223 .entry((run_id.to_string(), tool_use_id.clone()))
1224 .and_modify(|current| {
1225 if path < current {
1226 *current = path.clone();
1227 }
1228 })
1229 .or_insert_with(|| path.clone());
1230 index
1231 .first_tool_paths
1232 .entry(tool_use_id.clone())
1233 .and_modify(|current| {
1234 if path < current {
1235 *current = path.clone();
1236 }
1237 })
1238 .or_insert_with(|| path.clone());
1239 }
1240 index_nodes(&node.children, path, run_id, index);
1241 path.pop();
1242 }
1243}
1244
1245fn node_at_path<'a>(nodes: &'a [WorkflowNode], path: &[usize]) -> Option<&'a WorkflowNode> {
1246 #[cfg(test)]
1247 count_indexed_lookup(path);
1248 let (first, tail) = path.split_first()?;
1249 let mut node = nodes.get(*first)?;
1250 for child in tail {
1251 node = node.children.get(*child)?;
1252 }
1253 Some(node)
1254}
1255
1256fn node_at_path_mut<'a>(
1257 nodes: &'a mut [WorkflowNode],
1258 path: &[usize],
1259) -> Option<&'a mut WorkflowNode> {
1260 #[cfg(test)]
1261 count_indexed_lookup(path);
1262 let (first, tail) = path.split_first()?;
1263 let mut node = nodes.get_mut(*first)?;
1264 for child in tail {
1265 node = node.children.get_mut(*child)?;
1266 }
1267 Some(node)
1268}
1269
1270fn cascade_terminate(node: &mut WorkflowNode, status: NodeStatus, now: DateTime<Utc>) {
1271 if matches!(node.status, NodeStatus::Running | NodeStatus::Pending) {
1272 node.status = status;
1273 node.ended_at = Some(now);
1274 }
1275 for child in &mut node.children {
1276 cascade_terminate(child, status, now);
1277 }
1278}
1279
1280fn collect_flow_run_ids(node: &WorkflowNode, out: &mut Vec<String>) {
1281 match &node.kind {
1282 WorkflowNodeKind::Flow { run_id, .. } | WorkflowNodeKind::Subflow { run_id, .. } => {
1283 out.push(run_id.clone());
1284 }
1285 _ => {}
1286 }
1287 for child in &node.children {
1288 collect_flow_run_ids(child, out);
1289 }
1290}
1291
1292fn scope_id(run_id: &str, node_id: &str) -> String {
1293 format!("{run_id}::{node_id}")
1294}
1295
1296fn tool_node_id(run_id: &str, tool_use_id: &str) -> String {
1297 format!("tool:{run_id}:{tool_use_id}")
1298}
1299
1300fn parse_branch_index(node_id: &str) -> Option<usize> {
1301 let start = node_id.rfind(".branch[")?;
1302 let rest = &node_id[start + ".branch[".len()..];
1303 let end = rest.find(']')?;
1304 rest[..end].parse().ok()
1305}
1306
1307#[cfg(test)]
1308mod tests {
1309 use super::super::workflow_permission::{
1310 PermissionPerfCounters, perf_counters as permission_perf_counters,
1311 reset_perf_counters as reset_permission_perf_counters,
1312 };
1313 use super::*;
1314 use crate::event::FlowRunId;
1315 use crate::permission::{PermissionGroupId, PermissionRequestId};
1316 use crate::permission_audit::{
1317 PermissionAuditTarget, PermissionGroupAudit, PermissionGroupAuditOwner,
1318 PermissionPolicyReference, PermissionRequestAudit,
1319 };
1320
1321 fn start_projection() -> WorkflowProjection {
1322 projection_for_run("root")
1323 }
1324
1325 fn projection_for_run(run_id: &str) -> WorkflowProjection {
1326 let mut projection = WorkflowProjection::new(TurnId::now());
1327 projection.apply_stream_frame(&StreamFrame::FlowStart {
1328 run_id: run_id.into(),
1329 flow_name: "root".into(),
1330 parent_run_id: None,
1331 parent_node_id: None,
1332 });
1333 projection.apply_stream_frame(&StreamFrame::FlowNodeStart {
1334 run_id: run_id.into(),
1335 node_id: "dispatch".into(),
1336 kind: crate::nodegraph::NodeKind::ToolCall {
1337 path: "dispatch_all".into(),
1338 },
1339 label: "dispatch".into(),
1340 parent_node_id: None,
1341 });
1342 projection
1343 }
1344
1345 fn permission_payload(
1346 request_id: PermissionRequestId,
1347 run_id: &FlowRunId,
1348 tool_use_id: impl Into<String>,
1349 at: DateTime<Utc>,
1350 ) -> PermissionRequestAudit {
1351 PermissionRequestAudit {
1352 request_id: Some(request_id),
1353 revision: 1,
1354 session_id: "session".into(),
1355 requesting_run_id: run_id.clone(),
1356 parent_run_id: None,
1357 root_run_id: run_id.clone(),
1358 tool_use_id: tool_use_id.into(),
1359 tool: "fs.read".into(),
1360 call_intent: None,
1361 tier: crate::tool::Tier::Two,
1362 execution_boundary: None,
1363 provenance: Default::default(),
1364 target: PermissionAuditTarget::User,
1365 group_ids: Vec::new(),
1366 policy: PermissionPolicyReference {
1367 snapshot_id: "snapshot".into(),
1368 rule_id: "rule".into(),
1369 },
1370 escalation_path: Vec::new(),
1371 decision_id: None,
1372 actor: None,
1373 scope: None,
1374 reason: None,
1375 at,
1376 }
1377 }
1378
1379 fn add_tool(projection: &mut WorkflowProjection, run_id: &str, tool_use_id: &str) {
1380 projection.apply_stream_frame(&StreamFrame::ToolNode {
1381 run_id: run_id.into(),
1382 parent_node_id: "dispatch".into(),
1383 tool_use_id: tool_use_id.into(),
1384 tool: "fs.read".into(),
1385 args_preview: "{}".into(),
1386 call_intent: None,
1387 });
1388 }
1389
1390 fn permission_group(
1391 group_id: PermissionGroupId,
1392 owner: &FlowRunId,
1393 request_ids: Vec<PermissionRequestId>,
1394 at: DateTime<Utc>,
1395 ) -> PermissionGroupAudit {
1396 PermissionGroupAudit {
1397 group_id,
1398 owner: PermissionGroupAuditOwner::Flow {
1399 run_id: owner.clone(),
1400 },
1401 label: "batch".into(),
1402 request_ids,
1403 revision: 1,
1404 at,
1405 }
1406 }
1407
1408 fn without_timestamps(mut graph: WorkflowGraph) -> WorkflowGraph {
1409 fn clear(nodes: &mut [WorkflowNode]) {
1410 for node in nodes {
1411 node.started_at = None;
1412 node.ended_at = None;
1413 clear(&mut node.children);
1414 }
1415 }
1416 clear(&mut graph.root);
1417 graph
1418 }
1419
1420 #[test]
1421 fn indexed_tool_updates_are_depth_bounded_after_large_append() {
1422 let mut projection = start_projection();
1423 for idx in 0..10_000 {
1424 let delta = projection.apply_stream_frame(&StreamFrame::ToolNode {
1425 run_id: "root".into(),
1426 parent_node_id: "dispatch".into(),
1427 tool_use_id: format!("tool-{idx}"),
1428 tool: "fs.read".into(),
1429 args_preview: "{}".into(),
1430 call_intent: None,
1431 });
1432 assert!(delta.structural_changed);
1433 }
1434 assert_eq!(projection.index.node_paths.len(), 10_002);
1435 assert_eq!(projection.index.node_paths["tool:root:tool-9999"].len(), 3);
1436
1437 reset_perf_counters();
1438 let delta = projection.apply_stream_frame(&StreamFrame::ToolUseDone {
1439 tool: "fs.read".into(),
1440 id: "tool-9999".into(),
1441 ok: true,
1442 preview: "done".into(),
1443 });
1444 assert_eq!(
1445 perf_counters(),
1446 PerfCounters {
1447 indexed_lookups: 1,
1448 path_steps: 3,
1449 }
1450 );
1451 assert_eq!(delta.dirty_nodes, ["tool:root:tool-9999"]);
1452 assert_eq!(
1453 projection.find_node("tool:root:tool-9999").unwrap().status,
1454 NodeStatus::Ok
1455 );
1456 }
1457
1458 #[test]
1459 fn duplicate_structural_frames_preserve_revision_and_paths() {
1460 let mut projection = start_projection();
1461 let frame = StreamFrame::ToolNode {
1462 run_id: "root".into(),
1463 parent_node_id: "dispatch".into(),
1464 tool_use_id: "same".into(),
1465 tool: "fs.read".into(),
1466 args_preview: "{}".into(),
1467 call_intent: None,
1468 };
1469 let first = projection.apply_stream_frame(&frame);
1470 let revision = first.revision;
1471 let duplicate = projection.apply_stream_frame(&frame);
1472 assert!(!duplicate.changed());
1473 assert_eq!(duplicate.revision, revision);
1474 assert_eq!(
1475 projection
1476 .find_node("root::dispatch")
1477 .unwrap()
1478 .children
1479 .len(),
1480 1
1481 );
1482 }
1483
1484 #[test]
1485 fn append_only_mutations_preserve_existing_paths() {
1486 let mut projection = start_projection();
1487 projection.apply_stream_frame(&StreamFrame::ToolNode {
1488 run_id: "root".into(),
1489 parent_node_id: "dispatch".into(),
1490 tool_use_id: "first".into(),
1491 tool: "fs.read".into(),
1492 args_preview: "{}".into(),
1493 call_intent: None,
1494 });
1495 let path = projection.index.node_paths["tool:root:first"].clone();
1496
1497 projection.apply_stream_frame(&StreamFrame::ToolNode {
1498 run_id: "root".into(),
1499 parent_node_id: "dispatch".into(),
1500 tool_use_id: "second".into(),
1501 tool: "fs.read".into(),
1502 args_preview: "{}".into(),
1503 call_intent: None,
1504 });
1505
1506 assert_eq!(projection.index.node_paths["tool:root:first"], path);
1507 assert_eq!(
1508 projection.find_node("tool:root:first").unwrap().label,
1509 "fs.read"
1510 );
1511 }
1512
1513 #[test]
1514 fn batch_commits_one_revision_and_reports_all_dirty_nodes() {
1515 let run_id = crate::event::FlowRunId::now();
1516 let events = vec![
1517 Event::FlowStart {
1518 run_id: run_id.clone(),
1519 flow_name: "root".into(),
1520 parent_run_id: None,
1521 parent_node_id: None,
1522 spawned: false,
1523 },
1524 Event::FlowNodeStart {
1525 run_id: run_id.clone(),
1526 node_id: "dispatch".into(),
1527 kind: crate::nodegraph::NodeKind::ToolCall {
1528 path: "dispatch_all".into(),
1529 },
1530 label: "dispatch".into(),
1531 parent_node_id: None,
1532 },
1533 ];
1534 let mut projection = WorkflowProjection::new(TurnId::now());
1535 let delta = projection.apply_batch(&events);
1536 assert_eq!(delta.revision, 1);
1537 assert!(delta.structural_changed);
1538 assert_eq!(delta.dirty_nodes.len(), 2);
1539 assert!(
1540 projection
1541 .find_node(&scope_id(&run_id.0.to_string(), "dispatch"))
1542 .is_some()
1543 );
1544 }
1545
1546 #[test]
1547 fn serialization_preserves_graph_shape_and_rebuilds_indices() {
1548 let mut projection = start_projection();
1549 projection.apply_stream_frame(&StreamFrame::ToolNode {
1550 run_id: "root".into(),
1551 parent_node_id: "dispatch".into(),
1552 tool_use_id: "tool".into(),
1553 tool: "fs.read".into(),
1554 args_preview: "{}".into(),
1555 call_intent: None,
1556 });
1557 let graph_json = serde_json::to_value(projection.graph()).unwrap();
1558 let projection_json = serde_json::to_value(&projection).unwrap();
1559 assert_eq!(projection_json, graph_json);
1560
1561 let mut restored: WorkflowProjection = serde_json::from_value(projection_json).unwrap();
1562 assert!(restored.find_node("tool:root:tool").is_some());
1563 let delta = restored.apply_stream_frame(&StreamFrame::ToolUseDone {
1564 tool: "fs.read".into(),
1565 id: "tool".into(),
1566 ok: true,
1567 preview: "done".into(),
1568 });
1569 assert!(delta.changed());
1570 assert_eq!(
1571 restored.find_node("tool:root:tool").unwrap().status,
1572 NodeStatus::Ok
1573 );
1574 }
1575
1576 #[test]
1577 fn indexed_stream_application_matches_recursive_graph_semantics() {
1578 let run_id = crate::event::FlowRunId::now().0.to_string();
1579 let now = Utc::now();
1580 let tool_result = crate::message::Message {
1581 role: crate::message::MessageRole::Tool,
1582 parts: vec![crate::message::MessagePart::ToolResult {
1583 tool_use_id: "tool".into(),
1584 content: "contents".into(),
1585 is_error: false,
1586 }],
1587 turn_id: TurnId::now(),
1588 origin: crate::message::MessageOrigin::User,
1589 };
1590 let frames = vec![
1591 StreamFrame::FlowStart {
1592 run_id: run_id.clone(),
1593 flow_name: "root".into(),
1594 parent_run_id: None,
1595 parent_node_id: None,
1596 },
1597 StreamFrame::FlowNodeStart {
1598 run_id: run_id.clone(),
1599 node_id: "dispatch".into(),
1600 kind: crate::nodegraph::NodeKind::ToolCall {
1601 path: "dispatch_all".into(),
1602 },
1603 label: "dispatch".into(),
1604 parent_node_id: None,
1605 },
1606 StreamFrame::ToolNode {
1607 run_id: run_id.clone(),
1608 parent_node_id: "dispatch".into(),
1609 tool_use_id: "tool".into(),
1610 tool: "fs.read".into(),
1611 args_preview: "{}".into(),
1612 call_intent: None,
1613 },
1614 StreamFrame::ToolPendingApproval {
1615 run_id: run_id.clone(),
1616 tool_use_id: "tool".into(),
1617 tool_name: "fs.read".into(),
1618 args_preview: "{}".into(),
1619 level: "two".into(),
1620 preview: Some("read a file".into()),
1621 },
1622 StreamFrame::ToolResultMsg {
1623 flow_run_id: Some(run_id.clone()),
1624 message: tool_result,
1625 },
1626 StreamFrame::FlowNodeEnd {
1627 run_id: run_id.clone(),
1628 node_id: "dispatch".into(),
1629 status: FlowNodeStatus::Ok,
1630 output_preview: Some("done".into()),
1631 parent_node_id: None,
1632 },
1633 StreamFrame::FlowDone {
1634 run_id: run_id.clone(),
1635 flow_name: "root".into(),
1636 ok: true,
1637 cancelled: false,
1638 suicide: false,
1639 },
1640 ];
1641 let mut recursive = WorkflowGraph::new(TurnId::now());
1642 let mut indexed = WorkflowProjection::new(recursive.turn_id.clone());
1643
1644 for frame in frames {
1645 recursive.apply_stream_frame_at(&frame, Some(now));
1646 indexed.apply_stream_frame_at(&frame, Some(now));
1647 assert_eq!(
1648 without_timestamps(recursive.clone()),
1649 without_timestamps(indexed.graph().clone())
1650 );
1651 }
1652 }
1653
1654 #[test]
1655 fn scoped_tool_result_never_falls_back_to_a_different_run() {
1656 let mut projection = start_projection();
1657 projection.apply_stream_frame(&StreamFrame::ToolNode {
1658 run_id: "root".into(),
1659 parent_node_id: "dispatch".into(),
1660 tool_use_id: "shared".into(),
1661 tool: "fs.read".into(),
1662 args_preview: "{}".into(),
1663 call_intent: None,
1664 });
1665 let message = crate::message::Message {
1666 role: crate::message::MessageRole::Tool,
1667 parts: vec![crate::message::MessagePart::ToolResult {
1668 tool_use_id: "shared".into(),
1669 content: "wrong run".into(),
1670 is_error: false,
1671 }],
1672 turn_id: TurnId::now(),
1673 origin: crate::message::MessageOrigin::User,
1674 };
1675
1676 let delta = projection.apply_stream_frame(&StreamFrame::ToolResultMsg {
1677 flow_run_id: Some("other".into()),
1678 message,
1679 });
1680
1681 assert!(!delta.changed());
1682 assert_eq!(
1683 projection.find_node("tool:root:shared").unwrap().status,
1684 NodeStatus::Running
1685 );
1686 }
1687
1688 #[test]
1689 fn unscoped_tool_completion_preserves_depth_first_selection() {
1690 let mut graph = WorkflowGraph::new(TurnId::now());
1691 for run_id in ["first", "second"] {
1692 let mut flow = WorkflowNode {
1693 id: run_id.into(),
1694 kind: WorkflowNodeKind::Flow {
1695 run_id: run_id.into(),
1696 flow_name: run_id.into(),
1697 },
1698 label: run_id.into(),
1699 status: NodeStatus::Running,
1700 started_at: None,
1701 ended_at: None,
1702 output_preview: None,
1703 children: Vec::new(),
1704 parallelism: Parallelism::Serial,
1705 approval: None,
1706 llm_stats: None,
1707 };
1708 flow.children.push(WorkflowNode {
1709 id: tool_node_id(run_id, "shared"),
1710 kind: WorkflowNodeKind::ToolCall {
1711 tool_use_id: "shared".into(),
1712 tool: "fs.read".into(),
1713 args_preview: "{}".into(),
1714 call_intent: None,
1715 result_preview: None,
1716 },
1717 label: "fs.read".into(),
1718 status: NodeStatus::Running,
1719 started_at: None,
1720 ended_at: None,
1721 output_preview: None,
1722 children: Vec::new(),
1723 parallelism: Parallelism::Serial,
1724 approval: None,
1725 llm_stats: None,
1726 });
1727 graph.root.push(flow);
1728 }
1729 let mut projection = WorkflowProjection::from(graph);
1730
1731 projection.apply_stream_frame(&StreamFrame::ToolUseDone {
1732 tool: "fs.read".into(),
1733 id: "shared".into(),
1734 ok: true,
1735 preview: "done".into(),
1736 });
1737
1738 assert_eq!(
1739 projection.find_node("tool:first:shared").unwrap().status,
1740 NodeStatus::Ok
1741 );
1742 assert_eq!(
1743 projection.find_node("tool:second:shared").unwrap().status,
1744 NodeStatus::Running
1745 );
1746 }
1747
1748 #[test]
1749 fn indexed_permission_transition_is_bounded_after_large_rebuild() {
1750 const ENTRIES: usize = 10_000;
1751 let run_id = FlowRunId::now();
1752 let run_id_text = run_id.0.to_string();
1753 let at = Utc::now();
1754 let mut projection = projection_for_run(&run_id_text);
1755 for idx in 0..ENTRIES {
1756 add_tool(&mut projection, &run_id_text, &format!("tool-{idx}"));
1757 }
1758 let mut graph = projection.into_graph();
1759 let mut final_request = None;
1760 for idx in 0..ENTRIES {
1761 let request_id = PermissionRequestId::now();
1762 let identity = WorkflowPermissionIdentity::Canonical {
1763 request_id: request_id.clone(),
1764 };
1765 let payload = permission_payload(request_id, &run_id, format!("tool-{idx}"), at);
1766 graph.permission_requests.insert(
1767 identity.clone(),
1768 WorkflowPermissionRequest {
1769 payload: payload.clone(),
1770 state: WorkflowPermissionState::Pending,
1771 },
1772 );
1773 if idx + 1 == ENTRIES {
1774 final_request = Some((identity, payload));
1775 }
1776 }
1777 let mut projection = WorkflowProjection::from(graph);
1778 assert_eq!(projection.permissions.pending_tool_count(), ENTRIES);
1779 let (identity, mut payload) = final_request.unwrap();
1780 payload.reason = Some("approved".into());
1781 payload.at += chrono::Duration::seconds(1);
1782
1783 reset_perf_counters();
1784 reset_permission_perf_counters();
1785 let delta = projection.apply_permission_request_with_identity(
1786 identity,
1787 &payload,
1788 WorkflowPermissionState::Approved,
1789 );
1790 let workflow_counters = perf_counters();
1791 let permission_counters = permission_perf_counters();
1792
1793 assert_eq!(
1794 workflow_counters,
1795 PerfCounters {
1796 indexed_lookups: 1,
1797 path_steps: 3,
1798 }
1799 );
1800 assert_eq!(
1801 permission_counters,
1802 PermissionPerfCounters {
1803 winner_selections: 1,
1804 group_member_visits: 0,
1805 }
1806 );
1807 assert_eq!(delta.dirty_nodes, [format!("tool:{run_id_text}:tool-9999")]);
1808 assert_eq!(projection.permissions.pending_tool_count(), ENTRIES - 1);
1809 }
1810
1811 #[test]
1812 fn permission_before_tool_attaches_without_refreshing_other_nodes() {
1813 let run_id = FlowRunId::now();
1814 let run_id_text = run_id.0.to_string();
1815 let request_id = PermissionRequestId::now();
1816 let payload = permission_payload(request_id, &run_id, "late-tool", Utc::now());
1817 let mut projection = projection_for_run(&run_id_text);
1818
1819 projection.apply_permission_request(&payload, WorkflowPermissionState::Pending);
1820 reset_perf_counters();
1821 add_tool(&mut projection, &run_id_text, "late-tool");
1822
1823 assert_eq!(
1824 perf_counters(),
1825 PerfCounters {
1826 indexed_lookups: 1,
1827 path_steps: 2,
1828 }
1829 );
1830 let node_id = format!("tool:{run_id_text}:late-tool");
1831 assert!(matches!(
1832 projection.find_node(&node_id).unwrap().approval,
1833 Some(ApprovalState::Pending { .. })
1834 ));
1835 assert_eq!(
1836 projection
1837 .permission_request_for_node(&node_id)
1838 .unwrap()
1839 .payload
1840 .tool_use_id,
1841 "late-tool"
1842 );
1843 }
1844
1845 #[test]
1846 fn indexed_permission_winner_matches_recursive_canonical_and_legacy_semantics() {
1847 let run_id = FlowRunId::now();
1848 let run_id_text = run_id.0.to_string();
1849 let at = Utc::now();
1850 let canonical_id = PermissionRequestId::now();
1851 let legacy_request_id = PermissionRequestId::now();
1852 let canonical_identity = WorkflowPermissionIdentity::Canonical {
1853 request_id: canonical_id.clone(),
1854 };
1855 let legacy_identity = WorkflowPermissionIdentity::Legacy {
1856 seq: 7,
1857 run_id: run_id_text.clone(),
1858 tool_use_id: "shared".into(),
1859 };
1860 let mut canonical = permission_payload(canonical_id, &run_id, "shared", at);
1861 canonical.reason = Some("canonical".into());
1862 let mut legacy = permission_payload(legacy_request_id, &run_id, "shared", at);
1863 legacy.reason = Some("legacy".into());
1864 let mut indexed = projection_for_run(&run_id_text);
1865 add_tool(&mut indexed, &run_id_text, "shared");
1866 let mut recursive = indexed.graph().clone();
1867
1868 for (identity, payload, state) in [
1869 (
1870 canonical_identity.clone(),
1871 canonical.clone(),
1872 WorkflowPermissionState::Denied,
1873 ),
1874 (
1875 legacy_identity.clone(),
1876 legacy.clone(),
1877 WorkflowPermissionState::Denied,
1878 ),
1879 ] {
1880 recursive.apply_permission_request_with_identity(identity.clone(), &payload, state);
1881 indexed.apply_permission_request_with_identity(identity, &payload, state);
1882 }
1883 assert_eq!(indexed.graph(), &recursive);
1884 let node_id = format!("tool:{run_id_text}:shared");
1885 assert_eq!(
1886 indexed
1887 .permission_request_for_node(&node_id)
1888 .unwrap()
1889 .payload
1890 .reason
1891 .as_deref(),
1892 Some("canonical")
1893 );
1894
1895 legacy.at -= chrono::Duration::seconds(1);
1896 recursive.apply_permission_request_with_identity(
1897 legacy_identity.clone(),
1898 &legacy,
1899 WorkflowPermissionState::Pending,
1900 );
1901 indexed.apply_permission_request_with_identity(
1902 legacy_identity.clone(),
1903 &legacy,
1904 WorkflowPermissionState::Pending,
1905 );
1906 canonical.at += chrono::Duration::seconds(10);
1907 recursive.apply_permission_request_with_identity(
1908 canonical_identity.clone(),
1909 &canonical,
1910 WorkflowPermissionState::Approved,
1911 );
1912 indexed.apply_permission_request_with_identity(
1913 canonical_identity,
1914 &canonical,
1915 WorkflowPermissionState::Approved,
1916 );
1917 assert_eq!(indexed.graph(), &recursive);
1918 assert!(matches!(
1919 indexed.find_node(&node_id).unwrap().approval,
1920 Some(ApprovalState::Pending { .. })
1921 ));
1922
1923 recursive.apply_permission_request_with_identity(
1924 legacy_identity.clone(),
1925 &legacy,
1926 WorkflowPermissionState::Denied,
1927 );
1928 indexed.apply_permission_request_with_identity(
1929 legacy_identity,
1930 &legacy,
1931 WorkflowPermissionState::Denied,
1932 );
1933 assert_eq!(indexed.graph(), &recursive);
1934 assert_eq!(
1935 indexed.find_node(&node_id).unwrap().approval,
1936 Some(ApprovalState::Approved)
1937 );
1938 }
1939
1940 #[test]
1941 fn repeated_tool_use_ids_remain_scoped_by_run() {
1942 let run_a = FlowRunId::now();
1943 let run_b = FlowRunId::now();
1944 let run_a_text = run_a.0.to_string();
1945 let run_b_text = run_b.0.to_string();
1946 let mut projection = projection_for_run(&run_a_text);
1947 projection.apply_stream_frame(&StreamFrame::FlowStart {
1948 run_id: run_b_text.clone(),
1949 flow_name: "second".into(),
1950 parent_run_id: None,
1951 parent_node_id: None,
1952 });
1953 projection.apply_stream_frame(&StreamFrame::FlowNodeStart {
1954 run_id: run_b_text.clone(),
1955 node_id: "dispatch".into(),
1956 kind: crate::nodegraph::NodeKind::ToolCall {
1957 path: "dispatch_all".into(),
1958 },
1959 label: "dispatch".into(),
1960 parent_node_id: None,
1961 });
1962 add_tool(&mut projection, &run_a_text, "shared");
1963 add_tool(&mut projection, &run_b_text, "shared");
1964 let payload_a =
1965 permission_payload(PermissionRequestId::now(), &run_a, "shared", Utc::now());
1966 let payload_b =
1967 permission_payload(PermissionRequestId::now(), &run_b, "shared", Utc::now());
1968
1969 projection.apply_permission_request(&payload_a, WorkflowPermissionState::Pending);
1970 projection.apply_permission_request(&payload_b, WorkflowPermissionState::Approved);
1971
1972 assert!(matches!(
1973 projection
1974 .find_node(&format!("tool:{run_a_text}:shared"))
1975 .unwrap()
1976 .approval,
1977 Some(ApprovalState::Pending { .. })
1978 ));
1979 assert_eq!(
1980 projection
1981 .find_node(&format!("tool:{run_b_text}:shared"))
1982 .unwrap()
1983 .approval,
1984 Some(ApprovalState::Approved)
1985 );
1986 }
1987
1988 #[test]
1989 fn group_progress_updates_incrementally_and_rebuilds_from_graph_dto() {
1990 let run_id = FlowRunId::now();
1991 let run_id_text = run_id.0.to_string();
1992 let at = Utc::now();
1993 let group_id = PermissionGroupId(uuid::Uuid::now_v7());
1994 let request_a_id = PermissionRequestId::now();
1995 let request_b_id = PermissionRequestId::now();
1996 let mut request_a = permission_payload(request_a_id.clone(), &run_id, "a", at);
1997 let mut request_b = permission_payload(request_b_id.clone(), &run_id, "b", at);
1998 request_a.group_ids.push(group_id.clone());
1999 request_b.group_ids.push(group_id.clone());
2000 let group = permission_group(
2001 group_id.clone(),
2002 &run_id,
2003 vec![request_a_id.clone(), request_a_id, request_b_id],
2004 at,
2005 );
2006 let mut projection = projection_for_run(&run_id_text);
2007 add_tool(&mut projection, &run_id_text, "a");
2008 add_tool(&mut projection, &run_id_text, "b");
2009 projection.apply_permission_request(&request_a, WorkflowPermissionState::Pending);
2010 projection.apply_permission_request(&request_b, WorkflowPermissionState::Pending);
2011 projection.apply_permission_group(&group, false);
2012 assert_eq!(
2013 projection.permission_group_progress(&group_id),
2014 Some((0, 3))
2015 );
2016
2017 request_a.at += chrono::Duration::seconds(1);
2018 reset_permission_perf_counters();
2019 projection.apply_permission_request(&request_a, WorkflowPermissionState::Approved);
2020 assert_eq!(
2021 projection.permission_group_progress(&group_id),
2022 Some((2, 3))
2023 );
2024 assert_eq!(
2025 permission_perf_counters(),
2026 PermissionPerfCounters {
2027 winner_selections: 1,
2028 group_member_visits: 0,
2029 }
2030 );
2031
2032 let restored = WorkflowProjection::from(projection.into_graph());
2033 assert_eq!(restored.permission_group_progress(&group_id), Some((2, 3)));
2034 assert_eq!(restored.descendant_pending_permissions(&run_id_text), 1);
2035 }
2036
2037 #[test]
2038 fn interrupt_updates_canonical_and_legacy_pending_indices_then_becomes_noop() {
2039 let run_id = FlowRunId::now();
2040 let run_id_text = run_id.0.to_string();
2041 let canonical_id = PermissionRequestId::now();
2042 let canonical = permission_payload(canonical_id.clone(), &run_id, "canonical", Utc::now());
2043 let legacy = permission_payload(PermissionRequestId::now(), &run_id, "legacy", Utc::now());
2044 let mut projection = projection_for_run(&run_id_text);
2045 add_tool(&mut projection, &run_id_text, "canonical");
2046 add_tool(&mut projection, &run_id_text, "legacy");
2047 projection.apply_permission_request_with_identity(
2048 WorkflowPermissionIdentity::Canonical {
2049 request_id: canonical_id,
2050 },
2051 &canonical,
2052 WorkflowPermissionState::Pending,
2053 );
2054 projection.apply_permission_request_with_identity(
2055 WorkflowPermissionIdentity::Legacy {
2056 seq: 11,
2057 run_id: run_id_text.clone(),
2058 tool_use_id: "legacy".into(),
2059 },
2060 &legacy,
2061 WorkflowPermissionState::Pending,
2062 );
2063 assert_eq!(projection.descendant_pending_permissions(&run_id_text), 2);
2064 assert_eq!(projection.permissions.pending_tool_count(), 2);
2065
2066 let interrupted = projection.interrupt_pending_permissions();
2067 assert!(interrupted.changed());
2068 assert_eq!(interrupted.dirty_nodes.len(), 2);
2069 assert_eq!(projection.descendant_pending_permissions(&run_id_text), 0);
2070 assert_eq!(projection.permissions.pending_tool_count(), 0);
2071 assert!(
2072 projection
2073 .graph
2074 .permission_requests
2075 .values()
2076 .all(|request| {
2077 request.state == WorkflowPermissionState::Interrupted
2078 && request.payload.reason.as_deref()
2079 == Some("interrupted at end of persisted history")
2080 })
2081 );
2082
2083 let revision = projection.revision();
2084 let noop = projection.interrupt_pending_permissions();
2085 assert!(!noop.changed());
2086 assert_eq!(noop.revision, revision);
2087 assert_eq!(projection.revision(), revision);
2088 }
2089
2090 #[test]
2091 fn summary_tracks_incremental_structure_status_time_and_llm_usage() {
2092 let now = Utc::now();
2093 let mut projection = WorkflowProjection::new(TurnId::now());
2094 projection.apply_stream_frame_at(
2095 &StreamFrame::FlowStart {
2096 run_id: "root".into(),
2097 flow_name: "root".into(),
2098 parent_run_id: None,
2099 parent_node_id: None,
2100 },
2101 Some(now),
2102 );
2103 projection.apply_stream_frame_at(
2104 &StreamFrame::FlowNodeStart {
2105 run_id: "root".into(),
2106 node_id: "dispatch".into(),
2107 kind: crate::nodegraph::NodeKind::ToolCall {
2108 path: "dispatch_all".into(),
2109 },
2110 label: "dispatch".into(),
2111 parent_node_id: None,
2112 },
2113 Some(now + chrono::Duration::seconds(1)),
2114 );
2115 assert_eq!(projection.summary().counts().nodes, 2);
2116 assert_eq!(projection.summary().collapsed_leaf_paths(8), [vec![0, 0]]);
2117
2118 for (tool_use_id, tool) in [
2119 ("read", "fs.read"),
2120 ("agent", "flow.spawn"),
2121 ("edit", "fs.write"),
2122 ] {
2123 projection.apply_stream_frame_at(
2124 &StreamFrame::ToolNode {
2125 run_id: "root".into(),
2126 parent_node_id: "dispatch".into(),
2127 tool_use_id: tool_use_id.into(),
2128 tool: tool.into(),
2129 args_preview: "{}".into(),
2130 call_intent: None,
2131 },
2132 Some(now + chrono::Duration::seconds(2)),
2133 );
2134 }
2135 assert_eq!(
2136 projection.summary().counts(),
2137 WorkflowCounts {
2138 nodes: 5,
2139 agents: 1,
2140 tools: 3,
2141 edits: 1,
2142 }
2143 );
2144 assert_eq!(
2145 projection.summary().status(),
2146 WorkflowAggregateStatus::Running
2147 );
2148 assert!(
2149 projection
2150 .summary()
2151 .collapsed_leaf_paths(8)
2152 .iter()
2153 .all(|path| path.len() == 3)
2154 );
2155
2156 projection.apply_stream_frame_at(
2157 &StreamFrame::LlmCallStats {
2158 model: "reasoning-model".into(),
2159 provider: "openai-compatible".into(),
2160 context_call_purpose: crate::context_plan::ContextCallPurpose::General,
2161 context_call_scope: crate::context_plan::ContextCallScope::Root,
2162 input_tokens: 100,
2163 output_tokens: 20,
2164 cache_read: 300,
2165 cache_write: 40,
2166 ttft_ms: 50,
2167 tokens_per_second: 60.0,
2168 wallclock_ms: 70,
2169 run_id: Some("root".into()),
2170 node_id: Some("dispatch".into()),
2171 },
2172 Some(now + chrono::Duration::seconds(3)),
2173 );
2174 let aggregate = projection
2175 .summary()
2176 .llm_routes()
2177 .values()
2178 .next()
2179 .expect("LLM route summary");
2180 assert_eq!(aggregate.calls, 1);
2181 assert_eq!(aggregate.total_in, 440);
2182 assert_eq!(aggregate.total_out, 20);
2183
2184 projection.apply_stream_frame_at(
2185 &StreamFrame::FlowDone {
2186 run_id: "root".into(),
2187 flow_name: "root".into(),
2188 ok: true,
2189 cancelled: false,
2190 suicide: false,
2191 },
2192 Some(now + chrono::Duration::seconds(5)),
2193 );
2194 assert_eq!(projection.summary().status(), WorkflowAggregateStatus::Ok);
2195 assert_eq!(projection.summary().started_at(), Some(now));
2196 assert_eq!(
2197 projection.summary().ended_at(),
2198 Some(now + chrono::Duration::seconds(5))
2199 );
2200 assert_eq!(
2201 projection
2202 .summary()
2203 .elapsed_secs(now + chrono::Duration::seconds(100)),
2204 5
2205 );
2206
2207 let restored = WorkflowProjection::from(projection.graph().clone());
2208 assert_eq!(restored.summary().counts(), projection.summary().counts());
2209 assert_eq!(restored.summary().status(), projection.summary().status());
2210 assert_eq!(
2211 restored.summary().collapsed_leaf_paths(8),
2212 projection.summary().collapsed_leaf_paths(8)
2213 );
2214 assert_eq!(
2215 restored.summary().llm_routes(),
2216 projection.summary().llm_routes()
2217 );
2218 }
2219}