Select the complete pack with:
DATABASE=postgres OUTBOX=postgres MESSAGING=nats-jetstream make template-init \
MODULE=github.com/acme/orders CODEOWNER=@acme/backendOUTBOX=postgres requires both PostgreSQL and NATS JetStream. The generator
rejects an incomplete selection instead of producing a relay with no publisher.
OUTBOX=none, the default, removes the event contract, River appender and
worker, migration, command, tests, and this document.
The pack guarantees at-least-once publication:
- The feature's PostgreSQL adapter mutates business state and calls
Appender.Appendthrough the same caller-ownedpgx.Tx. - River inserts one
publish_domain_eventjob in that transaction. A rollback removes both writes; a commit exposes both. cmd/outbox-relayworks the River job through the concrete NATS producer.- A JetStream acknowledgement followed by process loss may run the job again.
Every attempt uses the same logical event ID for both
Message-IdandNats-Msg-Id.
Consumers must make non-idempotent effects duplicate-safe. The pack does not provide exactly-once delivery or generic ordering.
River owns job state, concurrency, retry scheduling, crash rescue, cleanup, operator retry, and maintenance. The template owns only typed event encoding, atomic insertion, subject routing, NATS mapping, and process composition.
Create the immutable event once, before any transaction callback that may be retried:
var orderUpdatedV1 = domainevent.Define[order.UpdatedV1]("order.updated", 1)
event, err := orderUpdatedV1.New(
eventID,
occurredAt,
order.UpdatedV1{
OrderID: orderID,
Revision: revision,
},
)
if err != nil {
return err
}Build the appender once in the service composition root. This is where the service owns its event-to-subject mapping; feature code never receives the subject:
outbox, err := natsjs.NewOutboxAppender(
cfg.Messaging.MaxPayloadBytes,
natsjs.Route{
Type: orderUpdatedV1.Type,
Version: orderUpdatedV1.Version,
Subject: "events.orders",
},
)
if err != nil {
return err
}The PostgreSQL adapter may declare only the method it consumes:
type outboxAppender interface {
Append(context.Context, pgx.Tx, domainevent.Event) error
}Its transaction contains the exact business call:
return postgres.InTx(ctx, pool, pgx.TxOptions{}, func(tx pgx.Tx) error {
if err := orders.Update(ctx, tx, change); err != nil {
return err
}
return outbox.Append(ctx, tx, event)
})The template does not construct this appender in cmd/service: it ships no
business event or repository that could consume it.
domainevent.Event carries only:
- one stable logical ID;
- event type and positive schema version;
- UTC occurrence time;
- JSON encoded from the typed payload.
It carries no source, broker subject, arbitrary metadata, ordering key, publication-attempt identity, retry policy, or trace carrier. The appender resolves the subject, enforces the configured NATS payload bound, and stores the payload bytes unchanged inside River's JSON job args.
Reusing an event ID with the same immutable job is idempotent while River retains
that job. Reusing it with different bytes returns
postgresoutbox.ErrEventIDConflict.
The River OpenTelemetry plugin stores W3C traceparent and tracestate in
job metadata. River's work span links to the producing operation; the NATS
outbox worker additionally restores that original context before calling the
producer, so the broker message retains the producing trace rather than the
relay process's local trace. Missing or malformed trace metadata never blocks
publication.
cmd/outbox-relay runs one River queue named outbox with at most 16 workers.
River polling is used instead of LISTEN because this repository applies a finite
statement_timeout to every pooled PostgreSQL connection; a long-lived
LISTEN session would be cancelled by that safety bound.
There are no APP__OUTBOX__* settings. River owns its retry and rescue
defaults. Cancelled and discarded jobs are retained indefinitely so unpublished
intent cannot disappear through cleanup; completed jobs use River's normal
retention. The relay gives River a code-owned 25-second drain inside the
existing process grace period; HTTP shutdown tuning does not change it.
Use River's job list/retry and queue pause APIs from an authenticated, service-owned operator tool when needed. The template ships no generic admin endpoint and performs no automatic discard.
Monitor durable intent before it reaches NATS with a low-frequency PostgreSQL query against the writer, or against a replica whose lag is inside the alert budget:
SELECT state,
count(*) AS jobs,
min(created_at) AS oldest_created_at
FROM river_job
WHERE queue = 'outbox'
AND state <> 'completed'
GROUP BY state
ORDER BY state;Alert when cancelled or discarded is non-zero, or when the oldest active
job exceeds the publication SLO. The state enum is the only grouping dimension;
event and River job IDs belong in an authenticated drill-down, never metric
labels. NATS Surveyor cannot replace this query because it sees only messages
that already reached the broker.
migrations/000008_river.sql installs the shared River v0.44.0 PostgreSQL
schema. Fresh template-init output removes the legacy
000001_postgres_outbox.sql before SQLC generation, so a new service starts
with River as its only outbox state. The River
river_migration ledger records main versions 1 through 7, so River's migrator
starts from the same baseline. Upgrade River modules together and append the
matching upstream migration delta; never edit an applied generated-service
migration.
The template source retains the former outbox_events,
outbox_ordering_heads, receipt, and redrive tables for existing adopters'
rollback lineage, but the River worker does not read them. A service that
deployed the older pack must drain or explicitly bridge its remaining rows
before switching workers. Dropping those tables or pending rows from an
existing database is a separate authorized production action.
go test -vet=off ./internal/domainevent ./internal/infra/postgresoutbox \
./internal/infra/natsjs ./cmd/outbox-relay/...
go test -vet=off -tags=integration ./test -run '^TestPostgresOutbox' -count=1
make sqlc-check migration-check test-outbox-race
ALLOW_HEAVY=1 make template-init-check