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}