mq.proto 5.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232
  1. syntax = "proto3";
  2. package messaging_pb;
  3. option go_package = "github.com/seaweedfs/seaweedfs/weed/pb/mq_pb";
  4. option java_package = "seaweedfs.mq";
  5. option java_outer_classname = "MessageQueueProto";
  6. //////////////////////////////////////////////////
  7. service SeaweedMessaging {
  8. // control plane
  9. rpc FindBrokerLeader (FindBrokerLeaderRequest) returns (FindBrokerLeaderResponse) {
  10. }
  11. // control plane for balancer
  12. rpc PublisherToPubBalancer (stream PublisherToPubBalancerRequest) returns (stream PublisherToPubBalancerResponse) {
  13. }
  14. rpc BalanceTopics (BalanceTopicsRequest) returns (BalanceTopicsResponse) {
  15. }
  16. // control plane for topic partitions
  17. rpc ListTopics (ListTopicsRequest) returns (ListTopicsResponse) {
  18. }
  19. rpc ConfigureTopic (ConfigureTopicRequest) returns (ConfigureTopicResponse) {
  20. }
  21. rpc LookupTopicBrokers (LookupTopicBrokersRequest) returns (LookupTopicBrokersResponse) {
  22. }
  23. // invoked by the balancer, running on each broker
  24. rpc AssignTopicPartitions (AssignTopicPartitionsRequest) returns (AssignTopicPartitionsResponse) {
  25. }
  26. rpc ClosePublishers(ClosePublishersRequest) returns (ClosePublishersResponse) {
  27. }
  28. rpc CloseSubscribers(CloseSubscribersRequest) returns (CloseSubscribersResponse) {
  29. }
  30. // subscriber connects to broker balancer, which coordinates with the subscribers
  31. rpc SubscriberToSubCoordinator (stream SubscriberToSubCoordinatorRequest) returns (stream SubscriberToSubCoordinatorResponse) {
  32. }
  33. // data plane for each topic partition
  34. rpc Publish (stream PublishRequest) returns (stream PublishResponse) {
  35. }
  36. rpc Subscribe (SubscribeRequest) returns (stream SubscribeResponse) {
  37. }
  38. }
  39. //////////////////////////////////////////////////
  40. message FindBrokerLeaderRequest {
  41. string filer_group = 1;
  42. }
  43. message FindBrokerLeaderResponse {
  44. string broker = 1;
  45. }
  46. message Topic {
  47. string namespace = 1;
  48. string name = 2;
  49. }
  50. message Partition {
  51. int32 ring_size = 1;
  52. int32 range_start = 2;
  53. int32 range_stop = 3;
  54. int64 unix_time_ns = 4;
  55. }
  56. //////////////////////////////////////////////////
  57. message BrokerStats {
  58. int32 cpu_usage_percent = 1;
  59. map<string, TopicPartitionStats> stats = 2;
  60. }
  61. message TopicPartitionStats {
  62. Topic topic = 1;
  63. Partition partition = 2;
  64. int32 consumer_count = 3;
  65. bool is_leader = 4;
  66. }
  67. message PublisherToPubBalancerRequest {
  68. message InitMessage {
  69. string broker = 1;
  70. }
  71. oneof message {
  72. InitMessage init = 1;
  73. BrokerStats stats = 2;
  74. }
  75. }
  76. message PublisherToPubBalancerResponse {
  77. }
  78. message BalanceTopicsRequest {
  79. }
  80. message BalanceTopicsResponse {
  81. }
  82. //////////////////////////////////////////////////
  83. message ConfigureTopicRequest {
  84. Topic topic = 1;
  85. int32 partition_count = 2;
  86. }
  87. message ConfigureTopicResponse {
  88. repeated BrokerPartitionAssignment broker_partition_assignments = 2;
  89. }
  90. message ListTopicsRequest {
  91. }
  92. message ListTopicsResponse {
  93. repeated Topic topics = 1;
  94. }
  95. message LookupTopicBrokersRequest {
  96. Topic topic = 1;
  97. bool is_for_publish = 2;
  98. }
  99. message LookupTopicBrokersResponse {
  100. Topic topic = 1;
  101. repeated BrokerPartitionAssignment broker_partition_assignments = 2;
  102. }
  103. message BrokerPartitionAssignment {
  104. Partition partition = 1;
  105. string leader_broker = 2;
  106. repeated string follower_brokers = 3;
  107. }
  108. message AssignTopicPartitionsRequest {
  109. Topic topic = 1;
  110. repeated BrokerPartitionAssignment broker_partition_assignments = 2;
  111. bool is_leader = 3;
  112. bool is_draining = 4;
  113. }
  114. message AssignTopicPartitionsResponse {
  115. }
  116. message SubscriberToSubCoordinatorRequest {
  117. message InitMessage {
  118. string consumer_group = 1;
  119. string consumer_instance_id = 2;
  120. Topic topic = 3;
  121. }
  122. message AckMessage {
  123. Partition partition = 1;
  124. int64 ts_ns = 2;
  125. }
  126. oneof message {
  127. InitMessage init = 1;
  128. AckMessage ack = 2;
  129. }
  130. }
  131. message SubscriberToSubCoordinatorResponse {
  132. message AssignedPartition {
  133. Partition partition = 1;
  134. int64 ts_ns = 2;
  135. }
  136. message Assignment {
  137. int64 generation = 1;
  138. repeated AssignedPartition assigned_partitions = 2;
  139. }
  140. oneof message {
  141. Assignment assignment = 1;
  142. }
  143. }
  144. //////////////////////////////////////////////////
  145. message DataMessage {
  146. bytes key = 1;
  147. bytes value = 2;
  148. int64 ts_ns = 3;
  149. }
  150. message PublishRequest {
  151. message InitMessage {
  152. Topic topic = 1;
  153. Partition partition = 2;
  154. int32 ack_interval = 3;
  155. }
  156. oneof message {
  157. InitMessage init = 1;
  158. DataMessage data = 2;
  159. }
  160. int64 sequence = 3;
  161. }
  162. message PublishResponse {
  163. int64 ack_sequence = 1;
  164. string error = 2;
  165. bool should_close = 3;
  166. }
  167. message SubscribeRequest {
  168. message InitMessage {
  169. string consumer_group = 1;
  170. string consumer_id = 2;
  171. string client_id = 3;
  172. Topic topic = 4;
  173. Partition partition = 5;
  174. oneof offset {
  175. int64 start_offset = 6;
  176. int64 start_timestamp_ns = 7;
  177. }
  178. string filter = 8;
  179. }
  180. message AckMessage {
  181. int64 sequence = 1;
  182. }
  183. oneof message {
  184. InitMessage init = 1;
  185. AckMessage ack = 2;
  186. }
  187. }
  188. message SubscribeResponse {
  189. message CtrlMessage {
  190. string error = 1;
  191. bool is_end_of_stream = 2;
  192. bool is_end_of_topic = 3;
  193. }
  194. oneof message {
  195. CtrlMessage ctrl = 1;
  196. DataMessage data = 2;
  197. }
  198. }
  199. message ClosePublishersRequest {
  200. Topic topic = 1;
  201. int64 unix_time_ns = 2;
  202. }
  203. message ClosePublishersResponse {
  204. }
  205. message CloseSubscribersRequest {
  206. Topic topic = 1;
  207. int64 unix_time_ns = 2;
  208. }
  209. message CloseSubscribersResponse {
  210. }