Skip to main content

atman_runtime/projection/
workflow.rs

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}