2020-09-29 11:15:29 +03:00
|
|
|
package http_test
|
|
|
|
|
|
|
|
import (
|
2020-10-30 23:27:33 +03:00
|
|
|
"context"
|
2020-09-29 11:15:29 +03:00
|
|
|
"sync"
|
|
|
|
"testing"
|
|
|
|
"time"
|
|
|
|
|
|
|
|
"github.com/google/uuid"
|
2021-01-16 02:49:58 +03:00
|
|
|
http "github.com/unistack-org/micro-broker-http/v3"
|
|
|
|
jsoncodec "github.com/unistack-org/micro-codec-json/v3"
|
2020-09-29 11:15:29 +03:00
|
|
|
"github.com/unistack-org/micro/v3/broker"
|
2021-01-29 16:47:26 +03:00
|
|
|
"github.com/unistack-org/micro/v3/register"
|
2020-09-29 11:15:29 +03:00
|
|
|
)
|
|
|
|
|
|
|
|
var (
|
|
|
|
// mock data
|
2021-01-29 16:47:26 +03:00
|
|
|
testData = map[string][]*register.Service{
|
2020-09-29 11:15:29 +03:00
|
|
|
"foo": {
|
|
|
|
{
|
|
|
|
Name: "foo",
|
|
|
|
Version: "1.0.0",
|
2021-01-29 16:47:26 +03:00
|
|
|
Nodes: []*register.Node{
|
2020-09-29 11:15:29 +03:00
|
|
|
{
|
|
|
|
Id: "foo-1.0.0-123",
|
|
|
|
Address: "localhost:9999",
|
|
|
|
},
|
|
|
|
{
|
|
|
|
Id: "foo-1.0.0-321",
|
|
|
|
Address: "localhost:9999",
|
|
|
|
},
|
|
|
|
},
|
|
|
|
},
|
|
|
|
{
|
|
|
|
Name: "foo",
|
|
|
|
Version: "1.0.1",
|
2021-01-29 16:47:26 +03:00
|
|
|
Nodes: []*register.Node{
|
2020-09-29 11:15:29 +03:00
|
|
|
{
|
|
|
|
Id: "foo-1.0.1-321",
|
|
|
|
Address: "localhost:6666",
|
|
|
|
},
|
|
|
|
},
|
|
|
|
},
|
|
|
|
{
|
|
|
|
Name: "foo",
|
|
|
|
Version: "1.0.3",
|
2021-01-29 16:47:26 +03:00
|
|
|
Nodes: []*register.Node{
|
2020-09-29 11:15:29 +03:00
|
|
|
{
|
|
|
|
Id: "foo-1.0.3-345",
|
|
|
|
Address: "localhost:8888",
|
|
|
|
},
|
|
|
|
},
|
|
|
|
},
|
|
|
|
},
|
|
|
|
}
|
|
|
|
)
|
|
|
|
|
2021-01-29 16:47:26 +03:00
|
|
|
func newTestRegister() register.Register {
|
2021-02-12 20:50:02 +03:00
|
|
|
return register.NewRegister()
|
2020-09-29 11:15:29 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
func sub(be *testing.B, c int) {
|
|
|
|
be.StopTimer()
|
2021-01-29 16:47:26 +03:00
|
|
|
m := newTestRegister()
|
2020-09-29 11:15:29 +03:00
|
|
|
|
2021-01-29 16:47:26 +03:00
|
|
|
b := http.NewBroker(broker.Codec(jsoncodec.NewCodec()), broker.Register(m))
|
2020-09-29 11:15:29 +03:00
|
|
|
topic := uuid.New().String()
|
|
|
|
|
|
|
|
if err := b.Init(); err != nil {
|
|
|
|
be.Fatalf("Unexpected init error: %v", err)
|
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Connect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
be.Fatalf("Unexpected connect error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
msg := &broker.Message{
|
|
|
|
Header: map[string]string{
|
|
|
|
"Content-Type": "application/json",
|
|
|
|
},
|
|
|
|
Body: []byte(`{"message": "Hello World"}`),
|
|
|
|
}
|
|
|
|
|
|
|
|
var subs []broker.Subscriber
|
|
|
|
done := make(chan bool, c)
|
|
|
|
|
|
|
|
for i := 0; i < c; i++ {
|
2020-10-30 23:27:33 +03:00
|
|
|
sub, err := b.Subscribe(context.TODO(), topic, func(p broker.Event) error {
|
2020-09-29 11:15:29 +03:00
|
|
|
done <- true
|
|
|
|
m := p.Message()
|
|
|
|
|
|
|
|
if string(m.Body) != string(msg.Body) {
|
|
|
|
be.Fatalf("Unexpected msg %s, expected %s", string(m.Body), string(msg.Body))
|
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}, broker.SubscribeGroup("shared"))
|
|
|
|
if err != nil {
|
|
|
|
be.Fatalf("Unexpected subscribe error: %v", err)
|
|
|
|
}
|
|
|
|
subs = append(subs, sub)
|
|
|
|
}
|
|
|
|
|
|
|
|
for i := 0; i < be.N; i++ {
|
|
|
|
be.StartTimer()
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Publish(context.TODO(), topic, msg); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
be.Fatalf("Unexpected publish error: %v", err)
|
|
|
|
}
|
|
|
|
<-done
|
|
|
|
be.StopTimer()
|
|
|
|
}
|
|
|
|
|
|
|
|
for _, sub := range subs {
|
2020-10-30 23:27:33 +03:00
|
|
|
sub.Unsubscribe(context.TODO())
|
2020-09-29 11:15:29 +03:00
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Disconnect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
be.Fatalf("Unexpected disconnect error: %v", err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func pub(be *testing.B, c int) {
|
|
|
|
be.StopTimer()
|
2021-01-29 16:47:26 +03:00
|
|
|
m := newTestRegister()
|
|
|
|
b := http.NewBroker(broker.Codec(jsoncodec.NewCodec()), broker.Register(m))
|
2020-09-29 11:15:29 +03:00
|
|
|
topic := uuid.New().String()
|
|
|
|
|
|
|
|
if err := b.Init(); err != nil {
|
|
|
|
be.Fatalf("Unexpected init error: %v", err)
|
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Connect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
be.Fatalf("Unexpected connect error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
msg := &broker.Message{
|
|
|
|
Header: map[string]string{
|
|
|
|
"Content-Type": "application/json",
|
|
|
|
},
|
|
|
|
Body: []byte(`{"message": "Hello World"}`),
|
|
|
|
}
|
|
|
|
|
|
|
|
done := make(chan bool, c*4)
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
sub, err := b.Subscribe(context.TODO(), topic, func(p broker.Event) error {
|
2020-09-29 11:15:29 +03:00
|
|
|
done <- true
|
|
|
|
m := p.Message()
|
|
|
|
if string(m.Body) != string(msg.Body) {
|
|
|
|
be.Fatalf("Unexpected msg %s, expected %s", string(m.Body), string(msg.Body))
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}, broker.SubscribeGroup("shared"))
|
|
|
|
if err != nil {
|
|
|
|
be.Fatalf("Unexpected subscribe error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
var wg sync.WaitGroup
|
|
|
|
ch := make(chan int, c*4)
|
|
|
|
be.StartTimer()
|
|
|
|
|
|
|
|
for i := 0; i < c; i++ {
|
|
|
|
go func() {
|
|
|
|
for range ch {
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Publish(context.TODO(), topic, msg); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
be.Fatalf("Unexpected publish error: %v", err)
|
|
|
|
}
|
|
|
|
select {
|
|
|
|
case <-done:
|
|
|
|
case <-time.After(time.Second):
|
|
|
|
}
|
|
|
|
wg.Done()
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
}
|
|
|
|
|
|
|
|
for i := 0; i < be.N; i++ {
|
|
|
|
wg.Add(1)
|
|
|
|
ch <- i
|
|
|
|
}
|
|
|
|
|
|
|
|
wg.Wait()
|
|
|
|
be.StopTimer()
|
2020-10-30 23:27:33 +03:00
|
|
|
sub.Unsubscribe(context.TODO())
|
2020-09-29 11:15:29 +03:00
|
|
|
close(ch)
|
|
|
|
close(done)
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Disconnect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
be.Fatalf("Unexpected disconnect error: %v", err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func TestBroker(t *testing.T) {
|
2021-01-29 16:47:26 +03:00
|
|
|
m := newTestRegister()
|
|
|
|
b := http.NewBroker(broker.Codec(jsoncodec.NewCodec()), broker.Register(m))
|
2020-09-29 11:15:29 +03:00
|
|
|
|
|
|
|
if err := b.Init(); err != nil {
|
|
|
|
t.Fatalf("Unexpected init error: %v", err)
|
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Connect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected connect error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
msg := &broker.Message{
|
|
|
|
Header: map[string]string{
|
|
|
|
"Content-Type": "application/json",
|
|
|
|
},
|
|
|
|
Body: []byte(`{"message": "Hello World"}`),
|
|
|
|
}
|
|
|
|
|
|
|
|
done := make(chan bool)
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
sub, err := b.Subscribe(context.TODO(), "test", func(p broker.Event) error {
|
2020-09-29 11:15:29 +03:00
|
|
|
m := p.Message()
|
|
|
|
|
|
|
|
if string(m.Body) != string(msg.Body) {
|
|
|
|
t.Fatalf("Unexpected msg %s, expected %s", string(m.Body), string(msg.Body))
|
|
|
|
}
|
|
|
|
|
|
|
|
close(done)
|
|
|
|
return nil
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
t.Fatalf("Unexpected subscribe error: %v", err)
|
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Publish(context.TODO(), "test", msg); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected publish error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
<-done
|
2020-10-30 23:27:33 +03:00
|
|
|
sub.Unsubscribe(context.TODO())
|
2020-09-29 11:15:29 +03:00
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Disconnect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected disconnect error: %v", err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func TestConcurrentSubBroker(t *testing.T) {
|
2021-01-29 16:47:26 +03:00
|
|
|
m := newTestRegister()
|
|
|
|
b := http.NewBroker(broker.Codec(jsoncodec.NewCodec()), broker.Register(m))
|
2020-09-29 11:15:29 +03:00
|
|
|
|
|
|
|
if err := b.Init(); err != nil {
|
|
|
|
t.Fatalf("Unexpected init error: %v", err)
|
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Connect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected connect error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
msg := &broker.Message{
|
|
|
|
Header: map[string]string{
|
|
|
|
"Content-Type": "application/json",
|
|
|
|
},
|
|
|
|
Body: []byte(`{"message": "Hello World"}`),
|
|
|
|
}
|
|
|
|
|
|
|
|
var subs []broker.Subscriber
|
|
|
|
var wg sync.WaitGroup
|
|
|
|
|
|
|
|
for i := 0; i < 10; i++ {
|
2020-10-30 23:27:33 +03:00
|
|
|
sub, err := b.Subscribe(context.TODO(), "test", func(p broker.Event) error {
|
2020-09-29 11:15:29 +03:00
|
|
|
defer wg.Done()
|
|
|
|
|
|
|
|
m := p.Message()
|
|
|
|
|
|
|
|
if string(m.Body) != string(msg.Body) {
|
|
|
|
t.Fatalf("Unexpected msg %s, expected %s", string(m.Body), string(msg.Body))
|
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
t.Fatalf("Unexpected subscribe error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
wg.Add(1)
|
|
|
|
subs = append(subs, sub)
|
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Publish(context.TODO(), "test", msg); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected publish error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
wg.Wait()
|
|
|
|
|
|
|
|
for _, sub := range subs {
|
2020-10-30 23:27:33 +03:00
|
|
|
sub.Unsubscribe(context.TODO())
|
2020-09-29 11:15:29 +03:00
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Disconnect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected disconnect error: %v", err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func TestConcurrentPubBroker(t *testing.T) {
|
2021-01-29 16:47:26 +03:00
|
|
|
m := newTestRegister()
|
|
|
|
b := http.NewBroker(broker.Codec(jsoncodec.NewCodec()), broker.Register(m))
|
2020-09-29 11:15:29 +03:00
|
|
|
|
|
|
|
if err := b.Init(); err != nil {
|
|
|
|
t.Fatalf("Unexpected init error: %v", err)
|
|
|
|
}
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Connect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected connect error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
msg := &broker.Message{
|
|
|
|
Header: map[string]string{
|
|
|
|
"Content-Type": "application/json",
|
|
|
|
},
|
|
|
|
Body: []byte(`{"message": "Hello World"}`),
|
|
|
|
}
|
|
|
|
|
|
|
|
var wg sync.WaitGroup
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
sub, err := b.Subscribe(context.TODO(), "test", func(p broker.Event) error {
|
2020-09-29 11:15:29 +03:00
|
|
|
defer wg.Done()
|
|
|
|
|
|
|
|
m := p.Message()
|
|
|
|
|
|
|
|
if string(m.Body) != string(msg.Body) {
|
|
|
|
t.Fatalf("Unexpected msg %s, expected %s", string(m.Body), string(msg.Body))
|
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
t.Fatalf("Unexpected subscribe error: %v", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
for i := 0; i < 10; i++ {
|
|
|
|
wg.Add(1)
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Publish(context.TODO(), "test", msg); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected publish error: %v", err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
wg.Wait()
|
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
sub.Unsubscribe(context.TODO())
|
2020-09-29 11:15:29 +03:00
|
|
|
|
2020-10-30 23:27:33 +03:00
|
|
|
if err := b.Disconnect(context.TODO()); err != nil {
|
2020-09-29 11:15:29 +03:00
|
|
|
t.Fatalf("Unexpected disconnect error: %v", err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkSub1(b *testing.B) {
|
|
|
|
sub(b, 1)
|
|
|
|
}
|
|
|
|
func BenchmarkSub8(b *testing.B) {
|
|
|
|
sub(b, 8)
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkSub32(b *testing.B) {
|
|
|
|
sub(b, 32)
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkSub64(b *testing.B) {
|
|
|
|
sub(b, 64)
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkSub128(b *testing.B) {
|
|
|
|
sub(b, 128)
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkPub1(b *testing.B) {
|
|
|
|
pub(b, 1)
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkPub8(b *testing.B) {
|
|
|
|
pub(b, 8)
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkPub32(b *testing.B) {
|
|
|
|
pub(b, 32)
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkPub64(b *testing.B) {
|
|
|
|
pub(b, 64)
|
|
|
|
}
|
|
|
|
|
|
|
|
func BenchmarkPub128(b *testing.B) {
|
|
|
|
pub(b, 128)
|
|
|
|
}
|