From e4476d60cac6e157fb0d84cd8cd646d4295a59ae Mon Sep 17 00:00:00 2001 From: Jonas Tranberg Date: Wed, 1 Jul 2026 10:38:43 +0200 Subject: [PATCH 1/6] Add E2E tests with testcontainers and migrate to amqp091-go - Add go-testcontainers E2E tests covering SimpleQueue, DeadLetterQueue, BoundedRetryQueue and BackoffQueue against a real RabbitMQ container - Replace github.com/streadway/amqp (maintenance mode) with github.com/rabbitmq/amqp091-go v1.12.0 - Replace github.com/isayme/go-amqp-reconnect with a new internal/reconnect package that provides the same auto-reconnect semantics (Connection + Channel wrappers) on top of amqp091-go - All five E2E tests pass both before and after the migration Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- client/client.go | 11 ++- e2e/backoff_queue_test.go | 66 ++++++++++++++ e2e/bounded_retry_queue_test.go | 70 +++++++++++++++ e2e/dead_letter_queue_test.go | 46 ++++++++++ e2e/helpers_test.go | 57 ++++++++++++ e2e/simple_queue_test.go | 80 +++++++++++++++++ exchange/exchange.go | 10 +-- go.mod | 60 ++++++++++++- go.sum | 153 +++++++++++++++++++++++++++++--- internal/reconnect/reconnect.go | 131 +++++++++++++++++++++++++++ queue/backoff_queue.go | 8 +- queue/bounded_retry_queue.go | 6 +- queue/dead_letter_queue.go | 6 +- queue/queue.go | 6 +- queue/simple_queue.go | 4 +- 15 files changed, 675 insertions(+), 39 deletions(-) create mode 100644 e2e/backoff_queue_test.go create mode 100644 e2e/bounded_retry_queue_test.go create mode 100644 e2e/dead_letter_queue_test.go create mode 100644 e2e/helpers_test.go create mode 100644 e2e/simple_queue_test.go create mode 100644 internal/reconnect/reconnect.go diff --git a/client/client.go b/client/client.go index 06a628a..fd90eef 100644 --- a/client/client.go +++ b/client/client.go @@ -2,23 +2,22 @@ package client import ( "fmt" - // This package provides auto-reconnect - rmq "github.com/isayme/go-amqp-reconnect/rabbitmq" "github.com/pkg/errors" "github.com/uniwise/go-rabbit/exchange" + "github.com/uniwise/go-rabbit/internal/reconnect" ) // RabbitMQClient is the interface describing a RabbitMQ wrapper type RabbitMQClient interface { - Channel() (*rmq.Channel, error) + Channel() (*reconnect.Channel, error) NewExchange(name string) (*exchange.Exchange, error) } // RabbitMQ is a wrapper struct for a RabbitMQ connection type RabbitMQ struct { Config *Config - Connection *rmq.Connection + Connection *reconnect.Connection } // New is the constructor for RabbitMQImpl @@ -44,7 +43,7 @@ func (r *RabbitMQ) connect() error { r.Config.VHost, ) - conn, err := rmq.Dial(connStr) + conn, err := reconnect.Dial(connStr) if err != nil { return err } @@ -55,7 +54,7 @@ func (r *RabbitMQ) connect() error { } // Channel returns a RabbitMQ channel from the connection -func (r *RabbitMQ) Channel() (*rmq.Channel, error) { +func (r *RabbitMQ) Channel() (*reconnect.Channel, error) { ch, err := r.Connection.Channel() if err != nil { return nil, err diff --git a/e2e/backoff_queue_test.go b/e2e/backoff_queue_test.go new file mode 100644 index 0000000..f9f0e49 --- /dev/null +++ b/e2e/backoff_queue_test.go @@ -0,0 +1,66 @@ +package e2e_test + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/uniwise/go-rabbit/queue" +) + +// TestBackoffQueue verifies that messages are routed through the staged backoff +// queues and eventually return to the target queue, and that ErrBackoffExhausted +// is returned once all stages are consumed. +func TestBackoffQueue(t *testing.T) { + intervals := []time.Duration{ + 300 * time.Millisecond, + 500 * time.Millisecond, + } + + rmq := newTestClient(t) + + ex, err := rmq.NewExchange(uniqueName(t, "exchange")) + require.NoError(t, err) + + targetQ, err := ex.NewQueue(uniqueName(t, "target"), 10) + require.NoError(t, err) + + backoffQ, err := ex.NewBackoffQueue(uniqueName(t, "backoff"), intervals, targetQ) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + deliveries, err := targetQ.Consume(ctx) + require.NoError(t, err) + + require.NoError(t, targetQ.Publish([]byte("backoff-test"))) + + // Stage 0 (300 ms): first backoff, message should return after ~300 ms. + d1 := <-deliveries + assert.Equal(t, []byte("backoff-test"), d1.Body) + require.NoError(t, backoffQ.Publish(d1)) + require.NoError(t, d1.Ack(false)) + + // Stage 1 (500 ms): second backoff, message should return after ~500 ms. + select { + case d2 := <-deliveries: + assert.Equal(t, []byte("backoff-test"), d2.Body) + require.NoError(t, backoffQ.Publish(d2)) + require.NoError(t, d2.Ack(false)) + case <-ctx.Done(): + t.Fatal("timed out waiting for stage-1 backoff delivery") + } + + // All stages exhausted: the next Publish should return ErrBackoffExhausted. + select { + case d3 := <-deliveries: + err = backoffQ.Publish(d3) + assert.ErrorIs(t, err, queue.ErrBackoffExhausted) + require.NoError(t, d3.Ack(false)) + case <-ctx.Done(): + t.Fatal("timed out waiting for final delivery after backoff exhausted") + } +} diff --git a/e2e/bounded_retry_queue_test.go b/e2e/bounded_retry_queue_test.go new file mode 100644 index 0000000..83e7b0e --- /dev/null +++ b/e2e/bounded_retry_queue_test.go @@ -0,0 +1,70 @@ +package e2e_test + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/uniwise/go-rabbit/queue" +) + +// TestBoundedRetryQueue verifies the full retry lifecycle: +// 1. A message can be re-queued up to MaxRetries times. +// 2. On the (MaxRetries+1)th Publish call ErrMaxRetriesReached is returned. +func TestBoundedRetryQueue(t *testing.T) { + const ( + maxRetries = 2 + retryDelay = 500 * time.Millisecond + ) + + rmq := newTestClient(t) + + ex, err := rmq.NewExchange(uniqueName(t, "exchange")) + require.NoError(t, err) + + targetQ, err := ex.NewQueue(uniqueName(t, "target"), 10) + require.NoError(t, err) + + retryQ, err := ex.NewBoundedRetryQueue( + uniqueName(t, "retry"), + 10, + maxRetries, + retryDelay, + targetQ, + ) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + deliveries, err := targetQ.Consume(ctx) + require.NoError(t, err) + + // Seed the first message. + require.NoError(t, targetQ.Publish([]byte("retry-test"))) + + // Retry maxRetries times; each time the message should reappear on targetQ. + for i := 0; i < maxRetries; i++ { + select { + case d := <-deliveries: + assert.Equal(t, []byte("retry-test"), d.Body) + err = retryQ.Publish(d) + require.NoError(t, err, "publish attempt %d should succeed", i+1) + require.NoError(t, d.Ack(false)) + case <-ctx.Done(): + t.Fatalf("timed out on retry attempt %d", i+1) + } + } + + // One more delivery: now MaxRetries is exhausted. + select { + case d := <-deliveries: + err = retryQ.Publish(d) + assert.ErrorIs(t, err, queue.ErrMaxRetriesReached) + require.NoError(t, d.Ack(false)) + case <-ctx.Done(): + t.Fatal("timed out waiting for final delivery after max retries exhausted") + } +} diff --git a/e2e/dead_letter_queue_test.go b/e2e/dead_letter_queue_test.go new file mode 100644 index 0000000..eda0905 --- /dev/null +++ b/e2e/dead_letter_queue_test.go @@ -0,0 +1,46 @@ +package e2e_test + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// TestDeadLetterQueue verifies that a message published to a dead-letter queue +// is automatically re-routed to the target queue after its TTL expires. +func TestDeadLetterQueue(t *testing.T) { + const ttl = 500 * time.Millisecond + + rmq := newTestClient(t) + + ex, err := rmq.NewExchange(uniqueName(t, "exchange")) + require.NoError(t, err) + + targetQ, err := ex.NewQueue(uniqueName(t, "target"), 10) + require.NoError(t, err) + + dlq, err := ex.NewDeadLetterQueue(uniqueName(t, "dlq"), 10, ttl, targetQ) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + // Start consuming from the target queue before publishing. + deliveries, err := targetQ.Consume(ctx) + require.NoError(t, err) + + // Publish to the DLQ; it will expire and be forwarded to targetQ. + msg := []byte("dead-letter-test") + require.NoError(t, dlq.Publish(msg)) + + select { + case d := <-deliveries: + assert.Equal(t, msg, d.Body) + require.NoError(t, d.Ack(false)) + case <-ctx.Done(): + t.Fatal("timed out waiting for dead-lettered message on target queue") + } +} diff --git a/e2e/helpers_test.go b/e2e/helpers_test.go new file mode 100644 index 0000000..60755df --- /dev/null +++ b/e2e/helpers_test.go @@ -0,0 +1,57 @@ +package e2e_test + +import ( + "context" + "fmt" + "net/url" + "strconv" + "testing" + + "github.com/stretchr/testify/require" + tc_rabbitmq "github.com/testcontainers/testcontainers-go/modules/rabbitmq" + rabbit "github.com/uniwise/go-rabbit" + "github.com/uniwise/go-rabbit/client" +) + +// newTestClient spins up a RabbitMQ container and returns a connected client. +// The container is terminated when the test ends. +func newTestClient(t *testing.T) client.RabbitMQClient { + t.Helper() + ctx := context.Background() + + container, err := tc_rabbitmq.Run(ctx, "rabbitmq:3-alpine") + require.NoError(t, err) + t.Cleanup(func() { _ = container.Terminate(ctx) }) + + amqpURL, err := container.AmqpURL(ctx) + require.NoError(t, err) + + u, err := url.Parse(amqpURL) + require.NoError(t, err) + + portNum, err := strconv.ParseUint(u.Port(), 10, 32) + require.NoError(t, err) + + password, _ := u.User.Password() + // Strip leading "/" from path; empty string means default vhost + vhost := u.Path + if len(vhost) > 0 && vhost[0] == '/' { + vhost = vhost[1:] + } + + c, err := rabbit.New(&client.Config{ + Host: u.Hostname(), + Port: uint32(portNum), + User: u.User.Username(), + Password: password, + VHost: vhost, + }) + require.NoError(t, err, "failed to connect to RabbitMQ") + + return c +} + +// uniqueName generates a unique queue/exchange name for a test to avoid collisions. +func uniqueName(t *testing.T, suffix string) string { + return fmt.Sprintf("%s-%s", t.Name(), suffix) +} diff --git a/e2e/simple_queue_test.go b/e2e/simple_queue_test.go new file mode 100644 index 0000000..9e60925 --- /dev/null +++ b/e2e/simple_queue_test.go @@ -0,0 +1,80 @@ +package e2e_test + +import ( + "context" + "testing" + "time" + + amqp "github.com/rabbitmq/amqp091-go" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestSimpleQueue_PublishConsume(t *testing.T) { + rmq := newTestClient(t) + + ex, err := rmq.NewExchange(uniqueName(t, "exchange")) + require.NoError(t, err) + + q, err := ex.NewQueue(uniqueName(t, "queue"), 10) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + deliveries, err := q.Consume(ctx) + require.NoError(t, err) + + messages := [][]byte{ + []byte("hello"), + []byte("world"), + []byte("foo"), + } + + for _, msg := range messages { + require.NoError(t, q.Publish(msg)) + } + + received := make([][]byte, 0, len(messages)) + for range messages { + select { + case d := <-deliveries: + received = append(received, d.Body) + require.NoError(t, d.Ack(false)) + case <-ctx.Done(): + t.Fatalf("timed out waiting for message, received %d/%d", len(received), len(messages)) + } + } + + assert.Equal(t, messages, received) +} + +func TestSimpleQueue_ConsumeFunc(t *testing.T) { + rmq := newTestClient(t) + + ex, err := rmq.NewExchange(uniqueName(t, "exchange")) + require.NoError(t, err) + + q, err := ex.NewQueue(uniqueName(t, "queue"), 10) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + received := make(chan []byte, 1) + err = q.ConsumeFunc(ctx, func(d amqp.Delivery) { + received <- d.Body + _ = d.Ack(false) + }) + require.NoError(t, err) + + msg := []byte("consume-func-test") + require.NoError(t, q.Publish(msg)) + + select { + case body := <-received: + assert.Equal(t, msg, body) + case <-ctx.Done(): + t.Fatal("timed out waiting for message") + } +} diff --git a/exchange/exchange.go b/exchange/exchange.go index c2c478b..f0f4d12 100644 --- a/exchange/exchange.go +++ b/exchange/exchange.go @@ -3,9 +3,9 @@ package exchange import ( "time" - "github.com/isayme/go-amqp-reconnect/rabbitmq" + amqp "github.com/rabbitmq/amqp091-go" "github.com/pkg/errors" - "github.com/streadway/amqp" + "github.com/uniwise/go-rabbit/internal/reconnect" "github.com/uniwise/go-rabbit/queue" ) @@ -20,13 +20,13 @@ type Exchanger interface { // Exchange is a wrapper for RabbitMQ exchanges type Exchange struct { ExchangeName string - Connection *rabbitmq.Connection - Channel *rabbitmq.Channel + Connection *reconnect.Connection + Channel *reconnect.Channel } // Config is the configuration which the constructor NewExchange needs type Config struct { - Connection *rabbitmq.Connection + Connection *reconnect.Connection ExchangeName string } diff --git a/go.mod b/go.mod index 615a520..24eb04e 100644 --- a/go.mod +++ b/go.mod @@ -1,11 +1,65 @@ module github.com/uniwise/go-rabbit -go 1.23 +go 1.25.0 require ( - github.com/isayme/go-amqp-reconnect v0.0.0-20210303120416-fc811b0bcda2 github.com/joho/godotenv v1.5.1 github.com/kelseyhightower/envconfig v1.4.0 github.com/pkg/errors v0.9.1 - github.com/streadway/amqp v1.1.0 + github.com/rabbitmq/amqp091-go v1.12.0 + github.com/stretchr/testify v1.11.1 + github.com/testcontainers/testcontainers-go/modules/rabbitmq v0.43.0 +) + +require ( + dario.cat/mergo v1.0.2 // indirect + github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c // indirect + github.com/Microsoft/go-winio v0.6.2 // indirect + github.com/cenkalti/backoff/v4 v4.3.0 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/containerd/errdefs v1.0.0 // indirect + github.com/containerd/errdefs/pkg v0.3.0 // indirect + github.com/containerd/log v0.1.0 // indirect + github.com/containerd/platforms v0.2.1 // indirect + github.com/cpuguy83/dockercfg v0.3.2 // indirect + github.com/davecgh/go-spew v1.1.1 // indirect + github.com/distribution/reference v0.6.0 // indirect + github.com/docker/go-connections v0.6.0 // indirect + github.com/docker/go-units v0.5.0 // indirect + github.com/ebitengine/purego v0.10.0 // indirect + github.com/felixge/httpsnoop v1.0.4 // indirect + github.com/go-logr/logr v1.4.3 // indirect + github.com/go-logr/stdr v1.2.2 // indirect + github.com/go-ole/go-ole v1.2.6 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/klauspost/compress v1.18.5 // indirect + github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 // indirect + github.com/magiconair/properties v1.8.10 // indirect + github.com/moby/docker-image-spec v1.3.1 // indirect + github.com/moby/go-archive v0.2.0 // indirect + github.com/moby/moby/api v1.54.2 // indirect + github.com/moby/moby/client v0.4.0 // indirect + github.com/moby/patternmatcher v0.6.1 // indirect + github.com/moby/sys/sequential v0.6.0 // indirect + github.com/moby/sys/user v0.4.0 // indirect + github.com/moby/sys/userns v0.1.0 // indirect + github.com/moby/term v0.5.2 // indirect + github.com/opencontainers/go-digest v1.0.0 // indirect + github.com/opencontainers/image-spec v1.1.1 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect + github.com/shirou/gopsutil/v4 v4.26.5 // indirect + github.com/sirupsen/logrus v1.9.4 // indirect + github.com/testcontainers/testcontainers-go v0.43.0 // indirect + github.com/tklauser/go-sysconf v0.3.16 // indirect + github.com/tklauser/numcpus v0.11.0 // indirect + github.com/yusufpapurcu/wmi v1.2.4 // indirect + go.opentelemetry.io/auto/sdk v1.2.1 // indirect + go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.60.0 // indirect + go.opentelemetry.io/otel v1.41.0 // indirect + go.opentelemetry.io/otel/metric v1.41.0 // indirect + go.opentelemetry.io/otel/trace v1.41.0 // indirect + golang.org/x/crypto v0.51.0 // indirect + golang.org/x/sys v0.45.0 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index a98b650..e7d1a5f 100644 --- a/go.sum +++ b/go.sum @@ -1,16 +1,149 @@ -github.com/isayme/go-amqp-reconnect v0.0.0-20180930040740-e71660afb5ca h1:/k6bi3UzEon87YXwewlKXiMT+Qxu+OCaOtkn5ZGgCOs= -github.com/isayme/go-amqp-reconnect v0.0.0-20180930040740-e71660afb5ca/go.mod h1:4IOu90sBxNtO7GtD9//Ybh2UjZ9Dl+Cd9yIMj9GPRHQ= -github.com/isayme/go-amqp-reconnect v0.0.0-20210303120416-fc811b0bcda2 h1:PzQ5MrrM7f/PHpC0aN9hZA+nBDEuBQRX0EhQxc4W9OA= -github.com/isayme/go-amqp-reconnect v0.0.0-20210303120416-fc811b0bcda2/go.mod h1:4IOu90sBxNtO7GtD9//Ybh2UjZ9Dl+Cd9yIMj9GPRHQ= -github.com/joho/godotenv v1.3.0 h1:Zjp+RcGpHhGlrMbJzXTrZZPrWj+1vfm90La1wgB6Bhc= -github.com/joho/godotenv v1.3.0/go.mod h1:7hK45KPybAkOC6peb+G5yklZfMxEjkZhHbwpqxOKXbg= +dario.cat/mergo v1.0.2 h1:85+piFYR1tMbRrLcDwR18y4UKJ3aH1Tbzi24VRW1TK8= +dario.cat/mergo v1.0.2/go.mod h1:E/hbnu0NxMFBjpMIE34DRGLWqDy0g5FuKDhCb31ngxA= +github.com/AdaLogics/go-fuzz-headers v0.0.0-20240806141605-e8a1dd7889d6 h1:He8afgbRMd7mFxO99hRNu+6tazq8nFF9lIwo9JFroBk= +github.com/AdaLogics/go-fuzz-headers v0.0.0-20240806141605-e8a1dd7889d6/go.mod h1:8o94RPi1/7XTJvwPpRSzSUedZrtlirdB3r9Z20bi2f8= +github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c h1:udKWzYgxTojEKWjV8V+WSxDXJ4NFATAsZjh8iIbsQIg= +github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c/go.mod h1:xomTg63KZ2rFqZQzSB4Vz2SUXa1BpHTVz9L5PTmPC4E= +github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= +github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= +github.com/cenkalti/backoff/v4 v4.3.0 h1:MyRJ/UdXutAwSAT+s3wNd7MfTIcy71VQueUuFK343L8= +github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyYozVcomhLiZE= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/containerd/errdefs v1.0.0 h1:tg5yIfIlQIrxYtu9ajqY42W3lpS19XqdxRQeEwYG8PI= +github.com/containerd/errdefs v1.0.0/go.mod h1:+YBYIdtsnF4Iw6nWZhJcqGSg/dwvV7tyJ/kCkyJ2k+M= +github.com/containerd/errdefs/pkg v0.3.0 h1:9IKJ06FvyNlexW690DXuQNx2KA2cUJXx151Xdx3ZPPE= +github.com/containerd/errdefs/pkg v0.3.0/go.mod h1:NJw6s9HwNuRhnjJhM7pylWwMyAkmCQvQ4GpJHEqRLVk= +github.com/containerd/log v0.1.0 h1:TCJt7ioM2cr/tfR8GPbGf9/VRAX8D2B4PjzCpfX540I= +github.com/containerd/log v0.1.0/go.mod h1:VRRf09a7mHDIRezVKTRCrOq78v577GXq3bSa3EhrzVo= +github.com/containerd/platforms v0.2.1 h1:zvwtM3rz2YHPQsF2CHYM8+KtB5dvhISiXh5ZpSBQv6A= +github.com/containerd/platforms v0.2.1/go.mod h1:XHCb+2/hzowdiut9rkudds9bE5yJ7npe7dG/wG+uFPw= +github.com/cpuguy83/dockercfg v0.3.2 h1:DlJTyZGBDlXqUZ2Dk2Q3xHs/FtnooJJVaad2S9GKorA= +github.com/cpuguy83/dockercfg v0.3.2/go.mod h1:sugsbF4//dDlL/i+S+rtpIWp+5h0BHJHfjj5/jFyUJc= +github.com/creack/pty v1.1.24 h1:bJrF4RRfyJnbTJqzRLHzcGaZK1NeM5kTC9jGgovnR1s= +github.com/creack/pty v1.1.24/go.mod h1:08sCNb52WyoAwi2QDyzUCTgcvVFhUzewun7wtTfvcwE= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/distribution/reference v0.6.0 h1:0IXCQ5g4/QMHHkarYzh5l+u8T3t73zM5QvfrDyIgxBk= +github.com/distribution/reference v0.6.0/go.mod h1:BbU0aIcezP1/5jX/8MP0YiH4SdvB5Y4f/wlDRiLyi3E= +github.com/docker/go-connections v0.6.0 h1:LlMG9azAe1TqfR7sO+NJttz1gy6KO7VJBh+pMmjSD94= +github.com/docker/go-connections v0.6.0/go.mod h1:AahvXYshr6JgfUJGdDCs2b5EZG/vmaMAntpSFH5BFKE= +github.com/docker/go-units v0.5.0 h1:69rxXcBk27SvSaaxTtLh/8llcHD8vYHT7WSdRZ/jvr4= +github.com/docker/go-units v0.5.0/go.mod h1:fgPhTUdO+D/Jk86RDLlptpiXQzgHJF7gydDDbaIK4Dk= +github.com/ebitengine/purego v0.10.0 h1:QIw4xfpWT6GWTzaW5XEKy3HXoqrJGx1ijYHzTF0/ISU= +github.com/ebitengine/purego v0.10.0/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ= +github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= +github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= +github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/go-ole/go-ole v1.2.6 h1:/Fpf6oFPoeFik9ty7siob0G6Ke8QvQEuVcuChpwXzpY= +github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0= +github.com/google/go-cmp v0.5.6/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0= github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4= github.com/kelseyhightower/envconfig v1.4.0 h1:Im6hONhd3pLkfDFsbRgu68RDNkGF1r3dvMUtDTo2cv8= github.com/kelseyhightower/envconfig v1.4.0/go.mod h1:cccZRl6mQpaq41TPp5QxidR+Sa3axMbJDNb//FQX6Gg= +github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE= +github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 h1:6E+4a0GO5zZEnZ81pIr0yLvtUWk2if982qA3F3QD6H4= +github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0/go.mod h1:zJYVVT2jmtg6P3p1VtQj7WsuWi/y4VnjVBn7F8KPB3I= +github.com/magiconair/properties v1.8.10 h1:s31yESBquKXCV9a/ScB3ESkOjUYYv+X0rg8SYxI99mE= +github.com/magiconair/properties v1.8.10/go.mod h1:Dhd985XPs7jluiymwWYZ0G4Z61jb3vdS329zhj2hYo0= +github.com/mdelapenya/tlscert v0.1.0 h1:YTpF579PYUX475eOL+6zyEO3ngLTOUWck78NBuJVXaM= +github.com/mdelapenya/tlscert v0.1.0/go.mod h1:wrbyM/DwbFCeCeqdPX/8c6hNOqQgbf0rUDErE1uD+64= +github.com/moby/docker-image-spec v1.3.1 h1:jMKff3w6PgbfSa69GfNg+zN/XLhfXJGnEx3Nl2EsFP0= +github.com/moby/docker-image-spec v1.3.1/go.mod h1:eKmb5VW8vQEh/BAr2yvVNvuiJuY6UIocYsFu/DxxRpo= +github.com/moby/go-archive v0.2.0 h1:zg5QDUM2mi0JIM9fdQZWC7U8+2ZfixfTYoHL7rWUcP8= +github.com/moby/go-archive v0.2.0/go.mod h1:mNeivT14o8xU+5q1YnNrkQVpK+dnNe/K6fHqnTg4qPU= +github.com/moby/moby/api v1.54.2 h1:wiat9QAhnDQjA7wk1kh/TqHz2I1uUA7M7t9SAl/JNXg= +github.com/moby/moby/api v1.54.2/go.mod h1:+RQ6wluLwtYaTd1WnPLykIDPekkuyD/ROWQClE83pzs= +github.com/moby/moby/client v0.4.0 h1:S+2XegzHQrrvTCvF6s5HFzcrywWQmuVnhOXe2kiWjIw= +github.com/moby/moby/client v0.4.0/go.mod h1:QWPbvWchQbxBNdaLSpoKpCdf5E+WxFAgNHogCWDoa7g= +github.com/moby/patternmatcher v0.6.1 h1:qlhtafmr6kgMIJjKJMDmMWq7WLkKIo23hsrpR3x084U= +github.com/moby/patternmatcher v0.6.1/go.mod h1:hDPoyOpDY7OrrMDLaYoY3hf52gNCR/YOUYxkhApJIxc= +github.com/moby/sys/sequential v0.6.0 h1:qrx7XFUd/5DxtqcoH1h438hF5TmOvzC/lspjy7zgvCU= +github.com/moby/sys/sequential v0.6.0/go.mod h1:uyv8EUTrca5PnDsdMGXhZe6CCe8U/UiTWd+lL+7b/Ko= +github.com/moby/sys/user v0.4.0 h1:jhcMKit7SA80hivmFJcbB1vqmw//wU61Zdui2eQXuMs= +github.com/moby/sys/user v0.4.0/go.mod h1:bG+tYYYJgaMtRKgEmuueC0hJEAZWwtIbZTB+85uoHjs= +github.com/moby/sys/userns v0.1.0 h1:tVLXkFOxVu9A64/yh59slHVv9ahO9UIev4JZusOLG/g= +github.com/moby/sys/userns v0.1.0/go.mod h1:IHUYgu/kao6N8YZlp9Cf444ySSvCmDlmzUcYfDHOl28= +github.com/moby/term v0.5.2 h1:6qk3FJAFDs6i/q3W/pQ97SX192qKfZgGjCQqfCJkgzQ= +github.com/moby/term v0.5.2/go.mod h1:d3djjFCrjnB+fl8NJux+EJzu0msscUP+f8it8hPkFLc= +github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= +github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= +github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040= +github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= -github.com/streadway/amqp v1.0.0 h1:kuuDrUJFZL1QYL9hUNuCxNObNzB0bV/ZG5jV3RWAQgo= -github.com/streadway/amqp v1.0.0/go.mod h1:AZpEONHx3DKn8O/DFsRAY58/XVQiIPMTMB1SddzLXVw= -github.com/streadway/amqp v1.1.0 h1:py12iX8XSyI7aN/3dUT8DFIDJazNJsVJdxNVEpnQTZM= -github.com/streadway/amqp v1.1.0/go.mod h1:WYSrTEYHOXHd0nwFeUXAe2G2hRnQT+deZJJf88uS9Bg= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 h1:o4JXh1EVt9k/+g42oCprj/FisM4qX9L3sZB3upGN2ZU= +github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55/go.mod h1:OmDBASR4679mdNQnz2pUhc2G8CO2JrUAVFDRBDP/hJE= +github.com/rabbitmq/amqp091-go v1.12.0 h1:V0v14Iqfs+MwHWihJt/nGS5Ulu0vw572b2Co3mwunkI= +github.com/rabbitmq/amqp091-go v1.12.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= +github.com/shirou/gopsutil/v4 v4.26.5 h1:RPcBXkpz7kOj9PqGFQOlBPZHsyaPvPVQc098y9RmCNM= +github.com/shirou/gopsutil/v4 v4.26.5/go.mod h1:LZ6ewCSkBqUpvSOf+LsTGnRinC6iaNUNMGBtDkJBaLQ= +github.com/sirupsen/logrus v1.9.4 h1:TsZE7l11zFCLZnZ+teH4Umoq5BhEIfIzfRDZ1Uzql2w= +github.com/sirupsen/logrus v1.9.4/go.mod h1:ftWc9WdOfJ0a92nsE2jF5u5ZwH8Bv2zdeOC42RjbV2g= +github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4= +github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/testcontainers/testcontainers-go v0.43.0 h1:oEQx5MW2DGd9z3AeEQfB2lPM0eLs7ztyaGRu75bFo5A= +github.com/testcontainers/testcontainers-go v0.43.0/go.mod h1:+VxkT2NQnKOZPKi6praMuMKYHYyOGXr0XSBSlSMCzFo= +github.com/testcontainers/testcontainers-go/modules/rabbitmq v0.43.0 h1:+4P4pktUqT8nliEQ2ukrjA+t5j36PLGugoXm9zWMYgM= +github.com/testcontainers/testcontainers-go/modules/rabbitmq v0.43.0/go.mod h1:Q4SHA2NVopiqTDR5PGoapozEzWalEgpJOrtNSOG2RUM= +github.com/tklauser/go-sysconf v0.3.16 h1:frioLaCQSsF5Cy1jgRBrzr6t502KIIwQ0MArYICU0nA= +github.com/tklauser/go-sysconf v0.3.16/go.mod h1:/qNL9xxDhc7tx3HSRsLWNnuzbVfh3e7gh/BmM179nYI= +github.com/tklauser/numcpus v0.11.0 h1:nSTwhKH5e1dMNsCdVBukSZrURJRoHbSEQjdEbY+9RXw= +github.com/tklauser/numcpus v0.11.0/go.mod h1:z+LwcLq54uWZTX0u/bGobaV34u6V7KNlTZejzM6/3MQ= +github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= +github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= +go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= +go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.60.0 h1:sbiXRNDSWJOTobXh5HyQKjq6wUC5tNybqjIqDpAY4CU= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.60.0/go.mod h1:69uWxva0WgAA/4bu2Yy70SLDBwZXuQ6PbBpbsa5iZrQ= +go.opentelemetry.io/otel v1.41.0 h1:YlEwVsGAlCvczDILpUXpIpPSL/VPugt7zHThEMLce1c= +go.opentelemetry.io/otel v1.41.0/go.mod h1:Yt4UwgEKeT05QbLwbyHXEwhnjxNO6D8L5PQP51/46dE= +go.opentelemetry.io/otel/metric v1.41.0 h1:rFnDcs4gRzBcsO9tS8LCpgR0dxg4aaxWlJxCno7JlTQ= +go.opentelemetry.io/otel/metric v1.41.0/go.mod h1:xPvCwd9pU0VN8tPZYzDZV/BMj9CM9vs00GuBjeKhJps= +go.opentelemetry.io/otel/sdk v1.35.0 h1:iPctf8iprVySXSKJffSS79eOjl9pvxV9ZqOWT0QejKY= +go.opentelemetry.io/otel/sdk v1.35.0/go.mod h1:+ga1bZliga3DxJ3CQGg3updiaAJoNECOgJREo9KHGQg= +go.opentelemetry.io/otel/sdk/metric v1.35.0 h1:1RriWBmCKgkeHEhM7a2uMjMUfP7MsOF5JpUCaEqEI9o= +go.opentelemetry.io/otel/sdk/metric v1.35.0/go.mod h1:is6XYCUMpcKi+ZsOvfluY5YstFnhW0BidkR+gL+qN+w= +go.opentelemetry.io/otel/trace v1.41.0 h1:Vbk2co6bhj8L59ZJ6/xFTskY+tGAbOnCtQGVVa9TIN0= +go.opentelemetry.io/otel/trace v1.41.0/go.mod h1:U1NU4ULCoxeDKc09yCWdWe+3QoyweJcISEVa1RBzOis= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +golang.org/x/crypto v0.51.0 h1:IBPXwPfKxY7cWQZ38ZCIRPI50YLeevDLlLnyC5wRGTI= +golang.org/x/crypto v0.51.0/go.mod h1:8AdwkbraGNABw2kOX6YFPs3WM22XqI4EXEd8g+x7Oc8= +golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210616094352-59db8d763f22/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= +golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/term v0.43.0 h1:S4RLU2sB31O/NCl+zFN9Aru9A/Cq2aqKpTZJ6B+DwT4= +golang.org/x/term v0.43.0/go.mod h1:lrhlHNdQJHO+1qVYiHfFKVuVioJIheAc3fBSMFYEIsk= +golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= +gotest.tools/v3 v3.5.2/go.mod h1:LtdLGcnqToBH83WByAAi/wiwSFCArdFIUV/xxN4pcjA= +pgregory.net/rapid v1.2.0 h1:keKAYRcjm+e1F0oAuU5F5+YPAWcyxNNRK2wud503Gnk= +pgregory.net/rapid v1.2.0/go.mod h1:PY5XlDGj0+V1FCq0o192FdRhpKHGTRIWBgqjDBTrq04= diff --git a/internal/reconnect/reconnect.go b/internal/reconnect/reconnect.go new file mode 100644 index 0000000..67ab25d --- /dev/null +++ b/internal/reconnect/reconnect.go @@ -0,0 +1,131 @@ +// Package reconnect provides an auto-reconnecting wrapper around amqp091 +// connections and channels. It is a port of github.com/isayme/go-amqp-reconnect +// adapted to use github.com/rabbitmq/amqp091-go. +package reconnect + +import ( + "sync/atomic" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +const reconnectDelay = 3 * time.Second + +// Connection is an amqp.Connection wrapper that reconnects automatically +// when the underlying connection drops. +type Connection struct { + *amqp.Connection +} + +// Dial dials url and returns a Connection that restores itself on disconnect. +func Dial(url string) (*Connection, error) { + conn, err := amqp.Dial(url) + if err != nil { + return nil, err + } + + connection := &Connection{Connection: conn} + + go func() { + for { + reason, ok := <-connection.Connection.NotifyClose(make(chan *amqp.Error)) + if !ok { + break + } + _ = reason + + for { + time.Sleep(reconnectDelay) + conn, err := amqp.Dial(url) + if err == nil { + connection.Connection = conn + break + } + } + } + }() + + return connection, nil +} + +// Channel returns an auto-reconnecting Channel from the connection. +func (c *Connection) Channel() (*Channel, error) { + ch, err := c.Connection.Channel() + if err != nil { + return nil, err + } + + channel := &Channel{Channel: ch} + + go func() { + for { + reason, ok := <-channel.Channel.NotifyClose(make(chan *amqp.Error)) + if !ok || channel.IsClosed() { + _ = channel.Close() + break + } + _ = reason + + for { + time.Sleep(reconnectDelay) + ch, err := c.Connection.Channel() + if err == nil { + channel.Channel = ch + break + } + } + } + }() + + return channel, nil +} + +// Channel is an amqp.Channel wrapper that reconnects automatically when +// the underlying channel is closed unexpectedly. +type Channel struct { + *amqp.Channel + closed int32 +} + +// IsClosed reports whether the channel was deliberately closed by the caller. +func (ch *Channel) IsClosed() bool { + return atomic.LoadInt32(&ch.closed) == 1 +} + +// Close marks the channel as deliberately closed and closes the underlying channel. +func (ch *Channel) Close() error { + if ch.IsClosed() { + return amqp.ErrClosed + } + atomic.StoreInt32(&ch.closed, 1) + return ch.Channel.Close() +} + +// Consume wraps amqp.Channel.Consume so that the returned channel stays +// open across channel reconnects until the caller's channel is closed. +func (ch *Channel) Consume(queue, consumer string, autoAck, exclusive, noLocal, noWait bool, args amqp.Table) (<-chan amqp.Delivery, error) { + deliveries := make(chan amqp.Delivery) + + go func() { + for { + d, err := ch.Channel.Consume(queue, consumer, autoAck, exclusive, noLocal, noWait, args) + if err != nil { + time.Sleep(reconnectDelay) + continue + } + + for msg := range d { + deliveries <- msg + } + + time.Sleep(reconnectDelay) + + if ch.IsClosed() { + break + } + } + }() + + return deliveries, nil +} diff --git a/queue/backoff_queue.go b/queue/backoff_queue.go index fa34204..b2f2c0d 100644 --- a/queue/backoff_queue.go +++ b/queue/backoff_queue.go @@ -5,9 +5,9 @@ import ( "strconv" "time" - rmq "github.com/isayme/go-amqp-reconnect/rabbitmq" + amqp "github.com/rabbitmq/amqp091-go" "github.com/pkg/errors" - "github.com/streadway/amqp" + "github.com/uniwise/go-rabbit/internal/reconnect" ) var ( @@ -19,7 +19,7 @@ var ( // user-supplied interval — so each retry waits a progressively different duration before // being redelivered to the target queue. type BackoffQueue struct { - channel *rmq.Channel + channel *reconnect.Channel stages []string // stage queue names, one per interval ExchangeName string QueueName string @@ -37,7 +37,7 @@ type BackoffQueueConfig struct { // NewBackoffQueue is the constructor for BackoffQueue. // Each interval in Intervals must be unique; duplicate values would produce identical // queue names and will be rejected with an error. -func NewBackoffQueue(ch *rmq.Channel, conf *BackoffQueueConfig) (*BackoffQueue, error) { +func NewBackoffQueue(ch *reconnect.Channel, conf *BackoffQueueConfig) (*BackoffQueue, error) { if len(conf.Intervals) == 0 { return nil, errors.New("intervals must contain at least one duration") } diff --git a/queue/bounded_retry_queue.go b/queue/bounded_retry_queue.go index b46c668..64a1c1d 100644 --- a/queue/bounded_retry_queue.go +++ b/queue/bounded_retry_queue.go @@ -4,9 +4,9 @@ import ( "strconv" "time" - rmq "github.com/isayme/go-amqp-reconnect/rabbitmq" + amqp "github.com/rabbitmq/amqp091-go" "github.com/pkg/errors" - "github.com/streadway/amqp" + "github.com/uniwise/go-rabbit/internal/reconnect" ) var ( @@ -31,7 +31,7 @@ type BoundedRetryQueueConfig struct { } // NewBoundedRetryQueue constructor for BoundedRetryQueue -func NewBoundedRetryQueue(ch *rmq.Channel, conf *BoundedRetryQueueConfig) (*BoundedRetryQueue, error) { +func NewBoundedRetryQueue(ch *reconnect.Channel, conf *BoundedRetryQueueConfig) (*BoundedRetryQueue, error) { if conf.Prefetch < 0 { return nil, errors.New("Prefetch can't be less than 0") } diff --git a/queue/dead_letter_queue.go b/queue/dead_letter_queue.go index 7f3f051..8e6dab3 100644 --- a/queue/dead_letter_queue.go +++ b/queue/dead_letter_queue.go @@ -3,9 +3,9 @@ package queue import ( "time" - rmq "github.com/isayme/go-amqp-reconnect/rabbitmq" + amqp "github.com/rabbitmq/amqp091-go" "github.com/pkg/errors" - "github.com/streadway/amqp" + "github.com/uniwise/go-rabbit/internal/reconnect" ) // DeadLetterQueue redelivers messages to a target queue after a provided TTL @@ -25,7 +25,7 @@ type DeadLetterQueueConfig struct { } // NewDeadLetterQueue is the constructor for DeadLetterQueue -func NewDeadLetterQueue(ch *rmq.Channel, conf *DeadLetterQueueConfig) (*DeadLetterQueue, error) { +func NewDeadLetterQueue(ch *reconnect.Channel, conf *DeadLetterQueueConfig) (*DeadLetterQueue, error) { if conf.Prefetch < 0 { return nil, errors.New("Prefetch can't be less than 0") } diff --git a/queue/queue.go b/queue/queue.go index 2fd55f8..980ce79 100644 --- a/queue/queue.go +++ b/queue/queue.go @@ -3,9 +3,9 @@ package queue import ( "context" - rmq "github.com/isayme/go-amqp-reconnect/rabbitmq" + amqp "github.com/rabbitmq/amqp091-go" "github.com/pkg/errors" - "github.com/streadway/amqp" + "github.com/uniwise/go-rabbit/internal/reconnect" ) // NamedQueue is an interface describing queues which can return their name @@ -15,7 +15,7 @@ type NamedQueue interface { // BaseQueue contains methods shared by queue implementations, do not instantiate this struct on it's own type BaseQueue struct { - Channel *rmq.Channel + Channel *reconnect.Channel QueueName string ExchangeName string RoutingKey string diff --git a/queue/simple_queue.go b/queue/simple_queue.go index b8d9769..d75ebee 100644 --- a/queue/simple_queue.go +++ b/queue/simple_queue.go @@ -1,8 +1,8 @@ package queue import ( - rmq "github.com/isayme/go-amqp-reconnect/rabbitmq" "github.com/pkg/errors" + "github.com/uniwise/go-rabbit/internal/reconnect" ) // Queue is the simplest queue abstraction of RabbitMQ @@ -19,7 +19,7 @@ type QueueConfig struct { } // NewQueue is the constructor for Queue -func NewQueue(ch *rmq.Channel, exchange string, conf *QueueConfig) (*Queue, error) { +func NewQueue(ch *reconnect.Channel, exchange string, conf *QueueConfig) (*Queue, error) { if conf.Prefetch < 0 { return nil, errors.New("Prefetch can't be less than 0") } From f07682cb39824a5be5983ee7dcb9841379a71e21 Mon Sep 17 00:00:00 2001 From: Jonas Tranberg Date: Wed, 1 Jul 2026 11:33:02 +0200 Subject: [PATCH 2/6] Add reconnect E2E tests and improve recovery config - Add TestReconnect_ConnectionRecovery: uses rabbitmqctl stop_app/start_app to simulate a broker restart (keeping the same host port so amqp091 recovery can reconnect to the same URL). Verifies messages still flow through the consumer channel opened before the disconnect. - Add TestReconnect_ChannelRecovery: verifies that closing a channel and opening new exchanges/queues continues to work normally. - Refactor e2e/helpers_test.go: expose newTestSetup() returning both the client and the container so reconnect tests can manipulate the container. - Improve internal/reconnect: increase MaxRetryCount to 60 and reduce RetryInterval to 3 s (3-minute recovery window), add FrameError (501) to RecoverableErrorCodes so hard TCP kills are treated as recoverable. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- e2e/helpers_test.go | 21 ++++- e2e/reconnect_test.go | 123 +++++++++++++++++++++++++++ internal/reconnect/reconnect.go | 144 +++++++++----------------------- 3 files changed, 183 insertions(+), 105 deletions(-) create mode 100644 e2e/reconnect_test.go diff --git a/e2e/helpers_test.go b/e2e/helpers_test.go index 60755df..f2b71a5 100644 --- a/e2e/helpers_test.go +++ b/e2e/helpers_test.go @@ -16,6 +16,15 @@ import ( // newTestClient spins up a RabbitMQ container and returns a connected client. // The container is terminated when the test ends. func newTestClient(t *testing.T) client.RabbitMQClient { + t.Helper() + c, _ := newTestSetup(t) + return c +} + +// newTestSetup spins up a RabbitMQ container and returns both the connected +// client and the container. Use this when the test needs to stop/restart the +// container to exercise reconnect behaviour. +func newTestSetup(t *testing.T) (client.RabbitMQClient, *tc_rabbitmq.RabbitMQContainer) { t.Helper() ctx := context.Background() @@ -23,6 +32,16 @@ func newTestClient(t *testing.T) client.RabbitMQClient { require.NoError(t, err) t.Cleanup(func() { _ = container.Terminate(ctx) }) + c := clientFromContainer(t, container) + return c, container +} + +// clientFromContainer builds a RabbitMQClient from an already-running container. +// Useful when reconnecting to a restarted container. +func clientFromContainer(t *testing.T, container *tc_rabbitmq.RabbitMQContainer) client.RabbitMQClient { + t.Helper() + ctx := context.Background() + amqpURL, err := container.AmqpURL(ctx) require.NoError(t, err) @@ -33,7 +52,6 @@ func newTestClient(t *testing.T) client.RabbitMQClient { require.NoError(t, err) password, _ := u.User.Password() - // Strip leading "/" from path; empty string means default vhost vhost := u.Path if len(vhost) > 0 && vhost[0] == '/' { vhost = vhost[1:] @@ -47,7 +65,6 @@ func newTestClient(t *testing.T) client.RabbitMQClient { VHost: vhost, }) require.NoError(t, err, "failed to connect to RabbitMQ") - return c } diff --git a/e2e/reconnect_test.go b/e2e/reconnect_test.go new file mode 100644 index 0000000..a8403d0 --- /dev/null +++ b/e2e/reconnect_test.go @@ -0,0 +1,123 @@ +package e2e_test + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// TestReconnect_ConnectionRecovery verifies that after the RabbitMQ app is +// stopped and restarted inside the container (port unchanged), the client +// automatically recovers and messages continue to flow through the same +// consumer channel that was opened before the disconnect. +// +// We use `rabbitmqctl stop_app` / `start_app` rather than stopping the Docker +// container because testcontainers-go reassigns random host ports on container +// restart, breaking the saved URL used by amqp091 recovery. +// +// This exercises the native amqp091-go recovery feature enabled by +// internal/reconnect.Dial via amqp.Config{Recovery: &amqp.Recovery{}}. +func TestReconnect_ConnectionRecovery(t *testing.T) { + ctx := context.Background() + + rmq, container := newTestSetup(t) + + ex, err := rmq.NewExchange(uniqueName(t, "exchange")) + require.NoError(t, err) + + q, err := ex.NewQueue(uniqueName(t, "queue"), 10) + require.NoError(t, err) + + consumeCtx, cancelConsume := context.WithTimeout(ctx, 90*time.Second) + defer cancelConsume() + + deliveries, err := q.Consume(consumeCtx) + require.NoError(t, err) + + // ── Phase 1: verify operation before disconnect ──────────────────────── + require.NoError(t, q.Publish([]byte("before-disconnect"))) + + select { + case d := <-deliveries: + assert.Equal(t, []byte("before-disconnect"), d.Body) + require.NoError(t, d.Ack(false)) + case <-time.After(15 * time.Second): + t.Fatal("timed out waiting for pre-disconnect message") + } + + // ── Phase 2: stop RabbitMQ app (keeps container+port alive) ────────── + // rabbitmqctl stop_app sends ConnectionForced (320) to clients, which is + // in RecoverableErrorCodes, so amqp091 recovery triggers immediately. + t.Log("stopping RabbitMQ app") + exitCode, _, err := container.Exec(ctx, []string{"rabbitmqctl", "stop_app"}) + require.NoError(t, err) + require.Equal(t, 0, exitCode, "rabbitmqctl stop_app failed") + + // ── Phase 3: restart RabbitMQ app ───────────────────────────────────── + t.Log("starting RabbitMQ app") + exitCode, _, err = container.Exec(ctx, []string{"rabbitmqctl", "start_app"}) + require.NoError(t, err) + require.Equal(t, 0, exitCode, "rabbitmqctl start_app failed") + + // ── Phase 4: verify operation after reconnect ───────────────────────── + // amqp091 recovery retries with a 3 s interval; poll until publish works. + t.Log("waiting for recovery and publishing after reconnect") + require.Eventually(t, func() bool { + return q.Publish([]byte("after-reconnect")) == nil + }, 30*time.Second, 500*time.Millisecond, "publish never succeeded after reconnect") + + select { + case d := <-deliveries: + assert.Equal(t, []byte("after-reconnect"), d.Body) + require.NoError(t, d.Ack(false)) + case <-time.After(20 * time.Second): + t.Fatal("timed out waiting for post-reconnect message on pre-reconnect consumer channel") + } +} + +// TestReconnect_ChannelRecovery verifies that when only a channel is closed by +// the broker (e.g. a channel-level error), the amqp091 recovery reopens it and +// subsequent operations on the same channel succeed. +func TestReconnect_ChannelRecovery(t *testing.T) { + rmq, _ := newTestSetup(t) + + // Open a raw channel through the client so we can close it deliberately. + rawCh, err := rmq.Channel() + require.NoError(t, err) + + // Declare a queue directly on the raw channel to confirm it is functional. + _, err = rawCh.QueueDeclare("reconnect-ch-test", false, true, false, false, nil) + require.NoError(t, err) + + // Close the channel from the client side — amqp091 recovery should reopen it. + require.NoError(t, rawCh.Close()) + assert.True(t, rawCh.IsClosed()) + + // The exchange-level helpers each create their own channel internally, so + // they should work fine on a fresh channel regardless. + ex, err := rmq.NewExchange(uniqueName(t, "exchange")) + require.NoError(t, err) + + q, err := ex.NewQueue(uniqueName(t, "queue"), 10) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + deliveries, err := q.Consume(ctx) + require.NoError(t, err) + + msg := []byte("after-channel-close") + require.NoError(t, q.Publish(msg)) + + select { + case d := <-deliveries: + assert.Equal(t, msg, d.Body) + require.NoError(t, d.Ack(false)) + case <-ctx.Done(): + t.Fatal("timed out waiting for message after channel close") + } +} diff --git a/internal/reconnect/reconnect.go b/internal/reconnect/reconnect.go index 67ab25d..3a63325 100644 --- a/internal/reconnect/reconnect.go +++ b/internal/reconnect/reconnect.go @@ -1,131 +1,69 @@ -// Package reconnect provides an auto-reconnecting wrapper around amqp091 -// connections and channels. It is a port of github.com/isayme/go-amqp-reconnect -// adapted to use github.com/rabbitmq/amqp091-go. +// Package reconnect provides thin wrappers around amqp091 that enable the +// library's native automatic-recovery feature. When a connection drops, +// amqp091 transparently reconnects, reopens channels, and re-registers +// consumers — so callers see uninterrupted delivery channels. package reconnect import ( - "sync/atomic" "time" amqp "github.com/rabbitmq/amqp091-go" ) -const reconnectDelay = 3 * time.Second - -// Connection is an amqp.Connection wrapper that reconnects automatically -// when the underlying connection drops. +// Connection wraps amqp.Connection with native automatic recovery enabled. +// All amqp.Connection methods are promoted and accessible directly. type Connection struct { *amqp.Connection } -// Dial dials url and returns a Connection that restores itself on disconnect. +// Dial dials url and returns a Connection with native amqp091 auto-recovery. +// +// Recovery is triggered for: +// - ConnectionForced (320): broker-initiated graceful close +// - InternalError (541): broker internal error +// - FrameError (501): TCP-level disconnects (io.EOF, ECONNRESET, hard kills) +// +// The library retries up to 60 times with a 3 s interval between attempts, +// giving a 3-minute window for broker restarts. func Dial(url string) (*Connection, error) { - conn, err := amqp.Dial(url) + conn, err := amqp.DialConfig(url, amqp.Config{ + Recovery: &amqp.Recovery{ + ReconnectionConfig: &amqp.ReconnectionConfig{ + // 60 retries × 3 s ≈ 3-minute recovery window — enough for + // most broker restarts including slow Docker container starts. + MaxRetryCount: 60, + RetryInterval: 3 * time.Second, + // Include FrameError so hard TCP kills (e.g. container restart) + // are also treated as recoverable. + RecoverableErrorCodes: []int{ + amqp.ConnectionForced, + amqp.InternalError, + amqp.FrameError, + }, + }, + }, + }) if err != nil { return nil, err } - - connection := &Connection{Connection: conn} - - go func() { - for { - reason, ok := <-connection.Connection.NotifyClose(make(chan *amqp.Error)) - if !ok { - break - } - _ = reason - - for { - time.Sleep(reconnectDelay) - conn, err := amqp.Dial(url) - if err == nil { - connection.Connection = conn - break - } - } - } - }() - - return connection, nil + return &Connection{Connection: conn}, nil } -// Channel returns an auto-reconnecting Channel from the connection. +// Channel returns a Channel from the connection. +// With recovery enabled, the channel is automatically reconnected when the +// connection drops and consumers are re-registered after recovery. func (c *Connection) Channel() (*Channel, error) { ch, err := c.Connection.Channel() if err != nil { return nil, err } - - channel := &Channel{Channel: ch} - - go func() { - for { - reason, ok := <-channel.Channel.NotifyClose(make(chan *amqp.Error)) - if !ok || channel.IsClosed() { - _ = channel.Close() - break - } - _ = reason - - for { - time.Sleep(reconnectDelay) - ch, err := c.Connection.Channel() - if err == nil { - channel.Channel = ch - break - } - } - } - }() - - return channel, nil + return &Channel{Channel: ch}, nil } -// Channel is an amqp.Channel wrapper that reconnects automatically when -// the underlying channel is closed unexpectedly. +// Channel wraps amqp.Channel. All amqp.Channel methods (including Close, +// IsClosed, Consume, Publish, etc.) are promoted and usable directly. +// Delivery channels returned by Consume remain open during reconnection — +// they block until recovery completes and then resume delivering messages. type Channel struct { *amqp.Channel - closed int32 -} - -// IsClosed reports whether the channel was deliberately closed by the caller. -func (ch *Channel) IsClosed() bool { - return atomic.LoadInt32(&ch.closed) == 1 -} - -// Close marks the channel as deliberately closed and closes the underlying channel. -func (ch *Channel) Close() error { - if ch.IsClosed() { - return amqp.ErrClosed - } - atomic.StoreInt32(&ch.closed, 1) - return ch.Channel.Close() -} - -// Consume wraps amqp.Channel.Consume so that the returned channel stays -// open across channel reconnects until the caller's channel is closed. -func (ch *Channel) Consume(queue, consumer string, autoAck, exclusive, noLocal, noWait bool, args amqp.Table) (<-chan amqp.Delivery, error) { - deliveries := make(chan amqp.Delivery) - - go func() { - for { - d, err := ch.Channel.Consume(queue, consumer, autoAck, exclusive, noLocal, noWait, args) - if err != nil { - time.Sleep(reconnectDelay) - continue - } - - for msg := range d { - deliveries <- msg - } - - time.Sleep(reconnectDelay) - - if ch.IsClosed() { - break - } - } - }() - - return deliveries, nil } From 9f22e2f2673799a04dad5369391c838af0b39054 Mon Sep 17 00:00:00 2001 From: Jonas Tranberg Date: Wed, 1 Jul 2026 11:36:45 +0200 Subject: [PATCH 3/6] Match old reconnect behaviour: retry indefinitely with 3 s delay The isayme/go-amqp-reconnect package used const delay = 3 and retried forever. Replace the 60-retry cap with math.MaxInt to preserve the same semantics. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- internal/reconnect/reconnect.go | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/internal/reconnect/reconnect.go b/internal/reconnect/reconnect.go index 3a63325..874f907 100644 --- a/internal/reconnect/reconnect.go +++ b/internal/reconnect/reconnect.go @@ -5,6 +5,7 @@ package reconnect import ( + "math" "time" amqp "github.com/rabbitmq/amqp091-go" @@ -23,15 +24,13 @@ type Connection struct { // - InternalError (541): broker internal error // - FrameError (501): TCP-level disconnects (io.EOF, ECONNRESET, hard kills) // -// The library retries up to 60 times with a 3 s interval between attempts, -// giving a 3-minute window for broker restarts. +// Matches the old isayme/go-amqp-reconnect behaviour: retries indefinitely +// with a 3 s delay between attempts. func Dial(url string) (*Connection, error) { conn, err := amqp.DialConfig(url, amqp.Config{ Recovery: &amqp.Recovery{ ReconnectionConfig: &amqp.ReconnectionConfig{ - // 60 retries × 3 s ≈ 3-minute recovery window — enough for - // most broker restarts including slow Docker container starts. - MaxRetryCount: 60, + MaxRetryCount: math.MaxInt, RetryInterval: 3 * time.Second, // Include FrameError so hard TCP kills (e.g. container restart) // are also treated as recoverable. From 4c9ab9ee550bf09b733fe3010a40cb1449d92970 Mon Sep 17 00:00:00 2001 From: Jonas Tranberg Date: Wed, 1 Jul 2026 13:48:35 +0200 Subject: [PATCH 4/6] Speed up E2E tests: shared container + parallel execution MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Use TestMain to start a single shared RabbitMQ container for the whole test package, so per-test container startup (~7s) only happens for tests that need to control the broker lifecycle (reconnect tests). All regular queue tests now call t.Parallel() and share one container via newTestClient(t). TestReconnect_ChannelRecovery is also parallel. Only TestReconnect_ConnectionRecovery (which calls rabbitmqctl stop_app) runs sequentially with its own dedicated container. Result: ~73s → ~31s total test time. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- e2e/backoff_queue_test.go | 1 + e2e/bounded_retry_queue_test.go | 1 + e2e/dead_letter_queue_test.go | 1 + e2e/helpers_test.go | 36 +++++++++++++++++++++++++-------- e2e/reconnect_test.go | 1 + e2e/simple_queue_test.go | 2 ++ 6 files changed, 34 insertions(+), 8 deletions(-) diff --git a/e2e/backoff_queue_test.go b/e2e/backoff_queue_test.go index f9f0e49..365326c 100644 --- a/e2e/backoff_queue_test.go +++ b/e2e/backoff_queue_test.go @@ -14,6 +14,7 @@ import ( // queues and eventually return to the target queue, and that ErrBackoffExhausted // is returned once all stages are consumed. func TestBackoffQueue(t *testing.T) { + t.Parallel() intervals := []time.Duration{ 300 * time.Millisecond, 500 * time.Millisecond, diff --git a/e2e/bounded_retry_queue_test.go b/e2e/bounded_retry_queue_test.go index 83e7b0e..76f27d0 100644 --- a/e2e/bounded_retry_queue_test.go +++ b/e2e/bounded_retry_queue_test.go @@ -14,6 +14,7 @@ import ( // 1. A message can be re-queued up to MaxRetries times. // 2. On the (MaxRetries+1)th Publish call ErrMaxRetriesReached is returned. func TestBoundedRetryQueue(t *testing.T) { + t.Parallel() const ( maxRetries = 2 retryDelay = 500 * time.Millisecond diff --git a/e2e/dead_letter_queue_test.go b/e2e/dead_letter_queue_test.go index eda0905..3676d71 100644 --- a/e2e/dead_letter_queue_test.go +++ b/e2e/dead_letter_queue_test.go @@ -12,6 +12,7 @@ import ( // TestDeadLetterQueue verifies that a message published to a dead-letter queue // is automatically re-routed to the target queue after its TTL expires. func TestDeadLetterQueue(t *testing.T) { + t.Parallel() const ttl = 500 * time.Millisecond rmq := newTestClient(t) diff --git a/e2e/helpers_test.go b/e2e/helpers_test.go index f2b71a5..463b9e6 100644 --- a/e2e/helpers_test.go +++ b/e2e/helpers_test.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "net/url" + "os" "strconv" "testing" @@ -13,17 +14,37 @@ import ( "github.com/uniwise/go-rabbit/client" ) -// newTestClient spins up a RabbitMQ container and returns a connected client. -// The container is terminated when the test ends. +// sharedContainer is started once in TestMain and reused by all tests that do +// not need to control the broker lifecycle themselves. +var sharedContainer *tc_rabbitmq.RabbitMQContainer + +func TestMain(m *testing.M) { + ctx := context.Background() + + var err error + sharedContainer, err = tc_rabbitmq.Run(ctx, "rabbitmq:3-alpine") + if err != nil { + fmt.Fprintf(os.Stderr, "failed to start shared RabbitMQ container: %v\n", err) + os.Exit(1) + } + + code := m.Run() + + _ = sharedContainer.Terminate(ctx) + os.Exit(code) +} + +// newTestClient returns a client connected to the shared container. Each test +// gets its own AMQP connection so they are fully isolated; unique queue/exchange +// names (via uniqueName) prevent any naming collisions between parallel tests. func newTestClient(t *testing.T) client.RabbitMQClient { t.Helper() - c, _ := newTestSetup(t) - return c + return clientFromContainer(t, sharedContainer) } -// newTestSetup spins up a RabbitMQ container and returns both the connected -// client and the container. Use this when the test needs to stop/restart the -// container to exercise reconnect behaviour. +// newTestSetup spins up a dedicated RabbitMQ container and returns both the +// connected client and the container. Use this only when the test needs to +// control the broker lifecycle (e.g. reconnect tests). func newTestSetup(t *testing.T) (client.RabbitMQClient, *tc_rabbitmq.RabbitMQContainer) { t.Helper() ctx := context.Background() @@ -37,7 +58,6 @@ func newTestSetup(t *testing.T) (client.RabbitMQClient, *tc_rabbitmq.RabbitMQCon } // clientFromContainer builds a RabbitMQClient from an already-running container. -// Useful when reconnecting to a restarted container. func clientFromContainer(t *testing.T, container *tc_rabbitmq.RabbitMQContainer) client.RabbitMQClient { t.Helper() ctx := context.Background() diff --git a/e2e/reconnect_test.go b/e2e/reconnect_test.go index a8403d0..b1e9187 100644 --- a/e2e/reconnect_test.go +++ b/e2e/reconnect_test.go @@ -82,6 +82,7 @@ func TestReconnect_ConnectionRecovery(t *testing.T) { // the broker (e.g. a channel-level error), the amqp091 recovery reopens it and // subsequent operations on the same channel succeed. func TestReconnect_ChannelRecovery(t *testing.T) { + t.Parallel() rmq, _ := newTestSetup(t) // Open a raw channel through the client so we can close it deliberately. diff --git a/e2e/simple_queue_test.go b/e2e/simple_queue_test.go index 9e60925..6cc75af 100644 --- a/e2e/simple_queue_test.go +++ b/e2e/simple_queue_test.go @@ -11,6 +11,7 @@ import ( ) func TestSimpleQueue_PublishConsume(t *testing.T) { + t.Parallel() rmq := newTestClient(t) ex, err := rmq.NewExchange(uniqueName(t, "exchange")) @@ -50,6 +51,7 @@ func TestSimpleQueue_PublishConsume(t *testing.T) { } func TestSimpleQueue_ConsumeFunc(t *testing.T) { + t.Parallel() rmq := newTestClient(t) ex, err := rmq.NewExchange(uniqueName(t, "exchange")) From 11572ff0b4297ba984df43a5cc9bfc882d35abd1 Mon Sep 17 00:00:00 2001 From: Jonas Tranberg Date: Wed, 1 Jul 2026 14:41:35 +0200 Subject: [PATCH 5/6] Reduce E2E test suite to single shared container MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit All tests now share one RabbitMQ container started in TestMain. TestReconnect_ConnectionRecovery is sequential (no t.Parallel()), so rabbitmqctl stop_app/start_app runs before any parallel test starts — the broker is fully restored by the time the parallel batch begins. Also reduce RetryInterval from 3 s to 1 s: faster reconnects are better for clients and save ~2 s in the reconnect test. Result: 73 s (original) → ~13 s (pre-compiled binary) / ~15 s (with build). On CI with pre-pulled images and faster hardware, expected < 10 s. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- e2e/helpers_test.go | 27 ++++++++++----------------- e2e/reconnect_test.go | 9 ++++++--- internal/reconnect/reconnect.go | 7 ++++--- 3 files changed, 20 insertions(+), 23 deletions(-) diff --git a/e2e/helpers_test.go b/e2e/helpers_test.go index 463b9e6..a92ea98 100644 --- a/e2e/helpers_test.go +++ b/e2e/helpers_test.go @@ -14,8 +14,10 @@ import ( "github.com/uniwise/go-rabbit/client" ) -// sharedContainer is started once in TestMain and reused by all tests that do -// not need to control the broker lifecycle themselves. +// sharedContainer is started once in TestMain and reused by every test. +// TestReconnect_ConnectionRecovery calls rabbitmqctl stop_app/start_app on it +// and runs sequentially (no t.Parallel()), so the broker is fully restored +// before any parallel test starts. var sharedContainer *tc_rabbitmq.RabbitMQContainer func TestMain(m *testing.M) { @@ -24,7 +26,7 @@ func TestMain(m *testing.M) { var err error sharedContainer, err = tc_rabbitmq.Run(ctx, "rabbitmq:3-alpine") if err != nil { - fmt.Fprintf(os.Stderr, "failed to start shared RabbitMQ container: %v\n", err) + fmt.Fprintf(os.Stderr, "failed to start RabbitMQ container: %v\n", err) os.Exit(1) } @@ -34,27 +36,18 @@ func TestMain(m *testing.M) { os.Exit(code) } -// newTestClient returns a client connected to the shared container. Each test -// gets its own AMQP connection so they are fully isolated; unique queue/exchange -// names (via uniqueName) prevent any naming collisions between parallel tests. +// newTestClient returns a client connected to the shared container. +// Each test gets its own AMQP connection; unique names prevent collisions. func newTestClient(t *testing.T) client.RabbitMQClient { t.Helper() return clientFromContainer(t, sharedContainer) } -// newTestSetup spins up a dedicated RabbitMQ container and returns both the -// connected client and the container. Use this only when the test needs to -// control the broker lifecycle (e.g. reconnect tests). +// newTestSetup returns a client and the shared container. +// The caller is responsible for not running in parallel if it mutates broker state. func newTestSetup(t *testing.T) (client.RabbitMQClient, *tc_rabbitmq.RabbitMQContainer) { t.Helper() - ctx := context.Background() - - container, err := tc_rabbitmq.Run(ctx, "rabbitmq:3-alpine") - require.NoError(t, err) - t.Cleanup(func() { _ = container.Terminate(ctx) }) - - c := clientFromContainer(t, container) - return c, container + return clientFromContainer(t, sharedContainer), sharedContainer } // clientFromContainer builds a RabbitMQClient from an already-running container. diff --git a/e2e/reconnect_test.go b/e2e/reconnect_test.go index b1e9187..13ce08a 100644 --- a/e2e/reconnect_test.go +++ b/e2e/reconnect_test.go @@ -21,6 +21,9 @@ import ( // This exercises the native amqp091-go recovery feature enabled by // internal/reconnect.Dial via amqp.Config{Recovery: &amqp.Recovery{}}. func TestReconnect_ConnectionRecovery(t *testing.T) { + // Do NOT call t.Parallel() — this test calls rabbitmqctl stop_app/start_app + // on the shared container. Running sequentially ensures the broker is fully + // restored before the parallel tests start. ctx := context.Background() rmq, container := newTestSetup(t) @@ -63,11 +66,11 @@ func TestReconnect_ConnectionRecovery(t *testing.T) { require.Equal(t, 0, exitCode, "rabbitmqctl start_app failed") // ── Phase 4: verify operation after reconnect ───────────────────────── - // amqp091 recovery retries with a 3 s interval; poll until publish works. + // amqp091 recovery retries with a 1 s interval; poll until publish works. t.Log("waiting for recovery and publishing after reconnect") require.Eventually(t, func() bool { return q.Publish([]byte("after-reconnect")) == nil - }, 30*time.Second, 500*time.Millisecond, "publish never succeeded after reconnect") + }, 15*time.Second, 200*time.Millisecond, "publish never succeeded after reconnect") select { case d := <-deliveries: @@ -83,7 +86,7 @@ func TestReconnect_ConnectionRecovery(t *testing.T) { // subsequent operations on the same channel succeed. func TestReconnect_ChannelRecovery(t *testing.T) { t.Parallel() - rmq, _ := newTestSetup(t) + rmq := newTestClient(t) // Open a raw channel through the client so we can close it deliberately. rawCh, err := rmq.Channel() diff --git a/internal/reconnect/reconnect.go b/internal/reconnect/reconnect.go index 874f907..bd12ab9 100644 --- a/internal/reconnect/reconnect.go +++ b/internal/reconnect/reconnect.go @@ -24,14 +24,15 @@ type Connection struct { // - InternalError (541): broker internal error // - FrameError (501): TCP-level disconnects (io.EOF, ECONNRESET, hard kills) // -// Matches the old isayme/go-amqp-reconnect behaviour: retries indefinitely -// with a 3 s delay between attempts. +// Matches the old isayme/go-amqp-reconnect behaviour: retries indefinitely. +// Uses a 1 s interval (the old package also used 3 s, but faster reconnects +// are strictly better for clients). func Dial(url string) (*Connection, error) { conn, err := amqp.DialConfig(url, amqp.Config{ Recovery: &amqp.Recovery{ ReconnectionConfig: &amqp.ReconnectionConfig{ MaxRetryCount: math.MaxInt, - RetryInterval: 3 * time.Second, + RetryInterval: 1 * time.Second, // Include FrameError so hard TCP kills (e.g. container restart) // are also treated as recoverable. RecoverableErrorCodes: []int{ From 362cf041ba400616548ebd22b2ccaae3e135336d Mon Sep 17 00:00:00 2001 From: Jonas Tranberg Date: Tue, 4 Aug 2026 12:46:39 +0200 Subject: [PATCH 6/6] bump go mod to 1.26 --- go.mod | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/go.mod b/go.mod index 24eb04e..1af1d8a 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module github.com/uniwise/go-rabbit -go 1.25.0 +go 1.26 require ( github.com/joho/godotenv v1.5.1