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