Skip to main content

atman_runtime/
workflow.rs

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