mq_broker.proto 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378
  1. syntax = "proto3";
  2. package messaging_pb;
  3. import "mq_schema.proto";
  4. option go_package = "github.com/seaweedfs/seaweedfs/weed/pb/mq_pb";
  5. option java_package = "seaweedfs.mq";
  6. option java_outer_classname = "MessageQueueProto";
  7. //////////////////////////////////////////////////
  8. service SeaweedMessaging {
  9. // control plane
  10. rpc FindBrokerLeader (FindBrokerLeaderRequest) returns (FindBrokerLeaderResponse) {
  11. }
  12. // control plane for balancer
  13. rpc PublisherToPubBalancer (stream PublisherToPubBalancerRequest) returns (stream PublisherToPubBalancerResponse) {
  14. }
  15. rpc BalanceTopics (BalanceTopicsRequest) returns (BalanceTopicsResponse) {
  16. }
  17. // control plane for topic partitions
  18. rpc ListTopics (ListTopicsRequest) returns (ListTopicsResponse) {
  19. }
  20. rpc ConfigureTopic (ConfigureTopicRequest) returns (ConfigureTopicResponse) {
  21. }
  22. rpc LookupTopicBrokers (LookupTopicBrokersRequest) returns (LookupTopicBrokersResponse) {
  23. }
  24. rpc GetTopicConfiguration (GetTopicConfigurationRequest) returns (GetTopicConfigurationResponse) {
  25. }
  26. rpc GetTopicPublishers (GetTopicPublishersRequest) returns (GetTopicPublishersResponse) {
  27. }
  28. rpc GetTopicSubscribers (GetTopicSubscribersRequest) returns (GetTopicSubscribersResponse) {
  29. }
  30. // invoked by the balancer, running on each broker
  31. rpc AssignTopicPartitions (AssignTopicPartitionsRequest) returns (AssignTopicPartitionsResponse) {
  32. }
  33. rpc ClosePublishers(ClosePublishersRequest) returns (ClosePublishersResponse) {
  34. }
  35. rpc CloseSubscribers(CloseSubscribersRequest) returns (CloseSubscribersResponse) {
  36. }
  37. // subscriber connects to broker balancer, which coordinates with the subscribers
  38. rpc SubscriberToSubCoordinator (stream SubscriberToSubCoordinatorRequest) returns (stream SubscriberToSubCoordinatorResponse) {
  39. }
  40. // data plane for each topic partition
  41. rpc PublishMessage (stream PublishMessageRequest) returns (stream PublishMessageResponse) {
  42. }
  43. rpc SubscribeMessage (stream SubscribeMessageRequest) returns (stream SubscribeMessageResponse) {
  44. }
  45. // The lead broker asks a follower broker to follow itself
  46. rpc PublishFollowMe (stream PublishFollowMeRequest) returns (stream PublishFollowMeResponse) {
  47. }
  48. rpc SubscribeFollowMe (stream SubscribeFollowMeRequest) returns (SubscribeFollowMeResponse) {
  49. }
  50. // SQL query support - get unflushed messages from broker's in-memory buffer (streaming)
  51. rpc GetUnflushedMessages (GetUnflushedMessagesRequest) returns (stream GetUnflushedMessagesResponse) {
  52. }
  53. }
  54. //////////////////////////////////////////////////
  55. message FindBrokerLeaderRequest {
  56. string filer_group = 1;
  57. }
  58. message FindBrokerLeaderResponse {
  59. string broker = 1;
  60. }
  61. //////////////////////////////////////////////////
  62. message BrokerStats {
  63. int32 cpu_usage_percent = 1;
  64. map<string, TopicPartitionStats> stats = 2;
  65. }
  66. message TopicPartitionStats {
  67. schema_pb.Topic topic = 1;
  68. schema_pb.Partition partition = 2;
  69. int32 publisher_count = 3;
  70. int32 subscriber_count = 4;
  71. string follower = 5;
  72. }
  73. message PublisherToPubBalancerRequest {
  74. message InitMessage {
  75. string broker = 1;
  76. }
  77. oneof message {
  78. InitMessage init = 1;
  79. BrokerStats stats = 2;
  80. }
  81. }
  82. message PublisherToPubBalancerResponse {
  83. }
  84. message BalanceTopicsRequest {
  85. }
  86. message BalanceTopicsResponse {
  87. }
  88. //////////////////////////////////////////////////
  89. message TopicRetention {
  90. int64 retention_seconds = 1; // retention duration in seconds
  91. bool enabled = 2; // whether retention is enabled
  92. }
  93. message ConfigureTopicRequest {
  94. schema_pb.Topic topic = 1;
  95. int32 partition_count = 2;
  96. schema_pb.RecordType record_type = 3;
  97. TopicRetention retention = 4;
  98. }
  99. message ConfigureTopicResponse {
  100. repeated BrokerPartitionAssignment broker_partition_assignments = 2;
  101. schema_pb.RecordType record_type = 3;
  102. TopicRetention retention = 4;
  103. }
  104. message ListTopicsRequest {
  105. }
  106. message ListTopicsResponse {
  107. repeated schema_pb.Topic topics = 1;
  108. }
  109. message LookupTopicBrokersRequest {
  110. schema_pb.Topic topic = 1;
  111. }
  112. message LookupTopicBrokersResponse {
  113. schema_pb.Topic topic = 1;
  114. repeated BrokerPartitionAssignment broker_partition_assignments = 2;
  115. }
  116. message BrokerPartitionAssignment {
  117. schema_pb.Partition partition = 1;
  118. string leader_broker = 2;
  119. string follower_broker = 3;
  120. }
  121. message GetTopicConfigurationRequest {
  122. schema_pb.Topic topic = 1;
  123. }
  124. message GetTopicConfigurationResponse {
  125. schema_pb.Topic topic = 1;
  126. int32 partition_count = 2;
  127. schema_pb.RecordType record_type = 3;
  128. repeated BrokerPartitionAssignment broker_partition_assignments = 4;
  129. int64 created_at_ns = 5;
  130. int64 last_updated_ns = 6;
  131. TopicRetention retention = 7;
  132. }
  133. message GetTopicPublishersRequest {
  134. schema_pb.Topic topic = 1;
  135. }
  136. message GetTopicPublishersResponse {
  137. repeated TopicPublisher publishers = 1;
  138. }
  139. message GetTopicSubscribersRequest {
  140. schema_pb.Topic topic = 1;
  141. }
  142. message GetTopicSubscribersResponse {
  143. repeated TopicSubscriber subscribers = 1;
  144. }
  145. message TopicPublisher {
  146. string publisher_name = 1;
  147. string client_id = 2;
  148. schema_pb.Partition partition = 3;
  149. int64 connect_time_ns = 4;
  150. int64 last_seen_time_ns = 5;
  151. string broker = 6;
  152. bool is_active = 7;
  153. int64 last_published_offset = 8;
  154. int64 last_acked_offset = 9;
  155. }
  156. message TopicSubscriber {
  157. string consumer_group = 1;
  158. string consumer_id = 2;
  159. string client_id = 3;
  160. schema_pb.Partition partition = 4;
  161. int64 connect_time_ns = 5;
  162. int64 last_seen_time_ns = 6;
  163. string broker = 7;
  164. bool is_active = 8;
  165. int64 current_offset = 9; // last acknowledged offset
  166. int64 last_received_offset = 10;
  167. }
  168. message AssignTopicPartitionsRequest {
  169. schema_pb.Topic topic = 1;
  170. repeated BrokerPartitionAssignment broker_partition_assignments = 2;
  171. bool is_leader = 3;
  172. bool is_draining = 4;
  173. }
  174. message AssignTopicPartitionsResponse {
  175. }
  176. message SubscriberToSubCoordinatorRequest {
  177. message InitMessage {
  178. string consumer_group = 1;
  179. string consumer_group_instance_id = 2;
  180. schema_pb.Topic topic = 3;
  181. // The consumer group instance will be assigned at most max_partition_count partitions.
  182. // If the number of partitions is less than the sum of max_partition_count,
  183. // the consumer group instance may be assigned partitions less than max_partition_count.
  184. // Default is 1.
  185. int32 max_partition_count = 4;
  186. // If consumer group instance changes, wait for rebalance_seconds before reassigning partitions
  187. // Exception: if adding a new consumer group instance and sum of max_partition_count equals the number of partitions,
  188. // the rebalance will happen immediately.
  189. // Default is 10 seconds.
  190. int32 rebalance_seconds = 5;
  191. }
  192. message AckUnAssignmentMessage {
  193. schema_pb.Partition partition = 1;
  194. }
  195. message AckAssignmentMessage {
  196. schema_pb.Partition partition = 1;
  197. }
  198. oneof message {
  199. InitMessage init = 1;
  200. AckAssignmentMessage ack_assignment = 2;
  201. AckUnAssignmentMessage ack_un_assignment = 3;
  202. }
  203. }
  204. message SubscriberToSubCoordinatorResponse {
  205. message Assignment {
  206. BrokerPartitionAssignment partition_assignment = 1;
  207. }
  208. message UnAssignment {
  209. schema_pb.Partition partition = 1;
  210. }
  211. oneof message {
  212. Assignment assignment = 1;
  213. UnAssignment un_assignment = 2;
  214. }
  215. }
  216. //////////////////////////////////////////////////
  217. message ControlMessage {
  218. bool is_close = 1;
  219. string publisher_name = 2;
  220. }
  221. message DataMessage {
  222. bytes key = 1;
  223. bytes value = 2;
  224. int64 ts_ns = 3;
  225. ControlMessage ctrl = 4;
  226. }
  227. message PublishMessageRequest {
  228. message InitMessage {
  229. schema_pb.Topic topic = 1;
  230. schema_pb.Partition partition = 2;
  231. int32 ack_interval = 3;
  232. string follower_broker = 4;
  233. string publisher_name = 5; // for debugging
  234. }
  235. oneof message {
  236. InitMessage init = 1;
  237. DataMessage data = 2;
  238. }
  239. }
  240. message PublishMessageResponse {
  241. int64 ack_sequence = 1;
  242. string error = 2;
  243. bool should_close = 3;
  244. }
  245. message PublishFollowMeRequest {
  246. message InitMessage {
  247. schema_pb.Topic topic = 1;
  248. schema_pb.Partition partition = 2;
  249. }
  250. message FlushMessage {
  251. int64 ts_ns = 1;
  252. }
  253. message CloseMessage {
  254. }
  255. oneof message {
  256. InitMessage init = 1;
  257. DataMessage data = 2;
  258. FlushMessage flush = 3;
  259. CloseMessage close = 4;
  260. }
  261. }
  262. message PublishFollowMeResponse {
  263. int64 ack_ts_ns = 1;
  264. }
  265. message SubscribeMessageRequest {
  266. message InitMessage {
  267. string consumer_group = 1;
  268. string consumer_id = 2;
  269. string client_id = 3;
  270. schema_pb.Topic topic = 4;
  271. schema_pb.PartitionOffset partition_offset = 5;
  272. schema_pb.OffsetType offset_type = 6;
  273. string filter = 10;
  274. string follower_broker = 11;
  275. int32 sliding_window_size = 12;
  276. }
  277. message AckMessage {
  278. int64 sequence = 1;
  279. bytes key = 2;
  280. }
  281. oneof message {
  282. InitMessage init = 1;
  283. AckMessage ack = 2;
  284. }
  285. }
  286. message SubscribeMessageResponse {
  287. message SubscribeCtrlMessage {
  288. string error = 1;
  289. bool is_end_of_stream = 2;
  290. bool is_end_of_topic = 3;
  291. }
  292. oneof message {
  293. SubscribeCtrlMessage ctrl = 1;
  294. DataMessage data = 2;
  295. }
  296. }
  297. message SubscribeFollowMeRequest {
  298. message InitMessage {
  299. schema_pb.Topic topic = 1;
  300. schema_pb.Partition partition = 2;
  301. string consumer_group = 3;
  302. }
  303. message AckMessage {
  304. int64 ts_ns = 1;
  305. }
  306. message CloseMessage {
  307. }
  308. oneof message {
  309. InitMessage init = 1;
  310. AckMessage ack = 2;
  311. CloseMessage close = 3;
  312. }
  313. }
  314. message SubscribeFollowMeResponse {
  315. int64 ack_ts_ns = 1;
  316. }
  317. message ClosePublishersRequest {
  318. schema_pb.Topic topic = 1;
  319. int64 unix_time_ns = 2;
  320. }
  321. message ClosePublishersResponse {
  322. }
  323. message CloseSubscribersRequest {
  324. schema_pb.Topic topic = 1;
  325. int64 unix_time_ns = 2;
  326. }
  327. message CloseSubscribersResponse {
  328. }
  329. //////////////////////////////////////////////////
  330. // SQL query support messages
  331. message GetUnflushedMessagesRequest {
  332. schema_pb.Topic topic = 1;
  333. schema_pb.Partition partition = 2;
  334. int64 start_buffer_index = 3; // Filter by buffer index (messages from buffers >= this index)
  335. }
  336. message GetUnflushedMessagesResponse {
  337. LogEntry message = 1; // Single message per response (streaming)
  338. string error = 2; // Error message if any
  339. bool end_of_stream = 3; // Indicates this is the final response
  340. }
  341. message LogEntry {
  342. int64 ts_ns = 1;
  343. bytes key = 2;
  344. bytes data = 3;
  345. uint32 partition_key_hash = 4;
  346. }