messenger

package module
v0.3.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 3, 2026 License: MIT Imports: 25 Imported by: 7

README

GoMessenger

CI Release Go Reference Go 1.27

Typed durable messaging for Go.

Build commands, local queries, and events with explicit transaction, delivery, retry, and broker semantics.

Transactional Outbox/Inbox, NATS JetStream, Kafka, bounded retries, DLQ/replay, tracing, and managed lifecycle.

Release status: v0.3.0 is the current release line. Source validation, dependency-ordered tag publication, and the clean published-consumer gate are separate release evidence. The real-service pilot remains pending, so controlled repository gates are not a production-readiness claim. The current checkout may accumulate follow-up changes beyond that release line.

What it is

GoMessenger is a typed facade for Go 1.27+ services that need both process-local messaging and durable one-way delivery. Commands, local request/reply queries, and events share explicit descriptors, middleware, lifecycle, and observations; one-way commands and events can additionally use transactional Outbox staging, NATS JetStream or Kafka delivery, and a durable SQL Inbox.

It keeps transport truth visible. An Outbox receipt means staged in the caller's transaction, a NATS receipt means JetStream returned PubAck, a Kafka receipt means the producer transaction committed, and a consumer ACK/offset follows the committed Inbox transaction. The root module remains standard library plus gobus; broker, SQL, Prometheus, and OpenTelemetry dependencies live in optional nested modules.

Why it exists

Use GoMessenger when a service needs at least one of these boundaries:

  • a business write and outgoing command/event must commit or roll back together;
  • broker redelivery must not repeat an already committed SQL handler effect;
  • retry, terminal failure, DLQ, and replay must be bounded and explicit;
  • local commands, local queries, and durable events should use stable typed descriptors without pretending that NATS and Kafka have identical semantics.

If all work is process-local, GoBus or direct function calls may be enough. If the problem is a durable multi-step workflow with timers and compensation, use a workflow engine. See the use-case comparison before choosing an abstraction.

Install

GoMessenger requires Go 1.27+. For local commands, queries, and events:

go get github.com/assurrussa/gomessenger@v0.3.0

For durable NATS JetStream delivery with Inbox and transactional Outbox integration:

go get github.com/assurrussa/gomessenger@v0.3.0 \
  github.com/assurrussa/gomessenger/adapters/inbox@v0.3.0 \
  github.com/assurrussa/gomessenger/adapters/nats@v0.3.0 \
  github.com/assurrussa/gomessenger/adapters/outbox@v0.3.0

For durable Kafka delivery with Inbox and transactional Outbox integration:

go get github.com/assurrussa/gomessenger@v0.3.0 \
  github.com/assurrussa/gomessenger/adapters/inbox@v0.3.0 \
  github.com/assurrussa/gomessenger/adapters/kafka@v0.3.0 \
  github.com/assurrussa/gomessenger/adapters/outbox@v0.3.0

Optional telemetry and CLI modules use the same release version:

go get github.com/assurrussa/gomessenger/observability@v0.3.0
go install github.com/assurrussa/gomessenger/tools/gomessengerctl@v0.3.0

These commands target the exact path-qualified v0.3.0 tags. The rest of this README tracks the current checkout and may describe unreleased APIs that are not present in that release line. Use the versioned Go Reference for the exact release API, or use the checkout workflow below when evaluating unreleased changes.

Keep every GoMessenger module in one consumer on the same version. The Outbox adapter requires Outbox v0.15.0; the host selects and installs its matching database backend separately. To evaluate the current checkout instead:

git clone https://github.com/assurrussa/gomessenger.git
cd gomessenger
GOWORK=off go test ./...

30-line local quickstart

The smallest local command is:

type ResizeMedia struct {
	JobID int64 `json:"jobId"`
}

resize := messenger.MustCommand("media.resize", 1, messenger.JSON[ResizeMedia]())
builder := messenger.NewBuilder(messenger.WithSource("urn:service:media-resizer"))
builder.HandleCommandFunc(resize, "media-worker", func(_ context.Context, payload ResizeMedia) error {
	fmt.Println("resize", payload.JobID)
	return nil
})
builder.RouteCommand(resize, messenger.NewLocalSyncRoute())

bus, _, err := builder.Build()
if err != nil {
	return err
}
resizeSender := messenger.BindSender(bus, resize)
_, err = resizeSender.Send(ctx, ResizeMedia{JobID: 42})

The equivalent runnable API is compiled as ExampleMessenger_Send; testdata/consumer separately compiles the complete public facade from an external Go module.

For the durable path, run the PostgreSQL + NATS demo:

make demo-durable-postgres-nats

It performs a business write and Outbox stage in one PostgreSQL transaction, relays through JetStream, retries an intentional handler failure, suppresses a distinct duplicate delivery in the Inbox, moves a permanent failure to the DLQ, and confirms replay. The example is a checkout-level demonstration with local GoMessenger replacements; it is not evidence of published-module resolution or production readiness.

The example also contains an opt-in open-loop NATS capacity experiment over HTTP -> business transaction + Outbox -> JetStream -> Inbox -> business projection:

make capacity-nats
make capacity-nats-site
make capacity-nats-site-single
make capacity-nats-site-batch-1
make capacity-nats-site-batch-100
make capacity-frontier
make capacity-frontier-matrix
make capacity-inbox-postgres

The legacy quick and site-shaped commands remain screening tools. The normalized frontier commands use PostgreSQL 18 and compare three fixed two-CPU topologies, four runtime-confirmed ingress/relay/consumer variants, small and mixed payloads, and stock versus separately labelled tuned database profiles. The PostgreSQL-only command isolates the real Inbox ProcessAttempt transaction without Outbox or NATS. They report unique committed business effects, canonical envelope bytes, Inbox/ACK latency, and PostgreSQL statement/WAL/I/O telemetry. See the example capacity contract. Results describe only the recorded checkout, host, and local Docker topology; they are not production benchmark claims.

Set OUTBOX_RESERVATION_BATCH_SIZE=16 (valid 1..1000) to A/B only the reservation/prefetch width. The default remains 1. Each Outbox worker still publishes and acknowledges those prefetched jobs sequentially. True relay batching is selected separately with OUTBOX_RELAY_MODE=batch and OUTBOX_BATCH_MAX_MESSAGES=1|100; true ingress batching requires OUTBOX_INGRESS_MODE=batch and the bulk HTTP path. Select the consumer with CONSUMER_MODE=single|batch. Batch mode accepts CONSUMER_BATCH_MAX_MESSAGES, CONSUMER_BATCH_MAX_BYTES, and CONSUMER_BATCH_MAX_WAIT; use batch + 1 as the same-path control and batch + 100 for real batching. Capacity report spec 2.1 records the runtime-confirmed ingress, relay, and consumer selections, actual batch calls and sizes, Outbox handler/publish/finalization latency and outcomes, image digests, limits, pool connection health, SQL calls, transactions/message, WAL/message, checkpoints, pprof, and PostgreSQL plans. It reports relay throughput from published envelopes, consumer throughput and MiB/s from committed projections, and separates staged - published Outbox lag from published - committed consumer lag. Warm-up and drain remain outside every throughput denominator.

Versioned capacity baselines, raw-evidence rules, and defensible claim boundaries are recorded in the performance evidence registry.

Transactional producer and relay batches

Producer batching is explicit and Outbox-only. NewBatchProducer stages one ordered set atomically in the caller's business transaction; direct broker and local routes return ErrUnsupportedCapability from the batch facade. Existing single-message producers and relay jobs do not change.

producer, err := outboxadapter.NewBatchProducer(outboxService, outboxadapter.ProducerConfig{
	Name: "orders-outbox",
})
if err != nil {
	return err
}

builder := messenger.NewBuilder(messenger.WithSource("urn:service:orders"))
builder.RouteEvent(orderCreated, producer)
app, _, err := builder.Build()
if err != nil {
	return err
}

receipts, err := app.PublishMessageBatch(ctx, orderCreated, []messenger.Outgoing[OrderCreated]{
	{Payload: first, Metadata: messenger.OutgoingMetadata{ID: firstID}},
	{Payload: second, Metadata: messenger.OutgoingMetadata{ID: secondID}},
})

Register outboxadapter.NewBatchRelayJob with Outbox RegisterBatchJob. Its zero outbox.BatchConfig means 100 jobs, 4 MiB of payload, and 25 ms; MaxMessages=1 is the same-path control. NATS confirms one asynchronous publish future per item. Kafka sends the valid subset in one multi-record transaction. Result and retry semantics are defined by ADR-0006.

Guarantees

The durable contract is at-least-once:

  • an Outbox route reports success only after the envelope is staged in the caller's database transaction;
  • a direct NATS route reports success only after a JetStream PubAck;
  • a direct Kafka route reports success only after its producer transaction commits;
  • a durable consumer acknowledges only after its Inbox transaction and handler commit;
  • stable message identity and an Inbox suppress a second execution of the same committed consumer transaction;
  • a batch consumer invokes one typed handler and one Inbox/business transaction for each real batch while retaining individual ACK/retry/defer/DLQ outcomes;
  • concurrency, envelopes, headers, handler attempts, retry delays, and shutdown waits are bounded by explicit configuration or documented limits.

Non-goals

  • No exactly-once external effects: HTTP calls, email, object storage, and writes to another database still need their own idempotency or durable hand-off.
  • No distributed queries: Query[Q,R] is process-local. Remote request/reply remains the separate, unimplemented contract in ADR-0003.
  • No workflow, saga, event-sourcing, service-discovery, or streaming-analytics engine.
  • No automatic ownership of host database connections, broker credentials, migrations, topology policy, supervision, or deployment.
  • No claim of universal broker abstraction or production-proven maturity; the real-service pilot in ADR-0002 is still pending.

Choose a route

Start with an explicit descriptor, choose the success boundary, build the messenger, and inject a narrow Sender, Querier, or Publisher into business code. Durable consumers are separate managed services.

Route A successful call means Important boundary
NewLocalSyncRoute the in-process handler completed no restart durability
NewLocalAsyncRoute bounded runtime work completed for queries or was accepted for one-way work no restart durability
natsadapter.NewRoute JetStream returned PubAck not atomic with a separate business database write
kafkaadapter.NewRoute the producer transaction committed not atomic with a separate business database write
outboxadapter.NewProducer the envelope was staged in the active host transaction provisional until that transaction commits
business code -> Messenger.Query -> local result handler -> typed result
             `-> Send/Publish -> local handler
                              `-> Outbox -> relay -> JetStream/Kafka -> consumer -> Inbox transaction -> handler -> ACK/offset

