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