loafer-awsx

module
v0.3.0 Latest Latest
Warning

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

Go to latest
Published: Jul 26, 2026 License: MIT

README

loafer-awsx

Go Reference CI Go Version License

A modern, idiomatic Go library for AWS SQS/SNS message processing, built on aws-sdk-go-v2 with generic type-safe handlers, a composable middleware pipeline, first-class log/slog logging, and built-in Prometheus and OpenTelemetry observability.

loafer-awsx organizes message processing into small, single-responsibility packages you can compose: build an AWS connection, declare routes, wrap them in a broker, and publish events with a producer. Everything is configured through functional options, and every component accepts the standard library *slog.Logger directly, no custom logger interface.

  • Module: github.com/silviolleite/loafer-awsx
  • Minimum Go version: Go 1.26 or later
  • AWS SDK: aws-sdk-go-v2 (SQS + SNS)

Table of Contents

  1. Architecture
  2. Installation
  3. Quickstart
  4. Examples
  5. Configuration Reference
  6. Acknowledgements
  7. License

Architecture

loafer-awsx is a library, not a service. Your application imports it, wires routes and a broker, and processes messages from AWS SQS while optionally publishing to AWS SNS. The library also exposes Prometheus metrics and OpenTelemetry spans for observability.

Context Diagram
flowchart LR
    dev([Developer])
    loafer[loafer-awsx]
    sqs[(AWS SQS)]
    sns[(AWS SNS)]
    prom[Prometheus]
    otel[OpenTelemetry]

    dev --> loafer
    loafer --> sqs
    loafer --> sns
    loafer --> prom
    loafer --> otel
Container Diagram

The broker orchestrates one consumer per route: each Route binds a queue to a handler, and its Consumer runs a worker pool that polls the matching SQS queue. Publishing runs alongside through the producer.

flowchart TB
    app([Application])
    broker[Broker]

    subgraph routeA[Route A]
        consumerA[Consumer / Workers]
    end
    subgraph routeB[Route B]
        consumerB[Consumer / Workers]
    end
    subgraph routeN[Route N]
        consumerN[Consumer / Workers]
    end

    sqsA[(SQS Queue A)]
    sqsB[(SQS Queue B)]
    sqsN[(SQS Queue N)]

    producer[Producer]
    sns[(AWS SNS)]

    app --> broker
    app --> producer

    broker --> routeA
    broker --> routeB
    broker --> routeN

    consumerA --> sqsA
    consumerB --> sqsB
    consumerN --> sqsN

    producer --> sns

Cross-cutting packages support this pipeline: conn builds the shared aws.Config, middleware wraps each route handler (global middleware outermost, route middleware innermost), typed adds generic type-safe handlers and producers, idgen generates FIFO IDs, and logger supplies the *slog.Logger used throughout.

Package responsibilities at a glance:

Package Responsibility
conn Factory for an aws.Config (region, credentials, endpoint, profile, retry).
logger Constructors for the standard library *slog.Logger (stdout + no-op).
middleware Handler, Middleware, Chain, and built-in Recovery, Logging, Metrics, OTel.
router Immutable Route value object binding a queue to a handler and options.
consumer SQS polling loop, worker-pool dispatch, visibility management, DLQ observability.
broker Lifecycle orchestrator that runs one consumer per route with graceful shutdown.
producer SNS single and batch publish for standard and FIFO topics.
typed Generic, type-safe handlers and producers via Codec[T].
idgen MessageGroupId / MessageDeduplicationId generation strategies.
errors Sentinel errors matchable with errors.Is.

Installation

Requires Go 1.26 or later.

go get github.com/silviolleite/loafer-awsx

Then import the packages you need, for example:

import (
    "github.com/silviolleite/loafer-awsx/broker"
    "github.com/silviolleite/loafer-awsx/conn"
    "github.com/silviolleite/loafer-awsx/logger"
    "github.com/silviolleite/loafer-awsx/router"
    "github.com/silviolleite/loafer-awsx/producer"
)

Quickstart

A typical setup builds a shared AWS config with conn.New, binds queue names to handlers with router.New, hands the routes to a broker, and calls broker.Run (which blocks until the context is canceled, then drains in-flight messages). Publishing works the same way: create a producer and call Publish or PublishBatch.

For complete, runnable programs covering the consumer, producer, FIFO, typed, and middleware setups, see the examples/ directory and its README.


Examples

Runnable, self-contained programs live in the examples/ directory, wired to run locally against LocalStack with infrastructure provisioned by Terraform. See examples/README.md for setup and run instructions (make up, make provision, make run-basic, and friends).

