* perf(jellycompat,sections): bound resume scan, batch leaf detail progress, widen section concurrency Three low-risk fixes from the section-fetch performance investigation (docs/superpowers/plans/2026-07-03-section-fetch-performance.md): - jellycompat: bound loadProgressPage at resumeScanMaxRows=300 so a single request never pages through more than that many in-progress rows. The cap is unconditional: it also covers the sparse-visible-set case (a heavy watcher whose recent rows are mostly dismissed/superseded, or a Series/Season-only request that matches no leaf in-progress row), where the page never fills and the loop would otherwise scan the entire history — previously an O(history) scan reaching tens of seconds. In the common case the loop exits far earlier, so the cap only bounds the pathological worst case; 300 leaves ample headroom to fill a ~20-item Continue Watching page. Beyond the cap the reported total is a clamped lower bound. Covered by TestLoadProgressPage_BoundsScanForSparseVisibleSet. - jellycompat: batch the leaf-item (movie/episode) progress lookup in GetItemDetailsByIDs via ListProgressWithCompletedHistory instead of a per-item GetProgressWithCompletedHistory (~100 sequential queries for a 50-item detail page). Series keep the per-item episode-rollup path (they own no progress row). Output is unchanged; a batch-lookup failure is now logged rather than silently dropping played state for the whole page. - sections: raise fetchAllMaxConcurrency 4 -> 6 to cut FetchAll wave count for large home layouts, staying within the default 20-conn pool. Part of the home/continue-watching latency work. AI-use disclosure: implemented with AI (Claude) assistance. * perf(jellycompat): keep Latest browse on the cross-library fast path under isPlayed /Items/Latest with isPlayed=false is the highest-frequency compat browse (~10.8k calls/day). The played overlay can't be pushed into SQL, so browse over-fetches and filters locally. The cross-library recently_added fast path (BrowseRecentlyAddedAcrossLibraries: one ~1ms index walk per library) was gated on Offset==0, so a heavy watcher who had already seen the newest items needed a 2nd chunk and fell through to BrowsePage — a whole-catalog MIN(first_seen_at) + GROUP BY HashAggregate over ~147k movies measured at ~755ms per call (0.8-1.6s observed end-to-end). Fetch the entire over-fetch budget (maxScannedRows) in a single merged fast-path walk instead of paging into BrowsePage, so the loop fills from one call. The clamp caveat (MaxLimit=1000 leaves a fall-through only for requestedLimit>200, off the Latest hot path) is documented inline. Part of the home/browse latency work. AI-use disclosure: implemented with AI (Claude) assistance. * fix(jellycompat): scope resume scan cap to resume path and bound the fast-path loop Addresses PR #292 review feedback: - Codex (P2): the resumeScanMaxRows cap was applied unconditionally in the general loop, which also paginates the completed (watched-items) list. Gate it on resumeFiltered so the completed path keeps exact TotalRecordCount and deep StartIndex pagination. Covered by TestLoadProgressPage_CompletedScanNotCapped. - CodeRabbit (Critical): the earlier raw-offset fast-path loop — the default Continue Watching shape and the sections-fallback route — had the same unbounded-scan bug and was not covered by the cap (the existing test forces EnableTotalRecordCount=true, routing around it). Bound it with the same resumeScanMaxRows guard. Covered by TestLoadProgressPage_BoundsFastPathScanForSparseVisibleSet. - CodeRabbit (Minor): tag the doc's fenced example blocks as text to satisfy markdownlint MD040. AI-use disclosure: implemented with AI (Claude) assistance. * perf(sections): cache shared user-agnostic home rails per access scope Home-screen rails that are identical for everyone who can see the same libraries (recently added, recently released, genre, trending on server, most watched, new to library, critically acclaimed, award winners, format showcase, seasonal, mood, trending discover, admin-curated lists, and library collections) were rebuilt from Postgres once per request, per user. Only the overlay on top of each row (watched flags, play position, presigned poster URLs) is actually per-user. Insert a process-global resolved-list cache at the FetchOne choke point in internal/sections. Each cacheable row is built once per access scope, held with a 15m TTL, and refreshed in the background 3m before expiry; singleflight collapses cold-miss stampedes into a single build. The per-user overlay still runs fresh in buildSectionsResponse, so no profile state is ever shared. Random and per-user rows (continue watching, next up, recommendations, hidden gems, forgotten favorites, activity feed, user collections) bypass the cache. The access-scope key captures every access boundary the fetch path enforces -- section identity (type + id + config hash) + item limit + accessible and disabled libraries + max content rating + excluded media types + name prefix + allowed-content-id allowlist -- and nothing per-user, so entries are safely shared. Empty membership is never cached (avoids freezing a transiently empty rail); background refreshes are bounded by a timeout. Scale (analytical, derived from the cache behavior -- not a measured latency): for the user-agnostic rows, Postgres section-query volume collapses from O(rows x concurrent requests) to O(rows x distinct access scopes) per 15m refresh window, because most users share a handful of access scopes. Illustrative -- 40 cacheable rows on a home screen, 1000 concurrent users falling into ~5 distinct access scopes: - before: ~40 x 1000 = ~40,000 section queries per wave of home loads - after: ~40 x 5 = ~200 builds per 15m window (plus one background refresh per row per scope), i.e. a warm home load runs zero section queries for these rows. That is a ~99% reduction in shared section-query load at that concurrency; the win grows with concurrency and shrinks as access-scope diversity rises. Design/plan doc added under docs/superpowers/plans/. * perf(jellycompat): serve per-library Latest via the cached recently-added section A jellyfin-compat per-library /Items/Latest rail is the same user-agnostic list as the native "recently added" library rail -- both order by mil.first_seen_at DESC. It was rebuilt on every request through directContentService.BrowseItems, missing the resolved-list cache entirely. Route per-library Latest for movies and series libraries through the native section fetch instead, so it reuses the shared cache. HandleLatest resolves the library's type once, and for a movies/series library builds a synthetic SectionRecentlyAdded with the same type + config + limit + access scope the native rail uses and calls FetchOne; the per-user overlay (favorites, progress, episode targets, presign) is extracted into buildLatestItemDTOs and shared by both the native and BrowseItems paths, so no overlay logic is duplicated. Cached *models.MediaItem values are read-only -- LocalizeItemModels deep-copies before any presign mutation. To let the two surfaces share one entry, resolvedListCacheKey no longer includes the arbitrary section ID: every cacheable section type derives its membership from type + config + limit + scope, never from its own ID (audited all 14 cacheable types plus the library-collection path; the sole s.ID read lives in the non-cacheable user-collection branch). A native recently-added rail and the compat Latest for the same library + scope now collapse to ONE cache entry, built once and reused. Access-scope isolation is unchanged -- the removed ID never carried access information, and every access boundary (libraries, rating cap, excluded types, content allow-list, name prefix) still keys the entry. Guardrails: the native path is restricted to movies and series libraries; every other library type (ebook, music, manga, mixed) is ignored and keeps its exact BrowseItems behavior -- important because an unfiltered recently-added fetch would otherwise surface non-video items to Jellyfin clients that only expect video. Deeper pages, played-filter and backdrop-required requests, a client asking for a type other than the library's own, and any FetchOne error also fall back to BrowseItems. Chosen over an alternative that gave the synthetic section a deterministic ID (which kept two separate cache entries): both returned identical data with similar complexity, so the shared-entry design won. * fix(sections,jellycompat): post-review fixes for the shared-list cache and Latest path Consolidates fixes from the branch's adversarial review and PR #292 review comments into one commit: - Latest fast path: fall back to BrowseItems when a request carries a genre, name-prefix, or person filter (the synthetic recently-added section cannot express these, so serving it unfiltered would return a wrong, broader set). Eligibility is decided by latestFastPathEligible and covered by a test. - Clamp the /Items/Latest page size to compatBrowseMaxLimit before building the section, matching the BrowseItems fallback, so a large client Limit can't drive an oversized recently-added fetch or explode the shared cache key with unbounded ItemLimit values. - Evict expired entries from the process-global resolvedListCache: resolvedListSet sweeps expired keys at most once per minute, bounding the map to scopes seen within one TTL window. Covered by TestResolvedListCacheEvictsExpiredEntries. - Log a short digest of the cache key (resolvedListLogKey) instead of the raw key in the background-refresh panic/error paths, since the key embeds user-controlled access-scope fields such as NamePrefix. Skipped review comments (verified already fixed or stale against current code): the resume fast-path scan bound and watched-items cap (04d2e795) and the docs fence-language tags (already addressed). Build, vet, and go test -race pass for internal/sections and internal/jellycompat. * perf(plugins): cache plugin installations in-memory, invalidated on lifecycle change ## Problem Every poster/image on a warm home rail re-read plugin_installations from Postgres to answer "is this plugin enabled?" and to acquire the plugin client (Source A: metadata chain buildProviders enabled-check; Source B: ensureClient -> loadInstallation). Plugin-resolved image URLs are never URL-cached, so the plugin source and the DB read behind it fired again on every identical warm request; 100% of images in the target library are plugin-backed. ## Solution - Guarded in-memory installation cache (map[int]*Installation + RWMutex) in plugins.Service. loadInstallation reads through it; the requireEnabled gate stays after the cache read so ErrInstallationDisabled semantics are unchanged. invalidateInstallationCache clears it and is self-registered as a lifecycle hook, so Service.OnLifecycleChange wipes it on install/enable/disable/update/ uninstall. - A generation counter closes an invalidate-vs-repopulate race: captured before installations.GetByID and re-checked under the write lock, so a row fetched before a lifecycle mutation is never written into a freshly cleared cache (would otherwise resurrect a just-disabled plugin). - Route the metadata chain enabled-check through the same cache via a structural InstallationEnabledChecker interface (nil-safe: falls back to the pool query when no checker is injected), wired in cmd/silo/main.go. ## Post-review fix (auto-update reliability blocker) AutoUpdateService mutated installations (new InstallPath, old dir deleted) on the default auto update policy without firing OnLifecycleChange, leaving the cache stale and breaking plugins with "stored plugin manifest mismatch" until restart. It now takes an onChange callback wired to Service.OnLifecycleChange and fires it once per Check run that mutated a row. ## Verification go build/vet, go test ./internal/plugins/... ./internal/metadata/... (-race). Tests: cache hit/invalidation, racing-invalidation guard, IsInstallationEnabled, auto-update fires onChange. ## AI-use disclosure Implemented with AI assistance (Claude). * perf(jellycompat): batch per-item presign, and enrich series on the cached Latest path ## Problem List rails presigned each item's poster/backdrop/logo/still image individually (~160 singular resolver calls for a 40-item page where 4 batched calls suffice), and ItemsHandler carried a near-verbatim duplicate of the batch presigner. ## Solution (batching) Promote the batch presigner to a shared package-level presignCompatListItems (presign_list.go) with a generic collectImagePaths[T]; convert the per-item loops (cached home/Latest rail, favorites, batch loaders, userdata favorites) to one batched PresignImageURLsWithExpiry per image type per page; batch the season/episode collections; delete the three duplicate presign helpers. URL output is unchanged (verified byte-for-byte). ## Post-review fix (series Latest data-parity regression) The native cached Latest fast path built items via compatListItemsFromModels + buildLatestItemDTOs and never ran the series watch-state rollup, so a series library's Latest lost Played / UnplayedItemCount and page 1 disagreed with the BrowseItems fallback. enrichSeriesUserData is promoted to the ContentService interface and called on the native path (reused, not duplicated). ## Verification go build/vet, go test ./internal/jellycompat/... ./internal/catalog/... Tests: bounded presign invocation counts + per-item URL mapping; series rollup populated on the native Latest path. ## AI-use disclosure Implemented with AI assistance (Claude). * perf(sections): gate personalized rails out of the shared cache; widen refresh lead ## Problem 1. The shared home-rail cache whitelisted custom_filter/genre sections by TYPE alone, but those route through fetchFiltered -> ParseQueryDefinition and can carry personalized (per-profile) rules/sorts (watched, favorited, in_watchlist, in_progress, last_watched; sorts progress/date_viewed/plays). Their membership is per-profile yet the cache key excludes userID/profileID, so a personalized rail built for one profile was served to others in the same access scope for up to 15m -- a cross-profile watchlist/watch-state leak. 2. The background-refresh lead was tuned so steady traffic is served a warm entry from a longer soft window. ## Solution - Add QueryDefinition.IsPersonalized() (reusing the existing QueryFieldRequiresProfile/QuerySortRequiresProfile helpers). isCacheableSectionType now parses the section QueryDefinition and refuses to cache custom_filter/genre when personalized; non-personalized definitions stay cacheable. Seasonal/mood/trending build their definitions server-side and stay unconditionally cacheable. - resolvedListRefreshLead 3m -> 10m (soft threshold builtAt+5min instead of builtAt+12min). ## Verification go build/vet, go test ./internal/sections/... ./internal/catalog/... (-race). Test: personalized custom_filter/genre not cacheable; non-personalized are. ## AI-use disclosure Implemented with AI assistance (Claude). * fix(sections,metadata): post-review fixes for shared cache and plugin chain staleness Addresses three review findings on PR #292: - sections: canonicalize section config JSON before hashing so configs differing only in whitespace/field order share a cache entry (native + jellycompat rail sharing). Added TestHashSectionConfigCanonicalizes. - metadata: invalidate the resolved-chain cache on plugin lifecycle changes; the installation-enabled check already reads the invalidated plugin cache, but resolveChainCached could serve a stale provider chain for up to chainCacheTTL after a provider's availability changed. - jellycompat: move ctx to the first parameter of presignCompatListItems for consistency with the other presign helpers. Skipped the episode-image presign batching nitpick: the resolver already dedupes+singleflights, so it is a Minor perf-only item not worth the two-pass refactor risk in this pass. * fix(sections,jellycompat): harden shared rail cache and Latest fast path per review Addresses the eight findings from the deep review of this PR: - Detach the blocking cold-miss rebuild from the singleflight leader's request context (context.WithoutCancel + the shared 30s build timeout) so one client disconnect no longer fails every collapsed waiter and leaves the entry uncached. - Stop client-controlled values minting unbounded cache entries: the compat Latest fast path now always fetches a fixed 100-row budget and slices to the requested limit (one entry per scope+library instead of one per Limit value), and an unrecognized MaxOfficialRating string disqualifies the fast path instead of entering the global cache key. - Add release_date to the sections item projection/scan so movies served via the Latest fast path keep PremiereDate (Jellyfin default-set field) in parity with the BrowseItems fallback. - Fall back to per-item progress lookups when the batched leaf progress query fails, restoring one-item-at-a-time degradation instead of blanking played state for the whole page. - Derive cache eligibility from a single source of truth: fetchSection and isCacheableSectionType now share the userAgnosticSectionFetcher table, whose no-userID/profileID signature makes a fetcher drop out of the cacheable set at compile time if it ever gains per-profile inputs. - Decide Latest fast-path eligibility off the actual browse params the fallback would receive, so any filter later added to buildBrowseParams automatically disqualifies the cached path; share one compatDefaultBrowseLimit constant between both paths. - Extract AccessFilter.WriteAccessScopeCacheKey as the shared, security- critical serializer for all access-scoped caches (resolved-list, editorial candidates, audiobook groups); the editorial key now captures ExcludedMediaTypes, which its loaders already applied in SQL. - Strip leaked agent-transcript markup from the section-fetch plan doc. go build ./..., go vet, gofmt clean; go test -race on internal/sections, internal/catalog, internal/jellycompat passes (TestBeginWebOperation* failures are the known pre-existing flakes). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Quick <31828688+Quick104@users.noreply.github.com> Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
2974 lines
107 KiB
Go
2974 lines
107 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/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/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/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"
|
|
"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/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/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/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) slog.Handler {
|
|
opts := &slog.HandlerOptions{Level: level}
|
|
if strings.EqualFold(format, "json") {
|
|
return slog.NewJSONHandler(os.Stderr, opts)
|
|
}
|
|
return slog.NewTextHandler(os.Stderr, opts)
|
|
}
|
|
|
|
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)
|
|
}
|
|
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.Warn("ensure operational log partitions; continuing in degraded mode", "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.Warn("postgres auto-tuning disabled", "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.Warn("postgres auto-tuning setting failed",
|
|
"name", failure.Name,
|
|
"value", failure.Value,
|
|
"error", failure.Err,
|
|
)
|
|
}
|
|
if err != nil {
|
|
slog.Warn("postgres auto-tuning failed",
|
|
"error", err,
|
|
"applied", result.Applied,
|
|
"failures", len(result.Failures),
|
|
)
|
|
return
|
|
}
|
|
|
|
slog.Info("postgres auto-tuning applied",
|
|
"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.Warn("postgres restart required to finish applying auto-tuned settings",
|
|
"settings", strings.Join(result.RestartRequired, ","),
|
|
)
|
|
}
|
|
if len(result.Reset) > 0 {
|
|
slog.Info("postgres auto-tuning reset stale settings",
|
|
"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.
|
|
func runCredentialBackfills(ctx context.Context, pool *pgxpool.Pool, cipher *secret.Cipher, settings *catalog.EncryptedSettingsRepo) {
|
|
settingsN, err := settings.BackfillSensitiveSettings(ctx)
|
|
if err != nil {
|
|
slog.Error("secret backfill: sensitive settings", "error", err)
|
|
}
|
|
columnsN, err := secret.BackfillColumns(ctx, pool, cipher, secret.ColumnBackfillTargets())
|
|
if err != nil {
|
|
slog.Error("secret backfill: credential columns", "error", err)
|
|
}
|
|
historyServersN, err := historyimport.NewRepository(pool, cipher).BackfillSessionServerSecrets(ctx)
|
|
if err != nil {
|
|
slog.Error("secret backfill: history import session server credentials", "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.Error("secret backfill: arr api keys", "error", err)
|
|
}
|
|
if total := settingsN + columnsN + historyServersN + arrN; total > 0 {
|
|
slog.Info("secret backfill: encrypted plaintext credentials at rest",
|
|
"settings", settingsN, "columns", columnsN, "history_session_servers", historyServersN, "arr_keys", arrN, "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")
|
|
flag.Parse()
|
|
|
|
ctx := context.Background()
|
|
|
|
// 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 *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))
|
|
baseHandler := buildBaseHandler(cfg.Server.LogFormat, logLevelVar)
|
|
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 "+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))
|
|
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
|
|
}
|
|
|
|
// 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() }()
|
|
}
|
|
|
|
deps := api.Dependencies{
|
|
Config: cfg,
|
|
LiveConfig: configWatcher.Config,
|
|
OnConfigChange: configWatcher.OnChange,
|
|
BootstrapSensitiveConfigured: bootstrapSensitiveConfigured,
|
|
BootstrapSensitiveValues: bootstrapSensitiveValues,
|
|
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, _, _ string) {
|
|
// 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()
|
|
},
|
|
}
|
|
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
|
|
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)
|
|
}
|
|
watchProviderService = watchsync.NewService(
|
|
watchsync.NewPostgresRepository(deps.DB, deps.SecretCipher),
|
|
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)
|
|
s.SetSearchIndexProvider(activeCatalogSearchProvider)
|
|
configWatcher.OnChange(func(_, updated *config.Config) {
|
|
s.SetWorkers(updated.Scanner.Workers)
|
|
})
|
|
s.SetLiteraryWorkLinker(literaryWorkService)
|
|
deps.Scanner = s
|
|
deps.ProbeEnsurer = scanner.NewPlaybackProbeEnsurer(fileRepo, ffprobePath, 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)
|
|
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 {
|
|
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 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.Warn("reload marker plugin providers failed", "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)
|
|
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.Warn("failed to reload plugin image resolvers", "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)
|
|
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)
|
|
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,
|
|
},
|
|
)
|
|
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.Warn("failed to clean legacy series group queue rows", "error", err)
|
|
} else if cleaned > 0 {
|
|
slog.Info("cleaned legacy series group queue rows", "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.Warn("failed to seed metadata queues", "error", err)
|
|
return
|
|
}
|
|
seedMovieQueue := func(folderID int) {
|
|
if movieQueueRepo == nil {
|
|
return
|
|
}
|
|
if err := movieQueueRepo.SyncForFolder(ctx, folderID); err != nil {
|
|
slog.Warn("failed to seed movie match queue", "folder_id", folderID, "error", err)
|
|
}
|
|
}
|
|
seedSeriesQueue := func(folderID int) {
|
|
if seriesQueueRepo == nil {
|
|
return
|
|
}
|
|
if err := seriesQueueRepo.SyncForFolder(ctx, folderID); err != nil {
|
|
slog.Warn("failed to seed series root queue", "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.Info("deferred init: metadata match queues seeded", "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()
|
|
}
|
|
|
|
// 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 {
|
|
notificationScopes := access.NewResolver(
|
|
auth.NewUserRepository(deps.DB),
|
|
userStoreProvider,
|
|
access.NewProfileTokenService(cfg.Auth.JWTSecret, 0),
|
|
)
|
|
notificationSystem = notifications.NewSystem(
|
|
deps.DB,
|
|
settingsRepo,
|
|
userStoreProvider,
|
|
notificationScopes,
|
|
auth.NewUserRepository(deps.DB),
|
|
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
|
|
if deps.DB != nil {
|
|
userRepo := auth.NewUserRepository(deps.DB)
|
|
sessionMgr.SetLimitProvider(func(ctx context.Context, userID int) (playback.SessionLimits, error) {
|
|
user, err := userRepo.GetByID(ctx, userID)
|
|
if err != nil {
|
|
return playback.SessionLimits{}, err
|
|
}
|
|
return playback.SessionLimits{
|
|
MaxStreams: user.MaxStreams,
|
|
MaxTranscodes: user.MaxTranscodes,
|
|
}, nil
|
|
})
|
|
}
|
|
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 err := watchProviderService.SweepOpenScrobbles(ctx); err != nil {
|
|
slog.Warn("failed to sweep open watch provider scrobbles", "error", err)
|
|
}
|
|
})
|
|
}
|
|
// Auto-remove fully-watched items from the watchlist (standalone behavior,
|
|
// default-on per profile), propagating removals to connected providers.
|
|
if itemRepo != nil && episodeRepo != nil && userStoreProvider != nil {
|
|
maintainer := watchlist.NewMaintainer(userStoreProvider, itemRepo, episodeRepo)
|
|
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
|
|
|
|
// 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)
|
|
}
|
|
// Reload trusted proxies
|
|
cidrs, loadErr := clientip.LoadTrustedCIDRs(context.Background(), settingsRepo)
|
|
if loadErr != nil {
|
|
slog.Warn("clientip config reload failed", "error", loadErr)
|
|
} else {
|
|
ipResolver.UpdateTrustedCIDRs(cidrs)
|
|
}
|
|
}
|
|
})
|
|
|
|
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)
|
|
}
|
|
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)
|
|
})
|
|
}
|
|
}
|
|
|
|
// 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)))
|
|
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))
|
|
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 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 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)
|
|
if userStoreProvider != nil {
|
|
reconcileResolver := access.NewResolver(
|
|
auth.NewUserRepository(deps.DB),
|
|
userStoreProvider,
|
|
access.NewProfileTokenService(cfg.Auth.JWTSecret, 0),
|
|
)
|
|
requestReconcileSvc.SetEntitlementResolver(mediarequests.NewAccessEntitlements(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))
|
|
}
|
|
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))
|
|
}
|
|
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)
|
|
}
|
|
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),
|
|
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
|
|
|
|
// White-label branding: one service shared by the API (public read + admin
|
|
// upload) and the frontend handler (index.html title, favicon, manifest).
|
|
// 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.
|
|
if settingsRepo != nil {
|
|
var brandingStore branding.AssetStore
|
|
if deps.S3Public != nil {
|
|
brandingStore = deps.S3Public
|
|
}
|
|
brandingSvc := branding.NewService(settingsRepo, brandingStore)
|
|
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),
|
|
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,
|
|
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.DetailSvc = detailSvc
|
|
compatDeps.FolderRepo = folderRepo
|
|
compatDeps.SessionMgr = sessionMgr
|
|
compatDeps.UserStoreProvider = userStoreProvider
|
|
compatDeps.WatchCompletionObserver = deps.WatchCompletionObserver
|
|
compatDeps.SettingsRepo = settingsRepo
|
|
compatDeps.PersonRepo = personRepo
|
|
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 {
|
|
compatDeps.AccessFilterFn = jellycompat.NewScopeAccessFilter(access.NewResolver(
|
|
userRepo,
|
|
userStoreProvider,
|
|
nil, // profile tokens unused: compat login already verifies PINs
|
|
))
|
|
}
|
|
}
|
|
|
|
compat := jellycompat.NewServerWithDependencies(compatDeps)
|
|
compatServer = compat
|
|
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
|
|
}
|
|
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.Warn("plugin image resolver capability has no valid schemes",
|
|
"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.Warn("skipping legacy metadata image resolver with invalid scheme",
|
|
"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.Info("reloaded plugin image resolvers", "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
|
|
}
|
|
}
|
|
|
|
type markerPluginCapabilityStore interface {
|
|
ListEnabled(ctx context.Context) ([]*plugins.Installation, error)
|
|
ListCapabilities(ctx context.Context, installationID int) ([]*plugins.Capability, error)
|
|
}
|
|
|
|
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
|
|
}
|
|
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"
|
|
}
|
|
}
|
|
|
|
// 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)
|
|
}
|