Files
sub2api-add/plugins/plugin-admin/main.go
T
Qiufeng 028b505c36
Business Plugins CI / check (plugin-admin) (push) Successful in 1m35s
Business Plugins CI / check (subscription-admin) (push) Successful in 1m29s
feat: add controlled plugin marketplace lifecycle
2026-08-28 00:51:34 +08:00

2974 lines
94 KiB
Go

// Command plugin-admin is the standalone Business Plugin V1 control plane.
// It deliberately lives outside Sub2API Core: Core remains the authority for
// identity and business data while this service owns plugin packages,
// revisions, lifecycle state, configuration metadata and menu intents.
package main
import (
"archive/zip"
"bytes"
"context"
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"crypto/sha256"
"embed"
"encoding/base64"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net"
"net/http"
"net/url"
"os"
"os/exec"
"os/signal"
"path"
"path/filepath"
"regexp"
"sort"
"strconv"
"strings"
"sync"
"syscall"
"time"
"git.awaioi.com/awaioi/sub2api-add/plugins/plugin-admin/internal/manifest"
)
const (
pluginID = "qiu.plugin-admin"
pluginVersion = "1.0.0"
sessionCookieName = "plugin_admin_session"
maxJSONBytes = 2 << 20
maxPackageBytes = 128 << 20
maxPackageFiles = 512
maxUncompressedBytes = 256 << 20
defaultSessionTTL = 30 * time.Minute
defaultSessionMaxTTL = 8 * time.Hour
defaultOperationLimit = 200
maxOperationBodyBytes = maxPackageBytes + (4 << 20)
)
//go:embed ui/*
var uiFS embed.FS
var requestIDPattern = regexp.MustCompile(`^[A-Za-z0-9._:-]{1,64}$`)
type coreEnvelope struct {
Code int `json:"code"`
Message string `json:"message"`
Data json.RawMessage `json:"data"`
status int
}
type coreError struct{ status int }
func (e *coreError) Error() string { return "core request failed" }
type coreClient struct {
base string
http *http.Client
admin map[string]struct{}
mu sync.Mutex
}
func newCoreClient(raw string) (*coreClient, error) {
raw = strings.TrimRight(strings.TrimSpace(raw), "/")
u, err := url.Parse(raw)
if err != nil || u.Host == "" || (u.Scheme != "http" && u.Scheme != "https") || u.User != nil || u.RawQuery != "" || u.Fragment != "" || (u.Path != "" && u.Path != "/") {
return nil, errors.New("CORE_BASE_URL must be an absolute origin without credentials, path, query or fragment")
}
if u.Scheme == "http" && !isLoopbackHost(u.Hostname()) {
return nil, errors.New("CORE_BASE_URL must use HTTPS unless Core is on loopback")
}
transport := http.DefaultTransport.(*http.Transport).Clone()
transport.Proxy = nil
return &coreClient{base: raw, http: &http.Client{Timeout: 10 * time.Second, Transport: transport, CheckRedirect: func(_ *http.Request, _ []*http.Request) error { return http.ErrUseLastResponse }}}, nil
}
func isLoopbackHost(host string) bool {
if strings.EqualFold(strings.TrimSuffix(host, "."), "localhost") {
return true
}
ip := net.ParseIP(host)
return ip != nil && ip.IsLoopback()
}
func (c *coreClient) call(ctx context.Context, method, requestPath string, body any, accessToken, requestID string) (coreEnvelope, error) {
if c == nil {
return coreEnvelope{}, errors.New("core client unavailable")
}
if !allowedCorePath(method, requestPath) {
return coreEnvelope{}, errors.New("core path is not in the control-plane allowlist")
}
var reader io.Reader
if body != nil {
payload, err := json.Marshal(body)
if err != nil {
return coreEnvelope{}, err
}
reader = bytes.NewReader(payload)
}
req, err := http.NewRequestWithContext(ctx, method, c.base+requestPath, reader)
if err != nil {
return coreEnvelope{}, err
}
req.Header.Set("Accept", "application/json")
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
if accessToken != "" {
req.Header.Set("Authorization", "Bearer "+accessToken)
}
requestID = normalizeRequestID(requestID)
req.Header.Set("X-Request-Id", requestID)
res, err := c.http.Do(req)
if err != nil {
return coreEnvelope{}, err
}
defer res.Body.Close()
var out coreEnvelope
if err := json.NewDecoder(io.LimitReader(res.Body, 4<<20)).Decode(&out); err != nil {
return coreEnvelope{}, err
}
out.status = res.StatusCode
if res.StatusCode >= 300 || out.Code != 0 {
status := res.StatusCode
if status < 400 && out.Code >= 400 {
status = out.Code
}
return out, &coreError{status: status}
}
return out, nil
}
func allowedCorePath(method, p string) bool {
if method == http.MethodGet && (p == "/api/v1/auth/me" || p == "/api/v1/settings/public" || p == "/api/v1/admin/settings") {
return true
}
return method == http.MethodPut && p == "/api/v1/admin/settings" || method == http.MethodPost && (p == "/api/v1/auth/login" || p == "/api/v1/auth/login/2fa" || p == "/api/v1/auth/refresh" || p == "/api/v1/auth/logout")
}
func (c *coreClient) login(ctx context.Context, body any, rid string) (coreEnvelope, error) {
return c.call(ctx, http.MethodPost, "/api/v1/auth/login", body, "", rid)
}
func (c *coreClient) login2FA(ctx context.Context, body any, rid string) (coreEnvelope, error) {
return c.call(ctx, http.MethodPost, "/api/v1/auth/login/2fa", body, "", rid)
}
func (c *coreClient) refresh(ctx context.Context, refresh, rid string) (coreEnvelope, error) {
return c.call(ctx, http.MethodPost, "/api/v1/auth/refresh", map[string]string{"refresh_token": refresh}, "", rid)
}
func (c *coreClient) logout(ctx context.Context, refresh, rid string) {
if refresh != "" {
_, _ = c.call(ctx, http.MethodPost, "/api/v1/auth/logout", map[string]string{"refresh_token": refresh}, "", rid)
}
}
func (c *coreClient) me(ctx context.Context, access, rid string) (coreEnvelope, error) {
return c.call(ctx, http.MethodGet, "/api/v1/auth/me", nil, access, rid)
}
func (c *coreClient) settings(ctx context.Context, access, rid string) (coreEnvelope, error) {
return c.call(ctx, http.MethodGet, "/api/v1/admin/settings", nil, access, rid)
}
func (c *coreClient) publicSettings(ctx context.Context, rid string) (coreEnvelope, error) {
return c.call(ctx, http.MethodGet, "/api/v1/settings/public", nil, "", rid)
}
func (c *coreClient) updateMenu(ctx context.Context, access, rid string, items []any) (coreEnvelope, error) {
return c.call(ctx, http.MethodPut, "/api/v1/admin/settings", map[string]any{"custom_menu_items": items}, access, rid)
}
type session struct {
AccessToken string `json:"-"`
RefreshToken string `json:"-"`
CSRFToken string `json:"csrf_token"`
User map[string]any `json:"user"`
CreatedAt time.Time `json:"created_at"`
LastSeen time.Time `json:"last_seen"`
}
type pendingLogin struct {
TempToken string
Expires time.Time
}
type revision struct {
ID string `json:"id"`
Version string `json:"version"`
Path string `json:"path"`
ArchiveSHA string `json:"archive_sha256"`
Manifest manifest.Manifest `json:"manifest"`
VerifiedAt time.Time `json:"verified_at"`
HealthyAt time.Time `json:"healthy_at,omitempty"`
Retired bool `json:"retired,omitempty"`
}
type pluginRecord struct {
Manifest manifest.Manifest `json:"manifest"`
State string `json:"state"`
ActiveRevision string `json:"active_revision,omitempty"`
PendingRevision string `json:"pending_revision,omitempty"`
PreviousState string `json:"previous_state,omitempty"`
Revisions []revision `json:"revisions"`
ConfigCipher string `json:"config_cipher,omitempty"`
Endpoint string `json:"endpoint,omitempty"`
LastError string `json:"last_error,omitempty"`
UpdatedAt time.Time `json:"updated_at"`
}
type operation struct {
ID string `json:"id"`
IdempotencyKey string `json:"idempotency_key"`
Kind string `json:"kind"`
PluginID string `json:"plugin_id"`
Revision string `json:"revision,omitempty"`
ActorID any `json:"actor_id,omitempty"`
RequestID string `json:"request_id"`
RequestHash string `json:"request_hash,omitempty"`
State string `json:"state"`
Error string `json:"error,omitempty"`
Warning string `json:"warning,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
}
type auditEvent struct {
Time time.Time `json:"time"`
Action string `json:"action"`
PluginID string `json:"plugin_id,omitempty"`
ActorID any `json:"actor_id,omitempty"`
Operation string `json:"operation_id,omitempty"`
RequestID string `json:"request_id"`
Result string `json:"result"`
}
// stagedRevisionPath keeps an uninstall reversible until the registry commit
// has succeeded. Revision directories are renamed within the installed root,
// then removed after the registry no longer references them.
type stagedRevisionPath struct {
original string
staged string
}
type registryData struct {
Plugins map[string]pluginRecord `json:"plugins"`
Operations []operation `json:"operations"`
Audit []auditEvent `json:"audit"`
}
func clonePluginRecord(p pluginRecord) pluginRecord {
p.Revisions = append([]revision(nil), p.Revisions...)
for i := range p.Revisions {
p.Revisions[i].Manifest.Capabilities = append([]string(nil), p.Revisions[i].Manifest.Capabilities...)
p.Revisions[i].Manifest.TestedCoreVersions = append([]string(nil), p.Revisions[i].Manifest.TestedCoreVersions...)
if p.Revisions[i].Manifest.Files != nil {
p.Revisions[i].Manifest.Files = make(map[string]string, len(p.Revisions[i].Manifest.Files))
for key, value := range p.Revisions[i].Manifest.Files {
p.Revisions[i].Manifest.Files[key] = value
}
}
}
p.Manifest.Capabilities = append([]string(nil), p.Manifest.Capabilities...)
p.Manifest.TestedCoreVersions = append([]string(nil), p.Manifest.TestedCoreVersions...)
if p.Manifest.Files != nil {
p.Manifest.Files = make(map[string]string, len(p.Manifest.Files))
for key, value := range p.Manifest.Files {
p.Manifest.Files[key] = value
}
}
return p
}
type registry struct {
path string
mu sync.Mutex
data registryData
}
func openRegistry(dir string) (*registry, error) {
if dir == "" {
dir = "./data"
}
if err := os.MkdirAll(dir, 0o700); err != nil {
return nil, err
}
r := &registry{path: filepath.Join(dir, "registry.json"), data: registryData{Plugins: map[string]pluginRecord{}, Operations: []operation{}, Audit: []auditEvent{}}}
raw, err := os.ReadFile(r.path)
if errors.Is(err, os.ErrNotExist) {
return r, nil
}
if err != nil {
return nil, err
}
if err := json.Unmarshal(raw, &r.data); err != nil {
return nil, fmt.Errorf("decode registry: %w", err)
}
if r.data.Plugins == nil {
r.data.Plugins = map[string]pluginRecord{}
}
if r.data.Operations == nil {
r.data.Operations = []operation{}
}
if r.data.Audit == nil {
r.data.Audit = []auditEvent{}
}
changed := false
for i := range r.data.Operations {
if r.data.Operations[i].State == "running" && time.Since(r.data.Operations[i].UpdatedAt) > 10*time.Minute {
r.data.Operations[i].State = "failed"
r.data.Operations[i].Error = "operation interrupted by control-plane restart"
r.data.Operations[i].UpdatedAt = time.Now().UTC()
changed = true
}
}
if changed {
if err := r.saveLocked(); err != nil {
return nil, err
}
}
return r, nil
}
func (r *registry) saveLocked() error {
raw, err := json.MarshalIndent(r.data, "", " ")
if err != nil {
return err
}
tmp := r.path + ".tmp-" + token(8)
if err := os.WriteFile(tmp, raw, 0o600); err != nil {
return err
}
if err := os.Rename(tmp, r.path); err != nil {
_ = os.Remove(tmp)
return err
}
return nil
}
func (r *registry) save() error { r.mu.Lock(); defer r.mu.Unlock(); return r.saveLocked() }
func (r *registry) addAudit(ev auditEvent) error {
r.mu.Lock()
defer r.mu.Unlock()
previous := append([]auditEvent(nil), r.data.Audit...)
r.data.Audit = append(r.data.Audit, ev)
if len(r.data.Audit) > 500 {
r.data.Audit = r.data.Audit[len(r.data.Audit)-500:]
}
if err := r.saveLocked(); err != nil {
// Do not report an audit event that was not persisted.
r.data.Audit = previous
return err
}
return nil
}
func (r *registry) operation(kind, pluginID, key string, actor any, rid, requestHash string) (operation, bool, bool, error) {
r.mu.Lock()
defer r.mu.Unlock()
for i := len(r.data.Operations) - 1; i >= 0; i-- {
op := r.data.Operations[i]
if op.Kind == kind && op.PluginID == pluginID && key != "" && op.IdempotencyKey == key && fmt.Sprint(op.ActorID) == fmt.Sprint(actor) {
return op, true, op.RequestHash != "" && requestHash != "" && op.RequestHash != requestHash, nil
}
}
previous := append([]operation(nil), r.data.Operations...)
op := operation{ID: token(16), IdempotencyKey: key, Kind: kind, PluginID: pluginID, ActorID: actor, RequestID: rid, RequestHash: requestHash, State: "running", CreatedAt: time.Now().UTC(), UpdatedAt: time.Now().UTC()}
r.data.Operations = append(r.data.Operations, op)
if len(r.data.Operations) > defaultOperationLimit {
r.data.Operations = r.data.Operations[len(r.data.Operations)-defaultOperationLimit:]
}
if err := r.saveLocked(); err != nil {
r.data.Operations = previous
return operation{}, false, false, err
}
return op, false, false, nil
}
func (r *registry) finishOperation(op operation, operationErr error) (operation, error) {
r.mu.Lock()
defer r.mu.Unlock()
previous := operation{}
found := false
for i := range r.data.Operations {
if r.data.Operations[i].ID == op.ID {
previous = r.data.Operations[i]
found = true
break
}
}
if !found {
return op, errors.New("operation not found")
}
op.State = "completed"
if operationErr != nil {
op.State = "failed"
op.Error = sanitizeError(operationErr)
}
op.UpdatedAt = time.Now().UTC()
for i := range r.data.Operations {
if r.data.Operations[i].ID == op.ID {
r.data.Operations[i] = op
break
}
}
if err := r.saveLocked(); err != nil {
if found {
for i := range r.data.Operations {
if r.data.Operations[i].ID == op.ID {
r.data.Operations[i] = previous
break
}
}
} else if len(r.data.Operations) > 0 {
r.data.Operations = r.data.Operations[:len(r.data.Operations)-1]
}
return op, err
}
return op, nil
}
func (r *registry) setOperationPlugin(operationID, pluginID string) error {
r.mu.Lock()
defer r.mu.Unlock()
for i := range r.data.Operations {
if r.data.Operations[i].ID != operationID {
continue
}
previous := r.data.Operations[i]
r.data.Operations[i].PluginID = pluginID
if err := r.saveLocked(); err != nil {
r.data.Operations[i] = previous
return err
}
return nil
}
return errors.New("operation not found")
}
func sanitizeError(err error) string {
if err == nil {
return ""
}
message := strings.TrimSpace(err.Error())
for _, secret := range []string{"Bearer ", "refresh_token", "password", "api_key", "secret"} {
if strings.Contains(strings.ToLower(message), strings.ToLower(secret)) {
return "operation failed"
}
}
if len(message) > 240 {
return message[:240]
}
return message
}
type app struct {
core *coreClient
registry *registry
root string
publicBasePath string
cookiePath string
cookieSecure bool
cookieSameSite http.SameSite
frameAncestors []string
allowUnsigned bool
trustedPublishers map[string][]byte
configKey []byte
sessions map[string]session
sessionLocks map[string]*sync.Mutex
pending map[string]pendingLogin
processes map[string]*exec.Cmd
pluginLocks map[string]*sync.Mutex
marketplaceConfig marketplaceService
mu sync.Mutex
}
func newApp(core *coreClient, r *registry, root string) *app {
key := sha256.Sum256([]byte(token(32)))
marketplace, _ := newMarketplaceService(filepath.Join(root, "marketplace", "index.json"), "", true)
return &app{core: core, registry: r, root: root, cookiePath: "/", cookieSameSite: http.SameSiteLaxMode, frameAncestors: []string{"'self'"}, configKey: key[:], sessions: map[string]session{}, sessionLocks: map[string]*sync.Mutex{}, pending: map[string]pendingLogin{}, processes: map[string]*exec.Cmd{}, pluginLocks: map[string]*sync.Mutex{}, marketplaceConfig: marketplace}
}
func (a *app) lockPlugin(id string) func() {
a.mu.Lock()
lock := a.pluginLocks[id]
if lock == nil {
lock = &sync.Mutex{}
a.pluginLocks[id] = lock
}
a.mu.Unlock()
lock.Lock()
return lock.Unlock
}
func requestID(r *http.Request) string {
if r == nil {
return token(12)
}
value := strings.TrimSpace(r.Header.Get("X-Request-Id"))
if requestIDPattern.MatchString(value) {
return value
}
return token(12)
}
func normalizeRequestID(value string) string {
value = strings.TrimSpace(value)
if requestIDPattern.MatchString(value) {
return value
}
return token(12)
}
func token(n int) string {
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
panic(err)
}
return hex.EncodeToString(b)
}
func decodeJSON(r *http.Request, dst any, limit int64) error {
if r == nil || r.Body == nil {
return errors.New("request body is required")
}
defer r.Body.Close()
raw, err := io.ReadAll(io.LimitReader(r.Body, limit+1))
if err != nil {
return err
}
if int64(len(raw)) > limit {
return errors.New("request body exceeds size limit")
}
if len(bytes.TrimSpace(raw)) == 0 {
return errors.New("request body is required")
}
dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields()
if err := dec.Decode(dst); err != nil {
return err
}
var trailing any
if err := dec.Decode(&trailing); err != io.EOF {
if err == nil {
return errors.New("request body must contain exactly one JSON value")
}
return err
}
return nil
}
func writeJSON(w http.ResponseWriter, status int, value any) {
w.Header().Set("Cache-Control", "no-store")
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(value)
}
func (a *app) setCookie(w http.ResponseWriter, id string, maxAge int) {
http.SetCookie(w, &http.Cookie{Name: sessionCookieName, Value: id, Path: a.cookiePath, HttpOnly: true, Secure: a.cookieSecure, SameSite: a.cookieSameSite, MaxAge: maxAge})
}
func (a *app) sessionFromRequest(r *http.Request) (string, session, bool) {
c, err := r.Cookie(sessionCookieName)
if err != nil || c.Value == "" {
return "", session{}, false
}
a.mu.Lock()
defer a.mu.Unlock()
s, ok := a.sessions[c.Value]
if !ok {
return "", session{}, false
}
now := time.Now()
if now.Sub(s.CreatedAt) > defaultSessionMaxTTL || now.Sub(s.LastSeen) > defaultSessionTTL {
delete(a.sessions, c.Value)
delete(a.sessionLocks, c.Value)
return "", session{}, false
}
s.LastSeen = now
a.sessions[c.Value] = s
return c.Value, s, true
}
func publicUser(user map[string]any) map[string]any {
out := map[string]any{}
for key, value := range user {
lower := strings.ToLower(key)
if strings.Contains(lower, "token") || strings.Contains(lower, "password") || strings.Contains(lower, "secret") || lower == "api_key" || lower == "refresh_token" {
continue
}
out[key] = value
}
return out
}
func envelopeData(e coreEnvelope) map[string]any {
var data map[string]any
if len(e.Data) > 0 && json.Unmarshal(e.Data, &data) == nil && data != nil {
return data
}
return map[string]any{}
}
func menuItemsFromSettings(data map[string]any) []any {
if raw, ok := data["custom_menu_items"].([]any); ok {
return raw
}
if encoded, ok := data["custom_menu_items"].(string); ok && strings.TrimSpace(encoded) != "" {
var items []any
if json.Unmarshal([]byte(encoded), &items) == nil {
return items
}
}
return nil
}
func isAdmin(user map[string]any) bool {
role, _ := user["role"].(string)
return strings.EqualFold(role, "admin") || strings.EqualFold(role, "administrator")
}
func (a *app) login(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
var in struct {
Email string `json:"email"`
Password string `json:"password"`
TurnstileToken string `json:"turnstile_token,omitempty"`
TencentCaptchaTicket string `json:"tencent_captcha_ticket,omitempty"`
TencentCaptchaRandstr string `json:"tencent_captcha_randstr,omitempty"`
}
if err := decodeJSON(r, &in, 1<<20); err != nil || strings.TrimSpace(in.Email) == "" || in.Password == "" {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid credentials"})
return
}
if a.core == nil {
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "core unavailable"})
return
}
out, err := a.core.login(r.Context(), in, requestID(r))
if err != nil {
a.coreError(w, err, "core login failed")
return
}
data := envelopeData(out)
if requires, _ := data["requires_2fa"].(bool); requires {
temp, _ := data["temp_token"].(string)
if temp == "" {
writeJSON(w, http.StatusBadGateway, map[string]string{"error": "2fa challenge missing"})
return
}
pending := token(16)
a.mu.Lock()
a.pending[pending] = pendingLogin{TempToken: temp, Expires: time.Now().Add(5 * time.Minute)}
a.mu.Unlock()
writeJSON(w, http.StatusOK, map[string]any{"requires_2fa": true, "pending_token": pending})
return
}
a.finishLogin(w, r, data)
}
func (a *app) login2FA(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
var in struct {
PendingToken string `json:"pending_token"`
TOTPCode string `json:"totp_code"`
}
if err := decodeJSON(r, &in, 1<<20); err != nil || in.PendingToken == "" || len(in.TOTPCode) != 6 {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid 2fa request"})
return
}
a.mu.Lock()
pending, ok := a.pending[in.PendingToken]
delete(a.pending, in.PendingToken)
a.mu.Unlock()
if !ok || time.Now().After(pending.Expires) {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "2fa session expired"})
return
}
out, err := a.core.login2FA(r.Context(), map[string]string{"temp_token": pending.TempToken, "totp_code": in.TOTPCode}, requestID(r))
if err != nil {
a.coreError(w, err, "2fa verification failed")
return
}
a.finishLogin(w, r, envelopeData(out))
}
func (a *app) finishLogin(w http.ResponseWriter, r *http.Request, data map[string]any) {
access, _ := data["access_token"].(string)
refresh, _ := data["refresh_token"].(string)
if access == "" {
writeJSON(w, http.StatusBadGateway, map[string]string{"error": "core token missing"})
return
}
me, err := a.core.me(r.Context(), access, requestID(r))
if err != nil {
a.core.logout(r.Context(), refresh, requestID(r))
writeJSON(w, http.StatusForbidden, map[string]string{"error": "admin verification failed"})
return
}
user := envelopeData(me)
if !isAdmin(user) {
a.core.logout(r.Context(), refresh, requestID(r))
writeJSON(w, http.StatusForbidden, map[string]string{"error": "admin role required"})
return
}
now := time.Now()
s := session{AccessToken: access, RefreshToken: refresh, CSRFToken: token(16), User: publicUser(user), CreatedAt: now, LastSeen: now}
id := token(32)
a.mu.Lock()
a.sessions[id] = s
a.sessionLocks[id] = &sync.Mutex{}
a.mu.Unlock()
a.setCookie(w, id, int(defaultSessionMaxTTL/time.Second))
writeJSON(w, http.StatusOK, map[string]any{"ok": true, "csrf_token": s.CSRFToken, "user": s.User})
}
func (a *app) removeSession(id string) {
a.mu.Lock()
delete(a.sessions, id)
delete(a.sessionLocks, id)
a.mu.Unlock()
}
func (a *app) logout(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
if id, s, ok := a.sessionFromRequest(r); ok {
if r.Header.Get("X-CSRF-Token") != s.CSRFToken {
writeJSON(w, http.StatusForbidden, map[string]string{"error": "csrf validation failed"})
return
}
a.removeSession(id)
a.core.logout(r.Context(), s.RefreshToken, requestID(r))
}
a.setCookie(w, "", -1)
writeJSON(w, http.StatusOK, map[string]bool{"ok": true})
}
func (a *app) refreshSession(ctx context.Context, id string, stale session) (session, bool) {
a.mu.Lock()
lock := a.sessionLocks[id]
a.mu.Unlock()
if lock == nil {
return session{}, false
}
lock.Lock()
defer lock.Unlock()
a.mu.Lock()
current, ok := a.sessions[id]
a.mu.Unlock()
if !ok {
return session{}, false
}
if current.AccessToken != stale.AccessToken {
return current, true
}
out, err := a.core.refresh(ctx, current.RefreshToken, token(12))
if err != nil {
return session{}, false
}
data := envelopeData(out)
access, _ := data["access_token"].(string)
if access == "" {
return session{}, false
}
current.AccessToken = access
if next, _ := data["refresh_token"].(string); next != "" {
current.RefreshToken = next
}
me, err := a.core.me(ctx, access, token(12))
if err != nil {
return session{}, false
}
current.User = publicUser(envelopeData(me))
if !isAdmin(envelopeData(me)) {
return current, false
}
current.LastSeen = time.Now()
a.mu.Lock()
a.sessions[id] = current
a.mu.Unlock()
return current, true
}
func (a *app) authenticate(w http.ResponseWriter, r *http.Request) (string, session, bool) {
id, s, ok := a.sessionFromRequest(r)
if !ok {
a.setCookie(w, "", -1)
writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "authentication required"})
return "", session{}, false
}
if r.Method != http.MethodGet && r.Header.Get("X-CSRF-Token") != s.CSRFToken {
writeJSON(w, http.StatusForbidden, map[string]string{"error": "csrf validation failed"})
return "", session{}, false
}
if a.core == nil {
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "core unavailable"})
return "", session{}, false
}
me, err := a.core.me(r.Context(), s.AccessToken, requestID(r))
if err != nil {
if ce, ok := err.(*coreError); ok && ce.status == http.StatusUnauthorized && s.RefreshToken != "" {
if next, refreshed := a.refreshSession(r.Context(), id, s); refreshed {
return id, next, true
} else if next.User != nil && !isAdmin(next.User) {
a.removeSession(id)
a.core.logout(r.Context(), s.RefreshToken, requestID(r))
writeJSON(w, http.StatusForbidden, map[string]string{"error": "admin role required"})
return "", session{}, false
}
}
a.removeSession(id)
a.core.logout(r.Context(), s.RefreshToken, requestID(r))
a.setCookie(w, "", -1)
if ce, ok := err.(*coreError); ok && ce.status == http.StatusUnauthorized {
writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "core session expired"})
} else {
a.coreError(w, err, "core session unavailable")
}
return "", session{}, false
}
user := envelopeData(me)
if !isAdmin(user) {
a.removeSession(id)
a.core.logout(r.Context(), s.RefreshToken, requestID(r))
writeJSON(w, http.StatusForbidden, map[string]string{"error": "admin role required"})
return "", session{}, false
}
s.User = publicUser(user)
a.mu.Lock()
a.sessions[id] = s
a.mu.Unlock()
return id, s, true
}
func (a *app) me(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
_, s, ok := a.authenticate(w, r)
if ok {
writeJSON(w, http.StatusOK, map[string]any{"user": s.User, "csrf_token": s.CSRFToken, "plugin_id": pluginID, "plugin_version": pluginVersion})
}
}
func (a *app) health(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
writeJSON(w, http.StatusOK, map[string]any{"status": "ok", "plugin_id": pluginID, "version": pluginVersion})
}
func (a *app) ready(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
if a.registry == nil || a.core == nil {
writeJSON(w, http.StatusServiceUnavailable, map[string]any{"status": "not_ready"})
return
}
writeJSON(w, http.StatusOK, map[string]any{"status": "ready", "plugin_id": pluginID, "version": pluginVersion})
}
func (a *app) apiPlugins(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
if _, _, ok := a.authenticate(w, r); !ok {
return
}
a.registry.mu.Lock()
items := make([]map[string]any, 0, len(a.registry.data.Plugins))
for _, p := range a.registry.data.Plugins {
items = append(items, a.publicPlugin(p))
}
a.registry.mu.Unlock()
sort.Slice(items, func(i, j int) bool { return items[i]["plugin_id"].(string) < items[j]["plugin_id"].(string) })
writeJSON(w, http.StatusOK, map[string]any{"items": items})
}
func (a *app) marketplace(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
if _, _, ok := a.authenticate(w, r); !ok {
return
}
index, _, err := a.marketplaceConfig.loadIndex(r.Context())
if err != nil {
writeJSON(w, http.StatusBadGateway, map[string]string{"error": "marketplace is unavailable"})
return
}
a.registry.mu.Lock()
installed := make(map[string]pluginRecord, len(a.registry.data.Plugins))
for id, plugin := range a.registry.data.Plugins {
installed[id] = clonePluginRecord(plugin)
}
a.registry.mu.Unlock()
items := make([]map[string]any, 0, len(index.Entries))
coreVersion := a.currentCoreVersion(r.Context())
for _, entry := range index.Entries {
var plugin *pluginRecord
if value, ok := installed[entry.PluginID]; ok {
copy := value
plugin = &copy
}
item := publicMarketplaceEntry(entry, plugin)
compatibility := manifest.Manifest{CoreAPIBaseline: entry.CoreAPIBaseline, TestedCoreVersions: entry.TestedCoreVersions}.EvaluateCompatibility(coreVersion)
item["compatibility"] = compatibility
items = append(items, item)
}
writeJSON(w, http.StatusOK, map[string]any{
"schema_version": index.SchemaVersion,
"source": index.Source,
"issued_at": index.IssuedAt,
"expires_at": index.ExpiresAt,
"items": items,
})
}
func publicMarketplaceEntry(entry marketplaceEntry, installed *pluginRecord) map[string]any {
item := map[string]any{
"plugin_id": entry.PluginID,
"name": entry.Name,
"version": entry.Version,
"description": entry.Description,
"publisher_key_id": entry.PublisherKeyID,
"core_api_baseline": entry.CoreAPIBaseline,
"tested_core_versions": entry.TestedCoreVersions,
"capabilities": entry.Capabilities,
"archive_sha256": entry.ArchiveSHA256,
"archive_size": entry.ArchiveSize,
"release_notes": entry.ReleaseNotes,
"published_at": entry.PublishedAt,
"installed": installed != nil,
"installed_state": "",
"installed_status": "",
"installed_version": "",
"active_revision": "",
}
if installed != nil {
item["installed_state"] = installed.State
item["installed_status"] = installationStatus(*installed)
item["installed_version"] = installed.Manifest.Version
item["active_revision"] = installed.ActiveRevision
}
return item
}
func (a *app) marketplaceInstall(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
raw, err := captureRequestBody(r, 32<<10)
if err != nil {
writeJSON(w, http.StatusRequestEntityTooLarge, map[string]string{"error": "request body exceeds size limit"})
return
}
var input struct {
PluginID string `json:"plugin_id"`
Version string `json:"version"`
}
decoder := json.NewDecoder(bytes.NewReader(raw))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&input); err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "plugin_id and version are required"})
return
}
var trailing any
if err := decoder.Decode(&trailing); err != io.EOF || !marketplaceIDPattern.MatchString(strings.TrimSpace(input.PluginID)) || !marketplaceVersionPattern.MatchString(marketplaceEntryVersion(marketplaceEntry{Version: input.Version})) {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "plugin_id and version are invalid"})
return
}
input.PluginID = strings.TrimSpace(input.PluginID)
input.Version = marketplaceEntryVersion(marketplaceEntry{Version: input.Version})
r.Body = io.NopCloser(bytes.NewReader(raw))
op, s, ok := a.mutationAuthWithHash(w, r, "marketplace_install", input.PluginID, operationHashWithBody(r, raw))
if !ok {
return
}
var installed pluginRecord
if index, origin, loadErr := a.marketplaceConfig.loadIndex(r.Context()); loadErr != nil {
err = loadErr
} else {
var entry *marketplaceEntry
for i := range index.Entries {
candidate := &index.Entries[i]
if candidate.PluginID == input.PluginID && marketplaceEntryVersion(*candidate) == input.Version {
entry = candidate
break
}
}
if entry == nil {
err = errors.New("plugin version is not available in the marketplace")
} else {
var archive []byte
archive, err = a.marketplaceConfig.archiveBytes(r.Context(), *entry, origin)
if err == nil {
var info packageInfo
info, err = a.inspectPackage(archive)
if err == nil {
if info.Manifest.PluginID != entry.PluginID || strings.TrimPrefix(info.Manifest.Version, "v") != input.Version {
err = errors.New("marketplace package metadata does not match the catalog")
} else if info.Manifest.Name != entry.Name || !sameCapabilities(info.Manifest.Capabilities, entry.Capabilities) {
err = errors.New("marketplace package metadata does not match the catalog")
} else if !sameCoreVersions(info.Manifest.TestedCoreVersions, entry.TestedCoreVersions) {
err = errors.New("marketplace package Core test metadata does not match the catalog")
} else if info.Manifest.Publisher.KeyID != entry.PublisherKeyID {
err = errors.New("marketplace package publisher does not match the catalog")
} else if strings.TrimPrefix(strings.TrimPrefix(info.Manifest.CoreAPIBaseline, "sub2api-"), "v") != strings.TrimPrefix(strings.TrimPrefix(entry.CoreAPIBaseline, "sub2api-"), "v") {
err = errors.New("marketplace package Core baseline does not match the catalog")
} else if compat := info.Manifest.EvaluateCompatibility(a.currentCoreVersion(r.Context())); !compat.Compatible {
err = errors.New("marketplace package is incompatible with the current Core")
}
}
if err == nil {
installed, err = a.installPackage(info)
}
}
}
}
if err == nil {
op.Revision = installed.ActiveRevision
}
finished, persistErr := a.finalizeOperation(op, err, auditEvent{Time: time.Now().UTC(), Action: "marketplace_install", PluginID: input.PluginID, ActorID: s.User["id"], RequestID: requestID(r)})
if persistErr != nil {
writeOperationPersistenceError(w, finished)
return
}
a.operationResponse(w, finished)
}
func (a *app) operationByID(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
if _, _, ok := a.authenticate(w, r); !ok {
return
}
id := strings.TrimPrefix(r.URL.Path, "/api/operations/")
if id == "" || strings.ContainsAny(id, "/\\") {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid operation id"})
return
}
a.registry.mu.Lock()
defer a.registry.mu.Unlock()
for _, op := range a.registry.data.Operations {
if op.ID == id {
writeJSON(w, http.StatusOK, op)
return
}
}
writeJSON(w, http.StatusNotFound, map[string]string{"error": "operation not found"})
}
func (a *app) publicPlugin(p pluginRecord) map[string]any {
revisions := make([]map[string]any, 0, len(p.Revisions))
for _, rev := range p.Revisions {
revisions = append(revisions, map[string]any{"id": rev.ID, "version": rev.Version, "archive_sha256": rev.ArchiveSHA, "verified_at": rev.VerifiedAt, "healthy_at": rev.HealthyAt})
}
compat := p.Manifest.EvaluateCompatibility(a.currentCoreVersion(context.Background()))
return map[string]any{"plugin_id": p.Manifest.PluginID, "name": p.Manifest.Name, "version": p.Manifest.Version, "capabilities": p.Manifest.SortedCapabilities(), "state": p.State, "installation_status": installationStatus(p), "active_revision": p.ActiveRevision, "pending_revision": p.PendingRevision, "revisions": revisions, "compatibility": compat, "endpoint": p.Endpoint, "last_error": p.LastError, "updated_at": p.UpdatedAt, "menu": p.Manifest.UI.Menu}
}
func installationStatus(p pluginRecord) string {
switch p.State {
case "healthy", "enabled":
return "installed"
case "disabled":
for _, revision := range p.Revisions {
if !revision.HealthyAt.IsZero() {
return "stopped"
}
}
return "staged"
case "incompatible":
return "staged"
case "starting", "draining", "upgrading", "rollback_pending":
return "transitioning"
default:
return "failed"
}
}
func (a *app) currentCoreVersion(ctx context.Context) string {
if value := strings.TrimSpace(os.Getenv("CORE_VERSION")); value != "" {
return value
}
if a.core == nil {
return ""
}
out, err := a.core.publicSettings(ctx, token(12))
if err != nil {
return ""
}
data := envelopeData(out)
if value, ok := data["version"].(string); ok {
return strings.TrimSpace(value)
}
return ""
}
func (a *app) pluginIDFromPath(r *http.Request) (string, string, bool) {
parts := strings.Split(strings.TrimPrefix(strings.Trim(r.URL.Path, "/"), "api/plugins/"), "/")
if len(parts) < 1 || parts[0] == "" || strings.ContainsAny(parts[0], "/\\") || strings.Contains(parts[0], "..") {
return "", "", false
}
action := ""
if len(parts) > 1 {
action = parts[1]
}
return parts[0], action, true
}
func (a *app) getPlugin(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
if _, _, ok := a.authenticate(w, r); !ok {
return
}
id, action, ok := a.pluginIDFromPath(r)
if !ok || action != "" {
writeJSON(w, http.StatusNotFound, map[string]string{"error": "plugin not found"})
return
}
a.registry.mu.Lock()
p, found := a.registry.data.Plugins[id]
a.registry.mu.Unlock()
if !found {
writeJSON(w, http.StatusNotFound, map[string]string{"error": "plugin not found"})
return
}
writeJSON(w, http.StatusOK, a.publicPlugin(p))
}
func (a *app) operationResponse(w http.ResponseWriter, op operation) {
writeJSON(w, http.StatusAccepted, map[string]any{"operation_id": op.ID, "state": op.State, "error": op.Error, "warning": op.Warning})
}
func (a *app) finalizeOperation(op operation, operationErr error, event auditEvent) (operation, error) {
finished, finishErr := a.registry.finishOperation(op, operationErr)
if finishErr != nil {
return finished, finishErr
}
event.Operation = finished.ID
event.Result = finished.State
auditErr := a.registry.addAudit(event)
if auditErr != nil {
return finished, auditErr
}
return finished, nil
}
func writeOperationPersistenceError(w http.ResponseWriter, op operation) {
writeJSON(w, http.StatusInternalServerError, map[string]any{
"operation_id": op.ID,
"state": op.State,
"error": "operation registry is unavailable",
})
}
func (a *app) mutationAuth(w http.ResponseWriter, r *http.Request, kind, plugin string) (operation, session, bool) {
body, err := captureRequestBody(r, maxOperationBodyBytes)
if err != nil {
writeJSON(w, http.StatusRequestEntityTooLarge, map[string]string{"error": "request body exceeds size limit"})
return operation{}, session{}, false
}
return a.mutationAuthWithHash(w, r, kind, plugin, operationHashWithBody(r, body))
}
func (a *app) mutationAuthWithHash(w http.ResponseWriter, r *http.Request, kind, plugin, hash string) (operation, session, bool) {
_, s, ok := a.authenticate(w, r)
if !ok {
return operation{}, session{}, false
}
key := strings.TrimSpace(r.Header.Get("Idempotency-Key"))
if key == "" || len(key) > 128 {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "Idempotency-Key is required"})
return operation{}, session{}, false
}
op, existing, conflict, err := a.registry.operation(kind, plugin, key, s.User["id"], requestID(r), hash)
if err != nil {
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "operation registry is unavailable"})
return operation{}, session{}, false
}
if conflict {
writeJSON(w, http.StatusConflict, map[string]string{"error": "Idempotency-Key was already used for a different request"})
return operation{}, session{}, false
}
if existing {
a.operationResponse(w, op)
return operation{}, session{}, false
}
return op, s, true
}
func operationHash(r *http.Request) string {
return operationHashWithBody(r, nil)
}
func operationHashWithBody(r *http.Request, body []byte) string {
if r == nil {
return ""
}
bodySum := sha256.Sum256(body)
value := r.Method + "\n" + r.URL.Path + "\n" + r.URL.RawQuery + "\n" + hex.EncodeToString(bodySum[:])
sum := sha256.Sum256([]byte(value))
return hex.EncodeToString(sum[:])
}
func (a *app) stageRevisionPaths(revisions []revision) ([]stagedRevisionPath, error) {
base, err := filepath.Abs(filepath.Join(a.root, "installed"))
if err != nil {
return nil, err
}
moves := make([]stagedRevisionPath, 0, len(revisions))
restore := func() {
for i := len(moves) - 1; i >= 0; i-- {
_ = os.Rename(moves[i].staged, moves[i].original)
}
}
for _, rev := range revisions {
if strings.TrimSpace(rev.Path) == "" {
continue
}
original, err := filepath.Abs(rev.Path)
if err != nil {
restore()
return nil, err
}
rel, err := filepath.Rel(base, original)
if err != nil || rel == "." || rel == ".." || strings.HasPrefix(rel, ".."+string(os.PathSeparator)) {
restore()
return nil, errors.New("revision path is outside the plugin install root")
}
if _, err := os.Lstat(original); errors.Is(err, os.ErrNotExist) {
continue
} else if err != nil {
restore()
return nil, err
}
staged := original + ".uninstall-" + token(8)
if err := os.Rename(original, staged); err != nil {
restore()
return nil, err
}
moves = append(moves, stagedRevisionPath{original: original, staged: staged})
}
return moves, nil
}
func restoreRevisionPaths(moves []stagedRevisionPath) error {
var firstErr error
for i := len(moves) - 1; i >= 0; i-- {
if _, err := os.Lstat(moves[i].staged); errors.Is(err, os.ErrNotExist) {
continue
}
if err := os.Rename(moves[i].staged, moves[i].original); err != nil && firstErr == nil {
firstErr = err
}
}
return firstErr
}
func removeStagedRevisionPaths(moves []stagedRevisionPath) error {
var firstErr error
for _, move := range moves {
if err := os.RemoveAll(move.staged); err != nil && firstErr == nil {
firstErr = err
}
}
return firstErr
}
// captureRequestBody makes the server-side request fingerprint include the
// actual payload while restoring the body for the handler that still needs to
// decode it. The limit prevents an idempotency check from becoming an upload
// amplification vector.
func captureRequestBody(r *http.Request, limit int64) ([]byte, error) {
if r == nil || r.Body == nil {
return nil, nil
}
raw, err := io.ReadAll(io.LimitReader(r.Body, limit+1))
_ = r.Body.Close()
r.Body = io.NopCloser(bytes.NewReader(raw))
if err != nil {
return raw, err
}
if int64(len(raw)) > limit {
return raw, errors.New("request body exceeds size limit")
}
return raw, nil
}
type packageInfo struct {
Manifest manifest.Manifest
Raw []byte
Archive []byte
Files map[string][]byte
}
func (a *app) inspectPackage(raw []byte) (packageInfo, error) {
if len(raw) == 0 || len(raw) > maxPackageBytes {
return packageInfo{}, errors.New("package exceeds size limit")
}
hash := sha256.Sum256(raw)
_ = hash
zr, err := zip.NewReader(bytes.NewReader(raw), int64(len(raw)))
if err != nil {
return packageInfo{}, errors.New("invalid plugin package")
}
if len(zr.File) > maxPackageFiles {
return packageInfo{}, errors.New("package contains too many files")
}
files := map[string][]byte{}
var total int64
for _, file := range zr.File {
name := strings.ReplaceAll(file.Name, "\\", "/")
if file.FileInfo().IsDir() {
continue
}
if file.Mode()&os.ModeSymlink != 0 {
return packageInfo{}, fmt.Errorf("symlink entries are not allowed: %s", file.Name)
}
if name == "" || strings.HasPrefix(name, "/") || strings.Contains(name, "..") || path.Clean(name) != name || strings.Contains(name, "\\") {
return packageInfo{}, fmt.Errorf("unsafe package path: %s", file.Name)
}
if _, exists := files[name]; exists {
return packageInfo{}, fmt.Errorf("duplicate package path: %s", name)
}
if file.UncompressedSize64 > maxUncompressedBytes {
return packageInfo{}, errors.New("package member exceeds size limit")
}
rc, err := file.Open()
if err != nil {
return packageInfo{}, err
}
remaining := maxUncompressedBytes - total
if remaining <= 0 {
_ = rc.Close()
return packageInfo{}, errors.New("package exceeds uncompressed size limit")
}
data, readErr := io.ReadAll(io.LimitReader(rc, remaining+1))
_ = rc.Close()
if readErr != nil {
return packageInfo{}, readErr
}
total += int64(len(data))
if int64(len(data)) > remaining || total > maxUncompressedBytes {
return packageInfo{}, errors.New("package exceeds uncompressed size limit")
}
files[name] = data
}
manifestRaw, ok := files["manifest.json"]
if !ok {
return packageInfo{}, errors.New("manifest.json is required")
}
manifestPath := filepath.Join(os.TempDir(), "plugin-manifest-"+token(8)+".json")
defer os.Remove(manifestPath)
if err := os.WriteFile(manifestPath, manifestRaw, 0o600); err != nil {
return packageInfo{}, err
}
m, _, err := manifest.Load(manifestPath)
if err != nil {
return packageInfo{}, err
}
if signature, ok := files["signature.json"]; ok {
pub, trusted := a.trustedPublishers[m.Publisher.KeyID]
if !trusted {
return packageInfo{}, errors.New("plugin publisher is not trusted")
}
if err := manifest.VerifyKeyID(signature, m.Publisher.KeyID); err != nil {
return packageInfo{}, err
}
if err := manifest.VerifySignature(manifestRaw, signature, []byte(base64.StdEncoding.EncodeToString(pub))); err != nil {
return packageInfo{}, err
}
} else if !a.allowUnsigned {
return packageInfo{}, errors.New("signed plugin package is required")
}
for name, expected := range m.Files {
data, exists := files[name]
if !exists {
return packageInfo{}, fmt.Errorf("manifest file is missing: %s", name)
}
sum := sha256.Sum256(data)
if hex.EncodeToString(sum[:]) != expected {
return packageInfo{}, fmt.Errorf("file hash mismatch: %s", name)
}
}
for name := range files {
if name == "manifest.json" || name == "signature.json" {
continue
}
if _, declared := m.Files[name]; !declared {
return packageInfo{}, fmt.Errorf("package file is not declared in manifest: %s", name)
}
}
if len(m.Files) == 0 {
return packageInfo{}, errors.New("manifest.files must declare package files")
}
if _, ok := m.Files[m.UI.Entrypoint]; !ok {
entry := strings.Trim(strings.TrimSpace(m.UI.Entrypoint), "/")
if entry == "admin" || entry == "admin/index.html" {
if _, ok := m.Files["ui/index.html"]; !ok {
return packageInfo{}, errors.New("ui entrypoint requires ui/index.html")
}
} else {
return packageInfo{}, errors.New("ui.entrypoint is missing from files")
}
}
if m.Backend.Command != "" {
if _, ok := m.Files[m.Backend.Command]; !ok {
return packageInfo{}, errors.New("backend.command is missing from files")
}
}
return packageInfo{Manifest: m, Raw: manifestRaw, Archive: raw, Files: files}, nil
}
func (a *app) installPackage(info packageInfo) (pluginRecord, error) {
unlock := a.lockPlugin(info.Manifest.PluginID)
defer unlock()
return a.installPackageMode(info, false)
}
func (a *app) installPackageMode(info packageInfo, preserveActive bool) (pluginRecord, error) {
m := info.Manifest
compat := m.EvaluateCompatibility(a.currentCoreVersion(context.Background()))
state := "disabled"
if compat.Status == "incompatible" {
state = "incompatible"
} else if compat.Status == "untested" {
state = "disabled"
}
archiveHash := sha256.Sum256(info.Archive)
revisionID := time.Now().UTC().Format("20060102T150405.000000000Z") + "-" + hex.EncodeToString(archiveHash[:4])
dir := filepath.Join(a.root, "installed", m.PluginID, revisionID)
staging := dir + ".staging-" + token(8)
if err := os.MkdirAll(staging, 0o700); err != nil {
return pluginRecord{}, err
}
committed := false
defer func() {
if !committed {
_ = os.RemoveAll(staging)
}
}()
for name, data := range info.Files {
dest := filepath.Join(staging, filepath.FromSlash(name))
if !strings.HasPrefix(filepath.Clean(dest), filepath.Clean(staging)+string(os.PathSeparator)) {
return pluginRecord{}, errors.New("package path escaped staging directory")
}
if err := os.MkdirAll(filepath.Dir(dest), 0o700); err != nil {
return pluginRecord{}, err
}
mode := os.FileMode(0o600)
if strings.HasPrefix(name, "service/") {
mode = 0o700
}
if err := os.WriteFile(dest, data, mode); err != nil {
return pluginRecord{}, err
}
}
if err := os.WriteFile(filepath.Join(staging, ".verified"), []byte(revisionID+"\n"), 0o600); err != nil {
return pluginRecord{}, err
}
if err := os.MkdirAll(filepath.Dir(dir), 0o700); err != nil {
return pluginRecord{}, err
}
if err := os.Rename(staging, dir); err != nil {
return pluginRecord{}, err
}
committed = true
p := pluginRecord{Manifest: m, State: state, ActiveRevision: revisionID, Revisions: []revision{{ID: revisionID, Version: m.Version, Path: dir, ArchiveSHA: hex.EncodeToString(archiveHash[:]), Manifest: m, VerifiedAt: time.Now().UTC()}}, UpdatedAt: time.Now().UTC()}
a.registry.mu.Lock()
old, existed := a.registry.data.Plugins[m.PluginID]
if existed {
old = clonePluginRecord(old)
}
if existed {
if !preserveActive {
a.registry.mu.Unlock()
_ = os.RemoveAll(dir)
return pluginRecord{}, errors.New("plugin is already installed; use upgrade")
}
for i := range old.Revisions {
if old.Revisions[i].Manifest.PluginID == "" {
old.Revisions[i].Manifest = old.Manifest
}
}
p.Revisions = append(old.Revisions, p.Revisions...)
if preserveActive {
p.Manifest = old.Manifest
p.ActiveRevision = old.ActiveRevision
p.PendingRevision = revisionID
p.PreviousState = old.State
p.State = "upgrading"
} else {
p.ActiveRevision = revisionID
}
// An endpoint belongs to a service mode. A command revision gets a new
// managed port; an external revision must be configured explicitly when
// switching from a command revision.
if old.Manifest.Backend.Command == "" && m.Backend.Command == "" {
p.Endpoint = old.Endpoint
}
p.ConfigCipher = old.ConfigCipher
if preserveActive {
p.LastError = ""
}
}
a.registry.data.Plugins[m.PluginID] = p
err := a.registry.saveLocked()
if err != nil {
if existed {
a.registry.data.Plugins[m.PluginID] = old
} else {
delete(a.registry.data.Plugins, m.PluginID)
}
if restoreErr := a.registry.saveLocked(); restoreErr != nil {
err = fmt.Errorf("registry save failed: %v; restore failed: %w", err, restoreErr)
}
}
a.registry.mu.Unlock()
if err != nil {
_ = os.RemoveAll(dir)
return pluginRecord{}, err
}
return p, nil
}
func readUpload(r *http.Request) ([]byte, error) {
if strings.HasPrefix(strings.ToLower(r.Header.Get("Content-Type")), "multipart/form-data") {
if err := r.ParseMultipartForm(maxPackageBytes); err != nil {
return nil, err
}
file, _, err := r.FormFile("package")
if err != nil {
file, _, err = r.FormFile("artifact")
}
if err != nil {
return nil, errors.New("multipart field package is required")
}
defer file.Close()
return io.ReadAll(io.LimitReader(file, maxPackageBytes+1))
}
return io.ReadAll(io.LimitReader(r.Body, maxPackageBytes+1))
}
func (a *app) install(w http.ResponseWriter, r *http.Request, upgrade bool, routeID string) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
expectedID := strings.TrimSpace(routeID)
operationKind := map[bool]string{false: "install", true: "upgrade"}[upgrade]
raw, err := readUpload(r)
if err == nil && len(raw) > maxPackageBytes {
err = errors.New("package exceeds size limit")
}
op, s, ok := a.mutationAuthWithHash(w, r, operationKind, expectedID, operationHashWithBody(r, raw))
if !ok {
return
}
var unlock func()
if upgrade && expectedID != "" {
unlock = a.lockPlugin(expectedID)
defer unlock()
}
var p pluginRecord
previousState := "disabled"
var oldRecord pluginRecord
var oldExists bool
if upgrade && expectedID != "" {
a.registry.mu.Lock()
if old, exists := a.registry.data.Plugins[expectedID]; exists && old.State != "" {
previousState = old.State
oldRecord = clonePluginRecord(old)
oldExists = true
} else {
err = errors.New("plugin to upgrade was not found")
}
a.registry.mu.Unlock()
}
if err == nil {
info, inspectErr := a.inspectPackage(raw)
err = inspectErr
if err == nil && expectedID != "" && info.Manifest.PluginID != expectedID {
err = errors.New("upgrade plugin_id does not match route")
}
if err == nil && expectedID == "" {
if setErr := a.registry.setOperationPlugin(op.ID, info.Manifest.PluginID); setErr != nil {
err = setErr
}
}
if err == nil {
if upgrade {
p, err = a.installPackageMode(info, true)
if err == nil {
candidate := p
candidate.ActiveRevision = p.PendingRevision
candidateRevision := revisionByID(&p, p.PendingRevision)
candidate.Manifest = candidateRevision.Manifest
candidateKey := p.Manifest.PluginID + "#" + p.PendingRevision
candidateEndpoint, probeErr := a.startRevision(&candidate, p.PendingRevision, candidateKey)
err = probeErr
if err == nil {
if previousState != "healthy" {
a.stopProcess(candidateKey)
}
if previousState == "healthy" {
p.Endpoint = candidateEndpoint
}
p.Manifest = candidateRevision.Manifest
p.ActiveRevision = p.PendingRevision
p.PendingRevision = ""
p.PreviousState = ""
p.State = previousState
if p.State == "upgrading" || p.State == "starting" || p.State == "draining" {
p.State = "disabled"
}
p.LastError = ""
a.registry.mu.Lock()
p.UpdatedAt = time.Now().UTC()
a.registry.data.Plugins[p.Manifest.PluginID] = p
saveErr := a.registry.saveLocked()
if saveErr != nil {
if oldExists {
a.registry.data.Plugins[p.Manifest.PluginID] = oldRecord
} else {
delete(a.registry.data.Plugins, p.Manifest.PluginID)
}
if restoreErr := a.registry.saveLocked(); restoreErr != nil {
saveErr = fmt.Errorf("registry save failed: %v; restore failed: %w", saveErr, restoreErr)
}
}
a.registry.mu.Unlock()
err = saveErr
if err != nil {
a.stopProcess(candidateKey)
} else if previousState == "healthy" {
// Keep the old process alive until the registry commit has
// succeeded, then atomically switch the process handle.
a.stopPlugin(p.Manifest.PluginID)
a.promoteProcess(candidateKey, p.Manifest.PluginID)
}
} else {
p.State = "rollback_pending"
p.LastError = sanitizeError(err)
a.registry.mu.Lock()
a.registry.data.Plugins[p.Manifest.PluginID] = p
saveErr := a.registry.saveLocked()
if saveErr != nil {
if oldExists {
a.registry.data.Plugins[p.Manifest.PluginID] = oldRecord
} else {
delete(a.registry.data.Plugins, p.Manifest.PluginID)
}
restoreErr := a.registry.saveLocked()
err = fmt.Errorf("candidate probe failed: %v; registry save failed: %v", err, saveErr)
if restoreErr != nil {
err = fmt.Errorf("%v; restore failed: %w", err, restoreErr)
}
}
a.registry.mu.Unlock()
}
}
} else {
if unlock == nil {
unlock = a.lockPlugin(info.Manifest.PluginID)
defer unlock()
}
p, err = a.installPackageMode(info, false)
}
}
}
if err == nil {
op.Revision = p.ActiveRevision
}
pluginAuditID := p.Manifest.PluginID
if pluginAuditID == "" {
pluginAuditID = expectedID
}
op, persistErr := a.finalizeOperation(op, err, auditEvent{Time: time.Now().UTC(), Action: op.Kind, PluginID: pluginAuditID, ActorID: s.User["id"], RequestID: requestID(r)})
if persistErr != nil {
writeOperationPersistenceError(w, op)
return
}
a.operationResponse(w, op)
}
func (a *app) removeOwnMenu(ctx context.Context, access, id string, rid string) error {
if a.core == nil {
return errors.New("core unavailable")
}
settings, err := a.core.settings(ctx, access, rid)
if err != nil {
return err
}
data := envelopeData(settings)
raw := menuItemsFromSettings(data)
next := make([]any, 0, len(raw))
for _, item := range raw {
object, ok := item.(map[string]any)
if ok && object["id"] == id {
continue
}
next = append(next, item)
}
if len(next) == len(raw) {
return nil
}
_, err = a.core.updateMenu(ctx, access, rid, next)
return err
}
func (a *app) restoreOwnMenu(ctx context.Context, access string, p pluginRecord, rid string) error {
if a.core == nil {
return errors.New("core unavailable")
}
itemURL := p.Manifest.UI.Menu.URL
if itemURL == "" {
cfg, err := a.decryptConfig(p.ConfigCipher)
if err != nil {
return err
}
itemURL, _ = cfg["public_url"].(string)
}
if itemURL == "" || validateMenuURL(itemURL) != nil {
return errors.New("menu URL is not configured")
}
settings, err := a.core.settings(ctx, access, rid)
if err != nil {
return err
}
current := menuItemsFromSettings(envelopeData(settings))
candidate := map[string]any{"id": p.Manifest.UI.Menu.ID, "label": p.Manifest.UI.Menu.Label, "url": itemURL, "visibility": "admin", "sort_order": p.Manifest.UI.Menu.SortOrder}
replaced := false
for i, raw := range current {
if object, ok := raw.(map[string]any); ok && object["id"] == p.Manifest.UI.Menu.ID {
current[i] = candidate
replaced = true
}
}
if !replaced {
current = append(current, candidate)
}
_, err = a.core.updateMenu(ctx, access, rid, current)
return err
}
func (a *app) withPlugin(w http.ResponseWriter, r *http.Request, kind string, fn func(*pluginRecord, session) error) {
id, _, ok := a.pluginIDFromPath(r)
if !ok {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid plugin id"})
return
}
op, s, ok := a.mutationAuth(w, r, kind, id)
if !ok {
return
}
unlock := a.lockPlugin(id)
defer unlock()
a.registry.mu.Lock()
p, found := a.registry.data.Plugins[id]
a.registry.mu.Unlock()
if found {
p = clonePluginRecord(p)
}
old := clonePluginRecord(p)
var err error
var warning string
if !found {
err = errors.New("plugin not found")
} else {
err = fn(&p, s)
}
if err == nil && kind != "uninstall" {
a.registry.mu.Lock()
p.UpdatedAt = time.Now().UTC()
a.registry.data.Plugins[id] = p
saveErr := a.registry.saveLocked()
if saveErr != nil {
a.registry.data.Plugins[id] = old
if restoreErr := a.registry.saveLocked(); restoreErr != nil {
saveErr = fmt.Errorf("registry save failed: %v; restore failed: %w", saveErr, restoreErr)
}
}
a.registry.mu.Unlock()
err = saveErr
if err != nil {
a.restorePluginRuntime(old, p)
if kind == "disable" {
_ = a.restoreOwnMenu(r.Context(), s.AccessToken, old, requestID(r))
}
}
} else if err == nil && kind == "uninstall" {
// Move revision directories to private tombstones before changing the
// registry. This keeps both the registry and files recoverable if either
// the rename or registry commit fails.
moves, stageErr := a.stageRevisionPaths(old.Revisions)
err = stageErr
if err == nil {
a.registry.mu.Lock()
delete(a.registry.data.Plugins, id)
saveErr := a.registry.saveLocked()
if saveErr != nil {
a.registry.data.Plugins[id] = old
if restoreErr := a.registry.saveLocked(); restoreErr != nil {
saveErr = fmt.Errorf("registry save failed: %v; restore failed: %w", saveErr, restoreErr)
}
}
a.registry.mu.Unlock()
err = saveErr
if err != nil {
if restoreErr := restoreRevisionPaths(moves); restoreErr != nil {
err = fmt.Errorf("%v; revision restore failed: %w", err, restoreErr)
}
_ = a.restoreOwnMenu(r.Context(), s.AccessToken, old, requestID(r))
} else if cleanupErr := removeStagedRevisionPaths(moves); cleanupErr != nil {
// The registry commit is already complete. Report a warning while
// keeping the deletion completed; a later startup pass removes the
// private tombstones without requiring a second uninstall operation.
warning = sanitizeError(fmt.Errorf("plugin files remain pending cleanup: %w", cleanupErr))
}
} else {
_ = a.restoreOwnMenu(r.Context(), s.AccessToken, old, requestID(r))
}
}
op.Warning = warning
op, persistErr := a.finalizeOperation(op, err, auditEvent{Time: time.Now().UTC(), Action: kind, PluginID: id, ActorID: s.User["id"], RequestID: requestID(r)})
if persistErr != nil {
writeOperationPersistenceError(w, op)
return
}
a.operationResponse(w, op)
}
func (a *app) restorePluginRuntime(old, current pluginRecord) {
if current.State == "healthy" && old.State != "healthy" {
a.stopPlugin(current.Manifest.PluginID)
return
}
if current.State != "healthy" && old.State == "healthy" {
a.stopPlugin(current.Manifest.PluginID)
candidate := old
if err := a.startPlugin(&candidate); err == nil {
return
}
}
if current.State == "healthy" && old.State == "healthy" && current.ActiveRevision != old.ActiveRevision {
a.stopPlugin(current.Manifest.PluginID)
candidate := old
_ = a.startPlugin(&candidate)
}
}
func (a *app) enable(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
a.withPlugin(w, r, "enable", func(p *pluginRecord, _ session) error {
compat := p.Manifest.EvaluateCompatibility(a.currentCoreVersion(r.Context()))
if !compat.Compatible {
p.State = "incompatible"
return errors.New("plugin is incompatible with current Core")
}
if compat.Status == "untested" && r.URL.Query().Get("accept_untested") != "true" {
return errors.New("untested Core version requires explicit acceptance")
}
p.State = "starting"
if err := a.startPlugin(p); err != nil {
a.stopPlugin(p.Manifest.PluginID)
p.State = "error"
p.LastError = sanitizeError(err)
return err
}
p.State = "healthy"
p.LastError = ""
for i := range p.Revisions {
if p.Revisions[i].ID == p.ActiveRevision {
p.Revisions[i].HealthyAt = time.Now().UTC()
}
}
return nil
})
}
func (a *app) disable(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
a.withPlugin(w, r, "disable", func(p *pluginRecord, s session) error {
if err := a.removeOwnMenu(r.Context(), s.AccessToken, p.Manifest.PluginID, requestID(r)); err != nil {
return err
}
p.State = "draining"
a.stopPlugin(p.Manifest.PluginID)
p.State = "disabled"
return nil
})
}
func (a *app) rollback(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
a.withPlugin(w, r, "rollback", func(p *pluginRecord, _ session) error {
targetState := p.State
if p.PreviousState != "" {
targetState = p.PreviousState
}
// A failed upgrade leaves the previous active revision untouched and
// stores the new revision as pending. Rolling back in that state means
// discarding the failed candidate, not attempting to boot it again.
if p.PendingRevision != "" && p.ActiveRevision != "" {
pendingID := p.PendingRevision
if targetState == "healthy" {
if err := a.checkPluginHealth(p); err != nil {
p.LastError = sanitizeError(err)
return err
}
}
kept := p.Revisions[:0]
for i := range p.Revisions {
if p.Revisions[i].ID == pendingID {
// Keep the candidate until a maintenance pass can remove its
// files after this registry update has committed.
p.Revisions[i].Retired = true
}
kept = append(kept, p.Revisions[i])
}
p.Revisions = kept
p.PendingRevision = ""
p.PreviousState = ""
p.State = targetState
if p.State == "starting" || p.State == "draining" || p.State == "upgrading" || p.State == "rollback_pending" {
p.State = "disabled"
}
p.LastError = ""
return nil
}
if len(p.Revisions) < 2 {
p.State = "disabled"
return errors.New("no rollback revision retained")
}
current := p.ActiveRevision
var target revision
found := false
for i := len(p.Revisions) - 1; i >= 0; i-- {
if p.Revisions[i].Retired || p.Revisions[i].ID == current || (p.PendingRevision != "" && p.Revisions[i].ID == p.PendingRevision) {
continue
}
target = p.Revisions[i]
if target.Manifest.PluginID == "" {
target.Manifest = p.Manifest
}
found = true
break
}
if !found {
return errors.New("rollback revision not found")
}
candidate := *p
candidate.Manifest = target.Manifest
candidate.ActiveRevision = target.ID
candidateKey := p.Manifest.PluginID + "#rollback"
endpoint, err := a.startRevision(&candidate, target.ID, candidateKey)
if err != nil {
p.State = "rollback_pending"
p.LastError = sanitizeError(err)
return err
}
if targetState == "healthy" {
a.stopPlugin(p.Manifest.PluginID)
a.promoteProcess(candidateKey, p.Manifest.PluginID)
p.Endpoint = endpoint
} else {
a.stopProcess(candidateKey)
}
p.ActiveRevision = target.ID
p.Manifest = target.Manifest
p.PendingRevision = ""
p.PreviousState = ""
p.State = targetState
if p.State == "starting" || p.State == "draining" || p.State == "upgrading" || p.State == "rollback_pending" {
p.State = "disabled"
}
p.LastError = ""
return nil
})
}
func (a *app) uninstall(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost && r.Method != http.MethodDelete {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
a.withPlugin(w, r, "uninstall", func(p *pluginRecord, s session) error {
if p.State == "healthy" || p.State == "starting" || p.State == "draining" {
return errors.New("disable plugin before uninstall")
}
if err := a.removeOwnMenu(r.Context(), s.AccessToken, p.Manifest.PluginID, requestID(r)); err != nil {
return err
}
a.stopPlugin(p.Manifest.PluginID)
return nil
})
}
func (a *app) config(w http.ResponseWriter, r *http.Request) {
id, action, ok := a.pluginIDFromPath(r)
if !ok || action != "config" {
writeJSON(w, http.StatusNotFound, map[string]string{"error": "not found"})
return
}
if r.Method == http.MethodGet {
if _, _, ok := a.authenticate(w, r); !ok {
return
}
a.registry.mu.Lock()
p, found := a.registry.data.Plugins[id]
a.registry.mu.Unlock()
if !found {
writeJSON(w, http.StatusNotFound, map[string]string{"error": "plugin not found"})
return
}
cfg, cfgErr := a.decryptConfig(p.ConfigCipher)
if cfgErr != nil {
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "plugin configuration is unavailable"})
return
}
public := map[string]any{}
for key, value := range cfg {
if isSecretKey(key) {
public[key] = map[string]any{"configured": strings.TrimSpace(fmt.Sprint(value)) != ""}
continue
}
public[key] = value
}
writeJSON(w, http.StatusOK, map[string]any{"config": public, "endpoint": p.Endpoint})
return
}
if r.Method != http.MethodPut {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
a.withPlugin(w, r, "config", func(p *pluginRecord, _ session) error {
var cfg map[string]any
if err := decodeJSON(r, &cfg, maxJSONBytes); err != nil || cfg == nil {
return errors.New("config must be a JSON object")
}
for key := range cfg {
normalized := strings.ReplaceAll(strings.ReplaceAll(strings.ToLower(strings.TrimSpace(key)), "-", "_"), " ", "_")
if strings.Contains(normalized, "core_token") || strings.Contains(normalized, "core_access_token") || strings.Contains(normalized, "core_refresh_token") || strings.Contains(normalized, "admin_key") || strings.Contains(normalized, "admin_api_key") {
return errors.New("credential field is not accepted")
}
}
cipherText, err := a.encryptConfig(cfg)
if err != nil {
return err
}
p.ConfigCipher = cipherText
if p.Manifest.Backend.Command == "" {
endpoint, ok := cfg["service_url"].(string)
if !ok || strings.TrimSpace(endpoint) == "" {
return errors.New("service_url is required for external plugins")
}
if err := validateServiceURL(endpoint); err != nil {
return err
}
endpoint = strings.TrimRight(strings.TrimSpace(endpoint), "/")
if p.State == "healthy" && endpoint != p.Endpoint {
return errors.New("disable plugin before changing service_url")
}
p.Endpoint = endpoint
} else {
// Managed command plugins receive their endpoint from the supervisor.
p.Endpoint = ""
}
return nil
})
}
func isSecretKey(key string) bool {
lower := strings.ToLower(strings.ReplaceAll(strings.ReplaceAll(strings.TrimSpace(key), "-", "_"), " ", "_"))
compact := strings.ReplaceAll(lower, "_", "")
for _, marker := range []string{"secret", "password", "token", "apikey", "privatekey", "credential", "authorization", "cookie", "session", "key"} {
if strings.Contains(compact, marker) {
return true
}
}
return false
}
func (a *app) cipherKey() []byte {
if len(a.configKey) == 32 {
return a.configKey
}
return nil
}
func (a *app) encryptConfig(value map[string]any) (string, error) {
if len(a.cipherKey()) != 32 {
return "", errors.New("PLUGIN_CONFIG_KEY is required")
}
plain, err := json.Marshal(value)
if err != nil {
return "", err
}
block, err := aes.NewCipher(a.cipherKey())
if err != nil {
return "", err
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return "", err
}
nonce := make([]byte, gcm.NonceSize())
if _, err := rand.Read(nonce); err != nil {
return "", err
}
return base64.RawStdEncoding.EncodeToString(gcm.Seal(nonce, nonce, plain, nil)), nil
}
func (a *app) decryptConfig(encoded string) (map[string]any, error) {
if encoded == "" {
return map[string]any{}, nil
}
raw, err := base64.RawStdEncoding.DecodeString(encoded)
if err != nil {
return nil, err
}
if len(a.cipherKey()) != 32 {
return nil, errors.New("PLUGIN_CONFIG_KEY is required")
}
block, err := aes.NewCipher(a.cipherKey())
if err != nil {
return nil, err
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, err
}
if len(raw) < gcm.NonceSize() {
return nil, errors.New("invalid encrypted config")
}
plain, err := gcm.Open(nil, raw[:gcm.NonceSize()], raw[gcm.NonceSize():], nil)
if err != nil {
return nil, err
}
var out map[string]any
if err := json.Unmarshal(plain, &out); err != nil {
return nil, err
}
return out, nil
}
func validateServiceURL(raw string) error {
u, err := url.Parse(strings.TrimSpace(raw))
if err != nil || u.Host == "" || (u.Scheme != "http" && u.Scheme != "https") || u.User != nil || u.RawQuery != "" || u.Fragment != "" {
return errors.New("service_url must be an absolute http(s) origin without credentials or query")
}
if !isLoopbackHost(u.Hostname()) {
return errors.New("service_url must point to a loopback plugin service")
}
return nil
}
func validateMenuURL(raw string) error {
u, err := url.Parse(strings.TrimSpace(raw))
if err != nil || u.Host == "" || (u.Scheme != "http" && u.Scheme != "https") || u.User != nil || u.RawQuery != "" || u.Fragment != "" {
return errors.New("menu URL must be an absolute http(s) origin")
}
if u.Scheme == "http" && !isLoopbackHost(u.Hostname()) {
return errors.New("menu URL must use HTTPS unless loopback")
}
return nil
}
func (a *app) startPlugin(p *pluginRecord) error {
endpoint, err := a.startRevision(p, p.ActiveRevision, p.Manifest.PluginID)
if err != nil {
return err
}
p.Endpoint = endpoint
return nil
}
func (a *app) startRevision(p *pluginRecord, revisionID, processKey string) (string, error) {
if revisionID == "" {
return "", errors.New("revision is required")
}
rev := revisionByID(p, revisionID)
pluginManifest := p.Manifest
if rev.Manifest.PluginID != "" {
pluginManifest = rev.Manifest
}
endpoint := ""
if pluginManifest.Backend.Command != "" {
commandPath := filepath.Join(a.root, "installed", pluginManifest.PluginID, revisionID, filepath.FromSlash(pluginManifest.Backend.Command))
if rev.Path != "" {
commandPath = filepath.Join(rev.Path, filepath.FromSlash(pluginManifest.Backend.Command))
}
if _, err := os.Stat(commandPath); err != nil {
return "", err
}
port, err := freeLoopbackPort()
if err != nil {
return "", err
}
cmd := exec.Command(commandPath)
cmd.Dir = filepath.Dir(commandPath)
cmd.Env = []string{"PATH=/usr/bin:/bin", "HOME=" + filepath.Dir(commandPath), "LANG=C", "PLUGIN_ID=" + pluginManifest.PluginID, pluginManifest.Backend.ListenEnv + "=" + port}
cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
if err := cmd.Start(); err != nil {
return "", err
}
endpoint = "http://127.0.0.1:" + port
a.mu.Lock()
a.processes[processKey] = cmd
a.mu.Unlock()
} else {
// External services are owned by their deployer. Never inherit a
// command-process endpoint or a stale endpoint across a mode switch.
endpoint = p.Endpoint
}
probe := *p
probe.Manifest = pluginManifest
probe.ActiveRevision = revisionID
probe.Endpoint = endpoint
var probeErr error
attempts := 1
if pluginManifest.Backend.Command != "" {
// A freshly spawned process can need a short interval before binding its
// port. Retry the probe within the bounded startup window.
attempts = 20
}
for attempt := 0; attempt < attempts; attempt++ {
probeErr = a.checkPluginHealth(&probe)
if probeErr == nil {
break
}
if attempt+1 < attempts {
time.Sleep(50 * time.Millisecond)
}
}
if probeErr != nil {
a.stopProcess(processKey)
return "", probeErr
}
return endpoint, nil
}
func revisionByID(p *pluginRecord, id string) revision {
for _, rev := range p.Revisions {
if rev.ID == id {
if rev.Manifest.PluginID == "" {
rev.Manifest = p.Manifest
}
return rev
}
}
return revision{ID: id, Manifest: p.Manifest}
}
func freeLoopbackPort() (string, error) {
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
return "", err
}
port := strconv.Itoa(listener.Addr().(*net.TCPAddr).Port)
if err := listener.Close(); err != nil {
return "", err
}
return port, nil
}
func (a *app) checkPluginHealth(p *pluginRecord) error {
if p.Endpoint == "" {
return errors.New("service_url is required before enabling")
}
u, err := url.Parse(p.Endpoint)
if err != nil {
return err
}
if err := validateServiceURL(p.Endpoint); err != nil {
return err
}
health := strings.TrimRight(p.Endpoint, "/") + "/" + strings.Trim(strings.TrimSpace(p.Manifest.Backend.HealthPath), "/")
req, err := http.NewRequest(http.MethodGet, health, nil)
if err != nil {
return err
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
req = req.WithContext(ctx)
client := &http.Client{Timeout: 5 * time.Second, Transport: &http.Transport{Proxy: nil}, CheckRedirect: func(_ *http.Request, _ []*http.Request) error { return http.ErrUseLastResponse }}
res, err := client.Do(req)
if err != nil {
return err
}
defer res.Body.Close()
body, err := io.ReadAll(io.LimitReader(res.Body, 64<<10))
if err != nil {
return err
}
_ = u
if res.StatusCode < 200 || res.StatusCode >= 300 {
return fmt.Errorf("plugin health check returned %d", res.StatusCode)
}
if version := responseVersion(body); version != "" && strings.TrimPrefix(version, "v") != strings.TrimPrefix(p.Manifest.Version, "v") {
return fmt.Errorf("plugin version mismatch: got %s", version)
}
ready := strings.TrimRight(p.Endpoint, "/") + "/" + strings.Trim(strings.TrimSpace(p.Manifest.Backend.ReadinessPath), "/")
readyReq, err := http.NewRequestWithContext(ctx, http.MethodGet, ready, nil)
if err != nil {
return err
}
readyRes, err := client.Do(readyReq)
if err != nil {
return err
}
defer readyRes.Body.Close()
if readyRes.StatusCode < 200 || readyRes.StatusCode >= 300 {
return fmt.Errorf("plugin readiness check returned %d", readyRes.StatusCode)
}
return nil
}
func responseVersion(body []byte) string {
var value map[string]any
if json.Unmarshal(body, &value) != nil {
return ""
}
version, _ := value["version"].(string)
return strings.TrimSpace(version)
}
func (a *app) stopPlugin(id string) {
a.stopProcess(id)
}
func (a *app) stopProcess(id string) {
a.mu.Lock()
cmd := a.processes[id]
delete(a.processes, id)
a.mu.Unlock()
if cmd != nil && cmd.Process != nil {
// Kill the process group so a plugin cannot leave a child server behind
// after disable, upgrade, rollback or control-plane shutdown.
if err := syscall.Kill(-cmd.Process.Pid, syscall.SIGTERM); err != nil {
_ = cmd.Process.Signal(syscall.SIGTERM)
}
finished := make(chan struct{})
go func() { _, _ = cmd.Process.Wait(); close(finished) }()
select {
case <-finished:
case <-time.After(10 * time.Second):
_ = cmd.Process.Kill()
<-finished
}
}
}
func (a *app) shutdown() {
a.mu.Lock()
ids := make([]string, 0, len(a.processes))
for id := range a.processes {
ids = append(ids, id)
}
a.mu.Unlock()
for _, id := range ids {
a.stopProcess(id)
}
}
func (a *app) promoteProcess(from, to string) {
a.mu.Lock()
if cmd, ok := a.processes[from]; ok {
a.processes[to] = cmd
delete(a.processes, from)
}
a.mu.Unlock()
}
// recoverPlugins reconciles persisted lifecycle state with the processes and
// external endpoints that exist after a control-plane restart. A healthy
// command plugin is started again; an external plugin is only probed. Any
// interrupted transition is failed closed instead of being advertised as
// healthy with no corresponding runtime.
func (a *app) recoverPlugins() error {
a.registry.mu.Lock()
// A prior delete may have committed the registry before the filesystem
// cleanup failed. Remove only tombstones whose original revision is no
// longer referenced; an interrupted delete still needs its tombstone to
// recover the registry record safely.
_ = a.cleanupUninstallTombstonesLocked()
ids := make([]string, 0, len(a.registry.data.Plugins))
for id := range a.registry.data.Plugins {
ids = append(ids, id)
}
a.registry.mu.Unlock()
sort.Strings(ids)
for _, id := range ids {
unlock := a.lockPlugin(id)
a.registry.mu.Lock()
p, ok := a.registry.data.Plugins[id]
a.registry.mu.Unlock()
if !ok {
unlock()
continue
}
old := clonePluginRecord(p)
var err error
switch p.State {
case "healthy", "enabled":
compat := p.Manifest.EvaluateCompatibility(a.currentCoreVersion(context.Background()))
if !compat.Compatible {
err = errors.New("plugin is incompatible with current Core")
} else if p.ActiveRevision == "" {
err = errors.New("active revision is missing")
} else {
candidate := p
endpoint, probeErr := a.startRevision(&candidate, p.ActiveRevision, p.Manifest.PluginID)
err = probeErr
if err == nil {
p.Endpoint = endpoint
p.State = "healthy"
p.LastError = ""
for i := range p.Revisions {
if p.Revisions[i].ID == p.ActiveRevision {
p.Revisions[i].HealthyAt = time.Now().UTC()
}
}
}
}
if err != nil {
p.State = "error"
p.LastError = sanitizeError(err)
p.Endpoint = ""
}
case "starting", "draining", "upgrading", "rollback_pending":
a.stopPlugin(p.Manifest.PluginID)
p.State = "error"
p.LastError = "control-plane restart interrupted plugin operation"
default:
unlock()
continue
}
p.UpdatedAt = time.Now().UTC()
a.registry.mu.Lock()
a.registry.data.Plugins[id] = p
saveErr := a.registry.saveLocked()
if saveErr != nil {
a.registry.data.Plugins[id] = old
if restoreErr := a.registry.saveLocked(); restoreErr != nil {
saveErr = fmt.Errorf("registry save failed: %v; restore failed: %w", saveErr, restoreErr)
}
}
a.registry.mu.Unlock()
unlock()
if saveErr != nil {
if p.Manifest.Backend.Command != "" {
a.stopPlugin(p.Manifest.PluginID)
}
return saveErr
}
}
return nil
}
func (a *app) cleanupUninstallTombstonesLocked() error {
live := map[string]struct{}{}
for _, plugin := range a.registry.data.Plugins {
for _, revision := range plugin.Revisions {
if revision.Path == "" {
continue
}
if absolute, err := filepath.Abs(revision.Path); err == nil {
live[filepath.Clean(absolute)] = struct{}{}
}
}
}
installedRoot := filepath.Join(a.root, "installed")
pluginDirs, err := os.ReadDir(installedRoot)
if errors.Is(err, os.ErrNotExist) {
return nil
}
if err != nil {
return err
}
var firstErr error
for _, pluginDir := range pluginDirs {
if !pluginDir.IsDir() {
continue
}
entries, readErr := os.ReadDir(filepath.Join(installedRoot, pluginDir.Name()))
if readErr != nil {
if firstErr == nil {
firstErr = readErr
}
continue
}
for _, entry := range entries {
if !entry.IsDir() || !strings.Contains(entry.Name(), ".uninstall-") {
continue
}
originalName := strings.SplitN(entry.Name(), ".uninstall-", 2)[0]
originalPath, pathErr := filepath.Abs(filepath.Join(installedRoot, pluginDir.Name(), originalName))
if pathErr == nil {
if _, referenced := live[filepath.Clean(originalPath)]; referenced {
if _, statErr := os.Lstat(originalPath); errors.Is(statErr, os.ErrNotExist) {
if restoreErr := os.Rename(filepath.Join(installedRoot, pluginDir.Name(), entry.Name()), originalPath); restoreErr != nil && firstErr == nil {
firstErr = restoreErr
}
} else if statErr == nil {
if removeErr := os.RemoveAll(filepath.Join(installedRoot, pluginDir.Name(), entry.Name())); removeErr != nil && firstErr == nil {
firstErr = removeErr
}
}
continue
}
}
if removeErr := os.RemoveAll(filepath.Join(installedRoot, pluginDir.Name(), entry.Name())); removeErr != nil && firstErr == nil {
firstErr = removeErr
}
}
}
return firstErr
}
func (a *app) audit(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
if _, _, ok := a.authenticate(w, r); !ok {
return
}
a.registry.mu.Lock()
items := append([]auditEvent(nil), a.registry.data.Audit...)
a.registry.mu.Unlock()
writeJSON(w, http.StatusOK, map[string]any{"items": items})
}
func (a *app) menu(w http.ResponseWriter, r *http.Request, apply bool) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
id, _, ok := a.pluginIDFromPath(r)
if !ok {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid plugin id"})
return
}
kind := "menu_preview"
if apply {
kind = "menu_apply"
}
op, s, ok := a.mutationAuth(w, r, kind, id)
if !ok {
return
}
unlock := a.lockPlugin(id)
defer unlock()
a.registry.mu.Lock()
p, found := a.registry.data.Plugins[id]
a.registry.mu.Unlock()
var err error
var current []any
if !found {
err = errors.New("plugin not found")
}
if err == nil && apply && (p.State != "healthy" && p.State != "enabled") {
err = errors.New("plugin must be healthy before menu injection")
}
if err == nil {
settings, callErr := a.core.settings(r.Context(), s.AccessToken, requestID(r))
if callErr != nil {
err = callErr
} else {
data := envelopeData(settings)
current = menuItemsFromSettings(data)
}
}
itemURL := p.Manifest.UI.Menu.URL
if itemURL == "" {
cfg, cfgErr := a.decryptConfig(p.ConfigCipher)
if cfgErr != nil {
err = errors.New("plugin configuration is unavailable")
} else if configured, ok := cfg["public_url"].(string); ok {
itemURL = configured
}
}
if itemURL == "" {
err = errors.New("menu URL is not configured")
} else if validateMenuURL(itemURL) != nil {
err = errors.New("menu URL must be an absolute http(s) URL")
}
next := append([]any(nil), current...)
if err == nil {
candidate := map[string]any{"id": p.Manifest.UI.Menu.ID, "label": p.Manifest.UI.Menu.Label, "url": itemURL, "visibility": "admin", "sort_order": p.Manifest.UI.Menu.SortOrder}
replaced := false
for i, raw := range next {
if object, ok := raw.(map[string]any); ok && object["id"] == p.Manifest.UI.Menu.ID {
next[i] = candidate
replaced = true
}
}
if !replaced {
next = append(next, candidate)
}
}
if apply && err == nil {
_, err = a.core.updateMenu(r.Context(), s.AccessToken, requestID(r), next)
}
op, persistErr := a.finalizeOperation(op, err, auditEvent{Time: time.Now().UTC(), Action: kind, PluginID: id, ActorID: s.User["id"], RequestID: requestID(r)})
if persistErr != nil {
writeOperationPersistenceError(w, op)
return
}
if apply {
a.operationResponse(w, op)
return
}
writeJSON(w, http.StatusOK, map[string]any{"operation_id": op.ID, "current": current, "next": next, "state": op.State, "error": op.Error})
}
func (a *app) menuGlobal(w http.ResponseWriter, r *http.Request, apply bool) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
var input struct {
PluginID string `json:"plugin_id"`
}
raw, bodyErr := captureRequestBody(r, 32<<10)
if bodyErr != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request body"})
return
}
dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields()
var trailing any
if err := dec.Decode(&input); err != nil || dec.Decode(&trailing) != io.EOF || strings.TrimSpace(input.PluginID) == "" || strings.ContainsAny(input.PluginID, "/\\") {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "plugin_id is required"})
return
}
// a.menu performs its own authenticated idempotency check and needs the
// original body available for its server-side request fingerprint.
r.Body = io.NopCloser(bytes.NewReader(raw))
pathAction := "menu-preview"
if apply {
pathAction = "menu-apply"
}
r.URL.Path = "/api/plugins/" + input.PluginID + "/" + pathAction
a.menu(w, r, apply)
}
func (a *app) coreError(w http.ResponseWriter, err error, fallback string) {
status := http.StatusBadGateway
if ce, ok := err.(*coreError); ok && ce.status >= 400 && ce.status < 600 {
status = ce.status
}
writeJSON(w, status, map[string]string{"error": fallback})
}
func (a *app) static(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/admin" {
http.Redirect(w, r, a.publicBasePath+"/admin/", http.StatusPermanentRedirect)
return
}
if r.URL.Path == "/" || r.URL.Path == "/admin/" {
data, err := uiFS.ReadFile("ui/index.html")
if err != nil {
http.Error(w, "ui unavailable", 500)
return
}
baseJSON, _ := json.Marshal(a.publicBasePath)
data = []byte(strings.ReplaceAll(string(data), "__PLUGIN_BASE_PATH_JSON__", string(baseJSON)))
w.Header().Set("Cache-Control", "no-store")
w.Header().Set("Content-Type", "text/html; charset=utf-8")
_, _ = w.Write(data)
return
}
for _, name := range []string{"app.js", "styles.css"} {
if r.URL.Path == "/"+name {
data, err := uiFS.ReadFile("ui/" + name)
if err != nil {
http.NotFound(w, r)
return
}
if strings.HasSuffix(name, ".js") {
w.Header().Set("Content-Type", "text/javascript; charset=utf-8")
} else {
w.Header().Set("Content-Type", "text/css; charset=utf-8")
}
_, _ = w.Write(data)
return
}
}
http.NotFound(w, r)
}
func (a *app) routes() http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("/healthz", a.health)
mux.HandleFunc("/readyz", a.ready)
mux.HandleFunc("/login", a.login)
mux.HandleFunc("/login/2fa", a.login2FA)
mux.HandleFunc("/logout", a.logout)
mux.HandleFunc("/api/me", a.me)
mux.HandleFunc("/api/audit", a.audit)
mux.HandleFunc("/api/marketplace", a.marketplace)
mux.HandleFunc("/api/marketplace/install", a.marketplaceInstall)
mux.HandleFunc("/api/menu-items/preview", func(w http.ResponseWriter, r *http.Request) { a.menuGlobal(w, r, false) })
mux.HandleFunc("/api/menu-items/apply", func(w http.ResponseWriter, r *http.Request) { a.menuGlobal(w, r, true) })
mux.HandleFunc("/api/operations/", a.operationByID)
mux.HandleFunc("/api/plugins", a.apiPlugins)
mux.HandleFunc("/api/plugins/install", func(w http.ResponseWriter, r *http.Request) { a.install(w, r, false, "") })
mux.HandleFunc("/api/plugins/", func(w http.ResponseWriter, r *http.Request) {
id, action, ok := a.pluginIDFromPath(r)
if !ok {
http.NotFound(w, r)
return
}
if action == "" {
if r.Method == http.MethodDelete {
a.uninstall(w, r)
} else {
a.getPlugin(w, r)
}
return
}
switch action {
case "install":
a.install(w, r, false, id)
case "upgrade":
a.install(w, r, true, id)
case "enable":
a.enable(w, r)
case "disable":
a.disable(w, r)
case "rollback":
a.rollback(w, r)
case "uninstall":
a.uninstall(w, r)
case "delete":
a.uninstall(w, r)
case "config":
a.config(w, r)
case "menu-preview":
a.menu(w, r, false)
case "menu-apply":
a.menu(w, r, true)
default:
http.NotFound(w, r)
}
})
mux.HandleFunc("/", a.static)
return a.securityHeaders(requestIDMiddleware(mux))
}
func requestIDMiddleware(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
rid := requestID(r)
w.Header().Set("X-Request-Id", rid)
next.ServeHTTP(w, r)
})
}
func (a *app) securityHeaders(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("X-Content-Type-Options", "nosniff")
w.Header().Set("Referrer-Policy", "no-referrer")
w.Header().Set("Content-Security-Policy", "default-src 'self'; script-src 'self'; style-src 'self'; connect-src 'self'; img-src 'self' data:; frame-ancestors "+strings.Join(a.frameAncestors, " ")+"; base-uri 'self'; form-action 'self'")
next.ServeHTTP(w, r)
})
}
func parseSameSite(value string) (http.SameSite, error) {
switch strings.ToLower(strings.TrimSpace(value)) {
case "", "lax":
return http.SameSiteLaxMode, nil
case "strict":
return http.SameSiteStrictMode, nil
case "none":
return http.SameSiteNoneMode, nil
default:
return 0, errors.New("PLUGIN_COOKIE_SAMESITE must be lax, strict or none")
}
}
func loadTrustedPublishers(raw string) map[string][]byte {
out := map[string][]byte{}
var values map[string]string
if json.Unmarshal([]byte(raw), &values) != nil {
return out
}
for key, encoded := range values {
if data, err := base64.StdEncoding.DecodeString(encoded); err == nil && len(data) == 32 {
out[key] = data
}
}
return out
}
func parseFrameAncestors(raw string) []string {
values := strings.Fields(raw)
if len(values) == 0 {
return []string{"'self'"}
}
allowed := make([]string, 0, len(values))
for _, value := range values {
if value == "'self'" || value == "'none'" {
allowed = append(allowed, value)
continue
}
u, err := url.Parse(value)
if err == nil && (u.Scheme == "http" || u.Scheme == "https") && u.Host != "" && u.User == nil && u.Path == "" && u.RawQuery == "" && u.Fragment == "" {
allowed = append(allowed, value)
}
}
if len(allowed) == 0 {
return []string{"'self'"}
}
return allowed
}
func deriveConfigKey(raw string) ([]byte, error) {
raw = strings.TrimSpace(raw)
if len(raw) < 32 {
return nil, errors.New("PLUGIN_CONFIG_KEY must contain at least 32 characters")
}
hash := sha256.Sum256([]byte(raw))
return hash[:], nil
}
func main() {
host := os.Getenv("PLUGIN_HOST")
if host == "" {
host = "127.0.0.1"
}
port := os.Getenv("PLUGIN_PORT")
if port == "" {
port = "8090"
}
coreURL := os.Getenv("CORE_BASE_URL")
if coreURL == "" {
coreURL = "http://127.0.0.1:8080"
}
core, err := newCoreClient(coreURL)
if err != nil {
slog.Error("invalid Core URL", "error", err)
os.Exit(2)
}
registryDir := os.Getenv("PLUGIN_REGISTRY_DIR")
if registryDir == "" {
registryDir = "./data"
}
reg, err := openRegistry(registryDir)
if err != nil {
slog.Error("open registry", "error", err)
os.Exit(2)
}
a := newApp(core, reg, registryDir)
environment := strings.ToLower(strings.TrimSpace(os.Getenv("PLUGIN_ENV")))
if environment == "" {
environment = "production"
}
a.publicBasePath = strings.TrimRight(strings.TrimSpace(os.Getenv("PLUGIN_PUBLIC_BASE_PATH")), "/")
a.cookiePath = os.Getenv("PLUGIN_COOKIE_PATH")
if a.cookiePath == "" {
a.cookiePath = "/"
}
a.cookieSecure = strings.EqualFold(os.Getenv("PLUGIN_COOKIE_SECURE"), "true")
if !isLoopbackHost(host) && !a.cookieSecure {
slog.Error("PLUGIN_COOKIE_SECURE must be true for non-loopback listeners")
os.Exit(2)
}
if sameSite, parseErr := parseSameSite(os.Getenv("PLUGIN_COOKIE_SAMESITE")); parseErr == nil {
a.cookieSameSite = sameSite
} else {
slog.Error("invalid cookie same site", "error", parseErr)
os.Exit(2)
}
if a.cookieSameSite == http.SameSiteNoneMode && !a.cookieSecure {
slog.Error("SameSite=None requires secure cookie")
os.Exit(2)
}
a.allowUnsigned = strings.EqualFold(os.Getenv("PLUGIN_ALLOW_UNSIGNED"), "true")
if a.allowUnsigned && (environment != "development" || !isLoopbackHost(host)) {
slog.Error("unsigned packages are allowed only in development on loopback")
os.Exit(2)
}
if key, keyErr := deriveConfigKey(os.Getenv("PLUGIN_CONFIG_KEY")); keyErr != nil {
slog.Error("invalid config encryption key", "error", keyErr)
os.Exit(2)
} else {
a.configKey = key
}
a.frameAncestors = parseFrameAncestors(os.Getenv("PLUGIN_FRAME_ANCESTORS"))
a.trustedPublishers = loadTrustedPublishers(os.Getenv("PLUGIN_TRUSTED_PUBLISHERS"))
marketplaceSource := strings.TrimSpace(os.Getenv("PLUGIN_MARKETPLACE_INDEX"))
if marketplaceSource == "" {
marketplaceSource = filepath.Join(registryDir, "marketplace", "index.json")
}
allowLoopbackMarketplace := environment == "development" && isLoopbackHost(host)
marketplace, marketplaceErr := newMarketplaceService(marketplaceSource, os.Getenv("PLUGIN_MARKETPLACE_ALLOWED_HOSTS"), allowLoopbackMarketplace)
if marketplaceErr != nil {
slog.Error("invalid marketplace configuration", "error", marketplaceErr)
os.Exit(2)
}
a.marketplaceConfig = marketplace
if err := a.recoverPlugins(); err != nil {
slog.Error("recover plugins", "error", err)
os.Exit(2)
}
addr := net.JoinHostPort(host, port)
srv := &http.Server{Addr: addr, Handler: a.routes(), ReadHeaderTimeout: 10 * time.Second, ReadTimeout: 30 * time.Second, WriteTimeout: 30 * time.Second, IdleTimeout: 60 * time.Second}
slog.Info("plugin-admin listening", "addr", addr)
stopSignals := make(chan os.Signal, 1)
signal.Notify(stopSignals, os.Interrupt, syscall.SIGTERM)
defer signal.Stop(stopSignals)
go func() {
<-stopSignals
shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
_ = srv.Shutdown(shutdownCtx)
a.shutdown()
}()
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
slog.Error("plugin-admin stopped", "error", err)
os.Exit(1)
}
}
// Keep strconv linked for manifests that use numeric listen settings in future
// revisions; this also makes the validation helper easy to extend without API churn.
var _ = strconv.IntSize