2019-10-03 18:19:02 +03:00
|
|
|
package service
|
|
|
|
|
|
|
|
import (
|
2020-01-30 14:39:00 +03:00
|
|
|
"github.com/micro/go-micro/v2/broker"
|
|
|
|
pb "github.com/micro/go-micro/v2/broker/service/proto"
|
2020-03-11 20:55:39 +03:00
|
|
|
"github.com/micro/go-micro/v2/logger"
|
2019-10-03 18:19:02 +03:00
|
|
|
)
|
|
|
|
|
|
|
|
type serviceSub struct {
|
|
|
|
topic string
|
|
|
|
queue string
|
|
|
|
handler broker.Handler
|
|
|
|
stream pb.Broker_SubscribeService
|
|
|
|
closed chan bool
|
|
|
|
options broker.SubscribeOptions
|
|
|
|
}
|
|
|
|
|
|
|
|
type serviceEvent struct {
|
|
|
|
topic string
|
2020-03-07 00:25:16 +03:00
|
|
|
err error
|
2019-10-03 18:19:02 +03:00
|
|
|
message *broker.Message
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *serviceEvent) Topic() string {
|
|
|
|
return s.topic
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *serviceEvent) Message() *broker.Message {
|
|
|
|
return s.message
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *serviceEvent) Ack() error {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2020-03-07 00:25:16 +03:00
|
|
|
func (s *serviceEvent) Error() error {
|
|
|
|
return s.err
|
|
|
|
}
|
|
|
|
|
2019-10-04 18:30:03 +03:00
|
|
|
func (s *serviceSub) isClosed() bool {
|
|
|
|
select {
|
|
|
|
case <-s.closed:
|
|
|
|
return true
|
|
|
|
default:
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *serviceSub) run() error {
|
2019-10-03 18:19:02 +03:00
|
|
|
exit := make(chan bool)
|
|
|
|
go func() {
|
|
|
|
select {
|
|
|
|
case <-exit:
|
|
|
|
case <-s.closed:
|
|
|
|
}
|
2019-10-04 18:30:03 +03:00
|
|
|
|
|
|
|
// close the stream
|
|
|
|
s.stream.Close()
|
2019-10-03 18:19:02 +03:00
|
|
|
}()
|
|
|
|
|
|
|
|
for {
|
|
|
|
// TODO: do not fail silently
|
|
|
|
msg, err := s.stream.Recv()
|
|
|
|
if err != nil {
|
2020-03-11 20:55:39 +03:00
|
|
|
if logger.V(logger.DebugLevel, logger.DefaultLogger) {
|
|
|
|
logger.Debugf("Streaming error for subcription to topic %s: %v", s.Topic(), err)
|
|
|
|
}
|
2019-10-04 18:30:03 +03:00
|
|
|
|
|
|
|
// close the exit channel
|
2019-10-03 18:19:02 +03:00
|
|
|
close(exit)
|
2019-10-04 18:30:03 +03:00
|
|
|
|
|
|
|
// don't return an error if we unsubscribed
|
|
|
|
if s.isClosed() {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// return stream error
|
|
|
|
return err
|
2019-10-03 18:19:02 +03:00
|
|
|
}
|
2019-10-04 18:30:03 +03:00
|
|
|
|
2020-03-07 00:25:16 +03:00
|
|
|
p := &serviceEvent{
|
2019-10-03 18:19:02 +03:00
|
|
|
topic: s.topic,
|
|
|
|
message: &broker.Message{
|
|
|
|
Header: msg.Header,
|
|
|
|
Body: msg.Body,
|
|
|
|
},
|
2020-03-07 00:25:16 +03:00
|
|
|
}
|
|
|
|
p.err = s.handler(p)
|
2019-10-03 18:19:02 +03:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *serviceSub) Options() broker.SubscribeOptions {
|
|
|
|
return s.options
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *serviceSub) Topic() string {
|
|
|
|
return s.topic
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *serviceSub) Unsubscribe() error {
|
|
|
|
select {
|
|
|
|
case <-s.closed:
|
|
|
|
return nil
|
|
|
|
default:
|
|
|
|
close(s.closed)
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|