diff --git a/README.md b/README.md index a9bc173..2424c76 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,20 @@ # `pgoutbox` - a transactional outbox for `pgx` -`pgoutbox` implements a simple [transactional outbox](https://microservices.io/patterns/data/transactional-outbox.html) for [`pgx`](https://github.com/jackc/pgx). New messages can be added to a Postgres table using `AddMessages` and can be flushed to a destination via `ProcessMessages`. +[![Go Reference](https://pkg.go.dev/badge/github.com/hatchet-dev/pgoutbox.svg)](https://pkg.go.dev/github.com/hatchet-dev/pgoutbox) + +`pgoutbox` implements a simple [transactional outbox](https://microservices.io/patterns/data/transactional-outbox.html) for [`pgx`](https://github.com/jackc/pgx). New messages can be added to a Postgres table within a transaction using `AddMessages` and can be flushed to a destination via `ProcessMessages`. + +## Why? + +While working on [Hatchet](https://github.com/hatchet-dev/hatchet) we needed a reliable and performant way to durably persist messages over a message boundary. In particular, we needed: + +- **Batched reads and writes.** `AddMessages` inserts a batch of messages in a single transaction, and `ProcessMessages` locks a batch, hands the whole batch to one `Flush` call, and deletes it in the same transaction (see [Atomic flush and delete](#atomic-flush-and-delete) and [Benchmarks](#benchmarks)). +- **Exclusive consumers with leasing semantics.** Exactly one instance across a fleet owns a topic under a renewing lease, and a standby takes over within seconds if the holder goes away (see [Exclusive consumers](#exclusive-consumers)). +- **Support for publishing across hundreds of thousands of topics.** Topics are plain strings that don't need to be declared up front, and they're tracked in a table rather than by a poller or worker pool per topic (see [Multiple topics and flushers](#multiple-topics-and-flushers) and [Message expiration](#message-expiration)). + +Without these particular requirements, a library like [River](https://github.com/riverqueue/river) would otherwise have been a good fit. `pgoutbox` is deliberately an outbox rather than a job queue: there are no retries, scheduling, priorities, or job history, and messages are deleted as soon as they're flushed. + +## Example usage Here's an example of flushing messages on `topic1` by simply printing them to the console: @@ -44,7 +58,7 @@ if err != nil { } ``` -## Schema +## Schema migrations By default, `NewOutbox` runs migrations and creates an outbox table in the schema `outbox.messages`. This can be overwritten via: @@ -97,7 +111,7 @@ outbox.ProcessMessages(ctx, "shipments") ## Atomic flush and delete -If your flusher writes to Postgres itself (e.g. into a relay table), use the transaction exposed by the `FlushContext` passed to `Flush`. It is the same transaction `ProcessMessages` uses to lock and delete messages, so your writes and the outbox delete commit or roll back together: +If your flusher writes to Postgres itself (e.g. into a separate table), use the transaction exposed by the `FlushContext` passed to `Flush`. It is the same transaction `ProcessMessages` uses to lock and delete messages, so your writes and the outbox delete commit or roll back together: ```go type relayFlusher struct{} @@ -117,7 +131,7 @@ func (f *relayFlusher) Flush(ctx pgoutbox.FlushContext, msgs []*sqlc.Message) er ## Continuous processing with `Subscribe` -Instead of calling `ProcessMessages` yourself, `Subscribe` runs it in a loop: it drains the topic, then waits until either the poll interval elapses or a new-message notification arrives (see below), and drains again. It blocks until its context is cancelled: +Instead of calling `ProcessMessages` yourself, `Subscribe` drains the topic automatically by listening for a poll or a new-message notification. It blocks until its context is cancelled: ```go go func() { @@ -131,7 +145,7 @@ go func() { }() ``` -Processing errors don't kill the loop — they're logged to the `WithLogger` logger and retried on the next wake-up. `Subscribe` returns an error immediately only if no flusher is registered for the topic or the subscription itself can't be established. +If processing fails, messages are logged to the `WithLogger` logger and retried again on the next wake-up. `Subscribe` returns an error immediately only if no flusher is registered for the topic or the subscription itself can't be established. ### Waking on new messages with `LISTEN`/`NOTIFY` @@ -146,9 +160,7 @@ if err != nil { outbox, err := pgoutbox.NewOutbox(ctx, pool, pgoutbox.WithPubSub(ps)) ``` -With a `PubSub` attached, `AddMessages` publishes a notification for each staged topic and `Subscribe` wakes on it instead of waiting out the poll interval. The pg-backed `PubSub` publishes *inside the `AddMessages` transaction*, so the notification is delivered exactly when the insert commits — and never for a transaction that rolls back. Postgres deduplicates identical notifications within a transaction, so any number of `AddMessages` calls for a topic in one transaction cost a single wake-up. - -A notification only makes sense once the staging transaction has committed, and a `PubSub` that doesn't implement `TxPublisher` has no way to defer a publish to commit time. For those transports, pass a `Notifier` to `AddMessages` and fire it after a successful commit: +You can swap your own pub/sub implementation as well. Note that a notification only makes sense once the inserting transaction has committed, so a `PubSub` that doesn't implement `TxPublisher` needs to defer notifying until after commit. In these cases, pass a `Notifier` to `AddMessages` and fire it after a successful commit: ```go var notifier pgoutbox.Notifier @@ -167,14 +179,9 @@ if err := tx.Commit(ctx); err != nil { notifier.Notify(ctx) // wakes the subscribers of both topics ``` -The `Notifier` accumulates one notification per `AddMessages` call it is passed to, so a single one can serve a whole transaction. With a `TxPublisher` transport like the pg-backed `PubSub`, `Notify` is a no-op (the notification already rode the transaction), so the pattern is transport-agnostic. Skipping it never loses messages — subscribers just fall back to the poll interval — and publish failures inside `Notify` are logged, not returned, for the same reason. +Delivery is best-effort by design; if notifications are lost, `pgoutbox` falls back to polling. -Delivery is best-effort by design: if a notification is lost (for example while the listener reconnects), polling picks the messages up within one poll interval. The listener occupies a single dedicated connection (hijacked out of the pool so it doesn't consume a pool slot) no matter how many topics are subscribed. - -Two details worth knowing: - -- All notifications travel over one NOTIFY channel, `pgoutbox_pubsub` by default. Two outboxes sharing a database (e.g. different schemas) should use distinct channels via `pgoutbox.WithNotifyChannel("my_channel")` to avoid waking each other's subscribers. -- `PubSub` is an interface, so you can bring your own transport (e.g. Redis, NATS) instead of `LISTEN`/`NOTIFY`. If your implementation also implements `TxPublisher` (detected once, at `NewOutbox`), notifications are published transactionally as described above; otherwise use `WithNotifier` and fire `Notify` after commit to publish best-effort. +The default LISTEN/NOTIFY channel is `pgoutbox_pubsub`, which can be configured via `pgoutbox.WithNotifyChannel("my_channel")` if necessary. ## Message expiration @@ -194,7 +201,7 @@ if err != nil { } ``` -The outbox always runs a background scanner goroutine that polls the `topics` table, but maintenance loops are only launched for topics that actually have an expiration configured. Topics don't need to be declared at startup: the library tracks every topic that receives a message in a `topics` table (via a Postgres trigger) and applies the default expiration automatically. +The outbox always runs a background scanner goroutine that polls the `topics` table, but maintenance loops are only launched for topics that actually have an expiration configured. Topics don't need to be declared at startup: the library tracks every topic that receives a message in a `topics` table and applies the default expiration automatically. Multiple outbox instances (e.g. replicas of the same service) coordinate cleanup using a per-topic maintenance lease, so only one instance runs the delete at a time. @@ -207,11 +214,9 @@ outbox, err := pgoutbox.NewOutbox(ctx, pool, ) ``` -Lease competition — another instance winning the cleanup race — is not logged. - ## Exclusive consumers -By default, any number of `ProcessMessages` callers can drain a topic concurrently (each call grabs a non-overlapping batch using `FOR UPDATE SKIP LOCKED`). Use `AcquireTopic` when you need exactly one active consumer at a time. +By default, any number of `ProcessMessages` callers can drain a topic concurrently (using `FOR UPDATE SKIP LOCKED`). Use `AcquireTopic` or `pgoutbox.WithExclusive()` passed to `Subscribe` when you need exactly one active consumer at a time. `AcquireTopic` blocks until this instance holds the exclusive lease for the topic, then returns. A background goroutine automatically renews the lease until the context is cancelled, at which point the lease expires and another instance can take over: @@ -240,9 +245,7 @@ An instance that never acquired the lease (or whose lease has expired) receives When the holder's context is cancelled, the lease expires naturally (within the lease duration, 30 s by default) and another instance's `AcquireTopic` call unblocks. -### Exclusive subscribers - -`Subscribe` composes with exclusive consumers: pass `WithExclusive()` and it manages the lease for you. +You can also pass `pgoutbox.WithExclusive()` to `Subscribe`: ```go // Exactly one instance across the fleet drains "orders" at a time; the rest @@ -256,19 +259,17 @@ With `WithExclusive()`, `Subscribe`: 2. re-acquires it automatically if the lease is ever lost mid-subscribe (for example, a heartbeat lapse during a database blip); and 3. releases it on return, so a waiting instance takes over immediately instead of waiting out the lease's grace period. -Several instances calling `Subscribe(..., WithExclusive())` on the same topic therefore form a failover group: one active consumer, the rest hot standbys. Nothing is missed during a handoff — the new holder's first drain pass covers any backlog that accumulated while the lease changed hands. - -Without `WithExclusive()`, subscribing to a topic whose exclusive lease is held elsewhere doesn't fail — every pass errors (visible via `WithLogger`) and is retried, so the subscriber sits idle until the lease frees up or is acquired. For an exclusive topic, either call `AcquireTopic` before `Subscribe` or pass `WithExclusive()`. +Several instances calling `Subscribe(..., WithExclusive())` on the same topic therefore form a failover group: one active consumer, the rest hot standbys. Without `WithExclusive()`, subscribing to a topic whose exclusive lease is held elsewhere doesn't fail — every pass errors (visible via `WithLogger`) and is retried, so the subscriber sits idle until the lease frees up or is acquired. For an exclusive topic, either call `AcquireTopic` before `Subscribe` or pass `WithExclusive()`. ## Benchmarks -You can run benchmarks locally; for example, to write and flush 100k messages, you can run: +On a local Macbook with an M3 Max core, `pgoutbox` reaches `223817 msgs/sec` with batching support. You can run benchmarks locally; for example, to write and flush 100k messages, you can run: ``` go test -bench=. -benchtime=100000x ``` -`BenchmarkOutbox_WriteAndPublishThroughput` drains each topic with a busy-polling `ProcessMessages` loop; `BenchmarkOutbox_SubscribeThroughput` drains each topic with a `Subscribe` call woken by `LISTEN`/`NOTIFY` (its poll interval is set far above the benchmark runtime, so throughput is carried entirely by notifications — and each producer commit pays the in-transaction `pg_notify`). Both run a matrix of 1 and 10 topics (messages spread round-robin, one consumer per topic) and producer batch sizes of 1, 10, and 100 messages per `AddMessages` transaction. On a local Macbook with an M3 Max core: +Full results from a local run: ``` $ go test -bench=. -benchtime=100000x @@ -301,5 +302,3 @@ BenchmarkOutbox_SubscribeThroughput/TxFlush/topics=10/batch=10-14 1 BenchmarkOutbox_SubscribeThroughput/Flush/topics=10/batch=100-14 100000 5352 ns/op 186840 msgs/sec BenchmarkOutbox_SubscribeThroughput/TxFlush/topics=10/batch=100-14 100000 39021 ns/op 25627 msgs/sec ``` - -Batching producer writes is the single biggest lever: staging 100 messages per transaction reaches ~150-220k msgs/sec on the cheap-flusher path, roughly 20× the one-message-per-transaction rate. The `batch=1` cells run long enough to be sensitive to background database noise (checkpoints, autovacuum), so expect swings between runs — the 1082 msgs/sec outlier above measures ~3300 msgs/sec in isolation. diff --git a/doc.go b/doc.go new file mode 100644 index 0000000..1482199 --- /dev/null +++ b/doc.go @@ -0,0 +1,48 @@ +// Package pgoutbox implements a simple transactional outbox for pgx. +// +// New messages are added to a Postgres table within a transaction using +// [Outbox.AddMessages] and are flushed to a destination via +// [Outbox.ProcessMessages]. Because the insert rides the caller's own +// transaction, a message is only ever visible to consumers if the business +// write it belongs to committed. +// +// Messages are grouped by topic, and each topic is drained by the [Flusher] +// registered for it with [Outbox.AddFlusher]. ProcessMessages locks a batch of +// messages, hands the whole batch to a single Flush call, and deletes the +// messages in the same transaction if the flush succeeds. A Flusher that +// writes to Postgres itself can use [FlushContext.Tx] so its writes and the +// outbox delete commit or roll back together. +// +// [Outbox.Subscribe] drains a topic continuously, waking on a poll interval +// or, when a [PubSub] such as [NewPGPubSub] is attached via [WithPubSub], the +// moment new messages commit. [Outbox.AcquireTopic] and [WithExclusive] give +// a topic exactly one active consumer across a fleet, backed by a renewing +// lease with automatic failover. [WithTopicExpiration] and +// [WithDefaultExpiration] delete old messages in the background. +// +// A minimal setup looks like: +// +// type printFlusher struct{} +// +// func (printFlusher) Flush(_ pgoutbox.FlushContext, msgs []*sqlc.Message) error { +// for _, m := range msgs { +// fmt.Printf("flushed id=%d topic=%s payload=%s\n", m.ID, m.Topic, string(m.Payload)) +// } +// return nil +// } +// +// outbox, err := pgoutbox.NewOutbox(ctx, pool) +// if err != nil { +// return err +// } +// outbox.AddFlusher("orders", printFlusher{}) +// +// // within a transaction +// err = outbox.AddMessages(ctx, tx, "orders", []pgoutbox.MessageOpts{{Payload: payload}}) +// +// // after the transaction commits +// _, err = outbox.ProcessMessages(ctx, "orders") +// +// See the README at https://github.com/hatchet-dev/pgoutbox for a longer +// walkthrough of each feature. +package pgoutbox diff --git a/migrate.go b/migrate.go index 28bc185..2e3a798 100644 --- a/migrate.go +++ b/migrate.go @@ -30,14 +30,13 @@ func validateSchemaName(schema string) error { return nil } -// Migrate runs the embedded pgoutbox migrations against the given pool. -// It is the explicit alternative to NewOutbox's auto-migration: callers -// that want to control when DDL runs (separate startup phase, release -// pipeline, etc.) should construct the outbox with WithAutoMigrate(false) -// and invoke Migrate themselves. +// Migrate runs the embedded pgoutbox migrations against pool, creating the +// schema if needed. NewOutbox runs them automatically by default. To run them +// yourself instead, for example as part of a separate release step, construct +// the outbox with WithAutoMigrate(false) and call Migrate explicitly. // -// Only WithSchema is consulted from opts; other options are accepted for -// API symmetry but ignored. +// Only WithSchema is consulted from opts. Other options are accepted so the +// same option list can be passed to NewOutbox, but they are ignored here. func Migrate(ctx context.Context, pool *pgxpool.Pool, opts ...OutboxOpt) error { o := defaultOpts() for _, f := range opts { diff --git a/no_op_flusher.go b/no_op_flusher.go index 2c5092d..07c8a7d 100644 --- a/no_op_flusher.go +++ b/no_op_flusher.go @@ -2,12 +2,16 @@ package pgoutbox import "github.com/hatchet-dev/pgoutbox/sqlc" +// NopFlusher is a Flusher that discards every message it is given. It is useful +// in tests, and for draining a topic whose messages are no longer needed. type NopFlusher struct{} +// NewNopFlusher returns a NopFlusher. func NewNopFlusher() *NopFlusher { return &NopFlusher{} } +// Flush discards msgs and returns nil. func (f *NopFlusher) Flush(_ FlushContext, _ []*sqlc.Message) error { return nil } diff --git a/outbox.go b/outbox.go index 77182ab..4da88a0 100644 --- a/outbox.go +++ b/outbox.go @@ -20,9 +20,10 @@ import ( ) // FlushContext is the context passed to Flusher.Flush. It embeds -// context.Context and exposes the transaction that ProcessMessages uses to -// lock and delete messages. Callers that want their writes to commit -// atomically with the outbox delete can enlist in that transaction via Tx(). +// context.Context, so flushers that don't need the transaction can treat it as +// a plain context. Tx returns the transaction ProcessMessages uses to lock and +// delete messages. Flushers that write to Postgres themselves can use it so +// their writes and the outbox delete commit or roll back together. type FlushContext interface { context.Context Tx() pgx.Tx @@ -35,68 +36,106 @@ type flushContext struct { func (f *flushContext) Tx() pgx.Tx { return f.tx } +// Flusher delivers a batch of messages to their destination. Register one per +// topic with Outbox.AddFlusher. ProcessMessages calls Flush with the messages it +// has locked for the topic and deletes them only if Flush returns nil. If Flush +// returns an error, the messages stay in the outbox and are retried on the next +// call. +// +// Flushers that write to Postgres can use ctx.Tx() so that their writes and the +// outbox delete commit or roll back together. type Flusher interface { Flush(ctx FlushContext, msgs []*sqlc.Message) error } +// MessageOpts describes a single message to add via Outbox.AddMessages. type MessageOpts struct { + // Payload is the opaque message body. It is stored as-is and handed back to + // the Flusher unchanged. Payload []byte } +// Outbox is a transactional outbox: messages are added to a Postgres table +// within the caller's transaction and later flushed to a destination by the +// Flusher registered for their topic. Create one with NewOutbox. +// +// Topics are plain strings and don't need to be declared up front. Any number +// of topics can be used, and every topic that receives a message is tracked +// automatically. type Outbox interface { + // AddFlusher registers the Flusher that ProcessMessages and Subscribe use to + // drain topic. Registering a flusher for a topic that already has one + // replaces it. AddFlusher(topic string, flusher Flusher) - // AddMessages stages msgs on the topic within the caller's transaction. - // When a PubSub is configured, it also arranges the new-message - // notification that wakes Subscribe callers: TxPublisher transports - // publish it on tx itself, and generic transports hand it to the Notifier - // passed via WithNotifier, for the caller to fire after commit (see - // Notifier). Skipping the option never loses messages, it only leaves - // generic transports waiting out Subscribe's poll interval. + // AddMessages adds msgs to the topic within the caller's transaction. The + // messages become visible to ProcessMessages once tx commits, and are + // discarded if it rolls back. + // + // When the outbox was built with WithPubSub, AddMessages also notifies + // subscribers of the topic. A TxPublisher such as NewPGPubSub publishes the + // notification inside tx, so it is delivered exactly when the transaction + // commits. Other PubSub implementations need to defer the notification + // until after commit: pass a Notifier via WithNotifier and call Notify once + // the transaction has committed. Skipping the notification never loses + // messages; subscribers pick them up on their next poll. AddMessages(ctx context.Context, tx pgx.Tx, topic string, msgs []MessageOpts, opts ...AddOpt) error - // ProcessMessages grabs a batch of messages for the given topic, flushes them using the registered Flusher for that - // topic, and deletes them from the outbox if the flush is successful. If the topic has an active exclusive consumer, - // the calling instance must hold the exclusive lease (via AcquireTopic) or an error is returned. + // ProcessMessages grabs a batch of messages for the topic, flushes them + // using the registered Flusher, and deletes them from the outbox in the + // same transaction if the flush succeeds. It returns the flushed messages, + // or nil if the topic was empty. The batch size defaults to 1000 and can be + // changed with WithBatchSize. + // + // Messages are locked with FOR UPDATE SKIP LOCKED, so any number of callers + // can drain a topic concurrently. If the topic has an active exclusive + // consumer, the caller must hold the lease via AcquireTopic; otherwise + // ErrExclusiveLeaseHeld or ErrExclusiveLeaseRequired is returned. ProcessMessages(ctx context.Context, topic string, opts ...ProcessOpt) ([]*sqlc.Message, error) - // Subscribe blocks and continuously drains the topic: it runs - // ProcessMessages until the topic is empty, then waits for the poll - // interval to elapse — or, when the outbox was built with WithPubSub, for - // a new-message notification — and drains again. Processing errors are - // logged to the WithLogger logger and retried on the next wake-up; as - // with ProcessMessages, topics with an active exclusive consumer require - // AcquireTopic first — either call it beforehand, or pass WithExclusive - // to have Subscribe acquire, re-acquire, and release the lease itself. - // Returns ctx.Err() when ctx ends, or an error immediately if no flusher - // is registered for the topic, the PubSub subscription cannot be - // established, or the WithExclusive initial acquisition fails. + // Subscribe drains the topic automatically by listening for a poll or a + // new-message notification. It runs ProcessMessages until the topic is + // empty, then waits for the poll interval (see WithPollInterval) or, when + // the outbox was built with WithPubSub, for a new-message notification, and + // drains again. It blocks until ctx is cancelled and then returns ctx.Err(). + // + // If processing fails, the error is logged to the WithLogger logger and the + // messages are retried on the next wake-up. Subscribe returns an error + // immediately only if no flusher is registered for the topic, the PubSub + // subscription can't be established, or the initial lease acquisition for + // WithExclusive fails. + // + // As with ProcessMessages, a topic with an active exclusive consumer + // requires the lease. Either call AcquireTopic before Subscribe or pass + // WithExclusive to have Subscribe manage the lease itself. Subscribe(ctx context.Context, topic string, opts ...SubscribeOpt) error - // AcquireTopic blocks until this instance holds the exclusive processing lease - // for the named topic, then returns. A background goroutine automatically renews - // the lease until ctx is cancelled or ReleaseTopic is called, at which point the - // lease expires naturally and another instance can take over. AcquireTopic must - // be called before ProcessMessages for any topic that has an active exclusive - // consumer. + // AcquireTopic blocks until this instance holds the exclusive lease for the + // topic, then returns. A background goroutine automatically renews the + // lease until ctx is cancelled or ReleaseTopic is called, at which point + // the lease expires and another instance can take over. + // + // Once a topic has an exclusive consumer, only the lease holder may call + // ProcessMessages for it. Other instances receive ErrExclusiveLeaseHeld, and + // an instance whose lease has expired receives ErrExclusiveLeaseRequired + // until it calls AcquireTopic again. AcquireTopic(ctx context.Context, topic string) error - // ReleaseTopic stops renewing and immediately expires the exclusive lease - // this instance holds for topic, letting another instance acquire it right - // away instead of waiting out the lease duration. It is a no-op if this - // instance does not currently hold the lease. As with a naturally expired - // lease, a subsequent ProcessMessages call still requires an explicit - // AcquireTopic first. + // ReleaseTopic stops renewing the exclusive lease this instance holds for + // topic and expires it immediately, so a waiting instance can take over + // right away instead of waiting out the lease duration. It is a no-op if + // this instance does not hold the lease. A later ProcessMessages call for + // the topic requires AcquireTopic again. ReleaseTopic(ctx context.Context, topic string) error } // ErrExclusiveLeaseHeld is returned by ProcessMessages when another outbox -// instance currently holds a valid exclusive lease for the topic. +// instance currently holds the exclusive lease for the topic. var ErrExclusiveLeaseHeld = errors.New("exclusive lease held by another instance") // ErrExclusiveLeaseRequired is returned by ProcessMessages when the topic has -// an exclusive-consumer record but this instance does not hold a live lease — -// either AcquireTopic was never called or the lease has since expired. +// an exclusive consumer but this instance does not hold a live lease, either +// because AcquireTopic was never called or because the lease has expired. var ErrExclusiveLeaseRequired = errors.New("exclusive lease required: call AcquireTopic first") // defaultBatchSize is the number of messages ProcessMessages will pull per @@ -110,35 +149,34 @@ type addOpts struct { notifier *Notifier } -// WithNotifier has AddMessages collect its post-commit notification into n -// instead of dropping it. Only generic (non-TxPublisher) PubSubs need it — -// they have no way to defer a publish to commit time, so the caller carries -// the notification past the transaction and fires it with Notify. One -// Notifier can be shared by every AddMessages call in a transaction and -// fired once after commit. +// WithNotifier has AddMessages record its new-message notification in n so the +// caller can publish it after the transaction commits. It is only needed for +// PubSub implementations that don't implement TxPublisher, since those can't +// defer publishing to commit time. One Notifier can be shared by every +// AddMessages call in a transaction and fired once with Notify after commit. func WithNotifier(n *Notifier) AddOpt { return func(opts *addOpts) { opts.notifier = n } } -// Notifier accumulates the new-message notifications of the AddMessages calls -// it is passed to (via WithNotifier), so they can be published once the -// staging transaction has committed. The zero value is ready to use; it is -// not safe for concurrent use, mirroring the pgx.Tx it accompanies. +// Notifier accumulates the new-message notifications from the AddMessages +// calls it is passed to via WithNotifier, so they can be published once the +// transaction has committed. The zero value is ready to use. Like the pgx.Tx it +// accompanies, it is not safe for concurrent use. type Notifier struct { hooks []func(context.Context) } -// Notify publishes the accumulated notifications. Invoke it once, after the -// transaction commits successfully; after a rollback, simply discard the -// Notifier. It is a no-op when there is nothing to publish — no PubSub -// configured, no messages staged, or a TxPublisher transport that already -// published on the transaction. Publishing is best-effort: failures are -// logged to the WithLogger logger, not returned, since durably staged -// messages are picked up by Subscribe's polling fallback regardless. Calling -// it more than once just repeats the wake-ups (harmless, like any spurious -// notification). +// Notify publishes the accumulated notifications. Call it once after the +// transaction commits successfully. After a rollback, discard the Notifier +// instead. Notify is a no-op when there is nothing to publish, including when +// the PubSub implements TxPublisher and already published on the transaction. +// +// Delivery is best-effort by design: failures are logged to the WithLogger +// logger rather than returned, since subscribers fall back to polling and pick +// the messages up regardless. Calling Notify more than once only repeats the +// wake-ups, which is harmless. func (n *Notifier) Notify(ctx context.Context) { for _, hook := range n.hooks { hook(ctx) @@ -156,9 +194,9 @@ func defaultProcessOpts() *processOpts { return &processOpts{batchSize: defaultBatchSize} } -// WithBatchSize sets the maximum number of messages ProcessMessages will -// acquire and hand to the Flusher in a single call. Must be > 0. Values -// above math.MaxInt32 are ignored and the default (1000) is used instead. +// WithBatchSize sets the maximum number of messages ProcessMessages locks and +// hands to the Flusher in a single call. Defaults to 1000. Values that are not +// positive or that exceed math.MaxInt32 are ignored. func WithBatchSize(n int) ProcessOpt { return func(opts *processOpts) { if n <= 0 || n > math.MaxInt32 { @@ -297,71 +335,75 @@ type outboxImpl struct { managed map[string]*managedTopic } +// OutboxOpt configures NewOutbox and Migrate. type OutboxOpt func(*outboxImplOpts) -func WithSchema(searchPath string) OutboxOpt { +// WithSchema sets the Postgres schema that holds the outbox tables. Defaults to +// "outbox". The schema is created if it doesn't exist when migrations run. +func WithSchema(schema string) OutboxOpt { return func(opts *outboxImplOpts) { - opts.schema = searchPath + opts.schema = schema } } -// WithAutoMigrate controls whether NewOutbox runs the embedded migrations on -// construction. Defaults to true. Set to false when the caller wants to run -// migrations explicitly via Migrate (for example, in a separate startup -// phase or release pipeline). +// WithAutoMigrate controls whether NewOutbox runs the embedded migrations. +// Defaults to true. Set it to false to run migrations yourself by calling +// Migrate, for example as part of a separate release step. func WithAutoMigrate(enabled bool) OutboxOpt { return func(opts *outboxImplOpts) { opts.autoMigrate = enabled } } -// WithTopicExpiration registers a TTL for the named topic. On Start, the TTL is -// written to the topics table so that any outbox instance can discover it. -// Messages older than ttl are eligible for deletion by the background maintenance -// goroutine launched by Start. Per-topic TTLs take precedence over -// WithDefaultExpiration. +// WithTopicExpiration sets how long messages on topic are kept before the +// background maintenance goroutines delete them. NewOutbox writes the +// expiration to the topics table so that every outbox instance can discover +// it. Per-topic expirations take precedence over WithDefaultExpiration. func WithTopicExpiration(topic string, ttl time.Duration) OutboxOpt { return func(opts *outboxImplOpts) { opts.expirations[topic] = ttl } } -// WithDefaultExpiration sets a fallback TTL used for topics that have no -// specific expiration configured via WithTopicExpiration. Any topic that -// appears in the topics table with a NULL expiration_nanos will be maintained -// using this TTL when Start is running. +// WithDefaultExpiration sets the expiration for topics that have none +// configured via WithTopicExpiration. Topics don't need to be declared at +// startup: every topic that receives a message is tracked in the topics table +// and picks up the default automatically. Without a default, topics that have +// no explicit expiration are never expired. func WithDefaultExpiration(ttl time.Duration) OutboxOpt { return func(opts *outboxImplOpts) { opts.defaultExpiration = ttl } } -// WithLogger attaches a zerolog logger that receives error-level messages from -// the background maintenance goroutines. Lease competition (another instance -// holding the lease) is not logged. If not set, maintenance errors are silent. +// WithLogger attaches a zerolog logger that receives errors from the +// background maintenance goroutines, from Subscribe's processing passes, and +// from Notifier.Notify. If not set, those errors are silent. func WithLogger(l zerolog.Logger) OutboxOpt { return func(opts *outboxImplOpts) { opts.logger = l } } -// WithPubSub attaches a PubSub used to cut end-to-end latency: AddMessages -// publishes a notification for each staged topic and Subscribe wakes on those -// notifications instead of waiting out its poll interval. Delivery is -// best-effort — Subscribe's polling remains the fallback for lost -// notifications. If ps also implements TxPublisher (NewPGPubSub does), the -// notification is published inside the AddMessages transaction and delivered -// exactly when it commits; otherwise pass a Notifier to AddMessages via -// WithNotifier and invoke Notify after committing. +// WithPubSub attaches a PubSub that wakes Subscribe callers the moment new +// messages commit, instead of waiting out the poll interval. NewPGPubSub +// provides an implementation built on Postgres LISTEN/NOTIFY. +// +// Delivery is best-effort by design; if a notification is lost, Subscribe falls +// back to polling. If ps implements TxPublisher, the notification is published +// inside the AddMessages transaction. Otherwise, pass a Notifier to AddMessages +// via WithNotifier and call Notify after commit. func WithPubSub(ps PubSub) OutboxOpt { return func(opts *outboxImplOpts) { opts.pubsub = ps } } -// NewOutbox creates an outbox backed by pool and starts the background -// maintenance goroutines. The goroutines run until ctx is cancelled; pass a -// context tied to your application lifetime (e.g. from signal.NotifyContext). +// NewOutbox creates an outbox backed by pool. By default it runs the embedded +// migrations, creating the outbox tables in the "outbox" schema (see WithSchema +// and WithAutoMigrate), and starts the background maintenance goroutines that +// delete expired messages. The goroutines run until ctx is cancelled, so pass a +// context tied to your application lifetime. func NewOutbox(ctx context.Context, pool *pgxpool.Pool, fs ...OutboxOpt) (Outbox, error) { opts := defaultOpts() diff --git a/pgpubsub.go b/pgpubsub.go index 20ceffe..7fc1e1c 100644 --- a/pgpubsub.go +++ b/pgpubsub.go @@ -43,10 +43,10 @@ func defaultPGPubSubOpts() *pgPubSubOpts { // PGPubSubOpt configures the PubSub returned by NewPGPubSub. type PGPubSubOpt func(*pgPubSubOpts) -// WithNotifyChannel overrides the Postgres NOTIFY channel the PubSub -// multiplexes over. All messages on a channel are broadcast to every listener -// of that channel, so two outboxes sharing a database (e.g. different -// schemas) should use distinct channels to avoid spurious wake-ups. +// WithNotifyChannel sets the Postgres LISTEN/NOTIFY channel that carries the +// PubSub's messages. Defaults to "pgoutbox_pubsub". Two outboxes sharing a +// database should use distinct channels so they don't wake each other's +// subscribers. func WithNotifyChannel(name string) PGPubSubOpt { return func(opts *pgPubSubOpts) { opts.channel = name @@ -54,8 +54,8 @@ func WithNotifyChannel(name string) PGPubSubOpt { } // WithNotifyLogger attaches a zerolog logger that receives errors from the -// background listener (connection failures, malformed payloads). If not set, -// those errors are silent. +// background listener, such as connection failures and malformed payloads. If +// not set, those errors are silent. func WithNotifyLogger(l zerolog.Logger) PGPubSubOpt { return func(opts *pgPubSubOpts) { opts.logger = l @@ -90,15 +90,16 @@ type pgPubSub struct { } // NewPGPubSub returns a PubSub backed by Postgres LISTEN/NOTIFY on the given -// pool. The background listener starts lazily on the first Sub call and runs -// until ctx is cancelled; pass a context tied to your application lifetime. +// pool. Pass it to WithPubSub to wake subscribers the moment new messages +// commit. The background listener starts on the first Sub call and runs until +// ctx is cancelled, so pass a context tied to your application lifetime. // -// The returned PubSub implements TxPublisher, so an outbox configured with it -// publishes new-message notifications transactionally: subscribers wake when -// the staging transaction commits, and not at all if it rolls back. +// The returned PubSub implements TxPublisher, so notifications are published +// inside the AddMessages transaction: subscribers wake when it commits, and not +// at all if it rolls back. // -// NOTIFY payloads are capped by Postgres at roughly 8000 bytes; Pub returns -// an error beyond that. The outbox's own notifications are empty. +// Postgres caps NOTIFY payloads at roughly 8000 bytes, and Pub returns an error +// beyond that. The outbox's own notifications carry no payload. func NewPGPubSub(ctx context.Context, pool *pgxpool.Pool, fs ...PGPubSubOpt) (PubSub, error) { opts := defaultPGPubSubOpts() diff --git a/pubsub.go b/pubsub.go index de2fee1..1aa1ff4 100644 --- a/pubsub.go +++ b/pubsub.go @@ -18,13 +18,15 @@ type PubSubMessage struct { } // PubSub is a minimal publish/subscribe transport for small notification -// messages. The outbox uses it (via WithPubSub) to wake Subscribe callers as -// soon as new messages are staged, instead of waiting out a poll interval. +// messages. The outbox uses it via WithPubSub to wake Subscribe callers as soon +// as new messages commit, instead of waiting out the poll interval. NewPGPubSub +// provides an implementation built on Postgres LISTEN/NOTIFY, and you can swap +// in your own transport by implementing this interface. // -// Delivery is expected to be best-effort: implementations may drop messages -// under load or while disconnected. The outbox tolerates both lost messages -// (Subscribe falls back to polling) and duplicate or spurious messages (an -// extra processing pass on an empty topic is a no-op). +// Delivery is best-effort by design. Implementations may drop messages under +// load or while disconnected: the outbox falls back to polling for lost +// messages, and an extra processing pass on an empty topic is a no-op, so +// duplicate or spurious messages are harmless. type PubSub interface { // Pub publishes payload to topic. Pub(ctx context.Context, topic string, payload []byte) error @@ -38,12 +40,12 @@ type PubSub interface { } // TxPublisher is an optional interface a PubSub can implement to publish -// within a pgx transaction. When the PubSub configured via WithPubSub -// implements it (detected once, at NewOutbox), AddMessages publishes its -// new-message notification inside the caller's transaction, so the -// notification is delivered exactly when the insert commits — and never for a -// transaction that rolls back. Without it, the notification is deferred to a -// Notifier the caller passes via WithNotifier and invokes after commit. +// within a pgx transaction. A notification only makes sense once the inserting +// transaction has committed. When the PubSub passed to WithPubSub implements +// TxPublisher, AddMessages publishes its notification inside the caller's +// transaction, so it is delivered exactly when the insert commits and never for +// a transaction that rolls back. A PubSub that doesn't implement it needs to +// defer notifying until after commit; see WithNotifier. type TxPublisher interface { PubInTx(ctx context.Context, tx pgx.Tx, topic string, payload []byte) error } diff --git a/subscribe.go b/subscribe.go index 8797cc9..6a57c2e 100644 --- a/subscribe.go +++ b/subscribe.go @@ -26,7 +26,8 @@ func defaultSubscribeOpts() *subscribeOpts { } // WithPollInterval sets how long Subscribe waits between processing passes -// when no new-message notification arrives. Must be > 0. +// when no new-message notification arrives. Defaults to 5 seconds. Values that +// are not positive are ignored. func WithPollInterval(d time.Duration) SubscribeOpt { return func(opts *subscribeOpts) { if d <= 0 { @@ -36,22 +37,23 @@ func WithPollInterval(d time.Duration) SubscribeOpt { } } -// WithProcessOpts forwards per-call ProcessMessages options (e.g. -// WithBatchSize) to every processing pass Subscribe makes. +// WithProcessOpts forwards ProcessMessages options such as WithBatchSize to +// every processing pass Subscribe makes. func WithProcessOpts(popts ...ProcessOpt) SubscribeOpt { return func(opts *subscribeOpts) { opts.processOpts = append(opts.processOpts, popts...) } } -// WithExclusive makes Subscribe manage the topic's exclusive-consumer lease -// for the duration of the call: it acquires the lease before the first -// processing pass (blocking, like AcquireTopic, while another instance holds -// it), re-acquires it if it is ever lost mid-subscribe, and releases it on -// return so a waiting instance can take over immediately instead of waiting -// out the lease's grace period. Several instances calling Subscribe with -// WithExclusive on the same topic therefore form a failover group: exactly -// one drains the topic while the rest block in line behind the lease. +// WithExclusive has Subscribe manage the topic's exclusive lease for the +// duration of the call. Subscribe then acquires the lease before its first +// processing pass, blocking while another instance holds it (like +// AcquireTopic); re-acquires it automatically if it is ever lost +// mid-subscribe; and releases it on return, so a waiting instance takes over +// immediately instead of waiting out the lease duration. +// +// Several instances calling Subscribe with WithExclusive on the same topic +// therefore form a failover group: one active consumer, the rest hot standbys. func WithExclusive() SubscribeOpt { return func(opts *subscribeOpts) { opts.exclusive = true