use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread::JoinHandle;
use std::time::Duration;
use super::{Caller, arming, check, row, window};
use crate::app::cadence::{self, CADENCE_YAML};
use crate::git_tree::Agent;
use crate::opslog;
use crate::state::SnapshotCell;
use crate::ui_state::Clock;
const TAIL: usize = 4000;
pub struct SentryCtx {
pub state_root: PathBuf,
pub cell: SnapshotCell,
pub clock: Arc<dyn Clock>,
pub caller: Box<dyn Caller>,
}
impl SentryCtx {
pub fn period(&self) -> Duration {
cadence::parse(&self.settings()).full_sweep
}
pub fn pass(&self) -> bool {
let settings = self.settings();
if arming::armed(&settings).is_empty() {
return false;
}
let snapshot = crate::state::latest_snapshot(&self.cell);
let checks = row::of_entries(&opslog::tail(&self.state_root, TAIL));
for workspace in &snapshot.workspaces {
let key = crate::nav::ws_key(&workspace.path);
let Some(watch) = arming::watch(&settings, &key) else {
continue;
};
let policy =
std::fs::read_to_string(self.state_root.join(&watch.prompt)).unwrap_or_default();
if policy.trim().is_empty() {
continue;
}
let Some(tree) = snapshot.trees.get(&workspace.path) else {
continue;
};
if let Some(agent) = tree
.agents
.iter()
.find(|a| due(&checks, &workspace.path, a))
{
self.fire(&workspace.path, agent, &watch, &policy, &checks);
return true;
}
}
false
}
fn fire(
&self,
workspace: &Path,
agent: &Agent,
watch: &arming::Watch,
policy: &str,
checks: &[row::Check],
) {
let standing = row::latest(checks, &crate::nav::ws_key(workspace), &agent.agent_id);
let evidence = window::gather(
workspace,
&agent.agent_id,
standing.as_ref().map(|c| c.sha.as_str()),
&agent.tip_oid,
);
let request = check::request(&evidence, standing.map(|c| c.verdict));
let ts = self.clock.stamp();
let entry = match check::run(&*self.caller, workspace, watch, policy, &request) {
Ok(answer) => row::entry(
ts,
&row::Check {
workspace: crate::nav::ws_key(workspace),
agent: agent.agent_id.clone(),
verdict: answer.reply.verdict,
sha: agent.tip_oid.clone(),
reason: answer.reply.reason,
model: watch.model.clone(),
input_tokens: answer.input_tokens,
output_tokens: answer.output_tokens,
},
),
Err(why) => row::failure(ts, workspace, &agent.agent_id, &why),
};
let _ = opslog::append(&self.state_root, &entry);
}
fn settings(&self) -> String {
std::fs::read_to_string(self.state_root.join(CADENCE_YAML)).unwrap_or_default()
}
}
fn due(checks: &[row::Check], workspace: &Path, agent: &Agent) -> bool {
row::latest(checks, &crate::nav::ws_key(workspace), &agent.agent_id).map(|c| c.sha)
!= Some(agent.tip_oid.clone())
}
pub struct Sentry {
stop: Arc<AtomicBool>,
handle: Option<JoinHandle<()>>,
}
impl Sentry {
pub fn spawn(ctx: SentryCtx) -> Self {
let stop = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&stop);
let handle = std::thread::spawn(move || {
while !flag.load(Ordering::Relaxed) {
ctx.pass();
std::thread::park_timeout(ctx.period());
}
});
Self {
stop,
handle: Some(handle),
}
}
}
impl Drop for Sentry {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(handle) = self.handle.take() {
handle.thread().unpark();
let _ = handle.join();
}
}
}
#[cfg(test)]
mod tests;