✨ Normalize storage metadata with a closed schema (#11987)

* ✨ 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

* ♻️ 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

* ♻️ 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

* ♻️ 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

* ♻️ 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

* ♻️ 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

* ♻️ 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

* 📚 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

* 📚 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

* ♻️ 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

* ♻️ 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

* ♻️ 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

* ♻️ 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

* 🐛 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
This commit is contained in:
Andrey Antukh 2026-10-01 07:19:21 +02:00 committed by GitHub
parent 0509e2b9d2
commit d67a00c1d5
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
15 changed files with 884 additions and 94 deletions

View File

@ -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.

View File

@ -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]]]))

View File

@ -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)

View File

@ -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]

View File

@ -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;

View File

@ -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';

View File

@ -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"})

View File

@ -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})))

View File

@ -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)))))

View File

@ -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

View File

@ -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))))

View File

@ -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
[]

View File

@ -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)})

View File

@ -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>" 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)))))

View File

@ -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))