Files
ai-gateway-go/cmd/gateway-api/main.go
T
LLMGuardX Dev 8000bccde3 0.11.6: 渠道权限管控(部门范围 + 用户级授权)
- 迁移 000047:channels.department_ids(空=全局) + channel_grants 用户级授权
  (source=manual/approval 区分来源)。
- 管理端:渠道部门范围配置 + 授权管理(列表/授予/撤销);渠道列表显示范围。
- 门户:我的渠道端点(/api/v1/portal/channels)按部门可见或明确授权返回,
  「个人渠道」页新增可使用渠道区(授权方式标识)。
- 审批流:资源申请中的渠道类型通过后自动写 channel_grants(source=approval),
  取代'批准记录即授权'的弱语义。
- 端到端验证:部门隔离(demo 无部门看不到)→手动授予→可见→撤销→不可见;
  审批通过自动授权。修复 JOIN 列歧义与 uuid/text 比较。
2026-08-13 15:03:39 +08:00

556 lines
28 KiB
Go

package main
import (
"context"
"errors"
"log/slog"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"aigateway.local/core/internal/agentnode"
"aigateway.local/core/internal/apikey"
"aigateway.local/core/internal/assistant"
"aigateway.local/core/internal/audit"
"aigateway.local/core/internal/channel"
"aigateway.local/core/internal/contentpolicy"
"aigateway.local/core/internal/factcheck"
"aigateway.local/core/internal/gateway"
"aigateway.local/core/internal/identity"
"aigateway.local/core/internal/memory"
"aigateway.local/core/internal/modelquota"
"aigateway.local/core/internal/operations"
"aigateway.local/core/internal/outbox"
"aigateway.local/core/internal/platform/cache"
"aigateway.local/core/internal/platform/config"
"aigateway.local/core/internal/platform/cryptox"
"aigateway.local/core/internal/platform/database"
"aigateway.local/core/internal/platform/health"
"aigateway.local/core/internal/platform/httpserver"
"aigateway.local/core/internal/platform/license"
"aigateway.local/core/internal/platform/storage"
"aigateway.local/core/internal/portal"
"aigateway.local/core/internal/pricing"
"aigateway.local/core/internal/provider"
providercontrolplane "aigateway.local/core/internal/provider/controlplane"
provideropenai "aigateway.local/core/internal/provider/openai"
providerruntime "aigateway.local/core/internal/provider/runtime"
"aigateway.local/core/internal/scheduler"
"aigateway.local/core/internal/shadow"
"aigateway.local/core/internal/trace"
"aigateway.local/core/internal/workbench"
)
var version = "dev"
func main() {
logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
cfg, err := config.Load()
if err != nil {
logger.Error("invalid configuration", "error", err)
os.Exit(1)
}
if err := cfg.ValidateRuntime(); err != nil {
logger.Error("invalid runtime configuration", "error", err)
os.Exit(1)
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
db, err := database.Open(ctx, cfg.Database)
if err != nil {
logger.Error("database initialization failed", "error", err)
os.Exit(1)
}
if db != nil {
defer db.Close()
}
criticalRedis, err := cache.Open(cfg.Redis.CriticalURL)
if err != nil {
logger.Error("critical redis initialization failed", "error", err)
os.Exit(1)
}
if criticalRedis != nil {
defer criticalRedis.Close()
}
cacheRedis, err := cache.Open(cfg.Redis.CacheURL)
if err != nil {
logger.Error("cache redis initialization failed", "error", err)
os.Exit(1)
}
if cacheRedis != nil {
defer cacheRedis.Close()
}
checker := health.Checker{Timeout: 1500 * time.Millisecond}
checker.Dependencies = append(checker.Dependencies,
health.Dependency{Name: "postgres", Required: true, Probe: func(probeCtx context.Context) error {
if db == nil {
return errors.New("not configured")
}
return db.Ping(probeCtx)
}},
health.Dependency{Name: "redis_critical", Required: true, Probe: func(probeCtx context.Context) error {
if criticalRedis == nil {
return errors.New("not configured")
}
return criticalRedis.Ping(probeCtx).Err()
}},
health.Dependency{Name: "redis_cache", Required: false, Probe: func(probeCtx context.Context) error {
if cacheRedis == nil {
return errors.New("not configured")
}
return cacheRedis.Ping(probeCtx).Err()
}},
)
adapter, err := provideropenai.New(cfg.Upstream.BaseURL, cfg.Upstream.APIKey)
if err != nil {
logger.Error("provider initialization failed", "error", err)
os.Exit(1)
}
var fallbackAdapter provider.Adapter
if cfg.Upstream.FallbackEnabled {
fallbackAdapter = adapter
}
credentialCipher, err := provider.NewCredentialCipher(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring,
)
if err != nil {
logger.Error("credential encryption initialization failed", "error", err)
os.Exit(1)
}
providerRepository := provider.NewRepository(db)
providerResolver := providerruntime.NewResolver(
providerRepository, credentialCipher, fallbackAdapter, cfg.Credentials.ProviderRefreshInterval, logger,
)
providerResolver.SetNotificationClient(criticalRedis)
go providerResolver.Run(ctx)
apiKeyRepository := apikey.NewRepository(db)
bootstrapAPIKey := ""
if cfg.Security.BootstrapAPIKeyEnabled {
bootstrapAPIKey = cfg.Security.BootstrapAPIKey
logger.Warn("bootstrap API key compatibility is enabled; monitor usage and disable after migration")
}
apiKeyAuthenticator := apikey.NewAuthenticator(apiKeyRepository, criticalRedis, bootstrapAPIKey)
apiKeyAuthenticator.SetLogger(logger)
proxy := gateway.NewDynamicProxy(providerResolver, apiKeyAuthenticator, cfg.Server.MaxBodyBytes, logger)
proxy.SetAllowPrivateProviderURLs(cfg.Credentials.AllowPrivateProviderURL)
proxy.SetAdmissionController(gateway.NewRedisAdmissionController(criticalRedis))
proxy.SetTokenQuotaController(gateway.NewRedisTokenQuotaController(criticalRedis))
// M8+ 模型级 Token 配额:企业模型总配额(所有 Key 共享),与 Key 级配额叠加。
modelQuotaService := modelquota.NewService(db, criticalRedis, logger)
if err := modelQuotaService.Reload(ctx); err != nil {
logger.Warn("model quota initial load failed; quota checks disabled until refresh", "error", err)
}
go modelQuotaService.Run(ctx, cfg.RuntimeData.PricingRefreshInterval)
proxy.SetModelQuotaController(modelQuotaService)
proxy.SetResiliencePolicy(gateway.ResiliencePolicy{
ResponseHeaderTimeout: cfg.Upstream.ResponseHeaderTimeout, MaxRetries: cfg.Upstream.MaxRetries,
RetryBackoff: cfg.Upstream.RetryBackoff, CircuitThreshold: cfg.Upstream.CircuitThreshold,
CircuitOpenDuration: cfg.Upstream.CircuitOpenDuration,
})
auditRecorder := audit.NewRecorder(db, logger, cfg.Audit.QueueSize, cfg.Audit.BatchSize, cfg.Audit.FlushInterval)
auditContext, stopAudit := context.WithCancel(context.Background())
auditStopped := make(chan struct{})
go func() {
auditRecorder.Run(auditContext)
close(auditStopped)
}()
proxy.SetAuditRecorder(auditRecorder)
contentPolicyEngine := contentpolicy.NewEngine(db, cfg.RuntimeData.ContentPolicyRefreshInterval, logger)
pricingService := pricing.NewService(db, cfg.RuntimeData.PricingRefreshInterval, logger)
if err := contentPolicyEngine.Reload(ctx); err != nil {
logger.Error("content policy initialization failed", "error", err)
os.Exit(1)
}
if err := pricingService.Reload(ctx); err != nil {
logger.Error("model pricing initialization failed", "error", err)
os.Exit(1)
}
go contentPolicyEngine.Run(ctx)
go pricingService.Run(ctx)
proxy.SetContentPolicyEngine(contentPolicyEngine)
proxy.SetOutputPolicyEngine(contentPolicyEngine)
proxy.SetPricingService(pricingService)
identityRepository := identity.NewRepository(db)
sessionStore := identity.NewSessionStore(criticalRedis, cfg.Auth.SessionTTL)
loginLimiter := identity.NewLoginLimiter(criticalRedis, cfg.Auth.LoginRateLimitMax, cfg.Auth.LoginRateLimitWindow, cfg.Auth.TrustedProxies)
totpCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "totp-secret",
)
if err != nil {
logger.Error("TOTP encryption initialization failed", "error", err)
os.Exit(1)
}
identityService := identity.NewService(identityRepository, sessionStore, loginLimiter, cfg.Auth, totpCipher)
idpCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "identity-provider-credentials",
)
if err != nil {
logger.Error("identity provider encryption initialization failed", "error", err)
os.Exit(1)
}
identityService.SetIdentityProviderCipher(idpCipher, cfg.Credentials.AllowPrivateProviderURL)
identityHandler := identity.NewHTTPHandler(identityService)
identityManagementHandler := identity.NewManagementHTTPHandler(identityService)
providerHandler := provider.NewAdminHTTPHandler(
providerRepository, credentialCipher, identityService, cfg.Credentials.AllowPrivateProviderURL,
)
providerOperations := providercontrolplane.NewService(
providerRepository, credentialCipher, cfg.Credentials.AllowPrivateProviderURL,
)
providerHandler.SetOperations(providerOperations)
providerHandler.SetChangeHook(func(changeCtx context.Context) error {
reloadErr := providerResolver.Reload(changeCtx)
notifyErr := providerResolver.Notify(changeCtx)
if reloadErr != nil || notifyErr != nil {
logger.Warn("provider change propagation was incomplete", "reload_error", reloadErr, "notify_error", notifyErr)
}
return errors.Join(reloadErr, notifyErr)
})
apiKeyHandler := apikey.NewAdminHTTPHandler(apiKeyRepository, apiKeyAuthenticator, identityService)
apiKeyHandler.SetUsageStore(apikey.NewUsageStore(criticalRedis))
auditHandler := audit.NewAdminHTTPHandler(audit.NewQueryService(db), identityService)
outboxHandler := outbox.NewAdminHTTPHandler(outbox.NewStore(db), identityService)
contentPolicyHandler := contentpolicy.NewAdminHTTPHandler(contentpolicy.NewStore(db), contentPolicyEngine, identityService)
pricingHandler := pricing.NewAdminHTTPHandler(pricingService, identityService)
factCheckHandler := factcheck.NewAdminHTTPHandler(factcheck.NewService(db), identityService)
workbenchService := workbench.NewService(db)
// M8 P2:本地 Ollama 向量化。EMBEDDINGS_ENABLED=false 时不构造 embedder,
// 知识库检索自动回退纯 FTS;Ollama 挂时入库降级(embedding 置 NULL)。
if cfg.Embeddings.Enabled {
workbenchService.SetEmbedder(workbench.NewOllamaEmbedder(workbench.OllamaEmbedderConfig{
BaseURL: cfg.Embeddings.BaseURL,
Model: cfg.Embeddings.Model,
Dim: cfg.Embeddings.Dim,
BatchSize: cfg.Embeddings.BatchSize,
Timeout: cfg.Embeddings.Timeout,
}))
logger.Info("knowledge embeddings enabled", "model", cfg.Embeddings.Model, "base_url", cfg.Embeddings.BaseURL)
} else {
logger.Info("knowledge embeddings disabled, knowledge retrieval uses postgres_fts only")
}
// 记忆管理(旗舰版):多层记忆 CRUD + 语义召回 + 授权;embedder 与知识库共用。
memoryService := memory.NewService(db, nil)
if cfg.Embeddings.Enabled {
memoryService.SetEmbedder(workbench.NewOllamaEmbedder(workbench.OllamaEmbedderConfig{
BaseURL: cfg.Embeddings.BaseURL, Model: cfg.Embeddings.Model,
Dim: cfg.Embeddings.Dim, BatchSize: cfg.Embeddings.BatchSize, Timeout: cfg.Embeddings.Timeout,
}))
}
memoryHandler := memory.NewHTTPHandler(memoryService, identityService)
toolCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "tool-request-headers",
)
if err != nil {
logger.Error("tool credential encryption initialization failed", "error", err)
os.Exit(1)
}
notificationCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "notification-signing-secret",
)
if err != nil {
logger.Error("notification encryption initialization failed", "error", err)
os.Exit(1)
}
envVarCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "user-env-var",
)
if err != nil {
logger.Error("user env var encryption initialization failed", "error", err)
os.Exit(1)
}
envVarService := workbench.NewEnvVarService(db, envVarCipher)
envVarHandler := workbench.NewEnvVarHTTPHandler(envVarService, identityService)
adminEnvVarHandler := workbench.NewAdminEnvVarHTTPHandler(envVarService, identityService)
toolService := workbench.NewToolService(workbenchService, toolCipher, cfg.Credentials.AllowPrivateToolURL)
notificationService := workbench.NewNotificationService(workbenchService, notificationCipher, cfg.Credentials.AllowPrivateWebhookURL)
workbenchHandler := workbench.NewAdminHTTPHandler(workbenchService, toolService, notificationService, identityService)
mcpServerCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "mcp-server-headers",
)
if err != nil {
logger.Error("MCP server credential encryption initialization failed", "error", err)
os.Exit(1)
}
mcpServerService := workbench.NewMCPServerService(workbenchService, mcpServerCipher, cfg.Credentials.AllowPrivateToolURL)
skillService := workbench.NewSkillService(workbenchService)
digitalEmployeeService := workbench.NewDigitalEmployeeService(workbenchService, skillService, toolService, mcpServerService)
marketplaceService := workbench.NewMarketplaceService(workbenchService, mcpServerService, skillService, digitalEmployeeService)
mcpClient := workbench.NewMCPClient(cfg.Credentials.AllowPrivateToolURL, 60*time.Second)
marketplaceHandler := workbench.NewMarketplaceAdminHTTPHandler(marketplaceService, mcpServerService, skillService, digitalEmployeeService, mcpClient, identityService)
// M8: 对象存储(MinIO)文件管理。文件体在 MinIO,元数据在 PostgreSQL。MinIO
// 不暴露主机端口,上传/下载全部经网关代理,凭据只留在 API 容器内。
objectStore, err := storage.NewClient(storage.Config{
Endpoint: cfg.ObjectStorage.Endpoint,
AccessKeyID: cfg.ObjectStorage.AccessKeyID,
SecretAccessKey: cfg.ObjectStorage.SecretAccessKey,
Bucket: cfg.ObjectStorage.Bucket,
Region: cfg.ObjectStorage.Region,
UseSSL: cfg.ObjectStorage.UseSSL,
MaxFileBytes: cfg.ObjectStorage.MaxFileBytes,
})
if err != nil {
logger.Error("object storage client initialization failed", "error", err)
os.Exit(1)
}
if err := objectStore.EnsureBucket(ctx); err != nil {
// 桶未就绪不致命:MinIO 起来后会自动建桶,当前上传请求会得到明确报错。
logger.Warn("object storage bucket is not ready; uploads fail until MinIO is reachable", "error", err)
}
fileService := workbench.NewFileService(workbenchService, objectStore)
filesAdminHandler := workbench.NewFilesAdminHTTPHandler(fileService, identityService)
filesPortalHandler := workbench.NewFilesPortalHTTPHandler(fileService, identityService)
// M8 P4:站内消息。未读数以 PostgreSQL 为权威源,inbox service 仅用 Redis PUBLISH
// 提醒订阅方;通知 worker 在同一消费循环内物化事件(见 gateway-notification-worker)。
inboxService := workbench.NewInboxService(workbenchService, criticalRedis, cfg.Inbox.Channel)
inboxAdminHandler := workbench.NewInboxAdminHTTPHandler(inboxService, identityService)
inboxPortalHandler := workbench.NewInboxPortalHTTPHandler(inboxService, identityService)
schedulerCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "scheduled-task-api-key",
)
if err != nil {
logger.Error("scheduled task encryption initialization failed", "error", err)
os.Exit(1)
}
schedulerService := scheduler.NewService(db, schedulerCipher)
schedulerHandler := scheduler.NewAdminHTTPHandler(schedulerService, identityService)
portalSchedulerHandler := scheduler.NewPortalHTTPHandler(schedulerService, identityService)
traceStore := trace.NewStore(db)
traceHandler := trace.NewAdminHTTPHandler(traceStore, identityService)
agentNodeStore := agentnode.NewStore(db)
agentNodeHandler := agentnode.NewHTTPHandler(agentNodeStore, identityService)
agentPolicyHandler := workbench.NewAgentPolicyHTTPHandler(workbench.NewAgentPolicyService(db), identityService)
applicationKeyCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "application-runtime-key",
)
if err != nil {
logger.Error("application runtime credential initialization failed", "error", err)
os.Exit(1)
}
channelCipher, err := cryptox.NewKeyring(
cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "channel-config",
)
if err != nil {
logger.Error("channel encryption initialization failed", "error", err)
os.Exit(1)
}
channelService := channel.NewService(db, cfg.Scheduler.GatewayBaseURL, channelCipher, logger)
channelHandler := channel.NewHTTPHandler(channelService, identityService)
channelInboundHandler := channel.NewInboundHTTPHandler(channelService)
shadowMiddleware := shadow.New(cfg.Shadow, logger)
governedGateway := shadowMiddleware.Wrap(proxy)
workbenchRuntime := workbench.NewRuntimeHTTPHandler(workbenchService, toolService, workbench.NewRetriever(workbenchService, workbenchService.Embedder()), apiKeyAuthenticator, governedGateway, workbench.MarketplaceDeps{
MCPServers: mcpServerService,
Skills: skillService,
Employees: digitalEmployeeService,
Market: marketplaceService,
MCPClient: mcpClient,
})
workbenchRuntime.SetLogger(logger)
workbenchRuntime.SetTraceStore(traceStore)
workbenchRuntime.SetEnvVarService(envVarService)
// Wire the fact-check engine: the admin fact-check settings/policies UI now
// actually governs application answers instead of being inert configuration.
factCheckEngine := factcheck.NewEngine(db, workbench.NewFactCheckRetriever(workbench.NewRetriever(workbenchService, workbenchService.Embedder())), logger)
workbenchRuntime.SetFactCheckEngine(factCheckEngine)
portalService := portal.NewService(db, workbenchService, toolService, identityService)
portalService.SetEnvVarService(envVarService)
portalService.SetApplicationRuntime(portal.NewRuntimeCredentials(db, apiKeyRepository, applicationKeyCipher), workbenchRuntime)
portalService.SetGateway(governedGateway)
portalService.SetMarketplace(marketplaceService)
portalService.SetChannelService(channelService)
portalHandler := portal.NewHTTPHandler(portalService, identityService)
portalAdminHandler := portal.NewAdminHTTPHandler(portalService, identityService)
// License 授权:文件校验 + 账号数管控 + 管理端查看/上传。
licenseManager, err := license.NewManager(cfg.License.FilePath, cfg.Credentials.MasterKey)
if err != nil {
logger.Warn("license initialization failed; running as community edition", "error", license.FormatError(err))
}
identityManagementHandler.SetLicenseManager(licenseManager)
licenseHandler := license.NewHTTPHandler(licenseManager, identityService)
startedAt := time.Now()
assistantService := assistant.NewService(db, providerResolver, "", logger)
assistantHandler := assistant.NewHTTPHandler(assistantService, identityService)
operationsHandler := operations.NewAdminHTTPHandler(db, identityService, version, startedAt, func(reloadCtx context.Context) error {
return errors.Join(providerResolver.Reload(reloadCtx), contentPolicyEngine.Reload(reloadCtx), pricingService.Reload(reloadCtx))
})
controlMux := http.NewServeMux()
controlMux.Handle("/api/v1/admin/providers", providerHandler)
controlMux.Handle("/api/v1/admin/providers/", providerHandler)
controlMux.Handle("/api/v1/admin/model-routes", providerHandler)
controlMux.Handle("/api/v1/admin/model-routes/", providerHandler)
controlMux.Handle("/api/v1/admin/api-keys", apiKeyHandler)
controlMux.Handle("/api/v1/admin/api-keys/", apiKeyHandler)
controlMux.Handle("/api/v1/admin/audit-events", auditHandler)
controlMux.Handle("/api/v1/admin/usage/", auditHandler)
controlMux.Handle("/api/v1/admin/outbox-events", outboxHandler)
controlMux.Handle("/api/v1/admin/outbox-events/", outboxHandler)
controlMux.Handle("/api/v1/admin/content-policies", contentPolicyHandler)
controlMux.Handle("/api/v1/admin/content-policies/", contentPolicyHandler)
controlMux.Handle("/api/v1/admin/model-prices", pricingHandler)
controlMux.Handle("/api/v1/admin/model-quotas", modelquota.NewHTTPHandler(modelQuotaService, identityService))
controlMux.Handle("/api/v1/admin/model-prices/", pricingHandler)
controlMux.Handle("/api/v1/admin/fact-check/", factCheckHandler)
controlMux.Handle("/api/v1/admin/prompt-categories", workbenchHandler)
controlMux.Handle("/api/v1/admin/prompt-categories/", workbenchHandler)
controlMux.Handle("/api/v1/admin/prompts", workbenchHandler)
controlMux.Handle("/api/v1/admin/prompts/", workbenchHandler)
controlMux.Handle("/api/v1/admin/knowledge-bases", workbenchHandler)
controlMux.Handle("/api/v1/admin/knowledge-bases/", workbenchHandler)
controlMux.Handle("/api/v1/admin/tools", workbenchHandler)
controlMux.Handle("/api/v1/admin/tools/", workbenchHandler)
controlMux.Handle("/api/v1/admin/tool-approvals", workbenchHandler)
controlMux.Handle("/api/v1/admin/tool-approvals/", workbenchHandler)
controlMux.Handle("/api/v1/admin/applications", workbenchHandler)
controlMux.Handle("/api/v1/admin/applications/", workbenchHandler)
controlMux.Handle("/api/v1/admin/marketplace-categories", marketplaceHandler)
controlMux.Handle("/api/v1/admin/marketplace-categories/", marketplaceHandler)
controlMux.Handle("/api/v1/admin/mcp-servers", marketplaceHandler)
controlMux.Handle("/api/v1/admin/mcp-servers/", marketplaceHandler)
controlMux.Handle("/api/v1/admin/skills", marketplaceHandler)
controlMux.Handle("/api/v1/admin/skills/", marketplaceHandler)
controlMux.Handle("/api/v1/admin/digital-employees", marketplaceHandler)
controlMux.Handle("/api/v1/admin/digital-employees/", marketplaceHandler)
controlMux.Handle("/api/v1/admin/marketplace/", marketplaceHandler)
controlMux.Handle("/api/v1/portal/marketplace", portalHandler)
controlMux.Handle("/api/v1/portal/marketplace/", portalHandler)
controlMux.Handle("/api/v1/admin/notification-channels", workbenchHandler)
controlMux.Handle("/api/v1/admin/notification-channels/", workbenchHandler)
controlMux.Handle("/api/v1/admin/notification-deliveries", workbenchHandler)
controlMux.Handle("/api/v1/admin/notification-deliveries/", workbenchHandler)
controlMux.Handle("/api/v1/admin/models", portalAdminHandler)
controlMux.Handle("/api/v1/admin/model-requests", portalAdminHandler)
controlMux.Handle("/api/v1/admin/model-requests/", portalAdminHandler)
controlMux.Handle("/api/v1/admin/system-info", operationsHandler)
controlMux.Handle("/api/v1/admin/monitoring/overview", operationsHandler)
controlMux.Handle("/api/v1/admin/reports/", operationsHandler)
controlMux.Handle("/api/v1/admin/tenants/", operationsHandler)
controlMux.Handle("/api/v1/admin/files", filesAdminHandler)
controlMux.Handle("/api/v1/admin/files/", filesAdminHandler)
controlMux.Handle("/api/v1/portal/files", filesPortalHandler)
controlMux.Handle("/api/v1/portal/files/", filesPortalHandler)
controlMux.Handle("/api/v1/admin/inbox", inboxAdminHandler)
controlMux.Handle("/api/v1/admin/inbox/", inboxAdminHandler)
controlMux.Handle("/api/v1/portal/inbox", inboxPortalHandler)
controlMux.Handle("/api/v1/portal/inbox/", inboxPortalHandler)
controlMux.Handle("/api/v1/admin/scheduled-tasks", schedulerHandler)
controlMux.Handle("/api/v1/portal/scheduled-tasks", portalSchedulerHandler)
controlMux.Handle("/api/v1/portal/scheduled-tasks/", portalSchedulerHandler)
controlMux.Handle("/api/v1/admin/scheduled-tasks/", schedulerHandler)
controlMux.Handle("/api/v1/admin/traces", traceHandler)
controlMux.Handle("/api/v1/admin/traces/", traceHandler)
controlMux.Handle("/api/v1/admin/agent-sessions", traceHandler)
controlMux.Handle("/api/v1/admin/agent-nodes", agentNodeHandler)
controlMux.Handle("/api/v1/admin/agent-nodes/", agentNodeHandler)
controlMux.Handle("/api/v1/admin/agent-tasks", agentNodeHandler)
controlMux.Handle("/api/v1/admin/agent-tasks/", agentNodeHandler)
controlMux.Handle("/api/v1/agent/nodes/", agentNodeHandler)
controlMux.Handle("/api/v1/admin/channels", channelHandler)
controlMux.Handle("/api/v1/admin/channels/", channelHandler)
controlMux.Handle("/api/v1/admin/assistant", assistantHandler)
controlMux.Handle("/api/v1/admin/assistant/", assistantHandler)
controlMux.Handle("/api/v1/admin/license", licenseHandler)
controlMux.Handle("/api/v1/admin/license/", licenseHandler)
controlMux.Handle("/api/v1/admin/reload", operationsHandler)
controlMux.Handle("/api/v1/admin/identities/", identityManagementHandler)
controlMux.Handle("/api/v1/admin/roles", identityManagementHandler)
controlMux.Handle("/api/v1/admin/roles/", identityManagementHandler)
controlMux.Handle("/api/v1/admin/departments", identityManagementHandler)
controlMux.Handle("/api/v1/admin/departments/", identityManagementHandler)
controlMux.Handle("/api/v1/admin/identity-providers", identityManagementHandler)
controlMux.Handle("/api/v1/admin/identity-providers/", identityManagementHandler)
controlMux.Handle("/api/v1/admin/saml-providers", identityManagementHandler)
controlMux.Handle("/api/v1/admin/saml-providers/", identityManagementHandler)
controlMux.Handle("/api/v1/admin/social-providers", identityManagementHandler)
controlMux.Handle("/api/v1/admin/social-providers/", identityManagementHandler)
controlMux.Handle("/api/v1/portal/applications", portalHandler)
controlMux.Handle("/api/v1/portal/apps/", portalHandler)
controlMux.Handle("/api/v1/portal/catalog", portalHandler)
controlMux.Handle("/api/v1/portal/cost", portalHandler)
controlMux.Handle("/api/v1/portal/docs-info", portalHandler)
controlMux.Handle("/api/v1/portal/knowledge", portalHandler)
controlMux.Handle("/api/v1/portal/logs", portalHandler)
controlMux.Handle("/api/v1/portal/logs/", portalHandler)
controlMux.Handle("/api/v1/portal/model-requests", portalHandler)
controlMux.Handle("/api/v1/portal/env-vars", envVarHandler)
controlMux.Handle("/api/v1/portal/env-vars/", envVarHandler)
controlMux.Handle("/api/v1/portal/memories", memoryHandler)
controlMux.Handle("/api/v1/portal/memories/", memoryHandler)
controlMux.Handle("/api/v1/portal/model-requests/", portalHandler)
controlMux.Handle("/api/v1/portal/password", portalHandler)
controlMux.Handle("/api/v1/portal/prompts", portalHandler)
controlMux.Handle("/api/v1/portal/prompts/", portalHandler)
controlMux.Handle("/api/v1/portal/stats", portalHandler)
controlMux.Handle("/api/v1/portal/tools", portalHandler)
controlMux.Handle("/api/v1/portal/chat/", portalHandler)
controlMux.Handle("/api/v1/portal/agent-policy", agentPolicyHandler)
controlMux.Handle("/api/v1/portal/personal-channels", portalHandler)
controlMux.Handle("/api/v1/portal/personal-channels/", portalHandler)
controlMux.Handle("/api/v1/portal/channels", portalHandler)
controlMux.Handle("/api/v1/portal/digital-employees", portalHandler)
controlMux.Handle("/api/v1/portal/digital-employees/", portalHandler)
controlMux.Handle("/api/v1/portal/resource-requests", portalHandler)
controlMux.Handle("/api/v1/portal/resource-requests/", portalHandler)
controlMux.Handle("/api/v1/admin/resource-requests", portalAdminHandler)
controlMux.Handle("/api/v1/admin/resource-requests/", portalAdminHandler)
controlMux.Handle("/api/v1/admin/env-vars", adminEnvVarHandler)
controlMux.Handle("/api/v1/admin/env-vars/", adminEnvVarHandler)
controlMux.Handle("/api/v1/", identityHandler)
publicMux := http.NewServeMux()
publicMux.Handle("/v1/prompts", workbenchRuntime)
publicMux.Handle("/v1/prompts/", workbenchRuntime)
publicMux.Handle("/v1/knowledge/", workbenchRuntime)
publicMux.Handle("/v1/tools", workbenchRuntime)
publicMux.Handle("/v1/tools/", workbenchRuntime)
publicMux.Handle("/v1/applications/", workbenchRuntime)
publicMux.Handle("/v1/skills/", workbenchRuntime)
publicMux.Handle("/v1/mcp-servers", workbenchRuntime)
publicMux.Handle("/v1/mcp-servers/", workbenchRuntime)
publicMux.Handle("/v1/channels/", channelInboundHandler)
publicMux.Handle("/v1/personal-channels/", portalHandler)
publicMux.Handle("/v1/agent/nodes/", agentNodeHandler)
publicMux.Handle("/v1/digital-employees/", workbenchRuntime)
publicMux.Handle("/v1/", governedGateway)
server := httpserver.New(httpserver.Dependencies{
Config: cfg, Logger: logger, Checker: checker, Gateway: publicMux, Control: controlMux,
Version: version, StartedAt: startedAt,
BootstrapUses: apiKeyAuthenticator.BootstrapUses,
ExtraMetrics: func() string { return auditRecorder.Prometheus() + shadowMiddleware.Prometheus() },
})
serverErrors := make(chan error, 1)
go func() {
logger.Info("gateway API started", "address", server.Addr, "version", version)
serverErrors <- server.ListenAndServe()
}()
select {
case <-ctx.Done():
logger.Info("shutdown requested")
case err := <-serverErrors:
if !errors.Is(err, http.ErrServerClosed) {
logger.Error("gateway API stopped unexpectedly", "error", err)
os.Exit(1)
}
}
shutdownCtx, cancel := context.WithTimeout(context.Background(), cfg.Server.ShutdownTimeout)
defer cancel()
if err := server.Shutdown(shutdownCtx); err != nil {
logger.Error("graceful shutdown failed", "error", err)
_ = server.Close()
os.Exit(1)
}
stopAudit()
select {
case <-auditStopped:
case <-shutdownCtx.Done():
logger.Warn("audit recorder drain timed out")
}
logger.Info("gateway API stopped")
}