For a durable flow, provision topology, apply Outbox and Inbox migrations, register the relay before producer traffic, start the managed runtimes, and expose readiness before accepting traffic. The practical usage guide shows the full composition and shutdown order.

Selective routing (Mixing Outbox and direct brokers)

Outbox is not mandatory for all messages. You can mix and match different routes in the same application based on descriptor requirements:

  • Transactional Outbox: Route critical state transitions (e.g., payments, orders, status changes) through outboxadapter.NewBatchProducer so messages commit atomically with your database writes.
  • Direct broker dispatch: Route high-throughput, low-latency, or non-transactional traffic (e.g., clickstream, telemetry, audit logs, ephemeral notifications) directly through natsadapter.NewRoute or kafkaadapter.NewRoute. This completely bypasses the database, avoids WAL/disk contention, and delivers at full broker throughput.
builder := messenger.NewBuilder(messenger.WithSource("urn:service:billing"))

// 1. Critical event: staged atomically in PostgreSQL transaction (ReceiptStaged)
builder.RouteEvent(OrderCreatedEvent, outboxProducer)

// 2. High-volume audit/telemetry: published directly to JetStream without DB (ReceiptBrokerConfirmed)
builder.RouteEvent(TelemetryMetricEvent, natsDirectRoute)

// 3. Local in-process command: executed in-memory via GoBus (ReceiptCompleted)
builder.RouteCommand(InvalidateLocalCache, messenger.NewLocalSyncRoute())
Custom routes and custom adapters

GoMessenger does not couple your application to the bundled database or broker implementations. The route boundary is a minimal Go interface:

type Route interface {
	Name() string
	Deliver(ctx context.Context, delivery Delivery) (Receipt, error)
}

For batched ingress, routes may additionally implement BatchRoute:

type BatchRoute interface {
	Name() string
	DeliverBatch(ctx context.Context, batch []Delivery) ([]Receipt, error)
}

Because Route and BatchRoute are standard interfaces, host applications can implement custom adapters and overriding strategies directly in their codebase without requiring any changes or new releases of GoMessenger:

  • High-throughput Append-only / CDC Outbox: The standard adapters/outbox provides transactional staging with a polling publisher relay and batching (~2,500–3,000 msg/s sustainable on single-instance PostgreSQL). For high-throughput platforms (e.g. 10,000–100,000+ msg/s) where polling queries and status updates (UPDATE ... status = 'sent') cause database contention or table bloat, applications can write minimal unindexed append-only rows (or stream via pgx.CopyFrom) within the business transaction, and let an external CDC pipeline (such as Debezium, Kafka Connect, or a PostgreSQL WAL logical replication streamer via pgoutput) forward raw events directly to Kafka or NATS.
  • Alternative storage backends: route outbox events into Redis Streams, ClickHouse, KeyDB, or Tarantool.
  • Dynamic / Conditional routing: inspect context.Context at runtime to route through Outbox when an active *sql.Tx is present, or fall back to direct broker publishing outside transactions.
  • Mock and testing routes: capture or assert dispatched messages in integration tests with zero infrastructure dependencies.

Modules and release status

GoMessenger requires Go 1.27 because the builder and messenger expose generic methods.

Module Responsibility
github.com/assurrussa/gomessenger descriptors, local queries, envelopes, local routes, runtime, manifest
.../adapters/outbox transactional staging and broker relay job
.../adapters/nats JetStream producers, consumers, topology, CloudEvents
.../adapters/kafka transactional Kafka producers, consumers, topology, retry/DLQ
.../adapters/inbox atomic PostgreSQL and SQLite consumer deduplication
.../observability Prometheus, OpenTelemetry spans, W3C Trace Context
.../tools/gomessengerctl manifest/topology validation, plan/apply, DLQ inspect/replay

The module set uses synchronized path-qualified v0.3.0 tags. Release completion requires every tag above plus the clean post-publication consumer probe; neither is inferred from source-only checks. Outbox root and its PostgreSQL/SQLite backend tags at v0.15.0 are the pinned durable-producer dependencies. During repository development go.work selects local GoMessenger modules; published consumers use matching path-qualified tags and no local replace directives. See the release process for dependency order and verification.

Minimal local query example

type FindArticle struct{ ID int64 }
type ArticleView struct {
	ID    int64
	Title string
}

findArticle := messenger.MustQuery[FindArticle, ArticleView](
	"article.find", 1, messenger.JSON[FindArticle](),
)
builder := messenger.NewBuilder(messenger.WithSource("urn:service:catalog"))
builder.HandleQueryFunc(findArticle, "article-reader", func(_ context.Context, query FindArticle) (ArticleView, error) {
	return ArticleView{ID: query.ID, Title: "CQRS in Go"}, nil
})
builder.RouteQuery(findArticle, messenger.NewLocalSyncRoute())

bus, _, err := builder.Build()
if err != nil {
	return err
}
reader := messenger.BindQuerier(bus, findArticle)
article, err := reader.Query(ctx, FindArticle{ID: 42})

Query[Q,R] uses the codec and schema only for request identity; R is an in-process type identity and is not written to the manifest. Every registered query must have exactly one handler and one built-in local route. The async route uses the caller context for admission, execution, and waiting, so cancellation returns ctx.Err() and accepted result delivery cannot block runtime drain. The runnable version is compiled as ExampleMessenger_Query.

Descriptors use explicit stable wire names, schema versions, content types, data encodings, and optional schema URIs. Go type names and package paths never become the wire contract implicitly. Native envelopes carry dataEncoding as json, text, or binary, so custom codecs do not infer their representation from a media-type prefix.

Local method-call benchmarks

Representative GitHub-hosted benchmark results for the synchronous local public path (Linux/amd64, Intel Xeon Platinum 8573C, Go 1.27, ten samples):

Public method Scenario Median time/op Approx. calls/s B/op allocs/op
Messenger.Send one local command handler 1.066 µs ~938k 1472 9
Messenger.Query one local query handler and typed result 937.6 ns ~1.07M 1296 10
Messenger.Publish one local event subscriber 1.086 µs ~921k 1472 9

The calls/s column is the reciprocal of the median single-thread time/op, not a concurrency or durable-throughput claim. These benchmarks exercise NewLocalSyncRoute with no-op application handlers. They do not include NATS, Kafka, Outbox, Inbox, SQL, network latency, retries, or telemetry exporters. Reproduce the sample with:

GOWORK=off go test -run '^$' -bench '^BenchmarkMessengerLocal' -benchmem -count=10 .

GitHub Actions retains the raw samples and benchstat comparison. Command and query results in this snapshot did not change significantly from the base commit.

Transactional outbox producer

Use adapters/outbox when a business write and message publication must commit or roll back together. Register the broker relay before the producer can stage its job:

natsRoute, err := natsadapter.NewRoute(natsConnection, natsadapter.RouteConfig{
	Name:      "nats.integration-events",
	Namespace: "prod",
	WireMode:  natsadapter.WireNative,
})
if err != nil {
	return err
}
relay, err := outboxadapter.NewRelayJob(natsRoute, outboxadapter.RelayJobConfig{})
if err != nil {
	return err
}
if err := outboxRuntime.Service().RegisterJob(relay); err != nil {
	return err
}

route, err := outboxadapter.NewProducer(
	outboxRuntime.Service(),
	outboxadapter.ProducerConfig{Name: "outbox.integration-events"},
)
if err != nil {
	return err
}
builder := messenger.NewBuilder(
	messenger.WithSource("urn:service:media-resizer"),
)
builder.RouteEvent(mediaResized, route)

bus, _, err := builder.Build()
if err != nil {
	return err
}

err = outboxRuntime.Transactor().RunInTx(ctx, func(txCtx context.Context) error {
	if err := mediaRepository.Save(txCtx, media); err != nil {
		return err
	}
	_, err := bus.Publish(txCtx, mediaResized, MediaResized{JobID: media.ID})
	return err
})

The repository write must use the transaction carried by txCtx. ReceiptStaged reports staging in that active transaction; a callback rollback still removes the row. The producer stores the canonical envelope under its MessageID; the relay later publishes those exact bytes and waits for PubAck. Repeating identical content resolves the same outbox tombstone, while reusing the identity with different content fails closed. Immediate messages use their immutable message time as availableAt, preventing retry-time fingerprint drift.

messenger.RetryAfter from the broker route becomes a persisted outbox.RetryAt; permanent envelope failures move directly to the outbox DLQ.

Durable JetStream consumer and inbox

A consumer owns one stable ConsumerID, explicit concurrency and retry bounds, and one SQL inbox:

if err := inboxpgsql.Migrate(ctx, database); err != nil {
	return err
}
store, err := inboxpgsql.New(database)
if err != nil {
	return err
}
consumer, err := natsadapter.NewEventConsumer(
	natsConnection,
	store,
	mediaResized,
	func(ctx context.Context, message messenger.Message[MediaResized]) error {
		tx, ok := inbox.SQLTxFromContext(ctx)
		if !ok {
			return errors.New("missing inbox transaction")
		}
		return projection.ApplyTx(ctx, tx, message.Metadata.ID.String(), message.Payload)
	},
	natsadapter.HandlerConfig{
		Stream: "MESSAGES", Namespace: "prod", ConsumerID: "media-projection", WireMode: natsadapter.WireNative,
		Concurrency: 8, Timeout: 30 * time.Second, FinalizationTimeout: 5 * time.Second,
		AckWait: 30 * time.Second, MaxAttempts: 10,
		DLQSubject: "prod.dlq",
	},
)
if err != nil {
	return err
}
consumerBuilder := messenger.NewBuilder(
	messenger.WithSource("urn:service:media-projection"),
)
consumerBuilder.Use("consumer.media-projection", consumer)

_, runtime, err := consumerBuilder.Build()
if err != nil {
	return err
}

For true consumer batching, select the separate constructor and return one keyed result for every message passed to the handler. Classify the complete batch before writing the successful subset:

