From cc8ae9d56b987f2dd362c3cb775e1dd4858879fb Mon Sep 17 00:00:00 2001 From: Yaroslav Markin Date: Mon, 7 Sep 2026 23:37:49 +0400 Subject: [PATCH 1/4] Add retired dispatch slots, slot reuse and pool counters to the native layer --- ext/kino/src/control.rs | 79 +++++++++++++++++++++-- ext/kino/src/lib.rs | 9 +++ ext/kino/src/queue.rs | 68 +++++++++++++++++++- ext/kino/src/registry.rs | 95 ++++++++++++++++++++++++++++ ext/kino/src/server.rs | 132 ++++++++++++++++++++++++++++++++++++++- 5 files changed, 371 insertions(+), 12 deletions(-) diff --git a/ext/kino/src/control.rs b/ext/kino/src/control.rs index 6ea82d0..800d506 100644 --- a/ext/kino/src/control.rs +++ b/ext/kino/src/control.rs @@ -18,6 +18,8 @@ pub struct WorkerStat { pub in_flight: usize, pub busy_ms: u64, pub quarantined: bool, + /// Sent home by the pool scaler; leaves at its next idle tick. + pub retired: bool, } /// Read every slot's per-worker sensors in one pass under the slots read @@ -46,6 +48,7 @@ pub fn collect_worker_status(server: &ServerInner) -> Vec { now.saturating_sub(started) }, quarantined, + retired: slot.retired.load(Ordering::Relaxed), } }) .collect() @@ -70,6 +73,10 @@ pub struct StatsSnapshot { pub worker_status: Vec, pub quarantined_count: usize, pub quarantine_replacements: u64, + pub max_workers: usize, + pub active_workers: usize, + pub scale_ups: u64, + pub scale_downs: u64, pub queue_histogram: crate::registry::QueueHistogramSnapshot, } @@ -94,6 +101,10 @@ impl StatsSnapshot { worker_status, quarantined_count, quarantine_replacements: server.quarantine_replacements.load(Ordering::Relaxed), + max_workers: server.topology.max_workers, + active_workers: server.active_workers.load(Ordering::Relaxed), + scale_ups: server.scale_ups.load(Ordering::Relaxed), + scale_downs: server.scale_downs.load(Ordering::Relaxed), queue_histogram: server.queue_histogram.snapshot(), } } @@ -114,9 +125,10 @@ pub fn stats_json(s: &StatsSnapshot) -> String { let mut out = String::with_capacity(256); write!( out, - r#"{{"mode":"{}","lanes":{},"workers":{},"threads":{},"batch":{},"respawns":{},"queued":{},"in_flight":{},"served":{},"rejected":{},"timeouts":{}"#, + r#"{{"mode":"{}","lanes":{},"workers":{},"threads":{},"batch":{},"respawns":{},"queued":{},"in_flight":{},"served":{},"rejected":{},"timeouts":{},"max_workers":{},"active_workers":{},"scale_ups":{},"scale_downs":{}"#, s.mode, s.lanes, s.workers, s.threads, s.batch, s.respawns, - s.queued, s.in_flight, s.served, s.rejected, s.timeouts + s.queued, s.in_flight, s.served, s.rejected, s.timeouts, + s.max_workers, s.active_workers, s.scale_ups, s.scale_downs ) .expect("writing to a String cannot fail"); if let Some(depths) = &s.lane_depths { @@ -134,8 +146,8 @@ pub fn stats_json(s: &StatsSnapshot) -> String { } write!( out, - r#"{{"index":{},"served":{},"in_flight":{},"busy_ms":{},"quarantined":{}}}"#, - w.index, w.served, w.in_flight, w.busy_ms, w.quarantined + r#"{{"index":{},"served":{},"in_flight":{},"busy_ms":{},"quarantined":{},"retired":{}}}"#, + w.index, w.served, w.in_flight, w.busy_ms, w.quarantined, w.retired ) .expect("writing to a String cannot fail"); } @@ -247,6 +259,34 @@ pub fn metrics_text(s: &StatsSnapshot) -> String { "Configured threads per worker.", s.threads, ); + metric( + &mut out, + "kino_max_workers", + "gauge", + "Worker pool ceiling (equals kino_workers for a fixed pool).", + s.max_workers, + ); + metric( + &mut out, + "kino_active_workers", + "gauge", + "Workers currently alive and serving.", + s.active_workers, + ); + metric( + &mut out, + "kino_scale_ups_total", + "counter", + "Workers added by the pool scaler.", + s.scale_ups, + ); + metric( + &mut out, + "kino_scale_downs_total", + "counter", + "Idle workers retired by the pool scaler.", + s.scale_downs, + ); metric( &mut out, "kino_ready", @@ -641,6 +681,10 @@ mod tests { worker_status: vec![], quarantined_count: 0, quarantine_replacements: 0, + max_workers: 32, + active_workers: 12, + scale_ups: 3, + scale_downs: 1, queue_histogram: crate::registry::QueueHistogramSnapshot { buckets: [0; crate::registry::QUEUE_BOUNDS_US.len()], overflow: 0, @@ -665,6 +709,10 @@ mod tests { r#""served":100"#, r#""rejected":5"#, r#""timeouts":6"#, + r#""max_workers":32"#, + r#""active_workers":12"#, + r#""scale_ups":3"#, + r#""scale_downs":1"#, r#""state":"ready""#, r#""version":""#, ] { @@ -690,6 +738,19 @@ mod tests { assert!(draining.contains("kino_ready 0")); } + #[test] + fn metrics_text_reports_the_elastic_pool() { + let text = metrics_text(&snapshot(crate::registry::STATE_READY)); + assert!(text.contains("# TYPE kino_max_workers gauge")); + assert!(text.contains("kino_max_workers 32")); + assert!(text.contains("# TYPE kino_active_workers gauge")); + assert!(text.contains("kino_active_workers 12")); + assert!(text.contains("# TYPE kino_scale_ups_total counter")); + assert!(text.contains("kino_scale_ups_total 3")); + assert!(text.contains("# TYPE kino_scale_downs_total counter")); + assert!(text.contains("kino_scale_downs_total 1")); + } + #[test] fn metrics_text_includes_lane_depth_samples_when_lanes_are_on() { let mut s = snapshot(crate::registry::STATE_READY); @@ -757,6 +818,7 @@ mod tests { in_flight: 1, busy_ms: 4, quarantined: false, + retired: false, }, WorkerStat { index: 1, @@ -764,10 +826,11 @@ mod tests { in_flight: 0, busy_ms: 0, quarantined: false, + retired: true, }, ]; let json = stats_json(&s); - assert!(json.contains(r#""worker_status":[{"index":0,"served":10,"in_flight":1,"busy_ms":4,"quarantined":false},{"index":1,"served":7,"in_flight":0,"busy_ms":0,"quarantined":false}]"#), "got {json}"); + assert!(json.contains(r#""worker_status":[{"index":0,"served":10,"in_flight":1,"busy_ms":4,"quarantined":false,"retired":false},{"index":1,"served":7,"in_flight":0,"busy_ms":0,"quarantined":false,"retired":true}]"#), "got {json}"); } #[test] @@ -786,6 +849,7 @@ mod tests { in_flight: 1, busy_ms: 4, quarantined: false, + retired: false, }, WorkerStat { index: 1, @@ -793,6 +857,7 @@ mod tests { in_flight: 0, busy_ms: 0, quarantined: false, + retired: true, }, ]; let text = metrics_text(&s); @@ -835,6 +900,7 @@ mod tests { in_flight: 1, busy_ms: 0, quarantined: true, + retired: false, }, WorkerStat { index: 1, @@ -842,6 +908,7 @@ mod tests { in_flight: 1, busy_ms: 5, quarantined: false, + retired: false, }, ]; let json = stats_json(&s); @@ -850,7 +917,7 @@ mod tests { "top-level count: {json}" ); assert!( - json.contains(r#"{"index":0,"served":3,"in_flight":1,"busy_ms":0,"quarantined":true}"#), + json.contains(r#"{"index":0,"served":3,"in_flight":1,"busy_ms":0,"quarantined":true,"retired":false}"#), "{json}" ); assert!(json.contains(r#""quarantined":false"#)); diff --git a/ext/kino/src/lib.rs b/ext/kino/src/lib.rs index 0aefb46..38e5cd6 100644 --- a/ext/kino/src/lib.rs +++ b/ext/kino/src/lib.rs @@ -49,6 +49,15 @@ fn init(ruby: &Ruby) -> Result<(), Error> { native.define_singleton_method("worker_stats", function!(server::worker_stats, 1))?; native.define_singleton_method("queue_time", function!(server::queue_time, 1))?; native.define_singleton_method("quarantine_slot", function!(server::quarantine_slot, 2))?; + native.define_singleton_method("retire_slot", function!(server::retire_slot, 2))?; + native.define_singleton_method("reset_slot", function!(server::reset_slot, 2))?; + native.define_singleton_method( + "set_active_workers", + function!(server::set_active_workers, 2), + )?; + native.define_singleton_method("record_scale_up", function!(server::record_scale_up, 1))?; + native.define_singleton_method("record_scale_down", function!(server::record_scale_down, 1))?; + native.define_singleton_method("pool_stats", function!(server::pool_stats, 1))?; native.define_singleton_method( "record_quarantine_replacement", function!(server::record_quarantine_replacement, 1), diff --git a/ext/kino/src/queue.rs b/ext/kino/src/queue.rs index cc3e9b1..b7618c0 100644 --- a/ext/kino/src/queue.rs +++ b/ext/kino/src/queue.rs @@ -25,6 +25,21 @@ pub const TICK: Duration = Duration::from_millis(50); /// the difference and doesn't need to. type Taken = Option; +/// What a take does when its bounded wait times out: keep waiting (None), +/// or end the loop because the pool scaler retired this slot. A retired +/// lane worker first empties its own lane, one item per timeout: the +/// dispatcher stopped feeding it when the flag went up (see +/// registry::ServerInner::retire_slot for the ordering), so whatever is +/// still there is the last of it, and nothing already assigned is +/// orphaned. Only reached with the queue momentarily empty, so the fast +/// path pays nothing for it. +fn after_timeout(slot: &WorkerSlot) -> Option { + if !slot.retired.load(Ordering::SeqCst) { + return None; + } + Some(slot.lane_rx.as_ref().and_then(|rx| rx.try_recv().ok())) +} + /// Block until one request arrives (GVL released, interruptible). /// No busy-poll before parking, deliberately: the wake-per-request futex /// cost is real (~20% of cycles at saturation, per perf), but a measured @@ -46,7 +61,7 @@ fn block_take(server: &ServerInner, slot: &Arc) -> Result Some(Some(ctx)), - Err(flume::RecvTimeoutError::Timeout) => None, + Err(flume::RecvTimeoutError::Timeout) => after_timeout(slot), Err(flume::RecvTimeoutError::Disconnected) => Some(None), })?; Ok(taken.flatten()) @@ -93,8 +108,10 @@ fn lane_take(server: &ServerInner, slot: &Arc) -> Result Some(Some(ctx)), // Periodic steal so a backlog behind a slow sibling can't - // outlive a tick. - Err(flume::RecvTimeoutError::Timeout) => steal().map(Some), + // outlive a tick; a retired worker leaves instead. + Err(flume::RecvTimeoutError::Timeout) => { + after_timeout(slot).or_else(|| steal().map(Some)) + } Err(flume::RecvTimeoutError::Disconnected) => Some(None), } }); @@ -235,3 +252,48 @@ pub fn respond_and_take( crate::request::respond_simple(ruby, request, status, headers, body)?; Worker::take_batch(ruby, &worker, max) } + +#[cfg(test)] +mod tests { + use super::*; + use crate::registry::test_server; + use crate::request::test_ctx; + + #[test] + fn a_take_keeps_waiting_after_a_timeout_unless_the_slot_is_retired() { + let server = test_server(false, 4); + server.register_worker(); + let slot = server.slots.read()[0].clone(); + + assert!(after_timeout(&slot).is_none()); + } + + #[test] + fn a_retired_shared_queue_worker_ends_its_loop_at_the_next_timeout() { + let server = test_server(false, 4); + server.register_worker(); + server.retire_slot(0); + let slot = server.slots.read()[0].clone(); + + assert!(matches!(after_timeout(&slot), Some(None))); + } + + #[test] + fn a_retired_lane_worker_drains_its_own_lane_before_leaving() { + let server = test_server(true, 4); + server.register_worker(); + let slot = server.slots.read()[0].clone(); + slot.lane_tx + .lock() + .as_ref() + .expect("lane open") + .send(test_ctx()) + .expect("lane has room"); + server.retire_slot(0); + + // One item still assigned to this lane: serve it, not orphan it. + assert!(matches!(after_timeout(&slot), Some(Some(_)))); + // Lane empty now: leave. + assert!(matches!(after_timeout(&slot), Some(None))); + } +} diff --git a/ext/kino/src/registry.rs b/ext/kino/src/registry.rs index 40361f7..622169d 100644 --- a/ext/kino/src/registry.rs +++ b/ext/kino/src/registry.rs @@ -21,7 +21,11 @@ pub const STATE_DRAINING: u8 = 2; /// "ractor" or "threaded", never "auto". pub struct Topology { pub mode: String, + /// The pool floor: the configured worker count. pub workers: usize, + /// The pool ceiling: equal to `workers` for a fixed pool. HTTP/2 + /// admission (SETTINGS_MAX_CONCURRENT_STREAMS) is advertised from it. + pub max_workers: usize, pub threads: usize, pub batch: usize, } @@ -113,6 +117,12 @@ pub struct ServerInner { pub respawns: AtomicU64, /// Replacements spawned by the quarantine monitor (Relaxed, advisory). pub quarantine_replacements: AtomicU64, + /// Workers currently alive and serving, reported by the Ruby pool as + /// it grows and shrinks (starts at the floor). Relaxed, advisory. + pub active_workers: AtomicUsize, + /// Pool scaler events (Relaxed, advisory). + pub scale_ups: AtomicU64, + pub scale_downs: AtomicU64, pub topology: Topology, pub https: bool, /// HTTP/2 serving (ALPN over TLS, prior-knowledge h2c on plaintext); @@ -158,6 +168,11 @@ pub struct WorkerSlot { /// Set by the quarantine monitor when this slot is abandoned as wedged: /// excluded from wedge detection, and its busy_ms is reported as 0. pub quarantined: std::sync::atomic::AtomicBool, + /// Set by the pool scaler to send this slot's worker home: the lane + /// dispatcher skips it, and the take loop ends at its next idle tick + /// (a request already taken finishes first). Cleared by `reset_slot` + /// when the slot is handed to a new worker. + pub retired: std::sync::atomic::AtomicBool, } /// Per-lane depth cap: small, so a slow handler can only ever delay this @@ -251,6 +266,7 @@ impl WorkerSlot { in_flight: AtomicUsize::new(0), last_started_ms: AtomicU64::new(0), quarantined: std::sync::atomic::AtomicBool::new(false), + retired: std::sync::atomic::AtomicBool::new(false), } } } @@ -312,6 +328,35 @@ impl ServerInner { slots.len() - 1 } + /// Send a slot's worker home once it is idle. The flag is raised under + /// the slot's lane lock, the same lock the lane dispatcher holds while + /// it checks the flag and sends: any dispatch that saw "not retired" + /// has landed in the lane before the worker can see the flag and + /// drain, so no request is orphaned. Unknown ids are ignored. + pub fn retire_slot(&self, worker_id: usize) { + if let Some(slot) = self.slots.read().get(worker_id) { + let _lane = slot.lane_tx.lock(); + slot.retired.store(true, Ordering::SeqCst); + } + } + + /// Return a retired slot to its fresh state so a new worker can take + /// it over (slots are never removed; recycling keeps the registry + /// from growing with every scale-up). The lane channel stays: the new + /// occupant simply starts taking from it. Unknown ids are ignored. + pub fn reset_slot(&self, worker_id: usize) { + if let Some(slot) = self.slots.read().get(worker_id) { + let _lane = slot.lane_tx.lock(); + slot.current.lock().clear(); + slot.retired.store(false, Ordering::SeqCst); + slot.parked.store(false, Ordering::SeqCst); + slot.interrupted.store(false, Ordering::SeqCst); + slot.served.store(0, Ordering::Relaxed); + slot.in_flight.store(0, Ordering::Relaxed); + slot.last_started_ms.store(0, Ordering::Relaxed); + } + } + pub fn slot( &self, ruby: &magnus::Ruby, @@ -350,9 +395,13 @@ pub fn test_server(lanes: bool, queue_depth: usize) -> Arc { state: std::sync::atomic::AtomicU8::new(STATE_BOOTING), respawns: AtomicU64::new(0), quarantine_replacements: AtomicU64::new(0), + active_workers: AtomicUsize::new(0), + scale_ups: AtomicU64::new(0), + scale_downs: AtomicU64::new(0), topology: Topology { mode: "threaded".to_string(), workers: 0, + max_workers: 0, threads: 0, batch: 1, }, @@ -491,6 +540,18 @@ mod tests { assert_eq!(server.topology.batch, 1); } + #[test] + fn pool_counters_start_at_the_configured_floor() { + let server = test_server(false, 4); + assert_eq!( + server.active_workers.load(Ordering::Relaxed), + server.topology.workers + ); + assert_eq!(server.topology.max_workers, server.topology.workers); + assert_eq!(server.scale_ups.load(Ordering::Relaxed), 0); + assert_eq!(server.scale_downs.load(Ordering::Relaxed), 0); + } + #[test] fn fresh_slot_has_zeroed_per_worker_sensors() { let server = test_server(false, 4); @@ -510,6 +571,40 @@ mod tests { assert_eq!(server.quarantine_replacements.load(Ordering::Relaxed), 0); } + #[test] + fn fresh_slot_is_not_retired() { + let server = test_server(false, 4); + server.register_worker(); + assert!(!server.slots.read()[0].retired.load(Ordering::Relaxed)); + } + + #[test] + fn retire_then_reset_returns_a_slot_to_its_fresh_state() { + let server = test_server(true, 4); + server.register_worker(); + let slot = server.slots.read()[0].clone(); + slot.parked.store(true, Ordering::Relaxed); + slot.interrupted.store(true, Ordering::Relaxed); + slot.served.store(5, Ordering::Relaxed); + slot.in_flight.store(1, Ordering::Relaxed); + slot.last_started_ms.store(9, Ordering::Relaxed); + slot.current.lock().push(std::sync::Weak::new()); + + server.retire_slot(0); + assert!(slot.retired.load(Ordering::Relaxed)); + + server.reset_slot(0); + assert!(!slot.retired.load(Ordering::Relaxed)); + assert!(!slot.parked.load(Ordering::Relaxed)); + assert!(!slot.interrupted.load(Ordering::Relaxed)); + assert_eq!(slot.served.load(Ordering::Relaxed), 0); + assert_eq!(slot.in_flight.load(Ordering::Relaxed), 0); + assert_eq!(slot.last_started_ms.load(Ordering::Relaxed), 0); + assert!(slot.current.lock().is_empty()); + // The lane survives reuse: the next occupant takes from it. + assert!(slot.lane_tx.lock().is_some()); + } + #[test] fn queue_histogram_buckets_by_wait() { let h = QueueHistogram::new(); diff --git a/ext/kino/src/server.rs b/ext/kino/src/server.rs index 5675373..074de65 100644 --- a/ext/kino/src/server.rs +++ b/ext/kino/src/server.rs @@ -68,6 +68,10 @@ pub fn server_start(ruby: &Ruby, config: magnus::RHash) -> Result<(u64, u16, Opt let log_requests: bool = cfg_opt(ruby, config, "log_requests")?.unwrap_or(false); let mode: String = cfg_opt(ruby, config, "mode")?.unwrap_or_else(|| "threaded".to_string()); let workers: usize = cfg_opt(ruby, config, "workers")?.unwrap_or(0); + // Absent or below the floor means a fixed pool. + let max_workers: usize = cfg_opt::(ruby, config, "max_workers")? + .unwrap_or(workers) + .max(workers); let threads: usize = cfg_opt(ruby, config, "threads")?.unwrap_or(0); let batch: usize = cfg_opt(ruby, config, "batch")?.unwrap_or(1); let acceptor = match (&tls_cert, &tls_key) { @@ -130,9 +134,13 @@ pub fn server_start(ruby: &Ruby, config: magnus::RHash) -> Result<(u64, u16, Opt state: std::sync::atomic::AtomicU8::new(registry::STATE_BOOTING), respawns: std::sync::atomic::AtomicU64::new(0), quarantine_replacements: std::sync::atomic::AtomicU64::new(0), + active_workers: std::sync::atomic::AtomicUsize::new(workers), + scale_ups: std::sync::atomic::AtomicU64::new(0), + scale_downs: std::sync::atomic::AtomicU64::new(0), topology: registry::Topology { mode, workers, + max_workers, threads, batch, }, @@ -383,6 +391,12 @@ fn advertised_streams(workers: usize, threads: usize) -> u32 { slots.clamp(8, 1024) as u32 } +/// The pool can grow to its ceiling to meet what an h2 balancer sends, so +/// the ceiling is the admission to advertise (the floor for a fixed pool). +fn stream_capacity(topology: ®istry::Topology) -> u32 { + advertised_streams(topology.max_workers, topology.threads) +} + fn conn_builder(http2: bool, max_streams: u32) -> auto::Builder { let mut builder = auto::Builder::new(TokioExecutor::new()); builder @@ -418,7 +432,7 @@ async fn serve_connection( I: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static, { let http2 = server.http2; - let max_streams = advertised_streams(server.topology.workers, server.topology.threads); + let max_streams = stream_capacity(&server.topology); let mut drain = server.shutdown_tx.subscribe(); let service = service_fn(move |req| handle_request(server.clone(), remote_addr, local_addr, req)); @@ -712,6 +726,11 @@ fn try_dispatch(server: &ServerInner, mut ctx: BoxedCtx) -> Dispatch { let guard = slot.lane_tx.lock(); let Some(tx) = guard.as_ref() else { continue }; any_open = true; + // Checked under the lane lock, which retire_slot also takes: + // see registry::ServerInner::retire_slot for the ordering. + if slot.retired.load(Ordering::SeqCst) { + continue; + } match tx.try_send(ctx) { Ok(()) => return Dispatch::Sent, Err(flume::TrySendError::Full(c)) | Err(flume::TrySendError::Disconnected(c)) => { @@ -900,7 +919,7 @@ pub fn queue_time(_ruby: &Ruby, server_id: u64) -> Result<(u64, f64), Error> { } /// One worker slot's [index, served, in_flight, busy_ms, quarantined] row. -pub type WorkerStatRow = (usize, u64, usize, u64, bool); +pub type WorkerStatRow = (usize, u64, usize, u64, bool, bool); /// Per-slot rows for Server#stats parity: [index, served, in_flight, /// busy_ms, quarantined] each. Empty when the server is gone. @@ -910,7 +929,16 @@ pub fn worker_stats(_ruby: &Ruby, server_id: u64) -> Result, }; Ok(crate::control::collect_worker_status(&server) .into_iter() - .map(|w| (w.index, w.served, w.in_flight, w.busy_ms, w.quarantined)) + .map(|w| { + ( + w.index, + w.served, + w.in_flight, + w.busy_ms, + w.quarantined, + w.retired, + ) + }) .collect()) } @@ -923,6 +951,63 @@ pub fn quarantine_slot(ruby: &Ruby, server_id: u64, worker_id: usize) -> Result< Ok(()) } +/// Send a slot's worker home once idle (pool scale-down); see +/// registry::ServerInner::retire_slot. +pub fn retire_slot(ruby: &Ruby, server_id: u64, worker_id: usize) -> Result<(), Error> { + if let Some(server) = registry::try_get(server_id) { + server.slot(ruby, worker_id)?; // unknown id is a caller bug: raise + server.retire_slot(worker_id); + } + Ok(()) +} + +/// Hand a retired slot to a new worker (pool scale-up reusing a slot). +pub fn reset_slot(ruby: &Ruby, server_id: u64, worker_id: usize) -> Result<(), Error> { + if let Some(server) = registry::try_get(server_id) { + server.slot(ruby, worker_id)?; + server.reset_slot(worker_id); + } + Ok(()) +} + +/// The Ruby pool reports how many workers are alive and serving. +pub fn set_active_workers(_ruby: &Ruby, server_id: u64, count: usize) -> Result<(), Error> { + if let Some(server) = registry::try_get(server_id) { + server.active_workers.store(count, Ordering::Relaxed); + } + Ok(()) +} + +/// One worker added by the pool scaler. +pub fn record_scale_up(_ruby: &Ruby, server_id: u64) -> Result<(), Error> { + if let Some(server) = registry::try_get(server_id) { + server.scale_ups.fetch_add(1, Ordering::Relaxed); + } + Ok(()) +} + +/// One worker retired by the pool scaler. +pub fn record_scale_down(_ruby: &Ruby, server_id: u64) -> Result<(), Error> { + if let Some(server) = registry::try_get(server_id) { + server.scale_downs.fetch_add(1, Ordering::Relaxed); + } + Ok(()) +} + +/// [active_workers, max_workers, scale_ups, scale_downs] for Server#stats. +/// Zeros when the server is gone. +pub fn pool_stats(_ruby: &Ruby, server_id: u64) -> Result<(usize, usize, u64, u64), Error> { + let Some(server) = registry::try_get(server_id) else { + return Ok((0, 0, 0, 0)); + }; + Ok(( + server.active_workers.load(Ordering::Relaxed), + server.topology.max_workers, + server.scale_ups.load(Ordering::Relaxed), + server.scale_downs.load(Ordering::Relaxed), + )) +} + /// One replacement spawned by the quarantine monitor. pub fn record_quarantine_replacement(_ruby: &Ruby, server_id: u64) -> Result<(), Error> { if let Some(server) = registry::try_get(server_id) { @@ -1355,6 +1440,20 @@ mod tests { assert_eq!(advertised_streams(8, 0), 200); } + #[test] + fn stream_capacity_is_advertised_from_the_pool_ceiling() { + // An elastic pool grows to meet load an h2 balancer sends, so the + // ceiling, not the floor, is the admission to advertise. + let topology = registry::Topology { + mode: "ractor".to_string(), + workers: 2, + max_workers: 8, + threads: 3, + batch: 1, + }; + assert_eq!(stream_capacity(&topology), 24); + } + #[tokio::test] async fn streams_beyond_the_advertised_cap_queue_instead_of_failing() { use http_body_util::Full; @@ -1566,4 +1665,31 @@ mod tests { assert!(matches!(try_dispatch(&server, test_ctx()), Dispatch::Sent)); assert_eq!(server.lane_depths(), Some(vec![0, 2])); } + + #[test] + fn dispatch_skips_retired_lanes() { + let server = test_server(true, 4); + server.register_worker(); + server.register_worker(); + server.retire_slot(0); + + // A retiring worker gets nothing new: both land on slot 1. + assert!(matches!(try_dispatch(&server, test_ctx()), Dispatch::Sent)); + assert!(matches!(try_dispatch(&server, test_ctx()), Dispatch::Sent)); + assert_eq!(server.lane_depths(), Some(vec![0, 2])); + } + + #[test] + fn dispatch_retries_rather_than_closing_when_only_retired_lanes_remain() { + let server = test_server(true, 4); + server.register_worker(); + server.retire_slot(0); + + // Retired is not draining: the caller keeps retrying until the + // queue timeout, so a replacement worker can still pick this up. + assert!(matches!( + try_dispatch(&server, test_ctx()), + Dispatch::Full(_) + )); + } } From 1fc1efe9f9badcaff4264d1ebc2c821c017c610c Mon Sep 17 00:00:00 2001 From: Yaroslav Markin Date: Mon, 7 Sep 2026 23:37:49 +0400 Subject: [PATCH 2/4] Add an experimental elastic worker pool with max_workers and scale_down_after --- lib/kino.rb | 2 + lib/kino/cli.rb | 10 +- lib/kino/configuration.rb | 16 +- lib/kino/pool_scaler.rb | 130 ++++++++++++++++ lib/kino/ractor_supervisor.rb | 109 ++++++++++++- lib/kino/server.rb | 143 ++++++++--------- lib/kino/templates/kino.rb.tt | 15 ++ lib/kino/threaded_pool.rb | 188 +++++++++++++++++++++++ sig/kino.rbs | 2 + spec/elastic_pool_spec.rb | 279 ++++++++++++++++++++++++++++++++++ spec/pool_scaler_spec.rb | 154 +++++++++++++++++++ spec/startup_info_spec.rb | 11 ++ 12 files changed, 970 insertions(+), 89 deletions(-) create mode 100644 lib/kino/pool_scaler.rb create mode 100644 lib/kino/threaded_pool.rb create mode 100644 spec/elastic_pool_spec.rb create mode 100644 spec/pool_scaler_spec.rb diff --git a/lib/kino.rb b/lib/kino.rb index e33e529..d73d6ce 100644 --- a/lib/kino.rb +++ b/lib/kino.rb @@ -58,7 +58,9 @@ def self.available_parallelism require_relative "kino/worker_hooks" require_relative "kino/worker" require_relative "kino/ractor_supervisor" +require_relative "kino/threaded_pool" require_relative "kino/quarantine_monitor" +require_relative "kino/pool_scaler" require_relative "kino/server" # Hand the frozen shareable singletons to the native layer: it sets them diff --git a/lib/kino/cli.rb b/lib/kino/cli.rb index 473aa81..48a221a 100644 --- a/lib/kino/cli.rb +++ b/lib/kino/cli.rb @@ -127,7 +127,7 @@ def action!(server) stats = server.stats puts dim("- ruby: #{RUBY_DESCRIPTION}") puts dim("- env: #{ENV["RAILS_ENV"] || ENV["RACK_ENV"] || "development"}") - puts dim("- mode: #{server.mode}, #{count(stats[:workers], "worker")} × #{count(stats[:threads], "thread")}") + puts dim("- mode: #{server.mode}, #{workers_label(stats)} × #{count(stats[:threads], "thread")}") puts dim("- pid: #{Process.pid}") puts dim("- listening: #{server.url}") puts dim("- control: #{server.control_url}") if server.control_url @@ -140,6 +140,14 @@ def count(number, noun) "#{number} #{noun}#{"s" unless number == 1}" end + # The pool as configured: "8 workers", or "8-32 workers" when it can + # grow. + def workers_label(stats) + return count(stats[:workers], "worker") if stats[:max_workers] == stats[:workers] + + "#{stats[:workers]}-#{stats[:max_workers]} workers" + end + # Roll credits when the process ends: normal exit or crash (at_exit # also runs after an uncaught exception; only a force-exit skips it). # @return [void] diff --git a/lib/kino/configuration.rb b/lib/kino/configuration.rb index 93e7983..99f7a63 100644 --- a/lib/kino/configuration.rb +++ b/lib/kino/configuration.rb @@ -10,6 +10,8 @@ class Configuration bind: "127.0.0.1", port: 0, workers: nil, # resolved to Kino.available_parallelism in #to_h + max_workers: nil, # nil = fixed pool; above workers = elastic pool + scale_down_after: nil, # resolved to 30 seconds in Server threads: nil, # resolved per mode in Server: 1 in :ractor, 3 in :threaded mode: :auto, queue_depth: 1024, @@ -146,6 +148,8 @@ def server_options # bind "0.0.0.0" # port 9292 # workers 8 # ractors (or thread groups in :threaded mode) + # max_workers 32 # elastic pool ceiling; unset = fixed pool + # scale_down_after 30 # seconds idle before an extra worker retires # threads 3 # threads per worker # mode :ractor # :auto | :ractor | :threaded # queue_depth 2048 @@ -172,9 +176,19 @@ def bind(host) = @config.set(:bind, host) # Port to listen on; 0 picks an ephemeral port. def port(port) = @config.set(:port, Integer(port)) - # Worker count (ractors in :ractor mode); defaults to CPU cores. + # Worker count (ractors in :ractor mode); defaults to CPU cores. The + # pool floor when max_workers is set. def workers(count) = @config.set(:workers, Integer(count)) + # Pool ceiling: under sustained queue pressure the pool grows past + # `workers`, one worker at a time, up to this many. Unset (the + # default) keeps the pool fixed at `workers`. + def max_workers(count) = @config.set(:max_workers, Integer(count)) + + # Seconds a worker above the floor must sit idle before it is + # retired (default 30). Only meaningful with max_workers. + def scale_down_after(seconds) = @config.set(:scale_down_after, seconds) + # Threads per worker (I/O concurrency inside one ractor); default is # mode-dependent: 1 in :ractor mode, 3 in :threaded. def threads(count) = @config.set(:threads, Integer(count)) diff --git a/lib/kino/pool_scaler.rb b/lib/kino/pool_scaler.rb new file mode 100644 index 0000000..7c22a72 --- /dev/null +++ b/lib/kino/pool_scaler.rb @@ -0,0 +1,130 @@ +# frozen_string_literal: true + +module Kino + # @private + # Grows and shrinks the worker pool between `floor` and `ceiling`. One + # thread on the main ractor polls the native queue and per-slot sensors + # every tick and drives a pool (RactorSupervisor or ThreadedPool) through + # three methods: `active_count`, `groups` (worker index => slot ids for + # every worker that may be retired), `grow`, and `retire(index)`. + # + # Policy: grow by one worker per tick once requests have been waiting in + # the queue on two consecutive ticks (a burst that clears within a tick + # is not pressure); retire one worker per tick, the one idle longest, + # once it has been idle for `scale_down_after`. "Idle" means no slot of + # the worker held a request at two consecutive samples and its served + # count did not move between them, so a worker serving short requests + # between samples is never mistaken for an idle one. + class PoolScaler + TICK = 0.1 + # Consecutive ticks with a non-empty queue before the pool grows. + PRESSURE_TICKS = 2 + + def initialize(server_id:, pool:, floor:, ceiling:, scale_down_after:, tick: TICK) + @server_id = server_id + @pool = pool + @floor = floor + @ceiling = ceiling + @scale_down_after = scale_down_after + @tick = tick + @pressure = 0 + @last_served = {} + @idle_since = {} + @running = false + @thread = nil + end + + def start + @running = true + @thread = Thread.new do + Thread.current.name = "pool-scaler" + run + end + self + end + + def stop + @running = false + @thread&.join(@tick * 2) + end + + # One policy step over one set of observations: the monotonic time, + # the queue depth, and worker_stats rows ([slot, served, in_flight, + # busy_ms, quarantined, retired]). Public so the policy is testable + # without a server. + def step(now, queued, rows) + @pressure = queued.positive? ? @pressure + 1 : 0 + if @pressure >= PRESSURE_TICKS + grow if @pool.active_count < @ceiling + return + end + + observe_idle(now, rows) + return unless @pool.active_count > @floor + + index, since = @idle_since.min_by { |_index, at| at } + return unless index && now - since >= @scale_down_after + + retire(index) + end + + private + + def run + tick while @running + rescue => e + Log.error("pool scaler crashed: #{e.class}: #{e.message}") + end + + def tick + queued, _in_flight = Native.queue_stats(@server_id) + step(Process.clock_gettime(Process::CLOCK_MONOTONIC), queued, Native.worker_stats(@server_id)) + rescue => e + # A bad tick must never kill the scaler. + Log.error("pool scaler tick error: #{e.class}: #{e.message}") + ensure + sleep @tick + end + + def grow + index = @pool.grow + return unless index + + Native.record_scale_up(@server_id) + Log.info("pool grew to #{@pool.active_count} workers (queue pressure)") + end + + def retire(index) + return unless @pool.retire(index) + + forget(index) + Native.record_scale_down(@server_id) + Log.info("pool shrank to #{@pool.active_count} workers (worker-#{index} idle)") + end + + # Per worker: idle when no slot holds a request at this sample and the + # served total did not move since the previous one. First sighting of + # a worker is never idle (it takes two samples to know). + def observe_idle(now, rows) + by_slot = rows.to_h { |row| [row[0], row] } + groups = @pool.groups + (@idle_since.keys - groups.keys).each { |index| forget(index) } + groups.each do |index, slot_ids| + served = slot_ids.sum { |id| by_slot.dig(id, 1) || 0 } + busy = slot_ids.any? { |id| (by_slot.dig(id, 2) || 0).positive? } + idle = !busy && @last_served[index] == served + @last_served[index] = served + if idle + @idle_since[index] ||= now + else + @idle_since.delete(index) + end + end + end + + def forget(index) + @idle_since.delete(index) + @last_served.delete(index) + end + end +end diff --git a/lib/kino/ractor_supervisor.rb b/lib/kino/ractor_supervisor.rb index 5858a1a..d05cae9 100644 --- a/lib/kino/ractor_supervisor.rb +++ b/lib/kino/ractor_supervisor.rb @@ -5,7 +5,12 @@ module Kino # Spawns worker ractors and keeps them alive. One supervisor thread per # ractor: it blocks in Ractor#value, and a crash (anything that kills the # ractor, Exception from app code included) wakes it to 500 the in-flight - # requests and respawn. Clean exits (queue drained) end supervision. + # requests and respawn. Clean exits (queue drained at shutdown, or the + # worker retired by the pool scaler) end supervision. + # + # Also the :ractor-mode pool behind PoolScaler: `grow` adds a worker, + # `retire` sends one home, `groups` lists the ones that may be retired, + # and `active_count` is what the control plane reports. class RactorSupervisor def initialize(server_id, app, workers:, threads:, batch: 1, hooks: nil, on_worker_exit: nil) @server_id = server_id @@ -21,6 +26,12 @@ def initialize(server_id, app, workers:, threads:, batch: 1, hooks: nil, on_work @worker_slots = {} @slot_to_worker = {} @replaced = {} + # Worker index => true while its ractor runs; => true once the + # scaler asked it to leave; and the slot ids retired workers gave + # back, for the next worker to take over. + @live = {} + @retiring = {} + @free_slots = [] # The first replacement's index; `replace` increments before using it, # so this starts one below the first free index (@workers). @next_worker_index = @workers - 1 @@ -28,6 +39,7 @@ def initialize(server_id, app, workers:, threads:, batch: 1, hooks: nil, on_work def start @supervisor_threads = Array.new(@workers) { |index| supervise(index) } + report_active self end @@ -52,6 +64,50 @@ def join @lock.synchronize { @supervisor_threads.dup }.each(&:join) end + # Workers alive and serving: the live ones minus those the quarantine + # monitor abandoned as wedged (their replacements count instead). + def active_count + @lock.synchronize { @live.count { |index, _| !@replaced.key?(index) } } + end + + # Worker index => slot ids for every worker the scaler may retire: + # live, not already leaving, not quarantined, and past its spawn (a + # worker marked live whose supervisor thread has not assigned slots + # yet is not listed until it has). + def groups + @lock.synchronize do + @live.keys + .reject { |index| @retiring.key?(index) || @replaced.key?(index) || !@worker_slots.key?(index) } + .to_h { |index| [index, @worker_slots[index].dup] } + end + end + + # Add one supervised worker; returns its index. + def grow + new_index = @lock.synchronize { @next_worker_index += 1 } + thread = supervise(new_index) + @lock.synchronize { @supervisor_threads << thread } + report_active + new_index + end + + # Send a worker home. Its slots stop receiving work now; the worker + # finishes what it holds, leaves at its next idle tick, and its slots + # come back to the free list once the ractor has exited. Returns + # false when there is no such live worker to retire. + def retire(worker_index) + slot_ids = @lock.synchronize do + next nil unless @live.key?(worker_index) && !@retiring.key?(worker_index) + + @retiring[worker_index] = true + @worker_slots[worker_index] + end + return false unless slot_ids + + slot_ids.each { |id| Native.retire_slot(@server_id, id) } + true + end + # Replace the ractor owning slot `worker_id`: spawn a fresh supervised # ractor, then quarantine the old ractor's slots. The old supervisor # thread stays blocked in ractor.value on the wedged ractor (it and the @@ -90,6 +146,9 @@ def replace(worker_id) private def supervise(index) + # Live from the moment it is asked for, not from when its thread gets + # around to spawning: `grow` reports the count right after this. + @lock.synchronize { @live[index] = true } Thread.new do Thread.current.name = "supervisor-#{index}" crashes = 0 @@ -97,8 +156,9 @@ def supervise(index) ractor, worker_ids = spawn_worker(index) begin ractor.value # blocks until the ractor terminates - HookFire.fire(@on_worker_exit, "on_worker_exit", index, nil) # clean exit: queue drained - break # clean exit: queue closed, workers drained + HookFire.fire(@on_worker_exit, "on_worker_exit", index, nil) # clean exit: drained or retired + exited(index, worker_ids) + break rescue Ractor::Error => e # The ractor died mid-flight. Anything it was serving will never # be answered by Ruby: 500 those clients NOW (not when GC gets @@ -106,11 +166,17 @@ def supervise(index) worker_ids.each { |id| Native.abort_inflight(@server_id, id) } cause = (e.respond_to?(:cause) && e.cause) ? e.cause : e HookFire.fire(@on_worker_exit, "on_worker_exit", index, cause) - break if draining? + if draining? + exited(index, nil) + break + end crashes += 1 Native.record_respawn(@server_id) Log.error("worker-#{index} crashed (#{cause.class}: #{cause.message}); respawning") + # A crashed worker that was on its way out respawns on fresh + # slots like any other; the scaler retires it again when idle. + @lock.synchronize { @retiring.delete(index) } # Policy (crash recovery): unlimited respawn # keeps the server up under rare crashes but turns a # crash-on-every-request bug into a busy loop. A circuit breaker @@ -121,11 +187,11 @@ def supervise(index) end end - # Fresh ractor + fresh native slots. Slots are never reused across - # respawns: stale interrupt kicks and dead weak refs go down with the - # old slot. + # Fresh ractor on fresh or recycled slots. Slots are never reused + # across crash respawns: stale interrupt kicks and dead weak refs go + # down with the old slot. They are reused after a clean retirement. def spawn_worker(worker_index) - worker_ids = Array.new(@threads) { Native.register_worker(@server_id) } + worker_ids = Array.new(@threads) { claim_slot } @lock.synchronize do @worker_slots[worker_index] = worker_ids worker_ids.each { |id| @slot_to_worker[id] = worker_index } @@ -147,6 +213,33 @@ def spawn_worker(worker_index) [ractor, worker_ids] end + # A slot a retired worker gave back, reset for its new occupant, or a + # fresh one. + def claim_slot + id = @lock.synchronize { @free_slots.pop } + return Native.register_worker(@server_id) unless id + + Native.reset_slot(@server_id, id) + id + end + + # Bookkeeping for a supervisor thread that is done: the worker is no + # longer live, a retired worker's slots go back to the free list, and + # the thread leaves the join set so a long-lived elastic pool does not + # accumulate dead threads. + def exited(index, worker_ids) + @lock.synchronize do + @live.delete(index) + @free_slots.concat(worker_ids) if @retiring.delete(index) && worker_ids + @supervisor_threads.delete(Thread.current) + end + report_active + end + + def report_active + Native.set_active_workers(@server_id, active_count) + end + def draining? @lock.synchronize { @draining } end diff --git a/lib/kino/server.rb b/lib/kino/server.rb index 533340a..4e1a454 100644 --- a/lib/kino/server.rb +++ b/lib/kino/server.rb @@ -68,6 +68,16 @@ def initialize(app, config_file: nil, **options) @bind = settings[:bind] @requested_port = settings[:port] @workers = Integer(settings[:workers]) + # The pool ceiling; equal to the floor for a fixed pool. + @max_workers = settings[:max_workers].nil? ? @workers : Integer(settings[:max_workers]) + if @max_workers < @workers + raise ArgumentError, "max_workers (#{@max_workers}) must be at least workers (#{@workers})" + end + @scale_down_after = settings[:scale_down_after].nil? ? 30.0 : Float(settings[:scale_down_after]) + raise ArgumentError, "scale_down_after must be positive" unless @scale_down_after.positive? + if !settings[:scale_down_after].nil? && !elastic? + Log.warn("scale_down_after has no effect unless max_workers is above workers") + end @on_error = validate_hook(settings[:on_error], :on_error) @after_worker_boot = validate_hook(settings[:after_worker_boot], :after_worker_boot) @after_request_complete = validate_hook(settings[:after_request_complete], :after_request_complete) @@ -81,8 +91,9 @@ def initialize(app, config_file: nil, **options) # The access log's GC and allocation figures come from the VM's # process-wide counters, so they are measured only where one # request at a time can own them: the GVL serializes :threaded - # mode, and a single ractor has nothing to race. - access_timing: !!settings[:log_requests] && (@mode == :threaded || @workers == 1) + # mode, and a single ractor (a pool that can never grow past one) + # has nothing to race. + access_timing: !!settings[:log_requests] && (@mode == :threaded || @max_workers == 1) ) # Default threads per mode: 1 in :ractor (threads inside a ractor # share its lock; a measured +17% on fast handlers; raise `workers` @@ -126,10 +137,10 @@ def initialize(app, config_file: nil, **options) else @workers * @threads end - @worker_threads = [] - @worker_threads_lock = Mutex.new @supervisor = nil + @threaded_pool = nil @quarantine_monitor = nil + @pool_scaler = nil @started = false end @@ -158,7 +169,8 @@ def start tls_cert: @tls&.fetch(:cert), tls_key: @tls&.fetch(:key), http2: @http2, lanes: @lanes, log_requests: @log_requests, - mode: @mode.to_s, workers: @workers, threads: @threads, batch: @batch, + mode: @mode.to_s, workers: @workers, max_workers: @max_workers, + threads: @threads, batch: @batch, control_bind: @control_bind, control_token: @control_token ) booted = true @@ -169,12 +181,15 @@ def start # lifetime so in-flight buffers survive even a worker ractor crash. @pin_keeper = Native.pin_keeper(@id) if @mode == :ractor + warn_scheduler_cap @supervisor = RactorSupervisor.new(@id, @app, workers: @workers, threads: @threads, batch: @batch, hooks: @worker_hooks, on_worker_exit: @on_worker_exit).start else - @worker_threads = (@workers * @threads).times.map { spawn_worker_thread } + @threaded_pool = ThreadedPool.new(@id, @app, threads: @threads, batch: @batch, + hooks: @worker_hooks, on_worker_exit: @on_worker_exit).start(@workers) end start_quarantine_monitor if @quarantine_timeout_ms + start_pool_scaler if elastic? Native.control_ready(@id) HookFire.fire(@after_boot, "after_boot") @started = true @@ -192,6 +207,7 @@ def start def shutdown(timeout: nil) return unless @started + @pool_scaler&.stop @quarantine_monitor&.stop deadline = monotonic_now + (timeout || @shutdown_timeout) Native.stop_accepting(@id) @@ -224,7 +240,6 @@ def shutdown(timeout: nil) # The runtime is gone, so hyper has dropped every pinned buffer; # the keeper (and the strings it marked) may now be collected. @pin_keeper = nil - @worker_threads.clear @started = false remove_pidfile if @pidfile nil @@ -233,7 +248,7 @@ def shutdown(timeout: nil) # Block until every worker has exited (i.e. until shutdown). # @return [void] def wait - @supervisor ? @supervisor.join : @worker_threads.each(&:join) + pool.join end # Production entry point: build the server and {#run} it. The `kino` @@ -298,19 +313,22 @@ def self.trap_signals(server) # lanes mode) once started def stats base = { - mode: @mode, lanes: @lanes, workers: @workers, threads: @threads, - batch: @batch, respawns: 0 + mode: @mode, lanes: @lanes, workers: @workers, max_workers: @max_workers, + threads: @threads, batch: @batch, respawns: 0, + active_workers: @workers, scale_ups: 0, scale_downs: 0 } return base unless @started queued, in_flight, served, rejected, timeouts, respawns, lane_depths = Native.server_stats(@id) base.merge!(queued:, in_flight:, served:, rejected:, timeouts:, respawns:) base[:lane_depths] = lane_depths if lane_depths + active_workers, _max_workers, scale_ups, scale_downs = Native.pool_stats(@id) + base.merge!(active_workers:, scale_ups:, scale_downs:) rows = Native.worker_stats(@id) - base[:worker_status] = rows.map do |index, served, in_flight, busy_ms, quarantined| - {index:, served:, in_flight:, busy_ms:, quarantined:} + base[:worker_status] = rows.map do |index, served, in_flight, busy_ms, quarantined, retired| + {index:, served:, in_flight:, busy_ms:, quarantined:, retired:} end - base[:quarantined] = rows.count { |_index, _served, _in_flight, _busy_ms, quarantined| quarantined } + base[:quarantined] = rows.count { |row| row[4] } count, sum_seconds = Native.queue_time(@id) base[:queue_time] = {count:, sum_seconds:} base @@ -318,63 +336,38 @@ def stats private - # Register a fresh dispatch slot and run a worker thread on it; returns - # the thread. Used at boot and by the quarantine replacer. - def spawn_worker_thread - worker_id = Native.register_worker(@id) - Thread.new do - # Named so log lines from inside say which worker spoke. - Thread.current.name = "worker-#{worker_id}" - error = nil - begin - Worker.run(@id, worker_id, @app, @batch, @worker_hooks) - rescue Exception => e # rubocop:disable Lint/RescueException -- a hard crash in a threaded worker thread - error = e - raise - ensure - HookFire.fire(@on_worker_exit, "on_worker_exit", worker_id, error) - end - end + # The worker pool this mode runs: the ractor supervisor, or the + # threaded pool. Both spawn, retire, replace and join workers behind + # the same methods. + def pool + @supervisor || @threaded_pool end - # Track a replacement thread spawned outside the initial pool assignment - # (the quarantine replacer) so shutdown's join/done?/kill sweeps see it. - def track_replacement_thread(thread) - @worker_threads_lock.synchronize { @worker_threads << thread } + # Both pools are quarantine replacers: replace(worker_id) spawns a + # replacement worker, then quarantines the wedged slot. + def start_quarantine_monitor + @quarantine_monitor = QuarantineMonitor.new( + server_id: @id, timeout_ms: @quarantine_timeout_ms, + max: @quarantine_max, replacer: pool + ).start end - # @private - # The :threaded-mode quarantine replacer: spawns a replacement worker - # thread, quarantines the wedged slot, then tracks the new thread so - # shutdown's join/done?/kill sweeps see it. Built from bound Method - # objects instead of a server reference, so it drives the server - # through those methods without send or instance_variable_get. - class ThreadedReplacer - def initialize(server_id:, spawner:, tracker:) - @server_id = server_id - @spawner = spawner - @tracker = tracker - end + # Ruby's M:N scheduler runs non-main ractors' Ruby code on at most + # RUBY_MAX_CPU native threads (default 8). Workers past that cap share + # timeslices instead of adding parallelism, and a fresh ractor can wait + # seconds for its first one while the others are CPU-bound. + def warn_scheduler_cap + cap = Integer(ENV.fetch("RUBY_MAX_CPU", "8"), exception: false) || 8 + return if @max_workers <= cap - def replace(worker_id) - thread = @spawner.call # spawn FIRST (may raise ThreadError) - Native.quarantine_slot(@server_id, worker_id) # quarantine after success - @tracker.call(thread) - true - end + Log.warn("#{@max_workers} ractor workers exceed RUBY_MAX_CPU=#{cap}: only #{cap} can run " \ + "Ruby code at once; set RUBY_MAX_CPU=#{@max_workers} for CPU-bound apps") end - private_constant :ThreadedReplacer - # A replacer.replace(worker_id) spawns a replacement worker, then - # quarantines the wedged slot, mode-appropriately. In :ractor the - # supervisor is the replacer; in :threaded a small object over - # spawn_worker_thread. - def start_quarantine_monitor - replacer = @supervisor || ThreadedReplacer.new(server_id: @id, spawner: method(:spawn_worker_thread), - tracker: method(:track_replacement_thread)) - @quarantine_monitor = QuarantineMonitor.new( - server_id: @id, timeout_ms: @quarantine_timeout_ms, - max: @quarantine_max, replacer: replacer + def start_pool_scaler + @pool_scaler = PoolScaler.new( + server_id: @id, pool: pool, floor: @workers, ceiling: @max_workers, + scale_down_after: @scale_down_after ).start end @@ -400,6 +393,11 @@ def monotonic_now Process.clock_gettime(Process::CLOCK_MONOTONIC) end + # A pool that can grow: the ceiling is above the floor. + def elastic? + @max_workers > @workers + end + # Default connection cap: most of the process open-file limit. A # connection flood's failure mode is descriptor exhaustion, and in # :ractor/:threaded mode the app's own sockets and files share this @@ -474,23 +472,11 @@ def remove_pidfile end def join_workers(deadline) - if @supervisor - @supervisor.shutdown([deadline - monotonic_now, 0].max) - else - threads = @worker_threads_lock.synchronize { @worker_threads.dup } - threads.each do |thread| - thread.join([deadline - monotonic_now, 0.01].max) - end - end + pool.shutdown([deadline - monotonic_now, 0].max) end def workers_done? - if @supervisor - @supervisor.done? - else - threads = @worker_threads_lock.synchronize { @worker_threads.dup } - threads.none?(&:alive?) - end + pool.done? end def kill_stragglers @@ -499,8 +485,7 @@ def kill_stragglers # by abort_all_inflight. The stuck ractor leaks until process exit. Log.error("shutdown deadline passed with stuck ractor workers") unless @supervisor.done? else - threads = @worker_threads_lock.synchronize { @worker_threads.dup } - threads.each { |thread| thread.kill if thread.alive? } + @threaded_pool.kill_stragglers end end diff --git a/lib/kino/templates/kino.rb.tt b/lib/kino/templates/kino.rb.tt index 9876632..b0c025b 100644 --- a/lib/kino/templates/kino.rb.tt +++ b/lib/kino/templates/kino.rb.tt @@ -40,6 +40,21 @@ # `workers` instead. # threads 1 +## Elastic pool (experimental) +# +# Let the pool grow past `workers` under load and shrink back when +# idle: one worker is added every 100 ms while requests wait in the +# queue, and a worker above `workers` retires after `scale_down_after` +# seconds idle. Helps apps that wait on databases or other services; +# pure CPU work gains nothing past the core count. Unset (the default) +# keeps the pool fixed. In :ractor mode, Ruby runs at most RUBY_MAX_CPU +# (default 8) ractors' Ruby code at once; set that variable to match +# `max_workers` on bigger boxes. +# max_workers 32 + +# Seconds an extra worker must sit idle before it is retired. +# scale_down_after 30 + ## Dispatch mode # # :auto - picks :ractor when your app supports it, else :threaded. diff --git a/lib/kino/threaded_pool.rb b/lib/kino/threaded_pool.rb new file mode 100644 index 0000000..30bf401 --- /dev/null +++ b/lib/kino/threaded_pool.rb @@ -0,0 +1,188 @@ +# frozen_string_literal: true + +module Kino + # @private + # The :threaded-mode worker pool: `workers` groups of `threads` plain + # Threads, each thread on its own dispatch slot. A group is one ractor's + # worth of capacity, so `workers` and `max_workers` mean the same thing + # in both modes. Behind PoolScaler here (`grow`, `retire`, `groups`, + # `active_count`), the quarantine monitor's replacer (`replace`), and + # the join/kill sweeps Server#shutdown runs. + class ThreadedPool + def initialize(server_id, app, threads:, batch: 1, hooks: nil, on_worker_exit: nil) + @server_id = server_id + @app = app + @threads = threads + @batch = batch + @hooks = hooks + @on_worker_exit = on_worker_exit + @lock = Mutex.new + # index => {slots:, threads:}; the groups asked to leave; the slot + # ids retired groups gave back; quarantine replacements (one thread + # each, standing in for a wedged slot: neither counted nor retired, + # so the wedged group keeps counting as the capacity it still is). + @groups = {} + @slot_to_group = {} + @retiring = {} + @free_slots = [] + @replacements = {} + @wedged = {} + @next_index = -1 + end + + def start(workers) + workers.times { spawn_group } + report_active + self + end + + # Groups alive and serving, quarantine replacements aside. + def active_count + @lock.synchronize { @groups.count { |index, _| !@replacements.key?(index) } } + end + + # Group index => slot ids for every group the scaler may retire: not + # already leaving, not wedged, not a quarantine replacement. + def groups + reap + @lock.synchronize do + @groups + .reject { |index, _| @retiring.key?(index) || @wedged.key?(index) || @replacements.key?(index) } + .transform_values { |group| group[:slots].dup } + end + end + + # Add one group; returns its index. + def grow + reap + index = spawn_group + report_active + index + end + + # Send a group home: its slots stop receiving work now, each thread + # finishes what it holds and leaves at its next idle tick, and the + # slots come back to the free list once every thread has exited. + def retire(index) + slots = @lock.synchronize do + next nil unless @groups.key?(index) && !@retiring.key?(index) + + @retiring[index] = true + @groups[index][:slots] + end + return false unless slots + + slots.each { |id| Native.retire_slot(@server_id, id) } + true + end + + # The quarantine replacer: spawn a replacement thread on a fresh slot + # FIRST (may raise ThreadError), then quarantine the wedged slot. + def replace(worker_id) + spawn_group(slots: 1, replacement: true) + Native.quarantine_slot(@server_id, worker_id) + @lock.synchronize do + wedged = @slot_to_group[worker_id] + @wedged[wedged] = true if wedged + end + true + end + + # Join every thread up to the (numeric) deadline. + def shutdown(timeout) + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + all_threads.each do |thread| + remaining = deadline - Process.clock_gettime(Process::CLOCK_MONOTONIC) + thread.join([remaining, 0.01].max) + end + end + + def done? + all_threads.none?(&:alive?) + end + + # Block until every thread exits on its own (drain elsewhere). + def join + all_threads.each(&:join) + end + + def kill_stragglers + all_threads.each { |thread| thread.kill if thread.alive? } + end + + private + + def all_threads + @lock.synchronize { @groups.values.flat_map { |group| group[:threads] } } + end + + # The group is in the table before its first thread starts, so a + # ThreadError partway through leaves nothing untracked for shutdown. + def spawn_group(slots: @threads, replacement: false) + group = {slots: [], threads: []} + index = @lock.synchronize do + @next_index += 1 + @groups[@next_index] = group + @replacements[@next_index] = true if replacement + @next_index + end + slots.times do + id = claim_slot + @lock.synchronize do + group[:slots] << id + @slot_to_group[id] = index + end + thread = spawn_thread(id) + @lock.synchronize { group[:threads] << thread } + end + index + end + + def spawn_thread(worker_id) + Thread.new do + # Named so log lines from inside say which worker spoke. + Thread.current.name = "worker-#{worker_id}" + error = nil + begin + Worker.run(@server_id, worker_id, @app, @batch, @hooks) + rescue Exception => e # rubocop:disable Lint/RescueException -- a hard crash in a threaded worker thread + error = e + raise + ensure + HookFire.fire(@on_worker_exit, "on_worker_exit", worker_id, error) + end + end + end + + # A slot a retired group gave back, reset for its new occupant, or a + # fresh one. + def claim_slot + id = @lock.synchronize { @free_slots.pop } + return Native.register_worker(@server_id) unless id + + Native.reset_slot(@server_id, id) + id + end + + # Retiring groups whose threads have all exited give their slots back + # and leave the table. Runs before every pool decision, so the count + # the control plane sees never lags by more than a scaler tick. + def reap + freed = @lock.synchronize do + done = @retiring.keys.select { |index| @groups[index][:threads].none?(&:alive?) } + done.each do |index| + group = @groups.delete(index) + @retiring.delete(index) + group[:slots].each { |id| @slot_to_group.delete(id) } + @free_slots.concat(group[:slots]) + end + done.any? + end + report_active if freed + end + + def report_active + Native.set_active_workers(@server_id, active_count) + end + end +end diff --git a/sig/kino.rbs b/sig/kino.rbs index 6ec5209..47e5923 100644 --- a/sig/kino.rbs +++ b/sig/kino.rbs @@ -110,6 +110,8 @@ module Kino def bind: (String host) -> untyped def port: (int port) -> untyped def workers: (int count) -> untyped + def max_workers: (int count) -> untyped + def scale_down_after: (Numeric seconds) -> untyped def threads: (int count) -> untyped def mode: (Symbol | String mode) -> untyped def queue_depth: (int depth) -> untyped diff --git a/spec/elastic_pool_spec.rb b/spec/elastic_pool_spec.rb new file mode 100644 index 0000000..cb640e9 --- /dev/null +++ b/spec/elastic_pool_spec.rb @@ -0,0 +1,279 @@ +# frozen_string_literal: true + +require "json" +require "tmpdir" + +RSpec.describe "elastic pool" do + let(:ok_app) { ->(_env) { [200, {"content-type" => "text/plain"}, ["ok"]] } } + + def monotonic = Process.clock_gettime(Process::CLOCK_MONOTONIC) + + def wait_for(timeout = 5) + deadline = monotonic + timeout + sleep 0.02 until yield || monotonic > deadline + end + + # Slow enough that a burst keeps the queue non-empty across several + # scaler ticks (100 ms), fast enough that specs stay short. + def slow_shareable_app + Ractor.shareable_proc do |env| + Kino.sleep(0.4) if env["PATH_INFO"] == "/slow" + [200, {"content-type" => "text/plain"}, ["ok"]] + end + end + + def burst(host, port, count, path = "/slow") + Array.new(count) { Thread.new { Net::HTTP.get_response(host, path, port) } } + end + + describe "configuration" do + it "loads max_workers and scale_down_after from the DSL" do + Dir.mktmpdir do |dir| + path = File.join(dir, "kino.rb") + File.write(path, "workers 2\nmax_workers 8\nscale_down_after 10\n") + config = Kino::Configuration.new.load_file(path) + + expect(config[:max_workers]).to eq(8) + expect(config[:scale_down_after]).to eq(10) + end + end + + it "defaults to a fixed pool" do + config = Kino::Configuration.new + + expect(config[:max_workers]).to be_nil + expect(config[:scale_down_after]).to be_nil + end + + it "rejects a ceiling below the floor" do + expect { Kino::Server.new(ok_app, mode: :threaded, workers: 4, max_workers: 2) } + .to raise_error(ArgumentError, /max_workers \(2\) must be at least workers \(4\)/) + end + + it "rejects a non-positive scale_down_after" do + expect { Kino::Server.new(ok_app, mode: :threaded, workers: 1, max_workers: 2, scale_down_after: 0) } + .to raise_error(ArgumentError, /scale_down_after/) + end + end + + describe "stats" do + it "reports the pool bounds and live count before and after start" do + server = Kino::Server.new(ok_app, mode: :threaded, workers: 1, threads: 1, max_workers: 3) + expect(server.stats).to include(workers: 1, max_workers: 3, active_workers: 1, + scale_ups: 0, scale_downs: 0) + + server.start + expect(server.stats).to include(workers: 1, max_workers: 3, active_workers: 1, + scale_ups: 0, scale_downs: 0) + expect(server.stats[:worker_status]).to all(include(retired: false)) + ensure + server&.shutdown(timeout: 0.2) + end + + it "reports the floor as the ceiling for a fixed pool" do + server = Kino::Server.new(ok_app, mode: :threaded, workers: 2, threads: 1) + + expect(server.stats).to include(max_workers: 2, active_workers: 2) + end + end + + describe "ractor mode" do + it "grows under queue pressure, shrinks back to the floor when idle, and recycles slots" do + server = Kino::Server.new(slow_shareable_app, mode: :ractor, workers: 1, threads: 1, + max_workers: 3, scale_down_after: 0.3).start + host, port = "127.0.0.1", server.port + + requests = burst(host, port, 4) + wait_for { server.stats[:active_workers] == 3 } + expect(server.stats).to include(active_workers: 3, scale_ups: 2) + requests.each { |t| expect(t.value.code).to eq("200") } + + wait_for { server.stats[:active_workers] == 1 } + expect(server.stats).to include(active_workers: 1, scale_downs: 2) + slots = server.stats[:worker_status].length + + requests = burst(host, port, 4) + wait_for { server.stats[:active_workers] == 3 } + requests.each { |t| expect(t.value.code).to eq("200") } + # Retired slots were handed to the new workers: the registry did not grow. + expect(server.stats[:worker_status].length).to eq(slots) + ensure + server&.shutdown(timeout: 1) + end + + it "fires on_worker_exit with a nil cause for each retired worker" do + exits = Queue.new + server = Kino::Server.new(slow_shareable_app, mode: :ractor, workers: 1, threads: 1, + max_workers: 3, scale_down_after: 0.3, + on_worker_exit: ->(index, cause) { exits << [index, cause] }).start + host, port = "127.0.0.1", server.port + + burst(host, port, 4).each(&:join) + wait_for { server.stats[:active_workers] == 1 } + + retired = Array.new(2) { exits.pop(timeout: 2) } + expect(retired.map(&:last)).to eq([nil, nil]) + expect(retired.map(&:first)).to all(be_a(Integer)) + ensure + server&.shutdown(timeout: 1) + end + end + + describe "threaded mode" do + let(:slow_app) do + lambda do |env| + sleep 0.4 if env["PATH_INFO"] == "/slow" + [200, {"content-type" => "text/plain"}, ["ok"]] + end + end + + it "grows under queue pressure, shrinks back to the floor when idle, and recycles slots" do + server = Kino::Server.new(slow_app, mode: :threaded, workers: 1, threads: 1, + max_workers: 3, scale_down_after: 0.3).start + host, port = "127.0.0.1", server.port + + requests = burst(host, port, 4) + wait_for { server.stats[:active_workers] == 3 } + expect(server.stats).to include(active_workers: 3, scale_ups: 2) + requests.each { |t| expect(t.value.code).to eq("200") } + + wait_for { server.stats[:active_workers] == 1 } + expect(server.stats).to include(active_workers: 1, scale_downs: 2) + slots = server.stats[:worker_status].length + + requests = burst(host, port, 4) + wait_for { server.stats[:active_workers] == 3 } + requests.each { |t| expect(t.value.code).to eq("200") } + expect(server.stats[:worker_status].length).to eq(slots) + ensure + server&.shutdown(timeout: 1) + end + + it "grows and retires whole workers when a worker is several threads" do + server = Kino::Server.new(slow_app, mode: :threaded, workers: 1, threads: 2, + max_workers: 2, scale_down_after: 0.3).start + host, port = "127.0.0.1", server.port + + requests = burst(host, port, 6) + wait_for { server.stats[:active_workers] == 2 } + expect(server.stats[:worker_status].length).to eq(4) # two workers x two slots + requests.each { |t| expect(t.value.code).to eq("200") } + + wait_for { server.stats[:active_workers] == 1 } + expect(server.stats[:worker_status].count { |w| w[:retired] }).to eq(2) + ensure + server&.shutdown(timeout: 1) + end + + it "fires after_worker_boot for grown workers and on_worker_exit with nil for retired ones" do + boots = Queue.new + exits = Queue.new + server = Kino::Server.new(slow_app, mode: :threaded, workers: 1, threads: 1, + max_workers: 3, scale_down_after: 0.3, + after_worker_boot: ->(id) { boots << id }, + on_worker_exit: ->(id, cause) { exits << [id, cause] }).start + host, port = "127.0.0.1", server.port + + burst(host, port, 4).each(&:join) + wait_for { server.stats[:active_workers] == 1 } + + expect(boots.size).to eq(3) + retired = Array.new(2) { exits.pop(timeout: 2) } + expect(retired.map(&:last)).to eq([nil, nil]) + ensure + server&.shutdown(timeout: 1) + end + + it "stays within the floor and the ceiling" do + server = Kino::Server.new(slow_app, mode: :threaded, workers: 2, threads: 1, + max_workers: 3, scale_down_after: 0.2).start + host, port = "127.0.0.1", server.port + + requests = burst(host, port, 8) + peak = 0 + until requests.none?(&:alive?) + peak = [peak, server.stats[:active_workers]].max + sleep 0.02 + end + expect(peak).to eq(3) + + wait_for { server.stats[:active_workers] == 2 } + sleep 0.5 # well past scale_down_after: still at the floor + expect(server.stats[:active_workers]).to eq(2) + ensure + server&.shutdown(timeout: 1) + end + end + + describe "lane dispatch" do + it "loses no request while the pool grows and shrinks underneath the dispatcher" do + server = Kino::Server.new(slow_shareable_app, mode: :ractor, lanes: true, workers: 1, threads: 1, + max_workers: 4, scale_down_after: 0.2).start + host, port = "127.0.0.1", server.port + + 3.times do + burst(host, port, 6).each { |t| expect(t.value.code).to eq("200") } + # Past scale_down_after: workers retire while these are dispatched. + sleep 0.35 + 10.times { expect(Net::HTTP.get_response(host, "/", port).code).to eq("200") } + end + + expect(server.stats[:scale_ups]).to be >= 3 + expect(server.stats[:scale_downs]).to be >= 1 + expect(server.stats[:rejected]).to eq(0) + ensure + server&.shutdown(timeout: 1) + end + end + + describe "control plane" do + it "reports the pool in /stats and /metrics" do + server = Kino::Server.new(ok_app, mode: :threaded, workers: 1, threads: 1, max_workers: 3, + control_bind: "127.0.0.1:0").start + + stats = JSON.parse(Net::HTTP.get_response("127.0.0.1", "/stats", server.control_port).body) + expect(stats).to include("workers" => 1, "max_workers" => 3, "active_workers" => 1, + "scale_ups" => 0, "scale_downs" => 0) + expect(stats["worker_status"]).to all(include("retired" => false)) + + metrics = Net::HTTP.get_response("127.0.0.1", "/metrics", server.control_port).body + expect(metrics).to include("kino_max_workers 3", "kino_active_workers 1", + "kino_scale_ups_total 0", "kino_scale_downs_total 0") + ensure + server&.shutdown(timeout: 0.2) + end + end + + describe "the Ruby scheduler cap" do + # Ruby runs non-main ractors' Ruby code on at most RUBY_MAX_CPU native + # threads (default 8); a pool past that gains no CPU parallelism. + def with_max_cpu(value) + previous = ENV["RUBY_MAX_CPU"] + ENV["RUBY_MAX_CPU"] = value + yield + ensure + ENV["RUBY_MAX_CPU"] = previous + end + + def start_and_stop(**opts) + app = Ractor.shareable_proc { |_env| [200, {"content-type" => "text/plain"}, ["ok"]] } + capture_native_stderr do + server = Kino::Server.new(app, mode: :ractor, threads: 1, **opts).start + server.shutdown(timeout: 0.2) + end + end + + it "warns at start when a ractor pool's ceiling exceeds the cap" do + err = with_max_cpu("2") { start_and_stop(workers: 1, max_workers: 3) } + + expect(err).to include("RUBY_MAX_CPU") + expect(err).to include("3") + end + + it "stays quiet when the pool fits under the cap" do + err = with_max_cpu("2") { start_and_stop(workers: 1, max_workers: 2) } + + expect(err).not_to include("RUBY_MAX_CPU") + end + end +end diff --git a/spec/pool_scaler_spec.rb b/spec/pool_scaler_spec.rb new file mode 100644 index 0000000..1adcb8c --- /dev/null +++ b/spec/pool_scaler_spec.rb @@ -0,0 +1,154 @@ +# frozen_string_literal: true + +# The pool seam the scaler drives: worker groups (index => slot ids), +# grow, retire. Records what the scaler asked for. +class ScalerFakePool + attr_reader :grown, :retired + + def initialize(groups) + @groups = groups + @grown = [] + @retired = [] + @next_index = groups.keys.max.to_i + 1 + end + + def active_count = @groups.size + + def groups = @groups.dup + + def grow + index = @next_index + @next_index += 1 + @groups[index] = [index] + @grown << index + index + end + + def retire(index) + @retired << index + @groups.delete(index) + true + end +end + +RSpec.describe Kino::PoolScaler do + # A worker_stats row: [slot, served, in_flight, busy_ms, quarantined, retired]. + def row(slot, served:, in_flight: 0) + [slot, served, in_flight, 0, false, false] + end + + def scaler(pool, floor: 1, ceiling: 4, scale_down_after: 30) + described_class.new(server_id: 0, pool: pool, floor: floor, ceiling: ceiling, + scale_down_after: scale_down_after) + end + + describe "scaling up" do + it "grows by one worker once the queue has been non-empty on two consecutive ticks" do + pool = ScalerFakePool.new({0 => [0]}) + s = scaler(pool) + busy = [row(0, served: 1, in_flight: 1)] + + s.step(0.0, 0, busy) + s.step(0.1, 3, busy) + expect(pool.grown).to be_empty + + s.step(0.2, 3, busy) + expect(pool.grown).to eq([1]) + end + + it "ignores a single-tick blip" do + pool = ScalerFakePool.new({0 => [0]}) + s = scaler(pool) + busy = [row(0, served: 1, in_flight: 1)] + + s.step(0.1, 5, busy) + s.step(0.2, 0, busy) + s.step(0.3, 5, busy) + + expect(pool.grown).to be_empty + end + + it "keeps growing one per tick while pressure lasts" do + pool = ScalerFakePool.new({0 => [0]}) + s = scaler(pool) + busy = [row(0, served: 1, in_flight: 1)] + + 4.times { |i| s.step(i / 10.0, 3, busy) } + + expect(pool.grown).to eq([1, 2, 3]) + end + + it "never grows past the ceiling" do + pool = ScalerFakePool.new({0 => [0], 1 => [1], 2 => [2]}) + s = scaler(pool, ceiling: 3) + busy = pool.groups.keys.map { |i| row(i, served: 1, in_flight: 1) } + + 5.times { |i| s.step(i / 10.0, 9, busy) } + + expect(pool.grown).to be_empty + end + end + + describe "scaling down" do + it "retires a worker above the floor once it has been idle for scale_down_after" do + pool = ScalerFakePool.new({0 => [0], 1 => [1]}) + s = scaler(pool, floor: 1, scale_down_after: 10) + + # Worker 0 keeps serving between ticks; worker 1 never does. Ticks + # are whole seconds so the idle arithmetic is exact. + (0..12).each do |t| + s.step(t, 0, [row(0, served: t + 1), row(1, served: 5)]) + expect(pool.retired).to be_empty if t < 11 + end + + expect(pool.retired).to eq([1]) + end + + it "does not count a worker as idle while its served count keeps moving" do + pool = ScalerFakePool.new({0 => [0], 1 => [1]}) + s = scaler(pool, floor: 1, scale_down_after: 5) + + (0..20).each do |t| + # in_flight is always 0 at the sampling instant, but the counter moves. + s.step(t, 0, [row(0, served: 1), row(1, served: t + 1)]) + end + + expect(pool.retired).to eq([0]) + end + + it "never retires below the floor" do + pool = ScalerFakePool.new({0 => [0], 1 => [1]}) + s = scaler(pool, floor: 2, scale_down_after: 2) + idle = [row(0, served: 1), row(1, served: 1)] + + 10.times { |t| s.step(t, 0, idle) } + + expect(pool.retired).to be_empty + end + + it "retires at most one worker per tick" do + pool = ScalerFakePool.new({0 => [0], 1 => [1], 2 => [2]}) + s = scaler(pool, floor: 1, scale_down_after: 2) + idle = [row(0, served: 1), row(1, served: 1), row(2, served: 1)] + + s.step(0, 0, idle) + s.step(1, 0, idle) # idle since 1 + s.step(3, 0, idle) # first retirement + expect(pool.retired.size).to eq(1) + + s.step(4, 0, idle) + expect(pool.retired.size).to eq(2) + end + + it "treats every slot of a multi-thread worker as one unit" do + pool = ScalerFakePool.new({0 => [0, 1], 1 => [2, 3]}) + s = scaler(pool, floor: 1, scale_down_after: 2) + # Worker 1's second slot stays busy: the worker is not idle. + rows = [row(0, served: 1), row(1, served: 1), row(2, served: 1), row(3, served: 1, in_flight: 1)] + + 6.times { |t| s.step(t, 0, rows) } + + expect(pool.retired).to eq([0]) + end + end +end diff --git a/spec/startup_info_spec.rb b/spec/startup_info_spec.rb index c1d8bd1..585f5d7 100644 --- a/spec/startup_info_spec.rb +++ b/spec/startup_info_spec.rb @@ -19,6 +19,17 @@ end end + it "prints the pool range when the pool is elastic" do + app = "run ->(_env) { [200, {}, []] }\n" + + Dir.mktmpdir("kino-startup") do |dir| + config = "workers 1\nmax_workers 3\nthreads 1\nmode :threaded\n" + with_cli_server(dir, config, app) do |_port, out| + expect(File.read(out)).to include("- mode: threaded, 1-3 workers × 1 thread") + end + end + end + it "names the control plane when one is bound" do app = "run ->(_env) { [200, {}, []] }\n" From f41565a5cf66a4f2b9e9ce564a0bd32d89baa474 Mon Sep 17 00:00:00 2001 From: Yaroslav Markin Date: Mon, 7 Sep 2026 23:37:49 +0400 Subject: [PATCH 3/4] Document the elastic pool --- CHANGELOG.md | 11 ++++++++++ README.md | 52 +++++++++++++++++++++++++++++++++++++++++++++ doc/architecture.md | 10 +++++++++ 3 files changed, 73 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6d1af8c..5c9e533 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,16 @@ ## [Unreleased] +- Experimental elastic worker pool. Set `max_workers` above `workers` + and the pool grows under load, one worker at a time, then shrinks + back to `workers` once the extra workers have sat idle for + `scale_down_after` seconds (default 30). Works in both modes; a + retiring worker finishes its request first. `stats`, `/stats` and + `/metrics` gain `max_workers`, `active_workers`, `scale_ups` and + `scale_downs`, and each `worker_status` row gains `retired`. Leave + `max_workers` unset and the pool is fixed, as before. +- Ractor mode warns at boot when `max_workers` (or `workers`) exceeds + `RUBY_MAX_CPU` (default 8), Ruby's cap on how many ractors run Ruby + code at once. Set the variable to your worker count to lift it. - Ractor-readiness fixes, so external Ractor audits of Kino pass: the env string caches root their strings through the lock-free pin slab instead of per-value GC registration (unsynchronized across ractors diff --git a/README.md b/README.md index 1bbd9af..ab26ccd 100644 --- a/README.md +++ b/README.md @@ -259,6 +259,8 @@ server = Kino::Server.new(app, bind: "127.0.0.1", # or "unix:///run/kino.sock" behind a proxy port: 9292, # 0 = ephemeral; read back via server.port workers: Kino.available_parallelism, # ractors (parallelism); the default + max_workers: nil, # experimental: grow past workers under load (see Elastic pool) + scale_down_after: 30, # seconds idle before an extra worker retires threads: 1, # per worker; ractor default 1, threaded default 3 mode: :auto, # :auto | :ractor | :threaded queue_depth: 1024, # bounded queue; overflow → 503 @@ -418,6 +420,56 @@ Kino fires four lifecycle hooks alongside `on_error`, split by firing context. A raising hook is logged and never kills a worker. +## Elastic pool (experimental) + +Size the pool for the quiet hours and let it grow for the busy ones: + +```ruby +# kino.rb +workers 4 # always running +max_workers 16 # reached only under load +scale_down_after 30 # seconds idle before an extra worker retires +``` + +Or `Kino::Server.new(app, workers: 4, max_workers: 16)`. Leave +`max_workers` unset and the pool is fixed at `workers`, as before. + +While requests wait in the queue, Kino adds a worker every 100 ms until +the queue clears or the pool hits `max_workers`. When the load passes, +workers above `workers` retire one at a time after `scale_down_after` +seconds idle, each finishing its current request first. Same behavior +in `:ractor` and `:threaded` mode. + +**Use it when your app waits**: on databases, upstream services, slow +clients. With `workers` at your core count, all workers can be blocked +on I/O while cores sit idle; a higher ceiling puts those cores to work, +and a ractor starts in microseconds, so the pool follows load closely. +Pure CPU work gains nothing past the core count. In `:ractor` mode Ruby +itself runs at most `RUBY_MAX_CPU` ractors' Ruby code at once (default +8), and Kino warns at boot when the pool can exceed it. On a bigger box: + +```sh +RUBY_MAX_CPU=16 kino +``` + +**Watch it breathe** in `server.stats`, `GET /stats` and `GET /metrics`: + +```sh +$ curl -s localhost:9293/stats | jq '{workers, max_workers, active_workers, scale_ups, scale_downs}' +{ + "workers": 4, + "max_workers": 16, + "active_workers": 9, + "scale_ups": 12, + "scale_downs": 7 +} +``` + +Prometheus gets `kino_max_workers`, `kino_active_workers`, +`kino_scale_ups_total` and `kino_scale_downs_total`. `after_worker_boot` +fires for every worker the pool adds, `on_worker_exit` (with a nil +cause) for every one it retires. + ## Stuck-worker quarantine `quarantine_timeout: seconds` (or `quarantine_timeout 60` in `kino.rb`) diff --git a/doc/architecture.md b/doc/architecture.md index 9cecd8e..71eee29 100644 --- a/doc/architecture.md +++ b/doc/architecture.md @@ -31,6 +31,16 @@ Puma-style two-level: `workers × threads`. - Identical machinery either way: the flume queue is MPMC, a "worker slot" is per-thread, and the worker loop (`lib/kino/worker.rb`) is shared verbatim. +- Elastic pool (`max_workers`): a scaler thread on the main ractor + samples queue depth and the per-slot sensors every 100 ms, adds one + worker per sample while requests wait, and retires the longest-idle + worker above the floor after `scale_down_after`. Retirement is a + per-slot flag raised under the slot's lane lock: the lane dispatcher + skips the slot, the take loop honors the flag at its next idle tick (a + request already taken finishes first, a lane worker drains its own + lane), and the slot is reset and reused by the next worker, so the + slot table never grows with churn. Both pools (the ractor supervisor + and the threaded pool) expose the same grow/retire/groups seam. - Experimental `lanes true` replaces the one shared queue with a small private queue per worker slot (awake-preferring dispatch, work stealing); see [benchmarks](benchmarks.md#lane-dispatch-experimental-lanes-true). From 54766ac9082cc89bf6f194376ee261b0791e24be Mon Sep 17 00:00:00 2001 From: Yaroslav Markin Date: Tue, 8 Sep 2026 09:14:48 +0400 Subject: [PATCH 4/4] Extract the monitor and slot bank shared by both pools, test slot retirement --- ext/kino/src/control.rs | 30 +++++++++ ext/kino/src/queue.rs | 40 +++++++++++- ext/kino/src/registry.rs | 111 +++++++++++++++++++++++---------- ext/kino/src/server.rs | 107 ++++++++++++++++++++++++++++--- lib/kino.rb | 2 + lib/kino/monitor.rb | 52 +++++++++++++++ lib/kino/pool_scaler.rb | 33 +--------- lib/kino/quarantine_monitor.rb | 42 ++----------- lib/kino/ractor_supervisor.rb | 40 ++++++------ lib/kino/server.rb | 8 +-- lib/kino/slot_bank.rb | 31 +++++++++ lib/kino/threaded_pool.rb | 33 ++++------ spec/monitor_spec.rb | 48 ++++++++++++++ 13 files changed, 416 insertions(+), 161 deletions(-) create mode 100644 lib/kino/monitor.rb create mode 100644 lib/kino/slot_bank.rb create mode 100644 spec/monitor_spec.rb diff --git a/ext/kino/src/control.rs b/ext/kino/src/control.rs index 800d506..45b4345 100644 --- a/ext/kino/src/control.rs +++ b/ext/kino/src/control.rs @@ -751,6 +751,36 @@ mod tests { assert!(text.contains("kino_scale_downs_total 1")); } + #[test] + fn worker_status_reports_retired_slots() { + let server = crate::registry::test_server(false, 4); + server.register_worker(); + server.register_worker(); + server.slots.read()[1].retire(); + + let status = collect_worker_status(&server); + + assert_eq!( + status.iter().map(|w| w.retired).collect::>(), + vec![false, true] + ); + assert!(status.iter().all(|w| w.busy_ms == 0)); + } + + #[test] + fn snapshot_reads_the_pool_counters() { + let server = crate::registry::test_server(false, 4); + server.active_workers.store(3, Ordering::Relaxed); + server.scale_ups.fetch_add(2, Ordering::Relaxed); + server.scale_downs.fetch_add(1, Ordering::Relaxed); + + let s = StatsSnapshot::take(&server); + + assert_eq!(s.active_workers, 3); + assert_eq!(s.max_workers, server.topology.max_workers); + assert_eq!((s.scale_ups, s.scale_downs), (2, 1)); + } + #[test] fn metrics_text_includes_lane_depth_samples_when_lanes_are_on() { let mut s = snapshot(crate::registry::STATE_READY); diff --git a/ext/kino/src/queue.rs b/ext/kino/src/queue.rs index b7618c0..ff07c6a 100644 --- a/ext/kino/src/queue.rs +++ b/ext/kino/src/queue.rs @@ -29,7 +29,7 @@ type Taken = Option; /// or end the loop because the pool scaler retired this slot. A retired /// lane worker first empties its own lane, one item per timeout: the /// dispatcher stopped feeding it when the flag went up (see -/// registry::ServerInner::retire_slot for the ordering), so whatever is +/// registry::WorkerSlot::retire for the ordering), so whatever is /// still there is the last of it, and nothing already assigned is /// orphaned. Only reached with the queue momentarily empty, so the fast /// path pays nothing for it. @@ -272,7 +272,7 @@ mod tests { fn a_retired_shared_queue_worker_ends_its_loop_at_the_next_timeout() { let server = test_server(false, 4); server.register_worker(); - server.retire_slot(0); + server.slots.read()[0].retire(); let slot = server.slots.read()[0].clone(); assert!(matches!(after_timeout(&slot), Some(None))); @@ -289,11 +289,45 @@ mod tests { .expect("lane open") .send(test_ctx()) .expect("lane has room"); - server.retire_slot(0); + server.slots.read()[0].retire(); // One item still assigned to this lane: serve it, not orphan it. assert!(matches!(after_timeout(&slot), Some(Some(_)))); // Lane empty now: leave. assert!(matches!(after_timeout(&slot), Some(None))); } + + #[test] + fn a_retired_lane_worker_drains_every_item_one_per_timeout() { + let server = test_server(true, 4); + server.register_worker(); + let slot = server.slots.read()[0].clone(); + for _ in 0..3 { + slot.lane_tx + .lock() + .as_ref() + .expect("lane open") + .send(test_ctx()) + .expect("lane has room"); + } + server.slots.read()[0].retire(); + + for _ in 0..3 { + assert!(matches!(after_timeout(&slot), Some(Some(_)))); + } + assert!(matches!(after_timeout(&slot), Some(None))); + } + + #[test] + fn a_reset_slot_keeps_waiting_again() { + let server = test_server(false, 4); + server.register_worker(); + let slot = server.slots.read()[0].clone(); + slot.retire(); + assert!(matches!(after_timeout(&slot), Some(None))); + + slot.reset(); + + assert!(after_timeout(&slot).is_none()); + } } diff --git a/ext/kino/src/registry.rs b/ext/kino/src/registry.rs index 622169d..c8028d6 100644 --- a/ext/kino/src/registry.rs +++ b/ext/kino/src/registry.rs @@ -170,8 +170,8 @@ pub struct WorkerSlot { pub quarantined: std::sync::atomic::AtomicBool, /// Set by the pool scaler to send this slot's worker home: the lane /// dispatcher skips it, and the take loop ends at its next idle tick - /// (a request already taken finishes first). Cleared by `reset_slot` - /// when the slot is handed to a new worker. + /// (a request already taken finishes first). Cleared by `reset` when + /// the slot is handed to a new worker. pub retired: std::sync::atomic::AtomicBool, } @@ -269,6 +269,32 @@ impl WorkerSlot { retired: std::sync::atomic::AtomicBool::new(false), } } + + /// Send this slot's worker home once it is idle. The flag is raised + /// under the lane lock, the same lock the lane dispatcher holds while + /// it checks the flag and sends: any dispatch that saw "not retired" + /// has landed in the lane before the worker can see the flag and + /// drain, so no request is orphaned. + pub fn retire(&self) { + let _lane = self.lane_tx.lock(); + self.retired.store(true, Ordering::SeqCst); + } + + /// Return a retired slot to its fresh state so a new worker can take + /// it over (slots are never removed; recycling keeps the registry + /// from growing with every scale-up). The lane channel stays, the new + /// occupant simply starts taking from it; the quarantine mark stays + /// too, a quarantined slot is never handed out. + pub fn reset(&self) { + let _lane = self.lane_tx.lock(); + self.current.lock().clear(); + self.retired.store(false, Ordering::SeqCst); + self.parked.store(false, Ordering::SeqCst); + self.interrupted.store(false, Ordering::SeqCst); + self.served.store(0, Ordering::Relaxed); + self.in_flight.store(0, Ordering::Relaxed); + self.last_started_ms.store(0, Ordering::Relaxed); + } } static REGISTRY: OnceLock>>> = OnceLock::new(); @@ -328,35 +354,6 @@ impl ServerInner { slots.len() - 1 } - /// Send a slot's worker home once it is idle. The flag is raised under - /// the slot's lane lock, the same lock the lane dispatcher holds while - /// it checks the flag and sends: any dispatch that saw "not retired" - /// has landed in the lane before the worker can see the flag and - /// drain, so no request is orphaned. Unknown ids are ignored. - pub fn retire_slot(&self, worker_id: usize) { - if let Some(slot) = self.slots.read().get(worker_id) { - let _lane = slot.lane_tx.lock(); - slot.retired.store(true, Ordering::SeqCst); - } - } - - /// Return a retired slot to its fresh state so a new worker can take - /// it over (slots are never removed; recycling keeps the registry - /// from growing with every scale-up). The lane channel stays: the new - /// occupant simply starts taking from it. Unknown ids are ignored. - pub fn reset_slot(&self, worker_id: usize) { - if let Some(slot) = self.slots.read().get(worker_id) { - let _lane = slot.lane_tx.lock(); - slot.current.lock().clear(); - slot.retired.store(false, Ordering::SeqCst); - slot.parked.store(false, Ordering::SeqCst); - slot.interrupted.store(false, Ordering::SeqCst); - slot.served.store(0, Ordering::Relaxed); - slot.in_flight.store(0, Ordering::Relaxed); - slot.last_started_ms.store(0, Ordering::Relaxed); - } - } - pub fn slot( &self, ruby: &magnus::Ruby, @@ -590,10 +587,10 @@ mod tests { slot.last_started_ms.store(9, Ordering::Relaxed); slot.current.lock().push(std::sync::Weak::new()); - server.retire_slot(0); + slot.retire(); assert!(slot.retired.load(Ordering::Relaxed)); - server.reset_slot(0); + slot.reset(); assert!(!slot.retired.load(Ordering::Relaxed)); assert!(!slot.parked.load(Ordering::Relaxed)); assert!(!slot.interrupted.load(Ordering::Relaxed)); @@ -605,6 +602,54 @@ mod tests { assert!(slot.lane_tx.lock().is_some()); } + #[test] + fn reset_keeps_a_quarantine_mark() { + let server = test_server(false, 4); + server.register_worker(); + let slot = server.slots.read()[0].clone(); + slot.quarantined.store(true, Ordering::Relaxed); + + slot.retire(); + slot.reset(); + + // A wedged slot stays flagged even if it ever came back around. + assert!(slot.quarantined.load(Ordering::Relaxed)); + } + + #[test] + fn retire_releases_the_lane_lock() { + let server = test_server(true, 4); + server.register_worker(); + let slot = server.slots.read()[0].clone(); + + slot.retire(); + + assert!(slot.lane_tx.try_lock().is_some()); + } + + #[test] + fn retire_waits_for_a_dispatcher_holding_the_lane_lock() { + let server = test_server(true, 4); + server.register_worker(); + let slot = server.slots.read()[0].clone(); + + // A dispatcher that has already read "not retired" and is sending. + let sending = slot.lane_tx.lock(); + let retiring = std::thread::spawn({ + let slot = slot.clone(); + move || slot.retire() + }); + std::thread::sleep(Duration::from_millis(30)); + assert!( + !slot.retired.load(Ordering::SeqCst), + "the flag must not go up under a dispatcher's lock" + ); + + drop(sending); + retiring.join().expect("retire thread"); + assert!(slot.retired.load(Ordering::SeqCst)); + } + #[test] fn queue_histogram_buckets_by_wait() { let h = QueueHistogram::new(); diff --git a/ext/kino/src/server.rs b/ext/kino/src/server.rs index 074de65..cf682e6 100644 --- a/ext/kino/src/server.rs +++ b/ext/kino/src/server.rs @@ -726,8 +726,8 @@ fn try_dispatch(server: &ServerInner, mut ctx: BoxedCtx) -> Dispatch { let guard = slot.lane_tx.lock(); let Some(tx) = guard.as_ref() else { continue }; any_open = true; - // Checked under the lane lock, which retire_slot also takes: - // see registry::ServerInner::retire_slot for the ordering. + // Checked under the lane lock, which WorkerSlot::retire also + // takes: see there for the ordering. if slot.retired.load(Ordering::SeqCst) { continue; } @@ -952,11 +952,10 @@ pub fn quarantine_slot(ruby: &Ruby, server_id: u64, worker_id: usize) -> Result< } /// Send a slot's worker home once idle (pool scale-down); see -/// registry::ServerInner::retire_slot. +/// registry::WorkerSlot::retire. An unknown id is a caller bug: raise. pub fn retire_slot(ruby: &Ruby, server_id: u64, worker_id: usize) -> Result<(), Error> { if let Some(server) = registry::try_get(server_id) { - server.slot(ruby, worker_id)?; // unknown id is a caller bug: raise - server.retire_slot(worker_id); + server.slot(ruby, worker_id)?.retire(); } Ok(()) } @@ -964,8 +963,7 @@ pub fn retire_slot(ruby: &Ruby, server_id: u64, worker_id: usize) -> Result<(), /// Hand a retired slot to a new worker (pool scale-up reusing a slot). pub fn reset_slot(ruby: &Ruby, server_id: u64, worker_id: usize) -> Result<(), Error> { if let Some(server) = registry::try_get(server_id) { - server.slot(ruby, worker_id)?; - server.reset_slot(worker_id); + server.slot(ruby, worker_id)?.reset(); } Ok(()) } @@ -1454,6 +1452,18 @@ mod tests { assert_eq!(stream_capacity(&topology), 24); } + #[test] + fn stream_capacity_of_a_fixed_pool_is_its_slot_count() { + let topology = registry::Topology { + mode: "threaded".to_string(), + workers: 4, + max_workers: 4, + threads: 3, + batch: 1, + }; + assert_eq!(stream_capacity(&topology), 12); + } + #[tokio::test] async fn streams_beyond_the_advertised_cap_queue_instead_of_failing() { use http_body_util::Full; @@ -1671,7 +1681,7 @@ mod tests { let server = test_server(true, 4); server.register_worker(); server.register_worker(); - server.retire_slot(0); + server.slots.read()[0].retire(); // A retiring worker gets nothing new: both land on slot 1. assert!(matches!(try_dispatch(&server, test_ctx()), Dispatch::Sent)); @@ -1683,7 +1693,7 @@ mod tests { fn dispatch_retries_rather_than_closing_when_only_retired_lanes_remain() { let server = test_server(true, 4); server.register_worker(); - server.retire_slot(0); + server.slots.read()[0].retire(); // Retired is not draining: the caller keeps retrying until the // queue timeout, so a replacement worker can still pick this up. @@ -1692,4 +1702,83 @@ mod tests { Dispatch::Full(_) )); } + + #[test] + fn dispatch_skips_a_retired_lane_even_when_it_is_the_only_awake_one() { + let server = test_server(true, 4); + server.register_worker(); + server.register_worker(); + server.slots.read()[0].retire(); + server.slots.read()[1].parked.store(true, Ordering::Relaxed); + + // The second pass (parked lanes allowed) still skips the retired one. + assert!(matches!(try_dispatch(&server, test_ctx()), Dispatch::Sent)); + assert_eq!(server.lane_depths(), Some(vec![0, 1])); + } + + #[test] + fn dispatch_keeps_retrying_when_the_other_lanes_are_closed() { + let server = test_server(true, 4); + server.register_worker(); + server.register_worker(); + server.slots.read()[0].lane_tx.lock().take(); // a crashed worker's lane + server.slots.read()[1].retire(); + + assert!(matches!( + try_dispatch(&server, test_ctx()), + Dispatch::Full(_) + )); + } + + #[test] + fn a_reset_lane_receives_dispatches_again() { + let server = test_server(true, 4); + server.register_worker(); + server.register_worker(); + server.slots.read()[0].retire(); + assert!(matches!(try_dispatch(&server, test_ctx()), Dispatch::Sent)); + assert_eq!(server.lane_depths(), Some(vec![0, 1])); + + server.slots.read()[0].reset(); + for _ in 0..2 { + assert!(matches!(try_dispatch(&server, test_ctx()), Dispatch::Sent)); + } + // Back in the round-robin. + assert_eq!(server.lane_depths(), Some(vec![1, 2])); + } + + #[test] + fn nothing_lands_in_a_lane_after_its_retirement_and_final_drain() { + let server = test_server(true, 4); + server.register_worker(); + server.register_worker(); + let stop = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let dispatcher = std::thread::spawn({ + let server = server.clone(); + let stop = stop.clone(); + move || { + while !stop.load(Ordering::Relaxed) { + let _ = try_dispatch(&server, test_ctx()); + // Lane 1 keeps draining, so dispatch never stalls on Full. + let slots = server.slots.read(); + if let Some(rx) = slots[1].lane_rx.as_ref() { + while rx.try_recv().is_ok() {} + } + } + } + }); + std::thread::sleep(Duration::from_millis(20)); + + // The scaler retires slot 0; its worker then drains what it holds. + let slot = server.slots.read()[0].clone(); + slot.retire(); + if let Some(rx) = slot.lane_rx.as_ref() { + while rx.try_recv().is_ok() {} + } + + std::thread::sleep(Duration::from_millis(50)); + stop.store(true, Ordering::Relaxed); + dispatcher.join().expect("dispatcher thread"); + assert_eq!(server.lane_depths(), Some(vec![0, 0])); + } } diff --git a/lib/kino.rb b/lib/kino.rb index d73d6ce..959dc29 100644 --- a/lib/kino.rb +++ b/lib/kino.rb @@ -57,8 +57,10 @@ def self.available_parallelism require_relative "kino/hook_fire" require_relative "kino/worker_hooks" require_relative "kino/worker" +require_relative "kino/slot_bank" require_relative "kino/ractor_supervisor" require_relative "kino/threaded_pool" +require_relative "kino/monitor" require_relative "kino/quarantine_monitor" require_relative "kino/pool_scaler" require_relative "kino/server" diff --git a/lib/kino/monitor.rb b/lib/kino/monitor.rb new file mode 100644 index 0000000..4195c60 --- /dev/null +++ b/lib/kino/monitor.rb @@ -0,0 +1,52 @@ +# frozen_string_literal: true + +module Kino + # @private + # A thread on the main ractor that calls `scan` every `tick` seconds + # until stopped. Main-ractor so it stays responsive when worker ractors + # are wedged. A scan that raises is logged and skipped, never fatal: + # monitors keep the server healthy, they must not take it down. + class Monitor + def initialize(name:, tick:) + @name = name + @tick = tick + @running = false + @thread = nil + end + + def start + @running = true + @thread = Thread.new do + Thread.current.name = @name + run + end + self + end + + def stop + @running = false + @thread&.join(@tick * 2) + end + + private + + def run + tick while @running + rescue => e + Log.error("#{@name} crashed: #{e.class}: #{e.message}") + end + + def tick + scan + rescue => e + Log.error("#{@name} tick error: #{e.class}: #{e.message}") + ensure + sleep @tick + end + + # One poll; subclasses define it. + def scan + raise NotImplementedError, "#{self.class} must define scan" + end + end +end diff --git a/lib/kino/pool_scaler.rb b/lib/kino/pool_scaler.rb index 7c22a72..a8f5fc8 100644 --- a/lib/kino/pool_scaler.rb +++ b/lib/kino/pool_scaler.rb @@ -15,37 +15,21 @@ module Kino # the worker held a request at two consecutive samples and its served # count did not move between them, so a worker serving short requests # between samples is never mistaken for an idle one. - class PoolScaler + class PoolScaler < Monitor TICK = 0.1 # Consecutive ticks with a non-empty queue before the pool grows. PRESSURE_TICKS = 2 def initialize(server_id:, pool:, floor:, ceiling:, scale_down_after:, tick: TICK) + super(name: "pool scaler", tick: tick) @server_id = server_id @pool = pool @floor = floor @ceiling = ceiling @scale_down_after = scale_down_after - @tick = tick @pressure = 0 @last_served = {} @idle_since = {} - @running = false - @thread = nil - end - - def start - @running = true - @thread = Thread.new do - Thread.current.name = "pool-scaler" - run - end - self - end - - def stop - @running = false - @thread&.join(@tick * 2) end # One policy step over one set of observations: the monotonic time, @@ -70,20 +54,9 @@ def step(now, queued, rows) private - def run - tick while @running - rescue => e - Log.error("pool scaler crashed: #{e.class}: #{e.message}") - end - - def tick + def scan queued, _in_flight = Native.queue_stats(@server_id) step(Process.clock_gettime(Process::CLOCK_MONOTONIC), queued, Native.worker_stats(@server_id)) - rescue => e - # A bad tick must never kill the scaler. - Log.error("pool scaler tick error: #{e.class}: #{e.message}") - ensure - sleep @tick end def grow diff --git a/lib/kino/quarantine_monitor.rb b/lib/kino/quarantine_monitor.rb index 3790311..0114f47 100644 --- a/lib/kino/quarantine_monitor.rb +++ b/lib/kino/quarantine_monitor.rb @@ -3,54 +3,22 @@ module Kino # @private # Polls per-slot busy_ms and, past the timeout, quarantines a wedged slot - # and asks the replacer to spawn a fresh worker. Runs one thread on the - # main ractor (uncontended by wedged worker ractors, so it stays - # responsive in :ractor mode). Never interrupts the wedged worker. - class QuarantineMonitor + # and asks the replacer to spawn a fresh worker. Never interrupts the + # wedged worker. + class QuarantineMonitor < Monitor def initialize(server_id:, timeout_ms:, max:, replacer:, tick: 0.5) + super(name: "quarantine monitor", tick: tick) @server_id = server_id @timeout_ms = timeout_ms @max = max @replacer = replacer - @tick = tick @outstanding = 0 @at_cap_logged = false - @running = false - @thread = nil - end - - def start - @running = true - @thread = Thread.new do - Thread.current.name = "quarantine" - run - end - self - end - - def stop - @running = false - @thread&.join(@tick * 2) end private - def run - tick while @running - rescue => e - Log.error("quarantine monitor crashed: #{e.class}: #{e.message}") - end - - def tick - scan_slots - rescue => e - # A bad tick must never kill the monitor. - Log.error("quarantine tick error: #{e.class}: #{e.message}") - ensure - sleep @tick - end - - def scan_slots + def scan Native.worker_stats(@server_id).each do |index, _served, _in_flight, busy_ms, quarantined| next if quarantined || busy_ms <= @timeout_ms diff --git a/lib/kino/ractor_supervisor.rb b/lib/kino/ractor_supervisor.rb index d05cae9..b872657 100644 --- a/lib/kino/ractor_supervisor.rb +++ b/lib/kino/ractor_supervisor.rb @@ -26,12 +26,12 @@ def initialize(server_id, app, workers:, threads:, batch: 1, hooks: nil, on_work @worker_slots = {} @slot_to_worker = {} @replaced = {} - # Worker index => true while its ractor runs; => true once the - # scaler asked it to leave; and the slot ids retired workers gave - # back, for the next worker to take over. + # Worker index => true while its ractor runs, and => true once the + # scaler asked it to leave. Retired workers hand their slots back to + # the bank for the next worker to take over. @live = {} @retiring = {} - @free_slots = [] + @bank = SlotBank.new(server_id) # The first replacement's index; `replace` increments before using it, # so this starts one below the first free index (@workers). @next_worker_index = @workers - 1 @@ -64,6 +64,12 @@ def join @lock.synchronize { @supervisor_threads.dup }.each(&:join) end + # Ractors cannot be force-killed; their clients were already freed by + # abort_all_inflight. A stuck ractor leaks until process exit. + def kill_stragglers + Log.error("shutdown deadline passed with stuck ractor workers") unless done? + end + # Workers alive and serving: the live ones minus those the quarantine # monitor abandoned as wedged (their replacements count instead). def active_count @@ -187,11 +193,10 @@ def supervise(index) end end - # Fresh ractor on fresh or recycled slots. Slots are never reused - # across crash respawns: stale interrupt kicks and dead weak refs go - # down with the old slot. They are reused after a clean retirement. + # Fresh ractor on slots from the bank: fresh ones, or ones a retired + # worker handed back. def spawn_worker(worker_index) - worker_ids = Array.new(@threads) { claim_slot } + worker_ids = Array.new(@threads) { @bank.claim } @lock.synchronize do @worker_slots[worker_index] = worker_ids worker_ids.each { |id| @slot_to_worker[id] = worker_index } @@ -213,26 +218,17 @@ def spawn_worker(worker_index) [ractor, worker_ids] end - # A slot a retired worker gave back, reset for its new occupant, or a - # fresh one. - def claim_slot - id = @lock.synchronize { @free_slots.pop } - return Native.register_worker(@server_id) unless id - - Native.reset_slot(@server_id, id) - id - end - # Bookkeeping for a supervisor thread that is done: the worker is no - # longer live, a retired worker's slots go back to the free list, and - # the thread leaves the join set so a long-lived elastic pool does not + # longer live, a retired worker's slots go back to the bank, and the + # thread leaves the join set so a long-lived elastic pool does not # accumulate dead threads. def exited(index, worker_ids) - @lock.synchronize do + retired = @lock.synchronize do @live.delete(index) - @free_slots.concat(worker_ids) if @retiring.delete(index) && worker_ids @supervisor_threads.delete(Thread.current) + @retiring.delete(index) end + @bank.release(worker_ids) if retired && worker_ids report_active end diff --git a/lib/kino/server.rb b/lib/kino/server.rb index 4e1a454..81a0b80 100644 --- a/lib/kino/server.rb +++ b/lib/kino/server.rb @@ -480,13 +480,7 @@ def workers_done? end def kill_stragglers - if @supervisor - # Ractors cannot be force-killed; their clients were already freed - # by abort_all_inflight. The stuck ractor leaks until process exit. - Log.error("shutdown deadline passed with stuck ractor workers") unless @supervisor.done? - else - @threaded_pool.kill_stragglers - end + pool.kill_stragglers end # Policy (mode resolution): when is an app safe for ractor diff --git a/lib/kino/slot_bank.rb b/lib/kino/slot_bank.rb new file mode 100644 index 0000000..d8855b6 --- /dev/null +++ b/lib/kino/slot_bank.rb @@ -0,0 +1,31 @@ +# frozen_string_literal: true + +module Kino + # @private + # Dispatch slots for a worker pool: fresh ones from the native registry, + # or ones that retired workers handed back, reset for their next + # occupant. The native side never removes a slot, so recycling is what + # keeps the slot table from growing as an elastic pool breathes. Only + # cleanly exited workers return slots; a crashed worker's slots are + # abandoned (stale interrupt kicks and dead weak refs go down with + # them). + class SlotBank + def initialize(server_id) + @server_id = server_id + @free = [] + @lock = Mutex.new + end + + def claim + id = @lock.synchronize { @free.pop } + return Native.register_worker(@server_id) unless id + + Native.reset_slot(@server_id, id) + id + end + + def release(ids) + @lock.synchronize { @free.concat(ids) } + end + end +end diff --git a/lib/kino/threaded_pool.rb b/lib/kino/threaded_pool.rb index 30bf401..46bfc31 100644 --- a/lib/kino/threaded_pool.rb +++ b/lib/kino/threaded_pool.rb @@ -17,14 +17,15 @@ def initialize(server_id, app, threads:, batch: 1, hooks: nil, on_worker_exit: n @hooks = hooks @on_worker_exit = on_worker_exit @lock = Mutex.new - # index => {slots:, threads:}; the groups asked to leave; the slot - # ids retired groups gave back; quarantine replacements (one thread - # each, standing in for a wedged slot: neither counted nor retired, - # so the wedged group keeps counting as the capacity it still is). + # index => {slots:, threads:}; the groups asked to leave; quarantine + # replacements (one thread each, standing in for a wedged slot: + # neither counted nor retired, so the wedged group keeps counting as + # the capacity it still is). Retired groups hand their slots back to + # the bank. @groups = {} @slot_to_group = {} @retiring = {} - @free_slots = [] + @bank = SlotBank.new(server_id) @replacements = {} @wedged = {} @next_index = -1 @@ -127,7 +128,7 @@ def spawn_group(slots: @threads, replacement: false) @next_index end slots.times do - id = claim_slot + id = @bank.claim @lock.synchronize do group[:slots] << id @slot_to_group[id] = index @@ -154,31 +155,23 @@ def spawn_thread(worker_id) end end - # A slot a retired group gave back, reset for its new occupant, or a - # fresh one. - def claim_slot - id = @lock.synchronize { @free_slots.pop } - return Native.register_worker(@server_id) unless id - - Native.reset_slot(@server_id, id) - id - end - # Retiring groups whose threads have all exited give their slots back # and leave the table. Runs before every pool decision, so the count # the control plane sees never lags by more than a scaler tick. def reap freed = @lock.synchronize do done = @retiring.keys.select { |index| @groups[index][:threads].none?(&:alive?) } - done.each do |index| + done.flat_map do |index| group = @groups.delete(index) @retiring.delete(index) group[:slots].each { |id| @slot_to_group.delete(id) } - @free_slots.concat(group[:slots]) + group[:slots] end - done.any? end - report_active if freed + return if freed.empty? + + @bank.release(freed) + report_active end def report_active diff --git a/spec/monitor_spec.rb b/spec/monitor_spec.rb new file mode 100644 index 0000000..00f89b3 --- /dev/null +++ b/spec/monitor_spec.rb @@ -0,0 +1,48 @@ +# frozen_string_literal: true + +# A monitor that counts its scans and can be told to raise on one. +class CountingMonitor < Kino::Monitor + attr_reader :scans + + def initialize(tick:, raise_on: nil) + super(name: "counting", tick: tick) + @scans = 0 + @raise_on = raise_on + end + + private + + def scan + @scans += 1 + raise "scan #{@scans} failed" if @scans == @raise_on + end +end + +RSpec.describe Kino::Monitor do + def monotonic = Process.clock_gettime(Process::CLOCK_MONOTONIC) + + it "scans every tick on its own thread until stopped" do + monitor = CountingMonitor.new(tick: 0.01).start + deadline = monotonic + 2 + sleep 0.01 until monitor.scans >= 3 || monotonic > deadline + monitor.stop + scans = monitor.scans + + expect(scans).to be >= 3 + sleep 0.05 + expect(monitor.scans).to eq(scans) + end + + it "logs a raising scan and keeps going" do + monitor = CountingMonitor.new(tick: 0.01, raise_on: 2) + err = capture_native_stderr do + monitor.start + deadline = monotonic + 2 + sleep 0.01 until monitor.scans >= 4 || monotonic > deadline + monitor.stop + end + + expect(monitor.scans).to be >= 4 + expect(err).to include("counting tick error: RuntimeError: scan 2 failed") + end +end