diff --git a/server/src/crate_server/readiness.gleam b/server/src/crate_server/readiness.gleam index faefb33..1401cc7 100644 --- a/server/src/crate_server/readiness.gleam +++ b/server/src/crate_server/readiness.gleam @@ -56,13 +56,34 @@ pub type Check { } type State { - State(last: Option(#(Bool, Int))) + State( + last: Option(#(Bool, Int)), + in_flight: Bool, + worker_replies: Subject(Msg), + ) } type Msg { Probe(Subject(Bool)) + ProbeResult(Bool, Int) } +/// Bounds how long a `Probe` message may block the actor waiting on a fresh +/// probe result before it gives up and answers from the cache instead. Kept +/// well under `call_budget_ms` so the actor always has room to reply before +/// its caller's own timeout fires, and well under the docker/Caddy +/// healthcheck timeouts below so a slow-but-healthy probe still answers +/// inside one health check even when it can't finish inside this budget. +const local_wait_ms = 500 + +/// The budget `cached()` gives callers to hear back from the actor. Deliberately +/// above the docker healthcheck's `--timeout=3s` and below Caddy's +/// `health_timeout 3s`-driven poll cadence: the actor itself never blocks +/// longer than `local_wait_ms`, so this budget is headroom for scheduling, +/// not for the probe itself. See deploy/Caddyfile and the Dockerfile +/// HEALTHCHECK. +const call_budget_ms = 4000 + pub fn ready() -> Check { Check(fn() { True }) } @@ -80,15 +101,35 @@ pub fn schema_check_sql() -> String { schema_check } -/// Coalesces calls through one actor, so only the first request after the TTL -/// executes `probe`; every other caller receives the cached result. +/// Coalesces calls through one actor, so only one `probe` runs at a time. +/// The actor never runs `probe` itself: it spawns a worker for it and waits +/// up to `local_wait_ms` for the result, so a slow probe never wedges the +/// actor's mailbox behind it. A probe that doesn't finish in time is +/// "unknown", not "broken": the actor answers with the last known state (or +/// `False` if none exists yet), and the worker's result folds into state +/// through the actor's own selector whenever it lands -- it must go through +/// that selector rather than a bare manual receive, or the actor's built-in +/// catch-all for unrecognised messages discards it the moment the loop goes +/// back to waiting. pub fn cached(probe: fn() -> Bool, now: fn() -> Int) -> Check { let assert Ok(started) = - actor.new(State(last: None)) + actor.new_with_initialiser(1000, fn(subject) { + let worker_replies = process.new_subject() + let selector = + process.new_selector() + |> process.select(subject) + |> process.select(worker_replies) + actor.initialised(State(last: None, in_flight: False, worker_replies:)) + |> actor.selecting(selector) + |> actor.returning(subject) + |> Ok + }) |> actor.on_message(fn(state, msg) { handle(state, msg, probe, now) }) |> actor.start let subject = started.data - Check(fn() { parallel.try_call(subject, 1000, Probe) |> result.unwrap(False) }) + Check(fn() { + parallel.try_call(subject, call_budget_ms, Probe) |> result.unwrap(False) + }) } pub fn is_ready(check: Check) -> Bool { @@ -103,34 +144,69 @@ fn handle( ) -> actor.Next(State, Msg) { case msg { Probe(reply) -> { - let #(available, next) = probe_or_cached(state.last, probe, now) + let #(available, next) = respond(state, probe, now) process.send(reply, available) - actor.continue(State(last: next)) + actor.continue(next) } + ProbeResult(available, checked_at) -> + actor.continue( + State(..state, last: Some(#(available, checked_at)), in_flight: False), + ) } } -fn probe_or_cached( - last: Option(#(Bool, Int)), +fn respond( + state: State, probe: fn() -> Bool, now: fn() -> Int, -) -> #(Bool, Option(#(Bool, Int))) { - case last { - Some(#(available, checked_at)) -> - case now() - checked_at < cache_ttl_seconds { - True -> #(available, last) - False -> run_probe(probe, now) +) -> #(Bool, State) { + case state.in_flight { + True -> #(fallback(state.last), state) + False -> + case needs_refresh(state.last, now) { + False -> #(fallback(state.last), state) + True -> refresh(state, probe, now) } - None -> run_probe(probe, now) } } -fn run_probe( +fn needs_refresh(last: Option(#(Bool, Int)), now: fn() -> Int) -> Bool { + case last { + None -> True + Some(#(_, checked_at)) -> now() - checked_at >= cache_ttl_seconds + } +} + +/// The best answer available without waiting on the database: the last +/// known state, or "not ready" if the database has never once answered. +fn fallback(last: Option(#(Bool, Int))) -> Bool { + case last { + Some(#(available, _)) -> available + None -> False + } +} + +fn refresh( + state: State, probe: fn() -> Bool, now: fn() -> Int, -) -> #(Bool, Option(#(Bool, Int))) { - let available = probe() - #(available, Some(#(available, now()))) +) -> #(Bool, State) { + let replies = state.worker_replies + process.spawn_unlinked(fn() { + process.send(replies, ProbeResult(probe(), now())) + }) + // Only ProbeResult ever lands on `replies`; the Probe arm below is + // unreachable but keeps this match total without a panic. + case process.receive(state.worker_replies, local_wait_ms) { + Ok(ProbeResult(available, checked_at)) -> #( + available, + State(..state, last: Some(#(available, checked_at)), in_flight: False), + ) + Ok(Probe(_)) | Error(Nil) -> #( + fallback(state.last), + State(..state, in_flight: True), + ) + } } fn schema_is_ready(conn: pog.Connection) -> Bool { diff --git a/server/test/health_test.gleam b/server/test/health_test.gleam index 2c38b0c..8ce3999 100644 --- a/server/test/health_test.gleam +++ b/server/test/health_test.gleam @@ -100,6 +100,84 @@ pub fn slow_failed_probe_is_cached_from_completion_test() { assert support.drain_count(calls) == 1 } +pub fn probe_timeout_before_any_success_reports_not_ready_test() { + let calls = process.new_subject() + let check = + readiness.cached( + fn() { + process.send(calls, Nil) + process.sleep(700) + True + }, + fn() { 0 }, + ) + + assert !readiness.is_ready(check) + + process.sleep(900) + + assert readiness.is_ready(check) + assert support.drain_count(calls) == 1 +} + +pub fn probe_timeout_falls_back_to_last_known_state_test() { + let seen = counter() + let #(_advance, now) = test_clock([0, 10, 15, 16, 20]) + let check = + readiness.cached( + fn() { + case seen() { + 0 -> True + _ -> { + process.sleep(700) + False + } + } + }, + now, + ) + + // First probe is fast: it establishes a known-good state. + assert readiness.is_ready(check) + + // The TTL has lapsed, so this call starts a refresh, but the new probe is + // slow. A timeout must fall back to the last known state, not go straight + // to "not ready". + assert readiness.is_ready(check) + + // Once the slow probe actually completes and fails, that becomes the + // reported state. + process.sleep(900) + assert !readiness.is_ready(check) +} + +type CountState { + CountState(n: Int) +} + +type CountMsg { + Bump(process.Subject(Int)) +} + +/// Returns how many times it has been called so far, starting at 0, via one +/// actor: readiness probes run on a spawned worker, so a plain closed-over +/// variable would race. +fn counter() -> fn() -> Int { + let assert Ok(started) = + actor.new(CountState(n: 0)) + |> actor.on_message(fn(state, msg) { + case msg { + Bump(reply) -> { + process.send(reply, state.n) + actor.continue(CountState(n: state.n + 1)) + } + } + }) + |> actor.start + let subject = started.data + fn() { parallel.try_call(subject, 1000, Bump) |> result.unwrap(-1) } +} + fn test_clock(ticks: List(Int)) -> #(fn() -> Nil, fn() -> Int) { let assert Ok(started) = actor.new(ClockState(ticks:))