batchConsumer, err := natsadapter.NewBatchEventConsumer(
	natsConnection,
	store,
	mediaResized,
	func(ctx context.Context, messages []messenger.Message[MediaResized]) (messenger.BatchResult, error) {
		result := messenger.BatchResult{Items: make([]messenger.BatchItemResult, len(messages))}
		successes := make([]messenger.Message[MediaResized], 0, len(messages))
		for i, message := range messages {
			itemErr := validateMedia(message.Payload)
			result.Items[i] = messenger.BatchItemResult{
				Key: messenger.BatchItemKey{Source: message.Metadata.Source, MessageID: message.Metadata.ID},
				Err: itemErr,
			}
			if itemErr == nil {
				successes = append(successes, message)
			}
		}
		tx, ok := inbox.SQLTxFromContext(ctx)
		if !ok {
			return messenger.BatchResult{}, errors.New("missing inbox transaction")
		}
		if err := projection.ApplyBatchTx(ctx, tx, successes); err != nil {
			return messenger.BatchResult{}, err
		}
		return result, nil
	},
	natsadapter.HandlerConfig{
		Stream: "MESSAGES", Namespace: "prod", ConsumerID: "media-batch-projection",
		WireMode: natsadapter.WireNative, Concurrency: 4, Timeout: 30 * time.Second,
		AckWait: 30 * time.Second, MaxAttempts: 10, DLQSubject: "prod.dlq",
	},
	messenger.BatchConfig{}, // 100 messages, 4 MiB canonical bytes, 25 ms
)

MaxMessages=1 is a control for the same batch path. DeferAfter requests an exact delayed retry without consuming an attempt; RetryAfter consumes one. An ordinary top-level error rolls back the whole handler transaction and retries the batch without changing item attempts. A top-level Permanent error or an inexact result fails the consumer closed. Kafka exposes matching NewBatchCommandConsumer and NewBatchEventConsumer constructors; each Kafka batch is confined to one contiguous topic-partition range.

Run the embedded additive inbox migrations explicitly with inboxpgsql.Migrate or inboxsqlite.Migrate before constructing consumers. They include the durable handler-attempt count and permanent-outcome state required by durable consumers. Hosts own database connections, migration ordering, NATS connections, credentials, and process supervision.

Each active consumer worker holds one SQL transaction for the full handler invocation. This is the atomic boundary for business writes and the Inbox completion marker; do not move the handler outside it. The adapter adds no connection semaphore: database/sql is the host-managed backpressure boundary. When several consumers share one sql.DB, configure SetMaxOpenConns to at least the sum of their HandlerConfig.Concurrency values, plus explicit headroom for application and maintenance queries. A smaller pool becomes the bottleneck before Go dispatch and can exhaust handler deadlines before application code runs.

Timeout bounds application handler execution. The Inbox transaction receives an additional FinalizationTimeout (5 seconds by default) to commit or roll back after that deadline. Increase it when a remote or otherwise slow database needs more finalization time; it does not extend the handler deadline.

Both SQL backends default to the gomessenger_ prefix and therefore preserve gomessenger_inbox, gomessenger_inbox_attempts, and gomessenger_inbox_attempt_generations. Pass the same options to migration and runtime construction when a host needs another namespace:

postgresInbox := []inboxpgsql.Option{
	inboxpgsql.WithSchema("messaging"), // the host creates the schema and grants access
	inboxpgsql.WithTablePrefix("site_"),
}
if err := inboxpgsql.Migrate(ctx, database, postgresInbox...); err != nil {
	return err
}
store, err := inboxpgsql.New(database, postgresInbox...)

WithTablePrefix is also available for SQLite; SQLite has no schema option. PostgreSQL qualifies every relation instead of relying on search_path, and its migrator never creates the configured schema. Changing the prefix or schema selects a separate Inbox: migrations do not rename, copy, or otherwise transfer existing deduplication history.

Provision source and DLQ capacity separately. DevStream keeps native source messages at the 1 MiB envelope bound; DevDLQStream reserves DefaultMaxDLQMessageBytes for the expanded JSON DLQ record. The NATS server/account max_payload must be at least the same value. Consumer.Run rejects a missing or undersized DLQ route before starting; the low-frequency Consumer.DeepHealth probe detects later topology drift without making ordinary readiness expensive.

Middleware, logging, and tracing

Global middleware wraps local query/command/event handlers and can be reused by durable NATS consumers. Registration order is execution order: the first middleware is outermost. A middleware may replace the context or short-circuit, but next may be called at most once. Typed one-way decorators use messenger.ChainHandler; typed query decorators use messenger.ChainQueryHandler and may return a cached or synthetic R. Global middleware cannot synthesize a typed result: successful completion without one returns ErrQueryResultMissing.

logger := messenger.AdaptSlog(slog.Default())
builder := messenger.NewBuilder(
	messenger.WithSource("urn:service:media-resizer"),
	messenger.WithLogger(logger),
	messenger.WithObserver(messenger.NewLoggingObserver(logger)),
	messenger.WithContextPropagator(observability.NewTraceContextPropagator()),
)
builder.UseMiddleware(func(
	ctx context.Context,
	metadata messenger.Metadata,
	handlerID string,
	next messenger.HandlerFunc,
) error {
	return next(ctx)
})

Logging is disabled by default. AdaptSlog(nil) is a safe no-op adapter; direct WithLogger(nil) is a configuration error. Observer registrations are additive. A panic in one observer is logged and isolated from the remaining observers. Kafka TransportConfig.Logger reports adapter-owned startup/readiness, producer and consumer transaction, abort/fencing, topology failures or applied changes, and retry partition deferrals. Core logging contains infrastructure state only and never logs record keys, payloads, message bodies, or arbitrary headers. WithClientLogger is a separate explicit opt-in to franz-go's own client logs.

Recovered handler panic values and stacks are dropped by default; ordinary errors satisfy the transport-neutral HandlerPanicError interface and can be classified with errors.As across independently versioned adapters. Configure WithPanicReporter (or the corresponding durable consumer field) only for a trusted diagnostic sink. Consumer observations and operational logs use the conservative DefaultFailureSanitizer unless the host explicitly supplies another sanitizer. In batch consumers, durable DLQ wire text strictly retains the built-in conservative sanitizer to ensure rebalance and finalization boundaries remain bounded, while configured host sanitizers apply to observations and logs. Custom sanitizers must execute promptly without blocking.

The observability propagator carries only W3C traceparent and tracestate. It works through native envelopes, CloudEvents structured/binary modes, and transactional Outbox storage. Baggage is intentionally not supported yet.

Runtime lifecycle

Run the managed services under the host supervisor. On a signal, close admission first and then wait for accepted work within a deadline:

signalCtx, stopSignals := signal.NotifyContext(
	context.Background(),
	os.Interrupt,
	syscall.SIGTERM,
)
defer stopSignals()

runErr := make(chan error, 1)
go func() {
	runErr <- runtime.Run(context.Background())
}()

select {
case err := <-runErr:
	return err
case <-signalCtx.Done():
}

runtime.BeginDrain()
shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
if err := runtime.Shutdown(shutdownCtx); err != nil {
	return err
}
return <-runErr

Runtime.Readiness is the lightweight admission/connectivity probe for the runtime and every attached service. Runtime.DeepHealth performs explicit topology diagnostics; Runtime.Liveness does not require readiness. The selected Outbox backend runtime has its own host-supervised lifecycle; GoMessenger does not open connections or supervise it implicitly. Never call Shutdown from a handler executing on the same runtime.

Failure semantics

  • Return messenger.Permanent(err) for a non-retryable handler failure.
  • Return messenger.RetryAfter(err, delay) for an explicit durable delay.
  • Other errors use bounded full-jitter exponential retry in the NATS and Kafka adapters.
  • MaxAttempts bounds handler invocations, not broker deliveries. A NotBefore deferral does not invoke the handler or consume an attempt.
  • AckWait is a broker redelivery deadline, not a handler timeout. It must be at least 100 ms and may be shorter than Timeout; active handlers refresh acknowledgement progress every AckWait / 3.
  • The NATS adapter publishes and confirms a DLQ record before broker-confirming acknowledgement of the original message. Failed terminal hand-off operations retry on the current delivery without invoking the application handler beyond MaxAttempts. A permanent outcome is persisted independently of the attempt count, so an interrupted hand-off cannot invoke that handler again after restart.
  • NotBefore becomes a retry delay until due; ExpiresAt becomes a permanent expired outcome.

Kafka

The independent Kafka module provides native-envelope transactional publish, Outbox relay, read-committed consumers, consumer-specific retry topics, atomic retry/DLQ offset hand-off, protected replay, static worker identity, and non-destructive topic planning. It intentionally does not reuse the NATS engine or support CloudEvents/query transport. Retry may be overtaken by later source records, so ordering is guaranteed only before the first failure. See the Kafka adapter guide and ADR-0004.

Level 2 roadmap

The next maturity level is an operational messaging platform, not a workflow framework. Work is ordered evidence-first: pilot and performance baseline, schema compatibility, broker capability declarations, partition-key and ordering semantics, batch capacity validation, then load, soak, and chaos validation. Saga engines, workflow orchestration, and generic distributed request/reply remain out of scope.

See the Level 2 roadmap for workstreams and exit criteria.

CloudEvents

The NATS adapter supports the native envelope for commands and events, plus CloudEvents 1.0 structured and binary modes for events. Native metadata is mapped explicitly; delivery-attempt state never enters the canonical envelope or fingerprint. An omitted CloudEvents time is derived deterministically from a UUIDv7 event ID; events without time and with a non-UUIDv7 ID are rejected so retries cannot change the canonical fingerprint. The required dataencoding extension carries json, text, or binary; missing, invalid, or descriptor-conflicting values are rejected before the payload codec runs.

Topology and CLI

Topology management is declarative and non-destructive. It may create missing resources or apply compatible additive changes. Retention/storage/replica changes, subject removal, reduced limits, or consumer delivery-contract drift are conflicts; the tool never deletes and recreates resources automatically. Compatible updates preserve broker settings outside the manifest's managed subset.

gomessengerctl manifest validate --file manifest.json
gomessengerctl topology validate --file topology.json
gomessengerctl topology plan --file topology.json --server nats://localhost:4222
gomessengerctl topology apply --file topology.json --server nats://localhost:4222
gomessengerctl dlq inspect --file record.json
gomessengerctl dlq replay --file record.json
gomessengerctl dlq replay --file record.json --confirm --server nats://localhost:4222
gomessengerctl kafka topology validate --file kafka-topology.json
gomessengerctl kafka topology plan --file kafka-topology.json --brokers localhost:9092 --instance-id ops-a
gomessengerctl kafka topology apply --file kafka-topology.json --brokers localhost:9092 --instance-id ops-a
gomessengerctl kafka dlq inspect --file kafka-record.json
gomessengerctl kafka dlq replay --file kafka-record.json
gomessengerctl kafka dlq replay --file kafka-record.json --confirm --brokers localhost:9092 --instance-id ops-a

