2021-02-12 16:33:16 +03:00
|
|
|
package broker
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"fmt"
|
|
|
|
"testing"
|
2021-07-23 12:03:18 +03:00
|
|
|
|
2023-04-11 22:20:37 +03:00
|
|
|
"go.unistack.org/micro/v4/metadata"
|
2021-02-12 16:33:16 +03:00
|
|
|
)
|
|
|
|
|
2021-07-22 22:53:44 +03:00
|
|
|
func TestMemoryBatchBroker(t *testing.T) {
|
|
|
|
b := NewBroker()
|
|
|
|
ctx := context.Background()
|
|
|
|
|
|
|
|
if err := b.Connect(ctx); err != nil {
|
|
|
|
t.Fatalf("Unexpected connect error %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
topic := "test"
|
|
|
|
count := 10
|
|
|
|
|
2023-07-29 00:40:58 +03:00
|
|
|
fn := func(evts []Message) error {
|
|
|
|
var err error
|
|
|
|
for _, evt := range evts {
|
|
|
|
if err = evt.Ack(); err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return nil
|
2021-07-22 22:53:44 +03:00
|
|
|
}
|
|
|
|
|
2023-07-29 00:40:58 +03:00
|
|
|
sub, err := b.Subscribe(ctx, topic, fn)
|
2021-07-22 22:53:44 +03:00
|
|
|
if err != nil {
|
|
|
|
t.Fatalf("Unexpected error subscribing %v", err)
|
|
|
|
}
|
|
|
|
|
2023-07-29 00:40:58 +03:00
|
|
|
msgs := make([]Message, 0, count)
|
2021-07-22 22:53:44 +03:00
|
|
|
for i := 0; i < count; i++ {
|
2023-07-29 00:40:58 +03:00
|
|
|
message := &memoryMessage{
|
|
|
|
header: map[string]string{
|
2021-07-23 12:03:18 +03:00
|
|
|
metadata.HeaderTopic: topic,
|
|
|
|
"foo": "bar",
|
|
|
|
"id": fmt.Sprintf("%d", i),
|
2021-07-22 22:53:44 +03:00
|
|
|
},
|
2023-07-29 00:40:58 +03:00
|
|
|
body: []byte(`"hello world"`),
|
2021-07-22 22:53:44 +03:00
|
|
|
}
|
|
|
|
msgs = append(msgs, message)
|
|
|
|
}
|
|
|
|
|
2023-07-29 00:40:58 +03:00
|
|
|
if err := b.Publish(ctx, msgs); err != nil {
|
2021-07-22 22:53:44 +03:00
|
|
|
t.Fatalf("Unexpected error publishing %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if err := sub.Unsubscribe(ctx); err != nil {
|
|
|
|
t.Fatalf("Unexpected error unsubscribing from %s: %v", topic, err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if err := b.Disconnect(ctx); err != nil {
|
|
|
|
t.Fatalf("Unexpected connect error %v", err)
|
|
|
|
}
|
|
|
|
}
|
2021-09-28 23:43:43 +03:00
|
|
|
|
2021-02-12 16:33:16 +03:00
|
|
|
func TestMemoryBroker(t *testing.T) {
|
|
|
|
b := NewBroker()
|
|
|
|
ctx := context.Background()
|
|
|
|
|
|
|
|
if err := b.Connect(ctx); err != nil {
|
|
|
|
t.Fatalf("Unexpected connect error %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
topic := "test"
|
|
|
|
count := 10
|
|
|
|
|
2023-07-29 00:40:58 +03:00
|
|
|
fn := func(p Message) error {
|
|
|
|
return p.Ack()
|
2021-02-12 16:33:16 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
sub, err := b.Subscribe(ctx, topic, fn)
|
|
|
|
if err != nil {
|
|
|
|
t.Fatalf("Unexpected error subscribing %v", err)
|
|
|
|
}
|
|
|
|
|
2023-07-29 00:40:58 +03:00
|
|
|
msgs := make([]Message, 0, count)
|
2021-02-12 16:33:16 +03:00
|
|
|
for i := 0; i < count; i++ {
|
2023-07-29 00:40:58 +03:00
|
|
|
message := &memoryMessage{
|
|
|
|
header: map[string]string{
|
2021-07-23 12:03:18 +03:00
|
|
|
metadata.HeaderTopic: topic,
|
|
|
|
"foo": "bar",
|
|
|
|
"id": fmt.Sprintf("%d", i),
|
2021-02-12 16:33:16 +03:00
|
|
|
},
|
2023-07-29 00:40:58 +03:00
|
|
|
body: []byte(`"hello world"`),
|
2021-02-12 16:33:16 +03:00
|
|
|
}
|
2021-07-22 22:53:44 +03:00
|
|
|
msgs = append(msgs, message)
|
2021-02-12 16:33:16 +03:00
|
|
|
}
|
|
|
|
|
2023-07-29 00:40:58 +03:00
|
|
|
if err := b.Publish(ctx, msgs); err != nil {
|
2021-07-22 22:53:44 +03:00
|
|
|
t.Fatalf("Unexpected error publishing %v", err)
|
|
|
|
}
|
|
|
|
|
2021-02-12 16:33:16 +03:00
|
|
|
if err := sub.Unsubscribe(ctx); err != nil {
|
|
|
|
t.Fatalf("Unexpected error unsubscribing from %s: %v", topic, err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if err := b.Disconnect(ctx); err != nil {
|
|
|
|
t.Fatalf("Unexpected connect error %v", err)
|
|
|
|
}
|
|
|
|
}
|