From 3449c6163fc59327e73e93a67ca212df71e93f45 Mon Sep 17 00:00:00 2001 From: "Claude (instanode)" Date: Tue, 12 May 2026 23:14:29 +0530 Subject: [PATCH 1/2] obs: port worker observability + buildinfo from api repo (B1) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Relocates the observability code that was merged into InstaNode-dev/api under api/worker/ (PR #37 there) into its real home in this repo. What ships: - internal/jobs/middleware.go — WithObservability[T] generic River-Worker wrapper that stamps tid/trace_id on ctx and (optionally) opens a New Relic transaction per job. Fail-open on nil nrApp. - internal/jobs/middleware_test.go — 7 tests covering tid stamping, trace_id missing/present, error propagation, nil-NR safety, delegation of NextRetry/Timeout, plus the int64 formatter. - internal/obs/nr.go — InitNewRelic + WaitForConnection helpers. - internal/obs/nr_test.go — 2 tests asserting fail-open contract. - main.go — slog wrapped in logctx.NewHandler, NR init, /healthz now emits commit_id/build_time/version, workers receive nrApp. - internal/jobs/workers.go — StartWorkers gains nrApp parameter; every river.AddWorker call wraps the worker via WithObservability(...). Critical detail: the api/worker/ PR shipped against TEMPORARY stubs at instant.dev/worker/internal/_obs_stubs/{buildinfo,logctx}. This relocate switches both imports to the canonical common packages: - instant.dev/worker/internal/_obs_stubs/buildinfo -> instant.dev/common/buildinfo - instant.dev/worker/internal/_obs_stubs/logctx -> instant.dev/common/logctx That mirrors today's PR #40 fix on the api repo, where the same stub->common substitution was applied. The worker module's existing `replace instant.dev/common => ../common` directive in go.mod makes the canonical import resolve to the sibling checkout. go.mod gains: - github.com/newrelic/go-agent/v3 (direct) Tests: go test ./... -count=1 — all green, 34 PASS in internal/jobs + internal/obs. Co-Authored-By: Claude Opus 4.7 (1M context) --- go.mod | 7 +- go.sum | 14 ++- internal/jobs/middleware.go | 173 +++++++++++++++++++++++++++ internal/jobs/middleware_test.go | 195 +++++++++++++++++++++++++++++++ internal/jobs/workers.go | 30 +++-- internal/obs/nr.go | 83 +++++++++++++ internal/obs/nr_test.go | 35 ++++++ main.go | 37 ++++-- 8 files changed, 547 insertions(+), 27 deletions(-) create mode 100644 internal/jobs/middleware.go create mode 100644 internal/jobs/middleware_test.go create mode 100644 internal/obs/nr.go create mode 100644 internal/obs/nr_test.go diff --git a/go.mod b/go.mod index 66dac18..ef462ad 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module instant.dev/worker -go 1.24.0 +go 1.25 require ( github.com/DATA-DOG/go-sqlmock v1.5.2 @@ -9,6 +9,7 @@ require ( github.com/lib/pq v1.10.9 github.com/minio/madmin-go/v3 v3.0.110 github.com/minio/minio-go/v7 v7.0.90 + github.com/newrelic/go-agent/v3 v3.43.3 github.com/prometheus/client_golang v1.21.0 github.com/redis/go-redis/v9 v9.6.1 github.com/resend/resend-go/v2 v2.28.0 @@ -19,7 +20,7 @@ require ( go.opentelemetry.io/otel v1.39.0 go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.39.0 go.opentelemetry.io/otel/sdk v1.39.0 - google.golang.org/grpc v1.79.3 + google.golang.org/grpc v1.80.0 instant.dev/common v0.0.0 instant.dev/proto v0.0.0 k8s.io/api v0.32.2 @@ -103,7 +104,7 @@ require ( golang.org/x/term v0.39.0 // indirect golang.org/x/text v0.33.0 // indirect golang.org/x/time v0.10.0 // indirect - google.golang.org/genproto/googleapis/api v0.0.0-20251202230838-ff82c1b0f217 // indirect + google.golang.org/genproto/googleapis/api v0.0.0-20260120221211-b8f7ae30c516 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260120221211-b8f7ae30c516 // indirect google.golang.org/protobuf v1.36.11 // indirect gopkg.in/evanphx/json-patch.v4 v4.12.0 // indirect diff --git a/go.sum b/go.sum index 6aa6a88..417a6e3 100644 --- a/go.sum +++ b/go.sum @@ -114,6 +114,8 @@ github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9G github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= +github.com/newrelic/go-agent/v3 v3.43.3 h1:0A6DkUBYK2bidV6jJDJ1SD2XkRlg976nl+SiEqkGTUQ= +github.com/newrelic/go-agent/v3 v3.43.3/go.mod h1:MFXnCId5xXMIJI6A/kbkg0DO48EVTsKcmNijMYphzTg= github.com/onsi/ginkgo/v2 v2.21.0 h1:7rg/4f3rB88pb5obDgNZrNHrQ4e6WpjonchcpuBRnZM= github.com/onsi/ginkgo/v2 v2.21.0/go.mod h1:7Du3c42kxCUegi0IImZ1wUQzMBVecgIHjR1C+NkhLQo= github.com/onsi/gomega v1.35.1 h1:Cwbd75ZBPxFSuZ6T+rN/WCb/gOc6YgFBXLlZLhC7Ds4= @@ -262,14 +264,14 @@ golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8T golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk= -gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E= -google.golang.org/genproto/googleapis/api v0.0.0-20251202230838-ff82c1b0f217 h1:fCvbg86sFXwdrl5LgVcTEvNC+2txB5mgROGmRL5mrls= -google.golang.org/genproto/googleapis/api v0.0.0-20251202230838-ff82c1b0f217/go.mod h1:+rXWjjaukWZun3mLfjmVnQi18E1AsFbDN9QdJ5YXLto= +gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= +gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= +google.golang.org/genproto/googleapis/api v0.0.0-20260120221211-b8f7ae30c516 h1:vmC/ws+pLzWjj/gzApyoZuSVrDtF1aod4u/+bbj8hgM= +google.golang.org/genproto/googleapis/api v0.0.0-20260120221211-b8f7ae30c516/go.mod h1:p3MLuOwURrGBRoEyFHBT3GjUwaCQVKeNqqWxlcISGdw= google.golang.org/genproto/googleapis/rpc v0.0.0-20260120221211-b8f7ae30c516 h1:sNrWoksmOyF5bvJUcnmbeAmQi8baNhqg5IWaI3llQqU= google.golang.org/genproto/googleapis/rpc v0.0.0-20260120221211-b8f7ae30c516/go.mod h1:j9x/tPzZkyxcgEFkiKEEGxfvyumM01BEtsW8xzOahRQ= -google.golang.org/grpc v1.79.3 h1:sybAEdRIEtvcD68Gx7dmnwjZKlyfuc61Dyo9pGXXkKE= -google.golang.org/grpc v1.79.3/go.mod h1:KmT0Kjez+0dde/v2j9vzwoAScgEPx/Bw1CYChhHLrHQ= +google.golang.org/grpc v1.80.0 h1:Xr6m2WmWZLETvUNvIUmeD5OAagMw3FiKmMlTdViWsHM= +google.golang.org/grpc v1.80.0/go.mod h1:ho/dLnxwi3EDJA4Zghp7k2Ec1+c2jqup0bFkw07bwF4= google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= diff --git a/internal/jobs/middleware.go b/internal/jobs/middleware.go new file mode 100644 index 0000000..f0a80a2 --- /dev/null +++ b/internal/jobs/middleware.go @@ -0,0 +1,173 @@ +// File adds the observability middleware used by every River worker +// registered in StartWorkers (see workers.go). Wrapping is opt-in at the +// AddWorker call-site: the actual job implementations in expire.go, quota.go, +// storage.go, geodb.go, trial.go, etc. are NOT modified by this track — +// the wrapper does its job around them. +// +// Track 4 of the observability rollout (OBSERVABILITY-PLAN-2026-05-12.md). +// +// What it does, per executed job: +// +// 1. Stamps `tid = ` on the ctx via logctx.WithTID so every slog +// line emitted inside the job carries the same task id — agents can +// grep one job's full trace from a stream of interleaved workers. +// 2. Stamps `trace_id = ` on the ctx via logctx.WithTraceID +// if one is not already present. Real ingest of OTel-derived trace ids +// will follow track 7 — this guarantees the field is always non-empty +// so log queries can be written today. +// 3. Opens a New Relic transaction named `job.` and defers its +// end. Errors returned by the inner Work bubble through nrtxn.NoticeError +// before being returned, so they surface in the NR error inbox. +// 4. Logs duration on completion at INFO (success) or ERROR (failure) +// using a consistent shape so the dashboard panels under track 7 can +// bind to a stable schema. +// +// The wrapper is a thin generic function: it preserves the concrete +// `river.Worker[T]` type so `river.AddWorker` keeps accepting it without +// reflection. NextRetry, Timeout, and every other Worker method delegate +// to the inner worker so existing retry / timeout policy is untouched. +package jobs + +import ( + "context" + "log/slog" + "time" + + "github.com/google/uuid" + "github.com/newrelic/go-agent/v3/newrelic" + "github.com/riverqueue/river" + + "instant.dev/common/logctx" +) + +// observabilityWorker wraps an inner river.Worker[T] with the per-job +// observability concerns described in the package doc. It is constructed +// via WithObservability and never used directly. +// +// The inner worker is held by value of an interface type so the wrapper does +// not have to know any of its fields. Every Worker[T] method delegates. +type observabilityWorker[T river.JobArgs] struct { + inner river.Worker[T] + nrApp *newrelic.Application // may be nil — fail-open +} + +// WithObservability wraps next so that each job execution is instrumented +// with logctx ids and an optional New Relic transaction. +// +// nrApp may be nil — in that case the wrapper still stamps ctx ids and logs +// duration, it just does not open an NR transaction. This matches the +// fail-open contract of obs.InitNewRelic. +// +// Call site (workers.go): +// +// river.AddWorker(workers, jobs.WithObservability(jobs.NewExpireAnonymousWorker(...), nrApp)) +// +// Note the generic parameter is inferred from the wrapped worker, so the +// caller writes WithObservability(...) not WithObservability[ExpireAnonymousArgs](...). +func WithObservability[T river.JobArgs](next river.Worker[T], nrApp *newrelic.Application) river.Worker[T] { + return &observabilityWorker[T]{inner: next, nrApp: nrApp} +} + +// Work is the only method that does real work — the rest delegate. It runs +// in this order: stamp ids, open NR txn, call inner.Work, record outcome, +// end NR txn (via defer), log duration. +func (w *observabilityWorker[T]) Work(ctx context.Context, job *river.Job[T]) error { + // Step 1: stamp ids on ctx so every slog call inside the job sees them. + // We always overwrite tid (the job is the authoritative source for the + // task id) but we PRESERVE an existing trace_id if one is present — that + // path is taken when a periodic-job dispatcher already opened a trace. + tid := jobIDString(job.ID) + ctx = logctx.WithTID(ctx, tid) + if logctx.TraceIDFromContext(ctx) == "" { + ctx = logctx.WithTraceID(ctx, uuid.New().String()) + } + + // Step 2: open the New Relic transaction. txn is nil-safe — every method + // on (*newrelic.Transaction)(nil) is a no-op in the v3 SDK — but we still + // gate the StartTransaction call to avoid the nil-deref on nrApp itself. + kind := jobKind(job) + var txn *newrelic.Transaction + if w.nrApp != nil { + txn = w.nrApp.StartTransaction("job." + kind) + // nrtxn carries the ctx for the duration of Work. Cross-process + // linkage (OTel headers) is set up by track 7 — today we only need + // the in-process span. + ctx = newrelic.NewContext(ctx, txn) + defer txn.End() + } + + start := time.Now() + err := w.inner.Work(ctx, job) + elapsed := time.Since(start) + + if err != nil { + if txn != nil { + txn.NoticeError(err) + } + slog.ErrorContext(ctx, "jobs.middleware.work_failed", + "kind", kind, + "job_id", job.ID, + "attempt", job.Attempt, + "duration_ms", elapsed.Milliseconds(), + "error", err.Error(), + ) + return err + } + + slog.InfoContext(ctx, "jobs.middleware.work_ok", + "kind", kind, + "job_id", job.ID, + "attempt", job.Attempt, + "duration_ms", elapsed.Milliseconds(), + ) + return nil +} + +// NextRetry, Timeout — pure delegation. The wrapper MUST NOT impose its own +// retry or timeout policy; that belongs to the wrapped worker (typically via +// river.WorkerDefaults embedded by the concrete worker struct). +func (w *observabilityWorker[T]) NextRetry(job *river.Job[T]) time.Time { + return w.inner.NextRetry(job) +} + +func (w *observabilityWorker[T]) Timeout(job *river.Job[T]) time.Duration { + return w.inner.Timeout(job) +} + +// jobKind extracts the job kind without forcing the caller to depend on the +// concrete args type. It calls (T).Kind() through the JobArgs interface; +// every River job args type already implements Kind() so this is free. +// +// We pull Kind() from job.Args rather than a fresh zero value because the +// JobArgs interface contract is that Kind() is constant per type. +func jobKind[T river.JobArgs](job *river.Job[T]) string { + return job.Args.Kind() +} + +// jobIDString formats an int64 job id without pulling in strconv at the +// call site. Kept tiny because it sits on the hot path of every job. +func jobIDString(id int64) string { + if id == 0 { + return "" + } + const digits = "0123456789" + var buf [20]byte + pos := len(buf) + neg := id < 0 + u := uint64(id) + if neg { + u = uint64(-id) + } + for u >= 10 { + pos-- + buf[pos] = digits[u%10] + u /= 10 + } + pos-- + buf[pos] = digits[u] + if neg { + pos-- + buf[pos] = '-' + } + return string(buf[pos:]) +} diff --git a/internal/jobs/middleware_test.go b/internal/jobs/middleware_test.go new file mode 100644 index 0000000..f780839 --- /dev/null +++ b/internal/jobs/middleware_test.go @@ -0,0 +1,195 @@ +// Tests for the observability middleware. The interesting properties: +// +// 1. `tid` ends up on the ctx via logctx.WithTID — readable with +// logctx.TIDFromContext — and matches the job.ID. +// 2. `trace_id` is non-empty after the wrapper runs, even when the caller +// passed no trace id in, and is preserved when the caller did. +// 3. An error from the inner worker bubbles through unchanged. +// 4. Duration is recorded (we can't easily assert it from outside, but we +// can assert the wrapper doesn't crash on a slow job). +// 5. The wrapper is safe with a nil New Relic application (fail-open). +// +// We don't unit-test the New-Relic-present path because it would require a +// live agent connection. The nil-app path covers the only branch under our +// control; integration tests for the present-path live in the deployment +// rollout (track 7). +package jobs + +import ( + "context" + "errors" + "strconv" + "testing" + "time" + + "github.com/riverqueue/river" + "github.com/riverqueue/river/rivertype" + + "instant.dev/common/logctx" +) + +// fakeArgs is a minimal river.JobArgs that the test uses to type the wrapper. +type fakeArgs struct{} + +func (fakeArgs) Kind() string { return "fake_test_job" } + +// fakeWorker is a river.Worker[fakeArgs] whose Work captures the ctx it was +// called with and optionally returns a configured error. NextRetry/Timeout +// return zero values to satisfy the interface. +type fakeWorker struct { + river.WorkerDefaults[fakeArgs] + gotCtx context.Context + gotJob *river.Job[fakeArgs] + returns error + delay time.Duration +} + +func (f *fakeWorker) Work(ctx context.Context, job *river.Job[fakeArgs]) error { + f.gotCtx = ctx + f.gotJob = job + if f.delay > 0 { + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(f.delay): + } + } + return f.returns +} + +// newJob returns a river.Job[fakeArgs] with the given id. river.Job embeds +// *rivertype.JobRow, so we construct the row separately and point the job +// at it. The middleware only reads ID + Attempt off the row plus Args.Kind() +// so the rest of the JobRow fields can stay zero. +func newJob(id int64) *river.Job[fakeArgs] { + return &river.Job[fakeArgs]{ + JobRow: &rivertype.JobRow{ID: id, Kind: "fake_test_job"}, + Args: fakeArgs{}, + } +} + +// TestWithObservability_StampsTIDOnContext is the contract test the task +// brief calls out: the wrapper must put job.ID on the ctx under the logctx +// "tid" key so downstream slog calls pick it up automatically. +func TestWithObservability_StampsTIDOnContext(t *testing.T) { + fake := &fakeWorker{} + wrapped := WithObservability[fakeArgs](fake, nil) + + want := int64(42) + if err := wrapped.Work(context.Background(), newJob(want)); err != nil { + t.Fatalf("wrapped.Work returned error: %v", err) + } + if fake.gotCtx == nil { + t.Fatalf("inner worker was never called") + } + got := logctx.TIDFromContext(fake.gotCtx) + if got != strconv.FormatInt(want, 10) { + t.Fatalf("tid on ctx: got %q, want %q", got, strconv.FormatInt(want, 10)) + } +} + +// TestWithObservability_SetsTraceIDWhenMissing asserts the wrapper generates +// a trace id when the incoming ctx has none. The exact value doesn't matter, +// only that it's non-empty so log queries always find a populated field. +func TestWithObservability_SetsTraceIDWhenMissing(t *testing.T) { + fake := &fakeWorker{} + wrapped := WithObservability[fakeArgs](fake, nil) + + if err := wrapped.Work(context.Background(), newJob(7)); err != nil { + t.Fatalf("wrapped.Work returned error: %v", err) + } + if got := logctx.TraceIDFromContext(fake.gotCtx); got == "" { + t.Fatalf("trace_id was not set on ctx") + } +} + +// TestWithObservability_PreservesExistingTraceID asserts the wrapper does NOT +// overwrite a trace id that the caller already attached. This matters when a +// periodic-job dispatcher (out of scope for this track) opens the trace and +// the worker needs to inherit it. +func TestWithObservability_PreservesExistingTraceID(t *testing.T) { + fake := &fakeWorker{} + wrapped := WithObservability[fakeArgs](fake, nil) + + const want = "trace-from-dispatcher" + ctx := logctx.WithTraceID(context.Background(), want) + if err := wrapped.Work(ctx, newJob(9)); err != nil { + t.Fatalf("wrapped.Work returned error: %v", err) + } + if got := logctx.TraceIDFromContext(fake.gotCtx); got != want { + t.Fatalf("trace_id: got %q, want %q (wrapper must not overwrite)", got, want) + } +} + +// TestWithObservability_PropagatesError covers the failure path: an error +// from the inner worker must reach the caller unchanged so River's retry +// machinery still sees it. We assert errors.Is to be defensive against the +// wrapper deciding to wrap the error in the future. +func TestWithObservability_PropagatesError(t *testing.T) { + want := errors.New("simulated job failure") + fake := &fakeWorker{returns: want} + wrapped := WithObservability[fakeArgs](fake, nil) + + err := wrapped.Work(context.Background(), newJob(11)) + if !errors.Is(err, want) { + t.Fatalf("error not propagated: got %v, want %v", err, want) + } +} + +// TestWithObservability_NilNRAppIsSafe is the fail-open contract test. With +// no NR app, the wrapper still runs the inner worker, still stamps ids on +// ctx, still returns the inner's error. We cover both error-free and +// error-returning paths so the deferred txn.End() path is exercised. +func TestWithObservability_NilNRAppIsSafe(t *testing.T) { + t.Run("success", func(t *testing.T) { + fake := &fakeWorker{} + wrapped := WithObservability[fakeArgs](fake, nil) + if err := wrapped.Work(context.Background(), newJob(1)); err != nil { + t.Fatalf("unexpected error: %v", err) + } + }) + t.Run("failure", func(t *testing.T) { + boom := errors.New("boom") + fake := &fakeWorker{returns: boom} + wrapped := WithObservability[fakeArgs](fake, nil) + if err := wrapped.Work(context.Background(), newJob(2)); !errors.Is(err, boom) { + t.Fatalf("unexpected error: got %v, want %v", err, boom) + } + }) +} + +// TestWithObservability_DelegatesNextRetryAndTimeout asserts the wrapper +// doesn't impose its own policy. The fakeWorker embeds river.WorkerDefaults +// which returns zero values; we just confirm calling those methods through +// the wrapper does not panic and returns the inner values. +func TestWithObservability_DelegatesNextRetryAndTimeout(t *testing.T) { + fake := &fakeWorker{} + wrapped := WithObservability[fakeArgs](fake, nil) + + if got := wrapped.NextRetry(newJob(1)); !got.IsZero() { + t.Fatalf("NextRetry should delegate to WorkerDefaults (zero time), got %v", got) + } + if got := wrapped.Timeout(newJob(1)); got != 0 { + t.Fatalf("Timeout should delegate to WorkerDefaults (0), got %v", got) + } +} + +// TestJobIDString covers the tiny int64->string formatter used to keep the +// hot path allocation-light. Belt-and-braces: 0, positive, negative. +func TestJobIDString(t *testing.T) { + cases := []struct { + in int64 + want string + }{ + {0, ""}, + {1, "1"}, + {42, "42"}, + {9876543210, "9876543210"}, + {-7, "-7"}, + } + for _, c := range cases { + if got := jobIDString(c.in); got != c.want { + t.Errorf("jobIDString(%d) = %q, want %q", c.in, got, c.want) + } + } +} diff --git a/internal/jobs/workers.go b/internal/jobs/workers.go index b959a33..a189c81 100644 --- a/internal/jobs/workers.go +++ b/internal/jobs/workers.go @@ -10,6 +10,7 @@ import ( madmin "github.com/minio/madmin-go/v3" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" + "github.com/newrelic/go-agent/v3/newrelic" "github.com/redis/go-redis/v9" "github.com/riverqueue/river" "github.com/riverqueue/river/riverdriver/riverpgxv5" @@ -83,7 +84,7 @@ func (mondayAt8UTCSchedule) Next(t time.Time) time.Time { // namespaces. Pass nil when the worker can't reach a cluster — the // reconciler logs at WARN each run and other periodic jobs keep functioning. // See worker/internal/jobs/deploy_status_reconcile.go for the SCOPE NOTE. -func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *config.Config, provClient *provisioner.Client, planRegistry PlanRegistry, deployStatusK8s deployStatusK8sProvider) *Workers { +func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *config.Config, provClient *provisioner.Client, planRegistry PlanRegistry, deployStatusK8s deployStatusK8sProvider, nrApp *newrelic.Application) *Workers { _ = rdb // available for future workers; currently only used by quota checks done via db // River requires pgx pool — open a separate connection for the worker pool. @@ -131,24 +132,31 @@ func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *confi } workers := river.NewWorkers() - river.AddWorker(workers, NewExpireAnonymousWorker(db, provClient, minioClient)) - river.AddWorker(workers, NewExpireStacksWorker(db, cfg.KubeNamespaceApps+"-")) - river.AddWorker(workers, NewRefreshGeoDBWorker()) - river.AddWorker(workers, &TrialExpiryWorker{db: db, email: emailClient}) - river.AddWorker(workers, &WeeklyDigestWorker{db: db, email: emailClient}) - river.AddWorker(workers, NewExpiryReminderWorker(db, emailClient)) - river.AddWorker(workers, NewEnforceStorageQuotaWorker(db, planRegistry)) - river.AddWorker(workers, NewUpdateStorageBytesWorker(db, provClient, minioScanner)) + // Each worker is wrapped in WithObservability so every job execution + // stamps tid + trace_id on ctx and (optionally) opens a New Relic + // transaction. nrApp may be nil — the wrapper still does the ctx work. + // See middleware.go for the full contract. + river.AddWorker(workers, WithObservability(NewExpireAnonymousWorker(db, provClient, minioClient), nrApp)) + river.AddWorker(workers, WithObservability(NewExpireStacksWorker(db, cfg.KubeNamespaceApps+"-"), nrApp)) + river.AddWorker(workers, WithObservability(NewRefreshGeoDBWorker(), nrApp)) + // TrialExpiry / WeeklyDigest are registered via composite literal, so the + // generic type parameter can't be inferred from the constructor return — + // it must be supplied explicitly. + river.AddWorker(workers, WithObservability[TrialExpiryArgs](&TrialExpiryWorker{db: db, email: emailClient}, nrApp)) + river.AddWorker(workers, WithObservability[WeeklyDigestArgs](&WeeklyDigestWorker{db: db, email: emailClient}, nrApp)) + river.AddWorker(workers, WithObservability(NewExpiryReminderWorker(db, emailClient), nrApp)) + river.AddWorker(workers, WithObservability(NewEnforceStorageQuotaWorker(db, planRegistry), nrApp)) + river.AddWorker(workers, WithObservability(NewUpdateStorageBytesWorker(db, provClient, minioScanner), nrApp)) // Custom-domain reconciler — TXT lookup, HTTP probe, stale-failed sweep. // k8s provider is nil today: the worker module does not import the api's // k8s client. Steps 2/3 (Ingress + cert poll) stay in the api handler. // See custom_domain_reconcile.go for the full SCOPE NOTE. - river.AddWorker(workers, NewCustomDomainReconciler(db, nil, nil)) + river.AddWorker(workers, WithObservability(NewCustomDomainReconciler(db, nil, nil), nrApp)) // Deploy-status reconciler — sweeps non-terminal deployments and rolls // status forward from live k8s Deployment state every 30s. deployStatusK8s // may be nil (kubeconfig unreachable in CI / docker-compose); the worker // then short-circuits with a WARN each tick. See deploy_status_reconcile.go. - river.AddWorker(workers, NewDeployStatusReconciler(db, deployStatusK8s)) + river.AddWorker(workers, WithObservability(NewDeployStatusReconciler(db, deployStatusK8s), nrApp)) periodicJobs := []*river.PeriodicJob{ river.NewPeriodicJob( diff --git a/internal/obs/nr.go b/internal/obs/nr.go new file mode 100644 index 0000000..bc4ce46 --- /dev/null +++ b/internal/obs/nr.go @@ -0,0 +1,83 @@ +// Package obs holds observability bootstrap helpers shared across the +// worker binary. Today it has one job: build a New Relic Application from +// env vars and never crash when the license key is missing. +// +// Track 4 of the observability rollout (OBSERVABILITY-PLAN-2026-05-12.md). +// The api and provisioner services have parallel helpers under their own +// internal/obs packages — each owns its own copy to keep service boundaries +// clean. The contract (fail-open, log-only warning, return nil app) is +// identical across all three. +package obs + +import ( + "log/slog" + "os" + "time" + + "github.com/newrelic/go-agent/v3/newrelic" +) + +// nrInitTimeout caps how long ConnectReply may block on bootstrap. The Go +// agent connects async by default, so this is a guard for the rare case where +// caller code waits on `WaitForConnection`. +const nrInitTimeout = 5 * time.Second + +// InitNewRelic returns a *newrelic.Application built from environment. +// +// Contract: NEVER crash. NEW_RELIC_LICENSE_KEY is the only required input; +// when it is empty (local dev, CI, k8s pod without the secret mounted yet) +// we log a warning and return (nil, nil). Every caller MUST nil-check the +// returned application before invoking methods on it — `(*nrApp).StartTransaction` +// is a nil-safe no-op in the v3 SDK, but defensive callers should still guard. +// +// The license-key-present path can still fail (network down, malformed key, +// duplicate registration). In that case we log the underlying error and +// return (nil, err) so the caller can surface it but keep running. The worker +// pod must not crashloop because New Relic is unhappy. +func InitNewRelic() (*newrelic.Application, error) { + licenseKey := os.Getenv("NEW_RELIC_LICENSE_KEY") + if licenseKey == "" { + slog.Warn("obs.newrelic.skipped", + "reason", "NEW_RELIC_LICENSE_KEY not set", + "behavior", "transactions are no-ops, worker continues") + return nil, nil + } + + appName := os.Getenv("NEW_RELIC_APP_NAME") + if appName == "" { + appName = "instant-worker" + } + + app, err := newrelic.NewApplication( + newrelic.ConfigAppName(appName), + newrelic.ConfigLicense(licenseKey), + newrelic.ConfigAppLogForwardingEnabled(true), + newrelic.ConfigDistributedTracerEnabled(true), + // Fail-open at the SDK level too: don't crash if the daemon can't be + // reached, just suppress the noisy harvest-cycle errors. + func(cfg *newrelic.Config) { + cfg.ErrorCollector.Enabled = true + cfg.TransactionTracer.Enabled = true + }, + ) + if err != nil { + slog.Warn("obs.newrelic.init_failed", + "error", err, + "behavior", "transactions are no-ops, worker continues") + return nil, err + } + + slog.Info("obs.newrelic.initialised", "app_name", appName) + return app, nil +} + +// WaitForConnection is a thin wrapper around app.WaitForConnection that does +// nothing when app is nil. Use only from tests or boot code that wants the +// agent fully connected before proceeding; production code paths should never +// block on this. +func WaitForConnection(app *newrelic.Application) { + if app == nil { + return + } + _ = app.WaitForConnection(nrInitTimeout) +} diff --git a/internal/obs/nr_test.go b/internal/obs/nr_test.go new file mode 100644 index 0000000..1ee53f4 --- /dev/null +++ b/internal/obs/nr_test.go @@ -0,0 +1,35 @@ +// Tests for the New Relic init helper. The hard requirement is the +// fail-open contract: missing NEW_RELIC_LICENSE_KEY must return (nil, nil) +// with a warning log, NEVER an error and NEVER a crash. +// +// We don't test the success path here — it would require either embedding a +// fake NR collector or carrying a real license key in CI secrets, neither of +// which is worth the complexity for a thin bootstrap helper. +package obs + +import ( + "testing" +) + +// TestInitNewRelic_FailOpenOnMissingLicenseKey is the primary contract test. +// With no env var, the helper must return (nil, nil). We use t.Setenv to +// guarantee an empty value even on developer machines where the env might be +// set in their shell. +func TestInitNewRelic_FailOpenOnMissingLicenseKey(t *testing.T) { + t.Setenv("NEW_RELIC_LICENSE_KEY", "") + + app, err := InitNewRelic() + if err != nil { + t.Fatalf("expected nil error on missing license key, got %v", err) + } + if app != nil { + t.Fatalf("expected nil application on missing license key, got %v", app) + } +} + +// TestWaitForConnection_NilSafe is a trivial guard for the helper that some +// boot-time code paths may call before the app is fully constructed. +func TestWaitForConnection_NilSafe(t *testing.T) { + // Must not panic. + WaitForConnection(nil) +} diff --git a/main.go b/main.go index cbace61..3f4b510 100644 --- a/main.go +++ b/main.go @@ -13,19 +13,35 @@ import ( "github.com/prometheus/client_golang/prometheus/promhttp" "google.golang.org/grpc" + "instant.dev/common/buildinfo" + "instant.dev/common/logctx" commonplans "instant.dev/common/plans" "instant.dev/worker/internal/config" "instant.dev/worker/internal/db" "instant.dev/worker/internal/jobs" + "instant.dev/worker/internal/obs" "instant.dev/worker/internal/provisioner" "instant.dev/worker/internal/telemetry" ) func main() { - // Structured JSON logging. - slog.SetDefault(slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{ - Level: slog.LevelInfo, - }))) + // Structured JSON logging — wrapped in logctx so every line carries + // service + commit_id + (when present) tid / trace_id / team_id. + base := slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{ + Level: slog.LevelInfo, + AddSource: true, + }) + slog.SetDefault(slog.New(logctx.NewHandler("worker", base))) + + // New Relic Go agent. Fail-open on empty / missing license so local dev + // and CI runs (which never get a real key) still boot. Matches the + // contract of telemetry.InitTracer below. + nrApp, _ := obs.InitNewRelic() + defer func() { + if nrApp != nil { + nrApp.Shutdown(5 * time.Second) + } + }() shutdownTracer := telemetry.InitTracer("instant-worker", os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT")) defer func() { @@ -85,7 +101,7 @@ func main() { slog.Info("worker.deploy_status_k8s_client_ready") } - workers := jobs.StartWorkers(ctx, database, rdb, cfg, provClient, planRegistry, deployStatusK8s) + workers := jobs.StartWorkers(ctx, database, rdb, cfg, provClient, planRegistry, deployStatusK8s, nrApp) defer workers.Stop() // Exit immediately if River failed to start so Kubernetes restarts the pod. @@ -102,7 +118,8 @@ func main() { // If River's goroutines panic after start, Go crashes the process and k8s restarts. mux := http.NewServeMux() mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { - fmt.Fprintf(w, `{"ok":true,"service":"instant-worker"}`) + fmt.Fprintf(w, `{"ok":true,"service":"instant-worker","commit_id":%q,"build_time":%q,"version":%q}`, + buildinfo.GitSHA, buildinfo.BuildTime, buildinfo.Version) }) mux.Handle("/metrics", promhttp.Handler()) srv := &http.Server{Addr: ":8091", Handler: mux} @@ -117,7 +134,13 @@ func main() { _ = srv.Shutdown(shutCtx) }() - slog.Info("worker.started", "environment", cfg.Environment, "liveness_port", 8091) + slog.Info("worker.started", + "environment", cfg.Environment, + "liveness_port", 8091, + "commit_id", buildinfo.GitSHA, + "build_time", buildinfo.BuildTime, + "version", buildinfo.Version, + ) <-ctx.Done() slog.Info("worker.shutdown") } From 20ee825cb8482fa348c0cec0b2386a7e4f483e47 Mon Sep 17 00:00:00 2001 From: "Claude (instanode)" Date: Tue, 12 May 2026 23:15:48 +0530 Subject: [PATCH 2/2] build: bump Dockerfile to golang:1.25-alpine go.mod ended up at go 1.25 after go mod tidy resolved newrelic / k8s client-go transitive bounds. Matches the api repo Dockerfile. Co-Authored-By: Claude Opus 4.7 (1M context) --- Dockerfile | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Dockerfile b/Dockerfile index 6e85937..bb2af79 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,4 +1,4 @@ -FROM golang:1.24-alpine AS builder +FROM golang:1.25-alpine AS builder WORKDIR /app COPY proto/ /proto/ COPY common/ /common/