Browse Source

adjust errors

pull/5890/head
chrislu 7 months ago
parent
commit
faffb2973c
  1. 4
      weed/mq/broker/broker_connect.go
  2. 5
      weed/mq/broker/broker_grpc_pub_balancer.go

4
weed/mq/broker/broker_connect.go

@ -55,9 +55,9 @@ func (b *MessageQueueBroker) BrokerConnectToBalancer(brokerBalancer string, stop
}) })
if err != nil { if err != nil {
if err == io.EOF { if err == io.EOF {
return err
// return err
} }
return fmt.Errorf("send stats message: %v", err)
return fmt.Errorf("send stats message to balancer %s: %v", brokerBalancer, err)
} }
// glog.V(3).Infof("sent stats: %+v", stats) // glog.V(3).Infof("sent stats: %+v", stats)

5
weed/mq/broker/broker_grpc_pub_balancer.go

@ -1,6 +1,7 @@
package broker package broker
import ( import (
"fmt"
"github.com/seaweedfs/seaweedfs/weed/mq/pub_balancer" "github.com/seaweedfs/seaweedfs/weed/mq/pub_balancer"
"github.com/seaweedfs/seaweedfs/weed/pb/mq_pb" "github.com/seaweedfs/seaweedfs/weed/pb/mq_pb"
"google.golang.org/grpc/codes" "google.golang.org/grpc/codes"
@ -14,7 +15,7 @@ func (b *MessageQueueBroker) PublisherToPubBalancer(stream mq_pb.SeaweedMessagin
} }
req, err := stream.Recv() req, err := stream.Recv()
if err != nil { if err != nil {
return err
return fmt.Errorf("receive init message: %v", err)
} }
// process init message // process init message
@ -33,7 +34,7 @@ func (b *MessageQueueBroker) PublisherToPubBalancer(stream mq_pb.SeaweedMessagin
for { for {
req, err := stream.Recv() req, err := stream.Recv()
if err != nil { if err != nil {
return err
return fmt.Errorf("receive stats message from %s: %v", initMessage.Broker, err)
} }
if !b.isLockOwner() { if !b.isLockOwner() {
return status.Errorf(codes.Unavailable, "not current broker balancer") return status.Errorf(codes.Unavailable, "not current broker balancer")

Loading…
Cancel
Save