Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 14 additions & 2 deletions rabbitmq/include/userver/urabbitmq/admin_channel.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -45,9 +45,21 @@ class AdminChannel final : IAdminInterface {
DeclareExchange(exchange, Exchange::Type::kFanOut, {}, deadline);
}

void DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline) override;
QueueDeclareResponse DeclareQueue(
const Queue& queue,
utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers,
engine::Deadline deadline
) override;

void DeclareQueue(const Queue& queue, engine::Deadline deadline) override { DeclareQueue(queue, {}, deadline); }
QueueDeclareResponse
DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline) override {
return DeclareQueue(queue, flags, {}, deadline);
}

QueueDeclareResponse DeclareQueue(const Queue& queue, engine::Deadline deadline) override {
return DeclareQueue(queue, {}, deadline);
}

void BindQueue(
const Exchange& exchange,
Expand Down
15 changes: 13 additions & 2 deletions rabbitmq/include/userver/urabbitmq/broker_interface.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -42,11 +42,22 @@ class IAdminInterface {
///
/// @param queue name of the queue
/// @param flags queue flags
/// @param headers metadata table of the queue
/// @param deadline execution deadline
virtual void DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline) = 0;
/// @returns the broker's `queue.declare-ok` reply (queue name, message and
/// consumer counts)
virtual QueueDeclareResponse DeclareQueue(
const Queue& queue,
utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers,
engine::Deadline deadline) = 0;

/// @brief overload of DeclareQueue
virtual QueueDeclareResponse
DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline) = 0;

/// @brief overload of DeclareQueue
virtual void DeclareQueue(const Queue& queue, engine::Deadline deadline) = 0;
virtual QueueDeclareResponse DeclareQueue(const Queue& queue, engine::Deadline deadline) = 0;

/// @brief Bind a queue to an exchange.
///
Expand Down
16 changes: 14 additions & 2 deletions rabbitmq/include/userver/urabbitmq/client.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -53,9 +53,21 @@ class Client
DeclareExchange(exchange, Exchange::Type::kFanOut, {}, deadline);
}

void DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline) override;
QueueDeclareResponse DeclareQueue(
const Queue& queue,
utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers,
engine::Deadline deadline
) override;

void DeclareQueue(const Queue& queue, engine::Deadline deadline) override { DeclareQueue(queue, {}, deadline); }
QueueDeclareResponse
DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline) override {
return DeclareQueue(queue, flags, {}, deadline);
}

QueueDeclareResponse DeclareQueue(const Queue& queue, engine::Deadline deadline) override {
return DeclareQueue(queue, {}, {}, deadline);
}