topology plan exits with code 3 when the printed plan contains a conflict. dlq inspect prints a safe summary and replayability status without handler error text, wire bytes, or header values. dlq replay is offline by default and prints a payload-free deterministic JSON plan. --confirm republishes the original subject, wire bytes, and bounded headers with a deterministic, DLQ-record-specific Nats-Msg-Id, then waits for JetStream PubAck. The target consumer starts one fresh bounded attempt generation while the original terminal generation remains protected; redeliveries of that replay retain the same generation and do not reset MaxAttempts. Replay does not support subject substitution, record deletion, or payload/header output. Its internal replay headers are reserved transport metadata and must not be injected by ordinary publishers; enforce that boundary with NATS publish permissions when producers can access subjects directly.

NATS DLQ inspection accepts normal v1 records and quarantine v2 records from the same subject. Quarantine reports its reason, original sizes, digest, and whether complete source content was omitted. It is explicitly non-replayable. PostgreSQL/SQLite terminal generations remain protected after handoff; hosts may opt into bounded retention with Store.PruneTerminalAttempts after choosing a safe cutoff. There is no automatic cleanup or default TTL. Upgrade the CLI first, apply additive migrations, drain old consumers, then start new consumers. See delivery guarantees and rollout.

Development and verification

make prepare
make check-workspace
make test-e2e
make test-integration
make test-batch-integration
make test-kafka
GOMESSENGER_POSTGRES_DSN='postgres://...' make test-postgres
make bench-all

make check-workspace covers every module, static lint, race and checkptr builds, a 90% root coverage gate, the checkout-consumer module, and the Docker-free durable pipeline E2E against the local go.work graph, including the sibling Outbox checkout. It is the maximal pre-publication source gate. make check runs the same gate with GOWORK=off; it intentionally remains blocked while a required GoMessenger or Outbox contract has not been published and pinned. release-readiness plus make test-consumer-release VERSION=vX.Y.Z is the separate no-replacement publication proof. make test-e2e reruns only the full Outbox-to-JetStream-to-Inbox path (using embedded JetStream for Docker-free local runs). make test-kafka is the local Docker entry point against official Apache Kafka 4.1.2 and 4.3.1 images in KRaft mode (ZooKeeper-less); hosted CI runs each version in an independent matrix job. Both NATS and Kafka adapters share complete functional E2E test parity under the race detector (Outbox staging, transactional relay, Inbox deduplication, retry tiers, DLQ, and replay). Fixed-rate capacity benchmarking (1,500 msg/s floor) is currently published for PostgreSQL + NATS. Run the full source gate before publishing the reviewed root tag. Then promote exact module requirements in reviewed dependency layers: root-dependent modules, Inbox-dependent transports, and finally the CLI and checkout fixtures. After the root, Inbox, NATS, and Kafka tags resolve through the Go proxy, finalize and check the complete graph:

make release-ready VERSION=vX.Y.Z OUTBOX_VERSION=v0.15.0
make release-readiness VERSION=vX.Y.Z OUTBOX_VERSION=v0.15.0
make check

release-ready is the final-layer command; do not run it before the dependency tags exist. The complete staged tagging and per-layer verification procedure is in the release guide.

A published release is verified separately, after all dependency-ordered module tags exist:

make test-consumer-release VERSION=vX.Y.Z

GitHub Actions runs static, race, checkptr, PostgreSQL, and Kafka integration shards; the aggregate Full gate requires all of them. A separate workflow compares base and head benchmarks with pinned benchstat, uploads raw samples, and reports when the base predates the Go module instead of failing the first comparison. It does not enforce an unstable cross-machine performance threshold.

See the practical usage guide, Kafka guide, contracts, architecture, durable pipeline E2E, use-case comparison, Level 2 roadmap, adoption and migration guide, release order, and the first real-project pilot decision. Distributed request/reply remains a separate, unimplemented boundary in ADR-0003.

Documentation

Overview

Package messenger provides typed commands, local queries, and events with explicit descriptors, GoBus dispatch, and transport-neutral one-way delivery contracts. Queries are process-local request/reply calls and never use the wire envelope or durable routes.

Index

Examples

Constants

View Source
const (
	// DefaultBatchMaxMessages is the zero-value BatchConfig message limit.
	DefaultBatchMaxMessages = 100
	// DefaultBatchMaxBytes is the zero-value BatchConfig canonical byte limit.
	DefaultBatchMaxBytes = 4 << 20
	// DefaultBatchMaxWait is the zero-value BatchConfig fill deadline.
	DefaultBatchMaxWait = 25 * time.Millisecond
)
View Source
const (
	// EnvelopeSpecVersion is the current native envelope contract.
	EnvelopeSpecVersion = "1.0"
	// DefaultMaxEnvelopeBytes is the default encoded envelope limit.
	DefaultMaxEnvelopeBytes = 1 << 20
	// DefaultMaxHeaders is the default number of application headers.
	DefaultMaxHeaders = 64
	// DefaultMaxHeaderBytes is the default aggregate application-header limit.
	DefaultMaxHeaderBytes = 16 << 10
)
View Source
const ManifestSpecVersion = "1.0"

ManifestSpecVersion is the current topology manifest contract.

Variables

View Source
var (
	// ErrInvalidDescriptor reports an invalid command, event, or query descriptor.
	ErrInvalidDescriptor = errors.New("messenger: invalid descriptor")
	// ErrInvalidMessage reports invalid outgoing metadata or envelope data.
	ErrInvalidMessage = errors.New("messenger: invalid message")
	// ErrDescriptorConflict reports two incompatible descriptors with one wire identity.
	ErrDescriptorConflict = errors.New("messenger: descriptor conflict")
	// ErrHandlerConflict reports a duplicate command handler or subscription ID.
	ErrHandlerConflict = errors.New("messenger: handler conflict")
	// ErrHandlerNotFound reports a missing required local handler.
	ErrHandlerNotFound = errors.New("messenger: handler not found")
	// ErrQueryResultMissing reports successful global middleware completion without a query result.
	ErrQueryResultMissing = errors.New("messenger: query result missing")
	// ErrRouteConflict reports more than one primary route for a descriptor.
	ErrRouteConflict = errors.New("messenger: route conflict")
	// ErrRouteNotFound reports that a descriptor has no outbound route.
	ErrRouteNotFound = errors.New("messenger: route not found")
	// ErrUnsupportedCapability reports a requested semantic guarantee that a route cannot provide.
	ErrUnsupportedCapability = errors.New("messenger: unsupported route capability")
	// ErrMessageExpired reports a message whose expiration boundary has been reached.
	ErrMessageExpired = errors.New("messenger: message expired")
	// ErrMessageNotReady reports a message whose not-before boundary is still in the future.
	ErrMessageNotReady = errors.New("messenger: message not ready")
	// ErrServiceConflict reports a duplicate managed service ID.
	ErrServiceConflict = errors.New("messenger: service conflict")
	// ErrRuntimeNotRunning reports an operation that requires a running runtime.
	ErrRuntimeNotRunning = errors.New("messenger: runtime not running")
	// ErrRuntimeRunning reports a second concurrent call to Runtime.Run.
	ErrRuntimeRunning = errors.New("messenger: runtime already running")
	// ErrRuntimeClosed reports use after a runtime has shut down.
	ErrRuntimeClosed = errors.New("messenger: runtime closed")
	// ErrEnvelopeTooLarge reports an envelope beyond the configured wire limit.
	ErrEnvelopeTooLarge = errors.New("messenger: envelope too large")
	// ErrInvalidBatchResult reports a missing, duplicate, unknown, or otherwise
	// inconsistent item in a batch handler result.
	ErrInvalidBatchResult = errors.New("messenger: invalid batch result")
)

Functions

func BoundedFailureText added in v0.2.0

func BoundedFailureText(sanitizer FailureSanitizer, err error, limit int) string

BoundedFailureText returns sanitized, valid UTF-8 text no longer than limit bytes. It never splits a UTF-8 code point.

func CanonicalizeEnvelope

func CanonicalizeEnvelope(data []byte) ([]byte, error)

CanonicalizeEnvelope parses and re-encodes an envelope using the native deterministic field order. Delivery metadata is never included.

func ContextWithMetadata

func ContextWithMetadata(ctx context.Context, metadata Metadata) context.Context

ContextWithMetadata installs immutable message lineage for an adapter or terminal handler entering the typed messenger boundary. As with standard context helpers, ctx must be non-nil.

func DecodeCommandPayload

func DecodeCommandPayload[T any](descriptor Command[T], data []byte) (T, error)

DecodeCommandPayload decodes codec bytes without an envelope.

func DecodeEventPayload

func DecodeEventPayload[T any](descriptor Event[T], data []byte) (T, error)

DecodeEventPayload decodes codec bytes without an envelope.

func DeferAfter added in v0.3.0

func DeferAfter(err error, delay time.Duration) error

DeferAfter asks a durable consumer to retry after an exact positive delay without consuming a handler attempt.

func DeferDelay added in v0.3.0

func DeferDelay(err error) (time.Duration, bool)

DeferDelay returns the exact no-attempt delay carried by err.

func EncodeCommandEnvelope

func EncodeCommandEnvelope[T any](descriptor Command[T], metadata Metadata, payload T) ([]byte, error)

EncodeCommandEnvelope encodes a typed command with already resolved metadata.

func EncodeEventEnvelope

func EncodeEventEnvelope[T any](descriptor Event[T], metadata Metadata, payload T) ([]byte, error)

EncodeEventEnvelope encodes a typed event with already resolved metadata.

func EnvelopeFingerprint

func EnvelopeFingerprint(data []byte) [sha256.Size]byte

EnvelopeFingerprint returns SHA-256 over canonical encoded envelope bytes.

func HandlerCompletionError added in v0.2.0

func HandlerCompletionError(ctx context.Context, handlerErr error) error

HandlerCompletionError prevents a handler that returns nil after its context deadline from committing its transaction. Handler deadlines remain cooperative: a handler must still observe ctx.Done to stop promptly.

func IsPermanent

func IsPermanent(err error) bool

IsPermanent reports whether err contains a Permanent marker.

func MarshalEnvelope

func MarshalEnvelope(metadata Metadata, payload []byte, encoding DataEncoding) ([]byte, error)

MarshalEnvelope validates and encodes an envelope with a codec payload.

func Permanent

func Permanent(err error) error

Permanent marks an error as non-retryable.

func ReportHandlerPanic added in v0.2.0

func ReportHandlerPanic(
	ctx context.Context,
	reporter PanicReporter,
	handlerID string,
	recovered any,
	stack []byte,
) error

ReportHandlerPanic sends sensitive details to the optional reporter and returns a safe error suitable for retries, observations, logs, and DLQ data.

func RetryAfter

func RetryAfter(err error, delay time.Duration) error

RetryAfter asks a durable transport to retry after an exact positive delay.

func RetryDelay

func RetryDelay(err error) (time.Duration, bool)

RetryDelay returns an explicitly requested retry delay.

