diff --git a/appview/ALERTS.md b/appview/ALERTS.md new file mode 100644 index 0000000..a5d7296 --- /dev/null +++ b/appview/ALERTS.md @@ -0,0 +1,224 @@ +# Effem AppView — Alerting + +Phase 1 Step 6. Pairs with `/_health` (see `health.go`) and the Prometheus +registry in `metrics/`. + +The goal is a **single, clean alert channel** with narrow, high-signal rules. +Over-alerting trains people to ignore pages; start with the minimal set below +and only add rules when you have evidence an incident class is not covered. + +## Signals we emit + +| Surface | Contract | Scrape/poll interval | +|---|---|---| +| `GET /_health` | HTTP 200 with `status: "ok"` when healthy; HTTP 503 with `status: "degraded"` or `"starting"` otherwise. Fields: `db`, `firehose_up`, `firehose_seq`, `firehose_age_s`. | External uptime monitor, every 30 s. | +| `GET /metrics` | Prometheus exposition format, behind admin scope. | Prometheus scraper with admin credentials, every 15 s. | + +`/_health` is public (no auth). `/metrics` is admin-only — provision a +dedicated admin token whose only job is scraping. + +## Tier 1 — uptime monitor (required for beta) + +Pick one external service that is **not** Railway itself. The point is to +detect the case where Railway (and therefore any self-hosted monitor) is +down. Options: + +- Better Stack (Better Uptime) +- UptimeRobot +- Checkly +- Grafana Cloud Synthetics + +Whichever you pick, configure: + +1. **Target**: `GET https:///_health` +2. **Interval**: 30 s +3. **Assertions**: + - HTTP status is `200` + - Response JSON has `status == "ok"` +4. **Alert after**: 2 consecutive failures (avoids paging on a single network + blip; keeps MTTD ≤ 1 min). +5. **Notification**: single channel (Slack/Discord/email). Keep the on-call + rotation short; one-person teams use their personal device. + +### Runbook link + +Every alert notification must include a link to `appview/ALERTS.md#on-call-runbook` +so whoever is paged has the playbook in one click. + +## Tier 2 — Prometheus rules (recommended for beta, required post-beta) + +Where Prometheus runs: + +- **Simplest**: Grafana Cloud free tier. Adds a `prometheus` remote-write or a + pull scrape from their hosted Prometheus. +- **Self-hosted**: a Prometheus container on Railway with persistent volume. + Cheaper but one more thing to operate. + +### Scrape config + +```yaml +scrape_configs: + - job_name: effem-appview + scrape_interval: 15s + metrics_path: /metrics + bearer_token: ${EFFEM_METRICS_TOKEN} + scheme: https + static_configs: + - targets: ['effem-appview.up.railway.app'] +``` + +Use a dedicated admin token for this scraper. Rotate it on the same cadence +as every other admin token (see Step 13 in `prod47.md`). + +### Alert rules + +```yaml +groups: + - name: effem-appview-beta + interval: 30s + rules: + - alert: EffemAppViewDown + expr: up{job="effem-appview"} == 0 + for: 2m + labels: { severity: page } + annotations: + summary: "Effem AppView unreachable for 2 minutes" + runbook: "https://github.com/SparrowTek/AtProto/blob/main/effem/AppView/appview/ALERTS.md#appview-down" + + - alert: EffemFirehoseDisconnected + expr: effem_firehose_connected == 0 + for: 2m + labels: { severity: page } + annotations: + summary: "Firehose consumer disconnected for 2 minutes" + runbook: "https://github.com/SparrowTek/AtProto/blob/main/effem/AppView/appview/ALERTS.md#firehose-disconnected" + + - alert: EffemFirehoseLag + expr: effem_firehose_lag_seconds > 60 + for: 5m + labels: { severity: page } + annotations: + summary: "Firehose has seen no events for 60+ seconds (stall or slow upstream)" + runbook: "https://github.com/SparrowTek/AtProto/blob/main/effem/AppView/appview/ALERTS.md#firehose-lag" + + - alert: EffemIndexerErrorSpike + expr: rate(effem_indexer_errors_total[5m]) > 0.1 + for: 10m + labels: { severity: warn } + annotations: + summary: "Indexer errors > 0.1/s sustained for 10 min" + description: "Collection breakdown: {{ $labels.collection }}/{{ $labels.action }}" + + - alert: EffemHTTP5xxRate + expr: | + sum(rate(effem_http_requests_total{status=~"5.."}[5m])) + / + sum(rate(effem_http_requests_total[5m])) + > 0.05 + for: 10m + labels: { severity: page } + annotations: + summary: "5xx error rate above 5% for 10 minutes" + + - alert: EffemRateLimitSaturation + expr: rate(effem_rate_limit_rejections_total[5m]) > 5 + for: 10m + labels: { severity: warn } + annotations: + summary: "Rate limiter rejecting > 5 req/s (possible abuse or legitimate surge)" + + - alert: EffemDBUnreachable + # From the /_health response converted into a gauge via blackbox_exporter + # or a dedicated probe. If Prometheus is configured as described above, + # this overlaps with EffemAppViewDown — keep only one. + expr: probe_http_status_code{instance=~".*/_health"} == 503 + for: 2m + labels: { severity: page } + annotations: + summary: "/_health reports degraded (DB or firehose failing)" +``` + +### Label discipline + +- `severity: page` → wakes someone up. Reserve for true outages. +- `severity: warn` → shows up in the on-call channel but does not page. + +## Notification routing + +A single Alertmanager config or Grafana alert rule set is enough at beta +scale. Route everything to one channel. Split channels only when volume +forces it. + +- Channel: Slack `#effem-alerts` (or Discord equivalent). +- Keep message format simple: alert name, severity, summary, runbook link. +- No @everyone. A clear name and a runbook link is enough. + +## On-call runbook + +Each alert must have a linked section below. If an alert fires and the +runbook is missing or wrong, fix it during the same incident — fresh pain is +the best time to write. + +### AppView down + +1. Confirm from a second network: `curl -sS -o /dev/null -w "%{http_code}\n" https:///_health`. +2. Check Railway dashboard: did the container crash? Recent deploys? +3. If the last deploy is the cause: roll back via Railway (pre-deploy + snapshot + previous image; see `database/BACKUPS.md §4`). +4. If DB is the cause (`db: false` in `/_health`): inspect Railway Postgres + status, swap to a restored database if necessary. +5. If the relay is the cause (firehose side): see **firehose-disconnected**. + +### Firehose disconnected + +1. Check the AppView logs for the last firehose log line: look for + "firehose disconnected, reconnecting" (expected transient) vs. a stack trace. +2. If the relay (bsky.network) is down, there is nothing to do but wait; + log the incident and continue. +3. If the AppView's WebSocket consistently fails to connect, verify + outbound WebSocket traffic is allowed from Railway. +4. If the cursor is too old (relay rejects the subscribe with a + `FutureCursor` or similar error): see + `firehose/RUNBOOK.md` (created in Phase 2 Step 9). + +### Firehose lag + +1. Check `effem_firehose_connected` — if 0, treat as **firehose-disconnected** instead. +2. Check relay health independently. `bsky.network` outages are real and + usually brief. +3. If connection is healthy but events are slow, check AppView CPU and + network saturation in Railway. +4. If the AppView is falling behind a busy firehose, consider bumping + `EFFEM_FIREHOSE_PARALLELISM` and redeploying. + +### Indexer error spike + +1. Check which `{collection, action}` is spiking in the alert payload. +2. Tail the AppView logs for `failed to index record` warnings — they + include the underlying error and repo path. +3. If a schema migration just shipped, rollback is the fastest + intervention; follow `database/BACKUPS.md §4`. +4. If a single repo is sending malformed records, the warning rate-limits + itself via the log sampler — no action beyond filing a Bluesky support + ticket. + +### HTTP 5xx rate + +1. Pull the top paths from `effem_http_requests_total{status=~"5.."}`. +2. Check DB health and PI upstream health — most 5xx on this AppView + originate from one of those two. +3. If PI is down, consider temporarily returning cached-but-stale data + instead of errors (tracked separately in the cache layer). + +## Maintenance + +- Review the rule set quarterly. Remove rules that have not fired in a + quarter and are not gating anything critical. +- Every alert that fires gets a note in this file's **Incident log** + section (below) if root cause is interesting. + +## Incident log + +| Date | Alert | Summary | Follow-up | +|---|---|---|---| +| _(first incident)_ | | | | diff --git a/appview/firehose.go b/appview/firehose.go index 411c815..9004511 100644 --- a/appview/firehose.go +++ b/appview/firehose.go @@ -19,6 +19,7 @@ import ( "github.com/bluesky-social/indigo/repomgr" "github.com/gorilla/websocket" "tangled.org/sparrowtek.com/effem-AppView/appview/database" + "tangled.org/sparrowtek.com/effem-AppView/appview/metrics" ) const ( @@ -60,10 +61,13 @@ func (srv *Server) RunFirehoseConsumer(ctx context.Context) error { rsc := &events.RepoStreamCallbacks{ RepoCommit: func(evt *comatproto.SyncSubscribeRepos_Commit) error { atomic.StoreInt64(&srv.lastSeq, evt.Seq) + srv.lastSeqTime.Store(time.Now()) return srv.handleCommit(ctx, evt) }, RepoIdentity: func(evt *comatproto.SyncSubscribeRepos_Identity) error { atomic.StoreInt64(&srv.lastSeq, evt.Seq) + srv.lastSeqTime.Store(time.Now()) + metrics.IncFirehoseEvent("identity", "") srv.logger.Debug("identity event", "did", evt.Did) return nil }, @@ -91,7 +95,11 @@ func (srv *Server) RunFirehoseConsumer(ctx context.Context) error { go srv.persistCursorLoop(ctx, cancel) srv.firehoseUp.Store(true) - defer srv.firehoseUp.Store(false) + metrics.SetFirehoseConnected(true) + defer func() { + srv.firehoseUp.Store(false) + metrics.SetFirehoseConnected(false) + }() srv.logger.Info("firehose consumer running") return events.HandleRepoStream(ctx, con, scheduler, srv.logger) @@ -125,6 +133,7 @@ func (srv *Server) handleCommit(ctx context.Context, evt *comatproto.SyncSubscri } srv.logger.Info("effem record event", "action", op.Action, "collection", collectionName, "rkey", rkey.String(), "repo", evt.Repo) + metrics.IncFirehoseEvent("commit", op.Action) switch repomgr.EventKind(op.Action) { case repomgr.EvtKindCreateRecord, repomgr.EvtKindUpdateRecord: @@ -135,22 +144,27 @@ func (srv *Server) handleCommit(ctx context.Context, evt *comatproto.SyncSubscri recCID, recCBOR, err := rr.GetRecordBytes(ctx, op.Path) if err != nil { srv.logger.Warn("failed to get record", "err", err, "path", op.Path) + metrics.IncIndexerError(collectionName, op.Action) continue } if op.Cid != nil && lexutil.LexLink(recCID) != *op.Cid { srv.logger.Warn("record CID mismatch", "path", op.Path, "carCID", recCID, "opCID", op.Cid) + metrics.IncIndexerError(collectionName, op.Action) continue } if recCBOR == nil { srv.logger.Warn("nil record payload", "path", op.Path) + metrics.IncIndexerError(collectionName, op.Action) continue } if err := srv.indexer.IndexRecord(ctx, evt.Repo, collectionName, rkey.String(), recCID.String(), *recCBOR); err != nil { srv.logger.Warn("failed to index record", "err", err, "path", op.Path) + metrics.IncIndexerError(collectionName, op.Action) } case repomgr.EvtKindDeleteRecord: if err := srv.indexer.DeleteRecord(ctx, evt.Repo, collectionName, rkey.String()); err != nil { srv.logger.Warn("failed to delete record", "err", err, "path", op.Path) + metrics.IncIndexerError(collectionName, op.Action) } default: continue diff --git a/appview/health.go b/appview/health.go new file mode 100644 index 0000000..7afc251 --- /dev/null +++ b/appview/health.go @@ -0,0 +1,97 @@ +package appview + +import ( + "context" + "net/http" + "sync/atomic" + "time" + + "github.com/labstack/echo/v4" +) + +// firehoseStaleAfter is how long the firehose may go without events before we +// mark the service degraded. The relay fires many events per second in normal +// operation, so a 30-second gap is already pathological. +const firehoseStaleAfter = 30 * time.Second + +// healthStatus is the JSON payload of GET /_health. Keep fields stable — +// external uptime monitors (see ALERTS.md) key off `status`. +type healthStatus struct { + Status string `json:"status"` + DB bool `json:"db"` + FirehoseUp bool `json:"firehose_up"` + FirehoseSeq int64 `json:"firehose_seq"` + FirehoseAgeS int `json:"firehose_age_s"` +} + +func (srv *Server) handleHealth(c echo.Context) error { + pingCtx, cancel := context.WithTimeout(c.Request().Context(), 2*time.Second) + defer cancel() + + sqlDB, dbErr := srv.db.DB() + if dbErr == nil { + dbErr = sqlDB.PingContext(pingCtx) + } + + status := evaluateHealth( + dbErr == nil, + srv.firehoseUp.Load(), + atomic.LoadInt64(&srv.lastSeq), + loadLastSeqTime(&srv.lastSeqTime), + time.Now(), + ) + + code := http.StatusOK + if status.Status != "ok" { + code = http.StatusServiceUnavailable + } + return c.JSON(code, status) +} + +// loadLastSeqTime reads the most recent firehose event time from an +// atomic.Value. The comma-ok assertion is critical: before the first event, +// Load returns nil and a naked type assertion would panic. +func loadLastSeqTime(v *atomic.Value) time.Time { + raw := v.Load() + if raw == nil { + return time.Time{} + } + t, ok := raw.(time.Time) + if !ok { + return time.Time{} + } + return t +} + +// evaluateHealth is a pure function over the observed state so it can be +// tested without a running DB or firehose. +func evaluateHealth(dbOK, firehoseUp bool, lastSeq int64, lastSeqTime, now time.Time) healthStatus { + out := healthStatus{ + DB: dbOK, + FirehoseUp: firehoseUp, + FirehoseSeq: lastSeq, + } + + var age time.Duration + if !lastSeqTime.IsZero() { + age = now.Sub(lastSeqTime) + if age < 0 { + age = 0 + } + out.FirehoseAgeS = int(age.Seconds()) + } + + switch { + case !dbOK: + out.Status = "degraded" + case !firehoseUp: + out.Status = "degraded" + case lastSeqTime.IsZero(): + out.Status = "starting" + case age > firehoseStaleAfter: + out.Status = "degraded" + default: + out.Status = "ok" + } + return out +} diff --git a/appview/health_test.go b/appview/health_test.go new file mode 100644 index 0000000..4b775a9 --- /dev/null +++ b/appview/health_test.go @@ -0,0 +1,88 @@ +package appview + +import ( + "sync/atomic" + "testing" + "time" +) + +func TestEvaluateHealth(t *testing.T) { + t.Parallel() + now := time.Date(2026, 4, 20, 12, 0, 0, 0, time.UTC) + fresh := now.Add(-1 * time.Second) + stale := now.Add(-(firehoseStaleAfter + time.Second)) + + cases := []struct { + name string + dbOK bool + firehoseUp bool + lastSeq int64 + lastSeqTime time.Time + wantStatus string + wantFirehoseS int + }{ + {"all ok", true, true, 42, fresh, "ok", 1}, + {"db down", false, true, 42, fresh, "degraded", 1}, + {"firehose not yet started", true, true, 0, time.Time{}, "starting", 0}, + {"firehose disconnected", true, false, 42, fresh, "degraded", 1}, + {"firehose stale events", true, true, 42, stale, "degraded", int(firehoseStaleAfter.Seconds()) + 1}, + {"zero lastSeqTime clamps age to zero", true, true, 0, time.Time{}, "starting", 0}, + } + for _, tc := range cases { + tc := tc + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + got := evaluateHealth(tc.dbOK, tc.firehoseUp, tc.lastSeq, tc.lastSeqTime, now) + if got.Status != tc.wantStatus { + t.Errorf("status: want %q, got %q", tc.wantStatus, got.Status) + } + if got.DB != tc.dbOK { + t.Errorf("db: want %v, got %v", tc.dbOK, got.DB) + } + if got.FirehoseUp != tc.firehoseUp { + t.Errorf("firehose_up: want %v, got %v", tc.firehoseUp, got.FirehoseUp) + } + if got.FirehoseSeq != tc.lastSeq { + t.Errorf("firehose_seq: want %d, got %d", tc.lastSeq, got.FirehoseSeq) + } + if got.FirehoseAgeS != tc.wantFirehoseS { + t.Errorf("firehose_age_s: want %d, got %d", tc.wantFirehoseS, got.FirehoseAgeS) + } + }) + } +} + +func TestEvaluateHealthClampsNegativeAge(t *testing.T) { + t.Parallel() + now := time.Now() + future := now.Add(1 * time.Second) + got := evaluateHealth(true, true, 1, future, now) + if got.FirehoseAgeS != 0 { + t.Fatalf("expected non-negative age, got %d", got.FirehoseAgeS) + } + if got.Status != "ok" { + t.Fatalf("expected ok, got %s", got.Status) + } +} + +func TestLoadLastSeqTimeHandlesNilAndWrongType(t *testing.T) { + t.Parallel() + + var v atomic.Value + if got := loadLastSeqTime(&v); !got.IsZero() { + t.Fatalf("empty atomic.Value should return zero time, got %v", got) + } + + v.Store("not a time") + if got := loadLastSeqTime(&v); !got.IsZero() { + t.Fatalf("wrong-type value should return zero time, got %v", got) + } + + want := time.Date(2026, 4, 20, 12, 0, 0, 0, time.UTC) + // atomic.Value panics on type change, so use a fresh one. + var v2 atomic.Value + v2.Store(want) + if got := loadLastSeqTime(&v2); !got.Equal(want) { + t.Fatalf("want %v, got %v", want, got) + } +} diff --git a/appview/httpmw/auth.go b/appview/httpmw/auth.go index 41a84de..4df3c6e 100644 --- a/appview/httpmw/auth.go +++ b/appview/httpmw/auth.go @@ -157,6 +157,11 @@ func shouldSkipAuth(r *http.Request) bool { if r.URL.Path == "/_health" { return true } + // /metrics must run through the authenticator so RequireScope("admin") can + // evaluate the token on the handler side. Other non-xrpc paths remain unauthed. + if r.URL.Path == "/metrics" { + return false + } return !strings.HasPrefix(r.URL.Path, "/xrpc/") } diff --git a/appview/httpmw/ratelimit.go b/appview/httpmw/ratelimit.go index 30399dc..998e34a 100644 --- a/appview/httpmw/ratelimit.go +++ b/appview/httpmw/ratelimit.go @@ -3,11 +3,13 @@ package httpmw import ( "fmt" "net/http" + "strings" "sync" "time" "github.com/labstack/echo/v4" "golang.org/x/time/rate" + "tangled.org/sparrowtek.com/effem-AppView/appview/metrics" ) type RateLimiterConfig struct { @@ -64,13 +66,15 @@ func (rl *PrincipalRateLimiter) Middleware(next echo.HandlerFunc) echo.HandlerFu } return func(c echo.Context) error { - if c.Request().Method == http.MethodOptions || c.Request().URL.Path == "/_health" { + p := c.Request().URL.Path + if c.Request().Method == http.MethodOptions || p == "/_health" || p == "/metrics" { return next(c) } key := rl.identifier(c) now := time.Now() if !rl.allow(key, now) { + metrics.IncRateLimitRejection(identifierKind(key)) c.Response().Header().Set("Retry-After", "1") return c.JSON(http.StatusTooManyRequests, map[string]string{ "error": "RateLimited", @@ -82,6 +86,21 @@ func (rl *PrincipalRateLimiter) Middleware(next echo.HandlerFunc) echo.HandlerFu } } +// identifierKind classifies a rate limit identifier for metrics labeling. +// Keep the returned set small — it ends up as a Prometheus label value. +func identifierKind(id string) string { + switch { + case strings.HasPrefix(id, "sub:"): + return "sub" + case id == "ip:unknown": + return "unknown" + case strings.HasPrefix(id, "ip:"): + return "ip" + default: + return "other" + } +} + func (rl *PrincipalRateLimiter) identifier(c echo.Context) string { if principal, ok := PrincipalFromContext(c); ok { if principal.Subject != "" { diff --git a/appview/metrics/metrics.go b/appview/metrics/metrics.go new file mode 100644 index 0000000..37d425e --- /dev/null +++ b/appview/metrics/metrics.go @@ -0,0 +1,211 @@ +// Package metrics wires Prometheus counters, histograms, and gauges for the +// Effem AppView. Everything lives on a dedicated registry so tests can count +// on a clean namespace and the default global registry stays empty. +// +// Cardinality is bounded by design: +// - path labels use Echo's matched route pattern, not the raw URL +// - collection labels are AT Protocol collection names, of which we index a +// handful +// - status labels are HTTP status codes, all in the 100–599 range +// - cache endpoint labels are the explicit names passed at the call site +// +// If you add a new label, make sure its value space is bounded before shipping. +package metrics + +import ( + "net/http" + "strconv" + "time" + + "github.com/labstack/echo/v4" + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" +) + +const namespace = "effem" + +var ( + registry = prometheus.NewRegistry() + + httpRequestsTotal = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: namespace, + Name: "http_requests_total", + Help: "Total HTTP requests served, partitioned by method, route pattern, and status.", + }, + []string{"method", "path", "status"}, + ) + + httpRequestDuration = prometheus.NewHistogramVec( + prometheus.HistogramOpts{ + Namespace: namespace, + Name: "http_request_duration_seconds", + Help: "HTTP request latency in seconds, partitioned by method and route pattern.", + Buckets: prometheus.DefBuckets, + }, + []string{"method", "path"}, + ) + + firehoseEventsTotal = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: namespace, + Name: "firehose_events_total", + Help: "AT Protocol firehose events observed, partitioned by kind (commit/identity) and action (create/update/delete).", + }, + []string{"kind", "action"}, + ) + + firehoseConnected = prometheus.NewGauge( + prometheus.GaugeOpts{ + Namespace: namespace, + Name: "firehose_connected", + Help: "1 if the AppView currently holds a firehose WebSocket, 0 otherwise.", + }, + ) + + indexerErrorsTotal = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: namespace, + Name: "indexer_errors_total", + Help: "Indexer errors by collection and operation (create/update/delete).", + }, + []string{"collection", "action"}, + ) + + piCacheHitsTotal = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: namespace, + Name: "pi_cache_hits_total", + Help: "Podcast Index cache hits, partitioned by endpoint.", + }, + []string{"endpoint"}, + ) + + piCacheMissesTotal = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: namespace, + Name: "pi_cache_misses_total", + Help: "Podcast Index cache misses (records fetched from upstream), partitioned by endpoint.", + }, + []string{"endpoint"}, + ) + + rateLimitRejectionsTotal = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: namespace, + Name: "rate_limit_rejections_total", + Help: "HTTP 429 responses emitted by the rate limiter, partitioned by identifier kind (sub/ip/unknown).", + }, + []string{"kind"}, + ) +) + +func init() { + registry.MustRegister( + httpRequestsTotal, + httpRequestDuration, + firehoseEventsTotal, + firehoseConnected, + indexerErrorsTotal, + piCacheHitsTotal, + piCacheMissesTotal, + rateLimitRejectionsTotal, + ) +} + +// Handler returns an http.Handler that serves the registered metrics in the +// Prometheus text exposition format. +func Handler() http.Handler { + return promhttp.HandlerFor(registry, promhttp.HandlerOpts{ + Registry: registry, + EnableOpenMetrics: true, + }) +} + +// RegisterFirehoseLag installs a dynamic gauge that computes the age of the +// most recent firehose event at scrape time. Passing fn lets the caller keep +// the underlying state (atomic.Value, mutex-guarded time, whatever) without +// this package depending on the Server type. +// +// Call once at startup. Calling a second time is a no-op returning an error +// surfaced through MustRegister so misuse shows up in tests. +func RegisterFirehoseLag(fn func() float64) { + gauge := prometheus.NewGaugeFunc( + prometheus.GaugeOpts{ + Namespace: namespace, + Name: "firehose_lag_seconds", + Help: "Wall-clock seconds since the most recent firehose event reached this AppView; 0 before the first event.", + }, + fn, + ) + registry.MustRegister(gauge) +} + +// HTTPMiddleware records request count and duration keyed by method and the +// matched Echo route pattern. If no route matched (404), the label is +// "unknown" so we don't emit a new label value per bogus URL. +func HTTPMiddleware() echo.MiddlewareFunc { + return func(next echo.HandlerFunc) echo.HandlerFunc { + return func(c echo.Context) error { + start := time.Now() + err := next(c) + duration := time.Since(start).Seconds() + + method := c.Request().Method + path := c.Path() + if path == "" { + path = "unknown" + } + status := c.Response().Status + if status == 0 { + // Echo leaves status at zero when the handler returns before + // writing; treat as 200 since Echo ends up writing that. + status = http.StatusOK + } + + httpRequestsTotal.WithLabelValues(method, path, strconv.Itoa(status)).Inc() + httpRequestDuration.WithLabelValues(method, path).Observe(duration) + return err + } + } +} + +// IncFirehoseEvent increments the event counter. Kind is "commit" or +// "identity"; action is "create"/"update"/"delete" for commits, or empty +// for identity events. +func IncFirehoseEvent(kind, action string) { + firehoseEventsTotal.WithLabelValues(kind, action).Inc() +} + +// SetFirehoseConnected flips the connection gauge. Pass true on connect, +// false on disconnect. +func SetFirehoseConnected(connected bool) { + if connected { + firehoseConnected.Set(1) + } else { + firehoseConnected.Set(0) + } +} + +// IncIndexerError increments the indexer error counter for a collection/action. +func IncIndexerError(collection, action string) { + indexerErrorsTotal.WithLabelValues(collection, action).Inc() +} + +// IncCacheHit records a Podcast Index cache hit for the given endpoint. +func IncCacheHit(endpoint string) { + piCacheHitsTotal.WithLabelValues(endpoint).Inc() +} + +// IncCacheMiss records a Podcast Index cache miss (an upstream fetch) for the +// given endpoint. +func IncCacheMiss(endpoint string) { + piCacheMissesTotal.WithLabelValues(endpoint).Inc() +} + +// IncRateLimitRejection records a 429. Kind classifies the identifier that +// tripped the bucket: "sub" for authenticated principals, "ip" for IP-based +// callers, "unknown" for the fallback bucket. +func IncRateLimitRejection(kind string) { + rateLimitRejectionsTotal.WithLabelValues(kind).Inc() +} diff --git a/appview/metrics/metrics_test.go b/appview/metrics/metrics_test.go new file mode 100644 index 0000000..12961fa --- /dev/null +++ b/appview/metrics/metrics_test.go @@ -0,0 +1,90 @@ +package metrics + +import ( + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/labstack/echo/v4" +) + +func TestHandlerServesRegisteredMetrics(t *testing.T) { + // Emit one observation for each metric so the exposition body is non-empty + // and we can assert the expected names appear. + IncFirehoseEvent("commit", "create") + SetFirehoseConnected(true) + IncIndexerError("xyz.effem.feed.comment", "create") + IncCacheHit("search_podcasts") + IncCacheMiss("search_podcasts") + IncRateLimitRejection("sub") + + // Exercise the HTTP middleware so the request counter/histogram have at + // least one sample to emit. + e := echo.New() + e.Use(HTTPMiddleware()) + e.GET("/ping", func(c echo.Context) error { + return c.NoContent(http.StatusOK) + }) + pingReq := httptest.NewRequest(http.MethodGet, "/ping", nil) + pingRec := httptest.NewRecorder() + e.ServeHTTP(pingRec, pingReq) + if pingRec.Code != http.StatusOK { + t.Fatalf("ping: want 200, got %d", pingRec.Code) + } + + req := httptest.NewRequest(http.MethodGet, "/metrics", nil) + rec := httptest.NewRecorder() + Handler().ServeHTTP(rec, req) + + if rec.Code != http.StatusOK { + t.Fatalf("want 200, got %d", rec.Code) + } + + body, err := io.ReadAll(rec.Body) + if err != nil { + t.Fatalf("read body: %v", err) + } + text := string(body) + + wantSubstrings := []string{ + "effem_http_requests_total", + "effem_http_request_duration_seconds", + "effem_firehose_events_total", + "effem_firehose_connected", + "effem_indexer_errors_total", + "effem_pi_cache_hits_total", + "effem_pi_cache_misses_total", + "effem_rate_limit_rejections_total", + } + for _, s := range wantSubstrings { + if !strings.Contains(text, s) { + t.Errorf("expected %q in /metrics output", s) + } + } +} + +func TestRegisterFirehoseLagExposesGauge(t *testing.T) { + // Install a lag source. Calling a second time on the global registry would + // panic on MustRegister, so assert both behaviors: first call succeeds and + // a panic is recovered on a second call. + defer func() { + if r := recover(); r == nil { + t.Fatal("expected panic on duplicate registration, got nil") + } + }() + + RegisterFirehoseLag(func() float64 { return 1.5 }) + + req := httptest.NewRequest(http.MethodGet, "/metrics", nil) + rec := httptest.NewRecorder() + Handler().ServeHTTP(rec, req) + body, _ := io.ReadAll(rec.Body) + if !strings.Contains(string(body), "effem_firehose_lag_seconds") { + t.Fatalf("expected firehose_lag_seconds in output, got:\n%s", body) + } + + // Second call must panic (MustRegister rejects duplicate collector). + RegisterFirehoseLag(func() float64 { return 9 }) +} diff --git a/appview/podcastindex/cache.go b/appview/podcastindex/cache.go index 832a1d7..e6a5888 100644 --- a/appview/podcastindex/cache.go +++ b/appview/podcastindex/cache.go @@ -6,9 +6,10 @@ import ( "fmt" "time" - "tangled.org/sparrowtek.com/effem-AppView/appview/database" "gorm.io/gorm" "gorm.io/gorm/clause" + "tangled.org/sparrowtek.com/effem-AppView/appview/database" + "tangled.org/sparrowtek.com/effem-AppView/appview/metrics" ) type CachedClient struct { @@ -20,20 +21,25 @@ func NewCachedClient(inner *Client, db *gorm.DB) *CachedClient { return &CachedClient{inner: inner, db: db} } +// getOrFetch serves a response from the pi_cache row when it is still valid, +// falling back to a fresh upstream call. Endpoint is a stable, low-cardinality +// label (not the cache key) so Prometheus cardinality stays bounded. func (cc *CachedClient) getOrFetch( - cacheKey string, + endpoint, cacheKey string, ttl time.Duration, fetcher func() (json.RawMessage, error), ) (json.RawMessage, error) { var cached database.PICache err := cc.db.Where("cache_key = ? AND expires_at > ?", cacheKey, time.Now()).First(&cached).Error if err == nil { + metrics.IncCacheHit(endpoint) return json.RawMessage(cached.Response), nil } if !errors.Is(err, gorm.ErrRecordNotFound) { return nil, err } + metrics.IncCacheMiss(endpoint) data, err := fetcher() if err != nil { return nil, err @@ -57,91 +63,91 @@ func (cc *CachedClient) getOrFetch( func (cc *CachedClient) SearchByTerm(query string, max int) (json.RawMessage, error) { key := fmt.Sprintf("search:podcasts:%s:%d", query, max) - return cc.getOrFetch(key, time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("search_podcasts", key, time.Hour, func() (json.RawMessage, error) { return cc.inner.SearchByTerm(query, max) }) } func (cc *CachedClient) SearchEpisodesByTerm(query string, max int) (json.RawMessage, error) { key := fmt.Sprintf("search:episodes:%s:%d", query, max) - return cc.getOrFetch(key, time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("search_episodes", key, time.Hour, func() (json.RawMessage, error) { return cc.inner.SearchEpisodesByTerm(query, max) }) } func (cc *CachedClient) GetPodcastByFeedID(feedID int64) (json.RawMessage, error) { key := fmt.Sprintf("podcast:%d", feedID) - return cc.getOrFetch(key, 24*time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("podcast_by_feed_id", key, 24*time.Hour, func() (json.RawMessage, error) { return cc.inner.GetPodcastByFeedID(feedID) }) } func (cc *CachedClient) GetEpisodesByFeedID(feedID int64, max int) (json.RawMessage, error) { key := fmt.Sprintf("episodes:%d:%d", feedID, max) - return cc.getOrFetch(key, 15*time.Minute, func() (json.RawMessage, error) { + return cc.getOrFetch("episodes_by_feed_id", key, 15*time.Minute, func() (json.RawMessage, error) { return cc.inner.GetEpisodesByFeedID(feedID, max) }) } func (cc *CachedClient) GetEpisodeByID(episodeID int64) (json.RawMessage, error) { key := fmt.Sprintf("episode:%d", episodeID) - return cc.getOrFetch(key, 6*time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("episode_by_id", key, 6*time.Hour, func() (json.RawMessage, error) { return cc.inner.GetEpisodeByID(episodeID) }) } func (cc *CachedClient) GetTrending(max int, lang string, categories string) (json.RawMessage, error) { key := fmt.Sprintf("trending:%d:%s:%s", max, lang, categories) - return cc.getOrFetch(key, 30*time.Minute, func() (json.RawMessage, error) { + return cc.getOrFetch("trending", key, 30*time.Minute, func() (json.RawMessage, error) { return cc.inner.GetTrending(max, lang, categories) }) } func (cc *CachedClient) GetCategories() (json.RawMessage, error) { key := "categories:list" - return cc.getOrFetch(key, 7*24*time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("categories", key, 7*24*time.Hour, func() (json.RawMessage, error) { return cc.inner.GetCategories() }) } func (cc *CachedClient) GetRecentEpisodes(max int) (json.RawMessage, error) { key := fmt.Sprintf("recent:episodes:%d", max) - return cc.getOrFetch(key, 15*time.Minute, func() (json.RawMessage, error) { + return cc.getOrFetch("recent_episodes", key, 15*time.Minute, func() (json.RawMessage, error) { return cc.inner.GetRecentEpisodes(max) }) } func (cc *CachedClient) GetPodcastByItunesID(itunesID int64) (json.RawMessage, error) { key := fmt.Sprintf("podcast:itunes:%d", itunesID) - return cc.getOrFetch(key, 24*time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("podcast_by_itunes_id", key, 24*time.Hour, func() (json.RawMessage, error) { return cc.inner.GetPodcastByItunesID(itunesID) }) } func (cc *CachedClient) GetEpisodesByFeedURL(feedURL string, max int) (json.RawMessage, error) { key := fmt.Sprintf("episodes:url:%s:%d", feedURL, max) - return cc.getOrFetch(key, 15*time.Minute, func() (json.RawMessage, error) { + return cc.getOrFetch("episodes_by_feed_url", key, 15*time.Minute, func() (json.RawMessage, error) { return cc.inner.GetEpisodesByFeedURL(feedURL, max) }) } func (cc *CachedClient) SearchByTitle(query string, max int) (json.RawMessage, error) { key := fmt.Sprintf("search:title:%s:%d", query, max) - return cc.getOrFetch(key, time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("search_title", key, time.Hour, func() (json.RawMessage, error) { return cc.inner.SearchByTitle(query, max) }) } func (cc *CachedClient) SearchByPerson(query string, max int) (json.RawMessage, error) { key := fmt.Sprintf("search:person:%s:%d", query, max) - return cc.getOrFetch(key, time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("search_person", key, time.Hour, func() (json.RawMessage, error) { return cc.inner.SearchByPerson(query, max) }) } func (cc *CachedClient) GetStats() (json.RawMessage, error) { key := "stats:current" - return cc.getOrFetch(key, 24*time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch("stats", key, 24*time.Hour, func() (json.RawMessage, error) { return cc.inner.GetStats() }) } diff --git a/appview/server.go b/appview/server.go index c0e8866..174578b 100644 --- a/appview/server.go +++ b/appview/server.go @@ -20,18 +20,20 @@ import ( "tangled.org/sparrowtek.com/effem-AppView/appview/handlers" "tangled.org/sparrowtek.com/effem-AppView/appview/httpmw" "tangled.org/sparrowtek.com/effem-AppView/appview/indexer" + "tangled.org/sparrowtek.com/effem-AppView/appview/metrics" "tangled.org/sparrowtek.com/effem-AppView/appview/podcastindex" ) type Server struct { - db *gorm.DB - echo *echo.Echo - pi *podcastindex.CachedClient - indexer *indexer.Indexer - config Config - logger *slog.Logger - lastSeq int64 - firehoseUp atomic.Bool + db *gorm.DB + echo *echo.Echo + pi *podcastindex.CachedClient + indexer *indexer.Indexer + config Config + logger *slog.Logger + lastSeq int64 + lastSeqTime atomic.Value // time.Time of the most recent firehose event + firehoseUp atomic.Bool } func NewServer(cfg Config) (*Server, error) { @@ -110,6 +112,7 @@ func NewServer(cfg Config) (*Server, error) { })) e.Use(httpmw.Authentication(authz, cfg.AuthRequired)) e.Use(rateLimiter.Middleware) + e.Use(metrics.HTTPMiddleware()) srv := &Server{ db: db, @@ -119,6 +122,18 @@ func NewServer(cfg Config) (*Server, error) { config: cfg, logger: logger, } + + // Expose firehose lag as a scrape-time gauge. Reading lastSeqTime via the + // closure means the gauge reflects the moment Prometheus scrapes, not a + // stale sampled value. + metrics.RegisterFirehoseLag(func() float64 { + t := loadLastSeqTime(&srv.lastSeqTime) + if t.IsZero() { + return 0 + } + return time.Since(t).Seconds() + }) + srv.registerRoutes() return srv, nil @@ -127,12 +142,12 @@ func NewServer(cfg Config) (*Server, error) { func (srv *Server) registerRoutes() { h := handlers.New(srv.db, srv.pi, srv.logger) - srv.echo.GET("/_health", func(c echo.Context) error { - if !srv.firehoseUp.Load() { - return c.JSON(http.StatusServiceUnavailable, map[string]string{"status": "firehose down"}) - } - return c.JSON(http.StatusOK, map[string]string{"status": "ok"}) - }) + srv.echo.GET("/_health", srv.handleHealth) + + // /metrics is gated by admin scope. The authenticator still runs because + // shouldSkipAuth explicitly keeps /metrics in scope; RequireScope then + // rejects any token that does not carry admin rights. + srv.echo.GET("/metrics", echo.WrapHandler(metrics.Handler()), httpmw.RequireScope("admin")) xrpc := srv.echo.Group("/xrpc", httpmw.RequireScope("read")) diff --git a/go.mod b/go.mod index 1f0f76b..450c7e0 100644 --- a/go.mod +++ b/go.mod @@ -6,6 +6,7 @@ require ( github.com/bluesky-social/indigo v0.0.0-20260211203311-b98f898303a4 github.com/gorilla/websocket v1.5.3 github.com/labstack/echo/v4 v4.12.0 + github.com/prometheus/client_golang v1.17.0 github.com/urfave/cli/v2 v2.27.6 golang.org/x/time v0.5.0 gorm.io/driver/postgres v1.5.11 @@ -74,7 +75,6 @@ require ( github.com/multiformats/go-varint v0.0.7 // indirect github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/polydawn/refmt v0.89.1-0.20221221234430-40501e09de1f // indirect - github.com/prometheus/client_golang v1.17.0 // indirect github.com/prometheus/client_model v0.5.0 // indirect github.com/prometheus/common v0.45.0 // indirect github.com/prometheus/procfs v0.12.0 // indirect