Files
DriverVault/API Server/internal/plugins/manager.go
T
tajniak81andClaude Opus 5 660af5736a Plugins: create the settings collection instead of waiting for it forever
01a8fec fixed the advice that led operators into this, but advice is not a
guard: a stack still running PB_BOOTSTRAP=false gets no app_settings
collection on upgrade, and the plugin panel sits at 503 while the retry
loop reads a collection that does not exist.

The fix is not to soften the reading. A missing collection stays "not
ready" rather than "no plugins configured", because the alternative lets
the first save write a fresh document over settings the server merely
failed to find - the failure this whole line of work exists to prevent.
Instead the server now fixes the cause: on a missing collection it creates
that collection and reads again.

Three pieces:

bootstrap.EnsureCollection creates one named collection from the desired
schema if absent, and nothing else. Deliberately narrower than Run - no
field reconcile elsewhere, no super-admin - so it is safe to call on a
deployment that turned the full bootstrap off. It creates the collection
the server cannot start without, not the schema the operator declined.

The store tells a missing collection apart from an outage. A 404 from a
list means the collection itself is gone: an existing but empty one answers
200 with no items. That is tagged errNoCollection, which wraps errNotReady
so every write is still refused, and IsMissingCollection narrows it. The
distinction matters because the remedies are opposites - creating
collections against a flaky database is exactly the wrong reflex, and a
test pins that an outage does not trigger it.

loadPlugins acts on the tag once, then re-reads. Failing to create is
reported as the original read error rather than the repair's, so the log
names the real problem.

Six tests: the tag and its negative in internal/plugins, and three in
internal/api against a fake PocketBase covering the collection being
created exactly once, an existing collection not being recreated, and an
outage creating nothing.

Docs from 01a8fec are corrected in the same pass - they said the panel
would answer 503 forever, which is no longer true. They now say what still
depends on the bootstrap (every other collection and field) and what does
not (app_settings alone).

go build, go vet and go test ./... pass; compose files still parse. Not
verified: no Docker CLI here, so the repair has not been exercised against
a real PocketBase, only the fake.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-21 17:07:18 +02:00

