saddle-runtime 0.3.27

Saddle managed asynchronous runtime and lifecycle
Documentation
//! Public OS/Tokio counters sampled on the existing managed executor.
//! No task, thread, request permit, or private scheduler is created here.
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>,
    /// The instant at which the gate state was read, relative to this sampler.
    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,
}

/// The caller reserves this value's storage and polls it on the same Runtime.
/// A snapshot does not claim cgroup exclusivity or OS CPU enforcement.
#[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);
        // Missing/stale evidence cannot become a low-pressure measurement.
        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;
    }
}

/// The meter is destroyed before its process reservation is refunded.
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();
            // A real process-clock/Runtime sampling gap, not a synthetic invalid flag.
            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);
            });
    }
}