2026-06-21 19:05:13 -07:00
|
|
|
syntax = "proto3";
|
|
|
|
|
|
|
|
|
|
package kafkamgmt.v1;
|
|
|
|
|
|
|
|
|
|
import "google/api/annotations.proto";
|
|
|
|
|
import "google/protobuf/timestamp.proto";
|
|
|
|
|
|
2026-06-21 19:40:12 -07:00
|
|
|
// Generated Go code is committed to this repo under gen/kafkamgmt/v1 and
|
|
|
|
|
// consumed directly via `go get`. Consumers do not run buf generate themselves.
|
|
|
|
|
option go_package = "forgejo.riotpiao.homelab.com/rock/kmsvc-proto/gen/kafkamgmt/v1;kafkamgmtv1";
|
2026-06-21 19:05:13 -07:00
|
|
|
|
|
|
|
|
// QueueService is the message-plane API for the Kafka Management Service.
|
|
|
|
|
// Queue lifecycle (create/delete/configure) is managed via the Queue CRD,
|
|
|
|
|
// not this service — see design.md §2a/§2b in kafaka_management_service.
|
|
|
|
|
service QueueService {
|
|
|
|
|
rpc SendMessage(SendMessageRequest) returns (SendMessageResponse) {
|
|
|
|
|
option (google.api.http) = {
|
|
|
|
|
post: "/v1/queues/{queue_name}/messages"
|
|
|
|
|
body: "*"
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
rpc SendMessageBatch(SendMessageBatchRequest) returns (SendMessageBatchResponse) {
|
|
|
|
|
option (google.api.http) = {
|
|
|
|
|
post: "/v1/queues/{queue_name}/messages:batch"
|
|
|
|
|
body: "*"
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
rpc ReceiveMessage(ReceiveMessageRequest) returns (ReceiveMessageResponse) {
|
|
|
|
|
option (google.api.http) = {
|
|
|
|
|
get: "/v1/queues/{queue_name}/messages"
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
rpc DeleteMessage(DeleteMessageRequest) returns (DeleteMessageResponse) {
|
|
|
|
|
option (google.api.http) = {
|
|
|
|
|
delete: "/v1/queues/{queue_name}/messages/{receipt_handle}"
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
rpc DeleteMessageBatch(DeleteMessageBatchRequest) returns (DeleteMessageBatchResponse) {
|
|
|
|
|
option (google.api.http) = {
|
|
|
|
|
post: "/v1/queues/{queue_name}/messages:batchDelete"
|
|
|
|
|
body: "*"
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
rpc ChangeMessageVisibility(ChangeMessageVisibilityRequest) returns (ChangeMessageVisibilityResponse) {
|
|
|
|
|
option (google.api.http) = {
|
|
|
|
|
patch: "/v1/queues/{queue_name}/messages/{receipt_handle}"
|
|
|
|
|
body: "*"
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message MessageAttributes {
|
|
|
|
|
map<string, string> values = 1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message SendMessageRequest {
|
|
|
|
|
string queue_name = 1;
|
|
|
|
|
bytes message_body = 2;
|
|
|
|
|
MessageAttributes message_attributes = 3;
|
|
|
|
|
string message_group_id = 4; // FIFO only
|
|
|
|
|
string message_deduplication_id = 5; // FIFO only
|
|
|
|
|
int32 delay_seconds = 6;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message SendMessageResponse {
|
|
|
|
|
string message_id = 1;
|
|
|
|
|
string sequence_number = 2; // FIFO only
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message SendMessageBatchEntry {
|
|
|
|
|
string id = 1;
|
|
|
|
|
bytes message_body = 2;
|
|
|
|
|
MessageAttributes message_attributes = 3;
|
|
|
|
|
string message_group_id = 4;
|
|
|
|
|
string message_deduplication_id = 5;
|
|
|
|
|
int32 delay_seconds = 6;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message SendMessageBatchRequest {
|
|
|
|
|
string queue_name = 1;
|
|
|
|
|
repeated SendMessageBatchEntry entries = 2;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message BatchResultEntry {
|
|
|
|
|
string id = 1;
|
|
|
|
|
string message_id = 2;
|
|
|
|
|
string error = 3;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message SendMessageBatchResponse {
|
|
|
|
|
repeated BatchResultEntry successful = 1;
|
|
|
|
|
repeated BatchResultEntry failed = 2;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message ReceiveMessageRequest {
|
|
|
|
|
string queue_name = 1;
|
|
|
|
|
int32 max_number_of_messages = 2; // <= 10
|
|
|
|
|
int32 wait_time_seconds = 3; // 0-20, long-poll
|
|
|
|
|
int32 visibility_timeout_seconds = 4; // override
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message Message {
|
|
|
|
|
string message_id = 1;
|
|
|
|
|
string receipt_handle = 2;
|
|
|
|
|
bytes body = 3;
|
|
|
|
|
MessageAttributes attributes = 4;
|
|
|
|
|
int32 receive_count = 5;
|
|
|
|
|
string message_group_id = 6;
|
|
|
|
|
google.protobuf.Timestamp enqueued_at = 7;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message ReceiveMessageResponse {
|
|
|
|
|
repeated Message messages = 1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message DeleteMessageRequest {
|
|
|
|
|
string queue_name = 1;
|
|
|
|
|
string receipt_handle = 2;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message DeleteMessageResponse {}
|
|
|
|
|
|
|
|
|
|
message DeleteMessageBatchEntry {
|
|
|
|
|
string id = 1;
|
|
|
|
|
string receipt_handle = 2;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message DeleteMessageBatchRequest {
|
|
|
|
|
string queue_name = 1;
|
|
|
|
|
repeated DeleteMessageBatchEntry entries = 2;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message DeleteMessageBatchResponse {
|
|
|
|
|
repeated BatchResultEntry successful = 1;
|
|
|
|
|
repeated BatchResultEntry failed = 2;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message ChangeMessageVisibilityRequest {
|
|
|
|
|
string queue_name = 1;
|
|
|
|
|
string receipt_handle = 2;
|
|
|
|
|
int32 visibility_timeout_seconds = 3;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
message ChangeMessageVisibilityResponse {}
|