Skip to main content

actl_core/
history.rs

1//! Bounded, content-minimal local history. Display snapshots remain ephemeral.
2use 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
34/// Each invocation owns a file; no shared append stream or heartbeat traffic.
35pub 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    /// For latency-sensitive callers that schedule retention on a separate worker.
51    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
159/// Run the existing cross-process-throttled retention policy outside UI callbacks.
160pub 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
169/// Only this directory's regular files are eligible. Desktop recovery records are separate.
170/// Files touched in the last hour are protected against concurrent writers.
171pub 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    /// References only: archiving does not copy or open screenshots or documents.
239    #[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
270/// Bounded newest-call query. Incomplete last lines are reported, never silently accepted.
271pub 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
279/// One time-ordered stream of every recorded event, across calls and files.
280/// The whole-task view that reconstructing a stop otherwise assembles by hand.
281pub 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
424/// Reports are caller assertions, not inferred task success. Unique immutable revisions.
425pub 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}