//// Readiness is intentionally a capability rather than an environment flag. //// One actor coalesces concurrent probe requests and caches the result for a //// short interval, so a public endpoint cannot amplify database work. import crate_server/catalog_index_postgres import crate_server/clock import crate_server/follow_index_postgres import crate_server/known_users_postgres import crate_server/oauth/sessions_postgres import crate_server/parallel import crate_server/shelf_index_postgres import gleam/dynamic/decode import gleam/erlang/process.{type Subject} import gleam/list import gleam/option.{type Option, None, Some} import gleam/otp/actor import gleam/result import gleam/string import pog const cache_ttl_seconds = 5 /// The tables every store's own `*_postgres` module says it needs (see each /// module's `required_tables`). This is the one place that list is /// assembled from every store, so it stays exhaustive; each module is the /// one place its own table names live, so adding a column to one of their /// `create table`/`alter table` statements never touches this file. pub fn required_tables() -> List(String) { known_users_postgres.required_tables() |> list.append(catalog_index_postgres.required_tables()) |> list.append(shelf_index_postgres.required_tables()) |> list.append(follow_index_postgres.required_tables()) |> list.append(sessions_postgres.required_tables()) } /// Returns the generated readiness query for regression tests. Deliberately /// narrower than a full column/type/not-null/unique-index audit: it checks /// that every required table exists and that the connected role holds /// baseline CRUD privileges on each, which is what actually gates the app /// from working and is cheap enough to run inside `local_wait_ms`. Table /// names come from `required_tables()`, are internal constants (never user /// input), so interpolating them is not injection -- the same reasoning /// `oauth/sessions_postgres` already relies on for its table-name /// interpolation. pub fn schema_check_sql() -> String { let names = required_tables() |> list.map(fn(name) { "'" <> name <> "'" }) |> string.join(", ") "select 1 where not exists ( select 1 from unnest(array[" <> names <> "]) as required(name) where to_regclass(current_schema() || '.' || required.name) is null or not has_table_privilege(current_user, current_schema() || '.' || required.name, 'SELECT') or not has_table_privilege(current_user, current_schema() || '.' || required.name, 'INSERT') or not has_table_privilege(current_user, current_schema() || '.' || required.name, 'UPDATE') or not has_table_privilege(current_user, current_schema() || '.' || required.name, 'DELETE') )" } pub type Check { Check(run: fn() -> Bool) } type State { State( last: Option(#(Bool, Int)), // The generation and start time of the refresh in flight, if any. The // generation lets a reply be matched to the attempt that asked for it, // so a late answer from an attempt we have since given up on cannot be // mistaken for the current one's. in_flight: Option(#(Int, Int)), next_generation: Int, worker_replies: Subject(Msg), ) } type Msg { Probe(Subject(Bool)) ProbeResult(Int, Bool) } /// 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 /// How long a stuck refresh is tolerated before a fresh one may start /// alongside it. Tied to Caddy's `health_interval 10s` (deploy/Caddyfile): /// once a probe has run longer than one full health-poll cycle without /// answering, waiting on it any further buys nothing, and it may be a /// worker that crashed without ever sending a reply rather than one that is /// merely slow -- either way, a probe that plainly overran must not disable /// all future probing for the actor's lifetime. const refresh_overrun_seconds = 10 /// How long a cached answer may be trusted before it must degrade to "not /// ready", regardless of what it says. Tied to the docker healthcheck's /// `interval: 30s` (docker-compose.yml), three Caddy poll cycles. Without /// this ceiling, a database that hangs rather than errors (so `in_flight` /// never clears) would have every later `Probe` fall back to the last /// known-good state forever: the exact inverse of the timeout-vs-failure /// bug this file exists to fix, and one that fails silently. const max_staleness_seconds = 30 pub fn ready() -> Check { Check(fn() { True }) } pub fn unavailable() -> Check { Check(fn() { False }) } pub fn postgres(conn: pog.Connection) -> Check { cached(fn() { schema_is_ready(conn) }, clock.now_seconds) } /// Coalesces calls through one actor, so normally 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, or if that state has gone stale past /// `max_staleness_seconds`), 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. If a refresh overruns `refresh_overrun_seconds` -- a /// hung query, or a worker that crashed without replying at all -- a fresh /// one is allowed to start rather than leaving probing wedged for good; the /// generation on each reply keeps a stale answer from either attempt being /// applied to the wrong one. pub fn cached(probe: fn() -> Bool, now: fn() -> Int) -> Check { let assert Ok(started) = 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: None, next_generation: 0, 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, call_budget_ms, Probe) |> result.unwrap(False) }) } pub fn is_ready(check: Check) -> Bool { check.run() } fn handle( state: State, msg: Msg, probe: fn() -> Bool, now: fn() -> Int, ) -> actor.Next(State, Msg) { case msg { Probe(reply) -> { let #(available, next) = respond(state, probe, now) process.send(reply, available) actor.continue(next) } ProbeResult(generation, available) -> actor.continue(accept(state, generation, available, now)) } } /// Folds a worker's answer into state only if it belongs to the refresh /// currently in flight; a reply from an attempt already abandoned to /// overrun is discarded rather than overwriting a newer one. fn accept( state: State, generation: Int, available: Bool, now: fn() -> Int, ) -> State { case state.in_flight { Some(#(current, _)) if current == generation -> State(..state, last: Some(#(available, now())), in_flight: None) _ -> state } } fn respond( state: State, probe: fn() -> Bool, now: fn() -> Int, ) -> #(Bool, State) { case state.in_flight { Some(#(_, started_at)) -> case now() - started_at > refresh_overrun_seconds { True -> refresh(state, probe, now) False -> #(fallback(state.last, now), state) } None -> case needs_refresh(state.last, now) { False -> #(fallback(state.last, now), state) True -> refresh(state, probe, now) } } } 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 or /// that state is older than `max_staleness_seconds` can vouch for. fn fallback(last: Option(#(Bool, Int)), now: fn() -> Int) -> Bool { case last { Some(#(available, checked_at)) -> case now() - checked_at > max_staleness_seconds { True -> False False -> available } None -> False } } fn refresh( state: State, probe: fn() -> Bool, now: fn() -> Int, ) -> #(Bool, State) { let generation = state.next_generation let replies = state.worker_replies process.spawn_unlinked(fn() { process.send(replies, ProbeResult(generation, probe())) }) let started = State( ..state, in_flight: Some(#(generation, now())), next_generation: generation + 1, ) // 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(reply_generation, available)) if reply_generation == generation -> #( available, State(..started, last: Some(#(available, now())), in_flight: None), ) Ok(_) | Error(Nil) -> #(fallback(state.last, now), started) } } fn schema_is_ready(conn: pog.Connection) -> Bool { let row = decode.success(Nil) case pog.query(schema_check_sql()) |> pog.returning(row) |> pog.execute(conn) { Ok(pog.Returned(rows: [_, ..], ..)) -> True _ -> False } }