use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::RecvTimeoutError;
use std::thread::JoinHandle;
use std::time::Duration;
use omgbase_store::{Config, Store};
use omgbase_sync::fs::RealFileSystem;
use omgbase_sync::registry::{AdapterRow, FS_ADAPTER, FS_ADAPTER_COMMAND, SourceRow};
use omgbase_sync::{
CheckpointResult, ExternalSource, Readiness, RepoRow, SyncSource, WatchEvent, WatchLease,
WriterLockOptions, freshness_sweep, reconcile_changes, wait_ready, with_writer_lock,
};
use crate::drain::DrainHandle;
use crate::{MinterSource, open_store, stamp};
pub const ADAPTER_ENV: &str = "OMGBASE_FS_ADAPTER";
pub const READY_PATIENCE: Duration = Duration::from_secs(30);
const POLL: Duration = Duration::from_millis(200);
pub struct WatchOptions {
pub omgbase_dir: PathBuf,
pub db_path: PathBuf,
pub repo: RepoRow,
pub minters: MinterSource,
pub clock: Option<String>,
pub adapter_override: Option<Vec<String>>,
pub ready_patience: Duration,
pub drain: Option<DrainHandle>,
}
pub enum Outcome {
Live(Watcher),
LeaseHeld,
NoSource,
AdapterUnavailable(String),
}
pub struct Watcher {
stop: Arc<AtomicBool>,
thread: Option<JoinHandle<()>>,
lease: Option<WatchLease>,
argv: Vec<String>,
}
impl Watcher {
#[must_use]
pub fn argv(&self) -> &[String] {
&self.argv
}
pub fn stop(mut self) {
self.stop.store(true, Ordering::SeqCst);
if let Some(t) = self.thread.take() {
let _ = t.join();
}
if let Some(mut lease) = self.lease.take() {
lease.release();
}
}
}
impl Drop for Watcher {
fn drop(&mut self) {
self.stop.store(true, Ordering::SeqCst);
if let Some(t) = self.thread.take() {
let _ = t.join();
}
if let Some(mut lease) = self.lease.take() {
lease.release();
}
}
}
#[must_use]
pub fn resolve_adapter(row: &AdapterRow, override_argv: Option<&[String]>) -> AdapterRow {
if row.name != FS_ADAPTER {
return row.clone();
}
let (command, leading): (String, &[String]) = match override_argv {
Some([command, rest @ ..]) => (command.clone(), rest),
_ => (FS_ADAPTER_COMMAND.to_owned(), &[]),
};
let mut args = leading.to_vec();
args.extend(row.args.iter().cloned());
AdapterRow {
name: row.name.clone(),
command,
args,
}
}
#[must_use]
pub fn adapter_override_from_env() -> Option<Vec<String>> {
let raw = std::env::var(ADAPTER_ENV).ok()?;
let argv: Vec<String> = raw.split_whitespace().map(str::to_owned).collect();
(!argv.is_empty()).then_some(argv)
}
fn fs_source(store: &omgbase_store::Store, repo_id: &str) -> Result<Option<SourceRow>, String> {
let sources = omgbase_sync::sources_for_repo(store, repo_id).map_err(|e| e.to_string())?;
Ok(sources.into_iter().find(|s| {
s.adapter == FS_ADAPTER
&& s.config
.get("root")
.and_then(|v| v.as_str())
.is_some_and(|r| !r.is_empty())
}))
}
fn prime(store: &mut Store, opts: &WatchOptions, config: &Config) {
let Some(root) = opts.repo.root_path.as_deref() else {
return;
};
let ts = stamp(opts.clock.as_ref());
let swept = with_writer_lock(&opts.omgbase_dir, WriterLockOptions::default(), || {
freshness_sweep(
store,
&opts.repo.repo_id,
&RealFileSystem,
Path::new(root),
&ts,
None,
config,
)
});
match swept {
Ok(s) if s.changed => eprintln!(
"[watch] primed: +{} -{}",
s.checkpoint.ingested.len(),
s.checkpoint.deleted.len()
),
Ok(_) => {}
Err(e) => eprintln!("[watch] priming sweep failed: {e}"),
}
}
pub fn start(opts: WatchOptions) -> Result<Outcome, String> {
let Some(lease) = WatchLease::try_acquire(&opts.omgbase_dir).map_err(|e| e.to_string())? else {
return Ok(Outcome::LeaseHeld);
};
let repo_id = opts.repo.repo_id.clone();
let config = Config::default();
let mut store = open_store(&opts.db_path, &opts.minters)?;
let Some(source) = fs_source(&store, &repo_id)? else {
prime(&mut store, &opts, &config);
return Ok(Outcome::NoSource);
};
let adapters = omgbase_sync::list_adapters(&store).map_err(|e| e.to_string())?;
let Some(row) = adapters.into_iter().find(|a| a.name == source.adapter) else {
prime(&mut store, &opts, &config);
return Ok(Outcome::AdapterUnavailable(format!(
"source `{}` names adapter `{}`, which the registry does not hold",
source.name, source.adapter
)));
};
let adapter = resolve_adapter(&row, opts.adapter_override.as_deref());
let argv = ExternalSource::argv(&source, &adapter);
let mut ext = match ExternalSource::spawn_source(&source, &adapter) {
Ok(s) => s,
Err(e) => {
prime(&mut store, &opts, &config);
return Ok(Outcome::AdapterUnavailable(format!(
"cannot start the `{}` adapter as `{}`: {e}",
source.adapter,
argv.join(" ")
)));
}
};
if !ext.capabilities().watch {
let _ = ext.close();
prime(&mut store, &opts, &config);
return Ok(Outcome::AdapterUnavailable(format!(
"adapter `{}` does not advertise `watch`",
argv.join(" ")
)));
}
let rx = match ext.watch() {
Ok(rx) => rx,
Err(e) => {
let _ = ext.close();
prime(&mut store, &opts, &config);
return Ok(Outcome::AdapterUnavailable(format!(
"adapter `{}` refused `watch`: {e}",
argv.join(" ")
)));
}
};
let (readiness, early) = wait_ready(&rx, opts.ready_patience);
match readiness {
Readiness::Ready => {}
Readiness::TimedOut => eprintln!(
"[watch] warning: adapter `{}` did not report ready within {:?}; proceeding as if ready (an edit made before now may be missed until the next sweep)",
argv.join(" "),
opts.ready_patience
),
Readiness::Ended => {
let _ = ext.close();
prime(&mut store, &opts, &config);
return Ok(Outcome::AdapterUnavailable(format!(
"adapter `{}` exited before reporting ready",
argv.join(" ")
)));
}
}
prime(&mut store, &opts, &config);
drop(store);
let stop = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&stop);
let thread_opts = opts;
let thread = std::thread::Builder::new()
.name("omgbase-watch".to_owned())
.spawn(move || {
run(&thread_opts, ext, &rx, early, &flag, &config);
})
.map_err(|e| format!("cannot spawn the watcher thread: {e}"))?;
Ok(Outcome::Live(Watcher {
stop,
thread: Some(thread),
lease: Some(lease),
argv,
}))
}
fn run(
opts: &WatchOptions,
mut source: ExternalSource,
rx: &std::sync::mpsc::Receiver<WatchEvent>,
early: Vec<Vec<String>>,
stop: &AtomicBool,
config: &Config,
) {
let mut store = match open_store(&opts.db_path, &opts.minters) {
Ok(s) => s,
Err(e) => {
eprintln!("[watch] error: {e}");
let _ = source.unwatch();
let _ = source.close();
return;
}
};
let checkpoint = |store: &mut Store, source: &mut ExternalSource, paths: &[String]| {
if paths.is_empty() {
return;
}
let ts = stamp(opts.clock.as_ref());
let result = with_writer_lock(&opts.omgbase_dir, WriterLockOptions::default(), || {
reconcile_changes(store, &opts.repo.repo_id, source, paths, &ts, None, config)
});
match result {
Ok(r) => on_checkpoint(&r, opts.drain.as_ref()),
Err(e) => eprintln!("[watch] error: {e}"),
}
};
for paths in early {
checkpoint(&mut store, &mut source, &paths);
}
while !stop.load(Ordering::SeqCst) {
let paths = match rx.recv_timeout(POLL) {
Ok(WatchEvent::Batch(p)) => p,
Ok(WatchEvent::Ready) | Err(RecvTimeoutError::Timeout) => continue,
Err(RecvTimeoutError::Disconnected) => {
eprintln!("[watch] the adapter's stream ended; watching stopped");
break;
}
};
checkpoint(&mut store, &mut source, &paths);
}
let _ = source.unwatch();
let _ = source.close();
}
fn on_checkpoint(r: &CheckpointResult, drain: Option<&DrainHandle>) {
if r.ingested.is_empty() && r.deleted.is_empty() {
return;
}
eprintln!(
"[watch] checkpoint: +{} -{}",
r.ingested.len(),
r.deleted.len()
);
if let Some(d) = drain {
d.schedule();
}
}
#[cfg(test)]
mod tests {
use super::*;
use omgbase_sync::Workspace;
use omgbase_sync::fs::TempDir;
fn options(tmp: &TempDir, root: Option<&str>, argv: Option<&[&str]>) -> WatchOptions {
let mut ws = Workspace::open(tmp.path()).unwrap();
let repo_id = omgbase_sync::ensure_repo(ws.store_mut(), "r", root).unwrap();
let (omgbase_dir, db_path) = (ws.omgbase_dir().to_path_buf(), ws.db_path().to_path_buf());
ws.close().unwrap();
WatchOptions {
omgbase_dir,
db_path,
repo: RepoRow {
repo_id,
slug: "r".into(),
root_path: root.map(str::to_owned),
},
minters: MinterSource::Random,
clock: None,
adapter_override: argv.map(|a| a.iter().map(|s| (*s).to_owned()).collect()),
ready_patience: Duration::from_millis(300),
drain: None,
}
}
fn count(db: &Path, table: &str) -> i64 {
let store = Store::open(db).unwrap();
store
.conn()
.query_row(&format!("SELECT count(*) FROM {table}"), [], |r| r.get(0))
.unwrap()
}
fn sh_adapter(dir: &Path, then: &str) -> String {
let script = r##"
printf '%s\n' '{"protocol":1,"capabilities":{"watch":true}}'
while IFS= read -r line; do
id="${line#*\"id\":}"; id="${id%%,*}"
case "$line" in
*'"watch"'*) printf '%s\n' "{\"id\":$id,\"result\":{\"ok\":true}}"; __THEN__ ;;
*'"fetch"'*) printf '%s\n' "{\"id\":$id,\"result\":{\"item\":{\"path\":\"a.md\",\"revision\":\"1\",\"content\":\"# A\\n\"}}}" ;;
*) printf '%s\n' "{\"id\":$id,\"result\":{\"ok\":true}}" ;;
esac
done
"##
.replace("__THEN__", then);
let path = dir.join("adapter.sh");
std::fs::write(&path, script).unwrap();
path.to_string_lossy().into_owned()
}
#[test]
fn a_sourceless_repo_has_nothing_to_watch() {
let tmp = TempDir::new("watch-sourceless");
let opts = options(&tmp, None, None);
let dir = opts.omgbase_dir.clone();
assert!(matches!(start(opts).unwrap(), Outcome::NoSource));
assert!(!WatchLease::live(&dir), "the lease is released");
}
#[test]
fn a_missing_adapter_degrades_to_no_watch() {
let tmp = TempDir::new("watch-noadapter");
let root = tmp.path().to_string_lossy().into_owned();
std::fs::write(tmp.path().join("a.md"), "# A\n").unwrap();
let opts = options(
&tmp,
Some(&root),
Some(&["definitely-not-an-omgbase-adapter-xyz", "--flag"]),
);
let (dir, db) = (opts.omgbase_dir.clone(), opts.db_path.clone());
match start(opts).unwrap() {
Outcome::AdapterUnavailable(msg) => {
assert!(
msg.contains("definitely-not-an-omgbase-adapter-xyz --flag --root"),
"{msg}"
);
}
_ => panic!("expected AdapterUnavailable"),
}
assert!(!WatchLease::live(&dir), "the lease is released");
assert_eq!(count(&db, "docs"), 1);
}
#[test]
fn a_live_lease_elsewhere_skips_the_watcher() {
let tmp = TempDir::new("watch-lease");
let root = tmp.path().to_string_lossy().into_owned();
let opts = options(&tmp, Some(&root), None);
let held = WatchLease::try_acquire(&opts.omgbase_dir).unwrap().unwrap();
assert!(matches!(start(opts).unwrap(), Outcome::LeaseHeld));
drop(held);
}
#[test]
fn a_non_watching_adapter_is_unavailable() {
let tmp = TempDir::new("watch-nowatch");
let root = tmp.path().to_string_lossy().into_owned();
let script = "echo '{\"protocol\":1,\"capabilities\":{}}'; cat >/dev/null";
let opts = options(&tmp, Some(&root), Some(&["sh", "-c", script]));
match start(opts).unwrap() {
Outcome::AdapterUnavailable(msg) => assert!(msg.contains("watch"), "{msg}"),
_ => panic!("expected AdapterUnavailable"),
}
}
#[test]
fn ready_is_awaited_and_an_early_batch_is_reconciled_after_the_prime() {
let tmp = TempDir::new("watch-ready");
let root = tmp.path().to_string_lossy().into_owned();
std::fs::write(tmp.path().join("a.md"), "# A\n").unwrap();
let script = sh_adapter(
tmp.path(),
r#"printf '%s\n' '{"event":"batch","paths":["a.md"]}' '{"event":"ready"}'"#,
);
let mut opts = options(&tmp, Some(&root), Some(&["sh", &script]));
opts.ready_patience = Duration::from_secs(10);
let (dir, db) = (opts.omgbase_dir.clone(), opts.db_path.clone());
let started = std::time::Instant::now();
let watcher = match start(opts).unwrap() {
Outcome::Live(w) => w,
Outcome::AdapterUnavailable(m) => panic!("{m}"),
_ => panic!("expected Live"),
};
assert!(
started.elapsed() < Duration::from_secs(5),
"ready arrived: no wait for patience"
);
assert_eq!(watcher.argv()[0], "sh");
assert!(WatchLease::live(&dir));
assert_eq!(count(&db, "docs"), 1, "the priming sweep ran before live");
watcher.stop();
assert!(!WatchLease::live(&dir), "the lease is released");
assert_eq!(
count(&db, "checkpoints"),
2,
"the sweep's checkpoint, then the early batch's"
);
assert_eq!(count(&db, "docs"), 1);
}
#[test]
fn a_silent_adapter_is_tolerated_after_patience() {
let tmp = TempDir::new("watch-silent");
let root = tmp.path().to_string_lossy().into_owned();
std::fs::write(tmp.path().join("a.md"), "# A\n").unwrap();
let script = sh_adapter(tmp.path(), ":");
let opts = options(&tmp, Some(&root), Some(&["sh", &script]));
let db = opts.db_path.clone();
let started = std::time::Instant::now();
let watcher = match start(opts).unwrap() {
Outcome::Live(w) => w,
Outcome::AdapterUnavailable(m) => panic!("{m}"),
_ => panic!("expected Live"),
};
assert!(started.elapsed() >= Duration::from_millis(300));
assert_eq!(count(&db, "docs"), 1);
watcher.stop();
assert_eq!(count(&db, "checkpoints"), 1);
}
#[test]
fn an_adapter_that_dies_before_ready_is_unavailable_but_the_sweep_runs() {
let tmp = TempDir::new("watch-dies");
let root = tmp.path().to_string_lossy().into_owned();
std::fs::write(tmp.path().join("a.md"), "# A\n").unwrap();
let script = "echo '{\"protocol\":1,\"capabilities\":{\"watch\":true}}'; read -r line; echo '{\"id\":1,\"result\":{\"ok\":true}}'; exit 0";
let opts = options(&tmp, Some(&root), Some(&["sh", "-c", script]));
let (dir, db) = (opts.omgbase_dir.clone(), opts.db_path.clone());
match start(opts).unwrap() {
Outcome::AdapterUnavailable(msg) => {
assert!(msg.contains("before reporting ready"), "{msg}")
}
_ => panic!("expected AdapterUnavailable"),
}
assert!(!WatchLease::live(&dir));
assert_eq!(count(&db, "docs"), 1);
}
#[test]
fn the_fs_launcher_never_reads_the_row_command() {
let row = AdapterRow {
name: "fs".into(),
command: "/registry/says/this".into(),
args: vec!["--v".into()],
};
let default = resolve_adapter(&row, None);
assert_eq!(default.command, FS_ADAPTER_COMMAND);
assert_eq!(default.args, ["--v"]);
assert_eq!(resolve_adapter(&row, Some(&[])), default, "empty override");
let o = resolve_adapter(&row, Some(&["node".to_owned(), "/x/bin.js".to_owned()]));
assert_eq!(o.command, "node");
assert_eq!(o.args, ["/x/bin.js", "--v"]);
assert_eq!(o.name, "fs");
let other = AdapterRow {
name: "git".into(),
command: "omgbase-git-adapter".into(),
args: vec![],
};
assert_eq!(resolve_adapter(&other, None), other);
assert_eq!(
resolve_adapter(&other, Some(&["node".to_owned()])),
other,
"the fs override is not a general one"
);
}
}