Files
sub2api-add/plugins/plugin-admin/main.go
T
Qiufeng ada4ab3c21
Business Plugins CI / check (plugin-admin) (push) Successful in 1m42s
Business Plugins CI / check (subscription-admin) (push) Successful in 1m30s
feat: complete unified plugin admin v1.1.0
2026-08-30 12:10:04 +08:00

3763 lines
122 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"
"html"
"io"
"log/slog"
"mime"
"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"
subscriptionPluginID = "qiu.subscription-admin"
pluginVersion = "1.1.0"
sessionCookieName = "plugin_admin_session"
maxJSONBytes = 2 << 20
maxPackageBytes = 128 << 20
maxPackageFiles = 512
maxUncompressedBytes = 256 << 20
defaultSessionTTL = 30 * time.Minute
defaultSessionMaxTTL = 8 * time.Hour
loginWindow = time.Minute
loginLimit = 10
maxPendingLogins = 1024
pendingLoginTTL = 5 * time.Minute
defaultOperationLimit = 200
maxOperationBodyBytes = maxPackageBytes + (4 << 20)
)
//go:embed ui/*
var uiFS embed.FS
var requestIDPattern = regexp.MustCompile(`^[A-Za-z0-9._:-]{1,64}$`)
var coreMenuIDPattern = regexp.MustCompile(`^[A-Za-z0-9_-]+$`)
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, clientIPs ...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)
if len(clientIPs) > 0 && net.ParseIP(strings.TrimSpace(clientIPs[0])) != nil {
clientIP := strings.TrimSpace(clientIPs[0])
req.Header.Set("X-Forwarded-For", clientIP)
req.Header.Set("X-Real-IP", clientIP)
}
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 {
parsed, err := url.Parse(p)
if err != nil || parsed.Host != "" || parsed.Scheme != "" || parsed.Fragment != "" || strings.Contains(parsed.Path, "\\") {
return false
}
pathName := parsed.Path
if method == http.MethodGet {
if pathName == "/api/v1/auth/me" || pathName == "/api/v1/settings/public" || pathName == "/api/v1/admin/settings" {
return len(parsed.Query()) == 0
}
return allowedSubscriptionReadPath(pathName) && validSubscriptionQuery(parsed.Query())
}
if len(parsed.Query()) != 0 {
return false
}
return method == http.MethodPut && pathName == "/api/v1/admin/settings" || method == http.MethodPost && (pathName == "/api/v1/auth/login" || pathName == "/api/v1/auth/login/2fa" || pathName == "/api/v1/auth/refresh" || pathName == "/api/v1/auth/logout")
}
func positiveID(value string) bool {
value = strings.TrimSpace(value)
if value == "" || strings.ContainsAny(value, "/\\") {
return false
}
n, err := strconv.ParseInt(value, 10, 64)
return err == nil && n > 0
}
func allowedSubscriptionReadPath(pathName string) bool {
if pathName == "/api/v1/admin/payment/plans" || pathName == "/api/v1/admin/subscriptions" {
return true
}
if strings.HasPrefix(pathName, "/api/v1/admin/subscriptions/") {
return positiveID(strings.TrimPrefix(pathName, "/api/v1/admin/subscriptions/"))
}
if strings.HasPrefix(pathName, "/api/v1/admin/users/") {
rest := strings.TrimPrefix(pathName, "/api/v1/admin/users/")
if positiveID(rest) {
return true
}
parts := strings.Split(rest, "/")
return len(parts) == 2 && positiveID(parts[0]) && parts[1] == "subscriptions"
}
return false
}
func validSubscriptionQuery(values url.Values) bool {
allowed := map[string]bool{"page": true, "page_size": true, "limit": true, "user_id": true, "group_id": true, "status": true, "platform": true, "sort_by": true, "sort_order": true}
for key, entries := range values {
if !allowed[key] {
return false
}
for _, value := range entries {
value = strings.TrimSpace(value)
if value == "" || len(value) > 100 {
return false
}
if key == "page" || key == "page_size" || key == "limit" || key == "user_id" || key == "group_id" {
n, err := strconv.ParseUint(value, 10, 63)
if err != nil || n == 0 || ((key == "page_size" || key == "limit") && n > 100) {
return false
}
}
}
}
return true
}
func sanitizeSubscriptionQuery(values url.Values) url.Values {
allowed := []string{"page", "page_size", "limit", "user_id", "group_id", "status", "platform", "sort_by", "sort_order"}
out := url.Values{}
for _, key := range allowed {
for _, value := range values[key] {
value = strings.TrimSpace(value)
if value == "" || len(value) > 100 {
continue
}
if (key == "page" || key == "page_size" || key == "limit" || key == "user_id" || key == "group_id") && !positiveID(value) {
continue
}
out.Add(key, value)
}
}
return out
}
func (c *coreClient) read(ctx context.Context, requestPath, accessToken, requestID string) (coreEnvelope, error) {
parsed, err := url.Parse(requestPath)
if err != nil || parsed.Host != "" || parsed.Scheme != "" || parsed.Fragment != "" || !allowedSubscriptionReadPath(parsed.Path) || !validSubscriptionQuery(parsed.Query()) {
return coreEnvelope{}, errors.New("core path is not in the subscription allowlist")
}
parsed.RawQuery = sanitizeSubscriptionQuery(parsed.Query()).Encode()
return c.call(ctx, http.MethodGet, parsed.EscapedPath()+func() string {
if parsed.RawQuery == "" {
return ""
}
return "?" + parsed.RawQuery
}(), nil, accessToken, requestID)
}
func (c *coreClient) login(ctx context.Context, body any, rid string, clientIP ...string) (coreEnvelope, error) {
return c.call(ctx, http.MethodPost, "/api/v1/auth/login", body, "", rid, clientIP...)
}
func (c *coreClient) login2FA(ctx context.Context, body any, rid string, clientIP ...string) (coreEnvelope, error) {
return c.call(ctx, http.MethodPost, "/api/v1/auth/login/2fa", body, "", rid, clientIP...)
}
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
ClientIP string
Identity string
}
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
publicURL 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
loginAttempts map[string]loginAttempt
trustProxy bool
clock func() time.Time
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{}, loginAttempts: map[string]loginAttempt{}, clock: time.Now, processes: map[string]*exec.Cmd{}, pluginLocks: map[string]*sync.Mutex{}, marketplaceConfig: marketplace}
}
type loginAttempt struct {
Started time.Time
Count int
}
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 (a *app) now() time.Time {
if a != nil && a.clock != nil {
return a.clock()
}
return time.Now()
}
// trustedClientIP only accepts forwarded client identity from a loopback peer
// when explicitly enabled. Direct/non-proxy requests use RemoteAddr, so an
// attacker cannot spoof the login limiter key with an arbitrary header.
func (a *app) trustedClientIP(r *http.Request) string {
if r == nil {
return ""
}
host, _, err := net.SplitHostPort(strings.TrimSpace(r.RemoteAddr))
if err != nil {
host = strings.TrimSpace(r.RemoteAddr)
}
peer := net.ParseIP(host)
if a != nil && a.trustProxy && peer != nil && peer.IsLoopback() {
for _, value := range strings.Split(r.Header.Get("X-Forwarded-For"), ",") {
if ip := net.ParseIP(strings.TrimSpace(value)); ip != nil {
return ip.String()
}
}
}
if peer != nil {
return peer.String()
}
return "unknown"
}
func (a *app) allowLoginAttempt(r *http.Request, identity string) bool {
if a == nil {
return false
}
key := a.trustedClientIP(r) + "|" + strings.ToLower(strings.TrimSpace(identity))
now := a.now()
a.mu.Lock()
defer a.mu.Unlock()
if a.loginAttempts == nil {
a.loginAttempts = map[string]loginAttempt{}
}
// Bound the map while opportunistically removing expired buckets.
for candidate, attempt := range a.loginAttempts {
if now.Sub(attempt.Started) >= loginWindow {
delete(a.loginAttempts, candidate)
}
}
attempt := a.loginAttempts[key]
if attempt.Started.IsZero() || now.Sub(attempt.Started) >= loginWindow {
attempt = loginAttempt{Started: now}
}
if attempt.Count >= loginLimit {
a.loginAttempts[key] = attempt
return false
}
attempt.Count++
a.loginAttempts[key] = attempt
return true
}
func (a *app) addPendingLogin(value pendingLogin) (string, bool) {
a.mu.Lock()
defer a.mu.Unlock()
now := a.now()
for id, pending := range a.pending {
if !pending.Expires.After(now) {
delete(a.pending, id)
}
}
if len(a.pending) >= maxPendingLogins {
return "", false
}
if a.pending == nil {
a.pending = map[string]pendingLogin{}
}
id := token(16)
a.pending[id] = value
return id, true
}
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 := a.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
}
// Core's custom_menu_items contract is narrower than the plugin ID contract:
// dots are valid namespace separators for plugin IDs, but Core accepts only
// ASCII letters, digits, hyphens, and underscores in menu IDs. Keep the
// namespace-shaped plugin ID in the registry and derive a stable Core-safe ID
// only at the integration boundary.
func coreMenuID(raw string) (string, error) {
raw = strings.TrimSpace(raw)
if raw == "" {
return "", errors.New("plugin menu id is required")
}
if len(raw) <= 32 && coreMenuIDPattern.MatchString(raw) {
return raw, nil
}
var b strings.Builder
lastSeparator := false
for _, r := range raw {
if (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') || r == '_' || r == '-' {
b.WriteRune(r)
lastSeparator = false
continue
}
if !lastSeparator {
b.WriteByte('-')
lastSeparator = true
}
}
value := strings.Trim(b.String(), "-_")
if value == "" {
return "", errors.New("plugin menu id cannot be represented as a Core menu id")
}
if len(value) > 32 {
sum := sha256.Sum256([]byte(raw))
suffix := "-" + hex.EncodeToString(sum[:])[:8]
prefixLen := 32 - len(suffix)
value = strings.TrimRight(value[:prefixLen], "-_") + suffix
}
if !coreMenuIDPattern.MatchString(value) || len(value) > 32 {
return "", errors.New("plugin menu id cannot be represented as a Core menu id")
}
return value, nil
}
func coreMenuIDSet(raw string) map[string]struct{} {
ids := map[string]struct{}{strings.TrimSpace(raw): {}}
if normalized, err := coreMenuID(raw); err == nil {
ids[normalized] = struct{}{}
}
return ids
}
func menuItemHasID(item any, ids map[string]struct{}) bool {
object, ok := item.(map[string]any)
if !ok {
return false
}
id, _ := object["id"].(string)
_, ok = ids[id]
return ok
}
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
}
clientIP := a.trustedClientIP(r)
if !a.allowLoginAttempt(r, in.Email) {
w.Header().Set("Retry-After", strconv.Itoa(int(loginWindow/time.Second)))
writeJSON(w, http.StatusTooManyRequests, map[string]string{"error": "too many login attempts"})
return
}
out, err := a.core.login(r.Context(), in, requestID(r), clientIP)
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, accepted := a.addPendingLogin(pendingLogin{TempToken: temp, Expires: a.now().Add(pendingLoginTTL), ClientIP: clientIP, Identity: in.Email})
if !accepted {
w.Header().Set("Retry-After", strconv.Itoa(int(pendingLoginTTL/time.Second)))
writeJSON(w, http.StatusTooManyRequests, map[string]string{"error": "too many pending login challenges"})
return
}
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 || !a.now().Before(pending.Expires) {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "2fa session expired"})
return
}
if pending.ClientIP != "" && pending.ClientIP != a.trustedClientIP(r) {
writeJSON(w, http.StatusForbidden, map[string]string{"error": "2fa session binding failed"})
return
}
if !a.allowLoginAttempt(r, "2fa:"+pending.Identity) {
w.Header().Set("Retry-After", strconv.Itoa(int(loginWindow/time.Second)))
writeJSON(w, http.StatusTooManyRequests, map[string]string{"error": "too many login attempts"})
return
}
out, err := a.core.login2FA(r.Context(), map[string]string{"temp_token": pending.TempToken, "totp_code": in.TOTPCode}, requestID(r), a.trustedClientIP(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 := a.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)
candidateRefresh := current.RefreshToken
if next, _ := data["refresh_token"].(string); next != "" {
candidateRefresh = next
}
access, _ := data["access_token"].(string)
if access == "" {
a.core.logout(ctx, candidateRefresh, token(12))
return session{}, false
}
// Core may rotate the refresh token. Keep the candidate separate until the
// replacement access token has been revalidated, so every token issued by a
// successful refresh can be revoked on a validation failure.
current.AccessToken = access
me, err := a.core.me(ctx, access, token(12))
if err != nil {
a.core.logout(ctx, candidateRefresh, token(12))
return session{}, false
}
current.User = publicUser(envelopeData(me))
if !isAdmin(envelopeData(me)) {
a.core.logout(ctx, candidateRefresh, token(12))
return current, false
}
current.RefreshToken = candidateRefresh
current.LastSeen = a.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 {
a.registry.mu.Lock()
installed, alreadyInstalled := a.registry.data.Plugins[input.PluginID]
a.registry.mu.Unlock()
if alreadyInstalled && compareMarketplaceVersions(entry.Version, installed.Manifest.Version) <= 0 {
err = errors.New("marketplace version is not newer than the installed plugin")
}
}
if err == nil {
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, "retired": rev.Retired})
}
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)
ids := coreMenuIDSet(id)
next := make([]any, 0, len(raw))
for _, item := range raw {
if menuItemHasID(item, ids) {
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, err := a.menuURL(p)
if err != nil {
return err
}
coreID, err := coreMenuID(p.Manifest.UI.Menu.ID)
if err != nil {
return err
}
settings, err := a.core.settings(ctx, access, rid)
if err != nil {
return err
}
current := menuItemsFromSettings(envelopeData(settings))
candidate := map[string]any{"id": coreID, "label": p.Manifest.UI.Menu.Label, "url": itemURL, "visibility": "admin", "sort_order": p.Manifest.UI.Menu.SortOrder}
ids := coreMenuIDSet(p.Manifest.UI.Menu.ID)
replaced := false
for i, raw := range current {
if menuItemHasID(raw, ids) {
current[i] = candidate
replaced = true
}
}
if !replaced {
current = append(current, candidate)
}
_, err = a.core.updateMenu(ctx, access, rid, current)
return err
}
// menuURL resolves the URL that Core should place in custom_menu_items. The
// subscription module is deliberately different from a generic plugin: its
// browser entry must always land in the Plugin Admin Shell, where the shared
// administrator session and module route are enforced. Its own service URL is
// an internal health/BFF dependency and must never become a second login entry.
func (a *app) menuURL(p pluginRecord) (string, error) {
if p.Manifest.PluginID == subscriptionPluginID {
base := strings.TrimRight(strings.TrimSpace(a.publicURL), "/")
if base == "" {
return "", errors.New("PLUGIN_PUBLIC_URL is required for the subscription menu")
}
if err := validatePublicURL(base); err != nil {
return "", err
}
return base + "/admin/#/modules/subscription/overview", nil
}
itemURL := strings.TrimSpace(p.Manifest.UI.Menu.URL)
if itemURL == "" {
cfg, err := a.decryptConfig(p.ConfigCipher)
if err != nil {
return "", err
}
itemURL, _ = cfg["public_url"].(string)
itemURL = strings.TrimSpace(itemURL)
}
if itemURL == "" || validateMenuURL(itemURL) != nil {
return "", errors.New("menu URL is not configured")
}
return itemURL, nil
}
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.UI.Menu.ID, requestID(r)); err != nil {
return err
}
p.State = "draining"
a.stopPlugin(p.Manifest.PluginID)
p.State = "disabled"
return nil
})
}
// decodeRollbackRequest keeps the API compatible with the old empty-body
// rollback (used to discard a failed pending candidate), while making an
// explicitly selected revision authoritative and strictly validated.
func decodeRollbackRequest(r *http.Request) (string, error) {
if r == nil || r.Body == nil {
return "", nil
}
raw, err := io.ReadAll(io.LimitReader(r.Body, 64<<10))
if err != nil {
return "", errors.New("invalid rollback request")
}
if len(bytes.TrimSpace(raw)) == 0 {
return "", nil
}
decoder := json.NewDecoder(bytes.NewReader(raw))
decoder.DisallowUnknownFields()
var input struct {
Revision string `json:"revision"`
}
if err := decoder.Decode(&input); err != nil {
return "", errors.New("rollback revision is invalid")
}
var trailing any
if err := decoder.Decode(&trailing); err != io.EOF {
return "", errors.New("rollback request must contain one JSON value")
}
input.Revision = strings.TrimSpace(input.Revision)
if input.Revision == "" || len(input.Revision) > 160 || strings.ContainsAny(input.Revision, "/\\") {
return "", errors.New("rollback revision is required")
}
return input.Revision, 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 {
requestedRevision, err := decodeRollbackRequest(r)
if err != nil {
return err
}
targetState := p.State
if p.PreviousState != "" {
targetState = p.PreviousState
}
current := p.ActiveRevision
pendingID := p.PendingRevision
var target revision
hasTarget := false
if requestedRevision != "" {
for _, candidate := range p.Revisions {
if candidate.ID != requestedRevision {
continue
}
if candidate.Retired {
return errors.New("rollback revision is retired")
}
if candidate.ID == current {
return errors.New("rollback revision is already active")
}
if candidate.ID == pendingID {
return errors.New("rollback revision is still pending")
}
target = candidate
if target.Manifest.PluginID == "" {
target.Manifest = p.Manifest
}
hasTarget = true
break
}
if !hasTarget {
return errors.New("rollback revision not found")
}
}
// 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. When
// the caller names an older revision, discard the candidate first and then
// continue with the explicitly selected target.
if pendingID != "" && current != "" {
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 = ""
if requestedRevision == "" {
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 !hasTarget && len(p.Revisions) < 2 {
p.State = "disabled"
return errors.New("no rollback revision retained")
}
if !hasTarget {
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
}
hasTarget = true
break
}
}
if !hasTarget {
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.UI.Menu.ID, 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
}
// validatePublicURL accepts the externally reachable base URL of Plugin Admin.
// It intentionally has no query, fragment, credentials, or path traversal;
// menuURL appends the one fixed module fragment after validation.
func validatePublicURL(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("PLUGIN_PUBLIC_URL must be an absolute http(s) URL without credentials, query or fragment")
}
if u.Scheme == "http" && !isLoopbackHost(u.Hostname()) {
return errors.New("PLUGIN_PUBLIC_URL must use HTTPS unless loopback")
}
if strings.Contains(u.Path, "\\") || strings.Contains(u.Path, "..") {
return errors.New("PLUGIN_PUBLIC_URL contains an unsafe path")
}
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
}
// Retired payloads are intentionally kept through the current operation
// commit so a failed registry write remains recoverable. Once recovery has
// persisted the authoritative record, clean those unreferenced payloads.
_ = a.cleanupRetiredRevisions(p)
}
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
}
// cleanupRetiredRevisions removes files for revisions that are no longer
// eligible for activation. Metadata is retained in the registry for audit and
// history, but retired payloads must not accumulate indefinitely on disk.
func (a *app) cleanupRetiredRevisions(p pluginRecord) error {
base, err := filepath.Abs(filepath.Join(a.root, "installed"))
if err != nil {
return err
}
var firstErr error
for _, rev := range p.Revisions {
if !rev.Retired || rev.Path == "" || rev.ID == p.ActiveRevision || rev.ID == p.PendingRevision {
continue
}
absolute, pathErr := filepath.Abs(rev.Path)
if pathErr != nil {
if firstErr == nil {
firstErr = pathErr
}
continue
}
rel, relErr := filepath.Rel(base, absolute)
if relErr != nil || rel == "." || rel == ".." || strings.HasPrefix(rel, ".."+string(os.PathSeparator)) {
if firstErr == nil {
firstErr = errors.New("retired revision path is outside the plugin install root")
}
continue
}
if removeErr := os.RemoveAll(absolute); 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})
}
// subscriptionStatus exposes module health without creating a second session.
// The browser receives only derived state; Core credentials remain in session.
func (a *app) subscriptionStatus(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 {
return
}
enabled := a.subscriptionEnabled()
writeJSON(w, http.StatusOK, map[string]any{
"module": "subscription",
"enabled": enabled,
"core_base_configured": a.core != nil,
"session_mode": "shared",
"credential_state": "server_managed",
"operator": s.User,
})
}
// captchaConfig exposes only the public browser-side fields required to render
// the Core login challenge. Provider secrets and the rest of Core settings
// never cross this boundary. The endpoint is intentionally unauthenticated so
// the first Plugin Admin login can discover whether a proof is required.
func (a *app) captchaConfig(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.core == nil {
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "core unavailable"})
return
}
out, err := a.core.publicSettings(r.Context(), requestID(r))
if err != nil {
a.coreError(w, err, "core settings unavailable")
return
}
data := envelopeData(out)
boolValue := func(key string) bool {
value, _ := data[key].(bool)
return value
}
stringValue := func(key string) string {
value, _ := data[key].(string)
value = strings.TrimSpace(value)
if len(value) > 512 {
return ""
}
return value
}
turnstileEnabled := boolValue("turnstile_enabled")
tencentEnabled := boolValue("tencent_captcha_enabled")
aliyunEnabled := boolValue("aliyun_captcha_enabled")
provider := ""
switch {
case turnstileEnabled:
provider = "turnstile"
case tencentEnabled:
provider = "tencent"
case aliyunEnabled:
provider = "aliyun"
}
writeJSON(w, http.StatusOK, map[string]any{
"enabled": turnstileEnabled || tencentEnabled || aliyunEnabled,
"provider": provider,
"turnstile_enabled": turnstileEnabled,
"turnstile_site_key": stringValue("turnstile_site_key"),
"tencent_captcha_enabled": tencentEnabled,
"tencent_captcha_app_id": stringValue("tencent_captcha_app_id"),
"tencent_captcha_region": stringValue("tencent_captcha_region"),
"aliyun_captcha_enabled": aliyunEnabled,
"aliyun_captcha_scene_id": stringValue("aliyun_captcha_scene_id"),
"aliyun_captcha_prefix": stringValue("aliyun_captcha_prefix"),
"aliyun_captcha_region": stringValue("aliyun_captcha_region"),
})
}
// subscriptionEnabled is the server-side source of truth for the optional
// subscription module. UI visibility is only a convenience; every module BFF
// endpoint calls requireSubscriptionEnabled so a stale/deep-linked browser
// cannot read Core subscription data after the module is disabled.
func (a *app) subscriptionEnabled() bool {
if a == nil || a.registry == nil {
return false
}
a.registry.mu.Lock()
defer a.registry.mu.Unlock()
plugin, found := a.registry.data.Plugins[subscriptionPluginID]
return found && (plugin.State == "healthy" || plugin.State == "enabled")
}
func (a *app) requireSubscriptionEnabled(w http.ResponseWriter) bool {
if a.subscriptionEnabled() {
return true
}
// Do not reveal whether the module was ever installed. The control plane
// still exposes /api/subscription/status to authenticated admins so the UI
// can render an accurate disabled state.
writeJSON(w, http.StatusNotFound, map[string]string{"error": "subscription module is not enabled"})
return false
}
func (a *app) subscriptionAudit(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
}
if !a.requireSubscriptionEnabled(w) {
return
}
if a.registry == nil {
writeJSON(w, http.StatusOK, map[string]any{"items": []auditEvent{}})
return
}
a.registry.mu.Lock()
items := make([]auditEvent, 0)
for _, item := range a.registry.data.Audit {
if strings.HasPrefix(item.Action, "subscription") || strings.Contains(item.Action, "/api/v1/admin/subscriptions") || strings.Contains(item.Action, "/api/v1/admin/payment/plans") || strings.Contains(item.Action, "/api/v1/admin/users/") {
items = append(items, item)
}
}
a.registry.mu.Unlock()
writeJSON(w, http.StatusOK, map[string]any{"items": items})
}
func sanitizedCoreValue(value any) any {
switch typed := value.(type) {
case []any:
out := make([]any, 0, len(typed))
for _, item := range typed {
out = append(out, sanitizedCoreValue(item))
}
return out
case map[string]any:
out := make(map[string]any, len(typed))
for key, item := range typed {
if isSecretKey(key) {
continue
}
out[key] = sanitizedCoreValue(item)
}
return out
default:
return value
}
}
func (a *app) writeSubscriptionEnvelope(w http.ResponseWriter, envelope coreEnvelope) {
var data any
if len(envelope.Data) > 0 {
if json.Unmarshal(envelope.Data, &data) != nil {
writeJSON(w, http.StatusBadGateway, map[string]string{"error": "invalid Core response"})
return
}
}
status := envelope.status
if status == 0 {
status = http.StatusOK
}
writeJSON(w, status, map[string]any{"code": envelope.Code, "message": envelope.Message, "data": sanitizedCoreValue(data)})
}
func (a *app) subscriptionProxy(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
return
}
if !validSubscriptionQuery(r.URL.Query()) {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid query parameters"})
return
}
if a.core == nil {
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "core unavailable"})
return
}
_, s, ok := a.authenticate(w, r)
if !ok {
return
}
if !a.requireSubscriptionEnabled(w) {
return
}
var corePath string
pathName := strings.TrimSuffix(r.URL.Path, "/")
switch {
case pathName == "/api/subscription/plans":
corePath = "/api/v1/admin/payment/plans"
case pathName == "/api/subscription/subscriptions":
corePath = "/api/v1/admin/subscriptions"
case strings.HasPrefix(pathName, "/api/subscription/subscriptions/"):
id := strings.TrimPrefix(pathName, "/api/subscription/subscriptions/")
if !positiveID(id) {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid subscription id"})
return
}
corePath = "/api/v1/admin/subscriptions/" + id
case strings.HasPrefix(pathName, "/api/subscription/users/"):
rest := strings.TrimPrefix(pathName, "/api/subscription/users/")
parts := strings.Split(rest, "/")
if len(parts) == 1 && positiveID(parts[0]) {
corePath = "/api/v1/admin/users/" + parts[0]
} else if len(parts) == 2 && parts[1] == "subscriptions" && positiveID(parts[0]) {
corePath = "/api/v1/admin/users/" + parts[0] + "/subscriptions"
} else {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid user path"})
return
}
default:
writeJSON(w, http.StatusNotFound, map[string]string{"error": "subscription endpoint not found"})
return
}
// authenticate() has already revalidated the Core access token. A failed
// read is intentionally returned as a generic gateway error to avoid leaking
// Core response details.
out, err := a.core.read(r.Context(), corePath+querySuffix(r.URL.Query()), s.AccessToken, requestID(r))
if err != nil {
a.coreError(w, err, "subscription data unavailable")
return
}
if a.registry != nil {
_ = a.registry.addAudit(auditEvent{Time: time.Now().UTC(), Action: "read:" + corePath, PluginID: "qiu.subscription-admin", ActorID: s.User["id"], RequestID: requestID(r)})
}
a.writeSubscriptionEnvelope(w, out)
}
func querySuffix(values url.Values) string {
if len(values) == 0 {
return ""
}
return "?" + sanitizeSubscriptionQuery(values).Encode()
}
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)
}
}
coreID := ""
if err == nil {
coreID, err = coreMenuID(p.Manifest.UI.Menu.ID)
}
itemURL := ""
if err == nil {
itemURL, err = a.menuURL(p)
}
next := append([]any(nil), current...)
if err == nil {
candidate := map[string]any{"id": coreID, "label": p.Manifest.UI.Menu.Label, "url": itemURL, "visibility": "admin", "sort_order": p.Manifest.UI.Menu.SortOrder}
ids := coreMenuIDSet(p.Manifest.UI.Menu.ID)
replaced := false
for i, raw := range next {
if menuItemHasID(raw, ids) {
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 staticMounts(base string) []string {
base = strings.TrimRight(strings.TrimSpace(base), "/")
mounts := []string{"/admin"}
if base != "" && base != "/" {
// Vue's relative ../ asset URLs resolve one level above /admin/.
// Keep that parent prefix as an asset alias for mounted deployments.
mounts = append(mounts, base+"/admin", base)
}
return mounts
}
// staticPath returns the path beneath a recognised UI mount. It rejects
// traversal before cleaning so encoded dot segments cannot escape ui/.
func staticPath(a *app, requestPath string) (string, bool, bool) {
escaped := requestPath
if escaped == "" {
escaped = "/"
}
decoded, err := url.PathUnescape(escaped)
if err != nil || strings.Contains(decoded, "\\") {
return "", false, true
}
for _, segment := range strings.Split(decoded, "/") {
if segment == "." || segment == ".." {
return "", false, true
}
}
clean := path.Clean("/" + decoded)
if clean == "/" {
return "", true, false
}
// The default shell references app.js/styles.css from the parent of
// /admin/. Vite builds commonly place additional files under /assets/.
if clean == "/app.js" || clean == "/styles.css" || strings.HasPrefix(clean, "/assets/") {
return strings.TrimPrefix(clean, "/"), true, false
}
for _, mount := range staticMounts(a.publicBasePath) {
if clean == mount || strings.HasPrefix(clean, mount+"/") {
rel := strings.TrimPrefix(clean, mount)
return strings.TrimPrefix(rel, "/"), true, false
}
}
return "", false, false
}
func (a *app) writeIndex(w http.ResponseWriter) {
data, err := uiFS.ReadFile("ui/index.html")
if err != nil {
http.Error(w, "ui unavailable", http.StatusInternalServerError)
return
}
baseJSON, _ := json.Marshal(a.publicBasePath)
content := strings.ReplaceAll(string(data), "__PLUGIN_BASE_PATH_JSON__", string(baseJSON))
// The Vue shell stores the mount in a meta attribute. Escape it as HTML so
// an operator-provided base path cannot break the document or inject markup.
content = strings.ReplaceAll(content, "__PLUGIN_BASE_PATH__", html.EscapeString(a.publicBasePath))
data = []byte(content)
w.Header().Set("Cache-Control", "no-store")
w.Header().Set("Content-Type", "text/html; charset=utf-8")
_, _ = w.Write(data)
}
func (a *app) static(w http.ResponseWriter, r *http.Request) {
rel, mounted, rejected := staticPath(a, r.URL.EscapedPath())
if rejected {
http.NotFound(w, r)
return
}
if !mounted {
http.NotFound(w, r)
return
}
if rel == "" {
if r.URL.Path == "/admin" || (a.publicBasePath != "" && r.URL.Path == strings.TrimRight(a.publicBasePath, "/")+"/admin") {
http.Redirect(w, r, strings.TrimRight(a.publicBasePath, "/")+"/admin/", http.StatusPermanentRedirect)
return
}
a.writeIndex(w)
return
}
data, err := uiFS.ReadFile("ui/" + rel)
if err == nil {
contentType := mime.TypeByExtension(path.Ext(rel))
if contentType == "" {
contentType = http.DetectContentType(data)
}
w.Header().Set("Content-Type", contentType)
_, _ = w.Write(data)
return
}
// Hash-routed SPA paths are client-side routes and have no file extension.
if path.Ext(rel) == "" {
a.writeIndex(w)
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/captcha-config", a.captchaConfig)
mux.HandleFunc("/api/me", a.me)
mux.HandleFunc("/api/audit", a.audit)
mux.HandleFunc("/api/subscription/status", a.subscriptionStatus)
mux.HandleFunc("/api/subscription/audit", a.subscriptionAudit)
mux.HandleFunc("/api/subscription/plans", a.subscriptionProxy)
mux.HandleFunc("/api/subscription/subscriptions", a.subscriptionProxy)
mux.HandleFunc("/api/subscription/subscriptions/", a.subscriptionProxy)
mux.HandleFunc("/api/subscription/users/", a.subscriptionProxy)
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)
handler := a.securityHeaders(requestIDMiddleware(mux))
base := strings.TrimRight(strings.TrimSpace(a.publicBasePath), "/")
if base == "" || base == "/" {
return handler
}
// Reverse proxies normally strip PLUGIN_PUBLIC_BASE_PATH. Stripping the
// exact configured prefix here as well keeps direct local previews and
// mounted production deployments behaviorally identical without exposing a
// generic path proxy.
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == base || strings.HasPrefix(r.URL.Path, base+"/") {
clone := r.Clone(r.Context())
clone.URL.Path = strings.TrimPrefix(r.URL.Path, base)
if clone.URL.Path == "" {
clone.URL.Path = "/"
}
clone.URL.RawPath = ""
handler.ServeHTTP(w, clone)
return
}
handler.ServeHTTP(w, r)
})
}
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")
// CAPTCHA providers are loaded only when Core reports them enabled. Keep
// the allowlist provider-specific; Core JWTs and application data remain
// same-origin and are never sent to these origins by the browser.
w.Header().Set("Content-Security-Policy", "default-src 'self'; worker-src 'self' blob:; script-src 'self' https://challenges.cloudflare.com https://o.alicdn.com https://*.alicdn.com https://turing.captcha.qcloud.com https://turing.captcha.gtimg.com https://ca.turing.captcha.qcloud.com https://global.turing.captcha.gtimg.com https://cloudcache.tencentcs.com; style-src 'self' 'unsafe-inline' https://*.captcha.gtimg.com https://o.alicdn.com https://*.alicdn.com; img-src 'self' data: blob: https:; font-src 'self' data: https://fonts.gstatic.com; connect-src 'self' https://challenges.cloudflare.com https://turing.captcha.qcloud.com https://turing.captcha.gtimg.com https://ca.turing.captcha.qcloud.com https://global.turing.captcha.gtimg.com https://cloudcache.tencentcs.com https://rce.tencentrio.com https://o.alicdn.com https://*.alicdn.com https://*.aliyuncs.com; frame-src 'self' https://challenges.cloudflare.com https://turing.captcha.qcloud.com https://turing.captcha.gtimg.com https://ca.turing.captcha.qcloud.com https://global.turing.captcha.gtimg.com https://cloudcache.tencentcs.com https://o.alicdn.com https://*.alicdn.com https://*.aliyuncs.com; 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.publicURL = strings.TrimRight(strings.TrimSpace(os.Getenv("PLUGIN_PUBLIC_URL")), "/")
if a.publicURL != "" {
if publicErr := validatePublicURL(a.publicURL); publicErr != nil {
slog.Error("invalid Plugin Admin public URL", "error", publicErr)
os.Exit(2)
}
}
a.trustProxy = strings.EqualFold(strings.TrimSpace(os.Getenv("PLUGIN_TRUST_PROXY")), "true")
if a.trustProxy && isLoopbackHost(host) == false {
slog.Error("PLUGIN_TRUST_PROXY requires a loopback plugin listener")
os.Exit(2)
}
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)
}
marketplace.indexSHA256 = strings.ToLower(strings.TrimSpace(os.Getenv("PLUGIN_MARKETPLACE_INDEX_SHA256")))
if marketplace.indexSHA256 != "" && !marketplaceSHA256Pattern.MatchString(marketplace.indexSHA256) {
slog.Error("invalid marketplace index sha256 pin")
os.Exit(2)
}
if marketplace.remoteURL != nil && !allowLoopbackMarketplace && marketplace.indexSHA256 == "" {
slog.Error("PLUGIN_MARKETPLACE_INDEX_SHA256 is required for remote production marketplace indexes")
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