use saddle_admission::pressure::{PressureGate, PressureState, PressureWindow};
use std::time::{Duration, Instant};
pub const SAMPLE_INTERVAL: Duration = Duration::from_millis(100);
const WINDOW: Duration = Duration::from_secs(1);
#[derive(Clone, Copy, Debug)]
pub struct PressureSnapshot {
pub os_error: Option<i32>,
pub window: Option<PressureWindow>,
pub decision_at: Duration,
pub state: PressureState,
pub workers: usize,
pub global_queue_depth: usize,
pub source: &'static str,
pub invalid_reason: Option<&'static str>,
}
#[derive(Clone, Copy)]
struct Counters {
at: Instant,
cpu: Duration,
busy: Duration,
}
#[derive(Clone, Copy, Debug)]
pub struct PressureReadError {
pub reason: &'static str,
pub os_error: Option<i32>,
}
impl From<&'static str> for PressureReadError {
fn from(reason: &'static str) -> Self {
Self {
reason,
os_error: None,
}
}
}
pub struct PressureSampler {
metrics: tokio::runtime::RuntimeMetrics,
origin: Instant,
next: Instant,
baseline: Counters,
previous: Counters,
delay: Duration,
cpu_cores: usize,
identity: u64,
needs_baseline: bool,
gate: PressureGate,
snapshot: PressureSnapshot,
}
impl PressureSampler {
pub fn new(cpu_cores: usize, identity: u64) -> Result<Self, PressureReadError> {
if cpu_cores == 0 || identity == 0 {
return Err("pressure.invalid_input".into());
}
let handle = tokio::runtime::Handle::try_current().map_err(|_| "pressure.no_runtime")?;
let metrics = handle.metrics();
if metrics.num_workers() == 0 {
return Err("pressure.no_workers".into());
}
let counters = read(&metrics)?;
Ok(Self {
snapshot: PressureSnapshot {
os_error: None,
window: None,
decision_at: Duration::ZERO,
state: PressureState::Cold,
workers: metrics.num_workers(),
global_queue_depth: metrics.global_queue_depth(),
source: "process_cpu_clock+tokio_workers+sampler_delay",
invalid_reason: None,
},
metrics,
origin: counters.at,
next: counters.at + SAMPLE_INTERVAL,
baseline: counters,
previous: counters,
delay: Duration::ZERO,
cpu_cores,
identity,
needs_baseline: false,
gate: PressureGate::default(),
})
}
pub fn next_sample(&self) -> Instant {
self.next
}
pub fn sample(&mut self) -> PressureSnapshot {
let now = Instant::now();
self.snapshot.decision_at = now.duration_since(self.origin);
if self.snapshot.window.is_some() {
self.snapshot.state = self.gate.state_at(now.duration_since(self.origin));
if self.snapshot.state == PressureState::Invalid
&& self.snapshot.invalid_reason.is_none()
{
self.snapshot.invalid_reason = Some("pressure.stale_window");
}
}
if now < self.next {
return self.snapshot;
}
self.delay = self.delay.max(now.duration_since(self.next));
self.next = now + SAMPLE_INTERVAL;
self.snapshot.global_queue_depth = self.metrics.global_queue_depth();
let counters = match read(&self.metrics) {
Ok(counters) => counters,
Err(reason) => {
self.invalid(now, reason.reason);
self.snapshot.os_error = reason.os_error;
self.needs_baseline = true;
return self.snapshot;
}
};
if self.needs_baseline
|| counters.at.duration_since(self.previous.at) > Duration::from_secs(2)
{
self.invalid(now, "pressure.missing_window");
self.baseline = counters;
self.previous = counters;
self.delay = Duration::ZERO;
self.needs_baseline = false;
return self.snapshot;
}
if counters.cpu < self.previous.cpu || counters.busy < self.previous.busy {
self.invalid(now, "pressure.counter_regression");
self.baseline = counters;
self.previous = counters;
self.delay = Duration::ZERO;
return self.snapshot;
}
self.previous = counters;
let elapsed = counters.at.duration_since(self.baseline.at);
if elapsed < WINDOW {
return self.snapshot;
}
let window = PressureWindow {
observed_at: counters.at.duration_since(self.origin),
runtime_identity: self.identity,
process_cpu_ratio: (counters.cpu - self.baseline.cpu).as_secs_f64()
/ elapsed.as_secs_f64()
/ self.cpu_cores as f64,
worker_busy_ratio: (counters.busy - self.baseline.busy).as_secs_f64()
/ elapsed.as_secs_f64()
/ self.metrics.num_workers() as f64,
scheduling_delay: self.delay,
valid: true,
};
self.snapshot.state = self.gate.observe(window);
self.snapshot.window = Some(window);
self.snapshot.invalid_reason = None;
self.snapshot.os_error = None;
self.baseline = counters;
self.delay = Duration::ZERO;
self.snapshot.decision_at = Instant::now().duration_since(self.origin);
self.snapshot
}
fn invalid(&mut self, now: Instant, reason: &'static str) {
let window = PressureWindow {
observed_at: now.duration_since(self.origin),
runtime_identity: self.identity,
process_cpu_ratio: 0.0,
worker_busy_ratio: 0.0,
scheduling_delay: self.delay,
valid: false,
};
self.snapshot.state = self.gate.observe(window);
self.snapshot.window = Some(window);
self.snapshot.invalid_reason = Some(reason);
self.snapshot.os_error = None;
}
}
pub struct ManagedPressureSampler {
sampler: PressureSampler,
_storage: saddle_admission::StoragePermit,
}
impl ManagedPressureSampler {
pub(crate) fn prepare(
process_storage: &saddle_admission::StoragePermit,
cpu_cores: usize,
) -> Result<Self, saddle_admission::AdmissionError> {
use saddle_admission::{AdmissionError, StorageDemand};
let storage = process_storage.try_reserve(StorageDemand::separate(&[(
std::alloc::Layout::new::<PressureSampler>(),
1,
)])?)?;
let sampler = PressureSampler::new(cpu_cores, 1).map_err(|error| {
AdmissionError::PressureUnavailable {
reason: error.reason,
os_error: error.os_error,
}
})?;
Ok(Self {
sampler,
_storage: storage,
})
}
pub fn next_sample(&self) -> Instant {
self.sampler.next_sample()
}
pub fn sample(&mut self) -> PressureSnapshot {
self.sampler.sample()
}
}
fn read(metrics: &tokio::runtime::RuntimeMetrics) -> Result<Counters, PressureReadError> {
use rustix::time::{ClockId, DynamicClockId, clock_gettime_dynamic};
let raw =
clock_gettime_dynamic(DynamicClockId::Known(ClockId::ProcessCPUTime)).map_err(|error| {
PressureReadError {
reason: "pressure.process_clock_unavailable",
os_error: Some(error.raw_os_error()),
}
})?;
let cpu = Duration::new(
raw.tv_sec
.try_into()
.map_err(|_| "pressure.invalid_cpu_counter")?,
raw.tv_nsec
.try_into()
.map_err(|_| "pressure.invalid_cpu_counter")?,
);
let busy = (0..metrics.num_workers())
.try_fold(Duration::ZERO, |sum, index| {
sum.checked_add(metrics.worker_total_busy_duration(index))
})
.ok_or("pressure.worker_counter_overflow")?;
Ok(Counters {
at: Instant::now(),
cpu,
busy,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn real_missing_window_is_recorded_with_original_sampling_source() {
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("real_missing_window_diagnostic_child")
.arg("--ignored")
.env("SADDLE_I027_PRESSURE_CHILD", "1")
.status().unwrap();
assert!(status.success(), "isolated real pressure diagnostic readback");
}
#[test]
#[ignore = "run by the isolated parent test"]
fn real_missing_window_diagnostic_child() {
if std::env::var_os("SADDLE_I027_PRESSURE_CHILD").is_none() { return; }
tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap().block_on(async {
let directory = std::env::temp_dir().join(format!("saddle-i027-pressure-{}", std::process::id()));
std::fs::create_dir_all(&directory).unwrap();
let mut writer = saddle_observability::EmergencyDiagnostics::start(
&saddle_observability::FileLoggingConfig::new(&directory, saddle_observability::Rotation::Daily)
).unwrap();
let output = writer.handle();
let observer = saddle_observability::Observer::with_writer(Default::default(), std::io::sink()).unwrap();
let mut sampler = PressureSampler::new(1, 7).unwrap();
std::thread::sleep(Duration::from_millis(2100));
let sample = sampler.sample();
assert_eq!(sample.invalid_reason, Some("pressure.missing_window"));
assert_eq!(sample.state, PressureState::Invalid);
assert!(!sample.window.unwrap().valid);
let submitted = observer.record_execution_pressure(Some(&output), saddle_observability::pressure::ExecutionPressureFacts {
os_error: sample.os_error, window: sample.window.unwrap(), decision_at: sample.decision_at,
source: sample.source, state: sample.state, workers: sample.workers,
global_queue_depth: sample.global_queue_depth, invalid_reason: sample.invalid_reason,
});
assert!(matches!(submitted, saddle_observability::DiagnosticSubmission::Enqueued));
drop(output);
let until = Instant::now() + Duration::from_secs(3);
while writer.shutdown() != saddle_observability::DiagnosticShutdown::Finished {
assert!(Instant::now() < until);
tokio::time::sleep(Duration::from_millis(10)).await;
}
let text = std::fs::read_to_string(directory.join("saddle.emergency.log")).unwrap();
let record: serde_json::Value = text.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.find(|row| row["event"] == "framework.execution_pressure").unwrap();
assert_eq!(record["source"], "process_cpu_clock+tokio_workers+sampler_delay");
assert_eq!(record["original_error"], "pressure.missing_window");
assert_eq!(record["window_valid"], false);
assert_eq!(record["valid"], false);
assert_eq!(record["state"], "invalid");
assert!(record["decision_monotonic_ms"].as_u64().unwrap() >= 2100);
});
}
#[test]
fn stale_real_window_is_invalid_at_decision_time() {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async {
let mut sampler = PressureSampler::new(1, 7).unwrap();
let window = PressureWindow {
observed_at: Duration::from_millis(100),
runtime_identity: 7,
process_cpu_ratio: 0.1,
worker_busy_ratio: 0.1,
scheduling_delay: Duration::ZERO,
valid: true,
};
assert_eq!(sampler.gate.observe(window), PressureState::Open);
sampler.snapshot.window = Some(window);
sampler.snapshot.state = PressureState::Open;
sampler.origin = Instant::now() - Duration::from_secs(3);
sampler.next = Instant::now() + Duration::from_secs(1);
let snapshot = sampler.sample();
assert_eq!(snapshot.state, PressureState::Invalid);
assert_eq!(snapshot.invalid_reason, Some("pressure.stale_window"));
assert!(snapshot.decision_at > snapshot.window.unwrap().observed_at);
});
}
}