All checks were successful
Release / Test backend (push) Successful in 43s
Release / Check ui (push) Successful in 43s
Release / Docker / caddy (push) Successful in 46s
Release / Docker / backend (push) Successful in 2m45s
Release / Docker / runner (push) Successful in 2m53s
Release / Docker / ui (push) Successful in 2m5s
Release / Gitea Release (push) Successful in 41s
248 lines
9.4 KiB
Go
248 lines
9.4 KiB
Go
// Command runner is the homelab worker binary.
|
|
//
|
|
// It polls PocketBase for pending scrape and audio tasks, executes them, and
|
|
// writes results back. It connects directly to PocketBase and MinIO using
|
|
// admin credentials loaded from environment variables.
|
|
//
|
|
// Usage:
|
|
//
|
|
// runner # start polling loop (blocks until SIGINT/SIGTERM)
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"os"
|
|
"os/signal"
|
|
"runtime"
|
|
"syscall"
|
|
"time"
|
|
|
|
"github.com/getsentry/sentry-go"
|
|
"github.com/libnovel/backend/internal/asynqqueue"
|
|
"github.com/libnovel/backend/internal/browser"
|
|
"github.com/libnovel/backend/internal/cfai"
|
|
"github.com/libnovel/backend/internal/config"
|
|
"github.com/libnovel/backend/internal/kokoro"
|
|
"github.com/libnovel/backend/internal/libretranslate"
|
|
"github.com/libnovel/backend/internal/meili"
|
|
"github.com/libnovel/backend/internal/novelfire"
|
|
"github.com/libnovel/backend/internal/otelsetup"
|
|
"github.com/libnovel/backend/internal/pockettts"
|
|
"github.com/libnovel/backend/internal/runner"
|
|
"github.com/libnovel/backend/internal/storage"
|
|
"github.com/libnovel/backend/internal/taskqueue"
|
|
)
|
|
|
|
// version and commit are set at build time via -ldflags.
|
|
var (
|
|
version = "dev"
|
|
commit = "unknown"
|
|
)
|
|
|
|
func main() {
|
|
if err := run(); err != nil {
|
|
fmt.Fprintf(os.Stderr, "runner: fatal: %v\n", err)
|
|
os.Exit(1)
|
|
}
|
|
}
|
|
|
|
func run() error {
|
|
cfg := config.Load()
|
|
|
|
// ── Sentry / GlitchTip error tracking ────────────────────────────────────
|
|
if dsn := os.Getenv("GLITCHTIP_DSN"); dsn != "" {
|
|
if err := sentry.Init(sentry.ClientOptions{
|
|
Dsn: dsn,
|
|
Release: version + "@" + commit,
|
|
TracesSampleRate: 0.1,
|
|
}); err != nil {
|
|
fmt.Fprintf(os.Stderr, "runner: sentry init warning: %v\n", err)
|
|
} else {
|
|
defer sentry.Flush(2 * time.Second)
|
|
}
|
|
}
|
|
|
|
// ── Logger ──────────────────────────────────────────────────────────────
|
|
log := buildLogger(cfg.LogLevel)
|
|
log.Info("runner starting",
|
|
"version", version,
|
|
"commit", commit,
|
|
"worker_id", cfg.Runner.WorkerID,
|
|
)
|
|
|
|
// ── Context: cancel on SIGINT / SIGTERM ─────────────────────────────────
|
|
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
|
defer stop()
|
|
|
|
// ── OpenTelemetry tracing + logs ─────────────────────────────────────────
|
|
otelShutdown, otelLog, err := otelsetup.Init(ctx, version)
|
|
if err != nil {
|
|
return fmt.Errorf("init otel: %w", err)
|
|
}
|
|
if otelShutdown != nil {
|
|
defer otelShutdown()
|
|
// Switch to the OTel-bridged logger so all structured log lines are
|
|
// forwarded to Loki with trace IDs attached.
|
|
log = otelLog
|
|
log.Info("otel tracing + logs enabled", "endpoint", os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT"))
|
|
}
|
|
|
|
// ── Storage ─────────────────────────────────────────────────────────────
|
|
store, err := storage.NewStore(ctx, cfg, log)
|
|
if err != nil {
|
|
return fmt.Errorf("init storage: %w", err)
|
|
}
|
|
|
|
// ── Browser / Scraper ───────────────────────────────────────────────────
|
|
workers := cfg.Runner.Workers
|
|
if workers <= 0 {
|
|
workers = runtime.NumCPU()
|
|
}
|
|
timeout := cfg.Runner.Timeout
|
|
if timeout <= 0 {
|
|
timeout = 90 * time.Second
|
|
}
|
|
|
|
browserClient := browser.NewDirectClient(browser.Config{
|
|
MaxConcurrent: workers,
|
|
Timeout: timeout,
|
|
})
|
|
novel := novelfire.New(browserClient, log)
|
|
|
|
// ── Kokoro ──────────────────────────────────────────────────────────────
|
|
var kokoroClient kokoro.Client
|
|
if cfg.Kokoro.URL != "" {
|
|
kokoroClient = kokoro.New(cfg.Kokoro.URL)
|
|
log.Info("kokoro TTS enabled", "url", cfg.Kokoro.URL)
|
|
} else {
|
|
log.Warn("KOKORO_URL not set — kokoro voice tasks will fail")
|
|
kokoroClient = &noopKokoro{}
|
|
}
|
|
|
|
// ── pocket-tts ──────────────────────────────────────────────────────────
|
|
var pocketTTSClient pockettts.Client
|
|
if cfg.PocketTTS.URL != "" {
|
|
pocketTTSClient = pockettts.New(cfg.PocketTTS.URL)
|
|
log.Info("pocket-tts enabled", "url", cfg.PocketTTS.URL)
|
|
} else {
|
|
log.Warn("POCKET_TTS_URL not set — pocket-tts voice tasks will fail")
|
|
}
|
|
|
|
// ── Cloudflare Workers AI ────────────────────────────────────────────────
|
|
var cfaiClient cfai.Client
|
|
if cfg.CFAI.AccountID != "" && cfg.CFAI.APIToken != "" {
|
|
cfaiClient = cfai.New(cfg.CFAI.AccountID, cfg.CFAI.APIToken, cfg.CFAI.Model)
|
|
log.Info("cloudflare AI TTS enabled", "model", cfg.CFAI.Model)
|
|
} else {
|
|
log.Info("CFAI_ACCOUNT_ID/CFAI_API_TOKEN not set — CF AI voice tasks will fail")
|
|
}
|
|
|
|
// ── LibreTranslate ──────────────────────────────────────────────────────
|
|
ltClient := libretranslate.New(cfg.LibreTranslate.URL, cfg.LibreTranslate.APIKey)
|
|
if ltClient != nil {
|
|
log.Info("libretranslate enabled", "url", cfg.LibreTranslate.URL)
|
|
} else {
|
|
log.Info("LIBRETRANSLATE_URL not set — machine translation disabled")
|
|
}
|
|
|
|
// ── Meilisearch ─────────────────────────────────────────────────────────
|
|
var searchIndex meili.Client
|
|
if cfg.Meilisearch.URL != "" {
|
|
if err := meili.Configure(cfg.Meilisearch.URL, cfg.Meilisearch.APIKey); err != nil {
|
|
log.Warn("meilisearch configure failed — search indexing disabled", "err", err)
|
|
searchIndex = meili.NoopClient{}
|
|
} else {
|
|
searchIndex = meili.New(cfg.Meilisearch.URL, cfg.Meilisearch.APIKey)
|
|
log.Info("meilisearch enabled", "url", cfg.Meilisearch.URL)
|
|
}
|
|
} else {
|
|
log.Info("MEILI_URL not set — search indexing disabled")
|
|
searchIndex = meili.NoopClient{}
|
|
}
|
|
|
|
// ── Runner ──────────────────────────────────────────────────────────────
|
|
rCfg := runner.Config{
|
|
WorkerID: cfg.Runner.WorkerID,
|
|
PollInterval: cfg.Runner.PollInterval,
|
|
MaxConcurrentScrape: cfg.Runner.MaxConcurrentScrape,
|
|
MaxConcurrentAudio: cfg.Runner.MaxConcurrentAudio,
|
|
MaxConcurrentTranslation: cfg.Runner.MaxConcurrentTranslation,
|
|
OrchestratorWorkers: workers,
|
|
MetricsAddr: cfg.Runner.MetricsAddr,
|
|
CatalogueRefreshInterval: cfg.Runner.CatalogueRefreshInterval,
|
|
CatalogueRequestDelay: cfg.Runner.CatalogueRequestDelay,
|
|
SkipInitialCatalogueRefresh: cfg.Runner.SkipInitialCatalogueRefresh,
|
|
RedisAddr: cfg.Redis.Addr,
|
|
RedisPassword: cfg.Redis.Password,
|
|
}
|
|
|
|
// In Asynq mode the Consumer is a thin wrapper: claim/heartbeat/reap are
|
|
// no-ops, but FinishAudioTask / FinishScrapeTask / FailTask write back to
|
|
// PocketBase as before.
|
|
var consumer taskqueue.Consumer = store
|
|
if cfg.Redis.Addr != "" {
|
|
log.Info("runner: asynq mode — using Redis for task dispatch", "addr", cfg.Redis.Addr)
|
|
consumer = asynqqueue.NewConsumer(store)
|
|
} else {
|
|
log.Info("runner: poll mode — using PocketBase for task dispatch")
|
|
}
|
|
|
|
deps := runner.Dependencies{
|
|
Consumer: consumer,
|
|
BookWriter: store,
|
|
BookReader: store,
|
|
AudioStore: store,
|
|
CoverStore: store,
|
|
TranslationStore: store,
|
|
SearchIndex: searchIndex,
|
|
Novel: novel,
|
|
Kokoro: kokoroClient,
|
|
PocketTTS: pocketTTSClient,
|
|
CFAI: cfaiClient,
|
|
LibreTranslate: ltClient,
|
|
Log: log,
|
|
}
|
|
r := runner.New(rCfg, deps)
|
|
|
|
return r.Run(ctx)
|
|
}
|
|
|
|
// ── Helpers ───────────────────────────────────────────────────────────────────
|
|
|
|
func buildLogger(level string) *slog.Logger {
|
|
var lvl slog.Level
|
|
switch level {
|
|
case "debug":
|
|
lvl = slog.LevelDebug
|
|
case "warn":
|
|
lvl = slog.LevelWarn
|
|
case "error":
|
|
lvl = slog.LevelError
|
|
default:
|
|
lvl = slog.LevelInfo
|
|
}
|
|
return slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: lvl}))
|
|
}
|
|
|
|
// noopKokoro is a no-op implementation used when KOKORO_URL is not set.
|
|
type noopKokoro struct{}
|
|
|
|
func (n *noopKokoro) GenerateAudio(_ context.Context, _, _ string) ([]byte, error) {
|
|
return nil, fmt.Errorf("kokoro not configured (KOKORO_URL is empty)")
|
|
}
|
|
|
|
func (n *noopKokoro) StreamAudioMP3(_ context.Context, _, _ string) (io.ReadCloser, error) {
|
|
return nil, fmt.Errorf("kokoro not configured (KOKORO_URL is empty)")
|
|
}
|
|
|
|
func (n *noopKokoro) StreamAudioWAV(_ context.Context, _, _ string) (io.ReadCloser, error) {
|
|
return nil, fmt.Errorf("kokoro not configured (KOKORO_URL is empty)")
|
|
}
|
|
|
|
func (n *noopKokoro) ListVoices(_ context.Context) ([]string, error) {
|
|
return nil, nil
|
|
}
|