use std::fs::File;
use std::io::{BufRead, BufReader};
use std::path::Path;
use crate::cm_turn_layout::event::TurnEvent;
use crate::cm_turn_layout::model::{SegmentKind, Turn};
use crate::cm_turn_layout::project::project_turn_web;
use crate::cm_turn_layout::reduce::TurnReducer;
#[derive(Debug, serde::Deserialize)]
struct SseReplayLine {
#[allow(dead_code)]
seq: u64,
#[allow(dead_code)]
job_id: u64,
data: String,
}
#[derive(Debug, serde::Deserialize)]
struct TurnReplayLine {
#[allow(dead_code)]
seq: u64,
#[allow(dead_code)]
job_id: u64,
event: String,
#[serde(default)]
detail: serde_json::Value,
#[serde(default)]
#[expect(dead_code, reason = "反序列化占位")]
title: String,
#[allow(dead_code)]
#[serde(default)]
replay_turn_seq: u64,
}
fn map_turn_replay_event_to_turn_events(line: &TurnReplayLine) -> Vec<TurnEvent> {
match line.event.as_str() {
"llm_response_done" => map_llm_response_done(&line.detail),
_ => Vec::new(),
}
}
fn map_llm_response_done(det: &serde_json::Value) -> Vec<TurnEvent> {
let mut out = Vec::new();
let tool_calls = det.get("tool_calls").and_then(|v| v.as_array());
let has_tools = tool_calls.is_some_and(|tcs| !tcs.is_empty());
if let Some(tcs) = tool_calls {
for tc in tcs {
let name = tc
.get("function")
.and_then(|f| f.get("name"))
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let tool_call_id = tc
.get("id")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
if !name.is_empty() {
out.push(TurnEvent::ToolCall {
tool_call_id,
name,
summary: String::new(),
});
}
}
}
if !has_tools {
out.push(TurnEvent::ToolPhaseEnd);
}
let _ = det.get("assistant_content");
out
}
fn parse_segment_start(data: &serde_json::Value) -> Option<TurnEvent> {
let segment_id = data
.get("segmentId")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let kind_str = data
.get("kind")
.and_then(|v| v.as_str())
.unwrap_or("commentary");
let kind = match kind_str {
"answer" => SegmentKind::Answer,
_ => SegmentKind::Commentary,
};
let before_tool_call_id = data
.get("beforeToolCallId")
.and_then(|v| v.as_str())
.filter(|s| !s.is_empty())
.map(|s| s.to_string());
Some(TurnEvent::SegmentStart {
segment_id,
kind,
before_tool_call_id,
})
}
fn parse_segment_end(data: &serde_json::Value) -> Option<TurnEvent> {
let segment_id = data
.get("segmentId")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
Some(TurnEvent::SegmentEnd { segment_id })
}
fn parse_timeline_log(data: &serde_json::Value) -> Option<TurnEvent> {
let kind = data.get("kind").and_then(|v| v.as_str()).unwrap_or("");
match kind {
"approval_decision" | "tool_result_summary" => {
let text = data
.get("title")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
if text.is_empty() {
None
} else {
Some(TurnEvent::TimelineAssistant { text })
}
}
"final_response" => {
None
}
_ => None,
}
}
fn map_single_sse_value(val: &serde_json::Value) -> Vec<TurnEvent> {
let Some(type_str) = val.get("type").and_then(|v| v.as_str()) else {
return Vec::new();
};
match type_str {
"TEXT_MESSAGE_CONTENT" => {
Vec::new()
}
"TOOL_CALL_START" => {
let name = val
.get("name")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let tool_call_id = val
.get("toolCallId")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
if name.is_empty() {
Vec::new()
} else {
vec![TurnEvent::ToolCall {
tool_call_id,
name,
summary: String::new(),
}]
}
}
"CUSTOM" => {
let custom_type = val.get("customType").and_then(|v| v.as_str());
let data = val.get("data");
match custom_type {
Some("turn_segment_start") => {
data.and_then(parse_segment_start).into_iter().collect()
}
Some("turn_segment_end") => data.and_then(parse_segment_end).into_iter().collect(),
Some("turn_tool_phase_end") => vec![TurnEvent::ToolPhaseEnd],
Some("timeline_log") => data.and_then(parse_timeline_log).into_iter().collect(),
_ => Vec::new(),
}
}
_ => Vec::new(),
}
}
fn map_sse_data_to_turn_events(data: &str) -> Vec<TurnEvent> {
let mut out = Vec::new();
for line in data.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
if let Ok(val) = serde_json::from_str::<serde_json::Value>(line) {
out.extend(map_single_sse_value(&val));
}
}
out
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum JsonlFormat {
SseReplay,
TurnReplay,
}
fn detect_format(first_line: &str) -> Result<JsonlFormat, String> {
let val: serde_json::Value =
serde_json::from_str(first_line).map_err(|e| format!("JSON 解析失败: {e}"))?;
if val.get("data").is_some() {
Ok(JsonlFormat::SseReplay)
} else if val.get("event").is_some() {
Ok(JsonlFormat::TurnReplay)
} else {
Err("无法识别 JSONL 格式:缺少 `data` 或 `event` 字段".to_string())
}
}
pub fn replay_sse_events_to_web_rows(
path: &Path,
) -> Result<Vec<crate::cm_turn_layout::project::ProjectedRow>, String> {
let turns = replay_all_turns(path)?;
let mut rows = Vec::new();
for (i, turn) in turns.iter().enumerate() {
if turns.len() > 1 {
rows.push(crate::cm_turn_layout::project::ProjectedRow {
kind: "turn_marker".to_string(),
text: format!("── Turn {} ──", i + 1),
tool_name: None,
tool_call_id: None,
});
}
rows.extend(project_turn_web(turn));
}
Ok(rows)
}
pub fn replay_sse_events_to_turn(path: &Path) -> Result<Turn, String> {
let all = replay_all_turns(path)?;
all.into_iter()
.last()
.ok_or_else(|| "JSONL 文件中无有效 Turn 事件".to_string())
}
fn process_replay_line(
line: &str,
format: JsonlFormat,
current_turn_seq: &mut u64,
current_turn: &mut Turn,
turns: &mut Vec<Turn>,
reducer: &TurnReducer,
) -> Result<usize, String> {
let events = match format {
JsonlFormat::SseReplay => {
let replay_line: SseReplayLine = serde_json::from_str(line)
.map_err(|e| format!("SSE replay JSON 解析失败: {e}\n 行: {line}"))?;
map_sse_data_to_turn_events(&replay_line.data)
}
JsonlFormat::TurnReplay => {
let replay_line: TurnReplayLine = serde_json::from_str(line)
.map_err(|e| format!("Turn replay JSON 解析失败: {e}\n 行: {line}"))?;
let this_seq = replay_line.replay_turn_seq;
if this_seq != *current_turn_seq && this_seq > 0 {
if *current_turn_seq > 0 {
crate::cm_turn_layout::close_open_commentary_segments(current_turn);
turns.push(std::mem::take(current_turn));
}
*current_turn_seq = this_seq;
}
map_turn_replay_event_to_turn_events(&replay_line)
}
};
let count = events.len();
for ev in events {
reducer.apply(current_turn, ev);
}
Ok(count)
}
fn detect_format_from_peekable_lines(
lines: &mut std::iter::Peekable<std::io::Lines<BufReader<File>>>,
) -> Result<JsonlFormat, String> {
let first = lines.peek().ok_or_else(|| "JSONL 文件为空".to_string())?;
let first_line = first.as_ref().map_err(|e| format!("读取首行失败: {e}"))?;
detect_format(first_line.trim())
}
fn finalize_last_replay_turn(event_count: usize, current_turn: Turn, turns: &mut Vec<Turn>) {
if event_count > 0 {
let mut current_turn = current_turn;
crate::cm_turn_layout::close_open_commentary_segments(&mut current_turn);
turns.push(current_turn);
}
}
pub fn replay_all_turns(path: &Path) -> Result<Vec<Turn>, String> {
let file = File::open(path).map_err(|e| format!("无法打开 {}: {e}", path.display()))?;
let reader = BufReader::new(file);
let mut lines = reader.lines().peekable();
let format = detect_format_from_peekable_lines(&mut lines)?;
let mut turns: Vec<Turn> = Vec::new();
let mut current_turn = Turn::default();
let mut current_turn_seq: u64 = 0;
let reducer = TurnReducer;
let mut event_count = 0usize;
for line_result in lines {
let line = line_result.map_err(|e| format!("读取行失败: {e}"))?;
let line = line.trim();
if line.is_empty() {
continue;
}
event_count += process_replay_line(
line,
format,
&mut current_turn_seq,
&mut current_turn,
&mut turns,
&reducer,
)?;
}
finalize_last_replay_turn(event_count, current_turn, &mut turns);
log::info!(
target: "crabmate-turn-layout",
"Replay ({format:?}): {event_count} TurnEvents → {} turns",
turns.len(),
);
Ok(turns)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn map_text_message_content_produces_no_event() {
let data = r#"{"type":"TEXT_MESSAGE_CONTENT","delta":"你好"}"#;
let events = map_sse_data_to_turn_events(data);
assert!(events.is_empty());
}
#[test]
fn map_turn_segment_start_to_segment_start() {
let data = r#"{"type":"CUSTOM","customType":"turn_segment_start","data":{"segmentId":"seg-before-tc1","kind":"commentary","beforeToolCallId":"tc1"}}"#;
let events = map_sse_data_to_turn_events(data);
assert_eq!(events.len(), 1);
assert_eq!(
events[0],
TurnEvent::SegmentStart {
segment_id: "seg-before-tc1".to_string(),
kind: SegmentKind::Commentary,
before_tool_call_id: Some("tc1".to_string()),
}
);
}
#[test]
fn map_timeline_log_final_response_produces_no_event() {
let data = r#"{"type":"CUSTOM","customType":"timeline_log","data":{"kind":"final_response","title":"终答","detail":"已完成创建。"}}"#;
let events = map_sse_data_to_turn_events(data);
assert!(events.is_empty());
}
#[test]
fn map_tool_call_start_to_tool_call() {
let data = r#"{"type":"TOOL_CALL_START","toolCallId":"tc-1","name":"read_file","parentMessageId":"msg-1"}"#;
let events = map_sse_data_to_turn_events(data);
assert_eq!(events.len(), 1);
assert_eq!(
events[0],
TurnEvent::ToolCall {
tool_call_id: "tc-1".to_string(),
name: "read_file".to_string(),
summary: String::new(),
}
);
}
#[test]
fn replay_sse_format() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("sse-replay-events.jsonl");
let jsonl = r#"{"seq":1,"job_id":1,"data":"{\"type\":\"CUSTOM\",\"customType\":\"turn_segment_start\",\"data\":{\"segmentId\":\"seg-1\",\"kind\":\"commentary\",\"beforeToolCallId\":null}}"}
{"seq":2,"job_id":1,"data":"{\"type\":\"CUSTOM\",\"customType\":\"turn_segment_end\",\"data\":{\"segmentId\":\"seg-1\"}}"}
{"seq":3,"job_id":1,"data":"{\"type\":\"CUSTOM\",\"customType\":\"turn_tool_phase_end\",\"data\":{\"phase\":\"tool_end\"}}"}
{"seq":4,"job_id":1,"data":"{\"type\":\"TEXT_MESSAGE_CONTENT\",\"delta\":\"完成。\"}"}
"#;
std::fs::write(&path, jsonl).expect("write jsonl");
let rows = replay_sse_events_to_web_rows(&path).expect("replay");
assert!(
rows.is_empty() || rows.iter().any(|r| r.kind == "assistant_batch_narration"),
"expected batch or empty rows, got: {rows:?}"
);
}
#[test]
fn replay_turn_replay_format() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("turn-replay-events.jsonl");
let jsonl = r#"{"seq":1,"job_id":1,"event":"llm_response_done","detail":{"llm_call_id":"llm-1","assistant_content":"","tool_calls":[{"id":"tc-1","function":{"name":"create_file"}}]},"title":"llm-1"}
{"seq":2,"job_id":1,"event":"tool_call_finished","detail":{"tool_call_id":"tc-1","name":"create_file"},"title":"create_file"}
{"seq":3,"job_id":1,"event":"llm_response_done","detail":{"llm_call_id":"llm-2","assistant_content":"完成。","tool_calls":[]},"title":"llm-2"}
"#;
std::fs::write(&path, jsonl).expect("write jsonl");
let rows = replay_sse_events_to_web_rows(&path).expect("replay");
assert!(
rows.iter().any(|r| r.kind == "tool"),
"expected tool row, got: {rows:?}"
);
}
#[test]
fn detect_format_sse_vs_turn() {
let sse_line = r#"{"seq":1,"job_id":1,"data":"{}"}"#;
let turn_line = r#"{"seq":1,"job_id":1,"event":"llm_response_done","detail":{}}"#;
assert_eq!(detect_format(sse_line).unwrap(), JsonlFormat::SseReplay);
assert_eq!(detect_format(turn_line).unwrap(), JsonlFormat::TurnReplay);
}
}