func SanitizeError added in v0.2.0

func SanitizeError(sanitizer FailureSanitizer, err error) error

SanitizeError preserves errors.Is/errors.As through Unwrap while exposing only sanitized text through Error.

func SanitizeFailure added in v0.2.0

func SanitizeFailure(sanitizer FailureSanitizer, err error) string

SanitizeFailure returns safe failure text. A nil or typed-nil sanitizer uses DefaultFailureSanitizer.

Types

type BatchConfig added in v0.3.0

type BatchConfig struct {
	MaxMessages int
	MaxBytes    int
	MaxWait     time.Duration
	Middlewares []BatchMiddleware
}

BatchConfig bounds one consumer batch. Its zero value resolves to 100 messages, 4 MiB of canonical envelope bytes, and 25 milliseconds.

func (BatchConfig) Normalize added in v0.3.0

func (c BatchConfig) Normalize(concurrency int) (BatchConfig, error)

Normalize applies zero-value defaults and validates process-wide bounds for the supplied positive batch concurrency.

type BatchHandler added in v0.3.0

type BatchHandler[T any] func(context.Context, []Message[T]) (BatchResult, error)

BatchHandler processes one broker-ordered batch of unique active messages. Implementations must classify the complete batch before performing business SQL and may write only for the successful subset.

func ChainBatchHandler added in v0.3.0

func ChainBatchHandler[T any](
	handler BatchHandler[T],
	middlewares ...BatchHandlerMiddleware[T],
) BatchHandler[T]

ChainBatchHandler applies typed batch middleware with the first item outermost. It returns nil when the chain is invalid.

type BatchHandlerFunc added in v0.3.0

type BatchHandlerFunc func(context.Context) (BatchResult, error)

BatchHandlerFunc is the transport-neutral terminal shape wrapped by batch middleware.

type BatchHandlerMiddleware added in v0.3.0

type BatchHandlerMiddleware[T any] func(BatchHandler[T]) BatchHandler[T]

BatchHandlerMiddleware wraps a typed batch handler.

type BatchItemKey added in v0.3.0

type BatchItemKey struct {
	Source    string
	MessageID MessageID
}

BatchItemKey is the consumer-independent logical identity returned by a BatchHandler. Consumer identity remains an Inbox and transport concern.

type BatchItemResult added in v0.3.0

type BatchItemResult struct {
	Key BatchItemKey
	Err error
}

BatchItemResult classifies one logical message from a BatchHandler input. A nil Err marks success.

type BatchMiddleware added in v0.3.0

type BatchMiddleware func(
	ctx context.Context,
	metadata []Metadata,
	handlerID string,
	next BatchHandlerFunc,
) (BatchResult, error)

BatchMiddleware wraps one batch invocation. Metadata is supplied as a defensive copy in broker order. The first registered middleware is the outermost wrapper.

type BatchPublisher added in v0.3.0

type BatchPublisher[T any] interface {
	PublishBatch(ctx context.Context, payloads []T) ([]Receipt, error)
	PublishMessageBatch(ctx context.Context, outgoing []Outgoing[T]) ([]Receipt, error)
}

BatchPublisher is the narrow typed DI surface for atomic event batches.

func BindBatchPublisher added in v0.3.0

func BindBatchPublisher[T any](messenger *Messenger, descriptor Event[T]) BatchPublisher[T]

BindBatchPublisher returns a narrow atomic batch facade bound to one event.

type BatchResult added in v0.3.0

type BatchResult struct {
	Items []BatchItemResult
}

BatchResult contains exactly one result for every logical message passed to a BatchHandler. Item order is irrelevant because results are keyed.

type BatchResultBuilder added in v0.3.0

type BatchResultBuilder[T any] struct {
	// contains filtered or unexported fields
}

BatchResultBuilder simplifies building a complete and valid BatchResult for a batch of messages. It initializes with every message in the batch marked as succeeded (nil error) and ensures that all input items are preserved in original order.

func NewBatchResultBuilder added in v0.3.0

func NewBatchResultBuilder[T any](messages []Message[T]) *BatchResultBuilder[T]

NewBatchResultBuilder initializes a builder for the supplied batch of messages. All items default to success (nil error).

func (*BatchResultBuilder[T]) Build added in v0.3.0

func (b *BatchResultBuilder[T]) Build() (BatchResult, error)

Build constructs the populated BatchResult containing one item result for every message in the original batch in input order, or returns an error if an unknown key was passed to the builder.

func (*BatchResultBuilder[T]) Error added in v0.3.0

func (b *BatchResultBuilder[T]) Error(message Message[T]) (error, bool)

Error returns the classified error for the message and reports whether the message was present in the batch.

func (*BatchResultBuilder[T]) ErrorKey added in v0.3.0

func (b *BatchResultBuilder[T]) ErrorKey(key BatchItemKey) (error, bool)

ErrorKey returns the classified error for key and reports whether the key was present in the batch.

func (*BatchResultBuilder[T]) Fail added in v0.3.0

func (b *BatchResultBuilder[T]) Fail(message Message[T], err error) *BatchResultBuilder[T]

Fail marks the message as failed with err.

func (*BatchResultBuilder[T]) FailKey added in v0.3.0

func (b *BatchResultBuilder[T]) FailKey(key BatchItemKey, err error) *BatchResultBuilder[T]

FailKey marks the message identified by key as failed with err.

func (*BatchResultBuilder[T]) HasErrors added in v0.3.0

func (b *BatchResultBuilder[T]) HasErrors() bool

HasErrors reports whether any message in the batch has a non-nil error.

func (*BatchResultBuilder[T]) OK added in v0.3.0

func (b *BatchResultBuilder[T]) OK(message Message[T]) *BatchResultBuilder[T]

OK marks the message as successfully processed.

func (*BatchResultBuilder[T]) OKKey added in v0.3.0

func (b *BatchResultBuilder[T]) OKKey(key BatchItemKey) *BatchResultBuilder[T]

OKKey marks the message identified by key as successfully processed.

type BatchRoute added in v0.3.0

type BatchRoute interface {
	Route
	DeliverBatch(ctx context.Context, deliveries []Delivery) ([]Receipt, error)
}

BatchRoute atomically delivers an ordered set of command/event deliveries. The durable Outbox route is the supported implementation; direct broker and local routes intentionally do not implement this capability.

type BatchSender added in v0.3.0

type BatchSender[T any] interface {
	SendBatch(ctx context.Context, payloads []T) ([]Receipt, error)
	SendMessageBatch(ctx context.Context, outgoing []Outgoing[T]) ([]Receipt, error)
}

BatchSender is the narrow typed DI surface for atomic command batches.

func BindBatchSender added in v0.3.0

func BindBatchSender[T any](messenger *Messenger, descriptor Command[T]) BatchSender[T]

BindBatchSender returns a narrow atomic batch facade bound to one command.

type Builder

type Builder struct {
	// contains filtered or unexported fields
}

Builder declares immutable descriptors, handlers, routes, and managed services. It is not safe for concurrent mutation.

func NewBuilder

func NewBuilder(options ...Option) *Builder

NewBuilder constructs an empty messenger builder.

func (*Builder) Build

func (b *Builder) Build() (*Messenger, *Runtime, error)

Build validates and freezes the declared topology.

func (*Builder) HandleCommand

func (b *Builder) HandleCommand[T any](descriptor Command[T], handlerID string, handler Handler[T])

HandleCommand registers the one local handler for a command descriptor. Validation errors are returned by Build.

func (*Builder) HandleCommandFunc

func (b *Builder) HandleCommandFunc[T any](
	descriptor Command[T],
	handlerID string,
	handler PayloadHandler[T],
)

HandleCommandFunc registers a payload-only command handler.

func (*Builder) HandleQuery

func (b *Builder) HandleQuery[Q, R any](
	descriptor Query[Q, R],
	handlerID string,
	handler QueryHandler[Q, R],
)

HandleQuery registers the one local handler for a query descriptor. Validation errors are returned by Build.

func (*Builder) HandleQueryFunc

func (b *Builder) HandleQueryFunc[Q, R any](
	descriptor Query[Q, R],
	handlerID string,
	handler QueryPayloadHandler[Q, R],
)

HandleQueryFunc registers a payload-only local query handler.

func (*Builder) RouteCommand

func (b *Builder) RouteCommand[T any](descriptor Command[T], route Route)

RouteCommand sets the command's one primary outbound route.

func (*Builder) RouteEvent

func (b *Builder) RouteEvent[T any](descriptor Event[T], route Route)

RouteEvent sets the event's one primary outbound route.

func (*Builder) RouteQuery

func (b *Builder) RouteQuery[Q, R any](descriptor Query[Q, R], route LocalQueryRoute)

RouteQuery sets the query's required local request/reply route.

func (*Builder) Subscribe

func (b *Builder) Subscribe[T any](descriptor Event[T], subscriptionID string, handler Handler[T])

Subscribe appends a named local event subscription. Validation errors are returned by Build.

func (*Builder) SubscribeFunc

func (b *Builder) SubscribeFunc[T any](
	descriptor Event[T],
	subscriptionID string,
	handler PayloadHandler[T],
)

SubscribeFunc appends a payload-only event subscription.

func (*Builder) Use

func (b *Builder) Use(serviceID string, service Service)

Use adds a named managed consumer or worker service to the returned Runtime.

func (*Builder) UseMiddleware

func (b *Builder) UseMiddleware(middlewares ...Middleware)

UseMiddleware appends global handler middleware. Validation errors are returned by Build.

type Codec

type Codec[T any] interface {
	Encode(value T) ([]byte, error)
	Decode(data []byte) (T, error)
	ContentType() string
	Encoding() DataEncoding
}

Codec encodes and decodes one descriptor payload type.

func Bytes

func Bytes() Codec[[]byte]

Bytes returns a binary codec that copies byte slices at the boundary.

func CustomCodec

func CustomCodec[T any](
	contentType string,
	encoding DataEncoding,
	encode func(T) ([]byte, error),
	decode func([]byte) (T, error),
) (Codec[T], error)

CustomCodec constructs a codec with an explicit content type and wire encoding.

func JSON

func JSON[T any]() Codec[T]

JSON returns the standard JSON codec for T.

func Text

func Text() Codec[string]

Text returns a UTF-8 text codec.

type Command

type Command[T any] struct {
	// contains filtered or unexported fields
}

Command is an immutable typed command descriptor.

func MustCommand

func MustCommand[T any](name string, schemaVersion int, codec Codec[T]) Command[T]

MustCommand constructs a command descriptor and panics when its declaration is invalid.

func NewCommand

