package runner // asynq_runner.go — Asynq-based task dispatch for the runner. // // When cfg.RedisAddr is set, Run() calls runAsynq() instead of runPoll(). // The Asynq server replaces the polling loop: it listens on Redis for tasks // enqueued by the backend Producer and delivers them immediately. // // Handlers in this file decode Asynq job payloads and call the existing // runScrapeTask / runAudioTask methods, keeping all execution logic in one place. import ( "context" "encoding/json" "fmt" "time" "github.com/hibiken/asynq" asynqmetrics "github.com/hibiken/asynq/x/metrics" "github.com/libnovel/backend/internal/asynqqueue" "github.com/libnovel/backend/internal/domain" ) // runAsynq starts an Asynq server that replaces the PocketBase poll loop. // It also starts the periodic catalogue refresh ticker. // Blocks until ctx is cancelled. func (r *Runner) runAsynq(ctx context.Context) error { redisOpt, err := r.redisConnOpt() if err != nil { return fmt.Errorf("runner: parse redis addr: %w", err) } srv := asynq.NewServer(redisOpt, asynq.Config{ // Allocate concurrency slots for each task type. // Total concurrency = scrape + audio slots. Concurrency: r.cfg.MaxConcurrentScrape + r.cfg.MaxConcurrentAudio, Queues: map[string]int{ asynqqueue.QueueDefault: 1, }, // Let Asynq handle retries with exponential back-off. RetryDelayFunc: asynq.DefaultRetryDelayFunc, // Log errors from handlers via the existing structured logger. ErrorHandler: asynq.ErrorHandlerFunc(func(_ context.Context, task *asynq.Task, err error) { r.deps.Log.Error("runner: asynq task failed", "type", task.Type(), "err", err, ) }), }) mux := asynq.NewServeMux() mux.HandleFunc(asynqqueue.TypeAudioGenerate, r.handleAudioTask) mux.HandleFunc(asynqqueue.TypeScrapeBook, r.handleScrapeTask) mux.HandleFunc(asynqqueue.TypeScrapeCatalogue, r.handleScrapeTask) // Register Asynq queue metrics with the default Prometheus registry so // the /metrics endpoint (metrics.go) can expose them. inspector := asynq.NewInspector(redisOpt) collector := asynqmetrics.NewQueueMetricsCollector(inspector) if err := r.metricsRegistry.Register(collector); err != nil { r.deps.Log.Warn("runner: could not register asynq prometheus collector", "err", err) } // Start the periodic catalogue refresh. catalogueTick := time.NewTicker(r.cfg.CatalogueRefreshInterval) defer catalogueTick.Stop() if !r.cfg.SkipInitialCatalogueRefresh { go r.runCatalogueRefresh(ctx) } else { r.deps.Log.Info("runner: skipping initial catalogue refresh (RUNNER_SKIP_INITIAL_CATALOGUE_REFRESH=true)") } r.deps.Log.Info("runner: asynq mode active", "redis_addr", r.cfg.RedisAddr) // Run catalogue refresh ticker in the background. go func() { for { select { case <-ctx.Done(): return case <-catalogueTick.C: go r.runCatalogueRefresh(ctx) } } }() // Start Asynq server (non-blocking). if err := srv.Start(mux); err != nil { return fmt.Errorf("runner: asynq server start: %w", err) } // Block until context is cancelled, then gracefully stop. <-ctx.Done() r.deps.Log.Info("runner: context cancelled, shutting down asynq server") srv.Shutdown() return nil } // redisConnOpt parses cfg.RedisAddr into an asynq.RedisConnOpt. // Supports full "redis://" / "rediss://" URLs and plain "host:port". func (r *Runner) redisConnOpt() (asynq.RedisConnOpt, error) { addr := r.cfg.RedisAddr // ParseRedisURI handles redis:// and rediss:// schemes. if len(addr) > 7 && (addr[:8] == "redis://" || addr[:9] == "rediss://") { return asynq.ParseRedisURI(addr) } // Plain "host:port" — use RedisClientOpt directly. return asynq.RedisClientOpt{ Addr: addr, Password: r.cfg.RedisPassword, }, nil } // handleScrapeTask is the Asynq handler for TypeScrapeBook and TypeScrapeCatalogue. func (r *Runner) handleScrapeTask(ctx context.Context, t *asynq.Task) error { var p asynqqueue.ScrapePayload if err := json.Unmarshal(t.Payload(), &p); err != nil { return fmt.Errorf("unmarshal scrape payload: %w", err) } task := domain.ScrapeTask{ ID: p.PBTaskID, Kind: p.Kind, TargetURL: p.TargetURL, FromChapter: p.FromChapter, ToChapter: p.ToChapter, } r.tasksRunning.Add(1) defer r.tasksRunning.Add(-1) r.runScrapeTask(ctx, task) return nil } // handleAudioTask is the Asynq handler for TypeAudioGenerate. func (r *Runner) handleAudioTask(ctx context.Context, t *asynq.Task) error { var p asynqqueue.AudioPayload if err := json.Unmarshal(t.Payload(), &p); err != nil { return fmt.Errorf("unmarshal audio payload: %w", err) } task := domain.AudioTask{ ID: p.PBTaskID, Slug: p.Slug, Chapter: p.Chapter, Voice: p.Voice, } r.tasksRunning.Add(1) defer r.tasksRunning.Add(-1) r.runAudioTask(ctx, task) return nil }