1use crate::{CtlError, ErrorCode};
3use serde::{Deserialize, Serialize};
4use std::fs::{self, File, OpenOptions};
5use std::io::Write;
6use std::path::{Path, PathBuf};
7
8const DAY: u64 = 86_400_000;
9const CALL_LIMIT: u64 = 1024 * 1024;
10
11pub fn valid_id(id: &str) -> bool {
12 !id.is_empty()
13 && id.len() <= 80
14 && id
15 .bytes()
16 .all(|b| b.is_ascii_alphanumeric() || b == b'-' || b == b'_')
17}
18
19pub fn task_from_env() -> Result<Option<String>, CtlError> {
20 match std::env::var("ACTL_TASK_ID") {
21 Ok(id) if valid_id(&id) => Ok(Some(id)),
22 Err(std::env::VarError::NotPresent) => Ok(None),
23 _ => Err(CtlError::protocol(
24 "ACTL_TASK_ID must contain 1-80 ASCII letters, digits, '-' or '_'",
25 )),
26 }
27}
28
29fn io_error(e: impl std::fmt::Display) -> CtlError {
30 CtlError::new(ErrorCode::Internal, format!("history: {e}"))
31}
32
33pub struct Journal {
35 file: File,
36 bytes: u64,
37 failed: bool,
38 call: String,
39 task: Option<String>,
40 seq: u64,
41 root: PathBuf,
42 dropped: u64,
43 failure_error: Option<std::io::Error>,
44}
45impl Journal {
46 pub fn open(root: &Path, call: &str, task: Option<String>) -> std::io::Result<Self> {
47 let result = Self::open_inner(root, call, task.clone());
48 if let Err(error) = &result {
49 crate::log_health::failure(root, call, task.as_deref(), "journal_open", 0, error);
50 }
51 result
52 }
53 fn open_inner(root: &Path, call: &str, task: Option<String>) -> std::io::Result<Self> {
54 if !valid_id(call) {
55 return Err(std::io::Error::other("invalid call id"));
56 }
57 let dir = root.join("history");
58 fs::create_dir_all(&dir)?;
59 cleanup(&dir, 7 * DAY, 64 * 1024 * 1024)?;
60 let file = OpenOptions::new()
61 .write(true)
62 .create_new(true)
63 .open(dir.join(format!("{call}.jsonl")))?;
64 Ok(Self {
65 file,
66 bytes: 0,
67 failed: false,
68 call: call.into(),
69 task,
70 seq: 0,
71 root: root.into(),
72 dropped: 0,
73 failure_error: None,
74 })
75 }
76
77 pub fn record(&mut self, event: &str, data: serde_json::Value) {
78 self.record_context(event, data, None);
79 }
80 pub fn record_context(
81 &mut self,
82 event: &str,
83 data: serde_json::Value,
84 context: Option<&crate::state::WorkflowProgress>,
85 ) {
86 if self.failed {
87 self.dropped += 1;
88 if let Some(error) = &self.failure_error {
89 crate::log_health::failure(
90 &self.root,
91 &self.call,
92 self.task.as_deref(),
93 "journal_write",
94 self.dropped,
95 error,
96 );
97 }
98 return;
99 }
100 self.seq += 1;
101 let mut value = serde_json::json!({"version":1,"call_id":self.call,"task_id":self.task,
102 "seq":self.seq,"ts_ms":crate::state::unix_ms(),"event":event,"data":data});
103 if let Some(context) = context {
104 value["run_id"] = serde_json::json!(context.run_id);
105 value["step_id"] = serde_json::json!(context.step_id);
106 }
107 let result = (|| -> std::io::Result<()> {
108 let mut line = serde_json::to_vec(&value)?;
109 line.push(b'\n');
110 if self.bytes + line.len() as u64 > CALL_LIMIT {
111 return Err(std::io::Error::new(
112 std::io::ErrorKind::FileTooLarge,
113 "per-call limit reached",
114 ));
115 }
116 self.file.write_all(&line)?;
117 self.file.flush()?;
118 self.bytes += line.len() as u64;
119 Ok(())
120 })();
121 if let Err(e) = result {
122 self.failed = true;
123 self.dropped += 1;
124 crate::log_health::failure(
125 &self.root,
126 &self.call,
127 self.task.as_deref(),
128 "journal_write",
129 self.dropped,
130 &e,
131 );
132 eprintln!("[actl-history] INTERNAL: record incomplete: {e}");
133 self.failure_error = Some(e);
134 }
135 }
136}
137
138pub fn cleanup(dir: &Path, age: u64, budget: u64) -> std::io::Result<()> {
141 let now = crate::state::unix_ms();
142 let mut files = Vec::new();
143 for entry in fs::read_dir(dir)? {
144 let entry = entry?;
145 if !entry.file_type()?.is_file() {
146 continue;
147 }
148 let path = entry.path();
149 if !matches!(
150 path.extension().and_then(|s| s.to_str()),
151 Some("json" | "jsonl")
152 ) {
153 continue;
154 }
155 let meta = entry.metadata()?;
156 let ts = meta
157 .modified()?
158 .duration_since(std::time::UNIX_EPOCH)
159 .unwrap_or_default()
160 .as_millis() as u64;
161 let active = path
162 .file_stem()
163 .and_then(|s| s.to_str())
164 .and_then(|id| {
165 let state = dir.parent()?.join("calls").join(format!("{id}.json"));
166 serde_json::from_slice::<crate::state::SessionState>(&fs::read(state).ok()?).ok()
167 })
168 .is_some_and(|s| {
169 s.phase == crate::state::Phase::Running && now.saturating_sub(s.ts_ms) < 10_000
170 });
171 if now.saturating_sub(ts) > age && !active {
172 fs::remove_file(path)?;
173 } else {
174 files.push((ts, meta.len(), path, active));
175 }
176 }
177 files.sort_by_key(|f| f.0);
178 let mut total: u64 = files.iter().map(|f| f.1).sum();
179 for (ts, size, path, active) in files {
180 if total <= budget {
181 break;
182 }
183 if active || now.saturating_sub(ts) < 3_600_000 {
184 continue;
185 }
186 fs::remove_file(path)?;
187 total = total.saturating_sub(size);
188 }
189 if total > budget {
190 return Err(std::io::Error::new(
191 std::io::ErrorKind::StorageFull,
192 "history capacity reached; recent records protected",
193 ));
194 }
195 Ok(())
196}
197
198#[derive(Debug, Serialize, Deserialize)]
199#[serde(deny_unknown_fields)]
200pub struct TaskReport {
201 pub task_id: String,
202 pub outcome: Outcome,
203 pub summary: String,
204 pub verification: String,
205 pub call_ids: Vec<String>,
206 pub feedback: Vec<Feedback>,
207 #[serde(default)]
209 pub evidence: Vec<String>,
210 #[serde(default, skip_serializing_if = "Vec::is_empty")]
211 pub artifacts: Vec<crate::reports::EvidenceRef>,
212}
213#[derive(Debug, Serialize, Deserialize)]
214#[serde(rename_all = "snake_case")]
215pub enum Outcome {
216 Completed,
217 Partial,
218 Stopped,
219 Failed,
220}
221#[derive(Debug, Serialize, Deserialize)]
222#[serde(deny_unknown_fields)]
223pub struct Feedback {
224 pub source: FeedbackSource,
225 pub text: String,
226 pub verified: bool,
227 #[serde(default, skip_serializing_if = "Option::is_none")]
228 pub kind: Option<crate::reports::StatementKind>,
229 #[serde(default, skip_serializing_if = "Vec::is_empty")]
230 pub evidence_ids: Vec<String>,
231}
232#[derive(Debug, Serialize, Deserialize)]
233#[serde(rename_all = "snake_case")]
234pub enum FeedbackSource {
235 Agent,
236 User,
237}
238
239pub fn read_history(
241 root: &Path,
242 task: Option<&str>,
243 limit: usize,
244) -> Result<serde_json::Value, CtlError> {
245 read_filtered(root, task, None, limit)
246}
247fn read_filtered(
248 root: &Path,
249 task: Option<&str>,
250 step: Option<&str>,
251 limit: usize,
252) -> Result<serde_json::Value, CtlError> {
253 if !(1..=100).contains(&limit) || task.is_some_and(|id| !valid_id(id)) {
254 return Err(CtlError::protocol(
255 "history requires limit 1-100 and a valid task ID",
256 ));
257 }
258 let dir = root.join("history");
259 if !dir.exists() {
260 return Ok(serde_json::json!({"calls":[],"incomplete_files":0}));
261 }
262 let mut files = Vec::new();
263 for entry in fs::read_dir(&dir).map_err(io_error)? {
264 let entry = entry.map_err(io_error)?;
265 if entry.file_type().map_err(io_error)?.is_file()
266 && entry.path().extension().is_some_and(|s| s == "jsonl")
267 {
268 files.push((
269 entry
270 .metadata()
271 .and_then(|m| m.modified())
272 .map_err(io_error)?,
273 entry.path(),
274 ));
275 }
276 }
277 files.sort_by_key(|f| std::cmp::Reverse(f.0));
278 let mut calls = Vec::new();
279 let mut incomplete = 0;
280 use std::io::{BufRead, BufReader, Read};
281 for (_, path) in files {
282 let file = File::open(path).map_err(io_error)?;
283 let mut lines = BufReader::new(file.take(CALL_LIMIT + 1)).lines();
284 let Some(first) = lines.next() else {
285 incomplete += 1;
286 continue;
287 };
288 let first = first.map_err(io_error)?;
289 let Ok(first) = serde_json::from_str::<serde_json::Value>(&first) else {
290 incomplete += 1;
291 continue;
292 };
293 if task.is_some_and(|id| first["task_id"].as_str() != Some(id)) {
294 continue;
295 }
296 let mut events = vec![first];
297 let mut valid = true;
298 for line in lines {
299 match serde_json::from_str::<serde_json::Value>(&line.map_err(io_error)?) {
300 Ok(event) => events.push(event),
301 Err(_) => {
302 valid = false;
303 break;
304 }
305 }
306 }
307 if step.is_some_and(|id| !events.iter().any(|e| e["step_id"].as_str() == Some(id))) {
308 continue;
309 }
310 let integrity = crate::log_integrity::check(&events, valid);
311 let finished = integrity == "complete";
312 if let Some(id) = step {
313 events.retain(|e| {
314 e["step_id"].as_str() == Some(id)
315 || matches!(e["event"].as_str(), Some("call_started" | "call_finished"))
316 });
317 }
318 if !valid || !finished {
319 incomplete += 1;
320 }
321 calls.push(serde_json::json!({"complete":finished,"integrity":integrity,"events":events}));
322 if calls.len() == limit {
323 break;
324 }
325 }
326 Ok(serde_json::json!({"calls":calls,"incomplete_files":incomplete}))
327}
328
329pub fn archive(root: &Path, input: &Path) -> Result<PathBuf, CtlError> {
331 let file = File::open(input).map_err(io_error)?;
332 use std::io::Read;
333 let mut bytes = Vec::new();
334 file.take(65_537)
335 .read_to_end(&mut bytes)
336 .map_err(io_error)?;
337 if bytes.len() > 65_536 {
338 return Err(CtlError::protocol("task report exceeds 64 KiB"));
339 }
340 let report: TaskReport = serde_json::from_slice(&bytes)
341 .map_err(|e| CtlError::protocol(format!("invalid task report: {e}")))?;
342 crate::reports::validate(&report)?;
343 if !valid_id(&report.task_id)
344 || report.call_ids.iter().any(|id| !valid_id(id))
345 || report.summary.trim().is_empty()
346 || report.verification.trim().is_empty()
347 {
348 return Err(CtlError::protocol(
349 "task report needs valid IDs, summary and verification",
350 ));
351 }
352 let dir = root.join("tasks");
353 fs::create_dir_all(&dir).map_err(io_error)?;
354 cleanup(&dir, 30 * DAY, 32 * 1024 * 1024).map_err(io_error)?;
355 let path = dir.join(format!(
356 "{}-{}.json",
357 report.task_id,
358 crate::snapshot::new_snapshot_id()
359 ));
360 let doc = serde_json::json!({"version":1,"archived_ms":crate::state::unix_ms(),"source":"caller","report":report});
361 let bytes = serde_json::to_vec_pretty(&doc).map_err(io_error)?;
362 let mut out = OpenOptions::new()
363 .write(true)
364 .create_new(true)
365 .open(&path)
366 .map_err(io_error)?;
367 if let Err(error) = out.write_all(&bytes).and_then(|_| out.sync_all()) {
368 drop(out);
369 let _ = fs::remove_file(&path);
370 return Err(io_error(error));
371 }
372 Ok(path)
373}
374
375pub fn query(
376 root: &Path,
377 task: Option<&str>,
378 step: Option<&str>,
379 limit: usize,
380 reports: bool,
381) -> Result<serde_json::Value, CtlError> {
382 if step.is_some_and(|s| !valid_id(s)) || ((step.is_some() || reports) && task.is_none()) {
383 return Err(CtlError::protocol(
384 "step/report queries require a task and valid step ID",
385 ));
386 }
387 let mut result = read_filtered(root, task, step, limit).map_err(|mut error| {
388 error.evidence =
389 Some(serde_json::json!({"logging_health":crate::log_health::read(root, task)}));
390 error
391 })?;
392 result["logging_health"] = crate::log_health::read(root, task);
393 if reports {
394 result["reports"] = crate::reports::read(root, task.unwrap_or_default(), limit)?;
395 }
396 Ok(result)
397}