Example Directory What it shows
Basic examples/basic/ Standard SQS queue consumption with a simple handler.
FIFO examples/fifo/ Ordered consumption in PerGroupID mode with custom group fields.
Typed examples/typed/ Generic type-safe handling via typed.WrapHandler + typed.JSONCodec.
Middleware examples/middleware/ Recovery, logging, Prometheus metrics, and OpenTelemetry tracing.
Producer examples/producer/ Single and batch publishing to standard and FIFO SNS topics.

Configuration Reference

Every component is configured through functional options. Invalid values are rejected at construction time (returning a descriptive error) rather than being silently accepted.

conn — AWS configuration

conn.New(ctx context.Context, opts ...conn.Option) (aws.Config, error)

Option Signature Default Description
WithRegion WithRegion(region string) — (required) AWS region. New returns ErrEmptyRegion when empty.
WithAccessKey WithAccessKey(key, secret string) Static credentials; take precedence over a profile.
WithSessionToken WithSessionToken(token string) Session token; applied only with static credentials.
WithProfile WithProfile(profile string) Shared config profile name.
WithEndpoint WithEndpoint(url string) Custom endpoint URL (LocalStack, etc.).
WithRetryCount WithRetryCount(n uint) 10 Maximum retry attempts.
router — route configuration

router.New(queueName string, handler middleware.Handler, opts ...router.Option) (*router.Route, error)

Returns ErrEmptyQueueName for an empty name and ErrNoHandler for a nil handler. Option failures are wrapped with ErrInvalidOption.

Option Signature Default Description
WithWorkerPoolSize WithWorkerPoolSize(n int) 5 Worker goroutines per route (must be > 0).
WithMaxMessages WithMaxMessages(n int32) 10 Messages per SQS receive call (range [1, 10]).
WithWaitTimeSeconds WithWaitTimeSeconds(n int32) 10 Long-poll wait in seconds (range [0, 20]).
WithVisibilityTimeout WithVisibilityTimeout(seconds int32) 30 Visibility timeout; values <= 11 are clamped up to 11.
WithExtensionLimit WithExtensionLimit(n int) 2 Max visibility extensions (must not be negative).
WithRunMode WithRunMode(mode router.Mode) Parallel Dispatch strategy: Parallel or PerGroupID.
WithCustomGroupFields WithCustomGroupFields(fields ...string) Fields forming the group key for PerGroupID.
WithMiddleware WithMiddleware(mws ...middleware.Middleware) Route-level middleware appended in order.
WithDLQ WithDLQ(maxReceiveCount int, opts ...router.DLQOption) Enables DLQ observability (must be > 0). See below.

Run modes (router.Mode): Parallel (random worker assignment) and PerGroupID (hash of MessageGroupId + custom group fields, preserving per-group ordering).

DLQ observability (router.DLQOption):

Option Signature Description
WithOnDLQ WithOnDLQ(fn func(ctx context.Context, msg middleware.Message)) Callback invoked when a message is treated as exhausted.

DLQ is observe-only. WithDLQ does not take a target ARN and the library never moves, publishes, or deletes messages for DLQ purposes. AWS SQS performs the actual redrive natively via the source queue's redrive policy. maxReceiveCount must mirror that policy; it is used only to detect when a message is exhausted so the consumer can emit an Error log, the loafer_messages_dlq_total metric, and the optional OnDLQ callback while leaving the message in the queue.

consumer — SQS polling

consumer.New(client consumer.SQSClient, route *router.Route, opts ...consumer.Option) (*consumer.Consumer, error)

Returns ErrNoSQSClient for a nil client and ErrNoRoute for a nil route. In most applications the broker creates consumers for you; use these options directly only when driving a consumer yourself.

Option Signature Default Description
WithLogger WithLogger(log *slog.Logger) no-op Structured logger (nil is ignored).
WithRetryTimeout WithRetryTimeout(d time.Duration) 5s Wait after a failed ReceiveMessage (non-positive ignored).
WithGlobalMiddleware WithGlobalMiddleware(mws ...middleware.Middleware) Outermost middleware, ahead of route middleware.
WithDLQMetric WithDLQMetric(inc consumer.DLQMetric) Increments loafer_messages_dlq_total; wire only with Metrics enabled.
broker — lifecycle orchestration

broker.New(sqsClient consumer.SQSClient, routes []*router.Route, opts ...broker.Option) (*broker.Broker, error)

Returns ErrNoRoute when no routes are provided. broker.Run(ctx) starts one consumer per route, blocks until the context is canceled, and drains in-flight messages within the shutdown timeout with no goroutine leaks.

