use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, mpsc};
use notify::Watcher;
use tokio::sync::{oneshot, watch};
use super::logger::{RegisteredSandboxLogger, SandboxLogger};
use crate::{MicrosandboxError, MicrosandboxResult};
#[derive(Clone)]
pub struct LogRegistry {
inner: Arc<RegistryInner>,
}
struct RegistryInner {
admin: tokio::sync::Mutex<()>,
admin_tx: mpsc::Sender<AdminCommand>,
dirs: Arc<Mutex<HashMap<PathBuf, DirEntry>>>,
registered_dirs: Arc<AtomicUsize>,
stats: Arc<RegistryStats>,
}
enum AdminCommand {
Watch {
dir: PathBuf,
reply: oneshot::Sender<notify::Result<()>>,
},
Unwatch {
dir: PathBuf,
},
}
struct DirEntry {
signal: Arc<watch::Sender<u64>>,
registrations: usize,
}
#[derive(Default)]
struct RegistryStats {
total_registrations: AtomicUsize,
wake_all_events: AtomicU64,
route_misses: AtomicU64,
watch_failures: AtomicU64,
}
#[derive(Debug, Clone, Copy)]
pub struct RegistryStatsSnapshot {
pub registered_dirs: usize,
pub total_registrations: usize,
pub wake_all_events: u64,
pub route_misses: u64,
pub watch_failures: u64,
}
pub struct LogRegistration {
inner: Arc<RegistryInner>,
dir: PathBuf,
}
pub(crate) struct LogSubscription {
rx: watch::Receiver<u64>,
_registration: LogRegistration,
}
impl LogRegistry {
pub fn new() -> MicrosandboxResult<Self> {
let dirs: Arc<Mutex<HashMap<PathBuf, DirEntry>>> = Arc::new(Mutex::new(HashMap::new()));
let stats = Arc::new(RegistryStats::default());
let cb_dirs = Arc::clone(&dirs);
let cb_stats = Arc::clone(&stats);
let watcher = notify::recommended_watcher(move |res: notify::Result<notify::Event>| {
route_event(&cb_dirs, &cb_stats, res);
})
.map_err(|e| MicrosandboxError::Custom(format!("log watch registry init failed: {e}")))?;
let (admin_tx, admin_rx) = mpsc::channel();
std::thread::Builder::new()
.name("log-watch-admin".into())
.spawn(move || admin_loop(watcher, admin_rx))
.map_err(|e| {
MicrosandboxError::Custom(format!("log watch admin thread spawn failed: {e}"))
})?;
Ok(Self {
inner: Arc::new(RegistryInner {
admin: tokio::sync::Mutex::new(()),
admin_tx,
dirs,
registered_dirs: Arc::new(AtomicUsize::new(0)),
stats,
}),
})
}
pub async fn register(
&self,
logger: SandboxLogger,
) -> MicrosandboxResult<RegisteredSandboxLogger> {
let dir = tokio::fs::canonicalize(logger.log_dir())
.await
.map_err(|_| MicrosandboxError::SandboxNotFound(logger.name().to_string()))?;
let inner = &self.inner;
let _admin = inner.admin.lock().await;
{
let mut map = lock(&inner.dirs);
if let Some(entry) = map.get_mut(&dir) {
entry.registrations += 1;
inner
.stats
.total_registrations
.fetch_add(1, Ordering::Relaxed);
return Ok(RegisteredSandboxLogger::new(
logger,
LogRegistration {
inner: Arc::clone(inner),
dir,
},
));
}
}
let (reply_tx, reply_rx) = oneshot::channel();
inner
.admin_tx
.send(AdminCommand::Watch {
dir: dir.clone(),
reply: reply_tx,
})
.map_err(|_| MicrosandboxError::Custom("log watch admin thread stopped".into()))?;
reply_rx
.await
.map_err(|_| MicrosandboxError::Custom("log watch admin thread dropped reply".into()))?
.map_err(|e| {
inner.stats.watch_failures.fetch_add(1, Ordering::Relaxed);
MicrosandboxError::Custom(format!(
"log watch subscribe failed for {}: {e}",
dir.display()
))
})?;
{
let (tx, _rx) = watch::channel(0u64);
lock(&inner.dirs).insert(
dir.clone(),
DirEntry {
signal: Arc::new(tx),
registrations: 1,
},
);
}
inner.registered_dirs.fetch_add(1, Ordering::Relaxed);
inner
.stats
.total_registrations
.fetch_add(1, Ordering::Relaxed);
Ok(RegisteredSandboxLogger::new(
logger,
LogRegistration {
inner: Arc::clone(inner),
dir,
},
))
}
pub fn stats(&self) -> RegistryStatsSnapshot {
let s = &self.inner.stats;
RegistryStatsSnapshot {
registered_dirs: self.inner.registered_dirs.load(Ordering::Relaxed),
total_registrations: s.total_registrations.load(Ordering::Relaxed),
wake_all_events: s.wake_all_events.load(Ordering::Relaxed),
route_misses: s.route_misses.load(Ordering::Relaxed),
watch_failures: s.watch_failures.load(Ordering::Relaxed),
}
}
}
impl LogRegistration {
pub(crate) fn subscribe(&self) -> LogSubscription {
let rx = lock(&self.inner.dirs)
.get(&self.dir)
.map(|entry| entry.signal.subscribe())
.expect("registration keeps its dir entry alive");
LogSubscription {
rx,
_registration: self.clone(),
}
}
}
impl Clone for LogRegistration {
fn clone(&self) -> Self {
{
let mut map = lock(&self.inner.dirs);
if let Some(entry) = map.get_mut(&self.dir) {
entry.registrations += 1;
self.inner
.stats
.total_registrations
.fetch_add(1, Ordering::Relaxed);
}
}
Self {
inner: Arc::clone(&self.inner),
dir: self.dir.clone(),
}
}
}
impl Drop for LogRegistration {
fn drop(&mut self) {
let mut map = lock(&self.inner.dirs);
let Some(entry) = map.get_mut(&self.dir) else {
return;
};
entry.registrations -= 1;
self.inner
.stats
.total_registrations
.fetch_sub(1, Ordering::Relaxed);
if entry.registrations > 0 {
return;
}
map.remove(&self.dir);
self.inner.registered_dirs.fetch_sub(1, Ordering::Relaxed);
let _ = self.inner.admin_tx.send(AdminCommand::Unwatch {
dir: self.dir.clone(),
});
}
}
impl LogSubscription {
pub(crate) async fn changed(&mut self) -> Result<(), watch::error::RecvError> {
self.rx.changed().await
}
}
fn admin_loop(mut watcher: notify::RecommendedWatcher, rx: mpsc::Receiver<AdminCommand>) {
while let Ok(command) = rx.recv() {
match command {
AdminCommand::Watch { dir, reply } => {
let result = watcher.watch(&dir, notify::RecursiveMode::NonRecursive);
let _ = reply.send(result);
}
AdminCommand::Unwatch { dir } => {
let _ = watcher.unwatch(&dir);
}
}
}
}
fn route_event(
dirs: &Mutex<HashMap<PathBuf, DirEntry>>,
stats: &RegistryStats,
res: notify::Result<notify::Event>,
) {
let event = match res {
Ok(event) => {
if event.need_rescan() {
wake_all(dirs, stats);
return;
}
event
}
Err(_) => {
wake_all(dirs, stats);
return;
}
};
use notify::EventKind;
if !matches!(
event.kind,
EventKind::Modify(_) | EventKind::Create(_) | EventKind::Remove(_)
) {
return;
}
let signals = collect_signals(dirs, stats, &event.paths);
for signal in signals {
signal.send_modify(|v| *v = v.wrapping_add(1));
}
}
fn collect_signals(
dirs: &Mutex<HashMap<PathBuf, DirEntry>>,
stats: &RegistryStats,
paths: &[PathBuf],
) -> Vec<Arc<watch::Sender<u64>>> {
let map = lock(dirs);
let mut signals = Vec::new();
for path in paths {
let entry = map
.get(path.as_path())
.or_else(|| path.parent().and_then(|parent| map.get(parent)));
match entry {
Some(entry) => signals.push(Arc::clone(&entry.signal)),
None => {
stats.route_misses.fetch_add(1, Ordering::Relaxed);
}
}
}
signals
}
fn wake_all(dirs: &Mutex<HashMap<PathBuf, DirEntry>>, stats: &RegistryStats) {
let signals: Vec<Arc<watch::Sender<u64>>> =
{ lock(dirs).values().map(|e| Arc::clone(&e.signal)).collect() };
for signal in &signals {
signal.send_modify(|v| *v = v.wrapping_add(1));
}
stats.wake_all_events.fetch_add(1, Ordering::Relaxed);
}
fn lock<T>(m: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
m.lock().unwrap_or_else(|e| e.into_inner())
}
#[cfg(test)]
mod tests {
use super::*;
fn map_with(
dirs: &[&std::path::Path],
) -> (
Arc<Mutex<HashMap<PathBuf, DirEntry>>>,
Vec<watch::Receiver<u64>>,
) {
let map = Arc::new(Mutex::new(HashMap::new()));
let mut receivers = Vec::new();
{
let mut guard = map.lock().unwrap();
for dir in dirs {
let (tx, rx) = watch::channel(0u64);
receivers.push(rx);
guard.insert(
dir.to_path_buf(),
DirEntry {
signal: Arc::new(tx),
registrations: 1,
},
);
}
}
(map, receivers)
}
fn synthetic_event(kind: notify::EventKind, paths: Vec<PathBuf>) -> notify::Event {
notify::Event {
kind,
paths,
attrs: Default::default(),
}
}
#[test]
fn routes_child_event_to_only_its_directory() {
let a = PathBuf::from("/tmp/msb-a/logs");
let b = PathBuf::from("/tmp/msb-b/logs");
let (map, receivers) = map_with(&[&a, &b]);
let stats = RegistryStats::default();
let event = synthetic_event(
notify::EventKind::Modify(notify::event::ModifyKind::Data(
notify::event::DataChange::Content,
)),
vec![a.join("exec.log")],
);
route_event(&map, &stats, Ok(event));
assert_eq!(*receivers[0].borrow(), 1, "dir a woke");
assert_eq!(*receivers[1].borrow(), 0, "dir b untouched");
assert_eq!(stats.route_misses.load(Ordering::Relaxed), 0);
}
#[test]
fn unmatched_path_counts_a_route_miss() {
let a = PathBuf::from("/tmp/msb-a/logs");
let (map, _receivers) = map_with(&[&a]);
let stats = RegistryStats::default();
let event = synthetic_event(
notify::EventKind::Create(notify::event::CreateKind::File),
vec![PathBuf::from("/tmp/msb-other/logs/exec.log")],
);
route_event(&map, &stats, Ok(event));
assert_eq!(stats.route_misses.load(Ordering::Relaxed), 1);
}
#[test]
fn error_wakes_all_directories() {
let a = PathBuf::from("/tmp/msb-a/logs");
let b = PathBuf::from("/tmp/msb-b/logs");
let (map, receivers) = map_with(&[&a, &b]);
let stats = RegistryStats::default();
route_event(&map, &stats, Err(notify::Error::generic("dropped events")));
assert_eq!(*receivers[0].borrow(), 1);
assert_eq!(*receivers[1].borrow(), 1);
assert_eq!(stats.wake_all_events.load(Ordering::Relaxed), 1);
}
#[test]
fn ignored_event_kinds_do_not_wake() {
let a = PathBuf::from("/tmp/msb-a/logs");
let (map, receivers) = map_with(&[&a]);
let stats = RegistryStats::default();
let event = synthetic_event(
notify::EventKind::Access(notify::event::AccessKind::Read),
vec![a.join("exec.log")],
);
route_event(&map, &stats, Ok(event));
assert_eq!(*receivers[0].borrow(), 0);
}
#[tokio::test]
async fn register_missing_dir_is_sandbox_not_found() {
let registry = LogRegistry::new().unwrap();
let logger = SandboxLogger::new(
"sbx".to_string(),
PathBuf::from("/tmp/msb-nonexistent-xyz/logs"),
);
match registry.register(logger).await {
Err(MicrosandboxError::SandboxNotFound(_)) => {}
Err(other) => panic!("expected SandboxNotFound, got {other:?}"),
Ok(_) => panic!("expected SandboxNotFound, got Ok"),
}
assert_eq!(registry.stats().registered_dirs, 0);
}
#[cfg_attr(
target_os = "macos",
ignore = "FSEvents watch() latency; runs on Linux CI"
)]
#[tokio::test]
async fn one_directory_two_registrations_share_one_descriptor() {
let registry = LogRegistry::new().unwrap();
let dir = tempfile::tempdir().unwrap();
let first = registry
.register(SandboxLogger::new("sbx".into(), dir.path().to_path_buf()))
.await
.unwrap();
let second = registry
.register(SandboxLogger::new("sbx".into(), dir.path().to_path_buf()))
.await
.unwrap();
assert_eq!(registry.stats().registered_dirs, 1);
assert_eq!(registry.stats().total_registrations, 2);
drop(first);
assert_eq!(registry.stats().registered_dirs, 1);
assert_eq!(registry.stats().total_registrations, 1);
drop(second);
assert_eq!(registry.stats().registered_dirs, 0);
let third = registry
.register(SandboxLogger::new("sbx".into(), dir.path().to_path_buf()))
.await
.unwrap();
assert_eq!(registry.stats().registered_dirs, 1);
drop(third);
assert_eq!(registry.stats().registered_dirs, 0);
}
#[cfg_attr(
target_os = "macos",
ignore = "FSEvents watch() latency; runs on Linux CI"
)]
#[tokio::test]
async fn stream_keeps_directory_registered_after_logger_drops() {
use crate::logs::LogStreamOptions;
let registry = LogRegistry::new().unwrap();
let dir = tempfile::tempdir().unwrap();
let logger = registry
.register(SandboxLogger::new("sbx".into(), dir.path().to_path_buf()))
.await
.unwrap();
let stream = logger
.stream(&LogStreamOptions {
follow: true,
..Default::default()
})
.await
.unwrap();
assert_eq!(registry.stats().registered_dirs, 1);
assert_eq!(registry.stats().total_registrations, 2);
drop(logger);
assert_eq!(registry.stats().registered_dirs, 1);
assert_eq!(registry.stats().total_registrations, 1);
drop(stream);
assert_eq!(registry.stats().registered_dirs, 0);
assert_eq!(registry.stats().total_registrations, 0);
}
}