From 2f729af21800523ce3cbb9c91c68b3924cac24d1 Mon Sep 17 00:00:00 2001 From: Furkan Akbulutlar <61627880+FrkAk@users.noreply.github.com> Date: Tue, 4 Aug 2026 22:58:29 +0200 Subject: [PATCH 1/2] feat: add database housekeeping and operational analytics (#283) --- .github/workflows/ci.yml | 18 +- bun.lock | 18 +- docker-compose.yml | 4 + docker/extensions.sql | 21 ++ docker/init-pg-cron.sql | 22 -- docker/rls-functions.sql | 119 +++++++++ eslint.config.mjs | 2 +- lib/db/raw/purge-expired-rows.ts | 35 +++ package.json | 11 +- scripts/apply-owner-rls.ts | 34 +-- scripts/db-stats.ts | 357 ++++++++++++++++++++++++++ scripts/smoke-workers.ts | 96 ++++++- scripts/verify-rls.ts | 50 +++- tests/data/housekeeping.test.ts | 413 +++++++++++++++++++++++++++++++ tests/scripts/db-stats.test.ts | 209 ++++++++++++++++ tests/setup/migrate.ts | 21 ++ worker-cf.ts | 80 ++++++ wrangler.jsonc | 4 + 18 files changed, 1441 insertions(+), 73 deletions(-) create mode 100644 docker/extensions.sql delete mode 100644 docker/init-pg-cron.sql create mode 100644 lib/db/raw/purge-expired-rows.ts create mode 100644 scripts/db-stats.ts create mode 100644 tests/data/housekeeping.test.ts create mode 100644 tests/scripts/db-stats.test.ts diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index ecfcddba..567dbe62 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -51,19 +51,11 @@ jobs: # next declares sharp ^0.34.5 (optional) and miniflare pins sharp # 0.34.5 exactly, so the fix is only reachable by forcing a major # override against both parents. Drop the ignore once either admits - # sharp >= 0.35.0. Tracked in PYZ-345. (fast-uri + postcss are fixed - # via package.json overrides; see the //overrides note there.) - # - # GHSA-mh99-v99m-4gvg: brace-expansion DoS, vulnerable range <= 5.0.7, - # spanning every major. Reached only through dev dependencies, via - # minimatch under eslint, eslint-plugin-import, typescript-eslint, and - # @opennextjs/cloudflare > glob. No runtime path: `bun audit --prod` - # does not report it. An override was rejected because 5.0.8 is - # the only patched release, so a single range drags the consumers - # pinned to 1.1.16 and 2.1.2 across two majors. Drop the ignore once - # upstream publishes patched 1.x and 2.x lines, or once the tree - # consolidates on brace-expansion >= 5.0.8. Tracked in PYZ-345. - run: bun audit --audit-level=high --ignore=GHSA-f88m-g3jw-g9cj --ignore=GHSA-mh99-v99m-4gvg + # sharp >= 0.35.0. Tracked in PYZ-345. (fast-uri, ip-address, postcss, + # and undici are fixed via package.json overrides; brace-expansion is + # fixed per-major in bun.lock, patched lines now exist for every + # major. See the //overrides note in package.json.) + run: bun audit --audit-level=high --ignore=GHSA-f88m-g3jw-g9cj - name: Validate PR title if: github.event_name == 'pull_request' env: diff --git a/bun.lock b/bun.lock index 706ebae5..0309f8c8 100644 --- a/bun.lock +++ b/bun.lock @@ -56,8 +56,10 @@ "@better-auth/core@1.6.23": "patches/@better-auth%2Fcore@1.6.23.patch", }, "overrides": { - "fast-uri": ">=3.1.4 <4", + "fast-uri": ">=3.1.5 <4", + "ip-address": ">=10.3.1 <11", "postcss": ">=8.5.18 <9", + "undici": ">=7.29.0 <8", }, "packages": { "@alloc/quick-lru": ["@alloc/quick-lru@5.2.0", "", {}, "sha512-UrcABB+4bUrFABwbluTIBErXwvbsU/V7TZWfmbgJfbkwiBuziS9gxdODUyuiecfdGQ85jglMW6juS3+z5TsKLw=="], @@ -804,7 +806,7 @@ "bowser": ["bowser@2.14.1", "", {}, "sha512-tzPjzCxygAKWFOJP011oxFHs57HzIhOEracIgAePE4pqB3LikALKnSzUyU4MGs9/iCEUuHlAJTjTc5M+u7YEGg=="], - "brace-expansion": ["brace-expansion@1.1.16", "", { "dependencies": { "balanced-match": "^1.0.0", "concat-map": "0.0.1" } }, "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw=="], + "brace-expansion": ["brace-expansion@1.1.18", "", { "dependencies": { "balanced-match": "^1.0.0", "concat-map": "0.0.1" } }, "sha512-Edep/X9fGqVNmzKBVsDYIOtD+z1tuezV70LBjdCst9Tqu76lsnvRiZ6oTic1n+/BIwX6QDGAO94PN4N2SADvtw=="], "braces": ["braces@3.0.3", "", { "dependencies": { "fill-range": "^7.1.1" } }, "sha512-yQbXgO/OSZVD2IsiLlro+7Hf6Q18EJrKSEsdoMzKePKXct3gvD8oLcOQdIzGupr5Fj+EDe8gO/lxc1BzfMpxvA=="], @@ -1028,7 +1030,7 @@ "fast-levenshtein": ["fast-levenshtein@2.0.6", "", {}, "sha512-DCXu6Ifhqcks7TZKY3Hxp3y6qphY5SJZmrWMDrKcERSOXWQdMhU9Ig/PYrzyw/ul9jOIyh0N4M0tbC5hodg8dw=="], - "fast-uri": ["fast-uri@3.1.4", "", {}, "sha512-8JnbkQ4juDyvYs4mgFGQqg4yCYtFDtUtmp2QIQq11ZZe5CFQ5wcqm1rqDgAh/QdMySuBnPzMUiJUNZG5N/AiQw=="], + "fast-uri": ["fast-uri@3.1.5", "", {}, "sha512-gHwA1O9LDIcKunMKhObS/HimwtehO1nPUECKAu5TpKgaO19fcWEl4bliWe1jWxVFvIXztJjjQ4L8XQ1EU9f7Jw=="], "fastq": ["fastq@1.20.1", "", { "dependencies": { "reusify": "^1.0.4" } }, "sha512-GGToxJ/w1x32s/D2EKND7kTil4n8OVk/9mycTc4VDza13lOvpUZTGX3mFSCtV9ksdGBVzvsyAVLM6mHFThxXxw=="], @@ -1156,7 +1158,7 @@ "internal-slot": ["internal-slot@1.1.0", "", { "dependencies": { "es-errors": "^1.3.0", "hasown": "^2.0.2", "side-channel": "^1.1.0" } }, "sha512-4gd7VpWNQNB4UKKCFFVcp1AVv+FMOgs9NKzjHKusc8jTMhd5eL1NqQqOpE0KzMds804/yHlglp3uxgluOqAPLw=="], - "ip-address": ["ip-address@10.2.0", "", {}, "sha512-/+S6j4E9AHvW9SWMSEY9Xfy66O5PWvVEJ08O0y5JGyEKQpojb0K0GKpz/v5HJ/G0vi3D2sjGK78119oXZeE0qA=="], + "ip-address": ["ip-address@10.4.0", "", {}, "sha512-oSK96Grm3aP6OrS263xVxbNDGVL7rzBtYdpGqlDG8iQdoenDoTs/nkki+DflYbAEE8Xl6o5YxhxlrKvI3nqKXQ=="], "ipaddr.js": ["ipaddr.js@1.9.1", "", {}, "sha512-0KI/607xoxSToH7GjN1FfSbLoU0+btTicjsQSWQlh/hZykN8KpmMf7uYwPW3R+akZ6R/w18ZlXSHBYXiYUPO3g=="], @@ -1726,7 +1728,7 @@ "unbox-primitive": ["unbox-primitive@1.1.0", "", { "dependencies": { "call-bound": "^1.0.3", "has-bigints": "^1.0.2", "has-symbols": "^1.1.0", "which-boxed-primitive": "^1.1.1" } }, "sha512-nWJ91DjeOkej/TA8pXQ3myruKpKEYgqvpw9lz4OPHj/NWFNluYrjbz9j01CJ8yKQd2g4jFoOkINCTW2I5LEEyw=="], - "undici": ["undici@7.28.0", "", {}, "sha512-cRZYrTDwWznlnRiPjggAGxZXanty6M8RV1ff8Wm4LWXBp7/IG8v5DnOm74DtUBp9OONpK75YlPnIjQqX0dBDtA=="], + "undici": ["undici@7.29.0", "", {}, "sha512-IDxfleLmmbSskfWSUATiN1nfn2rDuvnMOqb5CWR92iIfojA0Ud+ulOAAEQ57LPr9rWmsreUyf5lwyao+7GNNVw=="], "undici-types": ["undici-types@7.18.2", "", {}, "sha512-AsuCzffGHJybSaRrmr5eHr81mwJU3kjw6M+uprWvCXiNeN9SOGwQ3Jn8jb8m3Z6izVgknn1R0FTCEAP2QrLY/w=="], @@ -2042,7 +2044,7 @@ "@opennextjs/aws/esbuild/@esbuild/win32-x64": ["@esbuild/win32-x64@0.25.4", "", { "os": "win32", "cpu": "x64" }, "sha512-nOT2vZNw6hJ+z43oP1SPea/G/6AbN6X+bGNhNuq8NtRHy4wsMhw765IKLNmnjek7GvjWBYQ8Q5VBoYTFg9y1UQ=="], - "@typescript-eslint/typescript-estree/minimatch/brace-expansion": ["brace-expansion@5.0.7", "", { "dependencies": { "balanced-match": "^4.0.2" } }, "sha512-7oFy703dxfY3/NLxC1fh2SUCQ0H9rmAY+5EpDVfXjUTTs+HEwR2nYaqLv+GWcTsumwxPfiz6CzCNkwXwBUwqCA=="], + "@typescript-eslint/typescript-estree/minimatch/brace-expansion": ["brace-expansion@5.0.9", "", { "dependencies": { "balanced-match": "^4.0.2" } }, "sha512-ScQ4IuvIEF1TMlP7Zt+vjJ//9zlPb2SDcxWxM3bk8s6t6GGdJ7KO1dCcTidOPJKePW30LE/2cT7wCyPho9/Wxg=="], "ajv-formats/ajv/json-schema-traverse": ["json-schema-traverse@1.0.0", "", {}, "sha512-NM8/P9n3XjXhIZn1lLhkFaACTOURQXjWhV4BA/RnOv8xvgqtqpAX9IO4mRQxSx1Rlo4tqzeqb0sOlruaOy3dug=="], @@ -2054,7 +2056,7 @@ "form-data/mime-types/mime-db": ["mime-db@1.52.0", "", {}, "sha512-sPU4uV7dYlvtWJxwwxHD0PuihVNiE7TyAbQ5SWxDCB9mUYvOgroQOwYQQOKPJ8CIbE+1ETVlOoK1UC2nU3gYvg=="], - "glob/minimatch/brace-expansion": ["brace-expansion@5.0.7", "", { "dependencies": { "balanced-match": "^4.0.2" } }, "sha512-7oFy703dxfY3/NLxC1fh2SUCQ0H9rmAY+5EpDVfXjUTTs+HEwR2nYaqLv+GWcTsumwxPfiz6CzCNkwXwBUwqCA=="], + "glob/minimatch/brace-expansion": ["brace-expansion@5.0.9", "", { "dependencies": { "balanced-match": "^4.0.2" } }, "sha512-ScQ4IuvIEF1TMlP7Zt+vjJ//9zlPb2SDcxWxM3bk8s6t6GGdJ7KO1dCcTidOPJKePW30LE/2cT7wCyPho9/Wxg=="], "miniflare/sharp/@img/sharp-darwin-arm64": ["@img/sharp-darwin-arm64@0.35.2", "", { "optionalDependencies": { "@img/sharp-libvips-darwin-arm64": "1.3.1" }, "os": "darwin", "cpu": "arm64" }, "sha512-eEieHsMksAW4IiO5NzauESRl2D2qz3J/kwUxUrSfV06A93eEaRfMpHXyUb1mAqrR7i8U9A0GRqE9pjn6u1Jjpg=="], @@ -2212,7 +2214,7 @@ "wrap-ansi/strip-ansi/ansi-regex": ["ansi-regex@6.2.2", "", {}, "sha512-Bq3SmSpyFHaWjPk8If9yc6svM8c56dB5BAtW4Qbw5jHTwwXXcTLoRMkpDJp6VL0XzlWaCHTXrkFURMYmD0sLqg=="], - "@node-minify/core/glob/minimatch/brace-expansion": ["brace-expansion@2.1.2", "", { "dependencies": { "balanced-match": "^1.0.0" } }, "sha512-w5JZcKgdhDOgOwm8H+KgbosopHMuGcl6qbulwjtz3SM7I7P3yW1eAjzMPLrIE+NQ9vjgANKHWeMHnrT0OXW1oA=="], + "@node-minify/core/glob/minimatch/brace-expansion": ["brace-expansion@2.1.4", "", { "dependencies": { "balanced-match": "^1.0.0" } }, "sha512-hGfVzPxthbf3+2yjg/RBs60cB0FhqBS/zvdV/4wn4/BmN0bNMMHPc4V/BbFieqf1TKAGGAHnY4eSjajCl0f2Xg=="], "@node-minify/core/glob/path-scurry/lru-cache": ["lru-cache@10.4.3", "", {}, "sha512-JNAzZcXrCt42VGLuYz0zfAzDfAvJWW6AfYlDBQyDV5DClI2m5sAmK+OIO7s59XfsRsWHp02jAJrRadPRGTt6SQ=="], diff --git a/docker-compose.yml b/docker-compose.yml index 4251a1cc..8f7a471f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -3,6 +3,10 @@ services: db: image: postgres:18 restart: unless-stopped + # pg_stat_statements collects nothing unless preloaded; CREATE EXTENSION + # alone (docker/extensions.sql) only creates the SQL objects. Applies on + # the next `docker compose up -d`. + command: ["postgres", "-c", "shared_preload_libraries=pg_stat_statements"] environment: POSTGRES_USER: piyaz # Superuser: exempt from FORCE ROW LEVEL SECURITY, so this credential diff --git a/docker/extensions.sql b/docker/extensions.sql new file mode 100644 index 00000000..254d9219 --- /dev/null +++ b/docker/extensions.sql @@ -0,0 +1,21 @@ +-- ============================================================================= +-- Extensions every Piyaz database carries. Owner-only: CREATE EXTENSION needs +-- the database owner on Neon, never the least-privilege migration role, so +-- this is NOT a Drizzle migration. Applied by scripts/apply-owner-rls.ts +-- (db:rls:owner, hosted), the db:rls psql chain (self-host), and +-- tests/setup/migrate.ts (testcontainer). scripts/verify-rls.ts derives its +-- extension contract from this file, so a forgotten owner apply fails the +-- deploy. Idempotent. +-- +-- pg_stat_statements: CREATE EXTENSION succeeds without preloading, but +-- collecting/querying stats requires shared_preload_libraries — set for Neon +-- by the platform and for self-host by the postgres command in +-- docker-compose.yml. The read path is scripts/db-stats.ts (owner-only). +-- +-- Extensions live in their own schema, never public: the public schema is +-- owned by Drizzle, and `drizzle-kit push` (throwaway test DB) tries to drop +-- any non-Drizzle objects it finds there. +-- ============================================================================= + +CREATE SCHEMA IF NOT EXISTS extensions; +CREATE EXTENSION IF NOT EXISTS pg_stat_statements WITH SCHEMA extensions; diff --git a/docker/init-pg-cron.sql b/docker/init-pg-cron.sql deleted file mode 100644 index 7df905de..00000000 --- a/docker/init-pg-cron.sql +++ /dev/null @@ -1,22 +0,0 @@ --- Nightly purge of expired and revoked OAuth tokens. --- Apply against the application database. pg_cron schedules in UTC. --- Idempotent — safe to re-run. - -CREATE EXTENSION IF NOT EXISTS pg_cron; - -DO $$ -BEGIN - PERFORM cron.unschedule('purge-oauth-tokens') - WHERE EXISTS (SELECT 1 FROM cron.job WHERE jobname = 'purge-oauth-tokens'); -END $$; - -SELECT cron.schedule( - 'purge-oauth-tokens', - '0 3 * * *', - $$ - DELETE FROM piyaz_auth."oauthRefreshToken" - WHERE revoked IS NOT NULL OR "expiresAt" < now(); - DELETE FROM piyaz_auth."oauthAccessToken" - WHERE "expiresAt" < now(); - $$ -); diff --git a/docker/rls-functions.sql b/docker/rls-functions.sql index bd6fe26a..12db98e6 100644 --- a/docker/rls-functions.sql +++ b/docker/rls-functions.sql @@ -1317,3 +1317,122 @@ CREATE TRIGGER task_edges_touch_project_delete REFERENCING OLD TABLE AS changed_edges FOR EACH STATEMENT EXECUTE FUNCTION public.touch_projects_for_changed_task_edges(); + +-- --------------------------------------------------------------------------- +-- Nightly housekeeping sweep: bounded, idempotent deletion of expired auth +-- artifacts and stale invite codes. Called by the Cloudflare cron in +-- worker-cf.ts through service_role (JS caller: +-- lib/db/raw/purge-expired-rows.ts). SECURITY DEFINER because service_role +-- lacks DELETE on piyaz_auth."session" and all rights on +-- piyaz_auth."verification", and auth_role's 15s statement_timeout is unfit +-- for bulk deletes. +-- +-- The CONSTANT declarations are the retention matrix's single source of +-- truth. Dry runs and live runs share the same victim selection (the +-- data-modifying CTE always executes; NOT p_dry_run turns it into a no-op), +-- and the per-table LIMIT bounds one run — leftovers roll to the next. +-- Retained by policy: legal_acceptances, activity_events, note_revisions, +-- oauthConsent, jwks, piyaz_auth.invitation, account. +-- --------------------------------------------------------------------------- +CREATE OR REPLACE FUNCTION public.purge_expired_rows(p_dry_run boolean, p_batch_limit integer) +RETURNS TABLE (table_name text, row_count integer) +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, pg_temp +AS $$ +DECLARE + c_oauth_grace CONSTANT interval := interval '24 hours'; + c_session_grace CONSTANT interval := interval '7 days'; + c_verification_grace CONSTANT interval := interval '7 days'; + c_invite_grace CONSTANT interval := interval '30 days'; +BEGIN + IF p_batch_limit IS NULL OR p_batch_limit < 1 OR p_batch_limit > 50000 THEN + RAISE EXCEPTION 'p_batch_limit out of range: %', p_batch_limit; + END IF; + + WITH victims AS ( + SELECT ctid FROM piyaz_auth."oauthAccessToken" + WHERE "expiresAt" < now() - c_oauth_grace + LIMIT p_batch_limit + ), deleted AS ( + DELETE FROM piyaz_auth."oauthAccessToken" t + USING victims v + WHERE t.ctid = v.ctid AND NOT p_dry_run + RETURNING 1 + ) + SELECT CASE WHEN p_dry_run THEN (SELECT count(*) FROM victims) + ELSE (SELECT count(*) FROM deleted) END::integer + INTO row_count; + table_name := 'oauthAccessToken'; + RETURN NEXT; + + WITH victims AS ( + SELECT ctid FROM piyaz_auth."oauthRefreshToken" + WHERE revoked < now() - c_oauth_grace + OR "expiresAt" < now() - c_oauth_grace + LIMIT p_batch_limit + ), deleted AS ( + DELETE FROM piyaz_auth."oauthRefreshToken" t + USING victims v + WHERE t.ctid = v.ctid AND NOT p_dry_run + RETURNING 1 + ) + SELECT CASE WHEN p_dry_run THEN (SELECT count(*) FROM victims) + ELSE (SELECT count(*) FROM deleted) END::integer + INTO row_count; + table_name := 'oauthRefreshToken'; + RETURN NEXT; + + WITH victims AS ( + SELECT ctid FROM piyaz_auth."session" + WHERE "expiresAt" < now() - c_session_grace + LIMIT p_batch_limit + ), deleted AS ( + DELETE FROM piyaz_auth."session" t + USING victims v + WHERE t.ctid = v.ctid AND NOT p_dry_run + RETURNING 1 + ) + SELECT CASE WHEN p_dry_run THEN (SELECT count(*) FROM victims) + ELSE (SELECT count(*) FROM deleted) END::integer + INTO row_count; + table_name := 'session'; + RETURN NEXT; + + WITH victims AS ( + SELECT ctid FROM piyaz_auth."verification" + WHERE "expiresAt" < now() - c_verification_grace + LIMIT p_batch_limit + ), deleted AS ( + DELETE FROM piyaz_auth."verification" t + USING victims v + WHERE t.ctid = v.ctid AND NOT p_dry_run + RETURNING 1 + ) + SELECT CASE WHEN p_dry_run THEN (SELECT count(*) FROM victims) + ELSE (SELECT count(*) FROM deleted) END::integer + INTO row_count; + table_name := 'verification'; + RETURN NEXT; + + WITH victims AS ( + SELECT ctid FROM public.team_invite_code + WHERE revoked_at < now() - c_invite_grace + OR expires_at < now() - c_invite_grace + LIMIT p_batch_limit + ), deleted AS ( + DELETE FROM public.team_invite_code t + USING victims v + WHERE t.ctid = v.ctid AND NOT p_dry_run + RETURNING 1 + ) + SELECT CASE WHEN p_dry_run THEN (SELECT count(*) FROM victims) + ELSE (SELECT count(*) FROM deleted) END::integer + INTO row_count; + table_name := 'team_invite_code'; + RETURN NEXT; +END; +$$; +REVOKE EXECUTE ON FUNCTION public.purge_expired_rows(boolean, integer) FROM public; +REVOKE EXECUTE ON FUNCTION public.purge_expired_rows(boolean, integer) FROM app_user; +GRANT EXECUTE ON FUNCTION public.purge_expired_rows(boolean, integer) TO service_role; diff --git a/eslint.config.mjs b/eslint.config.mjs index d8421a02..efff6611 100644 --- a/eslint.config.mjs +++ b/eslint.config.mjs @@ -106,7 +106,7 @@ const eslintConfig = [ selector: "CallExpression[callee.object.name='serviceRoleDb'][callee.property.name=/^(select|insert|update|delete)$/]", message: - "serviceRoleDb. is BYPASSRLS. Allowed sites: lib/data/oauth-session.ts (oauth tables), lib/data/account.ts (clearOrgMembershipArtifacts, scrubLegalAcceptances, enumerateOwnedOrgsForDeletion), lib/data/membership.ts (admin lookups). Consider whether a SECURITY DEFINER function in docker/rls-functions.sql can replace this call site.", + "serviceRoleDb. is BYPASSRLS. Allowed sites: lib/data/oauth-session.ts (oauth tables), lib/data/account.ts (clearOrgMembershipArtifacts, scrubLegalAcceptances, enumerateOwnedOrgsForDeletion), lib/data/membership.ts (admin lookups), worker-cf.ts scheduled() (service-role handle from requestDbStore into lib/db/raw/purge-expired-rows.ts, which only EXECUTEs the SECURITY DEFINER public.purge_expired_rows). Consider whether a SECURITY DEFINER function in docker/rls-functions.sql can replace this call site.", }, { selector: "MemberExpression[object.name='db'][property.name='query']", diff --git a/lib/db/raw/purge-expired-rows.ts b/lib/db/raw/purge-expired-rows.ts new file mode 100644 index 00000000..4236948f --- /dev/null +++ b/lib/db/raw/purge-expired-rows.ts @@ -0,0 +1,35 @@ +import { sql } from "drizzle-orm"; +import { executeRaw, type ServiceRoleConn } from "@/lib/db/raw"; + +/** Rows the nightly sweep may delete per table per run; leftovers roll to the next run. */ +export const HOUSEKEEPING_BATCH_LIMIT = 5000; + +/** One per-table result row from `public.purge_expired_rows`. */ +export interface PurgeResultRow { + /** Table the row reports on (fixed order, one row per table per run). */ + table_name: string; + /** Rows deleted (live run) or rows that would be deleted (dry run). */ + row_count: number; +} + +/** + * Run the housekeeping sweep via `public.purge_expired_rows` (SECURITY + * DEFINER, service_role-only; retention windows live in + * `docker/rls-functions.sql`). Dry runs count the same victims a live run + * would delete without mutating anything. + * + * @param conn - BYPASSRLS service-role client. + * @param dryRun - True to count would-be deletions without deleting. + * @param batchLimit - Per-table row cap for this run (1..50000). + * @returns One result row per swept table, fixed order. + */ +export async function purgeExpiredRows( + conn: ServiceRoleConn, + dryRun: boolean, + batchLimit: number, +): Promise { + return executeRaw( + conn, + sql`SELECT table_name, row_count FROM public.purge_expired_rows(${dryRun}, ${batchLimit})`, + ); +} diff --git a/package.json b/package.json index 6c0c3758..06c3ae24 100644 --- a/package.json +++ b/package.json @@ -4,10 +4,12 @@ "private": true, "license": "AGPL-3.0-or-later", "type": "module", - "//overrides": "TEMP security pins for transitive CVEs. REMOVE each entry once its parent ships the fix and `bun audit --audit-level=high` passes without it. fast-uri >=3.1.4 <4 (GHSA-v2hh-gcrm-f6hx, via ajv under @modelcontextprotocol/sdk and eslint; ajv declares ^3.0.1, so the parent range already allows the fix). postcss >=8.5.18 <9 (GHSA-6g55-p6wh-862q, GHSA-r28c-9q8g-f849, via next, which pins postcss 8.4.31 EXACTLY: this entry contradicts its parent's declared range rather than nudging a loose one, so re-check it on every next bump and drop it as soon as next unpins). sharp is handled separately via a CI --ignore (no takeable fix). Tracked in PYZ-345.", + "//overrides": "TEMP security pins for transitive CVEs. REMOVE each entry once its parent ships the fix and `bun audit --audit-level=high` passes without it. fast-uri >=3.1.5 <4 (GHSA-v2hh-gcrm-f6hx, GHSA-7p8r-x3mc-p8w7, via ajv under @modelcontextprotocol/sdk and eslint; ajv declares ^3.0.1, so the parent range already allows the fix). ip-address >=10.3.1 <11 (GHSA-mwp4-54f8-5fhr, via express-rate-limit under @modelcontextprotocol/sdk; the parent declares ^10.2.0, so this only raises the floor). postcss >=8.5.18 <9 (GHSA-6g55-p6wh-862q, GHSA-r28c-9q8g-f849, via next, which pins postcss 8.4.31 EXACTLY: this entry contradicts its parent's declared range rather than nudging a loose one, so re-check it on every next bump and drop it as soon as next unpins). undici >=7.29.0 <8 (GHSA-4cwx-7wf7-3272, via miniflare under wrangler, which pins undici 7.28.0 EXACTLY: same caveat as postcss, re-check on every wrangler bump and drop once miniflare moves past 7.29.0). sharp is handled separately via a CI --ignore (no takeable fix). Tracked in PYZ-345.", "overrides": { - "fast-uri": ">=3.1.4 <4", - "postcss": ">=8.5.18 <9" + "fast-uri": ">=3.1.5 <4", + "ip-address": ">=10.3.1 <11", + "postcss": ">=8.5.18 <9", + "undici": ">=7.29.0 <8" }, "scripts": { "dev": "next dev --webpack", @@ -33,10 +35,11 @@ "db:setup": "docker compose --env-file .env.local up -d --wait && docker exec -i piyaz-db-1 psql -U piyaz -d piyaz < docker/init-auth.sql && docker exec piyaz-db-1 /docker-entrypoint-initdb.d/02-rls.sh && bun run db:migrate && bun run db:rls", "db:generate": "drizzle-kit generate", "db:migrate": "bun run scripts/assert-migration-journal.ts && drizzle-kit migrate", - "db:rls": "docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/grants.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/grants-auth.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/role-settings.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/rls-functions.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/rls-policies.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/storage.sql", + "db:rls": "docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/extensions.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/grants.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/grants-auth.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/role-settings.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/rls-functions.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/rls-policies.sql && docker exec -i piyaz-db-1 psql -v ON_ERROR_STOP=1 -U piyaz -d piyaz < docker/storage.sql", "db:rls:ci": "bun scripts/apply-public-rls.ts", "db:rls:verify": "bun scripts/verify-rls.ts", "db:rls:owner": "bun scripts/apply-owner-rls.ts", + "db:stats": "bun scripts/db-stats.ts", "db:baseline": "bun run scripts/baseline-self-host.ts", "db:grandfather-verified": "bun scripts/grandfather-verified-users.ts", "db:push": "bun run scripts/assert-throwaway-db.ts && drizzle-kit push", diff --git a/scripts/apply-owner-rls.ts b/scripts/apply-owner-rls.ts index 915c311c..1986fd23 100644 --- a/scripts/apply-owner-rls.ts +++ b/scripts/apply-owner-rls.ts @@ -1,10 +1,11 @@ /** - * Apply the owner-managed RLS SQL: the piyaz_auth grants - * (docker/grants-auth.sql), the request-path role settings - * (docker/role-settings.sql) and the SECURITY DEFINER helpers + triggers - * (docker/rls-functions.sql). These read or own piyaz_auth, or need ADMIN - * OPTION on a role, so they must run as the database owner, never the - * least-privilege migration role. Idempotent (CREATE OR REPLACE / GRANT / + * Apply the owner-managed RLS SQL: the extensions (docker/extensions.sql), + * the piyaz_auth grants (docker/grants-auth.sql), the request-path role + * settings (docker/role-settings.sql) and the SECURITY DEFINER helpers + + * triggers (docker/rls-functions.sql). These create extensions, read or own + * piyaz_auth, or need ADMIN OPTION on a role, so they must run as the + * database owner, never the least-privilege migration role. Idempotent + * (CREATE EXTENSION IF NOT EXISTS / CREATE OR REPLACE / GRANT / * ALTER ROLE ... SET). * * Reads DATABASE_OWNER_URL. Set this only in a trusted local shell, never as a @@ -16,22 +17,25 @@ import { join } from "node:path"; import postgres from "postgres"; const OWNER_RLS_FILES = [ + "extensions.sql", "grants-auth.sql", "role-settings.sql", "rls-functions.sql", ] as const; /** - * Read the owner connection string from the environment. + * Read the owner connection string from the environment. Shared by every + * owner-run script (`db:rls:owner`, `db:stats`) so the env-var name and the + * trusted-shell contract live in one place. * * @returns The database-owner DIRECT connection string. * @throws Error when DATABASE_OWNER_URL is unset. */ -function ownerUrl(): string { +export function ownerUrl(): string { const url = process.env.DATABASE_OWNER_URL; if (!url) { throw new Error( - "DATABASE_OWNER_URL is required to apply owner-managed RLS (database owner role).", + "DATABASE_OWNER_URL is required (database owner role; set it only in a trusted local shell).", ); } return url; @@ -64,9 +68,11 @@ async function applyOwnerRls(url: string): Promise { } } -try { - await applyOwnerRls(ownerUrl()); -} catch (err) { - console.error(err instanceof Error ? err.message : err); - process.exit(1); +if (import.meta.main) { + try { + await applyOwnerRls(ownerUrl()); + } catch (err) { + console.error(err instanceof Error ? err.message : err); + process.exit(1); + } } diff --git a/scripts/db-stats.ts b/scripts/db-stats.ts new file mode 100644 index 00000000..0ce21d95 --- /dev/null +++ b/scripts/db-stats.ts @@ -0,0 +1,357 @@ +/** + * Read-only operational analytics over pg_stat_statements: per-class volume + * (read / write / control / other), top-N reads and writes by total time and + * by calls, and the stats_reset observation window (Neon clears the view on + * every compute suspend, so the window line is load-bearing). + * + * Reads DATABASE_OWNER_URL — the same trusted-local-shell-only contract as + * scripts/apply-owner-rls.ts; Worker runtime roles have no stats access. + * Privacy: query text is printed only for read/write DML, which + * pg_stat_statements normalizes ($1 placeholders — no bind values, no tenant + * data). control/other statements are aggregate-only because utility + * statements are NOT normalized (`SET app.user_id = ''` carries a real + * user id) and track_utility cannot be disabled on Neon. Never paste this + * report into the public repo. + * + * Usage: `bun run db:stats [--top ]`. + */ +import postgres from "postgres"; +import { ownerUrl } from "./apply-owner-rls"; + +/** Statements shown per section when `--top` is not given. */ +const DEFAULT_TOP_N = 20; + +/** Display width for the normalized query column. */ +const QUERY_DISPLAY_WIDTH = 100; + +/** + * SQL CASE fragment classing every pg_stat_statements row. `\y` is the + * Postgres regex word boundary. Order matters: transaction/session control + * first, then plain DML writes, then writing CTEs, then reads; everything + * else (DDL, maintenance, other utility) falls through to 'other'. + */ +const CLASS_CASE = `CASE + WHEN query ~* '^\\s*(begin|commit|rollback|set|reset|show|deallocate|discard|listen|unlisten|savepoint|release|prepare|fetch|close)\\y' THEN 'control' + WHEN query ~* '^\\s*(insert|update|delete|merge)\\y' THEN 'write' + WHEN query ~* '^\\s*with\\y' AND query ~* '\\y(insert|update|delete|merge)\\y' THEN 'write' + WHEN query ~* '^\\s*(select|values|table|with)\\y' THEN 'read' + ELSE 'other' +END`; + +/** Scope every query to the connected database's statements. */ +const DBID_FILTER = `dbid = (SELECT oid FROM pg_database WHERE datname = current_database())`; + +/** One row of the per-class aggregate. */ +export interface ClassSummaryRow { + statement_class: string; + statements: string; + calls: string; + total_exec_time: number; +} + +/** One normalized statement row from pg_stat_statements. */ +export interface StatementRow { + calls: string; + total_exec_time: number; + mean_exec_time: number; + rows: string; + shared_blks_hit: string; + shared_blks_read: string; + shared_blks_dirtied: string; + temp_blks_read: string; + temp_blks_written: string; + wal_bytes: string; + query: string; +} + +/** Everything the report renders, fetched in one connection. */ +export interface DbStatsReport { + statsReset: Date | null; + classSummary: ClassSummaryRow[]; + readsByTime: StatementRow[]; + readsByCalls: StatementRow[]; + writesByTime: StatementRow[]; + writesByCalls: StatementRow[]; +} + +/** + * Collapse whitespace runs to single spaces and truncate with an ellipsis. + * + * @param query - Normalized statement text. + * @param width - Maximum output length including the ellipsis. + * @returns Single-line text no longer than `width`. + */ +export function sanitizeQueryText( + query: string, + width = QUERY_DISPLAY_WIDTH, +): string { + const collapsed = query.replace(/\s+/g, " ").trim(); + if (collapsed.length <= width) return collapsed; + return `${collapsed.slice(0, width - 1)}…`; +} + +/** + * Format milliseconds with two decimals. + * + * @param ms - Millisecond value. + * @returns Fixed-point string, e.g. `12.34`. + */ +export function formatMs(ms: number): string { + return ms.toFixed(2); +} + +/** + * Format a byte count with a binary unit suffix. + * + * @param value - Byte count (string when the driver returns int8/numeric). + * @returns Humanized size, e.g. `1.5 MiB`. + */ +export function formatBytes(value: string | number): string { + let n = Number(value); + if (!Number.isFinite(n)) return String(value); + for (const unit of ["B", "KiB", "MiB", "GiB"]) { + if (n < 1024) + return unit === "B" ? `${n} ${unit}` : `${n.toFixed(1)} ${unit}`; + n /= 1024; + } + return `${n.toFixed(1)} TiB`; +} + +/** + * Render rows as a padded text table. + * + * @param headers - Column headers. + * @param rows - Cell values, one array per row. + * @param leftAligned - Indexes of left-aligned columns (rest right-align). + * @returns Multi-line table string. + */ +export function renderTable( + headers: string[], + rows: string[][], + leftAligned: number[] = [], +): string { + const left = new Set(leftAligned); + const widths = headers.map((h, i) => + Math.max(h.length, ...rows.map((r) => r[i].length)), + ); + const line = (cells: string[]): string => + cells + .map((cell, i) => + left.has(i) ? cell.padEnd(widths[i]) : cell.padStart(widths[i]), + ) + .join(" ") + .trimEnd(); + return [line(headers), ...rows.map(line)].join("\n"); +} + +/** + * Render the per-class aggregate. Never includes query text: this is the + * only surface where control/other statements appear. + * + * @param rows - Class summary rows. + * @returns Table string, or a placeholder when nothing was recorded. + */ +export function formatClassSummary(rows: ClassSummaryRow[]): string { + if (rows.length === 0) return "(no statements recorded)"; + return renderTable( + ["class", "statements", "calls", "total ms"], + rows.map((r) => [ + r.statement_class, + r.statements, + r.calls, + formatMs(r.total_exec_time), + ]), + [0], + ); +} + +/** + * Render one top-N statement table (reads or writes only). + * + * @param rows - Statement rows. + * @returns Table string, or a placeholder when the section is empty. + */ +export function renderStatementsTable(rows: StatementRow[]): string { + if (rows.length === 0) return "(no statements recorded)"; + return renderTable( + [ + "calls", + "total ms", + "mean ms", + "rows", + "blk hit", + "blk read", + "blk dirty", + "tmp rd", + "tmp wr", + "wal", + "query", + ], + rows.map((r) => [ + r.calls, + formatMs(r.total_exec_time), + formatMs(r.mean_exec_time), + r.rows, + r.shared_blks_hit, + r.shared_blks_read, + r.shared_blks_dirtied, + r.temp_blks_read, + r.temp_blks_written, + formatBytes(r.wal_bytes), + sanitizeQueryText(r.query), + ]), + [10], + ); +} + +/** + * Assemble the full report text. + * + * @param report - Fetched report data. + * @returns Multi-section report string. + */ +export function formatDbStatsReport(report: DbStatsReport): string { + const window = report.statsReset + ? report.statsReset.toISOString() + : "unknown (stats_reset unavailable)"; + const sections = [ + `pg_stat_statements report (window since ${window})`, + `statement classes\n${formatClassSummary(report.classSummary)}`, + `top reads by total time\n${renderStatementsTable(report.readsByTime)}`, + `top reads by calls\n${renderStatementsTable(report.readsByCalls)}`, + `top writes by total time\n${renderStatementsTable(report.writesByTime)}`, + `top writes by calls\n${renderStatementsTable(report.writesByCalls)}`, + ]; + return sections.join("\n\n"); +} + +/** + * Parse the `--top ` argument. + * + * @param argv - Process arguments after the script path. + * @returns Statements per section (1..200). + * @throws Error on a malformed or out-of-range value. + */ +function parseTopN(argv: string[]): number { + const at = argv.indexOf("--top"); + if (at === -1) return DEFAULT_TOP_N; + const n = Number(argv[at + 1]); + if (!Number.isInteger(n) || n < 1 || n > 200) { + throw new Error( + `--top expects an integer between 1 and 200, got "${argv[at + 1]}"`, + ); + } + return n; +} + +/** + * Fetch one top-N statement list for a class and ordering. + * + * @param sql - Active owner client. + * @param statementClass - `read` or `write`. + * @param orderBy - `total_exec_time` or `calls`. + * @param topN - Row cap. + * @returns Statement rows. + */ +async function fetchTopStatements( + sql: ReturnType, + statementClass: "read" | "write", + orderBy: "total_exec_time" | "calls", + topN: number, +): Promise { + const rows = await sql.unsafe( + `SELECT calls::text, total_exec_time, mean_exec_time, rows::text, + shared_blks_hit::text, shared_blks_read::text, shared_blks_dirtied::text, + temp_blks_read::text, temp_blks_written::text, wal_bytes::text, + left(query, 400) AS query + FROM extensions.pg_stat_statements + WHERE ${DBID_FILTER} AND toplevel AND ${CLASS_CASE} = '${statementClass}' + ORDER BY ${orderBy} DESC + LIMIT $1`, + [topN], + ); + return rows as unknown as StatementRow[]; +} + +/** + * Fetch everything the report renders. + * + * @param sql - Active owner client. + * @param topN - Statements per top-N section. + * @returns The report data. + */ +async function fetchReport( + sql: ReturnType, + topN: number, +): Promise { + const classSummary = (await sql.unsafe( + `SELECT ${CLASS_CASE} AS statement_class, + count(*)::text AS statements, + sum(calls)::text AS calls, + sum(total_exec_time) AS total_exec_time + FROM extensions.pg_stat_statements + WHERE ${DBID_FILTER} AND toplevel + GROUP BY 1 + ORDER BY sum(total_exec_time) DESC`, + )) as unknown as ClassSummaryRow[]; + const [info] = (await sql.unsafe( + `SELECT stats_reset FROM extensions.pg_stat_statements_info`, + )) as unknown as Array<{ stats_reset: Date | null }>; + return { + statsReset: info?.stats_reset ?? null, + classSummary, + readsByTime: await fetchTopStatements(sql, "read", "total_exec_time", topN), + readsByCalls: await fetchTopStatements(sql, "read", "calls", topN), + writesByTime: await fetchTopStatements( + sql, + "write", + "total_exec_time", + topN, + ), + writesByCalls: await fetchTopStatements(sql, "write", "calls", topN), + }; +} + +/** + * Map a query failure to an actionable remediation hint. + * + * @param err - Caught error. + * @returns Hint line, or undefined for unrecognized failures. + */ +function failureHint(err: unknown): string | undefined { + const message = err instanceof Error ? err.message : String(err); + const code = (err as { code?: string }).code; + if (message.includes("shared_preload_libraries")) { + return "pg_stat_statements is installed but not preloaded. Self-host: docker-compose.yml sets shared_preload_libraries; restart the container (docker compose up -d)."; + } + if (code === "42P01" || /pg_stat_statements.*does not exist/i.test(message)) { + return "pg_stat_statements is not installed. Apply docker/extensions.sql: bun run db:rls (self-host) or bun run db:rls:owner (hosted)."; + } + return undefined; +} + +/** + * Fetch and print the report. + * + * @throws Error when the environment or database is not ready. + */ +async function main(): Promise { + const topN = parseTopN(process.argv.slice(2)); + const sql = postgres(ownerUrl(), { max: 1, onnotice: () => undefined }); + try { + console.log(formatDbStatsReport(await fetchReport(sql, topN))); + } finally { + await sql.end({ timeout: 5 }); + } +} + +if (import.meta.main) { + try { + await main(); + } catch (err) { + console.error(err instanceof Error ? err.message : err); + const hint = failureHint(err); + if (hint) console.error(` - ${hint}`); + process.exit(1); + } +} diff --git a/scripts/smoke-workers.ts b/scripts/smoke-workers.ts index dd60110c..dc585f62 100644 --- a/scripts/smoke-workers.ts +++ b/scripts/smoke-workers.ts @@ -138,6 +138,7 @@ function startWorker(): { "--env", "dev", "--local", + "--test-scheduled", "--ip", "127.0.0.1", "--port", @@ -225,6 +226,85 @@ async function runProbes(): Promise { return failures; } +/** + * Reason signature of the scheduled probe's expected failure. The smoke DB + * URLs are unreachable by design, so the sweep dies inside query execution + * and drizzle wraps it as a "Failed query" error naming the purge SELECT. + * Anything else (TypeError, ReferenceError, a dynamic-require bundling + * failure, a missing-frame throw) is unexpected: it fails the probe and + * stays visible to the error grep. + */ +const EXPECTED_HOUSEKEEPING_REASON = + "Failed query: SELECT table_name, row_count FROM public.purge_expired_rows"; + +/** + * Whether a log line is the scheduled probe's expected failure event: an + * `ok:false` `db_housekeeping` entry whose reason matches the unreachable-DB + * query-failure signature. Only lines matching this whitelist are exempt + * from the error grep. + * + * @param line - One worker log line. + * @returns True when the line is the expected unreachable-DB event. + */ +function isExpectedHousekeepingFailure(line: string): boolean { + return ( + line.includes('"event":"db_housekeeping"') && + line.includes('"ok":false') && + line.includes(EXPECTED_HOUSEKEEPING_REASON) + ); +} + +/** + * Trigger the scheduled handler through wrangler's scheduled test endpoint + * and wait for the housekeeping event in the worker log. The DB URLs are + * unreachable by design, so the expected outcome is the handler's own + * structured `ok:false` connection-failure log — which proves the handler, + * the request DB frame, and the raw-builder wiring survive the wrangler + * bundle without touching a database. A `db_housekeeping` failure carrying + * a TypeError fails the probe instead of passing it. + * + * @param output - Reader over everything the worker printed so far. + * @returns Human-readable failure lines, empty when the probe passed. + */ +async function probeScheduled(output: () => string): Promise { + const path = "/cdn-cgi/handler/scheduled"; + try { + const response = await fetch(`${BASE_URL}${path}`, { + signal: AbortSignal.timeout(20_000), + }); + if (response.status !== 200) { + return [ + `GET ${path} returned ${response.status}, expected 200 (the scheduled handler is not wired into the bundle)`, + ]; + } + } catch (error) { + return [`GET ${path} failed: ${String(error)}`]; + } + const deadline = Date.now() + 15_000; + while (Date.now() < deadline) { + const line = output() + .split("\n") + .find( + (entry) => + entry.includes('"event":"db_housekeeping"') && + entry.includes('"ok":false'), + ); + if (line) { + if (!isExpectedHousekeepingFailure(line)) { + return [ + `the scheduled handler logged an unexpected failure reason under the wrangler bundle: ${line.trim().slice(0, 300)}`, + ]; + } + console.log(` ok GET ${path} -> db_housekeeping event logged`); + return []; + } + await Bun.sleep(200); + } + return [ + "the scheduled probe never logged a db_housekeeping event (scheduled handler or DB-frame wiring is broken)", + ]; +} + /** * Build the bundle's verdict: boot it, probe it, and read its log. * @@ -239,6 +319,7 @@ async function main(): Promise { try { await waitForReady(); failures = await runProbes(); + failures.push(...(await probeScheduled(output))); } catch (error) { failures.push(String(error)); } finally { @@ -248,12 +329,17 @@ async function main(): Promise { // A probe can pass while the worker logs a runtime error on another path, // so the log itself is an assertion. This is what generalizes past the one - // bug that prompted the script. - const logged = output(); + // bug that prompted the script. Only the scheduled probe's expected + // unreachable-DB db_housekeeping event is exempt (whitelisted by reason + // signature); any other handler failure stays visible to both the probe + // and this grep. + const lines = output() + .split("\n") + .filter((entry) => !isExpectedHousekeepingFailure(entry)); for (const marker of ["✘ [ERROR]", "TypeError"]) { - if (!logged.includes(marker)) continue; - const line = logged.split("\n").find((entry) => entry.includes(marker)); - failures.push(`worker logged ${marker}: ${line?.trim().slice(0, 300)}`); + const line = lines.find((entry) => entry.includes(marker)); + if (!line) continue; + failures.push(`worker logged ${marker}: ${line.trim().slice(0, 300)}`); } if (failures.length > 0) { diff --git a/scripts/verify-rls.ts b/scripts/verify-rls.ts index 4985775f..55b6a837 100644 --- a/scripts/verify-rls.ts +++ b/scripts/verify-rls.ts @@ -1,7 +1,7 @@ /** * Verify the live database satisfies the public RLS contract: policies, - * FORCE RLS, the owner-managed functions AND triggers, lz4 compression, and - * the append-only REVOKE narrowing. Read-only: runs as the migration role + * FORCE RLS, the owner-managed extensions, functions AND triggers, lz4 + * compression, and the append-only REVOKE narrowing. Read-only: runs as the migration role * (system catalogs are world-readable). Exits non-zero with an actionable * message so a forgotten owner apply, grant drift, or policy drift blocks the * deploy instead of shipping broken RLS. @@ -16,6 +16,7 @@ interface ExpectedContract { forcedTables: string[]; functions: string[]; triggers: Array<{ table: string; trigger: string }>; + extensions: Array<{ name: string; schema: string | null }>; } /** @@ -53,14 +54,16 @@ function readDockerSql(file: string): string { } /** - * Extract the expected policies, FORCE-RLS tables, and SECURITY DEFINER - * function names from the hand-written docker SQL (the single source of truth). + * Extract the expected policies, FORCE-RLS tables, SECURITY DEFINER + * function names, and extensions from the hand-written docker SQL (the + * single source of truth). * * @returns The contract the live database must satisfy. */ function expectedContract(): ExpectedContract { const policiesSql = readDockerSql("rls-policies.sql"); const functionsSql = readDockerSql("rls-functions.sql"); + const extensionsSql = readDockerSql("extensions.sql"); const policies = [ ...policiesSql.matchAll(/CREATE\s+POLICY\s+"([^"]+)"\s+ON\s+"([^"]+)"/gi), @@ -88,11 +91,23 @@ function expectedContract(): ExpectedContract { ), ].map((m) => ({ trigger: m[1], table: m[2] })); + // Both halves are load-bearing: CREATE EXTENSION IF NOT EXISTS does NOT + // relocate an extension that already exists in another schema (it no-ops + // with a notice), and scripts/db-stats.ts addresses the views through the + // declared schema. A name-only check would pass on a database carrying a + // pre-existing public-schema install while the read path 42P01s. + const extensions = [ + ...extensionsSql.matchAll( + /CREATE\s+EXTENSION\s+IF\s+NOT\s+EXISTS\s+(\w+)(?:\s+WITH\s+SCHEMA\s+(\w+))?/gi, + ), + ].map((m) => ({ name: m[1], schema: m[2] ?? null })); + return { policies, forcedTables: [...new Set(forcedTables)], functions: [...new Set(functions)], triggers, + extensions, }; } @@ -144,6 +159,16 @@ async function findMissing( liveTriggers.map((r) => `${r.relname}.${r.tgname}`), ); + const liveExtensions = await sql<{ extname: string; nspname: string }[]>` + SELECT e.extname, n.nspname + FROM pg_extension e + JOIN pg_namespace n ON n.oid = e.extnamespace + `; + const liveExtensionSet = new Set( + liveExtensions.map((r) => `${r.extname}.${r.nspname}`), + ); + const liveExtensionNames = new Set(liveExtensions.map((r) => r.extname)); + const missing: string[] = []; for (const { table, policy } of expected.policies) { if (!livePolicySet.has(policyKey(table, policy))) { @@ -165,6 +190,16 @@ async function findMissing( missing.push(`trigger "${trigger}" on ${table}`); } } + for (const { name, schema } of expected.extensions) { + const present = schema + ? liveExtensionSet.has(`${name}.${schema}`) + : liveExtensionNames.has(name); + if (!present) { + missing.push( + schema ? `extension ${name} in schema ${schema}` : `extension ${name}`, + ); + } + } for (const { role, table, privilege } of REVOKED_PRIVILEGES) { const [row] = await sql<{ granted: boolean | null }[]>` @@ -246,8 +281,11 @@ async function verifyRls(url: string): Promise { throw new Error( `RLS contract not satisfied on the target database:\n${list}\n` + "Apply the owner-managed SQL as the database owner (db:rls:owner) " + - "for missing functions/triggers, re-run the public apply (db:rls:ci) " + - "for grants/policies/compression, then re-run the deploy.", + "for missing extensions/functions/triggers, re-run the public apply " + + "(db:rls:ci) for grants/policies/compression, then re-run the deploy. " + + "An extension already installed in a different schema needs an " + + "owner-run ALTER EXTENSION ... SET SCHEMA instead: CREATE EXTENSION " + + "IF NOT EXISTS cannot relocate it.", ); } } diff --git a/tests/data/housekeeping.test.ts b/tests/data/housekeeping.test.ts new file mode 100644 index 00000000..ce6639b7 --- /dev/null +++ b/tests/data/housekeeping.test.ts @@ -0,0 +1,413 @@ +import { test, expect, describe, afterEach } from "bun:test"; +import { truncateAll } from "@/tests/setup/schema"; +import { superuserPool } from "@/tests/setup/global"; +import { seedUserOrgProject } from "@/tests/setup/seed"; +import { captureAppUserError } from "@/tests/setup/expect-query"; +import { serviceRoleDb } from "@/lib/db/connection"; +import { + purgeExpiredRows, + type PurgeResultRow, +} from "@/lib/db/raw/purge-expired-rows"; + +/** Fixed table order `public.purge_expired_rows` reports in. */ +const TABLE_ORDER = [ + "oauthAccessToken", + "oauthRefreshToken", + "session", + "verification", + "team_invite_code", +]; + +/** + * Index result rows by table name. + * + * @param rows - Sweep result rows. + * @returns Map of table name to row count. + */ +function counts(rows: PurgeResultRow[]): Record { + return Object.fromEntries(rows.map((r) => [r.table_name, r.row_count])); +} + +/** + * Insert a bare user row and return its id. + * + * @param sql - Superuser client. + * @param suffix - Unique suffix for name/email. + * @returns The new user id. + */ +async function insertUser( + sql: ReturnType, + suffix: string, +): Promise { + const [u] = await sql<{ id: string }[]>` + INSERT INTO piyaz_auth."user" ("name", "email", "emailVerified", "updatedAt") + VALUES (${"User " + suffix}, ${"hk-" + suffix + "@test.local"}, true, now()) + RETURNING id + `; + return u.id; +} + +/** + * Insert a bare organization row and return its id. + * + * @param sql - Superuser client. + * @param suffix - Unique suffix for name/slug. + * @returns The new organization id. + */ +async function insertOrg( + sql: ReturnType, + suffix: string, +): Promise { + const [o] = await sql<{ id: string }[]>` + INSERT INTO piyaz_auth."organization" ("name", "slug", "createdAt") + VALUES (${"Org " + suffix}, ${"hk-org-" + suffix}, now()) + RETURNING id + `; + return o.id; +} + +/** + * Insert an oauth access token expiring at `now() + offset`. + * + * @param sql - Superuser client. + * @param expiresOffset - Signed interval, e.g. `-25 hours`. + * @returns The new row id. + */ +async function insertAccessToken( + sql: ReturnType, + expiresOffset: string, +): Promise { + const [t] = await sql<{ id: string }[]>` + INSERT INTO piyaz_auth."oauthAccessToken" ("token", "clientId", "scopes", "expiresAt") + VALUES ('at-' || gen_random_uuid()::text, 'hk-client', '{}', now() + ${expiresOffset}::interval) + RETURNING id + `; + return t.id; +} + +/** + * Insert an oauth refresh token; `revokedOffset` null keeps it unrevoked. + * + * @param sql - Superuser client. + * @param userId - Owning user id. + * @param expiresOffset - Signed interval for `expiresAt`. + * @param revokedOffset - Signed interval for `revoked`, or null. + * @returns The new row id. + */ +async function insertRefreshToken( + sql: ReturnType, + userId: string, + expiresOffset: string, + revokedOffset: string | null = null, +): Promise { + const [t] = await sql<{ id: string }[]>` + INSERT INTO piyaz_auth."oauthRefreshToken" + ("token", "clientId", "scopes", "userId", "expiresAt", "revoked") + VALUES ('rt-' || gen_random_uuid()::text, 'hk-client', '{}', ${userId}, + now() + ${expiresOffset}::interval, + now() + ${revokedOffset}::interval) + RETURNING id + `; + return t.id; +} + +/** + * Insert a session expiring at `now() + offset`. + * + * @param sql - Superuser client. + * @param userId - Owning user id. + * @param expiresOffset - Signed interval for `expiresAt`. + * @returns The new row id. + */ +async function insertSession( + sql: ReturnType, + userId: string, + expiresOffset: string, +): Promise { + const [s] = await sql<{ id: string }[]>` + INSERT INTO piyaz_auth."session" ("expiresAt", "token", "updatedAt", "userId") + VALUES (now() + ${expiresOffset}::interval, 'tok-' || gen_random_uuid()::text, now(), ${userId}) + RETURNING id + `; + return s.id; +} + +/** + * Insert a verification row expiring at `now() + offset`. + * + * @param sql - Superuser client. + * @param expiresOffset - Signed interval for `expiresAt`. + * @returns The new row id. + */ +async function insertVerification( + sql: ReturnType, + expiresOffset: string, +): Promise { + const [v] = await sql<{ id: string }[]>` + INSERT INTO piyaz_auth."verification" ("identifier", "value", "expiresAt") + VALUES ('hk-' || gen_random_uuid()::text, 'v', now() + ${expiresOffset}::interval) + RETURNING id + `; + return v.id; +} + +/** + * Insert a team invite code; null offsets leave the column NULL. + * + * @param sql - Superuser client. + * @param orgId - Owning organization id (one code per org). + * @param revokedOffset - Signed interval for `revoked_at`, or null. + * @param expiresOffset - Signed interval for `expires_at`, or null. + * @returns The new row id. + */ +async function insertInviteCode( + sql: ReturnType, + orgId: string, + revokedOffset: string | null, + expiresOffset: string | null, +): Promise { + const [c] = await sql<{ id: string }[]>` + INSERT INTO public.team_invite_code (organization_id, code, revoked_at, expires_at) + VALUES (${orgId}, 'code-' || gen_random_uuid()::text, + now() + ${revokedOffset}::interval, + now() + ${expiresOffset}::interval) + RETURNING id + `; + return c.id; +} + +/** + * Collect the surviving ids of a swept table. + * + * @param sql - Superuser client. + * @param table - Schema-qualified quoted table name. + * @returns Set of remaining row ids. + */ +async function remainingIds( + sql: ReturnType, + table: string, +): Promise> { + const rows = await sql.unsafe<{ id: string }[]>(`SELECT id FROM ${table}`); + return new Set(rows.map((r) => r.id)); +} + +afterEach(async () => { + await truncateAll(); +}); + +describe("purge_expired_rows boundaries", () => { + test("oauthAccessToken: deletes past the 24h grace, keeps within it", async () => { + const sql = superuserPool(); + await insertAccessToken(sql, "-25 hours"); + const keep = await insertAccessToken(sql, "-23 hours"); + const live = await insertAccessToken(sql, "30 days"); + + const res = await purgeExpiredRows(serviceRoleDb, false, 100); + expect(counts(res).oauthAccessToken).toBe(1); + expect(await remainingIds(sql, 'piyaz_auth."oauthAccessToken"')).toEqual( + new Set([keep, live]), + ); + }); + + test("oauthRefreshToken: revoked OR expired past the 24h grace", async () => { + const sql = superuserPool(); + const userId = await insertUser(sql, "refresh"); + await insertRefreshToken(sql, userId, "30 days", "-25 hours"); + const keepRevoked = await insertRefreshToken( + sql, + userId, + "30 days", + "-23 hours", + ); + await insertRefreshToken(sql, userId, "-25 hours"); + const live = await insertRefreshToken(sql, userId, "30 days"); + + const res = await purgeExpiredRows(serviceRoleDb, false, 100); + expect(counts(res).oauthRefreshToken).toBe(2); + expect(await remainingIds(sql, 'piyaz_auth."oauthRefreshToken"')).toEqual( + new Set([keepRevoked, live]), + ); + }); + + test("session: deletes past the 7d grace, keeps within it", async () => { + const sql = superuserPool(); + const userId = await insertUser(sql, "session"); + await insertSession(sql, userId, "-8 days"); + const keep = await insertSession(sql, userId, "-6 days"); + const live = await insertSession(sql, userId, "7 days"); + + const res = await purgeExpiredRows(serviceRoleDb, false, 100); + expect(counts(res).session).toBe(1); + expect(await remainingIds(sql, 'piyaz_auth."session"')).toEqual( + new Set([keep, live]), + ); + }); + + test("verification: deletes past the 7d grace, keeps within it", async () => { + const sql = superuserPool(); + await insertVerification(sql, "-8 days"); + const keep = await insertVerification(sql, "-6 days"); + const live = await insertVerification(sql, "7 days"); + + const res = await purgeExpiredRows(serviceRoleDb, false, 100); + expect(counts(res).verification).toBe(1); + expect(await remainingIds(sql, 'piyaz_auth."verification"')).toEqual( + new Set([keep, live]), + ); + }); + + test("team_invite_code: revoked or expired past the 30d grace", async () => { + const sql = superuserPool(); + await insertInviteCode(sql, await insertOrg(sql, "ic1"), "-31 days", null); + const keepRevoked = await insertInviteCode( + sql, + await insertOrg(sql, "ic2"), + "-29 days", + null, + ); + await insertInviteCode(sql, await insertOrg(sql, "ic3"), null, "-31 days"); + const active = await insertInviteCode( + sql, + await insertOrg(sql, "ic4"), + null, + "30 days", + ); + + const res = await purgeExpiredRows(serviceRoleDb, false, 100); + expect(counts(res).team_invite_code).toBe(2); + expect(await remainingIds(sql, "public.team_invite_code")).toEqual( + new Set([keepRevoked, active]), + ); + }); +}); + +describe("purge_expired_rows modes", () => { + test("dry run reports would-delete counts and mutates nothing", async () => { + const sql = superuserPool(); + const at = await insertAccessToken(sql, "-25 hours"); + const userId = await insertUser(sql, "dry"); + const rt = await insertRefreshToken(sql, userId, "-25 hours"); + const s = await insertSession(sql, userId, "-8 days"); + const v = await insertVerification(sql, "-8 days"); + const ic = await insertInviteCode( + sql, + await insertOrg(sql, "dry"), + "-31 days", + null, + ); + + const res = await purgeExpiredRows(serviceRoleDb, true, 100); + expect(counts(res)).toEqual({ + oauthAccessToken: 1, + oauthRefreshToken: 1, + session: 1, + verification: 1, + team_invite_code: 1, + }); + expect(await remainingIds(sql, 'piyaz_auth."oauthAccessToken"')).toEqual( + new Set([at]), + ); + expect(await remainingIds(sql, 'piyaz_auth."oauthRefreshToken"')).toEqual( + new Set([rt]), + ); + expect(await remainingIds(sql, 'piyaz_auth."session"')).toEqual( + new Set([s]), + ); + expect(await remainingIds(sql, 'piyaz_auth."verification"')).toEqual( + new Set([v]), + ); + expect(await remainingIds(sql, "public.team_invite_code")).toEqual( + new Set([ic]), + ); + }); + + test("batch limit caps each run and leftovers drain on the next", async () => { + const sql = superuserPool(); + for (let i = 0; i < 5; i++) await insertAccessToken(sql, "-25 hours"); + + const dry = await purgeExpiredRows(serviceRoleDb, true, 3); + expect(counts(dry).oauthAccessToken).toBe(3); + expect( + (await remainingIds(sql, 'piyaz_auth."oauthAccessToken"')).size, + ).toBe(5); + + const first = await purgeExpiredRows(serviceRoleDb, false, 3); + expect(counts(first).oauthAccessToken).toBe(3); + expect( + (await remainingIds(sql, 'piyaz_auth."oauthAccessToken"')).size, + ).toBe(2); + + const second = await purgeExpiredRows(serviceRoleDb, false, 3); + expect(counts(second).oauthAccessToken).toBe(2); + expect( + (await remainingIds(sql, 'piyaz_auth."oauthAccessToken"')).size, + ).toBe(0); + }); + + test("returns one zero row per table in fixed order on a clean database", async () => { + const res = await purgeExpiredRows(serviceRoleDb, false, 10); + expect(res.map((r) => r.table_name)).toEqual(TABLE_ORDER); + expect(res.every((r) => r.row_count === 0)).toBe(true); + }); + + test("rejects an out-of-range batch limit", async () => { + let caught: unknown; + try { + await purgeExpiredRows(serviceRoleDb, true, 0); + } catch (err) { + caught = err; + } + expect(caught).toBeDefined(); + const cause = (caught as { cause?: unknown }).cause; + expect(String(cause ?? caught)).toMatch(/out of range/); + }); +}); + +describe("purge_expired_rows exclusions", () => { + test("retains legal, activity, consent, and invitation rows regardless of age", async () => { + const sql = superuserPool(); + const f = await seedUserOrgProject("hk-retained", { legalCurrent: false }); + const [acceptance] = await sql<{ id: string }[]>` + INSERT INTO public.legal_acceptances ("user_id", "document_type", "document_version", "accepted_at") + VALUES (${f.userId}, 'terms', '2020-01-01', now() - interval '400 days') + RETURNING id + `; + const [event] = await sql<{ id: string }[]>` + INSERT INTO public.activity_events ("project_id", "type", "source", "summary", "created_at") + VALUES (${f.projectId}, 'task_created', 'web', 'hk seed', now() - interval '400 days') + RETURNING id + `; + const [consent] = await sql<{ id: string }[]>` + INSERT INTO piyaz_auth."oauthConsent" ("clientId", "userId", "referenceId", "scopes", "createdAt") + VALUES ('hk-client', ${f.userId}, ${f.organizationId}, '{}', now() - interval '400 days') + RETURNING id + `; + const [invite] = await sql<{ id: string }[]>` + INSERT INTO piyaz_auth."invitation" ("organizationId", "email", "status", "expiresAt", "inviterId") + VALUES (${f.organizationId}, 'hk-retained@test.local', 'pending', now() - interval '400 days', ${f.userId}) + RETURNING id + `; + + await purgeExpiredRows(serviceRoleDb, false, 1000); + const legal = + await sql`SELECT id FROM public.legal_acceptances WHERE id = ${acceptance.id}`; + const activity = + await sql`SELECT id FROM public.activity_events WHERE id = ${event.id}`; + const consents = + await sql`SELECT id FROM piyaz_auth."oauthConsent" WHERE id = ${consent.id}`; + const invites = + await sql`SELECT id FROM piyaz_auth."invitation" WHERE id = ${invite.id}`; + expect(legal.length).toBe(1); + expect(activity.length).toBe(1); + expect(consents.length).toBe(1); + expect(invites.length).toBe(1); + }); + + test("app_user cannot execute the sweep", async () => { + const f = await seedUserOrgProject("hk-priv", { legalCurrent: false }); + const { code } = await captureAppUserError(f.userId, async (tx) => { + await tx`SELECT * FROM public.purge_expired_rows(true, 10)`; + }); + expect(code).toBe("42501"); + }); +}); diff --git a/tests/scripts/db-stats.test.ts b/tests/scripts/db-stats.test.ts new file mode 100644 index 00000000..b8510483 --- /dev/null +++ b/tests/scripts/db-stats.test.ts @@ -0,0 +1,209 @@ +import { describe, expect, test } from "bun:test"; +import { + formatBytes, + formatClassSummary, + formatDbStatsReport, + formatMs, + renderStatementsTable, + renderTable, + sanitizeQueryText, + type ClassSummaryRow, + type DbStatsReport, + type StatementRow, +} from "../../scripts/db-stats"; + +/** + * Build a statement row with overridable fields. + * + * @param overrides - Field overrides. + * @returns A complete statement row. + */ +function statement(overrides: Partial = {}): StatementRow { + return { + calls: "10", + total_exec_time: 123.456, + mean_exec_time: 12.3456, + rows: "100", + shared_blks_hit: "5", + shared_blks_read: "2", + shared_blks_dirtied: "1", + temp_blks_read: "0", + temp_blks_written: "0", + wal_bytes: "2048", + query: "SELECT * FROM tasks WHERE id = $1", + ...overrides, + }; +} + +describe("sanitizeQueryText", () => { + test("collapses whitespace runs and newlines to single spaces", () => { + expect(sanitizeQueryText("SELECT *\n FROM tasks\t WHERE id = $1")).toBe( + "SELECT * FROM tasks WHERE id = $1", + ); + }); + + test("truncates with an ellipsis at the requested width", () => { + const out = sanitizeQueryText( + "SELECT column_a, column_b FROM somewhere", + 20, + ); + expect(out.length).toBe(20); + expect(out.endsWith("…")).toBe(true); + expect(out.startsWith("SELECT column_a")).toBe(true); + }); + + test("leaves short text untouched", () => { + expect(sanitizeQueryText("SELECT 1", 20)).toBe("SELECT 1"); + }); +}); + +describe("formatMs", () => { + test("renders two decimals", () => { + expect(formatMs(0)).toBe("0.00"); + expect(formatMs(1234.567)).toBe("1234.57"); + }); +}); + +describe("formatBytes", () => { + test("keeps small counts in bytes", () => { + expect(formatBytes(0)).toBe("0 B"); + expect(formatBytes("1023")).toBe("1023 B"); + }); + + test("scales binary units with one decimal", () => { + expect(formatBytes(1024)).toBe("1.0 KiB"); + expect(formatBytes(1536)).toBe("1.5 KiB"); + expect(formatBytes(5 * 1024 * 1024)).toBe("5.0 MiB"); + }); + + test("accepts driver string input", () => { + expect(formatBytes("2048")).toBe("2.0 KiB"); + }); + + test("passes non-numeric input through", () => { + expect(formatBytes("n/a")).toBe("n/a"); + }); +}); + +describe("renderTable", () => { + test("pads columns and honors alignment", () => { + const out = renderTable( + ["name", "count"], + [ + ["a", "1"], + ["long", "100"], + ], + [0], + ); + const lines = out.split("\n"); + expect(lines[0]).toBe("name count"); + expect(lines[1]).toBe("a 1"); + expect(lines[2]).toBe("long 100"); + }); +}); + +describe("formatClassSummary", () => { + const rows: ClassSummaryRow[] = [ + { + statement_class: "read", + statements: "12", + calls: "3400", + total_exec_time: 900.5, + }, + { + statement_class: "control", + statements: "3", + calls: "5000", + total_exec_time: 42.1, + }, + ]; + + test("renders one aggregate row per class without query text", () => { + const out = formatClassSummary(rows); + expect(out).toContain("read"); + expect(out).toContain("control"); + expect(out).toContain("900.50"); + expect(out).not.toContain("SELECT"); + expect(out).not.toContain("SET"); + }); + + test("renders a placeholder when nothing was recorded", () => { + expect(formatClassSummary([])).toBe("(no statements recorded)"); + }); +}); + +describe("renderStatementsTable", () => { + test("emits exactly the required columns", () => { + const out = renderStatementsTable([statement()]); + const header = out.split("\n")[0]; + for (const column of [ + "calls", + "total ms", + "mean ms", + "rows", + "blk hit", + "blk read", + "blk dirty", + "tmp rd", + "tmp wr", + "wal", + "query", + ]) { + expect(header).toContain(column); + } + expect(out).toContain("SELECT * FROM tasks WHERE id = $1"); + expect(out).toContain("2.0 KiB"); + }); + + test("renders a placeholder when the section is empty", () => { + expect(renderStatementsTable([])).toBe("(no statements recorded)"); + }); +}); + +describe("formatDbStatsReport", () => { + const report: DbStatsReport = { + statsReset: new Date("2026-08-01T03:00:00Z"), + classSummary: [ + { + statement_class: "control", + statements: "2", + calls: "900", + total_exec_time: 10, + }, + ], + readsByTime: [statement()], + readsByCalls: [], + writesByTime: [statement({ query: "UPDATE tasks SET title = $1" })], + writesByCalls: [], + }; + + test("always leads with the stats_reset observation window", () => { + expect(formatDbStatsReport(report)).toContain( + "window since 2026-08-01T03:00:00.000Z", + ); + expect(formatDbStatsReport({ ...report, statsReset: null })).toContain( + "window since unknown (stats_reset unavailable)", + ); + }); + + test("includes every section header and empty placeholders", () => { + const out = formatDbStatsReport(report); + for (const section of [ + "statement classes", + "top reads by total time", + "top reads by calls", + "top writes by total time", + "top writes by calls", + ]) { + expect(out).toContain(section); + } + expect(out).toContain("(no statements recorded)"); + }); + + test("prints query text only inside read/write sections", () => { + const out = formatDbStatsReport(report); + const classSection = out.split("top reads by total time")[0]; + expect(classSection).not.toContain("SELECT * FROM tasks"); + expect(out).toContain("UPDATE tasks SET title = $1"); + }); +}); diff --git a/tests/setup/migrate.ts b/tests/setup/migrate.ts index d97789ee..1f5a2111 100644 --- a/tests/setup/migrate.ts +++ b/tests/setup/migrate.ts @@ -36,6 +36,26 @@ async function provisionRoles(sql: ReturnType): Promise { await sql.unsafe(`REVOKE TEMPORARY ON DATABASE "${db}" FROM PUBLIC`); } +/** + * Apply `docker/extensions.sql` so the testcontainer carries the same + * extensions the owner apply (`db:rls:owner`) and self-host chain + * (`db:rls`) install. CREATE EXTENSION needs no preload and no tables, so + * it runs right after role provisioning; the test compose deliberately + * skips `shared_preload_libraries` because no test queries + * `pg_stat_statements`. Idempotent. + * + * @param sql - Active postgres client (must be the container superuser). + */ +async function applyExtensions( + sql: ReturnType, +): Promise { + const content = readFileSync( + join(process.cwd(), "docker", "extensions.sql"), + "utf8", + ); + await sql.unsafe(content); +} + /** * Apply `docker/grants.sql` (public) and `docker/grants-auth.sql` (piyaz_auth) * after `drizzle-kit push` so the `GRANT … ON ALL TABLES` statements land on @@ -171,6 +191,7 @@ export async function applyMigrations(url: string): Promise { // reset it so subsequent statements land in the `public` schema. await sql.unsafe("SET search_path TO public, piyaz_auth"); await provisionRoles(sql); + await applyExtensions(sql); } finally { await sql.end({ timeout: 5 }); } diff --git a/worker-cf.ts b/worker-cf.ts index c8d6fb87..1afe6175 100644 --- a/worker-cf.ts +++ b/worker-cf.ts @@ -18,10 +18,15 @@ import { CloudflareRateLimitBackend, type CloudflareRateLimitBinding, } from "./lib/api/rate-limit-cf"; +import { + HOUSEKEEPING_BATCH_LIMIT, + purgeExpiredRows, +} from "./lib/db/raw/purge-expired-rows"; import { scheduleRequestDbTeardown, withRequestDb, } from "./lib/db/request-scope.workers"; +import { requestDbStore } from "./lib/db/request-store"; import { broker, type DurableObjectNamespace, @@ -51,6 +56,7 @@ interface WorkerEnv { DATABASE_URL?: string; DATABASE_AUTH_URL?: string; DATABASE_SERVICE_ROLE_URL?: string; + HOUSEKEEPING_DRY_RUN?: string; RATE_LIMIT_API?: CloudflareRateLimitBinding; RATE_LIMIT_AUTH?: CloudflareRateLimitBinding; RATE_LIMIT_MCP?: CloudflareRateLimitBinding; @@ -68,6 +74,16 @@ interface WorkerCtx { passThroughOnException(): void; } +/** + * Minimal `ScheduledController` shape passed by workerd to `scheduled`. + * Same file-local-stub rationale as `WorkerEnv` above. + */ +interface ScheduledController { + cron: string; + scheduledTime: number; + noRetry(): void; +} + /** * Shape of the OpenNext-generated default export this entry delegates to. * The generated module carries no usable types (imported under @@ -288,6 +304,70 @@ const handler = { ctx.waitUntil(promise), ); }, + + /** + * Run the nightly housekeeping sweep. `wrangler.jsonc` owns the schedule + * (`triggers.crons`, one entry per env); the handler sweeps on whatever + * cron fired and logs it, so editing the schedule there cannot silently + * strand the job. Dispatch on `controller.cron` only when a second + * schedule actually lands. + * + * The sweep calls `public.purge_expired_rows` (SECURITY DEFINER) through + * the request frame's service-role handle, read via `requestDbStore` + * instead of `lib/db/connection.ts` — importing the latter would drag the + * Node postgres-js driver into this esbuild bundle (see + * `lib/db/request-store.ts`). Dry-run is fail-safe: anything but the + * literal `"false"` binding stays dry. Failures are logged as a + * structured `db_housekeeping` event with `ok:false` and not rethrown — + * cron invocations have no retry semantics, and the persisted Workers + * Logs are the alerting surface. Teardown is awaited directly: there is + * no response body to outlive. + * + * @param controller - Scheduled-event controller; `cron` is logged as fired. + * @param env - Worker bindings. + * @param _ctx - Execution context (unused; nothing outlives the handler). + */ + async scheduled( + controller: ScheduledController, + env: WorkerEnv, + _ctx: WorkerCtx, + ): Promise { + const startedAt = Date.now(); + const dryRun = env.HOUSEKEEPING_DRY_RUN !== "false"; + try { + const { result, teardown } = await withRequestDb(() => { + const frame = requestDbStore.getStore(); + if (!frame) throw new Error("db_housekeeping: no request DB frame"); + return purgeExpiredRows( + frame.serviceRoleDb, + dryRun, + HOUSEKEEPING_BATCH_LIMIT, + ); + }, dbBindings(env)); + await teardown(); + console.log( + JSON.stringify({ + event: "db_housekeeping", + cron: controller.cron, + dryRun, + ...Object.fromEntries(result.map((r) => [r.table_name, r.row_count])), + ms: Date.now() - startedAt, + ok: true, + }), + ); + } catch (err) { + console.error( + JSON.stringify({ + event: "db_housekeeping", + cron: controller.cron, + dryRun, + ms: Date.now() - startedAt, + ok: false, + reason: String(err), + }), + ); + } + }, }; export default handler; diff --git a/wrangler.jsonc b/wrangler.jsonc index a4b4e464..8d11b9f1 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -92,6 +92,7 @@ } ], "routes": [{ "pattern": "app.piyaz.ai", "custom_domain": true }], + "triggers": { "crons": ["0 3 * * *"] }, "workers_dev": false, "preview_urls": false, "limits": { "cpu_ms": 10000 }, @@ -128,6 +129,7 @@ "EMAIL_FROM_INFO": "noreply@piyaz.ai", "EMAIL_REPLY_TO": "hello@piyaz.ai", "REQUIRE_EMAIL_VERIFICATION": "true", + "HOUSEKEEPING_DRY_RUN": "true", "APP_NAME": "Piyaz", "BRAND_LOGO_URL": "https://app.piyaz.ai/piyaz-mark.png", "BRAND_COLOR": "#976b68", @@ -209,6 +211,7 @@ } ], "routes": [{ "pattern": "dev.app.piyaz.ai", "custom_domain": true }], + "triggers": { "crons": ["0 3 * * *"] }, "workers_dev": false, "preview_urls": false, "limits": { "cpu_ms": 10000 }, @@ -245,6 +248,7 @@ "EMAIL_FROM_INFO": "noreply@piyaz.ai", "EMAIL_REPLY_TO": "hello@piyaz.ai", "REQUIRE_EMAIL_VERIFICATION": "true", + "HOUSEKEEPING_DRY_RUN": "true", "APP_NAME": "Piyaz (dev)", "BRAND_LOGO_URL": "https://app.piyaz.ai/piyaz-mark.png", "BRAND_COLOR": "#976b68", From 2e1a4a42ecca47481f2a021a31f0edd120f1ffb8 Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Tue, 4 Aug 2026 21:02:21 +0000 Subject: [PATCH 2/2] chore(deps): bump actions/stale Bumps the all-actions group with 1 update in the / directory: [actions/stale](https://github.com/actions/stale). Updates `actions/stale` from 10 to 11 - [Release notes](https://github.com/actions/stale/releases) - [Changelog](https://github.com/actions/stale/blob/main/CHANGELOG.md) - [Commits](https://github.com/actions/stale/compare/v10...v11) --- updated-dependencies: - dependency-name: actions/stale dependency-version: '11' dependency-type: direct:production update-type: version-update:semver-major dependency-group: all-actions ... Signed-off-by: dependabot[bot] --- .github/workflows/stale.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/stale.yml b/.github/workflows/stale.yml index e2c6a4d0..5b5d3669 100644 --- a/.github/workflows/stale.yml +++ b/.github/workflows/stale.yml @@ -13,7 +13,7 @@ jobs: stale: runs-on: ubuntu-latest steps: - - uses: actions/stale@v10 + - uses: actions/stale@v11 with: days-before-issue-stale: 60 days-before-issue-close: 14