Skip to main content

atman_runtime/projection/
message_window.rs

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