Skip to main content

mj_core/targets/
diagnostics.rs

1//! Diagnostic snapshots only; never use these to decide lifecycle admission.
2
3use std::collections::BTreeMap;
4use std::sync::atomic::{AtomicU64, Ordering};
5use std::sync::{Mutex, OnceLock};
6use std::time::Instant;
7
8use super::CommandSpec;
9
10#[derive(Debug, Clone)]
11pub struct BlockingOperationSnapshot {
12    pub purpose: String,
13    pub program: String,
14    pub thread: String,
15    pub elapsed_ms: u64,
16}
17
18struct ActiveOperation {
19    purpose: String,
20    program: String,
21    thread: String,
22    started: Instant,
23}
24
25fn operations() -> &'static Mutex<BTreeMap<u64, ActiveOperation>> {
26    static OPERATIONS: OnceLock<Mutex<BTreeMap<u64, ActiveOperation>>> = OnceLock::new();
27    OPERATIONS.get_or_init(Mutex::default)
28}
29
30/// Names a synchronous wait before it starts, including SSH admission and
31/// master preparation. Arguments, environment and stdin are never retained.
32pub struct BlockingOperation(u64);
33
34impl BlockingOperation {
35    pub fn start(purpose: &str, program: &str) -> Self {
36        static NEXT_ID: AtomicU64 = AtomicU64::new(1);
37        let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
38        let thread = std::thread::current();
39        operations()
40            .lock()
41            .unwrap_or_else(std::sync::PoisonError::into_inner)
42            .insert(
43                id,
44                ActiveOperation {
45                    purpose: purpose.to_owned(),
46                    program: program.to_owned(),
47                    thread: format!("{} {:?}", thread.name().unwrap_or("unnamed"), thread.id()),
48                    started: Instant::now(),
49                },
50            );
51        Self(id)
52    }
53
54    pub(super) fn command(command: &CommandSpec) -> Self {
55        Self::start(&command.purpose, &command.program)
56    }
57}
58
59impl Drop for BlockingOperation {
60    fn drop(&mut self) {
61        operations()
62            .lock()
63            .unwrap_or_else(std::sync::PoisonError::into_inner)
64            .remove(&self.0);
65    }
66}
67
68/// Never wait for the registry lock from a stall reporter. `None` means the
69/// snapshot was contended, rather than that no work was running.
70pub fn active_blocking_operations() -> Option<Vec<BlockingOperationSnapshot>> {
71    let operations = match operations().try_lock() {
72        Ok(operations) => operations,
73        Err(std::sync::TryLockError::Poisoned(error)) => error.into_inner(),
74        Err(std::sync::TryLockError::WouldBlock) => return None,
75    };
76    Some(
77        operations
78            .values()
79            .map(|operation| BlockingOperationSnapshot {
80                purpose: operation.purpose.clone(),
81                program: operation.program.clone(),
82                thread: operation.thread.clone(),
83                elapsed_ms: operation.started.elapsed().as_millis() as u64,
84            })
85            .collect(),
86    )
87}
88
89#[cfg(test)]
90mod tests {
91    use super::*;
92
93    #[test]
94    fn diagnostics_keep_concurrent_operations_until_each_owner_finishes() {
95        let first = BlockingOperation::start("diagnostics-first", "ssh");
96        let second = BlockingOperation::start("diagnostics-second", "git");
97        let snapshot = || {
98            operations()
99                .lock()
100                .unwrap()
101                .values()
102                .filter(|operation| operation.purpose.starts_with("diagnostics-"))
103                .map(|operation| operation.purpose.clone())
104                .collect::<Vec<_>>()
105        };
106        assert_eq!(snapshot(), ["diagnostics-first", "diagnostics-second"]);
107        drop(first);
108        assert_eq!(snapshot(), ["diagnostics-second"]);
109        drop(second);
110        assert!(snapshot().is_empty());
111    }
112}