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 ────────────────────── 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 (v0.3.8) ─────────────── // 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 (v0.2.2) ───────────────── 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)) // ── Cluster Registry (v0.6.0) ──────────── // PG-backed node self-registration + heartbeat. No-op on SQLite. var clusterReg *cluster.Registry if database.IsPostgres() && stores.Cluster != nil { nodeID := cfg.ClusterNodeID if nodeID == "" { hostname, _ := os.Hostname() nodeID = fmt.Sprintf("%s-%d", hostname, os.Getpid()) } clusterReg = cluster.NewRegistry(nodeID, cfg.ClusterEndpoint, stores.Cluster, hub, cluster.RegistryConfig{ HeartbeatInterval: cfg.ClusterHeartbeatInterval, StaleThreshold: cfg.ClusterStaleThreshold, }) if err := clusterReg.Start(); err != nil { log.Printf("⚠ Cluster registry failed to start: %v", err) clusterReg = nil } else { defer clusterReg.Stop() } } // ── WebSocket Ticket Store (v0.32.0: PG-backed for cross-pod) ───── ticketAdapter := &events.TicketValidatorAdapter{Store: stores.Tickets} // ── Notification Service (v0.20.0) ─────── 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) base.GET("/health", func(c *gin.Context) { info := gin.H{ "status": "ok", "version": Version, "database": database.IsConnected(), "database_name": database.Name(), "schema_version": database.SchemaVersion(), } appendClusterHealth(info, clusterReg, stores) c.JSON(200, info) }) // 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) api.GET("/health", func(c *gin.Context) { info := gin.H{ "status": "ok", "version": Version, "schema_version": database.SchemaVersion(), "database": database.IsConnected(), "database_name": database.Name(), } if database.IsConnected() { info["registration_enabled"] = handlers.IsRegistrationEnabled(stores) } appendClusterHealth(info, clusterReg, stores) c.JSON(200, info) }) 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 (v0.2.2) ────────────── // 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 (v0.3.3) ───── 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 (v0.3.3) ────────── 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 (v0.28.8) ─────────── // Issue a single-use ticket for WebSocket auth. // Client fetches this, then connects with ?ticket=. 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 (v0.23.1) presence := handlers.NewPresenceHandler(stores) protected.POST("/presence/heartbeat", presence.Heartbeat) protected.GET("/presence", presence.Query) // User search (v0.23.2 — DM user picker) protected.GET("/users/search", presence.SearchUsers) // Workflows (v0.26.1 — team-owned staged processes) 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 (v0.3.2) 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 (v0.3.4) wfSignoffH := handlers.NewWorkflowSignoffHandler(wfEngine, stores) protected.POST("/instances/:iid/signoffs", wfSignoffH.Submit) protected.GET("/instances/:iid/signoffs", wfSignoffH.List) // Surface discovery (v0.25.0, v0.28.7: unified packages) pkgH := handlers.NewPackageHandler(stores) protected.GET("/surfaces", pkgH.ListEnabledSurfaces) // User package management (v0.30.0) 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 (v0.38.4) 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 (v0.6.1) docsDir := findDocsDir() docsH := handlers.NewDocsHandler(docsDir) protected.GET("/docs", docsH.ListDocs) protected.GET("/docs/:name", docsH.GetDoc) // Permission bootstrap (v0.37.1) — self-service resolved permissions permH := handlers.NewProfilePermissionsHandler(stores) protected.GET("/profile/permissions", permH.GetMyPermissions) // Boot payload (v0.37.15) — single-call SDK bootstrap bootH := handlers.NewProfileBootstrapHandler(stores) protected.GET("/profile/bootstrap", bootH.GetBootstrap) // Notifications (v0.20.0) 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 (v0.20.0 Phase 3) protected.GET("/notifications/preferences", notifH.ListPreferences) protected.PUT("/notifications/preferences/:type", notifH.SetPreference) protected.DELETE("/notifications/preferences/:type", notifH.DeletePreference) // ── Scheduled Tasks — User (v0.2.2) ──── 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 (v0.38.1) 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 (v0.30.0) teamPkgH := handlers.NewUserPackageHandler(stores, userPkgDir) teamScoped.POST("/packages/install", teamPkgH.InstallTeamPackage) teamScoped.DELETE("/packages/:id", teamPkgH.DeleteTeamPackage) // Team package settings — cascade overrides (v0.2.0) 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 (v0.31.2) 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 (v0.3.2) 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 (v0.3.4) 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 (v0.38.1) 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 (v0.28.6) 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 (v0.24.2) 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 (v0.20.0 Phase 3) 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 (v0.28.7 — replaces surfaces + extensions admin) 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 (v0.30.0) 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 (v0.30.2) wfPkgH := handlers.NewWorkflowPackageHandler(stores) admin.GET("/workflows/:id/export", wfPkgH.ExportWorkflowPackage) // ── Triggers (v0.2.2) ───────────────── 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 (v0.2.2) ── 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 (v0.6.0) ───────────────── if stores.Cluster != nil { clusterH := handlers.NewClusterHandler(stores) admin.GET("/cluster", clusterH.ListNodes) } // ── Backup/Restore (v0.6.1) ───────── 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 (v0.2.1) 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 (v0.2.0) }) // 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 \n") os.Exit(1) } }