func NewCommand[T any](name string, schemaVersion int, codec Codec[T]) (Command[T], error)

NewCommand constructs a command descriptor.

func (Command[T]) Info

func (d Command[T]) Info() DescriptorInfo

Info returns a copy of the command's wire identity.

func (Command[T]) WithSchema

func (d Command[T]) WithSchema(schema string) Command[T]

WithSchema returns a command descriptor with an explicit schema URI.

type ContextPropagator

type ContextPropagator interface {
	Inject(ctx context.Context, carrier map[string]string)
	Extract(ctx context.Context, carrier map[string]string) context.Context
}

ContextPropagator injects and extracts distributed context through immutable message headers.

func NoopContextPropagator

func NoopContextPropagator() ContextPropagator

NoopContextPropagator returns a propagator that leaves contexts and carriers unchanged.

type DataEncoding

type DataEncoding uint8

DataEncoding selects how encoded payload bytes appear in envelope JSON.

const (
	// DataJSON stores codec output directly in the envelope data field.
	DataJSON DataEncoding = iota + 1
	// DataText stores codec output as a JSON string in the envelope data field.
	DataText
	// DataBinary stores codec output in the envelope dataBase64 field.
	DataBinary
)

func (DataEncoding) MarshalJSON

func (e DataEncoding) MarshalJSON() ([]byte, error)

MarshalJSON encodes a data encoding as its stable wire name.

func (DataEncoding) String

func (e DataEncoding) String() string

String returns the stable wire name of the data encoding.

func (*DataEncoding) UnmarshalJSON

func (e *DataEncoding) UnmarshalJSON(data []byte) error

UnmarshalJSON decodes a stable data-encoding wire name.

type DeepHealthChecker added in v0.2.0

type DeepHealthChecker interface {
	DeepHealth(ctx context.Context) error
}

DeepHealthChecker optionally performs expensive topology and infrastructure validation outside the normal readiness probe path.

type Delivery

type Delivery interface {
	Metadata() Metadata
	HandlerCount() int
	MarshalEnvelope() ([]byte, error)
	Fingerprint() ([sha256.Size]byte, error)
	Invoke(ctx context.Context) error
}

Delivery is the transport-neutral route input. MarshalEnvelope is lazy, so local routes do not serialize payloads. Invoke is intended for local and terminal durable adapters.

type DescriptorInfo

type DescriptorInfo struct {
	Kind          Kind         `json:"kind"`
	Name          string       `json:"name"`
	SchemaVersion int          `json:"schemaVersion"`
	ContentType   string       `json:"contentType"`
	DataEncoding  DataEncoding `json:"dataEncoding"`
	Schema        string       `json:"schema,omitempty"`
}

DescriptorInfo is the transport-neutral public identity of a descriptor.

type Envelope

type Envelope struct {
	SpecVersion   string            `json:"specVersion"`
	ID            MessageID         `json:"id"`
	Kind          Kind              `json:"kind"`
	Name          string            `json:"name"`
	SchemaVersion int               `json:"schemaVersion"`
	Source        string            `json:"source"`
	Subject       string            `json:"subject,omitempty"`
	Time          time.Time         `json:"time"`
	CorrelationID MessageID         `json:"correlationId"`
	CausationID   MessageID         `json:"causationId,omitzero"`
	Key           string            `json:"key,omitempty"`
	ContentType   string            `json:"contentType"`
	DataEncoding  DataEncoding      `json:"dataEncoding"`
	Schema        string            `json:"schema,omitempty"`
	Headers       map[string]string `json:"headers,omitempty"`
	NotBefore     time.Time         `json:"notBefore,omitzero"`
	ExpiresAt     time.Time         `json:"expiresAt,omitzero"`
	Data          json.RawMessage   `json:"data,omitempty"`
	DataBase64    *string           `json:"dataBase64,omitempty"`
}

Envelope is the canonical native wire representation.

func UnmarshalEnvelope

func UnmarshalEnvelope(data []byte) (Envelope, error)

UnmarshalEnvelope parses and validates a native envelope.

func (Envelope) Metadata

func (e Envelope) Metadata() Metadata

Metadata returns a defensive copy of the envelope metadata.

func (Envelope) Payload

func (e Envelope) Payload() ([]byte, DataEncoding, error)

Payload returns decoded codec bytes and their envelope representation.

func (Envelope) Validate

func (e Envelope) Validate() error

Validate checks the native envelope invariants and configured default bounds.

type Event

type Event[T any] struct {
	// contains filtered or unexported fields
}

Event is an immutable typed event descriptor.

func MustEvent

func MustEvent[T any](name string, schemaVersion int, codec Codec[T]) Event[T]

MustEvent constructs an event descriptor and panics when its declaration is invalid.

func NewEvent

func NewEvent[T any](name string, schemaVersion int, codec Codec[T]) (Event[T], error)

NewEvent constructs an event descriptor.

func (Event[T]) Info

func (d Event[T]) Info() DescriptorInfo

Info returns a copy of the event's wire identity.

func (Event[T]) WithSchema

func (d Event[T]) WithSchema(schema string) Event[T]

WithSchema returns an event descriptor with an explicit schema URI.

type FailureSanitizer added in v0.2.0

type FailureSanitizer interface {
	SanitizeFailure(err error) string
}

FailureSanitizer converts an error to text safe for operational channels such as default logs and telemetry observations. In batch consumers, durable DLQ wire text strictly uses the conservative built-in sanitizer to protect rebalance and finalization bounds, while configured host sanitizers apply to observations and operational logs.

func DefaultFailureSanitizer added in v0.2.0

func DefaultFailureSanitizer() FailureSanitizer

DefaultFailureSanitizer returns the conservative built-in sanitizer. Hosts may opt in to richer text for observations and operational logs with an explicit FailureSanitizer implementation; durable batch DLQ wire payloads always retain the conservative sanitizer.

type FailureSanitizerFunc added in v0.2.0

type FailureSanitizerFunc func(error) string

FailureSanitizerFunc adapts a function to FailureSanitizer.

func (FailureSanitizerFunc) SanitizeFailure added in v0.2.0

func (f FailureSanitizerFunc) SanitizeFailure(err error) string

SanitizeFailure implements FailureSanitizer.

type Handler

type Handler[T any] func(context.Context, Message[T]) error

Handler processes one typed message.

func ChainHandler

func ChainHandler[T any](handler Handler[T], middlewares ...HandlerMiddleware[T]) Handler[T]

ChainHandler applies typed middleware with the first item outermost. It returns nil when handler, a middleware, or a middleware result is nil so the receiving Builder or durable consumer can reject the invalid chain.

func HandlePayload

func HandlePayload[T any](handler PayloadHandler[T]) Handler[T]

HandlePayload adapts a payload-only handler to the primary Handler contract.

type HandlerFunc

type HandlerFunc func(context.Context) error

HandlerFunc is the transport-neutral terminal handler shape used by global middleware.

type HandlerMiddleware

type HandlerMiddleware[T any] func(Handler[T]) Handler[T]

HandlerMiddleware wraps a typed handler.

type HandlerPanicError added in v0.2.0

type HandlerPanicError interface {
	error
	HandlerPanicID() string
}

HandlerPanicError is the transport-neutral safe view of a recovered handler or middleware panic. Independently versioned adapters implement this interface structurally without exposing the recovered value or stack.

type IDGenerator

type IDGenerator interface {
	New() (MessageID, error)
}

IDGenerator creates stable message identities.

func UUIDv7Generator

func UUIDv7Generator() IDGenerator

UUIDv7Generator returns the default cryptographically random UUIDv7 generator.

type Kind

type Kind string

Kind identifies the semantic message category.

const (
	// KindCommand identifies a command with one logical handler.
	KindCommand Kind = "command"
	// KindEvent identifies an event with zero or more subscriptions.
	KindEvent Kind = "event"
	// KindQuery identifies a local request/reply query with one handler.
	KindQuery Kind = "query"
)

type LivenessChecker added in v0.2.0

type LivenessChecker interface {
	Liveness(ctx context.Context) error
}

LivenessChecker optionally separates process liveness from readiness and transient broker or topology failures.

type LocalAsyncConfig

type LocalAsyncConfig struct {
	Capacity int
	Workers  int
	// DetachExecution is retained for source compatibility. Accepted one-way
	// jobs always detach execution from the caller's cancellation and deadline.
	// Query calls always retain the caller context for execution and waiting.
	//
	// Deprecated: caller context controls admission only.
	DetachExecution bool
}

LocalAsyncConfig bounds local asynchronous admission and execution.

type LocalAsyncRoute

type LocalAsyncRoute struct {
	// contains filtered or unexported fields
}

LocalAsyncRoute admits handler calls to a bounded GoBus async runtime.

func NewLocalAsyncRoute

func NewLocalAsyncRoute(name string, config LocalAsyncConfig) (*LocalAsyncRoute, error)

NewLocalAsyncRoute constructs a named bounded local asynchronous route.

func (*LocalAsyncRoute) BeginDrain

func (r *LocalAsyncRoute) BeginDrain()

BeginDrain rejects new work and drains accepted jobs.

func (*LocalAsyncRoute) Deliver

func (r *LocalAsyncRoute) Deliver(ctx context.Context, delivery Delivery) (Receipt, error)

Deliver implements Route and reports admission, not handler completion.

func (*LocalAsyncRoute) ManagedService

func (r *LocalAsyncRoute) ManagedService() (string, Service)

ManagedService exposes this route to Builder runtime aggregation.

func (*LocalAsyncRoute) Name

func (r *LocalAsyncRoute) Name() string

Name implements Route.

func (*LocalAsyncRoute) Readiness

func (r *LocalAsyncRoute) Readiness(context.Context) error

Readiness verifies that this route is accepting work.

func (*LocalAsyncRoute) Run

func (r *LocalAsyncRoute) Run(ctx context.Context) error

Run starts queue workers and blocks until cancellation or draining.

func (*LocalAsyncRoute) Shutdown

func (r *LocalAsyncRoute) Shutdown(ctx context.Context) error

Shutdown drains or force-cancels the route within ctx.

type LocalQueryRoute

type LocalQueryRoute interface {
	Name() string
	// contains filtered or unexported methods
}

LocalQueryRoute is the sealed local request/reply route contract. The built-in LocalSyncRoute and LocalAsyncRoute are its only implementations.

type LocalSyncRoute

type LocalSyncRoute struct {
	// contains filtered or unexported fields
}

LocalSyncRoute executes handlers synchronously through a private GoBus instance.

func NewLocalSyncRoute

func NewLocalSyncRoute() *LocalSyncRoute

NewLocalSyncRoute constructs a local synchronous route.

