KubeMQ
Client SDKsC++Reference

Queues

Queue messaging methods (Simple and Stream APIs) -- KubeMQ C++ SDK reference.

Guaranteed message delivery with acknowledgment. Two APIs are available: Simple (one-shot operations) and Stream (persistent streams with per-message settlement).

Simple API

SendQueueMessage

[[nodiscard]] StatusOr<QueueSendResult> SendQueueMessage(const QueueMessage& message);

Send a single message to a queue. Returns message ID and timestamps on success.

SendQueueMessages

[[nodiscard]] StatusOr<std::vector<QueueSendResult>> SendQueueMessages(
    const std::vector<QueueMessage>& messages);

Send a batch of messages. Returns a result for each message.

ReceiveQueueMessages

[[nodiscard]] StatusOr<ReceiveQueueMessagesResponse> ReceiveQueueMessages(
    const ReceiveQueueMessagesRequest& request);

Receive messages from a queue (simple, one-shot).

AckAllQueueMessages

[[nodiscard]] StatusOr<AckAllQueueMessagesResponse> AckAllQueueMessages(
    const AckAllQueueMessagesRequest& request);

Acknowledge all pending messages in a queue.

SendQueueMessageSimple

[[nodiscard]] StatusOr<QueueSendResult> SendQueueMessageSimple(
    const std::string& channel, const std::string& body,
    int expiration_seconds = 0, int delay_seconds = 0);

Convenience method with minimal parameters.

Stream API

QueueUpstream

[[nodiscard]] StatusOr<std::unique_ptr<QueueUpstreamHandle>> QueueUpstream(
    std::function<void(const QueueUpstreamResult&)> on_result,
    std::function<void(const Status&)> on_error);

Open an upstream stream for batch sending. Use handle->Send(request_id, messages) to send batches.

NewQueueDownstreamReceiver

[[nodiscard]] StatusOr<std::unique_ptr<QueueDownstreamReceiver>> NewQueueDownstreamReceiver();

Create a downstream receiver for polling and settling messages.

PollQueue

[[nodiscard]] StatusOr<PollResponse> PollQueue(const PollRequest& request);

Poll a queue for messages (convenience wrapper around downstream receiver).

Types

QueueMessage

Built via QueueMessage::Builder.

MethodTypeRequiredDescription
SetChannel(s)stringYesTarget queue channel
SetBody(s)stringNoMessage body
SetMetadata(s)stringNoString metadata
SetTags(m)map<string,string>NoKey-value tags
SetExpirationSeconds(n)intNoMessage TTL (0 = no expiry)
SetDelaySeconds(n)intNoDelay before delivery
SetMaxReceiveCount(n)intNoMax receives before DLQ
SetMaxReceiveQueue(s)stringNoDead-letter queue channel

PollRequest

FieldTypeDefaultDescription
channelstring--Queue channel to poll
max_itemsint1Max messages to receive
wait_timeout_secondsint5Wait timeout
auto_ackboolfalseAuto-acknowledge on receipt

PollResponse

MethodDescription
messages()Vector of QueueDownstreamMessage
is_error()Whether the response is an error
error()Error message
AckAll()Acknowledge all messages
NackAll()Negative-acknowledge all
ReQueueAll(ch)Re-queue all to another channel

QueueDownstreamMessage

MethodDescription
message()The underlying QueueMessage
sequence()Sequence number
transaction_id()Transaction identifier
Ack()Acknowledge this message
Nack()Negative-acknowledge
ReQueue(ch)Re-queue to another channel

Quick Usage

queues.cc
// Simple send
auto msg_or = kubemq::QueueMessage::Builder()
    .SetChannel("tasks")
    .SetBody("process-order")
    .Build();
auto result = client->SendQueueMessage(*msg_or);

// Poll with auto-ack
kubemq::PollRequest req;
req.channel = "tasks";
req.max_items = 10;
req.wait_timeout_seconds = 5;
req.auto_ack = true;
auto poll = client->PollQueue(req);

// Stream upstream
auto upstream = client->QueueUpstream(on_result, on_error);
(*upstream)->Send("req-1", messages);

See Also

Was this page helpful?

On this page