KubeMQ
Client SDKsC++How-to guidesQueues

Ack & Reject

Selectively acknowledge or reject individual KubeMQ Queue messages with the C++ SDK to control redelivery.

Overview

Ack and reject give you per-message control over queue delivery instead of an all-or-nothing batch outcome. When receiver->Poll fetches a batch with auto_ack = false, each message stays in an open transaction on the broker — invisible to other consumers — until the consumer explicitly settles it. That's what you need when one bad record in a batch shouldn't take the rest down with it.

Settlement happens through two calls on the received message: dm.Ack(), which permanently removes it from the queue, and dm.Nack(), which returns it to the queue for redelivery. Internally the broker tracks this against a receive count, which a dead-letter policy can use to stop retrying a poison message forever.

Gotchas: an unsettled message isn't gone — it snaps back to the queue once the transaction expires, so a slow consumer looks identical to a rejecting one; settle every message before that deadline, and never assume a batch is fully processed until you've called Ack() or Nack() on each one individually.

Prerequisites

  • KubeMQ server running on localhost:50000
  • C++ SDK installed (vcpkg or CMake FetchContent)
  • C++17 compiler (GCC 9+, Clang 9+, MSVC 2019+)

Code

main.cc
// Example: queues/ack_reject
//
// Demonstrates individual message acknowledgment and rejection using
// the queue downstream receiver. Messages can be individually acked
// (confirmed) or rejected (nacked) back to the queue.
//
// Channel: cpp-queues.ack-reject
// Client ID: cpp-queues-ack-reject-client
//
// Run with a KubeMQ server on localhost:50000
// (see https://docs.kubemq.io/deploy).

#include <kubemq/kubemq.h>

#include <iostream>
#include <string>

int main() {
    std::cout << "[1] Connecting to localhost:50000" << std::endl;

    kubemq::ClientOptions options;
    options.set_address("localhost", 50000);
    options.set_client_id("cpp-queues-ack-reject-client");

    auto client_result = kubemq::Client::Create(options);
    if (!client_result.ok()) {
        std::cerr << "[ERROR] Failed to create client: " << client_result.status().message()
                  << std::endl;
        return 1;
    }
    auto& client = *client_result;

    std::string channel = "cpp-queues.ack-reject";

    // Send two messages.
    std::cout << "[2] Sending 2 messages to queue " << channel << std::endl;
    for (int i = 1; i <= 2; i++) {
        auto msg_result = kubemq::QueueMessage::Builder()
                              .SetChannel(channel)
                              .SetBody("msg-" + std::to_string(i))
                              .Build();
        if (!msg_result.ok()) {
            std::cerr << "[ERROR] Build message " << i << ": " << msg_result.status().message()
                      << std::endl;
            return 1;
        }
        auto send_result = client->SendQueueMessage(*msg_result);
        if (!send_result.ok()) {
            std::cerr << "[ERROR] SendQueueMessage: " << send_result.status().message()
                      << std::endl;
            return 1;
        }
        if (send_result->is_error) {
            std::cerr << "[ERROR] Send failed: " << send_result->error << std::endl;
            return 1;
        }
    }
    std::cout << "[3] Sent 2 messages" << std::endl;

    // Poll messages with manual acknowledgment (AutoAck=false).
    std::cout << "[4] Creating downstream receiver" << std::endl;
    auto receiver_result = client->NewQueueDownstreamReceiver();
    if (!receiver_result.ok()) {
        std::cerr << "[ERROR] NewQueueDownstreamReceiver: " << receiver_result.status().message()
                  << std::endl;
        return 1;
    }
    auto& receiver = *receiver_result;

    kubemq::PollRequest poll_req;
    poll_req.channel = channel;
    poll_req.max_items = 10;
    poll_req.wait_timeout_seconds = 5;
    poll_req.auto_ack = false;

    std::cout << "[5] Polling with manual ack" << std::endl;
    auto resp_result = receiver->Poll(poll_req);
    if (!resp_result.ok()) {
        std::cerr << "[ERROR] Poll: " << resp_result.status().message() << std::endl;
        receiver->Close();
        auto close_status = client->Close();
        if (!close_status.ok()) {
            std::cerr << "[ERROR] Close failed: " << close_status.message() << std::endl;
        }
        return 1;
    }

    auto& messages = resp_result->messages();
    std::cout << "[6] Polled " << messages.size() << " messages" << std::endl;

    // Ack the first message, Nack (reject) the second.
    for (size_t i = 0; i < messages.size(); ++i) {
        auto& dm = messages[i];
        std::cout << "  [" << i + 1 << "] body=" << dm.message().body() << std::endl;
        if (i == 0) {
            // Acknowledge: message is removed from the queue.
            auto ack_status = dm.Ack();
            if (!ack_status.ok()) {
                std::cerr << "[ERROR] Ack failed: " << ack_status.message() << std::endl;
            } else {
                std::cout << "  [" << i + 1 << "] Acknowledged" << std::endl;
            }
        } else {
            // Nack (reject): message is returned to the queue for redelivery.
            auto nack_status = dm.Nack();
            if (!nack_status.ok()) {
                std::cerr << "[ERROR] Nack failed: " << nack_status.message() << std::endl;
            } else {
                std::cout << "  [" << i + 1 << "] Rejected (returned to queue)" << std::endl;
            }
        }
    }

    // Close the downstream receiver explicitly.
    // Note: The QueueDownstreamReceiver destructor also calls Close(), but explicit
    // cleanup is shown here for clarity and to match Go's defer pattern.
    receiver->Close();

    auto close_status = client->Close();
    if (!close_status.ok()) {
        std::cerr << "[ERROR] Close failed: " << close_status.message() << std::endl;
        return 1;
    }
    std::cout << "[7] Client closed" << std::endl;

    return 0;
}

How It Works

  • Demonstrates receiving messages and selectively acknowledging or rejecting them.
  • Uses the simple queue API with manual acknowledgment control.
  • Rejected messages can be re-delivered or moved to a dead-letter queue.
  • Provides fine-grained control over message processing outcomes.

Was this page helpful?

On this page