Documentation
¶
Index ¶
- Constants
- Variables
- func APIKeyGate(expectedKey, mcpPath string, next http.Handler) http.Handler
- func AuthGate(o AuthGateOptions, next http.Handler) http.Handler
- func DBHealthMiddleware(h *DBHealth) func(http.Handler) http.Handler
- func GzipMiddleware(mcpPath string) func(http.Handler) http.Handler
- func IsProtectedPath(path, mcpPath string) bool
- func IsWebSocketPath(path string) bool
- func MetricsMiddleware(metrics *telemetry.Metrics, next http.Handler) http.Handler
- func RecoverMiddleware(metrics *telemetry.Metrics, next http.Handler) http.Handler
- func RequireAPIKey(expectedKey string, next http.Handler) http.Handler
- func TenantMiddleware(cfg *config.Config) func(http.Handler) http.Handler
- func WSAllowedOriginHosts(allowed []string) []string
- func WebSocketGate(o WSGateOptions, next http.Handler) http.Handler
- type AggregateRuntime
- type AuthGateOptions
- type DBHealth
- type DBPinger
- type GraphEdge
- type GraphNode
- type NodeMetrics
- type RateLimiter
- type ReadinessThresholds
- type Server
- func (s *Server) BeginShutdown()
- func (s *Server) BroadcastLog(l storage.Log)
- func (s *Server) RegisterRoutes(mux *http.ServeMux)
- func (s *Server) SetAggregateDBProbe(fn func(context.Context) error)
- func (s *Server) SetAggregateEngine(e *aggregate.Engine)
- func (s *Server) SetAggregateRecoveryProbe(fn func() bool)
- func (s *Server) SetAggregateRuntimeProbe(fn func() AggregateRuntime)
- func (s *Server) SetDLQSaturationProbe(fn func() float64)
- func (s *Server) SetDiskPressureProbe(fn func() (string, bool))
- func (s *Server) SetGraph(g *graph.Graph)
- func (s *Server) SetGraphRAG(g *graphrag.GraphRAG)
- func (s *Server) SetPipelineSaturationProbe(fn func() float64)
- func (s *Server) SetReadinessThresholds(t ReadinessThresholds)
- func (s *Server) SetTopologyProvider(provider topology.Provider)
- type SystemGraphResponse
- type SystemSummary
- type WSGateOptions
Constants ¶
const ( // RequestedStartHeader restates the start the client requested. RequestedStartHeader = "OtelContext-Requested-Start" // EffectiveStartHeader carries the clamped start the response covers. EffectiveStartHeader = "OtelContext-Effective-Start" )
Range-clamp headers (#217). When an aggregate-mode request asks for more history than the engine's read-range cap allows, the handler clamps the range instead of refusing the request, and these headers say so: what the client asked for, and what the response actually covers.
const ( // WSSubprotocol is the only subprotocol the server ever echoes. Defined in // authn so the realtime hubs can negotiate it without importing this // package (api already imports realtime). WSSubprotocol = authn.WSSubprotocol // WSTenantParam is the query parameter an operator credential uses to // select the single tenant a socket is scoped to. WSTenantParam = "tenant" )
const DefaultExternalTenantHeader = "X-OtelContext-Tenant"
DefaultExternalTenantHeader is the dedicated identity header a front proxy injects when AUTH_TRUST_EXTERNAL=true. It is deliberately NOT X-Tenant-ID: that header stays client-controlled and untrusted, so a proxy that forgets to strip inbound copies of the dedicated header fails visibly rather than silently promoting client input to an identity.
const TenantHeader = "X-Tenant-ID"
TenantHeader is the canonical HTTP header carrying the tenant ID on read-side (query) requests. Ingest paths resolve tenant separately via gRPC metadata / OTLP resource attributes and do not go through this middleware.
Variables ¶
var AuthFailureHook func(reason string)
AuthFailureHook is an optional callback invoked whenever API-key auth rejects a request. Set by main.go to increment the APIAuthFailuresTotal metric. Left as a package-level function pointer (rather than a DI parameter) to avoid circular imports between api and telemetry. Safe to leave nil.
Reasons: "missing_header", "bad_scheme", "bad_key".
var OtelContextStartTime = time.Now()
Functions ¶
func APIKeyGate ¶
APIKeyGate wraps a handler so only requests matching IsProtectedPath require the key. Public paths flow through untouched, which is what keeps the UI bundle and health probes accessible without credentials.
func AuthGate ¶ added in v0.5.0
func AuthGate(o AuthGateOptions, next http.Handler) http.Handler
AuthGate authenticates /api/*, /v1/*, and the MCP endpoint.
Two credential classes, one gate:
- The operator key (API_KEY) authenticates and nothing more. Tenant resolution keeps its historical precedence — X-Tenant-ID header, then the OTLP resource attribute where trusted, then DEFAULT_TENANT — so an API_KEY-only deployment behaves exactly as it did before this middleware existed.
- A tenant key (API_TENANT_KEYS_FILE) authenticates AND binds. The bound tenant is pinned onto the request context and a client-supplied X-Tenant-ID is ignored and counted, never honoured.
When AUTH_TRUST_EXTERNAL=true and no bearer credential is presented, the proxy-injected identity header supplies a bound principal.
func DBHealthMiddleware ¶
DBHealthMiddleware returns 503 immediately when the DB poller reports unhealthy, for DB-dependent paths. Health/metrics/UI paths bypass the gate.
func GzipMiddleware ¶
GzipMiddleware compresses GET /api/* responses for clients that accept gzip. The 120-service system-graph JSON shrinks 5-8× on the wire.
Everything outside /api/* passes through untouched, which keeps the streaming surfaces safe by construction: WebSocket upgrades (/ws*) still see an http.Hijacker, OTLP ingest (/v1/*) and Prometheus scrapes (/metrics*) are write paths/exposition formats with their own encodings, and MCP SSE must flush uncompressed frames. Those prefixes — plus the configured MCP path — are also excluded explicitly in case an operator ever nests one under /api/.
Wire this innermost (directly around the mux): only handler output is compressed; error responses written by outer middleware (auth, rate limit) stay identity-encoded.
func IsProtectedPath ¶
IsProtectedPath reports whether a request path requires API-key authentication. Protected: /api/*, /v1/* (OTLP HTTP), and the MCP path. Unprotected: /live, /ready, /health*, /metrics* (Prometheus), /ws* (WebSocket), and the UI static bundle ("/" + assets).
func IsWebSocketPath ¶ added in v0.5.0
IsWebSocketPath reports whether a path belongs to the /ws* namespace.
func MetricsMiddleware ¶
MetricsMiddleware records OtelContext_http_requests_total and OtelContext_http_request_duration_seconds for every HTTP request.
func RecoverMiddleware ¶
RecoverMiddleware catches panics from downstream handlers/middleware, logs the stack trace, increments the panics-recovered metric, and responds with a generic 500. It must be installed as the OUTERMOST middleware (after MetricsMiddleware is wrapped around the stack) so panics anywhere below are caught. http.ErrAbortHandler is re-panicked to preserve net/http's sentinel-abort contract.
func RequireAPIKey ¶
RequireAPIKey returns middleware that requires an `Authorization: Bearer <key>` header matching the configured API key. When expectedKey is empty the middleware is a pass-through (auth disabled) — the caller is expected to log a warning at startup in that case.
The comparison is constant-time via subtle.ConstantTimeCompare to avoid timing side channels. On mismatch or missing header a 401 is returned with a JSON body `{"error":"unauthorized"}`.
func TenantMiddleware ¶
TenantMiddleware extracts the tenant ID from the X-Tenant-ID header (falling back to cfg.DefaultTenant — or "default" when cfg is nil or empty) and stashes it on the request context via storage.WithTenantContext so repository reads can scope their WHERE clause with storage.TenantFromContext.
The middleware is path-aware: only requests whose path begins with "/api/" are tenant-scoped. OTLP write endpoints ("/v1/..."), health probes ("/live", "/ready"), Prometheus scrape ("/metrics/..."), MCP, WebSocket and UI assets pass through untouched — these either resolve tenant separately (OTLP) or are tenant-agnostic/privileged.
func WSAllowedOriginHosts ¶ added in v0.5.0
WSAllowedOriginHosts reduces configured origins to bare hosts, which is the shape the WebSocket library's origin patterns expect.
func WebSocketGate ¶ added in v0.5.0
func WebSocketGate(o WSGateOptions, next http.Handler) http.Handler
WebSocketGate authenticates the /ws* handshake and scopes the connection to exactly one tenant before the upgrade happens.
Credential carriers, in order: `Authorization: Bearer <token>` (non-browser clients), then a `Sec-WebSocket-Protocol` entry of the form `auth.<base64url-token>` (browsers, which cannot set headers on a WebSocket). Query-string tokens are never accepted — they land in access logs, browser history, and Referer headers.
Scope: a tenant-key or proxy-injected identity binds the socket to its tenant and any client-selected tenant is ignored and counted. An operator credential selects exactly one tenant via ?tenant= or X-Tenant-ID, falling back to DEFAULT_TENANT. There is no merged all-tenant stream.
Types ¶
type AggregateRuntime ¶ added in v0.5.0
type AggregateRuntime struct {
// CommitFailureStreak and FinalizeFailureStreak are CONSECUTIVE failure
// counts: any success resets them.
CommitFailureStreak uint64
FinalizeFailureStreak uint64
// AdmissionRatio is the group-commit writer's admission occupancy as a
// fraction of its bounds — the fullest of pending bytes, pending deltas
// and parked waiters.
AdmissionRatio float64
// DeltaLogAgeSeconds is the age of the oldest un-finalized window,
// including the staleness of the sample it came from.
DeltaLogAgeSeconds float64
// DiskUsedBytes and DiskBudgetBytes are the aggregate tier's on-disk size
// against its share of the data budget. A zero budget disables the check.
DiskUsedBytes int64
DiskBudgetBytes int64
}
AggregateRuntime is the aggregate runtime health snapshot /ready consults. It is a plain value, sampled from counters the writer already maintains, so a readiness request never queries the store and never stacks behind the single SQLite writer.
func (AggregateRuntime) DiskRatio ¶ added in v0.5.0
func (a AggregateRuntime) DiskRatio() float64
DiskRatio is the aggregate tier's usage as a fraction of its budget.
type AuthGateOptions ¶ added in v0.5.0
type AuthGateOptions struct {
// Auth resolves bearer tokens. A disabled Authenticator makes the gate a
// pass-through — authentication arrives with configuration, not by default.
Auth *authn.Authenticator
// MCPPath is used by IsProtectedPath to gate the MCP endpoint.
MCPPath string
// ExternalTenantHeader overrides DefaultExternalTenantHeader.
ExternalTenantHeader string
}
AuthGateOptions configures the HTTP authentication middleware.
type DBHealth ¶
type DBHealth struct {
// contains filtered or unexported fields
}
DBHealth periodically pings the database and exposes the result via an atomic boolean. The HTTP middleware below short-circuits /api/* traffic with a 503 when the flag is false, preventing goroutine pile-up on pool acquisition when the DB is unreachable.
To avoid false-positive 503 windows under load — e.g. SQLite with MaxOpen=1 where a busy write transiently blocks the 2s health ping — the poller flips to unhealthy only after failureThreshold consecutive failed pings. A single successful ping clears the counter and restores healthy immediately.
func NewDBHealth ¶
NewDBHealth constructs a health poller. Default poll interval is 5s; ping timeout is 2s per attempt; failureThreshold defaults to 3 consecutive failed pings before the gate flips. Start() must be called to begin polling.
func (*DBHealth) SetFailureThreshold ¶
SetFailureThreshold overrides the number of consecutive failed pings before the middleware flips to 503. n <= 0 normalises to 1 (legacy behaviour: any single failure trips the gate).
func (*DBHealth) Shutdown ¶ added in v0.5.0
Shutdown signals the poller to exit and joins it within the application shutdown deadline.
type DBPinger ¶
DBPinger is the minimum DB interface DBHealth needs. *sql.DB satisfies it; tests can pass a stub.
type GraphEdge ¶
type GraphEdge struct {
Source string `json:"source"`
Target string `json:"target"`
CallCount int64 `json:"call_count"`
AvgLatencyMs float64 `json:"avg_latency_ms"`
ErrorRate float64 `json:"error_rate"`
Status string `json:"status"`
}
GraphEdge represents a call relationship between two services.
type GraphNode ¶
type GraphNode struct {
ID string `json:"id"`
Type string `json:"type"`
HealthScore float64 `json:"health_score"`
Status string `json:"status"`
Metrics NodeMetrics `json:"metrics"`
Alerts []string `json:"alerts"`
// Additive host projection (#288). Kind is service|host; a host node is
// excluded from the service summary.
Kind string `json:"kind,omitempty"`
HostCount int `json:"host_count,omitempty"`
Hosts []string `json:"hosts,omitempty"`
}
GraphNode represents a service in the system graph.
type NodeMetrics ¶
type NodeMetrics struct {
RequestRateRPS float64 `json:"request_rate_rps"`
ErrorRate float64 `json:"error_rate"`
AvgLatencyMs float64 `json:"avg_latency_ms"`
P99LatencyMs float64 `json:"p99_latency_ms"`
LatencyProvenance *latency.Provenance `json:"latency_provenance,omitempty"`
SpanCount1H int64 `json:"span_count_1h"`
}
NodeMetrics holds per-service observability metrics.
type RateLimiter ¶
type RateLimiter struct {
// contains filtered or unexported fields
}
RateLimiter is a per-IP token bucket rate limiter middleware.
func NewRateLimiter ¶
func NewRateLimiter(rps float64) *RateLimiter
NewRateLimiter creates a RateLimiter with the given requests-per-second limit.
func (*RateLimiter) Middleware ¶
func (rl *RateLimiter) Middleware(next http.Handler) http.Handler
Middleware returns an http.Handler that enforces rate limiting.
func (*RateLimiter) MiddlewareExcept ¶
func (rl *RateLimiter) MiddlewareExcept(skip func(path string) bool) func(http.Handler) http.Handler
MiddlewareExcept returns a handler chain that applies the rate limit only when skip(path) returns false. Intended to exempt OTLP ingestion paths from the per-IP API limiter.
type ReadinessThresholds ¶ added in v0.5.0
type ReadinessThresholds struct {
MaxCommitFailureStreak uint64
MaxFinalizeFailureStreak uint64
MaxAdmissionRatio float64
MaxDeltaLogAgeSeconds float64
MaxAggregateDiskRatio float64
}
ReadinessThresholds are the limits the aggregate runtime probes compare against. A non-positive limit disables that probe: an operator who disagrees with a default switches it off rather than patching the binary.
func DefaultReadinessThresholds ¶ added in v0.5.0
func DefaultReadinessThresholds() ReadinessThresholds
DefaultReadinessThresholds mirrors the config defaults so a Server built without explicit thresholds still probes rather than silently passing.
type Server ¶
type Server struct {
// contains filtered or unexported fields
}
Server handles HTTP API requests.
func NewServer ¶
func NewServer(repo *storage.Repository, hub *realtime.Hub, eventHub *realtime.EventHub, metrics *telemetry.Metrics) *Server
NewServer creates a new API server.
func (*Server) BeginShutdown ¶ added in v0.5.0
func (s *Server) BeginShutdown()
BeginShutdown removes this instance from readiness before admission stops. It is idempotent and leaves /live unchanged while the process drains.
func (*Server) BroadcastLog ¶
BroadcastLog sends a log entry to the buffered WebSocket hub.
func (*Server) RegisterRoutes ¶
RegisterRoutes registers API endpoints on the provided mux.
func (*Server) SetAggregateDBProbe ¶ added in v0.5.0
SetAggregateDBProbe registers the aggregate store reachability check. The callback must honour the context deadline the probe passes it.
func (*Server) SetAggregateEngine ¶ added in v0.5.0
SetAggregateEngine wires the aggregate query facade into the read path. Pass the engine only in AGGREGATE_MODE=aggregate; a nil engine, or an engine in any other mode, leaves every handler on the legacy path.
func (*Server) SetAggregateRecoveryProbe ¶ added in v0.5.0
SetAggregateRecoveryProbe registers a callback reporting whether the durable aggregate store has finished replaying its delta log. /ready returns 503 until it does. Pass nil (the default) when no aggregate store is configured.
func (*Server) SetAggregateRuntimeProbe ¶ added in v0.5.0
func (s *Server) SetAggregateRuntimeProbe(fn func() AggregateRuntime)
SetAggregateRuntimeProbe registers the aggregate runtime health sampler. Pass nil (the default) when no aggregate store is configured; every runtime check then reports "skipped" and readiness is unaffected.
func (*Server) SetDLQSaturationProbe ¶
SetDLQSaturationProbe registers a callback returning DLQ disk fullness as a fraction in [0.0, 1.0]. Used by /ready to flip to 503 when DLQ is at risk of FIFO-evicting unflushed batches. Pass nil to disable the check.
func (*Server) SetDiskPressureProbe ¶ added in v0.5.0
SetDiskPressureProbe registers a callback returning the disk watchdog's state label and whether readiness should pass. Pass nil (the default) when no watchdog is configured.
func (*Server) SetGraphRAG ¶
SetGraphRAG wires the GraphRAG instance for advanced queries.
func (*Server) SetPipelineSaturationProbe ¶
SetPipelineSaturationProbe registers a callback returning ingest pipeline queue fullness as a fraction in [0.0, 1.0]. Used by /ready to flip to 503 when the pipeline is at hard capacity (already returning 429/RESOURCE_EXHAUSTED to clients). Pass nil to disable the check.
func (*Server) SetReadinessThresholds ¶ added in v0.5.0
func (s *Server) SetReadinessThresholds(t ReadinessThresholds)
SetReadinessThresholds overrides the runtime-probe limits. A zero value in any field disables that probe.
func (*Server) SetTopologyProvider ¶ added in v0.5.0
SetTopologyProvider installs the construction-time mode owner used by both REST topology endpoints.
type SystemGraphResponse ¶
type SystemGraphResponse struct {
Timestamp time.Time `json:"timestamp"`
System SystemSummary `json:"system"`
Nodes []GraphNode `json:"nodes"`
Edges []GraphEdge `json:"edges"`
Source string `json:"source,omitempty"`
Coverage string `json:"coverage,omitempty"`
CoverageNote string `json:"coverage_note,omitempty"`
Epoch string `json:"epoch,omitempty"`
Revision uint64 `json:"revision,omitempty"`
Truncated bool `json:"truncated,omitempty"`
DroppedServices uint64 `json:"dropped_services,omitempty"`
DroppedOperations uint64 `json:"dropped_operations,omitempty"`
DroppedEdges uint64 `json:"dropped_edges,omitempty"`
DroppedMetrics uint64 `json:"dropped_metrics,omitempty"`
}
SystemGraphResponse is the full AI-consumable system graph.
type SystemSummary ¶
type SystemSummary struct {
TotalServices int `json:"total_services"`
Healthy int `json:"healthy"`
Degraded int `json:"degraded"`
Critical int `json:"critical"`
OverallHealthScore float64 `json:"overall_health_score"`
TotalErrorRate float64 `json:"total_error_rate"`
AvgLatencyMs float64 `json:"avg_latency_ms"`
UptimeSeconds float64 `json:"uptime_seconds"`
}
SystemSummary is the top-level system health summary.
type WSGateOptions ¶ added in v0.5.0
type WSGateOptions struct {
// Auth resolves credentials. Disabled Authenticator = pass-through, which
// keeps /ws* open in a default development deployment exactly as before.
Auth *authn.Authenticator
// DefaultTenant scopes an operator socket that selects no tenant.
DefaultTenant string
// AllowedOrigins is WS_ALLOWED_ORIGINS: exact origins ("https://app.example.com")
// or bare hosts ("app.example.com"). Empty means same-host only.
AllowedOrigins []string
// EnforceOrigin turns the origin policy on. Set when authentication is
// enabled or APP_ENV=production.
EnforceOrigin bool
// ExternalTenantHeader overrides DefaultExternalTenantHeader.
ExternalTenantHeader string
}
WSGateOptions configures the WebSocket handshake gate.