From ffeb2eb364d02750847c69ab65dba8c71d9c09b2 Mon Sep 17 00:00:00 2001 From: Niels Mokkenstorm Date: Sun, 9 Aug 2026 23:53:01 +0200 Subject: [PATCH] fix: bound cached readiness staleness, recover from a wedged or crashed probe, and derive the schema check --- server/src/crate_server/readiness.gleam | 203 ++++++++++++++++-------- server/test/health_test.gleam | 135 ++++++++-------- 2 files changed, 207 insertions(+), 131 deletions(-) diff --git a/server/src/crate_server/readiness.gleam b/server/src/crate_server/readiness.gleam index 1401cc7..b8e9c61 100644 --- a/server/src/crate_server/readiness.gleam +++ b/server/src/crate_server/readiness.gleam @@ -2,54 +2,60 @@ //// 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 -const schema_check = "with required_columns(table_name, column_name, type_name) as (values -('oauth_sessions','id','text'),('oauth_sessions','data','text'),('oauth_sessions','created_at','bigint'), -('discogs_creds','id','text'),('discogs_creds','data','text'),('discogs_creds','created_at','bigint'), -('known_users','did','text'),('known_users','handle','text'),('known_users','pds','text'),('known_users','first_seen','bigint'), -('catalog_releases','uri','text'),('catalog_releases','cid','text'),('catalog_releases','title','text'),('catalog_releases','artist_display','text'),('catalog_releases','genres','text[]'),('catalog_releases','styles','text[]'),('catalog_releases','released','text'),('catalog_releases','country','text'),('catalog_releases','cover_cid','text'),('catalog_releases','thumb_url','text'),('catalog_releases','discogs_id','text'),('catalog_releases','created_at','text'),('catalog_releases','publisher_did','text'),('catalog_releases','publisher_handle','text'),('catalog_releases','publisher_pds','text'),('catalog_releases','supersedes','text'),('catalog_releases','based_on','text'),('catalog_releases','format','text'),('catalog_releases','label','text'),('catalog_releases','master','text'),('catalog_releases','barcode','text'), -('catalog_adoptions','entry_uri','text'),('catalog_adoptions','did','text'),('catalog_adoptions','release_uri','text'),('catalog_adoptions','status','text'),('catalog_adoptions','created_at','text'),('catalog_adoptions','source','text'), -('catalog_edits','edit_uri','text'),('catalog_edits','did','text'),('catalog_edits','subject_uri','text'),('catalog_edits','entity','text'),('catalog_edits','created_at','text'),('catalog_edits','cid','text'),('catalog_edits','fields','text'),('catalog_edits','rationale','text'), -('jetstream_cursor','id','text'),('jetstream_cursor','time_us','bigint'), -('shelf_entry_events','event_uri','text'),('shelf_entry_events','entry_uri','text'),('shelf_entry_events','did','text'),('shelf_entry_events','rkey','text'),('shelf_entry_events','action','text'),('shelf_entry_events','created_at','text'),('shelf_entry_events','record_json','text'),('shelf_entry_events','indexed_at','text'), -('shelf_entries','entry_uri','text'),('shelf_entries','did','text'),('shelf_entries','status','text'),('shelf_entries','title','text'),('shelf_entries','artist_display','text'),('shelf_entries','year','integer'),('shelf_entries','format','text'),('shelf_entries','thumb_url','text'),('shelf_entries','cover_cid','text'),('shelf_entries','media_grade','text'),('shelf_entries','sleeve_grade','text'),('shelf_entries','rating','integer'),('shelf_entries','folder','text'),('shelf_entries','notes','text'),('shelf_entries','release_uri','text'),('shelf_entries','release_cid','text'),('shelf_entries','price_amount','integer'),('shelf_entries','price_currency','text'),('shelf_entries','counterparty','text'),('shelf_entries','source_client_agent','text'),('shelf_entries','source_origin','text'),('shelf_entries','source_origin_url','text'),('shelf_entries','source_external_provider','text'),('shelf_entries','source_external_id','text'),('shelf_entries','source_external_url','text'),('shelf_entries','source_record_uri','text'),('shelf_entries','source_record_cid','text'),('shelf_entries','created_at','text'),('shelf_entries','updated_at','text'),('shelf_entries','indexed_at','text'), -('shelf_index_seen','did','text'),('shelf_index_seen','seen_at','text'), -('follow_edges','follow_uri','text'),('follow_edges','did','text'),('follow_edges','subject','text'),('follow_edges','created_at','text'), -('follow_index_seen','did','text'),('follow_index_seen','seen_at','text') -), required_not_null(table_name, column_name) as (values -('oauth_sessions','id'),('oauth_sessions','data'),('oauth_sessions','created_at'),('discogs_creds','id'),('discogs_creds','data'),('discogs_creds','created_at'),('known_users','did'),('known_users','handle'),('known_users','pds'),('known_users','first_seen'), -('catalog_releases','uri'),('catalog_releases','cid'),('catalog_releases','title'),('catalog_releases','genres'),('catalog_releases','styles'),('catalog_releases','created_at'),('catalog_releases','publisher_did'),('catalog_releases','publisher_handle'),('catalog_releases','publisher_pds'), -('catalog_adoptions','entry_uri'),('catalog_adoptions','did'),('catalog_adoptions','release_uri'),('catalog_adoptions','status'),('catalog_adoptions','created_at'),('catalog_edits','edit_uri'),('catalog_edits','did'),('catalog_edits','subject_uri'),('catalog_edits','entity'),('catalog_edits','created_at'),('jetstream_cursor','id'),('jetstream_cursor','time_us'), -('shelf_entry_events','event_uri'),('shelf_entry_events','entry_uri'),('shelf_entry_events','did'),('shelf_entry_events','rkey'),('shelf_entry_events','action'),('shelf_entry_events','created_at'),('shelf_entry_events','record_json'),('shelf_entry_events','indexed_at'),('shelf_entries','entry_uri'),('shelf_entries','did'),('shelf_entries','status'),('shelf_entries','created_at'),('shelf_entries','updated_at'),('shelf_entries','indexed_at'),('shelf_index_seen','did'),('shelf_index_seen','seen_at'),('follow_edges','follow_uri'),('follow_edges','did'),('follow_edges','subject'),('follow_edges','created_at'),('follow_index_seen','did'),('follow_index_seen','seen_at') -), required_conflicts(table_name, column_name) as (values -('oauth_sessions','id'),('discogs_creds','id'),('known_users','did'),('catalog_releases','uri'),('catalog_adoptions','entry_uri'),('catalog_edits','edit_uri'),('jetstream_cursor','id'),('shelf_entry_events','event_uri'),('shelf_entries','entry_uri'),('shelf_index_seen','did'),('follow_edges','follow_uri'),('follow_index_seen','did') -), required_privileges(table_name, privilege) as (values -('oauth_sessions','SELECT'),('oauth_sessions','INSERT'),('oauth_sessions','UPDATE'),('oauth_sessions','DELETE'),('discogs_creds','SELECT'),('discogs_creds','INSERT'),('discogs_creds','UPDATE'),('discogs_creds','DELETE'),('known_users','SELECT'),('known_users','INSERT'),('known_users','UPDATE'), -('catalog_releases','SELECT'),('catalog_releases','INSERT'),('catalog_releases','UPDATE'),('catalog_releases','DELETE'),('catalog_adoptions','SELECT'),('catalog_adoptions','INSERT'),('catalog_adoptions','UPDATE'),('catalog_adoptions','DELETE'),('catalog_edits','SELECT'),('catalog_edits','INSERT'),('catalog_edits','UPDATE'),('catalog_edits','DELETE'),('jetstream_cursor','SELECT'),('jetstream_cursor','INSERT'),('jetstream_cursor','UPDATE'), -('shelf_entry_events','SELECT'),('shelf_entry_events','INSERT'),('shelf_entry_events','UPDATE'),('shelf_entry_events','DELETE'),('shelf_entries','SELECT'),('shelf_entries','INSERT'),('shelf_entries','UPDATE'),('shelf_entries','DELETE'),('shelf_index_seen','SELECT'),('shelf_index_seen','INSERT'),('shelf_index_seen','DELETE'), -('follow_edges','SELECT'),('follow_edges','INSERT'),('follow_edges','UPDATE'),('follow_edges','DELETE'),('follow_index_seen','SELECT'),('follow_index_seen','INSERT'),('follow_index_seen','DELETE') -), actual_columns as ( -select class.relname, attribute.attname, attribute.attnotnull, pg_catalog.format_type(attribute.atttypid, attribute.atttypmod) -from pg_catalog.pg_attribute attribute join pg_catalog.pg_class class on class.oid = attribute.attrelid join pg_catalog.pg_namespace namespace on namespace.oid = class.relnamespace -where namespace.nspname = current_schema() and attribute.attnum > 0 and not attribute.attisdropped -), actual_conflicts as ( -select class.relname, attribute.attname from pg_catalog.pg_index idx join pg_catalog.pg_class class on class.oid = idx.indrelid join pg_catalog.pg_namespace namespace on namespace.oid = class.relnamespace join pg_catalog.pg_attribute attribute on attribute.attrelid = class.oid and attribute.attnum = any(idx.indkey) -where namespace.nspname = current_schema() and idx.indisunique and idx.indisvalid and idx.indisready and idx.indpred is null and idx.indnkeyatts = 1 -) -select 1 where not exists (select 1 from required_columns required left join actual_columns actual on actual.relname = required.table_name and actual.attname = required.column_name where actual.attname is null or actual.format_type <> required.type_name) -and not exists (select 1 from required_not_null required left join actual_columns actual on actual.relname = required.table_name and actual.attname = required.column_name where actual.attname is null or not actual.attnotnull) -and not exists (select 1 from required_conflicts required left join actual_conflicts actual on actual.relname = required.table_name and actual.attname = required.column_name where actual.attname is null) -and not exists (select 1 from required_privileges where not has_table_privilege(current_user, table_name, privilege))" +/// 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) @@ -58,14 +64,19 @@ pub type Check { type State { State( last: Option(#(Bool, Int)), - in_flight: Bool, + // 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(Bool, Int) + ProbeResult(Int, Bool) } /// Bounds how long a `Probe` message may block the actor waiting on a fresh @@ -84,6 +95,24 @@ const local_wait_ms = 500 /// 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 }) } @@ -96,21 +125,21 @@ pub fn postgres(conn: pog.Connection) -> Check { cached(fn() { schema_is_ready(conn) }, clock.now_seconds) } -/// Returns the representative appview schema query for regression tests. -pub fn schema_check_sql() -> String { - schema_check -} - -/// 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 +/// 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), and the worker's result folds into state +/// `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. +/// 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) { @@ -119,7 +148,12 @@ pub fn cached(probe: fn() -> Bool, now: fn() -> Int) -> Check { process.new_selector() |> process.select(subject) |> process.select(worker_replies) - actor.initialised(State(last: None, in_flight: False, worker_replies:)) + actor.initialised(State( + last: None, + in_flight: None, + next_generation: 0, + worker_replies:, + )) |> actor.selecting(selector) |> actor.returning(subject) |> Ok @@ -148,10 +182,24 @@ fn handle( process.send(reply, available) actor.continue(next) } - ProbeResult(available, checked_at) -> - actor.continue( - State(..state, last: Some(#(available, checked_at)), in_flight: False), - ) + 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 } } @@ -161,10 +209,14 @@ fn respond( now: fn() -> Int, ) -> #(Bool, State) { case state.in_flight { - True -> #(fallback(state.last), state) - False -> + 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), state) + False -> #(fallback(state.last, now), state) True -> refresh(state, probe, now) } } @@ -178,10 +230,15 @@ fn needs_refresh(last: Option(#(Bool, Int)), now: fn() -> Int) -> Bool { } /// 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 { +/// 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, _)) -> available + Some(#(available, checked_at)) -> + case now() - checked_at > max_staleness_seconds { + True -> False + False -> available + } None -> False } } @@ -191,27 +248,35 @@ fn refresh( 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(probe(), now())) + 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(available, checked_at)) -> #( + Ok(ProbeResult(reply_generation, available)) + if reply_generation == generation + -> #( available, - State(..state, last: Some(#(available, checked_at)), in_flight: False), - ) - Ok(Probe(_)) | Error(Nil) -> #( - fallback(state.last), - State(..state, in_flight: True), + 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) |> pog.returning(row) |> pog.execute(conn) { + case + pog.query(schema_check_sql()) |> pog.returning(row) |> pog.execute(conn) + { Ok(pog.Returned(rows: [_, ..], ..)) -> True _ -> False } diff --git a/server/test/health_test.gleam b/server/test/health_test.gleam index 8ce3999..a387950 100644 --- a/server/test/health_test.gleam +++ b/server/test/health_test.gleam @@ -221,84 +221,95 @@ fn next_tick(ticks: List(Int)) -> #(Int, List(Int)) { } } -pub fn schema_probe_covers_required_table_columns_test() { - let sql = readiness.schema_check_sql() +pub fn required_tables_cover_every_store_test() { + let names = readiness.required_tables() [ - "required_columns(table_name, column_name, type_name)", - "pg_catalog.format_type(attribute.atttypid, attribute.atttypmod)", - "actual.format_type <> required.type_name", - "('oauth_sessions','created_at','bigint')", - "('discogs_creds','created_at','bigint')", - "('known_users','first_seen','bigint')", - "('catalog_releases','genres','text[]')", - "('catalog_releases','styles','text[]')", - "('catalog_releases','barcode','text')", - "('catalog_adoptions','source','text')", - "('catalog_edits','rationale','text')", - "('jetstream_cursor','time_us','bigint')", - "('shelf_entry_events','record_json','text')", - "('shelf_entries','source_record_cid','text')", - "('shelf_entries','price_amount','integer')", - "('shelf_entries','rating','integer')", - "('follow_edges','created_at','text')", - "('follow_index_seen','seen_at','text')", + "oauth_sessions", "discogs_creds", "known_users", "catalog_releases", + "catalog_adoptions", "catalog_edits", "jetstream_cursor", + "shelf_entry_events", "shelf_entries", "shelf_index_seen", "follow_edges", + "follow_index_seen", ] - |> list.each(fn(shape) { - assert string.contains(sql, shape) + |> list.each(fn(name) { + assert list.contains(names, name) }) + + assert list.length(names) == 12 } -pub fn schema_probe_requires_store_operation_privileges_test() { +pub fn schema_probe_requires_every_required_table_test() { let sql = readiness.schema_check_sql() - [ - "required_privileges(table_name, privilege)", - "has_table_privilege(current_user, table_name, privilege)", - "('oauth_sessions','DELETE')", - "('discogs_creds','UPDATE')", - "('known_users','INSERT')", - "('catalog_releases','DELETE')", - "('catalog_adoptions','UPDATE')", - "('catalog_edits','INSERT')", - "('jetstream_cursor','UPDATE')", - "('shelf_entry_events','DELETE')", - "('shelf_entries','UPDATE')", - "('shelf_index_seen','INSERT')", - "('follow_edges','DELETE')", - "('follow_index_seen','SELECT')", - ] - |> list.each(fn(clause) { - assert string.contains(sql, clause) + readiness.required_tables() + |> list.each(fn(name) { + assert string.contains(sql, "'" <> name <> "'") }) } -pub fn schema_probe_requires_not_null_and_conflict_targets_test() { +pub fn schema_probe_requires_existence_and_operation_privileges_test() { let sql = readiness.schema_check_sql() [ - "required_not_null(table_name, column_name)", - "actual.attnotnull", - "required_conflicts(table_name, column_name)", - "idx.indisunique", - "idx.indisvalid", - "idx.indisready", - "idx.indpred is null", - "idx.indnkeyatts = 1", - "('catalog_releases','genres')", - "('catalog_adoptions','entry_uri')", - "('catalog_edits','edit_uri')", - "('shelf_entry_events','record_json')", - "('shelf_entries','updated_at')", - "('oauth_sessions','id')", - "('discogs_creds','id')", - "('known_users','did')", - "('jetstream_cursor','id')", - "('shelf_index_seen','did')", - "('follow_edges','follow_uri')", - "('follow_index_seen','did')", + "to_regclass(current_schema() || '.' || required.name) is null", + "has_table_privilege(current_user, current_schema() || '.' || required.name, 'SELECT')", + "has_table_privilege(current_user, current_schema() || '.' || required.name, 'INSERT')", + "has_table_privilege(current_user, current_schema() || '.' || required.name, 'UPDATE')", + "has_table_privilege(current_user, current_schema() || '.' || required.name, 'DELETE')", ] |> list.each(fn(clause) { assert string.contains(sql, clause) }) } + +pub fn a_wedged_probe_eventually_reports_not_ready_test() { + let seen = counter() + let #(_advance, now) = test_clock([0, 0, 10, 10, 15, 25, 25, 35]) + let check = + readiness.cached( + fn() { + case seen() { + 0 -> True + _ -> { + process.sleep_forever() + True + } + } + }, + now, + ) + + // Establishes a known-good state. + assert readiness.is_ready(check) + + // The TTL lapses and the refresh hangs (the database never returns). The + // last known-good state is still within its staleness ceiling, so it is + // still trusted while the probe is stuck. + assert readiness.is_ready(check) + + // Once that cached state has been trusted longer than its staleness + // ceiling allows, a database that hangs rather than errors must stop + // advertising as ready, even though the original probe never finished. + assert !readiness.is_ready(check) +} + +pub fn a_crashing_worker_does_not_permanently_disable_probing_test() { + let seen = counter() + let #(_advance, now) = test_clock([0, 12, 12, 24]) + let check = + readiness.cached( + fn() { + case seen() { + 0 -> panic as "the database connection was severed mid-query" + _ -> True + } + }, + now, + ) + + // The very first probe's worker crashes before ever answering. + assert !readiness.is_ready(check) + + // Probing must recover once the crashed attempt has plainly overrun, + // rather than being wedged forever by a worker that died without a reply. + assert readiness.is_ready(check) +} -- 2.51.2