diff --git a/client/client.go b/client/client.go index 08526ac..b266994 100644 --- a/client/client.go +++ b/client/client.go @@ -2,25 +2,24 @@ 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 { Config() *Config - Connection() *rmq.Connection - Channel() (*rmq.Channel, error) + Connection() *reconnect.Connection + Channel() (*reconnect.Channel, error) NewExchange(name string) (*exchange.Exchange, error) } // RabbitMQ is a wrapper struct for a RabbitMQ connection type RabbitMQ struct { config *Config - connectionn *rmq.Connection + connectionn *reconnect.Connection } // New is the constructor for RabbitMQImpl @@ -46,7 +45,7 @@ func (r *RabbitMQ) connect() error { r.config.VHost, ) - conn, err := rmq.Dial(connStr) + conn, err := reconnect.Dial(connStr) if err != nil { return err } @@ -62,13 +61,13 @@ func (r *RabbitMQ) Config() *Config { } // Connection returns the underlying RabbitMQ connection -func (r *RabbitMQ) Connection() *rmq.Connection { +func (r *RabbitMQ) Connection() *reconnect.Connection { return r.connectionn } // Channel returns a RabbitMQ channel from the connection -func (r *RabbitMQ) Channel() (*rmq.Channel, error) { - ch, err := r.connectionn.Channel() +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..365326c --- /dev/null +++ b/e2e/backoff_queue_test.go @@ -0,0 +1,67 @@ +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) { + t.Parallel() + 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..76f27d0 --- /dev/null +++ b/e2e/bounded_retry_queue_test.go @@ -0,0 +1,71 @@ +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) { + t.Parallel() + 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..3676d71 --- /dev/null +++ b/e2e/dead_letter_queue_test.go @@ -0,0 +1,47 @@ +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) { + t.Parallel() + 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..a92ea98 --- /dev/null +++ b/e2e/helpers_test.go @@ -0,0 +1,87 @@ +package e2e_test + +import ( + "context" + "fmt" + "net/url" + "os" + "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" +) + +// 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) { + 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 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; unique names prevent collisions. +func newTestClient(t *testing.T) client.RabbitMQClient { + t.Helper() + return clientFromContainer(t, sharedContainer) +} + +// 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() + return clientFromContainer(t, sharedContainer), sharedContainer +} + +// clientFromContainer builds a RabbitMQClient from an already-running 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) + + 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() + 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/reconnect_test.go b/e2e/reconnect_test.go new file mode 100644 index 0000000..13ce08a --- /dev/null +++ b/e2e/reconnect_test.go @@ -0,0 +1,127 @@ +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) { + // 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) + + 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 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 + }, 15*time.Second, 200*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) { + t.Parallel() + rmq := newTestClient(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/e2e/simple_queue_test.go b/e2e/simple_queue_test.go new file mode 100644 index 0000000..6cc75af --- /dev/null +++ b/e2e/simple_queue_test.go @@ -0,0 +1,82 @@ +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) { + t.Parallel() + 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) { + t.Parallel() + 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..1af1d8a 100644 --- a/go.mod +++ b/go.mod @@ -1,11 +1,65 @@ module github.com/uniwise/go-rabbit -go 1.23 +go 1.26 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..bd12ab9 --- /dev/null +++ b/internal/reconnect/reconnect.go @@ -0,0 +1,69 @@ +// 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 ( + "math" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// 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 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) +// +// 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: 1 * 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 + } + return &Connection{Connection: conn}, nil +} + +// 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 + } + return &Channel{Channel: ch}, nil +} + +// 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 +} 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") }