Files
synapbus/internal/plugin/lifecycle.go
T
Algis DumbrisandClaude Opus 4.7 494e8f83e9 feat(plugin): plugin framework + plugintest + demo plugin + integration tests
Implements the compile-in plugin system designed in spec 019:

- internal/plugin/ (~1,300 LOC):
  * Tiny Plugin interface + 10 optional HasX capability sub-interfaces
  * Host struct with Logger / DB / Messenger / Channels / Attachments /
    Search / Secrets / Events / Config / DataDir / Tracer / Metrics /
    DefaultOwner / BaseURL
  * Registry with panic-on-duplicate, stability-level tracking, capability
    indexing (MCP tools, actions, panels, channel types, routes, event
    subscribers, CLI commands)
  * Migrator: per-plugin SHA-256-checksum'd migration chain, namespaced
    plugin_<name>_* table enforcement, idempotent re-apply
  * Three-phase lifecycle (Migrate → Init → Start) with panic-safe
    wrappers around every plugin call; failure per plugin isolated,
    core continues
  * YAML config loader that preserves unknown top-level keys on round-trip
  * Status store exposing /api/plugins/status JSON
  * Restart helpers (Noop + SignalRestarter); graceful reload is
    in-process for the demo

- internal/plugin/plugintest/ (~345 LOC):
  * NopHost(t) with in-memory modernc.org/sqlite
  * Run(t, plugin) full-lifecycle smoke helper
  * Assertions: HasTool, HasAction, HasPanel, HasChannelType,
    HasMigration, PluginStarted, PluginFailed
  * ScopedSecrets that returns ErrSecretNotFound for cross-plugin
    reads (satisfies SC-006)

- internal/plugins/demo/ (canonical showcase):
  * Plugin that exercises every HasX capability (migrations, actions,
    HTTP routes, web panel, lifecycle, config schema, stability)
  * Own SQL migration creating plugin_demo_notes
  * Embedded HTML panel that fetches notes via JS
  * 4 unit tests covering smoke, full capability registration, action
    handlers, and config-driven max_notes limit

