feat(metadata): reconcile artwork cache after public S3 provider changes (#349)

* feat(metadata): reconcile artwork cache after public S3 provider changes

Changing the public S3 provider previously broke every cached image
permanently: the DB keeps bucket-relative keys, the image cache pipeline
treats a cached path as its durable dedup marker and never re-enqueues,
and clients eat the 404s straight from S3 so the server never notices.

Add a storage identity fingerprint (s3.public_storage_identity, seeded
via SetIfAbsent at boot) and a reconcile_artwork_cache task whose
startup trigger only fires when the identity changed; manual runs
always sweep, doubling as bucket-data-loss recovery. The task probes a
random sample of cached objects, then either bulk-resets (near-total
miss) or per-row verifies. Missing provider-sourced artwork is reset to
its *_source_path so the existing enqueue loop re-caches it; surfaces
without a re-downloadable source (chapter thumbnails, collection
artwork, library posters, branding refs, embedded book covers) are
cleared so their owning pipelines refill them. Small upload-holding
tables are always per-row verified so bulk mode cannot blind-clear an
upload that survived migration, and transport errors never reset rows.

Users never see broken images during the transition: reset rows serve
the provider's original URL via the existing absolute-URL pass-through
and thumbhashes are preserved. The storage settings page now warns that
uploads cannot be re-downloaded when the identity fields are edited.

Part of #348

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(metadata): harden artwork reconcile per code review

Address the confirmed findings from the PR review:

- Fingerprint the key prefix case-sensitively and slash-trimmed exactly
  as s3client applies it (new exported NormalizeKeyPrefix): a case-only
  prefix edit is a real storage move and must reconcile; a slash-only
  edit is not and must not.
- Certify the storage fingerprint immediately after the artwork sweep
  succeeds and make the 4-object branding check non-fatal (reported in
  the task message), so a transient branding error cannot discard a
  completed catalog sweep and force it to repeat every boot.
- Fail closed on conditional-task preflight errors in the task manager
  (previously fail-open ran the task), and retry transient settings
  reads in ShouldRun since the startup trigger fires once per process.
- Track probe HEAD errors against a separate baseline so a flaky probe
  cannot consume the sweep's error budget.
- Probe before counting: bulk mode skips the per-surface count(*)
  full scans entirely, and probe sampling drops ORDER BY random()
  (plain LIMIT answers "is the cache in this bucket" just as well).
- Verify chapter thumbnails across a whole 500-file batch in one HEAD
  fan-out instead of per file, keeping the worker pool saturated.
- Replace the 10 inline non-provider-scheme ARRAY literals in the
  enqueue query with the shared nonProviderImageSchemesSQL constant.

Part of #348

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(metadata): guard bulk reset against degraded probes, certify only clean sweeps

Address bot review feedback on the reconcile hardening:

- A probe where more than half the HEAD requests error aborts the run:
  errored requests are excluded from the sample, so a partial outage
  could otherwise present a handful of surviving 404s as a ~100% miss
  rate and bulk-reset the catalog. Bulk mode additionally requires a
  minimum number of successful samples; thinned probes and tiny
  catalogs take the safe per-row verify path.
- Track sweep errors separately from probe/branding errors
  (stats.sweep_errors) and certify the storage fingerprint only when
  the sweep completed with zero of them — skipped rows were never
  verified, so the next startup retries. Applied resets stay durable.
- Give each ObjectExists attempt its own timeout so a stalled HEAD
  fails that attempt instead of pinning the retry loop to the run
  context.
- Report branding assets checked (not just cleared) in stats.Checked.
- Drop the dead settingsRepo/brandingSvc nil guards in cmd/silo and
  sync spec numbers with the implementation constants.

Part of #348

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Quick
2026-07-09 12:00:24 -04:00
committed by GitHub
co-authored by Claude Fable 5
parent 4d99597966
commit 04c4344f52
13 changed files with 1925 additions and 30 deletions
+35 -13
View File
@@ -1875,6 +1875,17 @@ func main() {
}
}
// White-label branding: one service shared by the API (public read + admin
// upload), the frontend handler (index.html title, favicon, manifest), and
// the artwork reconcile task. S3 is optional — pass a nil AssetStore (not
// the typed-nil *s3client.Client) when it isn't configured so text branding
// still works without it.
var brandingStore branding.AssetStore
if deps.S3Public != nil {
brandingStore = deps.S3Public
}
brandingSvc := branding.NewService(settingsRepo, brandingStore)
// Wire up task manager for admin task API.
if needsWorkers && deps.DB != nil {
triggerRepo := taskrepository.NewPgTriggerRepository(deps.DB)
@@ -1954,6 +1965,26 @@ func main() {
if metadataImageCacheProcessor != nil {
taskMgr.Register(tasks.NewCacheMetadataImagesTask(metadataImageCacheProcessor))
}
if deps.S3Public != nil {
identity := tasks.ArtworkStorageIdentity(cfg.S3.Public.Endpoint, cfg.S3.Public.Bucket, cfg.S3.Public.KeyPrefix)
// Seed the fingerprint on first boot so an unchanged storage
// identity never triggers a sweep. On the boot after a provider
// change the stored (old) identity survives this call and the
// startup trigger runs the reconcile.
if _, err := settingsRepo.SetIfAbsent(appCtx, tasks.ArtworkStorageIdentityKey, identity); err != nil {
slog.Warn("artwork reconcile: seeding storage identity failed", "error", err)
}
var brandingReconciler tasks.BrandingAssetReconciler
if brandingSvc != nil && brandingSvc.HasStorage() {
brandingReconciler = brandingSvc
}
taskMgr.Register(tasks.NewReconcileArtworkCacheTask(
metadata.NewArtworkCacheReconciler(deps.DB, deps.S3Public),
settingsRepo,
brandingReconciler,
identity,
))
}
if pluginAutoUpdater != nil {
taskMgr.Register(tasks.NewCheckPluginUpdatesTask(pluginAutoUpdater))
}
@@ -2232,19 +2263,10 @@ func main() {
deps.FrontendFS = distFS
server.WebDistFS = distFS
// White-label branding: one service shared by the API (public read + admin
// upload) and the frontend handler (index.html title, favicon, manifest).
// S3 is optional — pass a nil AssetStore (not the typed-nil *s3client.Client)
// when it isn't configured so text branding still works without it.
if settingsRepo != nil {
var brandingStore branding.AssetStore
if deps.S3Public != nil {
brandingStore = deps.S3Public
}
brandingSvc := branding.NewService(settingsRepo, brandingStore)
deps.BrandingService = brandingSvc
server.Branding = brandingSvc
}
// Expose the branding service (constructed before the task manager) to the
// API and the frontend handler.
deps.BrandingService = brandingSvc
server.Branding = brandingSvc
router := api.NewRouter(deps)
@@ -0,0 +1,265 @@
# S3 Provider Change: Artwork Cache Reconciliation Design
Commands and paths in this document assume the repository root is the cwd.
## Goal
When an admin changes the public S3 provider (endpoint, bucket, or key prefix), cached
artwork must heal automatically. Today the change silently breaks every cached image
forever: the database keeps pointing at object keys that only exist in the old bucket,
the caching pipeline permanently skips already-cached rows, and clients eat the 404s
directly from S3 so the server never even observes the breakage.
Design goals, in priority order:
1. **Users never see broken images.** During and after a provider change, every image
either serves from the new bucket, falls back to its provider source URL, or falls
back to its thumbhash placeholder / generated artwork.
2. **The admin makes no decisions and cannot get it wrong.** No "did you migrate your
data?" dialog. Silo verifies reality and does the right thing for migrated,
unmigrated, and partially migrated buckets alike.
3. **No wasted work.** An admin who migrated their objects to the new bucket must not
trigger a full catalog re-download (provider rate limits, bandwidth, time).
## Background: how caching works today
- The image cache pipeline (`internal/imagecache`, `internal/metadata/image_cache_processor.go`)
downloads provider artwork, generates variants (`original`, `w500`, `w300`, ...), and
uploads them to the **public** S3 bucket under bucket-relative keys such as
`tmdb/movies/550/poster/original.webp`.
- On success, the target row's path column (`media_items.poster_path`,
`episodes.still_path`, `people.photo_path`, ...) is rewritten from the provider URL to
the cached relative key. The original provider URL is preserved in the matching
`*_source_path` column. Thumbhashes are stored in their own columns and are
content-derived, so they remain valid across a re-cache of the same source.
- The recurring `cache_metadata_images` task (60s interval,
`internal/taskmanager/tasks/cache_metadata_images.go`) calls
`EnqueueExistingProviderArtwork` (`internal/metadata/image_cache_job_repo.go`), which
only enqueues rows whose destination path is still a provider URL (`LIKE '%://%'`) or
empty. **The cached path in the row is the durable dedup marker** — once set, the row
is never enqueued again.
- Serving: stored relative keys are resolved to presigned / read-endpoint URLs against
the *current* S3 config at request time (`internal/metadata/image_resolver.go`,
`PresignURLWithExpiry` in `internal/catalog/detail.go`). Absolute `http(s)://` paths
pass through to clients unchanged. Clients then fetch from S3 directly.
- All `s3.` settings keys are restart-required (`internal/config/restart_keys.go`), so a
provider change only takes effect at the next server start.
Consequences of a provider change with no migration: every resolved URL points into the
new bucket where nothing exists → broken images; nothing re-enqueues; recovery requires
hand-written SQL. If the admin *did* migrate objects, everything keeps working — which is
why a blind "wipe and re-cache on settings change" would be strictly worse than today.
## Approach overview
Three pieces:
1. **Storage identity fingerprint** persisted in `server_settings`. At startup a new
reconcile task compares the fingerprint against the live config; a mismatch means the
storage identity changed and a verification sweep is needed.
2. **`reconcile_artwork_cache` task** (task manager, startup trigger + manually
runnable). Probes the new bucket, verifies cached objects, and resets rows whose
objects are missing so the existing pipeline rebuilds them. Recovery action depends on
the image's source scheme (see table below).
3. **Admin messaging**: an informational note when saving changed S3 settings, live task
progress in the existing Tasks UI, and a completion notification summarizing what was
verified, re-queued, and unrecoverable.
There is deliberately **no new serving logic**: resetting a path back to its provider
source URL restores the pre-cache pass-through behavior, so images render (hotlinked)
immediately and flip back to cached S3 URLs as the queue drains.
## Storage identity fingerprint
- Value: normalized `endpoint|bucket|key_prefix` of the **public write** config
(`cfg.S3.Public.Endpoint`, `.Bucket`, `.KeyPrefix`), lowercased, trimmed, stored as a
plain (non-encrypted) `server_settings` row, e.g. key `s3.public_storage_identity`.
The read endpoint and URL-auth settings are excluded on purpose: changing how objects
are *served* (CDN, token auth) does not move the stored data.
- Bootstrap: if no fingerprint row exists (first run after upgrade), adopt the current
identity without reconciling. Admins who changed providers before this feature existed
can run the task manually.
- The fingerprint is updated to the new identity **only after a sweep completes**. An
interrupted sweep therefore re-triggers at the next startup; the sweep is idempotent.
## The `reconcile_artwork_cache` task
Task manager registration: key `reconcile_artwork_cache`, category Metadata, visible,
default trigger `TriggerTypeStartup`, manually runnable from the Tasks UI. On startup it
compares fingerprints and exits immediately (fast no-op) when they match. A manual run
skips the fingerprint check and always sweeps — this doubles as disaster recovery for
"my bucket was wiped / partially lost" with no provider change at all.
### Phase 1 — probe
The probe runs first (before any counting): HEAD (via the existing
`s3client.Client.ObjectExists`) a sample of ~200 stored cached keys across all
surfaces, taken with plain `LIMIT` sampling — `ORDER BY random()` would full-scan and
sort every surface, and the probe only has to answer "does this bucket hold the cache
at all". Because the DB stores the full key of the `original` variant (e.g.
`tmdb/movies/550/poster/original.webp`), the stored path is exactly the key to check —
no guessing. Progress-denominator `count(*)` queries run only in verify mode; bulk
mode reports `RowsAffected` and never pays for counts. Probe errors are reported but
tracked separately (`sweep_errors` vs `errors`) so they cannot consume the sweep's
error budget. Two reliability guards protect bulk mode from a degraded probe: a run
aborts outright when more than half the probe requests error (an outage is not a
miss), and bulk reset requires a minimum number of *successful* samples
(`artworkReconcileBulkMinSample`) — a probe thinned out by transport errors takes the
safe per-row path instead.
- **≥95% missing** → "bulk reset" mode: skip per-row verification and reset all cached
rows by SQL alone (the data plainly was not migrated). The threshold is not 100% so a
handful of coincidentally-present keys cannot force millions of pointless HEADs.
- **All present** → data was migrated; run the full sweep anyway (it is one cheap HEAD
per image and catches partially-missed prefixes) but expect near-zero resets.
- **Mixed** → full per-row sweep.
### Phase 2 — sweep and reset
Iterate every surface that stores cached public-bucket keys, in batches with keyset
pagination and a bounded HEAD concurrency (`artworkReconcileHeadWorkers`, 16 in
flight, each attempt under its own timeout), reporting progress
through the task's `ProgressReporter`:
| Surface | Cached columns | Recovery when object is missing |
| --- | --- | --- |
| `media_items` | `poster_path`, `backdrop_path`, `logo_path` | Reset column to its `*_source_path` when the source is a re-downloadable provider URL |
| `media_item_localizations` | `poster_path`, `backdrop_path`, `logo_path` | Same |
| `seasons`, `season_localizations` | `poster_path` | Same |
| `episodes` | `still_path` | Same |
| `people` | `photo_path` | Same |
| Chapter thumbnails (`media_files.chapters` JSONB) | per-chapter `thumbnail_path` | Clear path + thumbhash + retry state; the 6-hour `chapter_thumbnail_backfill` task regenerates from the media file |
| Collections (`library_collections`, `user_personal_collections`) | `poster_url` / `backdrop_url` | Clear (plus `poster_auto_generated` / `poster_from_template` flags) → UI falls back to the generated collage/poster; report. The original template reference is not stored, so template posters are not re-uploaded automatically |
| Server branding (`internal/branding`) | `server_settings` refs (`branding.*_ref`) | Verified by `branding.Service.ReconcileMissingAssets`; missing refs cleared → built-in defaults; report |
| Library artwork (`media_folders.poster_path`) | `library-posters/{id}.ext` | Clear → default tile; report (no auto-refill) |
| Audiobook / ebook covers (`local/...` keys in `media_items`) | `poster_path` | Clear + null `last_refreshed`; the enrichment sweeps re-extract embedded covers |
Profile avatars are stored in the **private** bucket (`deps.S3Private`), not the public
one, and are therefore out of scope for this reconcile.
Reset semantics for provider-sourced artwork (the overwhelming majority):
- Single guarded UPDATE per row/column:
`SET poster_path = poster_source_path WHERE content_id = $1 AND poster_path = $2` —
the `poster_path = $2` guard makes the sweep safe against concurrent metadata
refreshes (lost-update protection).
- **Thumbhash columns are left untouched.** Placeholders keep rendering throughout.
- Retained `metadata_image_cache_jobs` rows need no special handling:
`EnqueueExistingProviderArtwork` calls its upsert with `requeueSucceeded=true`, which
flips even `succeeded` rows back to `queued`. (An earlier draft of this spec called
for deleting job rows; implementation showed it is unnecessary.)
- Sources with schemes `upload://`, `s3://`, `file://`, `local://`, `generated://` are
never "reset to source" — those rows are cleared instead, and clearing `media_items`
artwork also nulls `last_refreshed` so the book-enrichment sweeps (which require
`last_refreshed IS NULL`) re-extract embedded covers.
### Phase 3 — rebuild (existing machinery, no new code)
Nothing re-downloads inside the reconcile task itself. After reset, the 60-second
`cache_metadata_images` loop sees provider URLs in destination columns again, enqueues
them, and the processor re-caches into the new bucket. The processor's per-variant
`ObjectExists` checks (`internal/imagecache/imagecache.go`) make rebuilds idempotent and
partial-migration-safe: variants that survived migration are skipped, missing ones are
uploaded.
### Phase 4 — report and finalize
- Update the fingerprint row to the new identity — but only when the sweep completed
with zero sweep errors. Skipped-on-error rows were never verified, so a run with
sweep errors reports its partial results (applied resets are durable) and leaves
the fingerprint stale for the next startup to retry.
- Task completion message: verified / re-queued / cleared counts, plus structured
result data (JSON) persisted with the execution record for the Tasks UI.
- No durable inbox notification: the notifications system has no free-form admin
message type (deliveries are content-keyed and per-profile), and adding one for this
feature is not worth the payload-renderer plumbing. The persisted task result plus
warn logs carry the report.
- Log unrecoverable items individually at `warn` with content IDs so an admin can
re-upload the specific artwork that was lost.
## Admin UX
1. Admin edits S3 settings → existing connection check (`admin_settings_checks.go`)
validates the new provider before saving, as today.
2. When the saved values change the storage identity, the settings UI adds one
informational callout next to the existing restart-required notice:
*"Artwork is cached in object storage. After restart, Silo verifies the cache against
the new storage and automatically re-caches anything missing. Uploaded artwork
(custom posters, avatars, branding) cannot be re-downloaded — migrate your bucket
contents if you want to keep them."*
No confirmation dialog, no migration question.
3. After restart the reconcile task appears in the Tasks UI with live progress, and
finishes with the summary + notification described above. It can be re-run manually
at any time.
## End-user UX
- Before the sweep touches a row: broken S3 URL, but the thumbhash placeholder renders.
The in-memory resolved-URL cache is empty after the (required) restart, so no stale
presigned URLs outlive the config change.
- After reset: the API serves the provider's original URL via the existing absolute-URL
pass-through — image present immediately, just hotlinked and un-resized.
- After re-cache: back to normal S3-served variants. Users should not be able to tell a
migration happened.
## Failure modes
- **New provider unreachable at startup**: probe HEADs fail with errors (not
"missing"). The task must distinguish `ObjectExists == false` from a transport error
and abort without resetting anything, leaving the fingerprint stale so it retries next
startup / next manual run. Never mass-reset on the basis of errors.
- **Sweep interrupted** (restart, crash): fingerprint was not updated; the task re-runs
at next startup. Already-reset rows are simply in the "provider URL" state the enqueue
loop handles; already-verified rows get re-HEADed. Idempotent.
- **Key-prefix-only change**: `ObjectExists` resolves keys under the current prefix, so
objects under the old prefix correctly count as missing and re-cache under the new
prefix.
- **Provider rate limits during rebuild**: unchanged from any large-catalog initial
cache; the existing job queue's retry/backoff behavior applies.
## Out of scope
- The **private** S3 bucket (subtitles, downloads, transcode artifacts). A private-
provider change has a different blast radius and mostly regenerable content; it
deserves its own pass. This spec's fingerprint/task mechanism is deliberately reusable
for it later.
- An S3→S3 object migration tool. Admins who want to keep uploaded artwork migrate with
standard tooling (`rclone`, `aws s3 sync`); Silo's job is to heal what it can and
report what it cannot.
- Serve-time 404 fallback. Clients fetch from S3 directly via presigned/read-endpoint
URLs, so the server never sees the misses; a proxying fallback would put every image
request through the server and is rejected on performance grounds.
## Implementation notes (as built)
- `internal/metadata/artwork_reconcile.go` — `ArtworkCacheReconciler`: table-driven
sweep surfaces, probe, bulk/verify modes, and the bespoke chapter-thumbnail sweep.
Small "precious upload" tables (collections, library posters) are marked
`alwaysVerify` so a bulk reset can never blind-clear an upload that survived
migration.
- `internal/taskmanager/tasks/reconcile_artwork_cache.go` — task registration,
`ArtworkStorageIdentity` fingerprint helper, and the `ShouldRun` gate
(`ScheduledConditionalTask` suppresses the startup trigger while the fingerprint
matches; manual runs always sweep).
- `internal/branding/service.go` — `ReconcileMissingAssets` verifies and clears the
four branding refs; composed into the task rather than the reconciler so `metadata`
does not import `branding`. It runs *after* the fingerprint is certified and its
errors are non-fatal (reported in the task message): a transient failure on a
4-object check must never discard a completed catalog sweep and force it to repeat.
- Fingerprint lives in `server_settings` under `s3.public_storage_identity`
(plaintext; encrypted values are GCM-bound to their key name, which complicates any
future rename for zero benefit). It is seeded with `SetIfAbsent` at wiring time in
`cmd/silo`, so first boot adopts the current identity without a sweep. Normalization
mirrors real S3 semantics: endpoint/bucket are case-folded, but the key prefix keeps
its case (S3 keys are case-sensitive — a case-only prefix edit is a real move) and is
slash-trimmed exactly like `s3client.NormalizeKeyPrefix` applies it.
- The task manager's conditional-task preflight fails **closed**: a `ShouldRun` error
skips the run (and the task's `ShouldRun` retries transient settings reads), so a
startup DB blip cannot fail open into a full catalog sweep.
- No schema migration: the fingerprint is a `server_settings` row and all reset
operations use existing columns.
- `s3client.Client` prepends `KeyPrefix` internally, so the reconciler passes DB-stored
logical keys straight to `ObjectExists`.
- Web UI: conditional callout in
`web/src/pages/admin-settings/StorageSettings.tsx`, shown while any of
`s3.public_endpoint` / `s3.public_bucket` / `s3.public_key_prefix` is dirty.
+33
View File
@@ -5,6 +5,7 @@ import (
"crypto/sha256"
"encoding/hex"
"errors"
"log/slog"
"path"
"github.com/Silo-Server/silo-server/internal/s3client"
@@ -138,6 +139,38 @@ func (s *Service) GetAsset(ctx context.Context, kind AssetKind) (data []byte, co
return data, contentTypeForExt(path.Ext(ref)), ref, nil
}
// ReconcileMissingAssets clears the ref of every configured branding asset
// whose stored object no longer exists (e.g. after the public S3 provider
// changed without migrating data), so the UI falls back to the built-in
// defaults instead of serving broken images. Returns how many configured
// assets were checked and how many of those were cleared.
func (s *Service) ReconcileMissingAssets(ctx context.Context) (checked, cleared int, err error) {
if s == nil || !s.HasStorage() {
return 0, 0, nil
}
for kind, spec := range assetSpecs {
ref, _ := s.settings.Get(ctx, spec.settingKey)
if ref == "" {
continue
}
checked++
key := spec.s3Prefix + "/" + ref
_, getErr := s.store.GetObject(ctx, s.store.Bucket(), key)
switch {
case getErr == nil:
case errors.Is(getErr, s3client.ErrNotFound):
if setErr := s.settings.Set(ctx, spec.settingKey, ""); setErr != nil {
return checked, cleared, setErr
}
cleared++
slog.Warn("branding: cleared asset whose stored object is missing", "kind", kind, "key", key)
default:
return checked, cleared, getErr
}
}
return checked, cleared, nil
}
func firstNonEmpty(values ...string) string {
for _, v := range values {
if v != "" {
+769
View File
@@ -0,0 +1,769 @@
package metadata
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"strings"
"sync"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// ArtworkObjectChecker is the S3 surface the reconciler needs: existence
// checks against the public asset bucket. Satisfied by *s3client.Client.
type ArtworkObjectChecker interface {
ObjectExists(ctx context.Context, bucket, key string) (bool, error)
Bucket() string
}
// nonProviderImageSchemesSQL mirrors isNonProviderImageScheme for use inside
// SQL predicates: source paths with these schemes cannot be re-downloaded.
const nonProviderImageSchemesSQL = `ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%']`
const (
artworkReconcileSampleTarget = 200
artworkReconcileBatchSize = 500
artworkReconcileHeadWorkers = 16
artworkReconcileHeadTimeout = 15 * time.Second
artworkReconcileErrorBudget = 200
// artworkReconcileBulkThreshold: when at least this fraction of sampled
// objects is missing, skip per-row verification for the large regenerable
// surfaces and reset every cached row.
artworkReconcileBulkThreshold = 0.95
// artworkReconcileBulkMinSample: bulk reset additionally requires this
// many *successful* probe samples. A probe degraded by transport errors
// (errored requests are excluded from the sample) must not bulk-reset the
// catalog off a handful of surviving 404s; small catalogs below this bar
// simply take the per-row verify path, which is cheap at that size.
artworkReconcileBulkMinSample = 25
)
// ArtworkReconcileStats summarizes one reconcile run.
type ArtworkReconcileStats struct {
Mode string `json:"mode"` // "verify" or "bulk_reset"
Sampled int `json:"sampled"`
SampleMissing int `json:"sample_missing"`
Checked int `json:"checked"`
Verified int `json:"verified"`
Requeued int `json:"requeued"` // reset to provider source; re-cached by the image cache pipeline
Cleared int `json:"cleared"` // no re-downloadable source; refilled by scans/enrichment or re-uploaded by an admin
Errors int `json:"errors"`
// SweepErrors is the subset of Errors from the sweep itself (skipped
// rows). Probe errors don't reduce sweep completeness — probed keys are
// re-checked by the sweep — so callers deciding whether the reconcile
// fully covered the catalog must look here, not at Errors.
SweepErrors int `json:"sweep_errors"`
}
// artworkSweepSurface describes one cached-path column the reconciler sweeps.
type artworkSweepSurface struct {
name string
table string
keyCols []string // pagination key expressions; must form a unique order
pathCol string
// sourceCol holds the original source the row can be reset to. Empty for
// surfaces without a re-downloadable source; their rows are always cleared.
sourceCol string
// clearSet is the SQL SET fragment applied when a row has no usable
// source: it must clear pathCol and whatever companion state the owning
// pipeline needs to refill the image.
clearSet string
// alwaysVerify forces per-row HEAD verification even in bulk-reset mode.
// Used for small tables holding admin/user uploads, where a blind reset
// would discard the last pointer to an object that survived migration.
alwaysVerify bool
}
func (s artworkSweepSurface) cachedPredicate() string {
return fmt.Sprintf(
`coalesce(%s, '') NOT IN ('', '-') AND %s NOT LIKE '%%://%%'`,
s.pathCol, s.pathCol,
)
}
func (s artworkSweepSurface) remoteSourcePredicate() string {
if s.sourceCol == "" {
return "FALSE"
}
return fmt.Sprintf(
`coalesce(%s, '') LIKE '%%://%%' AND lower(%s) NOT LIKE ALL (%s)`,
s.sourceCol, s.sourceCol, nonProviderImageSchemesSQL,
)
}
func (s artworkSweepSurface) resetSet() string {
return fmt.Sprintf(`%s = %s, updated_at = NOW()`, s.pathCol, s.sourceCol)
}
// artworkSweepSurfaces lists every cached-artwork destination in the public
// bucket that lives in a plain table column.
//
// The metadata surfaces are kept in sync with EnqueueExistingProviderArtwork:
// resetting a path column here is what makes that query pick the row up
// again. Clearing media_items artwork also nulls last_refreshed so the book
// enrichment sweeps (which require last_refreshed IS NULL) re-extract
// embedded covers.
//
// Chapter thumbnails (JSONB on media_files) and branding assets
// (server_settings refs) have bespoke sweeps and are not listed here.
func artworkSweepSurfaces() []artworkSweepSurface {
itemClear := func(pathCol string) string {
return fmt.Sprintf(`%s = '', last_refreshed = NULL, updated_at = NOW()`, pathCol)
}
plainClear := func(pathCol string) string {
return fmt.Sprintf(`%s = '', updated_at = NOW()`, pathCol)
}
return []artworkSweepSurface{
{name: "item posters", table: "media_items", keyCols: []string{"content_id"}, pathCol: "poster_path", sourceCol: "poster_source_path", clearSet: itemClear("poster_path")},
{name: "item backdrops", table: "media_items", keyCols: []string{"content_id"}, pathCol: "backdrop_path", sourceCol: "backdrop_source_path", clearSet: itemClear("backdrop_path")},
{name: "item logos", table: "media_items", keyCols: []string{"content_id"}, pathCol: "logo_path", sourceCol: "logo_source_path", clearSet: itemClear("logo_path")},
{name: "localized item posters", table: "media_item_localizations", keyCols: []string{"content_id", "language"}, pathCol: "poster_path", sourceCol: "poster_source_path", clearSet: plainClear("poster_path")},
{name: "localized item backdrops", table: "media_item_localizations", keyCols: []string{"content_id", "language"}, pathCol: "backdrop_path", sourceCol: "backdrop_source_path", clearSet: plainClear("backdrop_path")},
{name: "localized item logos", table: "media_item_localizations", keyCols: []string{"content_id", "language"}, pathCol: "logo_path", sourceCol: "logo_source_path", clearSet: plainClear("logo_path")},
{name: "season posters", table: "seasons", keyCols: []string{"content_id"}, pathCol: "poster_path", sourceCol: "poster_source_path", clearSet: plainClear("poster_path")},
{name: "localized season posters", table: "season_localizations", keyCols: []string{"season_content_id", "language"}, pathCol: "poster_path", sourceCol: "poster_source_path", clearSet: plainClear("poster_path")},
{name: "episode stills", table: "episodes", keyCols: []string{"content_id"}, pathCol: "still_path", sourceCol: "still_source_path", clearSet: plainClear("still_path")},
{name: "person photos", table: "people", keyCols: []string{"id::text"}, pathCol: "photo_path", sourceCol: "photo_source_path", clearSet: plainClear("photo_path")},
// Admin/user uploads: no re-downloadable source. Clearing falls back
// to the generated collage (admin collections), the generated poster
// (user collections), or the default tile (library posters); admins
// re-upload anything they want back. alwaysVerify protects surviving
// uploads from blind bulk resets.
{name: "collection posters", table: "library_collections", keyCols: []string{"id"}, pathCol: "poster_url", clearSet: `poster_url = '', poster_thumbhash = '', poster_auto_generated = FALSE, poster_from_template = FALSE, updated_at = NOW()`, alwaysVerify: true},
{name: "collection backdrops", table: "library_collections", keyCols: []string{"id"}, pathCol: "backdrop_url", clearSet: `backdrop_url = '', backdrop_thumbhash = '', updated_at = NOW()`, alwaysVerify: true},
{name: "user collection posters", table: "user_personal_collections", keyCols: []string{"id"}, pathCol: "poster_url", clearSet: `poster_url = '', poster_thumbhash = '', updated_at = NOW()`, alwaysVerify: true},
{name: "library posters", table: "media_folders", keyCols: []string{"id::text"}, pathCol: "poster_path", clearSet: `poster_path = ''`, alwaysVerify: true},
}
}
// ArtworkCacheReconciler verifies cached artwork keys against the public S3
// bucket and resets rows whose objects are missing, so the existing pipelines
// (image cache queue, book enrichment, chapter thumbnail backfill, collection
// collage generation) rebuild them in the currently configured storage.
type ArtworkCacheReconciler struct {
pool *pgxpool.Pool
s3 ArtworkObjectChecker
}
func NewArtworkCacheReconciler(pool *pgxpool.Pool, s3 ArtworkObjectChecker) *ArtworkCacheReconciler {
if pool == nil || s3 == nil {
return nil
}
return &ArtworkCacheReconciler{pool: pool, s3: s3}
}
// Run executes a full reconcile: probe, then either a bulk reset or a
// per-row verification sweep. It returns an error (leaving the storage
// fingerprint untouched at the caller) when storage cannot be reached or the
// error budget is exhausted, and never resets rows on the basis of transport
// errors.
func (r *ArtworkCacheReconciler) Run(ctx context.Context, progress func(percent float64, message string)) (ArtworkReconcileStats, error) {
stats := ArtworkReconcileStats{Mode: "verify"}
if r == nil || r.pool == nil || r.s3 == nil {
return stats, fmt.Errorf("artwork reconcile: not configured")
}
if progress == nil {
progress = func(float64, string) {}
}
surfaces := artworkSweepSurfaces()
// Probe before anything else: it decides the mode, and in bulk mode the
// per-surface count(*) queries (full scans on unindexable predicates)
// are never needed — bulk resets report their own RowsAffected.
progress(0, "Probing object storage")
if err := r.probe(ctx, surfaces, &stats); err != nil {
return stats, err
}
if stats.Sampled == 0 {
progress(100, "No cached artwork to verify")
return stats, nil
}
if shouldBulkReset(stats.Sampled, stats.SampleMissing) {
stats.Mode = "bulk_reset"
progress(5, fmt.Sprintf("Probe found %d/%d objects missing; resetting cached artwork", stats.SampleMissing, stats.Sampled))
steps := len(surfaces) + 1
for i, s := range surfaces {
pct := 5 + 90*float64(i+1)/float64(steps)
if s.alwaysVerify {
// Small upload-holding tables: never blind-reset; a surviving
// upload's row is the last pointer to its object.
if err := r.sweepSurface(ctx, s, &stats, func(done int) {
progress(pct, fmt.Sprintf("Verifying %s (%d rows)", s.name, done))
}); err != nil {
return stats, err
}
continue
}
if err := r.bulkResetSurface(ctx, s, &stats); err != nil {
return stats, err
}
progress(pct, fmt.Sprintf("Reset %s", s.name))
}
if err := r.bulkResetChapterThumbnails(ctx, &stats); err != nil {
return stats, err
}
progress(95, "Reset chapter thumbnails")
return stats, nil
}
// Verify mode: count cached rows once so progress has a denominator.
progress(2, "Counting cached artwork")
totals := make([]int, len(surfaces))
total := 0
for i, s := range surfaces {
n, err := r.countCached(ctx, s)
if err != nil {
return stats, err
}
totals[i] = n
total += n
}
chapterTotal, err := r.countChapterThumbnailFiles(ctx)
if err != nil {
return stats, err
}
total += chapterTotal
if total == 0 {
progress(100, "No cached artwork to verify")
return stats, nil
}
done := 0
report := func(surfaceName string) func(int) {
return func(surfaceDone int) {
pct := 5 + 90*float64(done+surfaceDone)/float64(total)
progress(pct, fmt.Sprintf("Verifying %s (%d/%d overall)", surfaceName, done+surfaceDone, total))
}
}
for i, s := range surfaces {
if totals[i] == 0 {
continue
}
if err := r.sweepSurface(ctx, s, &stats, report(s.name)); err != nil {
return stats, err
}
done += totals[i]
}
if chapterTotal > 0 {
if err := r.sweepChapterThumbnails(ctx, &stats, report("chapter thumbnails")); err != nil {
return stats, err
}
}
return stats, nil
}
// shouldBulkReset decides between a blind bulk reset and per-row
// verification. Probe HEADs are ground truth, so a near-total miss rate
// means the bucket plainly does not hold the cache; the threshold is below
// 1.0 only so a handful of coincidentally-present keys cannot force millions
// of pointless per-row checks. The minimum-sample bar keeps a probe thinned
// out by transport errors (or a tiny catalog) on the safe per-row path.
func shouldBulkReset(sampled, missing int) bool {
return sampled >= artworkReconcileBulkMinSample &&
float64(missing) >= artworkReconcileBulkThreshold*float64(sampled)
}
func (r *ArtworkCacheReconciler) countCached(ctx context.Context, s artworkSweepSurface) (int, error) {
var n int
q := fmt.Sprintf(`SELECT count(*) FROM %s WHERE %s`, s.table, s.cachedPredicate())
if err := r.pool.QueryRow(ctx, q).Scan(&n); err != nil {
return 0, fmt.Errorf("artwork reconcile: counting %s: %w", s.name, err)
}
return n, nil
}
// probe samples cached keys across all surfaces and HEADs them. A probe where
// every request errors aborts the run (storage unreachable ≠ objects missing).
func (r *ArtworkCacheReconciler) probe(ctx context.Context, surfaces []artworkSweepSurface, stats *ArtworkReconcileStats) error {
perSurface := artworkReconcileSampleTarget / (len(surfaces) + 1)
if perSurface < 1 {
perSurface = 1
}
// Plain LIMIT sampling (no ORDER BY random(), which would full-scan and
// sort every surface): the probe only has to answer "does the bucket
// hold this cache at all", and any N stored keys answer that. Partial
// migrations that skew the sample simply land in per-row verify mode,
// which handles them correctly anyway.
var keys []string
for _, s := range surfaces {
q := fmt.Sprintf(
`SELECT %s FROM %s WHERE %s LIMIT $1`,
s.pathCol, s.table, s.cachedPredicate(),
)
sampled, err := r.queryStrings(ctx, q, perSurface)
if err != nil {
return fmt.Errorf("artwork reconcile: sampling %s: %w", s.name, err)
}
keys = append(keys, sampled...)
}
chapterKeys, err := r.queryStrings(ctx, `
SELECT e->>'thumbnail_path'
FROM media_files, jsonb_array_elements(chapters) e
WHERE chapters IS NOT NULL AND coalesce(e->>'thumbnail_path', '') <> ''
LIMIT $1
`, perSurface)
if err != nil {
return fmt.Errorf("artwork reconcile: sampling chapter thumbnails: %w", err)
}
keys = append(keys, chapterKeys...)
if len(keys) == 0 {
return nil
}
present, missing, errored := r.headBatch(ctx, keys)
stats.Sampled = present + missing
stats.SampleMissing = missing
stats.Errors += errored
// A probe that mostly errors is not a probe of the cache, it is a probe
// of an outage: errored requests are excluded from the sample, so acting
// on the survivors could bulk-reset the catalog off a handful of 404s.
// Abort and leave the fingerprint stale; the next startup retries.
if errored*2 > len(keys) {
return fmt.Errorf("artwork reconcile: object storage unreliable: %d/%d probe requests failed", errored, len(keys))
}
return nil
}
func (r *ArtworkCacheReconciler) queryStrings(ctx context.Context, q string, args ...any) ([]string, error) {
rows, err := r.pool.Query(ctx, q, args...)
if err != nil {
return nil, err
}
defer rows.Close()
var out []string
for rows.Next() {
var v string
if err := rows.Scan(&v); err != nil {
return nil, err
}
out = append(out, v)
}
return out, rows.Err()
}
// headBatch checks the given keys with bounded concurrency and returns
// (present, missing, errored) counts.
func (r *ArtworkCacheReconciler) headBatch(ctx context.Context, keys []string) (present, missing, errored int) {
verdicts := r.headKeys(ctx, keys)
for _, v := range verdicts {
switch {
case v.err != nil:
errored++
case v.missing:
missing++
default:
present++
}
}
return present, missing, errored
}
type headVerdict struct {
missing bool
err error
}
// headKeys HEADs every key with bounded concurrency, preserving order.
func (r *ArtworkCacheReconciler) headKeys(ctx context.Context, keys []string) []headVerdict {
bucket := r.s3.Bucket()
verdicts := make([]headVerdict, len(keys))
var wg sync.WaitGroup
sem := make(chan struct{}, artworkReconcileHeadWorkers)
for i, key := range keys {
wg.Add(1)
go func(i int, key string) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
exists, err := r.objectExistsWithRetry(ctx, bucket, key)
verdicts[i] = headVerdict{missing: err == nil && !exists, err: err}
}(i, key)
}
wg.Wait()
return verdicts
}
func (r *ArtworkCacheReconciler) objectExistsWithRetry(ctx context.Context, bucket, key string) (bool, error) {
const maxAttempts = 3
var lastErr error
for attempt := 0; attempt < maxAttempts; attempt++ {
// Per-attempt deadline: a stalled HEAD must fail this attempt and
// move on, not hold the retry loop open until the run's context dies.
attemptCtx, cancel := context.WithTimeout(ctx, artworkReconcileHeadTimeout)
exists, err := r.s3.ObjectExists(attemptCtx, bucket, key)
cancel()
if err == nil {
return exists, nil
}
lastErr = err
if attempt == maxAttempts-1 {
break
}
timer := time.NewTimer(time.Duration(attempt+1) * 250 * time.Millisecond)
select {
case <-timer.C:
case <-ctx.Done():
timer.Stop()
return false, ctx.Err()
}
}
return false, lastErr
}
// bulkResetSurface resets every cached row without per-row verification. Rows
// with a re-downloadable provider source go back to that source (the enqueue
// loop re-caches them); rows without one are cleared so their owning pipeline
// can refill them.
func (r *ArtworkCacheReconciler) bulkResetSurface(ctx context.Context, s artworkSweepSurface, stats *ArtworkReconcileStats) error {
if s.sourceCol != "" {
requeue := fmt.Sprintf(
`UPDATE %s SET %s WHERE %s AND %s`,
s.table, s.resetSet(), s.cachedPredicate(), s.remoteSourcePredicate(),
)
tag, err := r.pool.Exec(ctx, requeue)
if err != nil {
return fmt.Errorf("artwork reconcile: bulk reset %s: %w", s.name, err)
}
stats.Requeued += int(tag.RowsAffected())
stats.Checked += int(tag.RowsAffected())
}
clearSQL := fmt.Sprintf(
`UPDATE %s SET %s WHERE %s AND NOT (%s)`,
s.table, s.clearSet, s.cachedPredicate(), s.remoteSourcePredicate(),
)
tag, err := r.pool.Exec(ctx, clearSQL)
if err != nil {
return fmt.Errorf("artwork reconcile: bulk clear %s: %w", s.name, err)
}
stats.Cleared += int(tag.RowsAffected())
stats.Checked += int(tag.RowsAffected())
return nil
}
// sweptRow is one candidate row in the per-row verification sweep.
type sweptRow struct {
keys []string
path string
remoteSource bool
}
func (r *ArtworkCacheReconciler) sweepSurface(ctx context.Context, s artworkSweepSurface, stats *ArtworkReconcileStats, onProgress func(done int)) error {
var cursor []string
done := 0
for {
rows, err := r.fetchSweepBatch(ctx, s, cursor)
if err != nil {
return err
}
if len(rows) == 0 {
return nil
}
cursor = rows[len(rows)-1].keys
if err := r.verifyAndReset(ctx, s, rows, stats); err != nil {
return err
}
if stats.SweepErrors > artworkReconcileErrorBudget {
return fmt.Errorf("artwork reconcile: aborting after %d sweep storage errors (errored rows were left untouched)", stats.SweepErrors)
}
done += len(rows)
onProgress(done)
}
}
func (r *ArtworkCacheReconciler) fetchSweepBatch(ctx context.Context, s artworkSweepSurface, cursor []string) ([]sweptRow, error) {
var b strings.Builder
args := make([]any, 0, len(cursor)+1)
fmt.Fprintf(&b, `SELECT %s, %s, (%s) FROM %s WHERE %s`,
strings.Join(s.keyCols, ", "), s.pathCol, s.remoteSourcePredicate(), s.table, s.cachedPredicate())
if len(cursor) > 0 {
placeholders := make([]string, len(cursor))
for i, v := range cursor {
args = append(args, v)
placeholders[i] = fmt.Sprintf("$%d", len(args))
}
fmt.Fprintf(&b, ` AND (%s) > (%s)`, strings.Join(s.keyCols, ", "), strings.Join(placeholders, ", "))
}
args = append(args, artworkReconcileBatchSize)
fmt.Fprintf(&b, ` ORDER BY %s LIMIT $%d`, strings.Join(s.keyCols, ", "), len(args))
rows, err := r.pool.Query(ctx, b.String(), args...)
if err != nil {
return nil, fmt.Errorf("artwork reconcile: fetching %s batch: %w", s.name, err)
}
defer rows.Close()
out := make([]sweptRow, 0, artworkReconcileBatchSize)
for rows.Next() {
row := sweptRow{keys: make([]string, len(s.keyCols))}
dest := make([]any, 0, len(s.keyCols)+2)
for i := range row.keys {
dest = append(dest, &row.keys[i])
}
dest = append(dest, &row.path, &row.remoteSource)
if err := rows.Scan(dest...); err != nil {
return nil, fmt.Errorf("artwork reconcile: scanning %s batch: %w", s.name, err)
}
out = append(out, row)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("artwork reconcile: iterating %s batch: %w", s.name, err)
}
return out, nil
}
func (r *ArtworkCacheReconciler) verifyAndReset(ctx context.Context, s artworkSweepSurface, batch []sweptRow, stats *ArtworkReconcileStats) error {
keys := make([]string, len(batch))
for i, row := range batch {
keys[i] = row.path
}
verdicts := r.headKeys(ctx, keys)
pkPredicate := keyEqualityPredicate(s.keyCols)
var pgBatch pgx.Batch
remoteByQueued := make([]bool, 0)
for i, v := range verdicts {
stats.Checked++
switch {
case v.err != nil:
stats.Errors++
stats.SweepErrors++
slog.Warn("artwork reconcile: object check failed; leaving row untouched",
"surface", s.name, "key", batch[i].path, "error", v.err)
case v.missing:
row := batch[i]
args := make([]any, 0, len(row.keys)+1)
for _, k := range row.keys {
args = append(args, k)
}
args = append(args, row.path)
var set string
if row.remoteSource {
set = s.resetSet()
} else {
set = s.clearSet
slog.Warn("artwork reconcile: cached image missing with no re-downloadable source; cleared",
"surface", s.name, "key", row.path, "row", strings.Join(row.keys, "/"))
}
pgBatch.Queue(fmt.Sprintf(`UPDATE %s SET %s WHERE %s AND %s = $%d`,
s.table, set, pkPredicate, s.pathCol, len(args)), args...)
remoteByQueued = append(remoteByQueued, row.remoteSource)
default:
stats.Verified++
}
}
if pgBatch.Len() == 0 {
return nil
}
results := r.pool.SendBatch(ctx, &pgBatch)
defer func() { _ = results.Close() }()
for _, remote := range remoteByQueued {
tag, err := results.Exec()
if err != nil {
return fmt.Errorf("artwork reconcile: resetting %s row: %w", s.name, err)
}
if tag.RowsAffected() == 0 {
// Row changed concurrently (metadata refresh, admin edit); leave it alone.
continue
}
if remote {
stats.Requeued++
} else {
stats.Cleared++
}
}
return nil
}
func keyEqualityPredicate(keyCols []string) string {
parts := make([]string, len(keyCols))
for i, col := range keyCols {
parts[i] = fmt.Sprintf("%s = $%d", col, i+1)
}
return strings.Join(parts, " AND ")
}
// --- Chapter thumbnails ---------------------------------------------------
//
// Chapter thumbnails live inside the media_files.chapters JSONB array
// (thumbnail_path / thumbnail_thumbhash per element). Clearing the path and
// the retry state makes the scheduled chapter_thumbnail_backfill task
// regenerate them from the media file.
const chapterThumbnailFilesPredicate = `chapters IS NOT NULL AND EXISTS (
SELECT 1 FROM jsonb_array_elements(chapters) e
WHERE coalesce(e->>'thumbnail_path', '') <> ''
)`
func (r *ArtworkCacheReconciler) countChapterThumbnailFiles(ctx context.Context) (int, error) {
var n int
q := `SELECT count(*) FROM media_files WHERE ` + chapterThumbnailFilesPredicate
if err := r.pool.QueryRow(ctx, q).Scan(&n); err != nil {
return 0, fmt.Errorf("artwork reconcile: counting chapter thumbnail files: %w", err)
}
return n, nil
}
func (r *ArtworkCacheReconciler) bulkResetChapterThumbnails(ctx context.Context, stats *ArtworkReconcileStats) error {
tag, err := r.pool.Exec(ctx, `
UPDATE media_files
SET chapters = (
SELECT jsonb_agg(
CASE WHEN coalesce(e->>'thumbnail_path', '') <> ''
THEN (e - 'thumbnail_retry_after' - 'thumbnail_failed_at' - 'thumbnail_last_error')
|| '{"thumbnail_path": "", "thumbnail_thumbhash": ""}'::jsonb
ELSE e
END
ORDER BY ord
)
FROM jsonb_array_elements(chapters) WITH ORDINALITY AS t(e, ord)
),
chapter_thumbnail_retry_after = NULL
WHERE `+chapterThumbnailFilesPredicate)
if err != nil {
return fmt.Errorf("artwork reconcile: bulk clearing chapter thumbnails: %w", err)
}
stats.Cleared += int(tag.RowsAffected())
stats.Checked += int(tag.RowsAffected())
return nil
}
// chapterFileRow is one media_files row in the chapter thumbnail sweep.
// Chapters are decoded as generic maps so fields this code does not know
// about survive a rewrite.
type chapterFileRow struct {
id int64
raw []byte
chapters []map[string]any
}
func (r *ArtworkCacheReconciler) sweepChapterThumbnails(ctx context.Context, stats *ArtworkReconcileStats, onProgress func(done int)) error {
cursor := int64(0)
done := 0
for {
rows, err := r.pool.Query(ctx, `
SELECT id, chapters FROM media_files
WHERE `+chapterThumbnailFilesPredicate+` AND id > $1
ORDER BY id LIMIT $2
`, cursor, artworkReconcileBatchSize)
if err != nil {
return fmt.Errorf("artwork reconcile: fetching chapter thumbnail batch: %w", err)
}
batch := make([]chapterFileRow, 0, artworkReconcileBatchSize)
for rows.Next() {
var f chapterFileRow
if err := rows.Scan(&f.id, &f.raw); err != nil {
rows.Close()
return fmt.Errorf("artwork reconcile: scanning chapter thumbnail batch: %w", err)
}
batch = append(batch, f)
}
rows.Close()
if err := rows.Err(); err != nil {
return fmt.Errorf("artwork reconcile: iterating chapter thumbnail batch: %w", err)
}
if len(batch) == 0 {
return nil
}
cursor = batch[len(batch)-1].id
if err := r.reconcileChapterBatch(ctx, batch, stats); err != nil {
return err
}
if stats.SweepErrors > artworkReconcileErrorBudget {
return fmt.Errorf("artwork reconcile: aborting after %d sweep storage errors (errored rows were left untouched)", stats.SweepErrors)
}
done += len(batch)
onProgress(done)
}
}
// reconcileChapterBatch verifies every chapter thumbnail across the whole
// batch in one HEAD fan-out — per-file checking would cap effective
// concurrency at one file's handful of chapters — then rewrites only the
// files whose arrays changed.
func (r *ArtworkCacheReconciler) reconcileChapterBatch(ctx context.Context, batch []chapterFileRow, stats *ArtworkReconcileStats) error {
type chapterRef struct{ file, chapter int }
var keys []string
var refs []chapterRef
for fi := range batch {
f := &batch[fi]
if err := json.Unmarshal(f.raw, &f.chapters); err != nil {
stats.Errors++
stats.SweepErrors++
slog.Warn("artwork reconcile: unparseable chapters JSON; skipping file", "file_id", f.id, "error", err)
f.chapters = nil
continue
}
for ci, ch := range f.chapters {
path, _ := ch["thumbnail_path"].(string)
if strings.TrimSpace(path) == "" {
continue
}
keys = append(keys, path)
refs = append(refs, chapterRef{file: fi, chapter: ci})
}
}
if len(keys) == 0 {
return nil
}
verdicts := r.headKeys(ctx, keys)
changed := make(map[int]bool, len(batch))
for vi, v := range verdicts {
stats.Checked++
ref := refs[vi]
switch {
case v.err != nil:
stats.Errors++
stats.SweepErrors++
slog.Warn("artwork reconcile: chapter thumbnail check failed; leaving chapter untouched",
"file_id", batch[ref.file].id, "key", keys[vi], "error", v.err)
case v.missing:
ch := batch[ref.file].chapters[ref.chapter]
ch["thumbnail_path"] = ""
ch["thumbnail_thumbhash"] = ""
delete(ch, "thumbnail_retry_after")
delete(ch, "thumbnail_failed_at")
delete(ch, "thumbnail_last_error")
changed[ref.file] = true
stats.Cleared++
default:
stats.Verified++
}
}
for fi := range batch {
if !changed[fi] {
continue
}
f := batch[fi]
updated, err := json.Marshal(f.chapters)
if err != nil {
return fmt.Errorf("artwork reconcile: encoding chapters for file %d: %w", f.id, err)
}
// Guard on the original JSON so a concurrent thumbnail-service write wins.
tag, err := r.pool.Exec(ctx, `
UPDATE media_files
SET chapters = $1::jsonb, chapter_thumbnail_retry_after = NULL
WHERE id = $2 AND chapters = $3::jsonb
`, updated, f.id, f.raw)
if err != nil {
return fmt.Errorf("artwork reconcile: updating chapters for file %d: %w", f.id, err)
}
if tag.RowsAffected() == 0 {
slog.Debug("artwork reconcile: chapters changed concurrently; skipped", "file_id", f.id)
}
}
return nil
}
+301
View File
@@ -0,0 +1,301 @@
package metadata
import (
"context"
"errors"
"fmt"
"os"
"strings"
"sync"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
// fakeObjectChecker treats every key as present unless listed in missing or
// erroring. Defaulting to present keeps sweeps over rows seeded by other
// tests in the shared test database side-effect free.
type fakeObjectChecker struct {
mu sync.Mutex
missing map[string]bool
erroring map[string]bool
errorAll bool
checked map[string]int
}
func (f *fakeObjectChecker) Bucket() string { return "test-bucket" }
func (f *fakeObjectChecker) ObjectExists(_ context.Context, _ string, key string) (bool, error) {
f.mu.Lock()
defer f.mu.Unlock()
if f.checked == nil {
f.checked = map[string]int{}
}
f.checked[key]++
if f.errorAll || f.erroring[key] {
return false, errors.New("simulated storage error")
}
return !f.missing[key], nil
}
func TestShouldBulkReset(t *testing.T) {
if shouldBulkReset(0, 0) {
t.Fatal("empty probe must not trigger bulk reset")
}
if shouldBulkReset(100, 94) {
t.Fatal("94% missing is below the bulk threshold")
}
if !shouldBulkReset(100, 95) {
t.Fatal("95% missing must trigger bulk reset")
}
// A probe thinned below the minimum successful-sample bar (transport
// errors, tiny catalog) must take the safe per-row path even at a 100%
// miss rate — a handful of surviving 404s is not a mandate to bulk-reset.
if shouldBulkReset(artworkReconcileBulkMinSample-1, artworkReconcileBulkMinSample-1) {
t.Fatal("below-minimum sample must not trigger bulk reset")
}
if !shouldBulkReset(artworkReconcileBulkMinSample, artworkReconcileBulkMinSample) {
t.Fatal("all-missing probe at the minimum sample size must trigger bulk reset")
}
}
func TestArtworkReconcileVerifySweep(t *testing.T) {
dsn := os.Getenv("SILO_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("SILO_TEST_DATABASE_URL is not set")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect test database: %v", err)
}
t.Cleanup(pool.Close)
suffix := time.Now().UnixNano()
id := func(name string) string { return fmt.Sprintf("arc-%s-%d", name, suffix) }
key := func(name string) string { return fmt.Sprintf("tmdb/movies/arc-%d/%s/original.webp", suffix, name) }
// Four items covering the sweep verdicts: intact, missing with a provider
// source, missing with an upload source, and an uncached provider URL.
seedItem := func(contentID, posterPath, posterSource string) {
if _, err := pool.Exec(ctx, `
INSERT INTO media_items (content_id, type, title, status, genres, poster_path, poster_source_path, last_refreshed)
VALUES ($1, 'movie', 'ARC Test', 'matched', '{}'::text[], $2, $3, NOW())
`, contentID, posterPath, posterSource); err != nil {
t.Fatalf("seed item %s: %v", contentID, err)
}
}
seedItem(id("intact"), key("intact"), "https://img.example/intact.jpg")
seedItem(id("missing"), key("missing"), "https://img.example/missing.jpg")
seedItem(id("upload"), key("upload"), "upload://admin/poster.jpg")
seedItem(id("uncached"), "https://img.example/direct.jpg", "https://img.example/direct.jpg")
var fileID int64
chapters := fmt.Sprintf(
`[{"index":0,"title":"One","thumbnail_path":%q,"thumbnail_thumbhash":"aGFzaA==","custom":"kept"},`+
`{"index":1,"title":"Two","thumbnail_path":%q,"thumbnail_thumbhash":"aGFzaA==","thumbnail_failed_at":"2026-01-01T00:00:00Z"}]`,
key("chapter-intact"), key("chapter-missing"),
)
var folderID int
if err := pool.QueryRow(ctx,
`INSERT INTO media_folders (type, name, enabled, poster_path) VALUES ('movies', 'ARC Folder', true, $1) RETURNING id`,
fmt.Sprintf("library-posters/arc-%d.png", suffix),
).Scan(&folderID); err != nil {
t.Fatalf("seed folder: %v", err)
}
if err := pool.QueryRow(ctx, `
INSERT INTO media_files (content_id, media_folder_id, file_path, chapters)
VALUES ($1, $2, $3, $4::jsonb) RETURNING id
`, id("intact"), folderID, fmt.Sprintf("/arc-%d/movie.mkv", suffix), chapters).Scan(&fileID); err != nil {
t.Fatalf("seed file: %v", err)
}
if _, err := pool.Exec(ctx, `
INSERT INTO library_collections (id, library_id, slug, title, collection_type, poster_url, poster_thumbhash, poster_from_template)
VALUES ($1, $2, $1, 'ARC Collection', 'manual', $3, 'aGFzaA==', TRUE)
`, id("coll"), folderID, key("coll")); err != nil {
t.Fatalf("seed collection: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(ctx, `DELETE FROM library_collections WHERE id = $1`, id("coll"))
_, _ = pool.Exec(ctx, `DELETE FROM media_files WHERE id = $1`, fileID)
_, _ = pool.Exec(ctx, `DELETE FROM media_folders WHERE id = $1`, folderID)
for _, name := range []string{"intact", "missing", "upload", "uncached"} {
_, _ = pool.Exec(ctx, `DELETE FROM media_items WHERE content_id = $1`, id(name))
}
})
checker := &fakeObjectChecker{missing: map[string]bool{
key("missing"): true,
key("upload"): true,
key("chapter-missing"): true,
key("coll"): true,
fmt.Sprintf("library-posters/arc-%d.png", suffix): true,
}}
stats, err := NewArtworkCacheReconciler(pool, checker).Run(ctx, nil)
if err != nil {
t.Fatalf("Run: %v", err)
}
if stats.Mode != "verify" {
t.Fatalf("Mode = %q, want verify (fake checker defaults to present)", stats.Mode)
}
var posterPath string
var lastRefreshed *time.Time
mustScanItem := func(contentID string) (string, *time.Time) {
if err := pool.QueryRow(ctx,
`SELECT poster_path, last_refreshed FROM media_items WHERE content_id = $1`, contentID,
).Scan(&posterPath, &lastRefreshed); err != nil {
t.Fatalf("read item %s: %v", contentID, err)
}
return posterPath, lastRefreshed
}
if got, _ := mustScanItem(id("intact")); got != key("intact") {
t.Fatalf("intact poster_path = %q, want untouched %q", got, key("intact"))
}
if got, _ := mustScanItem(id("missing")); got != "https://img.example/missing.jpg" {
t.Fatalf("missing poster_path = %q, want reset to provider source", got)
}
if got, refreshed := mustScanItem(id("upload")); got != "" || refreshed != nil {
t.Fatalf("upload poster_path = %q (last_refreshed %v), want cleared with last_refreshed NULL", got, refreshed)
}
if got, _ := mustScanItem(id("uncached")); got != "https://img.example/direct.jpg" {
t.Fatalf("uncached poster_path = %q, want untouched provider URL", got)
}
if checker.checked["https://img.example/direct.jpg"] != 0 {
t.Fatal("provider URLs must not be HEAD-checked")
}
var rawChapters string
var retryAfter *time.Time
if err := pool.QueryRow(ctx,
`SELECT chapters::text, chapter_thumbnail_retry_after FROM media_files WHERE id = $1`, fileID,
).Scan(&rawChapters, &retryAfter); err != nil {
t.Fatalf("read chapters: %v", err)
}
assertContains := func(s, substr, what string) {
t.Helper()
if !strings.Contains(s, substr) {
t.Fatalf("%s: %q not found in %s", what, substr, s)
}
}
assertContains(rawChapters, key("chapter-intact"), "intact chapter thumbnail kept")
assertContains(rawChapters, `"kept"`, "unknown chapter fields preserved")
if strings.Contains(rawChapters, key("chapter-missing")) || strings.Contains(rawChapters, "thumbnail_failed_at") {
t.Fatalf("missing chapter thumbnail not cleared: %s", rawChapters)
}
if retryAfter != nil {
t.Fatal("chapter_thumbnail_retry_after not cleared")
}
var collPoster, collHash string
var fromTemplate bool
if err := pool.QueryRow(ctx,
`SELECT poster_url, poster_thumbhash, poster_from_template FROM library_collections WHERE id = $1`, id("coll"),
).Scan(&collPoster, &collHash, &fromTemplate); err != nil {
t.Fatalf("read collection: %v", err)
}
if collPoster != "" || collHash != "" || fromTemplate {
t.Fatalf("collection artwork not fully cleared: url=%q hash=%q from_template=%v", collPoster, collHash, fromTemplate)
}
var folderPoster string
if err := pool.QueryRow(ctx, `SELECT poster_path FROM media_folders WHERE id = $1`, folderID).Scan(&folderPoster); err != nil {
t.Fatalf("read folder: %v", err)
}
if folderPoster != "" {
t.Fatalf("library poster not cleared: %q", folderPoster)
}
}
func TestArtworkReconcileLeavesRowsAloneOnStorageErrors(t *testing.T) {
dsn := os.Getenv("SILO_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("SILO_TEST_DATABASE_URL is not set")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect test database: %v", err)
}
t.Cleanup(pool.Close)
suffix := time.Now().UnixNano()
contentID := fmt.Sprintf("arc-err-%d", suffix)
okContentID := fmt.Sprintf("arc-err-ok-%d", suffix)
cachedKey := fmt.Sprintf("tmdb/movies/%s/poster/original.webp", contentID)
okKey := fmt.Sprintf("tmdb/movies/%s/poster/original.webp", okContentID)
seed := func(id, key string) {
if _, err := pool.Exec(ctx, `
INSERT INTO media_items (content_id, type, title, status, genres, poster_path, poster_source_path)
VALUES ($1, 'movie', 'ARC Err', 'matched', '{}'::text[], $2, 'https://img.example/err.jpg')
`, id, key); err != nil {
t.Fatalf("seed item %s: %v", id, err)
}
}
// The healthy sibling keeps the probe from concluding storage is
// unreachable (an all-errored probe aborts before any sweep runs), so
// the sweep-level skip-on-error behavior is what gets exercised.
seed(contentID, cachedKey)
seed(okContentID, okKey)
t.Cleanup(func() {
_, _ = pool.Exec(ctx, `DELETE FROM media_items WHERE content_id = ANY($1)`, []string{contentID, okContentID})
})
checker := &fakeObjectChecker{erroring: map[string]bool{cachedKey: true}}
stats, err := NewArtworkCacheReconciler(pool, checker).Run(ctx, nil)
if err != nil {
t.Fatalf("Run: %v", err)
}
if stats.Errors == 0 {
t.Fatal("expected the erroring key to be counted")
}
var posterPath string
if err := pool.QueryRow(ctx, `SELECT poster_path FROM media_items WHERE content_id = $1`, contentID).Scan(&posterPath); err != nil {
t.Fatalf("read item: %v", err)
}
if posterPath != cachedKey {
t.Fatalf("poster_path = %q, want untouched %q after storage error", posterPath, cachedKey)
}
}
func TestArtworkReconcileAbortsWhenStorageUnreachable(t *testing.T) {
dsn := os.Getenv("SILO_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("SILO_TEST_DATABASE_URL is not set")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect test database: %v", err)
}
t.Cleanup(pool.Close)
suffix := time.Now().UnixNano()
contentID := fmt.Sprintf("arc-down-%d", suffix)
cachedKey := fmt.Sprintf("tmdb/movies/%s/poster/original.webp", contentID)
if _, err := pool.Exec(ctx, `
INSERT INTO media_items (content_id, type, title, status, genres, poster_path, poster_source_path)
VALUES ($1, 'movie', 'ARC Down', 'matched', '{}'::text[], $2, 'https://img.example/down.jpg')
`, contentID, cachedKey); err != nil {
t.Fatalf("seed item: %v", err)
}
t.Cleanup(func() { _, _ = pool.Exec(ctx, `DELETE FROM media_items WHERE content_id = $1`, contentID) })
// Every probe HEAD errors: storage is unreachable, which must abort the
// run (missing ≠ unreachable) and leave the row untouched.
checker := &fakeObjectChecker{errorAll: true}
if _, err := NewArtworkCacheReconciler(pool, checker).Run(ctx, nil); err == nil {
t.Fatal("Run with unreachable storage returned nil error")
}
var posterPath string
if err := pool.QueryRow(ctx, `SELECT poster_path FROM media_items WHERE content_id = $1`, contentID).Scan(&posterPath); err != nil {
t.Fatalf("read item: %v", err)
}
if posterPath != cachedKey {
t.Fatalf("poster_path = %q, want untouched %q after unreachable storage", posterPath, cachedKey)
}
}
+13 -12
View File
@@ -546,7 +546,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
// cached relative path. The destination check makes the cached row itself the
// durable dedup marker, so pruning succeeded job rows does not cause the whole
// catalog to be re-downloaded once the rows age out.
rows, err := r.pool.Query(ctx, `
query := strings.ReplaceAll(`
WITH all_candidates AS (
SELECT
'poster'::text AS image_type,
@@ -563,7 +563,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
mi.imdb_id
FROM media_items mi
WHERE mi.poster_source_path LIKE '%://%'
AND lower(mi.poster_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(mi.poster_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (mi.poster_path LIKE '%://%' OR coalesce(mi.poster_path, '') = '')
UNION ALL
SELECT
@@ -581,7 +581,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
mi.imdb_id
FROM media_items mi
WHERE mi.backdrop_source_path LIKE '%://%'
AND lower(mi.backdrop_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(mi.backdrop_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (mi.backdrop_path LIKE '%://%' OR coalesce(mi.backdrop_path, '') = '')
UNION ALL
SELECT
@@ -599,7 +599,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
mi.imdb_id
FROM media_items mi
WHERE mi.logo_source_path LIKE '%://%'
AND lower(mi.logo_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(mi.logo_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (mi.logo_path LIKE '%://%' OR coalesce(mi.logo_path, '') = '')
UNION ALL
SELECT
@@ -618,7 +618,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
FROM media_item_localizations loc
JOIN media_items mi ON mi.content_id = loc.content_id
WHERE loc.poster_source_path LIKE '%://%'
AND lower(loc.poster_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(loc.poster_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (loc.poster_path LIKE '%://%' OR coalesce(loc.poster_path, '') = '')
UNION ALL
SELECT
@@ -637,7 +637,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
FROM media_item_localizations loc
JOIN media_items mi ON mi.content_id = loc.content_id
WHERE loc.backdrop_source_path LIKE '%://%'
AND lower(loc.backdrop_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(loc.backdrop_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (loc.backdrop_path LIKE '%://%' OR coalesce(loc.backdrop_path, '') = '')
UNION ALL
SELECT
@@ -656,7 +656,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
FROM media_item_localizations loc
JOIN media_items mi ON mi.content_id = loc.content_id
WHERE loc.logo_source_path LIKE '%://%'
AND lower(loc.logo_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(loc.logo_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (loc.logo_path LIKE '%://%' OR coalesce(loc.logo_path, '') = '')
UNION ALL
SELECT
@@ -675,7 +675,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
FROM seasons s
JOIN media_items mi ON mi.content_id = s.series_id
WHERE s.poster_source_path LIKE '%://%'
AND lower(s.poster_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(s.poster_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (s.poster_path LIKE '%://%' OR coalesce(s.poster_path, '') = '')
UNION ALL
SELECT
@@ -695,7 +695,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
JOIN seasons s ON s.content_id = loc.season_content_id
JOIN media_items mi ON mi.content_id = s.series_id
WHERE loc.poster_source_path LIKE '%://%'
AND lower(loc.poster_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(loc.poster_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (loc.poster_path LIKE '%://%' OR coalesce(loc.poster_path, '') = '')
UNION ALL
SELECT
@@ -714,7 +714,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
FROM episodes e
JOIN media_items mi ON mi.content_id = e.series_id
WHERE e.still_source_path LIKE '%://%'
AND lower(e.still_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(e.still_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (e.still_path LIKE '%://%' OR coalesce(e.still_path, '') = '')
UNION ALL
SELECT
@@ -732,7 +732,7 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
p.imdb_id
FROM people p
WHERE p.photo_source_path LIKE '%://%'
AND lower(p.photo_source_path) NOT LIKE ALL (ARRAY['s3://%', 'file://%', 'local://%', 'upload://%', 'generated://%'])
AND lower(p.photo_source_path) NOT LIKE ALL (@nonProviderSchemes)
AND (p.photo_path LIKE '%://%' OR coalesce(p.photo_path, '') = '')
),
candidates AS (
@@ -759,7 +759,8 @@ func (r *ImageCacheJobRepository) EnqueueExistingProviderArtwork(ctx context.Con
COALESCE(tvdb_id, '') AS tvdb_id,
COALESCE(imdb_id, '') AS imdb_id
FROM candidates
`, limit)
`, "@nonProviderSchemes", nonProviderImageSchemesSQL)
rows, err := r.pool.Query(ctx, query, limit)
if err != nil {
return 0, fmt.Errorf("enqueueing existing provider artwork: %w", err)
}
+6 -2
View File
@@ -104,7 +104,7 @@ func NewClient(cfg BucketConfig) *Client {
if tokenTTL <= 0 {
tokenTTL = 10800 // 3 hours
}
keyPrefix := normalizeKeyPrefix(cfg.KeyPrefix)
keyPrefix := NormalizeKeyPrefix(cfg.KeyPrefix)
return &Client{
s3Client: s3Client,
@@ -542,7 +542,11 @@ func newBytesReadSeeker(data []byte) *bytesReadSeeker {
return &bytesReadSeeker{data: data}
}
func normalizeKeyPrefix(prefix string) string {
// NormalizeKeyPrefix returns the canonical form of a bucket key prefix as the
// client applies it to object keys: whitespace- and slash-trimmed, case
// preserved (S3 keys are case-sensitive). Exported so storage-identity
// comparisons elsewhere normalize prefixes exactly the way this client does.
func NormalizeKeyPrefix(prefix string) string {
return strings.Trim(strings.TrimSpace(prefix), "/")
}
+9 -2
View File
@@ -148,10 +148,17 @@ func (m *TaskManager) triggerLoop(ctx context.Context, w *taskWorker) {
shouldRun, shouldRunErr := m.shouldRunScheduledTask(ctx, w)
if shouldRunErr != nil {
m.logger.WarnContext(ctx, "scheduled task preflight failed; running task",
// Fail closed: a preflight that cannot answer must not launch the
// task — for expensive conditional tasks a transient error would
// otherwise trigger the exact work the gate exists to suppress.
// Interval/daily triggers retry at their next firing; manual
// RunTask always bypasses the gate.
m.logger.WarnContext(ctx, "scheduled task preflight failed; skipping run",
"task", w.task.Key(), "error", shouldRunErr)
m.rearmTriggersFromNow(w)
continue
}
if shouldRunErr == nil && !shouldRun {
if !shouldRun {
m.rearmTriggersFromNow(w)
continue
}
+66 -1
View File
@@ -2,6 +2,8 @@ package taskmanager_test
import (
"context"
"errors"
"io"
"log/slog"
"reflect"
"sync"
@@ -146,6 +148,7 @@ func (t stubTask) DefaultTriggers() []taskmanager.TriggerConfig {
type conditionalStubTask struct {
stubTask
shouldRunCalled chan struct{}
shouldRunErr error
mu sync.Mutex
executeCalls int
}
@@ -155,7 +158,7 @@ func (t *conditionalStubTask) ShouldRun(context.Context) (bool, error) {
case t.shouldRunCalled <- struct{}{}:
default:
}
return false, nil
return false, t.shouldRunErr
}
func (t *conditionalStubTask) Execute(context.Context, taskmanager.ProgressReporter) error {
@@ -386,3 +389,65 @@ func TestTaskManagerTriggerSkipsConditionalTaskWithoutHistory(t *testing.T) {
beforeTrigger.Format(time.RFC3339Nano))
}
}
// A preflight that cannot answer must fail closed: skipping the run is
// recoverable at the next trigger, while running an expensive conditional
// task on a transient error is exactly what the gate exists to prevent.
func TestTaskManagerTriggerSkipsConditionalTaskOnPreflightError(t *testing.T) {
triggerRepo := &fakeTriggerRepository{
triggers: map[string][]taskmanager.TriggerConfig{
"conditional": {
{Type: taskmanager.TriggerTypeInterval, IntervalMs: int64(time.Hour / time.Millisecond)},
},
},
}
historyRepo := &recordingExecutionRepository{}
var triggers []*fakeTrigger
factory := func(cfg taskmanager.TriggerConfig) taskmanager.Trigger {
tr := &fakeTrigger{
cfg: cfg,
ch: make(chan struct{}, 1),
stopCh: make(chan struct{}),
}
triggers = append(triggers, tr)
return tr
}
manager := taskmanager.New(
triggerRepo,
historyRepo,
factory,
slog.New(slog.NewTextHandler(io.Discard, nil)),
)
task := &conditionalStubTask{
stubTask: stubTask{key: "conditional"},
shouldRunCalled: make(chan struct{}, 1),
shouldRunErr: errors.New("settings unavailable"),
}
manager.Register(task)
ctx, cancel := context.WithCancel(context.Background())
defer manager.Stop()
defer cancel()
manager.Start(ctx)
if len(triggers) != 1 {
t.Fatalf("triggers = %d, want 1", len(triggers))
}
beforeTrigger := time.Now()
triggers[0].ch <- struct{}{}
select {
case <-task.shouldRunCalled:
case <-time.After(time.Second):
t.Fatal("scheduled preflight was not called")
}
time.Sleep(25 * time.Millisecond)
if got := task.executeCount(); got != 0 {
t.Fatalf("Execute calls = %d, want 0 (preflight errors must fail closed)", got)
}
if !triggers[0].next.After(beforeTrigger) {
t.Fatalf("next run = %s, want rearmed after skipped preflight error",
triggers[0].next.Format(time.RFC3339Nano))
}
}
@@ -0,0 +1,186 @@
package tasks
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"strings"
"time"
"github.com/Silo-Server/silo-server/internal/metadata"
"github.com/Silo-Server/silo-server/internal/s3client"
"github.com/Silo-Server/silo-server/internal/taskmanager"
)
// ArtworkStorageIdentityKey is the server_settings key holding the storage
// identity fingerprint of the public S3 bucket the artwork cache was last
// reconciled against. Machine-managed; not an admin-editable setting.
const ArtworkStorageIdentityKey = "s3.public_storage_identity"
// ArtworkStorageIdentity builds the fingerprint of the public S3 storage the
// cached artwork lives in. Only fields that determine *where objects are
// stored* participate: the read endpoint and URL-auth settings affect how
// objects are served, not where they live, so changing them must not trigger
// a reconcile.
//
// Normalization mirrors how each field is actually used: endpoints (hostnames)
// and bucket names are case-insensitive, but the key prefix feeds into
// case-sensitive object keys, so it keeps its case and is normalized exactly
// like s3client applies it (slash- and whitespace-trimmed). A case-only prefix
// edit is a real storage move and must change the fingerprint; a slash-only
// edit is not and must not.
func ArtworkStorageIdentity(endpoint, bucket, keyPrefix string) string {
insensitive := func(v string) string { return strings.ToLower(strings.TrimSpace(v)) }
return insensitive(endpoint) + "|" + insensitive(bucket) + "|" + s3client.NormalizeKeyPrefix(keyPrefix)
}
// ArtworkReconcileSettingsStore is the server-settings surface the task needs.
// Satisfied by *catalog.ServerSettingsRepo and its encrypting decorator.
type ArtworkReconcileSettingsStore interface {
Get(ctx context.Context, key string) (string, error)
Set(ctx context.Context, key, value string) error
}
// ArtworkReconcileRunner runs a reconcile sweep. Satisfied by
// *metadata.ArtworkCacheReconciler.
type ArtworkReconcileRunner interface {
Run(ctx context.Context, progress func(percent float64, message string)) (metadata.ArtworkReconcileStats, error)
}
// BrandingAssetReconciler clears branding asset refs whose stored objects are
// missing. Satisfied by *branding.Service; may be nil when branding has no
// storage.
type BrandingAssetReconciler interface {
ReconcileMissingAssets(ctx context.Context) (checked, cleared int, err error)
}
// ReconcileArtworkCacheTask verifies cached artwork against the currently
// configured public object storage and resets whatever is missing so the
// image cache pipeline rebuilds it. Scheduled runs only fire when the storage
// identity changed since the last completed reconcile; manual runs always
// sweep, which doubles as recovery from bucket data loss.
type ReconcileArtworkCacheTask struct {
runner ArtworkReconcileRunner
settings ArtworkReconcileSettingsStore
branding BrandingAssetReconciler
identity string
}
func NewReconcileArtworkCacheTask(runner ArtworkReconcileRunner, settings ArtworkReconcileSettingsStore, branding BrandingAssetReconciler, identity string) *ReconcileArtworkCacheTask {
return &ReconcileArtworkCacheTask{runner: runner, settings: settings, branding: branding, identity: identity}
}
func (t *ReconcileArtworkCacheTask) Key() string { return "reconcile_artwork_cache" }
func (t *ReconcileArtworkCacheTask) Name() string { return "Reconcile Artwork Cache" }
func (t *ReconcileArtworkCacheTask) Description() string {
return "Verifies cached artwork against object storage and re-caches anything missing (runs automatically after the storage provider changes)"
}
func (t *ReconcileArtworkCacheTask) Category() taskmanager.TaskCategory {
return taskmanager.TaskCategoryMetadata
}
func (t *ReconcileArtworkCacheTask) IsHidden() bool { return false }
func (t *ReconcileArtworkCacheTask) DefaultTriggers() []taskmanager.TriggerConfig {
return []taskmanager.TriggerConfig{
{Type: taskmanager.TriggerTypeStartup},
}
}
// ShouldRun suppresses the startup trigger while the storage identity is
// unchanged. Manual RunTask calls bypass this and always sweep.
//
// The startup trigger fires exactly once per process, so a transient settings
// read failure here would postpone a needed reconcile until the next restart;
// retry briefly before giving up. (The task manager skips the run on a
// preflight error rather than failing open into a full sweep.)
func (t *ReconcileArtworkCacheTask) ShouldRun(ctx context.Context) (bool, error) {
if t.runner == nil || t.settings == nil {
return false, nil
}
var stored string
var err error
for attempt := 0; attempt < 3; attempt++ {
stored, err = t.settings.Get(ctx, ArtworkStorageIdentityKey)
if err == nil {
return stored != "" && stored != t.identity, nil
}
timer := time.NewTimer(time.Duration(attempt+1) * time.Second)
select {
case <-timer.C:
case <-ctx.Done():
timer.Stop()
return false, ctx.Err()
}
}
return false, fmt.Errorf("reading artwork storage identity: %w", err)
}
func (t *ReconcileArtworkCacheTask) Execute(ctx context.Context, progress taskmanager.ProgressReporter) error {
if t.runner == nil || t.settings == nil {
progress.Report(100, "Artwork reconcile is not configured")
return nil
}
stats, err := t.runner.Run(ctx, progress.Report)
if err != nil {
if data, marshalErr := json.Marshal(stats); marshalErr == nil {
progress.SetResultData(data)
}
return fmt.Errorf("reconciling artwork cache: %w", err)
}
// Only a clean, completed sweep certifies the current storage. Sweep
// errors mean rows were skipped unverified, so the fingerprint stays
// stale and the next startup retries; resets already applied this run
// are durable either way.
if stats.SweepErrors > 0 {
if data, marshalErr := json.Marshal(stats); marshalErr == nil {
progress.SetResultData(data)
}
return fmt.Errorf(
"artwork reconcile: %d rows skipped on storage errors (verified %d, re-queued %d, cleared %d); storage identity left uncertified so the next startup retries",
stats.SweepErrors, stats.Verified, stats.Requeued, stats.Cleared,
)
}
// Certify before the branding check: a transient failure on that
// 4-object pass must not discard a completed catalog sweep and force it
// to repeat every boot.
if setErr := t.settings.Set(ctx, ArtworkStorageIdentityKey, t.identity); setErr != nil {
return fmt.Errorf("persisting artwork storage identity: %w", setErr)
}
brandingNote := ""
if t.branding != nil {
brandingChecked, brandingCleared, brandingErr := t.branding.ReconcileMissingAssets(ctx)
stats.Cleared += brandingCleared
stats.Checked += brandingChecked
if brandingErr != nil {
stats.Errors++
brandingNote = fmt.Sprintf("; branding asset check failed: %v (re-run the task to retry)", brandingErr)
slog.Warn("artwork reconcile: branding asset check failed", "error", brandingErr)
}
}
if data, marshalErr := json.Marshal(stats); marshalErr == nil {
progress.SetResultData(data)
}
message := fmt.Sprintf(
"Verified %d cached images intact, re-queued %d for re-cache, cleared %d without a re-downloadable source",
stats.Verified, stats.Requeued, stats.Cleared,
)
if stats.Mode == "bulk_reset" {
message = fmt.Sprintf(
"Storage probe found %d/%d sampled objects missing; reset all cached artwork (re-queued %d, cleared %d)",
stats.SampleMissing, stats.Sampled, stats.Requeued, stats.Cleared,
)
}
if stats.Errors > 0 {
// SweepErrors is zero here (checked above), so these are probe or
// branding errors — reported, but they don't reduce sweep coverage.
message += fmt.Sprintf(", %d storage errors during probing", stats.Errors)
}
progress.Report(100, message+brandingNote)
return nil
}
@@ -0,0 +1,189 @@
package tasks
import (
"context"
"encoding/json"
"errors"
"strings"
"testing"
"github.com/Silo-Server/silo-server/internal/metadata"
)
type fakeSettingsStore struct {
values map[string]string
getErr error
}
func (f *fakeSettingsStore) Get(_ context.Context, key string) (string, error) {
if f.getErr != nil {
return "", f.getErr
}
return f.values[key], nil
}
func (f *fakeSettingsStore) Set(_ context.Context, key, value string) error {
if f.values == nil {
f.values = map[string]string{}
}
f.values[key] = value
return nil
}
type fakeReconcileRunner struct {
stats metadata.ArtworkReconcileStats
err error
runs int
}
func (f *fakeReconcileRunner) Run(context.Context, func(float64, string)) (metadata.ArtworkReconcileStats, error) {
f.runs++
return f.stats, f.err
}
type fakeBrandingReconciler struct {
checked int
cleared int
err error
}
func (f *fakeBrandingReconciler) ReconcileMissingAssets(context.Context) (int, int, error) {
return f.checked, f.cleared, f.err
}
type fakeProgress struct {
lastMessage string
resultData json.RawMessage
}
func (f *fakeProgress) Report(_ float64, message string) { f.lastMessage = message }
func (f *fakeProgress) SetResultData(data json.RawMessage) { f.resultData = data }
func TestArtworkStorageIdentityNormalizes(t *testing.T) {
// Endpoint and bucket are case-insensitive; whitespace is trimmed.
a := ArtworkStorageIdentity(" https://S3.Example.com ", "Assets", "silo/prod")
b := ArtworkStorageIdentity("https://s3.example.com", "assets", "silo/prod")
if a != b {
t.Fatalf("identity not normalized: %q != %q", a, b)
}
// The key prefix is slash-insensitive (the s3client trims slashes, so
// 'art' and '/art/' are the same storage location)...
if ArtworkStorageIdentity("e", "b", "art") != ArtworkStorageIdentity("e", "b", " /art/ ") {
t.Fatal("slash-only prefix differences must not change the identity")
}
// ...but case-SENSITIVE: S3 object keys are case-sensitive, so a
// case-only prefix edit is a real storage move and must reconcile.
if ArtworkStorageIdentity("e", "b", "Art") == ArtworkStorageIdentity("e", "b", "art") {
t.Fatal("case-only prefix differences are real storage moves and must change the identity")
}
if a == ArtworkStorageIdentity("https://s3.example.com", "assets", "") {
t.Fatal("key prefix must participate in the identity")
}
if a == ArtworkStorageIdentity("https://other.example.com", "assets", "silo/prod") {
t.Fatal("endpoint must participate in the identity")
}
}
func TestReconcileArtworkCacheShouldRun(t *testing.T) {
runner := &fakeReconcileRunner{}
store := &fakeSettingsStore{values: map[string]string{}}
task := NewReconcileArtworkCacheTask(runner, store, nil, "endpoint|bucket|prefix")
// No stored fingerprint: first boot, seeding happens at wiring time; the
// scheduled run must not sweep a catalog it has no baseline for.
if run, err := task.ShouldRun(context.Background()); err != nil || run {
t.Fatalf("ShouldRun with empty fingerprint = %v, %v; want false, nil", run, err)
}
store.values[ArtworkStorageIdentityKey] = "endpoint|bucket|prefix"
if run, err := task.ShouldRun(context.Background()); err != nil || run {
t.Fatalf("ShouldRun with matching fingerprint = %v, %v; want false, nil", run, err)
}
store.values[ArtworkStorageIdentityKey] = "old-endpoint|bucket|prefix"
if run, err := task.ShouldRun(context.Background()); err != nil || !run {
t.Fatalf("ShouldRun with changed fingerprint = %v, %v; want true, nil", run, err)
}
}
func TestReconcileArtworkCacheExecutePersistsFingerprintOnlyOnSuccess(t *testing.T) {
store := &fakeSettingsStore{values: map[string]string{ArtworkStorageIdentityKey: "old"}}
failing := &fakeReconcileRunner{err: errors.New("storage unreachable")}
task := NewReconcileArtworkCacheTask(failing, store, nil, "new")
if err := task.Execute(context.Background(), &fakeProgress{}); err == nil {
t.Fatal("Execute with failing runner returned nil error")
}
if got := store.values[ArtworkStorageIdentityKey]; got != "old" {
t.Fatalf("fingerprint after failed run = %q, want unchanged %q", got, "old")
}
ok := &fakeReconcileRunner{stats: metadata.ArtworkReconcileStats{Mode: "verify", Verified: 3, Requeued: 2, Cleared: 1}}
task = NewReconcileArtworkCacheTask(ok, store, nil, "new")
progress := &fakeProgress{}
if err := task.Execute(context.Background(), progress); err != nil {
t.Fatalf("Execute = %v, want nil", err)
}
if got := store.values[ArtworkStorageIdentityKey]; got != "new" {
t.Fatalf("fingerprint after successful run = %q, want %q", got, "new")
}
if progress.resultData == nil {
t.Fatal("Execute did not record result data")
}
}
func TestReconcileArtworkCacheExecuteDoesNotCertifyOnSweepErrors(t *testing.T) {
// Rows skipped on storage errors were never verified, so the sweep did
// not fully cover the catalog: the fingerprint must stay stale so the
// next startup retries.
store := &fakeSettingsStore{values: map[string]string{ArtworkStorageIdentityKey: "old"}}
runner := &fakeReconcileRunner{stats: metadata.ArtworkReconcileStats{
Mode: "verify", Verified: 10, Errors: 3, SweepErrors: 3,
}}
branding := &fakeBrandingReconciler{checked: 4}
task := NewReconcileArtworkCacheTask(runner, store, branding, "new")
if err := task.Execute(context.Background(), &fakeProgress{}); err == nil {
t.Fatal("Execute with sweep errors returned nil error")
}
if got := store.values[ArtworkStorageIdentityKey]; got != "old" {
t.Fatalf("fingerprint after sweep errors = %q, want unchanged %q", got, "old")
}
}
func TestReconcileArtworkCacheExecuteIncludesBranding(t *testing.T) {
store := &fakeSettingsStore{values: map[string]string{}}
runner := &fakeReconcileRunner{stats: metadata.ArtworkReconcileStats{Mode: "verify", Cleared: 1}}
task := NewReconcileArtworkCacheTask(runner, store, &fakeBrandingReconciler{checked: 4, cleared: 2}, "id")
progress := &fakeProgress{}
if err := task.Execute(context.Background(), progress); err != nil {
t.Fatalf("Execute = %v, want nil", err)
}
var stats metadata.ArtworkReconcileStats
if err := json.Unmarshal(progress.resultData, &stats); err != nil {
t.Fatalf("decode result data: %v", err)
}
if stats.Cleared != 3 {
t.Fatalf("Cleared = %d, want 3 (1 artwork + 2 branding)", stats.Cleared)
}
if stats.Checked != 4 {
t.Fatalf("Checked = %d, want 4 (all probed branding assets, not just cleared ones)", stats.Checked)
}
// A branding failure must NOT discard the completed sweep: the
// fingerprint is certified first and the failure is reported in the
// message instead, so the full catalog sweep never repeats over a
// transient error on a 4-object branding check.
fpStore := &fakeSettingsStore{values: map[string]string{}}
failing := NewReconcileArtworkCacheTask(runner, fpStore,
&fakeBrandingReconciler{err: errors.New("storage unreachable")}, "id")
failingProgress := &fakeProgress{}
if err := failing.Execute(context.Background(), failingProgress); err != nil {
t.Fatalf("Execute with failing branding reconcile = %v, want nil (non-fatal)", err)
}
if got := fpStore.values[ArtworkStorageIdentityKey]; got != "id" {
t.Fatalf("fingerprint after branding failure = %q, want certified %q", got, "id")
}
if !strings.Contains(failingProgress.lastMessage, "branding asset check failed") {
t.Fatalf("completion message %q does not surface the branding failure", failingProgress.lastMessage)
}
}
@@ -35,6 +35,7 @@ describe("StorageSettings", () => {
restartRequired: false,
sensitiveConfigured: [],
buildConnectionCheckRequest: vi.fn(),
isDirty: () => false,
});
const markup = renderToStaticMarkup(<StorageSettings />);
@@ -42,5 +43,34 @@ describe("StorageSettings", () => {
expect(markup).toContain("Public Assets");
expect(markup).toContain("Private Internal");
expect(markup).toContain("Check Connection");
expect(markup).not.toContain("Storage location change");
});
it("warns about the artwork cache when a public storage identity field is edited", () => {
useCheckAdminSettingsConnectionMock.mockReturnValue({
isPending: false,
mutateAsync: vi.fn(),
});
useSettingsFormMock.mockReturnValue({
isLoading: false,
getValue: (key: string) => {
if (key === "s3.public_url_auth") return "presigned";
return "";
},
setValue: vi.fn(),
dirtyCount: 1,
save: vi.fn(),
discard: vi.fn(),
isSaving: false,
restartRequired: false,
sensitiveConfigured: [],
buildConnectionCheckRequest: vi.fn(),
isDirty: (key: string) => key === "s3.public_bucket",
});
const markup = renderToStaticMarkup(<StorageSettings />);
expect(markup).toContain("Storage location change");
expect(markup).toContain("re-caches anything missing");
});
});
@@ -1,4 +1,5 @@
import { useMemo, useState } from "react";
import { AlertTriangle } from "lucide-react";
import type { ConnectionCheckResponse } from "@/api/types";
import { ConnectionCheckAction } from "@/components/admin/ConnectionCheckAction";
import { useCheckAdminSettingsConnection } from "@/hooks/queries/admin/settings";
@@ -25,6 +26,14 @@ const PUBLIC_S3_KEYS = [
"s3.public_token_ttl",
] as const;
// Changing any of these moves where cached artwork objects live; the server
// reconciles the artwork cache after a restart (see reconcile_artwork_cache).
const PUBLIC_S3_IDENTITY_KEYS = [
"s3.public_endpoint",
"s3.public_bucket",
"s3.public_key_prefix",
] as const;
const PRIVATE_S3_KEYS = [
"s3.private_endpoint",
"s3.private_region",
@@ -190,6 +199,20 @@ export default function StorageSettings() {
value={form.getValue("s3.public_key_prefix")}
onChange={(v) => form.setValue("s3.public_key_prefix", v)}
/>
{PUBLIC_S3_IDENTITY_KEYS.some((key) => form.isDirty(key)) && (
<div className="my-3 flex items-start gap-3 rounded-xl border border-amber-500/20 bg-amber-500/5 p-4">
<AlertTriangle className="mt-0.5 h-4 w-4 shrink-0 text-amber-500" />
<div className="text-[13px] leading-relaxed">
<p className="font-medium text-amber-500">Storage location change</p>
<p className="text-muted-foreground mt-1">
Artwork is cached in this bucket. After the server restarts, Silo verifies the
cache against the new storage and automatically re-caches anything missing.
Uploaded images (custom posters, collection artwork, branding) cannot be
re-downloaded — migrate your bucket contents if you want to keep them.
</p>
</div>
</div>
)}
<SettingField
label="Access Key"
type="password"