* feat(watchsync): add plugin-backed providers * fix(watchsync): address plugin review findings * fix(watchsync): harden plugin provider failures * feat(watchsync): complete plugin provider contract * fix(watchsync): address provider review feedback * fix(watchsync): keep device state host-private * fix(watchsync): build reconciliation index concurrently * fix(watchsync): preserve empty device state updates * chore(deps): use released watch-sync SDK --------- Co-authored-by: Quick <31828688+Quick104@users.noreply.github.com>
3444 lines
128 KiB
Go
3444 lines
128 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"flag"
|
|
"fmt"
|
|
"io"
|
|
"io/fs"
|
|
"log"
|
|
"log/slog"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"path/filepath"
|
|
"runtime/debug"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"syscall"
|
|
"time"
|
|
|
|
"github.com/go-chi/chi/v5"
|
|
chimiddleware "github.com/go-chi/chi/v5/middleware"
|
|
"github.com/google/uuid"
|
|
"github.com/hashicorp/go-hclog"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
"github.com/prometheus/client_golang/prometheus/promhttp"
|
|
|
|
pluginv1 "github.com/Silo-Server/silo-plugin-sdk/pkg/pluginproto/silo/plugin/v1"
|
|
sdkcapability "github.com/Silo-Server/silo-plugin-sdk/pkg/pluginsdk/capability"
|
|
|
|
"github.com/Silo-Server/silo-server/internal/access"
|
|
"github.com/Silo-Server/silo-server/internal/activitylog"
|
|
"github.com/Silo-Server/silo-server/internal/adminjob"
|
|
"github.com/Silo-Server/silo-server/internal/api"
|
|
"github.com/Silo-Server/silo-server/internal/api/handlers"
|
|
"github.com/Silo-Server/silo-server/internal/audiobooks"
|
|
"github.com/Silo-Server/silo-server/internal/audiobooks/podcastfeed"
|
|
"github.com/Silo-Server/silo-server/internal/auth"
|
|
"github.com/Silo-Server/silo-server/internal/autoscan"
|
|
"github.com/Silo-Server/silo-server/internal/branding"
|
|
"github.com/Silo-Server/silo-server/internal/cache"
|
|
"github.com/Silo-Server/silo-server/internal/catalog"
|
|
"github.com/Silo-Server/silo-server/internal/catalogseed"
|
|
"github.com/Silo-Server/silo-server/internal/chapterthumbs"
|
|
"github.com/Silo-Server/silo-server/internal/clientip"
|
|
"github.com/Silo-Server/silo-server/internal/config"
|
|
"github.com/Silo-Server/silo-server/internal/database"
|
|
"github.com/Silo-Server/silo-server/internal/diagnostics"
|
|
"github.com/Silo-Server/silo-server/internal/downloads"
|
|
"github.com/Silo-Server/silo-server/internal/ebooks"
|
|
evt "github.com/Silo-Server/silo-server/internal/events"
|
|
"github.com/Silo-Server/silo-server/internal/historyimport"
|
|
"github.com/Silo-Server/silo-server/internal/imagecache"
|
|
"github.com/Silo-Server/silo-server/internal/intromarkers"
|
|
"github.com/Silo-Server/silo-server/internal/jellycompat"
|
|
"github.com/Silo-Server/silo-server/internal/libraryingest"
|
|
"github.com/Silo-Server/silo-server/internal/literaryworks"
|
|
"github.com/Silo-Server/silo-server/internal/logfilter"
|
|
"github.com/Silo-Server/silo-server/internal/logredact"
|
|
"github.com/Silo-Server/silo-server/internal/logstream"
|
|
"github.com/Silo-Server/silo-server/internal/mail"
|
|
"github.com/Silo-Server/silo-server/internal/manga"
|
|
"github.com/Silo-Server/silo-server/internal/markers"
|
|
"github.com/Silo-Server/silo-server/internal/mdblist"
|
|
"github.com/Silo-Server/silo-server/internal/metadata"
|
|
|
|
// Built-in metadata providers self-register into the metadata package's
|
|
// builtin registry on import; buildProviders resolves their seeded chain
|
|
// entries in-process (no gRPC).
|
|
_ "github.com/Silo-Server/silo-server/internal/metadata/nfo"
|
|
"github.com/Silo-Server/silo-server/internal/models"
|
|
"github.com/Silo-Server/silo-server/internal/nodeconfig"
|
|
"github.com/Silo-Server/silo-server/internal/nodepool"
|
|
"github.com/Silo-Server/silo-server/internal/noderecipe"
|
|
"github.com/Silo-Server/silo-server/internal/nodesessions"
|
|
"github.com/Silo-Server/silo-server/internal/notifications"
|
|
"github.com/Silo-Server/silo-server/internal/opslog"
|
|
"github.com/Silo-Server/silo-server/internal/partman"
|
|
"github.com/Silo-Server/silo-server/internal/playback"
|
|
"github.com/Silo-Server/silo-server/internal/pluginhost"
|
|
"github.com/Silo-Server/silo-server/internal/plugins"
|
|
"github.com/Silo-Server/silo-server/internal/policy"
|
|
"github.com/Silo-Server/silo-server/internal/proxy"
|
|
"github.com/Silo-Server/silo-server/internal/ratelimit"
|
|
"github.com/Silo-Server/silo-server/internal/recommendations"
|
|
mediarequests "github.com/Silo-Server/silo-server/internal/requests"
|
|
"github.com/Silo-Server/silo-server/internal/s3client"
|
|
"github.com/Silo-Server/silo-server/internal/scanner"
|
|
"github.com/Silo-Server/silo-server/internal/scanqueue"
|
|
"github.com/Silo-Server/silo-server/internal/secret"
|
|
"github.com/Silo-Server/silo-server/internal/sections"
|
|
"github.com/Silo-Server/silo-server/internal/server"
|
|
"github.com/Silo-Server/silo-server/internal/settingscontract"
|
|
"github.com/Silo-Server/silo-server/internal/subtitles"
|
|
"github.com/Silo-Server/silo-server/internal/taskmanager"
|
|
taskrepository "github.com/Silo-Server/silo-server/internal/taskmanager/repository"
|
|
"github.com/Silo-Server/silo-server/internal/taskmanager/tasks"
|
|
"github.com/Silo-Server/silo-server/internal/taskmanager/triggers"
|
|
"github.com/Silo-Server/silo-server/internal/telemetry"
|
|
"github.com/Silo-Server/silo-server/internal/transcodenode"
|
|
"github.com/Silo-Server/silo-server/internal/usercollections"
|
|
"github.com/Silo-Server/silo-server/internal/userdb"
|
|
"github.com/Silo-Server/silo-server/internal/userstore"
|
|
"github.com/Silo-Server/silo-server/internal/userstore/pgstore"
|
|
"github.com/Silo-Server/silo-server/internal/watchlist"
|
|
"github.com/Silo-Server/silo-server/internal/watchstate"
|
|
"github.com/Silo-Server/silo-server/internal/watchsync"
|
|
watchmdblist "github.com/Silo-Server/silo-server/internal/watchsync/providers/mdblist"
|
|
"github.com/Silo-Server/silo-server/internal/watchsync/providers/simkl"
|
|
"github.com/Silo-Server/silo-server/internal/watchsync/providers/trakt"
|
|
"github.com/Silo-Server/silo-server/internal/worker"
|
|
"github.com/Silo-Server/silo-server/migrations"
|
|
siloweb "github.com/Silo-Server/silo-server/web"
|
|
)
|
|
|
|
// resolveNodeIdentity returns a stable node identifier used by the
|
|
// heartbeat writer, reconciler, and shutdown cleanup. Resolution order:
|
|
// SILO_NODE_NAME > NODE_NAME > os.Hostname().
|
|
func resolveNodeIdentity() string {
|
|
if v := os.Getenv("SILO_NODE_NAME"); v != "" {
|
|
return v
|
|
}
|
|
if v := os.Getenv("NODE_NAME"); v != "" {
|
|
return v
|
|
}
|
|
h, _ := os.Hostname()
|
|
return h
|
|
}
|
|
|
|
func resolvePluginCacheDir() string {
|
|
if v := strings.TrimSpace(os.Getenv("SILO_PLUGIN_CACHE_DIR")); v != "" {
|
|
return v
|
|
}
|
|
return filepath.Join(os.TempDir(), "silo-plugins")
|
|
}
|
|
|
|
func buildBaseHandler(format string, level slog.Leveler, otelHandler slog.Handler) slog.Handler {
|
|
opts := &slog.HandlerOptions{Level: level}
|
|
var console slog.Handler
|
|
if strings.EqualFold(format, "json") {
|
|
console = slog.NewJSONHandler(os.Stderr, opts)
|
|
} else {
|
|
console = slog.NewTextHandler(os.Stderr, opts)
|
|
}
|
|
if otelHandler == nil {
|
|
// Redact secrets before they reach stderr (the opslog DB path redacts
|
|
// separately when flattening rows).
|
|
return logredact.New(console)
|
|
}
|
|
// Fan out to the console and the OTel bridge. The OTel branch is level-gated
|
|
// by the shared level var so console and OTLP share one verbosity knob (see
|
|
// telemetry.LevelGated) — otherwise slog.MultiHandler.Enabled would OR the
|
|
// branches and export Debug records while stderr stays silent. The whole
|
|
// fan-out is wrapped in secret redaction so console and OTLP both emit
|
|
// masked output (the opslog DB path redacts separately).
|
|
return logredact.New(telemetry.FanOut(console, telemetry.LevelGated(otelHandler, level)))
|
|
}
|
|
|
|
func parseLogLevel(level string) slog.Level {
|
|
switch strings.ToLower(level) {
|
|
case "debug":
|
|
return slog.LevelDebug
|
|
case "warn", "warning":
|
|
return slog.LevelWarn
|
|
case "error":
|
|
return slog.LevelError
|
|
default:
|
|
return slog.LevelInfo
|
|
}
|
|
}
|
|
|
|
func mustGetSetting(store interface {
|
|
Get(context.Context, string) (string, error)
|
|
}, ctx context.Context, key, fallback string) string {
|
|
value, err := store.Get(ctx, key)
|
|
if err != nil || strings.TrimSpace(value) == "" {
|
|
return fallback
|
|
}
|
|
return value
|
|
}
|
|
|
|
func configureOperationalLogging(
|
|
ctx context.Context,
|
|
pool *pgxpool.Pool,
|
|
settingsRepo catalog.SettingsStore,
|
|
redisCfg config.RedisConfig,
|
|
logStreamHub *logstream.Hub,
|
|
filteredHandler slog.Handler,
|
|
nodeID string,
|
|
) (opslog.Writer, *opslog.Repo, *partman.Manager) {
|
|
if err := opslog.SeedDefaults(ctx, settingsRepo); err != nil {
|
|
log.Fatalf("seed opslog defaults: %v", err)
|
|
}
|
|
if err := diagnostics.SeedDefaults(ctx, settingsRepo); err != nil {
|
|
log.Fatalf("seed diagnostics defaults: %v", err)
|
|
}
|
|
opsPM := partman.NewManager(pool, "operational_logs", partman.Daily, 3)
|
|
if err := opsPM.EnsureFuturePartitions(ctx); err != nil {
|
|
// Non-fatal: a partition hiccup must not crash-loop the server (see the
|
|
// operational_logs partition incident). Writes fall back to the default
|
|
// partition and the periodic cleanup retries EnsureFuturePartitions.
|
|
slog.WarnContext(ctx, "ensure operational log partitions; continuing in degraded mode", "component", "app", "error", err)
|
|
}
|
|
|
|
var operationalWriter opslog.Writer
|
|
operationalConsumer := opslog.NewConsumer(pool, nil, logStreamHub)
|
|
if redisCfg.URL != "" {
|
|
redisClient, redisErr := cache.NewRedisClient(redisCfg)
|
|
if redisErr == nil && redisClient != nil {
|
|
operationalWriter = opslog.NewRedisWriter(redisClient)
|
|
operationalConsumer = opslog.NewConsumer(pool, redisClient, logStreamHub)
|
|
go operationalConsumer.RunRedis(ctx)
|
|
}
|
|
}
|
|
if operationalWriter == nil {
|
|
memWriter := opslog.NewMemoryWriter(10000)
|
|
operationalWriter = memWriter
|
|
go operationalConsumer.RunMemory(ctx, memWriter.Chan())
|
|
}
|
|
|
|
opsCaptureLevel := slog.LevelInfo
|
|
switch strings.ToLower(strings.TrimSpace(mustGetSetting(settingsRepo, ctx, "opslog.capture_level", "info"))) {
|
|
case "debug":
|
|
opsCaptureLevel = slog.LevelDebug
|
|
case "warn", "warning":
|
|
opsCaptureLevel = slog.LevelWarn
|
|
case "error":
|
|
opsCaptureLevel = slog.LevelError
|
|
}
|
|
|
|
slog.SetDefault(slog.New(opslog.NewHandler(filteredHandler, operationalWriter, opsCaptureLevel, nodeID)))
|
|
|
|
return operationalWriter, opslog.NewRepo(pool), opsPM
|
|
}
|
|
|
|
func maybeApplyPostgresTuning(ctx context.Context, pool *pgxpool.Pool, appMaxConnections int, mode string) {
|
|
switch strings.ToLower(strings.TrimSpace(mode)) {
|
|
case "", "integrated", "api":
|
|
default:
|
|
return
|
|
}
|
|
|
|
opts, err := database.LoadPostgresTuneOptionsFromEnv(appMaxConnections)
|
|
if err != nil {
|
|
slog.WarnContext(ctx, "postgres auto-tuning disabled", "component", "app", "error", err)
|
|
return
|
|
}
|
|
if !opts.Enabled {
|
|
return
|
|
}
|
|
|
|
tuneCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
|
|
defer cancel()
|
|
|
|
result, err := database.ApplyPostgresTuning(tuneCtx, pool, opts)
|
|
for _, failure := range result.Failures {
|
|
slog.WarnContext(ctx, "postgres auto-tuning setting failed", "component", "app",
|
|
"name", failure.Name,
|
|
"value", failure.Value,
|
|
"error", failure.Err,
|
|
)
|
|
}
|
|
if err != nil {
|
|
slog.WarnContext(ctx, "postgres auto-tuning failed", "component", "app",
|
|
"error", err,
|
|
"applied", result.Applied,
|
|
"failures", len(result.Failures),
|
|
)
|
|
return
|
|
}
|
|
|
|
slog.InfoContext(ctx, "postgres auto-tuning applied", "component", "app",
|
|
"profile", opts.Profile,
|
|
"postgres_major", result.PostgresMajorVersion,
|
|
"settings", result.Applied,
|
|
"resets", len(result.Reset),
|
|
"failures", len(result.Failures),
|
|
"memory_budget_bytes", opts.MemoryBudgetBytes,
|
|
"detected_memory_bytes", opts.DetectedMemoryBytes,
|
|
"memory_source", opts.MemorySource,
|
|
"memory_budget_percent", opts.MemoryBudgetPercent,
|
|
"cpus", opts.CPUs,
|
|
"connections", opts.Connections,
|
|
"storage", opts.Storage,
|
|
"db_size", result.DBSize,
|
|
"database_size_bytes", result.DatabaseSizeBytes,
|
|
)
|
|
if len(result.RestartRequired) > 0 {
|
|
slog.WarnContext(ctx, "postgres restart required to finish applying auto-tuned settings", "component", "app",
|
|
"settings", strings.Join(result.RestartRequired, ","),
|
|
)
|
|
}
|
|
if len(result.Reset) > 0 {
|
|
slog.InfoContext(ctx, "postgres auto-tuning reset stale settings", "component", "app",
|
|
"settings", strings.Join(result.Reset, ","),
|
|
)
|
|
}
|
|
}
|
|
|
|
// runCredentialBackfills sweeps any plaintext server-owned credential to
|
|
// ciphertext on the primary (migration-running) node. All passes are
|
|
// best-effort: a failed row leaves the prior plaintext (no new exposure) and
|
|
// still reads via the read-path pass-through, so a backfill error must never
|
|
// block boot. The sensitive-settings pass runs first so the arr
|
|
// resolve-then-encrypt pass sees consistent referenced settings.
|
|
// librarySettingsCleaner wires the per-user canonical settings cleanup the
|
|
// library delete job runs, or nil when the user store is unavailable — the
|
|
// executor treats a nil cleaner as "skip".
|
|
func librarySettingsCleaner(pool *pgxpool.Pool, stores userstore.UserStoreProvider) adminjob.LibrarySettingsCleaner {
|
|
if pool == nil || stores == nil {
|
|
return nil
|
|
}
|
|
return userstore.NewSettingValuesCleaner(auth.NewUserRepository(pool), stores)
|
|
}
|
|
|
|
func runCredentialBackfills(ctx context.Context, pool *pgxpool.Pool, cipher *secret.Cipher, settings *catalog.EncryptedSettingsRepo) {
|
|
settingsN, err := settings.BackfillSensitiveSettings(ctx)
|
|
if err != nil {
|
|
slog.ErrorContext(ctx, "secret backfill: sensitive settings", "component", "app", "error", err)
|
|
}
|
|
columnsN, err := secret.BackfillColumns(ctx, pool, cipher, secret.ColumnBackfillTargets())
|
|
if err != nil {
|
|
slog.ErrorContext(ctx, "secret backfill: credential columns", "component", "app", "error", err)
|
|
}
|
|
historyServersN, err := historyimport.NewRepository(pool, cipher).BackfillSessionServerSecrets(ctx)
|
|
if err != nil {
|
|
slog.ErrorContext(ctx, "secret backfill: history import session server credentials", "component", "app", "error", err)
|
|
}
|
|
// The arr resolver is the encrypting settings decorator: it decrypts a
|
|
// sensitive target (e.g. requests.radarr.api_key) or passes through a
|
|
// plaintext custom key, exactly replicating the deleted resolveAPIKey.
|
|
arrN, err := secret.BackfillReferencedColumns(ctx, pool, cipher, settings.Get, secret.ArrKeyBackfillTargets())
|
|
if err != nil {
|
|
slog.ErrorContext(ctx, "secret backfill: arr api keys", "component", "app", "error", err)
|
|
}
|
|
pluginConfigsN, err := plugins.NewRuntimeConfigStore(pool, cipher).BackfillEncryptedConfigs(ctx)
|
|
if err != nil {
|
|
slog.ErrorContext(ctx, "secret backfill: plugin runtime configs", "component", "app", "error", err)
|
|
}
|
|
if total := settingsN + columnsN + historyServersN + arrN + pluginConfigsN; total > 0 {
|
|
slog.InfoContext(ctx, "secret backfill: encrypted plaintext credentials at rest", "component", "app",
|
|
"settings", settingsN, "columns", columnsN, "history_session_servers", historyServersN,
|
|
"arr_keys", arrN, "plugin_configs", pluginConfigsN, "total", total)
|
|
}
|
|
}
|
|
|
|
func runCompatWebCommand(ctx context.Context, args []string) error {
|
|
if len(args) == 0 {
|
|
return fmt.Errorf("usage: silo compat-web {status|install|update|remove}")
|
|
}
|
|
command := args[0]
|
|
flags := flag.NewFlagSet("compat-web "+command, flag.ContinueOnError)
|
|
flags.SetOutput(io.Discard)
|
|
root := flags.String("dir", config.DefaultJellyfinWebInstallDir, "Jellyfin Web component install root")
|
|
version := flags.String("version", config.DefaultJellyfinWebVersion, "Jellyfin Web version without leading v")
|
|
source := flags.String("source", jellycompat.DefaultWebSourceURL, "upstream jellyfin-web git repository")
|
|
if err := flags.Parse(args[1:]); err != nil {
|
|
return err
|
|
}
|
|
|
|
switch command {
|
|
case "status":
|
|
status := jellycompat.WebComponentStatusForConfig(&config.Config{
|
|
JellyfinCompat: config.JellyfinCompatConfig{
|
|
Enabled: false,
|
|
WebVersion: *version,
|
|
WebInstallDir: *root,
|
|
WebDir: filepath.Join(*root, "current"),
|
|
},
|
|
}, map[string]string{
|
|
"jellyfin_compat.web_source_url": *source,
|
|
})
|
|
return json.NewEncoder(os.Stdout).Encode(status)
|
|
case "install", "update":
|
|
status, err := jellycompat.InstallWebComponent(ctx, jellycompat.WebComponentInstallOptions{
|
|
InstallRoot: *root,
|
|
SourceURL: *source,
|
|
Version: *version,
|
|
})
|
|
_ = json.NewEncoder(os.Stdout).Encode(status)
|
|
return err
|
|
case "remove":
|
|
return jellycompat.RemoveWebComponent(*root)
|
|
default:
|
|
return fmt.Errorf("unknown compat-web command %q", command)
|
|
}
|
|
}
|
|
|
|
func main() {
|
|
if len(os.Args) > 1 && os.Args[1] == "compat-web" {
|
|
if err := runCompatWebCommand(context.Background(), os.Args[2:]); err != nil {
|
|
log.Fatalf("compat-web: %v", err)
|
|
}
|
|
return
|
|
}
|
|
|
|
envFile := flag.String("env", ".env", "path to .env bootstrap file")
|
|
migrateOnly := flag.Bool("migrate-only", false, "apply database migrations and exit")
|
|
migrateStatus := flag.Bool("migrate-status", false, "show database migration status and exit")
|
|
migrateDownTo := flag.Int64("migrate-down-to", -1,
|
|
"roll back every migration newer than this version and exit (the version to KEEP)")
|
|
flag.Parse()
|
|
|
|
ctx := context.Background()
|
|
|
|
// Step 0: Validate the embedded settings contract before anything can
|
|
// depend on it. A malformed or self-inconsistent manifest is a build defect,
|
|
// not a runtime condition, so failing here — loudly, before the first
|
|
// request — is the whole point: the alternative is shipping an image whose
|
|
// contract disagrees with the clients that vendored it.
|
|
contract, err := settingscontract.Load()
|
|
if err != nil {
|
|
log.Fatalf("settings contract: %v", err)
|
|
}
|
|
contractETag, err := settingscontract.ETag()
|
|
if err != nil {
|
|
log.Fatalf("settings contract: %v", err)
|
|
}
|
|
slog.Info("settings contract loaded",
|
|
"revision", contract.Revision,
|
|
"definitions", len(contract.Definitions),
|
|
"etag", contractETag)
|
|
|
|
// Step 1: Bootstrap from .env
|
|
bc, err := config.LoadBootstrap(*envFile)
|
|
if err != nil {
|
|
log.Fatalf("bootstrap: %v", err)
|
|
}
|
|
|
|
// Construct the at-rest credential cipher from SECRET_KEY immediately after
|
|
// bootstrap, before any settings repo is built. It is threaded explicitly as
|
|
// a dependency into every repo that stores a server-owned secret — never a
|
|
// package-level global.
|
|
dataCipher, err := secret.New(bc.SecretKey)
|
|
if err != nil {
|
|
log.Fatalf("secret cipher: %v", err)
|
|
}
|
|
|
|
// Step 2: Connect to PostgreSQL (bootstrap pool with default max connections)
|
|
bootstrapDBCfg := config.DatabaseConfig{URL: bc.DatabaseURL, MaxConnections: 20}
|
|
pool, err := database.NewPool(ctx, bootstrapDBCfg)
|
|
if err != nil {
|
|
log.Fatalf("database pool: %v", err)
|
|
}
|
|
defer pool.Close()
|
|
slog.Info("connected to PostgreSQL")
|
|
|
|
if *migrateStatus {
|
|
migCtx, migCancel := database.MigrationContext(ctx)
|
|
statuses, statusErr := database.MigrationStatuses(migCtx, pool, migrations.FS, "sql")
|
|
migCancel()
|
|
if statusErr != nil {
|
|
log.Fatalf("failed to read migration status: %v", statusErr)
|
|
}
|
|
|
|
fmt.Printf("%-8s %8s %-25s %s\n", "STATE", "VERSION", "APPLIED_AT", "MIGRATION")
|
|
for _, status := range statuses {
|
|
appliedAt := "-"
|
|
if !status.AppliedAt.IsZero() {
|
|
appliedAt = status.AppliedAt.UTC().Format(time.RFC3339)
|
|
}
|
|
source := status.Source
|
|
if source != "" {
|
|
source = filepath.Base(source)
|
|
} else {
|
|
source = "-"
|
|
}
|
|
fmt.Printf("%-8s %8d %-25s %s\n", status.State, status.Version, appliedAt, source)
|
|
}
|
|
return
|
|
}
|
|
|
|
if *migrateDownTo >= 0 {
|
|
// Deliberately its own flag rather than a mode of --migrate-only: this
|
|
// discards data, and several of the migrations it reverses are Go ones
|
|
// the goose CLI cannot reach, so it is the only way to undo them
|
|
// short of restoring a backup.
|
|
migCtx, migCancel := database.MigrationContext(ctx)
|
|
migErr := database.MigrateDownTo(migCtx, pool, migrations.FS, "sql", *migrateDownTo)
|
|
migCancel()
|
|
if migErr != nil {
|
|
log.Fatalf("failed to roll back migrations: %v", migErr)
|
|
}
|
|
slog.Info("database migrations rolled back", "kept_through_version", *migrateDownTo)
|
|
return
|
|
}
|
|
|
|
if *migrateOnly {
|
|
migCtx, migCancel := database.MigrationContext(ctx)
|
|
migErr := database.RunMigrations(migCtx, pool, migrations.FS, "sql")
|
|
migCancel()
|
|
if migErr != nil {
|
|
log.Fatalf("failed to run migrations: %v", migErr)
|
|
}
|
|
slog.Info("database migrations applied")
|
|
return
|
|
}
|
|
|
|
// Run migrations only for integrated/api modes. Proxy and transcode nodes
|
|
// should never alter the schema — they may scale independently and would
|
|
// race or apply migrations before the primary node is deliberately upgraded.
|
|
// The same gate decides whether this node runs the credential-encryption
|
|
// backfills: only the primary (migration-running) node sweeps plaintext to
|
|
// ciphertext; secondary nodes read whatever the primary encrypted.
|
|
isPrimaryNode := bc.Mode == "integrated" || bc.Mode == "api" || bc.Mode == ""
|
|
if isPrimaryNode {
|
|
migCtx, migCancel := database.MigrationContext(ctx)
|
|
if migErr := database.RunMigrations(migCtx, pool, migrations.FS, "sql"); migErr != nil {
|
|
migCancel()
|
|
log.Fatalf("failed to run migrations: %v", migErr)
|
|
}
|
|
migCancel()
|
|
slog.Info("database migrations applied")
|
|
}
|
|
|
|
// Step 3: Load settings from DB. settingsRepo is the encrypting decorator so
|
|
// every consumer (config.LoadFromDB, admin, ABS, watchers) transparently sees
|
|
// plaintext while sensitive keys rest as ciphertext. The settings backfill
|
|
// (run after migrations, before this GetAll) is wired further below.
|
|
settingsRepo := catalog.NewEncryptedSettingsRepo(catalog.NewServerSettingsRepo(pool), dataCipher)
|
|
if isPrimaryNode {
|
|
runCredentialBackfills(ctx, pool, dataCipher, settingsRepo)
|
|
}
|
|
settings, err := settingsRepo.GetAll(ctx)
|
|
if err != nil {
|
|
log.Fatalf("loading settings: %v", err)
|
|
}
|
|
|
|
// Step 4: YAML import (one-time)
|
|
yamlPath := "silo.yaml"
|
|
if _, yamlErr := os.Stat(yamlPath); yamlErr == nil {
|
|
if settings["_yaml_imported"] == "" {
|
|
yamlSettings, importErr := config.YAMLToSettingsMap(yamlPath)
|
|
if importErr != nil {
|
|
log.Printf("WARN: could not import YAML config: %v", importErr)
|
|
} else {
|
|
for k, v := range yamlSettings {
|
|
if err := settingsRepo.Set(ctx, k, v); err != nil {
|
|
log.Printf("WARN: failed to import setting %s: %v", k, err)
|
|
}
|
|
}
|
|
if err := settingsRepo.Set(ctx, "_yaml_imported", "true"); err != nil {
|
|
slog.Warn("failed to set yaml import flag", "error", err)
|
|
}
|
|
log.Println("Imported config from silo.yaml — this file is no longer used")
|
|
settings, _ = settingsRepo.GetAll(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// Step 5: Auto-generate secrets
|
|
if settings["auth.jwt_secret"] == "" {
|
|
secret := make([]byte, 32)
|
|
if _, err := rand.Read(secret); err != nil {
|
|
log.Fatalf("generating jwt secret: %v", err)
|
|
}
|
|
encoded := base64.StdEncoding.EncodeToString(secret)
|
|
if err := settingsRepo.Set(ctx, "auth.jwt_secret", encoded); err != nil {
|
|
slog.Warn("failed to persist generated JWT secret", "error", err)
|
|
}
|
|
settings["auth.jwt_secret"] = encoded
|
|
}
|
|
if settings["jellyfin_compat.server_id"] == "" {
|
|
serverID := uuid.NewSHA1(uuid.NameSpaceURL, []byte("https://silo.local/jellycompat")).String()
|
|
if err := settingsRepo.Set(ctx, "jellyfin_compat.server_id", serverID); err != nil {
|
|
slog.Warn("failed to persist generated server ID", "error", err)
|
|
}
|
|
settings["jellyfin_compat.server_id"] = serverID
|
|
}
|
|
|
|
// Step 6: Build config from DB
|
|
cfg, err := config.LoadFromDB(settings)
|
|
if err != nil {
|
|
log.Fatalf("building config: %v", err)
|
|
}
|
|
|
|
// Step 7: Apply bootstrap overrides
|
|
cfg.Server.Listen = bc.Listen
|
|
cfg.Server.Mode = bc.Mode
|
|
cfg.Database.URL = bc.DatabaseURL
|
|
cfg.JellyfinCompat.Listen = bc.JFListen
|
|
if bc.RedisURL != "" {
|
|
cfg.Redis.URL = bc.RedisURL
|
|
}
|
|
|
|
// Step 8: Recreate pool if max_connections differs from bootstrap default
|
|
if cfg.Database.MaxConnections != bootstrapDBCfg.MaxConnections {
|
|
pool.Close()
|
|
pool, err = database.NewPool(ctx, cfg.Database)
|
|
if err != nil {
|
|
log.Fatalf("recreating pool with configured max_connections: %v", err)
|
|
}
|
|
}
|
|
// Re-wrap with the encrypting decorator so the recreated pool's settings repo
|
|
// still encrypts/decrypts — no raw settings repo may escape into later wiring.
|
|
settingsRepo = catalog.NewEncryptedSettingsRepo(catalog.NewServerSettingsRepo(pool), dataCipher)
|
|
nodeID := resolveNodeIdentity()
|
|
catalogSearchStartupSettings, err := catalog.CatalogSearchSettingsFromMap(settings)
|
|
if err != nil {
|
|
slog.Warn("catalog search: failed to load settings for startup wiring; using postgres", "err", err)
|
|
catalogSearchStartupSettings = catalog.DefaultCatalogSearchSettings()
|
|
}
|
|
activeCatalogSearchProvider := catalog.ActiveCatalogSearchProvider(catalogSearchStartupSettings)
|
|
|
|
// Step 9: Validate
|
|
if err := cfg.Validate(); err != nil {
|
|
log.Fatalf("config validation: %v", err)
|
|
}
|
|
|
|
// Step 10: Configure log level. The level var and quiet filter are
|
|
// shared with the operational-logging handler chain and hot-reloaded by
|
|
// the config watcher in integrated mode.
|
|
logLevelVar := new(slog.LevelVar)
|
|
logLevelVar.Set(parseLogLevel(cfg.Server.LogLevel))
|
|
|
|
// Bootstrap OpenTelemetry (logs + traces) before installing the log handler
|
|
// chain. Setup depends only on OTEL_* / SILO_OTEL_ENABLED env (not the DB),
|
|
// so it is safe to call here. When disabled, this is fully dormant: no
|
|
// providers are installed and telemetryShutdown is a no-op.
|
|
telemetryCfg := telemetry.LoadConfig(nodeID)
|
|
telemetryProviders, telemetryShutdown, err := telemetry.Setup(ctx, telemetryCfg)
|
|
if err != nil {
|
|
// Telemetry is best-effort: a malformed OTEL_* environment must not
|
|
// crash-loop the server. Setup installed no globals and returned no-op
|
|
// providers, so continue with telemetry disabled.
|
|
slog.ErrorContext(ctx, "telemetry setup failed; continuing with telemetry disabled", "component", "app", "error", err)
|
|
telemetryCfg.Enabled = false
|
|
}
|
|
defer func() {
|
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
if err := telemetryShutdown(shutdownCtx); err != nil {
|
|
slog.WarnContext(shutdownCtx, "telemetry shutdown error", "component", "app", "error", err)
|
|
}
|
|
}()
|
|
|
|
var otelLogHandler slog.Handler
|
|
if telemetryCfg.Enabled {
|
|
otelLogHandler = telemetry.NewOTelHandler(telemetryProviders.LoggerProvider)
|
|
}
|
|
|
|
baseHandler := buildBaseHandler(cfg.Server.LogFormat, logLevelVar, otelLogHandler)
|
|
quietFilter := logfilter.New(baseHandler, cfg.Server.LogQuiet)
|
|
slog.SetDefault(slog.New(quietFilter))
|
|
|
|
mode := cfg.Server.Mode
|
|
maybeApplyPostgresTuning(ctx, pool, cfg.Database.MaxConnections, mode)
|
|
slog.Info("silo starting", "mode", mode, "listen", cfg.Server.Listen, "log_level", cfg.Server.LogLevel, "node_id", nodeID)
|
|
|
|
appCtx, appCancel := context.WithCancel(ctx)
|
|
defer appCancel()
|
|
restartReqCh := make(chan struct{}, 1)
|
|
var restartRequested atomic.Bool
|
|
|
|
eventBus := cache.NewEventBus(cfg.Redis.URL)
|
|
logStreamHub := logstream.NewHub(nodeID, eventBus)
|
|
if err := logStreamHub.Start(appCtx); err != nil {
|
|
log.Fatalf("log stream hub start: %v", err)
|
|
}
|
|
realtimeHub := notifications.NewHub(nodeID, eventBus)
|
|
if err := realtimeHub.Start(appCtx); err != nil {
|
|
log.Fatalf("realtime hub start: %v", err)
|
|
}
|
|
eventsHub := realtimeHub.EventsHub()
|
|
scanRegistry := evt.NewScanRegistry()
|
|
operationalWriter, opsRepo, opsPM := configureOperationalLogging(appCtx, pool, settingsRepo, cfg.Redis, logStreamHub, quietFilter, nodeID)
|
|
defer func() {
|
|
if err := eventBus.Close(); err != nil {
|
|
slog.Warn("event bus close error", "error", err)
|
|
}
|
|
}()
|
|
|
|
// Proxy and transcode modes run with DB + Redis for hot-reload.
|
|
if mode == "proxy" || mode == "transcode" {
|
|
redisClient, err := cache.NewRedisClient(cfg.Redis)
|
|
if err != nil || redisClient == nil {
|
|
slog.Error("redis is required for this mode", "mode", mode, "error", err)
|
|
os.Exit(1)
|
|
}
|
|
|
|
bootstrap := nodeconfig.BootstrapOverrides{
|
|
Listen: cfg.Server.Listen,
|
|
Mode: cfg.Server.Mode,
|
|
DatabaseURL: cfg.Database.URL,
|
|
JFListen: cfg.JellyfinCompat.Listen,
|
|
RedisURL: bc.RedisURL,
|
|
}
|
|
watcher := nodeconfig.NewWatcher(pool, dataCipher, eventBus, bootstrap)
|
|
if err := watcher.Start(appCtx); err != nil {
|
|
slog.Error("config watcher start failed", "error", err)
|
|
os.Exit(1)
|
|
}
|
|
|
|
nodeURL := os.Getenv("NODE_URL")
|
|
nodeName := os.Getenv("NODE_NAME")
|
|
if nodeURL == "" {
|
|
nodeURL = "http://localhost" + cfg.Server.Listen
|
|
slog.Warn("NODE_URL not set, using listen address — session keys may collide across nodes")
|
|
}
|
|
if nodeName == "" {
|
|
nodeName = mode
|
|
}
|
|
|
|
tracker := nodesessions.NewTracker(redisClient, nodeURL, nodeName, mode)
|
|
tracker.StartRefresh(appCtx)
|
|
defer func() {
|
|
cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cleanupCancel()
|
|
tracker.Cleanup(cleanupCtx)
|
|
}()
|
|
|
|
var handler http.Handler
|
|
if mode == "proxy" {
|
|
srv := proxy.NewServer(watcher, tracker)
|
|
handler = srv.Handler()
|
|
} else {
|
|
srv := transcodenode.NewServer(watcher, tracker)
|
|
srv.SetFFmpegLogSink(playback.NewSlogFFmpegLogSink(slog.Default(), nodeID))
|
|
// Read jellycompat reconstruction recipes central wrote at transcode
|
|
// start, so this node can rebuild a Jellyfin transcode after its own
|
|
// restart (the node hop token is recipe-less). Shares the offload Redis.
|
|
srv.SetRecipeStore(noderecipe.NewStore(redisClient, 0))
|
|
// Reclaim orphaned transcode dirs at boot and hourly thereafter, bound
|
|
// to appCtx so it stops on shutdown.
|
|
srv.StartOrphanSweeper(appCtx)
|
|
handler = srv.Handler()
|
|
}
|
|
|
|
_ = operationalWriter
|
|
_ = opsRepo
|
|
startStandaloneServer(cfg.Server.Listen, handler)
|
|
return
|
|
}
|
|
|
|
// Hot-reload config watcher for integrated/api mode. Reloads on
|
|
// EventSettingsChanged (Redis) with a 60s poll fallback, so settings
|
|
// changes apply without restart even on Redis-less deployments. The
|
|
// watcher's config supersedes the startup snapshot from here on.
|
|
configWatcher := nodeconfig.NewWatcher(pool, dataCipher, eventBus, nodeconfig.BootstrapOverrides{
|
|
Listen: bc.Listen,
|
|
Mode: bc.Mode,
|
|
DatabaseURL: bc.DatabaseURL,
|
|
JFListen: bc.JFListen,
|
|
RedisURL: bc.RedisURL,
|
|
})
|
|
if err := configWatcher.Start(appCtx); err != nil {
|
|
log.Fatalf("config watcher start: %v", err)
|
|
}
|
|
cfg = configWatcher.Config()
|
|
|
|
// Apply server.log_level / server.log_quiet changes live. Both feed the
|
|
// shared level var and quiet filter inside the default logger chain.
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
logLevelVar.Set(parseLogLevel(updated.Server.LogLevel))
|
|
quietFilter.SetQuiet(updated.Server.LogQuiet)
|
|
})
|
|
|
|
// Determine which components to initialize based on mode.
|
|
needsS3 := mode == "integrated" || mode == "api"
|
|
needsScanner := mode == "integrated" || mode == "api"
|
|
needsUserDB := mode == "integrated" || mode == "api"
|
|
needsWorkers := mode == "integrated" || mode == "api"
|
|
|
|
bootstrapSensitiveConfigured := map[string]bool{}
|
|
bootstrapSensitiveValues := map[string]string{}
|
|
if bc.RedisURL != "" {
|
|
bootstrapSensitiveConfigured["redis.url"] = true
|
|
bootstrapSensitiveValues["redis.url"] = bc.RedisURL
|
|
}
|
|
if rawTrustedProxies := strings.TrimSpace(os.Getenv(clientip.EnvTrustedProxies)); rawTrustedProxies != "" {
|
|
normalizedTrustedProxies, normalizeErr := clientip.NormalizeCIDRList(rawTrustedProxies)
|
|
if normalizeErr != nil {
|
|
log.Fatalf("invalid %s: %v", clientip.EnvTrustedProxies, normalizeErr)
|
|
}
|
|
bootstrapSensitiveConfigured[clientip.SettingTrustedProxies] = true
|
|
bootstrapSensitiveValues[clientip.SettingTrustedProxies] = normalizedTrustedProxies
|
|
}
|
|
|
|
// Shared Redis client for components needing raw Redis beyond the event
|
|
// bus (websocket handshake tickets, session listing). Nil on Redis-less
|
|
// deployments; consumers fall back to in-process implementations.
|
|
apiRedisClient, apiRedisErr := cache.NewRedisClient(cfg.Redis)
|
|
if apiRedisErr != nil {
|
|
slog.Warn("redis client init failed; multi-node websocket tickets disabled", "error", apiRedisErr)
|
|
} else if apiRedisClient != nil {
|
|
defer func() { _ = apiRedisClient.Close() }()
|
|
}
|
|
|
|
// Assigned below once the trusted-proxy config is seeded; captured by the
|
|
// OnServerSettingUpdated closure, which only runs on admin requests after
|
|
// startup completes.
|
|
var ipResolver *clientip.Resolver
|
|
normalizedBootstrapRedisURL, bootstrapRedisURLErr := config.NormalizeRedisURL(bc.RedisURL)
|
|
redisBootstrapAvailable := (normalizedBootstrapRedisURL != "" && bootstrapRedisURLErr == nil) ||
|
|
(strings.TrimSpace(cfg.Redis.SentinelMaster) != "" && len(cfg.Redis.SentinelAddresses) > 0)
|
|
|
|
deps := api.Dependencies{
|
|
Config: cfg,
|
|
LiveConfig: configWatcher.Config,
|
|
OnConfigChange: configWatcher.OnChange,
|
|
BootstrapSensitiveConfigured: bootstrapSensitiveConfigured,
|
|
BootstrapSensitiveValues: bootstrapSensitiveValues,
|
|
RedisBootstrapAvailable: redisBootstrapAvailable,
|
|
AppContext: appCtx,
|
|
DB: pool,
|
|
SecretCipher: dataCipher,
|
|
EventBus: eventBus,
|
|
RedisClient: apiRedisClient,
|
|
LogStreamHub: logStreamHub,
|
|
RealtimeHub: realtimeHub,
|
|
EventsHub: eventsHub,
|
|
ScanRegistry: scanRegistry,
|
|
OpsLogRepo: opsRepo,
|
|
FFmpegLogSink: playback.NewSlogFFmpegLogSink(slog.Default(), nodeID),
|
|
PublicURL: os.Getenv("SILO_PUBLIC_URL"),
|
|
RequestServerRestart: func(context.Context) error {
|
|
if !restartRequested.CompareAndSwap(false, true) {
|
|
return handlers.ErrServerRestartAlreadyRequested
|
|
}
|
|
restartReqCh <- struct{}{}
|
|
return nil
|
|
},
|
|
OnServerSettingUpdated: func(_ context.Context, key, _ string) {
|
|
// Key-scoped reload for the client-IP trust boundary: unlike the
|
|
// whole-config watcher reload below, this cannot be blocked by an
|
|
// unrelated malformed setting failing config.LoadFromDB. Uses a
|
|
// fresh context — the setting is already persisted, so the reload
|
|
// must not be skipped because the admin request was canceled.
|
|
if key == clientip.SettingTrustedProxies && ipResolver != nil {
|
|
if cidrs, loadErr := clientip.LoadTrustedCIDRs(context.Background(), settingsRepo); loadErr != nil {
|
|
slog.WarnContext(context.Background(), "clientip config reload failed", "component", "app", "error", loadErr)
|
|
} else {
|
|
ipResolver.UpdateTrustedCIDRs(cidrs)
|
|
}
|
|
}
|
|
// Nudge the hot-reload watcher so same-process settings changes
|
|
// apply immediately even without Redis (the event bus is a no-op
|
|
// then, leaving only the 60s poll).
|
|
configWatcher.RequestReload()
|
|
},
|
|
}
|
|
accessGroupStore := access.NewGroupStore(pool)
|
|
audiobooksService := audiobooks.New(&audiobooksSettingsAdapter{repo: settingsRepo})
|
|
absCompatEnabled, err := audiobooksService.ABSCompatEnabled(appCtx)
|
|
if err != nil {
|
|
slog.Warn("Audiobookshelf compatibility disabled; failed to read setting", "err", err)
|
|
absCompatEnabled = false
|
|
}
|
|
adminJobCancelRegistry := adminjob.NewCancelRegistry()
|
|
deps.AdminJobCancelRegistry = adminJobCancelRegistry
|
|
if needsWorkers && deps.DB != nil {
|
|
deps.IntroRepository = intromarkers.NewRepository(deps.DB)
|
|
deps.IntroAnalyzer = intromarkers.NewAnalyzer(
|
|
deps.IntroRepository,
|
|
intromarkers.DefaultConfig(cfg.Playback.FFmpegPath),
|
|
slog.Default(),
|
|
)
|
|
}
|
|
if deps.DB != nil {
|
|
markerRegistry := markers.NewRegistry(slog.Default())
|
|
markerProviderConfig := markers.NewProviderConfigStore(deps.DB)
|
|
if err := markerProviderConfig.Reload(appCtx); err != nil {
|
|
slog.Warn("load marker provider config failed; falling back to registration-order fetch",
|
|
"error", err)
|
|
} else {
|
|
markerRegistry.UseConfigStore(markerProviderConfig)
|
|
if deps.EventBus != nil {
|
|
if err := deps.EventBus.Subscribe(appCtx, cache.ChannelAdmin, func(event cache.Event) {
|
|
if event.Type != cache.EventMarkerProviderConfigChanged {
|
|
return
|
|
}
|
|
if err := markerProviderConfig.Reload(appCtx); err != nil {
|
|
slog.Warn("reload marker provider config failed", "provider", event.Payload, "error", err)
|
|
}
|
|
}); err != nil {
|
|
slog.Warn("subscribe marker provider config reload failed", "error", err)
|
|
}
|
|
}
|
|
}
|
|
deps.MarkerProviderConfig = markerProviderConfig
|
|
deps.MarkerRegistry = markerRegistry
|
|
markerResolver := markers.NewDBExternalIDResolver(deps.DB)
|
|
deps.MarkerResolver = markerResolver
|
|
markerContributionStore := markers.NewContributionStore(deps.DB)
|
|
deps.MarkerContributionStore = markerContributionStore
|
|
deps.MarkerContributionService = markers.NewContributionService(
|
|
markerRegistry, markerResolver, markerProviderConfig, markerContributionStore, slog.Default(),
|
|
)
|
|
}
|
|
var watchProviderService *watchsync.Service
|
|
var watchProviderRegistry *watchsync.Registry
|
|
var watchProviderRepo *watchsync.PostgresRepository
|
|
if deps.DB != nil {
|
|
watchProviderRegistry = watchsync.NewRegistry()
|
|
if err := watchProviderRegistry.Register(trakt.NewProvider(nil, "")); err != nil {
|
|
log.Fatalf("register watch provider: %v", err)
|
|
}
|
|
if err := watchProviderRegistry.Register(simkl.NewProvider(nil, "")); err != nil {
|
|
log.Fatalf("register watch provider: %v", err)
|
|
}
|
|
if err := watchProviderRegistry.Register(watchmdblist.NewProvider(nil, "")); err != nil {
|
|
log.Fatalf("register watch provider: %v", err)
|
|
}
|
|
watchProviderRepo = watchsync.NewPostgresRepository(deps.DB, deps.SecretCipher)
|
|
watchProviderService = watchsync.NewService(watchProviderRepo, watchProviderRegistry)
|
|
deps.WatchProviderService = watchProviderService
|
|
}
|
|
|
|
// Initialize node pools for integrated/api modes.
|
|
if mode == "integrated" || mode == "api" {
|
|
nodeRepo := nodepool.NewRepository(pool)
|
|
deps.NodeRepo = nodeRepo
|
|
|
|
proxyPool := nodepool.NewProxyPool()
|
|
transcodePool := nodepool.NewTranscodePool()
|
|
|
|
proxyNodes, _ := nodeRepo.ListEnabled(context.Background(), nodepool.NodeTypeProxy)
|
|
transcodeNodes, _ := nodeRepo.ListEnabled(context.Background(), nodepool.NodeTypeTranscode)
|
|
proxyPool.SetNodes(proxyNodes)
|
|
transcodePool.SetNodes(transcodeNodes)
|
|
|
|
deps.ProxyPool = proxyPool
|
|
deps.TranscodePool = transcodePool
|
|
deps.NodePlanner = nodepool.NewPlanner(proxyPool, transcodePool)
|
|
|
|
healthChecker := nodepool.NewHealthChecker(proxyPool, transcodePool, nodeRepo)
|
|
healthChecker.Start(appCtx)
|
|
slog.Info("node pools initialized", "proxy_nodes", len(proxyNodes), "transcode_nodes", len(transcodeNodes))
|
|
|
|
// Subscribe to node pool change events for multi-instance reload.
|
|
_ = eventBus.Subscribe(appCtx, cache.ChannelAdmin, func(event cache.Event) {
|
|
if event.Type == cache.EventNodePoolChanged {
|
|
pNodes, pErr := nodeRepo.ListEnabled(context.Background(), nodepool.NodeTypeProxy)
|
|
tNodes, tErr := nodeRepo.ListEnabled(context.Background(), nodepool.NodeTypeTranscode)
|
|
if pErr != nil || tErr != nil {
|
|
slog.Warn("node pool reload from event failed, keeping current pools",
|
|
"proxy_err", pErr, "transcode_err", tErr)
|
|
return
|
|
}
|
|
proxyPool.SetNodes(pNodes)
|
|
transcodePool.SetNodes(tNodes)
|
|
slog.Info("node pools reloaded from event", "proxy", len(pNodes), "transcode", len(tNodes))
|
|
}
|
|
})
|
|
}
|
|
|
|
// Step 3: Create S3 clients (if needed).
|
|
if needsS3 {
|
|
configureS3Clients(cfg, &deps)
|
|
}
|
|
|
|
var literaryWorkService *literaryworks.Service
|
|
if deps.DB != nil {
|
|
literaryWorkService = literaryworks.NewService(literaryworks.NewRepository(deps.DB))
|
|
}
|
|
|
|
// Step 4: Create scanner (if needed).
|
|
if needsScanner && deps.DB != nil {
|
|
folderRepo := catalog.NewFolderRepository(deps.DB)
|
|
fileRepo := scanner.NewFileRepository(deps.DB)
|
|
deps.FolderRepo = folderRepo
|
|
deps.FileRepo = fileRepo
|
|
|
|
ffprobePath := scanner.FFprobePathFromFFmpeg(cfg.Playback.FFmpegPath)
|
|
s := scanner.NewScanner(fileRepo, ffprobePath, deps.S3Public, cfg.Scanner.Workers, cfg.Scanner.EmptyTrashAfterScan, cfg.Scanner.FileRemovalGrace)
|
|
s.SetSearchIndexProvider(activeCatalogSearchProvider)
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
s.SetWorkers(updated.Scanner.Workers)
|
|
})
|
|
s.SetLiteraryWorkLinker(literaryWorkService)
|
|
s.SetEbookEnrichmentQueue(ebooks.NewEnrichmentQueue(deps.DB))
|
|
deps.Scanner = s
|
|
deps.ProbeEnsurer = scanner.NewPlaybackProbeEnsurer(fileRepo, ffprobePath, cfg.Playback.FFmpegPath, 10*time.Second)
|
|
slog.Info("scanner initialized")
|
|
}
|
|
|
|
var chapterThumbService *chapterthumbs.Service
|
|
if deps.FileRepo != nil && deps.FolderRepo != nil && deps.S3Public != nil {
|
|
chapterThumbService = chapterthumbs.NewService(
|
|
deps.FileRepo,
|
|
deps.FolderRepo,
|
|
deps.ProbeEnsurer,
|
|
settingsRepo,
|
|
deps.S3Public,
|
|
nil,
|
|
deps.TranscodePool,
|
|
cfg.Playback.FFmpegPath,
|
|
cfg.Playback.HWAccel,
|
|
cfg.Playback.HWDevice,
|
|
cfg.Playback.ChapterThumbnailWorkers,
|
|
)
|
|
if chapterThumbService != nil {
|
|
chapterThumbService.Start(appCtx)
|
|
deps.ChapterThumbnailQueuer = chapterThumbService
|
|
}
|
|
}
|
|
|
|
var pluginHost *pluginhost.Host
|
|
var pluginService *plugins.Service
|
|
var pluginInstallationStore *plugins.InstallationStore
|
|
var pluginRuntimeConfigStore *plugins.RuntimeConfigStore
|
|
var pluginHTTPProxy *plugins.HTTPProxy
|
|
pluginAutoUpdateDone := make(chan struct{})
|
|
var pluginAutoUpdater *plugins.AutoUpdateService
|
|
if deps.DB != nil {
|
|
pluginCacheDir := resolvePluginCacheDir()
|
|
repositoryStore := plugins.NewRepositoryStore(deps.DB)
|
|
installationStore := plugins.NewInstallationStore(deps.DB)
|
|
runtimeConfigStore := plugins.NewRuntimeConfigStore(deps.DB, deps.SecretCipher)
|
|
catalogService := plugins.NewCatalogService(repositoryStore, plugins.CatalogServiceOptions{
|
|
SiloAPIVersion: plugins.DefaultSiloAPIVersion,
|
|
})
|
|
installer := plugins.NewInstaller(installationStore, plugins.InstallerOptions{
|
|
BaseDir: pluginCacheDir,
|
|
})
|
|
|
|
libDataSource := pluginhost.LibraryDataSourceFunc(
|
|
func(ctx context.Context, _ string) ([]pluginhost.LibraryRecord, error) {
|
|
// TODO: scope by userID when the requests plugin needs it (Plan B).
|
|
// For now, all callers see admin-scope.
|
|
if deps.FolderRepo == nil {
|
|
return nil, nil
|
|
}
|
|
folders, err := deps.FolderRepo.List(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out := make([]pluginhost.LibraryRecord, 0, len(folders))
|
|
for _, f := range folders {
|
|
out = append(out, pluginhost.LibraryRecord{
|
|
ID: strconv.Itoa(f.ID),
|
|
Name: f.Name,
|
|
MediaType: mapFolderTypeToMediaType(f.Type),
|
|
})
|
|
}
|
|
return out, nil
|
|
},
|
|
)
|
|
presenceItemRepo := catalog.NewItemRepository(deps.DB)
|
|
catalogPresence := pluginhost.NewCatalogPresence(
|
|
func(ctx context.Context, mediaType string, tmdbIDs []string) ([]pluginhost.LibraryPresenceRecord, error) {
|
|
rows, err := presenceItemRepo.LookupTMDBIDs(ctx, mediaType, tmdbIDs)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out := make([]pluginhost.LibraryPresenceRecord, 0, len(rows))
|
|
for _, r := range rows {
|
|
out = append(out, pluginhost.LibraryPresenceRecord{
|
|
ExternalID: r.TMDBID,
|
|
MediaID: r.MediaID,
|
|
LibraryID: r.LibraryID,
|
|
Title: r.Title,
|
|
})
|
|
}
|
|
return out, nil
|
|
},
|
|
)
|
|
pluginHost = pluginhost.NewHost(pluginhost.Config{
|
|
EventPublisher: eventsHub,
|
|
LibraryLister: pluginhost.NewLibraryLister(libDataSource),
|
|
CatalogPresence: catalogPresence,
|
|
InstalledPlugins: pluginhost.InstalledPluginListerFunc(
|
|
func(ctx context.Context) ([]pluginhost.InstalledPluginRecord, error) {
|
|
installations, err := installationStore.List(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out := make([]pluginhost.InstalledPluginRecord, 0, len(installations))
|
|
for _, installation := range installations {
|
|
// The reserved builtin row is not a plugin; keep it out
|
|
// of the host's installed-plugin listing.
|
|
if installation.IsBuiltin() {
|
|
continue
|
|
}
|
|
capabilities, err := installationStore.ListCapabilities(ctx, installation.ID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
descriptors := make([]*pluginv1.CapabilityDescriptor, 0, len(capabilities))
|
|
for _, capability := range capabilities {
|
|
descriptor, err := plugins.DecodeCapability(capability)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
descriptors = append(descriptors, descriptor)
|
|
}
|
|
out = append(out, pluginhost.InstalledPluginRecord{
|
|
InstallationID: installation.ID,
|
|
PluginID: installation.PluginID,
|
|
Version: installation.Version,
|
|
Enabled: installation.Enabled,
|
|
Capabilities: descriptors,
|
|
})
|
|
}
|
|
return out, nil
|
|
},
|
|
),
|
|
GlobalConfigSetter: pluginhost.GlobalConfigSetterFunc(
|
|
func(ctx context.Context, installationID int, key string, value map[string]any) error {
|
|
return runtimeConfigStore.PutGlobalConfig(ctx, installationID, key, value)
|
|
},
|
|
),
|
|
Logger: hclog.New(&hclog.LoggerOptions{
|
|
Name: "plugin-host",
|
|
Level: hclog.Info,
|
|
Output: os.Stderr,
|
|
}),
|
|
})
|
|
pluginService = plugins.NewService(
|
|
repositoryStore,
|
|
installationStore,
|
|
runtimeConfigStore,
|
|
catalogService,
|
|
installer,
|
|
plugins.NewHostAdapter(pluginHost),
|
|
)
|
|
if watchProviderRegistry != nil {
|
|
reloadWatchProviders := func(ctx context.Context) {
|
|
if err := reloadWatchSyncPluginProviders(ctx, watchProviderRegistry, installationStore, pluginService, watchProviderRepo); err != nil {
|
|
slog.WarnContext(ctx, "failed to reload watch sync plugin providers", "component", "app", "error", err)
|
|
}
|
|
}
|
|
pluginService.AddLifecycleHook(reloadWatchProviders)
|
|
reloadWatchProviders(appCtx)
|
|
}
|
|
if deps.MarkerRegistry != nil && deps.MarkerProviderConfig != nil {
|
|
markerPluginResolver := markers.NewPluginResolverAdapter(pluginService)
|
|
pluginService.AddLifecycleHook(func(ctx context.Context) {
|
|
if err := reloadMarkerPluginProviders(
|
|
ctx,
|
|
deps.MarkerRegistry,
|
|
deps.MarkerProviderConfig,
|
|
installationStore,
|
|
runtimeConfigStore,
|
|
settingsRepo,
|
|
markerPluginResolver,
|
|
); err != nil {
|
|
slog.WarnContext(ctx, "reload marker plugin providers failed", "component", "app", "error", err)
|
|
}
|
|
})
|
|
}
|
|
if err := pluginService.PreloadEnabled(appCtx); err != nil {
|
|
log.Fatalf("preload enabled plugins: %v", err)
|
|
}
|
|
slog.Info("plugin cache initialized", "base_dir", pluginCacheDir)
|
|
|
|
pluginAutoUpdater = plugins.NewAutoUpdateService(
|
|
repositoryStore,
|
|
installationStore,
|
|
catalogService,
|
|
installer,
|
|
pluginHost,
|
|
slog.Default(),
|
|
// Auto-updates rewrite installation rows (new version-specific
|
|
// InstallPath/Version) and delete the old install dir without going
|
|
// through pluginService. Wire OnLifecycleChange so the service's
|
|
// installation cache is invalidated and later plugin RPCs re-read
|
|
// the fresh row instead of a stale one.
|
|
pluginService.OnLifecycleChange,
|
|
)
|
|
go func() {
|
|
defer close(pluginAutoUpdateDone)
|
|
if err := pluginAutoUpdater.Run(appCtx); err != nil {
|
|
slog.Error("plugin auto-update failed", "error", err)
|
|
}
|
|
}()
|
|
pluginInstallationStore = installationStore
|
|
pluginRuntimeConfigStore = runtimeConfigStore
|
|
pluginHTTPProxy = plugins.NewHTTPProxyWithTypedResolver(pluginService, pluginInstallationStore)
|
|
if deps.DB != nil {
|
|
pluginHTTPProxy = pluginHTTPProxy.WithUserThemeLookup(plugins.NewPgUserThemeLookup(deps.DB))
|
|
pluginHTTPProxy = pluginHTTPProxy.WithUserIdentityLookup(plugins.NewPgUserIdentityLookup(deps.DB))
|
|
}
|
|
deps.PluginService = pluginService
|
|
deps.PluginHTTPProxy = pluginHTTPProxy
|
|
defer func() {
|
|
if pluginHost == nil {
|
|
return
|
|
}
|
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
if err := pluginHost.Shutdown(shutdownCtx); err != nil {
|
|
slog.Warn("failed to shut down plugin host", "error", err)
|
|
}
|
|
}()
|
|
} else {
|
|
close(pluginAutoUpdateDone)
|
|
}
|
|
if pluginService != nil && pluginInstallationStore != nil {
|
|
dispatcher := plugins.NewEventDispatcherWithTypedResolver(deps.EventBus, deps.EventsHub, pluginInstallationStore, pluginService, 4)
|
|
pluginService.SetEventDispatcher(dispatcher)
|
|
if err := dispatcher.Start(appCtx); err != nil {
|
|
log.Fatalf("plugin event dispatcher: %v", err)
|
|
}
|
|
defer dispatcher.Stop()
|
|
// Backfill the capability-subscriber index from the already-preloaded
|
|
// installations. PreloadEnabled ran earlier (before the dispatcher
|
|
// existed), so its rebuildDispatcherIndex was a no-op. Without this
|
|
// call, capability-scoped subscriptions never fire until the next
|
|
// lifecycle mutation.
|
|
pluginService.OnLifecycleChange(appCtx)
|
|
}
|
|
|
|
// backgroundInit collects non-critical startup work (catalog-size-dependent
|
|
// seeding, network-bound reconciliation) that must not block the HTTP
|
|
// listener. The steps run sequentially in a background goroutine once the
|
|
// server is ready to serve. Failures are logged, never fatal.
|
|
var backgroundInit []func(context.Context)
|
|
|
|
// Step 4b: Create metadata service and match worker (if needed).
|
|
var metadataService *metadata.MetadataService
|
|
var metadataImageCacheProcessor *metadata.ImageCacheProcessor
|
|
var personRefreshService *metadata.PersonRefreshService
|
|
var matchWorker *metadata.MatchWorker
|
|
var libraryIngestExecutor *libraryingest.Executor
|
|
var libraryScanQueue *scanqueue.Service
|
|
var itemRefreshExecutor *adminjob.ItemRefreshExecutor
|
|
var libraryRefreshExecutor *adminjob.LibraryRefreshExecutor
|
|
var itemRepo *catalog.ItemRepository
|
|
var skippedRootRepo *metadata.SkippedRootRepository
|
|
var movieQueueRepo *metadata.MovieMatchQueueRepository
|
|
var seriesQueueRepo *metadata.SeriesRootMatchQueueRepository
|
|
var matchQueueCoordinator *metadata.MatchQueueCoordinator
|
|
var rootClaimRepo *catalog.RootClaimRepository
|
|
var groupClaimRepo *catalog.GroupClaimRepository
|
|
var seasonRepo *catalog.SeasonRepository
|
|
var episodeRepo *catalog.EpisodeRepository
|
|
var audiobookEnricher *audiobooks.Enricher
|
|
var ebookEnricher *ebooks.Enricher
|
|
var mangaEnricher *manga.Enricher
|
|
if needsWorkers && deps.DB != nil && deps.FileRepo != nil {
|
|
chainRepo := metadata.NewChainRepository(deps.DB)
|
|
// Make every existing library chain aware of the built-in providers
|
|
// before serving: materialize legacy content_level='' chains per level,
|
|
// then append registered builtins disabled (idempotent; also the repair
|
|
// path after a stale chain-editor save drops a builtin row). Runs before
|
|
// the metadata service exists, so no chain cache to invalidate here.
|
|
syncCtx, syncCancel := context.WithTimeout(appCtx, 30*time.Second)
|
|
syncErr := metadata.SyncBuiltinProviderChains(syncCtx, chainRepo)
|
|
syncCancel()
|
|
if syncErr != nil {
|
|
log.Fatalf("sync builtin provider chains: %v", syncErr)
|
|
}
|
|
skippedRootRepo = metadata.NewSkippedRootRepository(deps.DB)
|
|
itemRepo = catalog.NewItemRepository(deps.DB).WithActiveSearchProvider(activeCatalogSearchProvider)
|
|
episodeRepo = catalog.NewEpisodeRepository(deps.DB)
|
|
seasonRepo = catalog.NewSeasonRepository(deps.DB)
|
|
personRepo := catalog.NewPersonRepository(deps.DB)
|
|
libraryRepo := catalog.NewLibraryItemRepository(deps.DB)
|
|
|
|
// Wait for plugin auto-update to finish before registering image resolvers.
|
|
<-pluginAutoUpdateDone
|
|
|
|
imageResolver := metadata.NewPluginImageResolver()
|
|
if pluginService != nil && pluginInstallationStore != nil {
|
|
reloadImageResolvers := func(ctx context.Context) {
|
|
if err := reloadPluginImageResolvers(ctx, pluginInstallationStore, imageResolver, pluginService); err != nil {
|
|
slog.WarnContext(ctx, "failed to reload plugin image resolvers", "component", "app", "error", err)
|
|
}
|
|
}
|
|
pluginService.AddLifecycleHook(reloadImageResolvers)
|
|
reloadImageResolvers(appCtx)
|
|
}
|
|
if deps.S3Public != nil {
|
|
presignTTL := cfg.S3.MetadataPresignExpiry
|
|
if presignTTL <= 0 {
|
|
presignTTL = 4 * time.Hour
|
|
}
|
|
imageResolver.SetS3Presigner(deps.S3Public, deps.S3Public.EffectivePresignTTL(presignTTL))
|
|
}
|
|
deps.ImageResolver = imageResolver
|
|
deps.PluginImageResolver = imageResolver
|
|
|
|
staleIDRepo := metadata.NewStaleMediaIDRepository(deps.DB)
|
|
providerIDRepo := catalog.NewProviderIDRepository(deps.DB)
|
|
movieQueueRepo = metadata.NewMovieMatchQueueRepository(deps.DB, deps.FileRepo)
|
|
seriesQueueRepo = metadata.NewSeriesRootMatchQueueRepository(deps.DB)
|
|
deps.MovieMatchQueueRepo = movieQueueRepo
|
|
deps.SeriesRootMatchQueueRepo = seriesQueueRepo
|
|
matchQueueCoordinator = metadata.NewMatchQueueCoordinator(movieQueueRepo, seriesQueueRepo)
|
|
backgroundInit = append(backgroundInit, func(ctx context.Context) {
|
|
if err := matchQueueCoordinator.WakeForChangedInputs(ctx); err != nil {
|
|
slog.WarnContext(ctx, "refresh metadata match queue inputs at startup failed", "component", "app", "error", err)
|
|
}
|
|
})
|
|
if pluginService != nil {
|
|
matchInputChanged := make(chan struct{}, 1)
|
|
go func() {
|
|
for {
|
|
select {
|
|
case <-appCtx.Done():
|
|
return
|
|
case <-matchInputChanged:
|
|
if err := matchQueueCoordinator.WakeForChangedInputs(appCtx); err != nil {
|
|
slog.WarnContext(appCtx, "wake metadata matches after plugin lifecycle change failed", "component", "app", "error", err)
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
pluginService.AddLifecycleHook(func(context.Context) {
|
|
// Queue fingerprint reconciliation may touch thousands of parked
|
|
// rows. Coalesce lifecycle bursts and keep plugin admin requests
|
|
// independent of that background database work.
|
|
select {
|
|
case matchInputChanged <- struct{}{}:
|
|
default:
|
|
}
|
|
})
|
|
}
|
|
rootClaimRepo = catalog.NewRootClaimRepository(deps.DB)
|
|
groupClaimRepo = catalog.NewGroupClaimRepository(deps.DB)
|
|
pluginResolver := metadata.NewPluginResolverAdapter(pluginService)
|
|
// Serve the metadata chain's plugin-installation enabled-check from the
|
|
// plugins service's in-memory installation cache. Declared as the
|
|
// interface type and only assigned when pluginService is non-nil so a
|
|
// nil *plugins.Service is passed as a genuine nil interface (not a
|
|
// typed-nil), letting buildProviders fall back to the pool query.
|
|
var installationEnabledChecker metadata.InstallationEnabledChecker
|
|
if pluginService != nil {
|
|
installationEnabledChecker = pluginService
|
|
}
|
|
metadataService = metadata.NewMetadataService(
|
|
chainRepo, pluginResolver, installationEnabledChecker,
|
|
itemRepo, providerIDRepo, episodeRepo, seasonRepo, libraryRepo, deps.FolderRepo,
|
|
personRepo,
|
|
deps.FileRepo, skippedRootRepo, staleIDRepo, rootClaimRepo,
|
|
)
|
|
// Drop the resolved-chain cache whenever a plugin is installed, enabled,
|
|
// disabled, updated, or uninstalled. The installation-enabled check is
|
|
// served from the plugins service's in-memory cache (invalidated on the
|
|
// same events), but resolveChainCached would otherwise keep serving a
|
|
// stale provider chain for up to chainCacheTTL after a provider's
|
|
// availability changes.
|
|
if pluginService != nil {
|
|
pluginService.AddLifecycleHook(func(context.Context) {
|
|
metadataService.InvalidateChainCache()
|
|
})
|
|
}
|
|
personRefreshService = metadata.NewPersonRefreshService(deps.DB, pluginResolver, personRepo)
|
|
personRefreshService.SetImageResolver(imageResolver)
|
|
|
|
// Wire the audiobook enricher. It uses the same plugin resolver and chain
|
|
// repo as the movie/TV pipeline, but resolves providers at
|
|
// content_level='audiobook' and sweeps items directly rather than via a queue.
|
|
audiobookEnricher = audiobooks.NewEnricher(
|
|
deps.DB,
|
|
chainRepo,
|
|
pluginResolver,
|
|
itemRepo,
|
|
personRepo,
|
|
providerIDRepo,
|
|
)
|
|
ebookEnricher = ebooks.NewEnricher(
|
|
deps.DB,
|
|
chainRepo,
|
|
pluginResolver,
|
|
itemRepo,
|
|
personRepo,
|
|
providerIDRepo,
|
|
)
|
|
audiobookEnricher.SetLiteraryWorkLinker(literaryWorkService)
|
|
ebookEnricher.SetLiteraryWorkLinker(literaryWorkService)
|
|
mangaEnricher = manga.NewEnricher(
|
|
deps.DB,
|
|
chainRepo,
|
|
pluginResolver,
|
|
itemRepo,
|
|
personRepo,
|
|
providerIDRepo,
|
|
)
|
|
|
|
// Always wire the image resolver so plugin-prefixed URLs (e.g.
|
|
// metadb://) can be resolved to presigned HTTP URLs in API responses.
|
|
metadataService.SetImageResolver(imageResolver)
|
|
|
|
// Wire the image cacher whenever object storage is available so explicit
|
|
// admin image applies can succeed even if automatic metadata caching is off.
|
|
if deps.S3Public != nil {
|
|
imageCacher := imagecache.New(deps.S3Public)
|
|
imageCacher.SetArtworkRevisionTracker(catalog.NewArtworkRevisionTracker(deps.DB))
|
|
metadataService.SetImageCacher(imageCacher)
|
|
imageCacheJobs := metadata.NewImageCacheJobRepository(deps.DB)
|
|
metadataService.SetImageCacheJobEnqueuer(imageCacheJobs)
|
|
metadataImageCacheProcessor = metadata.NewImageCacheProcessorWithTargets(
|
|
imageCacheJobs,
|
|
imageCacher,
|
|
imageResolver,
|
|
metadata.ImageCacheProcessorTargets{
|
|
Items: itemRepo,
|
|
Seasons: seasonRepo,
|
|
Episodes: episodeRepo,
|
|
ItemLocalizations: catalog.NewMediaItemLocalizationRepository(deps.DB),
|
|
SeasonLocalizations: catalog.NewSeasonLocalizationRepository(deps.DB),
|
|
People: personRepo,
|
|
},
|
|
)
|
|
// Local file:// artwork (NFO sidecars): confine reads to the owning
|
|
// library's roots and sweep stale hashed local/ prefixes on re-cache.
|
|
// The processor host must mount the libraries, like the metadata worker.
|
|
metadataImageCacheProcessor.SetLibraryRootResolver(deps.FolderRepo)
|
|
metadataImageCacheProcessor.SetImagePrefixDeleter(deps.S3Public)
|
|
metadataService.SetAutoCacheImages(cfg.Metadata.CacheImages)
|
|
metadataImageCacheProcessor.SetEnabled(cfg.Metadata.CacheImages)
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
metadataService.SetAutoCacheImages(updated.Metadata.CacheImages)
|
|
metadataImageCacheProcessor.SetEnabled(updated.Metadata.CacheImages)
|
|
})
|
|
if deps.Scanner != nil {
|
|
deps.Scanner.SetImageCacher(imageCacher)
|
|
}
|
|
if cfg.Metadata.CacheImages {
|
|
personRefreshService.SetImageCacher(imageCacher)
|
|
personRefreshService.SetImageCacheJobEnqueuer(imageCacheJobs)
|
|
slog.Info("metadata image caching enabled")
|
|
}
|
|
if audiobookEnricher != nil {
|
|
audiobookEnricher.SetImageCacher(imageCacher)
|
|
audiobookEnricher.SetImageCacheJobEnqueuer(imageCacheJobs)
|
|
audiobookEnricher.SetFFmpegPath(scanner.FFmpegPathFromFFprobe(scanner.FFprobePathFromFFmpeg(cfg.Playback.FFmpegPath)))
|
|
}
|
|
if ebookEnricher != nil {
|
|
ebookEnricher.SetImageCacher(imageCacher)
|
|
ebookEnricher.SetImageCacheJobEnqueuer(imageCacheJobs)
|
|
}
|
|
if mangaEnricher != nil {
|
|
mangaEnricher.SetImageCacher(imageCacher)
|
|
mangaEnricher.SetImageCacheJobEnqueuer(imageCacheJobs)
|
|
}
|
|
}
|
|
|
|
matchWorker = metadata.NewMatchWorker(metadataService, deps.FileRepo, cfg.Matcher.Workers, cfg.Matcher.BatchSize, 30*time.Second)
|
|
mwForReload := matchWorker
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
mwForReload.SetConcurrency(updated.Matcher.Workers, updated.Matcher.BatchSize)
|
|
})
|
|
matchWorker.SetRealtimeHub(deps.RealtimeHub)
|
|
if movieQueueRepo != nil {
|
|
matchWorker.SetMovieFileClaimer(movieQueueRepo)
|
|
}
|
|
if seriesQueueRepo != nil {
|
|
matchWorker.SetSeriesRootClaimer(seriesQueueRepo, cfg.Matcher.TVSeriesRootQueueEnabled())
|
|
backgroundInit = append(backgroundInit, func(ctx context.Context) {
|
|
if cleaned, err := seriesQueueRepo.CleanupLegacySeriesGroupQueue(ctx); err != nil {
|
|
slog.WarnContext(ctx, "failed to clean legacy series group queue rows", "component", "app", "error", err)
|
|
} else if cleaned > 0 {
|
|
slog.InfoContext(ctx, "cleaned legacy series group queue rows", "component", "app", "count", cleaned)
|
|
}
|
|
})
|
|
}
|
|
if deps.FolderRepo != nil {
|
|
backgroundInit = append(backgroundInit, func(ctx context.Context) {
|
|
start := time.Now()
|
|
enabledFolders, err := deps.FolderRepo.GetEnabled(ctx)
|
|
if err != nil {
|
|
slog.WarnContext(ctx, "failed to seed metadata queues", "component", "app", "error", err)
|
|
return
|
|
}
|
|
seedMovieQueue := func(folderID int) {
|
|
if movieQueueRepo == nil {
|
|
return
|
|
}
|
|
if err := movieQueueRepo.SyncForFolder(ctx, folderID); err != nil {
|
|
slog.WarnContext(ctx, "failed to seed movie match queue", "component", "app", "folder_id", folderID, "error", err)
|
|
}
|
|
}
|
|
seedSeriesQueue := func(folderID int) {
|
|
if seriesQueueRepo == nil {
|
|
return
|
|
}
|
|
if err := seriesQueueRepo.SyncForFolder(ctx, folderID); err != nil {
|
|
slog.WarnContext(ctx, "failed to seed series root queue", "component", "app", "folder_id", folderID, "error", err)
|
|
}
|
|
}
|
|
for _, folder := range enabledFolders {
|
|
if folder == nil {
|
|
continue
|
|
}
|
|
switch strings.ToLower(strings.TrimSpace(folder.Type)) {
|
|
case "movie", "movies":
|
|
seedMovieQueue(folder.ID)
|
|
case "series", "tv", "show", "tvshows":
|
|
seedSeriesQueue(folder.ID)
|
|
case "mixed":
|
|
seedSeriesQueue(folder.ID)
|
|
seedMovieQueue(folder.ID)
|
|
}
|
|
}
|
|
slog.InfoContext(ctx, "deferred init: metadata match queues seeded", "component", "app", "folders", len(enabledFolders), "duration", time.Since(start))
|
|
})
|
|
}
|
|
|
|
deps.SkippedRootRepo = skippedRootRepo
|
|
deps.StaleIDRepo = staleIDRepo
|
|
deps.PersonRepo = personRepo
|
|
deps.PersonRefreshQueue = worker.NewPersonRefreshWorker(
|
|
personRefreshService,
|
|
worker.DefaultPersonRefreshWorkerConfig(),
|
|
)
|
|
deps.PersonRefresher = personRefreshService
|
|
deps.Refresher = metadataService
|
|
deps.MetadataService = metadataService
|
|
slog.Info("metadata service initialized and running")
|
|
|
|
}
|
|
if deps.Scanner != nil {
|
|
if matchQueueCoordinator != nil {
|
|
deps.Scanner.SetMetadataQueueProducer(matchQueueCoordinator)
|
|
}
|
|
if movieQueueRepo != nil {
|
|
deps.Scanner.SetMovieQueueSyncer(movieQueueRepo)
|
|
}
|
|
if seriesQueueRepo != nil {
|
|
deps.Scanner.SetSeriesQueueSyncer(seriesQueueRepo)
|
|
}
|
|
}
|
|
if deps.Scanner != nil && matchWorker != nil && deps.FolderRepo != nil && skippedRootRepo != nil {
|
|
libraryIngestExecutor = libraryingest.NewExecutor(
|
|
deps.Scanner,
|
|
matchWorker,
|
|
deps.FolderRepo,
|
|
skippedRootRepo,
|
|
deps.EventBus,
|
|
deps.RealtimeHub,
|
|
)
|
|
deps.LibraryIngester = libraryIngestExecutor
|
|
if deps.DB != nil {
|
|
libraryScanQueue = scanqueue.NewService(
|
|
scanqueue.NewRepository(deps.DB),
|
|
deps.FolderRepo,
|
|
libraryIngestExecutor,
|
|
deps.EventsHub,
|
|
appCtx,
|
|
cfg.Scanner.MaxConcurrentLibraries,
|
|
cfg.Scanner.MaxConcurrentScoped,
|
|
)
|
|
// Started below, after the notification system has attached its
|
|
// availability detector to the executor: a scan resumed by the
|
|
// workers before that wiring would complete without recording
|
|
// episode availability, silently losing release notifications.
|
|
deps.LibraryScanQueue = libraryScanQueue
|
|
}
|
|
if deps.DB != nil && deps.FileRepo != nil && metadataService != nil {
|
|
itemRefreshResolver := adminjob.NewItemRefreshResolver(
|
|
itemRepo,
|
|
seasonRepo,
|
|
episodeRepo,
|
|
deps.FolderRepo,
|
|
deps.FileRepo,
|
|
)
|
|
libraryRefreshExecutor = adminjob.NewLibraryRefreshExecutor(
|
|
adminjob.NewPGLibraryRefreshItemLister(deps.DB),
|
|
deps.FolderRepo,
|
|
itemRefreshResolver,
|
|
libraryIngestExecutor,
|
|
metadataService,
|
|
deps.EventBus,
|
|
deps.RealtimeHub,
|
|
)
|
|
}
|
|
if metadataService != nil && deps.FileRepo != nil {
|
|
itemRefreshExecutor = adminjob.NewItemRefreshExecutor(
|
|
deps.FolderRepo,
|
|
deps.FileRepo,
|
|
rootClaimRepo,
|
|
groupClaimRepo,
|
|
skippedRootRepo,
|
|
seasonRepo,
|
|
episodeRepo,
|
|
libraryIngestExecutor,
|
|
metadataService,
|
|
deps.EventBus,
|
|
deps.RealtimeHub,
|
|
)
|
|
}
|
|
}
|
|
|
|
// Ensure PersonRepo is available for the router's DetailService.
|
|
if deps.DB != nil && deps.PersonRepo == nil {
|
|
deps.PersonRepo = catalog.NewPersonRepository(deps.DB)
|
|
}
|
|
|
|
// Step 5: Create user store provider (if needed).
|
|
var userStoreProvider userstore.UserStoreProvider
|
|
if needsUserDB {
|
|
switch cfg.UserDB.Backend {
|
|
case "sqlite":
|
|
poolConfig := userdb.PoolConfig{
|
|
MaxOpen: cfg.UserDB.PoolMaxOpen,
|
|
IdleTimeout: cfg.UserDB.IdleTimeout,
|
|
DataDir: "/var/lib/silo/userdb",
|
|
}
|
|
pool := userdb.NewUserDBPool(poolConfig)
|
|
userStoreProvider = userdb.NewSQLiteProvider(pool)
|
|
slog.Info("user store initialized", "backend", "sqlite", "max_open", poolConfig.MaxOpen)
|
|
default: // "postgres"
|
|
userStoreProvider = pgstore.NewPostgresProvider(deps.DB)
|
|
slog.Info("user store initialized", "backend", "postgres")
|
|
}
|
|
defer userStoreProvider.Close()
|
|
}
|
|
|
|
var policySystem *policy.System
|
|
if mode == "integrated" || mode == "api" {
|
|
policyDecisionLogger := policy.NewDecisionLogger(
|
|
deps.DB,
|
|
nodeID,
|
|
policy.WithDecisionLogLogger(slog.Default()),
|
|
)
|
|
policyDecisionLogger.SetVerbosity(cfg.Policy.DecisionLogVerbosity)
|
|
policyDecisionLogger.SetScopeSampleRate(cfg.Policy.DecisionLogScopeSampleRate)
|
|
policySystem = policy.NewSystem(
|
|
policy.NewPolicyStore(deps.DB),
|
|
deps.EventBus,
|
|
slog.Default(),
|
|
policy.WithSystemEvalTimeout(time.Duration(cfg.Policy.EvalTimeoutMS)*time.Millisecond),
|
|
policy.WithSystemDecisionLogger(policyDecisionLogger),
|
|
)
|
|
if err := policySystem.Start(appCtx); err != nil {
|
|
log.Fatalf("policy system start: %v", err)
|
|
}
|
|
deps.PolicySystem = policySystem
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
policySystem.SetEvalTimeout(time.Duration(updated.Policy.EvalTimeoutMS) * time.Millisecond)
|
|
if logger := policySystem.DecisionLogger(); logger != nil {
|
|
logger.SetVerbosity(updated.Policy.DecisionLogVerbosity)
|
|
logger.SetScopeSampleRate(updated.Policy.DecisionLogScopeSampleRate)
|
|
}
|
|
})
|
|
defer policySystem.Stop()
|
|
}
|
|
|
|
// User-facing release notifications. The system reads user state through
|
|
// the raw store provider; the provider handed to everything downstream is
|
|
// wrapped so every favorites/watchlist/progress mutation (REST handlers,
|
|
// jellycompat, imports, playback) feeds the interest index.
|
|
var notificationSystem *notifications.System
|
|
if deps.DB != nil && userStoreProvider != nil {
|
|
userRepo := auth.NewUserRepository(deps.DB)
|
|
profileTokens := access.NewProfileTokenService(cfg.Auth.JWTSecret, 0)
|
|
var notificationScopes notifications.ScopeResolver
|
|
if policySystem != nil {
|
|
notificationScopes = policy.NewViewerResolver(userRepo, userStoreProvider, profileTokens, policySystem.PDP(), accessGroupStore)
|
|
} else {
|
|
// Legacy resolver: proxy/test wiring without a policy system. Production integrated/api modes always take the policy path. Removed with the legacy cleanup phase.
|
|
notificationScopes = access.NewResolver(userRepo, userStoreProvider, profileTokens, accessGroupStore)
|
|
}
|
|
notificationSystem = notifications.NewSystem(
|
|
deps.DB,
|
|
settingsRepo,
|
|
userStoreProvider,
|
|
notificationScopes,
|
|
userRepo,
|
|
deps.EventsHub,
|
|
deps.RedisClient,
|
|
deps.SecretCipher,
|
|
mail.NewSMTPSender(settingsRepo),
|
|
)
|
|
userStoreProvider = notifications.WrapUserStoreProvider(userStoreProvider, notificationSystem)
|
|
deps.Notifications = notificationSystem
|
|
|
|
if libraryIngestExecutor != nil {
|
|
libraryIngestExecutor.SetAvailabilityDetector(notificationSystem.Detector)
|
|
}
|
|
if needsWorkers {
|
|
notificationSystem.Start(appCtx)
|
|
defer notificationSystem.Wait()
|
|
}
|
|
}
|
|
|
|
// Start the scan queue only now that the availability detector (when
|
|
// notifications are enabled) is attached to the ingest executor, so scans
|
|
// resumed at startup cannot complete before the detector exists.
|
|
if libraryScanQueue != nil {
|
|
libraryScanQueue.Start()
|
|
defer libraryScanQueue.Stop()
|
|
}
|
|
|
|
if userStoreProvider != nil && pluginService != nil {
|
|
deps.PluginUserConfig = plugins.NewUserConfigStore(userStoreProvider, pluginService)
|
|
}
|
|
|
|
// Step 6: Create playback session manager and wire into dependencies.
|
|
sessionMgr := playback.NewSessionManager(6, 2) // defaults from plan: max_streams=6, max_transcodes=2
|
|
var compatTerminalRecoveryReady <-chan struct{}
|
|
if userStoreProvider != nil {
|
|
deps.UserStoreProvider = userStoreProvider
|
|
}
|
|
if watchProviderService != nil {
|
|
historyRepo := historyimport.NewRepository(deps.DB, deps.SecretCipher)
|
|
historyIdentity := watchstate.NewStableIdentityResolver(itemRepo, episodeRepo, catalog.NewProviderIDRepository(deps.DB))
|
|
watchProviderService.
|
|
WithMatcher(historyimport.NewMatcher(historyRepo)).
|
|
WithWatchState(watchstate.NewService(userStoreProvider).WithStableIdentityResolver(historyIdentity)).
|
|
WithUserStoreProvider(userStoreProvider)
|
|
backgroundInit = append(backgroundInit, func(ctx context.Context) {
|
|
if compatTerminalRecoveryReady != nil {
|
|
select {
|
|
case <-compatTerminalRecoveryReady:
|
|
case <-ctx.Done():
|
|
return
|
|
}
|
|
}
|
|
if err := watchProviderService.SweepOpenScrobbles(ctx); err != nil {
|
|
slog.WarnContext(ctx, "failed to sweep open watch provider scrobbles", "component", "app", "error", err)
|
|
}
|
|
})
|
|
}
|
|
// Auto-remove fully-watched movies from the watchlist (standalone behavior,
|
|
// default-on per profile), propagating removals to connected providers.
|
|
// Series are never removed; watchlist read paths hide fully-watched ones
|
|
// (catalog.WatchlistVisibility) so newly added episodes bring them back.
|
|
if itemRepo != nil && userStoreProvider != nil {
|
|
maintainer := watchlist.NewMaintainer(userStoreProvider, itemRepo)
|
|
if watchProviderService != nil {
|
|
maintainer.WithListEventDispatcher(watchProviderService)
|
|
}
|
|
deps.WatchCompletionObserver = maintainer
|
|
}
|
|
deps.SessionMgr = sessionMgr
|
|
deps.PlaybackRealtimeHub = playback.NewRealtimeHub()
|
|
if chapterThumbService != nil && deps.S3Public != nil {
|
|
chapterThumbService.SetNotifier(
|
|
playback.NewChapterThumbnailNotifier(sessionMgr, deps.PlaybackRealtimeHub, deps.S3Public, 0),
|
|
)
|
|
}
|
|
|
|
// Build the reconciler early enough that playback handlers can trigger
|
|
// immediate session syncs after start/stop events.
|
|
nodeIdentity := resolveNodeIdentity()
|
|
|
|
var reconciler *worker.Reconciler
|
|
var heartbeatWriter *worker.HeartbeatWriter
|
|
if needsWorkers && deps.DB != nil {
|
|
sessionProvider := func() []worker.SessionSync {
|
|
sessions := sessionMgr.AllSessions()
|
|
syncs := make([]worker.SessionSync, len(sessions))
|
|
for i, s := range sessions {
|
|
syncs[i] = buildLiveSessionSync(s, nodeIdentity)
|
|
}
|
|
return syncs
|
|
}
|
|
reconciler = worker.NewReconciler(deps.DB, nodeIdentity, sessionProvider)
|
|
reconciler.EventBus = deps.EventBus
|
|
reconciler.EventsHub = deps.EventsHub
|
|
reconciler.PreSync = func() {
|
|
// Retire sessions that have not shown real playback activity
|
|
// recently enough to count as live. This keeps the in-memory
|
|
// limiter, transcode teardown, and synced admin view aligned.
|
|
if expired := sessionMgr.CleanStale(); len(expired) > 0 {
|
|
slog.Info("expired idle sessions", "count", len(expired))
|
|
}
|
|
}
|
|
deps.SessionSyncer = reconciler
|
|
|
|
nodeURL := fmt.Sprintf("http://%s%s", nodeIdentity, cfg.Server.Listen)
|
|
heartbeatWriter = worker.NewHeartbeatWriter(deps.DB, nodeIdentity, mode, nodeURL)
|
|
}
|
|
|
|
if deps.DB != nil {
|
|
adminStatsProvider, statsErr := handlers.NewAdminStatsProvider(appCtx, deps.DB, deps.EventBus)
|
|
if statsErr != nil {
|
|
log.Fatalf("failed to create admin stats provider: %v", statsErr)
|
|
}
|
|
defer adminStatsProvider.Close()
|
|
deps.AdminStatsProvider = adminStatsProvider
|
|
}
|
|
|
|
// Wire recommendations engine, worker, and ratings repo if enabled.
|
|
var recEngine *recommendations.Engine
|
|
var recWorker *recommendations.Worker
|
|
if cfg.Recommendations.Enabled && deps.DB != nil {
|
|
deps.RatingsRepo = catalog.NewRatingsRepo(deps.DB)
|
|
recEngine = recommendations.NewEngine(
|
|
deps.DB,
|
|
deps.RatingsRepo,
|
|
catalog.NewItemRepository(deps.DB),
|
|
catalog.NewPersonRepository(deps.DB),
|
|
userStoreProvider,
|
|
cfg.Recommendations,
|
|
)
|
|
deps.Recommender = recEngine
|
|
deps.CatalogSearchVectorizer = recEngine
|
|
|
|
var err error
|
|
recWorker, err = recommendations.NewWorker(
|
|
recEngine,
|
|
cfg.Recommendations.EmbeddingsCron,
|
|
cfg.Recommendations.TasteProfilesCron,
|
|
cfg.Recommendations.CowatchCron,
|
|
cfg.Recommendations.RecommendationsCron,
|
|
cfg.Recommendations.EmbeddingsJobTimeout,
|
|
)
|
|
if err != nil {
|
|
slog.Error("failed to create recommendation worker", "error", err)
|
|
} else {
|
|
deps.RecWorker = recWorker
|
|
}
|
|
}
|
|
|
|
// Client IP resolver with trusted proxy config.
|
|
if err := clientip.SeedDefaults(ctx, settingsRepo); err != nil {
|
|
log.Fatalf("seed clientip defaults: %v", err)
|
|
}
|
|
trustedCIDRs, err := clientip.LoadTrustedCIDRs(ctx, settingsRepo)
|
|
if err != nil {
|
|
log.Fatalf("load trusted CIDRs: %v", err)
|
|
}
|
|
ipResolver = clientip.NewResolver(trustedCIDRs)
|
|
deps.ClientIPResolver = ipResolver
|
|
// Hot-reload trusted proxies on settings changes via two complementary
|
|
// paths. The direct event-bus subscription re-reads only the clientip key,
|
|
// so a malformed unrelated setting (which fails the whole-config reload)
|
|
// cannot leave stale trust CIDRs on Redis-backed multi-instance deploys.
|
|
_ = eventBus.Subscribe(appCtx, cache.ChannelAdmin, func(event cache.Event) {
|
|
if event.Type != cache.EventSettingsChanged {
|
|
return
|
|
}
|
|
cidrs, loadErr := clientip.LoadTrustedCIDRs(context.Background(), settingsRepo)
|
|
if loadErr != nil {
|
|
slog.WarnContext(context.Background(), "clientip config reload failed", "component", "app", "error", loadErr)
|
|
return
|
|
}
|
|
ipResolver.UpdateTrustedCIDRs(cidrs)
|
|
})
|
|
// The config watcher covers the Redis-less poll/RequestReload path, so
|
|
// admin UI edits apply without a restart on single-node deployments too.
|
|
configWatcher.OnChange(func(old, updated *config.Config) {
|
|
if old != nil && old.ClientIP.TrustedProxies == updated.ClientIP.TrustedProxies {
|
|
return
|
|
}
|
|
raw := updated.ClientIP.TrustedProxies
|
|
if raw == "" {
|
|
raw = clientip.DefaultTrustedProxies
|
|
}
|
|
cidrs, parseErr := clientip.ParseCIDRs(raw)
|
|
if parseErr != nil {
|
|
slog.WarnContext(context.Background(), "clientip config reload failed", "component", "app", "error", parseErr)
|
|
return
|
|
}
|
|
ipResolver.UpdateTrustedCIDRs(cidrs)
|
|
})
|
|
|
|
// Step 6b: Create rate limiter.
|
|
if cfg.RateLimit.Enabled && deps.DB != nil {
|
|
var perKeyLimiter, globalLimiter ratelimit.RateLimiter
|
|
isMemory := true
|
|
|
|
if cfg.RateLimit.Backend == "redis" {
|
|
redisClient, redisErr := cache.NewRedisClient(cfg.Redis)
|
|
if redisErr != nil {
|
|
log.Fatalf("failed to create Redis client for rate limiting: %v", redisErr)
|
|
}
|
|
if redisClient != nil {
|
|
perKeyLimiter = ratelimit.NewRedisLimiter(redisClient)
|
|
globalLimiter = ratelimit.NewRedisLimiter(redisClient)
|
|
isMemory = false
|
|
defer redisClient.Close()
|
|
}
|
|
}
|
|
|
|
if isMemory {
|
|
perKeyLimiter = ratelimit.NewMemoryLimiter()
|
|
globalLimiter = ratelimit.NewMemoryLimiter()
|
|
}
|
|
defer perKeyLimiter.Close()
|
|
defer globalLimiter.Close()
|
|
|
|
rateLimitMW := ratelimit.NewMiddleware(perKeyLimiter, globalLimiter, settingsRepo, isMemory)
|
|
if err := rateLimitMW.Init(context.Background()); err != nil {
|
|
log.Fatalf("failed to init rate limiter: %v", err)
|
|
}
|
|
|
|
// Subscribe for multi-instance reload (only fires if EventBus is Redis-backed)
|
|
_ = eventBus.Subscribe(appCtx, cache.ChannelAdmin, func(event cache.Event) {
|
|
if event.Type == cache.EventSettingsChanged {
|
|
if reloadErr := rateLimitMW.Reload(context.Background()); reloadErr != nil {
|
|
slog.Warn("rate limit config reload from event failed", "error", reloadErr)
|
|
}
|
|
}
|
|
})
|
|
|
|
deps.RateLimitMW = rateLimitMW
|
|
}
|
|
|
|
// Activity log writer + consumer.
|
|
if err := activitylog.SeedDefaults(ctx, settingsRepo); err != nil {
|
|
log.Fatalf("seed activitylog defaults: %v", err)
|
|
}
|
|
|
|
// Seed default page sections for home and existing libraries.
|
|
sectionRepo := sections.NewRepository(pool)
|
|
var folders []*models.MediaFolder
|
|
if deps.FolderRepo != nil {
|
|
var listErr error
|
|
folders, listErr = deps.FolderRepo.List(ctx)
|
|
if listErr != nil {
|
|
log.Fatalf("list libraries for section defaults: %v", listErr)
|
|
}
|
|
}
|
|
if err := sectionRepo.SeedDefaults(ctx, "home", nil, sections.DefaultHomeSections(folders)); err != nil {
|
|
log.Fatalf("seed home section defaults: %v", err)
|
|
}
|
|
if deps.FolderRepo != nil {
|
|
for _, f := range folders {
|
|
id := f.ID
|
|
if seedErr := sectionRepo.SeedDefaults(ctx, "library", &id, sections.DefaultLibrarySectionsForType(&id, f.Type)); seedErr != nil {
|
|
slog.Warn("seed library section defaults", "library_id", id, "error", seedErr)
|
|
}
|
|
}
|
|
}
|
|
activityPM := partman.NewManager(pool, "activity_log", partman.Weekly, 2)
|
|
if err := activityPM.EnsureFuturePartitions(appCtx); err != nil {
|
|
// Non-fatal: see the operational_logs partition incident. Writes fall
|
|
// back to the default partition and periodic cleanup retries.
|
|
slog.Warn("ensure activity log partitions; continuing in degraded mode", "error", err)
|
|
}
|
|
policyPM := partman.NewManager(pool, "policy_decisions", partman.Daily, 3)
|
|
if err := policyPM.EnsureFuturePartitions(appCtx); err != nil {
|
|
// Non-fatal: decision logs fall back to the default partition and
|
|
// periodic cleanup retries partition creation.
|
|
slog.Warn("ensure policy decision log partitions; continuing in degraded mode", "error", err)
|
|
}
|
|
var activityWriter activitylog.Writer
|
|
activityConsumer := activitylog.NewConsumer(pool, nil, logStreamHub)
|
|
|
|
if cfg.Redis.URL != "" {
|
|
actRedisClient, actRedisErr := cache.NewRedisClient(cfg.Redis)
|
|
if actRedisErr == nil && actRedisClient != nil {
|
|
activityWriter = activitylog.NewRedisWriter(actRedisClient)
|
|
activityConsumer = activitylog.NewConsumer(pool, actRedisClient, logStreamHub)
|
|
go activityConsumer.RunRedis(appCtx)
|
|
defer actRedisClient.Close()
|
|
}
|
|
}
|
|
|
|
if activityWriter == nil {
|
|
memWriter := activitylog.NewMemoryWriter(10000)
|
|
activityWriter = memWriter
|
|
go activityConsumer.RunMemory(appCtx, memWriter.Chan())
|
|
}
|
|
deps.ActivityLogWriter = activityWriter
|
|
deps.ActivityLogRepo = activitylog.NewRepo(pool)
|
|
deps.NodeID = nodeID
|
|
|
|
// Create refresh worker early so the task manager can use it for FindCandidates.
|
|
var refreshWorker *worker.RefreshWorker
|
|
var personRefreshWorker *worker.PersonRefreshWorker
|
|
if needsWorkers && deps.DB != nil {
|
|
refreshWorker = worker.NewRefreshWorker(deps.DB)
|
|
if deps.PersonRefreshQueue != nil {
|
|
personRefreshWorker, _ = deps.PersonRefreshQueue.(*worker.PersonRefreshWorker)
|
|
}
|
|
}
|
|
|
|
// Construct collection service for both the router and the collection sync scheduler.
|
|
var collectionSyncScheduler *catalog.CollectionSyncScheduler
|
|
var userCollectionScheduler *usercollections.Scheduler
|
|
var trendingRefresher *sections.TrendingRefresher
|
|
if needsWorkers && deps.DB != nil {
|
|
collectionRepo := catalog.NewLibraryCollectionRepository(deps.DB)
|
|
collItemRepo := catalog.NewItemRepository(deps.DB)
|
|
libraryItemRepo := catalog.NewLibraryItemRepository(deps.DB)
|
|
collectionService := catalog.NewLibraryCollectionService(collectionRepo, collItemRepo, libraryItemRepo, nil)
|
|
collectionService.TMDBCollections = api.NewTMDBCollectionFetcher(cfg.TMDBAPIKey)
|
|
deps.CollectionService = collectionService
|
|
collectionSyncScheduler = catalog.NewCollectionSyncScheduler(collectionRepo, collectionService, slog.Default())
|
|
|
|
// The trending refresher reuses the section repo (to find used source/
|
|
// window combos), a snapshot repo, an item repo (external-ID matching),
|
|
// and the TMDB fetcher. The Trakt fetcher needs settingsRepo and is
|
|
// propagated onto deps.TrendingRefresher later in router.go.
|
|
trendingRefresher = sections.NewTrendingRefresher(
|
|
sectionRepo,
|
|
sections.NewTrendingSnapshotRepository(pool),
|
|
catalog.NewItemRepository(deps.DB),
|
|
collectionService.TMDBCollections,
|
|
collectionService.TraktCollections,
|
|
)
|
|
deps.TrendingRefresher = trendingRefresher
|
|
|
|
if deps.UserStoreProvider != nil {
|
|
userSync := usercollections.NewService(deps.UserStoreProvider, collItemRepo, libraryItemRepo, nil, slog.Default())
|
|
userSync.TMDBCollections = collectionService.TMDBCollections
|
|
// Trakt fetchers are wired in router.go (they need settingsRepo);
|
|
// router.go propagates them onto userSync once configured.
|
|
userCollectionScheduler = usercollections.NewScheduler(deps.DB, userSync, slog.Default())
|
|
deps.UserCollectionSync = userSync
|
|
deps.UserCollectionScheduler = userCollectionScheduler
|
|
deps.MDBListClient = mdblist.NewClient(cfg.MDBListAPIKey, nil)
|
|
mdblistForReload := deps.MDBListClient
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
mdblistForReload.SetAPIKey(updated.MDBListAPIKey)
|
|
})
|
|
}
|
|
}
|
|
|
|
// White-label branding: one service shared by the API (public read + admin
|
|
// upload), the frontend handler (index.html title, favicon, manifest), and
|
|
// the artwork reconcile task. S3 is optional — pass a nil AssetStore (not
|
|
// the typed-nil *s3client.Client) when it isn't configured so text branding
|
|
// still works without it.
|
|
var brandingStore branding.AssetStore
|
|
if deps.S3Public != nil {
|
|
brandingStore = deps.S3Public
|
|
}
|
|
brandingSvc := branding.NewService(settingsRepo, brandingStore)
|
|
|
|
// Wire up task manager for admin task API.
|
|
if needsWorkers && deps.DB != nil {
|
|
triggerRepo := taskrepository.NewPgTriggerRepository(deps.DB)
|
|
historyRepo := taskrepository.NewPgExecutionRepository(deps.DB)
|
|
taskMgr := taskmanager.New(triggerRepo, historyRepo, triggers.New, slog.Default())
|
|
if deps.EventsHub != nil {
|
|
taskMgr.AddObserver(evt.NewTaskObserver(deps.EventsHub))
|
|
}
|
|
|
|
if deps.FolderRepo != nil && deps.LibraryScanQueue != nil {
|
|
taskMgr.Register(tasks.NewScanLibrariesTask(deps.FolderRepo, deps.LibraryScanQueue, deps.EventBus))
|
|
}
|
|
taskMgr.Register(tasks.NewCleanupOrphanedMediaItemsTask(catalog.NewOrphanedProvisionalCleaner(deps.DB)))
|
|
taskMgr.Register(tasks.NewBackfillMediaItemAliasesTask(catalog.NewItemAliasRepository(deps.DB)))
|
|
if deps.S3Public != nil {
|
|
taskMgr.Register(tasks.NewCleanupArtworkRevisionsTask(
|
|
metadata.NewArtworkRevisionGarbageCollector(deps.DB, deps.S3Public),
|
|
))
|
|
}
|
|
catalogSearchIndexer := catalog.NewCatalogSearchIndexer(deps.DB, settingsRepo)
|
|
taskMgr.Register(tasks.NewSyncCatalogSearchIndexTask(catalogSearchIndexer))
|
|
taskMgr.Register(tasks.NewRebuildCatalogSearchIndexTask(catalogSearchIndexer))
|
|
if deps.IntroAnalyzer != nil {
|
|
taskMgr.Register(tasks.NewDetectIntroMarkersTask(deps.IntroAnalyzer, settingsRepo))
|
|
}
|
|
if deps.MarkerContributionService != nil && deps.MarkerProviderConfig != nil && deps.MarkerContributionStore != nil && deps.FileRepo != nil {
|
|
taskMgr.Register(tasks.NewContributeMarkersTask(
|
|
deps.MarkerContributionService, deps.MarkerProviderConfig, deps.MarkerContributionStore, deps.FileRepo,
|
|
))
|
|
}
|
|
if chapterBackfiller, ok := deps.ChapterThumbnailQueuer.(*chapterthumbs.Service); ok {
|
|
taskMgr.Register(tasks.NewChapterThumbnailBackfillTask(chapterBackfiller, 25))
|
|
}
|
|
taskMgr.Register(tasks.NewActivityLogCleanupTask(deps.DB, settingsRepo, activityPM))
|
|
taskMgr.Register(tasks.NewOperationalLogCleanupTask(deps.DB, settingsRepo, opsPM))
|
|
var diagnosticsStore diagnostics.ObjectStore
|
|
if deps.S3Private != nil {
|
|
diagnosticsStore = diagnostics.NewS3ObjectStore(deps.S3Private)
|
|
}
|
|
taskMgr.Register(tasks.NewClientDiagnosticsCleanupTask(
|
|
diagnostics.NewPostgresRepository(deps.DB),
|
|
settingsRepo,
|
|
diagnosticsStore,
|
|
))
|
|
taskMgr.Register(tasks.NewPolicyDecisionLogCleanupTask(deps.DB, settingsRepo, policyPM))
|
|
if deps.FileRepo != nil {
|
|
// Download prepare-to-file pipeline (Phase 3): a durable, leased encode
|
|
// queue hosted on the task manager. Built here (before Start) and shared
|
|
// with the API via deps so the download service can enqueue jobs.
|
|
artifactMgr := downloads.NewArtifactManager(
|
|
downloads.NewArtifactRepository(deps.DB),
|
|
downloads.NewRepository(deps.DB),
|
|
deps.FileRepo,
|
|
downloads.NewPlaybackPreparer(),
|
|
deps.NodeID,
|
|
func() *config.Config {
|
|
if deps.LiveConfig != nil {
|
|
if c := deps.LiveConfig(); c != nil {
|
|
return c
|
|
}
|
|
}
|
|
return deps.Config
|
|
},
|
|
func(ctx context.Context, d *downloads.Download) {
|
|
if deps.EventsHub == nil {
|
|
return
|
|
}
|
|
_ = deps.EventsHub.PublishJSON(ctx, evt.ChannelUserState, "download", map[string]any{
|
|
"download_id": d.ID,
|
|
"status": d.Status,
|
|
"media_item_id": d.ContentID,
|
|
"format": d.Format,
|
|
}, evt.PublishOptions{UserID: d.UserID, ProfileID: d.ProfileID})
|
|
},
|
|
)
|
|
encodeTask := tasks.NewEncodeDownloadArtifactsTask(artifactMgr)
|
|
artifactMgr.SetKick(func() { _ = taskMgr.RunTask(appCtx, encodeTask.Key()) })
|
|
taskMgr.Register(encodeTask)
|
|
deps.ArtifactManager = artifactMgr
|
|
}
|
|
if notificationSystem != nil {
|
|
taskMgr.Register(tasks.NewSeedContentAvailabilityTask(notificationSystem))
|
|
taskMgr.Register(tasks.NewRebuildReleaseInterestTask(notificationSystem))
|
|
taskMgr.Register(tasks.NewNotificationsRetentionTask(notificationSystem))
|
|
}
|
|
if userStoreProvider != nil {
|
|
taskMgr.Register(tasks.NewSettingMutationsRetentionTask(userstore.NewSettingMutationSweeper(
|
|
auth.NewUserRepository(deps.DB), userStoreProvider,
|
|
)))
|
|
}
|
|
if matchWorker != nil {
|
|
taskMgr.Register(tasks.NewMatchMediaTask(matchWorker))
|
|
}
|
|
if refreshWorker != nil && metadataService != nil {
|
|
taskMgr.Register(tasks.NewRefreshMetadataTask(refreshWorker, metadataService))
|
|
}
|
|
if metadataImageCacheProcessor != nil {
|
|
taskMgr.Register(tasks.NewCacheMetadataImagesTask(metadataImageCacheProcessor))
|
|
}
|
|
if deps.S3Public != nil {
|
|
identity := tasks.ArtworkStorageIdentity(cfg.S3.Public.Endpoint, cfg.S3.Public.Bucket, cfg.S3.Public.KeyPrefix)
|
|
// Seed the fingerprint on first boot so an unchanged storage
|
|
// identity never triggers a sweep. On the boot after a provider
|
|
// change the stored (old) identity survives this call and the
|
|
// startup trigger runs the reconcile.
|
|
if _, err := settingsRepo.SetIfAbsent(appCtx, tasks.ArtworkStorageIdentityKey, identity); err != nil {
|
|
slog.Warn("artwork reconcile: seeding storage identity failed", "error", err)
|
|
}
|
|
var brandingReconciler tasks.BrandingAssetReconciler
|
|
if brandingSvc != nil && brandingSvc.HasStorage() {
|
|
brandingReconciler = brandingSvc
|
|
}
|
|
taskMgr.Register(tasks.NewReconcileArtworkCacheTask(
|
|
metadata.NewArtworkCacheReconciler(deps.DB, deps.S3Public),
|
|
settingsRepo,
|
|
brandingReconciler,
|
|
identity,
|
|
))
|
|
}
|
|
if pluginAutoUpdater != nil {
|
|
taskMgr.Register(tasks.NewCheckPluginUpdatesTask(pluginAutoUpdater))
|
|
}
|
|
if collectionSyncScheduler != nil {
|
|
taskMgr.Register(tasks.NewSyncCollectionsTask(collectionSyncScheduler))
|
|
}
|
|
if trendingRefresher != nil {
|
|
taskMgr.Register(tasks.NewRefreshTrendingDiscoverTask(trendingRefresher))
|
|
}
|
|
if userCollectionScheduler != nil {
|
|
taskMgr.Register(tasks.NewSyncUserCollectionsTask(userCollectionScheduler))
|
|
}
|
|
if watchProviderService != nil {
|
|
taskMgr.Register(tasks.NewSyncWatchProvidersTask(watchProviderService))
|
|
}
|
|
requestReconcileSvc := mediarequests.NewService(
|
|
mediarequests.NewRepository(deps.DB, deps.SecretCipher),
|
|
nil,
|
|
mediarequests.NewCatalogPresence(
|
|
catalog.NewItemRepository(deps.DB),
|
|
catalog.NewProviderIDRepository(deps.DB),
|
|
),
|
|
)
|
|
requestReconcileSvc.SetRequesterIdentityResolver(plugins.RequesterIdentityFromLookup(plugins.NewPgUserIdentityLookup(deps.DB)))
|
|
api.AttachRequestRouter(requestReconcileSvc, pluginService)
|
|
requestReconcileSvc.SetGroupPolicyProvider(accessGroupStore)
|
|
if userStoreProvider != nil {
|
|
userRepo := auth.NewUserRepository(deps.DB)
|
|
profileTokens := access.NewProfileTokenService(cfg.Auth.JWTSecret, 0)
|
|
var reconcileResolver scopeResolver
|
|
if policySystem != nil {
|
|
reconcileResolver = policy.NewViewerResolver(userRepo, userStoreProvider, profileTokens, policySystem.PDP(), accessGroupStore)
|
|
} else {
|
|
// Legacy resolver: proxy/test wiring without a policy system. Production integrated/api modes always take the policy path. Removed with the legacy cleanup phase.
|
|
reconcileResolver = access.NewResolver(userRepo, userStoreProvider, profileTokens, accessGroupStore)
|
|
}
|
|
requestReconcileSvc.SetEntitlementResolver(scopeEntitlementResolver{resolver: reconcileResolver})
|
|
}
|
|
if notificationSystem != nil {
|
|
requestReconcileSvc.SetFulfillmentNotifier(notifications.NewRequestFulfillmentNotifier(notificationSystem))
|
|
}
|
|
taskMgr.Register(tasks.NewReconcileRequestsTask(requestReconcileSvc, 100))
|
|
if deps.FolderRepo != nil && deps.LibraryScanQueue != nil && pluginService != nil && pluginInstallationStore != nil {
|
|
autoscanRepo := autoscan.NewRepository(deps.DB, deps.SecretCipher)
|
|
if err := autoscanRepo.MarkInterruptedEvents(appCtx); err != nil {
|
|
slog.Warn("autoscan: failed to mark interrupted polls", "err", err)
|
|
}
|
|
autoscanSvc := api.BuildAutoscanService(
|
|
autoscanRepo,
|
|
pluginService,
|
|
pluginInstallationStore,
|
|
mediarequests.NewRepository(deps.DB, deps.SecretCipher),
|
|
deps.FolderRepo,
|
|
deps.LibraryScanQueue,
|
|
deps.RedisClient,
|
|
)
|
|
// The poll task's default interval seeds the schedule from the stored
|
|
// settings (DefaultPollIntervalSeconds); per-cycle gating still runs
|
|
// off the live settings inside PollOnce. Seed in MILLISECONDS as
|
|
// seconds*1000 — the SAME computation HandleUpdateSettings uses to
|
|
// reschedule — so startup and reschedule agree for sub-minute and
|
|
// non-60-multiple intervals (the old seconds/60 minutes path diverged).
|
|
var intervalMs int64 = 10 * 60 * 1000
|
|
if settings, serr := autoscanRepo.GetSettings(appCtx); serr == nil && settings.DefaultPollIntervalSeconds > 0 {
|
|
intervalMs = int64(settings.DefaultPollIntervalSeconds) * 1000
|
|
}
|
|
taskMgr.Register(tasks.NewAutoscanPollTask(autoscanSvc, intervalMs))
|
|
taskMgr.Register(tasks.NewAutoscanWebhookRetryTask(autoscanSvc))
|
|
}
|
|
reconcileProviderIDRepo := catalog.NewProviderIDRepository(deps.DB)
|
|
reconcileEpisodeRepo := catalog.NewEpisodeRepository(deps.DB)
|
|
historyResolver := watchstate.NewStableIdentityResolver(nil, reconcileEpisodeRepo, reconcileProviderIDRepo)
|
|
historyReconciler := watchstate.NewHistoryReconciler(deps.DB, historyResolver)
|
|
taskMgr.Register(tasks.NewRepairProviderIDIntegrityTask(metadata.NewProviderIDIntegrityRepairer(deps.DB), historyReconciler))
|
|
taskMgr.Register(tasks.NewReconcileWatchHistoryTask(historyReconciler))
|
|
taskMgr.Register(tasks.NewSyncPodcastFeedsTask(podcastfeed.New(), podcastfeed.NewDBStore(deps.DB)))
|
|
if audiobookEnricher != nil {
|
|
taskMgr.Register(tasks.NewSyncAudiobookMetadataTask(audiobookEnricher))
|
|
}
|
|
if ebookEnricher != nil {
|
|
taskMgr.Register(tasks.NewSyncEbookMetadataTask(ebookEnricher))
|
|
taskMgr.Register(tasks.NewBackfillEbookMetadataTask(ebookEnricher))
|
|
}
|
|
if mangaEnricher != nil {
|
|
taskMgr.Register(tasks.NewSyncMangaMetadataTask(mangaEnricher))
|
|
}
|
|
if pluginInstallationStore != nil && pluginRuntimeConfigStore != nil && pluginService != nil {
|
|
pluginTasks, err := plugins.NewTaskRegistryWithTypedResolver(pluginInstallationStore, pluginRuntimeConfigStore, pluginService).Tasks(appCtx)
|
|
if err != nil {
|
|
log.Fatalf("plugin task registry: %v", err)
|
|
}
|
|
for _, pluginTask := range pluginTasks {
|
|
taskMgr.Register(pluginTask)
|
|
}
|
|
}
|
|
|
|
taskMgr.Start(appCtx)
|
|
defer taskMgr.Stop()
|
|
deps.TaskManager = taskMgr
|
|
slog.Info("task manager started")
|
|
}
|
|
|
|
// Build the ABS-compatible REST + Socket.io handler when a DB pool is
|
|
// available. Routes are mounted at the root level by NewRouter (not under
|
|
// /api/v1/) so ABS clients resolve /login, /api/*, /abs/api/*, and
|
|
// /abs/socket.io/* without path prefix hacks.
|
|
if absCompatEnabled && deps.DB != nil {
|
|
absUserRepo := auth.NewUserRepository(deps.DB)
|
|
absSessionRepo := auth.NewSessionRepository(deps.DB)
|
|
absJWTService := auth.NewJWTService(
|
|
cfg.Auth.JWTSecret,
|
|
cfg.Auth.AccessTokenExpiry,
|
|
cfg.Auth.RefreshTokenExpiry,
|
|
)
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
absJWTService.SetExpiries(updated.Auth.AccessTokenExpiry, updated.Auth.RefreshTokenExpiry)
|
|
})
|
|
absAuthSvc := auth.NewService(
|
|
auth.NewLocalProvider(absUserRepo, absSessionRepo),
|
|
absJWTService,
|
|
absSessionRepo,
|
|
absUserRepo,
|
|
nil, // invite codes: not needed for ABS compat
|
|
nil, // settings: not needed here
|
|
nil, // user store: not needed here
|
|
)
|
|
absItemRepo := catalog.NewItemRepository(deps.DB)
|
|
absEpisodeRepo := catalog.NewEpisodeRepository(deps.DB)
|
|
absSeasonRepo := catalog.NewSeasonRepository(deps.DB)
|
|
absPersonRepo := catalog.NewPersonRepository(deps.DB)
|
|
var absFileFetcher catalog.FileVersionFetcher
|
|
if deps.FileRepo != nil {
|
|
absFileFetcher = deps.FileRepo
|
|
}
|
|
absDetailSvc := catalog.NewDetailService(absItemRepo, absEpisodeRepo, absSeasonRepo, absPersonRepo, absFileFetcher)
|
|
if deps.ImageResolver != nil {
|
|
absDetailSvc.SetImageResolver(deps.ImageResolver)
|
|
}
|
|
var absScopeResolver scopeResolver
|
|
if policySystem != nil {
|
|
absScopeResolver = policy.NewViewerResolver(absUserRepo, userStoreProvider, nil, policySystem.PDP(), accessGroupStore)
|
|
} else {
|
|
absScopeResolver = access.NewResolver(absUserRepo, userStoreProvider, nil, accessGroupStore)
|
|
}
|
|
absHDeps := audiobooks.ABSHandlerDeps{
|
|
Pool: deps.DB,
|
|
Items: absItemRepo,
|
|
Files: deps.FileRepo,
|
|
Settings: settingsRepo,
|
|
Auth: &audiobooks.SiloCredValidator{
|
|
Auth: absAuthSvc,
|
|
Pool: deps.DB,
|
|
},
|
|
AccessResolver: audiobooks.NewABSAccessResolver(absUserRepo, userStoreProvider, absScopeResolver, accessGroupStore),
|
|
Recs: recommendations.NewRepo(deps.DB),
|
|
Detail: absDetailSvc,
|
|
SessionMgr: sessionMgr,
|
|
SessionSyncer: deps.SessionSyncer,
|
|
}
|
|
absH := audiobooksService.BuildABSHandler(absHDeps)
|
|
deps.ABSHandler = absH
|
|
}
|
|
_ = audiobooksService
|
|
|
|
if deps.DB != nil && pluginInstallationStore != nil && pluginRuntimeConfigStore != nil && deps.PluginService != nil {
|
|
userRepo := auth.NewUserRepository(deps.DB)
|
|
sessionRepo := auth.NewSessionRepository(deps.DB)
|
|
authBindings, err := pluginRuntimeConfigStore.ListAuthBindings(appCtx)
|
|
if err != nil {
|
|
log.Fatalf("list plugin auth bindings: %v", err)
|
|
}
|
|
for _, binding := range authBindings {
|
|
if binding == nil || !binding.Enabled {
|
|
continue
|
|
}
|
|
installation, err := pluginInstallationStore.GetByID(appCtx, binding.InstallationID)
|
|
if err != nil {
|
|
log.Fatalf("load plugin auth installation %d: %v", binding.InstallationID, err)
|
|
}
|
|
if !installation.Enabled {
|
|
continue
|
|
}
|
|
displayName := binding.CapabilityID
|
|
mode := "credentials"
|
|
iconURL := ""
|
|
capabilities, err := pluginInstallationStore.ListCapabilities(appCtx, binding.InstallationID)
|
|
if err == nil {
|
|
for _, capability := range capabilities {
|
|
if capability != nil && capability.Type == "auth_provider.v1" && capability.ID == binding.CapabilityID {
|
|
if name, ok := capability.Metadata["display_name"].(string); ok && strings.TrimSpace(name) != "" {
|
|
displayName = name
|
|
}
|
|
// auth_modes ["oauth2"] flips the login button into
|
|
// an OAuth-style "Sign in with X" path. Mode is "oauth"
|
|
// when oauth2 is the only declared mode; "credentials"
|
|
// when password is supported alongside or alone.
|
|
if rawModes, ok := capability.Metadata["auth_modes"].([]any); ok {
|
|
hasPassword := false
|
|
hasOAuth := false
|
|
for _, m := range rawModes {
|
|
switch m {
|
|
case "password":
|
|
hasPassword = true
|
|
case "oauth2":
|
|
hasOAuth = true
|
|
}
|
|
}
|
|
if hasOAuth && !hasPassword {
|
|
mode = "oauth"
|
|
}
|
|
}
|
|
if url, ok := capability.Metadata["icon_url"].(string); ok {
|
|
iconURL = url
|
|
}
|
|
break
|
|
}
|
|
}
|
|
}
|
|
|
|
// Generic OIDC and similar multi-instance plugins ship one binary
|
|
// but install once per IdP. Their admin SPA writes display_name
|
|
// + icon_url_path to runtime config so each install renders its
|
|
// own brand on the login page. Manifest values are the fallback.
|
|
if runtimeConfigs, err := pluginRuntimeConfigStore.ListGlobalConfigs(appCtx, binding.InstallationID); err == nil {
|
|
for _, rc := range runtimeConfigs {
|
|
switch rc.Key {
|
|
case "display_name":
|
|
if v, ok := rc.Value["value"].(string); ok && strings.TrimSpace(v) != "" {
|
|
displayName = v
|
|
}
|
|
case "icon_url_path":
|
|
if v, ok := rc.Value["value"].(string); ok && strings.TrimSpace(v) != "" {
|
|
iconURL = fmt.Sprintf("/api/v1/plugins/%d/assets/%s", binding.InstallationID, strings.TrimLeft(v, "/"))
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
deps.AuthProviders = append(deps.AuthProviders, auth.RegisteredProvider{
|
|
Info: auth.LoginProviderInfo{
|
|
ID: fmt.Sprintf("plugin:%d:%s", binding.InstallationID, binding.CapabilityID),
|
|
DisplayName: displayName,
|
|
Mode: mode,
|
|
Default: binding.DefaultLogin,
|
|
IconURL: iconURL,
|
|
InstallationID: binding.InstallationID,
|
|
},
|
|
Provider: auth.NewPluginProvider(
|
|
auth.PluginProviderConfig{
|
|
InstallationID: binding.InstallationID,
|
|
CapabilityID: binding.CapabilityID,
|
|
DisplayName: displayName,
|
|
AutoProvision: binding.AutoProvision,
|
|
},
|
|
sessionRepo,
|
|
userRepo,
|
|
deps.DB,
|
|
deps.PluginService,
|
|
),
|
|
})
|
|
}
|
|
}
|
|
|
|
// Step 7: Build HTTP router with all dependencies.
|
|
// compatServer is populated after the compat server is constructed below;
|
|
// the closure captures the pointer so revocation calls reach the live instance.
|
|
var compatServer *jellycompat.Server
|
|
deps.OnUserSessionsRevoked = func(ctx context.Context, userID int) {
|
|
if compatServer != nil {
|
|
compatServer.SessionStore().DeleteByUserID(userID)
|
|
}
|
|
}
|
|
|
|
distFS, fsErr := fs.Sub(siloweb.DistFS, "dist")
|
|
if fsErr != nil {
|
|
log.Fatalf("failed to create frontend FS: %v", fsErr)
|
|
}
|
|
deps.FrontendFS = distFS
|
|
server.WebDistFS = distFS
|
|
|
|
// Expose the branding service (constructed before the task manager) to the
|
|
// API and the frontend handler.
|
|
deps.BrandingService = brandingSvc
|
|
server.Branding = brandingSvc
|
|
|
|
router := api.NewRouter(deps)
|
|
|
|
// Step 8: Expose Prometheus metrics endpoint (not behind auth).
|
|
metricsMux := http.NewServeMux()
|
|
metricsMux.Handle("/metrics", promhttp.Handler())
|
|
metricsMux.Handle("/api/", router)
|
|
// ABS-compat is NOT mounted on the main listener — see the "ABS compat
|
|
// listener" block below. It binds its own port so the discovery probes
|
|
// (/ping, /healthcheck, /status, /init, /login, /socket.io) own the URL
|
|
// space without collision with silo's SPA fallback. Mirrors how the
|
|
// Jellyfin compat server is set up at :8096.
|
|
metricsMux.Handle("/", server.FrontendHandler())
|
|
|
|
// Step 9: Start background workers (if needed).
|
|
var sessionCleaner *worker.SessionCleaner
|
|
var adminJobRunner *adminjob.Runner
|
|
|
|
if needsWorkers && deps.DB != nil {
|
|
if reconciler == nil {
|
|
log.Fatal("reconciler must be initialized before starting workers")
|
|
}
|
|
reconciler.Start()
|
|
defer reconciler.Stop()
|
|
|
|
if heartbeatWriter != nil {
|
|
heartbeatWriter.Start()
|
|
defer heartbeatWriter.Stop()
|
|
}
|
|
|
|
// RefreshWorker is kept as a RefreshCandidateFinder for the task manager's
|
|
// RefreshMetadataTask but no longer runs its own background loop.
|
|
// Scanning is handled exclusively by the task manager's ScanLibrariesTask.
|
|
if personRefreshWorker != nil {
|
|
personRefreshWorker.Start()
|
|
defer personRefreshWorker.Stop()
|
|
}
|
|
|
|
sessionCleaner = worker.NewSessionCleaner(deps.DB, cfg.UserDB.StaleGraceSeconds)
|
|
sessionCleaner.EventBus = deps.EventBus
|
|
sessionCleaner.EventsHub = deps.EventsHub
|
|
sessionCleaner.Start()
|
|
defer sessionCleaner.Stop()
|
|
|
|
var templateBundleApplyExecutor interface {
|
|
ExecuteTemplateBundleApply(context.Context, adminjob.TemplateBundleApplyRequest, func(int, int, string)) (any, error)
|
|
}
|
|
if deps.CollectionService != nil {
|
|
collectionRepo := catalog.NewLibraryCollectionRepository(deps.DB)
|
|
itemRepo := catalog.NewItemRepository(deps.DB)
|
|
collectionHandler := handlers.NewLibraryCollectionHandler(
|
|
collectionRepo,
|
|
deps.CollectionService,
|
|
itemRepo,
|
|
4*time.Hour,
|
|
nil,
|
|
deps.S3Public,
|
|
)
|
|
collectionHandler.FrontendFS = deps.FrontendFS
|
|
collectionHandler.SectionRepo = sectionRepo
|
|
collectionHandler.FolderRepo = deps.FolderRepo
|
|
if collectionHandler.FolderRepo == nil {
|
|
collectionHandler.FolderRepo = catalog.NewFolderRepository(deps.DB)
|
|
}
|
|
templateBundleApplyExecutor = collectionHandler
|
|
}
|
|
|
|
adminJobRunner = adminjob.NewRunner(
|
|
adminjob.NewRepository(deps.DB),
|
|
catalogseed.NewService(deps.DB, catalog.NewPersonRepository(deps.DB), recommendations.NewRepo(deps.DB)),
|
|
deps.S3Private,
|
|
itemRefreshExecutor,
|
|
libraryRefreshExecutor,
|
|
adminjob.NewLibraryDeleteExecutor(deps.FolderRepo, sectionRepo,
|
|
librarySettingsCleaner(deps.DB, userStoreProvider)),
|
|
adminjob.NewImageCacheCleanupExecutor(deps.S3Public),
|
|
templateBundleApplyExecutor,
|
|
deps.RealtimeHub,
|
|
)
|
|
adminJobRunner.SetCancelRegistry(adminJobCancelRegistry)
|
|
adminJobRunner.Start()
|
|
defer adminJobRunner.Stop()
|
|
|
|
// Start recommendation worker if enabled (reuse worker created above).
|
|
if recWorker != nil {
|
|
recWorker.Start()
|
|
defer recWorker.Stop()
|
|
|
|
// Check if this is first run (no embeddings yet).
|
|
embCount, _ := recommendations.NewRepo(deps.DB).EmbeddingCount(appCtx)
|
|
if embCount == 0 {
|
|
slog.Info("first run detected, triggering initial embedding")
|
|
recWorker.RunEmbeddingsNow()
|
|
}
|
|
}
|
|
|
|
slog.Info("background workers started")
|
|
}
|
|
|
|
// Step 10: Create and start the HTTP server.
|
|
srv := &http.Server{
|
|
Addr: cfg.Server.Listen,
|
|
Handler: metricsMux,
|
|
ReadTimeout: 30 * time.Second,
|
|
WriteTimeout: 120 * time.Second,
|
|
IdleTimeout: 120 * time.Second,
|
|
}
|
|
|
|
var compatSrv *http.Server
|
|
if (mode == "integrated" || mode == "api") && cfg.JellyfinCompat.Enabled && cfg.JellyfinCompat.Listen != "" {
|
|
compatDeps := jellycompat.Dependencies{
|
|
Config: cfg,
|
|
AppContext: appCtx,
|
|
LiveConfig: configWatcher.Config,
|
|
DB: deps.DB,
|
|
SecretCipher: dataCipher,
|
|
ClientIPResolver: ipResolver,
|
|
NodePlanner: deps.NodePlanner,
|
|
JWTSecret: cfg.Auth.JWTSecret,
|
|
RecWorker: recWorker,
|
|
FrontendFS: deps.FrontendFS,
|
|
// Hand remote-transcode recipes to the shared recipe store so a dedicated
|
|
// transcode node that restarts can rebuild a jellycompat session.
|
|
RecipeNodeStore: noderecipe.NewStore(apiRedisClient, 0),
|
|
SessionSyncer: deps.SessionSyncer,
|
|
}
|
|
|
|
// Wire direct dependencies when DB is available.
|
|
if deps.DB != nil {
|
|
browseRepo := catalog.NewBrowseRepository(deps.DB)
|
|
itemRepo := catalog.NewItemRepository(deps.DB)
|
|
seasonRepo := catalog.NewSeasonRepository(deps.DB)
|
|
episodeRepo := catalog.NewEpisodeRepository(deps.DB)
|
|
providerIDRepo := catalog.NewProviderIDRepository(deps.DB)
|
|
personRepo := catalog.NewPersonRepository(deps.DB)
|
|
folderRepo := deps.FolderRepo
|
|
|
|
var fileFetcher catalog.FileVersionFetcher
|
|
if deps.FileRepo != nil {
|
|
fileFetcher = deps.FileRepo
|
|
}
|
|
|
|
detailSvc := catalog.NewDetailService(itemRepo, episodeRepo, seasonRepo, personRepo, fileFetcher)
|
|
detailSvc.SetFolderRepository(folderRepo)
|
|
detailSvc.SetGroupClaimRepository(catalog.NewGroupClaimRepository(deps.DB))
|
|
detailSvc.SetProbeEnsurer(deps.ProbeEnsurer)
|
|
detailSvc.SetChapterThumbnailQueuer(deps.ChapterThumbnailQueuer)
|
|
if deps.ImageResolver != nil {
|
|
detailSvc.SetImageResolver(deps.ImageResolver)
|
|
}
|
|
|
|
compatDeps.BrowseRepo = browseRepo
|
|
compatDeps.ItemRepo = itemRepo
|
|
compatDeps.SeasonRepo = seasonRepo
|
|
compatDeps.EpisodeRepo = episodeRepo
|
|
compatDeps.ProviderIDRepo = providerIDRepo
|
|
compatDeps.StableIdentityResolver = watchstate.NewStableIdentityResolver(itemRepo, episodeRepo, providerIDRepo)
|
|
compatDeps.DetailSvc = detailSvc
|
|
compatDeps.FolderRepo = folderRepo
|
|
compatDeps.SessionMgr = sessionMgr
|
|
compatDeps.UserStoreProvider = userStoreProvider
|
|
compatDeps.WatchCompletionObserver = deps.WatchCompletionObserver
|
|
compatDeps.SettingsRepo = settingsRepo
|
|
compatDeps.PersonRepo = personRepo
|
|
if watchProviderService != nil {
|
|
compatDeps.WatchScrobbler = watchProviderService
|
|
}
|
|
compatSearchService := catalog.NewCatalogSearchService(
|
|
appCtx,
|
|
settingsRepo,
|
|
itemRepo,
|
|
catalog.NewSearchIndexEventRepository(deps.DB),
|
|
deps.CatalogSearchVectorizer,
|
|
)
|
|
if compatSearchService != nil {
|
|
compatSearchService.StartCoverageRefresh(appCtx)
|
|
compatDeps.CatalogSearchProvider = compatSearchService.Provider()
|
|
// Latch the resolved provider for the package-level enqueue
|
|
// helpers (idempotent with the API router's latch; this also
|
|
// covers modes that wire jellycompat without the router).
|
|
activeSearchProvider := catalog.SearchProviderPostgres
|
|
if _, ok := compatSearchService.Provider().(*catalog.MeilisearchSearchProvider); ok {
|
|
activeSearchProvider = catalog.SearchProviderMeilisearch
|
|
}
|
|
catalog.SetActiveSearchIndexProvider(activeSearchProvider)
|
|
}
|
|
|
|
if deps.S3Public != nil {
|
|
compatDeps.PosterPresigner = deps.S3Public
|
|
compatDeps.S3Client = deps.S3Public
|
|
compatDeps.S3Bucket = deps.S3Public.Bucket()
|
|
}
|
|
|
|
if deps.FileRepo != nil {
|
|
compatDeps.FileResolver = deps.FileRepo
|
|
}
|
|
|
|
compatDeps.SubtitleRepo = subtitles.NewPgRepository(deps.DB, deps.SecretCipher)
|
|
|
|
// Construct auth service for jellycompat login.
|
|
userRepo := auth.NewUserRepository(deps.DB)
|
|
compatDeps.APIKeyValidator = auth.NewAPIKeyRepository(deps.DB)
|
|
compatDeps.APIKeyUserLoader = userRepo
|
|
compatDeps.ScanQueue = deps.LibraryScanQueue
|
|
sessionRepo := auth.NewSessionRepository(deps.DB)
|
|
jwtService := auth.NewJWTService(
|
|
cfg.Auth.JWTSecret,
|
|
cfg.Auth.AccessTokenExpiry,
|
|
cfg.Auth.RefreshTokenExpiry,
|
|
)
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
jwtService.SetExpiries(updated.Auth.AccessTokenExpiry, updated.Auth.RefreshTokenExpiry)
|
|
})
|
|
provider := auth.NewLocalProvider(userRepo, sessionRepo)
|
|
compatDeps.AuthService = auth.NewService(provider, jwtService, sessionRepo, userRepo, nil, nil, nil)
|
|
|
|
// Access filter resolver for viewer-scoped library access.
|
|
// Backed by the shared access.Resolver so account-level library
|
|
// restrictions (users.library_ids), profile restrictions,
|
|
// user-disabled libraries, and rating/quality ceilings apply to
|
|
// the compat API exactly as they do to the native API.
|
|
if userStoreProvider != nil {
|
|
var compatScopeResolver jellycompat.ScopeResolver
|
|
if policySystem != nil {
|
|
compatScopeResolver = policy.NewViewerResolver(
|
|
userRepo,
|
|
userStoreProvider,
|
|
nil, // profile tokens unused: compat login already verifies PINs
|
|
policySystem.PDP(),
|
|
accessGroupStore,
|
|
)
|
|
} else {
|
|
// Legacy resolver: proxy/test wiring without a policy system. Production integrated/api modes always take the policy path. Removed with the legacy cleanup phase.
|
|
compatScopeResolver = access.NewResolver(
|
|
userRepo,
|
|
userStoreProvider,
|
|
nil, // profile tokens unused: compat login already verifies PINs
|
|
accessGroupStore,
|
|
)
|
|
}
|
|
compatDeps.AccessFilterFn = jellycompat.NewScopeAccessFilter(compatScopeResolver)
|
|
}
|
|
}
|
|
|
|
compat := jellycompat.NewServerWithDependencies(compatDeps)
|
|
compatServer = compat
|
|
compatTerminalRecoveryReady = compat.StartBackgroundTasks(context.Background())
|
|
compatSrv = compat.HTTPServer()
|
|
compatSrv.ReadTimeout = 30 * time.Second
|
|
compatSrv.WriteTimeout = 0
|
|
compatSrv.IdleTimeout = 120 * time.Second
|
|
}
|
|
|
|
// ABS-compat listener — dedicated http.Server bound to its own port
|
|
// (default :13378) that hosts the Audiobookshelf-compatible API.
|
|
// Mirrors the Jellyfin compat layout above. The ABS handler mounts
|
|
// onto a fresh chi router here so /ping, /healthcheck, /status, /login,
|
|
// /socket.io, etc. own the URL space at the root — no SPA fallback,
|
|
// no collision with silo's /api/v1.
|
|
var absSrv *http.Server
|
|
if (mode == "integrated" || mode == "api") && deps.ABSHandler != nil && cfg.AudiobookshelfCompat.Listen != "" {
|
|
absRouter := chi.NewRouter()
|
|
absRouter.Use(chimiddleware.Recoverer)
|
|
absRouter.Use(chimiddleware.Compress(5))
|
|
deps.ABSHandler.Mount(absRouter)
|
|
absSrv = &http.Server{
|
|
Addr: cfg.AudiobookshelfCompat.Listen,
|
|
Handler: absRouter,
|
|
ReadHeaderTimeout: 10 * time.Second,
|
|
ReadTimeout: 60 * time.Second,
|
|
WriteTimeout: 0,
|
|
IdleTimeout: 120 * time.Second,
|
|
}
|
|
}
|
|
|
|
// Run non-critical startup work in the background so it doesn't delay the
|
|
// HTTP listener from accepting connections. Steps run sequentially and stop
|
|
// early if the app context is cancelled (shutdown).
|
|
if len(backgroundInit) > 0 {
|
|
go func() {
|
|
start := time.Now()
|
|
for _, step := range backgroundInit {
|
|
if appCtx.Err() != nil {
|
|
return
|
|
}
|
|
func() {
|
|
defer func() {
|
|
if p := recover(); p != nil {
|
|
slog.Error("deferred startup init step panicked; continuing",
|
|
"panic", p, "stack", string(debug.Stack()))
|
|
}
|
|
}()
|
|
step(appCtx)
|
|
}()
|
|
}
|
|
slog.Info("deferred startup init completed", "steps", len(backgroundInit), "duration", time.Since(start))
|
|
}()
|
|
}
|
|
|
|
errCh := make(chan error, 3)
|
|
go func() {
|
|
slog.Info("HTTP server listening", "addr", cfg.Server.Listen)
|
|
if listenErr := srv.ListenAndServe(); listenErr != nil && listenErr != http.ErrServerClosed {
|
|
errCh <- fmt.Errorf("HTTP server error: %w", listenErr)
|
|
}
|
|
}()
|
|
if compatSrv != nil {
|
|
go func() {
|
|
slog.Info("Jellyfin compat server listening", "addr", compatSrv.Addr)
|
|
if listenErr := compatSrv.ListenAndServe(); listenErr != nil && listenErr != http.ErrServerClosed {
|
|
errCh <- fmt.Errorf("jellyfin compat server error: %w", listenErr)
|
|
}
|
|
}()
|
|
}
|
|
if absSrv != nil {
|
|
go func() {
|
|
slog.Info("ABS compat server listening", "addr", absSrv.Addr)
|
|
if listenErr := absSrv.ListenAndServe(); listenErr != nil && listenErr != http.ErrServerClosed {
|
|
errCh <- fmt.Errorf("abs compat server error: %w", listenErr)
|
|
}
|
|
}()
|
|
}
|
|
|
|
// Step 11: Wait for termination signal.
|
|
sigCh := make(chan os.Signal, 1)
|
|
signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT)
|
|
defer signal.Stop(sigCh)
|
|
|
|
select {
|
|
case sig := <-sigCh:
|
|
appCancel()
|
|
slog.Info("received signal, shutting down", "signal", sig)
|
|
case <-restartReqCh:
|
|
appCancel()
|
|
slog.Info("server restart requested, shutting down")
|
|
case serverErr := <-errCh:
|
|
appCancel()
|
|
slog.Error("server error, shutting down", "error", serverErr)
|
|
}
|
|
|
|
// Step 12: Graceful shutdown sequence.
|
|
slog.Info("beginning graceful shutdown")
|
|
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 30*time.Second)
|
|
defer shutdownCancel()
|
|
|
|
// 1. Stop accepting new requests.
|
|
if shutdownErr := srv.Shutdown(shutdownCtx); shutdownErr != nil {
|
|
slog.Error("HTTP shutdown error", "error", shutdownErr)
|
|
}
|
|
if compatSrv != nil {
|
|
if shutdownErr := compatSrv.Shutdown(shutdownCtx); shutdownErr != nil {
|
|
slog.Error("jellyfin compat shutdown error", "error", shutdownErr)
|
|
}
|
|
}
|
|
if absSrv != nil {
|
|
if shutdownErr := absSrv.Shutdown(shutdownCtx); shutdownErr != nil {
|
|
slog.Error("abs compat shutdown error", "error", shutdownErr)
|
|
}
|
|
}
|
|
|
|
// 2. Clean up stale sessions.
|
|
if sessionCleaner != nil {
|
|
cleaned, cleanErr := sessionCleaner.CleanStale(shutdownCtx)
|
|
if cleanErr != nil {
|
|
slog.Error("stale session cleanup error", "error", cleanErr)
|
|
} else if cleaned > 0 {
|
|
slog.Info("cleaned stale sessions", "count", cleaned)
|
|
}
|
|
}
|
|
|
|
// 2b. Remove this node's heartbeat and sessions from shared state.
|
|
if heartbeatWriter != nil {
|
|
if err := heartbeatWriter.CleanupSelf(shutdownCtx); err != nil {
|
|
slog.Error("heartbeat cleanup error", "error", err)
|
|
}
|
|
}
|
|
|
|
// 3. Close user store provider.
|
|
if userStoreProvider != nil {
|
|
if closeErr := userStoreProvider.Close(); closeErr != nil {
|
|
slog.Error("user store provider close error", "error", closeErr)
|
|
}
|
|
}
|
|
|
|
// 4. (match worker is now managed by the task manager — no separate cancel needed)
|
|
|
|
// Suppress unused variable warnings for workers used only in deferred calls.
|
|
_ = reconciler
|
|
_ = heartbeatWriter
|
|
_ = refreshWorker
|
|
_ = adminJobRunner
|
|
|
|
slog.Info("server stopped")
|
|
}
|
|
|
|
// startStandaloneServer runs a standalone HTTP server for proxy/transcode modes.
|
|
// It listens on the given address, handles graceful shutdown on SIGTERM/SIGINT.
|
|
func startStandaloneServer(addr string, handler http.Handler) {
|
|
srv := &http.Server{
|
|
Addr: addr,
|
|
Handler: handler,
|
|
ReadTimeout: 30 * time.Second,
|
|
WriteTimeout: 0, // no timeout for long streams
|
|
IdleTimeout: 120 * time.Second,
|
|
}
|
|
|
|
errCh := make(chan error, 1)
|
|
go func() {
|
|
slog.Info("HTTP server listening", "addr", addr)
|
|
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
errCh <- fmt.Errorf("HTTP server error: %w", err)
|
|
}
|
|
}()
|
|
|
|
sigCh := make(chan os.Signal, 1)
|
|
signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT)
|
|
|
|
select {
|
|
case sig := <-sigCh:
|
|
slog.Info("received signal, shutting down", "signal", sig)
|
|
case serverErr := <-errCh:
|
|
slog.Error("server error, shutting down", "error", serverErr)
|
|
}
|
|
|
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
|
defer cancel()
|
|
if err := srv.Shutdown(shutdownCtx); err != nil {
|
|
slog.Error("HTTP shutdown error", "error", err)
|
|
}
|
|
slog.Info("server stopped")
|
|
}
|
|
|
|
// newS3ClientIfConfigured creates an S3 client only if the bucket name is
|
|
// configured. Returns nil if the bucket is empty (not configured).
|
|
func newS3ClientIfConfigured(cfg s3client.BucketConfig) *s3client.Client {
|
|
if cfg.Bucket == "" {
|
|
return nil
|
|
}
|
|
return s3client.NewClient(cfg)
|
|
}
|
|
|
|
func configureS3Clients(cfg *config.Config, deps *api.Dependencies) {
|
|
if s3Public := newS3ClientIfConfigured(s3client.BucketConfig{
|
|
Endpoint: cfg.S3.Public.Endpoint,
|
|
PublicEndpoint: cfg.S3.Public.ReadEndpoint,
|
|
Region: cfg.S3.Public.Region,
|
|
Bucket: cfg.S3.Public.Bucket,
|
|
KeyPrefix: cfg.S3.Public.KeyPrefix,
|
|
AccessKey: cfg.S3.Public.AccessKey,
|
|
SecretKey: cfg.S3.Public.SecretKey,
|
|
PathStyle: cfg.S3.Public.PathStyle,
|
|
URLAuth: cfg.S3.Public.URLAuth,
|
|
TokenSecret: cfg.S3.Public.TokenSecret,
|
|
TokenParam: cfg.S3.Public.TokenParam,
|
|
TokenTTL: cfg.S3.Public.TokenTTL,
|
|
}); s3Public != nil {
|
|
deps.S3Public = s3Public
|
|
slog.Info("S3 public assets client configured", "bucket", s3Public.Bucket())
|
|
|
|
// Allow browsers to fetch presigned client-facing assets directly from S3.
|
|
// Skip for public/token auth (e.g. Cloudflare R2) where CORS is managed externally.
|
|
if !s3Public.UsesExternalAuth() {
|
|
corsCtx, corsCancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
if corsErr := s3Public.SetBucketCORS(corsCtx, s3Public.Bucket(), []string{"*"}); corsErr != nil {
|
|
slog.Warn("failed to set CORS on public assets bucket", "error", corsErr)
|
|
}
|
|
corsCancel()
|
|
}
|
|
}
|
|
|
|
if s3Private := newS3ClientIfConfigured(s3client.BucketConfig{
|
|
Endpoint: cfg.S3.Private.Endpoint,
|
|
Region: cfg.S3.Private.Region,
|
|
Bucket: cfg.S3.Private.Bucket,
|
|
KeyPrefix: cfg.S3.Private.KeyPrefix,
|
|
AccessKey: cfg.S3.Private.AccessKey,
|
|
SecretKey: cfg.S3.Private.SecretKey,
|
|
PathStyle: cfg.S3.Private.PathStyle,
|
|
}); s3Private != nil {
|
|
deps.S3Private = s3Private
|
|
slog.Info("S3 private internal client configured", "bucket", s3Private.Bucket())
|
|
if !s3Private.UsesExternalAuth() {
|
|
corsCtx, corsCancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
if corsErr := s3Private.SetBucketCORS(corsCtx, s3Private.Bucket(), []string{"*"}); corsErr != nil {
|
|
slog.Warn("failed to set CORS on private assets bucket", "error", corsErr)
|
|
}
|
|
corsCancel()
|
|
}
|
|
}
|
|
|
|
if s3UserDB := newS3ClientIfConfigured(s3client.BucketConfig{
|
|
Endpoint: cfg.S3.UserDB.Endpoint,
|
|
Region: cfg.S3.UserDB.Region,
|
|
Bucket: cfg.S3.UserDB.Bucket,
|
|
KeyPrefix: cfg.S3.UserDB.KeyPrefix,
|
|
AccessKey: cfg.S3.UserDB.AccessKey,
|
|
SecretKey: cfg.S3.UserDB.SecretKey,
|
|
PathStyle: cfg.S3.UserDB.PathStyle,
|
|
}); s3UserDB != nil {
|
|
deps.S3UserDB = s3UserDB
|
|
slog.Info("S3 user-db client configured", "bucket", s3UserDB.Bucket())
|
|
}
|
|
}
|
|
|
|
type pluginImageResolverCapabilityStore interface {
|
|
ListEnabled(ctx context.Context) ([]*plugins.Installation, error)
|
|
ListCapabilities(ctx context.Context, installationID int) ([]*plugins.Capability, error)
|
|
}
|
|
|
|
func reloadPluginImageResolvers(
|
|
ctx context.Context,
|
|
store pluginImageResolverCapabilityStore,
|
|
resolver *metadata.PluginImageResolver,
|
|
service *plugins.Service,
|
|
) error {
|
|
if resolver == nil {
|
|
return nil
|
|
}
|
|
if store == nil || service == nil {
|
|
resolver.ReplaceSources(nil)
|
|
return nil
|
|
}
|
|
|
|
installations, err := store.ListEnabled(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("list enabled plugin installations: %w", err)
|
|
}
|
|
sort.Slice(installations, func(i, j int) bool {
|
|
if installations[i] == nil {
|
|
return false
|
|
}
|
|
if installations[j] == nil {
|
|
return true
|
|
}
|
|
return installations[i].ID < installations[j].ID
|
|
})
|
|
|
|
var registrations []metadata.PluginImageResolverSourceRegistration
|
|
for _, installation := range installations {
|
|
if installation == nil {
|
|
continue
|
|
}
|
|
// Builtin installations resolve in-process metadata providers only;
|
|
// registering them here would claim their capability id as a gRPC
|
|
// image-resolver scheme with no binary behind it.
|
|
if installation.IsBuiltin() {
|
|
continue
|
|
}
|
|
capabilities, err := store.ListCapabilities(ctx, installation.ID)
|
|
if err != nil {
|
|
return fmt.Errorf("list image resolver capabilities for installation %d: %w", installation.ID, err)
|
|
}
|
|
sort.Slice(capabilities, func(i, j int) bool {
|
|
if capabilities[i] == nil {
|
|
return false
|
|
}
|
|
if capabilities[j] == nil {
|
|
return true
|
|
}
|
|
if capabilities[i].Type != capabilities[j].Type {
|
|
return capabilities[i].Type < capabilities[j].Type
|
|
}
|
|
return capabilities[i].ID < capabilities[j].ID
|
|
})
|
|
|
|
for _, capability := range capabilities {
|
|
if capability == nil {
|
|
continue
|
|
}
|
|
switch capability.Type {
|
|
case sdkcapability.ImageResolver:
|
|
schemes, priority := imageResolverCapabilityConfig(capability)
|
|
if len(schemes) == 0 {
|
|
slog.WarnContext(ctx, "plugin image resolver capability has no valid schemes", "component", "app",
|
|
"installation_id", installation.ID,
|
|
"capability_id", capability.ID)
|
|
continue
|
|
}
|
|
for _, scheme := range schemes {
|
|
source := metadata.NewPluginClientSource(installation.ID, capability.ID, func(
|
|
ctx context.Context, installationID int, capabilityID string,
|
|
) (metadata.PluginMetadataClient, error) {
|
|
return service.ImageResolverClient(ctx, installationID, capabilityID)
|
|
})
|
|
registrations = append(registrations, metadata.PluginImageResolverSourceRegistration{
|
|
Scheme: scheme,
|
|
Source: source,
|
|
Kind: metadata.PluginImageResolverSourceExplicit,
|
|
Priority: priority,
|
|
InstallationID: installation.ID,
|
|
CapabilityID: capability.ID,
|
|
})
|
|
}
|
|
case sdkcapability.MetadataProvider:
|
|
scheme := strings.TrimSpace(capability.ID)
|
|
if !metadata.ValidImageResolverScheme(scheme) {
|
|
slog.WarnContext(ctx, "skipping legacy metadata image resolver with invalid scheme", "component", "app",
|
|
"installation_id", installation.ID,
|
|
"capability_id", capability.ID)
|
|
continue
|
|
}
|
|
source := metadata.NewPluginClientSource(installation.ID, capability.ID, func(
|
|
ctx context.Context, installationID int, capabilityID string,
|
|
) (metadata.PluginMetadataClient, error) {
|
|
return service.MetadataProviderClient(ctx, installationID, capabilityID)
|
|
})
|
|
registrations = append(registrations, metadata.PluginImageResolverSourceRegistration{
|
|
Scheme: scheme,
|
|
Source: source,
|
|
Kind: metadata.PluginImageResolverSourceLegacy,
|
|
InstallationID: installation.ID,
|
|
CapabilityID: capability.ID,
|
|
})
|
|
}
|
|
}
|
|
}
|
|
|
|
resolver.ReplaceSources(registrations)
|
|
slog.InfoContext(ctx, "reloaded plugin image resolvers", "component", "app", "sources", len(registrations))
|
|
return nil
|
|
}
|
|
|
|
func imageResolverCapabilityConfig(capability *plugins.Capability) ([]string, int) {
|
|
if capability == nil {
|
|
return nil, 0
|
|
}
|
|
meta := capabilityMetadataFields(capability.Metadata)
|
|
return metadataStringList(meta["schemes"]), metadataInt(meta["priority"])
|
|
}
|
|
|
|
func capabilityMetadataFields(raw map[string]any) map[string]any {
|
|
if raw == nil {
|
|
return nil
|
|
}
|
|
if nested, ok := raw["metadata"]; ok {
|
|
switch typed := nested.(type) {
|
|
case map[string]any:
|
|
return typed
|
|
}
|
|
}
|
|
return raw
|
|
}
|
|
|
|
func metadataStringList(value any) []string {
|
|
var out []string
|
|
switch typed := value.(type) {
|
|
case []string:
|
|
for _, item := range typed {
|
|
if scheme := strings.TrimSpace(item); metadata.ValidImageResolverScheme(scheme) {
|
|
out = append(out, scheme)
|
|
}
|
|
}
|
|
case []any:
|
|
for _, item := range typed {
|
|
text, ok := item.(string)
|
|
if !ok {
|
|
continue
|
|
}
|
|
if scheme := strings.TrimSpace(text); metadata.ValidImageResolverScheme(scheme) {
|
|
out = append(out, scheme)
|
|
}
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func metadataInt(value any) int {
|
|
switch typed := value.(type) {
|
|
case int:
|
|
return typed
|
|
case int32:
|
|
return int(typed)
|
|
case int64:
|
|
return int(typed)
|
|
case float64:
|
|
return int(typed)
|
|
case json.Number:
|
|
n, _ := typed.Int64()
|
|
return int(n)
|
|
default:
|
|
return 0
|
|
}
|
|
}
|
|
|
|
var watchSyncPluginReloadMu sync.Mutex
|
|
|
|
type markerPluginCapabilityStore interface {
|
|
ListEnabled(ctx context.Context) ([]*plugins.Installation, error)
|
|
ListCapabilities(ctx context.Context, installationID int) ([]*plugins.Capability, error)
|
|
}
|
|
|
|
func reloadWatchSyncPluginProviders(
|
|
ctx context.Context,
|
|
registry *watchsync.Registry,
|
|
store markerPluginCapabilityStore,
|
|
service *plugins.Service,
|
|
repository watchsync.PluginCredentialRepository,
|
|
) error {
|
|
if registry == nil {
|
|
return nil
|
|
}
|
|
watchSyncPluginReloadMu.Lock()
|
|
defer watchSyncPluginReloadMu.Unlock()
|
|
|
|
var providers []watchsync.Provider
|
|
if store == nil || service == nil {
|
|
return registry.ReplacePluginProviders(providers)
|
|
}
|
|
installations, err := store.ListEnabled(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("list enabled watch sync plugin installations: %w", err)
|
|
}
|
|
sort.Slice(installations, func(i, j int) bool {
|
|
if installations[i] == nil {
|
|
return false
|
|
}
|
|
if installations[j] == nil {
|
|
return true
|
|
}
|
|
return installations[i].ID < installations[j].ID
|
|
})
|
|
for _, installation := range installations {
|
|
if installation == nil || installation.IsBuiltin() {
|
|
continue
|
|
}
|
|
capabilities, err := store.ListCapabilities(ctx, installation.ID)
|
|
if err != nil {
|
|
slog.WarnContext(ctx, "skip watch sync plugin with unreadable capabilities",
|
|
"component", "app",
|
|
"installation_id", installation.ID,
|
|
"error", err,
|
|
)
|
|
continue
|
|
}
|
|
for _, capability := range capabilities {
|
|
if capability == nil || capability.Type != sdkcapability.WatchSyncProvider {
|
|
continue
|
|
}
|
|
descriptor, err := plugins.DecodeCapability(capability)
|
|
if err != nil {
|
|
slog.WarnContext(ctx, "skip invalid watch sync plugin capability",
|
|
"component", "app",
|
|
"installation_id", installation.ID,
|
|
"capability_id", capability.ID,
|
|
"error", err,
|
|
)
|
|
continue
|
|
}
|
|
provider, err := watchsync.NewPluginProvider(watchsync.PluginProviderOptions{
|
|
InstallationID: installation.ID,
|
|
ProviderKey: fmt.Sprintf("plugin:%d:%s", installation.ID, capability.ID),
|
|
CapabilityID: capability.ID,
|
|
DisplayName: descriptor.GetDisplayName(),
|
|
Descriptor: descriptor.GetWatchSyncProvider(),
|
|
ResolveClient: func(callCtx context.Context, installationID int, capabilityID string) (watchsync.WatchSyncPluginClient, error) {
|
|
return service.WatchSyncProviderClient(callCtx, installationID, capabilityID)
|
|
},
|
|
ResolveConfig: func(callCtx context.Context, installationID int) (*pluginv1.WatchSyncProviderConfig, error) {
|
|
return service.WatchSyncProviderConfig(callCtx, installationID)
|
|
},
|
|
Repository: repository,
|
|
})
|
|
if err != nil {
|
|
slog.WarnContext(ctx, "skip unsupported watch sync plugin capability",
|
|
"component", "app",
|
|
"installation_id", installation.ID,
|
|
"capability_id", capability.ID,
|
|
"error", err,
|
|
)
|
|
continue
|
|
}
|
|
providers = append(providers, provider)
|
|
}
|
|
}
|
|
return registry.ReplacePluginProviders(providers)
|
|
}
|
|
|
|
type markerPluginRuntimeConfigStore interface {
|
|
ListGlobalConfigs(ctx context.Context, installationID int) ([]*plugins.RuntimeConfig, error)
|
|
PutGlobalConfig(ctx context.Context, installationID int, key string, value map[string]any) error
|
|
}
|
|
|
|
type markerLegacySettingsStore interface {
|
|
Get(ctx context.Context, key string) (string, error)
|
|
}
|
|
|
|
func reloadMarkerPluginProviders(
|
|
ctx context.Context,
|
|
registry *markers.Registry,
|
|
configStore *markers.ProviderConfigStore,
|
|
store markerPluginCapabilityStore,
|
|
runtimeConfigs markerPluginRuntimeConfigStore,
|
|
legacySettings markerLegacySettingsStore,
|
|
resolver *markers.PluginResolverAdapter,
|
|
) error {
|
|
if registry == nil {
|
|
return nil
|
|
}
|
|
var providers []markers.Provider
|
|
if store == nil || resolver == nil {
|
|
return registry.SetProviders(providers)
|
|
}
|
|
|
|
installations, err := store.ListEnabled(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("list enabled plugin installations: %w", err)
|
|
}
|
|
sort.Slice(installations, func(i, j int) bool {
|
|
if installations[i] == nil {
|
|
return false
|
|
}
|
|
if installations[j] == nil {
|
|
return true
|
|
}
|
|
return installations[i].ID < installations[j].ID
|
|
})
|
|
|
|
nextPriority := 1000
|
|
for _, installation := range installations {
|
|
if installation == nil {
|
|
continue
|
|
}
|
|
// Builtin installations expose no marker providers; defense in depth
|
|
// alongside the capability-type filter below.
|
|
if installation.IsBuiltin() {
|
|
continue
|
|
}
|
|
capabilities, err := store.ListCapabilities(ctx, installation.ID)
|
|
if err != nil {
|
|
return fmt.Errorf("list marker provider capabilities for installation %d: %w", installation.ID, err)
|
|
}
|
|
sort.Slice(capabilities, func(i, j int) bool {
|
|
if capabilities[i] == nil {
|
|
return false
|
|
}
|
|
if capabilities[j] == nil {
|
|
return true
|
|
}
|
|
return capabilities[i].ID < capabilities[j].ID
|
|
})
|
|
for _, capability := range capabilities {
|
|
if capability == nil || capability.Type != sdkcapability.MarkerProvider {
|
|
continue
|
|
}
|
|
descriptor, err := plugins.DecodeCapability(capability)
|
|
if err != nil {
|
|
return fmt.Errorf("decode marker provider capability %d/%s: %w", installation.ID, capability.ID, err)
|
|
}
|
|
metadataMap := markerCapabilityMetadata(descriptor)
|
|
provider, err := markers.NewPluginProvider(markers.PluginProviderOptions{
|
|
InstallationID: installation.ID,
|
|
CapabilityID: capability.ID,
|
|
DisplayName: firstNonEmptyMarkerText(descriptor.GetDisplayName(), capability.ID),
|
|
PluginID: installation.PluginID,
|
|
RequiredExternalIDs: markers.PluginRequiredExternalIDsFromMetadata(metadataMap),
|
|
}, resolver)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
providers = append(providers, provider)
|
|
|
|
priority := nextPriority
|
|
nextPriority++
|
|
if configuredPriority, ok := markers.PluginDefaultFetchPriorityFromMetadata(metadataMap); ok {
|
|
priority = configuredPriority
|
|
}
|
|
if configStore != nil {
|
|
defaultConfig := markers.ProviderConfig{
|
|
Provider: provider.ID(),
|
|
FetchEnabled: true,
|
|
FetchPriority: priority,
|
|
ContributeEnabled: false,
|
|
ContributeAutoLocal: false,
|
|
ContributeMinConfidence: 0.95,
|
|
}
|
|
if legacy, ok := legacyIntroDBProviderConfig(configStore, installation, capability, provider.ID()); ok {
|
|
defaultConfig = legacy
|
|
}
|
|
if err := configStore.Ensure(ctx, defaultConfig); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if err := copyLegacyIntroDBPluginConfig(ctx, runtimeConfigs, legacySettings, installation, capability); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return registry.SetProviders(providers)
|
|
}
|
|
|
|
func legacyIntroDBProviderConfig(
|
|
configStore *markers.ProviderConfigStore,
|
|
installation *plugins.Installation,
|
|
capability *plugins.Capability,
|
|
providerID string,
|
|
) (markers.ProviderConfig, bool) {
|
|
if configStore == nil ||
|
|
installation == nil ||
|
|
capability == nil ||
|
|
installation.PluginID != "silo.theintrodb" ||
|
|
capability.ID != "introdb" {
|
|
return markers.ProviderConfig{}, false
|
|
}
|
|
if _, exists := configStore.Get(providerID); exists {
|
|
return markers.ProviderConfig{}, false
|
|
}
|
|
legacy, ok := configStore.Get("introdb")
|
|
if !ok {
|
|
return markers.ProviderConfig{}, false
|
|
}
|
|
legacy.Provider = providerID
|
|
return legacy, true
|
|
}
|
|
|
|
func copyLegacyIntroDBPluginConfig(
|
|
ctx context.Context,
|
|
runtimeConfigs markerPluginRuntimeConfigStore,
|
|
legacySettings markerLegacySettingsStore,
|
|
installation *plugins.Installation,
|
|
capability *plugins.Capability,
|
|
) error {
|
|
if runtimeConfigs == nil ||
|
|
legacySettings == nil ||
|
|
installation == nil ||
|
|
capability == nil ||
|
|
installation.PluginID != "silo.theintrodb" ||
|
|
capability.ID != "introdb" {
|
|
return nil
|
|
}
|
|
configs, err := runtimeConfigs.ListGlobalConfigs(ctx, installation.ID)
|
|
if err != nil {
|
|
return fmt.Errorf("list TheIntroDB plugin config: %w", err)
|
|
}
|
|
for _, config := range configs {
|
|
if config != nil && config.Key == "account" {
|
|
return nil
|
|
}
|
|
}
|
|
apiKey, err := legacySettings.Get(ctx, "introdb.api_key")
|
|
if err != nil {
|
|
return fmt.Errorf("load legacy introdb.api_key: %w", err)
|
|
}
|
|
if strings.TrimSpace(apiKey) == "" {
|
|
return nil
|
|
}
|
|
if err := runtimeConfigs.PutGlobalConfig(ctx, installation.ID, "account", map[string]any{
|
|
"api_key": strings.TrimSpace(apiKey),
|
|
}); err != nil {
|
|
return fmt.Errorf("copy legacy introdb.api_key to plugin config: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func markerCapabilityMetadata(descriptor *pluginv1.CapabilityDescriptor) map[string]any {
|
|
if descriptor == nil || descriptor.GetMetadata() == nil {
|
|
return nil
|
|
}
|
|
return descriptor.GetMetadata().AsMap()
|
|
}
|
|
|
|
func firstNonEmptyMarkerText(values ...string) string {
|
|
for _, value := range values {
|
|
if strings.TrimSpace(value) != "" {
|
|
return strings.TrimSpace(value)
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
// mapFolderTypeToMediaType maps silo's MediaFolder.Type values
|
|
// ("movies", "series", "mixed") to the SDK's MediaType values
|
|
// ("movie", "tv", "mixed"). Unknown values map to "mixed".
|
|
func mapFolderTypeToMediaType(t string) string {
|
|
switch t {
|
|
case "movies":
|
|
return "movie"
|
|
case "series":
|
|
return "tv"
|
|
default:
|
|
return "mixed"
|
|
}
|
|
}
|
|
|
|
type scopeResolver interface {
|
|
Resolve(ctx context.Context, input access.ResolveInput) (access.Scope, error)
|
|
}
|
|
|
|
type scopeEntitlementResolver struct {
|
|
resolver scopeResolver
|
|
}
|
|
|
|
func (r scopeEntitlementResolver) MaxPlaybackQuality(ctx context.Context, userID int, profileID string) (string, error) {
|
|
scope, err := r.resolveScope(ctx, userID, profileID)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
return scope.MaxPlaybackQuality, nil
|
|
}
|
|
|
|
// MaxContentRating implements mediarequests.ContentRatingResolver so request
|
|
// discovery honors the profile's parental rating ceiling.
|
|
func (r scopeEntitlementResolver) MaxContentRating(ctx context.Context, userID int, profileID string) (string, error) {
|
|
scope, err := r.resolveScope(ctx, userID, profileID)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
return scope.MaxContentRating, nil
|
|
}
|
|
|
|
func (r scopeEntitlementResolver) resolveScope(ctx context.Context, userID int, profileID string) (access.Scope, error) {
|
|
return r.resolver.Resolve(ctx, access.ResolveInput{
|
|
UserID: userID,
|
|
ProfileID: profileID,
|
|
SkipPINVerification: true,
|
|
})
|
|
}
|
|
|
|
// audiobooksSettingsAdapter bridges catalog.ServerSettingsRepo (which
|
|
// exposes Get) to the audiobooks.SettingsReader interface (which
|
|
// requires GetString). The two signatures are identical modulo name.
|
|
type audiobooksSettingsAdapter struct {
|
|
repo catalog.SettingsStore
|
|
}
|
|
|
|
func (a *audiobooksSettingsAdapter) GetString(ctx context.Context, key string) (string, error) {
|
|
return a.repo.Get(ctx, key)
|
|
}
|