micro--go-micro
3e885308a0
Fixes #2988. Brings 'golangci-lint run ./...' to zero issues (was ~373): - errcheck: explicitly ignore fire-and-forget calls with '_ =' (and a small errcheck.exclude-functions list for response writes — json Encoder.Encode, http ResponseWriter.Write, fmt.Fprint*); genuine cases handled. - unused: remove dead code (unexported decls and dead test helpers) and the imports they orphaned. - staticcheck: ST1005 error strings, ST1016 receiver names, S1000/S1017/S1019/ S1023 simplifications, SA4004/SA4006/SA4010 dead code, SA1021 net.IP.Equal, SA6002 (store *[]byte in sync.Pool). - govet: fix a context leak (lostcancel) in internal/util/mdns and move t.Fatal/Fatalf out of goroutines (testinggoroutine) in tests. - ineffassign, unconvert: mechanical fixes. CI: the Lint workflow now runs a blocking full-tree 'golangci-lint run' on pushes and PRs (dropped only-new-issues now that the tree is clean). Verified: go build, go vet, test compilation, and unit tests for the behaviourally-touched packages all pass. Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL Co-authored-by: Claude <noreply@anthropic.com>
301 行
6.3 KiB
Go
301 行
6.3 KiB
Go
package rabbitmq_test
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"os"
|
|
"testing"
|
|
"time"
|
|
|
|
"go-micro.dev/v6/logger"
|
|
|
|
micro "go-micro.dev/v6"
|
|
broker "go-micro.dev/v6/broker"
|
|
rabbitmq "go-micro.dev/v6/broker/rabbitmq"
|
|
server "go-micro.dev/v6/server"
|
|
)
|
|
|
|
type Example struct{}
|
|
|
|
func init() {
|
|
rabbitmq.DefaultRabbitURL = "amqp://rabbitmq:rabbitmq@127.0.0.1:5672"
|
|
}
|
|
|
|
type TestEvent struct {
|
|
Name string `json:"name"`
|
|
Age int `json:"age"`
|
|
Time time.Time `json:"time"`
|
|
}
|
|
|
|
func (e *Example) Handler(ctx context.Context, r interface{}) error {
|
|
return nil
|
|
}
|
|
|
|
func TestDurable(t *testing.T) {
|
|
if tr := os.Getenv("TRAVIS"); len(tr) > 0 {
|
|
t.Skip()
|
|
}
|
|
brkrSub := broker.NewSubscribeOptions(
|
|
broker.Queue("queue.default"),
|
|
broker.DisableAutoAck(),
|
|
rabbitmq.DurableQueue(),
|
|
)
|
|
|
|
b := rabbitmq.NewBroker()
|
|
b.Init()
|
|
if err := b.Connect(); err != nil {
|
|
t.Logf("cant conect to broker, skip: %v", err)
|
|
t.Skip()
|
|
}
|
|
|
|
s := server.NewServer(server.Broker(b))
|
|
|
|
service := micro.NewService("test", micro.Server(s),
|
|
micro.Broker(b),
|
|
)
|
|
h := &Example{}
|
|
// Register a subscriber
|
|
micro.RegisterSubscriber(
|
|
"topic",
|
|
service.Server(),
|
|
h.Handler,
|
|
server.SubscriberContext(brkrSub.Context),
|
|
server.SubscriberQueue("queue.default"),
|
|
)
|
|
|
|
// service.Init()
|
|
|
|
if err := service.Run(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func TestWithoutExchange(t *testing.T) {
|
|
|
|
b := rabbitmq.NewBroker(rabbitmq.WithoutExchange())
|
|
b.Init()
|
|
if err := b.Connect(); err != nil {
|
|
t.Logf("cant conect to broker, skip: %v", err)
|
|
t.Skip()
|
|
}
|
|
|
|
s := server.NewServer(server.Broker(b))
|
|
|
|
service := micro.NewService("test", micro.Server(s),
|
|
micro.Broker(b),
|
|
)
|
|
brkrSub := broker.NewSubscribeOptions(
|
|
broker.Queue("direct.queue"),
|
|
broker.DisableAutoAck(),
|
|
rabbitmq.DurableQueue(),
|
|
)
|
|
// Register a subscriber
|
|
err := micro.RegisterSubscriber(
|
|
"direct.queue",
|
|
service.Server(),
|
|
func(ctx context.Context, evt *TestEvent) error {
|
|
logger.Logf(logger.InfoLevel, "receive event: %+v", evt)
|
|
return nil
|
|
},
|
|
server.SubscriberContext(brkrSub.Context),
|
|
server.SubscriberQueue("direct.queue"),
|
|
)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
go func() {
|
|
time.Sleep(5 * time.Second)
|
|
logger.Logf(logger.InfoLevel, "pub event")
|
|
jsonData, _ := json.Marshal(&TestEvent{
|
|
Name: "test",
|
|
Age: 16,
|
|
})
|
|
err := b.Publish("direct.queue", &broker.Message{
|
|
Body: jsonData,
|
|
},
|
|
rabbitmq.DeliveryMode(2),
|
|
rabbitmq.ContentType("application/json"))
|
|
if err != nil {
|
|
t.Errorf("%v", err)
|
|
}
|
|
}()
|
|
|
|
// service.Init()
|
|
|
|
if err := service.Run(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func TestFanoutExchange(t *testing.T) {
|
|
b := rabbitmq.NewBroker(rabbitmq.ExchangeType(rabbitmq.ExchangeTypeFanout), rabbitmq.ExchangeName("fanout.test"))
|
|
b.Init()
|
|
if err := b.Connect(); err != nil {
|
|
t.Logf("cant conect to broker, skip: %v", err)
|
|
t.Skip()
|
|
}
|
|
|
|
s := server.NewServer(server.Broker(b))
|
|
|
|
service := micro.NewService("test", micro.Server(s),
|
|
micro.Broker(b),
|
|
)
|
|
brkrSub := broker.NewSubscribeOptions(
|
|
broker.Queue("fanout.queue"),
|
|
broker.DisableAutoAck(),
|
|
rabbitmq.DurableQueue(),
|
|
)
|
|
// Register a subscriber
|
|
err := micro.RegisterSubscriber(
|
|
"fanout.queue",
|
|
service.Server(),
|
|
func(ctx context.Context, evt *TestEvent) error {
|
|
logger.Logf(logger.InfoLevel, "receive event: %+v", evt)
|
|
return nil
|
|
},
|
|
server.SubscriberContext(brkrSub.Context),
|
|
server.SubscriberQueue("fanout.queue"),
|
|
)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
go func() {
|
|
time.Sleep(5 * time.Second)
|
|
logger.Logf(logger.InfoLevel, "pub event")
|
|
jsonData, _ := json.Marshal(&TestEvent{
|
|
Name: "test",
|
|
Age: 16,
|
|
})
|
|
err := b.Publish("fanout.queue", &broker.Message{
|
|
Body: jsonData,
|
|
},
|
|
rabbitmq.DeliveryMode(2),
|
|
rabbitmq.ContentType("application/json"))
|
|
if err != nil {
|
|
t.Errorf("%v", err)
|
|
}
|
|
}()
|
|
|
|
// service.Init()
|
|
|
|
if err := service.Run(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func TestDirectExchange(t *testing.T) {
|
|
b := rabbitmq.NewBroker(rabbitmq.ExchangeType(rabbitmq.ExchangeTypeDirect), rabbitmq.ExchangeName("direct.test"))
|
|
b.Init()
|
|
if err := b.Connect(); err != nil {
|
|
t.Logf("cant conect to broker, skip: %v", err)
|
|
t.Skip()
|
|
}
|
|
|
|
s := server.NewServer(server.Broker(b))
|
|
|
|
service := micro.NewService("test", micro.Server(s),
|
|
micro.Broker(b),
|
|
)
|
|
brkrSub := broker.NewSubscribeOptions(
|
|
broker.Queue("direct.exchange.queue"),
|
|
broker.DisableAutoAck(),
|
|
rabbitmq.DurableQueue(),
|
|
)
|
|
// Register a subscriber
|
|
err := micro.RegisterSubscriber(
|
|
"direct.exchange.queue",
|
|
service.Server(),
|
|
func(ctx context.Context, evt *TestEvent) error {
|
|
logger.Logf(logger.InfoLevel, "receive event: %+v", evt)
|
|
return nil
|
|
},
|
|
server.SubscriberContext(brkrSub.Context),
|
|
server.SubscriberQueue("direct.exchange.queue"),
|
|
)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
go func() {
|
|
time.Sleep(5 * time.Second)
|
|
logger.Logf(logger.InfoLevel, "pub event")
|
|
jsonData, _ := json.Marshal(&TestEvent{
|
|
Name: "test",
|
|
Age: 16,
|
|
})
|
|
err := b.Publish("direct.exchange.queue", &broker.Message{
|
|
Body: jsonData,
|
|
},
|
|
rabbitmq.DeliveryMode(2),
|
|
rabbitmq.ContentType("application/json"))
|
|
if err != nil {
|
|
t.Errorf("%v", err)
|
|
}
|
|
}()
|
|
|
|
// service.Init()
|
|
|
|
if err := service.Run(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func TestTopicExchange(t *testing.T) {
|
|
b := rabbitmq.NewBroker()
|
|
b.Init()
|
|
if err := b.Connect(); err != nil {
|
|
t.Logf("cant conect to broker, skip: %v", err)
|
|
t.Skip()
|
|
}
|
|
|
|
s := server.NewServer(server.Broker(b))
|
|
|
|
service := micro.NewService("test", micro.Server(s),
|
|
micro.Broker(b),
|
|
)
|
|
brkrSub := broker.NewSubscribeOptions(
|
|
broker.Queue("topic.exchange.queue"),
|
|
broker.DisableAutoAck(),
|
|
rabbitmq.DurableQueue(),
|
|
)
|
|
// Register a subscriber
|
|
err := micro.RegisterSubscriber(
|
|
"my-test-topic",
|
|
service.Server(),
|
|
func(ctx context.Context, evt *TestEvent) error {
|
|
logger.Logf(logger.InfoLevel, "receive event: %+v", evt)
|
|
return nil
|
|
},
|
|
server.SubscriberContext(brkrSub.Context),
|
|
server.SubscriberQueue("topic.exchange.queue"),
|
|
)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
go func() {
|
|
time.Sleep(5 * time.Second)
|
|
logger.Logf(logger.InfoLevel, "pub event")
|
|
jsonData, _ := json.Marshal(&TestEvent{
|
|
Name: "test",
|
|
Age: 16,
|
|
})
|
|
err := b.Publish("my-test-topic", &broker.Message{
|
|
Body: jsonData,
|
|
},
|
|
rabbitmq.DeliveryMode(2),
|
|
rabbitmq.ContentType("application/json"))
|
|
if err != nil {
|
|
t.Errorf("%v", err)
|
|
}
|
|
}()
|
|
|
|
// service.Init()
|
|
|
|
if err := service.Run(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|