func (*LocalSyncRoute) Deliver

func (r *LocalSyncRoute) Deliver(ctx context.Context, delivery Delivery) (Receipt, error)

Deliver implements Route.

func (*LocalSyncRoute) Name

func (*LocalSyncRoute) Name() string

Name implements Route.

type LogAttr

type LogAttr struct {
	Key   string
	Value any
}

LogAttr is one structured logging attribute.

type LogLevel

type LogLevel uint8

LogLevel is the transport-neutral severity understood by Logger.

const (
	// LogDebug records diagnostic lifecycle details.
	LogDebug LogLevel = iota
	// LogInfo records ordinary lifecycle transitions.
	LogInfo
	// LogWarn records recoverable infrastructure failures.
	LogWarn
	// LogError records infrastructure failures requiring attention.
	LogError
)

type Logger

type Logger interface {
	Log(ctx context.Context, level LogLevel, message string, attrs ...LogAttr)
}

Logger is the minimal structured logging contract used by GoMessenger. Implementations must not retain ctx and should return quickly.

func AdaptSlog

func AdaptSlog(logger *slog.Logger) Logger

AdaptSlog adapts a standard slog logger. A nil logger returns a no-op adapter, which makes optional host wiring safe.

type Manifest

type Manifest struct {
	SpecVersion string               `json:"specVersion"`
	Source      string               `json:"source"`
	Descriptors []ManifestDescriptor `json:"descriptors"`
	Services    []string             `json:"services,omitempty"`
}

Manifest is a deterministic, secret-free description of runtime topology.

func (Manifest) Validate

func (m Manifest) Validate() error

Validate checks manifest structure without constructing a runtime.

type ManifestDescriptor

type ManifestDescriptor struct {
	DescriptorInfo
	Route      string   `json:"route,omitempty"`
	HandlerIDs []string `json:"handlerIds,omitempty"`
}

ManifestDescriptor describes one typed descriptor and its static route.

type Message

type Message[T any] struct {
	Metadata Metadata
	Payload  T
}

Message is the typed value passed to a handler.

func DecodeCommand

func DecodeCommand[T any](descriptor Command[T], data []byte) (Message[T], error)

DecodeCommand decodes and verifies a native command envelope.

func DecodeEvent

func DecodeEvent[T any](descriptor Event[T], data []byte) (Message[T], error)

DecodeEvent decodes and verifies a native event envelope.

type MessageID

type MessageID [16]byte

MessageID is a UUID-compatible 128-bit message identity.

func ParseMessageID

func ParseMessageID(value string) (MessageID, error)

ParseMessageID parses the canonical UUID text form.

func (MessageID) IsZero

func (id MessageID) IsZero() bool

IsZero reports whether no message identity is set.

func (MessageID) MarshalJSON

func (id MessageID) MarshalJSON() ([]byte, error)

MarshalJSON implements json.Marshaler.

func (MessageID) MarshalText

func (id MessageID) MarshalText() ([]byte, error)

MarshalText implements encoding.TextMarshaler.

func (MessageID) String

func (id MessageID) String() string

String returns the canonical lowercase UUID text form.

func (*MessageID) UnmarshalJSON

func (id *MessageID) UnmarshalJSON(data []byte) error

UnmarshalJSON implements json.Unmarshaler.

func (*MessageID) UnmarshalText

func (id *MessageID) UnmarshalText(text []byte) error

UnmarshalText implements encoding.TextUnmarshaler.

type Messenger

type Messenger struct {
	// contains filtered or unexported fields
}

Messenger sends commands, executes local queries, and publishes events through immutable descriptor bindings.

func (*Messenger) Manifest

func (m *Messenger) Manifest() Manifest

Manifest returns a defensive copy of the messenger topology.

func (*Messenger) MarshalManifest

func (m *Messenger) MarshalManifest() ([]byte, error)

MarshalManifest returns deterministic indented JSON suitable for gomessengerctl.

func (*Messenger) Publish

func (m *Messenger) Publish[T any](ctx context.Context, descriptor Event[T], payload T) (Receipt, error)

Publish publishes an event with generated metadata.

func (*Messenger) PublishBatch added in v0.3.0

func (m *Messenger) PublishBatch[T any](
	ctx context.Context,
	descriptor Event[T],
	payloads []T,
) ([]Receipt, error)

PublishBatch atomically stages event payloads through a BatchRoute.

func (*Messenger) PublishMessage

func (m *Messenger) PublishMessage[T any](
	ctx context.Context,
	descriptor Event[T],
	outgoing Outgoing[T],
) (Receipt, error)

PublishMessage publishes an event with explicit optional metadata.

func (*Messenger) PublishMessageBatch added in v0.3.0

func (m *Messenger) PublishMessageBatch[T any](
	ctx context.Context,
	descriptor Event[T],
	outgoing []Outgoing[T],
) ([]Receipt, error)

PublishMessageBatch validates all event messages and atomically stages them through the configured BatchRoute.

func (*Messenger) Query

func (m *Messenger) Query[Q, R any](
	ctx context.Context,
	descriptor Query[Q, R],
	payload Q,
) (R, error)

Query executes a typed local request/reply call through its configured route.

Example
package main

import (
	"context"
	"fmt"

	messenger "github.com/assurrussa/gomessenger"
)

func main() {
	type findArticle struct{ ID int64 }
	type articleView struct {
		ID    int64
		Title string
	}

	find := messenger.MustQuery[findArticle, articleView]("article.find", 1, messenger.JSON[findArticle]())
	builder := messenger.NewBuilder(messenger.WithSource("urn:service:catalog"))
	builder.HandleQueryFunc(find, "article-reader", func(_ context.Context, query findArticle) (articleView, error) {
		return articleView{ID: query.ID, Title: "CQRS in Go"}, nil
	})
	builder.RouteQuery(find, messenger.NewLocalSyncRoute())
	bus, _, err := builder.Build()
	if err != nil {
		panic(err)
	}

	reader := messenger.BindQuerier(bus, find)
	article, err := reader.Query(context.Background(), findArticle{ID: 42})
	if err != nil {
		panic(err)
	}
	fmt.Println(article.ID, article.Title)

}
Output:
42 CQRS in Go

func (*Messenger) Send

func (m *Messenger) Send[T any](ctx context.Context, descriptor Command[T], payload T) (Receipt, error)

Send sends a command with generated metadata.

Example
package main

import (
	"context"
	"fmt"

	messenger "github.com/assurrussa/gomessenger"
)

func main() {
	type resizeMedia struct {
		JobID int64 `json:"jobId"`
	}

	resize := messenger.MustCommand("media.resize", 1, messenger.JSON[resizeMedia]())
	builder := messenger.NewBuilder(messenger.WithSource("urn:service:media-resizer"))
	builder.HandleCommandFunc(resize, "media-worker", func(_ context.Context, payload resizeMedia) error {
		fmt.Println("handler", payload.JobID)
		return nil
	})
	builder.RouteCommand(resize, messenger.NewLocalSyncRoute())
	bus, _, err := builder.Build()
	if err != nil {
		panic(err)
	}
	resizeSender := messenger.BindSender(bus, resize)
	receipt, err := resizeSender.Send(context.Background(), resizeMedia{JobID: 42})
	if err != nil {
		panic(err)
	}
	fmt.Println(receipt.State)

}
Output:
handler 42
completed

func (*Messenger) SendBatch added in v0.3.0

func (m *Messenger) SendBatch[T any](
	ctx context.Context,
	descriptor Command[T],
	payloads []T,
) ([]Receipt, error)

SendBatch atomically stages command payloads through a BatchRoute.

func (*Messenger) SendMessage

func (m *Messenger) SendMessage[T any](
	ctx context.Context,
	descriptor Command[T],
	outgoing Outgoing[T],
) (Receipt, error)

SendMessage sends a command with explicit optional metadata.

func (*Messenger) SendMessageBatch added in v0.3.0

func (m *Messenger) SendMessageBatch[T any](
	ctx context.Context,
	descriptor Command[T],
	outgoing []Outgoing[T],
) ([]Receipt, error)

SendMessageBatch validates all command messages and atomically stages them through the configured BatchRoute.

type Metadata

type Metadata struct {
	ID            MessageID
	Kind          Kind
	Name          string
	SchemaVersion int
	Source        string
	Subject       string
	Time          time.Time
	CorrelationID MessageID
	CausationID   MessageID
	Key           string
	ContentType   string
	Schema        string
	Headers       map[string]string
	NotBefore     time.Time
	ExpiresAt     time.Time
}

Metadata is canonical message metadata independent of a transport attempt.

func MetadataFromContext

func MetadataFromContext(ctx context.Context) (Metadata, bool)

MetadataFromContext returns the currently handled message metadata.

type Middleware

type Middleware func(
	ctx context.Context,
	metadata Metadata,
	handlerID string,
	next HandlerFunc,
) error

Middleware wraps one local query/command/event or durable handler. The first registered middleware is the outermost wrapper and may short-circuit by not calling next.

type Observation

type Observation struct {
	Operation            Operation
	MessageID            MessageID
	Kind                 Kind
	Name                 string
	SchemaVersion        int
	Route                string
	HandlerID            string
	ConsumerID           string
	ServiceID            string
	Attempt              uint64
	Duplicate            bool
	RetryDelay           time.Duration
	BatchSize            int
	BatchBytes           int
	BatchHandlerMessages int
	BatchACKs            int
	BatchRetries         int
	BatchDeferrals       int
	BatchDLQs            int
	BatchFillDuration    time.Duration
	BatchHandlerDuration time.Duration
	State                ReceiptState
	StartedAt            time.Time
	Duration             time.Duration
	Err                  error
}

Observation contains bounded operational data. Observers decide which fields are safe for low-cardinality metric labels.

type Observer

type Observer interface {
	Observe(ctx context.Context, observation Observation)
}

Observer receives messaging lifecycle observations. Implementations must not retain ctx and should return quickly.

func NewLoggingObserver

func NewLoggingObserver(logger Logger) Observer

NewLoggingObserver reports sanitized observations through logger. Successful operations use Debug and failed operations use Error. A nil logger creates a no-op observer.

func NewSanitizedLoggingObserver added in v0.2.0

func NewSanitizedLoggingObserver(logger Logger, sanitizer FailureSanitizer) Observer

NewSanitizedLoggingObserver reports observations with an explicit failure sanitizer. A nil sanitizer uses DefaultFailureSanitizer.

type Operation

type Operation string

Operation identifies an observable messaging boundary.

