feat: wire main loop with periodic sync+scan and sync_interval config
Replace compute-only run() with an immediate sync+scan followed by a ticker-driven periodic loop, add the sync.interval config field (default 6h) with defaults and validation, and add tests for the scheduling logic.
This commit is contained in:
@@ -8,6 +8,7 @@ import (
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"naviwatcher/internal/config"
|
||||
"naviwatcher/internal/database"
|
||||
@@ -22,6 +23,10 @@ type App struct {
|
||||
db *database.DB
|
||||
mbClient *musicbrainz.MusicBrainzClient
|
||||
ndClient *navidrome.NavidromeClient
|
||||
|
||||
// syncFn, when non-nil, replaces the real syncAndScan call in
|
||||
// startPeriodicSync so tests can observe the loop without live clients.
|
||||
syncFn func(ctx context.Context) error
|
||||
}
|
||||
|
||||
// navidromeClientFactory constructs the Navidrome client. It is a package-level
|
||||
@@ -114,16 +119,50 @@ func (a *App) Close() {
|
||||
}
|
||||
|
||||
func (a *App) run(ctx context.Context) error {
|
||||
// Compute-only scanner hook: scan all monitored artists for missing
|
||||
// releases and log the count. Notifier/Web UI are out of scope for this
|
||||
// plan, so results are only logged. ScanAll is a blocking DB walk over
|
||||
// every monitored artist; it observes ctx cancellation and returns early.
|
||||
missing, err := scanner.ScanAll(ctx, a.db, a.cfg.Scanner.FuzzyThreshold)
|
||||
if err != nil {
|
||||
// Run an immediate sync+scan so the service produces results without
|
||||
// waiting a full interval, then kick off the periodic loop goroutine.
|
||||
// Business logic added in later tasks (notifier, web server) will be
|
||||
// wired as additional goroutines below.
|
||||
if err := a.doSync(ctx); err != nil {
|
||||
if ctx.Err() != nil {
|
||||
// Context cancelled (e.g. shutdown) — exit cleanly.
|
||||
return nil
|
||||
}
|
||||
log.Printf("Initial sync+scan failed: %v", err)
|
||||
}
|
||||
|
||||
a.startPeriodicSync(ctx)
|
||||
|
||||
<-ctx.Done()
|
||||
return nil
|
||||
}
|
||||
|
||||
// doSync runs the sync pipeline, using the injected syncFn when present (tests)
|
||||
// or the real syncAndScan otherwise.
|
||||
func (a *App) doSync(ctx context.Context) error {
|
||||
if a.syncFn != nil {
|
||||
return a.syncFn(ctx)
|
||||
}
|
||||
return a.syncAndScan(ctx)
|
||||
}
|
||||
|
||||
// syncAndScan runs the full data pipeline once: Navidrome artist sync, the
|
||||
// MusicBrainz discography pipeline (SyncAll), then the scanner over the
|
||||
// now-populated DB. It logs results and observes ctx cancellation.
|
||||
func (a *App) syncAndScan(ctx context.Context) error {
|
||||
if err := navidrome.SyncArtists(ctx, a.ndClient, a.db); err != nil {
|
||||
return fmt.Errorf("sync artists: %w", err)
|
||||
}
|
||||
|
||||
discography := musicbrainz.NewDiscographySyncer(a.mbClient)
|
||||
albums := musicbrainz.NewAlbumSyncer(func(ctx context.Context, db *database.DB) error {
|
||||
return navidrome.SyncAlbums(ctx, a.ndClient, db)
|
||||
})
|
||||
if err := musicbrainz.SyncAll(ctx, a.db, a.mbClient, discography, albums, a.cfg.MusicBrainz.CacheTTL); err != nil {
|
||||
return fmt.Errorf("sync all: %w", err)
|
||||
}
|
||||
|
||||
missing, err := scanner.ScanAll(ctx, a.db, a.cfg.Scanner.FuzzyThreshold)
|
||||
if err != nil {
|
||||
return fmt.Errorf("scan all: %w", err)
|
||||
}
|
||||
|
||||
@@ -131,10 +170,31 @@ func (a *App) run(ctx context.Context) error {
|
||||
for _, m := range missing {
|
||||
log.Printf(" missing: artist=%s rgid=%s title=%q", m.ArtistID, m.RGID, m.Title)
|
||||
}
|
||||
|
||||
// Main application loop — blocks until context is cancelled.
|
||||
// Business logic (notifier, web server) will be wired into separate
|
||||
// goroutines here in future tasks.
|
||||
<-ctx.Done()
|
||||
return nil
|
||||
}
|
||||
|
||||
// startPeriodicSync runs syncAndScan on a ticker at cfg.Sync.Interval. It
|
||||
// blocks until ctx is cancelled, then returns cleanly. Each tick runs in its
|
||||
// own goroutine so a slow sync does not block the ticker; a fresh interval is
|
||||
// still scheduled regardless.
|
||||
func (a *App) startPeriodicSync(ctx context.Context) {
|
||||
ticker := time.NewTicker(a.cfg.Sync.Interval)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
log.Println("Periodic sync stopped.")
|
||||
return
|
||||
case <-ticker.C:
|
||||
go func() {
|
||||
if err := a.doSync(ctx); err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
log.Printf("Periodic sync+scan failed: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,20 +7,22 @@ import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"naviwatcher/internal/config"
|
||||
"naviwatcher/internal/database"
|
||||
"naviwatcher/internal/navidrome"
|
||||
"naviwatcher/internal/scanner"
|
||||
)
|
||||
|
||||
func TestAppRun_ScanLogsMissingReleases(t *testing.T) {
|
||||
// Verify the compute-only run() hook scans monitored artists and returns
|
||||
// nil without starting notifier/web. Uses an in-memory DB with one
|
||||
// monitored artist that has one missing release (Animals) vs a local album
|
||||
// (The Wall). The context is left live so the scan actually executes; we
|
||||
// cancel shortly after to let run() return cleanly.
|
||||
// Verify run() performs the sync+scan pipeline once (via the injected
|
||||
// syncFn) and then blocks until ctx cancellation, returning nil. The
|
||||
// injected syncFn performs the scan and logs the missing release, mirroring
|
||||
// what syncAndScan does against live clients.
|
||||
db, err := database.New(":memory:")
|
||||
if err != nil {
|
||||
t.Fatalf("database.New() error: %v", err)
|
||||
@@ -49,21 +51,31 @@ func TestAppRun_ScanLogsMissingReleases(t *testing.T) {
|
||||
t.Fatalf("seed external release: %v", err)
|
||||
}
|
||||
|
||||
var buf bytes.Buffer
|
||||
log.SetOutput(&buf)
|
||||
defer log.SetOutput(os.Stderr)
|
||||
|
||||
app := &App{
|
||||
cfg: &config.Config{Scanner: config.ScannerConfig{FuzzyThreshold: 0.85}},
|
||||
db: db,
|
||||
cfg: &config.Config{
|
||||
Scanner: config.ScannerConfig{FuzzyThreshold: 0.85},
|
||||
Sync: config.SyncConfig{Interval: time.Hour},
|
||||
},
|
||||
db: db,
|
||||
syncFn: func(ctx context.Context) error {
|
||||
missing, err := scanner.ScanAll(ctx, db, 0.85)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, m := range missing {
|
||||
log.Printf(" missing: artist=%s rgid=%s title=%q", m.ArtistID, m.RGID, m.Title)
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
// Capture run()'s log output so we assert that run() ITSELF performed
|
||||
// the scan (not a separately re-run ScanAll). This guards against the
|
||||
// hook silently becoming a no-op while still passing.
|
||||
var buf bytes.Buffer
|
||||
log.SetOutput(&buf)
|
||||
defer log.SetOutput(os.Stderr)
|
||||
|
||||
// Run the (blocking) hook in a goroutine; cancel after it has had time to
|
||||
// perform the scan so run() returns nil via the ctx.Done() path.
|
||||
done := make(chan error, 1)
|
||||
@@ -83,6 +95,98 @@ func TestAppRun_ScanLogsMissingReleases(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestStartPeriodicSync_FiresOnTick(t *testing.T) {
|
||||
var calls int64
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(2)
|
||||
|
||||
app := &App{
|
||||
cfg: &config.Config{Sync: config.SyncConfig{Interval: 20 * time.Millisecond}},
|
||||
syncFn: func(ctx context.Context) error {
|
||||
atomic.AddInt64(&calls, 1)
|
||||
wg.Done()
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
go app.startPeriodicSync(ctx)
|
||||
|
||||
if !waitWG(&wg, 2*time.Second) {
|
||||
t.Fatal("expected syncFn to be called at least twice within timeout")
|
||||
}
|
||||
|
||||
if got := atomic.LoadInt64(&calls); got < 2 {
|
||||
t.Errorf("expected at least 2 sync calls, got %d", got)
|
||||
}
|
||||
|
||||
cancel()
|
||||
}
|
||||
|
||||
func TestStartPeriodicSync_CancelsCleanly(t *testing.T) {
|
||||
var calls int64
|
||||
app := &App{
|
||||
cfg: &config.Config{Sync: config.SyncConfig{Interval: time.Hour}},
|
||||
syncFn: func(ctx context.Context) error {
|
||||
atomic.AddInt64(&calls, 1)
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
done := make(chan struct{})
|
||||
|
||||
go func() {
|
||||
app.startPeriodicSync(ctx)
|
||||
close(done)
|
||||
}()
|
||||
|
||||
// With a 1h interval the ticker would never fire on its own; cancel should
|
||||
// return promptly.
|
||||
cancel()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
// clean exit
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("startPeriodicSync did not exit after ctx cancellation")
|
||||
}
|
||||
|
||||
if got := atomic.LoadInt64(&calls); got != 0 {
|
||||
t.Errorf("expected no sync calls with 1h interval, got %d", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDoSync_UsesInjectedSyncFn(t *testing.T) {
|
||||
// Verify doSync prefers an injected syncFn when present (so the periodic
|
||||
// loop and immediate run can be driven by tests without live clients),
|
||||
// and falls back to the real syncAndScan otherwise.
|
||||
db, err := database.New(":memory:")
|
||||
if err != nil {
|
||||
t.Fatalf("database.New() error: %v", err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
var called int32
|
||||
app := &App{
|
||||
cfg: &config.Config{Sync: config.SyncConfig{Interval: time.Hour}},
|
||||
db: db,
|
||||
syncFn: func(ctx context.Context) error {
|
||||
atomic.StoreInt32(&called, 1)
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
if err := app.doSync(context.Background()); err != nil {
|
||||
t.Fatalf("doSync returned error: %v", err)
|
||||
}
|
||||
if atomic.LoadInt32(&called) != 1 {
|
||||
t.Fatal("expected injected syncFn to be called")
|
||||
}
|
||||
}
|
||||
|
||||
func TestConfigIntegration(t *testing.T) {
|
||||
// Integration test: write a minimal valid config and load it via config.LoadConfig,
|
||||
// verifying the full path that main() uses.
|
||||
@@ -212,6 +316,7 @@ func TestAppRun_GracefulShutdown(t *testing.T) {
|
||||
MusicBrainz: config.MusicBrainzConfig{
|
||||
UserAgent: "NaviWatcher/1.0 ( test@example.com )",
|
||||
},
|
||||
Sync: config.SyncConfig{Interval: time.Hour},
|
||||
}
|
||||
|
||||
prevFactory := navidromeClientFactory
|
||||
@@ -260,3 +365,18 @@ musicbrainz:
|
||||
t.Fatal("expected config validation error for empty musicbrainz.user_agent, got nil")
|
||||
}
|
||||
}
|
||||
|
||||
// waitWG waits for wg with a timeout; returns true if it completed in time.
|
||||
func waitWG(wg *sync.WaitGroup, timeout time.Duration) bool {
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
wg.Wait()
|
||||
close(done)
|
||||
}()
|
||||
select {
|
||||
case <-done:
|
||||
return true
|
||||
case <-time.After(timeout):
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user