use std::collections::{HashMap, HashSet, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Instant;
use serde_json::Value;
use sha2::{Digest, Sha256};
use tokio::sync::{mpsc, watch};
use unicode_normalization::UnicodeNormalization;
use super::cache::{CacheEntry, ContentHashCache};
use super::claude_plugins;
use super::manifest::{expand_path, AgentPath, ConfigScope, Manifest, WatchStrategy};
use super::watcher::{glob_expand, is_excluded, watch_plan, WatchControl};
const MAX_PROJECT_ROOTS: usize = 64;
#[derive(Default)]
struct ProjectRoots(VecDeque<(PathBuf, PathBuf)>);
impl ProjectRoots {
fn register(&mut self, root: &Path) -> bool {
let key = crate::path_compat::dedup_key(root);
if let Some(i) = self.0.iter().position(|(k, _)| *k == key) {
if let Some(known) = self.0.remove(i) {
self.0.push_back(known);
}
return false;
}
self.0.push_back((key, root.to_path_buf()));
if self.0.len() > MAX_PROJECT_ROOTS {
self.0.pop_front();
}
true
}
fn paths(&self) -> Vec<PathBuf> {
self.0.iter().map(|(_, p)| p.clone()).collect()
}
}
use crate::cloud::CloudEvent;
use crate::config::{Config, ContentForwardMode};
use crate::core::logging::EventLogger;
use crate::privacy::{filter_event_with, PrivacyFilter};
#[derive(Debug, Clone)]
pub enum ConfigChangeRequest {
FsAdded(Vec<PathBuf>),
FsModified(Vec<PathBuf>),
FsRemoved(Vec<PathBuf>),
InitialInventory,
PeriodicRescan,
ManualRescan {
path_filter: Option<PathBuf>,
},
ProjectScopeRegister {
project_root: PathBuf,
},
DetectAgents,
Shutdown,
}
const DETECTION_INTERVAL: std::time::Duration = std::time::Duration::from_secs(60);
pub(crate) type AgentDetector = fn() -> HashSet<String>;
pub(crate) fn detected_agents() -> HashSet<String> {
crate::hooks::detect_agents()
.iter()
.map(|a| a.agent_type().to_string())
.collect()
}
#[derive(Debug, Clone, Copy)]
pub enum Severity {
Critical,
High,
Medium,
Low,
Info,
}
impl Severity {
pub fn as_str(&self) -> &'static str {
match self {
Severity::Critical => "critical",
Severity::High => "high",
Severity::Medium => "medium",
Severity::Low => "low",
Severity::Info => "info",
}
}
}
#[derive(Debug, Clone, Copy)]
pub enum EventSource {
InitScan,
Rescan,
FsWatcher,
}
impl EventSource {
pub fn as_str(&self) -> &'static str {
match self {
EventSource::InitScan => "init_scan",
EventSource::Rescan => "rescan",
EventSource::FsWatcher => "fs_watcher",
}
}
}
#[derive(Debug, Clone, Copy)]
pub enum ChangeKind {
Snapshot,
Added,
Modified,
Removed,
}
impl ChangeKind {
fn type_suffix(self) -> &'static str {
match self {
ChangeKind::Snapshot => "snapshot",
ChangeKind::Added => "added",
ChangeKind::Modified => "modified",
ChangeKind::Removed => "removed",
}
}
fn severity_key(self) -> &'static str {
match self {
ChangeKind::Snapshot | ChangeKind::Added => "added",
ChangeKind::Modified => "modified",
ChangeKind::Removed => "removed",
}
}
}
#[allow(clippy::too_many_arguments)]
pub async fn run(
manifest: Arc<Manifest>,
cache: Arc<ContentHashCache>,
privacy_filter: PrivacyFilter,
cloud_tx: Option<mpsc::Sender<CloudEvent>>,
event_logger: EventLogger,
config: Arc<Config>,
detector: Option<AgentDetector>,
watch: Option<WatchControl>,
mut request_rx: mpsc::Receiver<ConfigChangeRequest>,
request_tx: mpsc::Sender<ConfigChangeRequest>,
mut shutdown: watch::Receiver<bool>,
) {
let mut diff_counters: HashMap<(String, String), u64> = HashMap::new();
let mut projects = ProjectRoots::default();
let detect = || match detector {
Some(detector) => detector(),
None => manifest.agents.iter().map(|a| a.name.clone()).collect(),
};
let mut detected = detect();
let mut active = manifest.only_agents(&detected);
if detector.is_some() {
let detect_tx = request_tx.clone();
let mut detect_shutdown = shutdown.clone();
tokio::spawn(async move {
let mut interval = tokio::time::interval(DETECTION_INTERVAL);
interval.tick().await;
loop {
tokio::select! {
_ = interval.tick() => {
if detect_tx.send(ConfigChangeRequest::DetectAgents).await.is_err() {
break;
}
}
_ = wait_for_shutdown(&mut detect_shutdown) => break,
}
}
});
}
let rescan_tx = request_tx.clone();
let rescan_interval = std::time::Duration::from_secs(
config
.inventory_monitor
.periodic_rescan_interval_hours
.saturating_mul(3600),
);
if !rescan_interval.is_zero() {
let mut rescan_shutdown = shutdown.clone();
tokio::spawn(async move {
let mut interval = tokio::time::interval(rescan_interval);
interval.tick().await;
loop {
tokio::select! {
_ = interval.tick() => {
if rescan_tx
.send(ConfigChangeRequest::PeriodicRescan)
.await
.is_err()
{
break;
}
}
_ = wait_for_shutdown(&mut rescan_shutdown) => break,
}
}
});
}
loop {
let req = tokio::select! {
request = request_rx.recv() => request,
_ = wait_for_shutdown(&mut shutdown) => break,
};
let Some(req) = req else {
break;
};
match req {
ConfigChangeRequest::Shutdown => break,
ConfigChangeRequest::InitialInventory => {
let started = Instant::now();
crate::telemetry::capture_global(
crate::telemetry::Event::config_initial_scan_started("claude-code"),
);
let emitted = run_full_walk(
&active,
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
EventSource::InitScan,
None,
)
.await;
let duration_ms = started.elapsed().as_millis() as u64;
crate::telemetry::capture_global(
crate::telemetry::Event::config_initial_scan_completed(
"claude-code",
emitted,
duration_ms,
),
);
tracing::info!(
items_emitted = emitted,
duration_ms,
"config_monitor: initial inventory walk complete"
);
}
ConfigChangeRequest::PeriodicRescan => {
let started = Instant::now();
let (observed, changed) = run_rescan(
&active,
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
None,
)
.await;
let duration_ms = started.elapsed().as_millis() as u64;
crate::telemetry::capture_global(
crate::telemetry::Event::config_periodic_rescan_completed(
"claude-code",
observed,
changed,
duration_ms,
),
);
tracing::info!(
items_observed = observed,
items_changed = changed,
duration_ms,
"config_monitor: periodic rescan complete"
);
}
ConfigChangeRequest::ManualRescan { path_filter } => {
let started = Instant::now();
let (observed, changed) = run_rescan(
&active,
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
path_filter.as_deref(),
)
.await;
let duration_ms = started.elapsed().as_millis() as u64;
crate::telemetry::capture_global(
crate::telemetry::Event::config_periodic_rescan_completed(
"claude-code",
observed,
changed,
duration_ms,
),
);
}
ConfigChangeRequest::FsAdded(_)
| ConfigChangeRequest::FsModified(_)
| ConfigChangeRequest::FsRemoved(_) => {
let (kind, paths) = match req {
ConfigChangeRequest::FsAdded(p) => (ChangeKind::Added, p),
ConfigChangeRequest::FsModified(p) => (ChangeKind::Modified, p),
ConfigChangeRequest::FsRemoved(p) => (ChangeKind::Removed, p),
_ => continue,
};
let roots = projects.paths();
for path in &paths {
handle_fs_event(
kind,
path,
&active,
&roots,
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
)
.await;
}
if reshapes_tree(&paths, &manifest, &roots) {
let appeared: Vec<PathBuf> =
paths.iter().filter(|p| p.is_dir()).cloned().collect();
rearm(
watch.as_ref(),
&manifest,
&active,
&roots,
&appeared,
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
)
.await;
}
}
ConfigChangeRequest::DetectAgents => {
let now = detect();
if now != detected {
let appeared: HashSet<String> = now.difference(&detected).cloned().collect();
detected = now;
active = manifest.only_agents(&detected);
if !appeared.is_empty() {
let emitted = run_full_walk(
&manifest.only_agents(&appeared),
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
EventSource::Rescan,
None,
)
.await;
tracing::info!(
agents = ?appeared,
items_emitted = emitted,
"config_monitor: newly detected agents inventoried"
);
}
}
rearm(
watch.as_ref(),
&manifest,
&active,
&projects.paths(),
&[],
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
)
.await;
}
ConfigChangeRequest::ProjectScopeRegister { project_root } => {
if !projects.register(&project_root) {
continue;
}
let emitted = run_project_scope_scan(
&project_root,
&active,
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
)
.await;
rearm(
watch.as_ref(),
&manifest,
&active,
&projects.paths(),
&[],
&cache,
&privacy_filter,
cloud_tx.as_ref(),
&event_logger,
&config,
&mut diff_counters,
)
.await;
tracing::info!(
project_root = %project_root.display(),
items_emitted = emitted,
"config_monitor: project-scope registered"
);
}
}
}
}
async fn wait_for_shutdown(shutdown: &mut watch::Receiver<bool>) {
loop {
let stopped = *shutdown.borrow_and_update();
if stopped || shutdown.changed().await.is_err() {
return;
}
}
}
#[allow(clippy::too_many_arguments)]
async fn run_full_walk(
manifest: &Manifest,
cache: &ContentHashCache,
privacy_filter: &PrivacyFilter,
cloud_tx: Option<&mpsc::Sender<CloudEvent>>,
event_logger: &EventLogger,
config: &Config,
diff_counters: &mut HashMap<(String, String), u64>,
source: EventSource,
path_filter: Option<&Path>,
) -> usize {
let mut emitted = 0usize;
for agent in &manifest.agents {
for ap in &agent.paths {
if ap.is_project_scoped() {
continue;
}
for path in expand_paths_for_scan(ap) {
if let Some(filter) = path_filter {
if !crate::path_compat::dedup_key(&path)
.starts_with(crate::path_compat::dedup_key(filter))
{
continue;
}
}
if !path.exists() {
continue;
}
if is_excluded(&path, ap, manifest) {
continue;
}
let counter = diff_counters
.entry((agent.name.clone(), ap.kind.clone()))
.or_insert(0);
let counter_value = match source {
EventSource::InitScan => 0,
_ => {
*counter = counter.saturating_add(1);
*counter
}
};
for event in scan_path_to_event(
&path,
ap,
&agent.name,
ChangeKind::Snapshot,
source,
counter_value,
privacy_filter,
cache,
config,
) {
log_and_forward(&event, cloud_tx, event_logger).await;
emitted += 1;
}
}
}
}
emitted
}
#[allow(clippy::too_many_arguments)]
async fn run_rescan(
manifest: &Manifest,
cache: &ContentHashCache,
privacy_filter: &PrivacyFilter,
cloud_tx: Option<&mpsc::Sender<CloudEvent>>,
event_logger: &EventLogger,
config: &Config,
diff_counters: &mut HashMap<(String, String), u64>,
path_filter: Option<&Path>,
) -> (usize, usize) {
let mut observed = 0usize;
let mut changed = 0usize;
let mut seen: HashSet<PathBuf> = HashSet::new();
for agent in &manifest.agents {
for ap in &agent.paths {
if ap.is_project_scoped() {
continue;
}
for path in expand_paths_for_scan(ap) {
if let Some(filter) = path_filter {
if !crate::path_compat::dedup_key(&path)
.starts_with(crate::path_compat::dedup_key(filter))
{
continue;
}
}
if !path.exists() {
continue;
}
if is_excluded(&path, ap, manifest) {
continue;
}
seen.insert(crate::path_compat::dedup_key(&path));
observed += 1;
let mut prev_hashes: HashMap<Option<String>, [u8; 32]> = HashMap::new();
for entry in cache.entries_for_path(&path, &agent.name, &ap.kind) {
prev_hashes.insert(entry.subpath.clone(), entry.content_hash);
}
let counter = diff_counters
.entry((agent.name.clone(), ap.kind.clone()))
.or_insert(0);
*counter = counter.saturating_add(1);
let counter_value = *counter;
let events = scan_path_to_event_with_hash(
&path,
ap,
&agent.name,
ChangeKind::Snapshot,
EventSource::Rescan,
counter_value,
privacy_filter,
cache,
config,
);
if events.is_empty() {
continue;
}
let mut emitted_subpaths: HashSet<Option<String>> = HashSet::new();
for (event, new_hash, subpath) in events {
let prev = prev_hashes.get(&subpath).copied();
let is_change = prev.map(|p| p != new_hash).unwrap_or(true);
if is_change {
changed += 1;
}
emitted_subpaths.insert(subpath);
log_and_forward(&event, cloud_tx, event_logger).await;
}
for (sub, _) in prev_hashes
.iter()
.filter(|(s, _)| s.is_some() && !emitted_subpaths.contains(*s))
{
if let Some(sub_str) = sub {
cache.remove_subpath(&path, &agent.name, &ap.kind, Some(sub_str));
if let Some(removed_event) = build_removed_event_for_subpath(
&path,
Some(sub_str.as_str()),
&ap.kind,
&agent.name,
EventSource::Rescan,
*diff_counters
.entry((agent.name.clone(), ap.kind.clone()))
.and_modify(|c| *c = c.saturating_add(1))
.or_insert(1),
severity_for(&ap.kind, "removed"),
config,
) {
log_and_forward(&removed_event, cloud_tx, event_logger).await;
changed += 1;
}
}
}
}
}
}
let mut to_remove: Vec<CacheEntry> = Vec::new();
for entry in cache.snapshot() {
if seen.contains(&crate::path_compat::dedup_key(&entry.path)) {
continue;
}
if !manifest.agents.iter().any(|a| a.name == entry.agent) {
continue;
}
if entry.path.exists() && !dropped_plugin_file(&entry, manifest) {
continue;
}
to_remove.push(entry);
}
for entry in to_remove {
cache.remove_subpath(
&entry.path,
&entry.agent,
&entry.kind,
entry.subpath.as_deref(),
);
if let Some(event) = build_removed_event_for_subpath(
&entry.path,
entry.subpath.as_deref(),
&entry.kind,
&entry.agent,
EventSource::Rescan,
diff_counters
.entry((entry.agent.clone(), entry.kind.clone()))
.and_modify(|c| *c = c.saturating_add(1))
.or_insert(1)
.to_owned(),
severity_for(&entry.kind, "removed"),
config,
) {
log_and_forward(&event, cloud_tx, event_logger).await;
changed += 1;
}
}
(observed, changed)
}
#[allow(clippy::too_many_arguments)]
async fn run_project_scope_scan(
project_root: &Path,
manifest: &Manifest,
cache: &ContentHashCache,
privacy_filter: &PrivacyFilter,
cloud_tx: Option<&mpsc::Sender<CloudEvent>>,
event_logger: &EventLogger,
config: &Config,
diff_counters: &mut HashMap<(String, String), u64>,
) -> usize {
let mut emitted = 0usize;
for agent in &manifest.agents {
for ap in &agent.paths {
if !ap.is_project_scoped() {
continue;
}
for path in project_candidates(project_root, ap) {
if !path.exists() {
continue;
}
if is_excluded(&path, ap, manifest) {
continue;
}
let counter_value = {
let counter = diff_counters
.entry((agent.name.clone(), ap.kind.clone()))
.or_insert(0);
*counter = counter.saturating_add(1);
*counter
};
for event in scan_path_to_event(
&path,
ap,
&agent.name,
ChangeKind::Snapshot,
EventSource::InitScan,
counter_value,
privacy_filter,
cache,
config,
) {
log_and_forward(&event, cloud_tx, event_logger).await;
emitted += 1;
}
}
}
}
emitted
}
#[allow(clippy::too_many_arguments)]
async fn handle_fs_event(
kind: ChangeKind,
path: &Path,
manifest: &Manifest,
project_roots: &[PathBuf],
cache: &ContentHashCache,
privacy_filter: &PrivacyFilter,
cloud_tx: Option<&mpsc::Sender<CloudEvent>>,
event_logger: &EventLogger,
config: &Config,
diff_counters: &mut HashMap<(String, String), u64>,
) {
if !path.exists() {
sweep_removed_tree(
path,
manifest,
cache,
cloud_tx,
event_logger,
config,
diff_counters,
)
.await;
}
let claims = claims_for_path(path, manifest)
.into_iter()
.chain(project_claims_for_path(path, manifest, project_roots));
for (agent_name, ap) in claims {
if is_excluded(path, ap, manifest) {
continue;
}
handle_fs_claim(
kind,
path,
agent_name,
ap,
cache,
privacy_filter,
cloud_tx,
event_logger,
config,
diff_counters,
)
.await;
}
}
#[allow(clippy::too_many_arguments)]
async fn handle_fs_claim(
kind: ChangeKind,
path: &Path,
agent_name: &str,
ap: &AgentPath,
cache: &ContentHashCache,
privacy_filter: &PrivacyFilter,
cloud_tx: Option<&mpsc::Sender<CloudEvent>>,
event_logger: &EventLogger,
config: &Config,
diff_counters: &mut HashMap<(String, String), u64>,
) {
let known = !cache
.entries_for_path(path, agent_name, &ap.kind)
.is_empty();
let kind = match (kind, path.exists(), known) {
(ChangeKind::Snapshot, ..) => ChangeKind::Snapshot,
(_, false, _) => ChangeKind::Removed,
(_, true, true) => ChangeKind::Modified,
(_, true, false) => ChangeKind::Added,
};
let agent_owned = agent_name.to_string();
let counter = diff_counters
.entry((agent_owned.clone(), ap.kind.clone()))
.or_insert(0);
*counter = counter.saturating_add(1);
let counter_value = *counter;
let events: Vec<CloudEvent> = match kind {
ChangeKind::Removed => {
let removed_entries = cache.remove_all_under_path(path, agent_name, &ap.kind);
if removed_entries.is_empty() {
build_removed_event_for_subpath(
path,
None,
&ap.kind,
&agent_owned,
EventSource::FsWatcher,
counter_value,
severity_for(&ap.kind, "removed"),
config,
)
.into_iter()
.collect()
} else {
let mut out = Vec::with_capacity(removed_entries.len());
for (i, entry) in removed_entries.into_iter().enumerate() {
let c = if i == 0 {
counter_value
} else {
let counter = diff_counters
.entry((agent_owned.clone(), ap.kind.clone()))
.or_insert(0);
*counter = counter.saturating_add(1);
*counter
};
if let Some(event) = build_removed_event_for_subpath(
path,
entry.subpath.as_deref(),
&ap.kind,
&agent_owned,
EventSource::FsWatcher,
c,
severity_for(&ap.kind, "removed"),
config,
) {
out.push(event);
}
}
out
}
}
ChangeKind::Added | ChangeKind::Modified => {
let mut prev_hashes: HashMap<Option<String>, [u8; 32]> = HashMap::new();
for entry in cache.entries_for_path(path, agent_name, &ap.kind) {
prev_hashes.insert(entry.subpath.clone(), entry.content_hash);
}
let scanned = scan_path_to_event_with_hash(
path,
ap,
&agent_owned,
kind,
EventSource::FsWatcher,
counter_value,
privacy_filter,
cache,
config,
);
let mut emitted_subpaths: HashSet<Option<String>> = HashSet::new();
let mut out: Vec<CloudEvent> = Vec::with_capacity(scanned.len());
for (event, new_hash, subpath) in scanned {
emitted_subpaths.insert(subpath.clone());
let prev = prev_hashes.get(&subpath).copied();
if prev == Some(new_hash) {
continue;
}
out.push(event);
}
for (sub, _) in prev_hashes
.iter()
.filter(|(s, _)| s.is_some() && !emitted_subpaths.contains(*s))
{
if let Some(sub_str) = sub {
cache.remove_subpath(path, agent_name, &ap.kind, Some(sub_str));
let counter = diff_counters
.entry((agent_owned.clone(), ap.kind.clone()))
.or_insert(0);
*counter = counter.saturating_add(1);
let c = *counter;
if let Some(removed_event) = build_removed_event_for_subpath(
path,
Some(sub_str.as_str()),
&ap.kind,
&agent_owned,
EventSource::FsWatcher,
c,
severity_for(&ap.kind, "removed"),
config,
) {
out.push(removed_event);
}
}
}
out
}
ChangeKind::Snapshot => Vec::new(),
};
if events.is_empty() {
return;
}
crate::telemetry::capture_global(crate::telemetry::Event::config_change_detected(
&ap.kind,
severity_for(&ap.kind, kind.severity_key()).as_str(),
&agent_owned,
match kind {
ChangeKind::Added => "added",
ChangeKind::Modified => "modified",
ChangeKind::Removed => "removed",
ChangeKind::Snapshot => "snapshot",
},
));
for event in &events {
log_and_forward(event, cloud_tx, event_logger).await;
}
}
fn agent_path_matches(ap: &AgentPath, path: &Path) -> bool {
let expands_to = |template: &str| expand_path(template, None).is_ok_and(|e| e == path);
if ap.paths.iter().any(|p| expands_to(p))
|| ap.json_slice_paths.iter().any(|s| expands_to(&s.path))
{
return true;
}
let Some(g) = &ap.paths_glob else {
return false;
};
let Ok(expanded) = expand_path(g, None) else {
return false;
};
let matches = match (expanded.to_str().map(glob::Pattern::new), path.to_str()) {
(Some(Ok(pattern)), Some(path_str)) => pattern.matches(path_str),
_ => false,
};
matches && (!ap.claude_plugins || claude_plugins::admits(path))
}
fn dropped_plugin_file(entry: &CacheEntry, manifest: &Manifest) -> bool {
!claude_plugins::admits(&entry.path)
&& manifest
.agents
.iter()
.filter(|agent| agent.name == entry.agent)
.flat_map(|agent| &agent.paths)
.any(|ap| ap.claude_plugins && ap.kind == entry.kind)
}
fn claims_for_path<'a>(path: &Path, manifest: &'a Manifest) -> Vec<(&'a str, &'a AgentPath)> {
manifest
.agents
.iter()
.flat_map(|agent| agent.paths.iter().map(move |ap| (agent.name.as_str(), ap)))
.filter(|(_, ap)| !ap.is_project_scoped() && agent_path_matches(ap, path))
.collect()
}
fn project_claims_for_path<'a>(
path: &Path,
manifest: &'a Manifest,
project_roots: &[PathBuf],
) -> Vec<(&'a str, &'a AgentPath)> {
let options = glob::MatchOptions {
require_literal_separator: true,
..Default::default()
};
let rels: Vec<String> = project_roots
.iter()
.filter_map(|root| path.strip_prefix(root).ok())
.map(|rel| rel.to_string_lossy().replace('\\', "/"))
.collect();
let names = |ap: &AgentPath, rel: &str| {
ap.paths_relative.iter().any(|r| r == rel)
|| ap
.paths_glob_relative
.iter()
.any(|g| glob::Pattern::new(g).is_ok_and(|p| p.matches_with(rel, options)))
};
manifest
.agents
.iter()
.flat_map(|agent| agent.paths.iter().map(move |ap| (agent.name.as_str(), ap)))
.filter(|(_, ap)| ap.is_project_scoped() && rels.iter().any(|rel| names(ap, rel)))
.collect()
}
fn project_candidates(project_root: &Path, ap: &AgentPath) -> Vec<PathBuf> {
let mut candidates: Vec<PathBuf> = ap
.paths_relative
.iter()
.map(|rel| project_root.join(rel))
.collect();
for glob in &ap.paths_glob_relative {
candidates.extend(glob_expand(&project_root.join(glob)));
}
candidates.sort();
candidates.dedup();
candidates
}
fn reshapes_tree(paths: &[PathBuf], manifest: &Manifest, project_roots: &[PathBuf]) -> bool {
paths.iter().any(|p| {
p.is_dir()
|| (!p.exists()
&& claims_for_path(p, manifest).is_empty()
&& project_claims_for_path(p, manifest, project_roots).is_empty())
})
}
#[allow(clippy::too_many_arguments)]
async fn rearm(
watch: Option<&WatchControl>,
full_manifest: &Arc<Manifest>,
active: &Manifest,
project_roots: &[PathBuf],
appeared: &[PathBuf],
cache: &ContentHashCache,
privacy_filter: &PrivacyFilter,
cloud_tx: Option<&mpsc::Sender<CloudEvent>>,
event_logger: &EventLogger,
config: &Config,
diff_counters: &mut HashMap<(String, String), u64>,
) {
let Some(watch) = watch else {
return;
};
let plan_manifest = full_manifest.clone();
let plan_roots = project_roots.to_vec();
let Ok(plan) =
tokio::task::spawn_blocking(move || watch_plan(&plan_manifest, &plan_roots)).await
else {
return;
};
let applied = watch.apply(plan).await;
for gone in &applied.vanished {
sweep_removed_tree(
gone,
active,
cache,
cloud_tx,
event_logger,
config,
diff_counters,
)
.await;
}
if applied.armed.is_empty() && appeared.is_empty() {
return;
}
let armed: Vec<PathBuf> = applied
.armed
.iter()
.chain(appeared)
.map(|d| crate::path_compat::dedup_key(d))
.collect();
let mut candidates: Vec<PathBuf> = Vec::new();
for ap in active.agents.iter().flat_map(|a| a.paths.iter()) {
if ap.is_project_scoped() {
for root in project_roots {
candidates.extend(project_candidates(root, ap));
}
} else {
candidates.extend(expand_paths_for_scan(ap));
}
}
candidates.sort();
candidates.dedup();
candidates.retain(|c| {
let key = crate::path_compat::dedup_key(c);
c.is_file() && armed.iter().any(|dir| key.starts_with(dir))
});
for path in candidates {
handle_fs_event(
ChangeKind::Added,
&path,
active,
project_roots,
cache,
privacy_filter,
cloud_tx,
event_logger,
config,
diff_counters,
)
.await;
}
}
async fn sweep_removed_tree(
path: &Path,
manifest: &Manifest,
cache: &ContentHashCache,
cloud_tx: Option<&mpsc::Sender<CloudEvent>>,
event_logger: &EventLogger,
config: &Config,
diff_counters: &mut HashMap<(String, String), u64>,
) {
let gone = crate::path_compat::dedup_key(path);
let below: Vec<CacheEntry> = cache
.snapshot()
.into_iter()
.filter(|e| {
let key = crate::path_compat::dedup_key(&e.path);
key != gone
&& key.starts_with(&gone)
&& manifest.agents.iter().any(|a| a.name == e.agent)
})
.collect();
for entry in below {
cache.remove_subpath(
&entry.path,
&entry.agent,
&entry.kind,
entry.subpath.as_deref(),
);
let counter = diff_counters
.entry((entry.agent.clone(), entry.kind.clone()))
.or_insert(0);
*counter = counter.saturating_add(1);
if let Some(event) = build_removed_event_for_subpath(
&entry.path,
entry.subpath.as_deref(),
&entry.kind,
&entry.agent,
EventSource::FsWatcher,
*counter,
severity_for(&entry.kind, "removed"),
config,
) {
log_and_forward(&event, cloud_tx, event_logger).await;
}
}
}
fn expand_paths_for_scan(ap: &AgentPath) -> Vec<PathBuf> {
let mut out = Vec::new();
match ap.watch_strategy {
WatchStrategy::ExactFile | WatchStrategy::ExactFileWithSlice => {
for p in &ap.paths {
if let Ok(expanded) = expand_path(p, None) {
out.push(expanded);
}
}
for slice in &ap.json_slice_paths {
if let Ok(expanded) = expand_path(&slice.path, None) {
out.push(expanded);
}
}
}
WatchStrategy::Glob => {
if let Some(g) = &ap.paths_glob {
if let Ok(expanded) = expand_path(g, None) {
out.extend(glob_expand(&expanded));
}
}
}
WatchStrategy::ExactFileAndGlob => {
for p in &ap.paths {
if let Ok(expanded) = expand_path(p, None) {
out.push(expanded);
}
}
if let Some(g) = &ap.paths_glob {
if let Ok(expanded) = expand_path(g, None) {
out.extend(glob_expand(&expanded));
}
}
}
}
if ap.claude_plugins {
claude_plugins::retain_active(&mut out);
}
out.sort();
out.dedup();
out
}
#[allow(clippy::too_many_arguments)]
fn scan_path_to_event(
path: &Path,
ap: &AgentPath,
agent_name: &str,
kind: ChangeKind,
source: EventSource,
diffcounter: u64,
privacy_filter: &PrivacyFilter,
cache: &ContentHashCache,
config: &Config,
) -> Vec<CloudEvent> {
scan_path_to_event_with_hash(
path,
ap,
agent_name,
kind,
source,
diffcounter,
privacy_filter,
cache,
config,
)
.into_iter()
.map(|(event, _, _)| event)
.collect()
}
#[allow(clippy::too_many_arguments)]
fn scan_path_to_event_with_hash(
path: &Path,
ap: &AgentPath,
agent_name: &str,
kind: ChangeKind,
source: EventSource,
diffcounter: u64,
privacy_filter: &PrivacyFilter,
cache: &ContentHashCache,
config: &Config,
) -> Vec<(CloudEvent, [u8; 32], Option<String>)> {
let raw = match std::fs::read_to_string(path) {
Ok(s) => s,
Err(e) => {
tracing::debug!(
path = %path.display(),
error = %e,
"config_monitor: read failed during scan"
);
return Vec::new();
}
};
build_event_from_content(
path,
ap,
agent_name,
kind,
source,
diffcounter,
&raw,
privacy_filter,
cache,
config,
)
}
#[allow(clippy::too_many_arguments)]
fn build_event_from_content(
path: &Path,
ap: &AgentPath,
agent_name: &str,
kind: ChangeKind,
source: EventSource,
diffcounter: u64,
raw: &str,
privacy_filter: &PrivacyFilter,
cache: &ContentHashCache,
config: &Config,
) -> Vec<(CloudEvent, [u8; 32], Option<String>)> {
let slice_pointers: Vec<&str> =
if matches!(ap.watch_strategy, WatchStrategy::ExactFileWithSlice) {
ap.json_slice_paths
.iter()
.filter_map(|s| {
expand_path(&s.path, None)
.ok()
.and_then(|exp| (exp == path).then_some(s.json_pointer.as_str()))
})
.collect()
} else {
Vec::new()
};
if ap.kind == "mcp" && ap.slice_kind_subpath {
return build_mcp_fanout_events(
path,
ap,
agent_name,
kind,
source,
diffcounter,
raw,
&slice_pointers,
privacy_filter,
cache,
config,
);
}
let hashed = if ap.opaque {
Ok(content_hash_text(raw))
} else {
hash_for_kind(&ap.kind, raw, &slice_pointers)
};
let content_hash = match hashed {
Ok(h) => h,
Err(e) => {
tracing::warn!(
code = crate::error::ERR_INVENTORY_HASH_FAILED,
path = %path.display(),
error = ?e,
"config_monitor: hash pipeline failed"
);
return Vec::new();
}
};
let path_hash = sha256(path.to_string_lossy().as_bytes());
cache.insert(CacheEntry {
path: path.to_path_buf(),
subpath: None,
path_hash,
content_hash,
last_observed: Instant::now(),
kind: ap.kind.clone(),
agent: agent_name.to_string(),
});
let severity = severity_for(&ap.kind, kind.severity_key());
let data = match config.inventory_monitor.content_forward {
ContentForwardMode::HashOnly => Value::Null,
_ if ap.opaque => opaque_payload(&ap.kind, ap.scope, path),
ContentForwardMode::Filtered => {
filtered_payload(raw, &ap.kind, ap.scope, path, privacy_filter, config)
}
ContentForwardMode::FullUnfiltered => {
unfiltered_payload(raw, &ap.kind, ap.scope, path, config)
}
};
let event = build_modified_event(
path,
&ap.kind,
agent_name,
&content_hash,
&path_hash,
kind,
source,
diffcounter,
severity,
data,
config,
);
vec![(event, content_hash, None)]
}
fn resolve_mcp_servers(parsed: &Value, slice_pointers: &[&str]) -> serde_json::Map<String, Value> {
if slice_pointers.is_empty() {
return parsed
.pointer("/mcpServers")
.and_then(|v| v.as_object())
.cloned()
.or_else(|| parsed.as_object().cloned())
.unwrap_or_default();
}
let mut out = serde_json::Map::new();
for pointer in slice_pointers {
let Some(sliced) = parsed.pointer(pointer).and_then(|v| v.as_object()) else {
continue;
};
if pointer.rsplit('/').next() == Some("mcpServers") {
for (name, cfg) in sliced {
out.insert(name.clone(), cfg.clone());
}
} else {
for (block, value) in sliced {
let Some(nested) = value.pointer("/mcpServers").and_then(|v| v.as_object()) else {
continue;
};
for (name, cfg) in nested {
out.insert(format!("{block}::{name}"), cfg.clone());
}
}
}
}
out
}
#[allow(clippy::too_many_arguments)]
fn build_mcp_fanout_events(
path: &Path,
ap: &AgentPath,
agent_name: &str,
kind: ChangeKind,
source: EventSource,
diffcounter: u64,
raw: &str,
slice_pointers: &[&str],
privacy_filter: &PrivacyFilter,
cache: &ContentHashCache,
config: &Config,
) -> Vec<(CloudEvent, [u8; 32], Option<String>)> {
let parsed: Result<Value, String> = if ap.projected {
project_mcp_settings(raw)
} else {
serde_json::from_str(raw).map_err(|e| e.to_string())
};
let parsed = match parsed {
Ok(v) => v,
Err(e) => {
tracing::debug!(
path = %path.display(),
error = %e,
"config_monitor: mcp parse failed; falling back to non-fanout event"
);
return Vec::new();
}
};
let servers = resolve_mcp_servers(&parsed, slice_pointers);
if servers.is_empty() {
return Vec::new();
}
let severity = severity_for(&ap.kind, kind.severity_key());
let mut out = Vec::with_capacity(servers.len());
for (server_name, server_value) in servers {
let content_hash = hash_mcp_server_entry(&server_name, &server_value);
let path_hash = subpath_path_hash(path, Some(&server_name));
cache.insert(CacheEntry {
path: path.to_path_buf(),
subpath: Some(server_name.clone()),
path_hash,
content_hash,
last_observed: Instant::now(),
kind: ap.kind.clone(),
agent: agent_name.to_string(),
});
let data = match config.inventory_monitor.content_forward {
ContentForwardMode::HashOnly => Value::Null,
ContentForwardMode::Filtered => {
let mut payload = serde_json::json!({
"server_name": server_name,
"tools": [],
});
filter_event_with(&mut payload, privacy_filter);
attach_scope(&mut payload, &ap.kind, ap.scope);
payload
}
ContentForwardMode::FullUnfiltered => {
if std::env::var("OPENLATCH_TESTING").as_deref() == Ok("true") {
let mut payload = serde_json::json!({
"server_name": server_name,
"tools": [],
"raw": server_value,
});
attach_scope(&mut payload, &ap.kind, ap.scope);
payload
} else {
let mut payload = serde_json::json!({
"server_name": server_name,
"tools": [],
});
filter_event_with(&mut payload, privacy_filter);
attach_scope(&mut payload, &ap.kind, ap.scope);
payload
}
}
};
let event = build_modified_event(
path,
&ap.kind,
agent_name,
&content_hash,
&path_hash,
kind,
source,
diffcounter,
severity,
data,
config,
);
out.push((event, content_hash, Some(server_name)));
}
out
}
fn filtered_payload(
raw: &str,
kind: &str,
scope: Option<ConfigScope>,
path: &Path,
privacy_filter: &PrivacyFilter,
config: &Config,
) -> Value {
let inline_cap = config.inventory_monitor.max_inline_content_bytes as usize;
let mut payload = build_filtered_body(raw, kind, inline_cap, path, privacy_filter);
attach_scope(&mut payload, kind, scope);
payload
}
fn build_filtered_body(
raw: &str,
kind: &str,
inline_cap: usize,
path: &Path,
privacy_filter: &PrivacyFilter,
) -> Value {
match kind {
"rules" => {
let (frontmatter, body_only) = split_frontmatter(raw);
let cap = per_kind_body_cap(kind, inline_cap);
let body = truncate_body_in_place(&body_only, cap);
let mut body_v = Value::String(body);
filter_event_with(&mut body_v, privacy_filter);
let mut payload = serde_json::Map::new();
payload.insert("body".into(), body_v);
payload.insert("frontmatter".into(), frontmatter);
payload.insert(
"path".into(),
Value::String(truncate_path_for_platform(path)),
);
Value::Object(payload)
}
"skill" | "command" => {
let (frontmatter, body_only) = split_frontmatter(raw);
let cap = per_kind_body_cap(kind, inline_cap);
let body = truncate_body_in_place(&body_only, cap);
let mut body_v = Value::String(body);
filter_event_with(&mut body_v, privacy_filter);
let (name, description) = derive_name_and_description(&frontmatter, kind, path);
let mut payload = serde_json::Map::new();
payload.insert("body".into(), body_v);
payload.insert("frontmatter".into(), frontmatter);
payload.insert("name".into(), Value::String(name));
payload.insert("description".into(), Value::String(description));
Value::Object(payload)
}
"hooks" => build_hooks_body(raw, path, privacy_filter),
"mcp" => {
let server_name = derive_name_from_path(path);
serde_json::json!({
"server_name": server_name,
"tools": [],
})
}
_ => legacy_text_payload(kind, raw, privacy_filter),
}
}
fn build_hooks_body(raw: &str, path: &Path, privacy_filter: &PrivacyFilter) -> Value {
let parsed = match jsonc_parser::parse_to_serde_value(raw, &Default::default()) {
Ok(Some(v)) => Some((v, false)),
Ok(None) | Err(_) => serde_json::from_str::<Value>(raw)
.map(|v| (v, false))
.or_else(|_| toml::from_str::<Value>(raw).map(|v| (v, true)))
.ok(),
};
let hooks_obj = match parsed {
Some((doc, is_toml)) => match doc.pointer("/hooks").filter(|v| v.is_object()) {
Some(hooks) => hooks.clone(),
None if !is_toml && doc.is_object() => doc,
None => Value::Object(Default::default()),
},
None => {
tracing::debug!("config_monitor: hooks file parse failed; emitting empty hooks dict");
Value::Object(Default::default())
}
};
let mut wrapped = serde_json::json!({"hooks": hooks_obj});
filter_event_with(&mut wrapped, privacy_filter);
wrapped["path"] = Value::String(hooks_display_path(path));
wrapped
}
fn hooks_display_path(path: &Path) -> String {
let home = dirs::home_dir().map(|h| crate::path_compat::display_path(&h));
let shown = crate::path_compat::home_normalized(
&crate::path_compat::display_path(path),
home.as_deref(),
);
truncate_path_for_platform(Path::new(&shown))
}
const SUBPATH_HASH_SEP: char = '#';
fn subpath_path_hash(path: &Path, subpath: Option<&str>) -> [u8; 32] {
let path_str = path.to_string_lossy();
match subpath {
Some(sub) => sha256(format!("{path_str}{SUBPATH_HASH_SEP}{sub}").as_bytes()),
None => sha256(path_str.as_bytes()),
}
}
const PROJECTED_OUT_SERVER_FIELDS: &[&str] = &["env", "headers", "http_headers", "bearer_token"];
fn project_mcp_settings(raw: &str) -> Result<Value, String> {
#[derive(serde::Deserialize)]
struct Settings {
#[serde(default, rename = "mcpServers", alias = "mcp_servers")]
mcp_servers: std::collections::BTreeMap<String, ProjectedServer>,
}
let settings: Settings = serde_json::from_str(raw).or_else(|json| {
toml::from_str(raw).map_err(|toml| format!("json: {json}; toml: {toml}"))
})?;
let servers = settings
.mcp_servers
.into_iter()
.map(|(name, server)| (name, server.0))
.collect();
Ok(serde_json::json!({ "mcpServers": Value::Object(servers) }))
}
struct ProjectedServer(Value);
impl<'de> serde::Deserialize<'de> for ProjectedServer {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
struct Visitor;
impl<'de> serde::de::Visitor<'de> for Visitor {
type Value = ProjectedServer;
fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result {
formatter.write_str("an MCP server entry")
}
fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
where
A: serde::de::MapAccess<'de>,
{
let mut kept = serde_json::Map::new();
while let Some(key) = map.next_key::<String>()? {
if PROJECTED_OUT_SERVER_FIELDS.contains(&key.as_str()) {
map.next_value::<serde::de::IgnoredAny>()?;
} else {
kept.insert(key, map.next_value::<Value>()?);
}
}
Ok(ProjectedServer(Value::Object(kept)))
}
}
deserializer.deserialize_map(Visitor)
}
}
fn hash_mcp_server_entry(server_name: &str, server_value: &Value) -> [u8; 32] {
let nfc: String = server_name.nfc().collect();
let mut wrap = serde_json::Map::new();
wrap.insert(nfc, server_value.clone());
let synth = Value::Object(wrap);
let canonical = serde_json_canonicalizer::to_string(&synth)
.unwrap_or_else(|_| serde_json::to_string(&synth).unwrap_or_default());
sha256(canonical.as_bytes())
}
fn legacy_text_payload(kind: &str, raw: &str, privacy_filter: &PrivacyFilter) -> Value {
let mut as_string = Value::String(raw.to_string());
filter_event_with(&mut as_string, privacy_filter);
serde_json::json!({"kind": kind, "text": as_string})
}
const PLATFORM_RULES_BODY_MAX: usize = 512 * 1024;
const PLATFORM_SKILL_COMMAND_BODY_MAX: usize = 128 * 1024;
const PLATFORM_NAME_MAX: usize = 128;
const PLATFORM_PATH_MAX: usize = 1024;
const TRUNCATION_MARKER_RESERVE: usize = 64;
fn per_kind_body_cap(kind: &str, inline_cap: usize) -> usize {
let platform_max = match kind {
"rules" => PLATFORM_RULES_BODY_MAX,
"skill" | "command" => PLATFORM_SKILL_COMMAND_BODY_MAX,
_ => usize::MAX,
};
std::cmp::min(inline_cap, platform_max)
}
fn truncate_body_in_place(raw: &str, cap: usize) -> String {
if raw.len() <= cap {
return raw.to_string();
}
let half = cap.saturating_sub(TRUNCATION_MARKER_RESERVE) / 2;
let head_end = raw.char_indices().nth(half).map(|(i, _)| i).unwrap_or(half);
let tail_target = raw.len().saturating_sub(half);
let mut tail_start = tail_target;
while tail_start < raw.len() && !raw.is_char_boundary(tail_start) {
tail_start += 1;
}
let head = &raw[..head_end];
let tail = &raw[tail_start..];
let marker = format!(
"\n\n…<truncated {} bytes>…\n\n",
raw.len() - head.len() - tail.len()
);
format!("{head}{marker}{tail}")
}
fn derive_name_from_path(path: &Path) -> String {
let stem = path
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or("unnamed");
if stem.is_empty() {
return "unnamed".into();
}
if stem.len() <= PLATFORM_NAME_MAX {
return stem.to_string();
}
stem.chars().take(PLATFORM_NAME_MAX).collect()
}
fn derive_name_and_description(frontmatter: &Value, kind: &str, path: &Path) -> (String, String) {
let mut name = frontmatter
.get("name")
.and_then(|v| v.as_str())
.map(str::to_string)
.unwrap_or_default();
if name.is_empty() {
name = match kind {
"skill" => path
.parent()
.and_then(|p| p.file_name())
.and_then(|s| s.to_str())
.map(str::to_string)
.unwrap_or_else(|| derive_name_from_path(path)),
_ => derive_name_from_path(path),
};
}
if name.is_empty() {
name = "unnamed".to_string();
}
if name.chars().count() > PLATFORM_NAME_MAX {
name = name.chars().take(PLATFORM_NAME_MAX).collect();
}
let description = frontmatter
.get("description")
.and_then(|v| v.as_str())
.map(str::to_string)
.unwrap_or_default();
let description = if description.len() > 4096 {
description.chars().take(4096).collect()
} else {
description
};
(name, description)
}
fn split_frontmatter(raw: &str) -> (Value, String) {
let stripped = raw.strip_prefix('\u{FEFF}').unwrap_or(raw);
let after_optional_lf = stripped.strip_prefix('\n').unwrap_or(stripped);
let body = after_optional_lf;
let first_line_terminator = if body.starts_with("---\r\n") {
Some(5)
} else if body.starts_with("---\n") {
Some(4)
} else if body == "---" {
Some(3)
} else {
None
};
let Some(start) = first_line_terminator else {
return (Value::Object(Default::default()), raw.to_string());
};
let after_open = &body[start..];
let mut close_offset = None;
let mut cursor = 0;
for line in after_open.split_inclusive('\n') {
let stripped_line = line.trim_end_matches(['\n', '\r']);
if stripped_line == "---" {
close_offset = Some(cursor + line.len());
break;
}
cursor += line.len();
}
let Some(close) = close_offset else {
return (Value::Object(Default::default()), raw.to_string());
};
let frontmatter_block = &after_open[..cursor];
let body_remainder = &after_open[close..];
let body_remainder = body_remainder.strip_prefix('\n').unwrap_or(body_remainder);
let mut map = serde_json::Map::new();
for line in frontmatter_block.lines() {
let trimmed = line.trim();
if trimmed.is_empty() || trimmed.starts_with('#') {
continue;
}
let Some((key, value)) = trimmed.split_once(':') else {
continue;
};
let key = key.trim().to_string();
if key.is_empty() {
continue;
}
let mut value = value.trim().to_string();
if value.len() >= 2 {
let bytes = value.as_bytes();
let first = bytes[0];
let last = bytes[bytes.len() - 1];
if (first == b'"' && last == b'"') || (first == b'\'' && last == b'\'') {
value = value[1..value.len() - 1].to_string();
}
}
map.insert(key, Value::String(value));
}
(Value::Object(map), body_remainder.to_string())
}
fn opaque_payload(kind: &str, scope: Option<ConfigScope>, path: &Path) -> Value {
let name = truncate_path_for_platform(path);
let mut payload = match kind {
"hooks" => {
let stem = path
.file_stem()
.map(|s| s.to_string_lossy().into_owned())
.unwrap_or_default();
serde_json::json!({
"hooks": {stem: [{"type": "command", "command": name}]},
"path": hooks_display_path(path),
})
}
_ => serde_json::json!({"path": name}),
};
attach_scope(&mut payload, kind, scope);
payload
}
fn truncate_path_for_platform(path: &Path) -> String {
let s = path.to_string_lossy();
let count = s.chars().count();
if count <= PLATFORM_PATH_MAX {
return s.into_owned();
}
s.chars().skip(count - PLATFORM_PATH_MAX).collect()
}
fn unfiltered_payload(
raw: &str,
kind: &str,
scope: Option<ConfigScope>,
path: &Path,
config: &Config,
) -> Value {
if std::env::var("OPENLATCH_TESTING").as_deref() != Ok("true") {
tracing::warn!(
"config_monitor: full_unfiltered content_forward refused (OPENLATCH_TESTING != true)"
);
return filtered_payload(raw, kind, scope, path, &PrivacyFilter::new(&[]), config);
}
let mut payload = serde_json::json!({"kind": kind, "raw": raw});
attach_scope(&mut payload, kind, scope);
payload
}
fn attach_scope(payload: &mut Value, kind: &str, scope: Option<ConfigScope>) {
if kind != "mcp" {
return;
}
if let (Some(scope), Some(obj)) = (scope, payload.as_object_mut()) {
obj.insert(
"scope".to_string(),
Value::String(scope.as_str().to_string()),
);
}
}
#[allow(clippy::too_many_arguments)]
fn build_modified_event(
path: &Path,
kind: &str,
agent_name: &str,
content_hash: &[u8; 32],
path_hash: &[u8; 32],
change: ChangeKind,
source: EventSource,
diffcounter: u64,
severity: Severity,
data: Value,
config: &Config,
) -> CloudEvent {
let path_hash_hex = hex::encode(path_hash);
let content_hash_hex = hex::encode(content_hash);
let (resolved_path, is_symlink) = match path.symlink_metadata() {
Ok(m) if m.file_type().is_symlink() => (std::fs::canonicalize(path).ok(), true),
_ => (None, false),
};
let event_type = format!("ai.openlatch.config.{}", change.type_suffix());
let mut envelope = serde_json::json!({
"specversion": "1.0",
"id": crate::envelope::new_event_id(),
"source": agent_name,
"type": event_type,
"time": crate::envelope::current_timestamp(),
"datacontenttype": "application/json",
"configkind": kind,
"configsource": agent_name,
"configpathhash": path_hash_hex,
"configcontenthash": content_hash_hex,
"diffcounter": diffcounter,
"severityhint": severity.as_str(),
"eventsource": source.as_str(),
"configpath": path.display().to_string(),
"configissymlink": is_symlink,
"data": data,
});
if let Some(rp) = resolved_path {
envelope["configresolvedpath"] = serde_json::json!(crate::path_compat::display_path(&rp));
}
CloudEvent {
envelope,
agent_id: config.agent_id.clone().unwrap_or_default(),
}
}
#[allow(clippy::too_many_arguments)]
fn build_removed_event_for_subpath(
path: &Path,
subpath: Option<&str>,
kind: &str,
agent_name: &str,
source: EventSource,
diffcounter: u64,
severity: Severity,
config: &Config,
) -> Option<CloudEvent> {
let path_hash_hex = hex::encode(subpath_path_hash(path, subpath));
let envelope = serde_json::json!({
"specversion": "1.0",
"id": crate::envelope::new_event_id(),
"source": agent_name,
"type": "ai.openlatch.config.removed",
"time": crate::envelope::current_timestamp(),
"datacontenttype": "application/json",
"configkind": kind,
"configsource": agent_name,
"configpathhash": path_hash_hex,
"diffcounter": diffcounter,
"severityhint": severity.as_str(),
"eventsource": source.as_str(),
"configpath": path.display().to_string(),
"data": Value::Null,
});
Some(CloudEvent {
envelope,
agent_id: config.agent_id.clone().unwrap_or_default(),
})
}
async fn log_and_forward(
event: &CloudEvent,
cloud_tx: Option<&mpsc::Sender<CloudEvent>>,
event_logger: &EventLogger,
) {
if let Ok(line) = serde_json::to_string(&event.envelope) {
event_logger.log_backpressured(line).await;
}
if let Some(tx) = cloud_tx {
if tx.send(event.clone()).await.is_err() {
tracing::warn!(
code = crate::error::ERR_CLOUD_UNREACHABLE,
"config_monitor: cloud channel closed — config event dropped"
);
}
}
}
#[derive(thiserror::Error, Debug)]
pub enum HashError {
#[error("JSON parse failed: {0}")]
JsonParse(#[from] serde_json::Error),
#[error("JCS canonicalization failed: {0}")]
JcsCanonicalize(String),
#[error("JSONC parse failed: {0}")]
JsoncParse(String),
#[error("TOML parse failed: {0}")]
TomlParse(String),
}
fn hash_for_kind(kind: &str, raw: &str, slice_pointers: &[&str]) -> Result<[u8; 32], HashError> {
match kind {
"mcp" => {
if slice_pointers.is_empty() {
content_hash_json(raw)
} else {
content_hash_json_slices(raw, slice_pointers)
}
}
"hooks" => {
if slice_pointers.is_empty() {
content_hash_jsonc(raw).or_else(|_| content_hash_toml(raw))
} else {
content_hash_jsonc_slices(raw, slice_pointers)
}
}
_ => Ok(content_hash_text(raw)),
}
}
pub fn content_hash_json(text: &str) -> Result<[u8; 32], HashError> {
let nfc: String = text.nfc().collect();
let value: Value = serde_json::from_str(&nfc)?;
let canonical = serde_json_canonicalizer::to_string(&value)
.map_err(|e| HashError::JcsCanonicalize(e.to_string()))?;
Ok(sha256(canonical.as_bytes()))
}
pub fn content_hash_json_slices(text: &str, pointers: &[&str]) -> Result<[u8; 32], HashError> {
let nfc: String = text.nfc().collect();
let value: Value = serde_json::from_str(&nfc)?;
hash_slices(&value, pointers)
}
pub fn content_hash_jsonc(text: &str) -> Result<[u8; 32], HashError> {
let nfc: String = text.nfc().collect();
let parsed: Value = jsonc_parser::parse_to_serde_value(&nfc, &Default::default())
.map_err(|e| HashError::JsoncParse(e.to_string()))?;
let canonical = serde_json_canonicalizer::to_string(&parsed)
.map_err(|e| HashError::JcsCanonicalize(e.to_string()))?;
Ok(sha256(canonical.as_bytes()))
}
pub fn content_hash_toml(text: &str) -> Result<[u8; 32], HashError> {
let nfc: String = text.nfc().collect();
let parsed: Value = toml::from_str(&nfc).map_err(|e| HashError::TomlParse(e.to_string()))?;
let canonical = serde_json_canonicalizer::to_string(&parsed)
.map_err(|e| HashError::JcsCanonicalize(e.to_string()))?;
Ok(sha256(canonical.as_bytes()))
}
pub fn content_hash_jsonc_slices(text: &str, pointers: &[&str]) -> Result<[u8; 32], HashError> {
let nfc: String = text.nfc().collect();
let parsed: Value = jsonc_parser::parse_to_serde_value(&nfc, &Default::default())
.map_err(|e| HashError::JsoncParse(e.to_string()))?;
hash_slices(&parsed, pointers)
}
fn hash_slices(root: &Value, pointers: &[&str]) -> Result<[u8; 32], HashError> {
let subtrees: Vec<Value> = pointers
.iter()
.map(|ptr| root.pointer(ptr).cloned().unwrap_or(Value::Null))
.collect();
let synthetic = Value::Array(subtrees);
let canonical = serde_json_canonicalizer::to_string(&synthetic)
.map_err(|e| HashError::JcsCanonicalize(e.to_string()))?;
Ok(sha256(canonical.as_bytes()))
}
pub fn content_hash_text(text: &str) -> [u8; 32] {
let nfc: String = text.nfc().collect();
let normalized_eol = nfc.replace("\r\n", "\n").replace('\r', "\n");
let lines: Vec<&str> = normalized_eol.split('\n').map(|l| l.trim_end()).collect();
let normalized = lines.join("\n");
sha256(normalized.as_bytes())
}
fn sha256(bytes: &[u8]) -> [u8; 32] {
Sha256::digest(bytes).into()
}
pub fn severity_for(kind: &str, change_type: &str) -> Severity {
match (kind, change_type) {
("mcp", "added") => Severity::Critical,
("mcp", "modified") => Severity::High,
("mcp", "removed") => Severity::Medium,
("skill", "added") => Severity::High,
("skill", "modified") => Severity::High,
("skill", "removed") => Severity::Low,
("hooks", "added") => Severity::Critical,
("hooks", "modified") => Severity::Critical,
("hooks", "removed") => Severity::High,
("command", "added") => Severity::High,
("command", "modified") => Severity::Medium,
("command", "removed") => Severity::Low,
("rules", "added") => Severity::Medium,
("rules", "modified") => Severity::Medium,
("rules", "removed") => Severity::Info,
_ => Severity::Low,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn content_hash_json_canonicalizes_keys() {
let a = r#"{"b":1,"a":2}"#;
let b = r#"{"a":2,"b":1}"#;
assert_eq!(content_hash_json(a).unwrap(), content_hash_json(b).unwrap());
}
#[test]
fn content_hash_text_normalizes_eol() {
let a = "line1\nline2\n";
let b = "line1\r\nline2\r\n";
let c = "line1\rline2\r";
assert_eq!(content_hash_text(a), content_hash_text(b));
assert_eq!(content_hash_text(a), content_hash_text(c));
}
#[test]
fn content_hash_text_strips_trailing_whitespace_per_line() {
let a = "line1\nline2";
let b = "line1 \nline2\t";
assert_eq!(content_hash_text(a), content_hash_text(b));
}
#[test]
fn content_hash_jsonc_strips_comments() {
let with_comments = r#"{
// comment
"a": 1
}"#;
let plain = r#"{"a":1}"#;
assert_eq!(
content_hash_jsonc(with_comments).unwrap(),
content_hash_json(plain).unwrap()
);
}
#[test]
fn content_hash_json_slices_ignores_surrounding_keys() {
let with_theme = r#"{"theme":"dark","hooks":{"PreToolUse":[]}}"#;
let without_theme = r#"{"hooks":{"PreToolUse":[]}}"#;
let pointers = ["/hooks"];
assert_eq!(
content_hash_json_slices(with_theme, &pointers).unwrap(),
content_hash_json_slices(without_theme, &pointers).unwrap(),
);
}
#[test]
fn content_hash_json_slices_detects_slice_edits() {
let before = r#"{"hooks":{"PreToolUse":[]}}"#;
let after = r#"{"hooks":{"PreToolUse":[{"matcher":"*"}]}}"#;
let pointers = ["/hooks"];
assert_ne!(
content_hash_json_slices(before, &pointers).unwrap(),
content_hash_json_slices(after, &pointers).unwrap(),
);
}
#[test]
fn content_hash_jsonc_slices_strips_comments_and_slices() {
let jsonc_with_comments = r#"{
// ignore
"theme": "dark",
"hooks": { "PreToolUse": [] }
}"#;
let plain = r#"{"hooks":{"PreToolUse":[]}}"#;
let pointers = ["/hooks"];
assert_eq!(
content_hash_jsonc_slices(jsonc_with_comments, &pointers).unwrap(),
content_hash_json_slices(plain, &pointers).unwrap(),
);
}
#[test]
fn content_hash_json_slices_missing_pointer_is_deterministic() {
let a = r#"{"other":1}"#;
let b = r#"{"other":2}"#;
let pointers = ["/hooks"];
assert_eq!(
content_hash_json_slices(a, &pointers).unwrap(),
content_hash_json_slices(b, &pointers).unwrap(),
);
}
#[test]
fn content_hash_json_slices_multi_pointer_order_matters() {
let raw = r#"{"a":1,"b":2,"c":3}"#;
let forward = ["/a", "/b"];
let reversed = ["/b", "/a"];
assert_ne!(
content_hash_json_slices(raw, &forward).unwrap(),
content_hash_json_slices(raw, &reversed).unwrap(),
);
}
#[test]
fn severity_for_covers_known_pairs() {
assert!(matches!(severity_for("mcp", "added"), Severity::Critical));
assert!(matches!(severity_for("mcp", "modified"), Severity::High));
assert!(matches!(severity_for("rules", "removed"), Severity::Info));
assert!(matches!(severity_for("unknown", "added"), Severity::Low));
}
#[test]
fn change_kind_type_suffix_matches_namespace() {
assert_eq!(ChangeKind::Snapshot.type_suffix(), "snapshot");
assert_eq!(ChangeKind::Added.type_suffix(), "added");
assert_eq!(ChangeKind::Modified.type_suffix(), "modified");
assert_eq!(ChangeKind::Removed.type_suffix(), "removed");
}
#[test]
fn mcp_payload_omits_json_blob_to_match_platform_shape() {
let cfg = Config::defaults();
let filter = PrivacyFilter::new(&[]);
let raw = r#"{"mcpServers":{"a":{"command":"x"}}}"#;
let path = Path::new("/etc/claude/mcp.json");
let payload =
filtered_payload(raw, "mcp", Some(ConfigScope::Personal), path, &filter, &cfg);
assert!(
payload.get("kind").is_none(),
"platform forbids extra `kind` field"
);
assert!(payload.get("json").is_none());
assert!(payload.get("server_name").is_some());
assert_eq!(payload["scope"], "personal");
}
const CODEX_CONFIG_WITH_SECRETS: &str = r#"# my config, hand written
model = "gpt-5-codex" # trailing comment
model_provider = "openlatch"
[mcp_servers.ctx7]
command = "npx"
bearer_token = "SECRET-BEARER"
env = { API_KEY = "SECRET-ENV" }
[model_providers.openlatch]
name = "OpenLatch model relay"
base_url = "http://127.0.0.1:7600/v1"
http_headers = { Authorization = "SECRET-HEADER" }
[hooks.state."hooks.json:pre_tool_use:0:0"]
trusted_hash = "abc"
enabled = true
"#;
#[test]
fn toml_hooks_body_forwards_the_hooks_table_and_nothing_else() {
let filter = PrivacyFilter::new(&[]);
let body = build_hooks_body(
CODEX_CONFIG_WITH_SECRETS,
Path::new("/x/.codex/config.toml"),
&filter,
);
let hooks = body["hooks"].as_object().expect("the `hooks` wrapper");
assert_eq!(hooks.keys().collect::<Vec<_>>(), ["state"], "{body}");
assert_eq!(
hooks["state"]["hooks.json:pre_tool_use:0:0"]["trusted_hash"],
"abc"
);
let rendered = body.to_string();
for forbidden in [
"SECRET",
"mcp_servers",
"env",
"bearer_token",
"model_providers",
] {
assert!(!rendered.contains(forbidden), "{forbidden} in {rendered}");
}
}
#[test]
fn toml_hooks_body_without_a_hooks_table_is_empty() {
let raw = "model_provider = \"openlatch\"\n[mcp_servers.x]\nenv = { K = \"SECRET\" }\n";
let body = build_hooks_body(raw, Path::new("/x/config.toml"), &PrivacyFilter::new(&[]));
assert_eq!(body["hooks"], serde_json::json!({}));
assert!(!body.to_string().contains("SECRET"));
}
#[test]
fn codex_config_hooks_event_carries_no_secret_and_a_home_relative_path() {
let home = dirs::home_dir().expect("home dir");
let path = home.join(".codex").join("config.toml");
let ap = AgentPath {
kind: "hooks".to_string(),
scope: None,
paths: vec!["${CODEX_HOME}/config.toml".to_string()],
paths_relative: Vec::new(),
paths_glob: None,
paths_glob_relative: Vec::new(),
json_slice_paths: Vec::new(),
watch_strategy: WatchStrategy::ExactFile,
slice_kind_subpath: false,
opaque: false,
projected: false,
claude_plugins: false,
};
let mut cfg = Config::defaults();
cfg.inventory_monitor.content_forward = ContentForwardMode::Filtered;
let events = build_event_from_content(
&path,
&ap,
"codex-cli",
ChangeKind::Snapshot,
EventSource::InitScan,
0,
CODEX_CONFIG_WITH_SECRETS,
&PrivacyFilter::new(&[]),
&ContentHashCache::new(8),
&cfg,
);
assert_eq!(events.len(), 1);
let data = &events[0].0.envelope["data"];
let rendered = data.to_string();
for forbidden in ["SECRET", "mcp_servers", "bearer_token", "\"env\""] {
assert!(!rendered.contains(forbidden), "{forbidden} in {rendered}");
}
let shown = data["path"].as_str().expect("hooks data carries `path`");
assert!(shown.starts_with('~'), "{shown}");
assert!(shown.ends_with("config.toml"), "{shown}");
assert!(
!shown.contains(&*home.to_string_lossy()),
"never the absolute home: {shown}"
);
}
#[test]
fn hooks_display_path_writes_home_as_a_tilde() {
let home = dirs::home_dir().expect("home dir");
let shown = hooks_display_path(&home.join(".claude").join("settings.json"));
assert_eq!(
shown,
format!(
"~{}",
home.join(".claude")
.join("settings.json")
.to_string_lossy()
.strip_prefix(&*home.to_string_lossy())
.unwrap()
)
);
assert_eq!(
hooks_display_path(Path::new("/etc/x/hooks.json")),
"/etc/x/hooks.json"
);
}
#[test]
fn toml_hooks_body_degrades_on_malformed_toml() {
let filter = PrivacyFilter::new(&[]);
let body = build_hooks_body(
"[unclosed\nmodel = \n",
Path::new("/x/config.toml"),
&filter,
);
assert_eq!(
body,
serde_json::json!({"hooks": {}, "path": "/x/config.toml"}),
"unparseable input degrades, it does not panic and does not invent a body"
);
}
#[test]
fn toml_hooks_file_hashes_instead_of_dropping_the_event() {
let raw = "# my config, hand written\n\
model = \"gpt-5-codex\"\n\n\
[model_providers.openlatch]\n\
base_url = \"http://127.0.0.1:7600/v1\"\n";
let hash = hash_for_kind("hooks", raw, &[]).expect(
"a TOML hooks file must hash — an Err here drops the config event before any body \
builder runs",
);
let reformatted = "model=\"gpt-5-codex\"\n\
[model_providers.openlatch]\n\
base_url=\"http://127.0.0.1:7600/v1\"\n";
assert_eq!(
hash,
hash_for_kind("hooks", reformatted, &[]).expect("same document, reflowed"),
"comments and spacing are not content"
);
let edited = raw.replace("127.0.0.1:7600", "127.0.0.1:7601");
assert_ne!(
hash,
hash_for_kind("hooks", &edited, &[]).expect("still valid TOML"),
"a rerouted model relay base URL is exactly the change this file is watched for"
);
assert_eq!(
hash_for_kind("hooks", r#"{"hooks":{}} "#, &[]).expect("jsonc still hashes"),
content_hash_jsonc(r#"{"hooks":{}} "#).expect("jsonc leaf"),
);
assert!(hash_for_kind("hooks", "[unclosed\nmodel = \n", &[]).is_err());
}
fn opaque_hook_event(path: &Path, raw: &str, mode: ContentForwardMode) -> CloudEvent {
let mut cfg = Config::defaults();
cfg.inventory_monitor.content_forward = mode;
let ap = AgentPath {
kind: "hooks".to_string(),
scope: None,
paths: Vec::new(),
paths_relative: Vec::new(),
paths_glob: Some("${CLINE_ASSETS_DIR}/Hooks/*".to_string()),
paths_glob_relative: Vec::new(),
json_slice_paths: Vec::new(),
watch_strategy: WatchStrategy::Glob,
slice_kind_subpath: false,
opaque: true,
projected: false,
claude_plugins: false,
};
let mut events = build_event_from_content(
path,
&ap,
"cline",
ChangeKind::Snapshot,
EventSource::InitScan,
0,
raw,
&PrivacyFilter::new(&[]),
&ContentHashCache::new(8),
&cfg,
);
assert_eq!(
events.len(),
1,
"a script is not JSON, and must not be dropped"
);
let (event, hash, _) = events.remove(0);
assert_eq!(hash, content_hash_text(raw));
event
}
#[test]
fn opaque_hook_script_forwards_its_name_never_its_body() {
let raw = "#!/bin/sh\nSECRET=abc\nexec openlatch-hook \"$@\"\n";
let path = Path::new("/home/u/Documents/Cline/Hooks/PreToolUse");
let _testing =
crate::hooks::cline::EnvOverride::apply([("OPENLATCH_TESTING", Some("true".into()))]);
for mode in [
ContentForwardMode::Filtered,
ContentForwardMode::FullUnfiltered,
] {
let event = opaque_hook_event(path, raw, mode);
assert_eq!(
event.envelope["data"],
serde_json::json!({
"hooks": {"PreToolUse": [
{"type": "command", "command": path.display().to_string()}
]},
"path": path.display().to_string(),
}),
"{mode:?}"
);
assert_eq!(event.envelope["configkind"], "hooks");
let wire = serde_json::to_string(&event.envelope).unwrap();
assert!(
!wire.contains("abc"),
"{mode:?} leaked the script body: {wire}"
);
}
let event = opaque_hook_event(path, raw, ContentForwardMode::HashOnly);
assert_eq!(event.envelope["data"], Value::Null);
}
#[test]
fn opaque_hook_name_drops_the_extension() {
let path = Path::new("C:/Users/u/Documents/Cline/Hooks/PreToolUse.ps1");
let event = opaque_hook_event(path, "Write-Output 1", ContentForwardMode::Filtered);
assert!(event.envelope["data"]["hooks"]["PreToolUse"].is_array());
}
#[test]
fn hash_mcp_server_entry_distinguishes_servers_by_name() {
let raw = r#"{"context7":{"command":"npx"},"github":{"command":"docker"}}"#;
let parsed: Value = serde_json::from_str(raw).unwrap();
let servers = parsed.as_object().unwrap();
let h_ctx7 = hash_mcp_server_entry("context7", &servers["context7"]);
let h_gh = hash_mcp_server_entry("github", &servers["github"]);
assert_ne!(h_ctx7, h_gh);
}
#[test]
fn a_projected_toml_mcp_file_reads_mcp_servers_without_secrets() {
let raw = r#"
model = "gpt-5-codex"
model_provider = "openlatch"
[mcp_servers.ctx7]
command = "npx"
args = ["-y", "@upstash/context7-mcp"]
env = { API_KEY = "sk-SECRET" }
[mcp_servers.github]
url = "https://example.test/mcp"
bearer_token = "SECRET-TOKEN"
http_headers = { Authorization = "Bearer SECRET-HEADER" }
"#;
let projected = project_mcp_settings(raw).unwrap();
let rendered = projected.to_string();
assert!(!rendered.contains("SECRET"), "{rendered}");
let servers = projected["mcpServers"].as_object().unwrap();
assert_eq!(servers.keys().collect::<Vec<_>>(), ["ctx7", "github"]);
assert_eq!(servers["ctx7"]["command"], "npx");
assert_eq!(servers["github"]["url"], "https://example.test/mcp");
assert!(
project_mcp_settings("model = ").is_err(),
"neither JSON nor TOML"
);
}
#[test]
fn claims_for_path_is_one_claim_per_declaring_agent_and_kind() {
let manifest: Manifest = toml::from_str(
r#"
[[agent]]
name = "codex-cli"
[[agent.path]]
kind = "rules"
paths = ["/shared/AGENTS.md"]
watch_strategy = "exact_file"
[[agent.path]]
kind = "hooks"
paths = ["/codex/config.toml"]
watch_strategy = "exact_file"
[[agent.path]]
kind = "mcp"
paths = ["/codex/config.toml"]
watch_strategy = "exact_file"
[[agent]]
name = "cline"
[[agent.path]]
kind = "rules"
paths = ["/shared/AGENTS.md"]
watch_strategy = "exact_file"
"#,
)
.unwrap();
let owners = |m: &Manifest, path: &str| {
claims_for_path(Path::new(path), m)
.into_iter()
.map(|(agent, ap)| (agent.to_string(), ap.kind.clone()))
.collect::<Vec<_>>()
};
let pair = |a: &str, k: &str| (a.to_string(), k.to_string());
assert_eq!(
owners(&manifest, "/shared/AGENTS.md"),
[pair("codex-cli", "rules"), pair("cline", "rules")]
);
assert_eq!(
owners(&manifest, "/codex/config.toml"),
[pair("codex-cli", "hooks"), pair("codex-cli", "mcp")]
);
let cline_only = manifest.only_agents(&HashSet::from(["cline".to_string()]));
assert_eq!(
owners(&cline_only, "/shared/AGENTS.md"),
[pair("cline", "rules")]
);
assert!(owners(&manifest, "/elsewhere.md").is_empty());
}
#[test]
fn a_projected_mcp_file_never_materialises_env_or_headers() {
let settings = |secret: &str, command: &str| {
serde_json::json!({
"mcpServers": {
"github": {
"command": command,
"args": ["-y", "server-github"],
"env": {"GITHUB_TOKEN": secret},
"headers": {"Authorization": format!("Bearer {secret}")},
"disabled": false
}
},
"otherTopLevel": {"token": secret}
})
.to_string()
};
let projected = project_mcp_settings(&settings("ghp_SECRET", "npx")).unwrap();
let rendered = projected.to_string();
assert!(!rendered.contains("ghp_SECRET"), "{rendered}");
assert!(projected["mcpServers"]["github"].get("env").is_none());
assert!(projected["mcpServers"]["github"].get("headers").is_none());
assert_eq!(projected["mcpServers"]["github"]["command"], "npx");
let ap = AgentPath {
kind: "mcp".to_string(),
scope: Some(ConfigScope::Personal),
paths: vec!["${CLINE_DATA_DIR}/settings/cline_mcp_settings.json".to_string()],
paths_relative: Vec::new(),
paths_glob: None,
paths_glob_relative: Vec::new(),
json_slice_paths: Vec::new(),
watch_strategy: WatchStrategy::ExactFile,
slice_kind_subpath: true,
opaque: false,
projected: true,
claude_plugins: false,
};
let path = Path::new("/home/u/.cline/data/settings/cline_mcp_settings.json");
let events_for = |raw: &str| {
let mut cfg = Config::defaults();
cfg.inventory_monitor.content_forward = ContentForwardMode::FullUnfiltered;
build_event_from_content(
path,
&ap,
"cline",
ChangeKind::Snapshot,
EventSource::InitScan,
0,
raw,
&PrivacyFilter::new(&[]),
&ContentHashCache::new(8),
&cfg,
)
};
let first = events_for(&settings("ghp_SECRET", "npx"));
assert_eq!(first.len(), 1);
let body = first[0].0.envelope.to_string();
assert!(!body.contains("ghp_SECRET"), "{body}");
let rotated = events_for(&settings("ghp_ROTATED", "npx"));
assert_eq!(first[0].1, rotated[0].1, "a secret is not in the hash");
let edited = events_for(&settings("ghp_SECRET", "docker"));
assert_ne!(first[0].1, edited[0].1, "the server itself is");
}
#[test]
fn mcp_fanout_emits_one_envelope_per_server_with_distinct_hashes() {
let cache = ContentHashCache::new(64);
let cfg = Config::defaults();
let filter = PrivacyFilter::new(&[]);
let raw = r#"{"context7":{"command":"npx"},"github":{"command":"docker"}}"#;
let path = Path::new("/home/u/.claude/plugins/cache/x/y/z/.mcp.json");
let ap = AgentPath {
kind: "mcp".to_string(),
scope: Some(ConfigScope::Personal),
paths: Vec::new(),
paths_relative: Vec::new(),
paths_glob: None,
paths_glob_relative: Vec::new(),
json_slice_paths: Vec::new(),
watch_strategy: WatchStrategy::Glob,
slice_kind_subpath: true,
opaque: false,
projected: false,
claude_plugins: false,
};
let events = build_event_from_content(
path,
&ap,
"claude-code",
ChangeKind::Snapshot,
EventSource::InitScan,
0,
raw,
&filter,
&cache,
&cfg,
);
assert_eq!(events.len(), 2);
let mut subpaths: Vec<String> = events
.iter()
.map(|(_, _, s)| s.clone().expect("fan-out subpath set"))
.collect();
subpaths.sort();
assert_eq!(subpaths, vec!["context7", "github"]);
let path_hashes: HashSet<_> = events
.iter()
.map(|(e, _, _)| e.envelope["configpathhash"].as_str().unwrap().to_string())
.collect();
assert_eq!(path_hashes.len(), 2, "each server gets a unique path_hash");
let cached = cache.entries_for_path(path, "claude-code", "mcp");
assert_eq!(cached.len(), 2);
for (event, _, sub) in &events {
let data = &event.envelope["data"];
assert_eq!(data["server_name"].as_str(), sub.as_deref());
assert_eq!(data["tools"].as_array().unwrap().len(), 0);
assert_eq!(data["scope"], "personal");
}
}
#[test]
fn mcp_fanout_ignores_non_slice_keys_when_manifest_declares_pointers() {
let raw = r#"{
"numStartups": 42,
"userID": "abc",
"tipsHistory": {"x": 1},
"projects": {
"/repo/a": {"mcpServers": {"github": {"command": "docker"}}},
"/repo/b": {"mcpServers": {}}
}
}"#;
let parsed: Value = serde_json::from_str(raw).unwrap();
let sliced = resolve_mcp_servers(&parsed, &["/mcpServers", "/projects"]);
let mut names: Vec<&String> = sliced.keys().collect();
names.sort();
assert_eq!(names, vec!["/repo/a::github"]);
let unsliced = resolve_mcp_servers(&parsed, &[]);
assert_eq!(
unsliced.len(),
4,
"with no declared slices the whole document is still the server map"
);
}
#[test]
fn mcp_fanout_reads_top_level_mcpservers_slice() {
let raw = r#"{"mcpServers":{"context7":{"command":"npx"}},"userID":"abc"}"#;
let parsed: Value = serde_json::from_str(raw).unwrap();
let servers = resolve_mcp_servers(&parsed, &["/mcpServers", "/projects"]);
assert_eq!(servers.keys().collect::<Vec<_>>(), vec!["context7"]);
}
#[test]
fn rules_payload_matches_platform_shape() {
let cfg = Config::defaults();
let filter = PrivacyFilter::new(&[]);
let raw = "# CLAUDE.md\n\nBe terse.";
let path = Path::new("/repo/CLAUDE.md");
let payload = filtered_payload(
raw,
"rules",
Some(ConfigScope::Project),
path,
&filter,
&cfg,
);
assert_eq!(payload["path"], "/repo/CLAUDE.md");
assert_eq!(payload["body"], raw);
assert!(payload["frontmatter"].is_object());
assert!(
payload.get("kind").is_none(),
"extra `kind` field would 400 the envelope"
);
assert!(payload.get("text").is_none());
assert!(payload.get("scope").is_none(), "scope only attaches to mcp");
}
#[test]
fn skill_payload_derives_name_from_parent_dir_for_skill_md() {
let cfg = Config::defaults();
let filter = PrivacyFilter::new(&[]);
let raw = "Skill body content.";
let path = Path::new("/home/u/.claude/skills/code-review/SKILL.md");
let payload = filtered_payload(raw, "skill", None, path, &filter, &cfg);
assert_eq!(payload["name"], "code-review");
assert_eq!(payload["body"], raw);
assert_eq!(payload["description"], "");
assert!(payload["frontmatter"].is_object());
}
#[test]
fn skill_payload_pulls_name_and_description_from_frontmatter() {
let cfg = Config::defaults();
let filter = PrivacyFilter::new(&[]);
let raw =
"---\nname: playwright-cli\ndescription: \"Browser automation skill\"\n---\n\n# Body";
let path = Path::new("/home/u/.claude/skills/SKILL.md");
let payload = filtered_payload(raw, "skill", None, path, &filter, &cfg);
assert_eq!(payload["name"], "playwright-cli");
assert_eq!(payload["description"], "Browser automation skill");
let fm = payload["frontmatter"].as_object().unwrap();
assert_eq!(fm["name"], "playwright-cli");
assert_eq!(fm["description"], "Browser automation skill");
let body = payload["body"].as_str().unwrap();
assert!(body.starts_with("# Body"));
assert!(!body.contains("---"));
}
#[test]
fn command_payload_derives_name_from_file_stem() {
let cfg = Config::defaults();
let filter = PrivacyFilter::new(&[]);
let raw = "Run the command.";
let path = Path::new("/home/u/.claude/commands/security-review.md");
let payload = filtered_payload(raw, "command", None, path, &filter, &cfg);
assert_eq!(payload["name"], "security-review");
assert_eq!(payload["body"], raw);
}
#[test]
fn hooks_payload_lifts_hooks_slice_from_settings_json() {
let cfg = Config::defaults();
let filter = PrivacyFilter::new(&[]);
let raw = r#"{
"theme": "dark",
"hooks": {"PreToolUse": [{"matcher": "*"}]}
}"#;
let path = Path::new("/home/u/.claude/settings.json");
let payload = filtered_payload(raw, "hooks", None, path, &filter, &cfg);
assert!(payload.get("kind").is_none());
assert!(payload.get("text").is_none());
let hooks = payload["hooks"].as_object().unwrap();
assert!(hooks.contains_key("PreToolUse"));
assert!(!hooks.contains_key("theme"));
}
#[test]
fn hooks_payload_uses_root_when_no_hooks_slice() {
let cfg = Config::defaults();
let filter = PrivacyFilter::new(&[]);
let raw = r#"{"SessionStart": [{"matcher": "*"}]}"#;
let path = Path::new("/home/u/.claude/plugins/cache/x/y/z/hooks/hooks.json");
let payload = filtered_payload(raw, "hooks", None, path, &filter, &cfg);
let hooks = payload["hooks"].as_object().unwrap();
assert!(hooks.contains_key("SessionStart"));
}
#[test]
fn rules_truncated_payload_keeps_shape() {
let mut cfg = Config::defaults();
cfg.inventory_monitor.max_inline_content_bytes = 256;
let filter = PrivacyFilter::new(&[]);
let raw = "x".repeat(10_000);
let path = Path::new("/repo/CLAUDE.md");
let payload = filtered_payload(&raw, "rules", None, path, &filter, &cfg);
assert!(payload.get("kind").is_none());
assert!(payload.get("truncated").is_none());
let body = payload["body"].as_str().expect("body is string");
assert!(
body.len() <= 256 + 128,
"body should be capped near inline_cap"
);
assert!(
body.contains("…<truncated"),
"marker should appear in truncated body"
);
}
#[test]
fn build_modified_event_uses_agent_name_as_source() {
let cfg = Config::defaults();
let path = Path::new("/etc/claude/mcp.json");
let content_hash = [0u8; 32];
let path_hash = [0u8; 32];
let event = build_modified_event(
path,
"mcp",
"claude-code",
&content_hash,
&path_hash,
ChangeKind::Modified,
EventSource::FsWatcher,
1,
Severity::High,
Value::Null,
&cfg,
);
assert_eq!(event.envelope["source"].as_str(), Some("claude-code"));
assert_eq!(event.envelope["configsource"].as_str(), Some("claude-code"));
}
#[test]
fn build_removed_event_uses_agent_name_as_source() {
let cfg = Config::defaults();
let path = Path::new("/home/u/.codex/config.toml");
let event = build_removed_event_for_subpath(
path,
None,
"mcp",
"codex-cli",
EventSource::FsWatcher,
2,
Severity::Critical,
&cfg,
)
.expect("removed event must be built");
assert_eq!(event.envelope["source"].as_str(), Some("codex-cli"));
assert_eq!(event.envelope["configsource"].as_str(), Some("codex-cli"));
}
}