penpot/backend/test/backend_tests/storage_test.clj
Andrey Antukh d67a00c1d5
✨ 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
2026-10-01 07:19:21 +02:00

1239 lines
55 KiB
Clojure

;; 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-test
(:require
[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]
[datoteka.fs :as fs]
[datoteka.io :as io]
[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
NoSuchKeyException)
(software.amazon.awssdk.services.s3.presigner
S3Presigner)))
(t/use-fixtures :once th/state-init)
(t/use-fixtures :each (th/serial
th/database-reset
th/clean-storage))
(defn configure-storage-backend
"Given storage map, returns a storage configured with the appropriate
backend for assets."
[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"})]
(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 (= "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))
content (-> (sto/content "content")
(sto/wrap-with-hash "same-hash"))
object1 (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
::sto/touched-at (ct/in-future {:minutes 10})
:bucket "tempfile"
:content-type "text/plain"})
object2 (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
::sto/touched-at (ct/in-future {:minutes 10})
:bucket "tempfile"
:content-type "text/plain"})]
(t/is (not= (:id object1) (:id object2)))))
(t/deftest put-and-retrieve-expired-object
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
content (sto/content "content")
object (sto/put-object! storage {::sto/content content
::sto/expired-at (ct/in-future {:hours 1})
:content-type "text/plain"})]
(t/is (sto/object? object))
(t/is (ct/inst? (:expired-at object)))
(t/is (ct/is-after? (:expired-at object) (ct/now)))
(t/is (nil? (sto/get-object storage (:id object))))))
(t/deftest put-and-delete-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"})]
(t/is (sto/object? object))
(t/is (true? (sto/del-object! storage (:id object))))
;; retrieving the same object should be not nil because the
;; deletion is not immediate
(t/is (some? (sto/get-object-data storage object)))
(t/is (some? (sto/get-object-url storage object)))
(t/is (some? (sto/get-object-path storage object)))
;; But you can't retrieve the object again because in database is
;; 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))
content1 (sto/content "content1")
content2 (sto/content "content2")
content3 (sto/content "content3")
object1 (sto/put-object! storage {::sto/content content1
::sto/expired-at (ct/now)
:content-type "text/plain"})
object2 (sto/put-object! storage {::sto/content content2
::sto/expired-at (ct/in-future {:hours 2})
:content-type "text/plain"})
object3 (sto/put-object! storage {::sto/content content3
::sto/expired-at (ct/in-future {:hours 1})
:content-type "text/plain"})]
(binding [ct/*clock* (ct/fixed-clock (ct/in-future {:minutes 0}))]
(let [res (th/run-task! :storage-gc-deleted {})]
(t/is (= 1 (:deleted res)))))
(let [res (th/db-exec-one! ["select count(*) from storage_object;"])]
(t/is (= 2 (:count res))))
(binding [ct/*clock* (ct/fixed-clock (ct/in-future {:minutes 61}))]
(let [res (th/run-task! :storage-gc-deleted {})]
(t/is (= 1 (:deleted res)))))
(let [res (th/db-exec-one! ["select count(*) from storage_object;"])]
(t/is (= 1 (:count res))))))
(t/deftest touched-gc-task-1
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
prof (th/create-profile* 1)
proj (th/create-project* 1 {:profile-id (:id prof)
:team-id (:default-team-id prof)})
file (th/create-file* 1 {:profile-id (:id prof)
:project-id (:default-project-id prof)
:is-shared false})
mfile {:filename "sample.jpg"
:path (th/tempfile "backend_tests/test_files/sample.jpg")
:mtype "image/jpeg"
:size 312043}
params {::th/type :upload-file-media-object
::rpc/profile-id (:id prof)
:file-id (:id file)
:is-local true
:name "testfile"
:content mfile}
out1 (th/command! params)
out2 (th/command! params)]
(t/is (nil? (:error out1)))
(t/is (nil? (:error out2)))
(let [result-1 (:result out1)
result-2 (:result out2)]
(t/is (uuid? (:id result-1)))
(t/is (uuid? (:id result-2)))
(t/is (uuid? (:media-id result-1)))
(t/is (uuid? (:media-id result-2)))
(t/is (= (:media-id result-1) (:media-id result-2)))
(th/db-update! :file-media-object
{:deleted-at (ct/now)}
{:id (:id result-1)})
;; run the objects gc task for permanent deletion
(let [res (th/run-task! :objects-gc {})]
(t/is (= 1 (:processed res))))
;; check that we still have all the storage objects
(let [res (th/db-exec-one! ["select count(*) from storage_object"])]
(t/is (= 2 (:count res))))
;; now check if the storage objects are touched
(let [res (th/db-exec-one! ["select count(*) from storage_object where touched_at is not null"])]
(t/is (= 2 (:count res))))
;; run the touched gc task
(let [res (binding [ct/*clock* (ct/fixed-clock (ct/in-future {:hours 3}))]
(th/run-task! :storage-gc-touched {}))]
(t/is (= 2 (:freeze res)))
(t/is (= 0 (:delete res))))
;; now check that there are no touched objects
(let [res (th/db-exec-one! ["select count(*) from storage_object where touched_at is not null"])]
(t/is (= 0 (:count res))))
;; now check that all objects are marked to be deleted
(let [res (th/db-exec-one! ["select count(*) from storage_object where deleted_at is not null"])]
(t/is (= 0 (:count res)))))))
(defn- upload-font-chunked!
"Splits `font-bytes` into a single chunk, creates an upload session,
uploads the chunk, and returns the session-id UUID."
[prof ^bytes font-bytes mtype]
(let [tmp (fs/create-tempfile :dir "/tmp/penpot" :prefix "test-font-chunk-")
_ (io/write* tmp font-bytes)
mfile {:filename "chunk" :path tmp :mtype mtype :size (alength font-bytes)}
session-id (-> (th/command! {::th/type :create-upload-session
::rpc/profile-id (:id prof)
:total-chunks 1})
:result :session-id)
out (th/command! {::th/type :upload-chunk
::rpc/profile-id (:id prof)
:session-id session-id
:index 0
:content mfile})]
(assert (nil? (:error out)))
session-id))
(t/deftest touched-gc-task-2
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
prof (th/create-profile* 1 {:is-active true})
team-id (:default-team-id prof)
proj-id (:default-project-id prof)
font-id (uuid/custom 10 1)
proj (th/create-project* 1 {:profile-id (:id prof)
:team-id team-id})
file (th/create-file* 1 {:profile-id (:id prof)
:project-id proj-id
:is-shared false})
ttfdata (-> (io/resource "backend_tests/test_files/font-1.ttf")
(io/read*))
mfile {:filename "sample.jpg"
:path (th/tempfile "backend_tests/test_files/sample.jpg")
:mtype "image/jpeg"
:size 312043}
params1 {::th/type :upload-file-media-object
::rpc/profile-id (:id prof)
:file-id (:id file)
:is-local true
:name "testfile"
:content mfile}
session-id (upload-font-chunked! prof ttfdata "font/ttf")
params2 {::th/type :create-font-variant
::rpc/profile-id (:id prof)
:team-id team-id
:font-id font-id
:font-family "somefont"
:font-weight 400
:font-style "normal"
:uploads {"font/ttf" session-id}}
out1 (th/command! params1)
out2 (th/command! params2)]
;; (th/print-result! out)
(t/is (nil? (:error out1)))
(t/is (nil? (:error out2)))
;; run the touched gc task
(let [res (binding [ct/*clock* (ct/fixed-clock (ct/in-future {:hours 3}))]
(th/run-task! :storage-gc-touched {}))]
(t/is (= 5 (:freeze res)))
(t/is (= 1 (:delete res)))
(let [result-1 (:result out1)
result-2 (:result out2)]
(th/db-update! :team-font-variant
{:deleted-at (ct/now)}
{:id (:id result-2)})
;; run the objects gc task for permanent deletion
;; (processed = 2: the consumed upload session plus the font variant)
(let [res (th/run-task! :objects-gc {})]
(t/is (= 2 (:processed res))))
;; revert touched state to all storage objects
(th/db-exec-one! ["update storage_object set touched_at=?" (ct/now)])
;; Run the task again
(let [res (binding [ct/*clock* (ct/fixed-clock (ct/in-future {:hours 3}))]
(th/run-task! :storage-gc-touched {}))]
(t/is (= 2 (:freeze res)))
(t/is (= 4 (:delete res))))
;; now check that there are no touched objects
(let [res (th/db-exec-one! ["select count(*) from storage_object where touched_at is not null"])]
(t/is (= 0 (:count res))))
;; now check that all objects are marked to be deleted
(let [res (th/db-exec-one! ["select count(*) from storage_object where deleted_at is not null"])]
(t/is (= 4 (:count res))))))))
(t/deftest touched-gc-task-3
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
prof (th/create-profile* 1)
proj (th/create-project* 1 {:profile-id (:id prof)
:team-id (:default-team-id prof)})
file (th/create-file* 1 {:profile-id (:id prof)
:project-id (:default-project-id prof)
:is-shared false})
mfile {:filename "sample.jpg"
:path (th/tempfile "backend_tests/test_files/sample.jpg")
:mtype "image/jpeg"
:size 312043}
params {::th/type :upload-file-media-object
::rpc/profile-id (:id prof)
:file-id (:id file)
:is-local true
:name "testfile"
:content mfile}
out1 (th/command! params)
out2 (th/command! params)]
(t/is (nil? (:error out1)))
(t/is (nil? (:error out2)))
(let [result-1 (:result out1)
result-2 (:result out2)]
;; now we proceed to manually mark all storage objects touched
(th/db-exec! ["update storage_object set touched_at=?" (ct/now)])
;; run the touched gc task
(let [res (binding [ct/*clock* (ct/fixed-clock (ct/in-future {:hours 3}))]
(th/run-task! :storage-gc-touched {}))]
(t/is (= 2 (:freeze res)))
(t/is (= 0 (:delete res))))
;; check that we have all object in the db
(let [rows (th/db-exec! ["select * from storage_object"])]
(t/is (= 2 (count rows)))))
;; now we proceed to manually delete all file_media_object
(th/db-exec! ["update file_media_object set deleted_at = ?" (ct/now)])
(let [res (th/run-task! :objects-gc {})]
(t/is (= 2 (:processed res))))
;; run the touched gc task
(let [res (binding [ct/*clock* (ct/fixed-clock (ct/in-future {:hours 3}))]
(th/run-task! :storage-gc-touched {}))]
(t/is (= 0 (:freeze res)))
(t/is (= 2 (:delete res))))
;; check that we have all no objects
(let [rows (th/db-exec! ["select * from storage_object where deleted_at is null"])]
(t/is (= 0 (count rows))))))
(t/deftest tempfile-bucket-test
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
content1 (sto/content "content1")
now (ct/now)
object1 (sto/put-object! storage {::sto/content content1
::sto/touched-at (ct/plus now {:hours 1})
:bucket "tempfile"
:content-type "text/plain"})]
;; not eligible while the touched-at is in the future
(binding [ct/*clock* (ct/fixed-clock now)]
(let [res (th/run-task! :storage-gc-touched {})]
(t/is (= 0 (:freeze res)))
(t/is (= 0 (:delete res)))))
;; still not eligible: touched-at (now+1h) is beyond the threshold
(binding [ct/*clock* (ct/fixed-clock (ct/plus now {:hours 2}))]
(let [res (th/run-task! :storage-gc-touched {})]
(t/is (= 0 (:freeze res)))
(t/is (= 0 (:delete res)))))
;; eligible: marked for deletion immediately, without any extra delay
(let [clock (ct/plus now {:hours 3})]
(binding [ct/*clock* (ct/fixed-clock clock)]
(let [res (th/run-task! :storage-gc-touched {})]
(t/is (= 0 (:freeze res)))
(t/is (= 1 (:delete res)))))
(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 clock {:seconds 1})))))
;; removed on the next deleted gc run
(binding [ct/*clock* (ct/fixed-clock (ct/plus now {:hours 4}))]
(let [res (th/run-task! :storage-gc-deleted {})]
(t/is (= 1 (:deleted res)))))))
(t/deftest touched-gc-task-skip-delay
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
content (sto/content "content1")
now (ct/now)
object1 (sto/put-object! storage {::sto/content content
::sto/touched-at now
:bucket "tempfile"
:content-type "text/plain"})]
;; too recent: not processed without skip-delay
(binding [ct/*clock* (ct/fixed-clock now)]
(let [res (th/run-task! :storage-gc-touched {})]
(t/is (= 0 (:freeze res)))
(t/is (= 0 (:delete res)))))
;; processed immediately with skip-delay
(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)))))
;; and marked for deletion without any additional delay
(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))
content (sto/content "content1")
object (sto/put-object! storage {::sto/content content
:content-type "text/plain"})]
;; mark as deleted right now
(th/db-exec! ["update storage_object set deleted_at = ?" (ct/now)])
;; the deleted gc removes it on the next run
(let [res (th/run-task! :storage-gc-deleted {})]
(t/is (= 1 (:deleted res))))))
(t/deftest objects-gc-task-skip-delay
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
prof (th/create-profile* 1)
proj (th/create-project* 1 {:profile-id (:id prof)
:team-id (:default-team-id prof)})
file (th/create-file* 1 {:profile-id (:id prof)
:project-id (:default-project-id prof)
:is-shared false})
mfile {:filename "sample.jpg"
:path (th/tempfile "backend_tests/test_files/sample.jpg")
:mtype "image/jpeg"
:size 312043}
params {::th/type :upload-file-media-object
::rpc/profile-id (:id prof)
:file-id (:id file)
:is-local true
:name "testfile"
:content mfile}
out1 (th/command! params)
out2 (th/command! params)]
(t/is (nil? (:error out1)))
(t/is (nil? (:error out2)))
(let [result-1 (:result out1)
result-2 (:result out2)]
;; mark as deleted but in the future (not yet eligible)
(th/db-update! :file-media-object
{:deleted-at (ct/in-future {:days 1})}
{:id (:id result-1)})
;; without skip-delay the future deleted row is not processed
(let [res (th/run-task! :objects-gc {})]
(t/is (= 0 (:processed res))))
;; with skip-delay it is processed immediately
(let [res (th/run-task! :objects-gc {:skip-delay true})]
(t/is (= 1 (:processed res)))))))
(t/deftest put-object-write-failure-leaves-pending-row
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
;; Point the fs backend at a path that is actually a file so the
;; blob write fails.
blocked (fs/path "/tmp/penpot" (str "blocked-" (uuid/next)))
_ (spit (str blocked) "x")]
(try
(let [broken (assoc-in storage [::sto/backends :fs ::sto.fs/directory] (str blocked))
content (sto/content "content")
ex (try
(sto/put-object! broken {::sto/content content
:content-type "text/plain"})
nil
(catch Throwable cause cause))]
(t/is (some? ex))
;; the pending row stays behind and is reclaimed asynchronously
;; by the :storage-pending-gc task
(let [rows (th/db-query :storage-object {:status "pending"})]
(t/is (= 1 (count rows)))
(th/db-update! :storage-object
{:created-at (ct/in-past {:days 2})}
{:id (:id (first rows))})
(let [res (th/run-task! :storage-pending-gc {})]
(t/is (= 1 (:processed res))))
(let [row (th/db-exec-one! ["select count(*) from storage_object"])]
(t/is (= 0 (:count row))))))
(finally
(fs/delete blocked)))))
(t/deftest pending-gc-reclaims-unpromoted-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"})
path (sto/get-object-path storage object)]
;; valid objects are never reclaimed
(let [res (th/run-task! :storage-pending-gc {})]
(t/is (= 0 (:processed res))))
;; simulate a crash: the object was created but never promoted
(th/db-update! :storage-object {:status "pending"
:created-at (ct/in-past {:days 2})}
{:id (:id object)})
(t/is (fs/exists? path))
(let [res (th/run-task! :storage-pending-gc {})]
(t/is (= 1 (:processed res))))
;; both the row and the orphaned blob are removed
(let [row (th/db-exec-one! ["select count(*) from storage_object where id = ?" (:id object)])]
(t/is (= 0 (:count row))))
(t/is (not (fs/exists? path)))))
(t/deftest pending-objects-excluded-from-gc-touched
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
content (sto/content "content")
object (sto/put-object! storage {::sto/content content
::sto/touched-at (ct/now)
:content-type "text/plain"})]
;; mark it pending and touched in the past
(th/db-update! :storage-object {:status "pending"
:touched-at (ct/in-past {:days 1})}
{:id (:id object)})
(binding [ct/*clock* (ct/fixed-clock (ct/now))]
(let [res (th/run-task! :storage-gc-touched {})]
(t/is (= 0 (:freeze res)))
(t/is (= 0 (:delete res)))))
;; still present and not marked as deleted
(let [row (th/db-exec-one! ["select * from storage_object where id = ?" (:id object)])]
(t/is (some? row))
(t/is (nil? (:deleted-at row))))))
(t/deftest pending-objects-excluded-from-get
(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"})]
(t/is (some? (sto/get-object storage (:id object))))
(th/db-update! :storage-object {:status "pending"} {:id (:id object)})
(t/is (nil? (sto/get-object storage (:id object))))))
(t/deftest pending-objects-excluded-from-dedup
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
content (-> (sto/content "content")
(sto/wrap-with-hash "same-hash"))
object1 (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
:bucket "file-media-object"
:content-type "text/plain"})]
;; mark the only matching row as pending
(th/db-update! :storage-object {:status "pending"} {:id (:id object1)})
(let [object2 (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
:bucket "file-media-object"
:content-type "text/plain"})]
(t/is (not= (:id object1) (:id object2))))))
(t/deftest dedup-reuses-existing-blob
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
content (-> (sto/content "content")
(sto/wrap-with-hash "same-hash"))
object1 (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
:bucket "file-media-object"
:content-type "text/plain"})
object2 (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
:bucket "file-media-object"
:content-type "text/plain"})]
(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-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))
content (-> (sto/content "content")
(sto/wrap-with-hash "same-hash"))
object1 (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
:bucket "file-media-object"
:content-type "text/plain"})]
;; remove the physical blob to simulate a stale/broken object
(let [path (sto/get-object-path storage object1)]
(fs/delete path))
;; re-uploading identical content repairs the same reference in place
(let [object2 (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
:bucket "file-media-object"
:content-type "text/plain"})]
(t/is (= (:id object1) (:id object2)))
;; the row stays live: no tombstone and no extra row
(let [row (th/db-exec-one! ["select status, deleted_at from storage_object where id = ?" (:id object1)])]
(t/is (= "valid" (:status row)))
(t/is (nil? (:deleted-at row))))
(let [row (th/db-exec-one! ["select count(*) from storage_object"])]
(t/is (= 1 (:count row))))
;; the repaired blob is readable again under the original id
(t/is (= "content" (slurp (sto/get-object-data storage object2)))))))
(t/deftest gc-deleted-removes-broken-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"})]
;; mark as deleted and remove the physical blob
(th/db-update! :storage-object {:deleted-at (ct/in-past {:minutes 1})}
{:id (:id object)})
(let [path (sto/get-object-path storage object)]
(fs/delete path))
;; the deleted gc removes the row without error even though the blob is
;; missing (the physical deletion is best-effort)
(let [res (th/run-task! :storage-gc-deleted {})]
(t/is (= 1 (:deleted res))))
(let [row (th/db-exec-one! ["select count(*) from storage_object where id = ?" (:id object)])]
(t/is (= 0 (:count row))))))
(t/deftest pending-objects-excluded-from-gc-deleted
(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"})]
;; mark as pending + deleted in the past
(th/db-update! :storage-object {:status "pending"
:deleted-at (ct/in-past {:minutes 1})}
{:id (:id object)})
;; gc-deleted skips it because status != 'valid'
(let [res (th/run-task! :storage-gc-deleted {})]
(t/is (= 0 (:deleted res))))
;; row still exists (with deleted_at set — we set it above)
(let [row (th/db-exec-one! ["select count(*) from storage_object where id = ?"
(:id object)])]
(t/is (= 1 (:count row))))))
(t/deftest gc-deleted-gives-up-after-max-attempts
(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"})]
(th/db-update! :storage-object {:deleted-at (ct/in-past {:minutes 1})
:deletion_attempts 6}
{:id (:id object)})
(with-mocks [_mock {:target 'app.storage.impl/del-objects-in-bulk
:return (fn [_ ids] (set ids))}]
(let [res (th/run-task! :storage-gc-deleted {})]
(t/is (= 0 (:deleted res)))))
(let [row (th/db-exec-one! ["select count(*) from storage_object where id = ?" (:id object)])]
(t/is (= 0 (:count row))))))
(t/deftest dedup-reuses-existing-blob-with-touch
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
content (-> (sto/content "content")
(sto/wrap-with-hash "same-hash"))
t0 (ct/now)
params {::sto/deduplicate? true
::sto/touch true
:bucket "file-media-object"
:content-type "text/plain"}
object1 (binding [ct/*clock* (ct/fixed-clock t0)]
(sto/put-object! storage (assoc params ::sto/content content)))]
;; a touched hit reuses the object and updates its touched_at
(let [object2 (binding [ct/*clock* (ct/fixed-clock (ct/plus t0 {:hours 1}))]
(sto/put-object! storage (assoc params ::sto/content content)))]
(t/is (= (:id object1) (:id object2)))
(let [row (th/db-exec-one! ["select touched_at from storage_object where id = ?" (:id object1)])]
(t/is (ct/is-after? (:touched-at row) t0))))
;; with the blob removed, the touched hit repairs the stale row in
;; place: the same id is kept, the row is not deleted and touched_at
;; is left untouched (the touch flag only applies to healthy hits)
(let [path (sto/get-object-path storage object1)]
(fs/delete path))
(let [object3 (binding [ct/*clock* (ct/fixed-clock (ct/plus t0 {:hours 2}))]
(sto/put-object! storage (assoc params ::sto/content content)))]
(t/is (= (:id object1) (:id object3)))
(let [row (th/db-exec-one! ["select deleted_at, touched_at from storage_object where id = ?" (:id object1)])]
(t/is (nil? (:deleted-at row)))
;; the touch flag does not apply to repairs: touched_at was last
;; set by the healthy hit and is not bumped by the repair
(t/is (ct/is-before? (:touched-at row) (ct/plus t0 {:hours 2})))))))
(t/deftest put-object-repair-failure-leaves-row-intact
(let [storage (-> (:app.storage/storage th/*system*)
(configure-storage-backend))
content (-> (sto/content "content")
(sto/wrap-with-hash "same-hash"))
object (sto/put-object! storage {::sto/content content
::sto/deduplicate? true
:bucket "file-media-object"
:content-type "text/plain"})
path (sto/get-object-path storage object)
;; Point the fs backend at a path that is actually a file so the
;; blob write fails.
blocked (fs/path "/tmp/penpot" (str "blocked-" (uuid/next)))
_ (spit (str blocked) "x")]
(try
;; remove the physical blob to simulate a stale/broken object
(fs/delete path)
(let [broken (assoc-in storage [::sto/backends :fs ::sto.fs/directory] (str blocked))
ex (try
(sto/put-object! broken {::sto/content content
::sto/deduplicate? true
:bucket "file-media-object"
:content-type "text/plain"})
nil
(catch Throwable cause cause))]
(t/is (some? ex))
;; the failed repair leaves the original row exactly as it was:
;; live and valid, so a later upload can retry the healing
(let [row (th/db-exec-one! ["select status, deleted_at from storage_object where id = ?" (:id object)])]
(t/is (= "valid" (:status row)))
(t/is (nil? (:deleted-at row))))
(let [row (th/db-exec-one! ["select count(*) from storage_object"])]
(t/is (= 1 (:count row)))))
(finally
(fs/delete blocked)))))
(t/deftest upload-chunks-exclude-pending
(let [prof (th/create-profile* 1)
_ (th/create-project* 1 {:profile-id (:id prof)
:team-id (:default-team-id prof)})
file (th/create-file* 1 {:profile-id (:id prof)
:project-id (:default-project-id prof)
:is-shared false})
mfile {:filename "chunk"
:path (th/tempfile "backend_tests/test_files/sample.jpg")
:mtype "image/jpeg"
:size 312043}
session-id (-> (th/command! {::th/type :create-upload-session
::rpc/profile-id (:id prof)
:total-chunks 1})
:result :session-id)
out (th/command! {::th/type :upload-chunk
::rpc/profile-id (:id prof)
:session-id session-id
:index 0
:content mfile})]
(t/is (nil? (:error out)))
;; mark all the chunks of this session as pending (simulates rows that
;; were never promoted)
(th/db-exec! ["update storage_object set status = 'pending' where id in (select object_id from upload_session_chunk where session_id = ?)"
session-id])
;; assembling fails because no chunk is visible anymore
(let [assemble-out (th/command! {::th/type :assemble-file-media-object
::rpc/profile-id (:id prof)
:session-id session-id
:file-id (:id file)
:is-local true
:name "assembled-image"
:mtype "image/jpeg"})]
(t/is (some? (:error assemble-out))))))
(t/deftest upload-session-stalled-purge-lifecycle
;; Full lifecycle of a stalled session: objects-gc purges the session and
;; its mappings while touching the objects, touched-gc marks them deleted
;; and deleted-gc removes rows and blobs.
(let [prof (th/create-profile* 1)
_ (th/create-project* 1 {:profile-id (:id prof)
:team-id (:default-team-id prof)})
_ (th/create-file* 1 {:profile-id (:id prof)
:project-id (:default-project-id prof)
:is-shared false})
mfile {:filename "chunk"
:path (th/tempfile "backend_tests/test_files/sample.jpg")
:mtype "image/jpeg"
:size 312043}
session-id (-> (th/command! {::th/type :create-upload-session
::rpc/profile-id (:id prof)
:total-chunks 1})
:result :session-id)
out (th/command! {::th/type :upload-chunk
::rpc/profile-id (:id prof)
:session-id session-id
:index 0
:content mfile})]
(t/is (nil? (:error out)))
(t/is (= 1 (:count (th/db-exec-one! ["select count(*) from upload_session_chunk where session_id = ?"
session-id]))))
;; backdate the session so it counts as stalled
(th/db-exec! ["update upload_session set created_at = now() - interval '2 hours' where id = ?"
session-id])
;; objects-gc purges session and mappings, touching the objects
(let [res (th/run-task! :objects-gc {})]
(t/is (= 1 (:processed res))))
(t/is (= 0 (:count (th/db-exec-one! ["select count(*) from upload_session where id = ?"
session-id]))))
(t/is (= 0 (:count (th/db-exec-one! ["select count(*) from upload_session_chunk where session_id = ?"
session-id]))))
(t/is (= 1 (:count (th/db-exec-one! ["select count(*) from storage_object where touched_at is not null"]))))
;; touched-gc marks the orphaned object as deleted
(let [res (binding [ct/*clock* (ct/fixed-clock (ct/in-future {:hours 3}))]
(th/run-task! :storage-gc-touched {}))]
(t/is (= 0 (:freeze res)))
(t/is (= 1 (:delete res))))
;; deleted-gc removes the row and the blob (clock past the mark time)
(let [res (binding [ct/*clock* (ct/fixed-clock (ct/in-future {:hours 4}))]
(th/run-task! :storage-gc-deleted {}))]
(t/is (= 1 (:deleted res))))
(t/is (= 0 (:count (th/db-exec-one! ["select count(*) from storage_object"]))))))
(t/deftest upload-session-consumed-purge
;; An assembled session is marked as consumed and objects-gc purges it
;; right away, without waiting for the stalled threshold.
(let [prof (th/create-profile* 1)
_ (th/create-project* 1 {:profile-id (:id prof)
:team-id (:default-team-id prof)})
file (th/create-file* 1 {:profile-id (:id prof)
:project-id (:default-project-id prof)
:is-shared false})
mfile {:filename "chunk"
:path (th/tempfile "backend_tests/test_files/sample.jpg")
:mtype "image/jpeg"
:size 312043}
session-id (-> (th/command! {::th/type :create-upload-session
::rpc/profile-id (:id prof)
:total-chunks 1})
:result :session-id)
out (th/command! {::th/type :upload-chunk
::rpc/profile-id (:id prof)
:session-id session-id
:index 0
:content mfile})]
(t/is (nil? (:error out)))
(let [assemble-out (th/command! {::th/type :assemble-file-media-object
::rpc/profile-id (:id prof)
:session-id session-id
:file-id (:id file)
:is-local true
:name "assembled-image"
:mtype "image/jpeg"})]
(t/is (nil? (:error assemble-out))))
;; mappings stay, session row stays marked as consumed; objects-gc
;; purges both
(t/is (= 1 (:count (th/db-exec-one! ["select count(*) from upload_session_chunk where session_id = ?"
session-id]))))
(t/is (some? (:deleted-at (th/db-exec-one! ["select deleted_at from upload_session where id = ?"
session-id]))))
;; objects-gc purges the consumed session immediately
(let [res (th/run-task! :objects-gc {})]
(t/is (= 1 (:processed res))))
(t/is (= 0 (:count (th/db-exec-one! ["select count(*) from upload_session where id = ?"
session-id]))))
(t/is (= 0 (:count (th/db-exec-one! ["select count(*) from upload_session_chunk where session_id = ?"
session-id]))))))
(t/deftest upload-session-profile-purge
;; Sessions owned by a profile pending purge are drained first, so the
;; profile delete (which cascades to its sessions) never hits the chunk
;; RESTRICT foreign keys.
(let [prof (th/create-profile* 1)
mfile {:filename "chunk"
:path (th/tempfile "backend_tests/test_files/sample.jpg")
:mtype "image/jpeg"
:size 312043}
session-id (-> (th/command! {::th/type :create-upload-session
::rpc/profile-id (:id prof)
:total-chunks 1})
:result :session-id)
out (th/command! {::th/type :upload-chunk
::rpc/profile-id (:id prof)
:session-id session-id
:index 0
:content mfile})]
(t/is (nil? (:error out)))
;; soft-delete the profile; the live session is neither consumed nor stalled
(th/db-update! :profile {:deleted-at (ct/now)} {:id (:id prof)})
(th/run-task! :objects-gc {})
;; session and mappings are gone, profile row deletes cleanly
(t/is (= 0 (:count (th/db-exec-one! ["select count(*) from upload_session where id = ?"
session-id]))))
(t/is (= 0 (:count (th/db-exec-one! ["select count(*) from upload_session_chunk where session_id = ?"
session-id]))))
(t/is (= 0 (:count (th/db-exec-one! ["select count(*) from profile where id = ?"
(:id prof)]))))
;; and the chunk object was touched for the storage GC
(t/is (= 1 (:count (th/db-exec-one! ["select count(*) from storage_object where touched_at is not null"]))))))
(defn- fake-s3-backend
[]
{::sto/type :s3
::sto.s3/client (reify S3AsyncClient)
::sto.s3/presigner (reify S3Presigner)})
(t/deftest s3-exists-object-returns-true-on-found
(with-mocks [mock {:target 'app.storage.s3/head-object
:return (p/resolved {})}]
(t/is (true? (impl/exists-object? (fake-s3-backend) {:id (uuid/next)})))
(t/is (= 1 (:call-count @mock)))))
(t/deftest s3-exists-object-returns-false-on-missing-key
(with-mocks [mock {:target 'app.storage.s3/head-object
:return (p/rejected (-> (NoSuchKeyException/builder)
(.message "no key")
(.build)))}]
(t/is (false? (impl/exists-object? (fake-s3-backend) {:id (uuid/next)})))
;; a missing key is a definitive answer: no retries
(t/is (= 1 (:call-count @mock)))))
(t/deftest s3-exists-object-retries-transient-errors
(let [calls (atom 0)]
(with-mocks [_mock {:target 'app.storage.s3/head-object
:return (fn [& _]
(swap! calls inc)
(if (< @calls 3)
(p/rejected (RuntimeException. "boom"))
(p/resolved {})))}]
(t/is (true? (impl/exists-object? (fake-s3-backend) {:id (uuid/next)})))
(t/is (= 3 @calls)))))
(t/deftest s3-exists-object-throws-after-retries-exhausted
(with-mocks [mock {:target 'app.storage.s3/head-object
:return (p/rejected (RuntimeException. "boom"))}]
;; p/await returns the rejection wrapped in an ExecutionException
(let [ex (try
(impl/exists-object? (fake-s3-backend) {:id (uuid/next)})
nil
(catch Throwable cause cause))]
(t/is (some? ex))
(t/is (= "boom" (ex-message (ex-cause ex)))))
;; one initial attempt plus max-retries
(t/is (= 4 (:call-count @mock)))))