- cmd/plugindemo/ (~290 LOC):
  * Demo HTTP server wiring registry to chi
  * Mounts /api/plugins/status, /api/admin/plugins/{name}/{enable,disable},
    /api/actions/{name}, /api/plugins/<name>/* (per-plugin REST),
    /ui/plugins/<name>/ (per-plugin UI)
  * SIGHUP-triggered config reload + registry rebuild + mux swap
  * SIGTERM/SIGINT graceful shutdown

- test/integration/ (~357 LOC, build-tag "integration"):
  * 6 end-to-end tests against a spawned plugindemo binary
  * Enable/disable round-trip with data preservation
  * SIGHUP reload timing (measured 41 ms — SC-008 target is 2 s)
  * Action-404 on disabled plugin, panel-404 on disabled plugin
  * REST endpoints + UI panel reachable

Contract deviation: admin toggle endpoints moved from
/api/plugins/{name}/{enable,disable} to /api/admin/plugins/{name}/{...}
to avoid URL collision with chi per-plugin route mounts. rest.md updated.

Scope deferred to next session (mechanical follow-ups):
- Port internal/wiki/ to internal/plugins/wiki/
- Squash 26 migrations to schema/000_initial.sql
- Backup scripts for live kubic instance
- Remaining 9 plugin extractions
- Boundary-lint static analyzer
- Wire into cmd/synapbus/main.go

All unit + integration tests green. Chrome UI smoke test passes.
autonomous_summary.md carries the full verification record.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-04-19 07:14:27 +03:00

325 lines
8.4 KiB
Go

package plugin
import (
"context"
"errors"
"fmt"
"log/slog"
"sync"
"time"
)
// InitAll runs the three-phase boot sequence against an existing registry.
// 1. Migrate all enabled plugins (failure marks only that plugin failed).
// 2. Init each enabled plugin, collecting capabilities.
// 3. Start lifecycle-enabled plugins.
//
// The hostFactory produces a Host specialized for a given plugin name — it
// is responsible for scoping the logger, datadir, metrics, tracer, and
// secrets accessor. This indirection keeps the plugin package independent
// of concrete core types.
func (r *Registry) InitAll(ctx context.Context, hostFactory func(name string, cfg CapabilityContext) Host) error {
// Phase 1: Migrate
for _, name := range r.Enabled() {
p := r.Plugin(name)
if p == nil {
continue
}
// Use a temporary host to get the DB handle.
ctxHost := hostFactory(name, CapabilityContext{PhaseMigrate: true})
if ctxHost.DB == nil {
r.markFailed(name, fmt.Errorf("migrate: Host.DB is nil"))
continue
}
if _, err := ApplyMigrations(ctx, ctxHost.DB, p); err != nil {
r.markFailed(name, fmt.Errorf("migrate: %w", err))
continue
}
r.statuses.Update(name, func(e *StatusEntry) {
e.Status = StatusMigrated
// Populate migration version list for the status endpoint.
if versions, vErr := AppliedVersions(ctx, ctxHost.DB, name); vErr == nil {
e.MigrationVersions = versions
}
})
}
// Phase 2: Init
for _, name := range r.Enabled() {
if s, ok := r.statuses.Get(name); ok && s.Status == StatusFailed {
continue
}
p := r.Plugin(name)
host := hostFactory(name, CapabilityContext{PhaseInit: true})
if err := ensureDataDir(host.DataDir); err != nil {
r.markFailed(name, fmt.Errorf("prepare datadir: %w", err))
continue
}
if err := safeInit(ctx, p, host); err != nil {
r.markFailed(name, fmt.Errorf("init: %w", err))
continue
}
r.collectCapabilities(name, p)
r.statuses.Update(name, func(e *StatusEntry) {
e.Status = StatusInitialized
e.Capabilities = capabilityList(p)
e.ToolsRegistered = toolsList(p)
e.ActionsRegistered = actionsList(p)
})
}
// Phase 3: Start
for _, name := range r.Enabled() {
if s, ok := r.statuses.Get(name); ok && s.Status == StatusFailed {
continue
}
p := r.Plugin(name)
hl, ok := p.(HasLifecycle)
if !ok {
// Nothing to start; mark "started" to reflect "fully live".
r.statuses.Update(name, func(e *StatusEntry) {
e.Status = StatusStarted
e.StartedAt = time.Now()
})
continue
}
if err := safeStart(ctx, hl); err != nil {
r.markFailed(name, fmt.Errorf("start: %w", err))
continue
}
r.statuses.Update(name, func(e *StatusEntry) {
e.Status = StatusStarted
e.StartedAt = time.Now()
})
}
return nil
}
// CapabilityContext is a tiny struct the host factory can use to know which
// phase it is being invoked for. Allows the factory to populate only the
// subset needed for the phase (e.g. skip metrics registration during migrate).
type CapabilityContext struct {
PhaseMigrate bool
PhaseInit bool
}
// ShutdownAll calls Shutdown on lifecycle plugins in reverse start order.
func (r *Registry) ShutdownAll(ctx context.Context) {
enabled := r.Enabled()
for i := len(enabled) - 1; i >= 0; i-- {
name := enabled[i]
p := r.Plugin(name)
if hl, ok := p.(HasLifecycle); ok {
_ = safeShutdown(ctx, hl)
}
r.statuses.Update(name, func(e *StatusEntry) { e.Status = StatusStopped })
}
}
func (r *Registry) markFailed(name string, err error) {
r.statuses.Update(name, func(e *StatusEntry) {
e.Status = StatusFailed
if err != nil {
e.ErrorMessage = err.Error()
}
})
// Emit via slog at warn level. Host-level observers may do more.
slog.Warn("plugin failed", "plugin", name, "error", err)
}
func (r *Registry) collectCapabilities(name string, p Plugin) {
r.mu.Lock()
defer r.mu.Unlock()
if hmt, ok := p.(HasMCPTools); ok {
for _, t := range hmt.MCPTools() {
if _, exists := r.mcpTools[t.Name]; exists {
r.statuses.Update(name, func(e *StatusEntry) {
e.Status = StatusFailed
e.ErrorMessage = fmt.Sprintf("duplicate MCP tool %q", t.Name)
})
return
}
r.mcpTools[t.Name] = t
}
}
if ha, ok := p.(HasActions); ok {
for _, a := range ha.Actions() {
if _, exists := r.actions[a.Name]; exists {
r.statuses.Update(name, func(e *StatusEntry) {
e.Status = StatusFailed
e.ErrorMessage = fmt.Sprintf("duplicate action %q", a.Name)
})
return
}
r.actions[a.Name] = a
}
}
if hwp, ok := p.(HasWebPanels); ok {
for _, m := range hwp.WebPanels() {
r.panels = append(r.panels, m)
}
if h := hwp.PanelHandler(); h != nil {
r.panelHandlers[name] = h
}
}
if hct, ok := p.(HasChannelType); ok {
for _, ct := range hct.ChannelTypes() {
if _, exists := r.channelTypes[ct.Name]; exists {
r.statuses.Update(name, func(e *StatusEntry) {
e.Status = StatusFailed
e.ErrorMessage = fmt.Sprintf("duplicate channel type %q", ct.Name)
})
return
}
r.channelTypes[ct.Name] = ct
}
}
if hc, ok := p.(HasCLICommands); ok {
r.cliCommands = append(r.cliCommands, hc.CLICommands()...)
}
if hr, ok := p.(HasHTTPRoutes); ok {
mountPlugin := name
r.routeMounts = append(r.routeMounts, routeMount{
Plugin: mountPlugin,
Setup: func(rt Router) { hr.RegisterRoutes(rt) },
})
}
if he, ok := p.(HasEventHook); ok {
r.eventSubs = append(r.eventSubs, he)
}
}
func capabilityList(p Plugin) []string {
caps := []string{}
if _, ok := p.(HasMigrations); ok {
caps = append(caps, "migrations")
}
if _, ok := p.(HasMCPTools); ok {
caps = append(caps, "mcp_tools")
}
if _, ok := p.(HasActions); ok {
caps = append(caps, "actions")
}
if _, ok := p.(HasHTTPRoutes); ok {
caps = append(caps, "http_routes")
}
if _, ok := p.(HasWebPanels); ok {
caps = append(caps, "web_panel")
}
if _, ok := p.(HasCLICommands); ok {
caps = append(caps, "cli")
}
if _, ok := p.(HasChannelType); ok {
caps = append(caps, "channel_type")
}
if _, ok := p.(HasEventHook); ok {
caps = append(caps, "event_hook")
}
if _, ok := p.(HasLifecycle); ok {
caps = append(caps, "lifecycle")
}
if _, ok := p.(HasConfigSchema); ok {
caps = append(caps, "config_schema")
}
return caps
}
func toolsList(p Plugin) []string {
out := []string{}
if hmt, ok := p.(HasMCPTools); ok {
for _, t := range hmt.MCPTools() {
out = append(out, t.Name)
}
}
return out
}
func actionsList(p Plugin) []string {
out := []string{}
if ha, ok := p.(HasActions); ok {
for _, a := range ha.Actions() {
out = append(out, a.Name)
}
}
return out
}
// --- panic-safe wrappers ---
func safeInit(ctx context.Context, p Plugin, host Host) (err error) {
defer func() {
if r := recover(); r != nil {
err = fmt.Errorf("panic: %v", r)
}
}()
return p.Init(ctx, host)
}
func safeStart(ctx context.Context, hl HasLifecycle) (err error) {
defer func() {
if r := recover(); r != nil {
err = fmt.Errorf("panic: %v", r)
}
}()
return hl.Start(ctx)
}
func safeShutdown(ctx context.Context, hl HasLifecycle) (err error) {
defer func() {
if r := recover(); r != nil {
err = fmt.Errorf("panic: %v", r)
}
}()
return hl.Shutdown(ctx)
}
// --- in-process event bus (implements Events) ---
// NewEventBus returns a zero-configuration in-process bus suitable for tests
// and lightweight deployments. Subscribers are called sequentially per
// publish, each with panic-recover so one bad subscriber cannot break the
// others.
func NewEventBus() Events { return &localBus{subs: map[string][]func(context.Context, Event) error{}} }
type localBus struct {
mu sync.RWMutex
subs map[string][]func(context.Context, Event) error
}
func (b *localBus) Publish(ctx context.Context, e Event) error {
b.mu.RLock()
handlers := append([]func(context.Context, Event) error(nil), b.subs[e.Topic]...)
all := append(handlers, b.subs["*"]...)
b.mu.RUnlock()
var errs []error
for _, h := range all {
func() {
defer func() {
if r := recover(); r != nil {
errs = append(errs, fmt.Errorf("subscriber panic: %v", r))
}
}()
if err := h(ctx, e); err != nil {
errs = append(errs, err)
}
}()
}
return errors.Join(errs...)
}
func (b *localBus) Subscribe(topic string, fn func(ctx context.Context, e Event) error) (cancel func()) {
b.mu.Lock()
defer b.mu.Unlock()
b.subs[topic] = append(b.subs[topic], fn)
idx := len(b.subs[topic]) - 1
return func() {
b.mu.Lock()
defer b.mu.Unlock()
lst := b.subs[topic]
if idx < len(lst) {
b.subs[topic] = append(lst[:idx], lst[idx+1:]...)
}
}
}