Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
307 changes: 299 additions & 8 deletions crates/daemon/tests/upgrade.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::*;

Expand All @@ -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<bool>,
closed: Option<Instant>,
reopened: Option<Instant>,
count: u64,
last: Option<Instant>,
maximum_gap: Duration,
read_errors: Vec<String>,
}

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<AtomicBool>,
handle: Option<JoinHandle<MarkerSamples>>,
startup_problem: Option<String>,
}

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<MarkerSamples, String> {
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::<Vec<_>>();
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.
Expand Down Expand Up @@ -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");
Expand All @@ -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()
);
}

Expand Down