package main import ( "context" "encoding/json" "fmt" "log" "os" "strings" "time" "github.com/gin-gonic/gin" "git.gobha.me/xcaliber/chat-switchboard/compaction" "git.gobha.me/xcaliber/chat-switchboard/config" "git.gobha.me/xcaliber/chat-switchboard/crypto" "git.gobha.me/xcaliber/chat-switchboard/database" "git.gobha.me/xcaliber/chat-switchboard/events" "git.gobha.me/xcaliber/chat-switchboard/extraction" "git.gobha.me/xcaliber/chat-switchboard/handlers" "git.gobha.me/xcaliber/chat-switchboard/knowledge" "git.gobha.me/xcaliber/chat-switchboard/middleware" "git.gobha.me/xcaliber/chat-switchboard/providers" "git.gobha.me/xcaliber/chat-switchboard/roles" "git.gobha.me/xcaliber/chat-switchboard/storage" "git.gobha.me/xcaliber/chat-switchboard/store" postgres "git.gobha.me/xcaliber/chat-switchboard/store/postgres" "git.gobha.me/xcaliber/chat-switchboard/tools" "git.gobha.me/xcaliber/chat-switchboard/tools/search" ) 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() // Register LLM providers providers.Init() 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 stores = postgres.NewStores(database.DB) // Bootstrap admin from env (K8s secret) — upserts on every restart handlers.BootstrapAdmin(cfg, stores) // Seed additional users from env (dev/test only, skipped in production) handlers.SeedUsers(cfg, stores) // Seed builtin extensions from disk (idempotent, version-aware) handlers.SeedBuiltinExtensions(stores, "extensions/builtin") // Load search provider config from DB (defaults to DuckDuckGo if not set) loadSearchConfig(stores) } defer database.Close() // ── File Storage ───────────────────────── // 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) // ── Extraction Queue ──────────────────── // Filesystem-based queue for document text extraction (PDF, DOCX, etc.) // Nil if storage is disabled. var extQueue *extraction.Queue if objStore != nil && cfg.StoragePath != "" { q, err := extraction.NewQueue(cfg.StoragePath, cfg.ExtractionConcurrency) if err != nil { log.Printf("⚠ Extraction queue init failed: %v", err) } else { extQueue = q // Recover items stuck in "processing" from previous crash if recovered, err := extQueue.RecoverStale(30 * time.Minute); err != nil { log.Printf("⚠ Extraction recovery failed: %v", err) } else if recovered > 0 { log.Printf(" 📋 Recovered %d stale extraction items", recovered) } } } // Role resolver for model role dispatch (needs stores + vault) roleResolver := roles.NewResolver(stores, keyResolver) // ── Knowledge Base Pipeline ───────────── // Embedder + ingester for RAG document processing (v0.14.0). // Nil-safe: handler checks embedder.IsConfigured() before accepting uploads. kbEmbedder := knowledge.NewEmbedder(roleResolver).WithStores(stores) kbIngester := knowledge.NewIngester(stores, kbEmbedder, objStore, knowledge.DefaultConcurrency) defer kbIngester.Wait() // drain in-flight ingestions on shutdown // Register kb_search tool (late registration — needs stores + embedder) tools.RegisterKBSearch(stores, kbEmbedder) // Register note tools (late registration — needs stores + embedder for semantic search) tools.RegisterNoteTools(stores, kbEmbedder) r := gin.Default() r.Use(middleware.CORS()) // ── Base path group ────────────────────── base := r.Group(cfg.BasePath) // ── EventBus + WebSocket Hub ───────────── bus := events.NewBus() hub := events.NewHub(bus) // Health check (k8s probes hit this directly) base.GET("/health", func(c *gin.Context) { c.JSON(200, gin.H{ "status": "ok", "version": Version, "database": database.IsConnected(), "schema_version": database.SchemaVersion(), }) }) // WebSocket endpoint base.GET("/ws", middleware.Auth(cfg), hub.HandleWebSocket) // ── Auth routes (rate limited) ────────────── auth := handlers.NewAuthHandler(cfg, stores, uekCache) authLimiter := middleware.NewRateLimiter(1, 5) 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(), "providers": providers.List(), } if database.IsConnected() { info["registration_enabled"] = handlers.IsRegistrationEnabled(stores) } c.JSON(200, info) }) authGroup := api.Group("/auth") authGroup.Use(authLimiter.Limit()) { authGroup.POST("/register", auth.Register) authGroup.POST("/login", auth.Login) authGroup.POST("/refresh", auth.Refresh) authGroup.POST("/logout", auth.Logout) } // ── 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) // ── Protected routes ──────────────────── protected := api.Group("") protected.Use(middleware.Auth(cfg)) { // Channels channels := handlers.NewChannelHandler() protected.GET("/channels", channels.ListChannels) protected.POST("/channels", channels.CreateChannel) protected.GET("/channels/:id", channels.GetChannel) protected.PUT("/channels/:id", channels.UpdateChannel) protected.DELETE("/channels/:id", channels.DeleteChannel) // Messages msgs := handlers.NewMessageHandler(keyResolver, stores, hub, objStore) protected.GET("/channels/:id/messages", msgs.ListMessages) protected.POST("/channels/:id/messages", msgs.CreateMessage) // Message tree (forking) protected.GET("/channels/:id/path", msgs.GetActivePath) protected.PUT("/channels/:id/cursor", msgs.UpdateCursor) protected.POST("/channels/:id/messages/:msgId/edit", msgs.EditMessage) protected.POST("/channels/:id/messages/:msgId/regenerate", msgs.Regenerate) protected.GET("/channels/:id/messages/:msgId/siblings", msgs.ListSiblings) // Chat Completions comp := handlers.NewCompletionHandler(keyResolver, stores, hub, objStore) protected.POST("/chat/completions", comp.Complete) protected.GET("/tools", comp.ListTools) // Summarize & Continue (backed by compaction service) compactionSvc := compaction.NewService(stores, roleResolver) summarize := handlers.NewSummarizeHandler(compactionSvc) protected.POST("/channels/:id/summarize", summarize.Summarize) // Provider Configs (user-facing — replaces /api-configs) provCfg := handlers.NewProviderConfigHandler(stores, keyResolver) protected.GET("/api-configs", provCfg.ListConfigs) // backward compat protected.POST("/api-configs", provCfg.CreateConfig) protected.GET("/api-configs/:id", provCfg.GetConfig) protected.PUT("/api-configs/:id", provCfg.UpdateConfig) protected.DELETE("/api-configs/:id", provCfg.DeleteConfig) protected.GET("/api-configs/:id/models", provCfg.ListModels) protected.POST("/api-configs/:id/models/fetch", provCfg.FetchModels) // Models (unified resolver — replaces scattered endpoints) modelH := handlers.NewModelHandler(stores) protected.GET("/models/enabled", modelH.ListEnabledModels) protected.GET("/models", modelH.ListEnabledModels) // alias // Model Preferences modelPrefs := handlers.NewModelPrefsHandler(stores) protected.GET("/models/preferences", modelPrefs.GetPreferences) protected.PUT("/models/preferences", modelPrefs.SetPreference) protected.POST("/models/preferences/bulk", modelPrefs.BulkSetPreferences) // User Settings & Profile settings := handlers.NewSettingsHandler(uekCache) protected.GET("/profile", settings.GetProfile) protected.PUT("/profile", settings.UpdateProfile) protected.POST("/profile/password", settings.ChangePassword) protected.POST("/profile/avatar", settings.UploadAvatar) protected.DELETE("/profile/avatar", settings.DeleteAvatar) protected.GET("/settings", settings.GetSettings) protected.PUT("/settings", settings.UpdateSettings) // Usage (personal) usage := handlers.NewUsageHandler(stores) protected.GET("/usage", usage.PersonalUsage) // Personas (replaces /presets) personas := handlers.NewPersonaHandler(stores) protected.GET("/presets", personas.ListUserPersonas) // backward compat protected.POST("/presets", personas.CreateUserPersona) protected.PUT("/presets/:id", personas.UpdateUserPersona) protected.DELETE("/presets/:id", personas.DeleteUserPersona) protected.POST("/presets/:id/avatar", handlers.UploadPresetAvatar) protected.DELETE("/presets/:id/avatar", handlers.DeletePresetAvatar) // Notes notes := handlers.NewNoteHandler() protected.GET("/notes", notes.List) protected.POST("/notes", notes.Create) protected.GET("/notes/search", notes.Search) protected.GET("/notes/folders", notes.ListFolders) protected.POST("/notes/bulk-delete", notes.BulkDelete) protected.GET("/notes/:id", notes.Get) protected.PUT("/notes/:id", notes.Update) protected.DELETE("/notes/:id", notes.Delete) // Attachments (file upload/download) attachH := handlers.NewAttachmentHandler(stores, objStore, extQueue) protected.POST("/channels/:id/attachments", attachH.Upload) protected.GET("/channels/:id/attachments", attachH.ListByChannel) protected.GET("/attachments/:id", attachH.GetMetadata) protected.GET("/attachments/:id/download", attachH.Download) protected.DELETE("/attachments/:id", attachH.DeleteAttachment) // Hook: clean up storage files when channels are deleted handlers.SetChannelDeleteHook(attachH.CleanupChannelStorage) // Knowledge Bases (RAG — v0.14.0) kbH := handlers.NewKnowledgeBaseHandler(stores, objStore, kbIngester, kbEmbedder) protected.POST("/knowledge-bases", kbH.CreateKB) protected.GET("/knowledge-bases", kbH.ListKBs) protected.GET("/knowledge-bases/:id", kbH.GetKB) protected.PUT("/knowledge-bases/:id", kbH.UpdateKB) protected.DELETE("/knowledge-bases/:id", kbH.DeleteKB) protected.POST("/knowledge-bases/:id/documents", kbH.UploadDocument) protected.GET("/knowledge-bases/:id/documents", kbH.ListDocuments) protected.GET("/knowledge-bases/:id/documents/:docId/status", kbH.GetDocumentStatus) protected.DELETE("/knowledge-bases/:id/documents/:docId", kbH.DeleteDocument) protected.POST("/knowledge-bases/:id/search", kbH.SearchKB) protected.POST("/knowledge-bases/:id/rebuild", kbH.RebuildKB) protected.GET("/channels/:id/knowledge-bases", kbH.GetChannelKBs) protected.PUT("/channels/:id/knowledge-bases", kbH.SetChannelKBs) // Teams (user: my teams) teams := handlers.NewTeamHandler(keyResolver) protected.GET("/teams/mine", teams.MyTeams) // Team admin self-service teamScoped := protected.Group("/teams/:teamId") teamScoped.Use(middleware.RequireTeamAdmin()) { teamScoped.GET("/members", teams.ListMembers) teamScoped.POST("/members", teams.AddMember) teamScoped.PUT("/members/:memberId", teams.UpdateMember) teamScoped.DELETE("/members/:memberId", teams.RemoveMember) teamScoped.GET("/models", teams.ListAvailableModels) // Team providers teamScoped.GET("/providers", teams.ListTeamProviders) teamScoped.POST("/providers", teams.CreateTeamProvider) teamScoped.PUT("/providers/:id", teams.UpdateTeamProvider) teamScoped.DELETE("/providers/:id", teams.DeleteTeamProvider) teamScoped.GET("/providers/:id/models", teams.ListTeamProviderModels) // Team audit log (team admins only — RequireTeamAdmin on group) teamScoped.GET("/audit", teams.ListTeamAuditLog) teamScoped.GET("/audit/actions", teams.ListTeamAuditActions) // Team usage (team admins only — usage against team-owned providers) teamUsage := handlers.NewUsageHandler(stores) teamScoped.GET("/usage", teamUsage.TeamUsage) // Team personas teamPersonas := handlers.NewPersonaHandler(stores) teamScoped.GET("/presets", teamPersonas.ListTeamPersonas) teamScoped.POST("/presets", teamPersonas.CreateTeamPersona) teamScoped.DELETE("/presets/:id", teamPersonas.DeleteTeamPersona) // Team role overrides teamRoles := handlers.NewRolesHandler(stores, roleResolver) teamScoped.GET("/roles", teamRoles.ListTeamRoles) teamScoped.PUT("/roles/:role", teamRoles.UpdateTeamRole) teamScoped.DELETE("/roles/:role", teamRoles.DeleteTeamRole) } // Public global settings (non-admin users can read safe subset) adm := handlers.NewAdminHandler(stores, keyResolver, uekCache) 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)) admin.Use(middleware.RequireAdmin()) { adm := handlers.NewAdminHandler(stores, keyResolver, uekCache) // 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 Provider Configs admin.GET("/configs", adm.ListGlobalConfigs) admin.POST("/configs", adm.CreateGlobalConfig) admin.PUT("/configs/:id", adm.UpdateGlobalConfig) admin.DELETE("/configs/:id", adm.DeleteGlobalConfig) // Model Catalog admin.GET("/models", adm.ListModelConfigs) admin.POST("/models/fetch", adm.FetchModels) admin.PUT("/models/bulk", adm.BulkUpdateModels) admin.PUT("/models/:id", adm.UpdateModelConfig) admin.DELETE("/models/:id", adm.DeleteModelConfig) // Personas (admin global) personaAdm := handlers.NewPersonaHandler(stores) admin.GET("/presets", personaAdm.ListAdminPersonas) admin.POST("/presets", personaAdm.CreateAdminPersona) admin.PUT("/presets/:id", personaAdm.UpdateAdminPersona) admin.DELETE("/presets/:id", personaAdm.DeleteAdminPersona) admin.POST("/presets/:id/avatar", handlers.UploadPresetAvatar) admin.DELETE("/presets/:id/avatar", handlers.DeletePresetAvatar) // Teams (admin) teamAdm := handlers.NewTeamHandler(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) // Audit log admin.GET("/audit", adm.ListAuditLog) admin.GET("/audit/actions", adm.ListAuditActions) // Model Roles rolesH := handlers.NewRolesHandler(stores, roleResolver) admin.GET("/roles", rolesH.ListRoles) admin.GET("/roles/:role", rolesH.GetRole) admin.PUT("/roles/:role", rolesH.UpdateRole) admin.POST("/roles/:role/test", rolesH.TestRole) // Usage & Pricing usageH := handlers.NewUsageHandler(stores) admin.GET("/usage", usageH.AdminUsage) admin.GET("/usage/users/:id", usageH.AdminUserUsage) admin.GET("/usage/teams/:id", usageH.AdminTeamUsage) admin.GET("/pricing", usageH.ListPricing) admin.PUT("/pricing", usageH.UpsertPricing) admin.DELETE("/pricing/:provider/:model", usageH.DeletePricing) // Storage status storageH := handlers.NewStorageHandler(objStore) admin.GET("/storage/status", storageH.Status) // Vault admin.GET("/vault/status", adm.VaultStatus) // Storage management (orphan cleanup) attachAdm := handlers.NewAttachmentHandler(stores, objStore, extQueue) admin.GET("/storage/orphans", attachAdm.OrphanCount) admin.POST("/storage/cleanup", attachAdm.CleanupOrphans) admin.GET("/storage/extraction", attachAdm.ExtractionStatus) // 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) } } bp := cfg.BasePath if bp == "" { bp = "/" } log.Printf("🔀 Chat Switchboard API v%s starting on port %s", Version, cfg.Port) log.Printf(" Base path: %s", bp) log.Printf(" Schema: %s", database.SchemaVersion()) log.Printf(" Providers: %v", providers.List()) if objStore != nil { log.Printf(" Storage: %s", objStore.Backend()) if extQueue != nil { log.Printf(" Extraction: enabled (concurrency=%d)", cfg.ExtractionConcurrency) } else { log.Printf(" Extraction: disabled") } } else { log.Printf(" Storage: disabled") } log.Printf(" EventBus: ready, WebSocket on %s/ws", cfg.BasePath) if err := r.Run(":" + cfg.Port); err != nil { log.Fatalf("Failed to start server: %v", err) } } // ── Vault CLI Commands ────────────────────── // loadSearchConfig reads search provider config from global_config and applies it. // Falls back to DuckDuckGo (the default set in search package init) if not configured. func loadSearchConfig(stores store.Stores) { raw, err := stores.GlobalConfig.Get(context.Background(), "search_config") if err != nil || raw == nil { log.Println("🔍 Search: using default provider (DuckDuckGo)") return } b, _ := json.Marshal(raw) var cfg search.Config if err := json.Unmarshal(b, &cfg); err != nil { log.Printf("⚠️ Failed to parse search config: %v", err) return } if err := search.ApplyConfig(cfg); err != nil { log.Printf("⚠️ Failed to apply search config: %v", err) } } 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) } }