From c7ced96e86a6b6eecfd6af57d51d316d782cc4d8 Mon Sep 17 00:00:00 2001 From: Dave Thompson Date: Mon, 31 Aug 2026 22:34:08 -0400 Subject: [PATCH] feat(cli): add native filesystem watch demo Adds a diagnostic watch command backed by the native filesystem observer, with bounded descriptor storage and native retry-demo coverage. Generated by GPT-5.6 Sol via Codex under supervision of @3leapsdave Co-Authored-By: GPT-5.6 Sol Role: devlead Committer-of-Record: Dave Thompson [@3leapsdave] --- Cargo.lock | 2 + Makefile | 9 +- crates/waitprims-cli/Cargo.toml | 2 + crates/waitprims-cli/src/main.rs | 338 +++++++++++++++++++++- crates/waitprims-cli/tests/watch_cli.rs | 359 ++++++++++++++++++++++++ 5 files changed, 704 insertions(+), 6 deletions(-) create mode 100644 crates/waitprims-cli/tests/watch_cli.rs diff --git a/Cargo.lock b/Cargo.lock index 0d063d0..ab05e39 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1248,11 +1248,13 @@ dependencies = [ "clap", "jsonschema", "serde_json", + "sha2", "tokio", "tracing", "tracing-subscriber", "waitprims-async", "waitprims-core", + "waitprims-fs", "waitprims-testkit", ] diff --git a/Makefile b/Makefile index c343812..0f89db5 100644 --- a/Makefile +++ b/Makefile @@ -10,7 +10,7 @@ # make release-preflight Verify pre-tag requirements # make release-check Version consistency + cargo package (does not publish) -.PHONY: all help check test fmt fmt-check lint build clean version demo-follow demo-coalesce +.PHONY: all help check test fmt fmt-check lint build clean version demo-follow demo-coalesce demo-watch .PHONY: precommit prepush pr-final .PHONY: version-patch version-minor version-major version-set version-sync version-check .PHONY: release-check release-preflight release-guard-tag-version @@ -60,6 +60,7 @@ help: ## Show available targets @echo " build cargo build --workspace" @echo " demo-follow Locked offline held-follow CLI demo vs golden JSONL" @echo " demo-coalesce Locked offline held-coalesce CLI demo vs golden JSONL" + @echo " demo-watch Native filesystem CLI demo with bounded retry" @echo " clean cargo clean" @echo " precommit fmt-check, clippy" @echo " prepush fmt-check, clippy, locked tests, version-check" @@ -135,6 +136,10 @@ demo-coalesce: ## Locked offline held-coalesce CLI demo vs golden JSONL | cmp - fixtures/coalesce-demo/golden.jsonl @echo "[ok] demo-coalesce matched golden JSONL" +demo-watch: ## Native filesystem CLI demo with bounded create/remove retry + $(CARGO) test --locked -p waitprims-cli --test watch_cli native_watch_demo_uses_bounded_retry_and_visible_event_surface + @echo "[ok] demo-watch observed native filesystem JSONL" + clean: ## Remove build artifacts $(CARGO) clean @echo "[ok] Clean complete" @@ -145,7 +150,7 @@ precommit: fmt-check lint ## Fast checks for every commit prepush: fmt-check lint test version-check ## Thorough checks before push @echo "[ok] Pre-push checks passed" -pr-final: prepush demo-follow demo-coalesce ## Final PR merge-readiness gate +pr-final: prepush demo-follow demo-coalesce demo-watch ## Final PR merge-readiness gate @echo "[ok] PR final checks passed" # ----------------------------------------------------------------------------- diff --git a/crates/waitprims-cli/Cargo.toml b/crates/waitprims-cli/Cargo.toml index 60ed1f3..dd6b4f5 100644 --- a/crates/waitprims-cli/Cargo.toml +++ b/crates/waitprims-cli/Cargo.toml @@ -22,7 +22,9 @@ tracing.workspace = true tracing-subscriber.workspace = true waitprims-async.workspace = true waitprims-core.workspace = true +waitprims-fs.workspace = true waitprims-testkit.workspace = true +sha2.workspace = true [dev-dependencies] jsonschema.workspace = true diff --git a/crates/waitprims-cli/src/main.rs b/crates/waitprims-cli/src/main.rs index 17dd61e..4297df2 100644 --- a/crates/waitprims-cli/src/main.rs +++ b/crates/waitprims-cli/src/main.rs @@ -7,6 +7,9 @@ use std::fs; use std::path::{Path, PathBuf}; use std::process::ExitCode; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::Mutex; +use std::time::Instant; mod diagnostic; @@ -15,12 +18,16 @@ use std::time::Duration; use clap::{Parser, Subcommand, ValueEnum}; use tracing::info; use waitprims_async::{ - run_coalesce, run_first_match, run_follow, run_poll_cycle, Cancel, CoalescePolicy, + run_coalesce, run_first_match, run_follow, run_poll_cycle, Cancel, Clock, CoalescePolicy, }; use waitprims_core::{ bundled_entry_schema, bundled_message_schema, validate_message, validate_raw_documents, - AgentWaitMessage, Error, LiveWaitRequest, MessageType, PollCycleRequest, RegistrationSet, - ValidationError, + AgentWaitMessage, ContentDigest, DigestAlgorithm, Error, LiveWaitRequest, MessageType, + OpaqueRef, PayloadRef, PollCycleRequest, RegistrationSet, Timestamp, ValidationError, +}; +use waitprims_fs::{ + EventClock, EventDescriptor, EventRefSink, FilesystemPosture, FsObserver, SystemEventClock, + METHOD_FILE_WATCH, }; use waitprims_testkit::{FakeClock, Script, ScriptedObserver}; @@ -87,6 +94,18 @@ enum Command { #[arg(long, value_name = "PATH")] script: PathBuf, }, + /// Follow native local-filesystem events. + Watch { + /// Existing local directory that owns all watched subjects. + #[arg(long, value_name = "DIR")] + root: PathBuf, + /// Admitted `registration_set` JSON file. + #[arg(long, value_name = "PATH")] + registration_set: PathBuf, + /// Admitted `live_wait_request` JSON file. + #[arg(long, value_name = "PATH")] + request: PathBuf, + }, /// Replay a scripted held-coalesce session. Coalesce { /// Admitted `registration_set` JSON file. @@ -220,6 +239,17 @@ Diagnostic CLI. The library is the product; there is no daemon.", ExitCode::from(1) } }, + Some(Command::Watch { + root, + registration_set, + request, + }) => match run_native_watch(&root, ®istration_set, &request) { + Ok(()) => ExitCode::SUCCESS, + Err(err) => { + eprintln!("waitprims watch: {err}"); + ExitCode::from(1) + } + }, Some(Command::Coalesce { registration_set, request, @@ -369,6 +399,232 @@ fn run_held_follow(set_path: &Path, request_path: &Path, script_path: &Path) -> Ok(()) } +fn run_native_watch(root: &Path, set_path: &Path, request_path: &Path) -> Result<(), Error> { + let (root, set, request, source_instance_ref) = + admit_watch_inputs(root, set_path, request_path)?; + let event_sink = std::sync::Arc::new(CliEventSink::new( + set.aggregate_limits.max_events, + set.aggregate_limits.max_bytes, + )); + let observer = FsObserver::new( + source_instance_ref, + root, + FilesystemPosture::Local, + event_sink.clone(), + )?; + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_time() + .build() + .map_err(|_| Error::Contract { + path: "runtime", + constraint: "init", + })?; + let clock = RealtimeClock::new(); + let cancel = Cancel::new(); + let mut sink = diagnostic::JsonlSink::new(std::io::stdout()); + let end = runtime.block_on(run_follow( + &observer, + &clock, + &cancel, + &set, + &request, + |burst| { + let result = sink.emit_burst(&burst); + async move { result } + }, + ))?; + if event_sink.failed() { + return Err(Error::Contract { + path: "event_ref_sink", + constraint: "capacity", + }); + } + sink.emit_end(&end) +} + +fn admit_watch_inputs( + root: &Path, + set_path: &Path, + request_path: &Path, +) -> Result<(PathBuf, RegistrationSet, LiveWaitRequest, OpaqueRef), Error> { + reject_non_local_path(root)?; + reject_non_local_path(set_path)?; + reject_non_local_path(request_path)?; + let root = fs::canonicalize(root) + .map_err(|_| ValidationError::new("/root", "existing_local_directory_required"))?; + if !root.is_dir() { + return Err(ValidationError::new("/root", "existing_local_directory_required").into()); + } + + let set_raw = read_raw(set_path)?; + let request_raw = read_raw(request_path)?; + let admitted = validate_raw_documents([&set_raw, &request_raw])?; + let (set, request) = take_set_and_request(admitted)?; + let source_instance_ref = set + .registrations + .first() + .ok_or_else(|| ValidationError::new("/registrations", "nonempty_required"))? + .source_instance_ref + .clone(); + for registration in &set.registrations { + if registration.method_id.as_str() != METHOD_FILE_WATCH { + return Err( + ValidationError::new("/registrations/method_id", "file_watch_required").into(), + ); + } + if registration.source_instance_ref != source_instance_ref { + return Err(ValidationError::new( + "/registrations/source_instance_ref", + "single_source_required", + ) + .into()); + } + if registration.subject_kind.as_str() != "path" { + return Err( + ValidationError::new("/registrations/subject_kind", "path_required").into(), + ); + } + validate_watch_subject(registration.subject_id.as_str())?; + } + Ok((root, set, request, source_instance_ref)) +} + +fn validate_watch_subject(subject: &str) -> Result<(), Error> { + use std::path::Component; + + if subject.is_empty() + || subject.starts_with('/') + || subject.starts_with('\\') + || subject.contains('\\') + || subject.as_bytes().get(1) == Some(&b':') + { + return Err(ValidationError::new( + "/registrations/subject_id", + "portable_relative_path_required", + ) + .into()); + } + if Path::new(subject).components().any(|component| { + matches!( + component, + Component::ParentDir | Component::RootDir | Component::Prefix(_) + ) + }) { + return Err(ValidationError::new( + "/registrations/subject_id", + "portable_relative_path_required", + ) + .into()); + } + Ok(()) +} + +struct RealtimeClock { + timestamp: Timestamp, + instant: Instant, +} + +impl RealtimeClock { + fn new() -> Self { + Self { + timestamp: SystemEventClock.now(), + instant: Instant::now(), + } + } +} + +impl Clock for RealtimeClock { + fn now(&self) -> Timestamp { + self.timestamp.saturating_add(self.instant.elapsed()) + } + + async fn sleep_until(&self, deadline: &Timestamp) { + tokio::time::sleep(self.now().duration_until(deadline)).await; + } +} + +struct CliEventSink { + next: AtomicU64, + descriptors: Mutex, + max_events: u64, + max_bytes: u64, + failed: AtomicBool, + #[cfg(test)] + inspector: Option, +} + +#[derive(Default)] +struct CliDescriptorStore { + bytes: u64, + entries: std::collections::BTreeMap>, +} + +#[cfg(test)] +type DescriptorInspector = std::sync::Arc; + +impl CliEventSink { + fn new(max_events: u64, max_bytes: u64) -> Self { + Self { + next: AtomicU64::new(0), + descriptors: Mutex::new(CliDescriptorStore::default()), + max_events, + max_bytes, + failed: AtomicBool::new(false), + #[cfg(test)] + inspector: None, + } + } + + fn failed(&self) -> bool { + self.failed.load(Ordering::Acquire) + } +} + +impl Default for CliEventSink { + fn default() -> Self { + Self::new(u64::MAX, u64::MAX) + } +} + +impl EventRefSink for CliEventSink { + fn materialize(&self, descriptor: &EventDescriptor) -> waitprims_core::Result { + use sha2::{Digest, Sha256}; + + let bytes = descriptor.canonical_bytes(); + let digest = Sha256::digest(&bytes) + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + #[cfg(test)] + if let Some(inspector) = &self.inspector { + inspector( + descriptor, + &bytes, + &digest, + self.next.load(Ordering::Acquire), + ); + } + let sequence = self.next.fetch_add(1, Ordering::AcqRel) + 1; + let payload_ref = format!("ref:fs-cli-{sequence}"); + let mut store = self.descriptors.lock().expect("CLI descriptor store"); + let next_bytes = store.bytes.saturating_add(bytes.len() as u64); + if sequence > self.max_events || next_bytes > self.max_bytes { + self.failed.store(true, Ordering::Release); + return Err(ValidationError::new("/event_ref_sink", "capacity_exhausted").into()); + } + store.bytes = next_bytes; + store.entries.insert(payload_ref.clone(), bytes); + Ok(PayloadRef { + payload_ref: OpaqueRef::new(payload_ref), + content_digest: ContentDigest { + algorithm: DigestAlgorithm::Sha256, + value: digest, + }, + media_type: Some("application/json".to_string()), + }) + } +} + fn run_held_coalesce( set_path: &Path, request_path: &Path, @@ -611,8 +867,14 @@ fn read_raw(path: &Path) -> Result { #[cfg(test)] mod tests { - use super::{looks_like_uri, parse_duration, reject_non_local_path}; + use super::{ + looks_like_uri, parse_duration, reject_non_local_path, validate_watch_subject, + CliEventSink, DescriptorInspector, + }; use std::path::Path; + use std::sync::atomic::{AtomicBool, Ordering}; + use std::sync::Arc; + use waitprims_fs::{EventDescriptor, EventRefSink, FileEventClass}; #[test] fn millisecond_duration_is_preserved_not_truncated() { @@ -651,6 +913,74 @@ mod tests { } } + #[test] + fn cli_sink_inspects_descriptor_and_digest_before_ref_creation() { + use sha2::{Digest, Sha256}; + + let inspected = Arc::new(AtomicBool::new(false)); + let callback_inspected = inspected.clone(); + let inspector: DescriptorInspector = + Arc::new(move |descriptor, bytes, digest, refs_created| { + assert_eq!(refs_created, 0, "inspection must precede ref creation"); + assert_eq!(descriptor.class, FileEventClass::Create); + assert_eq!(descriptor.paths, ["watched/leaf"]); + assert_eq!(bytes, descriptor.canonical_bytes()); + let expected = Sha256::digest(bytes) + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + assert_eq!(digest, expected); + callback_inspected.store(true, Ordering::Release); + }); + let sink = CliEventSink { + inspector: Some(inspector), + ..CliEventSink::default() + }; + let payload = sink + .materialize(&EventDescriptor { + class: FileEventClass::Create, + paths: vec!["watched/leaf".to_string()], + }) + .expect("materialize"); + + assert!(inspected.load(Ordering::Acquire)); + assert_eq!(payload.payload_ref.as_str(), "ref:fs-cli-1"); + assert_eq!( + sink.descriptors + .lock() + .expect("descriptor store") + .entries + .get("ref:fs-cli-1") + .expect("stored descriptor"), + br#"{"class":"create","paths":["watched/leaf"]}"# + ); + } + + #[test] + fn cli_sink_capacity_exhaustion_is_fail_closed() { + let sink = CliEventSink::new(1, 1024); + let descriptor = EventDescriptor { + class: FileEventClass::Create, + paths: vec!["watched/leaf".to_string()], + }; + sink.materialize(&descriptor).expect("first descriptor"); + let error = sink + .materialize(&descriptor) + .expect_err("second descriptor must exceed event capacity"); + assert!(sink.failed()); + assert_eq!(error.to_string(), "/event_ref_sink: capacity_exhausted"); + } + + #[test] + fn watch_subjects_are_portable_relative_paths() { + for accepted in [".", "leaf", "directory/leaf"] { + validate_watch_subject(accepted).expect(accepted); + } + for rejected in ["", "/absolute", "../escape", r"directory\leaf", "C:/drive"] { + validate_watch_subject(rejected).expect_err(rejected); + } + } + #[test] fn uri_shaped_values_are_rejected() { for raw in [ diff --git a/crates/waitprims-cli/tests/watch_cli.rs b/crates/waitprims-cli/tests/watch_cli.rs new file mode 100644 index 0000000..ef37abc --- /dev/null +++ b/crates/waitprims-cli/tests/watch_cli.rs @@ -0,0 +1,359 @@ +use std::io::{BufRead, BufReader, Read, Write}; +use std::path::{Path, PathBuf}; +use std::process::{Command, Stdio}; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::mpsc; +use std::time::{Duration, Instant}; + +use waitprims_core::{registration_digest, validate_message}; +use waitprims_fs::{EventClock, SystemEventClock}; + +fn bin() -> Command { + Command::new(env!("CARGO_BIN_EXE_waitprims")) +} + +static TEMP_SEQUENCE: AtomicU64 = AtomicU64::new(1); + +struct TempRoot(PathBuf); + +impl TempRoot { + fn new() -> Self { + let sequence = TEMP_SEQUENCE.fetch_add(1, Ordering::Relaxed); + let path = std::env::temp_dir().join(format!( + "waitprims-watch-cli-{}-{sequence}", + std::process::id(), + )); + std::fs::create_dir(&path).expect("create temp root"); + Self(path) + } + + fn path(&self) -> &Path { + &self.0 + } +} + +impl Drop for TempRoot { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.0); + } +} + +fn write_watch_documents(root: &Path, methods_and_sources: &[(&str, &str)]) -> (PathBuf, PathBuf) { + let now = SystemEventClock.now(); + let lease = now.saturating_add(Duration::from_secs(12)); + let run_deadline = now.saturating_add(Duration::from_secs(3)); + let logical_deadline = now.saturating_add(Duration::from_secs(10)); + let registrations = methods_and_sources + .iter() + .enumerate() + .map(|(index, (method, source))| { + serde_json::json!({ + "registration_id": format!("reg:fs-demo-{index}"), + "method_id": method, + "subject_kind": "path", + "subject_id": "watched-leaf", + "baseline_policy": "latest", + "required": true, + "source_instance_ref": source, + "predicate_ref": "pred:file-any", + "capability_ref": "cap:fs-demo", + "lease_expires_at": lease, + "bounds": {"max_events": 4096, "max_bytes": 4096} + }) + }) + .collect::>(); + let digest = + registration_digest(&serde_json::to_string(®istrations).expect("registrations JSON")) + .expect("registration digest"); + let set = serde_json::json!({ + "capabilities": ["contract: agent-wait/v0"], + "message_type": "registration_set", + "message_id": "msg:fs-demo-set", + "correlation_id": "corr:fs-demo", + "created_at": now, + "actor_ref": "seat:fs-demo", + "waiter_id": "waiter:fs-demo", + "seat_ref": "seat:fs-demo", + "registration_revision": "regrev-fs-demo", + "registrations": registrations, + "principal_ref": "seat:fs-demo", + "logical_deadline": logical_deadline, + "authn_mode": "optional", + "aggregate_limits": {"max_events": 4096, "max_bytes": 1048576}, + "registration_digest": { + "canonicalization": "rfc8785", + "algorithm": "sha256", + "value": digest + } + }); + let request = serde_json::json!({ + "capabilities": ["contract: agent-wait/v0"], + "message_type": "live_wait_request", + "message_id": "msg:fs-demo-request", + "correlation_id": "corr:fs-demo", + "created_at": now, + "actor_ref": "seat:fs-demo", + "causation_id": "msg:fs-demo-set", + "waiter_id": "waiter:fs-demo", + "registration_set_ref": "msg:fs-demo-set", + "registration_revision": "regrev-fs-demo", + "logical_deadline": logical_deadline, + "run_deadline": run_deadline + }); + let set_path = root.join("registration_set.json"); + let request_path = root.join("live_wait_request.json"); + std::fs::write( + &set_path, + serde_json::to_vec_pretty(&set).expect("set JSON"), + ) + .expect("write set"); + std::fs::write( + &request_path, + serde_json::to_vec_pretty(&request).expect("request JSON"), + ) + .expect("write request"); + (set_path, request_path) +} + +fn rewrite_first_subject_kind(set_path: &Path, subject_kind: &str) { + let mut set: serde_json::Value = + serde_json::from_slice(&std::fs::read(set_path).expect("read set")).expect("set JSON"); + set["registrations"][0]["subject_kind"] = serde_json::json!(subject_kind); + let digest = registration_digest(&set["registrations"].to_string()).expect("digest"); + set["registration_digest"]["value"] = serde_json::json!(digest); + std::fs::write(set_path, serde_json::to_vec_pretty(&set).expect("set JSON")) + .expect("rewrite set"); +} + +#[test] +fn watch_help_shows_only_the_native_input_flags() { + let output = bin() + .args(["watch", "--help"]) + .output() + .expect("watch help"); + assert_eq!(output.status.code(), Some(0)); + let stdout = String::from_utf8_lossy(&output.stdout); + for flag in ["--root", "--registration-set", "--request"] { + assert!(stdout.contains(flag), "missing {flag}: {stdout}"); + } + for forbidden in ["--script", "--cancel", "--posture"] { + assert!( + !stdout.contains(forbidden), + "unexpected {forbidden}: {stdout}" + ); + } +} + +#[test] +fn watch_rejects_nonlocal_inputs_with_zero_stdout() { + let temp = TempRoot::new(); + let (set, request) = write_watch_documents(temp.path(), &[("file_watch", "source:fs-demo")]); + let cases = [ + ( + "https://example.invalid/root".to_string(), + set.to_string_lossy().into_owned(), + request.to_string_lossy().into_owned(), + ), + ( + temp.path().to_string_lossy().into_owned(), + "-".to_string(), + request.to_string_lossy().into_owned(), + ), + ( + temp.path().to_string_lossy().into_owned(), + set.to_string_lossy().into_owned(), + "file:///tmp/request.json".to_string(), + ), + ]; + for (root, set, request) in cases { + let output = bin() + .args([ + "watch", + "--root", + &root, + "--registration-set", + &set, + "--request", + &request, + ]) + .output() + .expect("watch rejection"); + assert_eq!(output.status.code(), Some(1)); + assert!(output.stdout.is_empty()); + let stderr = String::from_utf8_lossy(&output.stderr); + assert!(stderr.contains("local_path_required"), "{stderr}"); + assert!(!stderr.contains("example.invalid"), "{stderr}"); + } +} + +#[test] +fn watch_rejects_mixed_method_and_source_before_stdout() { + let cases = [ + vec![ + ("file_watch", "source:fs-demo"), + ("sms_inbound", "source:fs-demo"), + ], + vec![ + ("file_watch", "source:fs-demo"), + ("file_watch", "source:other"), + ], + ]; + for registrations in cases { + let temp = TempRoot::new(); + let (set, request) = write_watch_documents(temp.path(), ®istrations); + let output = bin() + .arg("watch") + .arg("--root") + .arg(temp.path()) + .arg("--registration-set") + .arg(set) + .arg("--request") + .arg(request) + .output() + .expect("watch mixed input"); + assert_eq!(output.status.code(), Some(1)); + assert!(output.stdout.is_empty()); + } +} + +#[test] +fn watch_rejects_non_path_subject_kind_before_stdout() { + let temp = TempRoot::new(); + let (set, request) = write_watch_documents(temp.path(), &[("file_watch", "source:fs-demo")]); + rewrite_first_subject_kind(&set, "inbox"); + let output = bin() + .arg("watch") + .arg("--root") + .arg(temp.path()) + .arg("--registration-set") + .arg(set) + .arg("--request") + .arg(request) + .output() + .expect("watch non-path subject kind"); + assert_eq!(output.status.code(), Some(1)); + assert!(output.stdout.is_empty()); + let stderr = String::from_utf8_lossy(&output.stderr); + assert!(stderr.contains("path_required"), "{stderr}"); + assert!(!stderr.contains("inbox"), "{stderr}"); +} + +#[test] +fn native_watch_demo_uses_bounded_retry_and_visible_event_surface() { + let temp = TempRoot::new(); + let (set, request) = write_watch_documents(temp.path(), &[("file_watch", "source:fs-demo")]); + let mut child = bin() + .args(["--log-level", "error", "watch", "--root"]) + .arg(temp.path()) + .arg("--registration-set") + .arg(set) + .arg("--request") + .arg(request) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .expect("spawn watch"); + let stdout = child.stdout.take().expect("stdout"); + let (sender, receiver) = mpsc::channel(); + let reader = std::thread::spawn(move || { + for line in BufReader::new(stdout).lines() { + if sender.send(line.expect("stdout line")).is_err() { + break; + } + } + }); + + let retry_deadline = Instant::now() + Duration::from_secs(6); + let leaf = temp.path().join("watched-leaf"); + std::fs::File::create(&leaf).expect("create watched leaf"); + let mut lines = Vec::new(); + let mut observed_burst = false; + while Instant::now() < retry_deadline && !observed_burst { + while let Ok(line) = receiver.try_recv() { + observed_burst |= line.contains("\"diagnostic_type\":\"follow_burst\""); + lines.push(line); + } + if observed_burst { + break; + } + let mut file = std::fs::OpenOptions::new() + .write(true) + .truncate(true) + .open(&leaf) + .expect("retry open"); + file.write_all(b"retry").expect("retry write"); + file.sync_all().expect("file barrier"); + let _ = std::fs::read_dir(temp.path()) + .expect("root barrier") + .collect::>>() + .expect("root entries"); + match receiver.recv_timeout(Duration::from_millis(25)) { + Ok(line) => { + observed_burst |= line.contains("\"diagnostic_type\":\"follow_burst\""); + lines.push(line); + } + Err(mpsc::RecvTimeoutError::Timeout) => {} + Err(mpsc::RecvTimeoutError::Disconnected) => break, + } + } + assert!( + observed_burst, + "bounded retry exhausted without a burst; lines={lines:?}" + ); + + let exit_deadline = Instant::now() + Duration::from_secs(6); + let status = loop { + if let Some(status) = child.try_wait().expect("child status") { + break status; + } + if Instant::now() >= exit_deadline { + let _ = child.kill(); + panic!("watch did not finish at its request deadline"); + } + std::thread::yield_now(); + }; + std::fs::remove_file(&leaf).expect("remove watched leaf after child exit"); + reader.join().expect("stdout reader"); + lines.extend(receiver.try_iter()); + let mut stderr = String::new(); + child + .stderr + .take() + .expect("stderr") + .read_to_string(&mut stderr) + .expect("read stderr"); + assert_eq!(status.code(), Some(0), "stderr={stderr} lines={lines:?}"); + + let records = lines + .iter() + .map(|line| serde_json::from_str::(line).expect("JSONL record")) + .collect::>(); + let bursts = records + .iter() + .filter(|record| record["diagnostic_type"] == "follow_burst") + .collect::>(); + assert!(!bursts.is_empty(), "{records:?}"); + assert_eq!( + records + .iter() + .filter(|record| record["diagnostic_type"] == "follow_end") + .count(), + 1 + ); + assert_eq!( + records.last().expect("last record")["diagnostic_type"], + "follow_end" + ); + let event = &bursts[0]["events"][0]; + assert_eq!(event["method_id"], "file_watch"); + assert_eq!(event["subject_id"], "watched-leaf"); + let combined = lines.join("\n"); + assert!( + !combined.contains(temp.path().to_string_lossy().as_ref()), + "{combined}" + ); + for (line, record) in lines.iter().zip(&records) { + assert!(record.get("message_type").is_none(), "{line}"); + validate_message(line).expect_err("diagnostic record must not admit as wire"); + } +}