Skip to main content

luft_core/
query.rs

1use anyhow::Result;
2use crate::contract::event::AgentEvent;
3use crate::contract::finding::Finding;
4use crate::state::{get_run_store, list_runs as list_run_dirs, RunCheckpoint};
5use std::path::{Path, PathBuf};
6
7/// Summary view of a run's checkpoint — the query DTO shared by the CLI.
8/// It lives in the query layer (not a presentation layer) so that
9/// the binary `commands` depend downward on `service`, not the reverse.
10#[derive(Debug, Clone, serde::Serialize)]
11pub struct StatusOutput {
12    pub run_id: String,
13    pub run_dir: String,
14    pub task: String,
15    pub status: String,
16    pub current_phase: u32,
17    pub completed_phases: usize,
18    pub total_started: usize,
19    pub completed_agents: usize,
20    pub running_agents: usize,
21    pub total_tokens: u64,
22    pub created_at: String,
23    pub updated_at: String,
24}
25
26impl From<(&str, &RunCheckpoint)> for StatusOutput {
27    fn from((run_dir, cp): (&str, &RunCheckpoint)) -> Self {
28        let created = chrono::DateTime::from_timestamp(cp.created_at as i64, 0)
29            .map(|dt| dt.to_rfc3339())
30            .unwrap_or_default();
31        let updated = chrono::DateTime::from_timestamp(cp.updated_at as i64, 0)
32            .map(|dt| dt.to_rfc3339())
33            .unwrap_or_default();
34
35        Self {
36            run_id: cp.run_id.to_string(),
37            run_dir: run_dir.to_string(),
38            task: cp.task.clone(),
39            status: format!("{:?}", cp.status).to_lowercase(),
40            current_phase: cp.current_phase,
41            completed_phases: cp.completed_phases.len(),
42            total_started: cp.started_agent_ids.len(),
43            completed_agents: cp.agent_results.len(),
44            running_agents: cp
45                .started_agent_ids
46                .len()
47                .saturating_sub(cp.agent_results.len()),
48            total_tokens: cp.total_tokens,
49            created_at: created,
50            updated_at: updated,
51        }
52    }
53}
54
55/// Fetch a run's checkpoint by reading `checkpoint.json` directly from disk.
56///
57/// This is a **read-only** query — it deliberately does *not* consult the
58/// in-memory `RunStore` cache, which is only populated for runs started in
59/// the current process (via `init_run` / `open_run` / `update_from_event`).
60/// Reading from disk is what makes `list_runs` / `get_run_status` work for
61/// runs created by *other* processes — e.g. the MCP server (started via
62/// `luft mcp serve`) querying runs created by `luft run`, or vice-versa.
63///
64/// Returns `Ok(None)` when the run directory or `checkpoint.json` is absent
65/// (unknown run id, or a run directory that never reached `init_run`).
66pub fn get_checkpoint(run_dir_name: &str, base_dir: &Path) -> Result<Option<RunCheckpoint>> {
67    read_checkpoint_from_disk(run_dir_name, base_dir)
68}
69
70/// Read a run's `checkpoint.json` from disk without touching the in-memory
71/// `RunStore` cache. Mirrors the disk-read pattern already used by
72/// `luft_core::journal::RunCreationMode::Auto`.
73fn read_checkpoint_from_disk(
74    run_dir_name: &str,
75    base_dir: &Path,
76) -> Result<Option<RunCheckpoint>> {
77    let run_dir = base_dir.join(run_dir_name);
78    if !run_dir.exists() {
79        return Ok(None);
80    }
81    // Surface an explicit io error when the run path exists but isn't a
82    // directory (e.g. a stray file where `.luft/runs/<name>` should be).
83    // The previous code path triggered this naturally via
84    // `RunStore::new` → `create_dir_all`; preserve that contract so
85    // callers can distinguish "no such run" (Ok(None)) from "corrupted
86    // run dir" (Err).
87    if !run_dir.is_dir() {
88        return Err(anyhow::anyhow!(
89            "Not a directory: run path '{}' exists as a non-directory file",
90            run_dir.display()
91        ));
92    }
93    let checkpoint_path: PathBuf = run_dir.join("checkpoint.json");
94    if !checkpoint_path.exists() {
95        return Ok(None);
96    }
97    let content = std::fs::read_to_string(&checkpoint_path)?;
98    let checkpoint: RunCheckpoint = serde_json::from_str(&content)
99        .map_err(|e| anyhow::anyhow!("parse checkpoint for '{run_dir_name}': {e}"))?;
100    Ok(Some(checkpoint))
101}
102
103pub fn list_runs(base_dir: &Path) -> Result<Vec<StatusOutput>> {
104    let run_dirs = list_run_dirs(base_dir)?;
105    let mut outputs = Vec::new();
106    for dir_name in run_dirs {
107        if let Ok(Some(checkpoint)) = get_checkpoint(&dir_name, base_dir) {
108            outputs.push(StatusOutput::from((&dir_name[..], &checkpoint)));
109        }
110    }
111    outputs.sort_by(|a, b| b.updated_at.cmp(&a.updated_at));
112    Ok(outputs)
113}
114
115pub fn get_status(run_dir_name: &str, base_dir: &Path) -> Result<Option<StatusOutput>> {
116    Ok(get_checkpoint(run_dir_name, base_dir)?
117        .as_ref()
118        .map(|cp| StatusOutput::from((run_dir_name, cp))))
119}
120
121/// Raw event log for a run (chronological). Callers own any slicing/formatting.
122pub fn get_events(run_dir_name: &str, base_dir: &Path) -> Result<Vec<AgentEvent>> {
123    let store = get_run_store(run_dir_name, base_dir)?;
124    Ok(store.get_event_log()?)
125}
126
127pub fn get_logs(run_dir_name: &str, base_dir: &Path, limit: Option<usize>) -> Result<Vec<String>> {
128    let logs: Vec<String> = get_events(run_dir_name, base_dir)?
129        .into_iter()
130        .take(limit.unwrap_or(1000))
131        .map(|e| serde_json::to_string(&e).unwrap_or_default())
132        .collect();
133    Ok(logs)
134}
135
136pub fn get_findings(run_dir_name: &str, base_dir: &Path) -> Result<Vec<Finding>> {
137    // Read from disk — same rationale as `get_checkpoint`: the in-memory
138    // `RunStore` cache is empty for runs created by other processes.
139    Ok(read_checkpoint_from_disk(run_dir_name, base_dir)?
140        .map(|cp| cp.findings)
141        .unwrap_or_default())
142}
143
144pub fn cancel_run(run_dir_name: &str, base_dir: &Path) -> Result<()> {
145    let store = get_run_store(run_dir_name, base_dir)?;
146    store.cancel()?;
147    Ok(())
148}
149
150pub enum ReportStatus {
151    Found(serde_json::Value),
152    NotFound,
153    RunFinished,
154}
155
156pub fn get_report(run_dir_name: &str, base_dir: &Path) -> Result<ReportStatus> {
157    let run_dir = base_dir.join(run_dir_name);
158    let events_path = run_dir.join("events.jsonl");
159    if !events_path.exists() {
160        if run_dir.exists() {
161            return Ok(ReportStatus::RunFinished);
162        } else {
163            return Ok(ReportStatus::NotFound);
164        }
165    }
166    let content = std::fs::read_to_string(&events_path)?;
167    for line in content.lines().rev() {
168        if let Ok(AgentEvent::RunDone { report, .. }) = serde_json::from_str::<AgentEvent>(line) {
169            return Ok(ReportStatus::Found(report));
170        }
171    }
172    Ok(ReportStatus::NotFound)
173}
174
175#[cfg(test)]
176mod tests {
177    use super::*;
178    use crate::contract::event::RunStatus;
179    use crate::contract::ids::TokenUsage;
180
181    use chrono::Utc;
182    use crate::contract::event::LogLevel;
183    use crate::contract::finding::Severity;
184    use crate::state::{AgentResultCache, CheckpointStatus, PhaseSummary};
185    use std::collections::HashMap;
186
187    #[test]
188    fn get_status_non_existing_run() {
189        let temp_dir = std::env::temp_dir().join("luft_svc_test_status_nonexist");
190        std::fs::create_dir_all(&temp_dir).unwrap();
191        let result = get_status("nonexistent_123", &temp_dir);
192        assert!(result.is_ok());
193        assert!(result.unwrap().is_none());
194        let _ = std::fs::remove_dir_all(&temp_dir);
195    }
196
197    #[test]
198    fn get_report_not_found_no_dir() {
199        let temp_dir = std::env::temp_dir().join("luft_svc_test_report");
200        let result = get_report("nonexistent_123", &temp_dir).unwrap();
201        assert!(matches!(result, ReportStatus::NotFound));
202    }
203
204    #[test]
205    fn get_report_run_finished_no_events() {
206        let temp_dir = std::env::temp_dir().join("luft_svc_test_report2");
207        let dir_name = "test_123";
208        let run_dir = temp_dir.join(dir_name);
209        std::fs::create_dir_all(&run_dir).unwrap();
210        let result = get_report(dir_name, &temp_dir).unwrap();
211        assert!(matches!(result, ReportStatus::RunFinished));
212        let _ = std::fs::remove_dir_all(&temp_dir);
213    }
214
215    #[test]
216    fn get_report_with_run_done() {
217        let temp_dir = std::env::temp_dir().join("luft_svc_test_report3");
218        let dir_name = "test_123";
219        let run_dir = temp_dir.join(dir_name);
220        std::fs::create_dir_all(&run_dir).unwrap();
221        let run_id = uuid::Uuid::now_v7();
222        let report = serde_json::json!({"summary": "done"});
223        let evt = AgentEvent::RunDone {
224            run_id,
225            status: RunStatus::Completed,
226            total_tokens: TokenUsage::default(),
227            report: report.clone(),
228            ts: chrono::Utc::now(),
229        };
230        std::fs::write(
231            run_dir.join("events.jsonl"),
232            serde_json::to_string(&evt).unwrap(),
233        )
234        .unwrap();
235        let result = get_report(dir_name, &temp_dir).unwrap();
236        match result {
237            ReportStatus::Found(data) => assert_eq!(data, report),
238            _ => panic!("expected Found"),
239        }
240        let _ = std::fs::remove_dir_all(&temp_dir);
241    }
242
243    #[test]
244    fn get_report_no_run_done_event() {
245        let temp_dir = std::env::temp_dir().join("luft_svc_test_report4");
246        let dir_name = "test_123";
247        let run_dir = temp_dir.join(dir_name);
248        std::fs::create_dir_all(&run_dir).unwrap();
249        let run_id = uuid::Uuid::now_v7();
250        let evt = AgentEvent::RunStarted {
251            run_id,
252            task: "t".into(),
253            ts: Utc::now(),
254        };
255        std::fs::write(
256            run_dir.join("events.jsonl"),
257            serde_json::to_string(&evt).unwrap(),
258        )
259        .unwrap();
260        let result = get_report(dir_name, &temp_dir).unwrap();
261        assert!(matches!(result, ReportStatus::NotFound));
262        let _ = std::fs::remove_dir_all(&temp_dir);
263    }
264
265    #[test]
266    fn status_output_from_checkpoint() {
267        let run_id = uuid::Uuid::now_v7();
268        let cp = RunCheckpoint {
269            run_id,
270            task: "test task".into(),
271            status: CheckpointStatus::Running,
272            current_phase: 1,
273            completed_phases: vec![],
274            agent_results: HashMap::new(),
275            findings: vec![],
276            total_tokens: 0,
277            created_at: 1719000000,
278            updated_at: 1719000100,
279            workflow_meta: None,
280            started_agent_ids: vec![],
281        };
282        let output = StatusOutput::from(("run_dir", &cp));
283        assert_eq!(output.run_id, run_id.to_string());
284        assert_eq!(output.run_dir, "run_dir");
285        assert_eq!(output.task, "test task");
286        assert_eq!(output.status, "running");
287        assert_eq!(output.current_phase, 1);
288        assert_eq!(output.completed_phases, 0);
289        assert_eq!(output.total_started, 0);
290        assert_eq!(output.completed_agents, 0);
291        assert_eq!(output.running_agents, 0);
292        assert_eq!(output.total_tokens, 0);
293        assert!(!output.created_at.is_empty());
294        assert!(!output.updated_at.is_empty());
295    }
296
297    #[test]
298    fn status_output_with_completed_agents() {
299        let run_id = uuid::Uuid::now_v7();
300        let agent_id = uuid::Uuid::now_v7();
301        let mut agent_results = HashMap::new();
302        agent_results.insert(
303            agent_id,
304            AgentResultCache {
305                agent_id,
306                phase_id: 1,
307                status: "ok".into(),
308                output: serde_json::json!({}),
309                findings: vec![],
310                tokens: 500,
311                completed_at: 1719000100,
312                cache_key_hash: None,
313                description: None,
314                role: None,
315            },
316        );
317        let cp = RunCheckpoint {
318            run_id,
319            task: "task".into(),
320            status: CheckpointStatus::Completed,
321            current_phase: 2,
322            completed_phases: vec![PhaseSummary {
323                phase_id: 1,
324                label: "phase 1".into(),
325                planned: 1,
326                ok: 1,
327                failed: 0,
328                description: None,
329                role: None,
330            }],
331            agent_results,
332            findings: vec![],
333            total_tokens: 500,
334            created_at: 1719000000,
335            updated_at: 1719000100,
336            workflow_meta: None,
337            started_agent_ids: vec![agent_id],
338        };
339        let output = StatusOutput::from(("run_dir", &cp));
340        assert_eq!(output.status, "completed");
341        assert_eq!(output.current_phase, 2);
342        assert_eq!(output.completed_phases, 1);
343        assert_eq!(output.total_started, 1);
344        assert_eq!(output.completed_agents, 1);
345        assert_eq!(output.running_agents, 0);
346        assert_eq!(output.total_tokens, 500);
347    }
348
349    #[test]
350    fn get_checkpoint_existing_run() {
351        let temp_dir = std::env::temp_dir().join("luft_svc_test_checkpoint");
352        std::fs::create_dir_all(&temp_dir).unwrap();
353        let dir_name = uuid::Uuid::now_v7().to_string();
354        let run_id = uuid::Uuid::now_v7();
355        let store = get_run_store(&dir_name, &temp_dir).unwrap();
356        store.init_run(run_id, "test task").unwrap();
357        let result = get_checkpoint(&dir_name, &temp_dir).unwrap();
358        assert!(result.is_some());
359        let cp = result.unwrap();
360        assert_eq!(cp.run_id, run_id);
361        assert_eq!(cp.task, "test task");
362        let _ = std::fs::remove_dir_all(&temp_dir);
363    }
364
365    // ── cross-process regression: disk-only checkpoint (no in-memory cache) ─
366    //
367    // This is the scenario that was broken before the fix: a run created by
368    // *another* process (e.g. `luft run`) leaves a valid `checkpoint.json` on
369    // disk, but the querying process (e.g. `luft mcp serve`, or `luft list`)
370    // has never called `init_run` for it — so the global `RunStore` cache is
371    // empty. Query tools must still find it by reading from disk.
372
373    /// Build a minimal `RunCheckpoint` with one finding, for the disk-only
374    /// regression tests below.
375    fn make_checkpoint(run_id: uuid::Uuid) -> RunCheckpoint {
376        RunCheckpoint {
377            run_id,
378            task: "cross-process task".into(),
379            status: CheckpointStatus::Completed,
380            current_phase: 1,
381            completed_phases: vec![],
382            agent_results: HashMap::new(),
383            findings: vec![Finding {
384                kind: "disk-only".into(),
385                severity: Severity::Medium,
386                title: "T".into(),
387                detail: "D".into(),
388                location: None,
389                evidence: vec![],
390                data: serde_json::json!({}),
391            }],
392            total_tokens: 42,
393            created_at: 1719000000,
394            updated_at: 1719000100,
395            workflow_meta: None,
396            started_agent_ids: vec![],
397        }
398    }
399
400    /// Write `checkpoint.json` directly to disk, bypassing the `RunStore`
401    /// cache entirely — mirroring what a *different* process would have left
402    /// behind.
403    fn write_checkpoint_disk_only(base_dir: &Path, dir_name: &str, cp: &RunCheckpoint) {
404        let run_dir = base_dir.join(dir_name);
405        std::fs::create_dir_all(&run_dir).unwrap();
406        let json = serde_json::to_string_pretty(cp).unwrap();
407        std::fs::write(run_dir.join("checkpoint.json"), json).unwrap();
408    }
409
410    #[test]
411    fn get_checkpoint_reads_disk_without_in_memory_cache() {
412        let temp_dir = std::env::temp_dir().join("luft_svc_test_disk_only_checkpoint");
413        let _ = std::fs::remove_dir_all(&temp_dir);
414        std::fs::create_dir_all(&temp_dir).unwrap();
415        let dir_name = format!("hello_{}", uuid::Uuid::now_v7());
416        let run_id = uuid::Uuid::now_v7();
417        write_checkpoint_disk_only(&temp_dir, &dir_name, &make_checkpoint(run_id));
418
419        // No `get_run_store` / `init_run` call has happened in this process —
420        // the global in-memory cache is empty for this run. Before the fix,
421        // this returned `None` ("run not found").
422        let cp = get_checkpoint(&dir_name, &temp_dir).unwrap().expect("checkpoint");
423        assert_eq!(cp.run_id, run_id);
424        assert_eq!(cp.task, "cross-process task");
425        assert_eq!(cp.status, CheckpointStatus::Completed);
426        assert_eq!(cp.total_tokens, 42);
427        let _ = std::fs::remove_dir_all(&temp_dir);
428    }
429
430    #[test]
431    fn get_status_reads_disk_without_in_memory_cache() {
432        let temp_dir = std::env::temp_dir().join("luft_svc_test_disk_only_status");
433        let _ = std::fs::remove_dir_all(&temp_dir);
434        std::fs::create_dir_all(&temp_dir).unwrap();
435        let dir_name = format!("hello_{}", uuid::Uuid::now_v7());
436        let run_id = uuid::Uuid::now_v7();
437        write_checkpoint_disk_only(&temp_dir, &dir_name, &make_checkpoint(run_id));
438
439        let status = get_status(&dir_name, &temp_dir)
440            .unwrap()
441            .expect("status");
442        assert_eq!(status.run_id, run_id.to_string());
443        assert_eq!(status.run_dir, dir_name);
444        assert_eq!(status.status, "completed");
445        assert_eq!(status.total_tokens, 42);
446        let _ = std::fs::remove_dir_all(&temp_dir);
447    }
448
449    #[test]
450    fn list_runs_reads_disk_without_in_memory_cache() {
451        let temp_dir = std::env::temp_dir().join("luft_svc_test_disk_only_list");
452        let _ = std::fs::remove_dir_all(&temp_dir);
453        std::fs::create_dir_all(&temp_dir).unwrap();
454        let dir_a = format!("run-a_{}", uuid::Uuid::now_v7());
455        let dir_b = format!("run-b_{}", uuid::Uuid::now_v7());
456        write_checkpoint_disk_only(&temp_dir, &dir_a, &make_checkpoint(uuid::Uuid::now_v7()));
457        write_checkpoint_disk_only(&temp_dir, &dir_b, &make_checkpoint(uuid::Uuid::now_v7()));
458
459        // Before the fix, this returned an empty list because neither run had
460        // an entry in the global in-memory `RunStore` cache.
461        let results = list_runs(&temp_dir).unwrap();
462        assert_eq!(results.len(), 2, "both disk-only runs must be listed");
463        let dirs: Vec<&str> = results.iter().map(|r| r.run_dir.as_str()).collect();
464        assert!(dirs.contains(&dir_a.as_str()));
465        assert!(dirs.contains(&dir_b.as_str()));
466        let _ = std::fs::remove_dir_all(&temp_dir);
467    }
468
469    #[test]
470    fn get_findings_reads_disk_without_in_memory_cache() {
471        let temp_dir = std::env::temp_dir().join("luft_svc_test_disk_only_findings");
472        let _ = std::fs::remove_dir_all(&temp_dir);
473        std::fs::create_dir_all(&temp_dir).unwrap();
474        let dir_name = format!("hello_{}", uuid::Uuid::now_v7());
475        let run_id = uuid::Uuid::now_v7();
476        write_checkpoint_disk_only(&temp_dir, &dir_name, &make_checkpoint(run_id));
477
478        let findings = get_findings(&dir_name, &temp_dir).unwrap();
479        assert_eq!(findings.len(), 1, "finding must be read from disk");
480        assert_eq!(findings[0].kind, "disk-only");
481        let _ = std::fs::remove_dir_all(&temp_dir);
482    }
483
484    #[test]
485    fn get_checkpoint_returns_none_when_checkpoint_json_missing() {
486        // Run directory exists but has no checkpoint.json (e.g. a half-written
487        // run). Must return Ok(None), not error.
488        let temp_dir = std::env::temp_dir().join("luft_svc_test_no_checkpoint_file");
489        let _ = std::fs::remove_dir_all(&temp_dir);
490        std::fs::create_dir_all(temp_dir.join("half_baked_run")).unwrap();
491        let result = get_checkpoint("half_baked_run", &temp_dir).unwrap();
492        assert!(result.is_none());
493        let _ = std::fs::remove_dir_all(&temp_dir);
494    }
495
496    #[test]
497    fn get_status_existing_run() {
498        let temp_dir = std::env::temp_dir().join("luft_svc_test_status_exist");
499        std::fs::create_dir_all(&temp_dir).unwrap();
500        let dir_name = uuid::Uuid::now_v7().to_string();
501        let run_id = uuid::Uuid::now_v7();
502        let store = get_run_store(&dir_name, &temp_dir).unwrap();
503        store.init_run(run_id, "test task").unwrap();
504        let result = get_status(&dir_name, &temp_dir).unwrap();
505        assert!(result.is_some());
506        let status = result.unwrap();
507        assert_eq!(status.run_id, run_id.to_string());
508        assert_eq!(status.task, "test task");
509        let _ = std::fs::remove_dir_all(&temp_dir);
510    }
511
512    #[test]
513    fn list_runs_empty() {
514        let temp_dir = std::env::temp_dir().join("luft_svc_test_list_empty");
515        let _ = std::fs::remove_dir_all(&temp_dir);
516        std::fs::create_dir_all(&temp_dir).unwrap();
517        let results = list_runs(&temp_dir).unwrap();
518        assert!(results.is_empty());
519        let _ = std::fs::remove_dir_all(&temp_dir);
520    }
521
522    #[test]
523    fn list_runs_with_runs() {
524        let temp_dir = std::env::temp_dir().join("luft_svc_test_list_runs");
525        let _ = std::fs::remove_dir_all(&temp_dir);
526        std::fs::create_dir_all(&temp_dir).unwrap();
527        let dir_a = uuid::Uuid::now_v7().to_string();
528        let dir_b = uuid::Uuid::now_v7().to_string();
529        let run_id_a = uuid::Uuid::now_v7();
530        let store_a = get_run_store(&dir_a, &temp_dir).unwrap();
531        store_a.init_run(run_id_a, "task a").unwrap();
532        let run_id_b = uuid::Uuid::now_v7();
533        let store_b = get_run_store(&dir_b, &temp_dir).unwrap();
534        store_b.init_run(run_id_b, "task b").unwrap();
535        let results = list_runs(&temp_dir).unwrap();
536        assert_eq!(results.len(), 2);
537        let ids: Vec<String> = results.iter().map(|r| r.run_id.clone()).collect();
538        assert!(ids.contains(&run_id_a.to_string()));
539        assert!(ids.contains(&run_id_b.to_string()));
540        assert!(results[0].updated_at >= results[1].updated_at);
541        let _ = std::fs::remove_dir_all(&temp_dir);
542    }
543
544    #[test]
545    fn get_events_success() {
546        let temp_dir = std::env::temp_dir().join("luft_svc_test_events");
547        std::fs::create_dir_all(&temp_dir).unwrap();
548        let dir_name = uuid::Uuid::now_v7().to_string();
549        let run_id = uuid::Uuid::now_v7();
550        let store = get_run_store(&dir_name, &temp_dir).unwrap();
551        store.init_run(run_id, "task").unwrap();
552        let event = AgentEvent::Log {
553            run_id,
554            agent_id: None,
555            level: LogLevel::Info,
556            msg: "test log".into(),
557        };
558        store.append_event(&event).unwrap();
559        let events = get_events(&dir_name, &temp_dir).unwrap();
560        assert_eq!(events.len(), 1);
561        assert!(matches!(events[0], AgentEvent::Log { .. }));
562        let _ = std::fs::remove_dir_all(&temp_dir);
563    }
564
565    #[test]
566    fn get_logs_without_limit() {
567        let temp_dir = std::env::temp_dir().join("luft_svc_test_logs_nolimit");
568        std::fs::create_dir_all(&temp_dir).unwrap();
569        let dir_name = uuid::Uuid::now_v7().to_string();
570        let run_id = uuid::Uuid::now_v7();
571        let store = get_run_store(&dir_name, &temp_dir).unwrap();
572        store.init_run(run_id, "task").unwrap();
573        for i in 0..3 {
574            store
575                .append_event(&AgentEvent::Log {
576                    run_id,
577                    agent_id: None,
578                    level: LogLevel::Info,
579                    msg: format!("log {}", i),
580                })
581                .unwrap();
582        }
583        let logs = get_logs(&dir_name, &temp_dir, None).unwrap();
584        assert_eq!(logs.len(), 3);
585        let _ = std::fs::remove_dir_all(&temp_dir);
586    }
587
588    #[test]
589    fn get_logs_with_limit() {
590        let temp_dir = std::env::temp_dir().join("luft_svc_test_logs_limit");
591        std::fs::create_dir_all(&temp_dir).unwrap();
592        let dir_name = uuid::Uuid::now_v7().to_string();
593        let run_id = uuid::Uuid::now_v7();
594        let store = get_run_store(&dir_name, &temp_dir).unwrap();
595        store.init_run(run_id, "task").unwrap();
596        for i in 0..5 {
597            store
598                .append_event(&AgentEvent::Log {
599                    run_id,
600                    agent_id: None,
601                    level: LogLevel::Info,
602                    msg: format!("log {}", i),
603                })
604                .unwrap();
605        }
606        let logs = get_logs(&dir_name, &temp_dir, Some(2)).unwrap();
607        assert_eq!(logs.len(), 2);
608        let _ = std::fs::remove_dir_all(&temp_dir);
609    }
610
611    #[test]
612    fn get_findings_success() {
613        let temp_dir = std::env::temp_dir().join("luft_svc_test_findings");
614        std::fs::create_dir_all(&temp_dir).unwrap();
615        let dir_name = uuid::Uuid::now_v7().to_string();
616        let run_id = uuid::Uuid::now_v7();
617        let store = get_run_store(&dir_name, &temp_dir).unwrap();
618        store.init_run(run_id, "task").unwrap();
619        let finding = Finding {
620            kind: "test_kind".into(),
621            severity: Severity::High,
622            title: "Test Finding".into(),
623            detail: "A detailed finding description".into(),
624            location: None,
625            evidence: vec![],
626            data: serde_json::json!({}),
627        };
628        let mut cp = store.get_checkpoint().unwrap();
629        cp.findings.push(finding);
630        store.save_checkpoint(&cp).unwrap();
631        let findings = get_findings(&dir_name, &temp_dir).unwrap();
632        assert_eq!(findings.len(), 1);
633        assert_eq!(findings[0].kind, "test_kind");
634        let _ = std::fs::remove_dir_all(&temp_dir);
635    }
636
637    #[test]
638    fn cancel_run_success() {
639        let temp_dir = std::env::temp_dir().join("luft_svc_test_cancel");
640        std::fs::create_dir_all(&temp_dir).unwrap();
641        let dir_name = uuid::Uuid::now_v7().to_string();
642        let run_id = uuid::Uuid::now_v7();
643        let store = get_run_store(&dir_name, &temp_dir).unwrap();
644        store.init_run(run_id, "task").unwrap();
645        cancel_run(&dir_name, &temp_dir).unwrap();
646        let cp = get_checkpoint(&dir_name, &temp_dir).unwrap().unwrap();
647        assert_eq!(cp.status, CheckpointStatus::Cancelled);
648        let _ = std::fs::remove_dir_all(&temp_dir);
649    }
650}
651