feat: round-11 abstractions — RuntimeConfig auto thread/core, ShutdownGuard RAII
This commit is contained in:
@@ -0,0 +1,28 @@
|
|||||||
|
# Implementation Spec: Round 11 — COMPLETE
|
||||||
|
|
||||||
|
## Goal
|
||||||
|
Auto thread/core allocation + race hardening (RAII shutdown).
|
||||||
|
|
||||||
|
## New Features
|
||||||
|
|
||||||
|
### 1. RuntimeConfig (mytheclipse-core, lifecycle)
|
||||||
|
File: `crates/mytheclipse/src/runtime_auto.rs`
|
||||||
|
- `RuntimeConfig::auto()` / `from_cores(n)` / `compact()` infer worker_threads,
|
||||||
|
max_blocking_threads, compute_threads, io_threads from host CPU topology
|
||||||
|
(std::thread::available_parallelism)
|
||||||
|
- `available_parallelism()` helper
|
||||||
|
- `build_rayon_pool(cfg)` gated on `compute` feature
|
||||||
|
- 3 tests
|
||||||
|
|
||||||
|
### 2. ShutdownGuard (mytheclipse-core, lifecycle)
|
||||||
|
File: `crates/mytheclipse/src/shutdown_guard.rs`
|
||||||
|
- RAII guard — runs completion callback exactly once on drop (panic-safe via
|
||||||
|
Mutex<Option<Box<FnOnce>>>), prevents double-shutdown race
|
||||||
|
- `new(cb)` + `finish()` (fire now + disarm)
|
||||||
|
- 3 tests (fires on drop, finish once, panic path)
|
||||||
|
|
||||||
|
## Files
|
||||||
|
- new: core/src/runtime_auto.rs, core/src/shutdown_guard.rs
|
||||||
|
- core/lib.rs: +module+export for both
|
||||||
|
|
||||||
|
Build: exit 0. Tests: 0 FAILED. Clippy: 0 new warnings.
|
||||||
@@ -70,6 +70,15 @@ pub mod lifecycle;
|
|||||||
|
|
||||||
#[cfg(feature = "lifecycle")]
|
#[cfg(feature = "lifecycle")]
|
||||||
pub mod bg_join;
|
pub mod bg_join;
|
||||||
|
#[cfg(feature = "lifecycle")]
|
||||||
|
pub mod runtime_auto;
|
||||||
|
#[cfg(feature = "lifecycle")]
|
||||||
|
pub use runtime_auto::{available_parallelism, RuntimeConfig};
|
||||||
|
|
||||||
|
#[cfg(feature = "lifecycle")]
|
||||||
|
pub mod shutdown_guard;
|
||||||
|
#[cfg(feature = "lifecycle")]
|
||||||
|
pub use shutdown_guard::ShutdownGuard;
|
||||||
|
|
||||||
#[cfg(all(feature = "observability", feature = "resiliency"))]
|
#[cfg(all(feature = "observability", feature = "resiliency"))]
|
||||||
pub mod middleware;
|
pub mod middleware;
|
||||||
|
|||||||
@@ -0,0 +1,122 @@
|
|||||||
|
//! Auto thread/core allocation (feature `lifecycle`).
|
||||||
|
//!
|
||||||
|
//! Provides [`RuntimeConfig`] — a small builder that infers sensible thread
|
||||||
|
//! counts from the host CPU topology (via [`std::thread::available_parallelism`])
|
||||||
|
//! so callers don't have to hand-tune `worker_threads` / `max_blocking_threads`.
|
||||||
|
//!
|
||||||
|
//! The same helpers let you size a [`rayon::ThreadPoolBuilder`] or any other
|
||||||
|
//! pool builder without pulling in `num_cpus`.
|
||||||
|
//!
|
||||||
|
//! ## Rationale
|
||||||
|
//!
|
||||||
|
//! Most async/runtime configs default to *one* worker thread per core, which
|
||||||
|
//! is usually fine — but compute-heavy or blocking-heavy workloads want a
|
||||||
|
//! separate, explicit breakdown. [`RuntimeConfig`] centralises that logic so
|
||||||
|
//! you configure it once and reuse the counts everywhere.
|
||||||
|
|
||||||
|
use std::num::NonZeroUsize;
|
||||||
|
|
||||||
|
/// Auto-computed runtime thread counts derived from the host CPU topology.
|
||||||
|
///
|
||||||
|
/// Each field is a concrete `usize` (never zero) so it can be fed directly
|
||||||
|
/// into a tokio/rayon/std-thread builder with no further `max(1)` guards.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub struct RuntimeConfig {
|
||||||
|
/// Worker threads for the main async runtime (default: one per core).
|
||||||
|
pub worker_threads: usize,
|
||||||
|
/// Extra threads reserved for blocking (tokio `max_blocking_threads`).
|
||||||
|
pub max_blocking_threads: usize,
|
||||||
|
/// Rayon (compute) pool size.
|
||||||
|
pub compute_threads: usize,
|
||||||
|
/// Suggested background/IO pool size.
|
||||||
|
pub io_threads: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl RuntimeConfig {
|
||||||
|
/// Sizes every pool to the host's available parallelism (one thread per
|
||||||
|
/// logical core), with a small reserved budget for blocking.
|
||||||
|
pub fn auto() -> Self {
|
||||||
|
let cores = available_parallelism();
|
||||||
|
Self {
|
||||||
|
worker_threads: cores,
|
||||||
|
max_blocking_threads: cores.saturating_add(cores / 2).max(4),
|
||||||
|
compute_threads: cores,
|
||||||
|
io_threads: cores.saturating_div(2).clamp(1, 8),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Sizes pools heuristically from `cores` (useful in tests or when you
|
||||||
|
/// want to override the host topology).
|
||||||
|
pub fn from_cores(cores: usize) -> Self {
|
||||||
|
let cores = cores.max(1);
|
||||||
|
Self {
|
||||||
|
worker_threads: cores,
|
||||||
|
max_blocking_threads: cores.saturating_add(cores / 2).max(4),
|
||||||
|
compute_threads: cores,
|
||||||
|
io_threads: cores.saturating_div(2).clamp(1, 8),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A compact runtime: minimal worker + compute threads for constrained
|
||||||
|
/// environments (embedded, small containers).
|
||||||
|
pub fn compact() -> Self {
|
||||||
|
let cores = available_parallelism();
|
||||||
|
Self {
|
||||||
|
worker_threads: cores.max(2),
|
||||||
|
max_blocking_threads: 2,
|
||||||
|
compute_threads: 1,
|
||||||
|
io_threads: 1,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Number of the host's logical CPU cores, falling back to 1 on error.
|
||||||
|
pub fn available_parallelism() -> usize {
|
||||||
|
std::thread::available_parallelism()
|
||||||
|
.map(NonZeroUsize::get)
|
||||||
|
.unwrap_or(1)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Builds a `rayon::ThreadPool` sized with [`RuntimeConfig::compute_threads`].
|
||||||
|
///
|
||||||
|
/// Returns `None` when rayon isn't enabled — call this in code that's gated
|
||||||
|
/// on the `compute` feature to avoid a compile error.
|
||||||
|
#[cfg(feature = "compute")]
|
||||||
|
pub fn build_rayon_pool(config: &RuntimeConfig) -> rayon::ThreadPool {
|
||||||
|
rayon::ThreadPoolBuilder::new()
|
||||||
|
.num_threads(config.compute_threads)
|
||||||
|
.thread_name(|i| format!("mytheclipse-compute-{i}"))
|
||||||
|
.build()
|
||||||
|
.expect("failed to build rayon compute pool")
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn auto_is_nonzero() {
|
||||||
|
let c = RuntimeConfig::auto();
|
||||||
|
assert!(c.worker_threads >= 1);
|
||||||
|
assert!(c.max_blocking_threads >= 4);
|
||||||
|
assert!(c.compute_threads >= 1);
|
||||||
|
assert!(c.io_threads >= 1);
|
||||||
|
assert!(c.io_threads <= 8);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn from_cores_clamps() {
|
||||||
|
let c = RuntimeConfig::from_cores(4);
|
||||||
|
assert_eq!(c.worker_threads, 4);
|
||||||
|
assert_eq!(c.max_blocking_threads, 6);
|
||||||
|
let one = RuntimeConfig::from_cores(0);
|
||||||
|
assert_eq!(one.worker_threads, 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn compact_is_minimal() {
|
||||||
|
let c = RuntimeConfig::compact();
|
||||||
|
assert_eq!(c.compute_threads, 1);
|
||||||
|
assert_eq!(c.io_threads, 1);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,96 @@
|
|||||||
|
//! RAII shutdown guard (feature `lifecycle`).
|
||||||
|
//!
|
||||||
|
//! [`ShutdownGuard`] gives background workers a *scope-based* way to announce
|
||||||
|
//! they have finished, without manually de-registering from a
|
||||||
|
//! [`crate::ShutdownManager`] — the guard fires its completion callback when
|
||||||
|
//! dropped (RAII), so even a `panic!` or an early `return` cannot leak a
|
||||||
|
//! dangling "task still running" registration.
|
||||||
|
//!
|
||||||
|
//! This complements [`crate::ShutdownManager::register`] (which tracks
|
||||||
|
//! `JoinHandle`s): a guard is useful when the work isn't behind a joinable
|
||||||
|
//! handle, or when you want to guarantee cleanup even on unwind.
|
||||||
|
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
|
/// An RAII guard that invokes a notification callback exactly once when
|
||||||
|
/// dropped.
|
||||||
|
///
|
||||||
|
/// The callback is wrapped in a `Mutex<Option<_>>` so it can be taken out and
|
||||||
|
/// run exactly once — even on a `panic!`-unwound drop — guaranteeing at-most-
|
||||||
|
/// once semantics (no double-shutdown race).
|
||||||
|
pub struct ShutdownGuard {
|
||||||
|
inner: Arc<Mutex<Option<Box<dyn FnOnce() + Send>>>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ShutdownGuard {
|
||||||
|
/// Creates a guard that will run `on_drop` (once, on drop) to mark this
|
||||||
|
/// unit of work as complete.
|
||||||
|
pub fn new(on_drop: impl FnOnce() + Send + 'static) -> Self {
|
||||||
|
Self {
|
||||||
|
inner: Arc::new(Mutex::new(Some(Box::new(on_drop)))),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Marks the work as complete immediately, invoking the callback once,
|
||||||
|
/// and disarms the guard so a later drop is a no-op.
|
||||||
|
pub fn finish(self) {
|
||||||
|
self.fire();
|
||||||
|
std::mem::forget(self);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn fire(&self) {
|
||||||
|
let cb = self.inner.lock().unwrap().take();
|
||||||
|
if let Some(cb) = cb {
|
||||||
|
cb();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Drop for ShutdownGuard {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
self.fire();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn fires_on_drop() {
|
||||||
|
let calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let c = Arc::clone(&calls);
|
||||||
|
{
|
||||||
|
let _guard = ShutdownGuard::new(move || {
|
||||||
|
c.fetch_add(1, Ordering::SeqCst);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
assert_eq!(calls.load(Ordering::SeqCst), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn finish_fires_once_and_disarms() {
|
||||||
|
let calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let c = Arc::clone(&calls);
|
||||||
|
let guard = ShutdownGuard::new(move || {
|
||||||
|
c.fetch_add(1, Ordering::SeqCst);
|
||||||
|
});
|
||||||
|
guard.finish();
|
||||||
|
assert_eq!(calls.load(Ordering::SeqCst), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn fires_exactly_once_despite_panic_path() {
|
||||||
|
let calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let c = Arc::clone(&calls);
|
||||||
|
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
|
||||||
|
let _guard = ShutdownGuard::new(move || {
|
||||||
|
c.fetch_add(1, Ordering::SeqCst);
|
||||||
|
});
|
||||||
|
panic!("boom");
|
||||||
|
}));
|
||||||
|
assert!(result.is_err());
|
||||||
|
assert_eq!(calls.load(Ordering::SeqCst), 1);
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user