const (
	// OperationDeliver covers outbound route delivery.
	OperationDeliver Operation = "deliver"
	// OperationExpire reports a local delivery skipped before execution because its deadline passed.
	OperationExpire Operation = "expire"
	// OperationHandle covers local or durable handler execution.
	OperationHandle Operation = "handle"
	// OperationBatchHandle covers one durable batch handler transaction.
	OperationBatchHandle Operation = "batch_handle"
	// OperationQuery covers a complete local request/reply call.
	OperationQuery Operation = "query"
	// OperationService covers managed service completion.
	OperationService Operation = "service"
	// OperationBrokerAck covers broker-confirmed acknowledgement of a consumed message.
	OperationBrokerAck Operation = "broker_ack"
	// OperationOffsetCommit covers transactional Kafka offset finalization.
	OperationOffsetCommit Operation = "offset_commit"
	// OperationRetryHandoff covers durable retry scheduling or broker hand-off.
	OperationRetryHandoff Operation = "retry_handoff"
	// OperationDLQHandoff covers durable terminal hand-off to a dead-letter destination.
	OperationDLQHandoff Operation = "dlq_handoff"
)

type Option

type Option func(*Builder)

Option configures a Builder.

func WithClock

func WithClock(clock func() time.Time) Option

WithClock overrides the UTC wall clock used for messages and receipts.

func WithContextPropagator

func WithContextPropagator(propagator ContextPropagator) Option

WithContextPropagator sets distributed-context injection for outgoing metadata. The default propagator is a no-op.

func WithIDGenerator

func WithIDGenerator(generator IDGenerator) Option

WithIDGenerator overrides UUIDv7 generation, primarily for deterministic tests.

func WithLogger

func WithLogger(logger Logger) Option

WithLogger sets the core structured logger. The default logger is a no-op.

func WithObserver

func WithObserver(observer Observer) Option

WithObserver appends a lifecycle observer.

func WithPanicReporter added in v0.2.0

func WithPanicReporter(reporter PanicReporter) Option

WithPanicReporter enables explicit handling of sensitive recovered-panic values and stacks. Without it, only a sanitized HandlerPanicError is emitted.

func WithRuntimeShutdownTimeout added in v0.2.0

func WithRuntimeShutdownTimeout(timeout time.Duration) Option

WithRuntimeShutdownTimeout sets the internal bound used when Run owns service shutdown after cancellation or an unexpected service return.

func WithSource

func WithSource(source string) Option

WithSource sets the required stable producer identity.

type Outgoing

type Outgoing[T any] struct {
	Payload  T
	Metadata OutgoingMetadata
}

Outgoing combines a typed payload with optional explicit metadata.

type OutgoingMetadata

type OutgoingMetadata struct {
	ID            MessageID
	Subject       string
	Time          time.Time
	CorrelationID MessageID
	CausationID   MessageID
	Key           string
	Headers       map[string]string
	NotBefore     time.Time
	ExpiresAt     time.Time
}

OutgoingMetadata customizes metadata generated for a new outgoing message. Source, kind, name, schema version, content type, and schema come from the builder and descriptor and cannot be overridden per call.

type PanicReport added in v0.2.0

type PanicReport struct {
	HandlerID string
	Value     any
	Stack     []byte
}

PanicReport contains sensitive diagnostics for an explicitly configured PanicReporter. Value and Stack must not be written to untrusted logs or DLQ records without host-side redaction.

type PanicReporter added in v0.2.0

type PanicReporter interface {
	ReportPanic(ctx context.Context, handlerID string, recovered any, stack []byte)
}

PanicReporter receives sensitive recovered-panic diagnostics. The default is to drop these details and return only HandlerPanicError to application code.

type PanicReporterFunc added in v0.2.0

type PanicReporterFunc func(context.Context, PanicReport)

PanicReporterFunc adapts a function to PanicReporter.

func (PanicReporterFunc) ReportPanic added in v0.2.0

func (f PanicReporterFunc) ReportPanic(
	ctx context.Context,
	handlerID string,
	recovered any,
	stack []byte,
)

ReportPanic implements PanicReporter.

type PayloadHandler

type PayloadHandler[T any] func(context.Context, T) error

PayloadHandler processes only the payload and ignores message metadata.

type Publisher

type Publisher[T any] interface {
	Publish(ctx context.Context, payload T) (Receipt, error)
	PublishMessage(ctx context.Context, outgoing Outgoing[T]) (Receipt, error)
}

Publisher is the ordinary generic DI interface for a bound event descriptor.

func BindPublisher

func BindPublisher[T any](messenger *Messenger, descriptor Event[T]) Publisher[T]

BindPublisher returns a narrow DI facade bound to one event descriptor.

type Querier

type Querier[Q, R any] interface {
	Query(ctx context.Context, payload Q) (R, error)
}

Querier is the ordinary generic DI interface for a bound query descriptor.

func BindQuerier

func BindQuerier[Q, R any](messenger *Messenger, descriptor Query[Q, R]) Querier[Q, R]

BindQuerier returns a narrow DI facade bound to one query descriptor.

type Query

type Query[Q, R any] struct {
	// contains filtered or unexported fields
}

Query is an immutable typed local query descriptor. Its codec describes the request Q only; R is a compile-time result identity and is never serialized.

func MustQuery

func MustQuery[Q, R any](name string, schemaVersion int, codec Codec[Q]) Query[Q, R]

MustQuery constructs a typed local query descriptor and panics when its declaration is invalid.

func NewQuery

func NewQuery[Q, R any](name string, schemaVersion int, codec Codec[Q]) (Query[Q, R], error)

NewQuery constructs a typed local query descriptor.

func (Query[Q, R]) Info

func (d Query[Q, R]) Info() DescriptorInfo

Info returns a copy of the query request identity.

func (Query[Q, R]) WithSchema

func (d Query[Q, R]) WithSchema(schema string) Query[Q, R]

WithSchema returns a query descriptor with an explicit request schema URI.

type QueryHandler

type QueryHandler[Q, R any] func(context.Context, Message[Q]) (R, error)

QueryHandler processes one typed query and returns its typed result.

func ChainQueryHandler

func ChainQueryHandler[Q, R any](
	handler QueryHandler[Q, R],
	middlewares ...QueryHandlerMiddleware[Q, R],
) QueryHandler[Q, R]

ChainQueryHandler applies typed query middleware with the first item outermost. It returns nil when the chain is invalid.

func HandleQueryPayload

func HandleQueryPayload[Q, R any](handler QueryPayloadHandler[Q, R]) QueryHandler[Q, R]

HandleQueryPayload adapts a payload-only query handler to QueryHandler.

type QueryHandlerMiddleware

type QueryHandlerMiddleware[Q, R any] func(QueryHandler[Q, R]) QueryHandler[Q, R]

QueryHandlerMiddleware wraps a typed query handler and may return a cached or synthetic result without invoking the wrapped handler.

type QueryPayloadHandler

type QueryPayloadHandler[Q, R any] func(context.Context, Q) (R, error)

QueryPayloadHandler processes only a query payload and returns its typed result.

type Receipt

type Receipt struct {
	MessageID MessageID    `json:"messageId"`
	Route     string       `json:"route"`
	State     ReceiptState `json:"state"`
	At        time.Time    `json:"at"`
}

Receipt describes the guarantee reached by one primary route.

type ReceiptState

type ReceiptState string

ReceiptState describes what a successful route call has guaranteed.

const (
	// ReceiptCompleted means local synchronous handlers completed.
	ReceiptCompleted ReceiptState = "completed"
	// ReceiptAccepted means bounded in-process async admission succeeded.
	ReceiptAccepted ReceiptState = "accepted"
	// ReceiptStaged means an outbox write succeeded in the current transaction.
	ReceiptStaged ReceiptState = "staged"
	// ReceiptBrokerConfirmed means the broker confirmed persistence.
	ReceiptBrokerConfirmed ReceiptState = "broker_confirmed"
	// ReceiptNoop means a local event had no subscribers.
	ReceiptNoop ReceiptState = "noop"
)

type Route

type Route interface {
	Name() string
	Deliver(ctx context.Context, delivery Delivery) (Receipt, error)
}

Route is one static primary delivery route.

type Runtime

type Runtime struct {
	// contains filtered or unexported fields
}

Runtime supervises the services declared on one immutable Builder. It never restarts a service automatically.

func (*Runtime) BeginDrain

func (r *Runtime) BeginDrain()

BeginDrain marks the runtime unready and asks every service to stop admission.

func (*Runtime) DeepHealth added in v0.2.0

func (r *Runtime) DeepHealth(ctx context.Context) error

DeepHealth performs explicit, potentially expensive service health and topology checks. It is intended for diagnostics or a low-frequency probe.

func (*Runtime) Liveness added in v0.2.0

func (r *Runtime) Liveness(ctx context.Context) error

Liveness checks that Runtime has not terminated and invokes optional service liveness checks without requiring readiness or topology access.

func (*Runtime) Readiness

func (r *Runtime) Readiness(ctx context.Context) error

Readiness checks runtime admission state and each service's lightweight readiness contract. Expensive topology validation belongs in DeepHealth.

func (*Runtime) Run

func (r *Runtime) Run(ctx context.Context) error

Run starts every service and blocks until cancellation, draining, or the first unexpected service return.

func (*Runtime) Shutdown

func (r *Runtime) Shutdown(ctx context.Context) error

Shutdown drains all services and waits for a concurrent Run to finish. Coordinate it outside handlers executing on this Runtime.

type Sender

type Sender[T any] interface {
	Send(ctx context.Context, payload T) (Receipt, error)
	SendMessage(ctx context.Context, outgoing Outgoing[T]) (Receipt, error)
}

Sender is the ordinary generic DI interface for a bound command descriptor.

func BindSender

func BindSender[T any](messenger *Messenger, descriptor Command[T]) Sender[T]

BindSender returns a narrow DI facade bound to one command descriptor.

type Service

type Service interface {
	Run(ctx context.Context) error
	Readiness(ctx context.Context) error
	BeginDrain()
	Shutdown(ctx context.Context) error
}

Service is a host-supervised managed consumer or worker lifecycle. BeginDrain must be non-blocking, stop admission, and cause a running Run call to return without requiring Shutdown to be invoked first. Shutdown waits for or force-cancels remaining work within its context.

type ServiceProvider

type ServiceProvider interface {
	ManagedService() (serviceID string, service Service)
}

ServiceProvider lets a route contribute one managed service to Builder.

Directories

Path Synopsis
adapters
inbox module
kafka module
nats module
outbox module
internal
batchruntime
Package batchruntime centralizes the transport-neutral batch handler contract shared by durable adapters.
Package batchruntime centralizes the transport-neutral batch handler contract shared by durable adapters.
tools

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL