use chrono::Utc;
use globset::{Glob, GlobSet, GlobSetBuilder};
use serde::Deserialize;
use sha2::{Digest, Sha256};
use std::{
collections::{HashMap, HashSet},
env,
error::Error,
fs,
path::{Path, PathBuf},
process::{Command, ExitCode},
};
const AGENTS: &[&str] = &[
"Naorix",
"EchoLoom",
"SealMind",
"WhisperGen",
"VaultCrafter",
"REEMSmith",
"Scribe-001",
"Overseer-7",
];
#[allow(dead_code)]
#[derive(Deserialize, Debug)]
struct AgentJson {
agent_name: String,
role: String,
intent: String,
reem_code: String,
timestamp: String,
cia_hash: String,
#[serde(default)]
exec_status: String,
#[serde(default)]
notes: String,
}
#[derive(Deserialize)]
struct Policy {
#[serde(default)]
policy: PolicyMeta,
#[serde(default)]
rule: Vec<Rule>, #[serde(rename = "rules")]
#[allow(dead_code)]
rules_vec: Option<Vec<Rule>>, }
#[derive(Deserialize, Default)]
struct PolicyMeta {
#[serde(default)]
min_agents_ok: Option<usize>,
}
#[derive(Deserialize, Clone)]
struct Rule {
name: String,
paths: Vec<String>,
required_agents: Vec<String>,
allowed_intents: Vec<String>,
}
fn read_policy() -> Result<Policy, Box<dyn Error>> {
let s = fs::read_to_string("reflex.toml")?;
let v: Policy = toml::from_str(&s)?;
Ok(v)
}
fn list_ts_dirs(agent_dir: &Path) -> Vec<i64> {
let mut out = vec![];
if let Ok(rd) = fs::read_dir(agent_dir) {
for e in rd.flatten() {
if e.file_type().map(|t| t.is_dir()).unwrap_or(false) {
if let Ok(name) = e.file_name().into_string() {
if let Ok(ts) = name.parse::<i64>() {
out.push(ts);
}
}
}
}
}
out
}
fn mode_ts(ts_lists: &[Vec<i64>]) -> Option<i64> {
let mut freq: HashMap<i64, u32> = HashMap::new();
for v in ts_lists {
for ts in v {
*freq.entry(*ts).or_default() += 1;
}
}
freq.into_iter().max_by_key(|(_, c)| *c).map(|(ts, _)| ts)
}
fn read_json_log(agent_name: &str, ts: i64) -> Result<AgentJson, Box<dyn Error>> {
let path = format!("vault_dev_log/{agent_name}_{ts}.json");
let data = fs::read_to_string(path)?;
Ok(serde_json::from_str(&data)?)
}
fn collect_stdout(agent_name: &str, ts: i64) -> Result<String, Box<dyn Error>> {
let p = format!("vault_dev_log/{agent_name}/{ts}/stdout.log");
Ok(fs::read_to_string(p).unwrap_or_default())
}
fn compute_seal(blobs: &[u8]) -> String {
let mut h = Sha256::new();
h.update(blobs);
format!("{:x}", h.finalize())
}
fn git_staged_paths() -> Vec<String> {
let out = Command::new("git")
.args(["diff", "--cached", "--name-only"])
.output();
match out {
Ok(o) if o.status.success() => String::from_utf8_lossy(&o.stdout)
.lines()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect(),
_ => vec![],
}
}
fn match_paths(paths: &[String], globs: &[String]) -> bool {
let mut b = GlobSetBuilder::new();
for g in globs {
b.add(Glob::new(g).unwrap());
}
let set: GlobSet = b.build().unwrap();
paths.iter().any(|p| set.is_match(p))
}
fn main() -> ExitCode {
if let Err(e) = run() {
eprintln!("watchtower error: {e}");
return ExitCode::from(1);
}
ExitCode::from(0)
}
fn run() -> Result<(), Box<dyn Error>> {
let require_all = env::args().any(|a| a == "--require-all");
let forced_ts =
env::args().find_map(|a| a.strip_prefix("--ts=").and_then(|v| v.parse::<i64>().ok()));
let policy = read_policy()?;
let min_agents_ok = policy.policy.min_agents_ok.unwrap_or(AGENTS.len());
let mut per_agent_ts: HashMap<&str, Vec<i64>> = HashMap::new();
for a in AGENTS {
let dir = PathBuf::from(format!("vault_dev_log/{a}"));
per_agent_ts.insert(a, list_ts_dirs(&dir));
}
let ts = if let Some(t) = forced_ts {
t
} else {
let lists: Vec<Vec<i64>> = per_agent_ts.values().cloned().collect();
if let Some(m) = mode_ts(&lists) {
m
} else if let Some(mx) = lists.iter().flatten().max().copied() {
mx
} else {
return Err("No runs found in vault_dev_log/*/<ts>".into());
}
};
println!("🔭 Watchtower checking run ts={ts}");
let mut bad = vec![];
let mut missing = vec![];
let mut seal_input: Vec<u8> = vec![];
let mut intents: HashMap<String, String> = HashMap::new();
let mut ok_agents: HashSet<String> = HashSet::new();
for a in AGENTS {
let dir = PathBuf::from(format!("vault_dev_log/{a}/{ts}"));
if !dir.exists() {
missing.push(a.to_string());
continue;
}
match read_json_log(a, ts) {
Ok(j) => {
intents.insert(j.agent_name.clone(), j.intent.clone());
if j.exec_status != "OK" && j.exec_status != "OK_WITH_WARNINGS" {
bad.push(format!("{} status={}", a, j.exec_status));
} else {
ok_agents.insert(j.agent_name.clone());
}
let json_path = format!("vault_dev_log/{a}_{ts}.json");
if let Ok(bytes) = fs::read(&json_path) {
seal_input.extend_from_slice(&bytes);
}
}
Err(e) => bad.push(format!("{a} json error: {e}")),
}
match collect_stdout(a, ts) {
Ok(s) => seal_input.extend_from_slice(s.as_bytes()),
Err(e) => bad.push(format!("{a} stdout error: {e}")),
}
}
if ok_agents.len() < min_agents_ok {
bad.push(format!(
"only {} agents OK; policy requires {}",
ok_agents.len(),
min_agents_ok
));
}
let staged = git_staged_paths();
if !staged.is_empty() {
for r in policy.rule {
if match_paths(&staged, &r.paths) {
for ra in &r.required_agents {
if !ok_agents.contains(ra) {
bad.push(format!("rule '{}' missing required agent {}", r.name, ra));
}
if let Some(int) = intents.get(ra) {
if !r.allowed_intents.iter().any(|ai| ai == int) {
bad.push(format!(
"rule '{}' intent '{}' from {} not in {:?}",
r.name, int, ra, r.allowed_intents
));
}
} else {
bad.push(format!("rule '{}' no intent logged for {}", r.name, ra));
}
}
}
}
}
if require_all && (!missing.is_empty() || !bad.is_empty()) {
eprintln!("❌ Missing agents: {missing:?}");
eprintln!("❌ Issues: {bad:?}");
return Err("strict verification failed".into());
}
let seal_hex = compute_seal(&seal_input);
let seal_dir = Path::new(".vaultseal");
fs::create_dir_all(seal_dir)?;
let meta = serde_json::json!({
"ts": ts, "seal": seal_hex,
"agents_present": ok_agents.into_iter().collect::<Vec<_>>(),
"agents_missing": missing, "issues": bad,
"intents": intents, "generated_at": Utc::now().to_rfc3339(),
"staged": staged,
});
fs::write(
seal_dir.join(format!("VaultSeal_{ts}.json")),
serde_json::to_string_pretty(&meta)?,
)?;
fs::write(seal_dir.join("LATEST"), format!("{ts}\n"))?;
println!("✅ VaultSeal(ts={ts}) = {seal_hex}");
Ok(())
}