235 lines
8.1 KiB
Go
235 lines
8.1 KiB
Go
package mq
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"net"
|
|
"testing"
|
|
"time"
|
|
|
|
kafkago "github.com/segmentio/kafka-go"
|
|
metadataapi "github.com/segmentio/kafka-go/protocol/metadata"
|
|
)
|
|
|
|
type kafkaRoundTripperFunc func(context.Context, net.Addr, kafkago.Request) (kafkago.Response, error)
|
|
|
|
func (f kafkaRoundTripperFunc) RoundTrip(ctx context.Context, addr net.Addr, request kafkago.Request) (kafkago.Response, error) {
|
|
return f(ctx, addr, request)
|
|
}
|
|
|
|
func TestDisabledKafkaIsSafeAndUnavailable(t *testing.T) {
|
|
client, err := NewKafka(KafkaConfig{})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if client.Connected() {
|
|
t.Fatal("disabled kafka reported connected")
|
|
}
|
|
if err = client.Publish(context.Background(), "orders", []byte("test"), AtLeastOnce, false); !errors.Is(err, ErrUnavailable) {
|
|
t.Fatalf("publish error = %v", err)
|
|
}
|
|
if err = client.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func TestEnabledKafkaRequiresConnectionSettings(t *testing.T) {
|
|
if _, err := NewKafka(KafkaConfig{Enabled: true}); err == nil {
|
|
t.Fatal("enabled kafka without brokers should fail")
|
|
}
|
|
if _, err := NewKafka(KafkaConfig{Enabled: true, Brokers: []string{"localhost:9092"}}); err == nil {
|
|
t.Fatal("enabled kafka without client and group ids should fail")
|
|
}
|
|
}
|
|
|
|
func TestKafkaConfigRejectsInvalidBrokerAndTLSSettings(t *testing.T) {
|
|
config := KafkaConfig{Enabled: true, Brokers: []string{"missing-port"}, ClientID: "kra", GroupID: "kra"}
|
|
if _, err := NewKafka(config); err == nil {
|
|
t.Fatal("invalid kafka broker was accepted")
|
|
}
|
|
config.Brokers = []string{"localhost:9092"}
|
|
config.TLSSkipVerify = true
|
|
if _, err := NewKafka(config); err == nil {
|
|
t.Fatal("tls skip verify without tls was accepted")
|
|
}
|
|
}
|
|
|
|
func TestKafkaRejectsUnsupportedMessageSemantics(t *testing.T) {
|
|
client, err := NewKafka(KafkaConfig{})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err = client.Publish(context.Background(), "orders", nil, AtMostOnce, true); err == nil {
|
|
t.Fatal("retained kafka message was accepted")
|
|
}
|
|
if err = client.Publish(context.Background(), "orders", nil, ExactlyOnce, false); err == nil {
|
|
t.Fatal("qos 2 kafka message was accepted")
|
|
}
|
|
if err = client.Subscribe(context.Background(), "orders", ExactlyOnce, func(context.Context, Message) {}); err == nil {
|
|
t.Fatal("qos 2 kafka subscription was accepted")
|
|
}
|
|
}
|
|
|
|
func TestKafkaTopicNotReadyClassification(t *testing.T) {
|
|
for _, err := range []error{
|
|
kafkago.UnknownTopicOrPartition,
|
|
kafkago.LeaderNotAvailable,
|
|
kafkago.NotLeaderForPartition,
|
|
kafkago.WriteErrors{kafkago.UnknownTopicOrPartition},
|
|
} {
|
|
if !kafkaTopicNotReady(err) {
|
|
t.Fatalf("error %v was not classified as topic-not-ready", err)
|
|
}
|
|
}
|
|
for _, err := range []error{
|
|
kafkago.TopicAuthorizationFailed,
|
|
errors.New("network failed"),
|
|
kafkago.WriteErrors{kafkago.UnknownTopicOrPartition, kafkago.TopicAuthorizationFailed},
|
|
} {
|
|
if kafkaTopicNotReady(err) {
|
|
t.Fatalf("error %v was classified as topic-not-ready", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
func TestKafkaStartOffset(t *testing.T) {
|
|
if kafkaStartOffset("latest") != -1 {
|
|
t.Fatal("latest offset was not mapped to kafka last offset")
|
|
}
|
|
if kafkaStartOffset("earliest") != -2 {
|
|
t.Fatal("earliest offset was not mapped to kafka first offset")
|
|
}
|
|
}
|
|
|
|
func TestNilKafkaClientIsSafe(t *testing.T) {
|
|
var client *Kafka
|
|
if err := client.Publish(context.Background(), "orders", nil, AtMostOnce, false); !errors.Is(err, ErrUnavailable) {
|
|
t.Fatalf("publish error = %v", err)
|
|
}
|
|
if err := client.Subscribe(context.Background(), "orders", AtMostOnce, func(context.Context, Message) {}); !errors.Is(err, ErrUnavailable) {
|
|
t.Fatalf("subscribe error = %v", err)
|
|
}
|
|
if err := client.Unsubscribe(context.Background(), "orders"); !errors.Is(err, ErrUnavailable) {
|
|
t.Fatalf("unsubscribe error = %v", err)
|
|
}
|
|
if err := client.Close(); err != nil {
|
|
t.Fatalf("close error = %v", err)
|
|
}
|
|
}
|
|
|
|
func TestKafkaReconnectDelayExpires(t *testing.T) {
|
|
client := &Kafka{config: KafkaConfig{ReconnectInterval: 20 * time.Millisecond}}
|
|
client.markUnavailable()
|
|
if client.Connected() || !client.Reconnecting() {
|
|
t.Fatalf("failure state connected=%t reconnecting=%t", client.Connected(), client.Reconnecting())
|
|
}
|
|
client.reconnectAt.Store(time.Now().Add(-time.Second).UnixNano())
|
|
if client.Reconnecting() {
|
|
t.Fatal("reconnect delay did not expire")
|
|
}
|
|
client.markConnected()
|
|
if !client.Connected() || client.Reconnecting() {
|
|
t.Fatalf("recovered state connected=%t reconnecting=%t", client.Connected(), client.Reconnecting())
|
|
}
|
|
}
|
|
|
|
func TestKafkaUnsubscribeDoesNotWaitForConsumerHandler(t *testing.T) {
|
|
consumeCtx, cancel := context.WithCancel(context.Background())
|
|
reader := kafkago.NewReader(kafkago.ReaderConfig{Brokers: []string{"127.0.0.1:9092"}, Topic: "orders", Partition: 0})
|
|
client := &Kafka{subscriptions: map[string]*kafkaSubscription{
|
|
"orders": {reader: reader, cancel: cancel},
|
|
}}
|
|
done := make(chan error, 1)
|
|
go func() { done <- client.Unsubscribe(context.Background(), "orders") }()
|
|
select {
|
|
case err := <-done:
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
case <-time.After(time.Second):
|
|
t.Fatal("unsubscribe waited for the consumer handler")
|
|
}
|
|
if consumeCtx.Err() == nil {
|
|
t.Fatal("consumer context was not canceled")
|
|
}
|
|
}
|
|
|
|
func TestKafkaUnsubscribeRejectsBlankTopics(t *testing.T) {
|
|
client := &Kafka{subscriptions: make(map[string]*kafkaSubscription)}
|
|
if err := client.Unsubscribe(context.Background(), " "); err == nil {
|
|
t.Fatal("blank kafka topic was accepted")
|
|
}
|
|
}
|
|
|
|
func TestProbeKafkaMetadataDoesNotAutoCreateTopic(t *testing.T) {
|
|
var captured *metadataapi.Request
|
|
client := &kafkago.Client{
|
|
Addr: kafkago.TCP("kafka-1:9092", "kafka-2:9092"),
|
|
Transport: kafkaRoundTripperFunc(func(_ context.Context, _ net.Addr, request kafkago.Request) (kafkago.Response, error) {
|
|
captured = request.(*metadataapi.Request)
|
|
return &metadataapi.Response{
|
|
Brokers: []metadataapi.ResponseBroker{{NodeID: 1, Host: "kafka-1", Port: 9092}},
|
|
Topics: []metadataapi.ResponseTopic{{
|
|
Name: "orders",
|
|
Partitions: []metadataapi.ResponsePartition{{
|
|
PartitionIndex: 0,
|
|
LeaderID: 1,
|
|
}},
|
|
}},
|
|
}, nil
|
|
}),
|
|
}
|
|
if err := probeKafkaMetadata(context.Background(), client, "orders"); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if captured == nil || len(captured.TopicNames) != 1 || captured.TopicNames[0] != "orders" {
|
|
t.Fatalf("metadata request = %#v", captured)
|
|
}
|
|
if captured.AllowAutoTopicCreation {
|
|
t.Fatal("subscription metadata lookup enabled automatic topic creation")
|
|
}
|
|
}
|
|
|
|
func TestProbeKafkaMetadataRequestsBrokersWithoutListingTopics(t *testing.T) {
|
|
var captured *metadataapi.Request
|
|
client := &kafkago.Client{
|
|
Addr: kafkago.TCP("kafka-1:9092"),
|
|
Transport: kafkaRoundTripperFunc(func(_ context.Context, _ net.Addr, request kafkago.Request) (kafkago.Response, error) {
|
|
captured = request.(*metadataapi.Request)
|
|
return &metadataapi.Response{Brokers: []metadataapi.ResponseBroker{{NodeID: 1, Host: "kafka-1", Port: 9092}}}, nil
|
|
}),
|
|
}
|
|
if err := probeKafkaMetadata(context.Background(), client, ""); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if captured == nil || captured.TopicNames == nil || len(captured.TopicNames) != 0 {
|
|
t.Fatalf("broker metadata request topics = %#v", captured)
|
|
}
|
|
}
|
|
|
|
func TestKafkaTopicMetadataErrorDoesNotMarkClusterUnavailable(t *testing.T) {
|
|
metadata := &kafkago.Client{
|
|
Addr: kafkago.TCP("kafka-1:9092"),
|
|
Transport: kafkaRoundTripperFunc(func(_ context.Context, _ net.Addr, _ kafkago.Request) (kafkago.Response, error) {
|
|
return &metadataapi.Response{
|
|
Brokers: []metadataapi.ResponseBroker{{NodeID: 1, Host: "kafka-1", Port: 9092}},
|
|
Topics: []metadataapi.ResponseTopic{{Name: "missing", ErrorCode: int16(kafkago.UnknownTopicOrPartition)}},
|
|
}, nil
|
|
}),
|
|
}
|
|
client := &Kafka{
|
|
config: defaultKafkaConfig(KafkaConfig{}),
|
|
dialer: &kafkago.Dialer{},
|
|
metadata: metadata,
|
|
subscriptions: make(map[string]*kafkaSubscription),
|
|
}
|
|
client.markUnavailable()
|
|
err := client.Subscribe(context.Background(), "missing", AtLeastOnce, func(context.Context, Message) {})
|
|
if err == nil {
|
|
t.Fatal("missing kafka topic was accepted")
|
|
}
|
|
if !client.Connected() || client.Reconnecting() {
|
|
t.Fatalf("topic error changed cluster state: connected=%t reconnecting=%t", client.Connected(), client.Reconnecting())
|
|
}
|
|
}
|