Documentation
¶
Overview ¶
Package broker defines Warren's messaging ports: one driver-neutral message envelope, a publisher, a subscriber, and a message handler.
Consumers written against these types run identically over Kafka, RabbitMQ, NATS, and the in-process broker. A consumer that touches a driver's record type loses that property, which is why the raw client is an explicit escape hatch and not the default path. The port carries zero implementations (invariant 5): the consumer middleware chain and the per-subscription options land with the runtime, and every driver runs the same exported contract suite.
Index ¶
- Constants
- func DeliveryHeaders(ctx context.Context) map[string]string
- func ExponentialBackoff(attempts int) app.RetryPolicy
- func InjectTrace(ctx context.Context, msgs []Message)
- func WithDeliveryHeaders(ctx context.Context, h map[string]string) context.Context
- type Message
- type MessageHandler
- type Middleware
- func ConcurrencyLimit(n int) Middleware
- func DeadLetter(pub Publisher, originTopic, dlqTopic string) Middleware
- func Deduplicate(subscription string, store inbox.Store, ttl time.Duration) Middleware
- func Drain() (Middleware, func(context.Context) error)
- func Recover() Middleware
- func RequireMessageID() Middleware
- func Retry(policy app.RetryPolicy) Middleware
- func TraceExtract() Middleware
- type Publisher
- type Redeliverer
- type SubscribeOption
- type Subscriber
Constants ¶
const CorrelationHeader = "correlation-id"
CorrelationHeader is the message header the correlation ID travels in. It is the broker-side counterpart of transport/http's X-Correlation-Id, and it is lower-case because a broker header map is a plain map of strings with no canonicalisation rules — a driver that round-trips keys verbatim and one that lower-cases them must agree, so the wire form is fixed here.
Variables ¶
This section is empty.
Functions ¶
func DeliveryHeaders ¶
DeliveryHeaders returns the headers of the message being processed, or nil when the context carries none. The returned map is the delivery's own — read it, don't mutate it.
func ExponentialBackoff ¶
func ExponentialBackoff(attempts int) app.RetryPolicy
ExponentialBackoff returns a bounded exponential app.RetryPolicy: attempts total attempts, base 100ms doubling per attempt, capped at 30s, with full jitter. Boot-time construction, request-path arithmetic only.
func InjectTrace ¶ added in v0.2.0
InjectTrace writes the trace context on ctx into each message's Headers, allocating a map only for a message that has none. A publisher adapter calls it as its first act, and it is what makes a span survive the trip through a broker into the consumer — the other half of the chain's TraceExtract stage.
It is a no-op when no telemetry is bound, so an uninstrumented service pays one nil check per publish.
The OUTBOX calls it when the row is WRITTEN, not when it is relayed: the relay runs long after the request's span ended, and a span parented to the relay's own context is a trace nobody can follow back to the request that caused it. It never OVERWRITES a header that is already there, for the same reason Correlating does not: the outbox stamps at Append, inside the request, and the relay publishes from a context whose span is its OWN. An overwriting injector would reparent every event to the drain that happened to carry it, which is a trace that leads back to a timer.
Types ¶
type Message ¶
type Message struct {
// ID is the idempotency key. Inbox dedupe is keyed on it, so it must be
// stable across redeliveries of the same fact.
ID string
// Type is the fact's name, such as "user.registered" — the same value a
// domain.Event reports from EventName.
Type string
// Key is the partition or routing key.
Key string
// Payload is the encoded body. This package does not define its format.
Payload []byte
// Headers carries metadata across the broker; trace context propagates
// here as strings, so a span survives the trip into the consumer without
// the core module knowing a telemetry SDK exists.
Headers map[string]string
// OccurredAt is when the fact happened, not when it was published.
OccurredAt time.Time
}
Message is the driver-neutral envelope every adapter translates to and from.
type MessageHandler ¶
MessageHandler processes one message. Returning nil acknowledges it; returning an error hands it to the retry and dead-letter middleware, which decide by the error's warren/errors code — the consumer column of the warren.md §2.6 table. A handler that maps a code to an ack decision itself has broken the ring: that decision belongs to the chain.
func Chain ¶
func Chain(h MessageHandler, mw ...Middleware) MessageHandler
Chain composes middleware around a handler; mw[0] is the outermost, the same convention as app.Chain. Like app.Chain it refuses nil at composition time — a boot-time panic naming the position instead of a request-time nil dereference (the AGENT.md-sanctioned guard).
func Pipeline ¶
func Pipeline(subscription, topic string, h MessageHandler, store inbox.Store, dlq Publisher, opts ...SubscribeOption) (MessageHandler, func(context.Context) error)
Pipeline assembles the full consumer chain around h for one subscription, in the fixed wrapping order:
Recover → Drain → TraceExtract → Deduplicate → DeadLetter → RequireMessageID → Retry → ConcurrencyLimit → handler
and returns the composed handler plus the drain-wait the lifecycle calls at shutdown step 3. Recover guards twice: innermost around the handler, so a handler panic becomes INTERNAL and takes the §2.6 path — retried, then dead-lettered — and outermost as a safety net, so a bug in a stage itself cannot kill the subscription either. topic is the subscription's topic — it names the default DLQ ("<topic>.dlq") and the origin-topic header on dead-lettered messages. A nil handler, store (unless WithoutDedupe), or dead-letter publisher panics here, at boot.
subscription names THIS subscription and must be unique in the process. It scopes the deduplication key, which is what lets two features consume one topic: keyed on the message id alone, whichever handler ran first marked the message seen and the second never saw it.
type Middleware ¶
type Middleware func(MessageHandler) MessageHandler
Middleware decorates a MessageHandler — the consumer ring's mirror of app.Middleware, one ring over. Written once, it applies to Kafka, Rabbit, NATS, and the in-process broker identically; that property is the entire messaging pitch, and it holds because no stage ever sees a driver type.
func ConcurrencyLimit ¶
func ConcurrencyLimit(n int) Middleware
ConcurrencyLimit caps concurrently executing handler invocations. It sits innermost — inside Retry — so the semaphore is held per attempt and a message sleeping in backoff starves nobody.
func DeadLetter ¶
func DeadLetter(pub Publisher, originTopic, dlqTopic string) Middleware
DeadLetter is the disposition stage: it maps the inner chain's final error by its outermost warren/errors code to the §2.6 consumer column. Terminal codes (INVALID, UNAUTHENTICATED, PERMISSION_DENIED, exhausted INTERNAL — and everything unknown, INTERNAL being the safe default) publish the original envelope to the dead-letter topic with forensic headers and ack. NOT_FOUND and CONFLICT ack. UNAVAILABLE nacks so the broker redelivers — never dead-lettered. A failed DLQ publish nacks: silent loss is the one forbidden outcome. Note what that nack implies: the chain holds no state, so the redelivery re-runs the FULL pipeline — the handler and its side effects included — before the DLQ publish is re-attempted. Handlers are idempotent in an at-least-once world; this is one more reason why.
func Deduplicate ¶
Deduplicate suppresses redeliveries of already-disposed messages, keyed on Message.ID: Seen before the handler (seen → ack without invoking it), MarkSeen only after the inner chain returns nil — a nacked message must not count as its own duplicate. A store error fails CLOSED: UNAVAILABLE nack, duplicates over loss, but never silently. A MarkSeen failure after success is acked — the work is done; refusing the ack would guarantee the duplicate it failed to prevent — but it is logged at ERROR, naming the subscription, because idempotency silently switching off is worse than the redelivery it was protecting against. A suppressed duplicate logs at DEBUG.
The mark does NOT join the handler's transaction, on any store: this stage runs outside the handler, so app.Transactional has committed by the time MarkSeen is called. warren.md §5.6 states the guarantee — at-least-once, for every store — and persistence/postgres.WithInbox says what durability does buy, which is reach, not atomicity.
func Drain ¶
func Drain() (Middleware, func(context.Context) error)
Drain returns the admission stage and the wait the lifecycle calls at shutdown step 3: wait refuses new deliveries (UNAVAILABLE nack — the broker redelivers them elsewhere) and returns when every in-flight message, its retries and its DLQ publish included, has finished.
func Recover ¶
func Recover() Middleware
Recover converts a panic into errors.Internal. Pipeline applies it twice: innermost — so a handler panic is classified INTERNAL and retried, then dead-lettered, per §2.6 — and outermost, so a bug in a stage itself cannot kill the subscription.
func RequireMessageID ¶ added in v0.2.0
func RequireMessageID() Middleware
RequireMessageID refuses a message whose ID cannot serve as a dedupe key, as INVALID — which is terminal, so the message is dead-lettered and preserved rather than dropped or redelivered for ever.
Two IDs are refused. An EMPTY one, because it is not a key: without this every message on the topic shares it, the first is handled, and every one after it is silently acked and destroyed. And one carrying a NUL byte, because the key must be storable by any inbox.Store — Postgres `text` rejects 0x00 outright, so such an ID would be a store error on a durable deployment and no error at all on the memory store, which is the worst pairing there is. Refusing it here makes the guarantee inboxtest states (keys hold no NUL) true by construction rather than by hope.
Pipeline installs it only when deduplication is on, because that is when the ID is load-bearing: it IS the key.
It sits INSIDE the dead-letter ring and outside Retry: retrying a message whose ID will never change achieves nothing.
func Retry ¶
func Retry(policy app.RetryPolicy) Middleware
Retry re-invokes the inner chain on the two §2.6 retry rows — UNAVAILABLE and INTERNAL, judged by the OUTERMOST code, a non-Warren error counting as INTERNAL — under the policy. Waits observe context cancellation, so shutdown never sits out a backoff. The final error is wrapped with the attempt count for the DLQ's forensic header; the code stays reachable.
func TraceExtract ¶
func TraceExtract() Middleware
TraceExtract continues the producer's trace from the delivery's headers, and seeds those headers on the context for anything else that wants them. It is the exact mirror of InjectTrace, and like it, it reaches the telemetry through app.TelemetryFromContext — so this package continues a distributed trace without knowing a telemetry SDK exists.
It used to do only the seeding half. broker.DeliveryHeaders had no non-test reader anywhere in the repository and app.Telemetry.Extract had exactly one caller, transport/http's edge — so trace context travelled INTO a message and never came back out. Every consumer span began a new root trace, silently, and a service that imported observability looked fully instrumented while answering none of the questions it was bought for.
It is a no-op when no telemetry is bound, so an uninstrumented service pays one nil check per delivery.
type Publisher ¶
Publisher sends messages to a topic. The outbox relay is its primary caller — it is variadic so a relay batch is one call; use cases publish through the unit of work, not directly.
func Correlating ¶ added in v0.2.0
Correlating returns a Publisher that copies the context's correlation ID into every message's headers, so the work a consumer does on the other side belongs to the request that caused it.
Without it the trail ends at the broker: a request logged under one ID published an event, and every line the consumer wrote while handling that event belonged to no request at all. Nothing in broker/ or outbox/ carried the ID, so the two halves of one causal chain could not be joined.
It never OVERWRITES a header that is already there. That rule is what makes the outbox work: outbox.Sink stamps the ID at Append, inside the request, and the relay publishes minutes later from a background context that has no correlation ID of its own. A decorator that overwrote would replace the one true value with nothing at exactly the moment it mattered.
It never stamps an EMPTY id either. A blank header reads as a correlation ID that is genuinely blank, and the consumer would seed one.
Correlating(nil) is nil, so wrapping whatever boot resolved is safe in a module that may have no broker configured.
type Redeliverer ¶ added in v0.2.0
type Redeliverer interface {
// Redelivers reports whether a nacked message comes back. It must return
// the same value for the lifetime of the driver.
Redelivers() bool
}
Redeliverer is implemented by a driver that can say whether nacking a message returns it for another attempt.
warren.md §2.6 gives UNAVAILABLE the consumer disposition "nack + backoff retry" and no dead-letter, which is right — but it rests on a PREMISE, and the premise is this method. A broker with a durable log and an acknowledgement protocol redelivers, so nacking is lossless and a DLQ would turn a transient blip into a queue of messages that would have succeeded. An in-process broker has neither, so a nack there is a DROP — and UNAVAILABLE, the code that means "try again", became the only lossy one while INVALID and INTERNAL were preserved.
A driver that does not implement this is assumed to redeliver, which is true of every durable broker and keeps the documented behaviour for them.
The answer must be constant for a driver ¶
It is asked ONCE, when the pipeline is composed, and it is asked of the publisher rather than of a subscription — so a driver's answer must hold for every subscription it serves. Kafka (offset rewind), RabbitMQ (nack with requeue) and JetStream (Nak) are all constantly true; the in-process broker is constantly false.
A driver that cannot promise that must expose TWO client values rather than one that answers differently per subscription — NATS Core and JetStream from a single connection is the case that forces this, and so is a RabbitMQ driver mixing requeue and no-requeue subscriptions. Moving this method onto Subscriber would make it per-subscription, and that is a breaking change to an exported interface, which is why the constraint is written down here rather than discovered later.
It also follows that the dead-letter publisher handed to Pipeline must be the SAME driver as the subscription. If it is a different broker, the pipeline asks the wrong one whether a nack comes back.
type SubscribeOption ¶
type SubscribeOption struct {
// contains filtered or unexported fields
}
SubscribeOption configures one subscription's pipeline. It is opaque: options are minted here, collected at the registration site, and applied by Pipeline at boot — drivers only forward them.
func WithConcurrency ¶
func WithConcurrency(n int) SubscribeOption
WithConcurrency caps the number of handler invocations executing concurrently. A message waiting out a retry backoff holds no slot. The default — omitting the option — is uncapped, bounded only by the driver's delivery parallelism; a cap of zero or less is a boot-time panic, because a subscription that can process no messages is a config typo, not intent.
func WithDeadLetter ¶
func WithDeadLetter(topic string) SubscribeOption
WithDeadLetter sets the topic exhausted and terminal messages are routed to. The default is "<topic>.dlq".
func WithDedupeTTL ¶
func WithDedupeTTL(d time.Duration) SubscribeOption
WithDedupeTTL sets how long a processed Message.ID is remembered. It must exceed the subscription's redelivery window. The default is 24 hours.
func WithRetry ¶
func WithRetry(p app.RetryPolicy) SubscribeOption
WithRetry sets the subscription's retry policy — the same app.RetryPolicy port the core Retrying middleware uses, so a policy written once is interchangeable with ExponentialBackoff at every site, inbound or consumer-side. One port, two rings, no library. The default is ExponentialBackoff(3).
func WithoutDedupe ¶
func WithoutDedupe() SubscribeOption
WithoutDedupe disables inbox dedupe for this subscription — the named opt-out from the on-by-default stance, for handlers that are naturally idempotent.
type Subscriber ¶
type Subscriber interface {
// Subscribe registers h for topic and RETURNS ONCE THE SUBSCRIPTION IS
// LIVE — once a Publish to that topic will reach h. It does not block
// for the lifetime of the subscription; the driver owns the delivery
// loop, which runs until ctx is cancelled.
//
// Returning early is the whole contract, and it exists because the
// alternative silently loses messages. Subscribe used to block, so every
// caller wrapped it in a goroutine:
//
// OnStart: func(context.Context) error {
// go func() { _ = sub.Subscribe(ctx, topic, h) }()
// return nil // "started" — but not yet subscribed
// }
//
// which reports success before the subscription exists. warren.md §1.3
// step 6 promises consumers start before publishers, and that promise was
// being defeated by the very code the framework generates: anything
// published in the gap went to a topic with no subscriber and was
// discarded — and the outbox relay, seeing a successful Publish, marked
// the record published. Silent, permanent event loss on the path the
// framework recommends. In tests it showed up as the first module test in
// a process passing and every later one receiving nothing at all.
//
// A driver that cannot register synchronously must block until it can:
// failing the boot because a broker is unreachable is the correct
// outcome, and it is the reason this returns an error.
//
// THE CONTEXT HANDED TO h FOR EVERY DELIVERY MUST CARRY ctx's VALUES.
// Cancellation is the driver's own — a shared poll loop may not stop
// fetching for one subscription because another was cancelled — but the
// values are the application's, and a driver that manufactures a delivery
// context from context.Background() severs them.
//
// This is not a formality. app.Telemetry rides the context, and
// TraceExtract and InjectTrace both reach it through
// app.TelemetryFromContext: a driver that drops the values makes every
// consumer span a new root and strips traceparent from every event the
// consumer raises, with nothing failing anywhere. The two shipped drivers
// disagreed about this until it was written down — memory passed the
// context through, kafka rebuilt it from Background — so brokertest now
// checks it, and a fix seeded through the context would otherwise have
// passed on the in-process broker and done nothing in production.
Subscribe(ctx context.Context, topic string, h MessageHandler) error
}
Subscriber consumes a topic, invoking the handler for each message. A subscription is a lifecycle component: it starts after its dependencies are ready and drains on shutdown — cancellation of ctx stops fetching and lets in-flight messages finish and ack. Never a goroutine someone forgot about.
Directories
¶
| Path | Synopsis |
|---|---|
|
Package brokertest is the contract suite every broker driver must pass: the in-process one, Kafka, RabbitMQ, NATS.
|
Package brokertest is the contract suite every broker driver must pass: the in-process one, Kafka, RabbitMQ, NATS. |
|
kafka
module
|
|
|
Package memory is the in-process broker: the default in tests, and the driver a modular monolith runs in production before its modules are extracted into services.
|
Package memory is the in-process broker: the default in tests, and the driver a modular monolith runs in production before its modules are extracted into services. |