Files
go-sip/internal/mq/current_pause_test.go
T

46 lines
1.6 KiB
Go

package mq
import (
"context"
"strings"
"testing"
"git.ipao.vip/rogee/go-sip/internal/tenant"
amqp "github.com/rabbitmq/amqp091-go"
)
func TestCurrentConsumerRequestStopDoesNotWaitForItself(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
consumer := &CurrentConsumer{cancel: cancel, done: make(chan struct{})}
consumer.RequestStop()
select {
case <-ctx.Done():
default:
t.Fatal("consumer cancellation was not requested")
}
consumer.RequestStop()
}
type currentRejectRecorder struct{ rejected int }
func (*currentRejectRecorder) Ack(uint64, bool) error { return nil }
func (*currentRejectRecorder) Nack(uint64, bool, bool) error { return nil }
func (a *currentRejectRecorder) Reject(uint64, bool) error { a.rejected++; return nil }
func TestCurrentMalformedControlStopsConsumerAfterReject(t *testing.T) {
id := "c046b893-8628-4589-ae50-619d049248a6"
route, err := tenant.CurrentControlRoute(id)
if err != nil {
t.Fatal(err)
}
broker := &CurrentBroker{dispatcherID: id, controlQueue: route.Queue}
ack := &currentRejectRecorder{}
deliveries := make(chan amqp.Delivery, 1)
deliveries <- amqp.Delivery{Body: []byte("{"), RoutingKey: route.BindingKey, Acknowledger: ack}
close(deliveries)
err = broker.consumeCurrentDeliveries(context.Background(), id, route.Queue, deliveries, func(context.Context, string, []byte) error { t.Fatal("invalid control reached handler"); return nil })
if err == nil || !strings.Contains(err.Error(), "invalid control") || ack.rejected != 1 {
t.Fatalf("malformed control was silently discarded before admission: rejected=%d err=%v", ack.rejected, err)
}
}