240 lines
7.0 KiB
Go
240 lines
7.0 KiB
Go
//go:build integration
|
|
|
|
package mq
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
|
"github.com/google/uuid"
|
|
amqp "github.com/rabbitmq/amqp091-go"
|
|
)
|
|
|
|
func TestLocalRabbitMQConfirmAckAndDeadLetter(t *testing.T) {
|
|
url := os.Getenv("RABBITMQ_URL")
|
|
if url == "" {
|
|
t.Skip("RABBITMQ_URL is not configured")
|
|
}
|
|
broker, err := OpenWithPrefetch(url, uuid.NewString(), 1)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer broker.Close()
|
|
|
|
tenantKey := fmt.Sprintf("integration-%d", time.Now().UnixNano())
|
|
queue, err := broker.DeclareTenantQueue(tenantKey)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
route, err := tenant.NewDispatcherRoute(broker.dispatcherID, tenantKey)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
routingKey := route.InboundKey
|
|
deadLetterQueue := route.DeadLetterQueue
|
|
body := []byte(`{"command_id":"integration-command"}`)
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
|
defer cancel()
|
|
publishCommand := func(ctx context.Context) error {
|
|
channel, err := broker.conn.Channel()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer channel.Close()
|
|
if err := channel.Confirm(false); err != nil {
|
|
return err
|
|
}
|
|
returned := channel.NotifyReturn(make(chan amqp.Return, 1))
|
|
confirmation, err := channel.PublishWithDeferredConfirmWithContext(ctx, DefaultExchange, routingKey, true, false, amqp.Publishing{
|
|
DeliveryMode: amqp.Persistent,
|
|
Body: body,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
acked, err := confirmation.WaitContext(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
select {
|
|
case result := <-returned:
|
|
return fmt.Errorf("command publication returned: code=%d", result.ReplyCode)
|
|
default:
|
|
}
|
|
if !acked {
|
|
return errors.New("command publication was negatively acknowledged")
|
|
}
|
|
return nil
|
|
}
|
|
acked := make(chan struct{}, 1)
|
|
consumeDone := make(chan error, 1)
|
|
go func() {
|
|
consumeDone <- broker.Consume(ctx, queue, func(_ context.Context, gotRoutingKey string, gotBody []byte) error {
|
|
if gotRoutingKey != routingKey || string(gotBody) != string(body) {
|
|
return fmt.Errorf("delivery mismatch: key=%q body=%q", gotRoutingKey, gotBody)
|
|
}
|
|
acked <- struct{}{}
|
|
return nil
|
|
})
|
|
}()
|
|
if err := publishCommand(ctx); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
select {
|
|
case <-acked:
|
|
case <-ctx.Done():
|
|
t.Fatal(ctx.Err())
|
|
}
|
|
cancel()
|
|
select {
|
|
case <-consumeDone:
|
|
case <-time.After(3 * time.Second):
|
|
t.Fatal("consumer did not stop")
|
|
}
|
|
|
|
permanentCtx, permanentCancel := context.WithTimeout(context.Background(), 15*time.Second)
|
|
defer permanentCancel()
|
|
permanentSeen := make(chan struct{}, 1)
|
|
permanentDone := make(chan error, 1)
|
|
go func() {
|
|
permanentDone <- broker.Consume(permanentCtx, queue, func(_ context.Context, _ string, _ []byte) error {
|
|
permanentSeen <- struct{}{}
|
|
return Permanent(errors.New("synthetic permanent command error"))
|
|
})
|
|
}()
|
|
if err := publishCommand(permanentCtx); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
select {
|
|
case <-permanentSeen:
|
|
case <-permanentCtx.Done():
|
|
t.Fatal(permanentCtx.Err())
|
|
}
|
|
permanentCancel()
|
|
select {
|
|
case <-permanentDone:
|
|
case <-time.After(3 * time.Second):
|
|
t.Fatal("permanent consumer did not stop")
|
|
}
|
|
|
|
dlqChannel, err := broker.conn.Channel()
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer dlqChannel.Close()
|
|
if _, err := dlqChannel.QueueDeclarePassive(deadLetterQueue, true, false, false, false, nil); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
dlqDeliveries, err := dlqChannel.Consume(deadLetterQueue, "sip-go-agent-dlq", false, false, false, false, nil)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
deadLettered := make(chan struct{}, 1)
|
|
dlqDone := make(chan error, 1)
|
|
go func() {
|
|
for delivery := range dlqDeliveries {
|
|
if string(delivery.Body) != string(body) {
|
|
dlqDone <- fmt.Errorf("dead-letter body mismatch: %q", delivery.Body)
|
|
return
|
|
}
|
|
if err := delivery.Ack(false); err != nil {
|
|
dlqDone <- err
|
|
return
|
|
}
|
|
deadLettered <- struct{}{}
|
|
dlqDone <- nil
|
|
return
|
|
}
|
|
dlqDone <- errors.New("dead-letter delivery channel closed")
|
|
}()
|
|
select {
|
|
case <-deadLettered:
|
|
_ = dlqChannel.Cancel("sip-go-agent-dlq", false)
|
|
case err := <-dlqDone:
|
|
t.Fatal(err)
|
|
case <-time.After(15 * time.Second):
|
|
t.Fatal("dead-letter message did not arrive")
|
|
}
|
|
select {
|
|
case err := <-dlqDone:
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
case <-time.After(3 * time.Second):
|
|
t.Fatal("dead-letter consumer did not stop")
|
|
}
|
|
}
|
|
|
|
func TestLocalRabbitMQV2PublishesOnlyConfirmedRoutedMessages(t *testing.T) {
|
|
url := os.Getenv("RABBITMQ_URL")
|
|
if url == "" {
|
|
t.Skip("RABBITMQ_URL is not configured")
|
|
}
|
|
broker, err := OpenWithPrefetch(url, uuid.NewString(), 1)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer broker.Close()
|
|
if _, err := broker.DeclareTenantQueue("publish-" + uuid.NewString()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var outbound string
|
|
for route := range broker.outbound {
|
|
outbound = route
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
|
defer cancel()
|
|
|
|
reject := func(exchange, route string, body []byte, reason string) {
|
|
t.Helper()
|
|
if err := broker.Publish(ctx, exchange, route, body); err == nil || !strings.Contains(err.Error(), reason) {
|
|
t.Fatalf("publish exchange=%q route=%q body-size=%d: got %v, want %q", exchange, route, len(body), err, reason)
|
|
}
|
|
}
|
|
reject(DefaultExchange, outbound, []byte(`{"event_id":"wrong-exchange"}`), "SaaS exchange")
|
|
reject(EventExchange, "", []byte(`{"event_id":"empty-route"}`), "declared route")
|
|
reject(EventExchange, outbound, nil, "body bytes")
|
|
reject(EventExchange, outbound, make([]byte, MaxMessageBytes+1), "body bytes")
|
|
reject(EventExchange, outbound, []byte(`{`), "decode outbound message identity")
|
|
reject(EventExchange, outbound, []byte(`{"event_id":"a","message_id":"b"}`), "exactly one")
|
|
reject(EventExchange, outbound, []byte(`{}`), "exactly one")
|
|
reject(EventExchange, outbound, []byte(fmt.Sprintf(`{"event_id":"%s"}`, strings.Repeat("x", 256))), "identity exceeds")
|
|
reject(EventExchange, outbound+".other", []byte(`{"event_id":"unknown-route"}`), "does not belong")
|
|
|
|
channel, err := broker.conn.Channel()
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer channel.Close()
|
|
for _, idField := range []string{"event_id", "message_id"} {
|
|
id := uuid.NewString()
|
|
body := []byte(fmt.Sprintf(`{"%s":"%s"}`, idField, id))
|
|
if err := broker.Publish(ctx, EventExchange, outbound, body); err != nil {
|
|
t.Fatalf("confirmed %s publication: %v", idField, err)
|
|
}
|
|
delivery, ok, err := channel.Get(SaaSQueue, true)
|
|
if err != nil || !ok || string(delivery.Body) != string(body) || delivery.MessageId != id || delivery.DeliveryMode != amqp.Persistent {
|
|
t.Fatalf("confirmed message lost identity or durability: delivery=%+v ok=%t err=%v", delivery, ok, err)
|
|
}
|
|
}
|
|
|
|
if err := channel.QueueUnbind(SaaSQueue, outbound, EventExchange, nil); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
reject(EventExchange, outbound, []byte(`{"event_id":"unroutable"}`), "publication returned")
|
|
if err := channel.QueueBind(SaaSQueue, outbound, EventExchange, false, nil); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := broker.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
reject(EventExchange, outbound, []byte(`{"event_id":"after-close"}`), "closed")
|
|
}
|