From d67a00c1d590d401450d7722e0f0988259f21c71 Mon Sep 17 00:00:00 2001 From: Andrey Antukh Date: Thu, 1 Oct 2026 07:19:21 +0200 Subject: [PATCH] :sparkles: Normalize storage metadata with a closed schema (#11987) * :sparkles: Add Malli schema for storage metadata with dual decode Phase 1 of the storage_object.metadata migration: reads accept both Transit and plain JSON (sniffed by the marker) and always return the normalized shape; writes validate against a closed per-bucket Malli schema and still serialize as Transit unless the new :storage-metadata-as-json config flag is set. The 0155 migration normalizes existing rows inside Transit (reference to bucket, default bucket, drop of chunk leftovers) and is idempotent; large instances should fake it and run the batched script instead. AI-assisted-by: muse-spark-1.3-contributor * :recycle: Address review findings on storage metadata Phase 1 Collapse the dead :reference leg of the gc-touched bucket fallback (the decode always sets :bucket on non-nil metadata, so it is only reachable with a NULL column) and fix its comment. Pin the write flag off in the transit-assuming metadata tests so the suite proves the same with the flag set, and add coverage for the flag rollback contract, JSON hash survival, NULL metadata in gc-touched, and the 0155 normalization statements. AI-assisted-by: muse-spark-1.3-contributor * :recycle: Backfill NULLs, canonical buckets, comment fix Backfill NULL metadata columns in 0155 via coalesce (the key-missing rule already matches them), derive valid-buckets from the Malli schema dispatch entries so the list lives in one place, and correct the lookup-bucket fallback comment to NULL columns. AI-assisted-by: muse-spark-1.3-contributor * :recycle: Defer corrupt metadata rows in storage gc-touched Decode touched rows individually so one non-map metadata value no longer aborts the whole chunk: corrupt rows are logged and deferred exactly one day in the same transaction, keeping their metadata intact for a later repair, while healthy rows process normally. AI-assisted-by: muse-spark-1.3-contributor * :recycle: Address storage metadata phase 1 review findings Address the review findings on the storage metadata phase 1 branch: - Fix put-and-delete-object: it stored the object with ::sto/expired-at, so the row was already deleted and del-object! returned false. Add delete-expired-object-returns-false to keep the expired-delete case covered. - Cache the Malli decoder and encoder per process. Building them compiles the closed multi-dispatch schema, and decode-metadata runs on every read path (get-object, dedup probes, GC batches). - Catch Exception instead of Throwable in try-decode-row so JVM Errors are not deferred as corrupt metadata. - Add penpot_storage_gc_poison_total, emitted from storage-gc-touched; wire ::mtx/metrics into its handler. - Anchor the encoding sniff to the start of the document so a plain JSON value that begins with a Transit-looking prefix is not read as Transit. - Cover every bucket on both encodings, a JSON roundtrip through the jsonb column, nil metadata, the canonical bucket set and a poison-only GC chunk. - Rename private check-metadata! to check-metadata. AI-assisted-by: deepseek-v4.1-flash * :recycle: Simplify the storage metadata schema to a single map Replace the per-bucket :multi dispatch with a single closed map: the bucket is validated with ::sm/one-of over metadata-buckets (now a plain set) and the remaining keys are typed optional fields. Per-bucket enforcement shrinks to a one-line :fn guard requiring :file-id and :id for file-data, whose ids the GC reads to resolve references. - Drop the dead (sm/register! ::metadata ...): nothing references the schema by keyword. - Define tempfile-bucket and upload-session-bucket in the schema and alias them from app.storage, removing duplicated literals. - Keep content-type required and the map closed, so an unknown bucket or key still fails fast on write. AI-assisted-by: deepseek-v4.1-flash * :recycle: Drop input coercion from encode-metadata encode-metadata no longer runs the json-transformer decoder before validation. On the write path its only effect was coercing string UUIDs to UUID, and every producer already passes native UUIDs (the RPC profile-id, uuid/random, or binfile ids decoded as ::sm/uuid). Reads keep decoding, so stored Transit or JSON values still come back as native types. - Replace encode-accepts-string-uuids with encode-rejects-string-uuids, pinning the stricter contract. - Pass native UUIDs in encode-writes-plain-json-with-flag. AI-assisted-by: deepseek-v4.1-flash * :books: Document each statement in the storage metadata migration Move the per-statement rules out of the header and add a comment to each UPDATE explaining what it does: drop chunk leftovers, promote the legacy "~:reference" to "~:bucket", drop residual "~:reference", and backfill the default bucket. The header keeps the scope, the encoding note, the `->` vs `?` note and the large-instance warning. AI-assisted-by: deepseek-v4.1-flash * :books: Unwrap wrapped lines in the backend storage memory One line per bullet or paragraph, as mem:memory-maintenance requires. Only formatting; no content change. AI-assisted-by: deepseek-v4.1-flash * :recycle: Defer storage GC poison rows in their own transaction process-chunk! no longer takes poison-ids; it only processes the healthy chunk. The deferral moves to defer-poison! and process-touched! runs it in its own transaction, separate from the freeze/delete work. The loop still drains while there is chunk or poison, so a batch made only of poison rows does not leave healthy rows behind the LIMIT 10 waiting for the next run. Add a regression test: ten poison rows plus one healthy row with a later touched_at are all handled in the same run. AI-assisted-by: deepseek-v4.1-flash * :recycle: Declare per-bucket metadata requirements in one map Replace the file-data-specific predicate with bucket-requirements, a map from bucket to the extra keys it must carry. metadata-buckets is derived from its keys and a single generic :fn enforces presence, so a new bucket and its contract are one entry. organization now requires :organization-id; file-data keeps requiring :file-id and :id. Update the http-assets test helper to set organization-id for its organization objects. AI-assisted-by: deepseek-v4.1-flash * :recycle: Drop the ! suffix from storage GC helpers Rename the internal helpers in app.storage.gc-touched (process-chunk, defer-poison, mark-freeze-in-bulk, ...) to drop the trailing !. AI-assisted-by: deepseek-v4.1-flash * :recycle: Drop the ! suffix from storage GC deleted helpers Rename the internal helpers in app.storage.gc-deleted (clean-deleted, delete-sobjects, delete-give-up, ...) to drop the trailing !. AI-assisted-by: deepseek-v4.1-flash * :bug: Fix dedup lookup for JSON-encoded storage metadata get-database-object-by-hash only matched the Transit keys, so once the :storage-metadata-as-json flag wrote plain JSON rows the dedup stopped finding them and duplicated blobs. Match both encodings with a UNION ALL of two indexable branches. - Add migration 0156 with the plain-key dedup index; the legacy 0068 index stays until Transit support is removed. - Cover it with a JSON dedup test and a Transit -> JSON cross test. AI-assisted-by: deepseek-v4.1-flash --- .serena/memories/backend/storage.md | 10 +- backend/src/app/config.clj | 4 + backend/src/app/main.clj | 9 +- backend/src/app/migrations.clj | 8 +- ...0155-normalize-storage-object-metadata.sql | 52 ++++ ...56-add-storage-object-json-dedup-index.sql | 11 + backend/src/app/storage.clj | 62 ++-- backend/src/app/storage/gc_deleted.clj | 24 +- backend/src/app/storage/gc_touched.clj | 111 +++++-- backend/src/app/storage/impl.clj | 4 +- backend/src/app/storage/schema.clj | 149 ++++++++++ .../test/backend_tests/http_assets_test.clj | 12 +- .../backend_tests/rpc_management_test.clj | 12 +- .../backend_tests/storage_metadata_test.clj | 276 ++++++++++++++++++ backend/test/backend_tests/storage_test.clj | 234 ++++++++++++++- 15 files changed, 884 insertions(+), 94 deletions(-) create mode 100644 backend/src/app/migrations/sql/0155-normalize-storage-object-metadata.sql create mode 100644 backend/src/app/migrations/sql/0156-add-storage-object-json-dedup-index.sql create mode 100644 backend/src/app/storage/schema.clj create mode 100644 backend/test/backend_tests/storage_metadata_test.clj diff --git a/.serena/memories/backend/storage.md b/.serena/memories/backend/storage.md index 43c925eeec..4a48cd8c71 100644 --- a/.serena/memories/backend/storage.md +++ b/.serena/memories/backend/storage.md @@ -54,8 +54,7 @@ ### Key Warning (from function notes): -The improved note in `import-storage-objects` and `handle-persistence` warns: -**Do not reuse the main database connection for storage operations within a transaction.** The storage upload process can fail mid-operation, leaving orphaned objects on the backend. If the outer transaction aborts, pending storage objects become unreconciliable because the storage subsystem registers its pending state in separate transactions. +The improved note in `import-storage-objects` and `handle-persistence` warns: **Do not reuse the main database connection for storage operations within a transaction.** The storage upload process can fail mid-operation, leaving orphaned objects on the backend. If the outer transaction aborts, pending storage objects become unreconciliable because the storage subsystem registers its pending state in separate transactions. ### Rule of Thumb for `sto/put-object!`: @@ -66,10 +65,7 @@ Since `put-object!` uses backend-specific operations (`impl/resolve-backend` + ` - Deduplication requires `::sto/deduplicate?`, a content hash, and bucket metadata. - The lookup matches hash, bucket, backend, and `deleted_at IS NULL`. - The lookup only considers rows with `status='valid'`; pending rows are invisible. -- A hit whose blob is missing is repaired in place: the same row/id is kept, - and `put-object!` rewrites the blob under that id. This heals all existing - references to the object. If the rewrite fails, the row is left live and - valid for a later retry. +- A hit whose blob is missing is repaired in place: the same row/id is kept, and `put-object!` rewrites the blob under that id. This heals all existing references to the object. If the rewrite fails, the row is left live and valid for a later retry. - The lookup does not include file ID, profile ID, team ID, or organization ID. - Objects can therefore share content across users and files within one bucket. - Deleted objects are not reused. @@ -130,6 +126,8 @@ Since `put-object!` uses backend-specific operations (`impl/resolve-backend` + ` - Logical storage operations (`app.storage`, `::mtx/metrics` required by the schema): - `penpot_storage_operations_total{op,bucket,backend}` — `put`, `repair`, `get-data`, `get-bytes`, `del`, `touch`, `exists`. All ops are success-only: `put`/`repair` emit after the backend write, `get-*` after the backend fetch opens, `touch`/`del` only when a row actually changed. `touch-object!`/`del-object!` take the object id (UUID) only — no object overload. Labels come from the updated row itself via `UPDATE ... RETURNING id, backend, metadata` (no extra `SELECT`); with no row matched they emit nothing. `del-object!` only matches live rows (`deleted_at IS NULL`): a repeated del returns `false` and emits nothing. Post-open stream read errors stay counted as attempts. `del` only marks `deleted_at`; physical deletion is a GC concern. `exists` is emitted per deduplication-hit probe, always paired with a `hit`/`repair` outcome (never on probe failure), not per user-facing existence check. - `penpot_storage_dedup_total{result,bucket}` — `hit`, `miss`, `repair`, `skip`. +- Storage GC poison (`app.storage.gc-touched`, `::mtx/metrics` required in the handler cfg): + - `penpot_storage_gc_poison_total` — touched rows whose metadata fails to decode. Each one is logged at error level and deferred exactly one day instead of aborting the chunk. No labels; incremented once per row per scan. - Asset serving (`app.http.assets`, `::mtx/metrics` required in the handler cfg): - `penpot_storage_asset_requests_total{route,backend,bucket,result}` — `route` is `by-id`, `by-file-media-id`, or `thumbnail`; `result` is `served` (<400), `not-found` (404), `unauthorized` (401/403), or `error` (everything else, including a nil/non-number status: every serve path must set `::yres/status`). Serve-path exceptions are counted by `serve-object-measured` and then rethrown. Permission-denied file-media requests and tempfile ownership mismatches both answer HTTP 404 (to avoid leaking existence) but are counted as `unauthorized`. Malformed UUIDs raise before any emission point and are never counted. Counts backend requests that trigger a browser GET to the object store (one per cache miss), so it is a proxy for object GETs, not an exact count. - The physical and logical counters intentionally overlap in coverage but differ in meaning; do not sum them. diff --git a/backend/src/app/config.clj b/backend/src/app/config.clj index 48081d6430..08c48476b5 100644 --- a/backend/src/app/config.clj +++ b/backend/src/app/config.clj @@ -304,6 +304,10 @@ [:objects-storage-s3-region {:optional true} :keyword] [:objects-storage-s3-endpoint {:optional true} ::sm/uri] + ;; Write storage_object.metadata as plain JSON instead of + ;; Transit-JSON. Unset by default (Phase 1: keep writing Transit). + [:storage-metadata-as-json {:optional true} ::sm/boolean] + ;; SSRF protection [:ssrf-allowed-hosts {:optional true} [::sm/set :string]] [:ssrf-extra-blocked-cidrs {:optional true} [::sm/set :string]]])) diff --git a/backend/src/app/main.clj b/backend/src/app/main.clj index 4133cc3a8b..de72b9a0ef 100644 --- a/backend/src/app/main.clj +++ b/backend/src/app/main.clj @@ -180,6 +180,12 @@ ::mdef/labels ["route" "backend" "bucket" "result"] ::mdef/type :counter} + :storage-gc-poison + {::mdef/name "penpot_storage_gc_poison_total" + ::mdef/help "Storage objects deferred by storage-gc-touched because their metadata is corrupt." + ::mdef/labels [] + ::mdef/type :counter} + :http-server-dispatch-timing {::mdef/name "penpot_http_server_dispatch_timing" ::mdef/help "Histogram of dispatch handler" @@ -271,7 +277,8 @@ ::sto/storage (ig/ref ::sto/storage)} ::sto.gc-touched/handler - {::db/pool (ig/ref ::db/pool)} + {::db/pool (ig/ref ::db/pool) + ::mtx/metrics (ig/ref ::mtx/metrics)} ::sto.pending-gc/handler {::db/pool (ig/ref ::db/pool) diff --git a/backend/src/app/migrations.clj b/backend/src/app/migrations.clj index fad1f91168..c458f45150 100644 --- a/backend/src/app/migrations.clj +++ b/backend/src/app/migrations.clj @@ -505,7 +505,13 @@ :fn (mg/resource "app/migrations/sql/0153-add-storage-object-status-and-deletion-attempts.sql")} {:name "0154-add-upload-session-chunk-table" - :fn (mg/resource "app/migrations/sql/0154-add-upload-session-chunk-table.sql")}]) + :fn (mg/resource "app/migrations/sql/0154-add-upload-session-chunk-table.sql")} + + {:name "0155-normalize-storage-object-metadata" + :fn (mg/resource "app/migrations/sql/0155-normalize-storage-object-metadata.sql")} + + {:name "0156-add-storage-object-json-dedup-index" + :fn (mg/resource "app/migrations/sql/0156-add-storage-object-json-dedup-index.sql")}]) (defn apply-migrations! [pool name migrations] diff --git a/backend/src/app/migrations/sql/0155-normalize-storage-object-metadata.sql b/backend/src/app/migrations/sql/0155-normalize-storage-object-metadata.sql new file mode 100644 index 0000000000..36ee9b2c97 --- /dev/null +++ b/backend/src/app/migrations/sql/0155-normalize-storage-object-metadata.sql @@ -0,0 +1,52 @@ +--- Normalize storage_object.metadata inside Transit (Phase 1). +--- +--- Brings every row to a canonical shape WITHOUT changing the encoding +--- family (still Transit JSON, "~:" keys), so old and new code keep +--- reading while it runs. Each statement below documents its own rule. +--- +--- NOTE: key-existence is tested with `-> ... IS [NOT] NULL` instead of +--- the `?` operator because the migration runner prepares statements +--- and a bare `?` is parsed as a bind placeholder. +--- +--- Idempotent: re-running changes zero rows. +--- NOTE for large instances (e.g. our production, ~33M rows): DO NOT run +--- this inline. Mark the migration as executed (fake) and run the batched +--- equivalent by hand: .agents/plans/2026-09-11-storage-normalize-phase1.sql + +-- Drop leftovers from the old chunked-upload metadata: no code reads +-- "~:upload-id" / "~:chunk-index" anymore (#11651 moved chunks to +-- upload_session_chunk and objects-gc purges the orphan rows), and the +-- Phase 1 readers already ignore them on decode. +UPDATE storage_object + SET metadata = metadata - '~:upload-id' - '~:chunk-index' + WHERE (metadata -> '~:upload-id') IS NOT NULL + OR (metadata -> '~:chunk-index') IS NOT NULL; + +-- Rows that carry only the legacy "~:reference": promote it to "~:bucket" +-- with the same value. The value comes as a Transit keyword ("~:xxx"), so +-- the "~:" prefix is stripped; plain strings are kept as they are. +UPDATE storage_object + SET metadata = (metadata - '~:reference') + || jsonb_build_object('~:bucket', + CASE WHEN metadata->>'~:reference' LIKE '~:%' + THEN substr(metadata->>'~:reference', 3) + ELSE metadata->>'~:reference' + END) + WHERE (metadata -> '~:reference') IS NOT NULL + AND (metadata -> '~:bucket') IS NULL; + +-- Rows that carry both keys: "~:bucket" wins and the residual +-- "~:reference" is dropped. +UPDATE storage_object + SET metadata = metadata - '~:reference' + WHERE (metadata -> '~:reference') IS NOT NULL + AND (metadata -> '~:bucket') IS NOT NULL; + +-- Rows with neither key get the historical default bucket +-- "file-media-object" (the value `lookup-bucket` used to fall back to). +-- NULL columns match this rule too (`->` on NULL is NULL) and are +-- backfilled by coalesce. +UPDATE storage_object + SET metadata = coalesce(metadata, '{}') || '{"~:bucket": "file-media-object"}' + WHERE (metadata -> '~:bucket') IS NULL + AND (metadata -> '~:reference') IS NULL; diff --git a/backend/src/app/migrations/sql/0156-add-storage-object-json-dedup-index.sql b/backend/src/app/migrations/sql/0156-add-storage-object-json-dedup-index.sql new file mode 100644 index 0000000000..0576318f50 --- /dev/null +++ b/backend/src/app/migrations/sql/0156-add-storage-object-json-dedup-index.sql @@ -0,0 +1,11 @@ +-- Dedup lookup index for the plain-JSON encoding of +-- storage_object.metadata ("hash" / "bucket"). +-- +-- The dedup lookup matches both metadata encodings while Transit rows +-- still exist (a UNION ALL of two indexable branches). The legacy "~:" +-- index (storage_object__hash_backend_bucket__idx, migration 0068) keeps +-- covering the Transit branch until Transit support is removed. + +CREATE INDEX storage_object__hash_backend_bucket_json__idx + ON storage_object ((metadata->>'hash'), (metadata->>'bucket'), backend) + WHERE deleted_at IS NULL AND status = 'valid'; diff --git a/backend/src/app/storage.clj b/backend/src/app/storage.clj index 52163bfc95..25c2d3482a 100644 --- a/backend/src/app/storage.clj +++ b/backend/src/app/storage.clj @@ -20,6 +20,7 @@ [app.storage.fs :as sfs] [app.storage.impl :as impl] [app.storage.s3 :as ss3] + [app.storage.schema :as stsch] [cuerdas.core :as str] [datoteka.fs :as fs] [integrant.core :as ig]) @@ -37,28 +38,18 @@ nil))) (def default-bucket - "file-media-object") + stsch/default-bucket) (def tempfile-bucket "Bucket name for temporary file uploads (10-minute expiry)." - "tempfile") + stsch/tempfile-bucket) (def upload-session-bucket "Bucket name for chunked-upload chunks." - "upload-session") + stsch/upload-session-bucket) (def valid-buckets - #{"file-media-object" - "team-font-variant" - "file-object-thumbnail" - "file-thumbnail" - "profile" - "organization" - tempfile-bucket - upload-session-bucket - "file-data" - "file-data-fragment" - "file-change"}) + stsch/metadata-buckets) ;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;; ;; Storage Module State @@ -118,18 +109,37 @@ params params)) +;; Matches both metadata encodings while Transit rows still exist. The +;; UNION ALL keeps each branch indexable: migration 0068 covers the legacy +;; "~:" keys and 0156 covers the plain JSON keys. Phase 3 drops the legacy +;; branch once no "~:" rows remain. +(def ^:private sql:get-database-object-by-hash + "( + select * from storage_object + where (metadata->>'~:hash') = ? + and (metadata->>'~:bucket') = ? + and backend = ? + and deleted_at is null + and status = 'valid' + limit 1 + ) union all ( + select * from storage_object + where (metadata->>'hash') = ? + and (metadata->>'bucket') = ? + and backend = ? + and deleted_at is null + and status = 'valid' + limit 1 + ) limit 1") + +;; NOTE: metadata is left encoded; row->storage-object is responsible for +;; decoding it. (defn- get-database-object-by-hash [connectable backend bucket hash] - (let [sql (str "select * from storage_object " - " where (metadata->>'~:hash') = ? " - " and (metadata->>'~:bucket') = ? " - " and backend = ?" - " and deleted_at is null" - " and status = 'valid'" - " limit 1")] - ;; NOTE: metadata is left encoded; row->storage-object is - ;; responsible for decoding it. - (db/exec-one! connectable [sql hash bucket (name backend)]))) + (let [backend (name backend)] + (db/exec-one! connectable [sql:get-database-object-by-hash + hash bucket backend + hash bucket backend]))) (defn- promote-object! [storage object] @@ -147,7 +157,7 @@ res)) (defn row->storage-object [res] - (let [mdata (or (some-> (:metadata res) (db/decode-transit-pgobject)) {})] + (let [mdata (or (some-> (:metadata res) (stsch/decode-metadata)) {})] (impl/storage-object (:id res) (:size res) @@ -289,7 +299,7 @@ {:id id :size (impl/get-size content) :backend (name backend) - :metadata (db/tjson mdata) + :metadata (stsch/encode-metadata mdata) :deleted-at expired-at :touched-at touched-at :status "pending"}) diff --git a/backend/src/app/storage/gc_deleted.clj b/backend/src/app/storage/gc_deleted.clj index d9dc908e0a..7048bd4d2e 100644 --- a/backend/src/app/storage/gc_deleted.clj +++ b/backend/src/app/storage/gc_deleted.clj @@ -51,7 +51,7 @@ "DELETE FROM storage_object WHERE id = ANY(?::uuid[])") -(defn- delete-sobjects! +(defn- delete-sobjects [conn ids] (let [ids (db/create-array conn "uuid" ids)] (-> (db/exec-one! conn [sql:delete-sobjects ids]) @@ -61,7 +61,7 @@ "DELETE FROM upload_session_chunk WHERE object_id = ANY(?::uuid[])") -(defn- delete-upload-session-chunks! +(defn- delete-upload-session-chunks "Remove the chunk mappings for the given storage object ids. This must run before the storage_object rows are deleted: the upload_session_chunk foreign keys are ON DELETE NO ACTION." @@ -75,7 +75,7 @@ deleted_at = NOW() + INTERVAL '1 day' WHERE id = ANY(?::uuid[])") -(defn- increment-attempts-and-defer! +(defn- increment-attempts-and-defer [conn ids] (let [ids (db/create-array conn "uuid" ids)] (db/exec-one! conn [sql:increment-attempts-and-defer ids]))) @@ -85,7 +85,7 @@ WHERE id = ANY(?::uuid[]) AND deletion_attempts >= ?") -(defn- delete-give-up! +(defn- delete-give-up [conn ids] (let [ids (db/create-array conn "uuid" ids)] (db/exec-one! conn [sql:delete-give-up ids max-attempts]))) @@ -93,7 +93,7 @@ (defn- process-chunk "Attempt to delete a chunk of storage objects from a specific backend. - This function runs inside the caller's transaction (clean-deleted!) — + This function runs inside the caller's transaction (clean-deleted) — it does NOT open its own transaction. The caller is responsible for ensuring the rows are locked via FOR UPDATE SKIP LOCKED before calling. @@ -121,18 +121,18 @@ ;; storage_object rows (NO ACTION foreign keys). It only affects ;; objects of the upload-session bucket; for any other bucket the ;; delete matches no rows. - (delete-upload-session-chunks! conn ok-ids) - (delete-sobjects! conn ok-ids)) + (delete-upload-session-chunks conn ok-ids) + (delete-sobjects conn ok-ids)) (when (seq fail-ids) - (increment-attempts-and-defer! conn fail-ids) + (increment-attempts-and-defer conn fail-ids) ;; NOTE: same NO ACTION ordering as above: the give-up DELETE below ;; removes storage_object rows, so chunk mappings must go first. ;; Deferred objects keep their rows; only the mapping of a ;; permanently given-up object disappears early, and that object is ;; already deleted-marked. - (delete-upload-session-chunks! conn fail-ids) - (let [given-up (delete-give-up! conn fail-ids)] + (delete-upload-session-chunks conn fail-ids) + (let [given-up (delete-give-up conn fail-ids)] (when (pos? (db/get-update-count given-up)) (l/wrn :hint "giving up on orphan blob after max attempts" :ids fail-ids @@ -160,7 +160,7 @@ [conn size] (db/exec! conn [sql:get-deleted-chunk (ct/now) size])) -(defn- clean-deleted! +(defn- clean-deleted [cfg] (loop [total 0] (let [deleted (db/tx-run! cfg @@ -184,6 +184,6 @@ (defmethod ig/init-key ::handler [_ cfg] (fn [_] - (let [total (clean-deleted! cfg)] + (let [total (clean-deleted cfg)] (l/inf :hint "task finished" :total total) {:deleted total}))) diff --git a/backend/src/app/storage/gc_touched.clj b/backend/src/app/storage/gc_touched.clj index 260739b02f..f73b6024f1 100644 --- a/backend/src/app/storage/gc_touched.clj +++ b/backend/src/app/storage/gc_touched.clj @@ -24,6 +24,7 @@ [app.common.logging :as l] [app.common.time :as ct] [app.db :as db] + [app.metrics :as mtx] [app.storage :as sto] [app.storage.impl :as impl] [integrant.core :as ig])) @@ -95,7 +96,7 @@ SET touched_at = NULL WHERE id = ANY(?::uuid[])") -(defn- mark-freeze-in-bulk! +(defn- mark-freeze-in-bulk [conn ids] (let [ids (db/create-array conn "uuid" ids)] (db/exec-one! conn [sql:mark-freeze-in-bulk ids]))) @@ -106,11 +107,21 @@ touched_at = NULL WHERE id = ANY(?::uuid[])") -(defn- mark-delete-in-bulk! +(defn- mark-delete-in-bulk [conn ids] (let [ids (db/create-array conn "uuid" ids)] (db/exec-one! conn [sql:mark-delete-in-bulk (ct/now) ids]))) +(def ^:private sql:defer-in-bulk + "UPDATE storage_object + SET touched_at = ? + WHERE id = ANY(?::uuid[])") + +(defn- defer-in-bulk + [conn ids timestamp] + (let [ids (db/create-array conn "uuid" ids)] + (db/exec-one! conn [sql:defer-in-bulk timestamp ids]))) + ;; NOTE: A getter that retrieves the key which will be used for group ;; ids; previously we have no value, then we introduced the ;; `:reference` prop, and then it is renamed to `:bucket` and now is @@ -125,12 +136,17 @@ ;; have value, it means :file-media-object. (defn- lookup-bucket - [{:keys [metadata]}] + [{:keys [id metadata]}] (or (some-> metadata :bucket) - (some-> metadata :reference d/name) - sto/default-bucket)) + (do + ;; Only reachable when the metadata column is NULL: the decode + ;; always sets :bucket on non-nil metadata (0155 also backfills + ;; NULL columns). Keep working, but make it visible. + (l/wrn :hint "storage object without bucket metadata, using fallback" + :id (str id)) + sto/default-bucket))) -(defn- process-objects! +(defn- process-objects [conn has-refs? bucket objects] (loop [to-freeze #{} to-delete #{} @@ -148,31 +164,37 @@ :bucket bucket) (recur to-freeze (conj to-delete id) (rest objects)))) (do - (some->> (seq to-freeze) (mark-freeze-in-bulk! conn)) - (some->> (seq to-delete) (mark-delete-in-bulk! conn)) + (some->> (seq to-freeze) (mark-freeze-in-bulk conn)) + (some->> (seq to-delete) (mark-delete-in-bulk conn)) [(count to-freeze) (count to-delete)])))) -(defn- process-bucket! +(defn- process-bucket [conn bucket objects] (cond - (= bucket "file-media-object") (process-objects! conn has-file-media-object-refs? bucket objects) - (= bucket "team-font-variant") (process-objects! conn has-team-font-variant-refs? bucket objects) - (= bucket "file-object-thumbnail") (process-objects! conn has-file-object-thumbnails-refs? bucket objects) - (= bucket "file-thumbnail") (process-objects! conn has-file-thumbnails-refs? bucket objects) - (= bucket "profile") (process-objects! conn has-profile-refs? bucket objects) - (= bucket "file-data") (process-objects! conn has-file-data-refs? bucket objects) - (= bucket sto/tempfile-bucket) (process-objects! conn (constantly false) sto/tempfile-bucket objects) - (= bucket sto/upload-session-bucket) (process-objects! conn (constantly false) sto/upload-session-bucket objects) - (= bucket "organization") (process-objects! conn (constantly false) bucket objects) + (= bucket "file-media-object") (process-objects conn has-file-media-object-refs? bucket objects) + (= bucket "team-font-variant") (process-objects conn has-team-font-variant-refs? bucket objects) + (= bucket "file-object-thumbnail") (process-objects conn has-file-object-thumbnails-refs? bucket objects) + (= bucket "file-thumbnail") (process-objects conn has-file-thumbnails-refs? bucket objects) + (= bucket "profile") (process-objects conn has-profile-refs? bucket objects) + (= bucket "file-data") (process-objects conn has-file-data-refs? bucket objects) + (= bucket sto/tempfile-bucket) (process-objects conn (constantly false) sto/tempfile-bucket objects) + (= bucket sto/upload-session-bucket) (process-objects conn (constantly false) sto/upload-session-bucket objects) + (= bucket "organization") (process-objects conn (constantly false) bucket objects) :else (ex/raise :type :internal :code :unexpected-unknown-reference :hint (dm/fmt "unknown reference '%'" bucket)))) -(defn process-chunk! +(defn- defer-poison + "Defer corrupt rows by one day in their own transaction, separate from + the healthy-chunk work." + [{:keys [::db/conn]} poison-ids] + (defer-in-bulk conn poison-ids (ct/plus (ct/now) {:days 1}))) + +(defn process-chunk [{:keys [::db/conn]} chunk] (reduce-kv (fn [[nfo ndo] bucket objects] - (let [[nfo' ndo'] (process-bucket! conn bucket objects)] + (let [[nfo' ndo'] (process-bucket conn bucket objects)] [(+ nfo nfo') (+ ndo ndo')])) [0 0] @@ -189,21 +211,45 @@ SKIP LOCKED LIMIT 10") +(defn- try-decode-row + "Decode a touched row, capturing corrupt metadata as poison instead + of aborting the whole chunk. Poison rows are deferred by the caller." + [row] + (try + [:ok (impl/decode-row row)] + ;; Exception, not Throwable: JVM Errors (OOM, StackOverflow) must + ;; not be swallowed as a deferrable data problem. + (catch Exception cause + (l/err :hint "storage object with corrupt metadata, deferring evaluation" + :id (str (:id row)) + :cause cause) + [:poison (:id row)]))) + (defn get-chunk [conn timestamp] - (->> (db/exec! conn [sql:get-touched-storage-objects timestamp]) - (map impl/decode-row) - (not-empty))) + (let [grouped (->> (db/exec! conn [sql:get-touched-storage-objects timestamp]) + (map try-decode-row) + (group-by first))] + {:chunk (not-empty (mapv second (:ok grouped))) + :poison (not-empty (mapv second (:poison grouped)))})) -(defn- process-touched! - [{:keys [::db/pool ::timestamp] :as cfg}] +(defn- process-touched + [{:keys [::db/pool ::mtx/metrics ::timestamp] :as cfg}] (loop [freezed 0 deleted 0] - (if-let [chunk (get-chunk pool timestamp)] - (let [[nfo ndo] (db/tx-run! cfg process-chunk! chunk)] - (recur (long (+ freezed nfo)) - (long (+ deleted ndo)))) - {:freeze freezed :delete deleted}))) + (let [{:keys [chunk poison]} (get-chunk pool timestamp)] + (when (seq poison) + (mtx/run! metrics :id :storage-gc-poison :inc (count poison)) + (db/tx-run! cfg defer-poison poison)) + ;; Keep draining after a poison-only batch: the deferred rows leave + ;; the selection and the next batch may hold healthy objects. + (if (or (seq chunk) (seq poison)) + (let [[nfo ndo] (if (seq chunk) + (db/tx-run! cfg process-chunk chunk) + [0 0])] + (recur (long (+ freezed nfo)) + (long (+ deleted ndo)))) + {:freeze freezed :delete deleted})))) ;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;; ;; HANDLER @@ -211,7 +257,8 @@ (defmethod ig/assert-key ::handler [_ params] - (assert (db/pool? (::db/pool params)) "expect valid storage")) + (assert (db/pool? (::db/pool params)) "expect valid storage") + (assert (mtx/metrics? (::mtx/metrics params)) "expect valid metrics")) (defmethod ig/expand-key ::handler [k v] @@ -223,5 +270,5 @@ (let [threshold (if (:skip-delay props) (ct/now) (ct/minus (ct/now) min-age))] - (process-touched! (assoc cfg ::timestamp threshold))))) + (process-touched (assoc cfg ::timestamp threshold))))) diff --git a/backend/src/app/storage/impl.clj b/backend/src/app/storage/impl.clj index 76fafe2cac..a6b3d2ca09 100644 --- a/backend/src/app/storage/impl.clj +++ b/backend/src/app/storage/impl.clj @@ -9,8 +9,8 @@ (:require [app.common.data.macros :as dm] [app.common.exceptions :as ex] - [app.db :as db] [app.storage :as-alias sto] + [app.storage.schema :as stsch] [buddy.core.codecs :as bc] [buddy.core.hash :as bh] [clojure.java.io :as jio] @@ -26,7 +26,7 @@ [{:keys [metadata] :as row}] (cond-> row (some? metadata) - (assoc :metadata (db/decode-transit-pgobject metadata)))) + (assoc :metadata (stsch/decode-metadata metadata)))) ;; --- API Definition diff --git a/backend/src/app/storage/schema.clj b/backend/src/app/storage/schema.clj new file mode 100644 index 0000000000..2ded471780 --- /dev/null +++ b/backend/src/app/storage/schema.clj @@ -0,0 +1,149 @@ +;; This Source Code Form is subject to the terms of the Mozilla Public +;; License, v. 2.0. If a copy of the MPL was not distributed with this +;; file, You can obtain one at http://mozilla.org/MPL/2.0/. +;; +;; Copyright (c) KALEIDOS SUBSIDIARY SL + +(ns app.storage.schema + "Schema, validation and encode/decode helpers for + `storage_object.metadata`. + + Phase 1: reads accept both the legacy Transit encoding and plain + JSON (sniffed by the `\"~:` marker); writes validate + normalize + but still serialize as Transit unless the + `:storage-metadata-as-json` config flag is set." + (:require + [app.common.data :as d] + [app.common.exceptions :as ex] + [app.common.schema :as sm] + [app.config :as cf] + [app.db :as db] + [cuerdas.core :as str]) + (:import + org.postgresql.util.PGobject)) + +(def default-bucket + "file-media-object") + +(def tempfile-bucket + "Bucket name for temporary file uploads (10-minute expiry)." + "tempfile") + +(def upload-session-bucket + "Bucket name for chunked-upload chunks." + "upload-session") + +(def bucket-requirements + "Canonical buckets and the extra keys each one must carry, so a new + bucket and its contract live in one place. Listed keys are declared as + optional in `schema:metadata` and only checked for presence; their type + is enforced by the map. Only buckets whose keys are read somewhere need + an entry: `file-data` ids are resolved by the GC, `organization-id` is + the logo's owner." + {"file-media-object" #{} + "team-font-variant" #{} + "file-object-thumbnail" #{} + "file-thumbnail" #{} + "profile" #{} + "organization" #{:organization-id} + tempfile-bucket #{} + upload-session-bucket #{} + "file-data" #{:file-id :id} + "file-data-fragment" #{} + "file-change" #{}}) + +(def metadata-buckets + "Canonical bucket set. `app.storage/valid-buckets` aliases it so the + list lives in exactly one place." + (set (keys bucket-requirements))) + +(defn- bucket-requirements-present? + [{:keys [bucket] :as mdata}] + (every? #(some? (get mdata %)) (get bucket-requirements bucket #{}))) + +(def schema:metadata + "Closed schema for `storage_object.metadata`. A single shape shared by + every bucket: `:bucket` is checked against `metadata-buckets` and the + rest are typed optional keys. `bucket-requirements` adds the per-bucket + presence checks." + [:and + [:map {:closed true} + [:bucket [::sm/one-of {:format :string} metadata-buckets]] + [:content-type :string] + [:hash {:optional true} :string] + [:profile-id {:optional true} ::sm/uuid] + [:organization-id {:optional true} ::sm/uuid] + [:file-id {:optional true} ::sm/uuid] + [:id {:optional true} ::sm/uuid]] + [:fn {:error/message "storage metadata is missing a required key for its bucket"} + bucket-requirements-present?]]) + +(defn- ->bucket + [v] + (cond + (string? v) v + (keyword? v) (d/name v) + :else (str v))) + +(defn- normalize-metadata + "Bring decoded (or incoming) metadata to its canonical shape: + `:bucket` is always present (`:reference` is mapped to it, keyword + or string, and defaults to `default-bucket`), and the legacy + `:reference`, `:upload-id` and `:chunk-index` keys are dropped." + [mdata] + (let [bucket (or (not-empty (->bucket (:bucket mdata))) + (not-empty (->bucket (:reference mdata))) + default-bucket)] + (-> mdata + (dissoc :reference :upload-id :chunk-index) + (assoc :bucket bucket)))) + +;; Decoder and encoder are built once per process: creating them +;; compiles the closed schema, and `decode-metadata` runs on every read +;; path (get-object, dedup probes, GC batches). +(def ^:private decode-metadata-fn + (sm/decode-fn schema:metadata (sm/json-transformer))) + +(def ^:private encode-metadata-fn + (sm/encoder schema:metadata (sm/json-transformer))) + +(defn- transit-encoded? + "Legacy metadata is Transit json-verbose, where every top-level key is + a keyword, so the document starts with `{\"~:`. Anchoring the check to + the start avoids reading a plain-JSON document as Transit just because + one of its values contains `\"~:`." + [value] + (and (string? value) + (str/starts-with? (str/trim value) "{\"~:"))) + +(defn decode-metadata + "Decode storage metadata from a PGobject, accepting both the legacy + Transit encoding (sniffed by the `\"~:` marker) and plain JSON. + Always returns the normalized shape with id fields as UUIDs." + [o] + (when (some? o) + (let [raw (if (transit-encoded? (.getValue ^PGobject o)) + (db/decode-transit-pgobject o) + (db/decode-json-pgobject o))] + (when-not (map? raw) + (ex/raise :type :internal + :code :invalid-storage-metadata + :hint "expected a map on storage object metadata")) + (decode-metadata-fn (normalize-metadata raw))))) + +(def ^:private check-metadata + (sm/check-fn schema:metadata + :hint "invalid storage object metadata" + :type :validation + :code :invalid-storage-metadata)) + +(defn encode-metadata + "Validate, normalize and encode storage metadata into a PGobject + ready for the `metadata` column. Serializes as Transit while the + `:storage-metadata-as-json` config flag is unset, as plain JSON + when it is set." + [mdata] + (let [mdata (check-metadata (normalize-metadata mdata))] + (if (cf/get :storage-metadata-as-json) + (db/json (encode-metadata-fn mdata)) + (db/tjson mdata)))) diff --git a/backend/test/backend_tests/http_assets_test.clj b/backend/test/backend_tests/http_assets_test.clj index 095d04f3e7..a5e60e9965 100644 --- a/backend/test/backend_tests/http_assets_test.clj +++ b/backend/test/backend_tests/http_assets_test.clj @@ -56,7 +56,17 @@ :bucket bucket :content-type "text/plain"} (some? profile-id) - (assoc :profile-id profile-id))))) + (assoc :profile-id profile-id) + + ;; file-data objects require file/id + ;; references for GC (has-file-data-refs?). + (= bucket "file-data") + (assoc :file-id (uuid/random) + :id (uuid/random)) + + ;; organization objects require their owner id. + (= bucket "organization") + (assoc :organization-id (uuid/random)))))) (defn- make-metrics [] diff --git a/backend/test/backend_tests/rpc_management_test.clj b/backend/test/backend_tests/rpc_management_test.clj index 570b858aaf..4b449af993 100644 --- a/backend/test/backend_tests/rpc_management_test.clj +++ b/backend/test/backend_tests/rpc_management_test.clj @@ -85,8 +85,7 @@ (configure-storage-backend)) sobject (sto/put-object! storage {::sto/content (sto/content "content") - :content-type "text/plain" - :other "data"}) + :content-type "text/plain"}) profile (th/create-profile* 1 {:is-active true}) project (th/create-project* 1 {:team-id (:default-team-id profile) :profile-id (:id profile)}) @@ -157,8 +156,7 @@ (let [storage (-> (:app.storage/storage th/*system*) (configure-storage-backend)) sobject (sto/put-object! storage {::sto/content (sto/content "content") - :content-type "text/plain" - :other "data"}) + :content-type "text/plain"}) profile (th/create-profile* 1 {:is-active true}) project (th/create-project* 1 {:team-id (:default-team-id profile) @@ -214,8 +212,7 @@ (configure-storage-backend)) sobject (sto/put-object! storage {::sto/content (sto/content "content") - :content-type "text/plain" - :other "data"}) + :content-type "text/plain"}) profile (th/create-profile* 1 {:is-active true}) project (th/create-project* 1 {:team-id (:default-team-id profile) @@ -278,8 +275,7 @@ (let [storage (-> (:app.storage/storage th/*system*) (configure-storage-backend)) sobject (sto/put-object! storage {::sto/content (sto/content "content") - :content-type "text/plain" - :other "data"}) + :content-type "text/plain"}) profile (th/create-profile* 1 {:is-active true}) project (th/create-project* 1 {:team-id (:default-team-id profile) :profile-id (:id profile)}) diff --git a/backend/test/backend_tests/storage_metadata_test.clj b/backend/test/backend_tests/storage_metadata_test.clj new file mode 100644 index 0000000000..4efb4c80df --- /dev/null +++ b/backend/test/backend_tests/storage_metadata_test.clj @@ -0,0 +1,276 @@ +;; This Source Code Form is subject to the terms of the Mozilla Public +;; License, v. 2.0. If a copy of the MPL was not distributed with this +;; file, You can obtain one at http://mozilla.org/MPL/2.0/. +;; +;; Copyright (c) KALEIDOS SUBSIDIARY SL + +(ns backend-tests.storage-metadata-test + (:require + [app.config :as cf] + [app.db :as db] + [app.storage.schema :as stsch] + [clojure.string :as str] + [clojure.test :as t]) + (:import + org.postgresql.util.PGobject)) + +(defn- pgobject + [^String value] + (doto (PGobject.) + (.setType "jsonb") + (.setValue value))) + +;; Raw Transit payloads shaped like the production rows sampled on +;; 2026-09-11 (verbose transit: "~:key" keys, "~u" uuid values, +;; "~:xxx" keyword values). +(def ^:private transit-media + "{\"~:hash\":\"blake2b:9f1c2e\",\"~:bucket\":\"file-media-object\",\"~:content-type\":\"image/png\"}") + +(def ^:private transit-tempfile + "{\"~:hash\":\"blake2b:9f1c2e\",\"~:bucket\":\"tempfile\",\"~:content-type\":\"application/zip\",\"~:profile-id\":\"~u86907e95-1cb8-8122-8008-4eb7ba07d89d\"}") + +(def ^:private transit-legacy-reference + "{\"~:reference\":\"~:file-media-object\",\"~:content-type\":\"image/png\",\"~:hash\":\"blake2b:9f1c2e\"}") + +(def ^:private transit-no-bucket + "{\"~:content-type\":\"image/svg+xml\"}") + +(def ^:private transit-chunk-leftovers + "{\"~:bucket\":\"tempfile\",\"~:content-type\":\"application/zip\",\"~:upload-id\":\"~u86907e95-1cb8-8122-8008-4eb7ba07d89d\",\"~:chunk-index\":3}") + +(def ^:private json-file-data + "{\"bucket\":\"file-data\",\"content-type\":\"application/octet-stream\",\"file-id\":\"86907e95-1cb8-8122-8008-4eb7ba07d89d\",\"id\":\"83df2f92-6bd4-4e6d-9c9a-3f6d2b1a4c55\"}") + +(t/deftest decode-transit-returns-native-types + (let [mdata (stsch/decode-metadata (pgobject transit-tempfile))] + (t/is (= "tempfile" (:bucket mdata))) + (t/is (= "application/zip" (:content-type mdata))) + (t/is (= "blake2b:9f1c2e" (:hash mdata))) + (t/is (uuid? (:profile-id mdata))) + (t/is (= (parse-uuid "86907e95-1cb8-8122-8008-4eb7ba07d89d") + (:profile-id mdata))))) + +(t/deftest decode-transit-maps-reference-to-bucket + (let [mdata (stsch/decode-metadata (pgobject transit-legacy-reference))] + (t/is (= "file-media-object" (:bucket mdata))) + (t/is (= "image/png" (:content-type mdata))) + (t/is (not (contains? mdata :reference))))) + +(t/deftest decode-transit-defaults-missing-bucket + (let [mdata (stsch/decode-metadata (pgobject transit-no-bucket))] + (t/is (= "file-media-object" (:bucket mdata))) + (t/is (= "image/svg+xml" (:content-type mdata))))) + +(t/deftest decode-transit-drops-chunk-leftovers + (let [mdata (stsch/decode-metadata (pgobject transit-chunk-leftovers))] + (t/is (= "tempfile" (:bucket mdata))) + (t/is (not (contains? mdata :upload-id))) + (t/is (not (contains? mdata :chunk-index))))) + +(t/deftest decode-json-coerces-uuids + (let [mdata (stsch/decode-metadata (pgobject json-file-data))] + (t/is (= "file-data" (:bucket mdata))) + (t/is (uuid? (:file-id mdata))) + (t/is (uuid? (:id mdata))) + (t/is (= (parse-uuid "86907e95-1cb8-8122-8008-4eb7ba07d89d") + (:file-id mdata))))) + +(t/deftest sniff-has-no-false-positives-on-values + ;; "~:" inside a value (not anchored to a quote) must not trigger + ;; the transit branch. + (let [mdata (stsch/decode-metadata + (pgobject "{\"bucket\":\"tempfile\",\"content-type\":\"text/plain;~:x\"}"))] + (t/is (= "tempfile" (:bucket mdata))) + (t/is (= "text/plain;~:x" (:content-type mdata))))) + +(t/deftest encode-writes-transit-by-default + ;; Pinned off: without the binding this test inherits the ambient + ;; config and proves nothing where the flag is set. + (binding [cf/config (assoc cf/config :storage-metadata-as-json nil)] + (let [encoded (stsch/encode-metadata {:bucket "file-media-object" + :content-type "image/png" + :hash "blake2b:9f1c2e"}) + value (.getValue ^PGobject encoded)] + (t/is (string? value)) + (t/is (str/includes? value "\"~:bucket\"")) + (t/is (not (str/includes? value "\"~:reference\"")))))) + +(t/deftest encode-normalizes-legacy-input + (binding [cf/config (assoc cf/config :storage-metadata-as-json nil)] + (let [encoded (stsch/encode-metadata {:reference :file-media-object + :content-type "image/png"}) + value (.getValue ^PGobject encoded) + decoded (stsch/decode-metadata encoded)] + (t/is (= "file-media-object" (:bucket decoded))) + (t/is (not (str/includes? value "\"~:reference\"")))))) + +(t/deftest encode-writes-plain-json-with-flag + (binding [cf/config (assoc cf/config :storage-metadata-as-json true)] + (let [file-id (parse-uuid "86907e95-1cb8-8122-8008-4eb7ba07d89d") + encoded (stsch/encode-metadata {:bucket "file-data" + :content-type "application/octet-stream" + :file-id file-id + :id (parse-uuid "83df2f92-6bd4-4e6d-9c9a-3f6d2b1a4c55")}) + value (.getValue ^PGobject encoded)] + (t/is (str/includes? value "\"bucket\"")) + (t/is (not (str/includes? value "\"~:"))) + ;; uuids travel as plain strings on the JSON encoding + (t/is (str/includes? value (str file-id))) + (let [decoded (stsch/decode-metadata encoded)] + (t/is (uuid? (:file-id decoded))) + (t/is (= file-id (:file-id decoded))))))) + +(t/deftest encode-rejects-string-uuids + ;; Encoding does not coerce input types: callers must pass native UUIDs + ;; (every producer does; reads still coerce on the way back). + (t/is (thrown-with-msg? clojure.lang.ExceptionInfo + #"invalid storage object metadata" + (stsch/encode-metadata {:bucket "tempfile" + :content-type "application/zip" + :profile-id "86907e95-1cb8-8122-8008-4eb7ba07d89d"})))) + +(t/deftest encode-rejects-unknown-bucket + (t/is (thrown-with-msg? clojure.lang.ExceptionInfo + #"invalid storage object metadata" + (stsch/encode-metadata {:bucket "no-such-bucket" + :content-type "image/png"})))) + +(t/deftest encode-rejects-extra-key + (t/is (thrown-with-msg? clojure.lang.ExceptionInfo + #"invalid storage object metadata" + (stsch/encode-metadata {:bucket "file-media-object" + :content-type "image/png" + :other "data"})))) + +(t/deftest encode-rejects-malformed-uuid + (t/is (thrown-with-msg? clojure.lang.ExceptionInfo + #"invalid storage object metadata" + (stsch/encode-metadata {:bucket "file-data" + :content-type "application/octet-stream" + :file-id "not-a-uuid" + :id "83df2f92-6bd4-4e6d-9c9a-3f6d2b1a4c55"})))) + +(t/deftest encode-rejects-missing-content-type + (t/is (thrown-with-msg? clojure.lang.ExceptionInfo + #"invalid storage object metadata" + (stsch/encode-metadata {:bucket "file-media-object"})))) + +(t/deftest encode-rejects-missing-file-data-ids + ;; Both ids are required for file-data; omitting one fails the bucket + ;; requirement check (the file-id here is a native UUID, so the map + ;; itself is valid). + (t/is (thrown-with-msg? clojure.lang.ExceptionInfo + #"invalid storage object metadata" + (stsch/encode-metadata {:bucket "file-data" + :content-type "application/octet-stream" + :file-id (parse-uuid "86907e95-1cb8-8122-8008-4eb7ba07d89d")})))) + +(t/deftest encode-rejects-missing-organization-id + (t/is (thrown-with-msg? clojure.lang.ExceptionInfo + #"invalid storage object metadata" + (stsch/encode-metadata {:bucket "organization" + :content-type "image/svg+xml"})))) + +(t/deftest decode-real-transit-payload + ;; Same legacy shape as above but produced by the real transit + ;; encoder, so the test breaks if the encoding ever drifts. + (let [mdata (stsch/decode-metadata + (db/tjson {:reference :team-font-variant + :content-type "font/woff2"}))] + (t/is (= "team-font-variant" (:bucket mdata))) + (t/is (not (contains? mdata :reference))))) + +(t/deftest transit-and-json-encodings-decode-to-same-shape + (let [mdata {:bucket "tempfile" + :content-type "application/zip" + :profile-id (parse-uuid "86907e95-1cb8-8122-8008-4eb7ba07d89d")} + transit (binding [cf/config (assoc cf/config :storage-metadata-as-json nil)] + (stsch/decode-metadata (stsch/encode-metadata mdata))) + json (binding [cf/config (assoc cf/config :storage-metadata-as-json true)] + (stsch/decode-metadata (stsch/encode-metadata mdata)))] + (t/is (= transit json)) + (t/is (= mdata json)))) + +(t/deftest flag-on-and-off-differ-only-in-encoding-family + ;; The Phase 2 rollback contract: flipping the flag changes how the + ;; same logical metadata hits the disk, never what it means. + (let [mdata {:bucket "file-media-object" + :content-type "image/png" + :hash "blake2b:9f1c2e"} + transit (binding [cf/config (assoc cf/config :storage-metadata-as-json nil)] + (.getValue ^PGobject (stsch/encode-metadata mdata))) + json (binding [cf/config (assoc cf/config :storage-metadata-as-json true)] + (.getValue ^PGobject (stsch/encode-metadata mdata)))] + (t/is (str/includes? transit "\"~:bucket\"")) + (t/is (str/includes? json "\"bucket\"")) + (t/is (not (str/includes? json "\"~:"))) + (t/is (= (stsch/decode-metadata (pgobject transit)) + (stsch/decode-metadata (pgobject json)))))) + +(t/deftest decode-json-keeps-hash-byte-exact + ;; Dedup matches on the hash string; it must survive the JSON + ;; roundtrip untouched. + (let [mdata (stsch/decode-metadata + (pgobject "{\"bucket\":\"file-media-object\",\"content-type\":\"image/png\",\"hash\":\"blake2b:9f1c2e\"}"))] + (t/is (= "blake2b:9f1c2e" (:hash mdata))))) + +(def ^:private bucket-samples + ;; One representative payload per bucket, with every required key and + ;; one optional key where the bucket declares any. + [["file-media-object" + {:bucket "file-media-object" :content-type "image/png" :hash "blake2b:9f1c2e"}] + ["team-font-variant" + {:bucket "team-font-variant" :content-type "font/woff2"}] + ["file-object-thumbnail" + {:bucket "file-object-thumbnail" :content-type "image/png"}] + ["file-thumbnail" + {:bucket "file-thumbnail" :content-type "image/png"}] + ["profile" + {:bucket "profile" :content-type "image/png"}] + ["organization" + {:bucket "organization" :content-type "image/svg+xml" + :organization-id #uuid "11111111-2222-3333-4444-555555555555"}] + ["tempfile" + {:bucket "tempfile" :content-type "application/zip" + :profile-id #uuid "86907e95-1cb8-8122-8008-4eb7ba07d89d"}] + ["upload-session" + {:bucket "upload-session" :content-type "application/octet-stream"}] + ["file-data" + {:bucket "file-data" :content-type "application/octet-stream" + :file-id #uuid "86907e95-1cb8-8122-8008-4eb7ba07d89d" + :id #uuid "83df2f92-6bd4-4e6d-9c9a-3f6d2b1a4c55"}] + ["file-data-fragment" + {:bucket "file-data-fragment" :content-type "application/octet-stream"}] + ["file-change" + {:bucket "file-change" :content-type "application/octet-stream"}]]) + +(t/deftest encode-decode-roundtrips-every-bucket + (binding [cf/config (assoc cf/config :storage-metadata-as-json nil)] + (doseq [[bucket mdata] bucket-samples] + (t/testing bucket + (t/is (= mdata (stsch/decode-metadata (stsch/encode-metadata mdata)))))))) + +(t/deftest json-encoding-roundtrips-every-bucket + (binding [cf/config (assoc cf/config :storage-metadata-as-json true)] + (doseq [[bucket mdata] bucket-samples] + (t/testing bucket + (t/is (= mdata (stsch/decode-metadata (stsch/encode-metadata mdata)))))))) + +(t/deftest metadata-buckets-is-the-canonical-set + (t/is (= #{"file-media-object" "team-font-variant" "file-object-thumbnail" + "file-thumbnail" "profile" "organization" "tempfile" + "upload-session" "file-data" "file-data-fragment" "file-change"} + stsch/metadata-buckets))) + +(t/deftest decode-metadata-returns-nil-for-nil + (t/is (nil? (stsch/decode-metadata nil)))) + +(t/deftest sniff-treats-leading-tilde-value-as-json + ;; A plain-JSON value starting with `~:` must not switch the reader to + ;; the Transit branch: the sniff is anchored to the first key. + (let [encoded (binding [cf/config (assoc cf/config :storage-metadata-as-json true)] + (stsch/encode-metadata {:bucket "tempfile" + :content-type "~:not-transit"})) + decoded (stsch/decode-metadata encoded)] + (t/is (= "~:not-transit" (:content-type decoded))) + (t/is (= "tempfile" (:bucket decoded))))) diff --git a/backend/test/backend_tests/storage_test.clj b/backend/test/backend_tests/storage_test.clj index 54953620a0..abdd8c1503 100644 --- a/backend/test/backend_tests/storage_test.clj +++ b/backend/test/backend_tests/storage_test.clj @@ -9,12 +9,16 @@ [app.common.exceptions :as ex] [app.common.time :as ct] [app.common.uuid :as uuid] + [app.config :as cf] [app.db :as db] + [app.metrics :as mtx] + [app.metrics.definition :as-alias mdef] [app.rpc :as-alias rpc] [app.storage :as sto] [app.storage.fs :as-alias sto.fs] [app.storage.impl :as impl] [app.storage.s3 :as-alias sto.s3] + [app.storage.schema :as stsch] [backend-tests.helpers :as th] [clojure.test :as t] [cuerdas.core :as str] @@ -23,6 +27,11 @@ [mockery.core :refer [with-mocks]] [promesa.core :as p]) (:import + (io.prometheus.client + Counter + Counter$Child) + (org.postgresql.util + PGobject) (software.amazon.awssdk.services.s3 S3AsyncClient) (software.amazon.awssdk.services.s3.model @@ -41,24 +50,51 @@ [storage] (assoc storage ::sto/backend :fs)) +(defn- poison-counter-value + "Read the process-wide gc-poison counter from the test system metrics." + [] + (let [metrics (:app.metrics/metrics th/*system*) + collector (mtx/get-collector metrics :storage-gc-poison) + instance (::mdef/instance collector)] + (.get ^Counter$Child (.labels ^Counter instance (into-array String []))))) + (t/deftest put-and-retrieve-object (let [storage (-> (:app.storage/storage th/*system*) (configure-storage-backend)) content (sto/content "content") object (sto/put-object! storage {::sto/content content - :content-type "text/plain" - :other "data"})] + :content-type "text/plain"})] (t/is (sto/object? object)) (t/is (fs/path? (sto/get-object-path storage object))) (t/is (nil? (:expired-at object))) (t/is (= :fs (:backend object))) - (t/is (= "data" (:other (meta object)))) + (t/is (= "file-media-object" (:bucket (meta object)))) (t/is (= "text/plain" (:content-type (meta object)))) (t/is (= "content" (slurp (sto/get-object-data storage object)))) (t/is (= "content" (slurp (sto/get-object-path storage object)))))) +(t/deftest put-and-retrieve-with-json-metadata-flag + ;; End-to-end through the jsonb column: with the flag on, the row is + ;; stored as plain JSON and still decodes to native types. + (binding [cf/config (assoc cf/config :storage-metadata-as-json true)] + (let [storage (-> (:app.storage/storage th/*system*) + (configure-storage-backend)) + profile (uuid/random) + object (sto/put-object! storage {::sto/content (sto/content "content") + :bucket "tempfile" + :content-type "application/zip" + :profile-id profile}) + row (th/db-exec-one! ["select metadata::text as metadata from storage_object where id = ?" + (:id object)])] + (t/is (not (str/includes? (:metadata row) "\"~:"))) + (t/is (str/includes? (:metadata row) "\"bucket\"")) + (let [loaded (sto/get-object storage (:id object))] + (t/is (= "tempfile" (:bucket (meta loaded)))) + (t/is (= "application/zip" (:content-type (meta loaded)))) + (t/is (= profile (:profile-id (meta loaded)))))))) + (t/deftest tempfile-objects-are-not-deduplicated (let [storage (-> (:app.storage/storage th/*system*) (configure-storage-backend)) @@ -94,8 +130,7 @@ (configure-storage-backend)) content (sto/content "content") object (sto/put-object! storage {::sto/content content - :content-type "text/plain" - :expired-at (ct/in-future {:seconds 1})})] + :content-type "text/plain"})] (t/is (sto/object? object)) (t/is (true? (sto/del-object! storage (:id object)))) @@ -109,6 +144,18 @@ ;; marked as deleted/expired. (t/is (nil? (sto/get-object storage (:id object)))))) +(t/deftest delete-expired-object-returns-false + ;; An object stored with `::sto/expired-at` is born with `deleted_at` + ;; set, so deleting it is a no-op and must report `false`. + (let [storage (-> (:app.storage/storage th/*system*) + (configure-storage-backend)) + content (sto/content "content") + object (sto/put-object! storage {::sto/content content + :content-type "text/plain" + ::sto/expired-at (ct/in-future {:hours 1})})] + (t/is (some? (:expired-at object))) + (t/is (false? (sto/del-object! storage (:id object)))))) + (t/deftest deleted-gc-task (let [storage (-> (:app.storage/storage th/*system*) (configure-storage-backend)) @@ -438,6 +485,146 @@ (let [row (th/db-exec-one! ["select deleted_at from storage_object where id = ?" (:id object1)])] (t/is (ct/is-before-or-equal? (:deleted-at row) (ct/plus now {:seconds 1})))))) +(t/deftest storage-gc-touched-null-metadata + ;; A NULL metadata column (predates any normalization) flows through + ;; the lookup fallback with a warning instead of breaking the GC loop. + (let [storage (-> (:app.storage/storage th/*system*) + (configure-storage-backend)) + content (sto/content "content") + object (sto/put-object! storage {::sto/content content + ::sto/touch true + :content-type "text/plain"})] + (th/db-exec! ["update storage_object set metadata = null where id = ?" (:id object)]) + (let [res (th/run-task! :storage-gc-touched {:skip-delay true})] + (t/is (= 0 (:freeze res))) + (t/is (= 1 (:delete res)))))) + +(t/deftest storage-gc-touched-defers-corrupt-metadata + ;; A non-map metadata row neither blocks the chunk nor gets + ;; collected: it is logged and deferred exactly one day, keeping its + ;; metadata intact for a later repair. + (let [now (ct/now) + before (poison-counter-value) + storage (-> (:app.storage/storage th/*system*) + (configure-storage-backend)) + healthy (sto/put-object! storage {::sto/content (sto/content "healthy") + ::sto/touched-at now + :content-type "text/plain"}) + poison (uuid/random)] + (th/db-exec! ["insert into storage_object (id, backend, metadata, touched_at) values (?, 'fs', '[]'::jsonb, ?)" + poison now]) + (binding [ct/*clock* (ct/fixed-clock now)] + (let [res (th/run-task! :storage-gc-touched {:skip-delay true})] + (t/is (= 0 (:freeze res))) + (t/is (= 1 (:delete res))))) + (let [row (th/db-exec-one! ["select metadata::text as metadata, touched_at from storage_object where id = ?" poison])] + (t/is (= "[]" (:metadata row))) + ;; inst-ms: timestamptz keeps micros, the frozen clock has nanos. + (t/is (= (inst-ms (ct/plus now {:days 1})) + (inst-ms (:touched-at row))))) + ;; One poison row observed, one increment. + (t/is (= (inc before) (poison-counter-value))))) + +(t/deftest storage-gc-touched-poison-only + ;; A chunk made only of poison rows still terminates: the rows are + ;; deferred, the next scan returns nothing, and no delete is counted. + (let [now (ct/now) + poison (uuid/random)] + (th/db-exec! ["insert into storage_object (id, backend, metadata, touched_at) values (?, 'fs', '[]'::jsonb, ?)" + poison now]) + (binding [ct/*clock* (ct/fixed-clock now)] + (let [res (th/run-task! :storage-gc-touched {:skip-delay true})] + (t/is (= 0 (:freeze res))) + (t/is (= 0 (:delete res))))) + (let [row (th/db-exec-one! ["select touched_at from storage_object where id = ?" poison])] + (t/is (= (inst-ms (ct/plus now {:days 1})) + (inst-ms (:touched-at row))))))) + +(t/deftest storage-gc-touched-drains-past-a-poison-only-batch + ;; A full batch of poison rows must not stop the run: they are deferred + ;; and the loop keeps draining, so a healthy row queued behind the + ;; LIMIT (later touched_at) still gets collected. + (let [now (ct/now) + storage (-> (:app.storage/storage th/*system*) + (configure-storage-backend)) + earlier (ct/minus now {:seconds 1}) + healthy (sto/put-object! storage {::sto/content (sto/content "healthy") + ::sto/touched-at now + :content-type "text/plain"}) + poisons (repeatedly 10 uuid/random)] + (doseq [id poisons] + (th/db-exec! ["insert into storage_object (id, backend, metadata, touched_at) values (?, 'fs', '[]'::jsonb, ?)" + id earlier])) + (binding [ct/*clock* (ct/fixed-clock now)] + (let [res (th/run-task! :storage-gc-touched {:skip-delay true})] + (t/is (= 0 (:freeze res))) + (t/is (= 1 (:delete res))))) + (doseq [id poisons] + (let [row (th/db-exec-one! ["select touched_at from storage_object where id = ?" id])] + (t/is (= (inst-ms (ct/plus now {:days 1})) + (inst-ms (:touched-at row)))))))) + +(def ^:private migration-0155-fixtures + ;; [id transit-metadata]: production-shaped legacy rows (nil payload + ;; means a NULL column). + [["11111111-1111-1111-1111-111111111111" + "{\"~:reference\":\"~:file-media-object\",\"~:content-type\":\"image/png\",\"~:hash\":\"blake2b:aaa\"}"] + ["22222222-2222-2222-2222-222222222222" + "{\"~:content-type\":\"image/svg+xml\"}"] + ["33333333-3333-3333-3333-333333333333" + "{\"~:bucket\":\"tempfile\",\"~:reference\":\"~:tempfile\",\"~:content-type\":\"application/zip\"}"] + ["44444444-4444-4444-4444-444444444444" + "{\"~:bucket\":\"tempfile\",\"~:content-type\":\"application/zip\",\"~:upload-id\":\"~u86907e95-1cb8-8122-8008-4eb7ba07d89d\",\"~:chunk-index\":3}"] + ["55555555-5555-5555-5555-555555555555" + nil]]) + +(defn- run-migration-0155! + ;; Re-runs the 0155 statements (not migratus: it already applied at + ;; bootstrap) over the fixture rows above. + [] + (let [sql (-> (io/resource "app/migrations/sql/0155-normalize-storage-object-metadata.sql") + (slurp)) + no-comments (->> (.split ^String sql "\n") + (remove #(.startsWith ^String (str/trim %) "--")) + (str/join "\n"))] + (doseq [stmt (->> (.split ^String no-comments ";") + (map str/trim) + (remove str/blank?))] + (th/db-exec! [stmt])))) + +(defn- get-metadata-by-id + [id] + (:metadata (th/db-exec-one! ["select metadata from storage_object where id = ?" + (parse-uuid id)]))) + +(t/deftest storage-migration-0155-normalizes-legacy-rows + (doseq [[id mdata] migration-0155-fixtures] + (th/db-exec! ["insert into storage_object (id, backend, metadata) values (?, 'fs', ?::jsonb)" + (parse-uuid id) mdata])) + (run-migration-0155!) + (let [mdata (fn [id] (stsch/decode-metadata (get-metadata-by-id id)))] + (t/is (= {:bucket "file-media-object" + :content-type "image/png" + :hash "blake2b:aaa"} + (mdata "11111111-1111-1111-1111-111111111111"))) + (t/is (= {:bucket "file-media-object" + :content-type "image/svg+xml"} + (mdata "22222222-2222-2222-2222-222222222222"))) + (t/is (= {:bucket "tempfile" + :content-type "application/zip"} + (mdata "33333333-3333-3333-3333-333333333333"))) + (t/is (= {:bucket "tempfile" + :content-type "application/zip"} + (mdata "44444444-4444-4444-4444-444444444444"))) + (t/is (= {:bucket "file-media-object"} + (mdata "55555555-5555-5555-5555-555555555555")))) + ;; second run changes nothing (idempotent) + (let [raw (fn [] (mapv #(.getValue ^PGobject (get-metadata-by-id %)) + (map first migration-0155-fixtures))) + before (raw)] + (run-migration-0155!) + (t/is (= before (raw))))) + (t/deftest storage-gc-deleted-immediate (let [storage (-> (:app.storage/storage th/*system*) (configure-storage-backend)) @@ -626,6 +813,43 @@ (let [row (th/db-exec-one! ["select count(*) from storage_object"])] (t/is (= 1 (:count row)))))) +(t/deftest dedup-reuses-json-encoded-blob + ;; With the JSON flag on, the dedup lookup must find the row it just + ;; wrote; a "~:"-only lookup used to miss it and duplicate the blob. + (binding [cf/config (assoc cf/config :storage-metadata-as-json true)] + (let [storage (-> (:app.storage/storage th/*system*) + (configure-storage-backend)) + content (-> (sto/content "json-content") + (sto/wrap-with-hash "json-hash")) + params {::sto/content content + ::sto/deduplicate? true + :bucket "file-media-object" + :content-type "text/plain"} + object1 (sto/put-object! storage params) + object2 (sto/put-object! storage params)] + (t/is (= (:id object1) (:id object2))) + (let [row (th/db-exec-one! ["select count(*) from storage_object"])] + (t/is (= 1 (:count row))))))) + +(t/deftest dedup-json-put-finds-transit-row + ;; Both encodings coexist during the transition: a JSON write must + ;; reuse a blob already stored as Transit. + (let [storage (-> (:app.storage/storage th/*system*) + (configure-storage-backend)) + content (-> (sto/content "mixed-content") + (sto/wrap-with-hash "mixed-hash")) + params {::sto/content content + ::sto/deduplicate? true + :bucket "file-media-object" + :content-type "text/plain"} + object1 (binding [cf/config (assoc cf/config :storage-metadata-as-json nil)] + (sto/put-object! storage params)) + object2 (binding [cf/config (assoc cf/config :storage-metadata-as-json true)] + (sto/put-object! storage params))] + (t/is (= (:id object1) (:id object2))) + (let [row (th/db-exec-one! ["select count(*) from storage_object"])] + (t/is (= 1 (:count row)))))) + (t/deftest dedup-repairs-stale-object (let [storage (-> (:app.storage/storage th/*system*) (configure-storage-backend))