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            agent_sessions: HashMap::new(),
276            findings: vec![],
277            total_tokens: 0,
278            created_at: 1719000000,
279            updated_at: 1719000100,
280            workflow_meta: None,
281            started_agent_ids: vec![],
282        };
283        let output = StatusOutput::from(("run_dir", &cp));
284        assert_eq!(output.run_id, run_id.to_string());
285        assert_eq!(output.run_dir, "run_dir");
286        assert_eq!(output.task, "test task");
287        assert_eq!(output.status, "running");
288        assert_eq!(output.current_phase, 1);
289        assert_eq!(output.completed_phases, 0);
290        assert_eq!(output.total_started, 0);
291        assert_eq!(output.completed_agents, 0);
292        assert_eq!(output.running_agents, 0);
293        assert_eq!(output.total_tokens, 0);
294        assert!(!output.created_at.is_empty());
295        assert!(!output.updated_at.is_empty());
296    }
297
298    #[test]
299    fn status_output_with_completed_agents() {
300        let run_id = uuid::Uuid::now_v7();
301        let agent_id = uuid::Uuid::now_v7();
302        let mut agent_results = HashMap::new();
303        agent_results.insert(
304            agent_id,
305            AgentResultCache {
306                agent_id,
307                phase_id: 1,
308                status: "ok".into(),
309                output: serde_json::json!({}),
310                findings: vec![],
311                tokens: 500,
312                completed_at: 1719000100,
313                cache_key_hash: None,
314                description: None,
315                role: None,
316            },
317        );
318        let cp = RunCheckpoint {
319            run_id,
320            task: "task".into(),
321            status: CheckpointStatus::Completed,
322            current_phase: 2,
323            completed_phases: vec![PhaseSummary {
324                phase_id: 1,
325                label: "phase 1".into(),
326                planned: 1,
327                ok: 1,
328                failed: 0,
329                description: None,
330                role: None,
331            }],
332            agent_results,
333            agent_sessions: HashMap::new(),
334            findings: vec![],
335            total_tokens: 500,
336            created_at: 1719000000,
337            updated_at: 1719000100,
338            workflow_meta: None,
339            started_agent_ids: vec![agent_id],
340        };
341        let output = StatusOutput::from(("run_dir", &cp));
342        assert_eq!(output.status, "completed");
343        assert_eq!(output.current_phase, 2);
344        assert_eq!(output.completed_phases, 1);
345        assert_eq!(output.total_started, 1);
346        assert_eq!(output.completed_agents, 1);
347        assert_eq!(output.running_agents, 0);
348        assert_eq!(output.total_tokens, 500);
349    }
350
351    #[test]
352    fn get_checkpoint_existing_run() {
353        let temp_dir = std::env::temp_dir().join("luft_svc_test_checkpoint");
354        std::fs::create_dir_all(&temp_dir).unwrap();
355        let dir_name = uuid::Uuid::now_v7().to_string();
356        let run_id = uuid::Uuid::now_v7();
357        let store = get_run_store(&dir_name, &temp_dir).unwrap();
358        store.init_run(run_id, "test task").unwrap();
359        let result = get_checkpoint(&dir_name, &temp_dir).unwrap();
360        assert!(result.is_some());
361        let cp = result.unwrap();
362        assert_eq!(cp.run_id, run_id);
363        assert_eq!(cp.task, "test task");
364        let _ = std::fs::remove_dir_all(&temp_dir);
365    }
366
367    // ── cross-process regression: disk-only checkpoint (no in-memory cache) ─
368    //
369    // This is the scenario that was broken before the fix: a run created by
370    // *another* process (e.g. `luft run`) leaves a valid `checkpoint.json` on
371    // disk, but the querying process (e.g. `luft mcp serve`, or `luft list`)
372    // has never called `init_run` for it — so the global `RunStore` cache is
373    // empty. Query tools must still find it by reading from disk.
374
375    /// Build a minimal `RunCheckpoint` with one finding, for the disk-only
376    /// regression tests below.
377    fn make_checkpoint(run_id: uuid::Uuid) -> RunCheckpoint {
378        RunCheckpoint {
379            run_id,
380            task: "cross-process task".into(),
381            status: CheckpointStatus::Completed,
382            current_phase: 1,
383            completed_phases: vec![],
384            agent_results: HashMap::new(),
385            agent_sessions: HashMap::new(),
386            findings: vec![Finding {
387                kind: "disk-only".into(),
388                severity: Severity::Medium,
389                title: "T".into(),
390                detail: "D".into(),
391                location: None,
392                evidence: vec![],
393                data: serde_json::json!({}),
394            }],
395            total_tokens: 42,
396            created_at: 1719000000,
397            updated_at: 1719000100,
398            workflow_meta: None,
399            started_agent_ids: vec![],
400        }
401    }
402
403    /// Write `checkpoint.json` directly to disk, bypassing the `RunStore`
404    /// cache entirely — mirroring what a *different* process would have left
405    /// behind.
406    fn write_checkpoint_disk_only(base_dir: &Path, dir_name: &str, cp: &RunCheckpoint) {
407        let run_dir = base_dir.join(dir_name);
408        std::fs::create_dir_all(&run_dir).unwrap();
409        let json = serde_json::to_string_pretty(cp).unwrap();
410        std::fs::write(run_dir.join("checkpoint.json"), json).unwrap();
411    }
412
413    #[test]
414    fn get_checkpoint_reads_disk_without_in_memory_cache() {
415        let temp_dir = std::env::temp_dir().join("luft_svc_test_disk_only_checkpoint");
416        let _ = std::fs::remove_dir_all(&temp_dir);
417        std::fs::create_dir_all(&temp_dir).unwrap();
418        let dir_name = format!("hello_{}", uuid::Uuid::now_v7());
419        let run_id = uuid::Uuid::now_v7();
420        write_checkpoint_disk_only(&temp_dir, &dir_name, &make_checkpoint(run_id));
421
422        // No `get_run_store` / `init_run` call has happened in this process —
423        // the global in-memory cache is empty for this run. Before the fix,
424        // this returned `None` ("run not found").
425        let cp = get_checkpoint(&dir_name, &temp_dir).unwrap().expect("checkpoint");
426        assert_eq!(cp.run_id, run_id);
427        assert_eq!(cp.task, "cross-process task");
428        assert_eq!(cp.status, CheckpointStatus::Completed);
429        assert_eq!(cp.total_tokens, 42);
430        let _ = std::fs::remove_dir_all(&temp_dir);
431    }
432
433    #[test]
434    fn get_status_reads_disk_without_in_memory_cache() {
435        let temp_dir = std::env::temp_dir().join("luft_svc_test_disk_only_status");
436        let _ = std::fs::remove_dir_all(&temp_dir);
437        std::fs::create_dir_all(&temp_dir).unwrap();
438        let dir_name = format!("hello_{}", uuid::Uuid::now_v7());
439        let run_id = uuid::Uuid::now_v7();
440        write_checkpoint_disk_only(&temp_dir, &dir_name, &make_checkpoint(run_id));
441
442        let status = get_status(&dir_name, &temp_dir)
443            .unwrap()
444            .expect("status");
445        assert_eq!(status.run_id, run_id.to_string());
446        assert_eq!(status.run_dir, dir_name);
447        assert_eq!(status.status, "completed");
448        assert_eq!(status.total_tokens, 42);
449        let _ = std::fs::remove_dir_all(&temp_dir);
450    }
451
452    #[test]
453    fn list_runs_reads_disk_without_in_memory_cache() {
454        let temp_dir = std::env::temp_dir().join("luft_svc_test_disk_only_list");
455        let _ = std::fs::remove_dir_all(&temp_dir);
456        std::fs::create_dir_all(&temp_dir).unwrap();
457        let dir_a = format!("run-a_{}", uuid::Uuid::now_v7());
458        let dir_b = format!("run-b_{}", uuid::Uuid::now_v7());
459        write_checkpoint_disk_only(&temp_dir, &dir_a, &make_checkpoint(uuid::Uuid::now_v7()));
460        write_checkpoint_disk_only(&temp_dir, &dir_b, &make_checkpoint(uuid::Uuid::now_v7()));
461
462        // Before the fix, this returned an empty list because neither run had
463        // an entry in the global in-memory `RunStore` cache.
464        let results = list_runs(&temp_dir).unwrap();
465        assert_eq!(results.len(), 2, "both disk-only runs must be listed");
466        let dirs: Vec<&str> = results.iter().map(|r| r.run_dir.as_str()).collect();
467        assert!(dirs.contains(&dir_a.as_str()));
468        assert!(dirs.contains(&dir_b.as_str()));
469        let _ = std::fs::remove_dir_all(&temp_dir);
470    }
471
472    #[test]
473    fn get_findings_reads_disk_without_in_memory_cache() {
474        let temp_dir = std::env::temp_dir().join("luft_svc_test_disk_only_findings");
475        let _ = std::fs::remove_dir_all(&temp_dir);
476        std::fs::create_dir_all(&temp_dir).unwrap();
477        let dir_name = format!("hello_{}", uuid::Uuid::now_v7());
478        let run_id = uuid::Uuid::now_v7();
479        write_checkpoint_disk_only(&temp_dir, &dir_name, &make_checkpoint(run_id));
480
481        let findings = get_findings(&dir_name, &temp_dir).unwrap();
482        assert_eq!(findings.len(), 1, "finding must be read from disk");
483        assert_eq!(findings[0].kind, "disk-only");
484        let _ = std::fs::remove_dir_all(&temp_dir);
485    }
486
487    #[test]
488    fn get_checkpoint_returns_none_when_checkpoint_json_missing() {
489        // Run directory exists but has no checkpoint.json (e.g. a half-written
490        // run). Must return Ok(None), not error.
491        let temp_dir = std::env::temp_dir().join("luft_svc_test_no_checkpoint_file");
492        let _ = std::fs::remove_dir_all(&temp_dir);
493        std::fs::create_dir_all(temp_dir.join("half_baked_run")).unwrap();
494        let result = get_checkpoint("half_baked_run", &temp_dir).unwrap();
495        assert!(result.is_none());
496        let _ = std::fs::remove_dir_all(&temp_dir);
497    }
498
499    #[test]
500    fn get_status_existing_run() {
501        let temp_dir = std::env::temp_dir().join("luft_svc_test_status_exist");
502        std::fs::create_dir_all(&temp_dir).unwrap();
503        let dir_name = uuid::Uuid::now_v7().to_string();
504        let run_id = uuid::Uuid::now_v7();
505        let store = get_run_store(&dir_name, &temp_dir).unwrap();
506        store.init_run(run_id, "test task").unwrap();
507        let result = get_status(&dir_name, &temp_dir).unwrap();
508        assert!(result.is_some());
509        let status = result.unwrap();
510        assert_eq!(status.run_id, run_id.to_string());
511        assert_eq!(status.task, "test task");
512        let _ = std::fs::remove_dir_all(&temp_dir);
513    }
514
515    #[test]
516    fn list_runs_empty() {
517        let temp_dir = std::env::temp_dir().join("luft_svc_test_list_empty");
518        let _ = std::fs::remove_dir_all(&temp_dir);
519        std::fs::create_dir_all(&temp_dir).unwrap();
520        let results = list_runs(&temp_dir).unwrap();
521        assert!(results.is_empty());
522        let _ = std::fs::remove_dir_all(&temp_dir);
523    }
524
525    #[test]
526    fn list_runs_with_runs() {
527        let temp_dir = std::env::temp_dir().join("luft_svc_test_list_runs");
528        let _ = std::fs::remove_dir_all(&temp_dir);
529        std::fs::create_dir_all(&temp_dir).unwrap();
530        let dir_a = uuid::Uuid::now_v7().to_string();
531        let dir_b = uuid::Uuid::now_v7().to_string();
532        let run_id_a = uuid::Uuid::now_v7();
533        let store_a = get_run_store(&dir_a, &temp_dir).unwrap();
534        store_a.init_run(run_id_a, "task a").unwrap();
535        let run_id_b = uuid::Uuid::now_v7();
536        let store_b = get_run_store(&dir_b, &temp_dir).unwrap();
537        store_b.init_run(run_id_b, "task b").unwrap();
538        let results = list_runs(&temp_dir).unwrap();
539        assert_eq!(results.len(), 2);
540        let ids: Vec<String> = results.iter().map(|r| r.run_id.clone()).collect();
541        assert!(ids.contains(&run_id_a.to_string()));
542        assert!(ids.contains(&run_id_b.to_string()));
543        assert!(results[0].updated_at >= results[1].updated_at);
544        let _ = std::fs::remove_dir_all(&temp_dir);
545    }
546
547    #[test]
548    fn get_events_success() {
549        let temp_dir = std::env::temp_dir().join("luft_svc_test_events");
550        std::fs::create_dir_all(&temp_dir).unwrap();
551        let dir_name = uuid::Uuid::now_v7().to_string();
552        let run_id = uuid::Uuid::now_v7();
553        let store = get_run_store(&dir_name, &temp_dir).unwrap();
554        store.init_run(run_id, "task").unwrap();
555        let event = AgentEvent::Log {
556            run_id,
557            agent_id: None,
558            level: LogLevel::Info,
559            msg: "test log".into(),
560        };
561        store.append_event(&event).unwrap();
562        let events = get_events(&dir_name, &temp_dir).unwrap();
563        assert_eq!(events.len(), 1);
564        assert!(matches!(events[0], AgentEvent::Log { .. }));
565        let _ = std::fs::remove_dir_all(&temp_dir);
566    }
567
568    #[test]
569    fn get_logs_without_limit() {
570        let temp_dir = std::env::temp_dir().join("luft_svc_test_logs_nolimit");
571        std::fs::create_dir_all(&temp_dir).unwrap();
572        let dir_name = uuid::Uuid::now_v7().to_string();
573        let run_id = uuid::Uuid::now_v7();
574        let store = get_run_store(&dir_name, &temp_dir).unwrap();
575        store.init_run(run_id, "task").unwrap();
576        for i in 0..3 {
577            store
578                .append_event(&AgentEvent::Log {
579                    run_id,
580                    agent_id: None,
581                    level: LogLevel::Info,
582                    msg: format!("log {}", i),
583                })
584                .unwrap();
585        }
586        let logs = get_logs(&dir_name, &temp_dir, None).unwrap();
587        assert_eq!(logs.len(), 3);
588        let _ = std::fs::remove_dir_all(&temp_dir);
589    }
590
591    #[test]
592    fn get_logs_with_limit() {
593        let temp_dir = std::env::temp_dir().join("luft_svc_test_logs_limit");
594        std::fs::create_dir_all(&temp_dir).unwrap();
595        let dir_name = uuid::Uuid::now_v7().to_string();
596        let run_id = uuid::Uuid::now_v7();
597        let store = get_run_store(&dir_name, &temp_dir).unwrap();
598        store.init_run(run_id, "task").unwrap();
599        for i in 0..5 {
600            store
601                .append_event(&AgentEvent::Log {
602                    run_id,
603                    agent_id: None,
604                    level: LogLevel::Info,
605                    msg: format!("log {}", i),
606                })
607                .unwrap();
608        }
609        let logs = get_logs(&dir_name, &temp_dir, Some(2)).unwrap();
610        assert_eq!(logs.len(), 2);
611        let _ = std::fs::remove_dir_all(&temp_dir);
612    }
613
614    #[test]
615    fn get_findings_success() {
616        let temp_dir = std::env::temp_dir().join("luft_svc_test_findings");
617        std::fs::create_dir_all(&temp_dir).unwrap();
618        let dir_name = uuid::Uuid::now_v7().to_string();
619        let run_id = uuid::Uuid::now_v7();
620        let store = get_run_store(&dir_name, &temp_dir).unwrap();
621        store.init_run(run_id, "task").unwrap();
622        let finding = Finding {
623            kind: "test_kind".into(),
624            severity: Severity::High,
625            title: "Test Finding".into(),
626            detail: "A detailed finding description".into(),
627            location: None,
628            evidence: vec![],
629            data: serde_json::json!({}),
630        };
631        let mut cp = store.get_checkpoint().unwrap();
632        cp.findings.push(finding);
633        store.save_checkpoint(&cp).unwrap();
634        let findings = get_findings(&dir_name, &temp_dir).unwrap();
635        assert_eq!(findings.len(), 1);
636        assert_eq!(findings[0].kind, "test_kind");
637        let _ = std::fs::remove_dir_all(&temp_dir);
638    }
639
640    #[test]
641    fn cancel_run_success() {
642        let temp_dir = std::env::temp_dir().join("luft_svc_test_cancel");
643        std::fs::create_dir_all(&temp_dir).unwrap();
644        let dir_name = uuid::Uuid::now_v7().to_string();
645        let run_id = uuid::Uuid::now_v7();
646        let store = get_run_store(&dir_name, &temp_dir).unwrap();
647        store.init_run(run_id, "task").unwrap();
648        cancel_run(&dir_name, &temp_dir).unwrap();
649        let cp = get_checkpoint(&dir_name, &temp_dir).unwrap().unwrap();
650        assert_eq!(cp.status, CheckpointStatus::Cancelled);
651        let _ = std::fs::remove_dir_all(&temp_dir);
652    }
653}
654