From 61a8c27b6787352bb656267252d1084c3da03a85 Mon Sep 17 00:00:00 2001 From: Kyle Crawshaw Date: Tue, 1 Sep 2026 11:41:32 -0400 Subject: [PATCH] feat(log): configurable json/logfmt/text output with request-scoped labels Operators can pick a collector-friendly encoding from [log], and request-path lines carry request_id (and W3C traceparent ids) without repeating those fields at each call site. Co-authored-by: Cursor --- README.md | 5 +- anteroom.example.toml | 21 ++++- charts/kyverno-policies/values.yaml | 7 ++ cmd/anteroom/main.go | 8 +- docs/docker.md | 4 +- docs/operating.md | 25 +++++- internal/config/config.go | 70 +++++++++++++++ internal/config/config_test.go | 32 +++++++ internal/gate/endpoints.go | 4 +- internal/gate/gate.go | 8 +- internal/gate/logging_test.go | 112 +++++++++++++++++++++++ internal/gate/pay.go | 51 +++++------ internal/gate/request.go | 2 + internal/logging/context.go | 128 ++++++++++++++++++++++++++ internal/logging/context_test.go | 86 ++++++++++++++++++ internal/logging/logging.go | 114 +++++++++++++++++++++++ internal/logging/logging_test.go | 135 ++++++++++++++++++++++++++++ internal/payment/callback.go | 28 +++--- 18 files changed, 783 insertions(+), 57 deletions(-) create mode 100644 internal/gate/logging_test.go create mode 100644 internal/logging/context.go create mode 100644 internal/logging/context_test.go create mode 100644 internal/logging/logging.go create mode 100644 internal/logging/logging_test.go diff --git a/README.md b/README.md index e7c89c7..d0d1a85 100644 --- a/README.md +++ b/README.md @@ -107,8 +107,9 @@ your own site a second later. Startup warnings are worth reading: an auto-generated HMAC key means passes die on restart and no second instance can verify them, which is fine for a first run and wrong in production. -`ANTEROOM_LISTEN`, `ANTEROOM_UPSTREAM`, `ANTEROOM_PAGES`, and `ANTEROOM_HMAC_KEY` -override the file, for containers and secret managers. +`ANTEROOM_LISTEN`, `ANTEROOM_UPSTREAM`, `ANTEROOM_PAGES`, `ANTEROOM_HMAC_KEY`, +`ANTEROOM_LOG_LEVEL`, and `ANTEROOM_LOG_FORMAT` override the file, for containers +and secret managers. Run with `-v` to log one line per request naming which rung of the ladder answered it (`pass-pow`, `wait-page`, `refusal`, `bypass-path`, …) with the status, size, diff --git a/anteroom.example.toml b/anteroom.example.toml index f7de02d..fd46d05 100644 --- a/anteroom.example.toml +++ b/anteroom.example.toml @@ -2,7 +2,8 @@ # # This file is the public configuration contract. # Only `upstream` is required. Environment variables override the file: -# ANTEROOM_LISTEN, ANTEROOM_ADMIN_LISTEN, ANTEROOM_UPSTREAM, ANTEROOM_HMAC_KEY, ... +# ANTEROOM_LISTEN, ANTEROOM_ADMIN_LISTEN, ANTEROOM_UPSTREAM, ANTEROOM_HMAC_KEY, +# ANTEROOM_LOG_LEVEL, ANTEROOM_LOG_FORMAT, ... listen = "127.0.0.1:8080" # bind address. Loopback by default: the documented # topology puts a TLS terminator in front, and a gate @@ -88,6 +89,20 @@ trusted_proxies = [] # CIDRs whose X-Forwarded-For is believed. #kid = "k1" #key = "base64..." +# --------------------------------------------------------------------------- +# Logging. Default is text at info — a terminal. json is the usual choice +# under Kubernetes; logfmt is the usual choice for Grafana Alloy / Promtail. +# `anteroom -v` still forces debug (one hit line per request) regardless of +# level. Request-scoped fields (request_id, and trace_id/span_id when the +# inbound request carried a W3C traceparent) attach automatically; they are +# not configured here. Static labels go on every line. +# --------------------------------------------------------------------------- +# [log] +# level = "info" # debug | info | warn | error +# format = "text" # json | logfmt | text +# [log.labels] +# service = "anteroom" + # --------------------------------------------------------------------------- # Challenge-activity log (optional). Omit this whole section to keep the gate # fully free of per-visitor state. When set, the admin listener (admin_listen @@ -269,7 +284,9 @@ allow_hosted_fetchers = true # Claude-User, ChatGPT-User, Google-Age # gate exposes counters for YOUR scraper and phones home to nobody. # # Run with -v to log one line per request naming which rung of the ladder answered -# it; useful when a request is being walled and you want to know why. +# it; useful when a request is being walled and you want to know why. Each line +# carries request_id (honoring inbound X-Request-ID, or generated) so it joins +# the upstream's logs. format = "json" under [log] if a collector will parse it. # # Before going live, read docs/operating.md — it lists what Anteroom breaks # without bypass rules (webhooks, OAuth callbacks, API clients, feeds, link diff --git a/charts/kyverno-policies/values.yaml b/charts/kyverno-policies/values.yaml index 4bfa5f6..93b2813 100644 --- a/charts/kyverno-policies/values.yaml +++ b/charts/kyverno-policies/values.yaml @@ -163,6 +163,13 @@ gateConfig: | renew_difficulty = 6 # renewals are cheap by design (mobile battery) inject = true # add the renewal script to proxied HTML + # Logging. Default is text at info. json is the usual choice in-cluster. + # [log] + # level = "info" # debug | info | warn | error. -v still forces debug. + # format = "json" # json | logfmt | text + # [log.labels] + # cluster = "prod" + # Trap 1 in docs/docker.md, and it applies with full force to # Kubernetes: WebCrypto and service workers require a secure # context — HTTPS or localhost. Reaching a NodePort at diff --git a/cmd/anteroom/main.go b/cmd/anteroom/main.go index a3883a1..a2d4870 100644 --- a/cmd/anteroom/main.go +++ b/cmd/anteroom/main.go @@ -10,7 +10,6 @@ import ( "flag" "fmt" "io" - "log/slog" "net" "net/http" "os" @@ -22,6 +21,7 @@ import ( "github.com/radiustechsystems/anteroom/internal/admin" "github.com/radiustechsystems/anteroom/internal/config" "github.com/radiustechsystems/anteroom/internal/gate" + "github.com/radiustechsystems/anteroom/internal/logging" ) // bindError supplies the context the standard library's message lacks: a @@ -59,11 +59,6 @@ func run() error { check := flag.Bool("healthcheck", false, "probe the local gate's health endpoint and exit 0 (healthy) or 1; for container HEALTHCHECK, which has no shell to run curl in") flag.Parse() - level := slog.LevelInfo - if *verbose { - level = slog.LevelDebug - } - lg := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: level})) cfg, err := config.Load(*cfgPath) if err != nil { return err @@ -73,6 +68,7 @@ func run() error { return healthcheck(cfg.Listen) } + lg := logging.New(os.Stderr, cfg.Log, *verbose) g, err := gate.New(cfg, lg) if err != nil { return err diff --git a/docs/docker.md b/docs/docker.md index 59f0dae..d95518c 100644 --- a/docs/docker.md +++ b/docs/docker.md @@ -113,7 +113,7 @@ which commit it is. ## The container contract -**Environment variables.** Only these four exist. Everything else is a config +**Environment variables.** Only these exist. Everything else is a config file setting, deliberately — the surface an operator can change without review is kept small. @@ -123,6 +123,8 @@ is kept small. | `ANTEROOM_LISTEN` | bind address; defaults to `:8080` | | `ANTEROOM_PAGES` | directory holding `header.html` and `footer.html` | | `ANTEROOM_HMAC_KEY` | the signing key; registers as `kid = "env"` | +| `ANTEROOM_LOG_LEVEL` | `debug`, `info`, `warn`, or `error`; `-v` still forces debug | +| `ANTEROOM_LOG_FORMAT` | `json`, `logfmt`, or `text` (the default) | **Ports.** The gate binds `:8080` as a non-root user. Publish it as `-p 80:8080` and the daemon owns the privileged half — no capability, no root, nothing to diff --git a/docs/operating.md b/docs/operating.md index 55c8461..67882c3 100644 --- a/docs/operating.md +++ b/docs/operating.md @@ -470,7 +470,7 @@ answered it — `own-endpoint`, `bypass-path`, `bypass-ip`, `bypass-crawler`, duration: ``` -level=DEBUG msg=hit method=GET path=/index.html decision=pass-pow status=200 bytes=142 dur=1.2ms ip=203.0.113.9 ua=Mozilla/5.0… +level=DEBUG msg=hit method=GET path=/index.html decision=pass-pow status=200 bytes=142 dur=1.2ms ip=203.0.113.9 ua=Mozilla/5.0… request_id=… ``` This is the fastest way to answer "why was this request walled?", which is @@ -480,6 +480,29 @@ ignore its logs. Note the line includes the client IP and user agent, so it is request-level data — appropriate for debugging, not for leaving on in production without deciding that is what you want. +Every request also carries a `request_id` (honoring inbound `X-Request-ID` or +`X-Correlation-ID`, otherwise generated) and, when the client sent a W3C +`traceparent`, `trace_id` and `span_id`. The same `X-Request-ID` is forwarded +upstream so the application's logs join. The gate does not mint spans it never +exports. + +Log encoding is configured under `[log]`. The default is `text` at `info` — a +terminal. Collectors want `json`; Grafana Alloy / Promtail want `logfmt`. +`anteroom -v` still forces debug regardless of `log.level`. Static labels +(`[log.labels]`) land on every line; request-scoped fields attach through +context and do not need repeating at each call site. + +```toml +[log] +level = "info" # debug | info | warn | error +format = "json" # json | logfmt | text +[log.labels] +service = "anteroom" +``` + +`ANTEROOM_LOG_LEVEL` and `ANTEROOM_LOG_FORMAT` override the file, the same way +`ANTEROOM_LISTEN` does. + ## Monitoring Set `admin_listen` to open an operator port alongside the gate: diff --git a/internal/config/config.go b/internal/config/config.go index e4f252f..88ce6cf 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -92,6 +92,7 @@ type Config struct { Activity *Activity `toml:"activity"` Bypass Bypass `toml:"bypass"` Triage Triage `toml:"triage"` + Log Log `toml:"log"` // KeyAutoGenerated is set when no hmac_keys were configured and one was // generated for this process. Fleets must not rely on it: every restart @@ -99,6 +100,23 @@ type Config struct { KeyAutoGenerated bool `toml:"-"` } +// Log is the process logger. Format is a Handler choice — every call site +// already emits slog key/value attrs — so json, logfmt, and text are encodings +// of the same records, not three logging APIs. +type Log struct { + // Level is the floor: debug, info, warn, or error. Default info. `anteroom + // -v` still forces debug, so per-request hit lines stay an explicit opt-in + // even when the file says info. + Level string `toml:"level"` + // Format is json, logfmt, or text (the default). json for collectors, + // logfmt for Grafana/Promtail, text for a terminal. + Format string `toml:"format"` + // Labels are static key/value pairs on every line (service, cluster, replica). + // Request-scoped fields (request_id, trace_id, span_id) are not configured + // here: they come from the request and attach through context. + Labels map[string]string `toml:"labels"` +} + type HMACKey struct { Kid string `toml:"kid"` Key string `toml:"key"` // base64 (std or url, padded or not) @@ -283,6 +301,10 @@ func defaults() Config { Difficulty: 14, RenewDifficulty: 6, Inject: true, + Log: Log{ + Level: "info", + Format: "text", + }, Triage: Triage{ JSONAccept: true, // Verified vendor-hosted user fetchers cannot complete PoW or x402, @@ -344,6 +366,12 @@ func applyEnv(cfg *Config) { if v := os.Getenv("ANTEROOM_HMAC_KEY"); v != "" { cfg.HMACKeys = []HMACKey{{Kid: "env", Key: v}} } + if v := os.Getenv("ANTEROOM_LOG_LEVEL"); v != "" { + cfg.Log.Level = v + } + if v := os.Getenv("ANTEROOM_LOG_FORMAT"); v != "" { + cfg.Log.Format = v + } } func (c *Config) validate() error { @@ -438,9 +466,51 @@ func (c *Config) validate() error { return err } } + if err := c.Log.validate(); err != nil { + return err + } return nil } +func (l Log) validate() error { + switch strings.ToLower(strings.TrimSpace(l.Level)) { + case "debug", "info", "warn", "warning", "error": + default: + return fmt.Errorf("config: log.level %q is not debug, info, warn, or error", l.Level) + } + switch strings.ToLower(strings.TrimSpace(l.Format)) { + case "json", "logfmt", "text": + default: + return fmt.Errorf("config: log.format %q is not json, logfmt, or text — json for collectors, logfmt for Grafana, text for a terminal", l.Format) + } + reserved := map[string]bool{"time": true, "ts": true, "level": true, "msg": true} + for k := range l.Labels { + if reserved[k] { + return fmt.Errorf("config: log.labels key %q is reserved (time, ts, level, msg)", k) + } + if !validLabelKey(k) { + return fmt.Errorf("config: log.labels key %q must start with a letter or underscore, then letters, digits, or underscores", k) + } + } + return nil +} + +func validLabelKey(k string) bool { + if k == "" { + return false + } + for i := 0; i < len(k); i++ { + c := k[i] + switch { + case c >= 'a' && c <= 'z', c >= 'A' && c <= 'Z', c == '_': + case i > 0 && c >= '0' && c <= '9': + default: + return false + } + } + return true +} + // checkFacilitatorURL validates a facilitator URL under key (a config path for // error messages). Empty is the caller's concern — for the global key it is an // error, for a rail it means "inherit". diff --git a/internal/config/config_test.go b/internal/config/config_test.go index e0fc6c9..a888da9 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -51,6 +51,9 @@ func TestLoadMinimal(t *testing.T) { if !cfg.Inject || !cfg.Triage.JSONAccept || !cfg.Triage.AllowHostedFetchers { t.Error("inject, json_accept, and allow_hosted_fetchers should default true") } + if cfg.Log.Level != "info" || cfg.Log.Format != "text" { + t.Errorf("log defaults: %+v", cfg.Log) + } // No admin listener unless asked for: it is unauthenticated, so silently // opening a port the operator never configured would be a surprise surface. if cfg.AdminListen != "" { @@ -375,6 +378,10 @@ price = "$0.01" {"activity ttl too long", minimal + "[activity]\nttl = \"48h\"\n", nil, "activity.ttl"}, {"activity max_ips out of range", minimal + "[activity]\nmax_ips = 5000000\n", nil, "activity.max_ips"}, {"activity unknown key", minimal + "[activity]\nttls = \"10m\"\n", nil, "unknown key"}, + {"bad log level", minimal + "[log]\nlevel = \"verbose\"\n", nil, "log.level"}, + {"bad log format", minimal + "[log]\nformat = \"yaml\"\n", nil, "log.format"}, + {"reserved log label", minimal + "[log.labels]\nmsg = \"nope\"\n", nil, "reserved"}, + {"bad log label key", minimal + "[log.labels]\n\"not a key\" = \"x\"\n", nil, "log.labels"}, } { t.Run(tc.name, func(t *testing.T) { _, err := Load(write(t, tc.body)) @@ -425,6 +432,8 @@ func TestEnvOverrides(t *testing.T) { t.Setenv("ANTEROOM_ADMIN_LISTEN", "127.0.0.1:9998") t.Setenv("ANTEROOM_UPSTREAM", "127.0.0.1:4000") t.Setenv("ANTEROOM_HMAC_KEY", "MDEyMzQ1Njc4OWFiY2RlZjAxMjM0NTY3ODlhYmNkZWY=") + t.Setenv("ANTEROOM_LOG_LEVEL", "debug") + t.Setenv("ANTEROOM_LOG_FORMAT", "json") cfg, err := Load(write(t, minimal)) if err != nil { t.Fatalf("Load: %v", err) @@ -438,6 +447,29 @@ func TestEnvOverrides(t *testing.T) { if cfg.KeyAutoGenerated || cfg.HMACKeys[0].Kid != "env" { t.Errorf("env key not used: %+v", cfg.HMACKeys) } + if cfg.Log.Level != "debug" || cfg.Log.Format != "json" { + t.Errorf("log env overrides not applied: %+v", cfg.Log) + } +} + +func TestLogSection(t *testing.T) { + cfg, err := Load(write(t, minimal+` +[log] +level = "warn" +format = "json" +[log.labels] +service = "anteroom" +cluster = "prod" +`)) + if err != nil { + t.Fatalf("Load: %v", err) + } + if cfg.Log.Level != "warn" || cfg.Log.Format != "json" { + t.Errorf("log section: %+v", cfg.Log) + } + if cfg.Log.Labels["service"] != "anteroom" || cfg.Log.Labels["cluster"] != "prod" { + t.Errorf("log.labels: %v", cfg.Log.Labels) + } } func TestParsePrice(t *testing.T) { diff --git a/internal/gate/endpoints.go b/internal/gate/endpoints.go index fa98339..7569cf2 100644 --- a/internal/gate/endpoints.go +++ b/internal/gate/endpoints.go @@ -151,7 +151,7 @@ func (g *Gate) serveChallenge(w http.ResponseWriter, r *http.Request) { } c, issuedAt, err := g.issuer.Issue(now, requestAudience(r), profile) if err != nil { - g.lg.Error("issuing challenge", "err", err) + g.lg.ErrorContext(r.Context(), "issuing challenge", "err", err) http.Error(w, "internal error", http.StatusInternalServerError) return } @@ -264,7 +264,7 @@ func (g *Gate) serveAnswer(w http.ResponseWriter, r *http.Request) { return } if err := g.setPassCookie(w, r, token.Pass{Kind: token.KindPoW, Scope: token.ScopeAll}, exp, rootAt); err != nil { - g.lg.Error("minting pass", "err", err) + g.lg.ErrorContext(r.Context(), "minting pass", "err", err) g.noteAnswer("error", r) w.WriteHeader(http.StatusInternalServerError) json.NewEncoder(w).Encode(answerResponse{Error: "internal error"}) diff --git a/internal/gate/gate.go b/internal/gate/gate.go index 7d73cb0..7874c50 100644 --- a/internal/gate/gate.go +++ b/internal/gate/gate.go @@ -312,7 +312,7 @@ func newProxy(u *url.URL, socket string, m *bypass.Matcher, lg *slog.Logger, ups FlushInterval: -1, ErrorHandler: func(w http.ResponseWriter, r *http.Request, err error) { upstreamErr.Inc() - lg.Error("upstream unreachable", "err", err, "path", r.URL.Path) + lg.ErrorContext(r.Context(), "upstream unreachable", "err", err, "path", r.URL.Path) http.Error(w, "upstream unreachable", http.StatusBadGateway) }, } @@ -390,7 +390,7 @@ func (g *Gate) ServeHTTP(w http.ResponseWriter, r *http.Request) { if ua := request.facts.userAgent; ua != "" { attrs = append(attrs, "ua", ua) } - g.lg.Debug("hit", attrs...) + g.lg.DebugContext(request.Context(), "hit", attrs...) } // serve is the ladder proper. It returns the name of the rung that answered, for @@ -571,11 +571,11 @@ func (g *Gate) forward(w http.ResponseWriter, r *http.Request) { // noteInjectionSkipped logs every skipped injection at debug level and warns // once per reason, avoiding per-request warning noise from persistent policies. func (g *Gate) noteInjectionSkipped(r *http.Request, reason string) { - g.lg.Debug("renewal script not injected", "reason", reason, "path", r.URL.Path) + g.lg.DebugContext(r.Context(), "renewal script not injected", "reason", reason, "path", r.URL.Path) if _, seen := g.skipReported.LoadOrStore(reason, struct{}{}); seen { return } - g.lg.Warn("renewal script not injected — visitors reading these pages will lapse and be re-challenged", + g.lg.WarnContext(r.Context(), "renewal script not injected — visitors reading these pages will lapse and be re-challenged", "reason", reason, "example_path", r.URL.Path, "consequence", "the pass is not renewed from this response, so the visitor is walled again when it expires", diff --git a/internal/gate/logging_test.go b/internal/gate/logging_test.go new file mode 100644 index 0000000..0c71ff8 --- /dev/null +++ b/internal/gate/logging_test.go @@ -0,0 +1,112 @@ +package gate + +import ( + "bytes" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "testing" + + "github.com/radiustechsystems/anteroom/internal/config" + "github.com/radiustechsystems/anteroom/internal/logging" +) + +func TestRequestIDIsForwardedUpstream(t *testing.T) { + var seen http.Header + up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + seen = r.Header.Clone() + w.Header().Set("Content-Type", "text/html") + io.WriteString(w, simplePage) + })) + t.Cleanup(up.Close) + + cfgPath := filepath.Join(t.TempDir(), "anteroom.toml") + if err := os.WriteFile(cfgPath, []byte("upstream = \""+up.URL+"\"\n"+fastCfg), 0o600); err != nil { + t.Fatal(err) + } + cfg, err := config.Load(cfgPath) + if err != nil { + t.Fatal(err) + } + g, err := New(cfg, logging.New(io.Discard, cfg.Log, false)) + if err != nil { + t.Fatal(err) + } + pass := solveAndGetCookie(t, g, nil) + + r := docReq("/page", pass) + r.Header.Set(logging.HeaderRequestID, "keep-me") + r.Header.Set(logging.HeaderTraceparent, "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01") + if res := do(g, r); res.Code != 200 { + t.Fatalf("status %d", res.Code) + } + if got := seen.Get(logging.HeaderRequestID); got != "keep-me" { + t.Errorf("upstream X-Request-ID = %q, want keep-me", got) + } + if got := seen.Get(logging.HeaderTraceparent); got != "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01" { + t.Errorf("upstream traceparent = %q", got) + } + + r = docReq("/page", pass) + if res := do(g, r); res.Code != 200 { + t.Fatalf("status %d", res.Code) + } + if got := seen.Get(logging.HeaderRequestID); got == "" || got == "keep-me" { + t.Errorf("generated X-Request-ID missing or reused: %q", got) + } +} + +func TestHitLineCarriesRequestContext(t *testing.T) { + up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/html") + io.WriteString(w, simplePage) + })) + t.Cleanup(up.Close) + + cfgPath := filepath.Join(t.TempDir(), "anteroom.toml") + if err := os.WriteFile(cfgPath, []byte("upstream = \""+up.URL+"\"\n"+fastCfg), 0o600); err != nil { + t.Fatal(err) + } + cfg, err := config.Load(cfgPath) + if err != nil { + t.Fatal(err) + } + var buf bytes.Buffer + g, err := New(cfg, logging.New(&buf, config.Log{Level: "info", Format: "json"}, true)) + if err != nil { + t.Fatal(err) + } + pass := solveAndGetCookie(t, g, nil) + buf.Reset() + + r := docReq("/page", pass) + r.Header.Set(logging.HeaderRequestID, "hit-req") + r.Header.Set(logging.HeaderTraceparent, "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01") + do(g, r) + + var rec map[string]any + for _, line := range bytes.Split(bytes.TrimSpace(buf.Bytes()), []byte("\n")) { + if !bytes.Contains(line, []byte(`"msg":"hit"`)) { + continue + } + if err := json.Unmarshal(line, &rec); err != nil { + t.Fatalf("json: %v\n%s", err, line) + } + break + } + if rec == nil { + t.Fatalf("no hit line in:\n%s", buf.String()) + } + if rec["request_id"] != "hit-req" { + t.Errorf("request_id = %v", rec["request_id"]) + } + if rec["trace_id"] != "0af7651916cd43dd8448eb211c80319c" || rec["span_id"] != "b7ad6b7169203331" { + t.Errorf("trace = %v %v", rec["trace_id"], rec["span_id"]) + } + if rec["decision"] != "pass-pow" { + t.Errorf("decision = %v", rec["decision"]) + } +} diff --git a/internal/gate/pay.go b/internal/gate/pay.go index c67015b..adac334 100644 --- a/internal/gate/pay.go +++ b/internal/gate/pay.go @@ -2,6 +2,7 @@ package gate import ( "cmp" + "context" "encoding/base64" "encoding/json" "fmt" @@ -247,7 +248,7 @@ func (g *Gate) servePaymentRequired(w http.ResponseWriter, r *http.Request, if enc, err := doc.Encode(); err == nil { w.Header().Set(payment.HeaderRequired, enc) } else { - g.lg.Error("encoding the payment offer", "err", err) + g.lg.ErrorContext(r.Context(), "encoding the payment offer", "err", err) } if retryAfter > 0 { w.Header().Set("Retry-After", fmt.Sprint(int(retryAfter.Seconds()))) @@ -295,7 +296,7 @@ func (g *Gate) servePaymentPending(w http.ResponseWriter, r *http.Request, res p }); err == nil { w.Header().Set(payment.HeaderResponse, b64Std(enc)) } else { - g.lg.Error("encoding the pending settlement receipt", "err", err) + g.lg.ErrorContext(r.Context(), "encoding the pending settlement receipt", "err", err) } retry := res.RetryAfter if retry < time.Second { @@ -430,7 +431,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR } p, err := payment.DecodePayload(header) if err != nil { - g.lg.Debug("malformed payment presentation", "err", err) + g.lg.DebugContext(r.Context(), "malformed payment presentation", "err", err) g.servePaymentRequired(w, r, rule, "PAYMENT-SIGNATURE could not be decoded", 0) return decisionPayMalformed } @@ -439,7 +440,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR req, err := payment.MatchRail(p, doc.Accepts) if err != nil { // The gate never offered this rail. Refused locally, no egress. - g.lg.Debug("payment presented on an unoffered rail", + g.lg.DebugContext(r.Context(), "payment presented on an unoffered rail", "network", p.Network(), "scheme", p.Scheme()) g.servePaymentRequired(w, r, rule, "payment presented on a network or scheme this resource does not accept", 0) @@ -454,7 +455,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR // authorization identifies this payment cannot be chosen by the payer. id, err := payment.ID(p, req) if err != nil { - g.lg.Debug("payment presentation could not be identified", "err", err) + g.lg.DebugContext(r.Context(), "payment presentation could not be identified", "err", err) g.servePaymentRequired(w, r, rule, "the authorization in this payment names no payer and nonce to identify it by", 0) return decisionPayUnidentified @@ -462,7 +463,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR state, lease, prior, err := g.grants.Begin(id) if err != nil { - g.lg.Error("payment state unavailable", "pay_id", id[:16], "err", err) + g.lg.ErrorContext(r.Context(), "payment state unavailable", "pay_id", id[:16], "err", err) g.servePaymentRequired(w, r, rule, "payment recovery state is temporarily unavailable; retry the same payment", 5*time.Second) return decisionPayStateUnavailable @@ -476,7 +477,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR } return g.servePaidGrant(w, r, rule, id, prior, true) case payment.BeginSpent: - g.lg.Debug("payment entitlement expired", "pay_id", id[:16]) + g.lg.DebugContext(r.Context(), "payment entitlement expired", "pay_id", id[:16]) g.servePaymentRequired(w, r, rule, "this payment has already been used; sign a fresh one", 0) return decisionPayReplay @@ -487,7 +488,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR } defer func() { if err := g.grants.Release(id, lease); err != nil { - g.lg.Error("releasing payment reservation", "pay_id", id[:16], "err", err) + g.lg.ErrorContext(r.Context(), "releasing payment reservation", "pay_id", id[:16], "err", err) } }() @@ -504,7 +505,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR limitKey = ip.String() } if !g.payLimit.Allow(limitKey) { - g.lg.Warn("payment presentation rate limited", "client", limitKey) + g.lg.WarnContext(r.Context(), "payment presentation rate limited", "client", limitKey) g.servePaymentRequired(w, r, rule, "too many payment attempts; slow down and retry", 10*time.Second) return decisionPayRateLimited @@ -516,7 +517,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR // New() and MatchRail only returns offered rails — so reaching this is // a gate bug, not a payment problem. Still 402, never 500: the free // door remains the honest answer while the bug is fixed. - g.lg.Error("no verifier for an offered rail — gate construction bug", "network", req.Network) + g.lg.ErrorContext(r.Context(), "no verifier for an offered rail — gate construction bug", "network", req.Network) g.servePaymentRequired(w, r, rule, "the gate cannot reach a settlement service right now; retry shortly or use the free challenge", 5*time.Second) return decisionPayInfra @@ -536,16 +537,16 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR ExpiresAt: exp.Unix(), }) if err != nil { - g.lg.Error("persisting a settled payment grant", "pay_id", id, "err", err) + g.lg.ErrorContext(r.Context(), "persisting a settled payment grant", "pay_id", id, "err", err) g.servePaymentRequired(w, r, rule, ambiguousGrantFailure, 2*time.Second) return decisionPayGrantFailed } if grant.Scope != route.scope || grant.Audience != requestAudience(r) { - g.lg.Error("settled payment raced with a different durable grant", "pay_id", id) + g.lg.ErrorContext(r.Context(), "settled payment raced with a different durable grant", "pay_id", id) g.servePaymentRequired(w, r, rule, ambiguousGrantFailure, 2*time.Second) return decisionPayGrantConflict } - g.countPaidValue(grant.Amount, req.Network, rule) + g.countPaidValue(r.Context(), grant.Amount, req.Network, rule) return g.servePaidGrant(w, r, rule, id, grant, false) case payment.Invalid: @@ -553,7 +554,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR if reason == "" { reason = "payment rejected" } - g.lg.Warn("payment rejected by the facilitator", + g.lg.WarnContext(r.Context(), "payment rejected by the facilitator", "reason", reason, "payer", res.Payer, "scope", rule.Name) g.servePaymentRequired(w, r, rule, reason, 0) return decisionPayRejected @@ -564,7 +565,7 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR // path — but unlike an ambiguous settle there is a transaction to name, // and naming it is the difference between "reconcile this" and "pay // again". - g.lg.Warn("settlement pending — the payer holds a broadcast transaction", + g.lg.WarnContext(r.Context(), "settlement pending — the payer holds a broadcast transaction", "pay_id", id[:16], "tx", res.Tx, "network", res.Network, "payer", res.Payer, "scope", rule.Name) g.servePaymentPending(w, r, res) @@ -577,13 +578,13 @@ func (g *Gate) servePayment(w http.ResponseWriter, r *http.Request, route *paidR } // Serve nothing, claim nothing. The payer's own retry is the recovery // path, and the payment ID is quotable because both sides can compute it. - g.lg.Error("settle ambiguous — reconcile against the facilitator or chain", + g.lg.ErrorContext(r.Context(), "settle ambiguous — reconcile against the facilitator or chain", "pay_id", id, "payer", res.Payer, "scope", rule.Name, "err", res.Err) g.servePaymentRequired(w, r, rule, res.Reason, retry) return decisionPayAmbiguous default: // Indeterminate - g.lg.Warn("facilitator unavailable; payments degraded, free path unaffected", + g.lg.WarnContext(r.Context(), "facilitator unavailable; payments degraded, free path unaffected", "err", res.Err, "scope", rule.Name) // Tell the client how long the gate will actually refuse to try. Advising // a retry sooner guarantees a wasted round trip, and synchronises every @@ -615,10 +616,10 @@ func (g *Gate) servePaidGrant(w http.ResponseWriter, r *http.Request, rule *conf if err := g.setPassCookie(w, r, token.Pass{ Kind: token.KindPaid, Scope: grant.Scope, - Payer: g.settlementField("payer", grant.Payer), - Tx: g.settlementField("tx", grant.Transaction), + Payer: g.settlementField(r.Context(), "payer", grant.Payer), + Tx: g.settlementField(r.Context(), "tx", grant.Transaction), }, exp, time.Time{}); err != nil { - g.lg.Error("minting a durable paid pass", "pay_id", id, "err", err) + g.lg.ErrorContext(r.Context(), "minting a durable paid pass", "pay_id", id, "err", err) g.servePaymentRequired(w, r, rule, ambiguousGrantFailure, 2*time.Second) return decisionPayGrantFailed } @@ -628,7 +629,7 @@ func (g *Gate) servePaidGrant(w http.ResponseWriter, r *http.Request, rule *conf if recovered { message = "payment grant recovered" } - g.lg.Info(message, + g.lg.InfoContext(r.Context(), message, "pay_id", id[:16], "scope", rule.Name, "payer", grant.Payer, "tx", grant.Transaction, "amount", grant.Amount, "network", grant.Network) @@ -652,7 +653,7 @@ func (g *Gate) servePaidGrant(w http.ResponseWriter, r *http.Request, rule *conf // aggregate rails with different decimals. The facilitator-reported amount is // preferred (Verify already refused any mismatch); when x402's optional amount // is absent, the rail's asking price is what the facilitator settled against. -func (g *Gate) countPaidValue(amount, network string, rule *config.Rule) { +func (g *Gate) countPaidValue(ctx context.Context, amount, network string, rule *config.Rule) { atomic, ok := new(big.Int).SetString(amount, 10) if !ok || atomic.Sign() < 0 { atomic = rule.PriceAtomic[network] @@ -670,7 +671,7 @@ func (g *Gate) countPaidValue(amount, network string, rule *config.Rule) { micros := new(big.Int).Mul(atomic, big.NewInt(1_000_000)) micros.Quo(micros, new(big.Int).Exp(big.NewInt(10), big.NewInt(int64(decimals)), nil)) if !micros.IsUint64() { - g.lg.Warn("settled payment value overflows the metric; not counted", + g.lg.WarnContext(ctx, "settled payment value overflows the metric; not counted", "amount", amount, "network", network) return } @@ -694,11 +695,11 @@ const maxSettlementField = 128 // it costs the audit trail for that grant and nothing else, and a TRUNCATED // hash would be worse than an absent one — it looks like evidence and matches // no transaction. -func (g *Gate) settlementField(name, v string) string { +func (g *Gate) settlementField(ctx context.Context, name, v string) string { if len(v) <= maxSettlementField { return v } - g.lg.Warn("facilitator returned an oversized settlement field; omitting it from the pass", + g.lg.WarnContext(ctx, "facilitator returned an oversized settlement field; omitting it from the pass", "field", name, "bytes", len(v), "cap", maxSettlementField, "consequence", "this grant is not traceable from its pass; the settlement log line still has it") return "" diff --git a/internal/gate/request.go b/internal/gate/request.go index 53a377b..08b7b91 100644 --- a/internal/gate/request.go +++ b/internal/gate/request.go @@ -7,6 +7,7 @@ import ( "strings" "github.com/radiustechsystems/anteroom/internal/hosted" + "github.com/radiustechsystems/anteroom/internal/logging" ) // requestFacts is the immutable, header-derived view of a request. Parsing it @@ -32,6 +33,7 @@ type gateRequest struct { } func (g *Gate) inspect(r *http.Request) *gateRequest { + r = logging.BindRequest(r) ip, err := g.match.ClientIP(r) if err != nil { ip = netip.Addr{} diff --git a/internal/logging/context.go b/internal/logging/context.go new file mode 100644 index 0000000..069d641 --- /dev/null +++ b/internal/logging/context.go @@ -0,0 +1,128 @@ +package logging + +import ( + "context" + "crypto/rand" + "encoding/hex" + "net/http" + "strings" +) + +// Header names the gate reads and, for X-Request-ID, restates on the request +// so the upstream sees the same identifier the logs carry. +const ( + HeaderRequestID = "X-Request-ID" + HeaderCorrelationID = "X-Correlation-ID" + HeaderTraceparent = "Traceparent" +) + +const maxRequestIDLen = 128 + +type ctxKey struct{} + +// Fields are the request-scoped labels attached to every log line that is +// written with this context. RequestID is always set. TraceID and SpanID are +// set only when the inbound request carried a valid W3C traceparent — the gate +// propagates traces, it does not mint spans it never exports. +type Fields struct { + RequestID string + TraceID string + SpanID string +} + +// FromContext returns the Fields bound by Bind, if any. +func FromContext(ctx context.Context) (Fields, bool) { + f, ok := ctx.Value(ctxKey{}).(Fields) + return f, ok +} + +// Bind stores f on ctx for the contextHandler to copy onto log records. +func Bind(ctx context.Context, f Fields) context.Context { + return context.WithValue(ctx, ctxKey{}, f) +} + +// BindRequest extracts Fields from r, generates a request ID if none was +// supplied, restates X-Request-ID on the request (so the reverse proxy +// forwards the identifier the logs use), and returns the request with those +// Fields on its context. +func BindRequest(r *http.Request) *http.Request { + f := FieldsFromRequest(r) + r.Header.Set(HeaderRequestID, f.RequestID) + return r.WithContext(Bind(r.Context(), f)) +} + +// FieldsFromRequest reads inbound correlation headers. An absent or unusable +// request ID is replaced with a generated one; a missing traceparent is left +// empty rather than invented. +func FieldsFromRequest(r *http.Request) Fields { + id := canonicalRequestID(r.Header.Get(HeaderRequestID)) + if id == "" { + id = canonicalRequestID(r.Header.Get(HeaderCorrelationID)) + } + if id == "" { + id = newRequestID() + } + traceID, spanID := parseTraceparent(r.Header.Get(HeaderTraceparent)) + return Fields{RequestID: id, TraceID: traceID, SpanID: spanID} +} + +func canonicalRequestID(s string) string { + s = strings.TrimSpace(s) + if s == "" || len(s) > maxRequestIDLen { + return "" + } + for i := 0; i < len(s); i++ { + c := s[i] + switch { + case c >= 'a' && c <= 'z', c >= 'A' && c <= 'Z', c >= '0' && c <= '9': + case c == '-', c == '_', c == '.': + default: + return "" + } + } + return s +} + +func newRequestID() string { + var b [16]byte + if _, err := rand.Read(b[:]); err != nil { + // crypto/rand failure is unheard of for 16 bytes; a degenerate + // fallback still gives every request a distinct-enough join key. + return "anteroom-unrandom" + } + return hex.EncodeToString(b[:]) +} + +// parseTraceparent reads a W3C traceparent. Version 00 layout is +// `{ver}-{32 hex trace}-{16 hex span}-{2 hex flags}`. All-zero ids are +// invalid. Future versions with the same four-field layout are accepted. +func parseTraceparent(s string) (traceID, spanID string) { + s = strings.TrimSpace(s) + parts := strings.Split(s, "-") + if len(parts) != 4 { + return "", "" + } + ver, tid, sid, flags := parts[0], parts[1], parts[2], parts[3] + if len(ver) != 2 || len(tid) != 32 || len(sid) != 16 || len(flags) != 2 { + return "", "" + } + if !isHex(ver) || !isHex(tid) || !isHex(sid) || !isHex(flags) { + return "", "" + } + if tid == "00000000000000000000000000000000" || sid == "00000000000000000000000000000000" { + return "", "" + } + return strings.ToLower(tid), strings.ToLower(sid) +} + +func isHex(s string) bool { + for i := 0; i < len(s); i++ { + c := s[i] + switch { + case c >= '0' && c <= '9', c >= 'a' && c <= 'f', c >= 'A' && c <= 'F': + default: + return false + } + } + return true +} diff --git a/internal/logging/context_test.go b/internal/logging/context_test.go new file mode 100644 index 0000000..d0e6ad1 --- /dev/null +++ b/internal/logging/context_test.go @@ -0,0 +1,86 @@ +package logging + +import ( + "net/http" + "net/http/httptest" + "testing" +) + +func TestFieldsFromRequestGeneratesID(t *testing.T) { + r := httptest.NewRequest(http.MethodGet, "/", nil) + f := FieldsFromRequest(r) + if len(f.RequestID) != 32 { + t.Errorf("generated request_id length = %d, want 32 hex chars, got %q", len(f.RequestID), f.RequestID) + } + if f.TraceID != "" || f.SpanID != "" { + t.Errorf("invented a trace: %+v", f) + } +} + +func TestFieldsFromRequestHonorsInboundID(t *testing.T) { + r := httptest.NewRequest(http.MethodGet, "/", nil) + r.Header.Set(HeaderRequestID, "abc-123") + f := FieldsFromRequest(r) + if f.RequestID != "abc-123" { + t.Errorf("request_id = %q", f.RequestID) + } +} + +func TestFieldsFromRequestFallsBackToCorrelationID(t *testing.T) { + r := httptest.NewRequest(http.MethodGet, "/", nil) + r.Header.Set(HeaderCorrelationID, "corr.1") + f := FieldsFromRequest(r) + if f.RequestID != "corr.1" { + t.Errorf("request_id = %q", f.RequestID) + } +} + +func TestFieldsFromRequestRejectsGarbageID(t *testing.T) { + r := httptest.NewRequest(http.MethodGet, "/", nil) + r.Header.Set(HeaderRequestID, "not a valid id\nwith newline") + f := FieldsFromRequest(r) + if f.RequestID == "not a valid id\nwith newline" || f.RequestID == "" { + t.Errorf("garbage id should be replaced, got %q", f.RequestID) + } +} + +func TestParseTraceparent(t *testing.T) { + const raw = "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01" + tid, sid := parseTraceparent(raw) + if tid != "0af7651916cd43dd8448eb211c80319c" || sid != "b7ad6b7169203331" { + t.Errorf("got %s %s", tid, sid) + } + tid, sid = parseTraceparent("00-00000000000000000000000000000000-b7ad6b7169203331-01") + if tid != "" || sid != "" { + t.Errorf("all-zero trace-id should be rejected, got %s %s", tid, sid) + } + tid, sid = parseTraceparent("not-a-traceparent") + if tid != "" || sid != "" { + t.Errorf("malformed should be rejected, got %s %s", tid, sid) + } +} + +func TestBindRequestRestatesHeader(t *testing.T) { + r := httptest.NewRequest(http.MethodGet, "/", nil) + r = BindRequest(r) + id := r.Header.Get(HeaderRequestID) + if id == "" { + t.Fatal("BindRequest did not set X-Request-ID") + } + f, ok := FromContext(r.Context()) + if !ok || f.RequestID != id { + t.Errorf("context Fields = %+v, header %q", f, id) + } + + r2 := httptest.NewRequest(http.MethodGet, "/", nil) + r2.Header.Set(HeaderRequestID, "keep-me") + r2.Header.Set(HeaderTraceparent, "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01") + r2 = BindRequest(r2) + if r2.Header.Get(HeaderRequestID) != "keep-me" { + t.Errorf("inbound id overwritten: %q", r2.Header.Get(HeaderRequestID)) + } + f, _ = FromContext(r2.Context()) + if f.TraceID != "0af7651916cd43dd8448eb211c80319c" || f.SpanID != "b7ad6b7169203331" { + t.Errorf("traceparent not parsed: %+v", f) + } +} diff --git a/internal/logging/logging.go b/internal/logging/logging.go new file mode 100644 index 0000000..6270e52 --- /dev/null +++ b/internal/logging/logging.go @@ -0,0 +1,114 @@ +// Package logging builds the process slog.Logger from [log] config and +// attaches request-scoped fields (request_id, trace_id, span_id) from +// context onto every record logged with that context. +// +// Format is a Handler choice. Call sites already emit slog key/value attrs, so +// json, logfmt, and text are encodings of the same records. +package logging + +import ( + "context" + "io" + "log/slog" + "os" + "sort" + "strings" + "time" + + "github.com/radiustechsystems/anteroom/internal/config" +) + +// New builds the process logger. verbose (the -v flag) forces debug regardless +// of cfg.Level, so per-request hit lines stay an explicit opt-in. +func New(w io.Writer, cfg config.Log, verbose bool) *slog.Logger { + if w == nil { + w = os.Stderr + } + level := parseLevel(cfg.Level) + if verbose { + level = slog.LevelDebug + } + opts := &slog.HandlerOptions{Level: level} + format := strings.ToLower(strings.TrimSpace(cfg.Format)) + if format == "logfmt" { + opts.ReplaceAttr = logfmtReplace + } + var h slog.Handler + switch format { + case "json": + h = slog.NewJSONHandler(w, opts) + default: // text, logfmt + h = slog.NewTextHandler(w, opts) + } + lg := slog.New(contextHandler{h}) + if len(cfg.Labels) == 0 { + return lg + } + keys := make([]string, 0, len(cfg.Labels)) + for k := range cfg.Labels { + keys = append(keys, k) + } + sort.Strings(keys) + args := make([]any, 0, len(keys)*2) + for _, k := range keys { + args = append(args, k, cfg.Labels[k]) + } + return lg.With(args...) +} + +func parseLevel(s string) slog.Level { + switch strings.ToLower(strings.TrimSpace(s)) { + case "debug": + return slog.LevelDebug + case "warn", "warning": + return slog.LevelWarn + case "error": + return slog.LevelError + default: + return slog.LevelInfo + } +} + +// logfmtReplace makes slog's text handler a collector-friendly logfmt: +// ts=RFC3339Nano (UTC) and lowercase level. Field names otherwise match text. +func logfmtReplace(_ []string, a slog.Attr) slog.Attr { + switch a.Key { + case slog.TimeKey: + a.Key = "ts" + if t, ok := a.Value.Any().(time.Time); ok { + a.Value = slog.StringValue(t.UTC().Format(time.RFC3339Nano)) + } + case slog.LevelKey: + a.Value = slog.StringValue(strings.ToLower(a.Value.String())) + } + return a +} + +// contextHandler copies request-scoped Fields off ctx onto every Record. +// Call sites must use the *Context slog methods or the fields stay off the line. +type contextHandler struct { + slog.Handler +} + +func (h contextHandler) Handle(ctx context.Context, r slog.Record) error { + if f, ok := FromContext(ctx); ok { + if f.RequestID != "" { + r.AddAttrs(slog.String("request_id", f.RequestID)) + } + if f.TraceID != "" { + r.AddAttrs(slog.String("trace_id", f.TraceID)) + } + if f.SpanID != "" { + r.AddAttrs(slog.String("span_id", f.SpanID)) + } + } + return h.Handler.Handle(ctx, r) +} + +func (h contextHandler) WithAttrs(attrs []slog.Attr) slog.Handler { + return contextHandler{Handler: h.Handler.WithAttrs(attrs)} +} + +func (h contextHandler) WithGroup(name string) slog.Handler { + return contextHandler{Handler: h.Handler.WithGroup(name)} +} diff --git a/internal/logging/logging_test.go b/internal/logging/logging_test.go new file mode 100644 index 0000000..dcc4565 --- /dev/null +++ b/internal/logging/logging_test.go @@ -0,0 +1,135 @@ +package logging + +import ( + "bytes" + "context" + "encoding/json" + "strings" + "testing" + + "github.com/radiustechsystems/anteroom/internal/config" +) + +func TestNewJSON(t *testing.T) { + var buf bytes.Buffer + lg := New(&buf, config.Log{Level: "info", Format: "json"}, false) + lg.Info("anteroom is up", "listen", "127.0.0.1:8080") + + var rec map[string]any + if err := json.Unmarshal(buf.Bytes(), &rec); err != nil { + t.Fatalf("json: %v\n%s", err, buf.String()) + } + if rec["msg"] != "anteroom is up" { + t.Errorf("msg = %v", rec["msg"]) + } + if rec["listen"] != "127.0.0.1:8080" { + t.Errorf("listen = %v", rec["listen"]) + } + if rec["level"] != "INFO" { + t.Errorf("level = %v", rec["level"]) + } +} + +func TestNewLogfmt(t *testing.T) { + var buf bytes.Buffer + lg := New(&buf, config.Log{Level: "info", Format: "logfmt"}, false) + lg.Info("anteroom is up", "listen", "127.0.0.1:8080") + line := buf.String() + for _, want := range []string{"ts=", "level=info", `msg="anteroom is up"`, "listen=127.0.0.1:8080"} { + if !strings.Contains(line, want) { + t.Errorf("logfmt line lacks %q:\n%s", want, line) + } + } + if strings.Contains(line, "time=") { + t.Errorf("logfmt still used time= rather than ts=:\n%s", line) + } +} + +func TestNewText(t *testing.T) { + var buf bytes.Buffer + lg := New(&buf, config.Log{Level: "info", Format: "text"}, false) + lg.Info("hit", "decision", "pass-pow") + line := buf.String() + if !strings.Contains(line, "msg=hit") || !strings.Contains(line, "decision=pass-pow") { + t.Errorf("text line:\n%s", line) + } +} + +func TestVerboseForcesDebug(t *testing.T) { + var buf bytes.Buffer + lg := New(&buf, config.Log{Level: "error", Format: "text"}, true) + lg.Debug("hit", "decision", "wait-page") + if !strings.Contains(buf.String(), "msg=hit") { + t.Errorf("-v should force debug, got:\n%s", buf.String()) + } +} + +func TestLabelsAreStableAndPresent(t *testing.T) { + var buf bytes.Buffer + lg := New(&buf, config.Log{ + Level: "info", + Format: "json", + Labels: map[string]string{"service": "anteroom", "cluster": "prod"}, + }, false) + lg.Info("up") + var rec map[string]any + if err := json.Unmarshal(buf.Bytes(), &rec); err != nil { + t.Fatal(err) + } + if rec["service"] != "anteroom" || rec["cluster"] != "prod" { + t.Errorf("labels missing: %v", rec) + } +} + +func TestContextFieldsAttachOnlyWithContextMethods(t *testing.T) { + var buf bytes.Buffer + lg := New(&buf, config.Log{Level: "info", Format: "json"}, false) + ctx := Bind(context.Background(), Fields{ + RequestID: "req-1", + TraceID: "0af7651916cd43dd8448eb211c80319c", + SpanID: "b7ad6b7169203331", + }) + lg.Info("without") + lg.InfoContext(ctx, "with") + + lines := bytes.Split(bytes.TrimSpace(buf.Bytes()), []byte("\n")) + if len(lines) != 2 { + t.Fatalf("want 2 lines, got %d:\n%s", len(lines), buf.String()) + } + var without, with map[string]any + if err := json.Unmarshal(lines[0], &without); err != nil { + t.Fatal(err) + } + if err := json.Unmarshal(lines[1], &with); err != nil { + t.Fatal(err) + } + if _, ok := without["request_id"]; ok { + t.Errorf("Info without context still has request_id: %v", without) + } + if with["request_id"] != "req-1" || with["trace_id"] != "0af7651916cd43dd8448eb211c80319c" || with["span_id"] != "b7ad6b7169203331" { + t.Errorf("InfoContext missing fields: %v", with) + } +} + +func TestWithPreservesContextHandler(t *testing.T) { + var buf bytes.Buffer + lg := New(&buf, config.Log{Level: "info", Format: "json", Labels: map[string]string{"service": "anteroom"}}, false) + ctx := Bind(context.Background(), Fields{RequestID: "req-2"}) + lg.With("listen", ":8080").InfoContext(ctx, "up") + var rec map[string]any + if err := json.Unmarshal(buf.Bytes(), &rec); err != nil { + t.Fatal(err) + } + if rec["service"] != "anteroom" || rec["listen"] != ":8080" || rec["request_id"] != "req-2" { + t.Errorf("child logger dropped attrs: %v", rec) + } +} + +func TestQuietByDefaultHidesDebug(t *testing.T) { + var buf bytes.Buffer + lg := New(&buf, config.Log{Level: "info", Format: "text"}, false) + lg.Debug("hit") + if buf.Len() != 0 { + t.Errorf("debug leaked at info: %s", buf.String()) + } +} diff --git a/internal/payment/callback.go b/internal/payment/callback.go index c434aaf..1332820 100644 --- a/internal/payment/callback.go +++ b/internal/payment/callback.go @@ -145,7 +145,7 @@ func (c *CallbackVerifier) Verify(ctx context.Context, p Payload, req Requiremen case c.egress <- struct{}{}: defer func() { <-c.egress }() default: - c.lg.Warn("facilitator egress saturated; shedding a payment presentation", + c.lg.WarnContext(ctx, "facilitator egress saturated; shedding a payment presentation", "in_flight", len(c.egress), "cap", maxInFlight, "consequence", "the client is told to retry; the free path is unaffected") return Result{ @@ -162,7 +162,7 @@ func (c *CallbackVerifier) Verify(ctx context.Context, p Payload, req Requiremen if err != nil { if res.transport { if opened := c.breaker.Fail(); opened { - c.lg.Error("facilitator circuit breaker opened", + c.lg.ErrorContext(ctx, "facilitator circuit breaker opened", "consequence", "payment presentations skip egress and fall back until it closes", "free_path", "unaffected") } @@ -184,7 +184,7 @@ func (c *CallbackVerifier) Verify(ctx context.Context, p Payload, req Requiremen if err != nil { if res.transport { if opened := c.breaker.Fail(); opened { - c.lg.Error("facilitator circuit breaker opened", + c.lg.ErrorContext(ctx, "facilitator circuit breaker opened", "consequence", "payment presentations skip egress and fall back until it closes") } } @@ -210,7 +210,7 @@ func (c *CallbackVerifier) Verify(ctx context.Context, p Payload, req Requiremen // broadcast leaves nothing to reconcile against. That is the // ambiguous case, with the same recovery: re-present the // identical payload and let settle idempotency answer. - c.lg.Error("facilitator reported settlement_pending with no transaction", + c.lg.ErrorContext(ctx, "facilitator reported settlement_pending with no transaction", "payer", sr.Payer, "network", sr.Network, "action", "treating as ambiguous; the payer may retry the identical payload") return Result{Verdict: Ambiguous, Reason: ambiguousReason, @@ -220,7 +220,7 @@ func (c *CallbackVerifier) Verify(ctx context.Context, p Payload, req Requiremen return Result{Verdict: Ambiguous, Reason: ambiguousReason, Err: fmt.Errorf("pending settlement reported network %q, requested %q", sr.Network, req.Network)} } - c.lg.Warn("settlement pending — broadcast, confirmation unknown", + c.lg.WarnContext(ctx, "settlement pending — broadcast, confirmation unknown", "tx", sr.Transaction, "network", sr.Network, "payer", sr.Payer, "action", "no pass minted, payment not claimed; the payer reconciles and re-presents") return Result{ @@ -233,7 +233,7 @@ func (c *CallbackVerifier) Verify(ctx context.Context, p Payload, req Requiremen if sr.Transaction == "" { // A settlement claim with no evidence must never mint a pass that // embeds a transaction. Ambiguous, never success. - c.lg.Error("facilitator claimed settle success without a transaction hash", + c.lg.ErrorContext(ctx, "facilitator claimed settle success without a transaction hash", "payer", sr.Payer, "network", sr.Network, "action", "treating as ambiguous; the payer may retry the identical payload") return Result{Verdict: Ambiguous, Reason: ambiguousReason, @@ -245,7 +245,7 @@ func (c *CallbackVerifier) Verify(ctx context.Context, p Payload, req Requiremen // gate asked for, whatever it says about success. Ambiguous rather than // invalid: something may well have moved on that other chain, and the // operator needs to know rather than have the payer told to sign again. - c.lg.Error("facilitator settled on a different network than requested", + c.lg.ErrorContext(ctx, "facilitator settled on a different network than requested", "requested", req.Network, "settled", sr.Network, "tx", sr.Transaction, "payer", sr.Payer, "action", "no pass minted; reconcile against the chain") return Result{Verdict: Ambiguous, Reason: ambiguousReason, @@ -332,7 +332,7 @@ func (c *CallbackVerifier) call(ctx context.Context, path string, p Payload, req if err != nil { cancel() last = fmt.Errorf("%s: %w", path, err) - c.lg.Warn("facilitator call failed", + c.lg.WarnContext(ctx, "facilitator call failed", "endpoint", path, "attempt", attempt, "elapsed", time.Since(start), "budget", budget, "err", err) if attempt < attempts && retryable(err) { @@ -341,7 +341,7 @@ func (c *CallbackVerifier) call(ctx context.Context, path string, p Payload, req return callResult{transport: true}, last } - res, err := c.decode(path, resp, out) + res, err := c.decode(ctx, path, resp, out) cancel() if err != nil { return res, err @@ -353,26 +353,26 @@ func (c *CallbackVerifier) call(ctx context.Context, path string, p Payload, req // decode turns a facilitator response into either a decoded document or a // classified failure. -func (c *CallbackVerifier) decode(path string, resp *http.Response, out any) (callResult, error) { +func (c *CallbackVerifier) decode(ctx context.Context, path string, resp *http.Response, out any) (callResult, error) { defer resp.Body.Close() raw, readErr := io.ReadAll(io.LimitReader(resp.Body, 64<<10)) switch { case resp.StatusCode >= 300 && resp.StatusCode < 400: loc := resp.Header.Get("Location") - c.lg.Error("facilitator redirected", + c.lg.ErrorContext(ctx, "facilitator redirected", "endpoint", path, "status", resp.StatusCode, "location", loc, "fix", "set the canonical facilitator URL in anteroom.toml; redirects are never followed on a POST") return callResult{transport: true}, fmt.Errorf("%s: redirected to %q", path, loc) case resp.StatusCode == http.StatusTooManyRequests: - c.lg.Warn("facilitator rate limited us", "endpoint", path, + c.lg.WarnContext(ctx, "facilitator rate limited us", "endpoint", path, "retry_after", resp.Header.Get("Retry-After")) return callResult{transport: true}, fmt.Errorf("%s: rate limited", path) case resp.StatusCode >= 500: - c.lg.Warn("facilitator server error", "endpoint", path, + c.lg.WarnContext(ctx, "facilitator server error", "endpoint", path, "status", resp.StatusCode, "body", snippet(raw)) return callResult{transport: true}, fmt.Errorf("%s: status %d", path, resp.StatusCode) @@ -385,7 +385,7 @@ func (c *CallbackVerifier) decode(path string, resp *http.Response, out any) (ca if err := json.Unmarshal(raw, out); err != nil { // A 4xx that does not decode is still a failure of this call, and on // /settle it is ambiguous — which the caller decides, not us. - c.lg.Warn("facilitator response did not decode", + c.lg.WarnContext(ctx, "facilitator response did not decode", "endpoint", path, "status", resp.StatusCode, "body", snippet(raw), "err", err) return callResult{}, fmt.Errorf("%s: malformed response: %w", path, err) }