diff --git a/crates/daemon/tests/upgrade.rs b/crates/daemon/tests/upgrade.rs index 851cd09..830ccd4 100644 --- a/crates/daemon/tests/upgrade.rs +++ b/crates/daemon/tests/upgrade.rs @@ -21,12 +21,14 @@ use devguard_daemon::config::HostConfig; use devguard_daemon::install::{self, InstallOptions, Manifest, ServiceManager, ServiceSpec}; use devguard_daemon::paths::{read_private, AuthorityPaths}; use devguard_daemon::upgrade::{self, DrainMode, UpgradeOptions, UpgradeOutcome, UpgradeReport}; -use serde_json::json; +use serde_json::{json, Value}; use std::fs; use std::io::{Seek, SeekFrom, Write}; use std::os::unix::fs::PermissionsExt; -use std::path::PathBuf; -use std::sync::atomic::AtomicBool; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{mpsc, Arc}; +use std::thread::{self, JoinHandle}; use std::time::{Duration, Instant}; use support::*; @@ -37,6 +39,278 @@ fn daemon_child() { } const MIB: u64 = 1024 * 1024; +const MARKER_SAMPLE_INTERVAL: Duration = Duration::from_millis(10); +const MARKER_WATCHER_STOP_LIMIT: Duration = Duration::from_secs(1); + +#[derive(Default)] +struct MarkerSamples { + initial_present: Option, + closed: Option, + reopened: Option, + count: u64, + last: Option, + maximum_gap: Duration, + read_errors: Vec, +} + +impl MarkerSamples { + fn observe(&mut self, marker: &Path) { + let observed = Instant::now(); + if let Some(last) = self.last { + self.maximum_gap = self.maximum_gap.max(observed.duration_since(last)); + } + self.last = Some(observed); + self.count += 1; + let present = match fs::symlink_metadata(marker) { + Ok(_) => Some(true), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Some(false), + Err(error) => { + let message = format!("cannot sample the admission marker: {error}"); + if !self.read_errors.contains(&message) { + self.read_errors.push(message); + } + None + } + }; + if self.count == 1 { + self.initial_present = present; + } + match present { + Some(true) if self.closed.is_none() => self.closed = Some(observed), + Some(false) if self.closed.is_some() && self.reopened.is_none() => { + self.reopened = Some(observed); + } + _ => {} + } + } +} + +/// Diagnostic sampling only: the marker is polled every 10 ms, so these +/// observations are not an exact clock for the close or reopen operations. +struct MarkerWatcher { + stop: Arc, + handle: Option>, + startup_problem: Option, +} + +impl MarkerWatcher { + fn start(marker: PathBuf) -> Self { + let stop = Arc::new(AtomicBool::new(false)); + let stopped = Arc::clone(&stop); + let (ready, started) = mpsc::channel(); + let handle = thread::spawn(move || { + let mut samples = MarkerSamples::default(); + let mut ready = Some(ready); + loop { + samples.observe(&marker); + if let Some(ready) = ready.take() { + let _ = ready.send(()); + } + if stopped.load(Ordering::Acquire) { + // This observation is guaranteed to start after the stop + // request was seen, reducing the chance of missing a + // marker removal just before upgrade returned. + samples.observe(&marker); + break; + } + thread::sleep(MARKER_SAMPLE_INTERVAL); + } + samples + }); + let startup_problem = started.recv_timeout(MARKER_WATCHER_STOP_LIMIT).err().map(|error| { + format!( + "marker watcher did not report its first sample within {MARKER_WATCHER_STOP_LIMIT:?}: {error}" + ) + }); + Self { + stop, + handle: Some(handle), + startup_problem, + } + } + + fn stop_and_join(&mut self) -> Result { + self.stop.store(true, Ordering::Release); + let deadline = Instant::now() + MARKER_WATCHER_STOP_LIMIT; + while self + .handle + .as_ref() + .is_some_and(|handle| !handle.is_finished()) + && Instant::now() < deadline + { + thread::sleep(Duration::from_millis(1)); + } + let finished = self.handle.as_ref().is_some_and(JoinHandle::is_finished); + if !finished { + // Dropping the handle detaches it; the stop flag remains set, so + // the thread exits at its next bounded sample. Never block a test + // indefinitely while collecting diagnostics. + self.handle.take(); + return Err(format!( + "marker watcher did not stop within {MARKER_WATCHER_STOP_LIMIT:?}" + )); + } + self.handle + .take() + .expect("a finished watcher has a handle") + .join() + .map_err(|_| "marker watcher panicked".to_string()) + } + + fn finish(mut self, upgrade_started: Instant, upgrade_elapsed: Duration) -> Value { + let returned = upgrade_started + upgrade_elapsed; + let mut missing = Vec::new(); + let mut diagnostics = self.startup_problem.take().into_iter().collect::>(); + let samples = match self.stop_and_join() { + Ok(samples) => samples, + Err(error) => { + diagnostics.push(error.clone()); + for field in [ + "initial_marker_present", + "sample_count", + "observed_max_sample_gap_ms", + "measurement_resolution_ms", + "marker_present_after_start_ms", + "marker_absent_after_start_ms", + "marker_observed_present_for_ms", + "marker_absence_sample_to_return_ms", + ] { + missing.push(json!({"field": field, "reason": error.clone()})); + } + return json!({ + "kind": "diagnostic_sampling", + "exact_latency_clock": false, + "sample_interval_ms": MARKER_SAMPLE_INTERVAL.as_millis(), + "measurement_resolution_ms": null, + "observed_max_sample_gap_ms": null, + "resolution_note": "polling and scheduler delays make transition observations no more precise than the observed sample gap", + "watcher_stop_limit_ms": MARKER_WATCHER_STOP_LIMIT.as_millis(), + "sample_count": null, + "initial_marker_present": null, + "total_upgrade_ms": upgrade_elapsed.as_millis(), + "marker_present_after_start_ms": null, + "marker_absent_after_start_ms": null, + "marker_observed_present_for_ms": null, + "marker_absence_sample_to_return_ms": null, + "missing": missing, + "diagnostics": diagnostics, + }); + } + }; + diagnostics.extend(samples.read_errors); + if samples.initial_present.is_none() { + missing.push(json!({ + "field": "initial_marker_present", + "reason": "the first marker metadata read failed" + })); + } + let observed_max_gap = (samples.count > 1).then_some(samples.maximum_gap); + if observed_max_gap.is_none() { + missing.push(json!({ + "field": "observed_max_sample_gap_ms", + "reason": "fewer than two marker samples completed" + })); + missing.push(json!({ + "field": "measurement_resolution_ms", + "reason": "fewer than two marker samples completed" + })); + } + let measurement_resolution = observed_max_gap.map(|gap| gap.max(MARKER_SAMPLE_INTERVAL)); + let present = match samples.closed { + Some(instant) => match instant.checked_duration_since(upgrade_started) { + Some(duration) => Some(duration), + None => { + missing.push(json!({ + "field": "marker_present_after_start_ms", + "reason": "the marker was already observed before upgrade timing began" + })); + None + } + }, + None => { + missing.push(json!({ + "field": "marker_present_after_start_ms", + "reason": "marker presence was not observed at the recorded sampling resolution" + })); + None + } + }; + let absent = match samples.reopened { + Some(instant) => match instant.checked_duration_since(upgrade_started) { + Some(duration) => Some(duration), + None => { + missing.push(json!({ + "field": "marker_absent_after_start_ms", + "reason": "marker absence was sampled before upgrade timing began" + })); + None + } + }, + None => { + missing.push(json!({ + "field": "marker_absent_after_start_ms", + "reason": "marker absence after observed presence was not sampled" + })); + None + } + }; + let present_for = samples + .reopened + .zip(samples.closed) + .and_then(|(absent, present)| absent.checked_duration_since(present)); + if present_for.is_none() { + missing.push(json!({ + "field": "marker_observed_present_for_ms", + "reason": "both ordered marker transitions were not sampled" + })); + } + let absence_to_return = match samples.reopened { + Some(absent) => match returned.checked_duration_since(absent) { + Some(duration) => Some(duration), + None => { + missing.push(json!({ + "field": "marker_absence_sample_to_return_ms", + "reason": "the marker-absence sample followed the return instant; diagnostic sampling cannot order them more precisely" + })); + None + } + }, + None => { + missing.push(json!({ + "field": "marker_absence_sample_to_return_ms", + "reason": "marker absence after observed presence was not sampled" + })); + None + } + }; + json!({ + "kind": "diagnostic_sampling", + "exact_latency_clock": false, + "sample_interval_ms": MARKER_SAMPLE_INTERVAL.as_millis(), + "measurement_resolution_ms": measurement_resolution.map(|duration| duration.as_millis()), + "observed_max_sample_gap_ms": observed_max_gap.map(|duration| duration.as_millis()), + "resolution_note": "polling and scheduler delays make transition observations no more precise than the observed sample gap", + "watcher_stop_limit_ms": MARKER_WATCHER_STOP_LIMIT.as_millis(), + "sample_count": samples.count, + "initial_marker_present": samples.initial_present, + "total_upgrade_ms": upgrade_elapsed.as_millis(), + "marker_present_after_start_ms": present.map(|duration| duration.as_millis()), + "marker_absent_after_start_ms": absent.map(|duration| duration.as_millis()), + "marker_observed_present_for_ms": present_for.map(|duration| duration.as_millis()), + "marker_absence_sample_to_return_ms": absence_to_return.map(|duration| duration.as_millis()), + "missing": missing, + "diagnostics": diagnostics, + }) + } +} + +impl Drop for MarkerWatcher { + fn drop(&mut self) { + if self.handle.is_some() { + let _ = self.stop_and_join(); + } + } +} /// An installed fixture service. The manager is dropped first, stopping the /// service before the base directory is removed. @@ -363,13 +637,17 @@ fn a_drain_that_does_not_finish_in_time_keeps_the_current_release_and_its_charge let running = committed(paths, "running"); installed.stage("b", "0.1.0-test-b"); let before = installed.pid(); + let watcher = MarkerWatcher::start(paths.admission_marker()); let started = Instant::now(); - let error = installed - .upgrade("0.1.0-test-b", Duration::from_secs(1), false) - .unwrap_err(); + let outcome = installed.upgrade("0.1.0-test-b", Duration::from_secs(1), false); + // Capture the total duration immediately when upgrade returns. The marker + // watcher is diagnostic sampling and is stopped only after this capture. + let elapsed = started.elapsed(); + let marker_timing = watcher.finish(started, elapsed); + let error = outcome.unwrap_err(); + eprintln!("drain-timeout upgrade returned in {elapsed:?}; marker timing: {marker_timing}"); assert_eq!(error.code, ErrorCode::ResourceUnavailable, "{error:?}"); assert!(error.message.contains("did not finish"), "{error:?}"); - assert!(started.elapsed() < Duration::from_secs(10)); // Nothing was replaced or backed up; the charge is kept and admission reopened. assert_eq!(installed.pid(), before); assert_eq!(installed.current(), "0.1.0-test-a"); @@ -387,7 +665,20 @@ fn a_drain_that_does_not_finish_in_time_keeps_the_current_release_and_its_charge assert_eq!(report.selection.current, "0.1.0-test-b"); record( "drain-timeout", - json!({"error": error, "kept": quiescence, "retried": report}), + json!({ + "error": error, + "kept": quiescence, + "retried": report, + "timing": marker_timing.clone(), + }), + ); + // Operability and safety assertions above must run even when a busy host + // exceeds this guard. Keep the existing bound until phase evidence shows + // whether it measures drain/reopen behavior or unrelated preparation. + assert!( + elapsed < Duration::from_secs(10), + "drain-timeout upgrade returned in {elapsed:?} ({} ms); marker timing: {marker_timing}", + elapsed.as_millis() ); }