120 lines
4.2 KiB
Rust
120 lines
4.2 KiB
Rust
use std::{
|
|
path::PathBuf,
|
|
sync::{Arc, Mutex},
|
|
thread,
|
|
time::Duration,
|
|
};
|
|
|
|
use notify::{EventKind, RecursiveMode, Watcher};
|
|
use tokio::sync::broadcast;
|
|
|
|
use crate::server::state::ServerState;
|
|
use tracing::{error, info};
|
|
|
|
const POLL_INTERVAL: Duration = Duration::from_secs(15);
|
|
|
|
/// A change observed by the watcher. Pushed onto the broadcast channel for
|
|
/// `/events` subscribers. The client treats these as wake-ups: correctness
|
|
/// always rests on the subsequent `/manifest` hash diff.
|
|
#[derive(Debug, Clone)]
|
|
pub struct ChangeEvent {
|
|
pub kind: ChangeKind,
|
|
}
|
|
|
|
#[derive(Debug, Clone, Copy)]
|
|
pub enum ChangeKind {
|
|
Create,
|
|
Modify,
|
|
Remove,
|
|
}
|
|
|
|
/// Server-side file watcher. Owns a background std thread driving `notify`
|
|
/// (inotify on Linux). On any filesystem event under `source` it:
|
|
/// 1. Rebuilds the shared [`ServerState`] from a fresh scan.
|
|
/// 2. Broadcasts a [`ChangeEvent`] so connected `/events` clients wake up.
|
|
///
|
|
/// The state rebuild is full-scan rather than incremental. This matches the
|
|
/// existing `LocalOriginFileWatcher` pattern, keeps the watcher simple, and
|
|
/// is cheap on a LAN-scale library. Hash-diff reconciliation on the client
|
|
/// absorbs any over-reporting.
|
|
pub struct ServerWatcher {
|
|
_events_tx: broadcast::Sender<ChangeEvent>,
|
|
_worker: Arc<Mutex<Option<thread::JoinHandle<()>>>>,
|
|
}
|
|
|
|
impl ServerWatcher {
|
|
/// Spawn the watcher. Returns the broadcast receiver that `/events`
|
|
/// handlers subscribe to.
|
|
pub fn spawn(source: PathBuf, state: ServerState) -> (Self, broadcast::Receiver<ChangeEvent>) {
|
|
let (events_tx, events_rx) = broadcast::channel(64);
|
|
let worker_tx = events_tx.clone();
|
|
let handle = thread::spawn(move || {
|
|
run_watcher_loop(source, state, worker_tx);
|
|
});
|
|
let watcher = ServerWatcher {
|
|
_events_tx: events_tx,
|
|
_worker: Arc::new(Mutex::new(Some(handle))),
|
|
};
|
|
return (watcher, events_rx);
|
|
}
|
|
}
|
|
|
|
fn run_watcher_loop(
|
|
source: PathBuf,
|
|
state: ServerState,
|
|
events_tx: broadcast::Sender<ChangeEvent>,
|
|
) {
|
|
let (tx, rx) = std::sync::mpsc::channel();
|
|
let mut watcher = match notify::recommended_watcher(tx) {
|
|
Ok(w) => w,
|
|
Err(e) => {
|
|
error!(error = %e, "server watcher: failed to create inotify watcher");
|
|
return;
|
|
}
|
|
};
|
|
|
|
if let Err(e) = watcher.watch(&source, RecursiveMode::Recursive) {
|
|
error!(source = %source.display(), error = %e, "server watcher: failed to watch source");
|
|
return;
|
|
}
|
|
|
|
info!(source = %source.display(), "server watcher: scanning for changes");
|
|
loop {
|
|
match rx.recv_timeout(POLL_INTERVAL) {
|
|
Ok(res) => match res {
|
|
Ok(event) => match event.kind {
|
|
EventKind::Create(_) | EventKind::Modify(_) | EventKind::Remove(_) => {
|
|
let kind = match event.kind {
|
|
EventKind::Create(_) => ChangeKind::Create,
|
|
EventKind::Modify(_) => ChangeKind::Modify,
|
|
EventKind::Remove(_) => ChangeKind::Remove,
|
|
_ => continue,
|
|
};
|
|
match state.replace_all(&source) {
|
|
Ok(true) => {
|
|
let _ = events_tx.send(ChangeEvent { kind });
|
|
}
|
|
Ok(false) => {}
|
|
Err(e) => {
|
|
error!(error = %e, "server watcher: state refresh failed");
|
|
}
|
|
}
|
|
}
|
|
_ => {}
|
|
},
|
|
Err(e) => error!(error = %e, "server watcher: inotify error"),
|
|
},
|
|
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => match state.replace_all(&source) {
|
|
Ok(true) => {
|
|
let _ = events_tx.send(ChangeEvent {
|
|
kind: ChangeKind::Modify,
|
|
});
|
|
}
|
|
Ok(false) => {}
|
|
Err(e) => error!(error = %e, "server watcher: poll rescan failed"),
|
|
},
|
|
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => break,
|
|
}
|
|
}
|
|
}
|