Skip to main content

luft_service/
query.rs

1use anyhow::Result;
2use luft_core::contract::event::AgentEvent;
3use luft_core::contract::finding::Finding;
4use luft_core::state::{get_run_store, list_runs as list_run_dirs, RunCheckpoint};
5use std::path::Path;
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, guarding against unknown run dirs: returns
56/// `Ok(None)` rather than letting `get_run_store` create the run directory for
57/// an id that was never started. The single existence-checked accessor shared
58/// by `list_runs` / `get_status` and the binary `status` command.
59pub fn get_checkpoint(run_dir_name: &str, base_dir: &Path) -> Result<Option<RunCheckpoint>> {
60    if !base_dir.join(run_dir_name).exists() {
61        return Ok(None);
62    }
63    let store = get_run_store(run_dir_name, base_dir)?;
64    Ok(store.get_checkpoint())
65}
66
67pub fn list_runs(base_dir: &Path) -> Result<Vec<StatusOutput>> {
68    let run_dirs = list_run_dirs(base_dir)?;
69    let mut outputs = Vec::new();
70    for dir_name in run_dirs {
71        if let Ok(Some(checkpoint)) = get_checkpoint(&dir_name, base_dir) {
72            outputs.push(StatusOutput::from((&dir_name[..], &checkpoint)));
73        }
74    }
75    outputs.sort_by(|a, b| b.updated_at.cmp(&a.updated_at));
76    Ok(outputs)
77}
78
79pub fn get_status(run_dir_name: &str, base_dir: &Path) -> Result<Option<StatusOutput>> {
80    Ok(get_checkpoint(run_dir_name, base_dir)?
81        .as_ref()
82        .map(|cp| StatusOutput::from((run_dir_name, cp))))
83}
84
85/// Raw event log for a run (chronological). Callers own any slicing/formatting.
86pub fn get_events(run_dir_name: &str, base_dir: &Path) -> Result<Vec<AgentEvent>> {
87    let store = get_run_store(run_dir_name, base_dir)?;
88    Ok(store.get_event_log()?)
89}
90
91pub fn get_logs(run_dir_name: &str, base_dir: &Path, limit: Option<usize>) -> Result<Vec<String>> {
92    let logs: Vec<String> = get_events(run_dir_name, base_dir)?
93        .into_iter()
94        .take(limit.unwrap_or(1000))
95        .map(|e| serde_json::to_string(&e).unwrap_or_default())
96        .collect();
97    Ok(logs)
98}
99
100pub fn get_findings(run_dir_name: &str, base_dir: &Path) -> Result<Vec<Finding>> {
101    let store = get_run_store(run_dir_name, base_dir)?;
102    Ok(store.get_findings())
103}
104
105pub fn cancel_run(run_dir_name: &str, base_dir: &Path) -> Result<()> {
106    let store = get_run_store(run_dir_name, base_dir)?;
107    store.cancel()?;
108    Ok(())
109}
110
111pub enum ReportStatus {
112    Found(serde_json::Value),
113    NotFound,
114    RunFinished,
115}
116
117pub fn get_report(run_dir_name: &str, base_dir: &Path) -> Result<ReportStatus> {
118    let run_dir = base_dir.join(run_dir_name);
119    let events_path = run_dir.join("events.jsonl");
120    if !events_path.exists() {
121        if run_dir.exists() {
122            return Ok(ReportStatus::RunFinished);
123        } else {
124            return Ok(ReportStatus::NotFound);
125        }
126    }
127    let content = std::fs::read_to_string(&events_path)?;
128    for line in content.lines().rev() {
129        if let Ok(AgentEvent::RunDone { report, .. }) = serde_json::from_str::<AgentEvent>(line) {
130            return Ok(ReportStatus::Found(report));
131        }
132    }
133    Ok(ReportStatus::NotFound)
134}
135
136#[cfg(test)]
137mod tests {
138    use super::*;
139    use luft_core::contract::event::RunStatus;
140    use luft_core::contract::ids::TokenUsage;
141
142    use chrono::Utc;
143    use luft_core::contract::event::LogLevel;
144    use luft_core::contract::finding::Severity;
145    use luft_core::state::{AgentResultCache, CheckpointStatus, PhaseSummary};
146    use std::collections::HashMap;
147
148    #[test]
149    fn get_status_non_existing_run() {
150        let temp_dir = std::env::temp_dir().join("luft_svc_test_status_nonexist");
151        std::fs::create_dir_all(&temp_dir).unwrap();
152        let result = get_status("nonexistent_123", &temp_dir);
153        assert!(result.is_ok());
154        assert!(result.unwrap().is_none());
155        let _ = std::fs::remove_dir_all(&temp_dir);
156    }
157
158    #[test]
159    fn get_report_not_found_no_dir() {
160        let temp_dir = std::env::temp_dir().join("luft_svc_test_report");
161        let result = get_report("nonexistent_123", &temp_dir).unwrap();
162        assert!(matches!(result, ReportStatus::NotFound));
163    }
164
165    #[test]
166    fn get_report_run_finished_no_events() {
167        let temp_dir = std::env::temp_dir().join("luft_svc_test_report2");
168        let dir_name = "test_123";
169        let run_dir = temp_dir.join(dir_name);
170        std::fs::create_dir_all(&run_dir).unwrap();
171        let result = get_report(dir_name, &temp_dir).unwrap();
172        assert!(matches!(result, ReportStatus::RunFinished));
173        let _ = std::fs::remove_dir_all(&temp_dir);
174    }
175
176    #[test]
177    fn get_report_with_run_done() {
178        let temp_dir = std::env::temp_dir().join("luft_svc_test_report3");
179        let dir_name = "test_123";
180        let run_dir = temp_dir.join(dir_name);
181        std::fs::create_dir_all(&run_dir).unwrap();
182        let run_id = uuid::Uuid::now_v7();
183        let report = serde_json::json!({"summary": "done"});
184        let evt = AgentEvent::RunDone {
185            run_id,
186            status: RunStatus::Completed,
187            total_tokens: TokenUsage::default(),
188            report: report.clone(),
189            ts: chrono::Utc::now(),
190        };
191        std::fs::write(
192            run_dir.join("events.jsonl"),
193            serde_json::to_string(&evt).unwrap(),
194        )
195        .unwrap();
196        let result = get_report(dir_name, &temp_dir).unwrap();
197        match result {
198            ReportStatus::Found(data) => assert_eq!(data, report),
199            _ => panic!("expected Found"),
200        }
201        let _ = std::fs::remove_dir_all(&temp_dir);
202    }
203
204    #[test]
205    fn get_report_no_run_done_event() {
206        let temp_dir = std::env::temp_dir().join("luft_svc_test_report4");
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 run_id = uuid::Uuid::now_v7();
211        let evt = AgentEvent::RunStarted {
212            run_id,
213            task: "t".into(),
214            ts: Utc::now(),
215        };
216        std::fs::write(
217            run_dir.join("events.jsonl"),
218            serde_json::to_string(&evt).unwrap(),
219        )
220        .unwrap();
221        let result = get_report(dir_name, &temp_dir).unwrap();
222        assert!(matches!(result, ReportStatus::NotFound));
223        let _ = std::fs::remove_dir_all(&temp_dir);
224    }
225
226    #[test]
227    fn status_output_from_checkpoint() {
228        let run_id = uuid::Uuid::now_v7();
229        let cp = RunCheckpoint {
230            run_id,
231            task: "test task".into(),
232            status: CheckpointStatus::Running,
233            current_phase: 1,
234            completed_phases: vec![],
235            agent_results: HashMap::new(),
236            findings: vec![],
237            total_tokens: 0,
238            created_at: 1719000000,
239            updated_at: 1719000100,
240            completed_spans: vec![],
241            workflow_meta: None,
242            started_agent_ids: vec![],
243        };
244        let output = StatusOutput::from(("run_dir", &cp));
245        assert_eq!(output.run_id, run_id.to_string());
246        assert_eq!(output.run_dir, "run_dir");
247        assert_eq!(output.task, "test task");
248        assert_eq!(output.status, "running");
249        assert_eq!(output.current_phase, 1);
250        assert_eq!(output.completed_phases, 0);
251        assert_eq!(output.total_started, 0);
252        assert_eq!(output.completed_agents, 0);
253        assert_eq!(output.running_agents, 0);
254        assert_eq!(output.total_tokens, 0);
255        assert!(!output.created_at.is_empty());
256        assert!(!output.updated_at.is_empty());
257    }
258
259    #[test]
260    fn status_output_with_completed_agents() {
261        let run_id = uuid::Uuid::now_v7();
262        let agent_id = uuid::Uuid::now_v7();
263        let mut agent_results = HashMap::new();
264        agent_results.insert(
265            agent_id,
266            AgentResultCache {
267                agent_id,
268                phase_id: 1,
269                status: "ok".into(),
270                output: serde_json::json!({}),
271                findings: vec![],
272                tokens: 500,
273                completed_at: 1719000100,
274                cache_key_hash: None,
275                description: None,
276                role: None,
277            },
278        );
279        let cp = RunCheckpoint {
280            run_id,
281            task: "task".into(),
282            status: CheckpointStatus::Completed,
283            current_phase: 2,
284            completed_phases: vec![PhaseSummary {
285                phase_id: 1,
286                label: "phase 1".into(),
287                planned: 1,
288                ok: 1,
289                failed: 0,
290                description: None,
291                role: None,
292            }],
293            agent_results,
294            findings: vec![],
295            total_tokens: 500,
296            created_at: 1719000000,
297            updated_at: 1719000100,
298            completed_spans: vec![],
299            workflow_meta: None,
300            started_agent_ids: vec![agent_id],
301        };
302        let output = StatusOutput::from(("run_dir", &cp));
303        assert_eq!(output.status, "completed");
304        assert_eq!(output.current_phase, 2);
305        assert_eq!(output.completed_phases, 1);
306        assert_eq!(output.total_started, 1);
307        assert_eq!(output.completed_agents, 1);
308        assert_eq!(output.running_agents, 0);
309        assert_eq!(output.total_tokens, 500);
310    }
311
312    #[test]
313    fn get_checkpoint_existing_run() {
314        let temp_dir = std::env::temp_dir().join("luft_svc_test_checkpoint");
315        std::fs::create_dir_all(&temp_dir).unwrap();
316        let dir_name = uuid::Uuid::now_v7().to_string();
317        let run_id = uuid::Uuid::now_v7();
318        let store = get_run_store(&dir_name, &temp_dir).unwrap();
319        store.init_run(run_id, "test task").unwrap();
320        let result = get_checkpoint(&dir_name, &temp_dir).unwrap();
321        assert!(result.is_some());
322        let cp = result.unwrap();
323        assert_eq!(cp.run_id, run_id);
324        assert_eq!(cp.task, "test task");
325        let _ = std::fs::remove_dir_all(&temp_dir);
326    }
327
328    #[test]
329    fn get_status_existing_run() {
330        let temp_dir = std::env::temp_dir().join("luft_svc_test_status_exist");
331        std::fs::create_dir_all(&temp_dir).unwrap();
332        let dir_name = uuid::Uuid::now_v7().to_string();
333        let run_id = uuid::Uuid::now_v7();
334        let store = get_run_store(&dir_name, &temp_dir).unwrap();
335        store.init_run(run_id, "test task").unwrap();
336        let result = get_status(&dir_name, &temp_dir).unwrap();
337        assert!(result.is_some());
338        let status = result.unwrap();
339        assert_eq!(status.run_id, run_id.to_string());
340        assert_eq!(status.task, "test task");
341        let _ = std::fs::remove_dir_all(&temp_dir);
342    }
343
344    #[test]
345    fn list_runs_empty() {
346        let temp_dir = std::env::temp_dir().join("luft_svc_test_list_empty");
347        let _ = std::fs::remove_dir_all(&temp_dir);
348        std::fs::create_dir_all(&temp_dir).unwrap();
349        let results = list_runs(&temp_dir).unwrap();
350        assert!(results.is_empty());
351        let _ = std::fs::remove_dir_all(&temp_dir);
352    }
353
354    #[test]
355    fn list_runs_with_runs() {
356        let temp_dir = std::env::temp_dir().join("luft_svc_test_list_runs");
357        let _ = std::fs::remove_dir_all(&temp_dir);
358        std::fs::create_dir_all(&temp_dir).unwrap();
359        let dir_a = uuid::Uuid::now_v7().to_string();
360        let dir_b = uuid::Uuid::now_v7().to_string();
361        let run_id_a = uuid::Uuid::now_v7();
362        let store_a = get_run_store(&dir_a, &temp_dir).unwrap();
363        store_a.init_run(run_id_a, "task a").unwrap();
364        let run_id_b = uuid::Uuid::now_v7();
365        let store_b = get_run_store(&dir_b, &temp_dir).unwrap();
366        store_b.init_run(run_id_b, "task b").unwrap();
367        let results = list_runs(&temp_dir).unwrap();
368        assert_eq!(results.len(), 2);
369        let ids: Vec<String> = results.iter().map(|r| r.run_id.clone()).collect();
370        assert!(ids.contains(&run_id_a.to_string()));
371        assert!(ids.contains(&run_id_b.to_string()));
372        assert!(results[0].updated_at >= results[1].updated_at);
373        let _ = std::fs::remove_dir_all(&temp_dir);
374    }
375
376    #[test]
377    fn get_events_success() {
378        let temp_dir = std::env::temp_dir().join("luft_svc_test_events");
379        std::fs::create_dir_all(&temp_dir).unwrap();
380        let dir_name = uuid::Uuid::now_v7().to_string();
381        let run_id = uuid::Uuid::now_v7();
382        let store = get_run_store(&dir_name, &temp_dir).unwrap();
383        store.init_run(run_id, "task").unwrap();
384        let event = AgentEvent::Log {
385            run_id,
386            agent_id: None,
387            level: LogLevel::Info,
388            msg: "test log".into(),
389        };
390        store.append_event(&event).unwrap();
391        let events = get_events(&dir_name, &temp_dir).unwrap();
392        assert_eq!(events.len(), 1);
393        assert!(matches!(events[0], AgentEvent::Log { .. }));
394        let _ = std::fs::remove_dir_all(&temp_dir);
395    }
396
397    #[test]
398    fn get_logs_without_limit() {
399        let temp_dir = std::env::temp_dir().join("luft_svc_test_logs_nolimit");
400        std::fs::create_dir_all(&temp_dir).unwrap();
401        let dir_name = uuid::Uuid::now_v7().to_string();
402        let run_id = uuid::Uuid::now_v7();
403        let store = get_run_store(&dir_name, &temp_dir).unwrap();
404        store.init_run(run_id, "task").unwrap();
405        for i in 0..3 {
406            store
407                .append_event(&AgentEvent::Log {
408                    run_id,
409                    agent_id: None,
410                    level: LogLevel::Info,
411                    msg: format!("log {}", i),
412                })
413                .unwrap();
414        }
415        let logs = get_logs(&dir_name, &temp_dir, None).unwrap();
416        assert_eq!(logs.len(), 3);
417        let _ = std::fs::remove_dir_all(&temp_dir);
418    }
419
420    #[test]
421    fn get_logs_with_limit() {
422        let temp_dir = std::env::temp_dir().join("luft_svc_test_logs_limit");
423        std::fs::create_dir_all(&temp_dir).unwrap();
424        let dir_name = uuid::Uuid::now_v7().to_string();
425        let run_id = uuid::Uuid::now_v7();
426        let store = get_run_store(&dir_name, &temp_dir).unwrap();
427        store.init_run(run_id, "task").unwrap();
428        for i in 0..5 {
429            store
430                .append_event(&AgentEvent::Log {
431                    run_id,
432                    agent_id: None,
433                    level: LogLevel::Info,
434                    msg: format!("log {}", i),
435                })
436                .unwrap();
437        }
438        let logs = get_logs(&dir_name, &temp_dir, Some(2)).unwrap();
439        assert_eq!(logs.len(), 2);
440        let _ = std::fs::remove_dir_all(&temp_dir);
441    }
442
443    #[test]
444    fn get_findings_success() {
445        let temp_dir = std::env::temp_dir().join("luft_svc_test_findings");
446        std::fs::create_dir_all(&temp_dir).unwrap();
447        let dir_name = uuid::Uuid::now_v7().to_string();
448        let run_id = uuid::Uuid::now_v7();
449        let store = get_run_store(&dir_name, &temp_dir).unwrap();
450        store.init_run(run_id, "task").unwrap();
451        let finding = Finding {
452            kind: "test_kind".into(),
453            severity: Severity::High,
454            title: "Test Finding".into(),
455            detail: "A detailed finding description".into(),
456            location: None,
457            evidence: vec![],
458            data: serde_json::json!({}),
459        };
460        let mut cp = store.get_checkpoint().unwrap();
461        cp.findings.push(finding);
462        store.save_checkpoint(&cp).unwrap();
463        let findings = get_findings(&dir_name, &temp_dir).unwrap();
464        assert_eq!(findings.len(), 1);
465        assert_eq!(findings[0].kind, "test_kind");
466        let _ = std::fs::remove_dir_all(&temp_dir);
467    }
468
469    #[test]
470    fn cancel_run_success() {
471        let temp_dir = std::env::temp_dir().join("luft_svc_test_cancel");
472        std::fs::create_dir_all(&temp_dir).unwrap();
473        let dir_name = uuid::Uuid::now_v7().to_string();
474        let run_id = uuid::Uuid::now_v7();
475        let store = get_run_store(&dir_name, &temp_dir).unwrap();
476        store.init_run(run_id, "task").unwrap();
477        cancel_run(&dir_name, &temp_dir).unwrap();
478        let cp = get_checkpoint(&dir_name, &temp_dir).unwrap().unwrap();
479        assert_eq!(cp.status, CheckpointStatus::Cancelled);
480        let _ = std::fs::remove_dir_all(&temp_dir);
481    }
482}