package main import ( "context" "encoding/json" "fmt" "log/slog" "net/http" "os" "os/signal" "sync" "syscall" "time" "github.com/itworx/pulse/internal/agentsource" "github.com/itworx/pulse/internal/agentstore" "github.com/itworx/pulse/internal/alert" "github.com/itworx/pulse/internal/alertapi" "github.com/itworx/pulse/internal/alertcontrol" "github.com/itworx/pulse/internal/alertcontrolapi" "github.com/itworx/pulse/internal/alertdefaults" "github.com/itworx/pulse/internal/alertopsapi" "github.com/itworx/pulse/internal/applicationapi" "github.com/itworx/pulse/internal/array" "github.com/itworx/pulse/internal/arrayapi" "github.com/itworx/pulse/internal/audit" "github.com/itworx/pulse/internal/auth" "github.com/itworx/pulse/internal/authapi" "github.com/itworx/pulse/internal/backup" "github.com/itworx/pulse/internal/backupapi" "github.com/itworx/pulse/internal/config" "github.com/itworx/pulse/internal/container" "github.com/itworx/pulse/internal/containerapi" "github.com/itworx/pulse/internal/correlation" "github.com/itworx/pulse/internal/dashboard" "github.com/itworx/pulse/internal/dashboardapi" "github.com/itworx/pulse/internal/database" "github.com/itworx/pulse/internal/disk" "github.com/itworx/pulse/internal/diskapi" "github.com/itworx/pulse/internal/eventapi" forecastdomain "github.com/itworx/pulse/internal/forecast" "github.com/itworx/pulse/internal/forecastapi" "github.com/itworx/pulse/internal/host" "github.com/itworx/pulse/internal/hostapi" "github.com/itworx/pulse/internal/incident" "github.com/itworx/pulse/internal/incidentapi" "github.com/itworx/pulse/internal/inventory" "github.com/itworx/pulse/internal/inventoryapi" "github.com/itworx/pulse/internal/live" "github.com/itworx/pulse/internal/livesampler" "github.com/itworx/pulse/internal/metriccatalog" "github.com/itworx/pulse/internal/metricquery" "github.com/itworx/pulse/internal/metricsapi" "github.com/itworx/pulse/internal/network" "github.com/itworx/pulse/internal/networkapi" "github.com/itworx/pulse/internal/observability" "github.com/itworx/pulse/internal/onboarding" "github.com/itworx/pulse/internal/onboardingapi" pooldomain "github.com/itworx/pulse/internal/pool" "github.com/itworx/pulse/internal/poolapi" "github.com/itworx/pulse/internal/problem" "github.com/itworx/pulse/internal/process" "github.com/itworx/pulse/internal/processapi" "github.com/itworx/pulse/internal/prometheus" "github.com/itworx/pulse/internal/queryplan" "github.com/itworx/pulse/internal/reverseproxy" "github.com/itworx/pulse/internal/reverseproxyapi" "github.com/itworx/pulse/internal/runtimeconfig" "github.com/itworx/pulse/internal/service" "github.com/itworx/pulse/internal/serviceapi" sharedomain "github.com/itworx/pulse/internal/share" "github.com/itworx/pulse/internal/shareapi" "github.com/itworx/pulse/internal/systemstatus" "github.com/itworx/pulse/internal/systemstatusapi" "github.com/itworx/pulse/internal/widget" "github.com/itworx/pulse/internal/widgetapi" "github.com/itworx/pulse/internal/workerruntime" "github.com/jackc/pgx/v5/pgxpool" ) func main() { logger := slog.New(slog.NewJSONHandler(os.Stdout, nil)) if err := run(logger); err != nil { logger.Error("pulse api failed", "error", err) os.Exit(1) } } func run(logger *slog.Logger) error { runtime, err := runtimeconfig.Load("api") if err != nil { return err } application, err := config.Load() if err != nil { return err } registry, err := metriccatalog.DefaultRegistry() if err != nil { return err } ctx := context.Background() internalMetrics := observability.NewRegistry(time.Now().UTC()) var pool *pgxpool.Pool var inventoryRepo *inventory.Repository var dashboardRepo dashboard.Repository var alertRepo alert.Store var alertControlStore alertcontrol.Store var alertOperationsStore alertopsapi.Store var incidentStore incident.Store var serviceProvider service.Provider = service.UnknownProvider{Reason: "source_unavailable"} var reverseProxyProvider reverseproxy.Provider = reverseproxy.DisabledProvider{SourceID: "reverse-proxy", SourceType: "connector", Reason: "connector_disabled"} var dependencyRepo *service.DependencyRepository var onboardingService onboarding.Service // agentReader is the read half of the pulse-agent telemetry transport. It stays nil // without a database, which keeps every monitoring surface Unknown instead of // inventing state (ADR-0008). var agentReader agentstore.Reader databaseReady := false if application.DatabaseURL != "" { pool, err = database.NewPool(ctx, database.Config{URL: application.DatabaseURL}) if err != nil { return err } defer pool.Close() if err := database.Ping(ctx, pool); err != nil { return err } databaseReady = true agentReader = agentstore.PostgresStore{Pool: pool} dashboardRepo = dashboard.Repository{Pool: pool} alertRepo = alert.Repository{Pool: pool, Registry: registry} if report, seedErr := alertdefaults.Seed(ctx, alertRepo, registry, "system-defaults"); seedErr != nil { return fmt.Errorf("seed alert defaults: %w", seedErr) } else { logger.Info("alert defaults reconciled", "added", report.Added, "existing", report.Existing) } alertControlStore = alertcontrol.Repository{Pool: pool} alertOperationsStore = alert.StateRepository{Pool: pool} incidentStore, err = incident.NewRepository(pool) if err != nil { return err } serviceProvider, err = service.NewPostgresProvider(pool, service.StatusPolicy{}) if err != nil { return err } dependencyRepo, err = service.NewDependencyRepository(pool, audit.PostgresStore{Pool: pool}) if err != nil { return err } inventoryRepo, err = inventory.NewRepository(pool) if err != nil { return err } onboardingService = onboarding.Service{State: onboarding.StateStore{Pool: pool}, Pool: pool, Dashboards: dashboardRepo, Alerts: alertRepo, AuthMode: application.AuthMode, OIDCIssuer: application.OIDCIssuer, OIDCClient: application.OIDCClientID, OIDCRedirect: application.OIDCRedirectURL, Prometheus: application.PrometheusURL != "", Unraid: application.UnraidURL != "" && application.UnraidAPIToken != ""} } var queryService *metricquery.Service var liveSampler live.Sampler var promSource *prometheus.Client if application.PrometheusURL != "" { source, sourceErr := prometheus.New(application.PrometheusURL, nil, prometheus.Limits{Timeout: application.PrometheusTimeout}) if sourceErr != nil { return sourceErr } promSource = source queryService = metricquery.NewService(queryplan.NewPlanner(registry, queryplan.Limits{}), source, nil) sampler, samplerErr := livesampler.New(registry, source, livesampler.Options{}) if samplerErr != nil { return samplerErr } liveSampler = sampler } if err != nil { return err } // publishSourceMetrics refreshes adapter counters into the internal registry just // before it is read, so operators see current Prometheus latency and error counts // rather than the values captured at process start. publishSourceMetrics := func() { if promSource != nil { promSource.PublishMetrics(internalMetrics) } } sessions := auth.NewSlidingSessionManager("pulse_session", application.SessionIdleTTL, application.SessionAbsoluteTTL, application.Environment == config.Production) backupManager := &backup.Manager{Pool: pool, Directory: application.BackupDirectory, Retention: application.BackupRetention} var backupObservation struct { sync.Mutex checkedAt time.Time latest time.Time verified time.Time err error } readBackupObservation := func(ctx context.Context, now time.Time) (time.Time, time.Time, error) { backupObservation.Lock() defer backupObservation.Unlock() if !backupObservation.checkedAt.IsZero() && now.Sub(backupObservation.checkedAt) < 5*time.Minute { return backupObservation.latest, backupObservation.verified, backupObservation.err } results, err := backupManager.List(ctx) latest := time.Time{} if len(results) > 0 { latest = results[0].Created } verified := time.Time{} if err == nil { verified = now } backupObservation.checkedAt, backupObservation.latest, backupObservation.verified, backupObservation.err = now, latest, verified, err return latest, verified, err } invalidateBackupObservation := func() { backupObservation.Lock() defer backupObservation.Unlock() backupObservation.checkedAt = time.Time{} } mux := service.HealthMuxWithReadiness(func() bool { return databaseReady }) mux.HandleFunc("/auth/test-login", func(response http.ResponseWriter, request *http.Request) { if application.Environment == config.Production || application.AuthMode != "mock" { problem.Write(response, request, http.StatusNotFound, "NOT_FOUND", "Not found", "The requested resource does not exist.", nil) return } principal := auth.Principal{Subject: "development-user", Role: auth.RoleAdministrator} if err := sessions.Issue(response, principal, time.Now().UTC()); err != nil { problem.Write(response, request, http.StatusInternalServerError, "SESSION_ERROR", "Session unavailable", "The session could not be created.", nil) return } if pool != nil { if err := audit.RecordSecurityAction(request.Context(), audit.PostgresStore{Pool: pool}, principal.Subject, "auth.test_login", "success", correlation.FromContext(request.Context())); err != nil { problem.Write(response, request, http.StatusServiceUnavailable, "AUDIT_UNAVAILABLE", "Authentication unavailable", "The authentication event could not be recorded.", nil) return } } response.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(response).Encode(map[string]string{"status": "authenticated", "mode": "mock-development"}) }) mux.HandleFunc("/session/logout", func(response http.ResponseWriter, request *http.Request) { sessions.Clear(response, request) response.WriteHeader(http.StatusNoContent) }) if application.AuthMode == "oidc" && application.OIDCIssuer != "" { oidcLogin, loginErr := authapi.New(authapi.Options{ OIDC: auth.OIDCConfig{ Issuer: application.OIDCIssuer, ClientID: application.OIDCClientID, ClientSecret: application.OIDCClientSecret, RedirectURL: application.OIDCRedirectURL, }, RoleMapping: roleMapping(application.OIDCRoleMapping), GroupsClaim: application.OIDCGroupsClaim, Sessions: sessions, Secure: application.Environment == config.Production, Logger: logger, // Failures land on the overview route, where the web app reads the reason // code from the query string, shows a localized notice and strips it from // the URL. There is deliberately no dedicated error page to maintain. ErrorPath: "/", Audit: func(ctx context.Context, actor, result string) error { if pool == nil { return nil } return audit.RecordSecurityAction(ctx, audit.PostgresStore{Pool: pool}, actor, "auth.login", result, correlation.FromContext(ctx)) }, }) if loginErr != nil { return loginErr } mux.Handle("/auth/login", oidcLogin.LoginHandler()) mux.Handle("/auth/callback", oidcLogin.CallbackHandler()) logger.Info("oidc login enabled", "mapped_claims", len(application.OIDCRoleMapping)) } // reportedJobs is the metadata-only view of the worker schedule. The API never runs // these jobs; it reads their recorded outcomes from job_runs so the status surface // reflects what the worker actually did instead of the hardcoded "not recorded" // placeholders it used before the worker runtime existed. reportedJobs := workerruntime.Schedule(workerruntime.ScheduleRuns{}) snapshot := func(requestContext context.Context) (systemstatus.Snapshot, error) { publishSourceMetrics() var auditEvents *int64 var statusOptions []systemstatus.Option now := time.Now().UTC() _, authenticated := auth.PrincipalFromContext(requestContext) statusOptions = append(statusOptions, systemstatus.WithAuthenticatedSession(authenticated)) sourceHealth := make([]systemstatus.SourceHealth, 0, 3) if promSource != nil { sourceHealth = append(sourceHealth, systemstatus.FromDatasource("prometheus", promSource.Health(requestContext), now)) } if agentReader != nil { unraidHealth, storageHealth, healthErr := readAgentSourceHealth(requestContext, agentReader, now) if healthErr != nil { return systemstatus.Snapshot{}, fmt.Errorf("read agent source health: %w", healthErr) } if unraidHealth.ReasonCode != agentsource.ReasonUnavailable || storageHealth.ReasonCode != agentsource.ReasonUnavailable { internalMetrics.SetGauge("pulse_unraid_configured", 1) } sourceHealth = append(sourceHealth, systemstatus.FromDatasource("unraid", unraidHealth, now), systemstatus.FromDatasource("storage", storageHealth, now)) } if len(sourceHealth) > 0 { statusOptions = append(statusOptions, systemstatus.WithSources(sourceHealth...)) } if application.BackupDirectory != "" { latestBackup, verifiedAt, backupErr := readBackupObservation(requestContext, now) if backupErr != nil { statusOptions = append(statusOptions, systemstatus.WithBackupObservation(time.Time{}, time.Time{}, backupErr)) } else if !latestBackup.IsZero() { statusOptions = append(statusOptions, systemstatus.WithBackupObservation(latestBackup, verifiedAt, nil)) } } if pool != nil && databaseReady { var count int64 if err := pool.QueryRow(requestContext, `SELECT count(*) FROM audit_events`).Scan(&count); err != nil { return systemstatus.Snapshot{}, fmt.Errorf("read audit event count: %w", err) } auditEvents = &count var migrationVersion string if err := pool.QueryRow(requestContext, `SELECT id FROM schema_migrations ORDER BY applied_at DESC, id DESC LIMIT 1`).Scan(&migrationVersion); err != nil { return systemstatus.Snapshot{}, fmt.Errorf("read migration version: %w", err) } statusOptions = append(statusOptions, systemstatus.WithMigrationVersion(migrationVersion)) jobs, jobsErr := workerruntime.ReadJobHealth(requestContext, pool, reportedJobs) if jobsErr != nil { // A failed read must not be reported as healthy. Omitting the option // leaves every job component Unknown, which is the honest answer when // the worker's recorded state cannot be established (ADR-0008). logger.Error("read worker job health", "error", jobsErr) } else { statusOptions = append(statusOptions, systemstatus.WithJobs(0, jobs...)) } } return systemstatus.Build(application, databaseReady, now, auditEvents, statusOptions...), nil } internalMetrics.SetGauge("pulse_database_ready", boolMetric(databaseReady)) internalMetrics.SetGauge("pulse_prometheus_configured", boolMetric(application.PrometheusURL != "")) internalMetrics.SetGauge("pulse_unraid_configured", boolMetric(application.UnraidURL != "" && application.UnraidAPIToken != "")) if onboardingService.State.Pool != nil { onboardingService.RuntimeCapabilities = func(requestContext context.Context) ([]onboarding.Capability, error) { status, statusErr := snapshot(requestContext) if statusErr != nil { return nil, statusErr } capabilities := make([]onboarding.Capability, 0, 2) for _, component := range status.Components { if component.ID != "prometheus" && component.ID != "unraid" { continue } state, detail := "unknown", "Geen actuele runtimewaarneming beschikbaar." switch component.State { case systemstatus.StateHealthy: state, detail = "ready", "Actuele telemetrie wordt ontvangen via de veilige runtimebron." case systemstatus.StateDegraded: state, detail = "incomplete", "De runtimebron vraagt aandacht." case systemstatus.StateDisabled: state, detail = "not-ready", "De runtimebron is niet geconfigureerd." } capabilities = append(capabilities, onboarding.Capability{ID: component.ID, State: state, Detail: detail}) } return capabilities, nil } } statusHandler := systemstatusapi.Handler{ Snapshot: snapshot, Diagnostics: func(requestContext context.Context) (systemstatusapi.Diagnostics, error) { status, err := snapshot(requestContext) if err != nil { return systemstatusapi.Diagnostics{}, err } return systemstatusapi.Diagnostics{ Status: status, Config: systemstatusapi.ConfigSummary{ Environment: string(application.Environment), Timezone: application.Timezone, Locale: application.DefaultLocale, AuthMode: application.AuthMode, PublicURLConfigured: application.PublicURL != "", DatabaseConfigured: application.DatabaseURL != "", PrometheusConfigured: application.PrometheusURL != "", UnraidConfigured: application.UnraidURL != "" && application.UnraidAPIToken != "", OIDCConfigured: application.OIDCIssuer != "" && application.OIDCClientID != "", }, Runtime: systemstatusapi.Runtime(), Metrics: internalMetrics.Exposition(time.Now().UTC()), }, nil }, } protectedStatus := withSession(sessions, auth.Require(auth.PermissionView, statusHandler)) protectedDiagnostics := withSession(sessions, auth.Require(auth.PermissionOperate, statusHandler)) mux.Handle("/api/v1/system/status", protectedStatus) mux.Handle("/api/v1/system/diagnostics", protectedDiagnostics) metricsExposition := internalMetrics.Handler() mux.Handle("/api/v1/system/metrics", withSession(sessions, auth.Require(auth.PermissionOperate, http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { publishSourceMetrics() metricsExposition.ServeHTTP(response, request) })))) mux.Handle("/api/v1/system/backups", withSession(sessions, auth.Require(auth.PermissionAdmin, backupapi.Handler{Manager: backupManager, OnCreated: invalidateBackupObservation, Audit: func(ctx context.Context, actor, result string) error { return audit.RecordSecurityAction(ctx, audit.PostgresStore{Pool: pool}, actor, "backup.create", result, correlation.FromContext(ctx)) }}))) onboardingHandler := onboardingapi.Handler{Service: onboardingService, Audit: audit.PostgresStore{Pool: pool}} mux.Handle("/api/v1/onboarding", withSession(sessions, auth.Require(auth.PermissionView, onboardingHandler))) metricsHandler := metricsapi.Handler{Registry: registry} widgetRegistry, err := widget.NewRegistry(widget.DefaultDefinitions()) if err != nil { return err } widgetHandler := widgetapi.Handler{Registry: widgetRegistry} mux.Handle("/api/v1/widgets/catalog", withSession(sessions, auth.Require(auth.PermissionView, widgetHandler))) mux.Handle("/api/v1/widgets/preview", withSession(sessions, auth.Require(auth.PermissionEdit, widgetHandler))) mux.Handle("/api/v1/metrics/catalog", withSession(sessions, auth.Require(auth.PermissionView, metricsHandler))) queryHandler := metricquery.Handler{Service: queryService} mux.Handle("/api/v1/metrics/query", withSession(sessions, auth.Require(auth.PermissionView, queryHandler))) mux.Handle("/api/v1/metrics/query-range", withSession(sessions, auth.Require(auth.PermissionView, queryHandler))) mux.Handle("/api/v1/metrics/inspect", withSession(sessions, auth.Require(auth.PermissionOperate, queryHandler))) livePlanner := queryplan.NewPlanner(registry, queryplan.Limits{}) liveRegistry := live.NewRegistry(liveSampler, live.RegistryOptions{}) liveHandler := live.Handler{Planner: &livePlanner, Registry: liveRegistry} mux.Handle("/api/v1/live", withSession(sessions, auth.Require(auth.PermissionView, liveHandler))) if alertRepo != nil { alertHandler := alertapi.Handler{Repository: alertRepo, Registry: registry, Audit: audit.PostgresStore{Pool: pool}} mux.Handle("/api/v1/alert-rules", withSession(sessions, alertHandler)) mux.Handle("/api/v1/alert-rules/", withSession(sessions, alertHandler)) operationsHandler := alertopsapi.Handler{Store: alertOperationsStore, Audit: audit.PostgresStore{Pool: pool}} protectedOperations := withSession(sessions, auth.Require(auth.PermissionView, operationsHandler)) mux.Handle("/api/v1/alerts", protectedOperations) mux.Handle("/api/v1/alerts/", protectedOperations) controlHandler := alertcontrolapi.Handler{Store: alertControlStore, Audit: audit.PostgresStore{Pool: pool}} protectedControls := withSession(sessions, auth.Require(auth.PermissionView, controlHandler)) mux.Handle("/api/v1/alert-silences", protectedControls) mux.Handle("/api/v1/alert-silences/", protectedControls) mux.Handle("/api/v1/maintenance-windows", protectedControls) mux.Handle("/api/v1/maintenance-windows/", protectedControls) } if incidentStore != nil { incidentHandler := incidentapi.Handler{Store: incidentStore, Audit: audit.PostgresStore{Pool: pool}} protectedIncidents := withSession(sessions, auth.Require(auth.PermissionView, incidentHandler)) mux.Handle("/api/v1/incidents", protectedIncidents) mux.Handle("/api/v1/incidents/", protectedIncidents) } if dashboardRepo.Pool != nil { dashboardHandler := dashboardapi.Handler{Repository: dashboardRepo, Audit: audit.PostgresStore{Pool: pool}} mux.Handle("/api/v1/dashboards", withSession(sessions, dashboardHandler)) mux.Handle("/api/v1/dashboards/", withSession(sessions, dashboardHandler)) } if inventoryRepo != nil { inventoryHandler := inventoryapi.Handler{Repository: inventoryRepo} protectedInventory := withSession(sessions, auth.Require(auth.PermissionView, inventoryHandler)) mux.Handle("/api/v1/entities", protectedInventory) mux.Handle("/api/v1/entities/", protectedInventory) } if pool != nil { eventsHandler := eventapi.Handler{Store: eventapi.PostgresStore{Pool: pool}} mux.Handle("/api/v1/events", withSession(sessions, auth.Require(auth.PermissionView, eventsHandler))) } // Monitoring surfaces are served from the bounded snapshots pulse-agent writes into // PostgreSQL. Without a database there is no transport at all, so each surface keeps // the empty adapter it had before, which resolves to Unknown rather than Healthy. agentWindows := agentsource.Windows{} if err := agentWindows.Validate(); err != nil { return err } var hostProvider host.Provider = host.UnknownProvider{SourceID: "host", SourceType: "agent", Reason: "source_unavailable"} var processProvider interface { Snapshot(context.Context) (process.Snapshot, error) } = process.Adapter{} var containerProvider container.Provider = container.Adapter{} var arrayProvider array.Provider = array.Adapter{} var diskProvider disk.Provider = disk.Adapter{} var poolProvider pooldomain.Provider = pooldomain.Adapter{} var shareProvider sharedomain.Provider = sharedomain.Adapter{} if agentReader != nil { hostProvider = agentsource.HostProvider{Reader: agentReader, Windows: agentWindows} processProvider = agentsource.ProcessProvider{Reader: agentReader, Windows: agentWindows} containerProvider = agentsource.ContainerProvider{Reader: agentReader, Windows: agentWindows} arrayProvider = agentsource.ArrayProvider{Reader: agentReader, Windows: agentWindows} diskProvider = agentsource.DiskProvider{Reader: agentReader, Windows: agentWindows} poolProvider = agentsource.PoolProvider{Reader: agentReader, Windows: agentWindows} shareProvider = agentsource.ShareProvider{Reader: agentReader, Windows: agentWindows} } // Applications have no capability of their own: they aggregate the container // inventory with the service probe results, and report Unknown when either input is // missing or stale. applicationProvider := agentsource.ApplicationProvider{Containers: containerProvider, Services: serviceProvider} hostHandler := hostapi.Handler{Provider: hostProvider} mux.Handle("/api/v1/host", withSession(sessions, auth.Require(auth.PermissionView, hostHandler))) processHandler := processapi.Handler{Provider: processProvider} mux.Handle("/api/v1/processes", withSession(sessions, auth.Require(auth.PermissionView, processHandler))) containerHandler := containerapi.Handler{Provider: containerProvider} mux.Handle("/api/v1/containers", withSession(sessions, auth.Require(auth.PermissionView, containerHandler))) mux.Handle("/api/v1/containers/", withSession(sessions, auth.Require(auth.PermissionView, containerHandler))) applicationHandler := applicationapi.Handler{Provider: applicationProvider} mux.Handle("/api/v1/applications", withSession(sessions, auth.Require(auth.PermissionView, applicationHandler))) mux.Handle("/api/v1/applications/", withSession(sessions, auth.Require(auth.PermissionView, applicationHandler))) arrayHandler := arrayapi.Handler{Provider: arrayProvider} mux.Handle("/api/v1/array", withSession(sessions, auth.Require(auth.PermissionView, arrayHandler))) diskHandler := diskapi.Handler{Provider: diskProvider} mux.Handle("/api/v1/disks", withSession(sessions, auth.Require(auth.PermissionView, diskHandler))) mux.Handle("/api/v1/disks/", withSession(sessions, auth.Require(auth.PermissionView, diskHandler))) poolHandler := poolapi.Handler{Provider: poolProvider} mux.Handle("/api/v1/pools", withSession(sessions, auth.Require(auth.PermissionView, poolHandler))) mux.Handle("/api/v1/pools/", withSession(sessions, auth.Require(auth.PermissionView, poolHandler))) shareHandler := shareapi.Handler{Provider: shareProvider} mux.Handle("/api/v1/shares", withSession(sessions, auth.Require(auth.PermissionView, shareHandler))) mux.Handle("/api/v1/shares/", withSession(sessions, auth.Require(auth.PermissionView, shareHandler))) forecastHandler := forecastapi.Handler{Provider: forecastdomain.StorageProvider{Shares: shareProvider, Pools: poolProvider, History: forecastdomain.PostgresHistory{Pool: pool}, Policy: forecastdomain.Policy{Enabled: true}}} mux.Handle("/api/v1/forecasts", withSession(sessions, auth.Require(auth.PermissionView, forecastHandler))) serviceHandler := serviceapi.Handler{Provider: serviceProvider, Dependencies: dependencyRepo, ReverseProxy: reverseProxyProvider} networkHandler := networkapi.Handler{Provider: network.Aggregator{Host: hostProvider, Services: serviceProvider}} mux.Handle("/api/v1/services", withSession(sessions, auth.Require(auth.PermissionView, serviceHandler))) mux.Handle("/api/v1/services/", withSession(sessions, auth.Require(auth.PermissionView, serviceHandler))) mux.Handle("/api/v1/topology", withSession(sessions, auth.Require(auth.PermissionView, serviceHandler))) mux.Handle("/api/v1/network", withSession(sessions, auth.Require(auth.PermissionView, networkHandler))) reverseProxyHandler := reverseproxyapi.Handler{Provider: reverseProxyProvider} mux.Handle("/api/v1/reverse-proxy", withSession(sessions, auth.Require(auth.PermissionView, reverseProxyHandler))) server := &http.Server{Addr: runtime.ListenAddress, Handler: observability.Middleware(internalMetrics, correlation.Middleware(mux)), ReadHeaderTimeout: 5 * time.Second} go func() { logger.Info("pulse api listening", "addr", runtime.ListenAddress, "environment", application.Environment) if err := server.ListenAndServe(); err != nil && err != http.ErrServerClosed { logger.Error("pulse api stopped", "error", err) } }() stopContext, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer stop() go alertcontrol.RunExpiryLoop(stopContext, alertControlStore, time.Minute, logger) service.WaitForStop(stopContext, nil) shutdownContext, cancel := context.WithTimeout(context.Background(), runtime.ShutdownAfter) defer cancel() if err := server.Shutdown(shutdownContext); err != nil { return err } logger.Info("pulse api stopped") return nil } func boolMetric(value bool) float64 { if value { return 1 } return 0 } // roleMapping converts the validated configuration mapping of identity provider // group claims onto the typed roles used by the authorization layer. Configuration // already rejects unknown role names, so no further validation is needed here. func roleMapping(configured map[string]string) map[string]auth.Role { if len(configured) == 0 { return nil } mapping := make(map[string]auth.Role, len(configured)) for claim, role := range configured { mapping[claim] = auth.Role(role) } return mapping } func withSession(manager *auth.SessionManager, next http.Handler) http.Handler { return http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { authentication, ok := manager.AuthenticateSession(response, request, time.Now().UTC()) if !ok { next.ServeHTTP(response, request) return } ctx, cancel := context.WithCancel(request.Context()) stop := context.AfterFunc(authentication.Context, cancel) defer func() { stop() cancel() }() next.ServeHTTP(response, request.WithContext(auth.WithPrincipal(ctx, authentication.Principal))) }) }