void BindQueue(
const Exchange& exchange,
Expand Down
11 changes: 11 additions & 0 deletions rabbitmq/include/userver/urabbitmq/typedefs.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,17 @@ enum class MessageType {
/// This is not JSON, but a convenient tree representation for AMQP field values.
using HeaderValue = formats::json::Value;

/// @brief Result of a `DeclareQueue` call, mirroring the broker's
/// `queue.declare-ok` response.
struct QueueDeclareResponse {
/// name of the declared queue (useful for server-named queues)
std::string name;
/// number of messages ready for delivery in the queue at the time of the reply
std::uint32_t message_count{0};
/// number of consumers subscribed to the queue at the time of the reply
std::uint32_t consumer_count{0};
};

/// @brief Structure holding an AMQP message body along with some of its
/// metadata fields. This struct is used to pass messages to the end user,
/// hiding the actual AMQP message object implementation.
Expand Down
61 changes: 61 additions & 0 deletions rabbitmq/src/tests/admin_rmqtest.cpp
Original file line number Diff line number Diff line change
@@ -1,5 +1,10 @@
#include "utils_rmqtest.hpp"

#include <string>
#include <unordered_map>

#include <userver/formats/json/value_builder.hpp>

USERVER_NAMESPACE_BEGIN

UTEST(AdminChannel, DeclareRemoveExchange) {
Expand Down Expand Up @@ -46,4 +51,60 @@ UTEST(AdminChannel, DeclareRemoveQueue) {
new_channel.RemoveQueue(queue, client.GetDeadline());
}

UTEST(AdminChannel, DeclareQueueWithMaxLengthArgument) {
// This test checks, whether optional args for queue are working
ClientWrapper client{};
auto channel = client->GetAdminChannel(client.GetDeadline());

constexpr std::int64_t kMaxLength = 3;
const std::unordered_map<std::string, urabbitmq::HeaderValue> arguments{
{"x-max-length", urabbitmq::HeaderValue::Builder{std::int64_t{kMaxLength}}.ExtractValue()},
};

channel.DeclareExchange(client.GetExchange(), urabbitmq::Exchange::Type::kFanOut, {}, client.GetDeadline());
channel.DeclareQueue(client.GetQueue(), {}, arguments, client.GetDeadline());
channel.BindQueue(client.GetExchange(), client.GetQueue(), client.GetRoutingKey(), client.GetDeadline());

constexpr int kPublishCount = kMaxLength + 10;
for (int i = 0; i < kPublishCount; ++i) {
client->PublishReliable(
client.GetExchange(), client.GetRoutingKey(), "message-" + std::to_string(i), client.GetDeadline()
);
}
const auto response = channel.DeclareQueue(client.GetQueue(), {}, arguments, client.GetDeadline());
EXPECT_EQ(response.message_count, kMaxLength);

// The default (drop-head) overflow drops the oldest messages, so the queue
// keeps only the newest kMaxLength ones. Reading them back in FIFO order must
// yield that tail: message-10, message-11, message-12.
for (int i = 0; i < kMaxLength; ++i) {
const auto message = client->Get(client.GetQueue(), urabbitmq::Queue::Flags::kNoAck, client.GetDeadline());
EXPECT_EQ(message, "message-" + std::to_string(kPublishCount - kMaxLength + i));
}
}

UTEST(AdminChannel, DeclareQueueArgumentMismatchFails) {
// Here we are testing the same queue with different optional args
ClientWrapper client{};

const auto declare_with_max_length = [&client](std::int64_t max_length) {
// A fresh channel per attempt: a PRECONDITION_FAILED breaks the channel
// it happens on, so reusing one would mask the cause of later failures.
auto channel = client->GetAdminChannel(client.GetDeadline());
const std::unordered_map<std::string, urabbitmq::HeaderValue> arguments{
{"x-max-length", urabbitmq::HeaderValue::Builder{std::int64_t{max_length}}.ExtractValue()},
};
channel.DeclareQueue(client.GetQueue(), {}, arguments, client.GetDeadline());
};

// First declaration creates the queue with x-max-length == 3.
declare_with_max_length(3);

// Re-declaring the same queue with a different argument value is a conflict.
EXPECT_ANY_THROW(declare_with_max_length(5));

// Re-declaring with the very same value is idempotent and must not throw.
EXPECT_NO_THROW(declare_with_max_length(3));
}

USERVER_NAMESPACE_END
7 changes: 5 additions & 2 deletions rabbitmq/src/urabbitmq/admin_channel.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,11 @@ void AdminChannel::DeclareExchange(
ConnectionHelper::DeclareExchange(*impl_, exchange, type, flags, deadline).Wait(deadline);
}

void AdminChannel::DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline) {
ConnectionHelper::DeclareQueue(*impl_, queue, flags, deadline).Wait(deadline);
QueueDeclareResponse AdminChannel::DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers, engine::Deadline deadline) {
QueueDeclareResponse response;
ConnectionHelper::DeclareQueue(*impl_, queue, flags, headers, response, deadline).Wait(deadline);
return response;
}

void AdminChannel::BindQueue(
Expand Down
8 changes: 6 additions & 2 deletions rabbitmq/src/urabbitmq/client.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,13 @@ void Client::DeclareExchange(
awaiter.Wait(deadline);
}

void Client::DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline) {
auto awaiter = ConnectionHelper::DeclareQueue(impl_->GetConnection(deadline), queue, flags, deadline);
QueueDeclareResponse Client::DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers, engine::Deadline deadline) {
QueueDeclareResponse response;
auto awaiter =
ConnectionHelper::DeclareQueue(impl_->GetConnection(deadline), queue, flags, headers, response, deadline);
awaiter.Wait(deadline);
return response;
}

void Client::BindQueue(
Expand Down
6 changes: 5 additions & 1 deletion rabbitmq/src/urabbitmq/connection_helper.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,13 @@ impl::ResponseAwaiter ConnectionHelper::DeclareQueue(
const ConnectionPtr& connection,
const Queue& queue,
utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers,
QueueDeclareResponse& response,
engine::Deadline deadline
) {
return WithSpan("declare_queue", [&] { return connection->GetChannel().DeclareQueue(queue, flags, deadline); });
return WithSpan("declare_queue", [&] {
return connection->GetChannel().DeclareQueue(queue, flags, headers, response, deadline);
});
}

impl::ResponseAwaiter ConnectionHelper::BindQueue(
Expand Down
2 changes: 2 additions & 0 deletions rabbitmq/src/urabbitmq/connection_helper.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ class ConnectionHelper final {
const ConnectionPtr& connection,
const Queue& queue,
utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers,
QueueDeclareResponse& response,
engine::Deadline deadline
);

Expand Down
8 changes: 7 additions & 1 deletion rabbitmq/src/urabbitmq/impl/amqp_channel.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -132,13 +132,19 @@ ResponseAwaiter AmqpChannel::DeclareExchange(
ResponseAwaiter AmqpChannel::DeclareQueue(
const Queue& queue,
utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers,
QueueDeclareResponse& response,
engine::Deadline deadline
) {
auto awaiter = conn_.GetAwaiter(deadline);

{
auto channel = conn_.GetChannel(deadline);
awaiter.GetWrapper()->Wrap(channel->declareQueue(queue.GetUnderlying(), Convert(flags)));
AMQP::Table table;
AddHeadersToTable(table, headers);
awaiter.GetWrapper()->WrapDeclareQueue(
channel->declareQueue(queue.GetUnderlying(), Convert(flags), table), response
);
}

return awaiter;
Expand Down
8 changes: 7 additions & 1 deletion rabbitmq/src/urabbitmq/impl/amqp_channel.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,13 @@ class AmqpChannel final {
engine::Deadline deadline
);

ResponseAwaiter DeclareQueue(const Queue& queue, utils::Flags<Queue::Flags> flags, engine::Deadline deadline);
ResponseAwaiter DeclareQueue(
const Queue& queue,
utils::Flags<Queue::Flags> flags,
const std::unordered_map<std::string, HeaderValue>& headers,
QueueDeclareResponse& response,
engine::Deadline deadline
);

ResponseAwaiter BindQueue(
const Exchange& exchange,
Expand Down
11 changes: 11 additions & 0 deletions rabbitmq/src/urabbitmq/impl/deferred_wrapper.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

#include <urabbitmq/make_shared_enabler.hpp>

#include <userver/urabbitmq/typedefs.hpp>
#include <userver/utils/assert.hpp>

USERVER_NAMESPACE_BEGIN
Expand Down Expand Up @@ -60,6 +61,16 @@ void DeferredWrapper::WrapGet(AMQP::DeferredGet& deferred, std::string& message)
.onError([wrap = shared_from_this()](const char* error) { wrap->Fail(error); });
}

void DeferredWrapper::WrapDeclareQueue(AMQP::DeferredQueue& deferred, QueueDeclareResponse& response) {
deferred
.onSuccess([wrap = shared_from_this(),
&response](const std::string& name, uint32_t message_count, uint32_t consumer_count) {
response = QueueDeclareResponse{name, message_count, consumer_count};
wrap->Ok();
})
.onError([wrap = shared_from_this()](const char* error) { wrap->Fail(error); });
}

} // namespace urabbitmq::impl

USERVER_NAMESPACE_END
7 changes: 7 additions & 0 deletions rabbitmq/src/urabbitmq/impl/deferred_wrapper.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,15 @@
namespace AMQP {
class Deferred;
class DeferredGet;
class DeferredQueue;
} // namespace AMQP

USERVER_NAMESPACE_BEGIN

namespace urabbitmq {
struct QueueDeclareResponse;
} // namespace urabbitmq

namespace urabbitmq::impl {

class AmqpConnection;
Expand All @@ -31,6 +36,8 @@ class DeferredWrapper : public std::enable_shared_from_this<DeferredWrapper> {

void WrapGet(AMQP::DeferredGet& deferred, std::string& message);

void WrapDeclareQueue(AMQP::DeferredQueue& deferred, QueueDeclareResponse& response);

static std::shared_ptr<DeferredWrapper> Create();

protected:
Expand Down
Loading