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