use std::path::{Path, PathBuf};
use clap::Args;
use serde_json::Value;
use crate::cli::output::{OutputConfig, OutputFormat};
use crate::config;
use crate::error::{OlError, ERR_INVALID_CONFIG};
#[derive(Args, Clone, Debug)]
pub struct ListArgs {
#[arg(long)]
pub kind: Option<String>,
#[arg(long)]
pub agent: Option<String>,
#[arg(long)]
pub severity: Option<String>,
}
#[derive(Args, Clone, Debug)]
pub struct LogArgs {
#[arg(long)]
pub since: Option<String>,
#[arg(long)]
pub severity: Option<String>,
#[arg(long)]
pub kind: Option<String>,
#[arg(long)]
pub agent: Option<String>,
#[arg(short = 'n', long, default_value = "100")]
pub tail: usize,
}
#[derive(Args, Clone, Debug)]
pub struct RescanArgs {
pub path: Option<PathBuf>,
}
#[derive(Args, Clone, Debug)]
pub struct StatusArgs {}
#[derive(Args, Clone, Debug)]
pub struct InspectArgs {
pub source_id: String,
}
#[derive(Args, Clone, Debug)]
pub struct ProjectsArgs {}
#[derive(Args, Clone, Debug)]
pub struct AckArgs {
pub alert_id: Option<String>,
}
pub fn run_inventory(
cmd: &crate::cli::InventoryCommands,
output: &OutputConfig,
) -> Result<(), OlError> {
match cmd {
crate::cli::InventoryCommands::List(args) => run_list(args, output),
crate::cli::InventoryCommands::Log(args) => run_log(args, output),
crate::cli::InventoryCommands::Rescan(args) => run_rescan(args, output),
crate::cli::InventoryCommands::Status(args) => run_status(args, output),
crate::cli::InventoryCommands::Inspect(args) => run_inspect(args, output),
crate::cli::InventoryCommands::Projects(args) => run_projects(args, output),
crate::cli::InventoryCommands::Ack(args) => run_ack(args, output),
}
}
fn run_list(args: &ListArgs, output: &OutputConfig) -> Result<(), OlError> {
let openlatch_dir = config::openlatch_dir();
let log_dir = openlatch_dir.join("logs");
let events = read_recent_config_events(&log_dir, 5_000)?;
use std::collections::BTreeMap;
let mut latest: BTreeMap<(String, String), Value> = BTreeMap::new();
for ev in events {
let Some(t) = ev.get("type").and_then(|v| v.as_str()) else {
continue;
};
if !t.starts_with("ai.openlatch.config.") {
continue;
}
let kind = ev
.get("configkind")
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string();
let agent = ev
.get("configsource")
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string();
let sev = ev
.get("severityhint")
.and_then(|v| v.as_str())
.unwrap_or("low")
.to_string();
if let Some(filter) = &args.kind {
if &kind != filter {
continue;
}
}
if let Some(filter) = &args.agent {
if &agent != filter {
continue;
}
}
if let Some(filter) = &args.severity {
if &sev != filter {
continue;
}
}
let path_hash = ev
.get("configpathhash")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
latest.insert((agent.clone(), path_hash.clone()), ev);
}
if matches!(output.format, OutputFormat::Json) {
let arr: Vec<Value> = latest.into_values().collect();
let count = arr.len();
output.print_json(&serde_json::json!({"sources": arr, "count": count}));
return Ok(());
}
if output.quiet {
return Ok(());
}
if latest.is_empty() {
println!("No configuration sources observed yet.");
return Ok(());
}
println!("{:<14} {:<10} {:<10} path", "agent", "kind", "severity");
for ((agent, _), ev) in &latest {
let kind = ev.get("configkind").and_then(|v| v.as_str()).unwrap_or("?");
let sev = ev
.get("severityhint")
.and_then(|v| v.as_str())
.unwrap_or("?");
let path = ev.get("configpath").and_then(|v| v.as_str()).unwrap_or("?");
println!("{:<14} {:<10} {:<10} {}", agent, kind, sev, path);
}
Ok(())
}
fn run_log(args: &LogArgs, output: &OutputConfig) -> Result<(), OlError> {
let openlatch_dir = config::openlatch_dir();
let log_dir = openlatch_dir.join("logs");
let mut events = read_recent_config_events(&log_dir, args.tail.max(1) * 4)?;
events.retain(|ev| {
if let Some(filter) = &args.kind {
if ev.get("configkind").and_then(|v| v.as_str()) != Some(filter.as_str()) {
return false;
}
}
if let Some(filter) = &args.agent {
if ev.get("configsource").and_then(|v| v.as_str()) != Some(filter.as_str()) {
return false;
}
}
if let Some(filter) = &args.severity {
if ev.get("severityhint").and_then(|v| v.as_str()) != Some(filter.as_str()) {
return false;
}
}
true
});
let start = events.len().saturating_sub(args.tail);
let tail = &events[start..];
if matches!(output.format, OutputFormat::Json) {
output.print_json(&serde_json::json!({"events": tail, "count": tail.len()}));
return Ok(());
}
if output.quiet {
return Ok(());
}
if tail.is_empty() {
println!("No matching configuration events.");
return Ok(());
}
for ev in tail {
let time = ev.get("time").and_then(|v| v.as_str()).unwrap_or("?");
let event_type = ev.get("type").and_then(|v| v.as_str()).unwrap_or("?");
let kind = ev.get("configkind").and_then(|v| v.as_str()).unwrap_or("?");
let sev = ev
.get("severityhint")
.and_then(|v| v.as_str())
.unwrap_or("?");
let path = ev.get("configpath").and_then(|v| v.as_str()).unwrap_or("?");
println!("{} {} {:<8} {:<8} {}", time, event_type, kind, sev, path);
}
Ok(())
}
fn run_rescan(args: &RescanArgs, output: &OutputConfig) -> Result<(), OlError> {
let cfg = config::Config::load(None, None, false)?;
let token = read_daemon_token()?;
let url = format!("http://127.0.0.1:{}/admin/inventory/rescan", cfg.port);
let body = serde_json::json!({"path": args.path});
let res = reqwest::blocking::Client::new()
.post(&url)
.bearer_auth(token)
.json(&body)
.timeout(std::time::Duration::from_secs(2))
.send();
match res {
Ok(r) if r.status().as_u16() == 202 => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(&serde_json::json!({"accepted": true}));
} else if !output.quiet {
println!("Rescan accepted.");
}
Ok(())
}
Ok(r) if r.status().as_u16() == 503 => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(
&serde_json::json!({"accepted": false, "reason": "monitor_unavailable"}),
);
} else if !output.quiet {
println!("Config monitor not running. Daemon may be stopped or monitor disabled.");
}
Ok(())
}
Ok(r) => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(
&serde_json::json!({"accepted": false, "status": r.status().as_u16()}),
);
} else if !output.quiet {
println!("Rescan request failed: HTTP {}", r.status());
}
Ok(())
}
Err(e) => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(&serde_json::json!({"accepted": false, "error": e.to_string()}));
} else if !output.quiet {
println!("Daemon unreachable: {e}");
}
Ok(())
}
}
}
fn run_status(_args: &StatusArgs, output: &OutputConfig) -> Result<(), OlError> {
let cfg = config::Config::load(None, None, false)?;
let token = read_daemon_token()?;
let url = format!("http://127.0.0.1:{}/admin/inventory/status", cfg.port);
let res = reqwest::blocking::Client::new()
.get(&url)
.bearer_auth(token)
.timeout(std::time::Duration::from_secs(2))
.send();
match res {
Ok(r) if r.status().is_success() => {
let body: Value = r
.json()
.unwrap_or(serde_json::json!({"error": "unparsable response"}));
if matches!(output.format, OutputFormat::Json) {
output.print_json(&body);
} else if !output.quiet {
let enabled = body
.get("enabled")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let cache_size = body.get("cache_size").and_then(|v| v.as_u64()).unwrap_or(0);
let manifest_loaded = body
.get("manifest_loaded")
.and_then(|v| v.as_bool())
.unwrap_or(false);
println!("Configuration plane monitor");
println!(" enabled: {enabled}");
println!(" manifest loaded: {manifest_loaded}");
println!(" cache entries: {cache_size}");
}
Ok(())
}
Ok(r) => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(
&serde_json::json!({"error": "non-success", "status": r.status().as_u16()}),
);
} else if !output.quiet {
println!("Daemon responded with HTTP {}", r.status());
}
Ok(())
}
Err(e) => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(&serde_json::json!({"error": e.to_string()}));
} else if !output.quiet {
println!("Daemon unreachable: {e}");
}
Ok(())
}
}
}
fn run_inspect(args: &InspectArgs, output: &OutputConfig) -> Result<(), OlError> {
let cfg = config::Config::load(None, None, false)?;
let token = read_daemon_token()?;
let url = format!(
"http://127.0.0.1:{}/admin/inventory/inspect/{}",
cfg.port, args.source_id
);
match daemon_get(&url, &token) {
Ok(body) => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(&body);
} else if !output.quiet {
let entry = body.get("cache_entry");
let alerts_len = body
.get("alerts")
.and_then(|v| v.as_array())
.map(|a| a.len())
.unwrap_or(0);
println!("Source: {}", args.source_id);
if let Some(e) = entry {
println!(
" agent: {}",
e.get("agent").and_then(|v| v.as_str()).unwrap_or("?")
);
println!(
" kind: {}",
e.get("kind").and_then(|v| v.as_str()).unwrap_or("?")
);
println!(
" content hash: {}",
e.get("content_hash")
.and_then(|v| v.as_str())
.unwrap_or("?")
);
println!(
" path: {}",
e.get("path").and_then(|v| v.as_str()).unwrap_or("?")
);
} else {
println!(" not currently in cache");
}
println!(" pending alerts: {alerts_len}");
}
Ok(())
}
Err(e) => {
print_daemon_error(output, &e);
Ok(())
}
}
}
fn run_projects(_args: &ProjectsArgs, output: &OutputConfig) -> Result<(), OlError> {
let cfg = config::Config::load(None, None, false)?;
let token = read_daemon_token()?;
let url = format!("http://127.0.0.1:{}/admin/inventory/projects", cfg.port);
match daemon_get(&url, &token) {
Ok(body) => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(&body);
} else if !output.quiet {
let projects = body.get("projects").and_then(|v| v.as_array());
match projects {
Some(arr) if !arr.is_empty() => {
for p in arr {
if let Some(s) = p.as_str() {
println!("{s}");
}
}
}
_ => println!("No projects observed yet."),
}
}
Ok(())
}
Err(e) => {
print_daemon_error(output, &e);
Ok(())
}
}
}
fn run_ack(args: &AckArgs, output: &OutputConfig) -> Result<(), OlError> {
let cfg = config::Config::load(None, None, false)?;
let token = read_daemon_token()?;
let url = format!("http://127.0.0.1:{}/admin/inventory/ack", cfg.port);
let body = serde_json::json!({"alert_id": args.alert_id});
match daemon_post(&url, &token, &body) {
Ok(resp) => {
if matches!(output.format, OutputFormat::Json) {
output.print_json(&resp);
} else if !output.quiet {
let n = resp
.get("acknowledged")
.and_then(|v| v.as_u64())
.unwrap_or(0);
if let Some(id) = &args.alert_id {
println!("Acknowledged {n} alert(s) matching {id}");
} else {
println!("Acknowledged {n} pending alert(s).");
}
}
Ok(())
}
Err(e) => {
print_daemon_error(output, &e);
Ok(())
}
}
}
fn daemon_get(url: &str, token: &str) -> Result<Value, String> {
let resp = reqwest::blocking::Client::new()
.get(url)
.bearer_auth(token)
.timeout(std::time::Duration::from_secs(2))
.send()
.map_err(|e| e.to_string())?;
if !resp.status().is_success() {
return Err(format!("daemon returned HTTP {}", resp.status()));
}
resp.json::<Value>().map_err(|e| e.to_string())
}
fn daemon_post(url: &str, token: &str, body: &Value) -> Result<Value, String> {
let resp = reqwest::blocking::Client::new()
.post(url)
.bearer_auth(token)
.json(body)
.timeout(std::time::Duration::from_secs(2))
.send()
.map_err(|e| e.to_string())?;
if !resp.status().is_success() {
return Err(format!("daemon returned HTTP {}", resp.status()));
}
resp.json::<Value>().map_err(|e| e.to_string())
}
fn print_daemon_error(output: &OutputConfig, msg: &str) {
if matches!(output.format, OutputFormat::Json) {
output.print_json(&serde_json::json!({"error": msg}));
} else if !output.quiet {
println!("Daemon request failed: {msg}");
}
}
fn read_daemon_token() -> Result<String, OlError> {
let openlatch_dir = config::openlatch_dir();
let path = openlatch_dir.join("daemon.token");
let raw = std::fs::read_to_string(&path).map_err(|e| {
OlError::new(ERR_INVALID_CONFIG, format!("Cannot read daemon.token: {e}"))
.with_suggestion("Run 'openlatch init' to (re-)generate the daemon token.")
})?;
Ok(raw.trim().to_string())
}
fn read_recent_config_events(log_dir: &Path, cap: usize) -> Result<Vec<Value>, OlError> {
if !log_dir.exists() {
return Ok(Vec::new());
}
let mut files: Vec<PathBuf> = std::fs::read_dir(log_dir)
.map_err(|e| OlError::new(ERR_INVALID_CONFIG, format!("Cannot read log dir: {e}")))?
.flatten()
.map(|e| e.path())
.filter(|p| {
p.file_name()
.and_then(|f| f.to_str())
.map(|s| s.starts_with("events-") && s.ends_with(".jsonl"))
.unwrap_or(false)
})
.collect();
files.sort();
let mut out: Vec<Value> = Vec::new();
for path in files.iter().rev() {
let raw = match std::fs::read_to_string(path) {
Ok(s) => s,
Err(_) => continue,
};
for line in raw.lines() {
if let Ok(v) = serde_json::from_str::<Value>(line) {
if let Some(t) = v.get("type").and_then(|x| x.as_str()) {
if t.starts_with("ai.openlatch.config.") {
out.push(v);
}
}
}
}
if out.len() >= cap {
break;
}
}
Ok(out)
}