Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
52 changes: 52 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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`)
Expand Down
10 changes: 10 additions & 0 deletions doc/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
109 changes: 103 additions & 6 deletions ext/kino/src/control.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -46,6 +48,7 @@ pub fn collect_worker_status(server: &ServerInner) -> Vec<WorkerStat> {
now.saturating_sub(started)
},
quarantined,
retired: slot.retired.load(Ordering::Relaxed),
}
})
.collect()
Expand All @@ -70,6 +73,10 @@ pub struct StatsSnapshot {
pub worker_status: Vec<WorkerStat>,
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,
}

Expand All @@ -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(),
}
}
Expand All @@ -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 {
Expand All @@ -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");
}
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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,
Expand All @@ -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":""#,
] {
Expand All @@ -690,6 +738,49 @@ 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 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<_>>(),
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);
Expand Down Expand Up @@ -757,17 +848,19 @@ mod tests {
in_flight: 1,
busy_ms: 4,
quarantined: false,
retired: false,
},
WorkerStat {
index: 1,
served: 7,
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]
Expand All @@ -786,13 +879,15 @@ mod tests {
in_flight: 1,
busy_ms: 4,
quarantined: false,
retired: false,
},
WorkerStat {
index: 1,
served: 7,
in_flight: 0,
busy_ms: 0,
quarantined: false,
retired: true,
},
];
let text = metrics_text(&s);
Expand Down Expand Up @@ -835,13 +930,15 @@ mod tests {
in_flight: 1,
busy_ms: 0,
quarantined: true,
retired: false,
},
WorkerStat {
index: 1,
served: 9,
in_flight: 1,
busy_ms: 5,
quarantined: false,
retired: false,
},
];
let json = stats_json(&s);
Expand All @@ -850,7 +947,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"#));
Expand Down
9 changes: 9 additions & 0 deletions ext/kino/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
Loading
Loading