Skip to main content

microsandbox_metrics/
registry.rs

1//! POSIX shared-memory registry: open/create, slot reservation, sample writes,
2//! and snapshot reads.
3//!
4//! All cross-process synchronization happens through atomics inside the
5//! mapped region. The seqlock pattern protects per-sample bytes against
6//! torn reads without requiring a kernel mutex.
7
8use std::ffi::CString;
9use std::ptr::NonNull;
10use std::sync::Arc;
11use std::sync::atomic::{AtomicU8, AtomicU16, AtomicU64, Ordering};
12use std::thread;
13use std::time::{Duration, Instant};
14
15use chrono::{DateTime, TimeZone, Utc};
16
17use crate::layout::{
18    DEFAULT_CAPACITY, HEADER_SIZE, HEADER_STATE_INITIALIZING, HEADER_STATE_READY,
19    HEADER_STATE_UNINIT, Header, NAME_BYTES, REGISTRY_MAGIC, REGISTRY_VERSION, SAMPLE_FLAG_CPU,
20    SAMPLE_FLAG_MEMORY_AVAILABLE, SAMPLE_FLAG_MEMORY_HOST_RESIDENT, SAMPLE_FLAG_MEMORY_USED,
21    SAMPLE_FLAG_UPPER_FREE, SAMPLE_FLAG_UPPER_HOST_ALLOCATED, SAMPLE_FLAG_UPPER_USED, SLOT_ACTIVE,
22    SLOT_FREE, SLOT_RESERVED, SLOT_SIZE, SLOT_STALE, Slot, registry_size,
23};
24use crate::snapshot::{LiveMetric, LiveMetricState};
25use crate::{MetricsError, MetricsResult};
26#[cfg(target_os = "windows")]
27use windows_sys::Win32::Foundation::{
28    CloseHandle, ERROR_ACCESS_DENIED, ERROR_ALREADY_EXISTS, ERROR_FILE_NOT_FOUND, GetLastError,
29    HANDLE, INVALID_HANDLE_VALUE, STILL_ACTIVE,
30};
31#[cfg(target_os = "windows")]
32use windows_sys::Win32::System::Memory::{
33    CreateFileMappingW, FILE_MAP_ALL_ACCESS, MEMORY_MAPPED_VIEW_ADDRESS, MapViewOfFile,
34    OpenFileMappingW, PAGE_READWRITE, UnmapViewOfFile,
35};
36#[cfg(target_os = "windows")]
37use windows_sys::Win32::System::Threading::{
38    GetExitCodeProcess, OpenProcess, PROCESS_QUERY_LIMITED_INFORMATION,
39};
40
41//--------------------------------------------------------------------------------------------------
42// Constants
43//--------------------------------------------------------------------------------------------------
44
45const INIT_WAIT_TIMEOUT: Duration = Duration::from_secs(5);
46const INIT_POLL_INTERVAL: Duration = Duration::from_millis(5);
47const QUIESCE_WAIT_TIMEOUT: Duration = Duration::from_millis(100);
48
49//--------------------------------------------------------------------------------------------------
50// Types
51//--------------------------------------------------------------------------------------------------
52
53/// Caller-supplied data used to reserve a slot.
54#[derive(Clone, Debug)]
55pub struct ReserveSlot<'a> {
56    /// Catalog sandbox id.
57    pub sandbox_id: i32,
58    /// Sandbox name. Must fit in the slot's fixed inline name buffer.
59    pub name: &'a str,
60    /// Configured guest memory limit in bytes.
61    pub memory_limit_bytes: u64,
62}
63
64/// Caller-supplied data used to transition a reservation to active.
65#[derive(Clone, Debug)]
66pub struct ActivateSlot {
67    /// Slot index returned by [`MetricsRegistry::reserve`].
68    pub slot: u32,
69    /// Generation returned by [`MetricsRegistry::reserve`].
70    pub generation: u64,
71    /// Catalog run id of the running sandbox.
72    pub run_id: i32,
73    /// PID of the runtime process.
74    pub pid: i32,
75    /// Wall-clock time at which the sandbox started.
76    pub started_at: DateTime<Utc>,
77}
78
79/// Mode passed to [`MetricsRegistry::release`].
80#[derive(Clone, Copy, Debug, PartialEq, Eq)]
81pub enum ReleaseMode {
82    /// Mark the slot stale so its last sample stays visible until reuse.
83    Stale,
84    /// Mark the slot immediately free.
85    Free,
86}
87
88/// Sample bytes to write into a slot.
89#[derive(Clone, Copy, Debug)]
90pub struct SampleWrite {
91    /// Wall-clock time the sample was captured.
92    pub sampled_at: DateTime<Utc>,
93    /// Guest vCPU usage as a percentage when sourced by the VMM.
94    pub cpu_percent: Option<f32>,
95    /// Cumulative guest vCPU execution time across all vCPUs when sourced by the VMM.
96    pub vcpu_time_ns: Option<u64>,
97    /// Guest-used memory in bytes when sourced by the guest.
98    pub memory_bytes: Option<u64>,
99    /// Guest-available memory in bytes when reported by the guest.
100    pub memory_available_bytes: Option<u64>,
101    /// Host-resident guest memory in bytes for capacity diagnostics.
102    pub memory_host_resident_bytes: Option<u64>,
103    /// Cumulative disk bytes read.
104    pub disk_read_bytes: u64,
105    /// Cumulative disk bytes written.
106    pub disk_write_bytes: u64,
107    /// Cumulative network bytes received.
108    pub net_rx_bytes: u64,
109    /// Cumulative network bytes transmitted.
110    pub net_tx_bytes: u64,
111    /// Guest-visible OCI upper filesystem used bytes when available.
112    pub upper_used_bytes: Option<u64>,
113    /// Guest-visible OCI upper filesystem free bytes when available.
114    pub upper_free_bytes: Option<u64>,
115    /// Host-allocated bytes for the writable upper image when available.
116    pub upper_host_allocated_bytes: Option<u64>,
117}
118
119/// Shared-memory registry.
120///
121/// Cloneable handle (`Arc`-backed). Dropping the last clone unmaps the region
122/// but never `shm_unlink`s — the registry outlives every process.
123#[derive(Clone)]
124pub struct MetricsRegistry {
125    inner: Arc<RegistryInner>,
126}
127
128/// Reservation token returned by [`MetricsRegistry::reserve`].
129#[derive(Clone, Copy, Debug)]
130pub struct SlotReservation {
131    /// Slot index assigned to the reservation.
132    pub slot: u32,
133    /// Generation stamp paired with this allocation. Carries through every
134    /// subsequent state transition to prevent stale writers from corrupting
135    /// a reused slot.
136    pub generation: u64,
137}
138
139/// Per-slot writer handle held by the runtime process.
140#[derive(Clone)]
141pub struct MetricsSlotWriter {
142    registry: MetricsRegistry,
143    slot: u32,
144    generation: u64,
145}
146
147struct RegistryInner {
148    // Kept alive so a future explicit `unlink` API has the resolved name.
149    // Currently only used for the open path; suppress the dead-code warning.
150    #[allow(dead_code)]
151    name: CString,
152    mapping: MappedRegion,
153    capacity: u32,
154}
155
156struct MappedRegion {
157    ptr: NonNull<u8>,
158    #[cfg(unix)]
159    len: usize,
160    #[cfg(target_os = "windows")]
161    handle: HANDLE,
162}
163
164//--------------------------------------------------------------------------------------------------
165// Trait Implementations
166//--------------------------------------------------------------------------------------------------
167
168// Safety: the pointer aliases the same shared-memory region across all
169// clones, and every read/write goes through atomics or the seqlock helpers.
170unsafe impl Send for RegistryInner {}
171unsafe impl Sync for RegistryInner {}
172
173impl Drop for MappedRegion {
174    fn drop(&mut self) {
175        // Unmap only. The named backing object is managed separately: POSIX
176        // registries intentionally outlive individual mappings, while Windows
177        // pagefile mappings live as long as at least one process has a handle
178        // or mapped view open.
179        #[cfg(unix)]
180        unsafe {
181            libc::munmap(self.ptr.as_ptr().cast(), self.len);
182        }
183
184        #[cfg(target_os = "windows")]
185        unsafe {
186            UnmapViewOfFile(MEMORY_MAPPED_VIEW_ADDRESS {
187                Value: self.ptr.as_ptr().cast(),
188            });
189            CloseHandle(self.handle);
190        }
191    }
192}
193
194//--------------------------------------------------------------------------------------------------
195// Methods
196//--------------------------------------------------------------------------------------------------
197
198impl MetricsRegistry {
199    /// Open the registry by name, creating it if it does not exist.
200    pub fn open_or_create(name: &str, capacity: u32) -> MetricsResult<Self> {
201        if capacity == 0 {
202            return Err(MetricsError::Custom(
203                "registry capacity must be non-zero".into(),
204            ));
205        }
206        let cname = CString::new(name)
207            .map_err(|_| MetricsError::Custom("registry name contains NUL byte".into()))?;
208        let map_len = registry_size(capacity);
209
210        // Loop around the create/open race: if another process wins creation
211        // between our open attempt and `O_EXCL`, retry the open path. Bound
212        // the loop so a pathological environment cannot spin forever.
213        const MAX_ATTEMPTS: u32 = 4;
214        for _ in 0..MAX_ATTEMPTS {
215            if let Some(reg) = try_open_existing(&cname, capacity, map_len)? {
216                return Ok(reg);
217            }
218            match create_and_init(&cname, capacity, map_len) {
219                Ok(reg) => return Ok(reg),
220                Err(MetricsError::AlreadyExists) => continue,
221                Err(other) => return Err(other),
222            }
223        }
224        Err(MetricsError::Custom(
225            "failed to open or create metrics registry after multiple attempts".into(),
226        ))
227    }
228
229    /// Open an existing registry. Errors if it has not yet been created.
230    pub fn open(name: &str) -> MetricsResult<Self> {
231        let cname = CString::new(name)
232            .map_err(|_| MetricsError::Custom("registry name contains NUL byte".into()))?;
233
234        // Two-pass: first map with the header-only length to discover the
235        // capacity, then remap with the full length.
236        let header_only_len = HEADER_SIZE;
237        let header_mapping = open_existing_region(&cname, header_only_len)?.ok_or_else(|| {
238            MetricsError::from(std::io::Error::new(
239                std::io::ErrorKind::NotFound,
240                "metrics registry does not exist",
241            ))
242        })?;
243        let header = unsafe { &*(header_mapping.ptr.as_ptr() as *const Header) };
244        let ready_result = wait_for_ready(header);
245        let capacity = match ready_result {
246            Ok(()) => match validate_header(header, None) {
247                Ok(()) => header.capacity,
248                Err(err) => return Err(err),
249            },
250            Err(WaitForReadyError::Stuck) => {
251                return Err(MetricsError::Custom(
252                    "metrics registry is still initializing".into(),
253                ));
254            }
255            Err(WaitForReadyError::Invalid(state)) => {
256                return Err(MetricsError::Custom(format!(
257                    "invalid registry header state: {state}"
258                )));
259            }
260        };
261        drop(header_mapping);
262
263        let map_len = registry_size(capacity);
264        let reg = try_open_existing(&cname, capacity, map_len)?
265            .ok_or_else(|| MetricsError::Custom("registry disappeared during open".into()))?;
266        Ok(reg)
267    }
268
269    /// Reserve a slot for an upcoming sandbox spawn.
270    pub fn reserve(&self, spec: ReserveSlot<'_>) -> MetricsResult<SlotReservation> {
271        if spec.name.len() > NAME_BYTES {
272            return Err(MetricsError::Custom(format!(
273                "sandbox name is too long for metrics slot: {} bytes (max {NAME_BYTES})",
274                spec.name.len()
275            )));
276        }
277
278        let capacity = self.inner.capacity;
279        // Scan slots once for a Free entry; fall back to a second pass that
280        // also reclaims Stale entries. Only when the registry is otherwise
281        // full do we reclaim Active slots whose owner PID is gone.
282        for pass in 0..3 {
283            for idx in 0..capacity {
284                let slot = self.slot(idx);
285                let mut current = slot.state.load(Ordering::Acquire);
286                let claimable = match (pass, current) {
287                    (_, SLOT_FREE) | (1, SLOT_STALE) => true,
288                    (2, SLOT_ACTIVE) => {
289                        let pid = slot.pid.load(Ordering::Acquire);
290                        if pid <= 0 || pid_is_alive(pid) {
291                            false
292                        } else {
293                            let generation = slot.generation.load(Ordering::Acquire);
294                            if self
295                                .release_inner(idx, generation, ReleaseMode::Free, true)
296                                .is_ok()
297                            {
298                                current = slot.state.load(Ordering::Acquire);
299                                current == SLOT_FREE
300                            } else {
301                                false
302                            }
303                        }
304                    }
305                    _ => false,
306                };
307                if !claimable {
308                    continue;
309                }
310                if slot
311                    .state
312                    .compare_exchange(current, SLOT_RESERVED, Ordering::AcqRel, Ordering::Acquire)
313                    .is_ok()
314                {
315                    let generation = self.next_generation();
316                    write_reservation_fields(slot, &spec, generation);
317                    return Ok(SlotReservation {
318                        slot: idx,
319                        generation,
320                    });
321                }
322            }
323        }
324        Err(MetricsError::Full)
325    }
326
327    /// Promote a reservation to an active writer.
328    pub fn activate_writer(&self, spec: ActivateSlot) -> MetricsResult<MetricsSlotWriter> {
329        let slot = self.try_slot(spec.slot)?;
330        let observed = slot.generation.load(Ordering::Acquire);
331        if observed != spec.generation {
332            return Err(MetricsError::GenerationMismatch {
333                expected: spec.generation,
334                actual: observed,
335            });
336        }
337
338        // Seq starts even; bump odd, write metadata, bump even again. We
339        // write inside the seqlock window so a reader that observes Active
340        // sees a coherent run_id/pid/started_at snapshot.
341        let begin = begin_write(slot);
342        // Re-check generation under the seqlock so a stale activator cannot
343        // resurrect a slot that was reused while we were preparing.
344        let observed_inside = slot.generation.load(Ordering::Acquire);
345        if observed_inside != spec.generation {
346            end_write(slot, begin);
347            return Err(MetricsError::GenerationMismatch {
348                expected: spec.generation,
349                actual: observed_inside,
350            });
351        }
352        slot.run_id.store(spec.run_id, Ordering::Relaxed);
353        slot.pid.store(spec.pid, Ordering::Relaxed);
354        slot.started_at_unix_ms
355            .store(spec.started_at.timestamp_millis(), Ordering::Relaxed);
356        end_write(slot, begin);
357
358        // Atomic Reserved → Active transition. Reject if anything moved the
359        // slot out of Reserved (an external release, a stale activator, or a
360        // reaper) between our generation check and now.
361        slot.state
362            .compare_exchange(
363                SLOT_RESERVED,
364                SLOT_ACTIVE,
365                Ordering::AcqRel,
366                Ordering::Acquire,
367            )
368            .map_err(|actual| {
369                MetricsError::Custom(format!(
370                    "cannot activate slot {}: state moved out of Reserved ({actual})",
371                    spec.slot
372                ))
373            })?;
374
375        Ok(MetricsSlotWriter {
376            registry: self.clone(),
377            slot: spec.slot,
378            generation: spec.generation,
379        })
380    }
381
382    /// Release a slot. Used both by the runtime exit observer and the host
383    /// reaper. Generation-checked so a stale caller cannot clear a reused slot.
384    ///
385    /// The release does three things, in order:
386    ///
387    /// 1. Bump the slot's generation so any writer that starts after release
388    ///    begins fails either its outer or inner generation check.
389    /// 2. Wait briefly for any already in-flight writer to finish its seqlock cycle.
390    ///    If `seq` stays odd past the wait budget, release fails and leaves
391    ///    the slot owned by the current generation; forced recovery is only
392    ///    used by dead-owner reclaim.
393    /// 3. Publish the new state. Readers observing `Free`/`Stale` either
394    ///    skip the slot or see a coherent terminal sample.
395    pub fn release(&self, slot_idx: u32, generation: u64, mode: ReleaseMode) -> MetricsResult<()> {
396        self.release_inner(slot_idx, generation, mode, false)
397    }
398
399    /// Release a reserved slot if it has not been activated yet.
400    ///
401    /// Returns `Ok(true)` when the reservation was cleared, `Ok(false)` when
402    /// the slot had already moved out of `Reserved`, and
403    /// [`MetricsError::GenerationMismatch`] when the reservation token is
404    /// stale.
405    pub fn release_reserved(&self, slot_idx: u32, generation: u64) -> MetricsResult<bool> {
406        let slot = self.try_slot(slot_idx)?;
407        let observed = slot.generation.load(Ordering::Acquire);
408        if observed != generation {
409            return Err(MetricsError::GenerationMismatch {
410                expected: generation,
411                actual: observed,
412            });
413        }
414
415        if slot
416            .state
417            .compare_exchange(
418                SLOT_RESERVED,
419                SLOT_FREE,
420                Ordering::AcqRel,
421                Ordering::Acquire,
422            )
423            .is_err()
424        {
425            return Ok(false);
426        }
427
428        let new_gen = self.next_generation();
429        let _ = slot.generation.compare_exchange(
430            generation,
431            new_gen,
432            Ordering::AcqRel,
433            Ordering::Acquire,
434        );
435        Ok(true)
436    }
437
438    fn release_inner(
439        &self,
440        slot_idx: u32,
441        generation: u64,
442        mode: ReleaseMode,
443        force_if_busy: bool,
444    ) -> MetricsResult<()> {
445        let slot = self.try_slot(slot_idx)?;
446        let observed = slot.generation.load(Ordering::Acquire);
447        if observed != generation {
448            return Err(MetricsError::GenerationMismatch {
449                expected: generation,
450                actual: observed,
451            });
452        }
453
454        // Invalidate any writer that might still be holding `self.generation`
455        // before we wait. A writer already inside the seqlock can finish; a
456        // writer starting after this store cannot enter a valid write cycle.
457        let new_gen = self.next_generation();
458        slot.generation.store(new_gen, Ordering::Release);
459
460        quiesce_seq(slot, force_if_busy)?;
461
462        let new_state = match mode {
463            ReleaseMode::Stale => SLOT_STALE,
464            ReleaseMode::Free => SLOT_FREE,
465        };
466        slot.state.store(new_state, Ordering::Release);
467        Ok(())
468    }
469
470    /// Snapshot every active or stale slot.
471    pub fn snapshot(&self) -> MetricsResult<Vec<LiveMetric>> {
472        self.snapshot_slots(true)
473    }
474
475    /// Snapshot every active slot that has written at least one sample.
476    pub fn active_snapshot(&self) -> MetricsResult<Vec<LiveMetric>> {
477        self.snapshot_slots(false)
478    }
479
480    fn snapshot_slots(&self, include_stale: bool) -> MetricsResult<Vec<LiveMetric>> {
481        let mut out = Vec::new();
482        for idx in 0..self.inner.capacity {
483            if let Some(metric) = self.read_slot_demoting(idx, include_stale) {
484                out.push(metric);
485            }
486        }
487        Ok(out)
488    }
489
490    /// Read a slot, demoting it to `Stale` first when its owner PID is dead.
491    ///
492    /// A runtime that crashes (or is SIGKILLed) never runs its exit observer,
493    /// leaving the slot `Active` forever with a frozen sample. Reserve-time
494    /// reclaim only fires under capacity pressure, so without this readers
495    /// keep reporting the dead sandbox as running. The demotion is
496    /// generation-checked: if the slot was reused between the read and the
497    /// release, the release fails and the (now unrelated) slot is skipped.
498    fn read_slot_demoting(&self, idx: u32, include_stale: bool) -> Option<LiveMetric> {
499        let (mut metric, generation) = self.read_slot(idx, true)?;
500        if metric.state == LiveMetricState::Active && metric.pid > 0 && !pid_is_alive(metric.pid) {
501            if self
502                .release_inner(idx, generation, ReleaseMode::Stale, true)
503                .is_err()
504            {
505                return None;
506            }
507            metric.state = LiveMetricState::Stale;
508        }
509        if metric.state == LiveMetricState::Stale && !include_stale {
510            return None;
511        }
512        Some(metric)
513    }
514
515    /// Lookup the active or stale slot for a sandbox id, if any.
516    pub fn get_by_sandbox_id(&self, sandbox_id: i32) -> MetricsResult<Option<LiveMetric>> {
517        self.get_by_sandbox_identity(sandbox_id, None)
518    }
519
520    /// Lookup the active or stale slot for a sandbox id, requiring the slot's
521    /// inline name to match when one is given.
522    ///
523    /// Prefers an `Active` slot over a `Stale` one: a sandbox that failed a
524    /// boot or restarted can briefly own two slots with the same id, and the
525    /// current run's slot must win over the preserved terminal sample.
526    ///
527    /// The name check guards against catalog id recycling: a removed
528    /// sandbox's stale slot can survive while its row id is reassigned to a
529    /// newly created sandbox, and id alone would then resolve to the ghost.
530    pub fn get_by_sandbox_identity(
531        &self,
532        sandbox_id: i32,
533        name: Option<&str>,
534    ) -> MetricsResult<Option<LiveMetric>> {
535        for want in [SLOT_ACTIVE, SLOT_STALE] {
536            for idx in 0..self.inner.capacity {
537                let slot = self.slot(idx);
538                if slot.state.load(Ordering::Acquire) != want {
539                    continue;
540                }
541                if slot.sandbox_id.load(Ordering::Acquire) != sandbox_id {
542                    continue;
543                }
544                // Re-verify identity from the coherent snapshot: the slot
545                // could have been released and reused for a different sandbox
546                // between the outer filter and the seqlock-protected read.
547                if let Some(metric) = self.read_slot_demoting(idx, true)
548                    && metric.sandbox_id == sandbox_id
549                    && name.is_none_or(|name| metric.name == name)
550                {
551                    return Ok(Some(metric));
552                }
553            }
554        }
555        Ok(None)
556    }
557
558    /// Lookup the active slot for a run id, if any.
559    pub fn get_by_run_id(&self, run_id: i32) -> MetricsResult<Option<LiveMetric>> {
560        for idx in 0..self.inner.capacity {
561            let slot = self.slot(idx);
562            let state = slot.state.load(Ordering::Acquire);
563            if state != SLOT_ACTIVE && state != SLOT_STALE {
564                continue;
565            }
566            if slot.run_id.load(Ordering::Acquire) != run_id {
567                continue;
568            }
569            if let Some(metric) = self.read_slot_demoting(idx, true)
570                && metric.run_id == run_id
571            {
572                return Ok(Some(metric));
573            }
574        }
575        Ok(None)
576    }
577
578    /// Number of slots in this registry.
579    pub fn capacity(&self) -> u32 {
580        self.inner.capacity
581    }
582
583    /// Release the slot owned by the given catalog identity, if any.
584    ///
585    /// Matches by run id first (most precise), falling back to sandbox id
586    /// when `run_id` is `None`. Returns the slot index that was released, or
587    /// `None` if no matching slot was found. The current slot generation is
588    /// looked up internally — callers do not have to track it across
589    /// catalog reads.
590    pub fn release_by_identity(
591        &self,
592        sandbox_id: i32,
593        run_id: Option<i32>,
594        mode: ReleaseMode,
595    ) -> MetricsResult<Option<u32>> {
596        for idx in 0..self.inner.capacity {
597            let slot = self.slot(idx);
598            let state = slot.state.load(Ordering::Acquire);
599            if state != SLOT_ACTIVE && state != SLOT_STALE {
600                continue;
601            }
602            let matches = match run_id {
603                Some(rid) => slot.run_id.load(Ordering::Acquire) == rid,
604                None => slot.sandbox_id.load(Ordering::Acquire) == sandbox_id,
605            };
606            if !matches {
607                continue;
608            }
609            let generation = slot.generation.load(Ordering::Acquire);
610            let force_if_busy = state == SLOT_ACTIVE && {
611                let pid = slot.pid.load(Ordering::Acquire);
612                pid > 0 && !pid_is_alive(pid)
613            };
614            self.release_inner(idx, generation, mode, force_if_busy)?;
615            return Ok(Some(idx));
616        }
617        Ok(None)
618    }
619
620    fn next_generation(&self) -> u64 {
621        let header = self.header();
622        // Generations start at 1 so the value `0` always means "unset".
623        header.global_generation.fetch_add(1, Ordering::AcqRel) + 1
624    }
625
626    fn header(&self) -> &Header {
627        unsafe { &*(self.inner.mapping.ptr.as_ptr() as *const Header) }
628    }
629
630    fn slot(&self, idx: u32) -> &Slot {
631        debug_assert!(idx < self.inner.capacity);
632        let base = self.inner.mapping.ptr.as_ptr();
633        let offset = HEADER_SIZE + (idx as usize) * SLOT_SIZE;
634        unsafe { &*(base.add(offset) as *const Slot) }
635    }
636
637    fn try_slot(&self, idx: u32) -> MetricsResult<&Slot> {
638        if idx >= self.inner.capacity {
639            return Err(MetricsError::Custom(format!(
640                "slot index {idx} out of range (capacity={})",
641                self.inner.capacity
642            )));
643        }
644        Ok(self.slot(idx))
645    }
646
647    /// Read one coherent (metric, generation) pair from a slot. The
648    /// generation is the value observed stable across the seqlock window,
649    /// letting callers perform generation-checked follow-up transitions.
650    fn read_slot(&self, idx: u32, include_stale: bool) -> Option<(LiveMetric, u64)> {
651        let slot = self.slot(idx);
652        // Try many times to obtain a coherent snapshot. A tight-loop writer
653        // can complete a full cycle in <100 ns, so we need a generous budget
654        // before giving up. 4096 attempts is still cheap (<1 ms in the worst
655        // case) and effectively unbounded in practice.
656        for _ in 0..4096 {
657            let state = slot.state.load(Ordering::Acquire);
658            if !slot_visible(state, include_stale) {
659                return None;
660            }
661            // Capture generation before reading fields so we can confirm the
662            // slot's identity stayed stable across the entire read.
663            let gen_before = slot.generation.load(Ordering::Acquire);
664
665            let s1 = slot.seq.load(Ordering::Acquire);
666            if s1 & 1 == 1 {
667                std::hint::spin_loop();
668                continue;
669            }
670
671            let sandbox_id = slot.sandbox_id.load(Ordering::Relaxed);
672            let run_id = slot.run_id.load(Ordering::Relaxed);
673            let pid = slot.pid.load(Ordering::Relaxed);
674            let started_at_ms = slot.started_at_unix_ms.load(Ordering::Relaxed);
675            let sampled_at_ms = slot.sampled_at_unix_ms.load(Ordering::Relaxed);
676            let sample_flags = slot.sample_flags.load(Ordering::Relaxed);
677            let memory_limit = slot.memory_limit_bytes.load(Ordering::Relaxed);
678            let vcpu_time_ns = slot.vcpu_time_ns.load(Ordering::Relaxed);
679            let cpu_bits = slot.cpu_percent_bits.load(Ordering::Relaxed);
680            let memory = slot.memory_bytes.load(Ordering::Relaxed);
681            let memory_available = slot.memory_available_bytes.load(Ordering::Relaxed);
682            let memory_host_resident = slot.memory_host_resident_bytes.load(Ordering::Relaxed);
683            let disk_r = slot.disk_read_bytes.load(Ordering::Relaxed);
684            let disk_w = slot.disk_write_bytes.load(Ordering::Relaxed);
685            let net_rx = slot.net_rx_bytes.load(Ordering::Relaxed);
686            let net_tx = slot.net_tx_bytes.load(Ordering::Relaxed);
687            let upper_used = slot.upper_used_bytes.load(Ordering::Relaxed);
688            let upper_free = slot.upper_free_bytes.load(Ordering::Relaxed);
689            let upper_host_allocated = slot.upper_host_allocated_bytes.load(Ordering::Relaxed);
690            let name = read_name(slot);
691
692            let s2 = slot.seq.load(Ordering::Acquire);
693            // Reject torn reads (s1 != s2), reads that landed mid-write
694            // (s2 odd), reads where the slot was reused under us
695            // (generation changed), or reads where state moved outside the
696            // caller's accepted state set.
697            let gen_after = slot.generation.load(Ordering::Acquire);
698            let state_after = slot.state.load(Ordering::Acquire);
699            let state_stable = slot_visible(state_after, include_stale);
700            if s1 != s2 || s2 & 1 == 1 || gen_before != gen_after || !state_stable {
701                std::hint::spin_loop();
702                continue;
703            }
704
705            // Skip slots whose owner has activated but not yet written a
706            // first sample — `sampled_at_unix_ms` is 0 from reservation
707            // zeroing, which would surface as a 1970-stamped metric.
708            if sampled_at_ms <= 0 {
709                return None;
710            }
711
712            let timestamp = ms_to_datetime(sampled_at_ms);
713            let started_at = ms_to_datetime(started_at_ms);
714            let uptime = timestamp
715                .signed_duration_since(started_at)
716                .to_std()
717                .unwrap_or_default();
718            let state = if state_after == SLOT_STALE {
719                LiveMetricState::Stale
720            } else {
721                LiveMetricState::Active
722            };
723
724            let metric = LiveMetric {
725                state,
726                sandbox_id,
727                run_id,
728                pid,
729                name,
730                timestamp,
731                uptime,
732                cpu_percent: f32::from_bits(cpu_bits),
733                vcpu_time_ns,
734                memory_bytes: memory,
735                memory_available_bytes: flag_value(
736                    sample_flags,
737                    SAMPLE_FLAG_MEMORY_AVAILABLE,
738                    memory_available,
739                ),
740                memory_host_resident_bytes: flag_value(
741                    sample_flags,
742                    SAMPLE_FLAG_MEMORY_HOST_RESIDENT,
743                    memory_host_resident,
744                ),
745                memory_limit_bytes: memory_limit,
746                disk_read_bytes: disk_r,
747                disk_write_bytes: disk_w,
748                net_rx_bytes: net_rx,
749                net_tx_bytes: net_tx,
750                upper_used_bytes: flag_value(sample_flags, SAMPLE_FLAG_UPPER_USED, upper_used),
751                upper_free_bytes: flag_value(sample_flags, SAMPLE_FLAG_UPPER_FREE, upper_free),
752                upper_host_allocated_bytes: flag_value(
753                    sample_flags,
754                    SAMPLE_FLAG_UPPER_HOST_ALLOCATED,
755                    upper_host_allocated,
756                ),
757            };
758            return Some((metric, gen_after));
759        }
760        None
761    }
762}
763
764impl MetricsSlotWriter {
765    /// Write a new sample into the owned slot.
766    ///
767    /// Returns [`MetricsError::GenerationMismatch`] if the slot was reclaimed
768    /// out from under this writer. Callers (the sampler) should stop on
769    /// mismatch instead of forcing the write.
770    pub fn write_sample(&self, sample: SampleWrite) -> MetricsResult<()> {
771        let slot = self.registry.try_slot(self.slot)?;
772        let observed = slot.generation.load(Ordering::Acquire);
773        if observed != self.generation {
774            return Err(MetricsError::GenerationMismatch {
775                expected: self.generation,
776                actual: observed,
777            });
778        }
779
780        let begin = begin_write(slot);
781
782        // Re-check generation while inside the seqlock window. Between the
783        // outer load and `begin_write`, an external caller can release the
784        // slot and a new reservation can claim it, bumping `generation`. If
785        // that happened, our stores would corrupt the new owner's freshly
786        // initialized fields. Close the seqlock without writing.
787        let observed_inside = slot.generation.load(Ordering::Acquire);
788        if observed_inside != self.generation {
789            end_write(slot, begin);
790            return Err(MetricsError::GenerationMismatch {
791                expected: self.generation,
792                actual: observed_inside,
793            });
794        }
795
796        slot.sampled_at_unix_ms
797            .store(sample.sampled_at.timestamp_millis(), Ordering::Relaxed);
798        let mut sample_flags = 0;
799        if sample.cpu_percent.is_some() && sample.vcpu_time_ns.is_some() {
800            sample_flags |= SAMPLE_FLAG_CPU;
801        }
802        if sample.memory_bytes.is_some() {
803            sample_flags |= SAMPLE_FLAG_MEMORY_USED;
804        }
805        if sample.memory_available_bytes.is_some() {
806            sample_flags |= SAMPLE_FLAG_MEMORY_AVAILABLE;
807        }
808        if sample.memory_host_resident_bytes.is_some() {
809            sample_flags |= SAMPLE_FLAG_MEMORY_HOST_RESIDENT;
810        }
811        if sample.upper_used_bytes.is_some() {
812            sample_flags |= SAMPLE_FLAG_UPPER_USED;
813        }
814        if sample.upper_free_bytes.is_some() {
815            sample_flags |= SAMPLE_FLAG_UPPER_FREE;
816        }
817        if sample.upper_host_allocated_bytes.is_some() {
818            sample_flags |= SAMPLE_FLAG_UPPER_HOST_ALLOCATED;
819        }
820        slot.sample_flags.store(sample_flags, Ordering::Relaxed);
821        slot.vcpu_time_ns
822            .store(sample.vcpu_time_ns.unwrap_or(0), Ordering::Relaxed);
823        slot.cpu_percent_bits.store(
824            sample.cpu_percent.unwrap_or(0.0).to_bits(),
825            Ordering::Relaxed,
826        );
827        slot.memory_bytes
828            .store(sample.memory_bytes.unwrap_or(0), Ordering::Relaxed);
829        slot.memory_available_bytes.store(
830            sample.memory_available_bytes.unwrap_or(0),
831            Ordering::Relaxed,
832        );
833        slot.memory_host_resident_bytes.store(
834            sample.memory_host_resident_bytes.unwrap_or(0),
835            Ordering::Relaxed,
836        );
837        slot.disk_read_bytes
838            .store(sample.disk_read_bytes, Ordering::Relaxed);
839        slot.disk_write_bytes
840            .store(sample.disk_write_bytes, Ordering::Relaxed);
841        slot.net_rx_bytes
842            .store(sample.net_rx_bytes, Ordering::Relaxed);
843        slot.net_tx_bytes
844            .store(sample.net_tx_bytes, Ordering::Relaxed);
845        slot.upper_used_bytes
846            .store(sample.upper_used_bytes.unwrap_or(0), Ordering::Relaxed);
847        slot.upper_free_bytes
848            .store(sample.upper_free_bytes.unwrap_or(0), Ordering::Relaxed);
849        slot.upper_host_allocated_bytes.store(
850            sample.upper_host_allocated_bytes.unwrap_or(0),
851            Ordering::Relaxed,
852        );
853        end_write(slot, begin);
854        Ok(())
855    }
856
857    /// Slot index owned by this writer.
858    pub fn slot(&self) -> u32 {
859        self.slot
860    }
861
862    /// Generation paired with this writer's slot reservation.
863    pub fn generation(&self) -> u64 {
864        self.generation
865    }
866
867    /// Convenience method for releasing the owned slot.
868    pub fn release(self, mode: ReleaseMode) -> MetricsResult<()> {
869        self.registry.release(self.slot, self.generation, mode)
870    }
871}
872
873//--------------------------------------------------------------------------------------------------
874// Functions
875//--------------------------------------------------------------------------------------------------
876
877/// Default capacity used by the host launcher when creating a registry.
878pub const fn default_capacity() -> u32 {
879    DEFAULT_CAPACITY
880}
881
882fn try_open_existing(
883    name: &std::ffi::CStr,
884    expected_capacity: u32,
885    map_len: usize,
886) -> MetricsResult<Option<MetricsRegistry>> {
887    let Some(header_mapping) = open_existing_region(name, HEADER_SIZE)? else {
888        return Ok(None);
889    };
890
891    let header_ref = unsafe { &*(header_mapping.ptr.as_ptr() as *const Header) };
892    match wait_for_ready(header_ref) {
893        Ok(()) => {}
894        Err(WaitForReadyError::Stuck) => {
895            return Err(MetricsError::Custom(
896                "metrics registry is still initializing".into(),
897            ));
898        }
899        Err(WaitForReadyError::Invalid(state)) => {
900            return Err(MetricsError::Custom(format!(
901                "invalid registry header state: {state}"
902            )));
903        }
904    }
905    validate_header(header_ref, Some(expected_capacity))?;
906    drop(header_mapping);
907
908    let Some(mapping) = open_existing_region(name, map_len)? else {
909        return Ok(None);
910    };
911
912    let header_ref = unsafe { &*(mapping.ptr.as_ptr() as *const Header) };
913    validate_header(header_ref, Some(expected_capacity))?;
914
915    let inner = RegistryInner {
916        name: name.to_owned(),
917        mapping,
918        capacity: header_ref.capacity,
919    };
920    Ok(Some(MetricsRegistry {
921        inner: Arc::new(inner),
922    }))
923}
924
925#[cfg(unix)]
926fn wait_for_size(fd: libc::c_int, min_size: usize) -> MetricsResult<()> {
927    let deadline = Instant::now() + INIT_WAIT_TIMEOUT;
928    loop {
929        let mut stat = std::mem::MaybeUninit::<libc::stat>::uninit();
930        if unsafe { libc::fstat(fd, stat.as_mut_ptr()) } != 0 {
931            return Err(std::io::Error::last_os_error().into());
932        }
933        let stat = unsafe { stat.assume_init() };
934        let size = stat.st_size as i128;
935        let wanted = min_size as i128;
936        if size >= wanted {
937            return Ok(());
938        }
939        if Instant::now() >= deadline {
940            return Err(MetricsError::Custom(format!(
941                "metrics registry size stayed below {min_size} bytes (size={})",
942                stat.st_size
943            )));
944        }
945        std::thread::sleep(INIT_POLL_INTERVAL);
946    }
947}
948
949fn create_and_init(
950    name: &std::ffi::CStr,
951    capacity: u32,
952    map_len: usize,
953) -> MetricsResult<MetricsRegistry> {
954    let mapping = create_region(name, map_len)?;
955
956    // Initialize header and slots while the state is still UNINIT/INITIALIZING.
957    let header = unsafe { &mut *(mapping.ptr.as_ptr() as *mut Header) };
958    // SAFETY: just-truncated memory is zero-filled. We're the exclusive
959    // initializer thanks to O_EXCL.
960    header
961        .state
962        .store(HEADER_STATE_INITIALIZING, Ordering::Release);
963    header.magic = REGISTRY_MAGIC;
964    header.version = REGISTRY_VERSION;
965    header.header_len = HEADER_SIZE as u32;
966    header.slot_len = SLOT_SIZE as u32;
967    header.capacity = capacity;
968    header.created_at_unix_ms = chrono::Utc::now().timestamp_millis();
969    header.global_generation.store(0, Ordering::Release);
970    // Slots are already zero-filled by ftruncate; SLOT_FREE == 0 and
971    // seq == 0 (even), so they are already in a valid initial state.
972    header.state.store(HEADER_STATE_READY, Ordering::Release);
973
974    let inner = RegistryInner {
975        name: name.to_owned(),
976        mapping,
977        capacity,
978    };
979    Ok(MetricsRegistry {
980        inner: Arc::new(inner),
981    })
982}
983
984#[cfg(unix)]
985fn open_existing_region(
986    name: &std::ffi::CStr,
987    map_len: usize,
988) -> MetricsResult<Option<MappedRegion>> {
989    let fd = unsafe { libc::shm_open(name.as_ptr(), libc::O_RDWR, 0) };
990    if fd < 0 {
991        let err = std::io::Error::last_os_error();
992        if err.raw_os_error() == Some(libc::ENOENT) {
993            return Ok(None);
994        }
995        return Err(err.into());
996    }
997
998    if let Err(err) = wait_for_size(fd, map_len) {
999        unsafe { libc::close(fd) };
1000        return Err(err);
1001    }
1002
1003    let ptr = unsafe {
1004        libc::mmap(
1005            std::ptr::null_mut(),
1006            map_len,
1007            libc::PROT_READ | libc::PROT_WRITE,
1008            libc::MAP_SHARED,
1009            fd,
1010            0,
1011        )
1012    };
1013    unsafe { libc::close(fd) };
1014    if ptr == libc::MAP_FAILED {
1015        return Err(std::io::Error::last_os_error().into());
1016    }
1017
1018    Ok(Some(MappedRegion {
1019        ptr: NonNull::new(ptr as *mut u8).expect("mmap returned non-null"),
1020        len: map_len,
1021    }))
1022}
1023
1024#[cfg(target_os = "windows")]
1025fn open_existing_region(
1026    name: &std::ffi::CStr,
1027    map_len: usize,
1028) -> MetricsResult<Option<MappedRegion>> {
1029    let name = windows_mapping_name(name)?;
1030    let handle = unsafe { OpenFileMappingW(FILE_MAP_ALL_ACCESS, 0, name.as_ptr()) };
1031    if handle.is_null() {
1032        let err = std::io::Error::last_os_error();
1033        if err.raw_os_error() == Some(ERROR_FILE_NOT_FOUND as i32) {
1034            return Ok(None);
1035        }
1036        return Err(err.into());
1037    }
1038
1039    map_windows_region(handle, map_len).map(Some)
1040}
1041
1042#[cfg(unix)]
1043fn create_region(name: &std::ffi::CStr, map_len: usize) -> MetricsResult<MappedRegion> {
1044    let fd = unsafe {
1045        libc::shm_open(
1046            name.as_ptr(),
1047            libc::O_RDWR | libc::O_CREAT | libc::O_EXCL,
1048            0o600,
1049        )
1050    };
1051    if fd < 0 {
1052        let err = std::io::Error::last_os_error();
1053        return match err.raw_os_error() {
1054            Some(libc::EEXIST) => Err(MetricsError::AlreadyExists),
1055            _ => Err(err.into()),
1056        };
1057    }
1058
1059    if unsafe { libc::ftruncate(fd, map_len as libc::off_t) } != 0 {
1060        let e = std::io::Error::last_os_error();
1061        unsafe {
1062            libc::close(fd);
1063            libc::shm_unlink(name.as_ptr());
1064        }
1065        return Err(e.into());
1066    }
1067
1068    let ptr = unsafe {
1069        libc::mmap(
1070            std::ptr::null_mut(),
1071            map_len,
1072            libc::PROT_READ | libc::PROT_WRITE,
1073            libc::MAP_SHARED,
1074            fd,
1075            0,
1076        )
1077    };
1078    unsafe { libc::close(fd) };
1079    if ptr == libc::MAP_FAILED {
1080        let e = std::io::Error::last_os_error();
1081        unsafe {
1082            libc::shm_unlink(name.as_ptr());
1083        }
1084        return Err(e.into());
1085    }
1086
1087    Ok(MappedRegion {
1088        ptr: NonNull::new(ptr as *mut u8).expect("mmap returned non-null"),
1089        len: map_len,
1090    })
1091}
1092
1093#[cfg(target_os = "windows")]
1094fn create_region(name: &std::ffi::CStr, map_len: usize) -> MetricsResult<MappedRegion> {
1095    let name = windows_mapping_name(name)?;
1096    let max_size = map_len as u64;
1097    let handle = unsafe {
1098        CreateFileMappingW(
1099            INVALID_HANDLE_VALUE,
1100            std::ptr::null(),
1101            PAGE_READWRITE,
1102            (max_size >> 32) as u32,
1103            max_size as u32,
1104            name.as_ptr(),
1105        )
1106    };
1107    if handle.is_null() {
1108        return Err(std::io::Error::last_os_error().into());
1109    }
1110    if unsafe { GetLastError() } == ERROR_ALREADY_EXISTS {
1111        unsafe { CloseHandle(handle) };
1112        return Err(MetricsError::AlreadyExists);
1113    }
1114
1115    map_windows_region(handle, map_len)
1116}
1117
1118#[cfg(all(unix, test))]
1119fn unlink_region(name: &std::ffi::CStr) {
1120    unsafe {
1121        libc::shm_unlink(name.as_ptr());
1122    }
1123}
1124
1125#[cfg(target_os = "windows")]
1126#[cfg(test)]
1127fn unlink_region(_name: &std::ffi::CStr) {}
1128
1129#[cfg(target_os = "windows")]
1130fn map_windows_region(handle: HANDLE, map_len: usize) -> MetricsResult<MappedRegion> {
1131    let ptr = unsafe { MapViewOfFile(handle, FILE_MAP_ALL_ACCESS, 0, 0, map_len) };
1132    if ptr.Value.is_null() {
1133        let err = std::io::Error::last_os_error();
1134        unsafe { CloseHandle(handle) };
1135        return Err(err.into());
1136    }
1137
1138    Ok(MappedRegion {
1139        ptr: NonNull::new(ptr.Value.cast()).expect("MapViewOfFile returned non-null"),
1140        handle,
1141    })
1142}
1143
1144#[cfg(target_os = "windows")]
1145fn windows_mapping_name(name: &std::ffi::CStr) -> MetricsResult<Vec<u16>> {
1146    use std::os::windows::ffi::OsStrExt;
1147
1148    let name = name
1149        .to_str()
1150        .map_err(|_| MetricsError::Custom("metrics registry name is not UTF-8".into()))?;
1151    let suffix: String = name
1152        .trim_start_matches('/')
1153        .chars()
1154        .map(|ch| match ch {
1155            '/' | '\\' => '_',
1156            other => other,
1157        })
1158        .collect();
1159    let object_name = format!("Local\\{suffix}");
1160
1161    Ok(std::ffi::OsStr::new(&object_name)
1162        .encode_wide()
1163        .chain(std::iter::once(0))
1164        .collect())
1165}
1166
1167/// Outcome of `wait_for_ready`.
1168enum WaitForReadyError {
1169    /// The header was still `UNINIT`/`INITIALIZING` when the wait expired.
1170    Stuck,
1171    /// The header carried an unrecognised state value.
1172    Invalid(u32),
1173}
1174
1175fn wait_for_ready(header: &Header) -> Result<(), WaitForReadyError> {
1176    let deadline = Instant::now() + INIT_WAIT_TIMEOUT;
1177    loop {
1178        let state = header.state.load(Ordering::Acquire);
1179        if state == HEADER_STATE_READY {
1180            return Ok(());
1181        }
1182        if state != HEADER_STATE_UNINIT && state != HEADER_STATE_INITIALIZING {
1183            return Err(WaitForReadyError::Invalid(state));
1184        }
1185        if Instant::now() >= deadline {
1186            return Err(WaitForReadyError::Stuck);
1187        }
1188        std::thread::sleep(INIT_POLL_INTERVAL);
1189    }
1190}
1191
1192fn validate_header(header: &Header, expected_capacity: Option<u32>) -> MetricsResult<()> {
1193    if header.magic != REGISTRY_MAGIC {
1194        return Err(MetricsError::Custom(format!(
1195            "invalid registry magic: 0x{:x}",
1196            header.magic
1197        )));
1198    }
1199    if header.version != REGISTRY_VERSION {
1200        return Err(MetricsError::Custom(format!(
1201            "incompatible registry version: {}",
1202            header.version
1203        )));
1204    }
1205    if header.header_len as usize != HEADER_SIZE {
1206        return Err(MetricsError::Custom(format!(
1207            "unexpected header length: {}",
1208            header.header_len
1209        )));
1210    }
1211    if header.slot_len as usize != SLOT_SIZE {
1212        return Err(MetricsError::Custom(format!(
1213            "unexpected slot length: {}",
1214            header.slot_len
1215        )));
1216    }
1217    if let Some(expected) = expected_capacity
1218        && header.capacity != expected
1219    {
1220        return Err(MetricsError::Custom(format!(
1221            "registry capacity mismatch: opened={}, expected={expected}",
1222            header.capacity
1223        )));
1224    }
1225    Ok(())
1226}
1227
1228fn write_reservation_fields(slot: &Slot, spec: &ReserveSlot<'_>, generation: u64) {
1229    // Reset the per-sample fields first so a reader seeing the slot before
1230    // activation doesn't observe stale counters from a prior owner.
1231    let begin = begin_write(slot);
1232    slot.sandbox_id.store(spec.sandbox_id, Ordering::Relaxed);
1233    slot.run_id.store(0, Ordering::Relaxed);
1234    slot.pid.store(0, Ordering::Relaxed);
1235    slot.started_at_unix_ms.store(0, Ordering::Relaxed);
1236    slot.sampled_at_unix_ms.store(0, Ordering::Relaxed);
1237    slot.sample_flags.store(0, Ordering::Relaxed);
1238    slot.memory_limit_bytes
1239        .store(spec.memory_limit_bytes, Ordering::Relaxed);
1240    slot.vcpu_time_ns.store(0, Ordering::Relaxed);
1241    slot.cpu_percent_bits.store(0, Ordering::Relaxed);
1242    slot.memory_bytes.store(0, Ordering::Relaxed);
1243    slot.memory_available_bytes.store(0, Ordering::Relaxed);
1244    slot.memory_host_resident_bytes.store(0, Ordering::Relaxed);
1245    slot.disk_read_bytes.store(0, Ordering::Relaxed);
1246    slot.disk_write_bytes.store(0, Ordering::Relaxed);
1247    slot.net_rx_bytes.store(0, Ordering::Relaxed);
1248    slot.net_tx_bytes.store(0, Ordering::Relaxed);
1249    slot.upper_used_bytes.store(0, Ordering::Relaxed);
1250    slot.upper_free_bytes.store(0, Ordering::Relaxed);
1251    slot.upper_host_allocated_bytes.store(0, Ordering::Relaxed);
1252    write_name(slot, spec.name);
1253    end_write(slot, begin);
1254
1255    // Publishing the generation last makes the activator's compare succeed
1256    // only after every reservation byte is visible.
1257    slot.generation.store(generation, Ordering::Release);
1258}
1259
1260fn write_name(slot: &Slot, name: &str) {
1261    let bytes = name.as_bytes();
1262    let len = bytes.len().min(NAME_BYTES);
1263    for (i, byte) in bytes.iter().take(len).enumerate() {
1264        slot.name_bytes[i].store(*byte, Ordering::Relaxed);
1265    }
1266    for i in len..NAME_BYTES {
1267        slot.name_bytes[i].store(0, Ordering::Relaxed);
1268    }
1269    slot.name_len.store(len as u16, Ordering::Relaxed);
1270}
1271
1272fn read_name(slot: &Slot) -> String {
1273    let len = (slot.name_len.load(Ordering::Relaxed) as usize).min(NAME_BYTES);
1274    let mut bytes = [0u8; NAME_BYTES];
1275    for (i, byte) in bytes.iter_mut().enumerate().take(len) {
1276        *byte = slot.name_bytes[i].load(Ordering::Relaxed);
1277    }
1278    String::from_utf8_lossy(&bytes[..len]).into_owned()
1279}
1280
1281fn flag_value(flags: u32, flag: u32, value: u64) -> Option<u64> {
1282    if flag_set(flags, flag) {
1283        Some(value)
1284    } else {
1285        None
1286    }
1287}
1288
1289fn flag_set(flags: u32, flag: u32) -> bool {
1290    flags & flag == flag
1291}
1292
1293fn begin_write(slot: &Slot) -> u64 {
1294    loop {
1295        let seq = slot.seq.load(Ordering::Acquire);
1296        if seq & 1 == 1 {
1297            std::hint::spin_loop();
1298            continue;
1299        }
1300        let begin = seq.wrapping_add(1);
1301        if slot
1302            .seq
1303            .compare_exchange(seq, begin, Ordering::AcqRel, Ordering::Acquire)
1304            .is_ok()
1305        {
1306            return begin;
1307        }
1308        std::hint::spin_loop();
1309    }
1310}
1311
1312fn end_write(slot: &Slot, begin: u64) {
1313    let prev = slot.seq.fetch_add(1, Ordering::AcqRel);
1314    debug_assert_eq!(prev, begin, "seqlock end did not pair with begin");
1315}
1316
1317/// Wait for any in-flight writer to leave the seqlock window before
1318/// surrendering the slot.
1319///
1320/// Normal release is conservative: if `seq` remains odd until the deadline,
1321/// the caller gets an error and the slot stays non-Free. Dead-owner reclaim
1322/// passes `force_if_busy = true` only after checking that the owner PID no
1323/// longer exists, which makes restoring even parity safe enough for reuse.
1324fn quiesce_seq(slot: &Slot, force_if_busy: bool) -> MetricsResult<()> {
1325    const SPIN_BEFORE_YIELD: u32 = 4096;
1326    let deadline = Instant::now() + QUIESCE_WAIT_TIMEOUT;
1327    let mut spins = 0u32;
1328    loop {
1329        let s = slot.seq.load(Ordering::Acquire);
1330        if s & 1 == 0 {
1331            return Ok(());
1332        }
1333        if Instant::now() >= deadline {
1334            if !force_if_busy {
1335                return Err(MetricsError::Custom(
1336                    "metrics slot writer still in progress".into(),
1337                ));
1338            }
1339            let _ = slot.seq.compare_exchange(
1340                s,
1341                s.wrapping_add(1),
1342                Ordering::AcqRel,
1343                Ordering::Acquire,
1344            );
1345            continue;
1346        }
1347        if spins < SPIN_BEFORE_YIELD {
1348            spins += 1;
1349            std::hint::spin_loop();
1350        } else {
1351            thread::yield_now();
1352        }
1353    }
1354}
1355
1356#[cfg(unix)]
1357fn pid_is_alive(pid: i32) -> bool {
1358    if pid <= 0 {
1359        return false;
1360    }
1361    if unsafe { libc::kill(pid, 0) } == 0 {
1362        return true;
1363    }
1364    matches!(
1365        std::io::Error::last_os_error().raw_os_error(),
1366        Some(libc::EPERM)
1367    )
1368}
1369
1370#[cfg(target_os = "windows")]
1371fn pid_is_alive(pid: i32) -> bool {
1372    if pid <= 0 {
1373        return false;
1374    }
1375
1376    let handle = unsafe { OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid as u32) };
1377    if handle.is_null() {
1378        return matches!(
1379            std::io::Error::last_os_error().raw_os_error(),
1380            Some(code) if code == ERROR_ACCESS_DENIED as i32
1381        );
1382    }
1383
1384    let mut exit_code = 0;
1385    let ok = unsafe { GetExitCodeProcess(handle, &mut exit_code) };
1386    unsafe { CloseHandle(handle) };
1387    ok != 0 && exit_code == STILL_ACTIVE as u32
1388}
1389
1390fn ms_to_datetime(ms: i64) -> DateTime<Utc> {
1391    if ms <= 0 {
1392        return DateTime::<Utc>::UNIX_EPOCH;
1393    }
1394    Utc.timestamp_millis_opt(ms)
1395        .single()
1396        .unwrap_or_else(|| Utc.timestamp_opt(0, 0).unwrap())
1397}
1398
1399fn slot_visible(state: u32, include_stale: bool) -> bool {
1400    state == SLOT_ACTIVE || (include_stale && state == SLOT_STALE)
1401}
1402
1403// `AtomicU8`, `AtomicU16`, and `AtomicU64` are referenced to keep the
1404// reservation/write helpers self-contained.
1405fn _assert_atomic_traits() {
1406    fn assert_send_sync<T: Send + Sync>() {}
1407    assert_send_sync::<AtomicU8>();
1408    assert_send_sync::<AtomicU16>();
1409    assert_send_sync::<AtomicU64>();
1410}
1411
1412//--------------------------------------------------------------------------------------------------
1413// Tests
1414//--------------------------------------------------------------------------------------------------
1415
1416#[cfg(test)]
1417mod tests {
1418    use std::sync::Barrier;
1419    use std::thread;
1420
1421    use super::*;
1422
1423    fn unique_name(tag: &str) -> String {
1424        // macOS shm_open names are capped at ~31 bytes; keep the test name
1425        // short enough to fit while remaining unique per test invocation.
1426        let nanos = std::time::SystemTime::now()
1427            .duration_since(std::time::UNIX_EPOCH)
1428            .unwrap()
1429            .as_nanos();
1430        format!("/msb-mtt-{tag}-{}", nanos & 0xffff_ffff)
1431    }
1432
1433    fn cleanup(name: &str) {
1434        let cname = CString::new(name).unwrap();
1435        unlink_region(&cname);
1436    }
1437
1438    /// A PID that is guaranteed alive for the duration of the test, so
1439    /// readers don't demote the slot as dead-owner.
1440    fn alive_pid() -> i32 {
1441        std::process::id() as i32
1442    }
1443
1444    #[test]
1445    fn reserve_activate_write_and_snapshot_roundtrip() {
1446        let name = unique_name("rt");
1447        let reg = MetricsRegistry::open_or_create(&name, 16).unwrap();
1448
1449        let res = reg
1450            .reserve(ReserveSlot {
1451                sandbox_id: 7,
1452                name: "alpine",
1453                memory_limit_bytes: 256 * 1024 * 1024,
1454            })
1455            .unwrap();
1456        let started_at = Utc::now() - chrono::Duration::seconds(2);
1457        let writer = reg
1458            .activate_writer(ActivateSlot {
1459                slot: res.slot,
1460                generation: res.generation,
1461                run_id: 99,
1462                pid: alive_pid(),
1463                started_at,
1464            })
1465            .unwrap();
1466
1467        let sample = SampleWrite {
1468            sampled_at: Utc::now(),
1469            cpu_percent: Some(12.5),
1470            vcpu_time_ns: Some(123_456),
1471            memory_bytes: Some(1024 * 1024),
1472            memory_available_bytes: Some(255 * 1024 * 1024),
1473            memory_host_resident_bytes: Some(2 * 1024 * 1024),
1474            disk_read_bytes: 4096,
1475            disk_write_bytes: 8192,
1476            net_rx_bytes: 100,
1477            net_tx_bytes: 200,
1478            upper_used_bytes: Some(64 * 1024),
1479            upper_free_bytes: Some(128 * 1024),
1480            upper_host_allocated_bytes: Some(16 * 1024),
1481        };
1482        writer.write_sample(sample).unwrap();
1483
1484        let snap = reg.snapshot().unwrap();
1485        assert_eq!(snap.len(), 1);
1486        let item = &snap[0];
1487        assert_eq!(item.sandbox_id, 7);
1488        assert_eq!(item.run_id, 99);
1489        assert_eq!(item.pid, alive_pid());
1490        assert_eq!(item.name, "alpine");
1491        assert!((item.cpu_percent - 12.5).abs() < 1e-6);
1492        assert_eq!(item.vcpu_time_ns, 123_456);
1493        assert_eq!(item.memory_bytes, 1024 * 1024);
1494        assert_eq!(item.memory_available_bytes, Some(255 * 1024 * 1024));
1495        assert_eq!(item.memory_host_resident_bytes, Some(2 * 1024 * 1024));
1496        assert_eq!(item.memory_limit_bytes, 256 * 1024 * 1024);
1497        assert_eq!(item.disk_read_bytes, 4096);
1498        assert_eq!(item.disk_write_bytes, 8192);
1499        assert_eq!(item.net_rx_bytes, 100);
1500        assert_eq!(item.net_tx_bytes, 200);
1501        assert_eq!(item.upper_used_bytes, Some(64 * 1024));
1502        assert_eq!(item.upper_free_bytes, Some(128 * 1024));
1503        assert_eq!(item.upper_host_allocated_bytes, Some(16 * 1024));
1504
1505        // Lookup by sandbox + run id.
1506        assert_eq!(
1507            reg.get_by_sandbox_id(7).unwrap().map(|m| m.run_id),
1508            Some(99)
1509        );
1510        assert_eq!(
1511            reg.get_by_run_id(99).unwrap().map(|m| m.sandbox_id),
1512            Some(7)
1513        );
1514
1515        writer.release(ReleaseMode::Free).unwrap();
1516        assert!(reg.snapshot().unwrap().is_empty());
1517        cleanup(&name);
1518    }
1519
1520    #[test]
1521    fn reserve_rejects_name_that_exceeds_inline_slot_capacity() {
1522        let name = unique_name("long");
1523        let reg = MetricsRegistry::open_or_create(&name, 16).unwrap();
1524
1525        let err = reg
1526            .reserve(ReserveSlot {
1527                sandbox_id: 7,
1528                name: &"x".repeat(NAME_BYTES + 1),
1529                memory_limit_bytes: 1,
1530            })
1531            .unwrap_err();
1532
1533        assert_eq!(
1534            err.to_string(),
1535            "sandbox name is too long for metrics slot: 129 bytes (max 128)"
1536        );
1537        cleanup(&name);
1538    }
1539
1540    #[test]
1541    fn active_snapshot_excludes_stale_slots() {
1542        let name = unique_name("act");
1543        let reg = MetricsRegistry::open_or_create(&name, 2).unwrap();
1544
1545        let res = reg
1546            .reserve(ReserveSlot {
1547                sandbox_id: 7,
1548                name: "alpine",
1549                memory_limit_bytes: 256 * 1024 * 1024,
1550            })
1551            .unwrap();
1552        let writer = reg
1553            .activate_writer(ActivateSlot {
1554                slot: res.slot,
1555                generation: res.generation,
1556                run_id: 99,
1557                pid: alive_pid(),
1558                started_at: Utc::now(),
1559            })
1560            .unwrap();
1561        writer
1562            .write_sample(SampleWrite {
1563                sampled_at: Utc::now(),
1564                cpu_percent: Some(12.5),
1565                vcpu_time_ns: Some(0),
1566                memory_bytes: Some(1024 * 1024),
1567                memory_available_bytes: None,
1568                memory_host_resident_bytes: None,
1569                disk_read_bytes: 4096,
1570                disk_write_bytes: 8192,
1571                net_rx_bytes: 100,
1572                net_tx_bytes: 200,
1573                upper_used_bytes: None,
1574                upper_free_bytes: None,
1575                upper_host_allocated_bytes: None,
1576            })
1577            .unwrap();
1578
1579        assert_eq!(reg.active_snapshot().unwrap().len(), 1);
1580        writer.release(ReleaseMode::Stale).unwrap();
1581        assert_eq!(reg.snapshot().unwrap().len(), 1);
1582        assert!(reg.active_snapshot().unwrap().is_empty());
1583
1584        cleanup(&name);
1585    }
1586
1587    #[test]
1588    fn coherent_reads_under_writer_pressure() {
1589        let name = unique_name("coh");
1590        let reg = MetricsRegistry::open_or_create(&name, 8).unwrap();
1591
1592        let res = reg
1593            .reserve(ReserveSlot {
1594                sandbox_id: 1,
1595                name: "x",
1596                memory_limit_bytes: 1,
1597            })
1598            .unwrap();
1599        let writer = reg
1600            .activate_writer(ActivateSlot {
1601                slot: res.slot,
1602                generation: res.generation,
1603                run_id: 1,
1604                pid: alive_pid(),
1605                started_at: Utc::now(),
1606            })
1607            .unwrap();
1608
1609        let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
1610        let stop_w = Arc::clone(&stop);
1611        let writer_clone = writer.clone();
1612        let writer_thread = thread::spawn(move || {
1613            let mut counter: u64 = 0;
1614            while !stop_w.load(Ordering::Relaxed) {
1615                counter = counter.wrapping_add(1);
1616                let sample = SampleWrite {
1617                    sampled_at: Utc::now(),
1618                    cpu_percent: Some((counter % 100) as f32),
1619                    vcpu_time_ns: Some(counter),
1620                    memory_bytes: Some(counter),
1621                    memory_available_bytes: Some(counter),
1622                    memory_host_resident_bytes: Some(counter),
1623                    disk_read_bytes: counter,
1624                    disk_write_bytes: counter,
1625                    net_rx_bytes: counter,
1626                    net_tx_bytes: counter,
1627                    upper_used_bytes: None,
1628                    upper_free_bytes: None,
1629                    upper_host_allocated_bytes: None,
1630                };
1631                writer_clone.write_sample(sample).unwrap();
1632                // Yield occasionally so the reader has a quiescent window;
1633                // production writers sample at ~1 Hz, not in a tight loop.
1634                std::thread::yield_now();
1635            }
1636        });
1637
1638        // Coherence check: any snapshot we manage to read must have all five
1639        // counter fields equal. The seqlock retries internally; on extreme
1640        // contention we may briefly observe `None`, which is acceptable.
1641        let mut successful_reads = 0;
1642        for _ in 0..10_000 {
1643            let snap = reg.snapshot().unwrap();
1644            if let Some(item) = snap.first() {
1645                assert_eq!(item.memory_bytes, item.disk_read_bytes);
1646                assert_eq!(item.memory_bytes, item.disk_write_bytes);
1647                assert_eq!(item.memory_bytes, item.net_rx_bytes);
1648                assert_eq!(item.memory_bytes, item.net_tx_bytes);
1649                assert_eq!(item.memory_available_bytes, Some(item.memory_bytes));
1650                assert_eq!(item.memory_host_resident_bytes, Some(item.memory_bytes));
1651                successful_reads += 1;
1652            }
1653        }
1654        stop.store(true, Ordering::Relaxed);
1655        writer_thread.join().unwrap();
1656
1657        assert!(
1658            successful_reads > 100,
1659            "expected non-trivial number of successful reads under contention, got {successful_reads}"
1660        );
1661        cleanup(&name);
1662    }
1663
1664    #[test]
1665    fn cloned_writers_serialize_sample_updates() {
1666        let name = unique_name("cw");
1667        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
1668
1669        let res = reg
1670            .reserve(ReserveSlot {
1671                sandbox_id: 1,
1672                name: "x",
1673                memory_limit_bytes: 1,
1674            })
1675            .unwrap();
1676        let writer = reg
1677            .activate_writer(ActivateSlot {
1678                slot: res.slot,
1679                generation: res.generation,
1680                run_id: 1,
1681                pid: alive_pid(),
1682                started_at: Utc::now(),
1683            })
1684            .unwrap();
1685
1686        const WRITERS: usize = 4;
1687        const ITERATIONS: u64 = 4_000;
1688        let barrier = Arc::new(Barrier::new(WRITERS + 1));
1689        let mut handles = Vec::new();
1690        for worker in 0..WRITERS {
1691            let writer = writer.clone();
1692            let barrier = Arc::clone(&barrier);
1693            handles.push(thread::spawn(move || {
1694                barrier.wait();
1695                for n in 0..ITERATIONS {
1696                    let counter = ((worker as u64) << 48) | n | 1;
1697                    writer
1698                        .write_sample(SampleWrite {
1699                            sampled_at: Utc::now(),
1700                            cpu_percent: Some((counter % 100) as f32),
1701                            vcpu_time_ns: Some(counter),
1702                            memory_bytes: Some(counter),
1703                            memory_available_bytes: Some(counter),
1704                            memory_host_resident_bytes: Some(counter),
1705                            disk_read_bytes: counter,
1706                            disk_write_bytes: counter,
1707                            net_rx_bytes: counter,
1708                            net_tx_bytes: counter,
1709                            upper_used_bytes: None,
1710                            upper_free_bytes: None,
1711                            upper_host_allocated_bytes: None,
1712                        })
1713                        .unwrap();
1714                    if n % 64 == 0 {
1715                        thread::yield_now();
1716                    }
1717                }
1718            }));
1719        }
1720
1721        barrier.wait();
1722        let mut successful_reads = 0;
1723        while handles.iter().any(|handle| !handle.is_finished()) {
1724            let snap = reg.snapshot().unwrap();
1725            if let Some(item) = snap.first() {
1726                assert_eq!(item.memory_bytes, item.disk_read_bytes);
1727                assert_eq!(item.memory_bytes, item.disk_write_bytes);
1728                assert_eq!(item.memory_bytes, item.net_rx_bytes);
1729                assert_eq!(item.memory_bytes, item.net_tx_bytes);
1730                assert_eq!(item.memory_available_bytes, Some(item.memory_bytes));
1731                assert_eq!(item.memory_host_resident_bytes, Some(item.memory_bytes));
1732                successful_reads += 1;
1733            }
1734            thread::yield_now();
1735        }
1736        for handle in handles {
1737            handle.join().unwrap();
1738        }
1739
1740        assert!(
1741            successful_reads > 0,
1742            "expected live reads while cloned writers were active, got {successful_reads}"
1743        );
1744        cleanup(&name);
1745    }
1746
1747    #[test]
1748    fn concurrent_reservations_do_not_collide() {
1749        let name = unique_name("res");
1750        let reg = MetricsRegistry::open_or_create(&name, 64).unwrap();
1751
1752        const WORKERS: usize = 8;
1753        const PER_WORKER: usize = 6;
1754        let barrier = Arc::new(Barrier::new(WORKERS));
1755        let mut handles = Vec::new();
1756        for w in 0..WORKERS {
1757            let reg = reg.clone();
1758            let barrier = Arc::clone(&barrier);
1759            handles.push(thread::spawn(move || {
1760                barrier.wait();
1761                let mut claimed = Vec::new();
1762                for n in 0..PER_WORKER {
1763                    let res = reg
1764                        .reserve(ReserveSlot {
1765                            sandbox_id: (w * PER_WORKER + n) as i32 + 1,
1766                            name: "concur",
1767                            memory_limit_bytes: 1,
1768                        })
1769                        .unwrap();
1770                    claimed.push((res.slot, res.generation));
1771                }
1772                claimed
1773            }));
1774        }
1775
1776        let mut all = Vec::new();
1777        for h in handles {
1778            all.extend(h.join().unwrap());
1779        }
1780        assert_eq!(all.len(), WORKERS * PER_WORKER);
1781        let mut slot_indices: Vec<u32> = all.iter().map(|(s, _)| *s).collect();
1782        slot_indices.sort();
1783        slot_indices.dedup();
1784        assert_eq!(slot_indices.len(), WORKERS * PER_WORKER);
1785        cleanup(&name);
1786    }
1787
1788    #[test]
1789    fn generation_mismatch_blocks_stale_writes() {
1790        let name = unique_name("gen");
1791        // Capacity 1 forces the second reservation to reuse the stale slot.
1792        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
1793        let res = reg
1794            .reserve(ReserveSlot {
1795                sandbox_id: 1,
1796                name: "x",
1797                memory_limit_bytes: 1,
1798            })
1799            .unwrap();
1800        let writer = reg
1801            .activate_writer(ActivateSlot {
1802                slot: res.slot,
1803                generation: res.generation,
1804                run_id: 1,
1805                pid: alive_pid(),
1806                started_at: Utc::now(),
1807            })
1808            .unwrap();
1809        // Stale-release the slot, then reuse it for a different sandbox.
1810        writer.clone().release(ReleaseMode::Stale).unwrap();
1811        let res2 = reg
1812            .reserve(ReserveSlot {
1813                sandbox_id: 2,
1814                name: "y",
1815                memory_limit_bytes: 1,
1816            })
1817            .unwrap();
1818        assert_eq!(res2.slot, res.slot);
1819        assert_ne!(res2.generation, res.generation);
1820
1821        let err = writer
1822            .write_sample(SampleWrite {
1823                sampled_at: Utc::now(),
1824                cpu_percent: Some(0.0),
1825                vcpu_time_ns: Some(0),
1826                memory_bytes: Some(0),
1827                memory_available_bytes: None,
1828                memory_host_resident_bytes: None,
1829                disk_read_bytes: 0,
1830                disk_write_bytes: 0,
1831                net_rx_bytes: 0,
1832                net_tx_bytes: 0,
1833                upper_used_bytes: None,
1834                upper_free_bytes: None,
1835                upper_host_allocated_bytes: None,
1836            })
1837            .unwrap_err();
1838        assert!(matches!(err, MetricsError::GenerationMismatch { .. }));
1839        cleanup(&name);
1840    }
1841
1842    #[test]
1843    fn release_by_identity_does_not_clear_new_reservation() {
1844        let name = unique_name("rid");
1845        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
1846        let res = reg
1847            .reserve(ReserveSlot {
1848                sandbox_id: 1,
1849                name: "x",
1850                memory_limit_bytes: 1,
1851            })
1852            .unwrap();
1853        let writer = reg
1854            .activate_writer(ActivateSlot {
1855                slot: res.slot,
1856                generation: res.generation,
1857                run_id: 10,
1858                pid: alive_pid(),
1859                started_at: Utc::now(),
1860            })
1861            .unwrap();
1862        writer
1863            .write_sample(SampleWrite {
1864                sampled_at: Utc::now(),
1865                cpu_percent: Some(0.0),
1866                vcpu_time_ns: Some(0),
1867                memory_bytes: Some(1),
1868                memory_available_bytes: None,
1869                memory_host_resident_bytes: None,
1870                disk_read_bytes: 0,
1871                disk_write_bytes: 0,
1872                net_rx_bytes: 0,
1873                net_tx_bytes: 0,
1874                upper_used_bytes: None,
1875                upper_free_bytes: None,
1876                upper_host_allocated_bytes: None,
1877            })
1878            .unwrap();
1879        writer.release(ReleaseMode::Stale).unwrap();
1880
1881        let next = reg
1882            .reserve(ReserveSlot {
1883                sandbox_id: 1,
1884                name: "x",
1885                memory_limit_bytes: 1,
1886            })
1887            .unwrap();
1888        assert_eq!(
1889            reg.release_by_identity(1, None, ReleaseMode::Free).unwrap(),
1890            None
1891        );
1892
1893        reg.activate_writer(ActivateSlot {
1894            slot: next.slot,
1895            generation: next.generation,
1896            run_id: 11,
1897            pid: alive_pid(),
1898            started_at: Utc::now(),
1899        })
1900        .unwrap();
1901        cleanup(&name);
1902    }
1903
1904    #[test]
1905    fn release_reserved_only_clears_unactivated_slot() {
1906        let name = unique_name("rr");
1907        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
1908        let res = reg
1909            .reserve(ReserveSlot {
1910                sandbox_id: 1,
1911                name: "x",
1912                memory_limit_bytes: 1,
1913            })
1914            .unwrap();
1915
1916        assert!(reg.release_reserved(res.slot, res.generation).unwrap());
1917        let next = reg
1918            .reserve(ReserveSlot {
1919                sandbox_id: 2,
1920                name: "y",
1921                memory_limit_bytes: 1,
1922            })
1923            .unwrap();
1924        assert_eq!(next.slot, res.slot);
1925        assert_ne!(next.generation, res.generation);
1926
1927        let writer = reg
1928            .activate_writer(ActivateSlot {
1929                slot: next.slot,
1930                generation: next.generation,
1931                run_id: 20,
1932                pid: alive_pid(),
1933                started_at: Utc::now(),
1934            })
1935            .unwrap();
1936        assert!(!reg.release_reserved(next.slot, next.generation).unwrap());
1937        writer
1938            .write_sample(SampleWrite {
1939                sampled_at: Utc::now(),
1940                cpu_percent: Some(0.0),
1941                vcpu_time_ns: Some(0),
1942                memory_bytes: Some(1),
1943                memory_available_bytes: None,
1944                memory_host_resident_bytes: None,
1945                disk_read_bytes: 0,
1946                disk_write_bytes: 0,
1947                net_rx_bytes: 0,
1948                net_tx_bytes: 0,
1949                upper_used_bytes: None,
1950                upper_free_bytes: None,
1951                upper_host_allocated_bytes: None,
1952            })
1953            .unwrap();
1954        assert_eq!(reg.snapshot().unwrap().len(), 1);
1955        writer.release(ReleaseMode::Free).unwrap();
1956        cleanup(&name);
1957    }
1958
1959    #[test]
1960    fn full_registry_returns_full_error() {
1961        let name = unique_name("full");
1962        let reg = MetricsRegistry::open_or_create(&name, 2).unwrap();
1963        let _ = reg
1964            .reserve(ReserveSlot {
1965                sandbox_id: 1,
1966                name: "a",
1967                memory_limit_bytes: 1,
1968            })
1969            .unwrap();
1970        let _ = reg
1971            .reserve(ReserveSlot {
1972                sandbox_id: 2,
1973                name: "b",
1974                memory_limit_bytes: 1,
1975            })
1976            .unwrap();
1977        let err = reg
1978            .reserve(ReserveSlot {
1979                sandbox_id: 3,
1980                name: "c",
1981                memory_limit_bytes: 1,
1982            })
1983            .unwrap_err();
1984        assert!(matches!(err, MetricsError::Full));
1985        cleanup(&name);
1986    }
1987
1988    #[test]
1989    fn reserve_reclaims_dead_active_slot_under_pressure() {
1990        let name = unique_name("dead");
1991        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
1992        let res = reg
1993            .reserve(ReserveSlot {
1994                sandbox_id: 1,
1995                name: "x",
1996                memory_limit_bytes: 1,
1997            })
1998            .unwrap();
1999        let writer = reg
2000            .activate_writer(ActivateSlot {
2001                slot: res.slot,
2002                generation: res.generation,
2003                run_id: 10,
2004                pid: i32::MAX,
2005                started_at: Utc::now(),
2006            })
2007            .unwrap();
2008        writer
2009            .write_sample(SampleWrite {
2010                sampled_at: Utc::now(),
2011                cpu_percent: Some(0.0),
2012                vcpu_time_ns: Some(0),
2013                memory_bytes: Some(1),
2014                memory_available_bytes: None,
2015                memory_host_resident_bytes: None,
2016                disk_read_bytes: 0,
2017                disk_write_bytes: 0,
2018                net_rx_bytes: 0,
2019                net_tx_bytes: 0,
2020                upper_used_bytes: None,
2021                upper_free_bytes: None,
2022                upper_host_allocated_bytes: None,
2023            })
2024            .unwrap();
2025
2026        let next = reg
2027            .reserve(ReserveSlot {
2028                sandbox_id: 2,
2029                name: "y",
2030                memory_limit_bytes: 1,
2031            })
2032            .unwrap();
2033        assert_eq!(next.slot, res.slot);
2034        assert_ne!(next.generation, res.generation);
2035
2036        let err = writer
2037            .write_sample(SampleWrite {
2038                sampled_at: Utc::now(),
2039                cpu_percent: Some(0.0),
2040                vcpu_time_ns: Some(0),
2041                memory_bytes: Some(999),
2042                memory_available_bytes: None,
2043                memory_host_resident_bytes: None,
2044                disk_read_bytes: 0,
2045                disk_write_bytes: 0,
2046                net_rx_bytes: 0,
2047                net_tx_bytes: 0,
2048                upper_used_bytes: None,
2049                upper_free_bytes: None,
2050                upper_host_allocated_bytes: None,
2051            })
2052            .unwrap_err();
2053        assert!(matches!(err, MetricsError::GenerationMismatch { .. }));
2054        cleanup(&name);
2055    }
2056
2057    #[test]
2058    fn activated_but_unsampled_slot_is_not_visible() {
2059        // Until a sampler writes its first sample, readers should not see a
2060        // freshly-activated slot — otherwise consumers observe a 1970-stamped
2061        // metric for a sandbox that hasn't reported yet.
2062        let name = unique_name("pre");
2063        let reg = MetricsRegistry::open_or_create(&name, 2).unwrap();
2064        let res = reg
2065            .reserve(ReserveSlot {
2066                sandbox_id: 1,
2067                name: "x",
2068                memory_limit_bytes: 1,
2069            })
2070            .unwrap();
2071        let _writer = reg
2072            .activate_writer(ActivateSlot {
2073                slot: res.slot,
2074                generation: res.generation,
2075                run_id: 1,
2076                pid: alive_pid(),
2077                started_at: Utc::now(),
2078            })
2079            .unwrap();
2080        assert!(reg.snapshot().unwrap().is_empty());
2081        assert!(reg.get_by_sandbox_id(1).unwrap().is_none());
2082        assert!(reg.get_by_run_id(1).unwrap().is_none());
2083        cleanup(&name);
2084    }
2085
2086    #[test]
2087    fn sample_without_guest_source_flags_uses_default_values() {
2088        let name = unique_name("inv");
2089        let reg = MetricsRegistry::open_or_create(&name, 2).unwrap();
2090        let res = reg
2091            .reserve(ReserveSlot {
2092                sandbox_id: 1,
2093                name: "x",
2094                memory_limit_bytes: 1,
2095            })
2096            .unwrap();
2097        let writer = reg
2098            .activate_writer(ActivateSlot {
2099                slot: res.slot,
2100                generation: res.generation,
2101                run_id: 1,
2102                pid: alive_pid(),
2103                started_at: Utc::now(),
2104            })
2105            .unwrap();
2106
2107        writer
2108            .write_sample(SampleWrite {
2109                sampled_at: Utc::now(),
2110                cpu_percent: None,
2111                vcpu_time_ns: None,
2112                memory_bytes: None,
2113                memory_available_bytes: Some(1),
2114                memory_host_resident_bytes: Some(1),
2115                disk_read_bytes: 1,
2116                disk_write_bytes: 1,
2117                net_rx_bytes: 1,
2118                net_tx_bytes: 1,
2119                upper_used_bytes: None,
2120                upper_free_bytes: None,
2121                upper_host_allocated_bytes: None,
2122            })
2123            .unwrap();
2124
2125        let live = reg.get_by_run_id(1).unwrap().unwrap();
2126        assert_eq!(live.cpu_percent, 0.0);
2127        assert_eq!(live.vcpu_time_ns, 0);
2128        assert_eq!(live.memory_bytes, 0);
2129        assert_eq!(live.memory_available_bytes, Some(1));
2130        assert_eq!(live.memory_host_resident_bytes, Some(1));
2131        assert_eq!(live.disk_read_bytes, 1);
2132        assert_eq!(live.disk_write_bytes, 1);
2133        assert_eq!(live.net_rx_bytes, 1);
2134        assert_eq!(live.net_tx_bytes, 1);
2135        assert!(reg.get_by_sandbox_id(1).unwrap().is_some());
2136
2137        writer
2138            .write_sample(SampleWrite {
2139                sampled_at: Utc::now(),
2140                cpu_percent: Some(0.0),
2141                vcpu_time_ns: Some(0),
2142                memory_bytes: Some(0),
2143                memory_available_bytes: Some(1),
2144                memory_host_resident_bytes: Some(1),
2145                disk_read_bytes: 1,
2146                disk_write_bytes: 1,
2147                net_rx_bytes: 1,
2148                net_tx_bytes: 1,
2149                upper_used_bytes: None,
2150                upper_free_bytes: None,
2151                upper_host_allocated_bytes: None,
2152            })
2153            .unwrap();
2154
2155        assert!(reg.get_by_run_id(1).unwrap().is_some());
2156        cleanup(&name);
2157    }
2158
2159    #[test]
2160    fn write_sample_inner_generation_recheck_rejects_stale_writer() {
2161        // Force a release+reserve race between the outer generation check
2162        // and the seqlock window: do the release+reserve manually, then call
2163        // write_sample on the original writer and confirm it errors out.
2164        let name = unique_name("igen");
2165        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
2166        let res = reg
2167            .reserve(ReserveSlot {
2168                sandbox_id: 1,
2169                name: "x",
2170                memory_limit_bytes: 1,
2171            })
2172            .unwrap();
2173        let writer = reg
2174            .activate_writer(ActivateSlot {
2175                slot: res.slot,
2176                generation: res.generation,
2177                run_id: 10,
2178                pid: alive_pid(),
2179                started_at: Utc::now(),
2180            })
2181            .unwrap();
2182        // Stale-release and re-reserve the slot for a different sandbox.
2183        reg.release(res.slot, res.generation, ReleaseMode::Stale)
2184            .unwrap();
2185        let res2 = reg
2186            .reserve(ReserveSlot {
2187                sandbox_id: 2,
2188                name: "y",
2189                memory_limit_bytes: 1,
2190            })
2191            .unwrap();
2192        assert_eq!(res2.slot, res.slot);
2193        assert_ne!(res2.generation, res.generation);
2194        // The original writer must refuse to write a sample.
2195        let err = writer
2196            .write_sample(SampleWrite {
2197                sampled_at: Utc::now(),
2198                cpu_percent: Some(0.0),
2199                vcpu_time_ns: Some(0),
2200                memory_bytes: Some(999_999),
2201                memory_available_bytes: None,
2202                memory_host_resident_bytes: None,
2203                disk_read_bytes: 0,
2204                disk_write_bytes: 0,
2205                net_rx_bytes: 0,
2206                net_tx_bytes: 0,
2207                upper_used_bytes: None,
2208                upper_free_bytes: None,
2209                upper_host_allocated_bytes: None,
2210            })
2211            .unwrap_err();
2212        assert!(matches!(err, MetricsError::GenerationMismatch { .. }));
2213        cleanup(&name);
2214    }
2215
2216    #[test]
2217    fn release_by_identity_recovers_dead_owner_odd_seq() {
2218        let name = unique_name("rdead");
2219        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
2220        let res = reg
2221            .reserve(ReserveSlot {
2222                sandbox_id: 1,
2223                name: "x",
2224                memory_limit_bytes: 1,
2225            })
2226            .unwrap();
2227        let writer = reg
2228            .activate_writer(ActivateSlot {
2229                slot: res.slot,
2230                generation: res.generation,
2231                run_id: 10,
2232                pid: i32::MAX,
2233                started_at: Utc::now(),
2234            })
2235            .unwrap();
2236        writer
2237            .write_sample(SampleWrite {
2238                sampled_at: Utc::now(),
2239                cpu_percent: Some(0.0),
2240                vcpu_time_ns: Some(0),
2241                memory_bytes: Some(1),
2242                memory_available_bytes: None,
2243                memory_host_resident_bytes: None,
2244                disk_read_bytes: 0,
2245                disk_write_bytes: 0,
2246                net_rx_bytes: 0,
2247                net_tx_bytes: 0,
2248                upper_used_bytes: None,
2249                upper_free_bytes: None,
2250                upper_host_allocated_bytes: None,
2251            })
2252            .unwrap();
2253
2254        let slot_ref = reg.slot(res.slot);
2255        let prev = slot_ref.seq.fetch_add(1, Ordering::AcqRel);
2256        assert_eq!(prev & 1, 0);
2257
2258        assert_eq!(
2259            reg.release_by_identity(1, Some(10), ReleaseMode::Free)
2260                .unwrap(),
2261            Some(res.slot)
2262        );
2263        assert_eq!(
2264            slot_ref.seq.load(Ordering::Acquire) & 1,
2265            0,
2266            "dead-owner identity release must leave seq even"
2267        );
2268        assert!(reg.snapshot().unwrap().is_empty());
2269        cleanup(&name);
2270    }
2271
2272    #[test]
2273    fn clean_release_refuses_odd_seq_until_dead_owner_reclaim() {
2274        // Simulate a writer SIGKILL'd between `begin_write` and `end_write`:
2275        // seq is left odd. Clean release must not force recovery because a
2276        // preempted writer could still resume. Under reservation pressure,
2277        // the dead-owner reclaim path may force seq even after proving the
2278        // owner PID is gone.
2279        let name = unique_name("sigk");
2280        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
2281        let res = reg
2282            .reserve(ReserveSlot {
2283                sandbox_id: 1,
2284                name: "x",
2285                memory_limit_bytes: 1,
2286            })
2287            .unwrap();
2288        let writer = reg
2289            .activate_writer(ActivateSlot {
2290                slot: res.slot,
2291                generation: res.generation,
2292                run_id: 10,
2293                pid: i32::MAX,
2294                started_at: Utc::now(),
2295            })
2296            .unwrap();
2297        writer
2298            .write_sample(SampleWrite {
2299                sampled_at: Utc::now(),
2300                cpu_percent: Some(0.0),
2301                vcpu_time_ns: Some(0),
2302                memory_bytes: Some(1),
2303                memory_available_bytes: None,
2304                memory_host_resident_bytes: None,
2305                disk_read_bytes: 0,
2306                disk_write_bytes: 0,
2307                net_rx_bytes: 0,
2308                net_tx_bytes: 0,
2309                upper_used_bytes: None,
2310                upper_free_bytes: None,
2311                upper_host_allocated_bytes: None,
2312            })
2313            .unwrap();
2314
2315        // Inject the SIGKILL-mid-write state: bump seq once to leave it odd.
2316        let slot_ref = reg.slot(res.slot);
2317        let prev = slot_ref.seq.fetch_add(1, Ordering::AcqRel);
2318        assert_eq!(prev & 1, 0, "seq should have been even before injection");
2319        assert_eq!(
2320            slot_ref.seq.load(Ordering::Acquire) & 1,
2321            1,
2322            "seq is now odd, simulating a writer killed mid-write"
2323        );
2324
2325        // Clean release invalidates the generation but refuses to force the
2326        // seqlock back to even while the owner might still be alive.
2327        let gen_before_release = slot_ref.generation.load(Ordering::Acquire);
2328        let err = reg
2329            .release(res.slot, gen_before_release, ReleaseMode::Free)
2330            .unwrap_err();
2331        assert_eq!(err.to_string(), "metrics slot writer still in progress");
2332        assert_eq!(
2333            slot_ref.seq.load(Ordering::Acquire) & 1,
2334            1,
2335            "clean release must leave seq odd when it cannot prove owner death"
2336        );
2337        let gen_after_release = slot_ref.generation.load(Ordering::Acquire);
2338        assert_ne!(
2339            gen_after_release, gen_before_release,
2340            "failed release still bumps generation to invalidate stranded writers"
2341        );
2342
2343        // The stranded writer must now fail on its next write_sample.
2344        let err = writer
2345            .write_sample(SampleWrite {
2346                sampled_at: Utc::now(),
2347                cpu_percent: Some(0.0),
2348                vcpu_time_ns: Some(0),
2349                memory_bytes: Some(999),
2350                memory_available_bytes: None,
2351                memory_host_resident_bytes: None,
2352                disk_read_bytes: 0,
2353                disk_write_bytes: 0,
2354                net_rx_bytes: 0,
2355                net_tx_bytes: 0,
2356                upper_used_bytes: None,
2357                upper_free_bytes: None,
2358                upper_host_allocated_bytes: None,
2359            })
2360            .unwrap_err();
2361        assert!(matches!(err, MetricsError::GenerationMismatch { .. }));
2362
2363        // The registry is full, so the next reservation must reclaim the dead
2364        // active slot, restore even seq, and make the slot reusable.
2365        let res2 = reg
2366            .reserve(ReserveSlot {
2367                sandbox_id: 2,
2368                name: "y",
2369                memory_limit_bytes: 1,
2370            })
2371            .unwrap();
2372        assert_eq!(res2.slot, res.slot);
2373        assert_eq!(
2374            slot_ref.seq.load(Ordering::Acquire) & 1,
2375            0,
2376            "dead-owner reclaim must leave seq even"
2377        );
2378        let writer2 = reg
2379            .activate_writer(ActivateSlot {
2380                slot: res2.slot,
2381                generation: res2.generation,
2382                run_id: 20,
2383                pid: alive_pid(),
2384                started_at: Utc::now(),
2385            })
2386            .unwrap();
2387        writer2
2388            .write_sample(SampleWrite {
2389                sampled_at: Utc::now(),
2390                cpu_percent: Some(1.0),
2391                vcpu_time_ns: Some(0),
2392                memory_bytes: Some(42),
2393                memory_available_bytes: None,
2394                memory_host_resident_bytes: None,
2395                disk_read_bytes: 0,
2396                disk_write_bytes: 0,
2397                net_rx_bytes: 0,
2398                net_tx_bytes: 0,
2399                upper_used_bytes: None,
2400                upper_free_bytes: None,
2401                upper_host_allocated_bytes: None,
2402            })
2403            .unwrap();
2404        let live = reg
2405            .get_by_sandbox_id(2)
2406            .unwrap()
2407            .expect("recovered slot must produce a coherent live sample");
2408        assert_eq!(live.run_id, 20);
2409        assert_eq!(live.memory_bytes, 42);
2410        cleanup(&name);
2411    }
2412
2413    #[test]
2414    fn reopen_existing_registry_reuses_slots() {
2415        let name = unique_name("reopen");
2416        let reg = MetricsRegistry::open_or_create(&name, 4).unwrap();
2417        let res = reg
2418            .reserve(ReserveSlot {
2419                sandbox_id: 11,
2420                name: "alpine",
2421                memory_limit_bytes: 1,
2422            })
2423            .unwrap();
2424        let writer = reg
2425            .activate_writer(ActivateSlot {
2426                slot: res.slot,
2427                generation: res.generation,
2428                run_id: 22,
2429                pid: alive_pid(),
2430                started_at: Utc::now(),
2431            })
2432            .unwrap();
2433        // Write a sample so the slot is visible to readers (readers skip
2434        // freshly-activated slots whose `sampled_at_unix_ms` is still 0).
2435        writer
2436            .write_sample(SampleWrite {
2437                sampled_at: Utc::now(),
2438                cpu_percent: Some(1.0),
2439                vcpu_time_ns: Some(0),
2440                memory_bytes: Some(0),
2441                memory_available_bytes: None,
2442                memory_host_resident_bytes: None,
2443                disk_read_bytes: 0,
2444                disk_write_bytes: 0,
2445                net_rx_bytes: 0,
2446                net_tx_bytes: 0,
2447                upper_used_bytes: None,
2448                upper_free_bytes: None,
2449                upper_host_allocated_bytes: None,
2450            })
2451            .unwrap();
2452
2453        let reopened = MetricsRegistry::open(&name).unwrap();
2454        let found = reopened
2455            .get_by_sandbox_id(11)
2456            .unwrap()
2457            .expect("slot is visible after reopen");
2458        assert_eq!(found.run_id, 22);
2459        assert_eq!(found.name, "alpine");
2460        cleanup(&name);
2461    }
2462
2463    fn write_minimal_sample(writer: &MetricsSlotWriter, memory_bytes: u64) {
2464        writer
2465            .write_sample(SampleWrite {
2466                sampled_at: Utc::now(),
2467                cpu_percent: Some(1.0),
2468                vcpu_time_ns: Some(0),
2469                memory_bytes: Some(memory_bytes),
2470                memory_available_bytes: None,
2471                memory_host_resident_bytes: None,
2472                disk_read_bytes: 0,
2473                disk_write_bytes: 0,
2474                net_rx_bytes: 0,
2475                net_tx_bytes: 0,
2476                upper_used_bytes: None,
2477                upper_free_bytes: None,
2478                upper_host_allocated_bytes: None,
2479            })
2480            .unwrap();
2481    }
2482
2483    #[test]
2484    fn snapshot_demotes_dead_owner_to_stale() {
2485        let name = unique_name("demote");
2486        let reg = MetricsRegistry::open_or_create(&name, 1).unwrap();
2487        let res = reg
2488            .reserve(ReserveSlot {
2489                sandbox_id: 1,
2490                name: "x",
2491                memory_limit_bytes: 1,
2492            })
2493            .unwrap();
2494        let writer = reg
2495            .activate_writer(ActivateSlot {
2496                slot: res.slot,
2497                generation: res.generation,
2498                run_id: 10,
2499                pid: i32::MAX,
2500                started_at: Utc::now(),
2501            })
2502            .unwrap();
2503        write_minimal_sample(&writer, 7);
2504
2505        // The owner PID is dead, so the active view must not report it —
2506        // and the read itself demotes the slot to Stale.
2507        assert!(reg.active_snapshot().unwrap().is_empty());
2508
2509        let snap = reg.snapshot().unwrap();
2510        assert_eq!(snap.len(), 1);
2511        assert_eq!(snap[0].state, LiveMetricState::Stale);
2512        assert_eq!(snap[0].memory_bytes, 7, "terminal sample is preserved");
2513
2514        // The demoted slot is reclaimable without capacity pressure.
2515        let next = reg
2516            .reserve(ReserveSlot {
2517                sandbox_id: 2,
2518                name: "y",
2519                memory_limit_bytes: 1,
2520            })
2521            .unwrap();
2522        assert_eq!(next.slot, res.slot);
2523        cleanup(&name);
2524    }
2525
2526    #[test]
2527    fn get_by_sandbox_id_prefers_active_over_stale_slot() {
2528        // A failed boot or restart can leave a stale slot with the same
2529        // sandbox id at a lower index than the current run's active slot;
2530        // the active one must win.
2531        let name = unique_name("pref");
2532        let reg = MetricsRegistry::open_or_create(&name, 2).unwrap();
2533
2534        let old = reg
2535            .reserve(ReserveSlot {
2536                sandbox_id: 9,
2537                name: "x",
2538                memory_limit_bytes: 1,
2539            })
2540            .unwrap();
2541        let old_writer = reg
2542            .activate_writer(ActivateSlot {
2543                slot: old.slot,
2544                generation: old.generation,
2545                run_id: 90,
2546                pid: alive_pid(),
2547                started_at: Utc::now(),
2548            })
2549            .unwrap();
2550        write_minimal_sample(&old_writer, 1);
2551        old_writer.release(ReleaseMode::Stale).unwrap();
2552
2553        let new = reg
2554            .reserve(ReserveSlot {
2555                sandbox_id: 9,
2556                name: "x",
2557                memory_limit_bytes: 1,
2558            })
2559            .unwrap();
2560        assert!(new.slot > old.slot, "stale slot must sit at a lower index");
2561        let new_writer = reg
2562            .activate_writer(ActivateSlot {
2563                slot: new.slot,
2564                generation: new.generation,
2565                run_id: 91,
2566                pid: alive_pid(),
2567                started_at: Utc::now(),
2568            })
2569            .unwrap();
2570        write_minimal_sample(&new_writer, 2);
2571
2572        let found = reg.get_by_sandbox_id(9).unwrap().unwrap();
2573        assert_eq!(found.state, LiveMetricState::Active);
2574        assert_eq!(found.run_id, 91);
2575        cleanup(&name);
2576    }
2577
2578    #[test]
2579    fn get_by_sandbox_identity_skips_ghost_slot_with_recycled_id() {
2580        // Catalog row ids are recycled after removal: a deleted sandbox's
2581        // stale slot can share an id with a newly created sandbox. The
2582        // name-qualified lookup must skip the ghost.
2583        let name = unique_name("ghost");
2584        let reg = MetricsRegistry::open_or_create(&name, 2).unwrap();
2585
2586        let ghost = reg
2587            .reserve(ReserveSlot {
2588                sandbox_id: 4,
2589                name: "old-removed",
2590                memory_limit_bytes: 1,
2591            })
2592            .unwrap();
2593        let ghost_writer = reg
2594            .activate_writer(ActivateSlot {
2595                slot: ghost.slot,
2596                generation: ghost.generation,
2597                run_id: 40,
2598                pid: alive_pid(),
2599                started_at: Utc::now(),
2600            })
2601            .unwrap();
2602        write_minimal_sample(&ghost_writer, 1);
2603        ghost_writer.release(ReleaseMode::Stale).unwrap();
2604
2605        let real = reg
2606            .reserve(ReserveSlot {
2607                sandbox_id: 4,
2608                name: "fresh",
2609                memory_limit_bytes: 1,
2610            })
2611            .unwrap();
2612        let real_writer = reg
2613            .activate_writer(ActivateSlot {
2614                slot: real.slot,
2615                generation: real.generation,
2616                run_id: 41,
2617                pid: alive_pid(),
2618                started_at: Utc::now(),
2619            })
2620            .unwrap();
2621        write_minimal_sample(&real_writer, 2);
2622        real_writer.release(ReleaseMode::Stale).unwrap();
2623
2624        let found = reg
2625            .get_by_sandbox_identity(4, Some("fresh"))
2626            .unwrap()
2627            .unwrap();
2628        assert_eq!(found.name, "fresh");
2629        assert_eq!(found.run_id, 41);
2630
2631        // Unqualified lookup still answers (first match), and a name with no
2632        // slot returns none.
2633        assert!(reg.get_by_sandbox_id(4).unwrap().is_some());
2634        assert!(
2635            reg.get_by_sandbox_identity(4, Some("never-existed"))
2636                .unwrap()
2637                .is_none()
2638        );
2639        cleanup(&name);
2640    }
2641
2642    #[test]
2643    fn live_metric_reports_slot_state() {
2644        let name = unique_name("state");
2645        let reg = MetricsRegistry::open_or_create(&name, 2).unwrap();
2646        let res = reg
2647            .reserve(ReserveSlot {
2648                sandbox_id: 5,
2649                name: "x",
2650                memory_limit_bytes: 1,
2651            })
2652            .unwrap();
2653        let writer = reg
2654            .activate_writer(ActivateSlot {
2655                slot: res.slot,
2656                generation: res.generation,
2657                run_id: 50,
2658                pid: alive_pid(),
2659                started_at: Utc::now(),
2660            })
2661            .unwrap();
2662        write_minimal_sample(&writer, 1);
2663
2664        let live = reg.get_by_sandbox_id(5).unwrap().unwrap();
2665        assert_eq!(live.state, LiveMetricState::Active);
2666
2667        writer.release(ReleaseMode::Stale).unwrap();
2668        let stale = reg.get_by_sandbox_id(5).unwrap().unwrap();
2669        assert_eq!(stale.state, LiveMetricState::Stale);
2670        assert!(reg.active_snapshot().unwrap().is_empty());
2671        cleanup(&name);
2672    }
2673}