From b89c369620a4fec37860e1afbe9da87835ff4cf8 Mon Sep 17 00:00:00 2001 From: Niels Mokkenstorm Date: Sun, 9 Aug 2026 09:35:16 +0200 Subject: [PATCH] feat: gate traffic on appview readiness --- Dockerfile | 3 +- deploy/AUTOPUSH.md | 9 +- deploy/Caddyfile | 9 +- deploy/crate-pull.service | 3 +- deploy/docker-compose.yml | 2 +- deploy/publish.sh | 4 +- deploy/setup-autopull.sh | 7 +- docker-compose.yml | 2 +- server/src/crate_server.gleam | 2 + server/src/crate_server/config_env.gleam | 28 ++- server/src/crate_server/context.gleam | 2 + server/src/crate_server/handlers/health.gleam | 16 ++ server/src/crate_server/readiness.gleam | 142 +++++++++++ server/src/crate_server/router.gleam | 4 +- server/src/crate_server/wiring.gleam | 3 + server/test/health_test.gleam | 226 ++++++++++++++++++ server/test/support.gleam | 2 + 17 files changed, 450 insertions(+), 14 deletions(-) create mode 100644 server/src/crate_server/handlers/health.gleam create mode 100644 server/src/crate_server/readiness.gleam create mode 100644 server/test/health_test.gleam diff --git a/Dockerfile b/Dockerfile index b1b831b..c734942 100644 --- a/Dockerfile +++ b/Dockerfile @@ -37,6 +37,7 @@ RUN apt-get update \ && rm -rf /var/lib/apt/lists/* COPY --from=build /build/server/build/erlang-shipment ./ +COPY deploy/Caddyfile /app/deploy/Caddyfile RUN useradd --system --no-create-home --uid 10001 app \ && chown -R app /app @@ -44,6 +45,6 @@ USER app EXPOSE 8080 HEALTHCHECK --interval=30s --timeout=3s --start-period=10s --retries=3 \ - CMD curl -fsS "http://localhost:${PORT:-8080}/healthz" || exit 1 + CMD curl -fsS "http://localhost:${PORT:-8080}/readyz" || exit 1 # gleam's generated entrypoint puts SPDX headers above the shebang, so exec-form CMD gets ENOEXEC. CMD ["/bin/sh", "./entrypoint.sh", "run"] diff --git a/deploy/AUTOPUSH.md b/deploy/AUTOPUSH.md index 4d52a36..19cba00 100644 --- a/deploy/AUTOPUSH.md +++ b/deploy/AUTOPUSH.md @@ -8,7 +8,9 @@ droplet, `crate-pull.timer` polls once a minute; when the digest behind push, so a red build never becomes an image. Nothing SSHes into production to deploy, and production holds no CI -credentials. `deploy/publish.sh` does the same two steps by hand for +credentials. The pull unit extracts the versioned Caddyfile from the pulled +app image and recreates Caddy only when that file changed, before replacing +the app container. `deploy/publish.sh` does the same two steps by hand for break-glass: build and push from your machine, then start the timer's unit immediately instead of waiting for the next poll. @@ -26,7 +28,10 @@ deploy/setup-autopull.sh ``` Installs `crate-pull.{service,timer}` into `/etc/systemd/system`, enables -the timer, points `.env` at the `latest` tag, and prints the schedule. +the timer, installs the initial Caddyfile, points `.env` at the `latest` tag, +and prints the schedule. Subsequent image pulls extract the versioned +Caddyfile, validates it in the Caddy image, and force-recreates Caddy only +when it changes. Invalid configuration leaves the running proxy untouched. Idempotent. If the ghcr package is private (GitHub creates new packages private, and the diff --git a/deploy/Caddyfile b/deploy/Caddyfile index 0216759..56a63c6 100644 --- a/deploy/Caddyfile +++ b/deploy/Caddyfile @@ -1,3 +1,10 @@ crate.mokkenstorm.dev { - reverse_proxy app:8080 + @readiness path /readyz + respond @readiness 404 + + reverse_proxy app:8080 { + health_uri /readyz + health_interval 10s + health_timeout 3s + } } diff --git a/deploy/crate-pull.service b/deploy/crate-pull.service index 83a7f57..a1462e3 100644 --- a/deploy/crate-pull.service +++ b/deploy/crate-pull.service @@ -6,7 +6,8 @@ After=docker.service [Service] Type=oneshot WorkingDirectory=/opt/crate -ExecStart=/usr/bin/docker compose pull --quiet app +ExecStartPre=/usr/bin/docker compose pull --quiet app +ExecStartPre=/bin/sh -c 'image=$(/usr/bin/docker compose config --images | sed -n "1p"); /usr/bin/docker run --rm --entrypoint cat "$image" /app/deploy/Caddyfile > /tmp/crate.Caddyfile && if ! cmp -s /tmp/crate.Caddyfile /opt/crate/Caddyfile; then /usr/bin/docker run --rm -v /tmp/crate.Caddyfile:/etc/caddy/Caddyfile:ro caddy:2 caddy validate --config /etc/caddy/Caddyfile --adapter caddyfile && install -m 644 /tmp/crate.Caddyfile /opt/crate/Caddyfile && /usr/bin/docker compose up -d --force-recreate caddy; fi' ExecStart=/usr/bin/docker compose up -d app # Each pull leaves the previous :latest dangling; nothing else prunes here. ExecStartPost=-/usr/bin/docker image prune -f diff --git a/deploy/docker-compose.yml b/deploy/docker-compose.yml index f5c2fea..b713a73 100644 --- a/deploy/docker-compose.yml +++ b/deploy/docker-compose.yml @@ -17,7 +17,7 @@ services: db: condition: service_healthy healthcheck: - test: ["CMD", "curl", "-fsS", "http://localhost:8080/healthz"] + test: ["CMD", "curl", "-fsS", "http://localhost:8080/readyz"] interval: 30s timeout: 3s retries: 3 diff --git a/deploy/publish.sh b/deploy/publish.sh index fb55943..945a7d8 100755 --- a/deploy/publish.sh +++ b/deploy/publish.sh @@ -41,8 +41,8 @@ docker buildx build \ echo "rolling the droplet onto $tag" ssh -o ConnectTimeout=20 "$droplet" " set -e - systemctl start at-record-pull.service - cd /opt/at-record && docker compose ps + systemctl start crate-pull.service + cd /opt/crate && docker compose ps " echo "published $tag; verify:" diff --git a/deploy/setup-autopull.sh b/deploy/setup-autopull.sh index 47b4214..ea39ff2 100755 --- a/deploy/setup-autopull.sh +++ b/deploy/setup-autopull.sh @@ -14,14 +14,19 @@ set -euo pipefail droplet="${CRATE_DROPLET:-root@206.189.15.37}" echo "installing the pull timer on $droplet" -tar cf - -C deploy crate-pull.service crate-pull.timer \ +tar cf - -C deploy crate-pull.service crate-pull.timer Caddyfile \ | ssh -o ConnectTimeout=20 "$droplet" " set -e tar xf - -C /etc/systemd/system + install -m 644 /etc/systemd/system/Caddyfile /tmp/crate.Caddyfile + docker run --rm -v /tmp/crate.Caddyfile:/etc/caddy/Caddyfile:ro caddy:2 caddy validate --config /etc/caddy/Caddyfile --adapter caddyfile + install -m 644 /tmp/crate.Caddyfile /opt/crate/Caddyfile + rm /etc/systemd/system/Caddyfile chown root:root /etc/systemd/system/crate-pull.{service,timer} chmod 644 /etc/systemd/system/crate-pull.{service,timer} systemctl daemon-reload systemctl enable --now crate-pull.timer + cd /opt/crate && docker compose up -d --force-recreate caddy " # The timer re-resolves whatever tag the container runs; a pinned sha would diff --git a/docker-compose.yml b/docker-compose.yml index 8562b1b..2b6bcc7 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -26,7 +26,7 @@ services: condition: service_healthy restart: unless-stopped healthcheck: - test: ["CMD", "curl", "-fsS", "http://localhost:8080/healthz"] + test: ["CMD", "curl", "-fsS", "http://localhost:8080/readyz"] interval: 30s timeout: 3s retries: 3 diff --git a/server/src/crate_server.gleam b/server/src/crate_server.gleam index f0498d5..6a17939 100644 --- a/server/src/crate_server.gleam +++ b/server/src/crate_server.gleam @@ -37,6 +37,7 @@ pub fn main() -> Nil { catalog_index, shelf_index, follow_index, + readiness, ) = config_env.stores() let assert Ok(identity_cache) = identity_cache.start( @@ -94,6 +95,7 @@ pub fn main() -> Nil { follow_index:, shelf_index:, identity_cache:, + readiness:, oauth:, ) diff --git a/server/src/crate_server/config_env.gleam b/server/src/crate_server/config_env.gleam index 4d1c1bc..8ed5ccd 100644 --- a/server/src/crate_server/config_env.gleam +++ b/server/src/crate_server/config_env.gleam @@ -16,6 +16,7 @@ import crate_server/known_users_postgres import crate_server/oauth/sessions.{type Store} import crate_server/oauth/sessions_memory import crate_server/oauth/sessions_postgres +import crate_server/readiness import crate_server/sealed_store import crate_server/shelf_index import crate_server/shelf_index_postgres @@ -128,9 +129,10 @@ pub fn stores() -> #( catalog_index.Store, shelf_index.Store, follow_index.Store, + readiness.Check, ) { let key = store_key() - let #(session_store, discogs_store, known, index, shelf, follows) = + let #(session_store, discogs_store, known, index, shelf, follows, readiness) = raw_stores() #( sealed_store.wrap(session_store, key), @@ -139,6 +141,7 @@ pub fn stores() -> #( index, shelf, follows, + readiness, ) } @@ -162,6 +165,7 @@ fn raw_stores() -> #( catalog_index.Store, shelf_index.Store, follow_index.Store, + readiness.Check, ) { case envoy.get("DATABASE_URL") { Ok(url) -> @@ -194,6 +198,7 @@ fn postgres_stores( catalog_index.Store, shelf_index.Store, follow_index.Store, + readiness.Check, ), String, ) { @@ -212,7 +217,15 @@ fn postgres_stores( use index <- result.try(catalog_index_postgres.table_store(conn)) use shelf <- result.try(shelf_index_postgres.table_store(conn)) use follows <- result.map(follow_index_postgres.table_store(conn)) - #(session_store, discogs_store, known, index, shelf, follows) + #( + session_store, + discogs_store, + known, + index, + shelf, + follows, + readiness.postgres(conn), + ) } fn memory_stores() -> #( @@ -222,6 +235,7 @@ fn memory_stores() -> #( catalog_index.Store, shelf_index.Store, follow_index.Store, + readiness.Check, ) { let assert Ok(session_store) = sessions_memory.start() let assert Ok(discogs_store) = sessions_memory.start() @@ -229,7 +243,15 @@ fn memory_stores() -> #( let assert Ok(index) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() let assert Ok(follows) = follow_index.start() - #(session_store, discogs_store, known, index, shelf, follows) + #( + session_store, + discogs_store, + known, + index, + shelf, + follows, + readiness.ready(), + ) } /// App-level Discogs auth from env. Absent is fine: search still works, just diff --git a/server/src/crate_server/context.gleam b/server/src/crate_server/context.gleam index 18ec4d1..37583f2 100644 --- a/server/src/crate_server/context.gleam +++ b/server/src/crate_server/context.gleam @@ -15,6 +15,7 @@ import crate_server/oauth/authed import crate_server/oauth/config.{type Config} import crate_server/oauth/session_store import crate_server/oauth/sessions.{type OauthSession, type Store} +import crate_server/readiness import crate_server/shelf_index import gleam/json import gleam/option.{type Option, None, Some} @@ -61,6 +62,7 @@ pub type Context { /// The resolution layer's read port: candidate release rows (chain edges /// included) plus a public adoption tally. See `catalog/source`. variant_source: catalog_source.Source, + readiness: readiness.Check, oauth: Config, ) } diff --git a/server/src/crate_server/handlers/health.gleam b/server/src/crate_server/handlers/health.gleam new file mode 100644 index 0000000..f0721f3 --- /dev/null +++ b/server/src/crate_server/handlers/health.gleam @@ -0,0 +1,16 @@ +//// Deployment probes: liveness proves the process can serve HTTP, while +//// readiness also proves the appview backing store remains available. + +import crate_server/readiness.{type Check} +import wisp.{type Response} + +pub fn liveness() -> Response { + wisp.json_response("{\"status\":\"ok\"}", 200) +} + +pub fn readiness(check: Check) -> Response { + case readiness.is_ready(check) { + True -> wisp.json_response("{\"status\":\"ready\"}", 200) + False -> wisp.json_response("{\"status\":\"not ready\"}", 503) + } +} diff --git a/server/src/crate_server/readiness.gleam b/server/src/crate_server/readiness.gleam new file mode 100644 index 0000000..faefb33 --- /dev/null +++ b/server/src/crate_server/readiness.gleam @@ -0,0 +1,142 @@ +//// 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/clock +import crate_server/parallel +import gleam/dynamic/decode +import gleam/erlang/process.{type Subject} +import gleam/option.{type Option, None, Some} +import gleam/otp/actor +import gleam/result +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))" + +pub type Check { + Check(run: fn() -> Bool) +} + +type State { + State(last: Option(#(Bool, Int))) +} + +type Msg { + Probe(Subject(Bool)) +} + +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) +} + +/// Returns the representative appview schema query for regression tests. +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. +pub fn cached(probe: fn() -> Bool, now: fn() -> Int) -> Check { + let assert Ok(started) = + actor.new(State(last: None)) + |> 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) }) +} + +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) = probe_or_cached(state.last, probe, now) + process.send(reply, available) + actor.continue(State(last: next)) + } + } +} + +fn probe_or_cached( + last: Option(#(Bool, Int)), + 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) + } + None -> run_probe(probe, now) + } +} + +fn run_probe( + probe: fn() -> Bool, + now: fn() -> Int, +) -> #(Bool, Option(#(Bool, Int))) { + let available = probe() + #(available, Some(#(available, now()))) +} + +fn schema_is_ready(conn: pog.Connection) -> Bool { + let row = decode.success(Nil) + case pog.query(schema_check) |> pog.returning(row) |> pog.execute(conn) { + Ok(pog.Returned(rows: [_, ..], ..)) -> True + _ -> False + } +} diff --git a/server/src/crate_server/router.gleam b/server/src/crate_server/router.gleam index 9054fda..c1277b8 100644 --- a/server/src/crate_server/router.gleam +++ b/server/src/crate_server/router.gleam @@ -11,6 +11,7 @@ import crate_server/handlers/discogs import crate_server/handlers/edit_inbox import crate_server/handlers/feed import crate_server/handlers/graph +import crate_server/handlers/health import crate_server/handlers/oauth import crate_server/handlers/shelf import crate_server/oauth/client_metadata @@ -30,7 +31,8 @@ pub fn handle_request(req: Request, ctx: Context) -> Response { from: ctx.web.static_directory, ) case request.path_segments(req), req.method { - ["healthz"], Get -> wisp.json_response("{\"status\":\"ok\"}", 200) + ["healthz"], Get -> health.liveness() + ["readyz"], Get -> health.readiness(ctx.readiness) ["debug", "index"], Get -> debug_index.status(req, ctx) ["client-metadata.json"], Get -> json_doc(client_metadata.document(ctx.web.base_url)) diff --git a/server/src/crate_server/wiring.gleam b/server/src/crate_server/wiring.gleam index 374b7f8..c056548 100644 --- a/server/src/crate_server/wiring.gleam +++ b/server/src/crate_server/wiring.gleam @@ -16,6 +16,7 @@ import crate_server/identity_cache.{type Cache, type Identity, Identity} import crate_server/identity_resolver.{Resolver} import crate_server/musicbrainz_client import crate_server/promotion +import crate_server/readiness.{type Check} import crate_server/shelf_index import gleam/dynamic/decode import gleam/option.{None} @@ -34,6 +35,7 @@ pub fn context( follow_index follow_index: follow_index.Store, shelf_index shelf_index: shelf_index.Store, identity_cache identity_cache: Cache, + readiness readiness: Check, oauth oauth, ) -> Context { let catalog = catalog_deps(client, constellation_host, resolver) @@ -53,6 +55,7 @@ pub fn context( follow_index:, shelf_index:, variant_source: variant_source(catalog_index), + readiness:, oauth:, ) } diff --git a/server/test/health_test.gleam b/server/test/health_test.gleam new file mode 100644 index 0000000..2c38b0c --- /dev/null +++ b/server/test/health_test.gleam @@ -0,0 +1,226 @@ +//// Liveness only proves that the HTTP process can answer. Readiness also +//// proves the appview's configured backing store remains usable. + +import crate_server/handlers/health +import crate_server/parallel +import crate_server/readiness +import crate_server/router +import gleam/erlang/process +import gleam/http +import gleam/list +import gleam/otp/actor +import gleam/result +import gleam/string +import support +import wisp/simulate + +type ClockState { + ClockState(ticks: List(Int)) +} + +type ClockMsg { + Advance + Now(process.Subject(Int)) +} + +pub fn liveness_is_always_ok_test() { + let response = health.liveness() + + assert response.status == 200 + assert simulate.read_body(response) == "{\"status\":\"ok\"}" +} + +pub fn readiness_is_ok_when_the_appview_store_is_available_test() { + let response = health.readiness(readiness.ready()) + + assert response.status == 200 + assert simulate.read_body(response) == "{\"status\":\"ready\"}" +} + +pub fn readiness_is_unavailable_when_the_appview_store_is_down_test() { + let response = health.readiness(readiness.unavailable()) + + assert response.status == 503 + assert string.contains(simulate.read_body(response), "not ready") +} + +pub fn health_routes_keep_liveness_and_readiness_separate_test() { + let cfg = + support.stub_config_with( + support.unreachable_client(), + "https://resolver.test", + "http://localhost:8080", + ) + let context = support.stub_context(cfg) + + let liveness = + router.handle_request(simulate.request(http.Get, "/healthz"), context) + let readiness = + router.handle_request(simulate.request(http.Get, "/readyz"), context) + + assert liveness.status == 200 + assert readiness.status == 200 +} + +pub fn cached_readiness_coalesces_probes_inside_its_ttl_test() { + let calls = process.new_subject() + let check = + readiness.cached( + fn() { + process.send(calls, Nil) + True + }, + fn() { 100 }, + ) + + list.repeat(Nil, 3) + |> list.each(fn(_) { + assert readiness.is_ready(check) + }) + + assert support.drain_count(calls) == 1 +} + +pub fn slow_failed_probe_is_cached_from_completion_test() { + let calls = process.new_subject() + let #(advance, now) = test_clock([100, 101, 105]) + let check = + readiness.cached( + fn() { + advance() + process.send(calls, Nil) + False + }, + now, + ) + + assert !readiness.is_ready(check) + assert !readiness.is_ready(check) + + assert support.drain_count(calls) == 1 +} + +fn test_clock(ticks: List(Int)) -> #(fn() -> Nil, fn() -> Int) { + let assert Ok(started) = + actor.new(ClockState(ticks:)) + |> actor.on_message(clock_handle) + |> actor.start + let subject = started.data + #( + fn() { + process.send(subject, Advance) + Nil + }, + fn() { parallel.try_call(subject, 1000, Now) |> result.unwrap(0) }, + ) +} + +fn clock_handle( + state: ClockState, + msg: ClockMsg, +) -> actor.Next(ClockState, ClockMsg) { + case msg { + Advance -> actor.continue(ClockState(ticks: drop_tick(state.ticks))) + Now(reply) -> { + let #(time, ticks) = next_tick(state.ticks) + process.send(reply, time) + actor.continue(ClockState(ticks:)) + } + } +} + +fn drop_tick(ticks: List(Int)) -> List(Int) { + case ticks { + [_tick, ..remaining] -> remaining + [] -> [] + } +} + +fn next_tick(ticks: List(Int)) -> #(Int, List(Int)) { + case ticks { + [tick, ..remaining] -> #(tick, remaining) + [] -> #(0, []) + } +} + +pub fn schema_probe_covers_required_table_columns_test() { + let sql = readiness.schema_check_sql() + + [ + "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')", + ] + |> list.each(fn(shape) { + assert string.contains(sql, shape) + }) +} + +pub fn schema_probe_requires_store_operation_privileges_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) + }) +} + +pub fn schema_probe_requires_not_null_and_conflict_targets_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')", + ] + |> list.each(fn(clause) { + assert string.contains(sql, clause) + }) +} diff --git a/server/test/support.gleam b/server/test/support.gleam index 244e978..20624f6 100644 --- a/server/test/support.gleam +++ b/server/test/support.gleam @@ -22,6 +22,7 @@ import crate_server/oauth/keys import crate_server/oauth/sessions import crate_server/oauth/sessions_memory import crate_server/oauth/store +import crate_server/readiness import crate_server/shelf_index import crate_server/wiring import gleam/dict @@ -259,6 +260,7 @@ pub fn stub_context_with( shelf_index: fresh_shelf_index(), follow_index: fresh_follow_index(), variant_source: empty_variant_source(), + readiness: readiness.ready(), oauth: cfg, ) } -- 2.51.2