Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 25 additions & 8 deletions packages/apigateway/cmd/apigateway/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,9 @@ func main() {
log.Info("Server exited gracefully")
}

// otelSetupTimeout bounds OTel initialization; see its use in run().
const otelSetupTimeout = 10 * time.Second

func run(log *slog.Logger) error {
// Load configuration
cfg, err := config.LoadConfig()
Expand All @@ -85,8 +88,18 @@ func run(log *slog.Logger) error {

ctx := context.Background()

// Initialize OTel metrics + tracing (warn and continue with no-ops on failure)
providers, otelShutdown, otelErr := otel.Setup(ctx, log, "desirelines-api-gateway")
// Initialize OTel metrics + tracing (warn and continue with no-ops on failure).
//
// Bounded: Setup() passes its context to GCP resource detection (a
// metadata-server call) and to trace-exporter construction, so an
// unreachable metadata service or a stuck exporter would otherwise hang
// startup indefinitely — before the HTTP server is listening, so Cloud Run
// sees a container that never becomes ready. On timeout Setup errors and the
// branch below falls back to NoopProviders: serve traffic unobserved rather
// than not at all.
otelCtx, otelCancel := context.WithTimeout(ctx, otelSetupTimeout)
providers, otelShutdown, otelErr := otel.Setup(otelCtx, log, "desirelines-api-gateway")
otelCancel()
if otelErr != nil {
log.Warn("OTel disabled, using no-op providers", "error", otelErr)
providers = otel.NoopProviders()
Expand Down Expand Up @@ -275,6 +288,10 @@ func initDependencies(ctx context.Context, cfg *config.Config, log *slog.Logger,
// the apigateway. Old metric `auth/firebase_verify.duration` predates the
// span and is no longer written; query the new name going forward.
authHist := newDurationHistogram(meter, log, "desirelines.io/auth/verify_id_token.duration", "Firebase ID token verification duration")
// Separate from authHist: the allowlist re-check hits a different backend
// (Firestore, or its TTL cache) and was previously traced but never metered,
// so there was no time-series to alert a degraded allowlist on.
accessHist := newDurationHistogram(meter, log, "desirelines.io/auth/check_access.duration", "Allowlist access-check duration")
oauthHist := newDurationHistogram(meter, log, "desirelines.io/strava/oauth_exchange.duration", "Strava OAuth exchange duration")
httpHist := newDurationHistogram(meter, log, "desirelines.io/http/request.duration", "HTTP request duration")
deps.httpHistogram = httpHist
Expand All @@ -300,12 +317,12 @@ func initDependencies(ctx context.Context, cfg *config.Config, log *slog.Logger,
// 4–7. Auth setup: Firebase (via emulator in local dev) + Strava (mock in local dev)
if cfg.Environment.IsLocal() && os.Getenv("FIREBASE_AUTH_EMULATOR_HOST") != "" {
// Local dev: real Firebase auth (via emulator) + mock Strava
if authErr := initLocalDevAuth(ctx, cfg, deps, log, authHist, tracer); authErr != nil {
if authErr := initLocalDevAuth(ctx, cfg, deps, log, authHist, accessHist, tracer); authErr != nil {
return nil, authErr
}
} else if !cfg.Environment.IsLocal() {
// Production/staging: real Firebase + real Strava
if authErr := initFirebaseAuth(ctx, cfg, deps, log, authHist, oauthHist, tracer); authErr != nil {
if authErr := initFirebaseAuth(ctx, cfg, deps, log, authHist, accessHist, oauthHist, tracer); authErr != nil {
return nil, authErr
}
} else {
Expand Down Expand Up @@ -536,7 +553,7 @@ func newFirebaseAuthClient(ctx context.Context, projectID string) (*firebaseauth

// initFirebaseAuth initializes Firebase, Firestore, and OAuth dependencies.
// Extracted from initDependencies for cyclomatic complexity.
func initFirebaseAuth(ctx context.Context, cfg *config.Config, deps *Dependencies, log *slog.Logger, authHist, oauthHist otelmetric.Float64Histogram, tracer trace.Tracer) error {
func initFirebaseAuth(ctx context.Context, cfg *config.Config, deps *Dependencies, log *slog.Logger, authHist, accessHist, oauthHist otelmetric.Float64Histogram, tracer trace.Tracer) error {
authClient, err := newFirebaseAuthClient(ctx, cfg.GCPProjectID)
if err != nil {
return err
Expand All @@ -555,7 +572,7 @@ func initFirebaseAuth(ctx context.Context, cfg *config.Config, deps *Dependencie
allowChecker = allowlist.NewCachingChecker(allowChecker, cfg.AllowlistCacheTTL, 0)
}
log.Info("API allowlist cache configured", "ttl", cfg.AllowlistCacheTTL)
deps.authMiddleware = middleware.NewAuthMiddlewareWithAccessCheck(authClient, allowChecker, log, authHist, tracer)
deps.authMiddleware = middleware.NewAuthMiddlewareWithAccessCheck(authClient, allowChecker, log, authHist, accessHist, tracer)

authHandler, err := initAuthHandler(cfg, authClient, firestoreClient, allowChecker, log, oauthHist)
if err != nil {
Expand Down Expand Up @@ -600,7 +617,7 @@ func randomSecret(n int) (string, error) {
return base64.StdEncoding.EncodeToString(b), nil
}

func initLocalDevAuth(ctx context.Context, cfg *config.Config, deps *Dependencies, log *slog.Logger, authHist otelmetric.Float64Histogram, tracer trace.Tracer) error {
func initLocalDevAuth(ctx context.Context, cfg *config.Config, deps *Dependencies, log *slog.Logger, authHist, accessHist otelmetric.Float64Histogram, tracer trace.Tracer) error {
log.Info("Local dev auth: Firebase emulator + mock Strava")

// Firebase Admin SDK auto-detects FIREBASE_AUTH_EMULATOR_HOST
Expand Down Expand Up @@ -646,7 +663,7 @@ func initLocalDevAuth(ctx context.Context, cfg *config.Config, deps *Dependencie
// Exercise the same per-request authorization path locally. The mock checker
// always allows, but keeping it wired prevents local development from
// silently diverging from the production middleware chain.
deps.authMiddleware = middleware.NewAuthMiddlewareWithAccessCheck(authClient, mockAllowlist, log, authHist, tracer)
deps.authMiddleware = middleware.NewAuthMiddlewareWithAccessCheck(authClient, mockAllowlist, log, authHist, accessHist, tracer)

handler, err := auth.NewHandler(&auth.HandlerConfig{
Strava: mockStrava,
Expand Down
19 changes: 16 additions & 3 deletions packages/apigateway/middleware/auth.go
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,12 @@ type AuthMiddleware struct {
access AccessChecker
logger *slog.Logger
histogram metric.Float64Histogram
tracer otelTrace.Tracer
// accessHistogram times the allowlist re-check. Separate from `histogram`
// (Firebase verification) because the two answer different questions and a
// shared instrument would blur a slow allowlist into token-verify latency.
// Nil when there is no access checker, or from callers that don't meter.
accessHistogram metric.Float64Histogram
tracer otelTrace.Tracer
}

// NewAuthMiddleware creates authentication middleware with a pre-initialized token verifier.
Expand All @@ -88,8 +93,10 @@ func NewAuthMiddleware(verifier TokenVerifier, logger *slog.Logger, histogram me
// re-checks the user's current authorization after token verification. Production
// and local composition roots should use this constructor; the shorter constructor
// remains useful for focused verifier tests and examples.
func NewAuthMiddlewareWithAccessCheck(verifier TokenVerifier, access AccessChecker, logger *slog.Logger, histogram metric.Float64Histogram, tracer otelTrace.Tracer) *AuthMiddleware {
return newAuthMiddleware(verifier, access, logger, histogram, tracer)
func NewAuthMiddlewareWithAccessCheck(verifier TokenVerifier, access AccessChecker, logger *slog.Logger, histogram, accessHistogram metric.Float64Histogram, tracer otelTrace.Tracer) *AuthMiddleware {
m := newAuthMiddleware(verifier, access, logger, histogram, tracer)
m.accessHistogram = accessHistogram
return m
}

func newAuthMiddleware(verifier TokenVerifier, access AccessChecker, logger *slog.Logger, histogram metric.Float64Histogram, tracer otelTrace.Tracer) *AuthMiddleware {
Expand Down Expand Up @@ -158,8 +165,14 @@ func (m *AuthMiddleware) Middleware(next http.Handler) http.Handler {

if m.access != nil {
accessCtx, cancel := context.WithTimeout(r.Context(), accessCheckTimeout)
// Span + histogram pairing, matching auth.verify_id_token above: the
// span shows one request's latency, the histogram gives the
// alertable time-series. The check was previously traced only, so a
// degraded allowlist backend had no metric to alert on.
spanCtx, accessSpanDone := otel.StartSpan(accessCtx, m.tracer, "auth.check_access")
accessDone := otel.RecordDuration(spanCtx, m.accessHistogram)
allowed, accessErr := m.access.IsAllowed(spanCtx, token.UID)
accessDone(accessErr)
accessSpanDone(accessErr)
cancel()
if accessErr != nil {
Expand Down
2 changes: 1 addition & 1 deletion packages/apigateway/middleware/auth_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,7 @@ func TestAuthMiddleware_AccessCheck(t *testing.T) {
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
nextCalled := false
am := NewAuthMiddlewareWithAccessCheck(verifier, tt.checker, logger, nil, nil)
am := NewAuthMiddlewareWithAccessCheck(verifier, tt.checker, logger, nil, nil, nil)
handler := am.Middleware(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
nextCalled = true
w.WriteHeader(http.StatusOK)
Expand Down
44 changes: 40 additions & 4 deletions packages/dispatcher/adapters/strava/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -183,7 +183,14 @@ var _ ports.StravaClient = (*Client)(nil)
// NewClient creates a new Strava API client.
// OAuth client credentials must be injected by the caller (composition root).
// Per-user tokens are read from the TokenStore on each request.
func NewClient(clientID, clientSecret string, tokenStore ports.TokenStore, logger *slog.Logger, histogram metric.Float64Histogram, tracer trace.Tracer) *Client {
func NewClient(
clientID, clientSecret string,
tokenStore ports.TokenStore,
logger *slog.Logger,
histogram metric.Float64Histogram,
breakerStateCounter metric.Int64Counter,
tracer trace.Tracer,
) *Client {
return &Client{
httpClient: &http.Client{
Timeout: httpClientTimeout,
Expand All @@ -197,15 +204,22 @@ func NewClient(clientID, clientSecret string, tokenStore ports.TokenStore, logge
logger: logger,
histogram: histogram,
tracer: tracer,
breaker: newStravaBreaker(logger, breakerOpenTimeout),
breaker: newStravaBreaker(logger, breakerOpenTimeout, breakerStateCounter),
}
}

// newStravaBreaker builds the circuit breaker shared across all Strava
// outbound calls on a Client. `timeout` parameterizes the open-state
// duration so tests can use a short value; production calls pass
// breakerOpenTimeout (30s).
func newStravaBreaker(logger *slog.Logger, timeout time.Duration) *gobreaker.CircuitBreaker[[]byte] {
//
// stateCounter may be nil (tests, or a failed instrument construction); the
// state-change hook nil-guards it, matching the handler's counter convention.
func newStravaBreaker(
logger *slog.Logger,
timeout time.Duration,
stateCounter metric.Int64Counter,
) *gobreaker.CircuitBreaker[[]byte] {
return gobreaker.NewCircuitBreaker[[]byte](gobreaker.Settings{
Name: "strava-api",
Timeout: timeout,
Expand All @@ -218,6 +232,22 @@ func newStravaBreaker(logger *slog.Logger, timeout time.Duration) *gobreaker.Cir
"from", from.String(),
"to", to.String(),
)
if stateCounter == nil {
return
}
// gobreaker gives the hook no context, and a state change is not
// request-scoped anyway — it is a property of the breaker, not of
// whichever unlucky call tripped it. Background() is correct here.
//
// Both `from` and `to` are recorded: "how many times did we open"
// needs `to`, but distinguishing a genuine recovery
// (half-open -> closed) from a flap (open -> half-open -> open)
// needs the pair.
stateCounter.Add(context.Background(), 1, metric.WithAttributes(
attribute.String("breaker", name),
attribute.String("from", from.String()),
attribute.String("to", to.String()),
))
},
IsSuccessful: isStravaCallSuccessful,
})
Expand Down Expand Up @@ -373,7 +403,13 @@ func (c *Client) verifyGrantWithTokens(ctx context.Context, ownerID int64, token
// token is nominally live, /athlete is a cheap, read-only proof that it still
// belongs to this owner. A 401 falls through to the stronger refresh check.
if tokens.AccessToken != "" && time.Now().Unix() < tokens.ExpiresAt {
if verifyErr := c.doVerifyCurrentAthlete(ctx, ownerID, tokens.AccessToken); verifyErr == nil {
// Timed like every other outbound Strava op. Without this the deauth
// path was absent from strava.api.duration entirely, so the histogram
// under-counted real Strava traffic and a slow /athlete was invisible.
done := otel.RecordDuration(ctx, c.histogram, attribute.String("operation", "verify_athlete"))
verifyErr := c.doVerifyCurrentAthlete(ctx, ownerID, tokens.AccessToken)
done(verifyErr)
if verifyErr == nil {
return ports.GrantActive, nil
} else if !isAuthError(verifyErr) {
return ports.GrantUnknown, fmt.Errorf("verify current athlete %d: %w", ownerID, verifyErr)
Expand Down
78 changes: 77 additions & 1 deletion packages/dispatcher/adapters/strava/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ import (
"github.com/andy-esch/desirelines/packages/shared/otel"
"github.com/andy-esch/desirelines/packages/shared/stravatoken"
"github.com/sony/gobreaker/v2"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/metric/metricdata"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
"go.opentelemetry.io/otel/sdk/trace/tracetest"
)
Expand Down Expand Up @@ -163,7 +165,7 @@ func newTestClient(server *httptest.Server, tokenStore ports.TokenStore) *Client
logger: logger,
histogram: noopHist,
tracer: noopProviders.Tracer,
breaker: newStravaBreaker(logger, testBreakerTimeout),
breaker: newStravaBreaker(logger, testBreakerTimeout, nil),
}
}

Expand Down Expand Up @@ -1185,6 +1187,80 @@ func TestJitterBackoff(t *testing.T) {
})
}

// TestCircuitBreaker_EmitsStateChangeMetric pins the alerting signal added for
// audit 2026-08-19-dispatcher L1: before it, OnStateChange only called
// logger.Warn, so an open breaker was discoverable solely by reading logs after
// the fact and could never fire an alert.
//
// The from/to pair is asserted, not just `to`: distinguishing a genuine recovery
// (half-open -> closed) from a flap (open -> half-open -> open) needs both.
func TestCircuitBreaker_EmitsStateChangeMetric(t *testing.T) {
reader := sdkmetric.NewManualReader()
provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
counter, err := provider.Meter("test").Int64Counter("desirelines.io/strava/breaker_state_change")
if err != nil {
t.Fatalf("create counter: %v", err)
}

server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
if _, writeErr := w.Write([]byte(`{"error":"down"}`)); writeErr != nil {
t.Errorf("failed to write response: %v", writeErr)
}
}))
defer server.Close()

tokenStore := &portstest.MockTokenStore{
Tokens: map[int64]*stravatoken.Data{
testOwnerID: {AccessToken: "t", RefreshToken: "r", ExpiresAt: futureExpiry()},
},
}
client := newTestClient(server, tokenStore)
// Rebuild the breaker with the counter attached; newTestClient passes nil.
client.breaker = newStravaBreaker(gcplog.NewNoOpLogger(), testBreakerTimeout, counter)

for range breakerFailureThreshold {
if _, fetchErr := client.FetchActivity(context.Background(), testOwnerID, testActivityID); fetchErr == nil {
t.Fatal("expected error while server is 500")
}
}
if state := client.breaker.State(); state != gobreaker.StateOpen {
t.Fatalf("breaker state = %v, want %v", state, gobreaker.StateOpen)
}

var rm metricdata.ResourceMetrics
if collectErr := reader.Collect(context.Background(), &rm); collectErr != nil {
t.Fatalf("collect metrics: %v", collectErr)
}

got := breakerStateChanges(rm)
if got["closed->open"] != 1 {
t.Errorf("closed->open transitions = %d, want 1 (all: %v)", got["closed->open"], got)
}
}

// breakerStateChanges flattens the breaker counter into "from->to" => count.
func breakerStateChanges(rm metricdata.ResourceMetrics) map[string]int64 {
out := map[string]int64{}
for _, sm := range rm.ScopeMetrics {
for _, m := range sm.Metrics {
if m.Name != "desirelines.io/strava/breaker_state_change" {
continue
}
sum, ok := m.Data.(metricdata.Sum[int64])
if !ok {
continue
}
for _, dp := range sum.DataPoints {
from, _ := dp.Attributes.Value("from")
to, _ := dp.Attributes.Value("to")
out[from.AsString()+"->"+to.AsString()] += dp.Value
}
}
}
return out
}

// TestCircuitBreaker_TripsAfterConsecutiveFailures verifies the breaker
// opens after the configured threshold of consecutive failing
// operations and short-circuits subsequent calls before any HTTP
Expand Down
20 changes: 18 additions & 2 deletions packages/dispatcher/cmd/dispatcher/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,16 @@ import (
)

const (
// otelSetupTimeout bounds OTel initialization. Setup() passes its context to
// GCP resource detection (a metadata-server call) and to trace-exporter
// construction, so an unreachable metadata service or a stuck exporter would
// otherwise hang process startup indefinitely — before the HTTP server is
// listening, so Cloud Run sees a container that never becomes ready. On
// timeout Setup returns an error and the existing otelErr branch falls back
// to NoopProviders, which is the right trade: serve traffic unobserved
// rather than not at all.
otelSetupTimeout = 10 * time.Second

// startupTimeout is the maximum time allowed for initializing dependencies
// (e.g., PubSub gRPC connection). Prevents indefinite hang if GCP metadata
// service is unreachable.
Expand Down Expand Up @@ -77,7 +87,9 @@ func run(log *slog.Logger) error {
}

// Initialize OTel metrics + tracing (warn and continue with no-ops on failure)
providers, otelShutdown, otelErr := otel.Setup(context.Background(), log, "desirelines-dispatcher")
otelCtx, otelCancel := context.WithTimeout(context.Background(), otelSetupTimeout)
providers, otelShutdown, otelErr := otel.Setup(otelCtx, log, "desirelines-dispatcher")
otelCancel()
if otelErr != nil {
log.Warn("OTel disabled, using no-op providers", "error", otelErr)
providers = otel.NoopProviders()
Expand Down Expand Up @@ -239,6 +251,10 @@ func initDependencies(cfg *config.Config, log *slog.Logger, meter metric.Meter,
// 1. Create OTel instruments first so they can be injected into adapters.
// Errors are non-fatal; instruments will be no-op on failure.
stravaHist := newHistogram(meter, log, "desirelines.io/strava/api.duration", "Strava API call duration")
// The breaker previously announced open/half-open transitions only via
// logger.Warn, so an open breaker was invisible to alerting — findable only
// by reading logs after the fact.
stravaBreakerCounter := newCounter(meter, log, "desirelines.io/strava/breaker_state_change", "Strava circuit breaker state transitions (labeled from/to)")
firestoreHist := newHistogram(meter, log, "desirelines.io/firestore/operation.duration", "Firestore operation duration")
pubsubHist := newHistogram(meter, log, "desirelines.io/pubsub/publish.duration", "PubSub publish duration")
webhookCounter := newCounter(meter, log, "desirelines.io/webhook/events", "Webhook events processed")
Expand Down Expand Up @@ -343,7 +359,7 @@ func initDependencies(cfg *config.Config, log *slog.Logger, meter metric.Meter,
return nil, fmt.Errorf("strava client_secret: %w", err)
}

stravaClient := strava.NewClient(stravaClientID, stravaClientSecret, tokenStore, log, stravaHist, tracer)
stravaClient := strava.NewClient(stravaClientID, stravaClientSecret, tokenStore, log, stravaHist, stravaBreakerCounter, tracer)

// Rate limiter: 5 req/s, burst 10 (Strava sends a few events/day normally)
// Uses Background context (not startupCtx) because the cleanup goroutine must
Expand Down
1 change: 1 addition & 0 deletions packages/shared/otel/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ var extendedDurationInstrumentNames = []string{
"desirelines.io/firestore/operation.duration",
"desirelines.io/pubsub/publish.duration",
"desirelines.io/auth/verify_id_token.duration",
"desirelines.io/auth/check_access.duration",
"desirelines.io/strava/oauth_exchange.duration",
}

Expand Down
1 change: 1 addition & 0 deletions packages/shared/otel/provider_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ func TestExtendedDurationViews_MatchEachListedInstrument(t *testing.T) {
"desirelines.io/firestore/operation.duration",
"desirelines.io/pubsub/publish.duration",
"desirelines.io/auth/verify_id_token.duration",
"desirelines.io/auth/check_access.duration",
"desirelines.io/strava/oauth_exchange.duration",
}

Expand Down
Loading