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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Dockerfile
Original file line number Diff line number Diff line change
@@ -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/
Expand Down
7 changes: 4 additions & 3 deletions go.mod
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
14 changes: 8 additions & 6 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down Expand Up @@ -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=
Expand Down
173 changes: 173 additions & 0 deletions internal/jobs/middleware.go
Original file line number Diff line number Diff line change
@@ -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 = <job.ID>` 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 = <uuid.New()>` 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.<JobKind>` 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:])
}
Loading