Merge pull request #45 from JMR-dev/ticket-17-otel-observability

#17 Observability: OpenTelemetry logging + tracing + alerting
This commit was merged in pull request #45.
This commit is contained in:
Jason Ross
2026-07-02 17:16:49 -05:00
committed by GitHub
15 changed files with 2382 additions and 8 deletions
+134
View File
@@ -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{}
+182
View File
@@ -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)
}
}
+79 -4
View File
@@ -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
+233
View File
@@ -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))
}
}
+32 -2
View File
@@ -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
}
+141
View File
@@ -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) }
+99
View File
@@ -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)
+356
View File
@@ -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)
+229
View File
@@ -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
}
+502
View File
@@ -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
}
+239
View File
@@ -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
View File
@@ -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{})))
}
+6
View File
@@ -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
+120
View File
@@ -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
View File
@@ -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