api

package
v0.6.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 21, 2026 License: MIT Imports: 36 Imported by: 0

Documentation

Index

Constants

View Source
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.

View Source
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"
)
View Source
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.

View Source
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

View Source
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".

View Source
var OtelContextStartTime = time.Now()

Functions

func APIKeyGate

func APIKeyGate(expectedKey, mcpPath string, next http.Handler) http.Handler

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

func DBHealthMiddleware(h *DBHealth) func(http.Handler) http.Handler

DBHealthMiddleware returns 503 immediately when the DB poller reports unhealthy, for DB-dependent paths. Health/metrics/UI paths bypass the gate.

func GzipMiddleware

func GzipMiddleware(mcpPath string) func(http.Handler) http.Handler

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

func IsProtectedPath(path, mcpPath string) bool

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

func IsWebSocketPath(path string) bool

IsWebSocketPath reports whether a path belongs to the /ws* namespace.

func MetricsMiddleware

func MetricsMiddleware(metrics *telemetry.Metrics, next http.Handler) http.Handler

MetricsMiddleware records OtelContext_http_requests_total and OtelContext_http_request_duration_seconds for every HTTP request.

func RecoverMiddleware

func RecoverMiddleware(metrics *telemetry.Metrics, next http.Handler) http.Handler

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

func RequireAPIKey(expectedKey string, next http.Handler) http.Handler

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

func TenantMiddleware(cfg *config.Config) func(http.Handler) http.Handler

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

func WSAllowedOriginHosts(allowed []string) []string

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

func NewDBHealth(db DBPinger, driver string, metrics *telemetry.Metrics) *DBHealth

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) Healthy

func (h *DBHealth) Healthy() bool

Healthy reports the most recent ping result.

func (*DBHealth) SetFailureThreshold

func (h *DBHealth) SetFailureThreshold(n int)

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

func (h *DBHealth) Shutdown(ctx context.Context) error

Shutdown signals the poller to exit and joins it within the application shutdown deadline.

func (*DBHealth) Start

func (h *DBHealth) Start(ctx context.Context)

Start launches the background poller.

func (*DBHealth) Stop

func (h *DBHealth) Stop()

Stop preserves the previous two-second bounded lifecycle surface.

type DBPinger

type DBPinger interface {
	PingContext(ctx context.Context) error
}

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

func (s *Server) BroadcastLog(l storage.Log)

BroadcastLog sends a log entry to the buffered WebSocket hub.

func (*Server) RegisterRoutes

func (s *Server) RegisterRoutes(mux *http.ServeMux)

RegisterRoutes registers API endpoints on the provided mux.

func (*Server) SetAggregateDBProbe added in v0.5.0

func (s *Server) SetAggregateDBProbe(fn func(context.Context) error)

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

func (s *Server) SetAggregateEngine(e *aggregate.Engine)

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

func (s *Server) SetAggregateRecoveryProbe(fn func() bool)

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

func (s *Server) SetDLQSaturationProbe(fn func() float64)

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

func (s *Server) SetDiskPressureProbe(fn func() (string, bool))

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) SetGraph

func (s *Server) SetGraph(g *graph.Graph)

SetGraph wires the in-memory service graph into the API server.

func (*Server) SetGraphRAG

func (s *Server) SetGraphRAG(g *graphrag.GraphRAG)

SetGraphRAG wires the GraphRAG instance for advanced queries.

func (*Server) SetPipelineSaturationProbe

func (s *Server) SetPipelineSaturationProbe(fn func() float64)

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

func (s *Server) SetTopologyProvider(provider topology.Provider)

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.

Directories

Path Synopsis
Package views provides explicit JSON view models for the HTTP API.
Package views provides explicit JSON view models for the HTTP API.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL