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