use serde::Deserialize;
use std::io::Read;
use std::path::PathBuf;
use std::process::{Command, Stdio};
use std::thread;
use std::time::{Duration, Instant};
use super::{AicxIntent, AicxSearchResult, AicxSteerResult, OracleStatus};
pub const AICX_BINARY_ENV: &str = "LOCT_AICX_BINARY";
pub const DEFAULT_TIMEOUT: Duration = Duration::from_secs(15);
pub const AICX_TIMEOUT_ENV: &str = "LOCT_AICX_TIMEOUT_SECS";
const POLL_INTERVAL: Duration = Duration::from_millis(20);
fn invocation_timeout() -> Duration {
std::env::var(AICX_TIMEOUT_ENV)
.ok()
.and_then(|raw| raw.parse::<u64>().ok())
.filter(|n| *n > 0)
.map(Duration::from_secs)
.unwrap_or(DEFAULT_TIMEOUT)
}
fn aicx_binary() -> PathBuf {
std::env::var(AICX_BINARY_ENV)
.map(PathBuf::from)
.unwrap_or_else(|_| PathBuf::from("aicx"))
}
pub fn is_aicx_available() -> bool {
static PROBE: std::sync::OnceLock<bool> = std::sync::OnceLock::new();
*PROBE.get_or_init(is_aicx_available_uncached)
}
fn is_aicx_available_uncached() -> bool {
let bin = aicx_binary();
let mut child = match Command::new(&bin)
.arg("--version")
.stdout(Stdio::null())
.stderr(Stdio::null())
.stdin(Stdio::null())
.spawn()
{
Ok(c) => c,
Err(_) => return false,
};
let started = Instant::now();
let probe_timeout = Duration::from_secs(2);
loop {
match child.try_wait() {
Ok(Some(status)) => return status.success(),
Ok(None) => {
if started.elapsed() > probe_timeout {
let _ = child.kill();
let _ = child.wait();
return false;
}
std::thread::sleep(POLL_INTERVAL);
}
Err(_) => {
let _ = child.kill();
let _ = child.wait();
return false;
}
}
}
}
pub(super) fn debug_log(msg: impl AsRef<str>) {
if std::env::var("LOCT_DEBUG")
.ok()
.filter(|v| !v.is_empty())
.is_some()
{
eprintln!("[loctree::aicx] {}", msg.as_ref());
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum AicxRunFailure {
Timeout,
Failed,
}
pub(super) fn run_aicx(args: &[&str]) -> Option<String> {
run_aicx_outcome(args, None).ok()
}
pub(super) fn run_aicx_outcome(
args: &[&str],
timeout_cap: Option<Duration>,
) -> Result<String, AicxRunFailure> {
let bin = aicx_binary();
let mut cmd = Command::new(&bin);
cmd.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.stdin(Stdio::null());
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
cmd.process_group(0);
}
let mut child = match cmd.spawn() {
Ok(c) => c,
Err(e) => {
debug_log(format!(
"spawn failed for {} (args={:?}): {}",
bin.display(),
args,
e
));
return Err(AicxRunFailure::Failed);
}
};
let (tx_out, rx_out) = std::sync::mpsc::channel();
let _stdout_handle = child.stdout.take().map(|mut handle| {
thread::spawn(move || {
let mut buf = Vec::new();
let _ = handle.read_to_end(&mut buf);
let _ = tx_out.send(buf);
})
});
let (tx_err, rx_err) = std::sync::mpsc::channel();
let _stderr_handle = child.stderr.take().map(|mut handle| {
thread::spawn(move || {
let mut buf = Vec::new();
let _ = handle.read_to_end(&mut buf);
let _ = tx_err.send(buf);
})
});
let started = Instant::now();
let base_timeout = invocation_timeout();
let timeout = timeout_cap
.map(|cap| cap.min(base_timeout))
.unwrap_or(base_timeout);
let exit_status = loop {
match child.try_wait() {
Ok(Some(status)) => break status,
Ok(None) => {
if started.elapsed() > timeout {
#[cfg(unix)]
{
let _ = Command::new("kill")
.arg("-9")
.arg(format!("-{}", child.id()))
.status();
}
#[cfg(not(unix))]
{
let _ = child.kill();
}
let _ = child.wait();
debug_log(format!("timeout after {:?} (args={:?})", timeout, args));
return Err(AicxRunFailure::Timeout);
}
std::thread::sleep(POLL_INTERVAL);
}
Err(e) => {
#[cfg(unix)]
{
let _ = Command::new("kill")
.arg("-9")
.arg(format!("-{}", child.id()))
.status();
}
#[cfg(not(unix))]
{
let _ = child.kill();
}
let _ = child.wait();
debug_log(format!("try_wait error (args={:?}): {}", args, e));
return Err(AicxRunFailure::Failed);
}
}
};
let stdout_buf = rx_out
.recv_timeout(Duration::from_millis(100))
.unwrap_or_default();
let stderr_buf = rx_err
.recv_timeout(Duration::from_millis(100))
.unwrap_or_default();
if !exit_status.success() {
debug_log(format!(
"aicx exited with {:?} (args={:?}): stderr={}",
exit_status,
args,
String::from_utf8_lossy(&stderr_buf).trim()
));
return Err(AicxRunFailure::Failed);
}
Ok(String::from_utf8_lossy(&stdout_buf).into_owned())
}
#[derive(Debug, Deserialize)]
struct IntentWire {
kind: String,
summary: String,
project: String,
agent: String,
date: String,
#[serde(default)]
timestamp: Option<String>,
session_id: String,
source_chunk: String,
#[serde(default)]
frame_kind: Option<String>,
}
#[derive(Debug, Deserialize)]
struct IntentEnvelope {
#[serde(default)]
oracle_status: Option<OracleStatus>,
#[serde(default)]
items: Vec<IntentWire>,
}
pub(super) fn parse_intents(stdout: &str) -> Vec<AicxIntent> {
let trimmed = stdout.trim();
if trimmed.is_empty() {
return Vec::new();
}
let (wire, oracle_status): (Vec<IntentWire>, Option<OracleStatus>) =
match serde_json::from_str::<Vec<IntentWire>>(trimmed) {
Ok(v) => (v, None),
Err(array_err) => match serde_json::from_str::<IntentEnvelope>(trimmed) {
Ok(envelope) => (envelope.items, envelope.oracle_status),
Err(envelope_err) => {
debug_log(format!(
"intents JSON parse error: array={}; envelope={}",
array_err, envelope_err
));
return Vec::new();
}
},
};
wire.into_iter()
.map(|w| AicxIntent {
kind: w.kind,
text: w.summary,
agent: w.agent,
date: w.date,
timestamp: w.timestamp,
session_id: w.session_id,
project: w.project,
source_chunk_path: w.source_chunk,
frame_kind: w.frame_kind,
oracle_status: oracle_status.clone(),
})
.collect()
}
#[derive(Debug, Deserialize)]
struct SearchEnvelope {
#[serde(default)]
oracle_status: Option<OracleStatus>,
#[serde(default)]
items: Vec<SearchWire>,
}
#[derive(Debug, Deserialize)]
struct SearchWire {
#[serde(default)]
score: i64,
#[serde(default)]
label: Option<String>,
project: String,
agent: String,
date: String,
#[serde(default)]
timestamp: Option<String>,
#[serde(default)]
frame_kind: Option<String>,
#[serde(default)]
session: Option<String>,
#[serde(default)]
matches: Vec<String>,
path: String,
}
pub(super) fn parse_search(stdout: &str) -> Vec<AicxSearchResult> {
let trimmed = stdout.trim();
if trimmed.is_empty() {
return Vec::new();
}
let envelope: SearchEnvelope = match serde_json::from_str(trimmed) {
Ok(v) => v,
Err(e) => {
debug_log(format!("search JSON parse error: {}", e));
return Vec::new();
}
};
let oracle_status = envelope.oracle_status;
envelope
.items
.into_iter()
.map(|w| AicxSearchResult {
score: w.score,
label: w.label,
project: w.project,
agent: w.agent,
date: w.date,
timestamp: w.timestamp,
frame_kind: w.frame_kind,
session: w.session,
matches: w.matches,
path: w.path,
oracle_status: oracle_status.clone(),
})
.collect()
}
pub(super) fn parse_steer(stdout: &str) -> Vec<AicxSteerResult> {
let mut out = Vec::new();
let mut lines = stdout.lines().peekable();
while let Some(header) = lines.next() {
let header = header.trim();
if header.is_empty() {
continue;
}
let parts: Vec<&str> = header.split('|').map(str::trim).collect();
if parts.len() < 4 {
continue;
}
let project = parts[0].to_string();
let agent = parts[1].to_string();
let date = parts[2].to_string();
let kind = parts[3].to_string();
let meta_line = match lines.next() {
Some(l) => l.trim().to_string(),
None => break,
};
let (run_id, prompt_id, model) = parse_steer_meta(&meta_line);
let path_line = match lines.next() {
Some(l) => l.trim().to_string(),
None => break,
};
if path_line.is_empty() {
continue;
}
out.push(AicxSteerResult {
project,
agent,
date,
kind,
run_id,
prompt_id,
model,
source_chunk_path: path_line,
});
if let Some(peek) = lines.peek()
&& peek.trim().is_empty()
{
lines.next();
}
}
out
}
fn parse_steer_meta(line: &str) -> (Option<String>, Option<String>, Option<String>) {
let mut run_id = None;
let mut prompt_id = None;
let mut model = None;
for chunk in line.split(" ") {
let chunk = chunk.trim();
if let Some(rest) = chunk.strip_prefix("run_id:") {
run_id = sentinel(rest);
} else if let Some(rest) = chunk.strip_prefix("prompt_id:") {
prompt_id = sentinel(rest);
} else if let Some(rest) = chunk.strip_prefix("model:") {
model = sentinel(rest);
}
}
(run_id, prompt_id, model)
}
fn sentinel(raw: &str) -> Option<String> {
let trimmed = raw.trim();
if trimmed.is_empty() || trimmed == "-" {
None
} else {
Some(trimmed.to_string())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_intents_handles_canonical_payload() {
let payload = r#"[
{
"kind": "intent",
"summary": "ship cut5",
"evidence": [],
"project": "loctree-suite",
"agent": "claude",
"date": "2026-04-28",
"timestamp": "2026-04-28T00:00:00Z",
"session_id": "abc123",
"count": null,
"source_chunk": "/Users/x/.aicx/store/foo.md"
}
]"#;
let parsed = parse_intents(payload);
assert_eq!(parsed.len(), 1);
let intent = &parsed[0];
assert_eq!(intent.kind, "intent");
assert_eq!(intent.text, "ship cut5");
assert_eq!(intent.session_id, "abc123");
assert_eq!(intent.source_chunk_path, "/Users/x/.aicx/store/foo.md");
assert_eq!(intent.timestamp.as_deref(), Some("2026-04-28T00:00:00Z"));
assert!(intent.frame_kind.is_none());
}
#[test]
fn parse_intents_handles_oracle_envelope() {
let payload = r#"{
"oracle_status": {
"backend": "filesystem_fuzzy",
"index_kind": "none",
"fallback_reason": "fallback_filesystem_fuzzy: content index unavailable",
"store_root": "/Users/x/.aicx",
"indexed_count": 0,
"scanned_count": 1,
"candidate_count": 1,
"source_paths_verified": true,
"stale_or_unknown": true
},
"results": 1,
"items": [
{
"kind": "decision",
"summary": "keep oracle fallback explicit",
"evidence": [],
"project": "loctree-suite",
"agent": "codex",
"date": "2026-05-04",
"timestamp": "2026-05-04T12:00:00Z",
"session_id": "oracle-1",
"source_chunk": "/Users/x/.aicx/store/oracle.md",
"frame_kind": "assistant"
}
]
}"#;
let parsed = parse_intents(payload);
assert_eq!(parsed.len(), 1);
let intent = &parsed[0];
assert_eq!(intent.kind, "decision");
assert_eq!(intent.text, "keep oracle fallback explicit");
assert_eq!(intent.agent, "codex");
assert_eq!(intent.session_id, "oracle-1");
assert_eq!(intent.frame_kind.as_deref(), Some("assistant"));
}
#[test]
fn parse_intents_returns_empty_on_garbage() {
assert!(parse_intents("").is_empty());
assert!(parse_intents("not json at all { [").is_empty());
assert!(parse_intents("null").is_empty());
}
#[test]
fn parse_search_handles_envelope() {
let payload = r#"{
"results": 1,
"scanned": 100,
"items": [
{
"score": 85,
"label": "HIGH",
"project": "Loctree/loctree-suite",
"agent": "claude",
"date": "2026-04-28",
"timestamp": "2026-04-28T00:00:00Z",
"frame_kind": "tool_call",
"session": "df9d7e52",
"cwd": "/tmp",
"matches": ["snippet a", "snippet b"],
"path": "/Users/x/.aicx/foo.md"
}
]
}"#;
let parsed = parse_search(payload);
assert_eq!(parsed.len(), 1);
let row = &parsed[0];
assert_eq!(row.score, 85);
assert_eq!(row.label.as_deref(), Some("HIGH"));
assert_eq!(row.frame_kind.as_deref(), Some("tool_call"));
assert_eq!(row.matches.len(), 2);
assert_eq!(row.path, "/Users/x/.aicx/foo.md");
}
#[test]
fn parse_search_returns_empty_on_garbage() {
assert!(parse_search("").is_empty());
assert!(parse_search("garbage}").is_empty());
}
#[test]
fn parse_steer_handles_three_line_blocks() {
let payload = "Loctree/loctree-suite | claude | 2026-04-19 | conversations\n run_id: - prompt_id: - model: -\n /Users/polyversai/.aicx/store/foo.md\n\nLoctree/loctree-suite | codex | 2026-04-25 | reports\n run_id: mrbl-001 prompt_id: cut3-task model: opus-4.7\n /Users/polyversai/.aicx/store/bar.md\n";
let parsed = parse_steer(payload);
assert_eq!(parsed.len(), 2);
assert_eq!(parsed[0].agent, "claude");
assert_eq!(parsed[0].kind, "conversations");
assert!(parsed[0].run_id.is_none());
assert!(parsed[0].prompt_id.is_none());
assert!(parsed[0].model.is_none());
assert_eq!(parsed[1].agent, "codex");
assert_eq!(parsed[1].kind, "reports");
assert_eq!(parsed[1].run_id.as_deref(), Some("mrbl-001"));
assert_eq!(parsed[1].prompt_id.as_deref(), Some("cut3-task"));
assert_eq!(parsed[1].model.as_deref(), Some("opus-4.7"));
assert_eq!(
parsed[1].source_chunk_path,
"/Users/polyversai/.aicx/store/bar.md"
);
}
#[test]
fn parse_steer_returns_empty_for_blank_input() {
assert!(parse_steer("").is_empty());
assert!(parse_steer("\n\n\n").is_empty());
}
#[test]
fn parse_steer_skips_truncated_blocks() {
let payload = "Loctree/loctree-suite | claude | 2026-04-19 | conversations\n run_id: - prompt_id: - model: -\n";
let parsed = parse_steer(payload);
assert!(parsed.is_empty());
}
use crate::aicx::{OracleBackend, OracleIndexKind};
#[test]
fn parse_search_propagates_filesystem_fuzzy_oracle_status() {
let payload = r#"{
"oracle_status": {
"source_layer": "layer_1_canonical_corpus",
"backend": "filesystem_fuzzy",
"index_kind": "none",
"fallback_reason": "fallback_filesystem_fuzzy: content index unavailable",
"derived_view": "none_filesystem_scan",
"store_root": "/Users/x/.aicx",
"indexed_count": 0,
"scanned_count": 127,
"candidate_count": 1,
"source_paths_verified": true,
"stale_or_unknown": true,
"loctree_scope_safe": false,
"loctree_scope_note": "unsafe_for_scope_narrowing; use as routing evidence, then read canonical chunks"
},
"results": 1,
"scanned": 127,
"items": [
{
"score": 88,
"label": "HIGH",
"project": "Loctree/loctree-suite",
"agent": "codex",
"date": "2026-04-28",
"matches": ["snippet"],
"path": "/tmp/fuzzy.md"
}
]
}"#;
let parsed = parse_search(payload);
assert_eq!(parsed.len(), 1);
let row = &parsed[0];
let status = row
.oracle_status
.as_ref()
.expect("filesystem_fuzzy envelope must propagate oracle_status");
assert_eq!(status.backend, OracleBackend::FilesystemFuzzy);
assert_eq!(status.index_kind, OracleIndexKind::None);
assert_eq!(
status.fallback_reason.as_deref(),
Some("fallback_filesystem_fuzzy: content index unavailable"),
"fallback_reason must propagate verbatim"
);
assert!(status.stale_or_unknown);
assert!(!status.loctree_scope_safe);
assert_eq!(status.scanned_count, 127);
assert_eq!(status.candidate_count, 1);
assert_eq!(status.retrieval_mode(), "filesystem_fuzzy_fallback");
}
#[test]
fn parse_search_propagates_content_semantic_oracle_status() {
let payload = r#"{
"oracle_status": {
"source_layer": "layer_2_embedded_semantic",
"backend": "content_semantic",
"index_kind": "content_chunks",
"fallback_reason": null,
"derived_view": "embedded_semantic_top_k",
"store_root": "/Users/x/.aicx",
"indexed_count": 5000,
"scanned_count": 5000,
"candidate_count": 5,
"source_paths_verified": true,
"stale_or_unknown": false,
"loctree_scope_safe": true,
"loctree_scope_note": "safe_as_semantic_oracle"
},
"results": 1,
"scanned": 5000,
"items": [
{
"score": 95,
"label": "HIGH",
"project": "Loctree/loctree-suite",
"agent": "claude",
"date": "2026-04-30",
"matches": ["embedded match"],
"path": "/tmp/embed.md"
}
]
}"#;
let parsed = parse_search(payload);
let status = parsed[0]
.oracle_status
.as_ref()
.expect("content_semantic envelope must propagate oracle_status");
assert_eq!(status.backend, OracleBackend::ContentSemantic);
assert_eq!(status.index_kind, OracleIndexKind::ContentChunks);
assert_eq!(status.fallback_reason, None);
assert!(!status.stale_or_unknown);
assert!(status.loctree_scope_safe);
assert_eq!(status.retrieval_mode(), "embedded_semantic");
}
#[test]
fn parse_search_propagates_canonical_corpus_oracle_status() {
let payload = r#"{
"oracle_status": {
"source_layer": "layer_1_canonical_corpus",
"backend": "canonical_corpus",
"index_kind": "canonical_chunks",
"derived_view": "canonical_chunk_scan_no_semantic_index",
"store_root": "/Users/x/.aicx",
"indexed_count": 0,
"scanned_count": 200,
"candidate_count": 2,
"source_paths_verified": true,
"stale_or_unknown": false,
"loctree_scope_safe": true
},
"results": 0,
"scanned": 200,
"items": [
{
"score": 70,
"label": "MEDIUM",
"project": "x",
"agent": "claude",
"date": "2026-04-29",
"path": "/tmp/canon.md"
}
]
}"#;
let parsed = parse_search(payload);
let status = parsed[0].oracle_status.as_ref().unwrap();
assert_eq!(status.backend, OracleBackend::CanonicalCorpus);
assert_eq!(status.index_kind, OracleIndexKind::CanonicalChunks);
assert_eq!(status.retrieval_mode(), "canonical_corpus");
}
#[test]
fn parse_search_legacy_envelope_without_oracle_status() {
let payload = r#"{
"results": 1,
"scanned": 10,
"items": [
{
"score": 50,
"label": "LOW",
"project": "x",
"agent": "claude",
"date": "2026-04-01",
"path": "/tmp/legacy.md"
}
]
}"#;
let parsed = parse_search(payload);
assert_eq!(parsed.len(), 1);
assert!(
parsed[0].oracle_status.is_none(),
"legacy envelope without oracle_status must yield None, not a parse error"
);
}
#[test]
fn parse_search_unknown_backend_falls_through_to_unknown_variant() {
let payload = r#"{
"oracle_status": {
"backend": "future_quantum_oracle",
"index_kind": "none"
},
"items": [
{
"score": 1,
"label": "LOW",
"project": "x",
"agent": "claude",
"date": "2026-04-01",
"path": "/tmp/future.md"
}
]
}"#;
let parsed = parse_search(payload);
let status = parsed[0].oracle_status.as_ref().unwrap();
assert_eq!(status.backend, OracleBackend::Unknown);
assert_eq!(status.retrieval_mode(), "unknown");
}
#[test]
fn parse_intents_propagates_oracle_status_to_every_row() {
let payload = r#"{
"oracle_status": {
"backend": "filesystem_fuzzy",
"index_kind": "none",
"fallback_reason": "fallback_filesystem_fuzzy: content index unavailable",
"stale_or_unknown": true
},
"results": 2,
"items": [
{
"kind": "decision",
"summary": "first",
"project": "x",
"agent": "claude",
"date": "2026-04-29",
"session_id": "s1",
"source_chunk": "/tmp/a.md"
},
{
"kind": "intent",
"summary": "second",
"project": "x",
"agent": "codex",
"date": "2026-04-30",
"session_id": "s2",
"source_chunk": "/tmp/b.md"
}
]
}"#;
let parsed = parse_intents(payload);
assert_eq!(parsed.len(), 2);
for intent in &parsed {
let status = intent.oracle_status.as_ref().unwrap_or_else(|| {
panic!("intent {} must inherit envelope oracle_status", intent.text)
});
assert_eq!(status.backend, OracleBackend::FilesystemFuzzy);
assert_eq!(
status.fallback_reason.as_deref(),
Some("fallback_filesystem_fuzzy: content index unavailable")
);
assert!(status.stale_or_unknown);
}
}
#[test]
fn parse_intents_bare_array_yields_no_oracle_status() {
let payload = r#"[
{
"kind": "intent",
"summary": "legacy",
"project": "x",
"agent": "claude",
"date": "2026-04-29",
"session_id": "s1",
"source_chunk": "/tmp/a.md"
}
]"#;
let parsed = parse_intents(payload);
assert_eq!(parsed.len(), 1);
assert!(
parsed[0].oracle_status.is_none(),
"bare-array wire must NOT synthesise an Unknown oracle_status"
);
}
}