604 lines
18 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package plugins
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"log"
"net/http"
"sort"
"strings"
"sync"
"time"
)
// secretMask is what a set secret value is echoed back as. On save, a field that
// still equals the mask is left unchanged (mirrors the pb-config password flow).
const secretMask = "••••••••"
// record is the persisted state for one plugin. For builtins, Kind/BaseURL are
// omitted (the descriptor comes from the registry); external plugins set them.
type record struct {
Kind string `json:"kind,omitempty"`
BaseURL string `json:"baseURL,omitempty"`
Provider string `json:"provider,omitempty"`
Enabled bool `json:"enabled"`
Config map[string]string `json:"config,omitempty"`
}
// View is the plugin shape returned to the panel (secrets masked).
type View struct {
Descriptor
Enabled bool `json:"enabled"`
Config map[string]string `json:"config"`
BaseURL string `json:"baseURL,omitempty"`
Health *Health `json:"health,omitempty"`
}
// Manager owns the plugin registry, persisted state, and live instances.
type Manager struct {
store Store
mu sync.Mutex
loaded bool // settings have been read; until then saves are refused
factories map[string]Factory
records map[string]*record
live map[string]Plugin
health map[string]*Health
client *http.Client
}
// NewManager builds a Manager over the given Store. It starts unloaded: call
// Load before serving, and keep calling it until it succeeds if the store is
// not reachable yet (see Ready).
func NewManager(store Store) *Manager {
return &Manager{
store: store,
factories: builtinFactories(),
records: map[string]*record{},
live: map[string]Plugin{},
health: map[string]*Health{},
client: &http.Client{Timeout: 12 * time.Second},
}
}
// Ready reports whether the settings have been read from the store. While it is
// false the Manager knows of no configured plugins and refuses every mutation,
// so an unreachable store cannot cause settings to be overwritten or lost.
func (m *Manager) Ready() bool {
m.mu.Lock()
defer m.mu.Unlock()
return m.loaded
}
// Load reads the settings from the store and initialises every enabled plugin.
// An absent document is fine (nothing configured yet); an unreachable store is
// not, and leaves the Manager unloaded so a caller can retry.
func (m *Manager) Load(ctx context.Context) error {
// Read outside the lock: the production store is a network call.
data, found, err := m.store.Load(ctx)
if err != nil {
return classify(err, errNotReady)
}
recs := map[string]*record{}
if found {
if err := decodeRecords(data, recs); err != nil {
// Stay unloaded on purpose. A document we cannot parse must never be
// replaced by one built from an empty map — that would turn a read
// problem into permanent data loss. Saves refuse until it is fixed.
return fmt.Errorf("plugin settings in %s are unreadable: %w", m.store.Describe(), err)
}
}
m.mu.Lock()
defer m.mu.Unlock()
// Load is called again by the boot retry, so make it idempotent: the
// previous generation of instances must not be left running beside the new.
for name, p := range m.live {
_ = p.Shutdown(ctx)
delete(m.live, name)
}
m.records = recs
m.loaded = true
for name, rec := range m.records {
if !rec.Enabled {
continue
}
p := construct(name, m.factories[name], rec)
if p == nil {
log.Printf("plugins: cannot construct %q (unknown builtin?)", name)
continue
}
if err := p.Init(ctx, rec.Config); err != nil {
log.Printf("plugins: init %q failed: %v", name, err)
continue
}
m.live[name] = p
}
return nil
}
// decodeRecords parses a settings document into into. It is deliberately strict
// about two shapes that would otherwise take the server down:
//
// - An empty document, or a literal "null", decodes to a nil map. Assigning
// that straight to m.records made the next save panic with "assignment to
// entry in nil map"; both now mean "nothing configured".
// - A null entry ({"toyota": null}) leaves a nil *record that Load's enable
// loop would dereference. Those entries are dropped.
func decodeRecords(data []byte, into map[string]*record) error {
if len(bytes.TrimSpace(data)) == 0 {
return nil
}
var recs map[string]*record
if err := json.Unmarshal(data, &recs); err != nil {
return err
}
for name, rec := range recs {
if rec == nil {
continue
}
into[name] = rec
}
return nil
}
// encodeRecords renders the settings document written to the store.
func encodeRecords(recs map[string]*record) ([]byte, error) {
data, err := json.MarshalIndent(recs, "", " ")
if err != nil {
return nil, fmt.Errorf("%w: %v", errPersist, err)
}
return append(data, '\n'), nil
}
// construct builds a plugin instance from a builtin factory or an external record.
func construct(name string, f Factory, rec *record) Plugin {
if f != nil {
return f()
}
if rec != nil && rec.Kind == KindExternal {
return newExternalPlugin(name, rec.BaseURL, rec.Provider)
}
return nil
}
// descriptorFor returns a plugin's descriptor without needing a live instance.
func (m *Manager) descriptorFor(name string, rec *record) Descriptor {
if p := m.live[name]; p != nil {
return p.Descriptor()
}
if f := m.factories[name]; f != nil {
return f().Descriptor()
}
if rec != nil && rec.Kind == KindExternal {
return newExternalPlugin(name, rec.BaseURL, rec.Provider).Descriptor()
}
return Descriptor{Name: name}
}
// maskConfig echoes config back with secret fields masked when set.
func maskConfig(d Descriptor, cfg map[string]string) map[string]string {
out := map[string]string{}
for k, v := range cfg {
out[k] = v
}
for _, f := range d.ConfigFields {
if f.Secret && out[f.Key] != "" {
out[f.Key] = secretMask
}
}
return out
}
// List returns every known plugin (registry persisted), sorted by name.
func (m *Manager) List() []View {
m.mu.Lock()
defer m.mu.Unlock()
names := map[string]bool{}
for n := range m.factories {
names[n] = true
}
for n := range m.records {
names[n] = true
}
out := make([]View, 0, len(names))
for name := range names {
rec := m.records[name]
d := m.descriptorFor(name, rec)
v := View{Descriptor: d, Health: m.health[name]}
if rec != nil {
v.Enabled = rec.Enabled
v.BaseURL = rec.BaseURL
v.Config = maskConfig(d, rec.Config)
} else {
v.Config = map[string]string{}
}
out = append(out, v)
}
sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name })
return out
}
// Get returns a single plugin view (ok=false when unknown).
func (m *Manager) Get(name string) (View, bool) {
for _, v := range m.List() {
if v.Name == name {
return v, true
}
}
return View{}, false
}
// Upsert enables/disables a plugin and merges its config, then (re)initialises or
// shuts down the live instance to match. Secrets left at the mask are preserved.
func (m *Manager) Upsert(ctx context.Context, name string, enabled bool, incoming map[string]string) (View, error) {
m.mu.Lock()
if !m.loaded {
m.mu.Unlock()
return View{}, errNotReady
}
_, isBuiltin := m.factories[name]
rec := m.records[name]
if !isBuiltin && (rec == nil || rec.Kind != KindExternal) {
m.mu.Unlock()
return View{}, errUnknown
}
isNew := rec == nil
if isNew {
rec = &record{}
m.records[name] = rec
}
d := m.descriptorFor(name, rec)
merged := map[string]string{}
for k, v := range rec.Config {
merged[k] = v
}
// Apply incoming values, honouring the secret-mask keep-current rule.
secretKeys := map[string]bool{}
for _, f := range d.ConfigFields {
if f.Secret {
secretKeys[f.Key] = true
}
}
for k, v := range incoming {
if secretKeys[k] && v == secretMask {
continue // keep existing secret
}
merged[k] = strings.TrimSpace(v)
}
// Validate required fields when enabling.
if enabled {
for _, f := range d.ConfigFields {
if f.Required && merged[f.Key] == "" {
if isNew {
delete(m.records, name) // don't leave a blank record behind
}
m.mu.Unlock()
return View{}, errors.New("missing required setting: " + f.Label)
}
}
}
prevEnabled, prevConfig := rec.Enabled, rec.Config
rec.Enabled = enabled
rec.Config = merged
if err := m.saveLocked(ctx); err != nil {
// Roll back, so the panel shows what the store actually holds. Keeping
// the change in memory is what made a failed save look like a
// successful one — right up until the next restart brought every
// plugin back disabled.
if isNew {
delete(m.records, name)
} else {
rec.Enabled, rec.Config = prevEnabled, prevConfig
}
m.mu.Unlock()
return View{}, err
}
// Reconcile the live instance.
if old := m.live[name]; old != nil {
_ = old.Shutdown(ctx)
delete(m.live, name)
}
var initErr error
if enabled {
p := construct(name, m.factories[name], rec)
if p != nil {
if err := p.Init(ctx, merged); err != nil {
initErr = err
} else {
m.live[name] = p
}
}
}
m.mu.Unlock()
v, _ := m.Get(name)
return v, initErr
}
// RegisterExternal adds a new external (remote HTTP) plugin at runtime — the
// "add a plugin without a rebuild" path. It starts disabled.
func (m *Manager) RegisterExternal(ctx context.Context, name, baseURL, provider string) error {
name = strings.TrimSpace(name)
baseURL = strings.TrimRight(strings.TrimSpace(baseURL), "/")
if name == "" || baseURL == "" {
return errors.New("name and baseURL are required")
}
if !strings.HasPrefix(baseURL, "http://") && !strings.HasPrefix(baseURL, "https://") {
baseURL = "http://" + baseURL
}
m.mu.Lock()
defer m.mu.Unlock()
if !m.loaded {
return errNotReady
}
if _, dup := m.factories[name]; dup {
return errors.New("a builtin plugin already uses that name")
}
if _, dup := m.records[name]; dup {
return errors.New("a plugin with that name already exists")
}
m.records[name] = &record{Kind: KindExternal, BaseURL: baseURL, Provider: provider}
if err := m.saveLocked(ctx); err != nil {
delete(m.records, name)
return err
}
return nil
}
// Remove deletes an external plugin registration. Builtins can only be disabled.
func (m *Manager) Remove(ctx context.Context, name string) error {
m.mu.Lock()
defer m.mu.Unlock()
if !m.loaded {
return errNotReady
}
rec := m.records[name]
if rec == nil || rec.Kind != KindExternal {
return errors.New("only external plugins can be removed")
}
// Persist before tearing the instance down, so a failed write leaves a
// still-registered plugin still running rather than a half-removed one.
delete(m.records, name)
if err := m.saveLocked(ctx); err != nil {
m.records[name] = rec
return err
}
if p := m.live[name]; p != nil {
_ = p.Shutdown(ctx)
delete(m.live, name)
}
delete(m.health, name)
return nil
}
// HealthCheck probes a plugin now, building a transient instance if it is not
// currently live (so disabled plugins can still be tested). Result is cached.
func (m *Manager) HealthCheck(ctx context.Context, name string) (Health, error) {
m.mu.Lock()
p := m.live[name]
transient := false
var cfg map[string]string
if p == nil {
rec := m.records[name]
if rec != nil {
cfg = rec.Config
}
p = construct(name, m.factories[name], rec)
transient = true
}
m.mu.Unlock()
if p == nil {
return Health{}, errUnknown
}
if transient {
_ = p.Init(ctx, cfg)
defer func() { _ = p.Shutdown(context.Background()) }()
}
h := p.HealthCheck(ctx)
m.mu.Lock()
hc := h
m.health[name] = &hc
m.mu.Unlock()
return h, nil
}
// HealthCheckWith probes a plugin against a caller-resolved config rather than
// the stored global config. It builds a transient instance, Inits it with cfg,
// probes, and tears it down — so a per-user cascade (see internal/api/
// integrations.go) can health-check under the credentials in force for that
// caller without disturbing the global instance or its cached health.
func (m *Manager) HealthCheckWith(ctx context.Context, name string, cfg map[string]string) (Health, error) {
m.mu.Lock()
rec := m.records[name]
p := construct(name, m.factories[name], rec)
m.mu.Unlock()
if p == nil {
return Health{}, errUnknown
}
_ = p.Init(ctx, cfg)
defer func() { _ = p.Shutdown(context.Background()) }()
return p.HealthCheck(ctx), nil
}
// InvokeWith runs a capability against a caller-resolved config. Like
// HealthCheckWith, it uses a transient instance Inited with cfg so per-user
// credentials drive the call. Returns the plugin's raw JSON result.
func (m *Manager) InvokeWith(ctx context.Context, name string, cfg map[string]string, action string, payload json.RawMessage) (json.RawMessage, error) {
m.mu.Lock()
rec := m.records[name]
p := construct(name, m.factories[name], rec)
m.mu.Unlock()
if p == nil {
return nil, errUnknown
}
_ = p.Init(ctx, cfg)
defer func() { _ = p.Shutdown(context.Background()) }()
return p.Invoke(ctx, action, payload)
}
// BatchCall is one capability invocation inside an InvokeBatchWith request.
type BatchCall struct {
ID string // caller-chosen id, echoed back on the result
Action string // capability id
Params json.RawMessage // action params; may be nil
}
// BatchResult is the outcome of one BatchCall. Exactly one of Result/Err is set.
type BatchResult struct {
ID string
Result json.RawMessage
Err error
}
// batchConcurrency caps how many calls of one batch are in flight at once, so a
// snapshot of a whole vehicle doesn't arrive at the upstream as a burst.
const batchConcurrency = 4
// InvokeBatchWith runs several capabilities against one caller-resolved config,
// sharing a single transient instance. A connector that authenticates lazily
// (Toyota's OAuth login on first request) would otherwise repeat that login for
// every action, because InvokeWith builds and tears down an instance per call;
// sharing the instance logs in once for the whole batch.
//
// Calls run concurrently, so a plugin's Invoke must be safe for concurrent use —
// which the contract already implies, since the live instance is shared by every
// HTTP request. Results come back in request order, each carrying its own error;
// a non-nil error return means the batch never started (unknown plugin).
func (m *Manager) InvokeBatchWith(ctx context.Context, name string, cfg map[string]string, calls []BatchCall) ([]BatchResult, error) {
m.mu.Lock()
rec := m.records[name]
p := construct(name, m.factories[name], rec)
m.mu.Unlock()
if p == nil {
return nil, errUnknown
}
_ = p.Init(ctx, cfg)
defer func() { _ = p.Shutdown(context.Background()) }()
out := make([]BatchResult, len(calls))
sem := make(chan struct{}, batchConcurrency)
var wg sync.WaitGroup
for i, c := range calls {
wg.Add(1)
go func(i int, c BatchCall) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
res, err := p.Invoke(ctx, c.Action, c.Params)
out[i] = BatchResult{ID: c.ID, Result: res, Err: err}
}(i, c)
}
wg.Wait()
return out, nil
}
// RawConfig returns a copy of a plugin's stored (global) config and its enabled
// flag. ok is false for an unknown plugin. This is the top layer (L1) of the
// per-user cascade: the config a superadmin set in the panel, which lower layers
// inherit blank fields from. Secrets are returned in clear — callers must mask
// before returning anything to a client.
func (m *Manager) RawConfig(name string) (cfg map[string]string, enabled, ok bool) {
m.mu.Lock()
defer m.mu.Unlock()
_, isBuiltin := m.factories[name]
rec := m.records[name]
if !isBuiltin && rec == nil {
return nil, false, false
}
out := map[string]string{}
if rec != nil {
for k, v := range rec.Config {
out[k] = v
}
enabled = rec.Enabled
}
return out, enabled, true
}
// Shutdown tears down every live plugin instance. Wire into graceful shutdown.
func (m *Manager) Shutdown(ctx context.Context) {
m.mu.Lock()
defer m.mu.Unlock()
for name, p := range m.live {
_ = p.Shutdown(ctx)
delete(m.live, name)
}
}
// saveLocked writes the current records to the store. Caller must hold m.mu.
//
// Every failure is classified as errPersist here rather than relying on the
// Store to have done it. That distinction drives the HTTP status: an
// unclassified save failure would fall through to the "saved, but the plugin
// failed to start" branch and be reported as a 200, which is precisely how a
// failed save used to masquerade as a successful one.
func (m *Manager) saveLocked(ctx context.Context) error {
data, err := encodeRecords(m.records)
if err != nil {
return err
}
if err := m.store.Save(ctx, data); err != nil {
return classify(err, errPersist)
}
return nil
}
// classify tags err with sentinel unless it already carries it, so callers can
// branch on IsPersist/IsNotReady whatever a Store returned.
func classify(err error, sentinel error) error {
if errors.Is(err, sentinel) {
return err
}
return fmt.Errorf("%w: %v", sentinel, err)
}
var (
errUnknown = errors.New("unknown plugin")
errPersist = errors.New("plugin settings could not be saved")
errNotReady = errors.New("plugin settings are not loaded yet")
// errNoCollection is a not-ready that waiting alone will never resolve: the
// collection the settings live in does not exist. It wraps errNotReady, so
// the safety behaviour (refuse every write) is unchanged; the extra tag only
// lets the caller fix the cause instead of retrying forever.
errNoCollection = fmt.Errorf("%w: the settings collection does not exist", errNotReady)
)
// IsUnknown reports whether err came from addressing a plugin that doesn't exist.
func IsUnknown(err error) bool { return errors.Is(err, errUnknown) }
// IsPersist reports whether err means the change never reached the store. Such a
// change has been rolled back in memory: it must be reported as a failure, or
// the caller sees a save that silently vanishes on the next restart.
func IsPersist(err error) bool { return errors.Is(err, errPersist) }
// IsNotReady reports whether err means the settings have not been read yet —
// the store was unreachable at boot and is still being retried. Callers should
// answer 503 rather than present the plugin list as empty.
func IsNotReady(err error) bool { return errors.Is(err, errNotReady) }
// IsMissingCollection reports whether err means the settings collection does not
// exist. It implies IsNotReady, and narrows it: retrying cannot help, so the
// caller should create the collection and read again.
func IsMissingCollection(err error) bool { return errors.Is(err, errNoCollection) }