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 ¶
- Constants
- Variables
- func BoundedFailureText(sanitizer FailureSanitizer, err error, limit int) string
- func CanonicalizeEnvelope(data []byte) ([]byte, error)
- func ContextWithMetadata(ctx context.Context, metadata Metadata) context.Context
- func DecodeCommandPayload[T any](descriptor Command[T], data []byte) (T, error)
- func DecodeEventPayload[T any](descriptor Event[T], data []byte) (T, error)
- func DeferAfter(err error, delay time.Duration) error
- func DeferDelay(err error) (time.Duration, bool)
- func EncodeCommandEnvelope[T any](descriptor Command[T], metadata Metadata, payload T) ([]byte, error)
- func EncodeEventEnvelope[T any](descriptor Event[T], metadata Metadata, payload T) ([]byte, error)
- func EnvelopeFingerprint(data []byte) [sha256.Size]byte
- func HandlerCompletionError(ctx context.Context, handlerErr error) error
- func IsPermanent(err error) bool
- func MarshalEnvelope(metadata Metadata, payload []byte, encoding DataEncoding) ([]byte, error)
- func Permanent(err error) error
- func ReportHandlerPanic(ctx context.Context, reporter PanicReporter, handlerID string, recovered any, ...) error
- func RetryAfter(err error, delay time.Duration) error
- func RetryDelay(err error) (time.Duration, bool)
- func SanitizeError(sanitizer FailureSanitizer, err error) error
- func SanitizeFailure(sanitizer FailureSanitizer, err error) string
- type BatchConfig
- type BatchHandler
- type BatchHandlerFunc
- type BatchHandlerMiddleware
- type BatchItemKey
- type BatchItemResult
- type BatchMiddleware
- type BatchPublisher
- type BatchResult
- type BatchResultBuilder
- func (b *BatchResultBuilder[T]) Build() (BatchResult, error)
- func (b *BatchResultBuilder[T]) Error(message Message[T]) (error, bool)
- func (b *BatchResultBuilder[T]) ErrorKey(key BatchItemKey) (error, bool)
- func (b *BatchResultBuilder[T]) Fail(message Message[T], err error) *BatchResultBuilder[T]
- func (b *BatchResultBuilder[T]) FailKey(key BatchItemKey, err error) *BatchResultBuilder[T]
- func (b *BatchResultBuilder[T]) HasErrors() bool
- func (b *BatchResultBuilder[T]) OK(message Message[T]) *BatchResultBuilder[T]
- func (b *BatchResultBuilder[T]) OKKey(key BatchItemKey) *BatchResultBuilder[T]
- type BatchRoute
- type BatchSender
- type Builder
- func (b *Builder) Build() (*Messenger, *Runtime, error)
- func (b *Builder) HandleCommand[T any](descriptor Command[T], handlerID string, handler Handler[T])
- func (b *Builder) HandleCommandFunc[T any](descriptor Command[T], handlerID string, handler PayloadHandler[T])
- func (b *Builder) HandleQuery[Q, R any](descriptor Query[Q, R], handlerID string, handler QueryHandler[Q, R])
- func (b *Builder) HandleQueryFunc[Q, R any](descriptor Query[Q, R], handlerID string, handler QueryPayloadHandler[Q, R])
- func (b *Builder) RouteCommand[T any](descriptor Command[T], route Route)
- func (b *Builder) RouteEvent[T any](descriptor Event[T], route Route)
- func (b *Builder) RouteQuery[Q, R any](descriptor Query[Q, R], route LocalQueryRoute)
- func (b *Builder) Subscribe[T any](descriptor Event[T], subscriptionID string, handler Handler[T])
- func (b *Builder) SubscribeFunc[T any](descriptor Event[T], subscriptionID string, handler PayloadHandler[T])
- func (b *Builder) Use(serviceID string, service Service)
- func (b *Builder) UseMiddleware(middlewares ...Middleware)
- type Codec
- type Command
- type ContextPropagator
- type DataEncoding
- type DeepHealthChecker
- type Delivery
- type DescriptorInfo
- type Envelope
- type Event
- type FailureSanitizer
- type FailureSanitizerFunc
- type Handler
- type HandlerFunc
- type HandlerMiddleware
- type HandlerPanicError
- type IDGenerator
- type Kind
- type LivenessChecker
- type LocalAsyncConfig
- type LocalAsyncRoute
- func (r *LocalAsyncRoute) BeginDrain()
- func (r *LocalAsyncRoute) Deliver(ctx context.Context, delivery Delivery) (Receipt, error)
- func (r *LocalAsyncRoute) ManagedService() (string, Service)
- func (r *LocalAsyncRoute) Name() string
- func (r *LocalAsyncRoute) Readiness(context.Context) error
- func (r *LocalAsyncRoute) Run(ctx context.Context) error
- func (r *LocalAsyncRoute) Shutdown(ctx context.Context) error
- type LocalQueryRoute
- type LocalSyncRoute
- type LogAttr
- type LogLevel
- type Logger
- type Manifest
- type ManifestDescriptor
- type Message
- type MessageID
- type Messenger
- func (m *Messenger) Manifest() Manifest
- func (m *Messenger) MarshalManifest() ([]byte, error)
- func (m *Messenger) Publish[T any](ctx context.Context, descriptor Event[T], payload T) (Receipt, error)
- func (m *Messenger) PublishBatch[T any](ctx context.Context, descriptor Event[T], payloads []T) ([]Receipt, error)
- func (m *Messenger) PublishMessage[T any](ctx context.Context, descriptor Event[T], outgoing Outgoing[T]) (Receipt, error)
- func (m *Messenger) PublishMessageBatch[T any](ctx context.Context, descriptor Event[T], outgoing []Outgoing[T]) ([]Receipt, error)
- func (m *Messenger) Query[Q, R any](ctx context.Context, descriptor Query[Q, R], payload Q) (R, error)
- func (m *Messenger) Send[T any](ctx context.Context, descriptor Command[T], payload T) (Receipt, error)
- func (m *Messenger) SendBatch[T any](ctx context.Context, descriptor Command[T], payloads []T) ([]Receipt, error)
- func (m *Messenger) SendMessage[T any](ctx context.Context, descriptor Command[T], outgoing Outgoing[T]) (Receipt, error)
- func (m *Messenger) SendMessageBatch[T any](ctx context.Context, descriptor Command[T], outgoing []Outgoing[T]) ([]Receipt, error)
- type Metadata
- type Middleware
- type Observation
- type Observer
- type Operation
- type Option
- func WithClock(clock func() time.Time) Option
- func WithContextPropagator(propagator ContextPropagator) Option
- func WithIDGenerator(generator IDGenerator) Option
- func WithLogger(logger Logger) Option
- func WithObserver(observer Observer) Option
- func WithPanicReporter(reporter PanicReporter) Option
- func WithRuntimeShutdownTimeout(timeout time.Duration) Option
- func WithSource(source string) Option
- type Outgoing
- type OutgoingMetadata
- type PanicReport
- type PanicReporter
- type PanicReporterFunc
- type PayloadHandler
- type Publisher
- type Querier
- type Query
- type QueryHandler
- type QueryHandlerMiddleware
- type QueryPayloadHandler
- type Receipt
- type ReceiptState
- type Route
- type Runtime
- type Sender
- type Service
- type ServiceProvider
Examples ¶
Constants ¶
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 )
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 )
const ManifestSpecVersion = "1.0"
ManifestSpecVersion is the current topology manifest contract.
Variables ¶
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 ¶
CanonicalizeEnvelope parses and re-encodes an envelope using the native deterministic field order. Delivery metadata is never included.
func ContextWithMetadata ¶
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 ¶
DecodeCommandPayload decodes codec bytes without an envelope.
func DecodeEventPayload ¶
DecodeEventPayload decodes codec bytes without an envelope.
func DeferAfter ¶ added in v0.3.0
DeferAfter asks a durable consumer to retry after an exact positive delay without consuming a handler attempt.
func DeferDelay ¶ added in v0.3.0
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 ¶
EncodeEventEnvelope encodes a typed event with already resolved metadata.
func EnvelopeFingerprint ¶
EnvelopeFingerprint returns SHA-256 over canonical encoded envelope bytes.
func HandlerCompletionError ¶ added in v0.2.0
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 ¶
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 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 ¶
RetryAfter asks a durable transport to retry after an exact positive delay.
func RetryDelay ¶
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
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
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 ¶
NewBuilder constructs an empty messenger builder.
func (*Builder) HandleCommand ¶
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 ¶
RouteCommand sets the command's one primary outbound route.
func (*Builder) RouteEvent ¶
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 ¶
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) 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.
type Command ¶
type Command[T any] struct { // contains filtered or unexported fields }
Command is an immutable typed command descriptor.
func MustCommand ¶
MustCommand constructs a command descriptor and panics when its declaration is invalid.
func NewCommand ¶
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 ¶
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
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 ¶
UnmarshalEnvelope parses and validates a native envelope.
type Event ¶
type Event[T any] struct { // contains filtered or unexported fields }
Event is an immutable typed event descriptor.
func MustEvent ¶
MustEvent constructs an event descriptor and panics when its declaration is invalid.
func (Event[T]) Info ¶
func (d Event[T]) Info() DescriptorInfo
Info returns a copy of the event's wire identity.
func (Event[T]) WithSchema ¶
WithSchema returns an event descriptor with an explicit schema URI.
type FailureSanitizer ¶ added in v0.2.0
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
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 ¶
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 ¶
HandlerFunc is the transport-neutral terminal handler shape used by global middleware.
type HandlerMiddleware ¶
HandlerMiddleware wraps a typed handler.
type HandlerPanicError ¶ added in v0.2.0
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 ¶
IDGenerator creates stable message identities.
func UUIDv7Generator ¶
func UUIDv7Generator() IDGenerator
UUIDv7Generator returns the default cryptographically random UUIDv7 generator.
type LivenessChecker ¶ added in v0.2.0
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 ¶
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) Readiness ¶
func (r *LocalAsyncRoute) Readiness(context.Context) error
Readiness verifies that this route is accepting work.
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.
type LogLevel ¶
type LogLevel uint8
LogLevel is the transport-neutral severity understood by Logger.
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.
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.
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 ¶
Message is the typed value passed to a handler.
func DecodeCommand ¶
DecodeCommand decodes and verifies a native command envelope.
type MessageID ¶
type MessageID [16]byte
MessageID is a UUID-compatible 128-bit message identity.
func ParseMessageID ¶
ParseMessageID parses the canonical UUID text form.
func (MessageID) MarshalJSON ¶
MarshalJSON implements json.Marshaler.
func (MessageID) MarshalText ¶
MarshalText implements encoding.TextMarshaler.
func (*MessageID) UnmarshalJSON ¶
UnmarshalJSON implements json.Unmarshaler.
func (*MessageID) UnmarshalText ¶
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) MarshalManifest ¶
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.
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.
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 ¶
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 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 ¶
WithLogger sets the core structured logger. The default logger is a no-op.
func WithObserver ¶
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
WithRuntimeShutdownTimeout sets the internal bound used when Run owns service shutdown after cancellation or an unexpected service return.
func WithSource ¶
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
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 ¶
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.
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 ¶
MustQuery constructs a typed local query descriptor and panics when its declaration is invalid.
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 ¶
WithSchema returns a query descriptor with an explicit request schema URI.
type QueryHandler ¶
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 ¶
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
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
Liveness checks that Runtime has not terminated and invokes optional service liveness checks without requiring readiness or topology access.
func (*Runtime) Readiness ¶
Readiness checks runtime admission state and each service's lightweight readiness contract. Expensive topology validation belongs in DeepHealth.
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.
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 ¶
ServiceProvider lets a route contribute one managed service to Builder.
Source Files
¶
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. |
|
observability
module
|
|
|
tools
|
|
|
gomessengerctl
module
|