This repository has been archived on 2026-04-03. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
core/server/main.go
Jeffrey Smith 349ff5c80a
All checks were successful
CI/CD / detect-changes (pull_request) Successful in 4s
CI/CD / test-frontend (pull_request) Successful in 5s
CI/CD / test-go-pg (pull_request) Successful in 2m46s
CI/CD / test-sqlite (pull_request) Successful in 2m53s
CI/CD / build-and-deploy (pull_request) Successful in 1m42s
Fix health tab: node identity indicator + cluster card overflow
Add node_id to metrics snapshot so the health tab shows which instance
the user is viewing from. Hoist nodeID computation so it works on both
SQLite and Postgres deployments. Fix cluster node card overflow with
text-overflow ellipsis and proper flex constraints.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-31 13:40:18 +00:00

994 lines
38 KiB
Go

package main
import (
"bytes"
"context"
_ "embed"
"fmt"
"log"
"net/http"
"os"
"path/filepath"
"strings"
"time"
"github.com/gin-gonic/gin"
"github.com/prometheus/client_golang/prometheus/promhttp"
"switchboard-core/auth"
"switchboard-core/cluster"
"switchboard-core/config"
"switchboard-core/crypto"
"switchboard-core/database"
"switchboard-core/events"
"switchboard-core/sandbox"
"switchboard-core/handlers"
"switchboard-core/logging"
"switchboard-core/metrics"
"switchboard-core/middleware"
"switchboard-core/notifications"
"switchboard-core/pages"
"switchboard-core/storage"
"switchboard-core/store"
"switchboard-core/triggers"
"switchboard-core/workflow"
postgres "switchboard-core/store/postgres"
sqliteStore "switchboard-core/store/sqlite"
)
//
//go:embed static/openapi.yaml
var openapiSpec []byte
//go:embed static/swagger.html
var swaggerHTML []byte
func main() {
// ── Subcommand dispatch ──────────────────
// Usage: switchboard vault rekey
if len(os.Args) > 2 && os.Args[1] == "vault" {
runVaultCommand(os.Args[2])
return
}
if len(os.Args) > 1 && os.Args[1] == "version" {
fmt.Println("switchboard", Version)
return
}
// ── Server startup ──────────────────────
startTime := time.Now()
cfg := config.Load()
// Structured logging — must be first so all subsequent
// log output goes through slog.
logging.Init(cfg.LogFormat, cfg.LogLevel)
var stores store.Stores
uekCache := crypto.NewUEKCache()
var keyResolver *crypto.KeyResolver
var objStore storage.ObjectStore
if err := database.Connect(cfg); err != nil {
log.Printf("⚠ Database unavailable: %v", err)
log.Println(" Running in unmanaged mode (no persistence)")
} else {
// Schema check: init if fresh, upgrade if behind, proceed if current
if err := database.Migrate(); err != nil {
log.Fatalf("❌ Schema migration failed: %v", err)
}
// Vault: enforce encryption key + backfill plaintext keys
if err := crypto.EnforceEncryptionKey(database.DB, cfg.EncryptionKey); err != nil {
log.Fatalf("❌ Vault check failed: %v", err)
}
if cfg.EncryptionKey != "" {
if err := crypto.BackfillEncryptedKeys(database.DB, cfg.EncryptionKey); err != nil {
log.Fatalf("❌ Vault backfill failed: %v", err)
}
}
// Derive env key for key resolver (nil if not set)
var envKey []byte
if cfg.EncryptionKey != "" {
var err error
envKey, err = crypto.DeriveKeyFromEnv(cfg.EncryptionKey)
if err != nil {
log.Fatalf("❌ Failed to derive encryption key: %v", err)
}
log.Println(" 🔐 API key encryption active")
}
keyResolver = crypto.NewKeyResolver(envKey, uekCache)
// Initialize store layer
if database.IsSQLite() {
stores = sqliteStore.NewStores(database.DB)
} else {
stores = postgres.NewStores(database.DB)
}
metrics.StartDBCollector(database.DB, 15*time.Second)
// Bootstrap admin from env (K8s secret) — upserts on every restart
handlers.BootstrapAdmin(cfg, stores, uekCache)
// Seed additional users from env (dev/test only, skipped in production)
handlers.SeedUsers(cfg, stores, uekCache)
}
defer database.Close()
// Auto-detects PVC if STORAGE_PATH is writable. Explicit STORAGE_BACKEND
// overrides auto-detection. nil objStore = storage features disabled.
var s3Cfg *storage.S3Config
if cfg.S3Bucket != "" {
s3Cfg = &storage.S3Config{
Endpoint: cfg.S3Endpoint,
Bucket: cfg.S3Bucket,
Region: cfg.S3Region,
AccessKey: cfg.S3AccessKey,
SecretKey: cfg.S3SecretKey,
Prefix: cfg.S3Prefix,
ForcePathStyle: cfg.S3ForcePathStyle,
}
}
if s, err := storage.Init(cfg.StorageBackend, cfg.StoragePath, s3Cfg); err != nil {
if cfg.StorageBackend != "" {
// Explicit backend requested but failed — fatal
log.Fatalf("❌ Storage init failed: %v", err)
}
log.Printf("⚠ Storage init failed: %v", err)
} else {
objStore = s
}
handlers.SetStorageConfigured(objStore != nil)
// ── EventBus (created early — needed by role resolver) ──
bus := events.NewBus()
// Wire Postgres LISTEN/NOTIFY for cross-pod event fan-out (multi-replica).
// No-op when running SQLite — in-process Bus is sufficient for single-pod.
events.StartPGBroadcast(bus)
// Sandboxed interpreter for extension scripts. Runner assembles
// modules based on granted permissions. Notifier attached below
// after notification service init.
starlarkRunner := sandbox.NewRunner(
sandbox.New(sandbox.DefaultConfig()),
stores,
)
starlarkRunner.SetConnectionResolver(handlers.NewConnectionResolverAdapter(stores, keyResolver))
starlarkRunner.SetDB(database.DB, database.IsPostgres())
if cfg.StoragePath != "" {
starlarkRunner.SetPackagesDir(cfg.StoragePath + "/packages")
}
starlarkRunner.SetBus(bus)
if os.Getenv("EXT_ALLOW_PRIVATE_IPS") == "true" {
starlarkRunner.SetAllowPrivateIPs(true)
log.Printf(" ⚠️ Extension SSRF protection relaxed: private IPs allowed")
}
// ── Bundled Packages ───────────────
// Auto-install bundled .pkg archives on first run.
// Skipped if SKIP_BUNDLED_PACKAGES=true or packages already exist.
if !cfg.SkipBundledPackages && stores.Packages != nil {
bundledPkgDir := ""
if cfg.StoragePath != "" {
bundledPkgDir = cfg.StoragePath + "/packages"
}
handlers.InstallBundledPackages(cfg.BundledPackagesDir, bundledPkgDir, cfg.BundledPackages, stores, starlarkRunner)
}
// ── Trigger Engine ─────────────────
triggerEngine := triggers.New(stores, starlarkRunner, bus)
if err := triggerEngine.Start(context.Background()); err != nil {
log.Printf(" ⚠️ Trigger engine start: %v", err)
}
triggers.SetGlobalEngine(triggerEngine)
defer triggerEngine.Stop()
r := gin.New()
r.Use(middleware.RequestID())
r.Use(middleware.Prometheus())
r.Use(middleware.Logger())
r.Use(middleware.Recovery())
userCache := middleware.NewUserStatusCache()
r.Use(middleware.CORS(cfg))
// ── Base path group ──────────────────────
base := r.Group(cfg.BasePath)
// ── WebSocket Hub ─────────────────────────
hub := events.NewHub(bus, middleware.GetAllowedOrigins(cfg))
// ── Node Identity ────────────────
// Used by cluster registry (PG) and metrics endpoint (all deployments).
nodeID := cfg.ClusterNodeID
if nodeID == "" {
hostname, _ := os.Hostname()
nodeID = fmt.Sprintf("%s-%d", hostname, os.Getpid())
}
// ── Cluster Registry ────────────
// PG-backed node self-registration + heartbeat. No-op on SQLite.
var clusterReg *cluster.Registry
if database.IsPostgres() && stores.Cluster != nil {
clusterReg = cluster.NewRegistry(nodeID, cfg.ClusterEndpoint, stores.Cluster, hub, cluster.RegistryConfig{
HeartbeatInterval: cfg.ClusterHeartbeatInterval,
StaleThreshold: cfg.ClusterStaleThreshold,
})
clusterReg.SetSandboxStats(sandbox.SandboxStats)
clusterReg.SetTriggerFireCount(triggerEngine.FireCount)
clusterReg.SetExtensionCount(func() int {
if stores.Packages == nil {
return 0
}
pkgs, err := stores.Packages.List(context.Background())
if err != nil {
return 0
}
count := 0
for _, p := range pkgs {
if p.Status == "active" {
count++
}
}
return count
})
if err := clusterReg.Start(); err != nil {
log.Printf("⚠ Cluster registry failed to start: %v", err)
clusterReg = nil
} else {
defer clusterReg.Stop()
}
}
// ── WebSocket Ticket Store ─────
ticketAdapter := &events.TicketValidatorAdapter{Store: stores.Tickets}
// ── Notification Service ───────
var notifSvc *notifications.Service
if stores.Notifications != nil {
notifSvc = notifications.NewService(stores.Notifications, hub).
WithPrefs(stores.NotifPrefs).
WithUsers(stores.Users)
// Load SMTP config from platform settings for email transport (Phase 3)
if stores.GlobalConfig != nil {
if smtpCfg, err := notifications.LoadSMTPConfig(stores.GlobalConfig, keyResolver); err == nil && smtpCfg != nil {
transport := notifications.NewEmailTransport(*smtpCfg)
instanceName := "Switchboard Core"
if brandCfg, err := stores.GlobalConfig.Get(context.Background(), "branding"); err == nil {
if name, ok := brandCfg["instance_name"].(string); ok && name != "" {
instanceName = name
}
}
notifSvc.WithEmail(transport, instanceName)
log.Println("[notifications] email transport enabled")
}
}
notifSvc.StartCleanup()
notifications.SetDefault(notifSvc)
starlarkRunner.SetNotifier(notifSvc)
// Subscribe to role.fallback events → generate notifications for admins
bus.Subscribe("role.fallback", notifications.RoleFallbackHandler(notifSvc, stores))
defer notifSvc.StopCleanup()
}
// Health check (k8s probes hit this directly)
buildHealthResponse := func() gin.H {
info := gin.H{
"status": "ok",
"version": Version,
"database": database.IsConnected(),
"database_name": database.Name(),
"schema_version": database.SchemaVersion(),
}
if database.IsConnected() {
info["registration_enabled"] = handlers.IsRegistrationEnabled(stores)
}
appendClusterHealth(info, clusterReg, stores)
return info
}
base.GET("/health", func(c *gin.Context) {
c.JSON(200, buildHealthResponse())
})
// Liveness: process is alive and serving (no dependency checks).
base.GET("/healthz/live", func(c *gin.Context) {
c.JSON(200, gin.H{"status": "ok"})
})
// Readiness: process can serve traffic (PG reachable).
// Failing readiness pulls the pod from the Service — new requests
// route to healthy replicas.
base.GET("/healthz/ready", func(c *gin.Context) {
if database.DB == nil {
c.JSON(503, gin.H{"error": "database not initialized"})
return
}
ctx, cancel := context.WithTimeout(c.Request.Context(), 2*time.Second)
defer cancel()
if err := database.DB.PingContext(ctx); err != nil {
c.JSON(503, gin.H{"error": "database unavailable"})
return
}
c.JSON(200, gin.H{"status": "ok"})
})
base.GET("/metrics", gin.WrapH(promhttp.Handler()))
base.GET("/api/docs", func(c *gin.Context) {
c.Data(http.StatusOK, "text/html; charset=utf-8", swaggerHTML)
})
base.GET("/api/docs/openapi.yaml", func(c *gin.Context) {
// Replace ${VERSION} placeholder with actual version from VERSION file
patched := bytes.Replace(openapiSpec, []byte("${VERSION}"), []byte(Version), 1)
c.Data(http.StatusOK, "application/yaml", patched)
})
base.GET("/api/docs/openapi.json", func(c *gin.Context) {
spec, err := handlers.BuildOpenAPISpec(stores, openapiSpec, Version)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "failed to build spec"})
return
}
c.JSON(http.StatusOK, spec)
})
// WebSocket endpoint
base.GET("/ws", middleware.WsAuth(cfg, stores.Users, userCache, ticketAdapter), hub.HandleWebSocket)
// ── Auth routes (rate limited) ──────────────
authMode, err := auth.ParseMode(cfg.AuthMode)
if err != nil {
log.Fatalf("❌ Invalid AUTH_MODE=%q: %v", cfg.AuthMode, err)
}
var authProvider auth.Provider
switch authMode {
case auth.ModeBuiltin:
authProvider = auth.NewBuiltinProvider()
case auth.ModeMTLS:
authProvider = auth.NewMTLSProvider(auth.MTLSConfig{
HeaderDN: cfg.MTLSHeaderDN,
HeaderVerify: cfg.MTLSHeaderVerify,
HeaderFingerprint: cfg.MTLSHeaderFingerprint,
AutoActivate: cfg.MTLSAutoActivate,
DefaultTeam: cfg.MTLSDefaultTeam,
})
case auth.ModeOIDC:
var err error
authProvider, err = auth.NewOIDCProvider(auth.OIDCConfig{
IssuerURL: cfg.OIDCIssuerURL,
ExternalIssuerURL: cfg.OIDCExternalIssuerURL,
ClientID: cfg.OIDCClientID,
ClientSecret: cfg.OIDCClientSecret,
RedirectURL: cfg.OIDCRedirectURL,
AutoActivate: cfg.OIDCAutoActivate,
DefaultTeam: cfg.OIDCDefaultTeam,
RolesClaim: cfg.OIDCRolesClaim,
GroupsClaim: cfg.OIDCGroupsClaim,
AdminRole: cfg.OIDCAdminRole,
})
if err != nil {
log.Fatalf("❌ OIDC provider init failed: %v", err)
}
}
log.Printf(" 🔑 Auth mode: %s", authMode)
authH := handlers.NewAuthHandler(cfg, stores, uekCache, authProvider)
authLimiter := middleware.NewRateLimiter(stores.RateLimits, 5, 8)
api := base.Group("/api/v1")
{
// Health (routable through ingress — same shape as /health)
api.GET("/health", func(c *gin.Context) {
c.JSON(200, buildHealthResponse())
})
authGroup := api.Group("/auth")
authGroup.Use(authLimiter.Limit())
{
authGroup.POST("/register", authH.Register)
authGroup.POST("/login", authH.Login)
authGroup.POST("/refresh", authH.Refresh)
authGroup.POST("/logout", authH.Logout)
// OIDC routes (only active when AUTH_MODE=oidc)
authGroup.GET("/oidc/login", authH.OIDCLogin)
authGroup.GET("/oidc/callback", authH.OIDCCallback)
}
// ── Public extension assets ────────────────
// Script tags can't send Authorization headers, so asset serving must be public.
// Extension JS is the same for all users — no user-specific data.
extH := handlers.NewExtensionHandler(stores)
api.GET("/extensions/:id/assets/*path", extH.ServeExtensionAsset)
// ── Webhook Inbound ──────────────
// Public routes — HMAC-verified in handler, no auth middleware.
api.POST("/hooks/:package_id/:slug", triggerEngine.HandleWebhook)
api.GET("/hooks/:package_id/:slug", triggerEngine.HandleWebhook)
// ── Workflow Engine (shared by public + protected routes) ──
wfEngine := workflow.NewEngine(stores, bus, starlarkRunner)
// ── Public Workflow Entry ─────
publicWfH := handlers.NewWorkflowPublicHandler(wfEngine, stores)
publicWf := api.Group("/public/workflows")
publicWf.Use(authLimiter.Limit())
{
publicWf.POST("/:id/start", publicWfH.StartPublic)
publicWf.GET("/resume/:token", publicWfH.ResumePublic)
publicWf.POST("/advance/:token", publicWfH.AdvancePublic)
}
// ── Workflow Scanner ──────────
wfScanner := workflow.NewScanner(stores, bus)
wfScanner.Start()
defer wfScanner.Stop()
// ── Protected routes ────────────────────
protected := api.Group("")
protected.Use(middleware.Auth(cfg, stores.Users, userCache))
protected.Use(middleware.ValidatePathParams())
{
// ── WebSocket Ticket ───────────
// Issue a single-use ticket for WebSocket auth.
// Client fetches this, then connects with ?ticket=<opaque>.
protected.POST("/ws/ticket", func(c *gin.Context) {
userID := c.GetString("user_id")
ticket, err := stores.Tickets.Issue(c.Request.Context(), userID)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "failed to issue ticket"})
return
}
c.JSON(http.StatusOK, gin.H{"ticket": ticket})
})
// Presence
presence := handlers.NewPresenceHandler(stores)
protected.POST("/presence/heartbeat", presence.Heartbeat)
protected.GET("/presence", presence.Query)
// User search
protected.GET("/users/search", presence.SearchUsers)
// Workflows
wfH := handlers.NewWorkflowHandler(stores)
protected.GET("/workflows", wfH.List)
protected.POST("/workflows", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.Create)
protected.GET("/workflows/:id", wfH.Get)
protected.PATCH("/workflows/:id", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.Update)
protected.DELETE("/workflows/:id", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.Delete)
protected.GET("/workflows/:id/stages", wfH.ListStages)
protected.POST("/workflows/:id/stages", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.CreateStage)
protected.PUT("/workflows/:id/stages/:sid", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.UpdateStage)
protected.DELETE("/workflows/:id/stages/:sid", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.DeleteStage)
protected.PATCH("/workflows/:id/stages/reorder", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.ReorderStages)
protected.POST("/workflows/:id/publish", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.Publish)
protected.POST("/workflows/:id/clone", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfH.Clone)
protected.GET("/workflows/:id/versions/:version", wfH.GetVersion)
// Workflow instances + assignments
wfInstH := handlers.NewWorkflowInstanceHandler(wfEngine, stores)
protected.POST("/workflows/:id/instances", middleware.RequirePermission(auth.PermWorkflowSubmit, stores), wfInstH.Start)
protected.GET("/workflows/:id/instances", wfInstH.ListInstances)
protected.GET("/workflows/:id/instances/:iid", wfInstH.GetInstance)
protected.POST("/workflows/:id/instances/:iid/advance", middleware.RequirePermission(auth.PermWorkflowSubmit, stores), wfInstH.Advance)
protected.POST("/workflows/:id/instances/:iid/cancel", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfInstH.Cancel)
wfAssignH := handlers.NewWorkflowAssignmentHandler(wfEngine, stores)
protected.POST("/assignments/:id/claim", wfAssignH.Claim)
protected.POST("/assignments/:id/unclaim", wfAssignH.Unclaim)
protected.POST("/assignments/:id/complete", wfAssignH.Complete)
protected.POST("/assignments/:id/cancel", middleware.RequirePermission(auth.PermWorkflowCreate, stores), wfAssignH.Cancel)
protected.GET("/assignments/mine", wfAssignH.ListMine)
// Workflow signoffs
wfSignoffH := handlers.NewWorkflowSignoffHandler(wfEngine, stores)
protected.POST("/instances/:iid/signoffs", wfSignoffH.Submit)
protected.GET("/instances/:iid/signoffs", wfSignoffH.List)
// Surface discovery
pkgH := handlers.NewPackageHandler(stores)
protected.GET("/surfaces", pkgH.ListEnabledSurfaces)
// User package management
userPkgDir := ""
if cfg.StoragePath != "" {
userPkgDir = cfg.StoragePath + "/packages"
}
userPkgH := handlers.NewUserPackageHandler(stores, userPkgDir)
protected.GET("/packages", userPkgH.ListVisiblePackages)
protected.POST("/packages/install", userPkgH.InstallPersonalPackage)
protected.DELETE("/packages/:id", userPkgH.DeletePersonalPackage)
// Connection Type Discovery
connTypeH := handlers.NewConnectionTypeHandler(stores)
protected.GET("/connection-types", connTypeH.ListConnectionTypes)
// Extension Connections (personal scope
connH := handlers.NewConnectionHandler(stores, keyResolver)
protected.GET("/connections", connH.ListConnections)
protected.POST("/connections", connH.CreateConnection)
protected.GET("/connections/resolve", connH.ResolveConnection)
protected.GET("/connections/:id", connH.GetConnection)
protected.PUT("/connections/:id", connH.UpdateConnection)
protected.DELETE("/connections/:id", connH.DeleteConnection)
// User Settings & Profile
settings := handlers.NewSettingsHandler(stores, uekCache)
protected.GET("/profile", settings.GetProfile)
protected.PUT("/profile", settings.UpdateProfile)
protected.POST("/profile/password", settings.ChangePassword)
protected.GET("/settings", settings.GetSettings)
protected.PUT("/settings", settings.UpdateSettings)
// Documentation
docsDir := findDocsDir()
docsH := handlers.NewDocsHandler(docsDir)
protected.GET("/docs", docsH.ListDocs)
protected.GET("/docs/:name", docsH.GetDoc)
// Permission bootstrap — self-service resolved permissions
permH := handlers.NewProfilePermissionsHandler(stores)
protected.GET("/profile/permissions", permH.GetMyPermissions)
// Boot payload — single-call SDK bootstrap
bootH := handlers.NewProfileBootstrapHandler(stores)
protected.GET("/profile/bootstrap", bootH.GetBootstrap)
// Notifications
notifH := handlers.NewNotificationHandler(stores, hub)
protected.GET("/notifications", notifH.List)
protected.GET("/notifications/unread-count", notifH.UnreadCount)
protected.PATCH("/notifications/:id/read", notifH.MarkRead)
protected.POST("/notifications/mark-all-read", notifH.MarkAllRead)
protected.DELETE("/notifications/:id", notifH.Delete)
// Notification preferences
protected.GET("/notifications/preferences", notifH.ListPreferences)
protected.PUT("/notifications/preferences/:type", notifH.SetPreference)
protected.DELETE("/notifications/preferences/:type", notifH.DeletePreference)
// ── Scheduled Tasks — User ────
userSchedH := handlers.NewScheduleHandler(stores, triggerEngine)
protected.GET("/schedules", userSchedH.ListMySchedules)
protected.POST("/schedules", userSchedH.CreateSchedule)
protected.GET("/schedules/:id", userSchedH.GetSchedule)
protected.PUT("/schedules/:id", userSchedH.UpdateSchedule)
protected.DELETE("/schedules/:id", userSchedH.DeleteSchedule)
protected.POST("/schedules/:id/run", userSchedH.RunSchedule)
protected.GET("/schedules/:id/logs", userSchedH.ListScheduleLogs)
// Teams (user: my teams)
teams := handlers.NewTeamHandler(stores, keyResolver)
protected.GET("/teams/mine", teams.MyTeams)
// Groups (user: my groups
groupH := handlers.NewGroupHandler(stores)
protected.GET("/groups/mine", groupH.MyGroups)
// Team admin self-service
teamScoped := protected.Group("/teams/:teamId")
teamScoped.Use(middleware.RequireTeamAdmin(stores.Teams))
{
teamScoped.GET("/members", teams.ListMembers)
teamScoped.POST("/members", teams.AddMember)
teamScoped.PUT("/members/:memberId", teams.UpdateMember)
teamScoped.DELETE("/members/:memberId", teams.RemoveMember)
teamScoped.GET("/roles", teams.ListRoles)
teamScoped.PUT("/roles", teams.UpdateRoles)
// Team groups (team admins manage team-scoped groups)
teamScoped.GET("/groups", groupH.ListTeamGroups)
// Team connections
teamScoped.GET("/connections", teams.ListTeamConnections)
teamScoped.POST("/connections", teams.CreateTeamConnection)
teamScoped.PUT("/connections/:id", teams.UpdateTeamConnection)
teamScoped.DELETE("/connections/:id", teams.DeleteTeamConnection)
// Team package management
teamPkgH := handlers.NewUserPackageHandler(stores, userPkgDir)
teamScoped.POST("/packages/install", teamPkgH.InstallTeamPackage)
teamScoped.DELETE("/packages/:id", teamPkgH.DeleteTeamPackage)
// Team package settings — cascade overrides
teamPkgSettingsH := handlers.NewTeamPackageSettingsHandler(stores)
teamScoped.GET("/packages/:id/settings", teamPkgSettingsH.GetTeamPackageSettings)
teamScoped.PUT("/packages/:id/settings", teamPkgSettingsH.UpdateTeamPackageSettings)
teamScoped.DELETE("/packages/:id/settings", teamPkgSettingsH.DeleteTeamPackageSettings)
// Team audit log (team admins only — RequireTeamAdmin on group)
teamScoped.GET("/audit", teams.ListTeamAuditLog)
teamScoped.GET("/audit/actions", teams.ListTeamAuditActions)
// Team workflows — self-service
teamWfH := handlers.NewWorkflowHandler(stores)
teamScoped.GET("/workflows", teamWfH.ListTeamWorkflows)
teamScoped.POST("/workflows", teamWfH.CreateTeamWorkflow)
teamScoped.GET("/workflows/:id", teamWfH.GetTeamWorkflow)
teamScoped.PATCH("/workflows/:id", teamWfH.UpdateTeamWorkflow)
teamScoped.DELETE("/workflows/:id", teamWfH.DeleteTeamWorkflow)
teamScoped.GET("/workflows/:id/stages", teamWfH.ListTeamWorkflowStages)
teamScoped.POST("/workflows/:id/stages", teamWfH.CreateTeamWorkflowStage)
teamScoped.PUT("/workflows/:id/stages/:sid", teamWfH.UpdateTeamWorkflowStage)
teamScoped.DELETE("/workflows/:id/stages/:sid", teamWfH.DeleteTeamWorkflowStage)
teamScoped.PATCH("/workflows/:id/stages/reorder", teamWfH.ReorderTeamWorkflowStages)
teamScoped.POST("/workflows/:id/publish", teamWfH.PublishTeamWorkflow)
teamScoped.POST("/workflows/:id/clone", teamWfH.CloneTeamWorkflow)
teamScoped.POST("/workflows/:id/adopt", teamWfH.AdoptTeamWorkflow)
teamScoped.GET("/workflows/available", teamWfH.ListGlobalWorkflows)
teamScoped.GET("/workflows/:id/versions/:version", teamWfH.GetTeamWorkflowVersion)
// Team workflow instances + assignments
teamWfInstH := handlers.NewWorkflowInstanceHandler(wfEngine, stores)
teamScoped.POST("/workflows/:id/instances", teamWfInstH.Start)
teamScoped.GET("/workflows/:id/instances", teamWfInstH.ListInstances)
teamScoped.GET("/workflows/:id/instances/:iid", teamWfInstH.GetInstance)
teamScoped.POST("/workflows/:id/instances/:iid/advance", teamWfInstH.Advance)
teamScoped.POST("/workflows/:id/instances/:iid/cancel", teamWfInstH.Cancel)
teamWfAssignH := handlers.NewWorkflowAssignmentHandler(wfEngine, stores)
teamScoped.GET("/assignments", teamWfAssignH.ListByTeam)
// Team workflow signoffs
teamWfSignoffH := handlers.NewWorkflowSignoffHandler(wfEngine, stores)
teamScoped.POST("/instances/:iid/signoffs", teamWfSignoffH.Submit)
teamScoped.GET("/instances/:iid/signoffs", teamWfSignoffH.List)
}
// Public global settings (non-admin users can read safe subset)
adm := handlers.NewAdminHandler(stores, keyResolver, uekCache, objStore)
protected.GET("/settings/public", adm.PublicSettings)
// Extensions (user-facing)
protected.GET("/extensions", extH.ListUserExtensions)
protected.POST("/extensions/:id/settings", extH.UpdateUserExtensionSettings)
protected.GET("/extensions/:id/manifest", extH.GetExtensionManifest)
protected.GET("/extensions/tools", extH.ListBrowserToolSchemas)
}
// ── Admin routes ────────────────────────
admin := api.Group("/admin")
admin.Use(middleware.Auth(cfg, stores.Users, userCache))
admin.Use(middleware.RequireAdmin(stores))
admin.Use(middleware.ValidatePathParams())
{
adm := handlers.NewAdminHandler(stores, keyResolver, uekCache, objStore)
adm.OnUserChanged(userCache.Evict)
// User management
admin.GET("/users", adm.ListUsers)
admin.POST("/users", adm.CreateUser)
admin.PUT("/users/:id/role", adm.UpdateUserRole)
admin.PUT("/users/:id/active", adm.ToggleUserActive)
admin.POST("/users/:id/reset-password", adm.ResetPassword)
admin.POST("/users/:id/vault/reset", adm.ResetVault)
admin.DELETE("/users/:id", adm.DeleteUser)
// Global settings
admin.GET("/settings", adm.ListGlobalSettings)
admin.GET("/settings/:key", adm.GetGlobalSetting)
admin.PUT("/settings/:key", adm.UpdateGlobalSetting)
// Stats
admin.GET("/stats", adm.GetStats)
// Global Connections
admin.GET("/connections", adm.ListGlobalConnections)
admin.POST("/connections", adm.CreateGlobalConnection)
admin.PUT("/connections/:id", adm.UpdateGlobalConnection)
admin.DELETE("/connections/:id", adm.DeleteGlobalConnection)
// Teams (admin)
teamAdm := handlers.NewTeamHandler(stores, keyResolver)
admin.GET("/teams", teamAdm.ListTeams)
admin.POST("/teams", teamAdm.CreateTeam)
admin.GET("/teams/:id", teamAdm.GetTeam)
admin.PUT("/teams/:id", teamAdm.UpdateTeam)
admin.DELETE("/teams/:id", teamAdm.DeleteTeam)
admin.GET("/teams/:id/members", teamAdm.ListMembers)
admin.POST("/teams/:id/members", teamAdm.AddMember)
admin.PUT("/teams/:id/members/:memberId", teamAdm.UpdateMember)
admin.DELETE("/teams/:id/members/:memberId", teamAdm.RemoveMember)
// Admin broadcast
adminNotifH := handlers.NewNotificationHandler(stores, hub)
admin.POST("/notifications/broadcast", adminNotifH.Broadcast)
// Audit log
admin.GET("/audit", adm.ListAuditLog)
admin.GET("/audit/actions", adm.ListAuditActions)
// Groups (admin
groupAdm := handlers.NewGroupHandler(stores)
admin.GET("/groups", groupAdm.ListGroups)
admin.POST("/groups", groupAdm.CreateGroup)
admin.GET("/groups/:id", groupAdm.GetGroup)
admin.PUT("/groups/:id", groupAdm.UpdateGroup)
admin.DELETE("/groups/:id", groupAdm.DeleteGroup)
admin.GET("/groups/:id/members", groupAdm.ListMembers)
admin.POST("/groups/:id/members", groupAdm.AddMember)
admin.DELETE("/groups/:id/members/:userId", groupAdm.RemoveMember)
// Permissions
admin.GET("/permissions", groupAdm.ListPermissions)
admin.GET("/users/:id/permissions", groupAdm.GetUserPermissions)
// Resource Grants (admin
admin.GET("/grants/:type/:id", groupAdm.GetResourceGrant)
admin.PUT("/grants/:type/:id", groupAdm.SetResourceGrant)
admin.DELETE("/grants/:type/:id", groupAdm.DeleteResourceGrant)
// Storage status
storageH := handlers.NewStorageHandler(objStore)
admin.GET("/storage/status", storageH.Status)
// Vault
admin.GET("/vault/status", adm.VaultStatus)
// Email / SMTP test
emailAdm := handlers.NewAdminEmailHandler(stores)
admin.POST("/notifications/test-email", emailAdm.TestEmail)
// Extensions (admin)
extAdm := handlers.NewExtensionHandler(stores)
admin.GET("/extensions", extAdm.AdminListExtensions)
admin.POST("/extensions", extAdm.AdminInstallExtension)
admin.PUT("/extensions/:id", extAdm.AdminUpdateExtension)
admin.DELETE("/extensions/:id", extAdm.AdminUninstallExtension)
// Extension permissions (admin
extPermH := handlers.NewExtPermHandler(stores)
admin.GET("/extensions/:id/permissions", extPermH.ListPackagePermissions)
admin.GET("/extensions/:id/review", extPermH.ReviewPackage)
admin.POST("/extensions/:id/permissions/:perm/grant", extPermH.GrantPermission)
admin.POST("/extensions/:id/permissions/:perm/revoke", extPermH.RevokePermission)
admin.POST("/extensions/:id/permissions/grant-all", extPermH.GrantAllPermissions)
// Extension secrets (admin
extSecH := handlers.NewExtSecretsHandler(stores)
admin.GET("/extensions/:id/secrets", extSecH.GetSecrets)
admin.PUT("/extensions/:id/secrets", extSecH.SetSecrets)
admin.DELETE("/extensions/:id/secrets", extSecH.DeleteSecrets)
// Package lifecycle management
packagesDir := ""
if cfg.StoragePath != "" {
packagesDir = cfg.StoragePath + "/packages"
}
pkgAdm := handlers.NewPackageHandler(stores, packagesDir)
pkgAdm.SetSandbox(sandbox.New(sandbox.DefaultConfig()))
// Package registry — must be registered before /packages/:id
registryH := handlers.NewRegistryHandler(stores, packagesDir, pkgAdm)
admin.GET("/packages/registry", registryH.BrowseRegistry)
admin.POST("/packages/registry/install", registryH.InstallFromRegistry)
admin.GET("/packages", pkgAdm.ListPackages)
admin.GET("/packages/:id", pkgAdm.GetPackage)
admin.POST("/packages/install", pkgAdm.InstallPackage)
admin.PUT("/packages/:id/enable", pkgAdm.EnablePackage)
admin.PUT("/packages/:id/disable", pkgAdm.DisablePackage)
admin.DELETE("/packages/:id", pkgAdm.DeletePackage)
admin.GET("/packages/:id/settings", pkgAdm.GetPackageSettings)
admin.PUT("/packages/:id/settings", pkgAdm.UpdatePackageSettings)
admin.GET("/packages/:id/dependencies", pkgAdm.ListDependencies)
admin.GET("/packages/:id/consumers", pkgAdm.ListConsumers)
admin.POST("/packages/:id/test-tool", pkgAdm.TestTool)
admin.POST("/packages/:id/update", pkgAdm.UpdatePackage)
admin.GET("/packages/:id/export", pkgAdm.ExportPackage)
admin.GET("/dependencies", pkgAdm.ListAllDependencies)
// Workflow package export
wfPkgH := handlers.NewWorkflowPackageHandler(stores)
admin.GET("/workflows/:id/export", wfPkgH.ExportWorkflowPackage)
// ── Triggers ─────────────────
trigH := handlers.NewTriggerHandler(stores, triggerEngine)
admin.GET("/triggers", trigH.ListTriggers)
admin.GET("/triggers/:id", trigH.GetTrigger)
admin.PUT("/triggers/:id/enable", trigH.EnableTrigger)
admin.PUT("/triggers/:id/disable", trigH.DisableTrigger)
admin.DELETE("/triggers/:id", trigH.DeleteTrigger)
admin.GET("/triggers/:id/logs", trigH.ListTriggerLogs)
admin.GET("/packages/:id/triggers", trigH.ListPackageTriggers)
// ── Scheduled Tasks — Admin ──
schedH := handlers.NewScheduleHandler(stores, triggerEngine)
admin.GET("/schedules", schedH.AdminListSchedules)
admin.PUT("/schedules/:id/enable", schedH.AdminEnableSchedule)
admin.PUT("/schedules/:id/disable", schedH.AdminDisableSchedule)
admin.DELETE("/schedules/:id", schedH.AdminDeleteSchedule)
// ── Cluster ─────────────────
if stores.Cluster != nil {
clusterH := handlers.NewClusterHandler(stores)
admin.GET("/cluster", clusterH.ListNodes)
}
// ── Metrics ─────────────────
metricsCollector := metrics.NewCollector(
nodeID, database.DB, hub, bus, stores,
sandbox.SandboxStats,
triggerEngine.FireCount,
startTime,
)
metricsH := handlers.NewMetricsHandler(metricsCollector)
admin.GET("/metrics", metricsH.GetMetrics)
// ── Backup/Restore ─────────
backupH := handlers.NewBackupHandler(stores, packagesDir, cfg.StoragePath)
admin.POST("/backup", backupH.CreateBackup)
admin.GET("/backups", backupH.ListBackups)
admin.GET("/backups/:name", backupH.DownloadBackup)
admin.DELETE("/backups/:name", backupH.DeleteBackup)
admin.POST("/restore", backupH.RestoreBackup)
// Surface aliases (backward compat — same handlers)
admin.GET("/surfaces", pkgAdm.ListPackages)
admin.GET("/surfaces/:id", pkgAdm.GetPackage)
admin.POST("/surfaces/install", pkgAdm.InstallPackage)
admin.PUT("/surfaces/:id/enable", pkgAdm.EnablePackage)
admin.PUT("/surfaces/:id/disable", pkgAdm.DisablePackage)
admin.DELETE("/surfaces/:id", pkgAdm.DeletePackage)
}
}
// ── Page Routes ──────────────────────────
// Served without auth — same rationale as extension assets (script tags can't send headers).
// In split deployment, nginx serves these from /data/packages/ directly.
if cfg.StoragePath != "" {
packagesStaticDir := cfg.StoragePath + "/packages"
base.GET("/surfaces/:id/*path", func(c *gin.Context) {
id := c.Param("id")
filePath := c.Param("path")
// Security: prevent path traversal
clean := filepath.Clean(filepath.Join(packagesStaticDir, id, filePath))
if !strings.HasPrefix(clean, filepath.Clean(packagesStaticDir)) {
c.Status(http.StatusForbidden)
return
}
c.File(clean)
})
}
pages.SetVersion(Version)
pageEngine := pages.New(cfg, stores)
// Root redirect → configurable default surface
base.GET("/", pageEngine.DefaultSurfaceRedirect())
// Login page — no auth required
base.GET("/login", pageEngine.RenderLogin())
// Replaces manual per-surface route blocks.
pageEngine.RegisterPageRoutes(base, pages.PageRouteMiddleware{
Authenticated: middleware.AuthOrRedirect(cfg, stores.Users, userCache),
Admin: []gin.HandlerFunc{
middleware.AuthOrRedirect(cfg, stores.Users, userCache),
middleware.RequireAdminPage(stores),
},
Session: middleware.AuthOrRedirect(cfg, stores.Users, userCache), // TODO: session middleware for workflow visitors
})
// Mounted at /s/:slug/api/* with JWT auth (returns 401, not redirect).
{
extAPIH := handlers.NewExtAPIHandler(stores, starlarkRunner)
extAPI := base.Group("/s/:slug/api")
extAPI.Use(middleware.Auth(cfg, stores.Users, userCache))
extAPI.Any("/*path", extAPIH.Handle)
}
bp := cfg.BasePath
if bp == "" {
bp = "/"
}
log.Printf("🔀 Switchboard Core v%s starting on port %s", Version, cfg.Port)
log.Printf(" Base path: %s", bp)
log.Printf(" Schema: %s", database.SchemaVersion())
if objStore != nil {
log.Printf(" Storage: %s", objStore.Backend())
} else {
log.Printf(" Storage: disabled")
}
log.Printf(" EventBus: ready, WebSocket on %s/ws", cfg.BasePath)
log.Printf(" Pages: template engine active (%d surfaces registered)", len(pageEngine.Surfaces()))
if err := r.Run(":" + cfg.Port); err != nil {
log.Fatalf("Failed to start server: %v", err)
}
}
// appendClusterHealth adds cluster info to a health response if the registry is active.
func appendClusterHealth(info gin.H, reg *cluster.Registry, stores store.Stores) {
if reg == nil || stores.Cluster == nil {
return
}
nodes, err := stores.Cluster.ListNodes(context.Background())
if err != nil {
return
}
var peers []string
var heartbeatAge time.Duration
for _, n := range nodes {
if n.NodeID != reg.NodeID() {
peers = append(peers, n.NodeID)
} else {
heartbeatAge = time.Since(n.Heartbeat)
}
}
info["node_id"] = reg.NodeID()
info["cluster"] = gin.H{
"size": len(nodes),
"peers": peers,
"heartbeat_age_ms": heartbeatAge.Milliseconds(),
}
}
// findDocsDir locates the docs/ directory. Checks common locations.
func findDocsDir() string {
candidates := []string{
"docs",
"../docs",
"/app/docs",
}
for _, dir := range candidates {
if info, err := os.Stat(dir); err == nil && info.IsDir() {
return dir
}
}
return ""
}
// ── Vault CLI Commands ──────────────────────
func runVaultCommand(subcmd string) {
cfg := config.Load()
switch strings.ToLower(subcmd) {
case "rekey":
if cfg.DatabaseURL == "" {
log.Fatal("DATABASE_URL is required for vault operations")
}
if err := database.Connect(cfg); err != nil {
log.Fatalf("❌ Database connection failed: %v", err)
}
defer database.Close()
oldKey := os.Getenv("ENCRYPTION_KEY")
newKey := os.Getenv("NEW_ENCRYPTION_KEY")
if err := crypto.Rekey(database.DB, oldKey, newKey); err != nil {
log.Fatalf("❌ Vault rekey failed: %v", err)
}
case "status":
if cfg.DatabaseURL == "" {
log.Fatal("DATABASE_URL is required for vault operations")
}
if err := database.Connect(cfg); err != nil {
log.Fatalf("❌ Database connection failed: %v", err)
}
defer database.Close()
status, err := crypto.VaultStatus(database.DB, cfg.EncryptionKey)
if err != nil {
log.Fatalf("❌ Failed to get vault status: %v", err)
}
fmt.Printf("Encryption key set: %v\n", status.EncryptionKeySet)
fmt.Printf("Encrypted keys: %d\n", status.EncryptedKeys)
fmt.Printf("Vault users (active): %d\n", status.VaultUsers)
default:
fmt.Fprintf(os.Stderr, "Unknown vault command: %s\n", subcmd)
fmt.Fprintf(os.Stderr, "Usage: switchboard vault <rekey|status>\n")
os.Exit(1)
}
}