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