1use std::fs::{self, File, OpenOptions};
2use std::io::{BufRead, BufReader, Write};
3use std::path::{Path, PathBuf};
4
5use kimetsu_core::KimetsuResult;
6use kimetsu_core::event::Event;
7use kimetsu_core::ids::RunId;
8use kimetsu_core::paths::ProjectPaths;
9
10#[derive(Debug, Clone)]
11pub struct RunPaths {
12 pub run_dir: PathBuf,
13 pub trace_jsonl: PathBuf,
14 pub artifacts_dir: PathBuf,
15 pub patch_plans_dir: PathBuf,
16 pub final_report: PathBuf,
17 pub run_log: PathBuf,
18}
19
20impl RunPaths {
21 pub fn new(paths: &ProjectPaths, run_id: RunId) -> Self {
22 let run_dir = paths.runs_dir.join(run_id.to_string());
23 Self {
24 trace_jsonl: run_dir.join("trace.jsonl"),
25 artifacts_dir: run_dir.join("artifacts"),
26 patch_plans_dir: run_dir.join("patch_plans"),
27 final_report: run_dir.join("final_report.md"),
28 run_log: run_dir.join("kimetsu.log"),
29 run_dir,
30 }
31 }
32
33 pub fn create_dirs(&self) -> KimetsuResult<()> {
34 fs::create_dir_all(&self.artifacts_dir)?;
35 fs::create_dir_all(&self.patch_plans_dir)?;
36 Ok(())
37 }
38}
39
40pub struct TraceWriter {
41 file: File,
42}
43
44impl TraceWriter {
45 pub fn create(paths: &ProjectPaths, run_id: RunId) -> KimetsuResult<(Self, RunPaths)> {
46 let run_paths = RunPaths::new(paths, run_id);
47 run_paths.create_dirs()?;
48 let file = OpenOptions::new()
49 .create(true)
50 .append(true)
51 .open(&run_paths.trace_jsonl)?;
52 Ok((Self { file }, run_paths))
53 }
54
55 pub fn append(&mut self, event: &Event, fsync: bool) -> KimetsuResult<()> {
56 serde_json::to_writer(&mut self.file, event)?;
57 self.file.write_all(b"\n")?;
58 self.file.flush()?;
59 if fsync {
60 self.file.sync_data()?;
61 }
62 Ok(())
63 }
64}
65
66pub fn read_trace(trace_jsonl: &Path) -> KimetsuResult<Vec<Event>> {
67 let file = File::open(trace_jsonl)?;
68 let mut reader = BufReader::new(file);
69 let mut events = Vec::new();
70 let mut line = String::new();
71 let mut line_number = 0usize;
72
73 loop {
74 line.clear();
75 let bytes_read = reader.read_line(&mut line)?;
76 if bytes_read == 0 {
77 break;
78 }
79 line_number += 1;
80
81 let trimmed = line.trim();
82 if trimmed.is_empty() {
83 continue;
84 }
85
86 match serde_json::from_str::<Event>(trimmed) {
87 Ok(event) => events.push(event),
88 Err(err) => {
89 if !line.ends_with('\n') {
90 eprintln!(
91 "warning: ignoring invalid trailing JSONL line in {}: {err}",
92 trace_jsonl.display()
93 );
94 break;
95 }
96
97 return Err(format!(
98 "invalid JSONL at {}:{}: {err}",
99 trace_jsonl.display(),
100 line_number
101 )
102 .into());
103 }
104 }
105 }
106
107 Ok(events)
108}
109
110pub fn discover_traces(paths: &ProjectPaths) -> KimetsuResult<Vec<PathBuf>> {
111 if !paths.runs_dir.exists() {
112 return Ok(Vec::new());
113 }
114
115 let mut traces = Vec::new();
116 for entry in fs::read_dir(&paths.runs_dir)? {
117 let entry = entry?;
118 if !entry.file_type()?.is_dir() {
119 continue;
120 }
121
122 let trace = entry.path().join("trace.jsonl");
123 if trace.exists() {
124 traces.push(trace);
125 }
126 }
127
128 traces.sort();
129 Ok(traces)
130}
131
132pub fn read_all_traces(paths: &ProjectPaths) -> KimetsuResult<Vec<Event>> {
133 let mut events = Vec::new();
134 for trace in discover_traces(paths)? {
135 events.extend(read_trace(&trace)?);
136 }
137
138 events.sort_by(|left, right| {
139 left.event_id
140 .0
141 .cmp(&right.event_id.0)
142 .then_with(|| left.ts.cmp(&right.ts))
143 });
144 events.dedup_by_key(|event| event.event_id);
145 Ok(events)
146}