Files
silo-server/internal/taskmanager/triggers/interval.go
14ffc91dfb [codex] Expand provider image cache queue (#176)
* 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>
2026-06-18 10:07:58 -04:00

99 lines
1.9 KiB
Go

package triggers
import (
"sync"
"time"
"github.com/Silo-Server/silo-server/internal/taskmanager"
)
// IntervalTrigger fires every N milliseconds, measured from completion of the
// previous run (not from start). If no previous run exists, it fires
// interval-after-Start().
type IntervalTrigger struct {
cfg taskmanager.TriggerConfig
interval time.Duration
ch chan struct{}
nextRun time.Time
timer *time.Timer
stopCh chan struct{}
mu sync.Mutex
}
func NewIntervalTrigger(cfg taskmanager.TriggerConfig) *IntervalTrigger {
return &IntervalTrigger{
cfg: cfg,
interval: time.Duration(cfg.IntervalMs) * time.Millisecond,
ch: make(chan struct{}, 1),
}
}
func (t *IntervalTrigger) Start(lastResult *taskmanager.ExecutionResult) {
t.mu.Lock()
defer t.mu.Unlock()
// Drain any stale signal from a previous timer fire.
select {
case <-t.ch:
default:
}
var base time.Time
if lastResult != nil && !lastResult.CompletedAt.IsZero() {
base = lastResult.CompletedAt
} else {
base = time.Now()
}
t.nextRun = base.Add(t.interval)
delay := time.Until(t.nextRun)
if delay < 0 {
delay = 0
t.nextRun = time.Now()
}
stopCh := make(chan struct{})
timer := time.NewTimer(delay)
t.stopCh = stopCh
t.timer = timer
go func() {
select {
case <-stopCh:
if !timer.Stop() {
select {
case <-timer.C:
default:
}
}
return
case <-timer.C:
select {
case t.ch <- struct{}{}:
default:
}
}
}()
}
func (t *IntervalTrigger) Stop() {
t.mu.Lock()
defer t.mu.Unlock()
if t.stopCh != nil {
select {
case <-t.stopCh:
default:
close(t.stopCh)
}
}
}
func (t *IntervalTrigger) NextRunTime() time.Time {
t.mu.Lock()
defer t.mu.Unlock()
return t.nextRun
}
func (t *IntervalTrigger) Config() taskmanager.TriggerConfig { return t.cfg }
func (t *IntervalTrigger) C() <-chan struct{} { return t.ch }