pub mod alerts;
pub mod cache;
pub mod enrich;
pub mod manifest;
pub mod monitor;
pub mod watcher;
use std::future::Future;
use std::sync::{Arc, Mutex, RwLock};
use std::time::Duration;
use serde::Serialize;
use tokio::sync::{mpsc, watch};
pub use alerts::{run_long_poll, PendingAlert, PendingAlerts};
pub use cache::{CacheEntry, ContentHashCache};
pub use enrich::{
enrich_session_start, try_register_project_from_cwd, DeclaredSource, EnrichError,
};
pub use manifest::{
AgentManifest, AgentPath, ConfigScope, JsonSlicePath, Manifest, ManifestError, WatchStrategy,
};
pub use monitor::{ChangeKind, ConfigChangeRequest, EventSource, Severity};
use crate::cloud::CloudEvent;
use crate::config::Config;
use crate::core::logging::EventLogger;
use crate::core::supervision::task::{TaskHealth, TaskState};
use crate::privacy::PrivacyFilter;
const REQUEST_CHANNEL_SIZE: usize = 256;
const SHUTDOWN_DRAIN: Duration = Duration::from_secs(4);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum ConfigMonitorState {
Disabled,
Pending,
Running,
Failed,
Stopped,
}
#[derive(Debug, Clone, Serialize)]
pub struct ConfigMonitorSnapshot {
pub state: ConfigMonitorState,
pub manifest_loaded: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
pub struct ConfigMonitorRuntime {
snapshot: RwLock<ConfigMonitorSnapshot>,
request_tx: Mutex<Option<mpsc::Sender<ConfigChangeRequest>>>,
}
impl ConfigMonitorRuntime {
pub fn disabled() -> Arc<Self> {
Arc::new(Self {
snapshot: RwLock::new(ConfigMonitorSnapshot {
state: ConfigMonitorState::Disabled,
manifest_loaded: false,
error: None,
}),
request_tx: Mutex::new(None),
})
}
pub fn pending() -> (Arc<Self>, mpsc::Receiver<ConfigChangeRequest>) {
let (request_tx, request_rx) = mpsc::channel(REQUEST_CHANNEL_SIZE);
let _ = request_tx.try_send(ConfigChangeRequest::InitialInventory);
(
Arc::new(Self {
snapshot: RwLock::new(ConfigMonitorSnapshot {
state: ConfigMonitorState::Pending,
manifest_loaded: false,
error: None,
}),
request_tx: Mutex::new(Some(request_tx)),
}),
request_rx,
)
}
pub fn snapshot(&self) -> ConfigMonitorSnapshot {
self.snapshot
.read()
.map(|snapshot| snapshot.clone())
.unwrap_or(ConfigMonitorSnapshot {
state: ConfigMonitorState::Failed,
manifest_loaded: false,
error: Some("monitor state unavailable".to_string()),
})
}
pub fn request_tx(&self) -> Option<mpsc::Sender<ConfigChangeRequest>> {
self.request_tx.lock().ok().and_then(|tx| tx.clone())
}
pub fn mark_manifest_loaded(&self) {
if let Ok(mut snapshot) = self.snapshot.write() {
snapshot.manifest_loaded = true;
}
}
fn mark_running(&self) {
if let Ok(mut snapshot) = self.snapshot.write() {
snapshot.state = ConfigMonitorState::Running;
snapshot.error = None;
}
}
fn mark_failed(&self, error: String) {
if let Ok(mut snapshot) = self.snapshot.write() {
snapshot.state = ConfigMonitorState::Failed;
snapshot.error = Some(error);
}
self.clear_sender();
}
fn mark_stopped(&self) {
if let Ok(mut snapshot) = self.snapshot.write() {
snapshot.state = ConfigMonitorState::Stopped;
}
self.clear_sender();
}
fn clear_sender(&self) {
if let Ok(mut tx) = self.request_tx.lock() {
*tx = None;
}
}
}
fn scrub_error(error: &str, privacy_filter: &PrivacyFilter) -> String {
let mut value = serde_json::Value::String(error.to_string());
crate::privacy::filter_value(&mut value, privacy_filter);
value.as_str().unwrap_or_default().to_string()
}
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;
}
}
}
pub(crate) async fn run_startup<I, IFut, H, R, RFut>(
runtime: Arc<ConfigMonitorRuntime>,
health: Arc<TaskHealth>,
privacy_filter: PrivacyFilter,
mut shutdown: watch::Receiver<bool>,
initialize: I,
run: R,
) where
I: FnOnce() -> IFut + Send + 'static,
IFut: Future<Output = Result<H, String>> + Send,
H: Send + 'static,
R: FnOnce(H, watch::Receiver<bool>) -> RFut + Send + 'static,
RFut: Future<Output = Result<(), String>> + Send,
{
if *shutdown.borrow() {
runtime.mark_stopped();
health.set_state(TaskState::Stopped);
return;
}
let initialized = tokio::select! {
result = initialize() => result,
_ = wait_for_shutdown(&mut shutdown) => {
runtime.mark_stopped();
health.set_state(TaskState::Stopped);
return;
}
};
let handle = match initialized {
Ok(handle) => handle,
Err(error) => {
let safe_error = scrub_error(&error, &privacy_filter);
runtime.mark_failed(safe_error.clone());
health.mark_failed(&safe_error);
tracing::error!(
code = crate::error::ERR_INVENTORY_INIT_FAILED,
error = %error,
"config monitor failed to start"
);
return;
}
};
if *shutdown.borrow() {
runtime.mark_stopped();
health.set_state(TaskState::Stopped);
return;
}
runtime.mark_running();
health.set_state(TaskState::Running);
tracing::info!("config monitor active");
if let Err(error) = run(handle, shutdown.clone()).await {
let safe_error = scrub_error(&error, &privacy_filter);
runtime.mark_failed(safe_error.clone());
health.mark_failed(&safe_error);
tracing::error!(
code = crate::error::ERR_INVENTORY_INIT_FAILED,
error = %error,
"config monitor stopped unexpectedly"
);
} else {
runtime.mark_stopped();
health.set_state(TaskState::Stopped);
}
}
pub struct ConfigMonitor {
manifest: Arc<Manifest>,
cache: Arc<ContentHashCache>,
privacy_filter: PrivacyFilter,
cloud_tx: Option<mpsc::Sender<CloudEvent>>,
event_logger: EventLogger,
config: Arc<Config>,
}
impl ConfigMonitor {
pub fn new(
manifest: Arc<Manifest>,
cache: Arc<ContentHashCache>,
privacy_filter: PrivacyFilter,
cloud_tx: Option<mpsc::Sender<CloudEvent>>,
event_logger: EventLogger,
config: Arc<Config>,
) -> Self {
Self {
manifest,
cache,
privacy_filter,
cloud_tx,
event_logger,
config,
}
}
pub async fn spawn(self) -> anyhow::Result<ConfigMonitorHandle> {
let (request_tx, request_rx) = mpsc::channel::<ConfigChangeRequest>(REQUEST_CHANNEL_SIZE);
request_tx
.send(ConfigChangeRequest::InitialInventory)
.await
.map_err(|error| anyhow::anyhow!("failed to queue initial inventory: {error}"))?;
self.spawn_with_requests(request_tx, request_rx).await
}
pub async fn spawn_with_requests(
self,
request_tx: mpsc::Sender<ConfigChangeRequest>,
request_rx: mpsc::Receiver<ConfigChangeRequest>,
) -> anyhow::Result<ConfigMonitorHandle> {
let watcher_manifest = self.manifest.clone();
let watcher_tx = request_tx.clone();
let debounce_ms = self.config.inventory_monitor.watcher_debounce_ms;
let watchers = tokio::task::spawn_blocking(move || {
watcher::spawn_watchers(&watcher_manifest, watcher_tx, debounce_ms)
})
.await
.map_err(|error| anyhow::anyhow!("watcher initialization task failed: {error}"))??;
let (shutdown_tx, shutdown_rx) = watch::channel(false);
let join = tokio::spawn(monitor::run(
self.manifest,
self.cache,
self.privacy_filter,
self.cloud_tx,
self.event_logger,
self.config,
request_rx,
request_tx.clone(),
shutdown_rx,
));
Ok(ConfigMonitorHandle {
join,
request_tx,
shutdown_tx,
_watchers: watchers,
})
}
}
pub struct ConfigMonitorHandle {
join: tokio::task::JoinHandle<()>,
pub request_tx: mpsc::Sender<ConfigChangeRequest>,
shutdown_tx: watch::Sender<bool>,
_watchers: Vec<watcher::WatcherGuard>,
}
impl ConfigMonitorHandle {
pub async fn run_until_shutdown(
mut self,
mut shutdown: watch::Receiver<bool>,
) -> Result<(), String> {
tokio::select! {
joined = &mut self.join => {
match joined {
Ok(()) => Err("monitor loop exited before daemon shutdown".to_string()),
Err(error) => Err(format!("monitor task failed: {error}")),
}
}
_ = wait_for_shutdown(&mut shutdown) => {
let _ = self.shutdown_tx.send(true);
if tokio::time::timeout(SHUTDOWN_DRAIN, &mut self.join).await.is_err() {
self.join.abort();
return Err("monitor did not stop within the shutdown drain window".to_string());
}
Ok(())
}
}
}
}
impl Drop for ConfigMonitorHandle {
fn drop(&mut self) {
let _ = self.shutdown_tx.send(true);
self.join.abort();
}
}
#[cfg(test)]
mod startup_tests {
use super::*;
use crate::core::supervision::task::RestartPolicy;
use tokio::sync::oneshot;
fn pending_runtime() -> (
Arc<ConfigMonitorRuntime>,
mpsc::Receiver<ConfigChangeRequest>,
Arc<TaskHealth>,
) {
let (runtime, request_rx) = ConfigMonitorRuntime::pending();
let health = Arc::new(TaskHealth::new(
"config-monitor-test",
RestartPolicy::Always,
));
(runtime, request_rx, health)
}
#[tokio::test]
async fn delayed_success_queues_requests_and_only_then_marks_running() {
let (runtime, request_rx, health) = pending_runtime();
let (release_tx, release_rx) = oneshot::channel::<()>();
let (running_tx, running_rx) = oneshot::channel::<()>();
let (shutdown_tx, shutdown_rx) = watch::channel(false);
let runtime_for_task = runtime.clone();
let health_for_task = health.clone();
let task = tokio::spawn(run_startup(
runtime_for_task,
health_for_task,
PrivacyFilter::new(&[]),
shutdown_rx,
move || async move {
release_rx.await.map_err(|error| error.to_string())?;
Ok(request_rx)
},
move |mut requests, mut shutdown| async move {
assert!(matches!(
requests.recv().await,
Some(ConfigChangeRequest::InitialInventory)
));
let queued = requests.recv().await;
assert!(matches!(
queued,
Some(ConfigChangeRequest::ManualRescan { .. })
));
let _ = running_tx.send(());
wait_for_shutdown(&mut shutdown).await;
Ok(())
},
));
runtime
.request_tx()
.expect("pending sender")
.send(ConfigChangeRequest::ManualRescan { path_filter: None })
.await
.expect("early request queues");
assert_eq!(runtime.snapshot().state, ConfigMonitorState::Pending);
assert_eq!(health.state(), TaskState::Starting);
release_tx.send(()).expect("release initializer");
running_rx.await.expect("runner observed queued request");
assert_eq!(runtime.snapshot().state, ConfigMonitorState::Running);
assert_eq!(health.state(), TaskState::Running);
shutdown_tx.send(true).expect("signal shutdown");
task.await.expect("startup task joins");
}
#[tokio::test]
async fn failed_startup_is_visible_and_scrubbed() {
let (runtime, _request_rx, health) = pending_runtime();
let (_shutdown_tx, shutdown_rx) = watch::channel(false);
run_startup(
runtime.clone(),
health.clone(),
PrivacyFilter::new(&[]),
shutdown_rx,
|| async { Err::<(), _>("watcher rejected Bearer top-secret".to_string()) },
|(), _| async { Ok(()) },
)
.await;
let snapshot = runtime.snapshot();
assert_eq!(snapshot.state, ConfigMonitorState::Failed);
assert!(snapshot.error.unwrap_or_default().contains("[BEARER:***]"));
assert_eq!(health.state(), TaskState::Failed);
assert!(runtime.request_tx().is_none());
}
#[tokio::test]
async fn shutdown_while_pending_prevents_late_running_transition() {
let (runtime, _request_rx, health) = pending_runtime();
let (shutdown_tx, shutdown_rx) = watch::channel(false);
let (started_tx, started_rx) = oneshot::channel::<()>();
let (release_tx, release_rx) = oneshot::channel::<()>();
let runtime_for_task = runtime.clone();
let health_for_task = health.clone();
let task = tokio::spawn(run_startup(
runtime_for_task,
health_for_task,
PrivacyFilter::new(&[]),
shutdown_rx,
move || async move {
let _ = started_tx.send(());
release_rx.await.map_err(|error| error.to_string())?;
Ok(())
},
|(), _| async { Ok(()) },
));
started_rx.await.expect("initializer started");
shutdown_tx.send(true).expect("signal shutdown");
task.await.expect("startup task joins promptly");
let _ = release_tx.send(());
assert_eq!(runtime.snapshot().state, ConfigMonitorState::Stopped);
assert_eq!(health.state(), TaskState::Stopped);
assert!(runtime.request_tx().is_none());
}
}