#17 Observability: OpenTelemetry logging + tracing + alerting #45
@@ -14,12 +14,15 @@
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"mime"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/telemetry"
|
||||
)
|
||||
|
||||
// MaxBodyBytes is the hard cap on the request body: 256 KiB (262,144 bytes),
|
||||
@@ -54,7 +57,34 @@ func NewHandler(sink Sink) *Handler {
|
||||
// sink failure -> 503
|
||||
//
|
||||
// Error bodies are generic ({"error":"..."}) and never reflect request content.
|
||||
//
|
||||
// Observability (#17): when a telemetry provider rides in the request context
|
||||
// (injected by the Worker), each request emits an "ingest.request" span and a
|
||||
// correlated structured log classifying the outcome (accepted / rejected /
|
||||
// error). Instrumentation is observe-only — it wraps the response writer to read
|
||||
// the status code and never alters the HTTP contract above — and is skipped
|
||||
// entirely (zero overhead, identical behaviour) when no provider is present,
|
||||
// which is the default until an OTLP endpoint is configured.
|
||||
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
tel := telemetry.FromContext(r.Context())
|
||||
if !tel.Enabled() {
|
||||
h.serve(w, r)
|
||||
return
|
||||
}
|
||||
ctx, span := tel.StartSpan(r.Context(), "ingest.request",
|
||||
telemetry.WithSpanKind(telemetry.SpanKindServer),
|
||||
telemetry.WithAttributes(telemetry.String("http.request.method", r.Method)))
|
||||
defer span.End()
|
||||
|
||||
sw := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
|
||||
h.serve(sw, r.WithContext(ctx))
|
||||
finishIngestSpan(ctx, tel, span, sw.status)
|
||||
}
|
||||
|
||||
// serve is the transport contract from ServeHTTP's doc, unchanged. It is split
|
||||
// out so ServeHTTP can wrap it with observe-only instrumentation without
|
||||
// touching a single response path.
|
||||
func (h *Handler) serve(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
w.Header().Set("Allow", http.MethodPost)
|
||||
writeError(w, http.StatusMethodNotAllowed, "method not allowed")
|
||||
@@ -129,6 +159,110 @@ func writeError(w http.ResponseWriter, status int, reason string) {
|
||||
writeJSON(w, status, map[string]string{"error": reason})
|
||||
}
|
||||
|
||||
// statusRecorder wraps an http.ResponseWriter to capture the status code the
|
||||
// handler wrote, so the ingest span/log can record the outcome. It is a pure
|
||||
// pass-through: it changes no bytes and no headers.
|
||||
type statusRecorder struct {
|
||||
http.ResponseWriter
|
||||
status int
|
||||
wrote bool
|
||||
}
|
||||
|
||||
func (s *statusRecorder) WriteHeader(code int) {
|
||||
if !s.wrote {
|
||||
s.status = code
|
||||
s.wrote = true
|
||||
}
|
||||
s.ResponseWriter.WriteHeader(code)
|
||||
}
|
||||
|
||||
func (s *statusRecorder) Write(b []byte) (int, error) {
|
||||
if !s.wrote {
|
||||
// An implicit 200 (no explicit WriteHeader) — not a path this handler
|
||||
// takes, but recorded correctly for completeness.
|
||||
s.status = http.StatusOK
|
||||
s.wrote = true
|
||||
}
|
||||
return s.ResponseWriter.Write(b)
|
||||
}
|
||||
|
||||
// Ingest outcome vocabulary, recorded as the "ingest.outcome" attribute so a
|
||||
// backend can group/alert by it.
|
||||
const (
|
||||
outcomeAccepted = "accepted"
|
||||
outcomeRejected = "rejected"
|
||||
outcomeRateLimited = "rate_limited"
|
||||
outcomeError = "error"
|
||||
)
|
||||
|
||||
// alertIngestError is the alert.type for the ingest error signal (issue #17):
|
||||
// the "elevated ingest error rate" a backend alert keys on. It marks the 5xx
|
||||
// (e.g. 503 storage-unavailable) responses that indicate a Worker-side fault, as
|
||||
// opposed to ordinary client rejections (4xx), which are expected and not
|
||||
// alerted.
|
||||
const alertIngestError = "ingest.error"
|
||||
|
||||
// finishIngestSpan classifies the written status and records the span status,
|
||||
// attributes, and a correlated log record. Client rejections (4xx) are recorded
|
||||
// at INFO; server errors (5xx) set the span to Error and emit the alertable
|
||||
// ingest.error signal.
|
||||
func finishIngestSpan(ctx context.Context, tel *telemetry.Telemetry, span *telemetry.Span, status int) {
|
||||
outcome, reason := classifyIngest(status)
|
||||
span.SetAttributes(
|
||||
telemetry.Int("http.response.status_code", status),
|
||||
telemetry.String("ingest.outcome", outcome),
|
||||
)
|
||||
if reason != "" {
|
||||
span.SetAttributes(telemetry.String("ingest.reason", reason))
|
||||
}
|
||||
|
||||
switch outcome {
|
||||
case outcomeAccepted:
|
||||
span.SetStatus(telemetry.StatusOK, "")
|
||||
tel.Info(ctx, "ingest request accepted",
|
||||
telemetry.String("ingest.outcome", outcome),
|
||||
telemetry.Int("http.response.status_code", status))
|
||||
case outcomeError:
|
||||
span.SetStatus(telemetry.StatusError, reason)
|
||||
attrs := append(telemetry.Alert(alertIngestError),
|
||||
telemetry.String("ingest.outcome", outcome),
|
||||
telemetry.String("ingest.reason", reason),
|
||||
telemetry.Int("http.response.status_code", status))
|
||||
tel.Error(ctx, "ingest request failed", attrs...)
|
||||
default: // rejected / rate_limited: expected client outcomes, not alerts.
|
||||
tel.Info(ctx, "ingest request "+outcome,
|
||||
telemetry.String("ingest.outcome", outcome),
|
||||
telemetry.String("ingest.reason", reason),
|
||||
telemetry.Int("http.response.status_code", status))
|
||||
}
|
||||
}
|
||||
|
||||
// classifyIngest maps an HTTP status to the ingest outcome vocabulary. Reasons
|
||||
// are derived at status granularity (the response paths never echo request
|
||||
// content, so neither do these).
|
||||
func classifyIngest(status int) (outcome, reason string) {
|
||||
switch {
|
||||
case status >= 200 && status < 300:
|
||||
return outcomeAccepted, ""
|
||||
case status == http.StatusTooManyRequests: // 429 — enforced at the edge today
|
||||
return outcomeRateLimited, "rate_limited"
|
||||
case status == http.StatusServiceUnavailable: // 503
|
||||
return outcomeError, "storage_unavailable"
|
||||
case status >= 500:
|
||||
return outcomeError, "server_error"
|
||||
case status == http.StatusRequestEntityTooLarge: // 413
|
||||
return outcomeRejected, "payload_too_large"
|
||||
case status == http.StatusUnsupportedMediaType: // 415
|
||||
return outcomeRejected, "unsupported_media_type"
|
||||
case status == http.StatusMethodNotAllowed: // 405
|
||||
return outcomeRejected, "method_not_allowed"
|
||||
case status >= 400:
|
||||
return outcomeRejected, "invalid_request"
|
||||
default:
|
||||
return outcomeAccepted, ""
|
||||
}
|
||||
}
|
||||
|
||||
// compile-time assurance that the stub sinks satisfy Sink.
|
||||
var (
|
||||
_ Sink = NopSink{}
|
||||
|
||||
@@ -0,0 +1,182 @@
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/telemetry"
|
||||
)
|
||||
|
||||
// doWithTelemetry runs one request through h with a telemetry provider injected
|
||||
// into the request context, returning the recorder and the exporter.
|
||||
func doWithTelemetry(t *testing.T, h http.Handler, sink Sink, method, contentType, body string) (*httptest.ResponseRecorder, *telemetry.MemoryExporter) {
|
||||
t.Helper()
|
||||
exp := telemetry.NewMemoryExporter()
|
||||
tel := telemetry.New(exp)
|
||||
|
||||
var r io.Reader
|
||||
if body != "" {
|
||||
r = strings.NewReader(body)
|
||||
}
|
||||
req := httptest.NewRequest(method, "/v1/reports", r)
|
||||
if contentType != "" {
|
||||
req.Header.Set("Content-Type", contentType)
|
||||
}
|
||||
req = req.WithContext(telemetry.NewContext(req.Context(), tel))
|
||||
rec := httptest.NewRecorder()
|
||||
h.ServeHTTP(rec, req)
|
||||
return rec, exp
|
||||
}
|
||||
|
||||
// requireSpan returns the single "ingest.request" span, failing if absent.
|
||||
func requireIngestSpan(t *testing.T, exp *telemetry.MemoryExporter) telemetry.SpanData {
|
||||
t.Helper()
|
||||
spans := exp.SpansByName("ingest.request")
|
||||
if len(spans) != 1 {
|
||||
t.Fatalf("got %d ingest.request spans, want 1", len(spans))
|
||||
}
|
||||
return spans[0]
|
||||
}
|
||||
|
||||
func attrString(t *testing.T, attrs []telemetry.KeyValue, key string) string {
|
||||
t.Helper()
|
||||
v, ok := telemetry.Attr(attrs, key)
|
||||
if !ok {
|
||||
t.Fatalf("attribute %q not present in %v", key, attrs)
|
||||
}
|
||||
s, ok := v.(string)
|
||||
if !ok {
|
||||
t.Fatalf("attribute %q = %v, not a string", key, v)
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
func TestIngestAcceptedEmitsSpanAndLog(t *testing.T) {
|
||||
rec, exp := doWithTelemetry(t, NewHandler(&MemorySink{}), nil, http.MethodPost, "application/json", validBody)
|
||||
if rec.Code != http.StatusAccepted {
|
||||
t.Fatalf("status = %d, want 202 (contract must be unchanged)", rec.Code)
|
||||
}
|
||||
span := requireIngestSpan(t, exp)
|
||||
if span.Kind != telemetry.SpanKindServer {
|
||||
t.Errorf("span kind = %d, want server", span.Kind)
|
||||
}
|
||||
if span.Status.Code != telemetry.StatusOK {
|
||||
t.Errorf("span status = %d, want OK", span.Status.Code)
|
||||
}
|
||||
if got := attrString(t, span.Attributes, "ingest.outcome"); got != "accepted" {
|
||||
t.Errorf("outcome = %q, want accepted", got)
|
||||
}
|
||||
if v, _ := telemetry.Attr(span.Attributes, "http.response.status_code"); v != int64(202) {
|
||||
t.Errorf("status_code attr = %v, want 202", v)
|
||||
}
|
||||
|
||||
logs := exp.Logs()
|
||||
if len(logs) != 1 {
|
||||
t.Fatalf("got %d logs, want 1", len(logs))
|
||||
}
|
||||
if logs[0].Severity != telemetry.SeverityInfo {
|
||||
t.Errorf("log severity = %v, want INFO", logs[0].Severity)
|
||||
}
|
||||
// The log must correlate to the span.
|
||||
if logs[0].SpanContext != span.SpanContext {
|
||||
t.Errorf("log not correlated to span: %+v vs %+v", logs[0].SpanContext, span.SpanContext)
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngestRejectedEmitsSpanAndLog(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
method string
|
||||
contentType string
|
||||
body string
|
||||
wantStatus int
|
||||
wantReason string
|
||||
}{
|
||||
{"wrong content type", http.MethodPost, "text/plain", validBody, http.StatusUnsupportedMediaType, "unsupported_media_type"},
|
||||
{"malformed json", http.MethodPost, "application/json", `{not json`, http.StatusBadRequest, "invalid_request"},
|
||||
{"wrong method", http.MethodGet, "application/json", validBody, http.StatusMethodNotAllowed, "method_not_allowed"},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
rec, exp := doWithTelemetry(t, NewHandler(&MemorySink{}), nil, tc.method, tc.contentType, tc.body)
|
||||
if rec.Code != tc.wantStatus {
|
||||
t.Fatalf("status = %d, want %d (contract must be unchanged)", rec.Code, tc.wantStatus)
|
||||
}
|
||||
span := requireIngestSpan(t, exp)
|
||||
if got := attrString(t, span.Attributes, "ingest.outcome"); got != "rejected" {
|
||||
t.Errorf("outcome = %q, want rejected", got)
|
||||
}
|
||||
if got := attrString(t, span.Attributes, "ingest.reason"); got != tc.wantReason {
|
||||
t.Errorf("reason = %q, want %q", got, tc.wantReason)
|
||||
}
|
||||
// A client rejection is NOT a span error and NOT an alert.
|
||||
if span.Status.Code == telemetry.StatusError {
|
||||
t.Errorf("rejection should not set span status Error")
|
||||
}
|
||||
logs := exp.Logs()
|
||||
if len(logs) != 1 || logs[0].Severity != telemetry.SeverityInfo {
|
||||
t.Fatalf("want 1 INFO log for a rejection, got %+v", logs)
|
||||
}
|
||||
if _, ok := telemetry.Attr(logs[0].Attributes, telemetry.AttrAlert); ok {
|
||||
t.Errorf("client rejection must not carry an alert marker")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// failSinkErr is a Sink that always fails, to drive the 503 error path.
|
||||
type failSinkErr struct{}
|
||||
|
||||
func (failSinkErr) Store(context.Context, []byte) error { return errors.New("boom") }
|
||||
|
||||
// TestIngestErrorEmitsAlertSignal is the "elevated ingest error" signal: a 503
|
||||
// storage failure sets the span to Error and emits an ERROR log carrying the
|
||||
// alert marker + alert.type=ingest.error that a backend alert keys on.
|
||||
func TestIngestErrorEmitsAlertSignal(t *testing.T) {
|
||||
rec, exp := doWithTelemetry(t, NewHandler(failSinkErr{}), nil, http.MethodPost, "application/json", validBody)
|
||||
if rec.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("status = %d, want 503", rec.Code)
|
||||
}
|
||||
span := requireIngestSpan(t, exp)
|
||||
if span.Status.Code != telemetry.StatusError {
|
||||
t.Errorf("span status = %d, want Error", span.Status.Code)
|
||||
}
|
||||
if got := attrString(t, span.Attributes, "ingest.outcome"); got != "error" {
|
||||
t.Errorf("outcome = %q, want error", got)
|
||||
}
|
||||
|
||||
logs := exp.Logs()
|
||||
if len(logs) != 1 {
|
||||
t.Fatalf("got %d logs, want 1", len(logs))
|
||||
}
|
||||
l := logs[0]
|
||||
if l.Severity != telemetry.SeverityError {
|
||||
t.Errorf("severity = %v, want ERROR", l.Severity)
|
||||
}
|
||||
if v, ok := telemetry.Attr(l.Attributes, telemetry.AttrAlert); !ok || v != true {
|
||||
t.Errorf("missing alert marker: %v (present=%v)", v, ok)
|
||||
}
|
||||
if v, _ := telemetry.Attr(l.Attributes, telemetry.AttrAlertType); v != alertIngestError {
|
||||
t.Errorf("alert.type = %v, want %q", v, alertIngestError)
|
||||
}
|
||||
if l.SpanContext != span.SpanContext {
|
||||
t.Errorf("alert log not correlated to span")
|
||||
}
|
||||
}
|
||||
|
||||
// TestIngestNoTelemetryIsInert confirms the default path (no provider in ctx)
|
||||
// emits nothing and preserves the response — the behaviour existing/parallel
|
||||
// tests rely on.
|
||||
func TestIngestNoTelemetryIsInert(t *testing.T) {
|
||||
h := NewHandler(&MemorySink{})
|
||||
// No telemetry in context.
|
||||
rec := do(t, h, http.MethodPost, "application/json", validBody)
|
||||
if rec.Code != http.StatusAccepted {
|
||||
t.Fatalf("status = %d, want 202", rec.Code)
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"strings"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/crypto"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/telemetry"
|
||||
)
|
||||
|
||||
// DefaultMaxPerRun is ADR #6 §3.1's per-run issue cap: at most 50 issues per
|
||||
@@ -15,6 +16,16 @@ import (
|
||||
// rest stay pending for the next run (schedule/lifecycle list oldest-first).
|
||||
const DefaultMaxPerRun = 50
|
||||
|
||||
// Alert types (issue #17) emitted as structured OTEL signals so a backend alert
|
||||
// can key on them once an OTLP endpoint is chosen (TBD).
|
||||
const (
|
||||
// alertCapHit fires when the per-run cap defers reports (the #14 follow-up:
|
||||
// this previously only logged; it is now also a structured, alertable signal).
|
||||
alertCapHit = "publish.cap_hit"
|
||||
// alertRunFailed fires when the run finishes with one or more failures.
|
||||
alertRunFailed = "publish.run_failed"
|
||||
)
|
||||
|
||||
// PendingGetter fetches a pending report's still-encrypted frame by id. It is the
|
||||
// read half of lifecycle.Manager (#10) that the publish job needs; a narrow
|
||||
// interface keeps this package host-testable and free of the Wasm-only storage
|
||||
@@ -188,7 +199,20 @@ func newPublisher(gh issueCreator, kr *crypto.Keyring, getter PendingGetter, opt
|
||||
// as failed and the maintainer is alerted. A failed report is not marked
|
||||
// published (onPublished not called), so it is retried on a later run.
|
||||
func (p *Publisher) Publish(ctx context.Context, ids []string) error {
|
||||
// Observability (#17): a run-level span plus per-report child spans and a
|
||||
// correlated log per report. When no telemetry provider rides in ctx (the
|
||||
// default until an OTLP endpoint is configured) every call here is a no-op and
|
||||
// behaviour is unchanged. The span is a child of the schedule.run span when
|
||||
// invoked from the weekly trigger.
|
||||
tel := telemetry.FromContext(ctx)
|
||||
ctx, span := tel.StartSpan(ctx, "publish.run", telemetry.WithSpanKind(telemetry.SpanKindInternal))
|
||||
defer span.End()
|
||||
span.SetAttributes(telemetry.Int("publish.pending", len(ids)))
|
||||
|
||||
if len(ids) == 0 {
|
||||
span.SetAttributes(telemetry.Int("publish.attempted", 0), telemetry.Int("publish.published", 0))
|
||||
span.SetStatus(telemetry.StatusOK, "")
|
||||
tel.Info(ctx, "publish run: no pending reports")
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -198,6 +222,7 @@ func (p *Publisher) Publish(ctx context.Context, ids []string) error {
|
||||
batch = batch[:p.maxPerRun] // oldest-first; remainder drains next run
|
||||
capHit = true
|
||||
}
|
||||
span.SetAttributes(telemetry.Int("publish.attempted", len(batch)))
|
||||
|
||||
// Ensure the three ADR #6 labels exist once per run (idempotent). Best-effort:
|
||||
// a failure is surfaced but does not abort publishing — issue creation still
|
||||
@@ -206,6 +231,8 @@ func (p *Publisher) Publish(ctx context.Context, ids []string) error {
|
||||
if err := p.gh.EnsureLabels(ctx, p.labels); err != nil {
|
||||
p.logf("publish: ensure labels (continuing, labels applied best-effort): %v", err)
|
||||
errs = append(errs, fmt.Errorf("ensure labels: %w", err))
|
||||
tel.Warn(ctx, "publish: ensure labels failed (continuing best-effort)",
|
||||
telemetry.String("error", err.Error()))
|
||||
}
|
||||
|
||||
names := LabelNames(p.labels)
|
||||
@@ -215,24 +242,72 @@ func (p *Publisher) Publish(ctx context.Context, ids []string) error {
|
||||
errs = append(errs, err)
|
||||
break
|
||||
}
|
||||
if err := p.publishOne(ctx, id, names); err != nil {
|
||||
rctx, rspan := tel.StartSpan(ctx, "publish.report",
|
||||
telemetry.WithSpanKind(telemetry.SpanKindInternal),
|
||||
telemetry.WithAttributes(telemetry.String("report.id", id)))
|
||||
if err := p.publishOne(rctx, id, names); err != nil {
|
||||
p.logf("publish: report %s failed, left pending: %v", id, err)
|
||||
errs = append(errs, fmt.Errorf("report %s: %w", id, err))
|
||||
rspan.SetAttributes(telemetry.String("report.outcome", "failed"))
|
||||
rspan.SetStatus(telemetry.StatusError, err.Error())
|
||||
tel.Error(rctx, "publish: report failed, left pending",
|
||||
telemetry.String("report.id", id),
|
||||
telemetry.String("report.outcome", "failed"),
|
||||
telemetry.String("error", err.Error()))
|
||||
rspan.End()
|
||||
continue
|
||||
}
|
||||
published++
|
||||
rspan.SetAttributes(telemetry.String("report.outcome", "published"))
|
||||
rspan.SetStatus(telemetry.StatusOK, "")
|
||||
tel.Info(rctx, "publish: report published",
|
||||
telemetry.String("report.id", id),
|
||||
telemetry.String("report.outcome", "published"))
|
||||
rspan.End()
|
||||
}
|
||||
|
||||
if capHit {
|
||||
deferred := len(ids) - len(batch)
|
||||
// ADR #6 §3.1: alert the maintainer that the cap was hit and reports were
|
||||
// deferred. The log line is the alert channel until dedicated alerting lands.
|
||||
// deferred. The log line remains for humans...
|
||||
p.logf("publish: per-run cap %d reached: published %d of %d pending this run; %d deferred to next run",
|
||||
p.maxPerRun, published, len(ids), len(ids)-len(batch))
|
||||
p.maxPerRun, published, len(ids), deferred)
|
||||
// ...and (#17, folding in the #14 follow-up) the cap-hit is now also emitted
|
||||
// as a structured, alertable OTEL signal rather than only a log line.
|
||||
span.AddEvent(alertCapHit,
|
||||
telemetry.Int("publish.cap", p.maxPerRun),
|
||||
telemetry.Int("publish.deferred", deferred))
|
||||
span.SetAttributes(telemetry.Bool("publish.cap_hit", true), telemetry.Int("publish.deferred", deferred))
|
||||
capAttrs := append(telemetry.Alert(alertCapHit),
|
||||
telemetry.Int("publish.cap", p.maxPerRun),
|
||||
telemetry.Int("publish.published", published),
|
||||
telemetry.Int("publish.pending", len(ids)),
|
||||
telemetry.Int("publish.deferred", deferred))
|
||||
tel.Warn(ctx, "publish: per-run cap reached; reports deferred to next run", capAttrs...)
|
||||
}
|
||||
p.logf("publish: run complete: %d published, %d failed, of %d attempted",
|
||||
published, len(batch)-published, len(batch))
|
||||
|
||||
return errors.Join(errs...)
|
||||
joined := errors.Join(errs...)
|
||||
failed := len(batch) - published
|
||||
span.SetAttributes(
|
||||
telemetry.Int("publish.published", published),
|
||||
telemetry.Int("publish.failed", failed),
|
||||
)
|
||||
if joined != nil {
|
||||
span.SetStatus(telemetry.StatusError, joined.Error())
|
||||
runAttrs := append(telemetry.Alert(alertRunFailed),
|
||||
telemetry.Int("publish.published", published),
|
||||
telemetry.Int("publish.failed", failed),
|
||||
telemetry.Int("publish.attempted", len(batch)))
|
||||
tel.Error(ctx, "publish run failed", runAttrs...)
|
||||
} else {
|
||||
span.SetStatus(telemetry.StatusOK, "")
|
||||
tel.Info(ctx, "publish run complete",
|
||||
telemetry.Int("publish.published", published),
|
||||
telemetry.Int("publish.failed", failed))
|
||||
}
|
||||
return joined
|
||||
}
|
||||
|
||||
// publishOne runs the full pipeline for a single report id. It returns an error
|
||||
|
||||
@@ -0,0 +1,233 @@
|
||||
package publish
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/telemetry"
|
||||
)
|
||||
|
||||
// telemetryCtx returns a context carrying a fresh telemetry provider plus its
|
||||
// in-memory exporter for assertions.
|
||||
func telemetryCtx() (context.Context, *telemetry.MemoryExporter) {
|
||||
exp := telemetry.NewMemoryExporter()
|
||||
tel := telemetry.New(exp)
|
||||
return telemetry.NewContext(context.Background(), tel), exp
|
||||
}
|
||||
|
||||
// TestPublishEmitsRunAndReportSpans: a successful run emits one publish.run span
|
||||
// and one publish.report child span per report, each with a correlated log.
|
||||
func TestPublishEmitsRunAndReportSpans(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
ids := []string{"id-1", "id-2", "id-3"}
|
||||
frames := map[string][]byte{}
|
||||
for _, id := range ids {
|
||||
frames[id] = seal(t, kr, reportJSON(t, sampleReport(id)))
|
||||
}
|
||||
mc := &mockCreator{}
|
||||
pub := newPublisher(mc, kr, &fakeGetter{frames: frames})
|
||||
|
||||
ctx, exp := telemetryCtx()
|
||||
if err := pub.Publish(ctx, ids); err != nil {
|
||||
t.Fatalf("Publish: %v", err)
|
||||
}
|
||||
|
||||
run, ok := exp.SpanByName("publish.run")
|
||||
if !ok {
|
||||
t.Fatal("no publish.run span")
|
||||
}
|
||||
if run.Status.Code != telemetry.StatusOK {
|
||||
t.Errorf("run status = %d, want OK", run.Status.Code)
|
||||
}
|
||||
if v, _ := telemetry.Attr(run.Attributes, "publish.published"); v != int64(3) {
|
||||
t.Errorf("publish.published = %v, want 3", v)
|
||||
}
|
||||
if v, _ := telemetry.Attr(run.Attributes, "publish.failed"); v != int64(0) {
|
||||
t.Errorf("publish.failed = %v, want 0", v)
|
||||
}
|
||||
|
||||
reports := exp.SpansByName("publish.report")
|
||||
if len(reports) != len(ids) {
|
||||
t.Fatalf("got %d publish.report spans, want %d", len(reports), len(ids))
|
||||
}
|
||||
for _, rs := range reports {
|
||||
// Each report span is a child of the run span, sharing its trace.
|
||||
if rs.SpanContext.TraceID != run.SpanContext.TraceID {
|
||||
t.Errorf("report span trace %s != run trace %s", rs.SpanContext.TraceID, run.SpanContext.TraceID)
|
||||
}
|
||||
if rs.ParentSpanID != run.SpanContext.SpanID {
|
||||
t.Errorf("report span parent = %s, want run span %s", rs.ParentSpanID, run.SpanContext.SpanID)
|
||||
}
|
||||
if v, _ := telemetry.Attr(rs.Attributes, "report.outcome"); v != "published" {
|
||||
t.Errorf("report outcome = %v, want published", v)
|
||||
}
|
||||
}
|
||||
|
||||
// One "published" log per report, all correlated to a report span.
|
||||
published := 0
|
||||
for _, l := range exp.Logs() {
|
||||
if v, _ := telemetry.Attr(l.Attributes, "report.outcome"); v == "published" {
|
||||
published++
|
||||
if !l.SpanContext.IsValid() {
|
||||
t.Errorf("published log not correlated to a span")
|
||||
}
|
||||
}
|
||||
}
|
||||
if published != len(ids) {
|
||||
t.Errorf("got %d published logs, want %d", published, len(ids))
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishCapHitEmitsAlertSignal is the #14 follow-up folded into #17: when
|
||||
// the per-run cap defers reports, a structured, alertable OTEL signal is emitted
|
||||
// (a WARN log with alert.type=publish.cap_hit plus counts, and a run-span
|
||||
// attribute/event) — not merely a human log line.
|
||||
func TestPublishCapHitEmitsAlertSignal(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
const n = DefaultMaxPerRun + 5
|
||||
ids := make([]string, n)
|
||||
frames := map[string][]byte{}
|
||||
for i := range ids {
|
||||
id := fmt.Sprintf("id-%03d", i)
|
||||
ids[i] = id
|
||||
frames[id] = seal(t, kr, reportJSON(t, sampleReport(id)))
|
||||
}
|
||||
mc := &mockCreator{}
|
||||
pub := newPublisher(mc, kr, &fakeGetter{frames: frames})
|
||||
|
||||
ctx, exp := telemetryCtx()
|
||||
if err := pub.Publish(ctx, ids); err != nil {
|
||||
t.Fatalf("Publish: %v", err)
|
||||
}
|
||||
|
||||
// The run span records the cap-hit.
|
||||
run, _ := exp.SpanByName("publish.run")
|
||||
if v, _ := telemetry.Attr(run.Attributes, "publish.cap_hit"); v != true {
|
||||
t.Errorf("run span publish.cap_hit = %v, want true", v)
|
||||
}
|
||||
if v, _ := telemetry.Attr(run.Attributes, "publish.deferred"); v != int64(5) {
|
||||
t.Errorf("run span publish.deferred = %v, want 5", v)
|
||||
}
|
||||
foundEvent := false
|
||||
for _, ev := range run.Events {
|
||||
if ev.Name == alertCapHit {
|
||||
foundEvent = true
|
||||
}
|
||||
}
|
||||
if !foundEvent {
|
||||
t.Errorf("run span missing %q event", alertCapHit)
|
||||
}
|
||||
|
||||
// The alertable log record.
|
||||
var capLog *telemetry.LogRecord
|
||||
for i := range exp.Logs() {
|
||||
l := exp.Logs()[i]
|
||||
if v, _ := telemetry.Attr(l.Attributes, telemetry.AttrAlertType); v == alertCapHit {
|
||||
capLog = &l
|
||||
break
|
||||
}
|
||||
}
|
||||
if capLog == nil {
|
||||
t.Fatal("no cap-hit alert log emitted")
|
||||
}
|
||||
if capLog.Severity != telemetry.SeverityWarn {
|
||||
t.Errorf("cap-hit log severity = %v, want WARN", capLog.Severity)
|
||||
}
|
||||
if v, ok := telemetry.Attr(capLog.Attributes, telemetry.AttrAlert); !ok || v != true {
|
||||
t.Errorf("cap-hit log missing alert marker")
|
||||
}
|
||||
if v, _ := telemetry.Attr(capLog.Attributes, "publish.deferred"); v != int64(5) {
|
||||
t.Errorf("cap-hit log publish.deferred = %v, want 5", v)
|
||||
}
|
||||
if v, _ := telemetry.Attr(capLog.Attributes, "publish.cap"); v != int64(DefaultMaxPerRun) {
|
||||
t.Errorf("cap-hit log publish.cap = %v, want %d", v, DefaultMaxPerRun)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishRunFailureEmitsAlertSignal: a per-report failure makes the run span
|
||||
// Error and emits the alert.type=publish.run_failed signal, while the failed
|
||||
// report gets its own Error span + correlated ERROR log.
|
||||
func TestPublishRunFailureEmitsAlertSignal(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
wrongKr := mustKeyring(t) // different key -> decrypt failure on "bad"
|
||||
ids := []string{"good-1", "bad", "good-2"}
|
||||
frames := map[string][]byte{
|
||||
"good-1": seal(t, kr, reportJSON(t, sampleReport("good-1"))),
|
||||
"bad": seal(t, wrongKr, reportJSON(t, sampleReport("bad"))),
|
||||
"good-2": seal(t, kr, reportJSON(t, sampleReport("good-2"))),
|
||||
}
|
||||
mc := &mockCreator{}
|
||||
pub := newPublisher(mc, kr, &fakeGetter{frames: frames})
|
||||
|
||||
ctx, exp := telemetryCtx()
|
||||
if err := pub.Publish(ctx, ids); err == nil {
|
||||
t.Fatal("expected a surfaced error for the decrypt failure")
|
||||
}
|
||||
|
||||
run, _ := exp.SpanByName("publish.run")
|
||||
if run.Status.Code != telemetry.StatusError {
|
||||
t.Errorf("run span status = %d, want Error", run.Status.Code)
|
||||
}
|
||||
if v, _ := telemetry.Attr(run.Attributes, "publish.failed"); v != int64(1) {
|
||||
t.Errorf("publish.failed = %v, want 1", v)
|
||||
}
|
||||
|
||||
// A run-failed alert log.
|
||||
foundRunAlert := false
|
||||
for _, l := range exp.Logs() {
|
||||
if v, _ := telemetry.Attr(l.Attributes, telemetry.AttrAlertType); v == alertRunFailed {
|
||||
foundRunAlert = true
|
||||
if l.Severity != telemetry.SeverityError {
|
||||
t.Errorf("run-failed log severity = %v, want ERROR", l.Severity)
|
||||
}
|
||||
}
|
||||
}
|
||||
if !foundRunAlert {
|
||||
t.Errorf("no %q alert log emitted", alertRunFailed)
|
||||
}
|
||||
|
||||
// The failed report span is present and marked failed.
|
||||
var failedSpan *telemetry.SpanData
|
||||
for i, rs := range exp.SpansByName("publish.report") {
|
||||
if v, _ := telemetry.Attr(rs.Attributes, "report.id"); v == "bad" {
|
||||
s := exp.SpansByName("publish.report")[i]
|
||||
failedSpan = &s
|
||||
}
|
||||
}
|
||||
if failedSpan == nil {
|
||||
t.Fatal("no publish.report span for the failed report")
|
||||
}
|
||||
if failedSpan.Status.Code != telemetry.StatusError {
|
||||
t.Errorf("failed report span status = %d, want Error", failedSpan.Status.Code)
|
||||
}
|
||||
if v, _ := telemetry.Attr(failedSpan.Attributes, "report.outcome"); v != "failed" {
|
||||
t.Errorf("failed report outcome = %v, want failed", v)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishEmptyBatchTraced: an empty run still emits a publish.run span (OK,
|
||||
// zero attempted) so every scheduled run is observable.
|
||||
func TestPublishEmptyBatchTraced(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
mc := &mockCreator{}
|
||||
pub := newPublisher(mc, kr, &fakeGetter{})
|
||||
ctx, exp := telemetryCtx()
|
||||
if err := pub.Publish(ctx, nil); err != nil {
|
||||
t.Fatalf("Publish(nil): %v", err)
|
||||
}
|
||||
run, ok := exp.SpanByName("publish.run")
|
||||
if !ok {
|
||||
t.Fatal("empty run should still emit a publish.run span")
|
||||
}
|
||||
if v, _ := telemetry.Attr(run.Attributes, "publish.attempted"); v != int64(0) {
|
||||
t.Errorf("publish.attempted = %v, want 0", v)
|
||||
}
|
||||
// No per-report spans, and the mock was never touched (behaviour unchanged).
|
||||
if len(exp.SpansByName("publish.report")) != 0 {
|
||||
t.Errorf("empty run emitted report spans")
|
||||
}
|
||||
if mc.ensureCalls != 0 || len(mc.created) != 0 {
|
||||
t.Errorf("empty batch did work: ensureCalls=%d created=%d", mc.ensureCalls, len(mc.created))
|
||||
}
|
||||
}
|
||||
@@ -31,8 +31,15 @@ import (
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/telemetry"
|
||||
)
|
||||
|
||||
// alertScheduleRunFailed is the alert.type (issue #17) emitted as a structured
|
||||
// OTEL signal when a gated weekly run fails to list or publish, so a backend
|
||||
// alert can key on a failed scheduled run once an OTLP endpoint is chosen (TBD).
|
||||
const alertScheduleRunFailed = "schedule.run_failed"
|
||||
|
||||
// PendingLister is the seam for "list the reports awaiting publication". It is
|
||||
// satisfied by *lifecycle.Manager (#10) via its ListPending method; tests inject
|
||||
// a fake. Kept as a one-method interface so this package does not import the
|
||||
@@ -68,13 +75,36 @@ func Run(ctx context.Context, scheduledFor time.Time, lister PendingLister, publ
|
||||
// sibling cron will fire at the correct Central hour. Do nothing.
|
||||
return false, nil
|
||||
}
|
||||
// Observability (#17): trace the gated run end to end (list + publish). When
|
||||
// no telemetry provider rides in ctx (the default until an OTLP endpoint is
|
||||
// configured) this is a no-op and behaviour is unchanged; when present, the
|
||||
// span becomes the parent of the publisher's publish.run span. A list/publish
|
||||
// failure sets the span to Error and emits the alertable run-failed signal.
|
||||
tel := telemetry.FromContext(ctx)
|
||||
ctx, span := tel.StartSpan(ctx, "schedule.run",
|
||||
telemetry.WithSpanKind(telemetry.SpanKindInternal),
|
||||
telemetry.WithAttributes(telemetry.String("schedule.scheduled_for", scheduledFor.UTC().Format(time.RFC3339))))
|
||||
defer span.End()
|
||||
|
||||
ids, err := lister.ListPending(ctx)
|
||||
if err != nil {
|
||||
return true, fmt.Errorf("schedule: list pending: %w", err)
|
||||
wrapped := fmt.Errorf("schedule: list pending: %w", err)
|
||||
span.SetStatus(telemetry.StatusError, wrapped.Error())
|
||||
tel.Error(ctx, "schedule: list pending failed",
|
||||
append(telemetry.Alert(alertScheduleRunFailed), telemetry.String("error", wrapped.Error()))...)
|
||||
return true, wrapped
|
||||
}
|
||||
span.SetAttributes(telemetry.Int("schedule.pending", len(ids)))
|
||||
tel.Info(ctx, "schedule: weekly trigger fired", telemetry.Int("schedule.pending", len(ids)))
|
||||
|
||||
if err := publisher.Publish(ctx, ids); err != nil {
|
||||
return true, fmt.Errorf("schedule: publish: %w", err)
|
||||
wrapped := fmt.Errorf("schedule: publish: %w", err)
|
||||
span.SetStatus(telemetry.StatusError, wrapped.Error())
|
||||
tel.Error(ctx, "schedule: publish run failed",
|
||||
append(telemetry.Alert(alertScheduleRunFailed), telemetry.String("error", wrapped.Error()))...)
|
||||
return true, wrapped
|
||||
}
|
||||
span.SetStatus(telemetry.StatusOK, "")
|
||||
return true, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,141 @@
|
||||
package schedule_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/schedule"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/telemetry"
|
||||
)
|
||||
|
||||
func telemetryCtx() (context.Context, *telemetry.MemoryExporter) {
|
||||
exp := telemetry.NewMemoryExporter()
|
||||
return telemetry.NewContext(context.Background(), telemetry.New(exp)), exp
|
||||
}
|
||||
|
||||
// TestScheduleRunEmitsSpan: a gated run emits a schedule.run span (OK) plus a
|
||||
// "weekly trigger fired" log recording the pending count.
|
||||
func TestScheduleRunEmitsSpan(t *testing.T) {
|
||||
want := []string{"a", "b"}
|
||||
lister := &fakeLister{ids: want}
|
||||
pub := &recordingPublisher{}
|
||||
|
||||
ctx, exp := telemetryCtx()
|
||||
ran, err := schedule.Run(ctx, fireInstant, lister, pub)
|
||||
if err != nil || !ran {
|
||||
t.Fatalf("Run = (%v, %v), want (true, nil)", ran, err)
|
||||
}
|
||||
|
||||
span, ok := exp.SpanByName("schedule.run")
|
||||
if !ok {
|
||||
t.Fatal("no schedule.run span")
|
||||
}
|
||||
if span.Status.Code != telemetry.StatusOK {
|
||||
t.Errorf("span status = %d, want OK", span.Status.Code)
|
||||
}
|
||||
if v, _ := telemetry.Attr(span.Attributes, "schedule.pending"); v != int64(2) {
|
||||
t.Errorf("schedule.pending = %v, want 2", v)
|
||||
}
|
||||
if _, ok := telemetry.Attr(span.Attributes, "schedule.scheduled_for"); !ok {
|
||||
t.Errorf("span missing schedule.scheduled_for")
|
||||
}
|
||||
|
||||
firedLog := false
|
||||
for _, l := range exp.Logs() {
|
||||
if l.Body == "schedule: weekly trigger fired" {
|
||||
firedLog = true
|
||||
if l.SpanContext != span.SpanContext {
|
||||
t.Errorf("fired log not correlated to schedule.run span")
|
||||
}
|
||||
}
|
||||
}
|
||||
if !firedLog {
|
||||
t.Errorf("no weekly-trigger-fired log")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScheduleGateClosedEmitsNothing: the sibling (no-op) fire starts no span and
|
||||
// emits no telemetry — matching its do-nothing contract.
|
||||
func TestScheduleGateClosedEmitsNothing(t *testing.T) {
|
||||
lister := &fakeLister{ids: []string{"x"}}
|
||||
pub := &recordingPublisher{}
|
||||
ctx, exp := telemetryCtx()
|
||||
|
||||
ran, err := schedule.Run(ctx, notFireInstant, lister, pub)
|
||||
if err != nil || ran {
|
||||
t.Fatalf("Run = (%v, %v), want (false, nil)", ran, err)
|
||||
}
|
||||
if len(exp.Spans()) != 0 || len(exp.Logs()) != 0 {
|
||||
t.Errorf("gate-closed fire emitted telemetry: %d spans, %d logs", len(exp.Spans()), len(exp.Logs()))
|
||||
}
|
||||
}
|
||||
|
||||
// TestScheduleListErrorEmitsAlert: a list failure sets the span to Error, emits
|
||||
// the alert.type=schedule.run_failed signal, and still returns the wrapped error.
|
||||
func TestScheduleListErrorEmitsAlert(t *testing.T) {
|
||||
sentinel := errors.New("r2 list failed")
|
||||
lister := &fakeLister{err: sentinel}
|
||||
pub := &recordingPublisher{}
|
||||
ctx, exp := telemetryCtx()
|
||||
|
||||
ran, err := schedule.Run(ctx, fireInstant, lister, pub)
|
||||
if !ran || !errors.Is(err, sentinel) {
|
||||
t.Fatalf("Run = (%v, %v), want (true, wraps sentinel)", ran, err)
|
||||
}
|
||||
|
||||
span, _ := exp.SpanByName("schedule.run")
|
||||
if span.Status.Code != telemetry.StatusError {
|
||||
t.Errorf("span status = %d, want Error", span.Status.Code)
|
||||
}
|
||||
found := false
|
||||
for _, l := range exp.Logs() {
|
||||
if v, _ := telemetry.Attr(l.Attributes, telemetry.AttrAlertType); v == "schedule.run_failed" {
|
||||
found = true
|
||||
if l.Severity != telemetry.SeverityError {
|
||||
t.Errorf("alert log severity = %v, want ERROR", l.Severity)
|
||||
}
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Errorf("no schedule.run_failed alert log")
|
||||
}
|
||||
if len(pub.batches) != 0 {
|
||||
t.Errorf("publisher called despite a list error")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScheduleRunParentsPublisher: the ctx handed to the publisher carries the
|
||||
// schedule.run span, so a publish.run span started from it is a child — the trace
|
||||
// links the scheduled run to the publish run end to end.
|
||||
func TestScheduleRunParentsPublisher(t *testing.T) {
|
||||
ctx, exp := telemetryCtx()
|
||||
var childTrace telemetry.TraceID
|
||||
var childParent telemetry.SpanID
|
||||
spy := publisherFunc(func(pctx context.Context, _ []string) error {
|
||||
tel := telemetry.FromContext(pctx)
|
||||
_, span := tel.StartSpan(pctx, "publish.run")
|
||||
childTrace = span.SpanContext().TraceID
|
||||
sc, _ := telemetry.SpanContextFromContext(pctx)
|
||||
childParent = sc.SpanID
|
||||
span.End()
|
||||
return nil
|
||||
})
|
||||
|
||||
ran, err := schedule.Run(ctx, fireInstant, &fakeLister{ids: []string{"a"}}, spy)
|
||||
if err != nil || !ran {
|
||||
t.Fatalf("Run = (%v, %v), want (true, nil)", ran, err)
|
||||
}
|
||||
run, _ := exp.SpanByName("schedule.run")
|
||||
if childTrace != run.SpanContext.TraceID {
|
||||
t.Errorf("publisher's span trace %s != schedule.run trace %s", childTrace, run.SpanContext.TraceID)
|
||||
}
|
||||
if childParent != run.SpanContext.SpanID {
|
||||
t.Errorf("ctx active span in publisher = %s, want schedule.run span %s", childParent, run.SpanContext.SpanID)
|
||||
}
|
||||
}
|
||||
|
||||
// publisherFunc adapts a function to schedule.Publisher.
|
||||
type publisherFunc func(context.Context, []string) error
|
||||
|
||||
func (f publisherFunc) Publish(ctx context.Context, ids []string) error { return f(ctx, ids) }
|
||||
@@ -0,0 +1,99 @@
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// MemoryExporter is an in-memory [Exporter] for host tests: it retains every
|
||||
// span and log record it is handed so tests can assert on them. It is safe for
|
||||
// concurrent use and never fails. It must not be used in production (it grows
|
||||
// without bound).
|
||||
type MemoryExporter struct {
|
||||
mu sync.Mutex
|
||||
spans []SpanData
|
||||
logs []LogRecord
|
||||
}
|
||||
|
||||
// NewMemoryExporter returns an empty MemoryExporter.
|
||||
func NewMemoryExporter() *MemoryExporter { return &MemoryExporter{} }
|
||||
|
||||
// ExportSpans records spans.
|
||||
func (m *MemoryExporter) ExportSpans(_ context.Context, spans []SpanData) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.spans = append(m.spans, spans...)
|
||||
return nil
|
||||
}
|
||||
|
||||
// ExportLogs records log records.
|
||||
func (m *MemoryExporter) ExportLogs(_ context.Context, logs []LogRecord) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.logs = append(m.logs, logs...)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Shutdown is a no-op.
|
||||
func (m *MemoryExporter) Shutdown(context.Context) error { return nil }
|
||||
|
||||
// Spans returns a snapshot copy of the exported spans.
|
||||
func (m *MemoryExporter) Spans() []SpanData {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
return append([]SpanData(nil), m.spans...)
|
||||
}
|
||||
|
||||
// Logs returns a snapshot copy of the exported log records.
|
||||
func (m *MemoryExporter) Logs() []LogRecord {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
return append([]LogRecord(nil), m.logs...)
|
||||
}
|
||||
|
||||
// Reset clears all captured spans and logs.
|
||||
func (m *MemoryExporter) Reset() {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.spans = nil
|
||||
m.logs = nil
|
||||
}
|
||||
|
||||
// SpanByName returns the first captured span with the given name and whether one
|
||||
// was found. A convenience for assertions.
|
||||
func (m *MemoryExporter) SpanByName(name string) (SpanData, bool) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
for _, s := range m.spans {
|
||||
if s.Name == name {
|
||||
return s, true
|
||||
}
|
||||
}
|
||||
return SpanData{}, false
|
||||
}
|
||||
|
||||
// SpansByName returns all captured spans with the given name, in export order.
|
||||
func (m *MemoryExporter) SpansByName(name string) []SpanData {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
var out []SpanData
|
||||
for _, s := range m.spans {
|
||||
if s.Name == name {
|
||||
out = append(out, s)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// Attr returns the value of attribute key among attrs and whether it was
|
||||
// present. It is a helper for tests asserting on span/log attributes.
|
||||
func Attr(attrs []KeyValue, key string) (any, bool) {
|
||||
for _, kv := range attrs {
|
||||
if kv.Key == key {
|
||||
return kv.Value, true
|
||||
}
|
||||
}
|
||||
return nil, false
|
||||
}
|
||||
|
||||
var _ Exporter = (*MemoryExporter)(nil)
|
||||
@@ -0,0 +1,356 @@
|
||||
package telemetry
|
||||
|
||||
// OTLP/HTTP exporter with JSON encoding.
|
||||
//
|
||||
// This is the production sink: it POSTs spans and log records to an OTLP/HTTP
|
||||
// endpoint as JSON (the OTLP-over-HTTP "application/json" encoding, the ProtoJSON
|
||||
// mapping of opentelemetry-proto). JSON is chosen over protobuf because it needs
|
||||
// no code generation or protobuf runtime — just encoding/json and net/http, both
|
||||
// of which compile and run under this project's TinyGo/Wasm Worker target (see
|
||||
// internal/publish, which drives the GitHub API the same way).
|
||||
//
|
||||
// Export is synchronous and best-effort: one POST per span-batch and per
|
||||
// log-batch. Volumes here are tiny (a handful of spans per ingest request; one
|
||||
// run plus per-report spans once a week), so batching/async delivery is a
|
||||
// deliberate non-goal for this first cut — an OTLP collector or a batching
|
||||
// wrapper can be added later without touching the instrumented code.
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Scope identifies the instrumentation scope (the library emitting the
|
||||
// telemetry) in the OTLP payload.
|
||||
type Scope struct {
|
||||
Name string
|
||||
Version string
|
||||
}
|
||||
|
||||
// OTLPConfig configures an [OTLPHTTPExporter].
|
||||
type OTLPConfig struct {
|
||||
// Endpoint is the base OTLP/HTTP URL, e.g. "https://otlp.example.com". The
|
||||
// signal paths "/v1/traces" and "/v1/logs" are appended. This is the value
|
||||
// issue #17 leaves TBD; it comes from the OTEL_EXPORTER_OTLP_ENDPOINT Worker
|
||||
// var and is never hardcoded.
|
||||
Endpoint string
|
||||
// Headers are added to every request (e.g. an auth header). In the Worker
|
||||
// these come from a secret, never a committed value.
|
||||
Headers map[string]string
|
||||
// Resource attributes (e.g. service.name) describe the emitting service.
|
||||
Resource []KeyValue
|
||||
// Scope names the instrumentation library.
|
||||
Scope Scope
|
||||
// HTTPClient overrides the client (tests inject one; default &http.Client{}).
|
||||
HTTPClient *http.Client
|
||||
// Timeout bounds each export request when HTTPClient is not supplied
|
||||
// (default 10s). Ignored if HTTPClient is set.
|
||||
Timeout time.Duration
|
||||
}
|
||||
|
||||
// OTLPHTTPExporter exports spans and logs to an OTLP/HTTP endpoint as JSON.
|
||||
type OTLPHTTPExporter struct {
|
||||
client *http.Client
|
||||
tracesURL string
|
||||
logsURL string
|
||||
headers map[string]string
|
||||
resource []KeyValue
|
||||
scope Scope
|
||||
}
|
||||
|
||||
// NewOTLPHTTPExporter builds an exporter from cfg. It does no I/O.
|
||||
func NewOTLPHTTPExporter(cfg OTLPConfig) *OTLPHTTPExporter {
|
||||
client := cfg.HTTPClient
|
||||
if client == nil {
|
||||
timeout := cfg.Timeout
|
||||
if timeout <= 0 {
|
||||
timeout = 10 * time.Second
|
||||
}
|
||||
client = &http.Client{Timeout: timeout}
|
||||
}
|
||||
base := strings.TrimRight(cfg.Endpoint, "/")
|
||||
headers := make(map[string]string, len(cfg.Headers))
|
||||
for k, v := range cfg.Headers {
|
||||
headers[k] = v
|
||||
}
|
||||
return &OTLPHTTPExporter{
|
||||
client: client,
|
||||
tracesURL: base + "/v1/traces",
|
||||
logsURL: base + "/v1/logs",
|
||||
headers: headers,
|
||||
resource: cfg.Resource,
|
||||
scope: cfg.Scope,
|
||||
}
|
||||
}
|
||||
|
||||
// ExportSpans POSTs spans to {endpoint}/v1/traces as OTLP/JSON.
|
||||
func (e *OTLPHTTPExporter) ExportSpans(ctx context.Context, spans []SpanData) error {
|
||||
if len(spans) == 0 {
|
||||
return nil
|
||||
}
|
||||
body, err := json.Marshal(e.buildTracesPayload(spans))
|
||||
if err != nil {
|
||||
return fmt.Errorf("otlp: marshal traces: %w", err)
|
||||
}
|
||||
return e.post(ctx, e.tracesURL, body)
|
||||
}
|
||||
|
||||
// ExportLogs POSTs logs to {endpoint}/v1/logs as OTLP/JSON.
|
||||
func (e *OTLPHTTPExporter) ExportLogs(ctx context.Context, logs []LogRecord) error {
|
||||
if len(logs) == 0 {
|
||||
return nil
|
||||
}
|
||||
body, err := json.Marshal(e.buildLogsPayload(logs))
|
||||
if err != nil {
|
||||
return fmt.Errorf("otlp: marshal logs: %w", err)
|
||||
}
|
||||
return e.post(ctx, e.logsURL, body)
|
||||
}
|
||||
|
||||
// Shutdown is a no-op; the exporter holds no long-lived resources.
|
||||
func (e *OTLPHTTPExporter) Shutdown(context.Context) error { return nil }
|
||||
|
||||
// post sends body to url with the configured headers and returns an error on a
|
||||
// transport failure or a non-2xx status.
|
||||
func (e *OTLPHTTPExporter) post(ctx context.Context, url string, body []byte) error {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return fmt.Errorf("otlp: build request: %w", err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
for k, v := range e.headers {
|
||||
req.Header.Set(k, v)
|
||||
}
|
||||
resp, err := e.client.Do(req)
|
||||
if err != nil {
|
||||
return fmt.Errorf("otlp: post %s: %w", url, err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
// Drain so the connection can be reused.
|
||||
_, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 1<<16))
|
||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||
return fmt.Errorf("otlp: post %s: unexpected status %d", url, resp.StatusCode)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// --- OTLP/JSON wire types (ProtoJSON mapping of opentelemetry-proto) ---
|
||||
|
||||
type otlpTracesPayload struct {
|
||||
ResourceSpans []otlpResourceSpans `json:"resourceSpans"`
|
||||
}
|
||||
|
||||
type otlpResourceSpans struct {
|
||||
Resource otlpResource `json:"resource"`
|
||||
ScopeSpans []otlpScopeSpans `json:"scopeSpans"`
|
||||
}
|
||||
|
||||
type otlpScopeSpans struct {
|
||||
Scope otlpScope `json:"scope"`
|
||||
Spans []otlpSpan `json:"spans"`
|
||||
}
|
||||
|
||||
type otlpLogsPayload struct {
|
||||
ResourceLogs []otlpResourceLogs `json:"resourceLogs"`
|
||||
}
|
||||
|
||||
type otlpResourceLogs struct {
|
||||
Resource otlpResource `json:"resource"`
|
||||
ScopeLogs []otlpScopeLogs `json:"scopeLogs"`
|
||||
}
|
||||
|
||||
type otlpScopeLogs struct {
|
||||
Scope otlpScope `json:"scope"`
|
||||
LogRecords []otlpLogRecord `json:"logRecords"`
|
||||
}
|
||||
|
||||
type otlpResource struct {
|
||||
Attributes []otlpKeyValue `json:"attributes,omitempty"`
|
||||
}
|
||||
|
||||
type otlpScope struct {
|
||||
Name string `json:"name,omitempty"`
|
||||
Version string `json:"version,omitempty"`
|
||||
}
|
||||
|
||||
type otlpSpan struct {
|
||||
TraceID string `json:"traceId"`
|
||||
SpanID string `json:"spanId"`
|
||||
ParentSpanID string `json:"parentSpanId,omitempty"`
|
||||
Name string `json:"name"`
|
||||
Kind int `json:"kind,omitempty"`
|
||||
StartTimeUnixNano string `json:"startTimeUnixNano"`
|
||||
EndTimeUnixNano string `json:"endTimeUnixNano"`
|
||||
Attributes []otlpKeyValue `json:"attributes,omitempty"`
|
||||
Events []otlpEvent `json:"events,omitempty"`
|
||||
Status *otlpStatus `json:"status,omitempty"`
|
||||
}
|
||||
|
||||
type otlpEvent struct {
|
||||
TimeUnixNano string `json:"timeUnixNano"`
|
||||
Name string `json:"name"`
|
||||
Attributes []otlpKeyValue `json:"attributes,omitempty"`
|
||||
}
|
||||
|
||||
type otlpStatus struct {
|
||||
Message string `json:"message,omitempty"`
|
||||
Code int `json:"code"`
|
||||
}
|
||||
|
||||
type otlpLogRecord struct {
|
||||
TimeUnixNano string `json:"timeUnixNano"`
|
||||
SeverityNumber int `json:"severityNumber,omitempty"`
|
||||
SeverityText string `json:"severityText,omitempty"`
|
||||
Body *otlpAnyValue `json:"body,omitempty"`
|
||||
Attributes []otlpKeyValue `json:"attributes,omitempty"`
|
||||
TraceID string `json:"traceId,omitempty"`
|
||||
SpanID string `json:"spanId,omitempty"`
|
||||
}
|
||||
|
||||
type otlpKeyValue struct {
|
||||
Key string `json:"key"`
|
||||
Value otlpAnyValue `json:"value"`
|
||||
}
|
||||
|
||||
// otlpAnyValue is the OTLP AnyValue. Exactly one field is non-nil. Per proto3
|
||||
// JSON, int64 is encoded as a string.
|
||||
type otlpAnyValue struct {
|
||||
StringValue *string `json:"stringValue,omitempty"`
|
||||
BoolValue *bool `json:"boolValue,omitempty"`
|
||||
IntValue *string `json:"intValue,omitempty"`
|
||||
DoubleValue *float64 `json:"doubleValue,omitempty"`
|
||||
}
|
||||
|
||||
func (e *OTLPHTTPExporter) buildTracesPayload(spans []SpanData) otlpTracesPayload {
|
||||
out := make([]otlpSpan, 0, len(spans))
|
||||
for _, s := range spans {
|
||||
os := otlpSpan{
|
||||
TraceID: s.SpanContext.TraceID.String(),
|
||||
SpanID: s.SpanContext.SpanID.String(),
|
||||
Name: s.Name,
|
||||
Kind: int(s.Kind),
|
||||
StartTimeUnixNano: nanos(s.StartTime),
|
||||
EndTimeUnixNano: nanos(s.EndTime),
|
||||
Attributes: toOTLPAttrs(s.Attributes),
|
||||
}
|
||||
if !s.ParentSpanID.IsZero() {
|
||||
os.ParentSpanID = s.ParentSpanID.String()
|
||||
}
|
||||
if len(s.Events) > 0 {
|
||||
os.Events = make([]otlpEvent, 0, len(s.Events))
|
||||
for _, ev := range s.Events {
|
||||
os.Events = append(os.Events, otlpEvent{
|
||||
TimeUnixNano: nanos(ev.Time),
|
||||
Name: ev.Name,
|
||||
Attributes: toOTLPAttrs(ev.Attributes),
|
||||
})
|
||||
}
|
||||
}
|
||||
if s.Status.Code != StatusUnset || s.Status.Message != "" {
|
||||
os.Status = &otlpStatus{Code: int(s.Status.Code), Message: s.Status.Message}
|
||||
}
|
||||
out = append(out, os)
|
||||
}
|
||||
return otlpTracesPayload{ResourceSpans: []otlpResourceSpans{{
|
||||
Resource: otlpResource{Attributes: toOTLPAttrs(e.resource)},
|
||||
ScopeSpans: []otlpScopeSpans{{Scope: e.otlpScope(), Spans: out}},
|
||||
}}}
|
||||
}
|
||||
|
||||
func (e *OTLPHTTPExporter) buildLogsPayload(logs []LogRecord) otlpLogsPayload {
|
||||
out := make([]otlpLogRecord, 0, len(logs))
|
||||
for _, l := range logs {
|
||||
body := l.Body
|
||||
rec := otlpLogRecord{
|
||||
TimeUnixNano: nanos(l.Time),
|
||||
SeverityNumber: int(l.Severity),
|
||||
SeverityText: l.Severity.String(),
|
||||
Body: &otlpAnyValue{StringValue: &body},
|
||||
Attributes: toOTLPAttrs(l.Attributes),
|
||||
}
|
||||
if l.SpanContext.IsValid() {
|
||||
rec.TraceID = l.SpanContext.TraceID.String()
|
||||
rec.SpanID = l.SpanContext.SpanID.String()
|
||||
}
|
||||
out = append(out, rec)
|
||||
}
|
||||
return otlpLogsPayload{ResourceLogs: []otlpResourceLogs{{
|
||||
Resource: otlpResource{Attributes: toOTLPAttrs(e.resource)},
|
||||
ScopeLogs: []otlpScopeLogs{{Scope: e.otlpScope(), LogRecords: out}},
|
||||
}}}
|
||||
}
|
||||
|
||||
func (e *OTLPHTTPExporter) otlpScope() otlpScope {
|
||||
return otlpScope{Name: e.scope.Name, Version: e.scope.Version}
|
||||
}
|
||||
|
||||
// toOTLPAttrs converts attributes to their OTLP KeyValue form.
|
||||
func toOTLPAttrs(attrs []KeyValue) []otlpKeyValue {
|
||||
if len(attrs) == 0 {
|
||||
return nil
|
||||
}
|
||||
out := make([]otlpKeyValue, 0, len(attrs))
|
||||
for _, kv := range attrs {
|
||||
out = append(out, otlpKeyValue{Key: kv.Key, Value: toAnyValue(kv.Value)})
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// toAnyValue maps a Go attribute value to an OTLP AnyValue.
|
||||
func toAnyValue(v any) otlpAnyValue {
|
||||
switch x := v.(type) {
|
||||
case string:
|
||||
return otlpAnyValue{StringValue: &x}
|
||||
case bool:
|
||||
return otlpAnyValue{BoolValue: &x}
|
||||
case int64:
|
||||
s := strconv.FormatInt(x, 10)
|
||||
return otlpAnyValue{IntValue: &s}
|
||||
case int:
|
||||
s := strconv.FormatInt(int64(x), 10)
|
||||
return otlpAnyValue{IntValue: &s}
|
||||
case float64:
|
||||
return otlpAnyValue{DoubleValue: &x}
|
||||
default:
|
||||
s := fmt.Sprintf("%v", v)
|
||||
return otlpAnyValue{StringValue: &s}
|
||||
}
|
||||
}
|
||||
|
||||
// nanos formats t as OTLP's Unix-nanoseconds-as-string. A zero time is "0".
|
||||
func nanos(t time.Time) string {
|
||||
if t.IsZero() {
|
||||
return "0"
|
||||
}
|
||||
return strconv.FormatInt(t.UnixNano(), 10)
|
||||
}
|
||||
|
||||
// ParseHeaders parses the OTEL_EXPORTER_OTLP_HEADERS format — a comma-separated
|
||||
// list of key=value pairs, e.g. "Authorization=Bearer abc,x-tenant=libremail" —
|
||||
// into a header map. Whitespace around keys and values is trimmed; malformed
|
||||
// entries (no '=') are skipped.
|
||||
func ParseHeaders(s string) map[string]string {
|
||||
out := map[string]string{}
|
||||
for _, pair := range strings.Split(s, ",") {
|
||||
pair = strings.TrimSpace(pair)
|
||||
if pair == "" {
|
||||
continue
|
||||
}
|
||||
k, v, ok := strings.Cut(pair, "=")
|
||||
k = strings.TrimSpace(k)
|
||||
if !ok || k == "" {
|
||||
continue
|
||||
}
|
||||
out[k] = strings.TrimSpace(v)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
var _ Exporter = (*OTLPHTTPExporter)(nil)
|
||||
@@ -0,0 +1,229 @@
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func sampleSpan() SpanData {
|
||||
start := time.Unix(1_700_000_000, 0).UTC()
|
||||
return SpanData{
|
||||
Name: "publish.report",
|
||||
SpanContext: SpanContext{TraceID: TraceID{15: 0xab}, SpanID: SpanID{7: 0x01}},
|
||||
ParentSpanID: SpanID{7: 0x09},
|
||||
Kind: SpanKindInternal,
|
||||
StartTime: start,
|
||||
EndTime: start.Add(2 * time.Second),
|
||||
Attributes: []KeyValue{
|
||||
String("report.id", "id-1"),
|
||||
Int("publish.deferred", 5),
|
||||
Bool("alert", true),
|
||||
Float("ratio", 0.5),
|
||||
},
|
||||
Status: Status{Code: StatusError, Message: "decrypt failed"},
|
||||
}
|
||||
}
|
||||
|
||||
// TestBuildTracesPayloadJSON asserts the OTLP/JSON trace shape: hex ids,
|
||||
// nanosecond string timestamps, int-as-string, and the resource/scope nesting.
|
||||
func TestBuildTracesPayloadJSON(t *testing.T) {
|
||||
exp := NewOTLPHTTPExporter(OTLPConfig{
|
||||
Endpoint: "https://otlp.example.com/",
|
||||
Resource: []KeyValue{String("service.name", "libremail-bug-report-ingest")},
|
||||
Scope: Scope{Name: "lib", Version: "1.0"},
|
||||
})
|
||||
b, err := json.Marshal(exp.buildTracesPayload([]SpanData{sampleSpan()}))
|
||||
if err != nil {
|
||||
t.Fatalf("marshal: %v", err)
|
||||
}
|
||||
got := string(b)
|
||||
|
||||
// Decode back generically and spot-check the load-bearing fields.
|
||||
var payload struct {
|
||||
ResourceSpans []struct {
|
||||
Resource struct {
|
||||
Attributes []map[string]any `json:"attributes"`
|
||||
} `json:"resource"`
|
||||
ScopeSpans []struct {
|
||||
Scope struct {
|
||||
Name string `json:"name"`
|
||||
Version string `json:"version"`
|
||||
} `json:"scope"`
|
||||
Spans []map[string]any `json:"spans"`
|
||||
} `json:"scopeSpans"`
|
||||
} `json:"resourceSpans"`
|
||||
}
|
||||
if err := json.Unmarshal(b, &payload); err != nil {
|
||||
t.Fatalf("round-trip decode: %v (json=%s)", err, got)
|
||||
}
|
||||
if len(payload.ResourceSpans) != 1 || len(payload.ResourceSpans[0].ScopeSpans) != 1 {
|
||||
t.Fatalf("unexpected nesting: %s", got)
|
||||
}
|
||||
ss := payload.ResourceSpans[0].ScopeSpans[0]
|
||||
if ss.Scope.Name != "lib" || ss.Scope.Version != "1.0" {
|
||||
t.Errorf("scope = %+v", ss.Scope)
|
||||
}
|
||||
if len(ss.Spans) != 1 {
|
||||
t.Fatalf("want 1 span, got %d", len(ss.Spans))
|
||||
}
|
||||
sp := ss.Spans[0]
|
||||
// trace_id and span_id are hex strings in OTLP/JSON (a documented exception
|
||||
// to proto3's base64 default).
|
||||
if sp["traceId"] != "000000000000000000000000000000ab" {
|
||||
t.Errorf("traceId = %v, want hex", sp["traceId"])
|
||||
}
|
||||
if sp["spanId"] != "0000000000000001" {
|
||||
t.Errorf("spanId = %v, want hex", sp["spanId"])
|
||||
}
|
||||
if sp["parentSpanId"] != "0000000000000009" {
|
||||
t.Errorf("parentSpanId = %v", sp["parentSpanId"])
|
||||
}
|
||||
if sp["startTimeUnixNano"] != "1700000000000000000" {
|
||||
t.Errorf("startTimeUnixNano = %v", sp["startTimeUnixNano"])
|
||||
}
|
||||
if sp["endTimeUnixNano"] != "1700000002000000000" {
|
||||
t.Errorf("endTimeUnixNano = %v", sp["endTimeUnixNano"])
|
||||
}
|
||||
// status.code error == 2
|
||||
status, _ := sp["status"].(map[string]any)
|
||||
if status == nil || status["code"] != float64(2) {
|
||||
t.Errorf("status = %v, want code 2", sp["status"])
|
||||
}
|
||||
// int attribute encoded as string.
|
||||
if !strings.Contains(got, `"intValue":"5"`) {
|
||||
t.Errorf("expected intValue as string in %s", got)
|
||||
}
|
||||
if !strings.Contains(got, `"boolValue":true`) {
|
||||
t.Errorf("expected boolValue true in %s", got)
|
||||
}
|
||||
if !strings.Contains(got, `"doubleValue":0.5`) {
|
||||
t.Errorf("expected doubleValue 0.5 in %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildLogsPayloadJSON(t *testing.T) {
|
||||
exp := NewOTLPHTTPExporter(OTLPConfig{Endpoint: "https://otlp.example.com"})
|
||||
rec := LogRecord{
|
||||
Time: time.Unix(1_700_000_001, 0).UTC(),
|
||||
Severity: SeverityError,
|
||||
Body: "publish run failed",
|
||||
Attributes: Alert("publish.run_failed"),
|
||||
SpanContext: SpanContext{TraceID: TraceID{15: 0x02}, SpanID: SpanID{7: 0x03}},
|
||||
}
|
||||
b, err := json.Marshal(exp.buildLogsPayload([]LogRecord{rec}))
|
||||
if err != nil {
|
||||
t.Fatalf("marshal: %v", err)
|
||||
}
|
||||
got := string(b)
|
||||
for _, want := range []string{
|
||||
`"resourceLogs"`,
|
||||
`"severityNumber":17`,
|
||||
`"severityText":"ERROR"`,
|
||||
`"body":{"stringValue":"publish run failed"}`,
|
||||
`"traceId":"00000000000000000000000000000002"`,
|
||||
`"spanId":"0000000000000003"`,
|
||||
`"alert.type"`,
|
||||
} {
|
||||
if !strings.Contains(got, want) {
|
||||
t.Errorf("logs JSON missing %q in %s", want, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestOTLPHTTPExporterRoundTrip drives the exporter against a local httptest
|
||||
// server (NOT a real OTLP backend) to prove it POSTs JSON with the auth header
|
||||
// to the right signal paths.
|
||||
func TestOTLPHTTPExporterRoundTrip(t *testing.T) {
|
||||
var mu sync.Mutex
|
||||
seen := map[string]string{} // path -> body
|
||||
var authHeader, contentType string
|
||||
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
body, _ := io.ReadAll(r.Body)
|
||||
mu.Lock()
|
||||
seen[r.URL.Path] = string(body)
|
||||
authHeader = r.Header.Get("Authorization")
|
||||
contentType = r.Header.Get("Content-Type")
|
||||
mu.Unlock()
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
exp := NewOTLPHTTPExporter(OTLPConfig{
|
||||
Endpoint: srv.URL,
|
||||
Headers: map[string]string{"Authorization": "Bearer secret-token"},
|
||||
Resource: []KeyValue{String("service.name", "svc")},
|
||||
HTTPClient: srv.Client(),
|
||||
})
|
||||
|
||||
if err := exp.ExportSpans(context.Background(), []SpanData{sampleSpan()}); err != nil {
|
||||
t.Fatalf("ExportSpans: %v", err)
|
||||
}
|
||||
if err := exp.ExportLogs(context.Background(), []LogRecord{{
|
||||
Time: time.Now(), Severity: SeverityInfo, Body: "hi",
|
||||
}}); err != nil {
|
||||
t.Fatalf("ExportLogs: %v", err)
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if _, ok := seen["/v1/traces"]; !ok {
|
||||
t.Errorf("no POST to /v1/traces; saw %v", keys(seen))
|
||||
}
|
||||
if _, ok := seen["/v1/logs"]; !ok {
|
||||
t.Errorf("no POST to /v1/logs; saw %v", keys(seen))
|
||||
}
|
||||
if authHeader != "Bearer secret-token" {
|
||||
t.Errorf("Authorization = %q, want the configured secret header", authHeader)
|
||||
}
|
||||
if contentType != "application/json" {
|
||||
t.Errorf("Content-Type = %q, want application/json", contentType)
|
||||
}
|
||||
if !strings.Contains(seen["/v1/traces"], `"publish.report"`) {
|
||||
t.Errorf("traces body missing span name: %s", seen["/v1/traces"])
|
||||
}
|
||||
}
|
||||
|
||||
// TestOTLPExporterNon2xxIsError ensures a non-2xx response surfaces as an error
|
||||
// (which the provider routes to its error handler, never to the caller).
|
||||
func TestOTLPExporterNon2xxIsError(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
}))
|
||||
defer srv.Close()
|
||||
exp := NewOTLPHTTPExporter(OTLPConfig{Endpoint: srv.URL, HTTPClient: srv.Client()})
|
||||
if err := exp.ExportSpans(context.Background(), []SpanData{sampleSpan()}); err == nil {
|
||||
t.Error("expected an error on a 500 response")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseHeaders(t *testing.T) {
|
||||
got := ParseHeaders("Authorization=Bearer abc, x-tenant = libremail ,,bad")
|
||||
if got["Authorization"] != "Bearer abc" {
|
||||
t.Errorf("Authorization = %q", got["Authorization"])
|
||||
}
|
||||
if got["x-tenant"] != "libremail" {
|
||||
t.Errorf("x-tenant = %q", got["x-tenant"])
|
||||
}
|
||||
if len(got) != 2 {
|
||||
t.Errorf("got %d headers, want 2: %v", len(got), got)
|
||||
}
|
||||
if len(ParseHeaders("")) != 0 {
|
||||
t.Error("empty string should yield no headers")
|
||||
}
|
||||
}
|
||||
|
||||
func keys(m map[string]string) []string {
|
||||
out := make([]string, 0, len(m))
|
||||
for k := range m {
|
||||
out = append(out, k)
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -0,0 +1,502 @@
|
||||
// Package telemetry is a minimal OpenTelemetry (OTEL) tracing + structured
|
||||
// logging shim for the LibreMail bug-report ingest Worker (issue #17).
|
||||
//
|
||||
// # Why a hand-rolled shim instead of the OpenTelemetry-Go SDK
|
||||
//
|
||||
// The Worker is compiled to Wasm by TinyGo. The full OpenTelemetry-Go SDK
|
||||
// (go.opentelemetry.io/otel/sdk + the OTLP exporters) pulls in a large
|
||||
// dependency tree — google.golang.org/protobuf, grpc, reflection-heavy option
|
||||
// plumbing — that bloats the Wasm binary and does not reliably compile under
|
||||
// TinyGo. This package instead implements just the slice of OTEL the Worker
|
||||
// needs: spans, structured log records, and OTLP/HTTP export, using only the
|
||||
// standard library (net/http, encoding/json, crypto/rand — all already proven
|
||||
// to compile and run under this project's js/wasm target, see internal/publish
|
||||
// and internal/crypto). The wire format is OTLP so any OTLP-compatible backend
|
||||
// can ingest it; only the SDK is hand-rolled, not the protocol.
|
||||
//
|
||||
// # Behaviour-preserving by construction
|
||||
//
|
||||
// Instrumentation is threaded through context. Instrumented code pulls an
|
||||
// optional [*Telemetry] from the context with [FromContext]; when it is absent
|
||||
// (or its exporter is nil) every method is a no-op and the surrounding code
|
||||
// behaves exactly as it did before. This is what lets the ingest handler, the
|
||||
// publish job, and the schedule gate be instrumented without changing any of
|
||||
// their public constructor/function signatures (a hard requirement so parallel
|
||||
// work that constructs those components via their current APIs keeps compiling).
|
||||
//
|
||||
// # Host-testable core
|
||||
//
|
||||
// The package carries no build constraints, so it is unit-tested with the
|
||||
// standard toolchain against [MemoryExporter], an in-memory test exporter. The
|
||||
// production OTLP/HTTP exporter ([OTLPHTTPExporter]) is exercised in tests
|
||||
// against an httptest server, never a real backend.
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// SpanKind mirrors the OTLP span kind enum (opentelemetry-proto SpanKind).
|
||||
type SpanKind int
|
||||
|
||||
const (
|
||||
// SpanKindInternal is an internal operation within an application.
|
||||
SpanKindInternal SpanKind = 1
|
||||
// SpanKindServer is a synchronous inbound (server-side) request span.
|
||||
SpanKindServer SpanKind = 2
|
||||
// SpanKindClient is a synchronous outbound (client-side) request span.
|
||||
SpanKindClient SpanKind = 3
|
||||
)
|
||||
|
||||
// StatusCode mirrors the OTLP Status.StatusCode enum.
|
||||
type StatusCode int
|
||||
|
||||
const (
|
||||
// StatusUnset is the default, un-set span status.
|
||||
StatusUnset StatusCode = 0
|
||||
// StatusOK marks a span the application explicitly considers successful.
|
||||
StatusOK StatusCode = 1
|
||||
// StatusError marks a span the application considers failed; backends can
|
||||
// alert on it.
|
||||
StatusError StatusCode = 2
|
||||
)
|
||||
|
||||
// Status is a span's completion status.
|
||||
type Status struct {
|
||||
Code StatusCode
|
||||
Message string
|
||||
}
|
||||
|
||||
// Severity mirrors the OTLP LogRecord SeverityNumber values used by this Worker.
|
||||
type Severity int
|
||||
|
||||
const (
|
||||
// SeverityInfo is OTEL severity number 9 (INFO).
|
||||
SeverityInfo Severity = 9
|
||||
// SeverityWarn is OTEL severity number 13 (WARN).
|
||||
SeverityWarn Severity = 13
|
||||
// SeverityError is OTEL severity number 17 (ERROR).
|
||||
SeverityError Severity = 17
|
||||
)
|
||||
|
||||
// String returns the OTEL severity text ("INFO"/"WARN"/"ERROR").
|
||||
func (s Severity) String() string {
|
||||
switch {
|
||||
case s >= SeverityError:
|
||||
return "ERROR"
|
||||
case s >= SeverityWarn:
|
||||
return "WARN"
|
||||
default:
|
||||
return "INFO"
|
||||
}
|
||||
}
|
||||
|
||||
// Alert attribute keys. A telemetry backend can raise an alert on the presence
|
||||
// of AttrAlert=true (any alert-worthy signal) or on a specific AttrAlertType
|
||||
// value. These are the hooks issue #17's alerting keys on once an OTLP backend
|
||||
// is chosen (that choice is TBD); the signals are emitted here so the wiring is
|
||||
// ready.
|
||||
const (
|
||||
// AttrAlert is a boolean attribute; true marks an alert-worthy signal.
|
||||
AttrAlert = "alert"
|
||||
// AttrAlertType is a string attribute naming the specific signal, e.g.
|
||||
// "publish.cap_hit" or "ingest.error".
|
||||
AttrAlertType = "alert.type"
|
||||
)
|
||||
|
||||
// Alert returns the attribute pair that marks a log record or span as an
|
||||
// alert-worthy signal of the given type. Callers append their own context
|
||||
// attributes (counts, ids, ...) after it.
|
||||
func Alert(alertType string) []KeyValue {
|
||||
return []KeyValue{Bool(AttrAlert, true), String(AttrAlertType, alertType)}
|
||||
}
|
||||
|
||||
// KeyValue is a single attribute. Value holds one of string, int64, bool, or
|
||||
// float64; the OTLP exporter maps it to the matching AnyValue variant. Using a
|
||||
// bare any keeps construction ergonomic and lets tests compare Value directly.
|
||||
type KeyValue struct {
|
||||
Key string
|
||||
Value any
|
||||
}
|
||||
|
||||
// String, Int, Int64, Bool, and Float build typed attributes.
|
||||
func String(k, v string) KeyValue { return KeyValue{Key: k, Value: v} }
|
||||
func Int(k string, v int) KeyValue { return KeyValue{Key: k, Value: int64(v)} }
|
||||
func Int64(k string, v int64) KeyValue {
|
||||
return KeyValue{Key: k, Value: v}
|
||||
}
|
||||
func Bool(k string, v bool) KeyValue { return KeyValue{Key: k, Value: v} }
|
||||
func Float(k string, v float64) KeyValue { return KeyValue{Key: k, Value: v} }
|
||||
|
||||
// TraceID is a 16-byte W3C trace id.
|
||||
type TraceID [16]byte
|
||||
|
||||
// SpanID is an 8-byte W3C span id.
|
||||
type SpanID [8]byte
|
||||
|
||||
// String returns the lowercase hex encoding used by OTLP/JSON.
|
||||
func (t TraceID) String() string { return hex.EncodeToString(t[:]) }
|
||||
|
||||
// String returns the lowercase hex encoding used by OTLP/JSON.
|
||||
func (s SpanID) String() string { return hex.EncodeToString(s[:]) }
|
||||
|
||||
// IsZero reports whether the id is all-zero (unset).
|
||||
func (t TraceID) IsZero() bool { return t == TraceID{} }
|
||||
|
||||
// IsZero reports whether the id is all-zero (unset).
|
||||
func (s SpanID) IsZero() bool { return s == SpanID{} }
|
||||
|
||||
// SpanContext is the trace-correlation identity of a span: its trace id and its
|
||||
// own span id. It is what rides in context so children can parent themselves and
|
||||
// log records can correlate to the active span.
|
||||
type SpanContext struct {
|
||||
TraceID TraceID
|
||||
SpanID SpanID
|
||||
}
|
||||
|
||||
// IsValid reports whether both ids are set.
|
||||
func (sc SpanContext) IsValid() bool { return !sc.TraceID.IsZero() && !sc.SpanID.IsZero() }
|
||||
|
||||
// Event is a timestamped event recorded on a span.
|
||||
type Event struct {
|
||||
Name string
|
||||
Time time.Time
|
||||
Attributes []KeyValue
|
||||
}
|
||||
|
||||
// SpanData is a completed span handed to an [Exporter] at End.
|
||||
type SpanData struct {
|
||||
Name string
|
||||
SpanContext SpanContext
|
||||
ParentSpanID SpanID
|
||||
Kind SpanKind
|
||||
StartTime time.Time
|
||||
EndTime time.Time
|
||||
Attributes []KeyValue
|
||||
Events []Event
|
||||
Status Status
|
||||
}
|
||||
|
||||
// LogRecord is a structured log record handed to an [Exporter]. SpanContext is
|
||||
// the active span at emit time (zero if none), so records correlate to traces.
|
||||
type LogRecord struct {
|
||||
Time time.Time
|
||||
Severity Severity
|
||||
Body string
|
||||
Attributes []KeyValue
|
||||
SpanContext SpanContext
|
||||
}
|
||||
|
||||
// Exporter is the sink for finished spans and log records. The Worker wires
|
||||
// [OTLPHTTPExporter]; host tests wire [MemoryExporter]. Implementations must be
|
||||
// safe for concurrent use.
|
||||
type Exporter interface {
|
||||
ExportSpans(ctx context.Context, spans []SpanData) error
|
||||
ExportLogs(ctx context.Context, logs []LogRecord) error
|
||||
// Shutdown releases any resources; exporters that hold none return nil.
|
||||
Shutdown(ctx context.Context) error
|
||||
}
|
||||
|
||||
// Telemetry is the tracer + logger provider. A nil *Telemetry, or one built with
|
||||
// a nil exporter, is a valid no-op provider: every method is safe to call and
|
||||
// does nothing, which is the default when OTEL is not configured (the OTLP
|
||||
// endpoint is TBD, issue #17).
|
||||
type Telemetry struct {
|
||||
exporter Exporter
|
||||
now func() time.Time
|
||||
newTrace func() TraceID
|
||||
newSpan func() SpanID
|
||||
onErr func(error)
|
||||
}
|
||||
|
||||
// Option customises a [Telemetry].
|
||||
type Option func(*Telemetry)
|
||||
|
||||
// WithClock injects the time source (tests use a fixed clock for determinism).
|
||||
func WithClock(now func() time.Time) Option {
|
||||
return func(t *Telemetry) {
|
||||
if now != nil {
|
||||
t.now = now
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// WithIDGenerator injects the trace/span id sources (tests use counters so ids
|
||||
// are deterministic). Either function may be nil to keep the default.
|
||||
func WithIDGenerator(newTrace func() TraceID, newSpan func() SpanID) Option {
|
||||
return func(t *Telemetry) {
|
||||
if newTrace != nil {
|
||||
t.newTrace = newTrace
|
||||
}
|
||||
if newSpan != nil {
|
||||
t.newSpan = newSpan
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// WithErrorHandler sets a handler for export errors. Telemetry is best-effort:
|
||||
// export failures never propagate to the instrumented operation, but the Worker
|
||||
// can wire this to a log so misconfiguration is visible. The default swallows.
|
||||
func WithErrorHandler(fn func(error)) Option {
|
||||
return func(t *Telemetry) {
|
||||
if fn != nil {
|
||||
t.onErr = fn
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// New returns a Telemetry exporting to exp. A nil exp yields a no-op provider
|
||||
// (so callers can always construct one and let configuration decide whether it
|
||||
// does anything).
|
||||
func New(exp Exporter, opts ...Option) *Telemetry {
|
||||
t := &Telemetry{
|
||||
exporter: exp,
|
||||
now: time.Now,
|
||||
newTrace: randomTraceID,
|
||||
newSpan: randomSpanID,
|
||||
onErr: func(error) {},
|
||||
}
|
||||
for _, o := range opts {
|
||||
o(t)
|
||||
}
|
||||
return t
|
||||
}
|
||||
|
||||
// Enabled reports whether this provider will actually record anything. It is
|
||||
// safe on a nil receiver. Instrumented hot paths can branch on it to skip work
|
||||
// entirely when telemetry is off.
|
||||
func (t *Telemetry) Enabled() bool { return t != nil && t.exporter != nil }
|
||||
|
||||
// StartSpan begins a span named name and returns a context carrying its
|
||||
// [SpanContext] (so nested StartSpan calls and log records correlate) plus the
|
||||
// span. When the provider is disabled it returns ctx unchanged and a nil *Span
|
||||
// whose methods are all no-ops — so callers need no telemetry-on/off branching.
|
||||
func (t *Telemetry) StartSpan(ctx context.Context, name string, opts ...SpanOption) (context.Context, *Span) {
|
||||
if !t.Enabled() {
|
||||
return ctx, nil
|
||||
}
|
||||
cfg := spanConfig{kind: SpanKindInternal}
|
||||
for _, o := range opts {
|
||||
o(&cfg)
|
||||
}
|
||||
parent, _ := SpanContextFromContext(ctx)
|
||||
traceID := parent.TraceID
|
||||
if traceID.IsZero() {
|
||||
traceID = t.newTrace()
|
||||
}
|
||||
sc := SpanContext{TraceID: traceID, SpanID: t.newSpan()}
|
||||
s := &Span{
|
||||
tel: t,
|
||||
sc: sc,
|
||||
parent: parent.SpanID,
|
||||
name: name,
|
||||
kind: cfg.kind,
|
||||
start: t.now(),
|
||||
attrs: append([]KeyValue(nil), cfg.attrs...),
|
||||
}
|
||||
return ContextWithSpanContext(ctx, sc), s
|
||||
}
|
||||
|
||||
// Log emits a structured log record at severity sev, correlated to the active
|
||||
// span in ctx (if any). It is safe (and a no-op) on a disabled provider.
|
||||
func (t *Telemetry) Log(ctx context.Context, sev Severity, body string, attrs ...KeyValue) {
|
||||
if !t.Enabled() {
|
||||
return
|
||||
}
|
||||
sc, _ := SpanContextFromContext(ctx)
|
||||
rec := LogRecord{
|
||||
Time: t.now(),
|
||||
Severity: sev,
|
||||
Body: body,
|
||||
Attributes: append([]KeyValue(nil), attrs...),
|
||||
SpanContext: sc,
|
||||
}
|
||||
if err := t.exporter.ExportLogs(context.Background(), []LogRecord{rec}); err != nil {
|
||||
t.onErr(err)
|
||||
}
|
||||
}
|
||||
|
||||
// Info, Warn, and Error are severity-specific [Telemetry.Log] wrappers.
|
||||
func (t *Telemetry) Info(ctx context.Context, body string, attrs ...KeyValue) {
|
||||
t.Log(ctx, SeverityInfo, body, attrs...)
|
||||
}
|
||||
func (t *Telemetry) Warn(ctx context.Context, body string, attrs ...KeyValue) {
|
||||
t.Log(ctx, SeverityWarn, body, attrs...)
|
||||
}
|
||||
func (t *Telemetry) Error(ctx context.Context, body string, attrs ...KeyValue) {
|
||||
t.Log(ctx, SeverityError, body, attrs...)
|
||||
}
|
||||
|
||||
// SpanOption customises a span at start.
|
||||
type SpanOption func(*spanConfig)
|
||||
|
||||
type spanConfig struct {
|
||||
kind SpanKind
|
||||
attrs []KeyValue
|
||||
}
|
||||
|
||||
// WithSpanKind sets the span kind (default [SpanKindInternal]).
|
||||
func WithSpanKind(k SpanKind) SpanOption {
|
||||
return func(c *spanConfig) { c.kind = k }
|
||||
}
|
||||
|
||||
// WithAttributes sets initial span attributes.
|
||||
func WithAttributes(attrs ...KeyValue) SpanOption {
|
||||
return func(c *spanConfig) { c.attrs = append(c.attrs, attrs...) }
|
||||
}
|
||||
|
||||
// Span is an in-progress span. All methods are safe on a nil *Span (the value
|
||||
// returned by a disabled provider's StartSpan), so instrumented code never needs
|
||||
// to nil-check.
|
||||
type Span struct {
|
||||
tel *Telemetry
|
||||
sc SpanContext
|
||||
parent SpanID
|
||||
name string
|
||||
kind SpanKind
|
||||
start time.Time
|
||||
|
||||
mu sync.Mutex
|
||||
attrs []KeyValue
|
||||
events []Event
|
||||
status Status
|
||||
ended bool
|
||||
}
|
||||
|
||||
// SpanContext returns the span's trace-correlation identity. On a nil span it
|
||||
// returns the zero value.
|
||||
func (s *Span) SpanContext() SpanContext {
|
||||
if s == nil {
|
||||
return SpanContext{}
|
||||
}
|
||||
return s.sc
|
||||
}
|
||||
|
||||
// SetAttributes adds attributes to the span.
|
||||
func (s *Span) SetAttributes(attrs ...KeyValue) {
|
||||
if s == nil || len(attrs) == 0 {
|
||||
return
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.attrs = append(s.attrs, attrs...)
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// SetStatus sets the span's completion status.
|
||||
func (s *Span) SetStatus(code StatusCode, msg string) {
|
||||
if s == nil {
|
||||
return
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.status = Status{Code: code, Message: msg}
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// RecordError marks the span failed and records the error message as an event
|
||||
// and status. A nil error or span is ignored.
|
||||
func (s *Span) RecordError(err error) {
|
||||
if s == nil || err == nil {
|
||||
return
|
||||
}
|
||||
s.AddEvent("exception", String("exception.message", err.Error()))
|
||||
s.SetStatus(StatusError, err.Error())
|
||||
}
|
||||
|
||||
// AddEvent records a timestamped event on the span.
|
||||
func (s *Span) AddEvent(name string, attrs ...KeyValue) {
|
||||
if s == nil {
|
||||
return
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.events = append(s.events, Event{Name: name, Time: s.tel.now(), Attributes: append([]KeyValue(nil), attrs...)})
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// End finishes the span and exports it. It is idempotent; a nil span is ignored.
|
||||
func (s *Span) End() {
|
||||
if s == nil {
|
||||
return
|
||||
}
|
||||
s.mu.Lock()
|
||||
if s.ended {
|
||||
s.mu.Unlock()
|
||||
return
|
||||
}
|
||||
s.ended = true
|
||||
data := SpanData{
|
||||
Name: s.name,
|
||||
SpanContext: s.sc,
|
||||
ParentSpanID: s.parent,
|
||||
Kind: s.kind,
|
||||
StartTime: s.start,
|
||||
EndTime: s.tel.now(),
|
||||
Attributes: append([]KeyValue(nil), s.attrs...),
|
||||
Events: append([]Event(nil), s.events...),
|
||||
Status: s.status,
|
||||
}
|
||||
s.mu.Unlock()
|
||||
if err := s.tel.exporter.ExportSpans(context.Background(), []SpanData{data}); err != nil {
|
||||
s.tel.onErr(err)
|
||||
}
|
||||
}
|
||||
|
||||
// --- context propagation ---
|
||||
|
||||
type telemetryKeyType struct{}
|
||||
type spanContextKeyType struct{}
|
||||
|
||||
var telemetryKey telemetryKeyType
|
||||
var spanContextKey spanContextKeyType
|
||||
|
||||
// NewContext returns a context carrying the telemetry provider. Instrumented
|
||||
// code deeper in the call tree retrieves it with [FromContext]. Passing a nil
|
||||
// provider is fine — it yields the no-op behaviour.
|
||||
func NewContext(ctx context.Context, t *Telemetry) context.Context {
|
||||
return context.WithValue(ctx, telemetryKey, t)
|
||||
}
|
||||
|
||||
// FromContext returns the telemetry provider carried by ctx, or nil if none.
|
||||
// A nil result is a valid no-op provider (all methods are nil-safe).
|
||||
func FromContext(ctx context.Context) *Telemetry {
|
||||
if ctx == nil {
|
||||
return nil
|
||||
}
|
||||
t, _ := ctx.Value(telemetryKey).(*Telemetry)
|
||||
return t
|
||||
}
|
||||
|
||||
// ContextWithSpanContext returns a context whose active span is sc (used as the
|
||||
// parent of the next StartSpan and the correlation id of Log calls).
|
||||
func ContextWithSpanContext(ctx context.Context, sc SpanContext) context.Context {
|
||||
return context.WithValue(ctx, spanContextKey, sc)
|
||||
}
|
||||
|
||||
// SpanContextFromContext returns the active [SpanContext] in ctx.
|
||||
func SpanContextFromContext(ctx context.Context) (SpanContext, bool) {
|
||||
if ctx == nil {
|
||||
return SpanContext{}, false
|
||||
}
|
||||
sc, ok := ctx.Value(spanContextKey).(SpanContext)
|
||||
return sc, ok
|
||||
}
|
||||
|
||||
// --- default id generation ---
|
||||
|
||||
// randomTraceID / randomSpanID draw ids from crypto/rand. crypto/rand.Read is
|
||||
// used elsewhere in build-tag-free code that compiles into the Wasm Worker
|
||||
// (internal/crypto, internal/storage), so it is available under TinyGo/Wasm.
|
||||
func randomTraceID() TraceID {
|
||||
var t TraceID
|
||||
_, _ = rand.Read(t[:])
|
||||
return t
|
||||
}
|
||||
|
||||
func randomSpanID() SpanID {
|
||||
var s SpanID
|
||||
_, _ = rand.Read(s[:])
|
||||
return s
|
||||
}
|
||||
@@ -0,0 +1,239 @@
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// fixedClock returns a clock function yielding a fixed instant.
|
||||
func fixedClock(t time.Time) func() time.Time { return func() time.Time { return t } }
|
||||
|
||||
// seqIDs returns deterministic id generators: trace ids 0x..01, 0x..02, ...;
|
||||
// span ids likewise, so tests can assert exact correlation.
|
||||
func seqIDs() (func() TraceID, func() SpanID) {
|
||||
var tn, sn byte
|
||||
return func() TraceID {
|
||||
tn++
|
||||
return TraceID{15: tn}
|
||||
}, func() SpanID {
|
||||
sn++
|
||||
return SpanID{7: sn}
|
||||
}
|
||||
}
|
||||
|
||||
func testProvider(t *testing.T) (*Telemetry, *MemoryExporter) {
|
||||
t.Helper()
|
||||
exp := NewMemoryExporter()
|
||||
nt, ns := seqIDs()
|
||||
tel := New(exp,
|
||||
WithClock(fixedClock(time.Unix(1_700_000_000, 0).UTC())),
|
||||
WithIDGenerator(nt, ns))
|
||||
return tel, exp
|
||||
}
|
||||
|
||||
func TestStartSpanExportsSpanData(t *testing.T) {
|
||||
tel, exp := testProvider(t)
|
||||
ctx, span := tel.StartSpan(context.Background(), "ingest.request",
|
||||
WithSpanKind(SpanKindServer),
|
||||
WithAttributes(String("http.request.method", "POST")))
|
||||
span.SetAttributes(Int("http.response.status_code", 202))
|
||||
span.SetStatus(StatusOK, "")
|
||||
span.End()
|
||||
_ = ctx
|
||||
|
||||
spans := exp.Spans()
|
||||
if len(spans) != 1 {
|
||||
t.Fatalf("got %d spans, want 1", len(spans))
|
||||
}
|
||||
s := spans[0]
|
||||
if s.Name != "ingest.request" {
|
||||
t.Errorf("name = %q, want ingest.request", s.Name)
|
||||
}
|
||||
if s.Kind != SpanKindServer {
|
||||
t.Errorf("kind = %d, want %d", s.Kind, SpanKindServer)
|
||||
}
|
||||
if !s.SpanContext.IsValid() {
|
||||
t.Errorf("span context invalid: %+v", s.SpanContext)
|
||||
}
|
||||
if !s.ParentSpanID.IsZero() {
|
||||
t.Errorf("root span should have zero parent, got %s", s.ParentSpanID)
|
||||
}
|
||||
if s.Status.Code != StatusOK {
|
||||
t.Errorf("status = %d, want OK", s.Status.Code)
|
||||
}
|
||||
if v, ok := Attr(s.Attributes, "http.request.method"); !ok || v != "POST" {
|
||||
t.Errorf("method attr = %v (present=%v), want POST", v, ok)
|
||||
}
|
||||
if v, ok := Attr(s.Attributes, "http.response.status_code"); !ok || v != int64(202) {
|
||||
t.Errorf("status_code attr = %v (present=%v), want int64 202", v, ok)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNestedSpanSharesTraceAndParents(t *testing.T) {
|
||||
tel, exp := testProvider(t)
|
||||
ctx, parent := tel.StartSpan(context.Background(), "publish.run")
|
||||
_, child := tel.StartSpan(ctx, "publish.report")
|
||||
child.End()
|
||||
parent.End()
|
||||
|
||||
spans := exp.Spans()
|
||||
if len(spans) != 2 {
|
||||
t.Fatalf("got %d spans, want 2", len(spans))
|
||||
}
|
||||
// Export order: child ended first.
|
||||
c, p := spans[0], spans[1]
|
||||
if c.SpanContext.TraceID != p.SpanContext.TraceID {
|
||||
t.Errorf("child/parent trace ids differ: %s vs %s", c.SpanContext.TraceID, p.SpanContext.TraceID)
|
||||
}
|
||||
if c.ParentSpanID != p.SpanContext.SpanID {
|
||||
t.Errorf("child parent = %s, want parent span id %s", c.ParentSpanID, p.SpanContext.SpanID)
|
||||
}
|
||||
if p.SpanContext.SpanID == c.SpanContext.SpanID {
|
||||
t.Errorf("child and parent share a span id: %s", c.SpanContext.SpanID)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLogCorrelatesWithActiveSpan(t *testing.T) {
|
||||
tel, exp := testProvider(t)
|
||||
ctx, span := tel.StartSpan(context.Background(), "ingest.request")
|
||||
tel.Info(ctx, "ingest request accepted", String("ingest.outcome", "accepted"))
|
||||
span.End()
|
||||
|
||||
logs := exp.Logs()
|
||||
if len(logs) != 1 {
|
||||
t.Fatalf("got %d logs, want 1", len(logs))
|
||||
}
|
||||
l := logs[0]
|
||||
if l.Severity != SeverityInfo {
|
||||
t.Errorf("severity = %v, want INFO", l.Severity)
|
||||
}
|
||||
if l.Body != "ingest request accepted" {
|
||||
t.Errorf("body = %q", l.Body)
|
||||
}
|
||||
if l.SpanContext != span.SpanContext() {
|
||||
t.Errorf("log span context %+v != span %+v (not correlated)", l.SpanContext, span.SpanContext())
|
||||
}
|
||||
if v, ok := Attr(l.Attributes, "ingest.outcome"); !ok || v != "accepted" {
|
||||
t.Errorf("outcome attr = %v (present=%v)", v, ok)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLogWithoutSpanHasNoCorrelation(t *testing.T) {
|
||||
tel, exp := testProvider(t)
|
||||
tel.Warn(context.Background(), "standalone warning")
|
||||
logs := exp.Logs()
|
||||
if len(logs) != 1 {
|
||||
t.Fatalf("got %d logs, want 1", len(logs))
|
||||
}
|
||||
if logs[0].SpanContext.IsValid() {
|
||||
t.Errorf("log outside a span should not correlate, got %+v", logs[0].SpanContext)
|
||||
}
|
||||
if logs[0].Severity != SeverityWarn {
|
||||
t.Errorf("severity = %v, want WARN", logs[0].Severity)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecordError(t *testing.T) {
|
||||
tel, exp := testProvider(t)
|
||||
_, span := tel.StartSpan(context.Background(), "op")
|
||||
span.RecordError(errors.New("boom"))
|
||||
span.End()
|
||||
s := exp.Spans()[0]
|
||||
if s.Status.Code != StatusError {
|
||||
t.Errorf("status = %d, want Error", s.Status.Code)
|
||||
}
|
||||
if len(s.Events) != 1 || s.Events[0].Name != "exception" {
|
||||
t.Fatalf("want one exception event, got %+v", s.Events)
|
||||
}
|
||||
if v, ok := Attr(s.Events[0].Attributes, "exception.message"); !ok || v != "boom" {
|
||||
t.Errorf("exception.message = %v (present=%v)", v, ok)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDisabledProviderIsNoop is the behaviour-preserving guarantee: a nil
|
||||
// provider (what FromContext returns when telemetry is not configured) and a
|
||||
// New(nil) provider must be fully no-op and never panic, and StartSpan must
|
||||
// return the context unchanged so downstream behaviour is identical.
|
||||
func TestDisabledProviderIsNoop(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
tel *Telemetry
|
||||
}{
|
||||
{"nil provider", nil},
|
||||
{"nil exporter", New(nil)},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
if tc.tel.Enabled() {
|
||||
t.Fatal("provider should report disabled")
|
||||
}
|
||||
base := context.Background()
|
||||
ctx, span := tc.tel.StartSpan(base, "op", WithSpanKind(SpanKindServer))
|
||||
if ctx != base {
|
||||
t.Error("disabled StartSpan must return the context unchanged")
|
||||
}
|
||||
if span != nil {
|
||||
t.Error("disabled StartSpan must return a nil span")
|
||||
}
|
||||
// None of these must panic on the nil span / disabled provider.
|
||||
span.SetAttributes(String("k", "v"))
|
||||
span.SetStatus(StatusError, "x")
|
||||
span.RecordError(errors.New("e"))
|
||||
span.AddEvent("ev")
|
||||
span.End()
|
||||
if sc := span.SpanContext(); sc.IsValid() {
|
||||
t.Error("nil span should have an invalid span context")
|
||||
}
|
||||
tc.tel.Info(ctx, "msg")
|
||||
tc.tel.Warn(ctx, "msg")
|
||||
tc.tel.Error(ctx, "msg")
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestFromContextRoundTrip(t *testing.T) {
|
||||
tel, _ := testProvider(t)
|
||||
ctx := NewContext(context.Background(), tel)
|
||||
if got := FromContext(ctx); got != tel {
|
||||
t.Errorf("FromContext = %p, want %p", got, tel)
|
||||
}
|
||||
if got := FromContext(context.Background()); got != nil {
|
||||
t.Errorf("FromContext on a bare context = %p, want nil", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAlertHelper(t *testing.T) {
|
||||
kvs := Alert("publish.cap_hit")
|
||||
if len(kvs) != 2 {
|
||||
t.Fatalf("Alert returned %d attrs, want 2", len(kvs))
|
||||
}
|
||||
if v, ok := Attr(kvs, AttrAlert); !ok || v != true {
|
||||
t.Errorf("alert marker = %v (present=%v), want true", v, ok)
|
||||
}
|
||||
if v, ok := Attr(kvs, AttrAlertType); !ok || v != "publish.cap_hit" {
|
||||
t.Errorf("alert.type = %v (present=%v)", v, ok)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSeverityString(t *testing.T) {
|
||||
for sev, want := range map[Severity]string{
|
||||
SeverityInfo: "INFO",
|
||||
SeverityWarn: "WARN",
|
||||
SeverityError: "ERROR",
|
||||
} {
|
||||
if got := sev.String(); got != want {
|
||||
t.Errorf("Severity(%d).String() = %q, want %q", sev, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestEndIsIdempotent(t *testing.T) {
|
||||
tel, exp := testProvider(t)
|
||||
_, span := tel.StartSpan(context.Background(), "op")
|
||||
span.End()
|
||||
span.End() // second End must not double-export
|
||||
if got := len(exp.Spans()); got != 1 {
|
||||
t.Errorf("got %d spans after double End, want 1", got)
|
||||
}
|
||||
}
|
||||
+6
-1
@@ -25,5 +25,10 @@ func main() {
|
||||
// The maintainer admin API (#11) is wired with workerAdminBackend, which reads
|
||||
// the admin shared secret (ADMIN_TOKEN) from Secrets Store and drives the
|
||||
// lifecycle Manager over R2 per request (see admin.go).
|
||||
workers.Serve(handler.New(storage.NewWorkerSink(), workerAdminBackend{}))
|
||||
//
|
||||
// withWorkerTelemetry (#17) injects the OTEL telemetry provider into each
|
||||
// request's context so the ingest handler emits spans + structured logs. It is
|
||||
// a no-op until an OTLP endpoint is configured (the endpoint is TBD), so this
|
||||
// wrapper does not change behaviour by default.
|
||||
workers.Serve(withWorkerTelemetry(handler.New(storage.NewWorkerSink(), workerAdminBackend{})))
|
||||
}
|
||||
|
||||
@@ -82,6 +82,12 @@ func runWeeklyTrigger(ctx context.Context) error {
|
||||
return fmt.Errorf("schedule: build publisher: %w", err)
|
||||
}
|
||||
|
||||
// Observability (#17): carry the OTEL telemetry provider in the context so the
|
||||
// gated run emits a schedule.run span, the publisher emits a publish.run span
|
||||
// with per-report child spans, and a failed run / cap-hit surfaces as an
|
||||
// alertable signal. A no-op until an OTLP endpoint is configured (TBD).
|
||||
ctx = scheduledContext(ctx)
|
||||
|
||||
// schedule.Run re-checks the gate (open here) and hands the pending ids to the
|
||||
// publisher, which decrypts, formats, and creates one labeled GitHub issue per
|
||||
// report. Cross-run de-dup is wired here (#15): the publisher's onPublished hook
|
||||
|
||||
@@ -0,0 +1,120 @@
|
||||
//go:build js && wasm
|
||||
|
||||
// This file wires OpenTelemetry (#17) into the Cloudflare Worker: it builds the
|
||||
// OTLP/HTTP exporter from Worker config and injects the telemetry provider into
|
||||
// the fetch request context so internal/ingest (and, via scheduled_wasm.go,
|
||||
// internal/schedule + internal/publish) emit spans and structured logs.
|
||||
//
|
||||
// The OTLP endpoint/backend is deliberately TBD (issue #17): it is read from a
|
||||
// plain Worker var and the auth header from a Secrets Store secret, never
|
||||
// hardcoded. When the endpoint var is unset the provider is nil and all
|
||||
// instrumentation is a no-op, so the Worker runs exactly as before until an
|
||||
// endpoint is configured.
|
||||
//
|
||||
// Compiled only into the js/wasm Worker; excluded from host builds and tests.
|
||||
// The non-trivial telemetry logic lives in the build-tag-free, host-tested
|
||||
// internal/telemetry package; this file is the thin runtime adapter that reads
|
||||
// the Cloudflare bindings.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"net/http"
|
||||
"sync"
|
||||
|
||||
"github.com/syumai/workers/cloudflare"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/storage"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/telemetry"
|
||||
)
|
||||
|
||||
// OTEL config names, declared in wrangler.jsonc. The endpoint and service name
|
||||
// are plain (non-secret) vars; the headers (which carry auth) are a Secrets
|
||||
// Store secret read via binding.get(), like the encryption keyring.
|
||||
const (
|
||||
// otelEndpointVar is the base OTLP/HTTP URL, e.g. "https://otlp.example.com".
|
||||
// TBD (#17): empty/unset => telemetry disabled (no-op).
|
||||
otelEndpointVar = "OTEL_EXPORTER_OTLP_ENDPOINT"
|
||||
// otelServiceVar overrides the reported service.name.
|
||||
otelServiceVar = "OTEL_SERVICE_NAME"
|
||||
// otelHeadersSecret is the Secrets Store secret holding the OTLP auth
|
||||
// headers in "key=value,key2=value2" form (e.g. "Authorization=Bearer ...").
|
||||
otelHeadersSecret = "OTEL_EXPORTER_OTLP_HEADERS"
|
||||
|
||||
otelServiceDefault = "libremail-bug-report-ingest"
|
||||
otelScopeName = "github.com/JMR-dev/LibreMail-Bug-Report-Ingest"
|
||||
)
|
||||
|
||||
var (
|
||||
telMu sync.Mutex
|
||||
telProvider *telemetry.Telemetry // cached for the isolate once built
|
||||
telBuilt bool
|
||||
)
|
||||
|
||||
// workerTelemetry lazily builds the telemetry provider from Worker config and
|
||||
// caches it for the isolate lifetime. It returns nil (a valid no-op provider)
|
||||
// when OTEL_EXPORTER_OTLP_ENDPOINT is unset — the endpoint is TBD (#17).
|
||||
//
|
||||
// It must be called from within a request or scheduled handler: the headers
|
||||
// secret read is async and per-request on the Workers runtime (mirroring how the
|
||||
// encryption keyring and admin token are read). The auth header is optional — if
|
||||
// the secret is unavailable the exporter is still built (some backends accept
|
||||
// unauthenticated ingest, or auth may live in the endpoint URL); the failure is
|
||||
// logged, never fatal.
|
||||
func workerTelemetry() *telemetry.Telemetry {
|
||||
telMu.Lock()
|
||||
defer telMu.Unlock()
|
||||
if telBuilt {
|
||||
return telProvider
|
||||
}
|
||||
|
||||
endpoint := cloudflare.Getenv(otelEndpointVar)
|
||||
if endpoint == "" {
|
||||
// Endpoint TBD/unset: leave telemetry disabled. Do not cache, so a later
|
||||
// deploy that sets the var (new isolate) picks it up.
|
||||
return nil
|
||||
}
|
||||
|
||||
headers := map[string]string{}
|
||||
if raw, err := storage.ReadSecret(otelHeadersSecret); err != nil {
|
||||
log.Printf("telemetry: OTLP headers secret %q unavailable (%v); exporting without auth headers", otelHeadersSecret, err)
|
||||
} else {
|
||||
headers = telemetry.ParseHeaders(string(raw))
|
||||
}
|
||||
|
||||
service := cloudflare.Getenv(otelServiceVar)
|
||||
if service == "" {
|
||||
service = otelServiceDefault
|
||||
}
|
||||
|
||||
exp := telemetry.NewOTLPHTTPExporter(telemetry.OTLPConfig{
|
||||
Endpoint: endpoint,
|
||||
Headers: headers,
|
||||
Resource: []telemetry.KeyValue{telemetry.String("service.name", service)},
|
||||
Scope: telemetry.Scope{Name: otelScopeName},
|
||||
})
|
||||
telProvider = telemetry.New(exp, telemetry.WithErrorHandler(func(err error) {
|
||||
// Best-effort: a telemetry export failure must never break the request.
|
||||
log.Printf("telemetry: export failed: %v", err)
|
||||
}))
|
||||
telBuilt = true
|
||||
return telProvider
|
||||
}
|
||||
|
||||
// withWorkerTelemetry wraps the fetch handler so every request carries the
|
||||
// telemetry provider in its context (context propagation). When telemetry is
|
||||
// disabled the provider is nil and the wrapped handler behaves identically.
|
||||
func withWorkerTelemetry(next http.Handler) http.Handler {
|
||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
ctx := telemetry.NewContext(r.Context(), workerTelemetry())
|
||||
next.ServeHTTP(w, r.WithContext(ctx))
|
||||
})
|
||||
}
|
||||
|
||||
// scheduledContext returns a context carrying the telemetry provider for the
|
||||
// weekly publish run, so internal/schedule and internal/publish emit their run
|
||||
// and per-report spans/logs.
|
||||
func scheduledContext(ctx context.Context) context.Context {
|
||||
return telemetry.NewContext(ctx, workerTelemetry())
|
||||
}
|
||||
+24
-1
@@ -66,6 +66,19 @@
|
||||
"binding": "GITHUB_TOKEN",
|
||||
"store_id": "<store-id>",
|
||||
"secret_name": "github-token"
|
||||
},
|
||||
// OpenTelemetry OTLP auth headers (#17). The OTLP exporter (internal/telemetry)
|
||||
// sends these headers on every /v1/traces and /v1/logs POST — typically an auth
|
||||
// token, e.g. "Authorization=Bearer <token>" or "x-api-key=<key>"; multiple are
|
||||
// comma-separated ("k1=v1,k2=v2"). Kept as a Secrets Store secret so the token
|
||||
// is never committed. OPTIONAL: if the OTLP backend needs no auth (or auth is in
|
||||
// the endpoint URL) this can be omitted — the exporter then sends no auth header.
|
||||
// The Worker reads it lazily per run via worker/telemetry_wasm.go
|
||||
// otelHeadersSecret. Replace "<store-id>" at deploy time.
|
||||
{
|
||||
"binding": "OTEL_EXPORTER_OTLP_HEADERS",
|
||||
"store_id": "<store-id>",
|
||||
"secret_name": "otel-exporter-otlp-headers"
|
||||
}
|
||||
],
|
||||
|
||||
@@ -73,8 +86,18 @@
|
||||
// job (#14) files issues on; it is configurable here without a code change and
|
||||
// defaults to "JMR-dev/LibreMail" in worker/scheduled_wasm.go if unset. Read via
|
||||
// cloudflare.Getenv(GITHUB_REPO).
|
||||
//
|
||||
// OpenTelemetry (#17): OTEL_EXPORTER_OTLP_ENDPOINT is the base OTLP/HTTP URL the
|
||||
// exporter POSTs traces + logs to ("/v1/traces" and "/v1/logs" are appended).
|
||||
// The backend/endpoint is deliberately TBD — left EMPTY here, which disables
|
||||
// telemetry (the Worker behaves exactly as before) until a collector/backend is
|
||||
// chosen and this is set (a non-secret URL; the auth token lives in the
|
||||
// OTEL_EXPORTER_OTLP_HEADERS secret above). OTEL_SERVICE_NAME overrides the
|
||||
// reported service.name (defaults to "libremail-bug-report-ingest").
|
||||
"vars": {
|
||||
"GITHUB_REPO": "JMR-dev/LibreMail"
|
||||
"GITHUB_REPO": "JMR-dev/LibreMail",
|
||||
"OTEL_EXPORTER_OTLP_ENDPOINT": "",
|
||||
"OTEL_SERVICE_NAME": "libremail-bug-report-ingest"
|
||||
},
|
||||
|
||||
// Cron Triggers for the weekly publish job (#13): "Friday 17:00 America/Chicago
|
||||
|
||||
Reference in New Issue
Block a user