Option Signature Default Description
WithLogger WithLogger(log *slog.Logger) logger.New() Structured logger, forwarded to every consumer.
WithRetryTimeout WithRetryTimeout(d time.Duration) 5s Per-consumer wait after a failed receive.
WithShutdownTimeout WithShutdownTimeout(d time.Duration) unbounded Max wait for in-flight messages on shutdown. Unset waits until consumers finish; set a duration to bound it.
WithMiddleware WithMiddleware(mws ...middleware.Middleware) Global middleware applied outermost to all routes.

Middleware ordering: broker-level (global) middleware is applied outermost and route-level middleware innermost (closest to the handler).

producer — SNS publishing

producer.New(client producer.SNSClient, opts ...producer.Option) (*producer.Producer, error)

Returns ErrNoSNSClient for a nil client. Publish returns ErrEmptyInput for a nil/empty input; PublishBatch returns ErrEmptyInput for an empty batch and ErrMaxBatchSize for more than 10 entries.

Option Signature Description
WithGroupIDGenerator WithGroupIDGenerator(gen idgen.GroupIDGenerator) Auto-generate MessageGroupId for FIFO topics when not set.
WithDeduplicationIDGenerator WithDeduplicationIDGenerator(gen idgen.DeduplicationIDGenerator) Auto-generate MessageDeduplicationId for FIFO topics when not set.

Helpers: producer.BuildTopicARN(region, accountID, topicName string) string.

Auto-generation only applies to FIFO topics (ARNs ending in .fifo) and only when the corresponding ID is empty. Standard topics never receive auto-generated IDs, because SNS rejects them on non-FIFO topics.

typed — generic type-safe handlers
  • typed.Codec[T] — interface with Encode(T) ([]byte, error) and Decode([]byte) (T, error).
  • typed.JSONCodec[T] — JSON implementation of Codec[T].
  • typed.WrapHandler[T](codec Codec[T], fn func(ctx, msg T) error) middleware.Handler — adapts a typed handler into a standard Handler; a decode error is returned to the consumer.
  • typed.NewProducer[T](p *producer.Producer, codec Codec[T]) *typed.Producer[T] — a typed producer that encodes before publishing.
  • typed.Producer[T].Publish(ctx, topicARN, value T, opts ...typed.PublishOption) (string, error).
Publish option Signature Description
WithGroupID WithGroupID(id string) Sets MessageGroupId.
WithDeduplicationID WithDeduplicationID(id string) Sets MessageDeduplicationId.
WithAttributes WithAttributes(attrs map[string]string) Sets message attributes.
middleware package
  • middleware.Handlerfunc(ctx context.Context, msg middleware.Message) error.
  • middleware.Middlewarefunc(Handler) Handler.
  • middleware.Chain(mws ...Middleware) Middleware — composes middleware; the first is outermost.
  • middleware.Recovery(log *slog.Logger) Middleware — recovers panics, logs the stack, returns ErrPanic.
  • middleware.Logging(log *slog.Logger) Middleware — logs receipt, duration, and outcome.
  • middleware.Metrics(routeName string, opts ...MetricsOption) Middleware — Prometheus counters, histogram, and inflight gauge.
  • middleware.OTel(routeName string, opts ...OTelOption) Middleware — an OpenTelemetry span per message.
Option Signature Default Description
WithMetricsRegisterer WithMetricsRegisterer(r prometheus.Registerer) prometheus.DefaultRegisterer Custom Prometheus registerer.
WithTracerProvider WithTracerProvider(tp trace.TracerProvider) global provider Custom OpenTelemetry tracer provider.

Metrics emitted: loafer_messages_received_total, loafer_messages_processed_total (labeled by status), loafer_messages_errors_total, loafer_message_processing_duration_seconds (histogram), loafer_messages_inflight (gauge), and loafer_messages_dlq_total — all labeled by route.

logger — standard library slog constructors
  • logger.New() *slog.Logger — structured, leveled output to stdout via slog.TextHandler.
  • logger.NewNoOp() *slog.Logger — a discard-backed logger (silent).

The library uses *slog.Logger everywhere and defines no custom logger interface. A *slog.Logger produced by a third-party bridge (for example zap via zapslog, or zerolog via a slog handler) is accepted directly, no adapter required.

idgen — ID generation

Interfaces idgen.GroupIDGenerator and idgen.DeduplicationIDGenerator both expose Generate(ctx context.Context, fields map[string]string) (string, error). A single concrete generator satisfies both.

