Files
silo-server/internal/plugins/service.go
T
91e1164090 feat(metadata): local NFO metadata and sidecar artwork (builtin chain provider) (#390)
* feat(metadata): register builtin NFO provider and broaden parsing

Phases A and B of the #216 local-NFO work, implemented test-first.

Registration & hint-first identity (Phase A):
- Migration seeds a reserved kind='builtin' silo.builtin installation
  and an 'nfo' metadata capability (default_enabled=false, priority 1
  for movie/series) with a partial unique index and documented Down.
- In-process builtin provider registry (internal/metadata/builtin.go);
  buildProviders returns the registered provider for builtin rows.
- Guard rails keep the reserved row out of every plugin surface (user
  plugin-settings, installations list, image resolvers, preload,
  auto-update, store Delete, mutation handlers -> 409); silo.builtin is
  a reserved manifest id.
- Startup sync materializes legacy content_level='' chains per level,
  then appends builtin capabilities disabled via
  AppendProviderToAllChains (idempotent); resolveEnabledProvidersBy
  priority now respects default_enabled=false.
- NFO uniqueids seed the trusted-hint machinery via IdentityHintProvider
  with per-mode conflict policy (stored IDs win on scheduled refresh,
  NFO wins on manual refresh, Identify skips NFO); ID-less candidates
  are excluded from provider-priority tie-breaks and nfo never counts
  as corroboration.
- Web chain-editor empty-state gate is now server-derived so builtin
  providers are reachable on plugin-less servers.

Parser breadth & sidecar hardening (Phase B):
- Parser covers the practical Kodi/Jellyfin field set for <movie> and
  <tvshow>: original title, tagline, runtime, dates, content rating,
  genres/studios/countries/tags, multi-source ratings with scale
  normalization, cast with roles/order, director/credits. Empty
  collections stay nil so merge early-returns apply.
- findNFO parses candidates and falls through on read/parse failure or
  root-type mismatch, so a stray movie.nfo cannot shadow tvshow.nfo;
  GetMetadata gains the same ContentType guard Search has.
- New FieldReleaseDates lock gates Year/ReleaseDate/First+LastAirDate
  in merge (Go) and the edit-metadata dialog (web), closing the gap
  where a manual refresh re-applied NFO dates over admin corrections.
- Merge-contract tests pin NFO fill semantics, genres whole-list
  first-provider-wins, and NFO edits propagating on manual refresh only.
- Docs: new admin wiki page (supported fields, merge semantics,
  naming-supplies-structure contract), index bullet, sidecar wording
  revision, v1-scope feature-detection note.

Zero behavior change while the provider is disabled (default); pinned
by CI-mode and DB-gated test suites.

Part of #216

AI-use disclosure: implemented with Claude Code (Fable 5) via
spec-driven TDD and agent-assisted implementation.

* feat(metadata): ingest local sidecar artwork and read series-depth NFO

Phases C and D of the #216 local-NFO work, implemented test-first, plus
the mixed-library use-case pins. Together these deliver the headline
case: a series absent from every remote database (e.g. a fitness
library) scans into a fully presented show -> named seasons -> titled
episodes tree from NFO files and sidecar art alone.

Local sidecar artwork through the S3 image cache (Phase C):
- The NFO provider implements ImageProvider: poster/backdrop/logo
  sidecar discovery with a fixed precedence map, symlink/non-regular
  rejection, an 8 MiB cap, and file:// source URLs at rating 0. Generic
  filenames apply only via the sidecar search paths, so a shared
  folder.jpg in a flat multi-movie directory applies to none.
- file:// becomes a live local source scheme: routed into *_source_path
  (never *_path), accepted by every image enqueue gate, attributed as
  provider "local", excluded from cached-path detection.
- The image-cache processor caches local files with lexical-on-logical
  confinement to the library roots, open-handle reads with re-checks,
  the same variant widths as remote art, and stable (7-day) failure
  classification. Keys land under
  local/{contentType}/{contentID}/{hash8}/{imageType}; superseded
  prefixes are cleaned on re-cache and item deletion.
- applyIfBetter gains a local exemption so rating-0 local art can fill
  matched items without being stickily displaced; ImageRequest carries
  additive sidecar path context.

Series depth (Phase D):
- SeasonsRequest/EpisodesRequest carry additive local path context
  (series roots, per-season directories, per-episode file paths),
  derived from naming at match time and reconstructed on refresh.
- season.nfo supplies season name/plot; NFO season numbers are advisory
  (directory-derived number wins with a Warn - naming owns structure).
  <episodedetails> gains aired/runtime/ratings; <basename>.nfo titles
  episodes and <basename>-thumb.ext supplies thumbs; filename SxxEyy
  wins over NFO numbers.
- Episode NFOs work without a season.nfo (provider seasons unioned with
  on-disk seasons); SynthesizeFallbackEpisodes always runs after persist
  so NFO-less episodes keep synthesized rows. Season/episode file:// art
  rides the Phase C pipeline unchanged.
- Migration adds season:1/episode:1 to the builtin NFO capability's
  default_priority (still default_enabled=false).

Mixed sports-library use case (tests only, no product change):
- Pins the classification contract for one library holding movie-shaped
  and show-shaped content (WWE PPV events as movies next to a "WWE
  SmackDown" show, NASCAR/F1/FIFA with partial TVDB/TMDB data): naming
  decides movie-vs-series per file before any provider runs; the NFO
  supplies metadata/identity but never flips type (ContentType guard);
  the per-root Type override is the correction path.
- NFO-driven type classification at scan time is recorded as an explicit
  deferred open question.

Part of #216

AI-use disclosure: implemented with Claude Code (Fable 5) via
spec-driven TDD and agent-assisted implementation.

* docs(metadata): document local NFO metadata architecture

Add a single as-built architecture page
(docs/architecture/local-nfo-metadata.md) for the #216 local-NFO
feature: the builtin registration model, hint-first identity semantics,
the file:// -> S3 artwork pipeline and its deployment constraint, series
depth, the mixed-library classification contract, and known limitations.

This replaces the working implementation plan, the per-phase specs, and
the narrow sidecar-artwork note, which were planning drafts and are left
untracked; admin-facing behavior remains in the wiki.

Part of #216

AI-use disclosure: planned, drafted, and consolidated with Claude Code
(Fable 5) using multi-agent exploration and adversarial review.

* fix(metadata): address PR review findings on NFO builtin provider

Fold in the valid, low-risk fixes surfaced by automated review on #390:

- imagecache: extract validateCacheRequest so CacheBytes (the local
  sidecar season/episode path) enforces the same episode-requires-season
  guard as Cache, preventing distinct episodes' art from colliding under
  one S3 key.
- image_cache_processor: close the sidecar symlink-swap window by
  rejecting the opened handle unless os.SameFile matches the Lstat'd
  file, so a leaf swapped to a symlink can't pull an out-of-root target
  into the public cache.
- plugins: guard the reserved builtin installation row in the store's
  Update, matching Delete, so its version/enabled/capabilities can never
  be rewritten even if a mutation slips past the HTTP layer.
- cmd/silo: bound SyncBuiltinProviderChains with a 30s timeout so a stuck
  DB round-trip fails fast at startup instead of hanging.
- metadata: panic instead of silently no-op'ing on an invalid
  RegisterBuiltinProvider call (init-time programmer error).
- docs: correct the media-folder-and-naming NFO paragraph to state
  season/episode NFOs and sidecar artwork are actively read.

---------

Co-authored-by: Quick104 <31828688+Quick104@users.noreply.github.com>
2026-07-16 17:55:36 -04:00

947 lines
30 KiB
Go

package plugins
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"log/slog"
"os"
"path/filepath"
"runtime/debug"
"slices"
"strconv"
"strings"
"sync"
"sync/atomic"
"golang.org/x/sync/singleflight"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/structpb"
pluginv1 "github.com/Silo-Server/silo-plugin-sdk/pkg/pluginproto/silo/plugin/v1"
"github.com/Silo-Server/silo-server/internal/pluginhost"
)
type pluginClient interface {
Manifest() *pluginv1.PluginManifest
MetadataProvider(capabilityID string) (*pluginhost.MetadataProviderClient, error)
ImageResolver(capabilityID string) (*pluginhost.ImageResolverClient, error)
MarkerProvider(capabilityID string) (*pluginhost.MarkerProviderClient, error)
MediaAnalyzer(capabilityID string) (*pluginhost.MediaAnalyzerClient, error)
ScheduledTask(capabilityID string) (*pluginhost.ScheduledTaskClient, error)
ScanSource(capabilityID string) (*pluginhost.ScanSourceClient, error)
RequestRouter(capabilityID string) (*pluginhost.RequestRouterClient, error)
EventConsumer(capabilityID string) (*pluginhost.EventConsumerClient, error)
AuthProvider(capabilityID string) (*pluginhost.AuthProviderClient, error)
HTTPRoutes(capabilityID string) (*pluginhost.HTTPRoutesClient, error)
}
type Host interface {
Start(ctx context.Context, req pluginhost.StartRequest) (pluginClient, error)
Client(installationID int) (pluginClient, error)
Stop(installationID int) error
Shutdown(ctx context.Context) error
}
type serviceInstallationStore interface {
archiveStore
GetByID(ctx context.Context, id int) (*Installation, error)
List(ctx context.Context) ([]*Installation, error)
ListEnabled(ctx context.Context) ([]*Installation, error)
ListByPluginID(ctx context.Context, pluginID string) ([]*Installation, error)
Update(ctx context.Context, id int, input UpdateInstallationInput) error
ListCapabilities(ctx context.Context, installationID int) ([]*Capability, error)
}
type serviceConfigStore interface {
ListGlobalConfigs(ctx context.Context, installationID int) ([]*RuntimeConfig, error)
PutGlobalConfig(ctx context.Context, installationID int, key string, value map[string]any) error
}
type Service struct {
repositories *RepositoryStore
installations serviceInstallationStore
configs serviceConfigStore
catalog *CatalogService
installer *Installer
archiveCache *ArchiveCache
host Host
testConfigSeq atomic.Int64
dispatcher *EventDispatcher
lifecycleMu sync.RWMutex
lifecycleHooks []func(context.Context)
launchGroup singleflight.Group
// installationCache memoizes plugin_installations rows keyed by ID so the
// hot plugin-RPC path (ensureClient -> loadInstallation) and the metadata
// chain enabled-check answer from memory instead of a per-call DB read. It
// is wiped wholesale by invalidateInstallationCache, registered as a
// lifecycle hook, so a cached row is at most one lifecycle event stale.
//
// installationCacheGen guards the read-through against an invalidate that
// races an in-flight GetByID: the generation is captured before the store
// read and re-checked under the write lock, so a row fetched before a
// lifecycle mutation is never written back into a freshly-cleared cache.
installationCacheMu sync.RWMutex
installationCache map[int]*Installation
installationCacheGen uint64
}
// SetEventDispatcher wires the EventDispatcher into the Service. The
// dispatcher reference is retained so future hooks can act on lifecycle
// changes; the current dispatcher implementation is fully driven by
// per-event store reads and needs no notification on install/enable/disable.
func (s *Service) SetEventDispatcher(d *EventDispatcher) { s.dispatcher = d }
// AddLifecycleHook registers a callback invoked after plugin install, enable,
// disable, uninstall, or preload lifecycle changes.
func (s *Service) AddLifecycleHook(hook func(context.Context)) {
if s == nil || hook == nil {
return
}
s.lifecycleMu.Lock()
s.lifecycleHooks = append(s.lifecycleHooks, hook)
s.lifecycleMu.Unlock()
}
// OnLifecycleChange is invoked by API handlers and the installer after every
// plugin lifecycle mutation. Hooks should be best-effort and log their own
// errors so plugin admin operations are not failed by secondary cache refreshes.
func (s *Service) OnLifecycleChange(ctx context.Context) {
if s == nil {
return
}
s.lifecycleMu.RLock()
hooks := append([]func(context.Context){}, s.lifecycleHooks...)
s.lifecycleMu.RUnlock()
for _, hook := range hooks {
func(hook func(context.Context)) {
defer func() {
if recovered := recover(); recovered != nil {
slog.ErrorContext(ctx,
"plugin lifecycle hook panicked; continuing", "component", "plugins",
"panic", recovered,
"stack", string(debug.Stack()),
)
}
}()
hook(ctx)
}(hook)
}
}
func NewService(
repositories *RepositoryStore,
installations *InstallationStore,
configs *RuntimeConfigStore,
catalog *CatalogService,
installer *Installer,
host Host,
) *Service {
svc := &Service{
repositories: repositories,
installations: installations,
configs: configs,
catalog: catalog,
installer: installer,
archiveCache: NewArchiveCache(installations),
host: host,
}
// Self-register the installation-cache invalidation so it can never be
// silently forgotten by a new caller: every OnLifecycleChange (install /
// enable / disable / update / uninstall) wipes the cache, keeping the
// memoized rows correct without any external wiring.
svc.AddLifecycleHook(func(context.Context) { svc.invalidateInstallationCache() })
return svc
}
func (s *Service) FetchCatalog(ctx context.Context) ([]CatalogEntry, error) {
entries, err := s.catalog.Fetch(ctx)
if err != nil {
return nil, err
}
return catalogEntriesForDiscovery(entries), nil
}
func catalogEntriesForDiscovery(entries []CatalogEntry) []CatalogEntry {
selected := make(map[string]CatalogEntry, len(entries))
for _, entry := range entries {
if entry.Manifest == nil {
continue
}
pluginID := entry.Manifest.GetPluginId()
if isApprovedCommunityPlugin(pluginID) && entry.SourceKind != RepositorySourceApprovedCommunity {
continue
}
existing, ok := selected[pluginID]
if !ok || catalogEntryPreferredForDiscovery(entry, existing) {
selected[pluginID] = entry
}
}
result := make([]CatalogEntry, 0, len(selected))
for _, entry := range selected {
result = append(result, entry)
}
slices.SortFunc(result, func(left, right CatalogEntry) int {
return strings.Compare(left.Manifest.GetPluginId(), right.Manifest.GetPluginId())
})
return result
}
func catalogEntryPreferredForDiscovery(candidate, current CatalogEntry) bool {
candidatePrecedence := repositorySourcePrecedence(candidate.SourceKind)
currentPrecedence := repositorySourcePrecedence(current.SourceKind)
if candidatePrecedence != currentPrecedence {
return candidatePrecedence < currentPrecedence
}
if candidate.RepositoryID != current.RepositoryID {
return candidate.RepositoryID < current.RepositoryID
}
return compareVersions(candidate.Manifest.GetVersion(), current.Manifest.GetVersion()) > 0
}
func repositorySourcePrecedence(sourceKind string) int {
switch sourceKind {
case RepositorySourceSilo:
return 0
case RepositorySourceApprovedCommunity:
return 1
default:
return 2
}
}
func (s *Service) InstallLocal(ctx context.Context, req InstallArchiveRequest) (*InstallResult, error) {
if req.ArchivePath == "" {
return nil, fmt.Errorf("archive path is required")
}
data, err := os.ReadFile(req.ArchivePath)
if err != nil {
return nil, fmt.Errorf("read archive %q: %w", req.ArchivePath, err)
}
_, _, manifest, err := openPluginArchive(data)
if err != nil {
return nil, err
}
existing, err := s.existingInstallationByPluginID(ctx, manifest.GetPluginId())
if err != nil {
return nil, err
}
var result *InstallResult
if existing != nil {
if err := s.stopInstallationIfRunning(existing); err != nil {
return nil, err
}
result, err = s.installer.ReplaceLocal(ctx, existing, req)
} else {
result, err = s.installer.InstallLocal(ctx, req)
}
if err != nil {
return nil, err
}
s.OnLifecycleChange(ctx)
return result, nil
}
func (s *Service) InstallRemote(ctx context.Context, req InstallArchiveRequest) (*InstallResult, error) {
result, err := s.installer.InstallRemote(ctx, req)
if err != nil {
return nil, err
}
s.OnLifecycleChange(ctx)
return result, nil
}
func (s *Service) InstallCatalog(ctx context.Context, req InstallCatalogRequest) (*InstallResult, error) {
target, err := s.catalog.ResolveInstall(ctx, req)
if err != nil {
return nil, err
}
repositoryID := target.RepositoryID
existing, err := s.existingInstallationByPluginID(ctx, req.PluginID)
if err != nil {
return nil, err
}
var result *InstallResult
if target.LegacyArchive {
archiveReq := InstallArchiveRequest{
ArchiveURL: target.ArchiveURL,
RepositoryID: &repositoryID,
}
if existing == nil {
result, err = s.installer.InstallRemote(ctx, archiveReq)
} else {
if err = s.stopInstallationIfRunning(existing); err != nil {
return nil, err
}
result, err = s.installer.ReplaceRemote(ctx, existing, archiveReq)
}
} else {
binaryReq := InstallBinaryRequest{
BinaryURL: target.ArchiveURL,
Checksum: target.Checksum,
RepositoryID: &repositoryID,
}
if existing == nil {
result, err = s.installer.InstallBinary(ctx, binaryReq)
} else {
if err = s.stopInstallationIfRunning(existing); err != nil {
return nil, err
}
result, err = s.installer.ReplaceBinary(ctx, existing, binaryReq)
}
}
if err != nil {
return nil, err
}
s.OnLifecycleChange(ctx)
return result, nil
}
// UpdateToAvailableVersion updates a plugin to its available_version.
// Returns the updated installation after the update completes.
func (s *Service) UpdateToAvailableVersion(ctx context.Context, installationID int) (*Installation, error) {
installation, err := s.installations.GetByID(ctx, installationID)
if err != nil {
return nil, err
}
if installation.AvailableVersion == nil || *installation.AvailableVersion == "" {
return nil, fmt.Errorf("no update available for plugin %q", installation.PluginID)
}
targetVersion := *installation.AvailableVersion
if installation.RepositoryID == nil || *installation.RepositoryID == 0 {
return nil, fmt.Errorf("plugin %q has no repository_id, cannot update from catalog", installation.PluginID)
}
_, err = s.InstallCatalog(ctx, InstallCatalogRequest{
RepositoryID: *installation.RepositoryID,
PluginID: installation.PluginID,
Version: targetVersion,
})
if err != nil {
return nil, fmt.Errorf("update plugin %q to %s: %w", installation.PluginID, targetVersion, err)
}
// Clear available_version now that we've updated.
empty := ""
if err := s.installations.Update(ctx, installationID, UpdateInstallationInput{
AvailableVersion: &empty,
}); err != nil {
slog.WarnContext(ctx, "failed to clear available_version after update", "component", "plugins",
"installation_id", installationID, "error", err)
}
// Reload and return the updated installation.
return s.installations.GetByID(ctx, installationID)
}
func (s *Service) InstallBinary(ctx context.Context, req InstallBinaryRequest) (*InstallResult, error) {
var result *InstallResult
var err error
if req.Manifest != nil && req.Manifest.GetPluginId() != "" {
existing, existErr := s.existingInstallationByPluginID(ctx, req.Manifest.GetPluginId())
if existErr != nil {
return nil, existErr
}
if existing != nil {
if err = s.stopInstallationIfRunning(existing); err != nil {
return nil, err
}
result, err = s.installer.ReplaceBinary(ctx, existing, req)
if err != nil {
return nil, err
}
s.OnLifecycleChange(ctx)
return result, nil
}
}
result, err = s.installer.InstallBinary(ctx, req)
if err != nil {
return nil, err
}
s.OnLifecycleChange(ctx)
return result, nil
}
func (s *Service) InstallBinaryUpload(ctx context.Context, binaryData []byte) (*InstallResult, error) {
if len(binaryData) == 0 {
return nil, fmt.Errorf("binary data is required")
}
manifest, err := loadManifestFromBinary(ctx, binaryData)
if err != nil {
return nil, err
}
checksum := sha256.Sum256(binaryData)
actualChecksum := hex.EncodeToString(checksum[:])
var result *InstallResult
if s.installations == nil {
result, err = s.installer.installBinary(ctx, binaryData, actualChecksum, manifest, nil)
if err != nil {
return nil, err
}
s.OnLifecycleChange(ctx)
return result, nil
}
existing, err := s.installations.ListByPluginID(ctx, manifest.GetPluginId())
if err != nil {
return nil, fmt.Errorf("list existing plugin installations for %q: %w", manifest.GetPluginId(), err)
}
if len(existing) == 0 {
result, err = s.installer.installBinary(ctx, binaryData, actualChecksum, manifest, nil)
if err != nil {
return nil, err
}
s.OnLifecycleChange(ctx)
return result, nil
}
if len(existing) > 1 {
return nil, fmt.Errorf("multiple existing installations found for plugin %q", manifest.GetPluginId())
}
oldInstallation := existing[0]
if err := s.stopInstallationIfRunning(oldInstallation); err != nil {
return nil, err
}
result, err = s.installer.replaceBinary(ctx, oldInstallation, binaryData, actualChecksum, manifest)
if err != nil {
return nil, err
}
s.OnLifecycleChange(ctx)
return result, nil
}
func (s *Service) existingInstallationByPluginID(ctx context.Context, pluginID string) (*Installation, error) {
if s.installations == nil || pluginID == "" {
return nil, nil
}
existing, err := s.installations.ListByPluginID(ctx, pluginID)
if err != nil {
return nil, fmt.Errorf("list existing plugin installations for %q: %w", pluginID, err)
}
if len(existing) == 0 {
return nil, nil
}
if len(existing) > 1 {
return nil, fmt.Errorf("multiple existing installations found for plugin %q", pluginID)
}
return existing[0], nil
}
func (s *Service) stopInstallationIfRunning(existing *Installation) error {
if existing == nil || !existing.Enabled || s.host == nil {
return nil
}
if err := s.host.Stop(existing.ID); err != nil && !errors.Is(err, pluginhost.ErrClientNotFound) {
return fmt.Errorf("stop existing plugin installation %d: %w", existing.ID, err)
}
return nil
}
func (s *Service) PreloadEnabled(ctx context.Context) error {
if s.installations == nil {
return nil
}
installations, err := s.installations.ListEnabled(ctx)
if err != nil {
return err
}
for _, installation := range installations {
if installation == nil {
continue
}
// Builtin installations have no archive or binary; skip them explicitly
// instead of leaning on the tolerated ErrArchiveNotFound branch below
// (any other load error here is fatal to startup).
if installation.IsBuiltin() {
continue
}
if _, err := s.ensureLoadedInstallation(ctx, installation); err != nil {
if errors.Is(err, ErrArchiveNotFound) {
slog.WarnContext(ctx,
"plugin preload skipped: archive not found for enabled installation", "component", "plugins",
"installation_id", installation.ID,
"plugin_id", installation.PluginID,
"version", installation.Version,
)
continue
}
return fmt.Errorf("preload plugin installation %d: %w", installation.ID, err)
}
}
s.OnLifecycleChange(ctx)
return nil
}
func (s *Service) Start(ctx context.Context, installationID int) (pluginClient, error) {
installation, manifest, err := s.ensureInstallationCache(ctx, installationID, true)
if err != nil {
return nil, err
}
configEntries, err := s.globalConfigEntries(ctx, installation.ID)
if err != nil {
return nil, err
}
return s.host.Start(ctx, pluginhost.StartRequest{
InstallationID: installation.ID,
BinaryPath: installation.InstallPath,
Manifest: manifest,
Config: configEntries,
})
}
func (s *Service) Stop(installationID int) error {
if s.host == nil {
return nil
}
return s.host.Stop(installationID)
}
func (s *Service) MediaAnalyzerClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.MediaAnalyzerClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.MediaAnalyzer(capabilityID)
}
func (s *Service) MetadataProviderClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.MetadataProviderClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.MetadataProvider(capabilityID)
}
func (s *Service) ImageResolverClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.ImageResolverClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.ImageResolver(capabilityID)
}
func (s *Service) MarkerProviderClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.MarkerProviderClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.MarkerProvider(capabilityID)
}
func (s *Service) ScheduledTaskClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.ScheduledTaskClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.ScheduledTask(capabilityID)
}
func (s *Service) ScanSourceClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.ScanSourceClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.ScanSource(capabilityID)
}
func (s *Service) RequestRouterClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.RequestRouterClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.RequestRouter(capabilityID)
}
func (s *Service) ScanSourceClientByPluginID(
ctx context.Context,
pluginID string,
capabilityID string,
) (*pluginhost.ScanSourceClient, error) {
if s == nil || s.installations == nil {
return nil, fmt.Errorf("scan source plugin resolver is not configured")
}
installations, err := s.installations.ListByPluginID(ctx, pluginID)
if err != nil {
return nil, err
}
if len(installations) == 0 {
return nil, fmt.Errorf("scan source plugin %q is not installed", pluginID)
}
if len(installations) > 1 {
return nil, fmt.Errorf("scan source plugin %q is ambiguous across %d installations", pluginID, len(installations))
}
var matches []*Installation
for _, installation := range installations {
if installation == nil {
continue
}
capabilities, err := s.installations.ListCapabilities(ctx, installation.ID)
if err != nil {
return nil, err
}
for _, capability := range capabilities {
if capability == nil {
continue
}
if capability.Type == "scan_source.v1" && capability.ID == capabilityID {
matches = append(matches, installation)
break
}
}
}
if len(matches) == 0 {
return nil, fmt.Errorf("scan source capability %q is not installed for plugin %q", capabilityID, pluginID)
}
if !matches[0].Enabled {
return nil, fmt.Errorf("scan source plugin %q is disabled", pluginID)
}
return s.ScanSourceClient(ctx, matches[0].ID, capabilityID)
}
func (s *Service) EventConsumerClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.EventConsumerClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.EventConsumer(capabilityID)
}
func (s *Service) AuthProviderClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.AuthProviderClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.AuthProvider(capabilityID)
}
func (s *Service) HTTPRoutesClient(
ctx context.Context,
installationID int,
capabilityID string,
) (*pluginhost.HTTPRoutesClient, error) {
client, err := s.ensureClient(ctx, installationID)
if err != nil {
return nil, err
}
return client.HTTPRoutes(capabilityID)
}
func (s *Service) RouteDescriptors(ctx context.Context, installationID int) ([]*pluginv1.HttpRouteDescriptor, error) {
manifest, err := s.manifestForInstallation(ctx, installationID, true)
if err != nil {
return nil, err
}
return append([]*pluginv1.HttpRouteDescriptor(nil), manifest.GetHttpRoutes()...), nil
}
func (s *Service) ResolveAssetPath(ctx context.Context, installationID int, assetPath string) (string, error) {
installation, manifest, err := s.ensureInstallationCache(ctx, installationID, true)
if err != nil {
return "", err
}
for _, asset := range manifest.GetAssets() {
if asset.GetPath() == assetPath {
resolved := filepath.Join(filepath.Dir(installation.InstallPath), assetPath)
if _, err := os.Stat(resolved); err != nil {
return "", fmt.Errorf("plugin asset %q: %w", assetPath, err)
}
return resolved, nil
}
}
return "", fmt.Errorf("plugin asset %q not found", assetPath)
}
func (s *Service) UserConfigSchema(ctx context.Context, installationID int) ([]*pluginv1.ConfigSchema, error) {
manifest, err := s.manifestForInstallation(ctx, installationID, false)
if err != nil {
return nil, err
}
return append([]*pluginv1.ConfigSchema(nil), manifest.GetUserConfigSchema()...), nil
}
func (s *Service) ManifestForInstallation(
ctx context.Context,
installationID int,
) (*pluginv1.PluginManifest, error) {
return s.manifestForInstallation(ctx, installationID, false)
}
// ensureClient returns a running client for the installation, collapsing
// concurrent first-use of a cold installation into a single launch so a burst of
// callers does not spawn redundant plugin processes (Host.Start releases its lock
// during the slow launch and cannot dedupe). After the flight completes the key
// is freed, so subsequent callers re-run and hit the now-warm cache.
func (s *Service) ensureClient(ctx context.Context, installationID int) (pluginClient, error) {
v, err, _ := s.launchGroup.Do(strconv.Itoa(installationID), func() (any, error) {
// Isolate the shared launch from the leader caller's cancellation: other
// waiters depend on this in-flight launch, so a single caller's canceled
// request must not tear it down. Values (tracing, auth) are preserved.
return s.doEnsureClient(context.WithoutCancel(ctx), installationID)
})
if err != nil {
return nil, err
}
return v.(pluginClient), nil
}
func (s *Service) doEnsureClient(ctx context.Context, installationID int) (pluginClient, error) {
installation, err := s.loadInstallation(ctx, installationID, true)
if err != nil {
return nil, err
}
client, err := s.host.Client(installationID)
if err == nil {
installedManifest, manifestErr := LoadManifestFile(InstalledManifestPath(installation.InstallPath))
if manifestErr != nil {
slog.WarnContext(ctx, "plugin installed manifest unavailable; reusing healthy client", "component", "plugins",
"installation_id", installation.ID,
"plugin_id", installation.PluginID,
"version", installation.Version,
"error", manifestErr,
)
return client, nil
}
cachedManifest := client.Manifest()
if cachedManifest != nil && proto.Equal(cachedManifest, installedManifest) {
return client, nil
}
slog.WarnContext(ctx, "plugin client manifest drift detected; restarting", "component", "plugins",
"installation_id", installation.ID,
"plugin_id", installation.PluginID,
"cached_plugin_id", manifestPluginID(cachedManifest),
"installed_plugin_id", installedManifest.GetPluginId(),
"cached_version", manifestVersion(cachedManifest),
"installed_version", installedManifest.GetVersion(),
)
if stopErr := s.host.Stop(installationID); stopErr != nil && !errors.Is(stopErr, pluginhost.ErrClientNotFound) {
return nil, fmt.Errorf("stop stale plugin installation %d: %w", installationID, stopErr)
}
return s.Start(ctx, installationID)
}
if errors.Is(err, pluginhost.ErrPluginUnhealthy) {
if stopErr := s.host.Stop(installationID); stopErr != nil && !errors.Is(stopErr, pluginhost.ErrClientNotFound) {
return nil, fmt.Errorf("stop unhealthy plugin installation %d: %w", installationID, stopErr)
}
return s.Start(ctx, installationID)
}
if errors.Is(err, pluginhost.ErrClientNotFound) {
return s.Start(ctx, installationID)
}
return nil, err
}
func (s *Service) manifestForInstallation(ctx context.Context, installationID int, requireEnabled bool) (*pluginv1.PluginManifest, error) {
_, manifest, err := s.ensureInstallationCache(ctx, installationID, requireEnabled)
return manifest, err
}
func (s *Service) ensureInstallationCache(
ctx context.Context,
installationID int,
requireEnabled bool,
) (*Installation, *pluginv1.PluginManifest, error) {
installation, err := s.loadInstallation(ctx, installationID, requireEnabled)
if err != nil {
return nil, nil, err
}
manifest, err := s.ensureLoadedInstallation(ctx, installation)
if err != nil {
return nil, nil, err
}
return installation, manifest, nil
}
func (s *Service) loadInstallation(ctx context.Context, installationID int, requireEnabled bool) (*Installation, error) {
installation, err := s.cachedInstallation(ctx, installationID)
if err != nil {
return nil, err
}
// The requireEnabled gate is applied after the cache read so the cache
// stores the row regardless of its enabled state and ErrInstallationDisabled
// semantics are unchanged.
if requireEnabled && !installation.Enabled {
return nil, ErrInstallationDisabled
}
// requireEnabled marks paths that intend to launch or serve the plugin
// (start, manifest routes/assets, gRPC clients, HTTP proxy). The reserved
// builtin row has no binary behind it and must never reach those paths;
// not-found gives the proxy and API a clean 4xx. Reads with
// requireEnabled=false (IsInstallationEnabled for the metadata chain,
// generic listings) still see the row.
if requireEnabled && installation.IsBuiltin() {
return nil, ErrInstallationNotFound
}
return installation, nil
}
// cachedInstallation returns the plugin_installations row for installationID
// from the in-memory cache, loading it from the store on a miss. The returned
// *Installation is shared and must be treated as read-only by callers; it is
// evicted wholesale by invalidateInstallationCache on every lifecycle change.
func (s *Service) cachedInstallation(ctx context.Context, installationID int) (*Installation, error) {
s.installationCacheMu.RLock()
cached, ok := s.installationCache[installationID]
gen := s.installationCacheGen
s.installationCacheMu.RUnlock()
if ok {
return cached, nil
}
installation, err := s.installations.GetByID(ctx, installationID)
if err != nil {
return nil, err
}
s.installationCacheMu.Lock()
// Only publish the fetched row if no invalidation happened while GetByID was
// in flight. Otherwise the row may pre-date a just-committed lifecycle change
// (e.g. a disable), and writing it would resurrect stale state until the next
// event. On a generation mismatch we still return the freshly-read row to the
// caller but leave the cache untouched.
if s.installationCacheGen == gen {
if s.installationCache == nil {
s.installationCache = make(map[int]*Installation)
}
s.installationCache[installationID] = installation
}
s.installationCacheMu.Unlock()
return installation, nil
}
// invalidateInstallationCache clears the in-memory installation cache. It is
// registered as a lifecycle hook (see NewService) so OnLifecycleChange evicts
// stale rows after every install / enable / disable / update / uninstall.
func (s *Service) invalidateInstallationCache() {
if s == nil {
return
}
s.installationCacheMu.Lock()
s.installationCache = nil
s.installationCacheGen++
s.installationCacheMu.Unlock()
}
// IsInstallationEnabled reports whether the given plugin installation is
// enabled, served from the in-memory installation cache. It backs the metadata
// chain's enabled-check (internal/metadata/chain.go) so provider construction
// no longer issues a per-capability SELECT on the hot path.
func (s *Service) IsInstallationEnabled(ctx context.Context, installationID int) (bool, error) {
installation, err := s.loadInstallation(ctx, installationID, false)
if err != nil {
return false, err
}
return installation.Enabled, nil
}
// InstallationKind returns the installation's kind ("plugin" or "builtin") from
// the same in-memory cache IsInstallationEnabled reads, so metadata chain
// resolution can identify builtin rows without a per-capability DB query.
func (s *Service) InstallationKind(ctx context.Context, installationID int) (string, error) {
installation, err := s.loadInstallation(ctx, installationID, false)
if err != nil {
return "", err
}
return installation.Kind, nil
}
func (s *Service) ensureLoadedInstallation(
ctx context.Context,
installation *Installation,
) (*pluginv1.PluginManifest, error) {
if s.archiveCache == nil {
return LoadManifestFile(InstalledManifestPath(installation.InstallPath))
}
return s.archiveCache.Ensure(ctx, installation)
}
func (s *Service) globalConfigEntries(ctx context.Context, installationID int) ([]*pluginv1.ConfigEntry, error) {
if s.configs == nil {
return nil, nil
}
configs, err := s.configs.ListGlobalConfigs(ctx, installationID)
if err != nil {
return nil, fmt.Errorf("list plugin runtime configs for installation %d: %w", installationID, err)
}
entries := make([]*pluginv1.ConfigEntry, 0, len(configs))
for _, config := range configs {
if config == nil {
continue
}
value := config.Value
if value == nil {
value = map[string]any{}
}
structValue, err := structpb.NewStruct(value)
if err != nil {
return nil, fmt.Errorf(
"encode runtime config %q for installation %d: %w",
config.Key,
installationID,
err,
)
}
entries = append(entries, &pluginv1.ConfigEntry{
Key: config.Key,
Value: structValue,
})
}
return entries, nil
}