* feat(metadata): expand provider image cache queue * fix(metadata): harden provider image cache queue Addresses bug-review feedback from Codex/CodeRabbit on the metadata image cache pipeline. All findings validated against the code before fixing; false positives (rows/connection deadlock, PhotoSourcePath merge coupling) were confirmed non-issues and left unchanged. - Honor metadata.cache_images for the background processor. The cache_metadata_images task was registered whenever S3 was configured, so merely enabling object storage downloaded the entire provider-artwork catalog even with caching disabled. Add ImageCacheProcessor.SetEnabled, gate RunOnce/RunUntilIdle on it, and wire it (with hot reload) from cfg.Metadata.CacheImages in main.go. - Guard terminal job updates with lease ownership. EnqueueBatch can repurpose a running row with a new source; MarkSucceeded/MarkFailed keyed on id alone let a stale worker finalize the replacement job and drop the new artwork. Thread locked_by through and add status='running' AND locked_by=$n guards. - Avoid uploading stale jobs onto the live artwork key. Verify the target still references the job's source (CurrentTargetSourcePath) before CacheImage, so a job whose source an admin/refresh already replaced cannot overwrite the deterministic storage object. - COALESCE nullable external IDs in EnqueueExistingProviderArtwork. A NULL tmdb_id/tvdb_id/imdb_id on any candidate failed the scan and aborted the whole cache run; matches the existing item_repo pattern. - Stop re-downloading the catalog every 30 days. Discovery now skips targets whose *_path is already a cached relative path, making the cached row the durable dedup marker instead of the prunable job row. - Decouple catalog sweeps from queue draining. RunOnce no longer runs discovery per batch; RunUntilIdle sweeps only when the queue drains and throttles full sweeps to every 15m, so idle installs stop full-scanning every entity table each minute. - Requeue claimed-but-unstarted jobs on cancellation. Acquire the semaphore before spawning workers and RequeueClaimed any jobs not yet started, instead of leaving them locked until the 15m lease expires. - Skip the backoff sleep after the final upload attempt in putObjectWithRetry (saves ~1.5s on permanent failures). - Add the s3/file/local/upload/generated exclusion to the seasons and episodes backfill in migration 20260617184537 for consistency with the later migration (the bad backfill was inert downstream, but the asymmetry is removed). 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
98 lines
1.7 KiB
Go
98 lines
1.7 KiB
Go
package triggers
|
|
|
|
import (
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/Silo-Server/silo-server/internal/taskmanager"
|
|
)
|
|
|
|
const startupDelay = 5 * time.Second
|
|
|
|
// StartupTrigger fires once, shortly after Start() is called. Subsequent
|
|
// calls to Start() (e.g. re-arming after execution) are no-ops.
|
|
type StartupTrigger struct {
|
|
cfg taskmanager.TriggerConfig
|
|
ch chan struct{}
|
|
delay time.Duration
|
|
nextRun time.Time
|
|
timer *time.Timer
|
|
stopCh chan struct{}
|
|
fired bool
|
|
mu sync.Mutex
|
|
}
|
|
|
|
func NewStartupTrigger(cfg taskmanager.TriggerConfig) *StartupTrigger {
|
|
return &StartupTrigger{
|
|
cfg: cfg,
|
|
ch: make(chan struct{}, 1),
|
|
delay: startupDelay,
|
|
}
|
|
}
|
|
|
|
func (s *StartupTrigger) Start(_ *taskmanager.ExecutionResult) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
// Drain any stale signal from a previous timer fire.
|
|
select {
|
|
case <-s.ch:
|
|
default:
|
|
}
|
|
|
|
if s.fired {
|
|
s.nextRun = time.Time{}
|
|
return
|
|
}
|
|
s.fired = true
|
|
|
|
stopCh := make(chan struct{})
|
|
s.nextRun = time.Now().Add(s.delay)
|
|
timer := time.NewTimer(s.delay)
|
|
s.stopCh = stopCh
|
|
s.timer = timer
|
|
|
|
go func() {
|
|
select {
|
|
case <-stopCh:
|
|
if !timer.Stop() {
|
|
select {
|
|
case <-timer.C:
|
|
default:
|
|
}
|
|
}
|
|
return
|
|
case <-timer.C:
|
|
s.mu.Lock()
|
|
s.nextRun = time.Time{}
|
|
s.mu.Unlock()
|
|
select {
|
|
case s.ch <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
}()
|
|
}
|
|
|
|
func (s *StartupTrigger) Stop() {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if s.stopCh != nil {
|
|
select {
|
|
case <-s.stopCh:
|
|
default:
|
|
close(s.stopCh)
|
|
}
|
|
}
|
|
s.nextRun = time.Time{}
|
|
}
|
|
|
|
func (s *StartupTrigger) NextRunTime() time.Time {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
return s.nextRun
|
|
}
|
|
|
|
func (s *StartupTrigger) Config() taskmanager.TriggerConfig { return s.cfg }
|
|
func (s *StartupTrigger) C() <-chan struct{} { return s.ch }
|