broker_server.go 3.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116
  1. package broker
  2. import (
  3. "time"
  4. "github.com/seaweedfs/seaweedfs/weed/cluster"
  5. "github.com/seaweedfs/seaweedfs/weed/pb/mq_pb"
  6. "github.com/seaweedfs/seaweedfs/weed/wdclient"
  7. "google.golang.org/grpc"
  8. "github.com/seaweedfs/seaweedfs/weed/pb"
  9. "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
  10. "github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
  11. )
  12. type MessageQueueBrokerOption struct {
  13. Masters map[string]pb.ServerAddress
  14. FilerGroup string
  15. DataCenter string
  16. Rack string
  17. DefaultReplication string
  18. MaxMB int
  19. Ip string
  20. Port int
  21. Cipher bool
  22. }
  23. type MessageQueueBroker struct {
  24. mq_pb.UnimplementedSeaweedMessagingServer
  25. option *MessageQueueBrokerOption
  26. grpcDialOption grpc.DialOption
  27. MasterClient *wdclient.MasterClient
  28. filers map[pb.ServerAddress]struct{}
  29. currentFiler pb.ServerAddress
  30. }
  31. func NewMessageBroker(option *MessageQueueBrokerOption, grpcDialOption grpc.DialOption) (mqBroker *MessageQueueBroker, err error) {
  32. mqBroker = &MessageQueueBroker{
  33. option: option,
  34. grpcDialOption: grpcDialOption,
  35. MasterClient: wdclient.NewMasterClient(grpcDialOption, option.FilerGroup, cluster.BrokerType, pb.NewServerAddress(option.Ip, option.Port, 0), option.DataCenter, option.Rack, option.Masters),
  36. filers: make(map[pb.ServerAddress]struct{}),
  37. }
  38. mqBroker.MasterClient.SetOnPeerUpdateFn(mqBroker.OnBrokerUpdate)
  39. go mqBroker.MasterClient.KeepConnectedToMaster()
  40. existingNodes := cluster.ListExistingPeerUpdates(mqBroker.MasterClient.GetMaster(), grpcDialOption, option.FilerGroup, cluster.FilerType)
  41. for _, newNode := range existingNodes {
  42. mqBroker.OnBrokerUpdate(newNode, time.Now())
  43. }
  44. return mqBroker, nil
  45. }
  46. func (broker *MessageQueueBroker) OnBrokerUpdate(update *master_pb.ClusterNodeUpdate, startFrom time.Time) {
  47. if update.NodeType != cluster.FilerType {
  48. return
  49. }
  50. address := pb.ServerAddress(update.Address)
  51. if update.IsAdd {
  52. broker.filers[address] = struct{}{}
  53. if broker.currentFiler == "" {
  54. broker.currentFiler = address
  55. }
  56. } else {
  57. delete(broker.filers, address)
  58. if broker.currentFiler == address {
  59. for filer := range broker.filers {
  60. broker.currentFiler = filer
  61. break
  62. }
  63. }
  64. }
  65. }
  66. func (broker *MessageQueueBroker) GetFiler() pb.ServerAddress {
  67. return broker.currentFiler
  68. }
  69. func (broker *MessageQueueBroker) WithFilerClient(streamingMode bool, fn func(filer_pb.SeaweedFilerClient) error) error {
  70. return pb.WithFilerClient(streamingMode, broker.GetFiler(), broker.grpcDialOption, fn)
  71. }
  72. func (broker *MessageQueueBroker) AdjustedUrl(location *filer_pb.Location) string {
  73. return location.Url
  74. }
  75. func (broker *MessageQueueBroker) GetDataCenter() string {
  76. return ""
  77. }
  78. func (broker *MessageQueueBroker) withMasterClient(streamingMode bool, master pb.ServerAddress, fn func(client master_pb.SeaweedClient) error) error {
  79. return pb.WithMasterClient(streamingMode, master, broker.grpcDialOption, false, func(client master_pb.SeaweedClient) error {
  80. return fn(client)
  81. })
  82. }
  83. func (broker *MessageQueueBroker) withBrokerClient(streamingMode bool, server pb.ServerAddress, fn func(client mq_pb.SeaweedMessagingClient) error) error {
  84. return pb.WithBrokerClient(streamingMode, server, broker.grpcDialOption, func(client mq_pb.SeaweedMessagingClient) error {
  85. return fn(client)
  86. })
  87. }