Constructor Description
NewKeyBased(opts ...Option) Deterministic ID from sorted field values, hashed with the configured algorithm. Returns ErrEmptyFields when no field is selected.
NewRandom() Random UUID v4 on every call; ignores fields.
NewComposite(opts ...Option) Joins selected field values with a separator (no hashing).
NewCompositeWithSuffix(opts ...Option) Like NewComposite, plus a random numeric suffix from the configured range to spread load across partitions.
Option Signature Default Description
WithHashAlgorithm WithHashAlgorithm(algorithm idgen.HashAlgorithm) SHA256 Digest for key-based hashing: SHA256 or FNV64.
WithSeparator WithSeparator(separator string) ":" Separator joining key/value pairs.
WithFields WithFields(fields ...string) all fields Whitelist of fields to include.
WithSuffixRange WithSuffixRange(min, max int) [1, 20] Inclusive suffix range for NewCompositeWithSuffix.
errors — sentinel errors

The errors package exports sentinels matchable with errors.Is, including ErrNoRoute, ErrNoSQSClient, ErrNoHandler, ErrGetMessage, ErrQueueResolve, ErrNoSNSClient, ErrEmptyInput, ErrMaxBatchSize, ErrEmptyRegion, ErrEmptyQueueName, ErrInvalidOption, ErrEmptyFields, and ErrPanic. errors.New(text string) error and errors.Wrap(sentinel, err error) error help build and combine errors while preserving errors.Is matching against both causes.


Acknowledgements

This project was inspired by JustCodes/loafer-go.


License

See LICENSE.

Directories

Path Synopsis
Package broker provides the top-level orchestrator that creates and manages Consumer instances for multiple routes, offering coordinated startup, graceful shutdown, and fail-fast behavior.
Package broker provides the top-level orchestrator that creates and manages Consumer instances for multiple routes, offering coordinated startup, graceful shutdown, and fail-fast behavior.
Package conn provides a factory for AWS SDK v2 configuration.
Package conn provides a factory for AWS SDK v2 configuration.
Package consumer implements the SQS polling loop, worker-pool dispatch, visibility-timeout management, and the message commit/backoff lifecycle for a single queue.
Package consumer implements the SQS polling loop, worker-pool dispatch, visibility-timeout management, and the message commit/backoff lifecycle for a single queue.
Package errors defines the sentinel errors used across loafer-go v3 and a Wrap helper that preserves errors.Is matchability.
Package errors defines the sentinel errors used across loafer-go v3 and a Wrap helper that preserves errors.Is matchability.
examples
basic command
Command basic demonstrates standard SQS queue consumption with loafer-go v3.
Command basic demonstrates standard SQS queue consumption with loafer-go v3.
fifo command
Command fifo demonstrates ordered consumption from an SQS FIFO queue with loafer-go v3.
Command fifo demonstrates ordered consumption from an SQS FIFO queue with loafer-go v3.
middleware command
Command middleware demonstrates observability middleware with loafer-go v3.
Command middleware demonstrates observability middleware with loafer-go v3.
producer command
Command producer demonstrates publishing messages to SNS topics with loafer-go v3.
Command producer demonstrates publishing messages to SNS topics with loafer-go v3.
typed command
Command typed demonstrates strongly-typed message handling with loafer-go v3.
Command typed demonstrates strongly-typed message handling with loafer-go v3.
Package fake provides configurable test doubles for the core loafer-go v3 interfaces used across package tests: Message (consumer.Message and middleware.Message), SQSClient (consumer.SQSClient), and SNSClient (producer.SNSClient).
Package fake provides configurable test doubles for the core loafer-go v3 interfaces used across package tests: Message (consumer.Message and middleware.Message), SQSClient (consumer.SQSClient), and SNSClient (producer.SNSClient).
Package idgen generates MessageGroupId and MessageDeduplicationId values using deterministic (key-based, composite) or random (UUID) strategies.
Package idgen generates MessageGroupId and MessageDeduplicationId values using deterministic (key-based, composite) or random (UUID) strategies.
Package logger provides constructors for the standard library *slog.Logger used throughout loafer-go v3.
Package logger provides constructors for the standard library *slog.Logger used throughout loafer-go v3.
Package middleware defines the Handler and Middleware types, the Chain combinator, and the built-in middlewares (Recovery, Logging, Metrics, and OpenTelemetry) used to add cross-cutting concerns to message processing.
Package middleware defines the Handler and Middleware types, the Chain combinator, and the built-in middlewares (Recovery, Logging, Metrics, and OpenTelemetry) used to add cross-cutting concerns to message processing.
Package producer publishes messages to AWS SNS topics.
Package producer publishes messages to AWS SNS topics.
Package router defines a Route as the binding between a queue name, a handler, a middleware chain, and route-level configuration.
Package router defines a Route as the binding between a queue name, a handler, a middleware chain, and route-level configuration.
Package typed provides generic, type-safe handlers and producers built on a Codec interface, eliminating manual JSON marshaling boilerplate.
Package typed provides generic, type-safe handlers and producers built on a Codec interface, eliminating manual JSON marshaling boilerplate.

Jump to

Keyboard shortcuts

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