mirror of
https://github.com/penpot/penpot.git
synced 2026-10-03 17:26:16 +00:00
* ✨ 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
1239 lines
55 KiB
Clojure
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)))))
|