This is an automated email from the ASF dual-hosted git repository.
numinnex pushed a commit to branch rework_topic_commands
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/rework_topic_commands by this
push:
new edaf3630f fix CI
edaf3630f is described below
commit edaf3630ff69d3e4d8b812076fe156b5f489b2d3
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Fri Aug 14 11:17:35 2026 +0200
fix CI
---
examples/go/getting-started/producer/main.go | 3 +-
foreign/cpp/tests/e2e/client.cpp | 56 ++++----
foreign/cpp/tests/e2e/consumer_group.cpp | 132 +++++++++---------
foreign/cpp/tests/e2e/partition.cpp | 44 +++---
foreign/cpp/tests/e2e/stream.cpp | 16 +--
foreign/cpp/tests/e2e/topic.cpp | 149 +++++++++++----------
.../common/provision/ResourceProvisioner.java | 1 -
.../iggy/client/async/tcp/TopicsTcpClient.java | 17 +--
.../apache/iggy/client/blocking/TopicsClient.java | 14 +-
.../client/blocking/http/TopicsHttpClient.java | 3 +-
.../iggy/client/blocking/tcp/TopicsTcpClient.java | 4 +-
.../client/async/AsyncClientIntegrationTest.java | 21 +--
.../client/async/AsyncConnectionPoolAuthTest.java | 1 -
.../iggy/client/async/AsyncConsumerGroupsTest.java | 16 +--
.../async/tcp/TopicsTcpClientPayloadTest.java | 17 +--
.../iggy/client/blocking/IntegrationTest.java | 1 -
.../iggy/client/blocking/TopicsClientBaseTest.java | 1 -
.../client/blocking/tcp/TlsConnectionTest.java | 9 +-
18 files changed, 221 insertions(+), 284 deletions(-)
diff --git a/examples/go/getting-started/producer/main.go
b/examples/go/getting-started/producer/main.go
index 129e454fe..74bd36686 100644
--- a/examples/go/getting-started/producer/main.go
+++ b/examples/go/getting-started/producer/main.go
@@ -83,8 +83,7 @@ func initSystem(ctx context.Context, client iggcon.Client) {
1,
iggcon.CompressionAlgorithmNone,
iggcon.IggyExpiryNeverExpire,
- 0,
- nil); err != nil {
+ 0); err != nil {
log.Printf("WARN: Topic already exists and will not be created
again or error: %v", err)
}
log.Println("Topic was created.")
diff --git a/foreign/cpp/tests/e2e/client.cpp b/foreign/cpp/tests/e2e/client.cpp
index b0d67381b..3185376c9 100644
--- a/foreign/cpp/tests/e2e/client.cpp
+++ b/foreign/cpp/tests/e2e/client.cpp
@@ -572,8 +572,8 @@ TEST_F(LowLevelE2E_Client,
FlushUnsavedBufferThrowsForExistingPartition) {
ASSERT_NO_THROW(client->create_stream(stream_name));
auto stream = client->get_stream(make_string_identifier(stream_name));
TrackStream(stream.id);
- ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
rust::Vec<iggy::ffi::IggyMessageToSend> messages;
messages.push_back(iggy::ffi::make_message(to_payload("flush-me"),
rust::Vec<iggy::ffi::HeaderEntry>()));
@@ -595,8 +595,8 @@ TEST_F(LowLevelE2E_Client,
FlushUnsavedBufferThrowsForExistingEmptyPartition) {
ASSERT_NO_THROW(client->create_stream(stream_name));
auto stream = client->get_stream(make_string_identifier(stream_name));
TrackStream(stream.id);
- ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
ASSERT_THROW(client->flush_unsaved_buffer(make_numeric_identifier(stream.id),
make_numeric_identifier(0), 0, true),
std::exception);
@@ -627,8 +627,8 @@ TEST_F(LowLevelE2E_Client,
FlushUnsavedBufferOnNonExistentStreamThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
ASSERT_THROW(
client->flush_unsaved_buffer(make_string_identifier(GetRandomName()),
make_numeric_identifier(0), 0, true),
@@ -643,8 +643,8 @@ TEST_F(LowLevelE2E_Client,
FlushUnsavedBufferOnNonExistentTopicThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
ASSERT_THROW(client->flush_unsaved_buffer(make_string_identifier(stream_name),
make_string_identifier(GetRandomName()), 0, true),
@@ -660,8 +660,8 @@ TEST_F(LowLevelE2E_Client,
FlushUnsavedBufferAfterStreamDeletedThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
auto stream = client->get_stream(make_string_identifier(stream_name));
TrackStream(stream.id);
- ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
const std::uint32_t saved_stream_id = stream.id;
ASSERT_NO_THROW(client->delete_stream(make_numeric_identifier(saved_stream_id)));
@@ -681,8 +681,8 @@ TEST_F(LowLevelE2E_Client,
FlushUnsavedBufferAfterTopicDeletedThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
auto stream = client->get_stream(make_string_identifier(stream_name));
TrackStream(stream.id);
- ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
ASSERT_NO_THROW(client->delete_topic(make_numeric_identifier(stream.id),
make_string_identifier(topic_name)));
ASSERT_THROW(
@@ -700,8 +700,8 @@ TEST_F(LowLevelE2E_Client, FlushUnsavedBufferTwiceThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
auto stream = client->get_stream(make_string_identifier(stream_name));
TrackStream(stream.id);
- ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
rust::Vec<iggy::ffi::IggyMessageToSend> messages;
messages.push_back(iggy::ffi::make_message(to_payload("flush-twice"),
rust::Vec<iggy::ffi::HeaderEntry>()));
@@ -723,8 +723,8 @@ TEST_F(LowLevelE2E_Client,
FlushUnsavedBufferWithInvalidPartitionIdsThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
auto stream = client->get_stream(make_string_identifier(stream_name));
TrackStream(stream.id);
- ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream.id),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
const std::uint32_t invalid_partition_ids[] = {1u, 9999u,
static_cast<std::uint32_t>(-1)};
for (const std::uint32_t invalid_partition_id : invalid_partition_ids) {
@@ -771,8 +771,8 @@ TEST_F(LowLevelE2E_Client,
DeleteSegmentsOnNonExistentStreamThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
ASSERT_THROW(
client->delete_segments(make_string_identifier(missing_stream_name),
make_string_identifier(topic_name), 0, 1),
@@ -788,8 +788,8 @@ TEST_F(LowLevelE2E_Client,
DeleteSegmentsOnNonExistentTopicThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
ASSERT_THROW(
client->delete_segments(make_string_identifier(stream_name),
make_string_identifier(missing_topic_name), 0, 1),
@@ -804,8 +804,8 @@ TEST_F(LowLevelE2E_Client,
DeleteSegmentsOnNonExistentPartitionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
ASSERT_THROW(
client->delete_segments(make_string_identifier(stream_name),
make_string_identifier(topic_name), 999, 1),
@@ -820,8 +820,8 @@ TEST_F(LowLevelE2E_Client,
DeleteSegmentsWithZeroCountIsNoOp) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
std::uint32_t stream_id = 0;
std::uint32_t topic_id = 0;
@@ -901,8 +901,8 @@ TEST_F(LowLevelE2E_Client,
DeleteSegmentsWhenOnlyActiveSegmentRemainsIsNoOp) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "never_expire", 0,
+ "server_default"));
std::uint32_t stream_id = 0;
std::uint32_t topic_id = 0;
@@ -1167,8 +1167,8 @@ TEST_F(LowLevelE2E_Client,
GetMeReflectsConsumerGroupMembershipChanges) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
const auto stream_details =
client->get_stream(make_string_identifier(stream_name));
ASSERT_EQ(stream_details.topics.size(), 1u);
diff --git a/foreign/cpp/tests/e2e/consumer_group.cpp
b/foreign/cpp/tests/e2e/consumer_group.cpp
index 3a50fe4d2..ead2d7e41 100644
--- a/foreign/cpp/tests/e2e/consumer_group.cpp
+++ b/foreign/cpp/tests/e2e/consumer_group.cpp
@@ -36,8 +36,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
CreateConsumerGroupSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW({
const auto group =
client->create_consumer_group(make_string_identifier(stream_name),
@@ -59,8 +59,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
CreateConsumerGroupOnNonExistentResourcesThrow
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_THROW(client->create_consumer_group(make_string_identifier(missing_stream_name),
make_string_identifier(topic_name), GetRandomName()),
@@ -79,8 +79,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
CreateConsumerGroupTwiceOnSameInputThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -97,8 +97,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
CreateConsumerGroupWithInvalidNamesThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
const std::string invalid_names[] = {"", std::string(256, 'a')};
for (const std::string &invalid_name : invalid_names) {
@@ -118,8 +118,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
CreateConsumerGroupAfterStreamDeletionThrows)
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->delete_stream(make_string_identifier(stream_name)));
ForgetTrackedStream(stream_name);
@@ -167,8 +167,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupReturnsSameInfoAsCreateConsume
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
const auto created_group =
client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name);
@@ -194,8 +194,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupsReturnsCreatedGroups) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), first_group_name));
@@ -248,8 +248,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
JoinConsumerGroupSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -304,8 +304,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
JoinConsumerGroupOnNonExistentResourcesThrows)
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), created_group_name));
TrackConsumerGroup(stream_name, topic_name, created_group_name);
@@ -332,8 +332,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
JoinConsumerGroupAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -354,8 +354,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
JoinConsumerGroupAfterTopicDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -375,8 +375,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
JoinConsumerGroupReflectsInGetConsumerGroup) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::ConsumerGroupDetails created_group;
ASSERT_NO_THROW({
@@ -407,8 +407,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
JoinConsumerGroupTwiceKeepsSingleMember) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -434,8 +434,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
JoinConsumerGroupFromTwoClientsIncreasesMember
ASSERT_NO_THROW(first->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(first->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
- 0, "server_default"));
+ ASSERT_NO_THROW(first->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default", 0,
+ "server_default"));
ASSERT_NO_THROW(first->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -460,8 +460,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
JoinConsumerGroupThenLeaveRestoresMembersCount
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -492,8 +492,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
LeaveConsumerGroupReducesMembersCount) {
ASSERT_NO_THROW(first->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(first->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
- 0, "server_default"));
+ ASSERT_NO_THROW(first->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default", 0,
+ "server_default"));
ASSERT_NO_THROW(first->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -565,8 +565,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
LeaveConsumerGroupOnNonExistentResourcesThrows
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), created_group_name));
TrackConsumerGroup(stream_name, topic_name, created_group_name);
@@ -595,8 +595,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
LeaveConsumerGroupAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -619,8 +619,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
LeaveConsumerGroupAfterTopicDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -643,8 +643,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
LeaveConsumerGroupTwiceThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -667,8 +667,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
LeaveConsumerGroupWithoutJoiningThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -688,8 +688,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupsReflectsJoinedGroupMembersCou
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::ConsumerGroupDetails joined_group;
iggy::ffi::ConsumerGroupDetails other_group;
@@ -758,8 +758,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupsIsStableAcrossBackToBackCalls
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), first_group_name));
@@ -795,8 +795,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupsReturnsCorrectNumberOfGroups)
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), deleted_group_name));
@@ -847,8 +847,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupsAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), first_group_name));
TrackConsumerGroup(stream_name, topic_name, first_group_name);
@@ -874,8 +874,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupsAfterTopicDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), first_group_name));
TrackConsumerGroup(stream_name, topic_name, first_group_name);
@@ -937,8 +937,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupOnNonExistentResourcesThrows)
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), created_group_name));
TrackConsumerGroup(stream_name, topic_name, created_group_name);
@@ -965,8 +965,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
GetConsumerGroupAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
ASSERT_NO_THROW(client->delete_stream(make_string_identifier(stream_name)));
@@ -986,8 +986,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
DeleteConsumerGroupSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -1047,8 +1047,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
DeleteConsumerGroupOnNonExistentResourcesThrow
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), created_group_name));
TrackConsumerGroup(stream_name, topic_name, created_group_name);
@@ -1074,8 +1074,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
DeleteConsumerGroupTwiceThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -1097,8 +1097,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
DeleteConsumerGroupAfterStreamDeletionThrows)
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
TrackConsumerGroup(stream_name, topic_name, group_name);
@@ -1121,8 +1121,8 @@ TEST_F(LowLevelE2E_ConsumerGroup,
DeleteConsumerGroupAndRecreateWithSameNameSucc
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->create_consumer_group(make_string_identifier(stream_name),
make_string_identifier(topic_name), group_name));
diff --git a/foreign/cpp/tests/e2e/partition.cpp
b/foreign/cpp/tests/e2e/partition.cpp
index 8edba3ad7..d58ed97fb 100644
--- a/foreign/cpp/tests/e2e/partition.cpp
+++ b/foreign/cpp/tests/e2e/partition.cpp
@@ -36,8 +36,8 @@ TEST_F(LowLevelE2E_Partition, CreatePartitionsSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(
client->create_partitions(make_string_identifier(stream_name),
make_string_identifier(topic_name), 43));
@@ -79,8 +79,8 @@ TEST_F(LowLevelE2E_Partition,
CreatePartitionsOnNonExistentResourcesThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_THROW(
client->create_partitions(make_string_identifier(missing_stream_name),
make_string_identifier(topic_name), 1),
@@ -99,8 +99,8 @@ TEST_F(LowLevelE2E_Partition,
CreatePartitionsWithInvalidIdentifiersThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Identifier invalid_stream_kind_id;
invalid_stream_kind_id.kind = "invalid";
@@ -201,8 +201,8 @@ TEST_F(LowLevelE2E_Partition,
CreatePartitionsWithNumericIdentifiersSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
const auto stream_details =
client->get_stream(make_string_identifier(stream_name));
ASSERT_EQ(stream_details.topics.size(), 1u);
@@ -228,8 +228,8 @@ TEST_F(LowLevelE2E_Partition, DeletePartitionsSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 44, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 44, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(
client->delete_partitions(make_string_identifier(stream_name),
make_string_identifier(topic_name), 43));
@@ -311,8 +311,8 @@ TEST_F(LowLevelE2E_Partition,
DeletePartitionsBeforeCreatingAdditionalPartitions
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(
client->delete_partitions(make_string_identifier(stream_name),
make_string_identifier(topic_name), 1));
@@ -333,8 +333,8 @@ TEST_F(LowLevelE2E_Partition,
DeletePartitionsFromTopicWithZeroPartitionsThrows)
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 0, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 0, "none", "server_default",
+ 0, "server_default"));
ASSERT_THROW(client->delete_partitions(make_string_identifier(stream_name),
make_string_identifier(topic_name), 1),
std::exception);
@@ -377,8 +377,8 @@ TEST_F(LowLevelE2E_Partition,
DeletePartitionsOnNonExistentResourcesThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none", "server_default",
+ 0, "server_default"));
ASSERT_THROW(
client->delete_partitions(make_string_identifier(missing_stream_name),
make_string_identifier(topic_name), 1),
@@ -397,8 +397,8 @@ TEST_F(LowLevelE2E_Partition,
DeletePartitionsWithInvalidIdentifiersThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Identifier invalid_stream_kind_id;
invalid_stream_kind_id.kind = "invalid";
@@ -438,8 +438,8 @@ TEST_F(LowLevelE2E_Partition,
DeletePartitionsTwiceForSameTopicSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 45, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 45, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(
client->delete_partitions(make_string_identifier(stream_name),
make_string_identifier(topic_name), 20));
ASSERT_NO_THROW(
@@ -462,8 +462,8 @@ TEST_F(LowLevelE2E_Partition,
DeletePartitionsAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none", "server_default",
+ 0, "server_default"));
const auto stream_details =
client->get_stream(make_string_identifier(stream_name));
ASSERT_EQ(stream_details.topics.size(), 1u);
diff --git a/foreign/cpp/tests/e2e/stream.cpp b/foreign/cpp/tests/e2e/stream.cpp
index 6a045230f..b1073b365 100644
--- a/foreign/cpp/tests/e2e/stream.cpp
+++ b/foreign/cpp/tests/e2e/stream.cpp
@@ -289,8 +289,8 @@ TEST_F(LowLevelE2E_Stream, UpdateStreamOnlyChangesName) {
ForgetTrackedStream(stream_name);
TrackStream(stream_id);
- ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream_id),
topic_name, 2, "none", "never_expire",
- 0, "server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_numeric_identifier(stream_id),
topic_name, 2, "none", "never_expire", 0,
+ "server_default"));
rust::Vec<iggy::ffi::IggyMessageToSend> messages;
for (std::uint32_t i = 0; i < 3; ++i) {
@@ -724,8 +724,8 @@ TEST_F(LowLevelE2E_Stream,
PurgeStreamPreservesStreamMetadata) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
first_topic_name, 2, "gzip",
- "duration", 1000, "1GiB"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
first_topic_name, 2, "gzip", "duration",
+ 1000, "1GiB"));
ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
second_topic_name, 3, "none",
"never_expire", 0, "server_default"));
@@ -971,8 +971,8 @@ TEST_F(LowLevelE2E_Stream,
PurgeStreamThenSendMessagesAgainSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
const auto created_stream =
client->get_stream(make_string_identifier(stream_name));
ASSERT_EQ(created_stream.topics.size(), 1u);
@@ -1007,8 +1007,8 @@ TEST_F(LowLevelE2E_Stream,
PurgeStreamTwiceKeepsStreamEmptyAndTopicsIntact) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
const auto created_stream =
client->get_stream(make_string_identifier(stream_name));
ASSERT_EQ(created_stream.topics.size(), 1u);
diff --git a/foreign/cpp/tests/e2e/topic.cpp b/foreign/cpp/tests/e2e/topic.cpp
index b329ede5d..d3031af07 100644
--- a/foreign/cpp/tests/e2e/topic.cpp
+++ b/foreign/cpp/tests/e2e/topic.cpp
@@ -59,13 +59,13 @@ TEST_F(LowLevelE2E_Topic,
CreateTopicWithAllOptionCombinations) {
for (const auto &expiry_option : expiry_options) {
for (const auto &max_topic_size : max_topic_sizes) {
const std::string topic_name = GetRandomName();
- SCOPED_TRACE(
- "compression=" + compression_algorithm + ", expiry_kind="
+ expiry_option.kind +
- ", expiry_value=" + std::to_string(expiry_option.value) +
", max_topic_size=" + max_topic_size);
+ SCOPED_TRACE("compression=" + compression_algorithm + ",
expiry_kind=" + expiry_option.kind +
+ ", expiry_value=" +
std::to_string(expiry_option.value) +
+ ", max_topic_size=" + max_topic_size);
ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1,
- compression_algorithm,
expiry_option.kind,
- expiry_option.value,
max_topic_size));
+ compression_algorithm,
expiry_option.kind, expiry_option.value,
+ max_topic_size));
++expected_topics_count;
expected_topic_names.insert(topic_name);
}
@@ -100,7 +100,8 @@ TEST_F(LowLevelE2E_Topic,
CreateTopicWithBoundaryPartitionsCountValues) {
ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
zero_partitions_topic_name, 0, "none",
"server_default", 0,
"server_default"));
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
max_partitions_topic_name, 1000, "none", "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
max_partitions_topic_name, 1000, "none",
+ "server_default", 0,
"server_default"));
ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
overflow_topic_name, 1001, "none",
"server_default", 0, "server_default"),
std::exception);
@@ -134,8 +135,8 @@ TEST_F(LowLevelE2E_Topic,
CreateTopicWithInvalidNamesThrows) {
};
for (const auto &topic_name : illegal_topic_names) {
SCOPED_TRACE(topic_name);
- ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"),
+ ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"),
std::exception);
}
@@ -153,10 +154,10 @@ TEST_F(LowLevelE2E_Topic, CreateDuplicateTopicThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
- ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
- 0, "server_default"),
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
+ ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default", 0,
+ "server_default"),
std::exception);
}
@@ -212,9 +213,9 @@ TEST_F(LowLevelE2E_Topic,
CreateTopicWithMaxTopicSizeBelowSegmentSizeThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
- 0, "1024"),
- std::exception);
+ ASSERT_THROW(
+ client->create_topic(make_string_identifier(stream_name), topic_name,
1, "none", "server_default", 0, "1024"),
+ std::exception);
}
TEST_F(LowLevelE2E_Topic, CreateTopicOnNonExistentStreamThrows) {
@@ -224,8 +225,8 @@ TEST_F(LowLevelE2E_Topic,
CreateTopicOnNonExistentStreamThrows) {
iggy::ffi::Client *client = GetLoggedInClient();
- ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
- 0, "server_default"),
+ ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default", 0,
+ "server_default"),
std::exception);
}
@@ -241,8 +242,8 @@ TEST_F(LowLevelE2E_Topic,
CreateTopicAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->delete_stream(make_string_identifier(stream_name)));
ForgetTrackedStream(stream_name);
- ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
- 0, "server_default"),
+ ASSERT_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default", 0,
+ "server_default"),
std::exception);
}
@@ -269,8 +270,8 @@ TEST_F(LowLevelE2E_Topic,
CreateTopicWithInvalidStreamIdentifierThrows) {
invalid_numeric_id.kind = "numeric";
invalid_numeric_id.length = 1;
invalid_numeric_id.value.push_back(1);
- ASSERT_THROW(client->create_topic(std::move(invalid_numeric_id),
second_topic_name, 1, "none", "server_default",
- 0, "server_default"),
+ ASSERT_THROW(client->create_topic(std::move(invalid_numeric_id),
second_topic_name, 1, "none", "server_default", 0,
+ "server_default"),
std::exception);
}
@@ -306,8 +307,8 @@ TEST_F(LowLevelE2E_Topic, DeleteTopicAfterCreate) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->delete_topic(make_string_identifier(stream_name),
make_string_identifier(topic_name)));
@@ -351,8 +352,8 @@ TEST_F(LowLevelE2E_Topic, DeleteTopicTwiceThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->delete_topic(make_string_identifier(stream_name),
make_string_identifier(topic_name)));
ASSERT_THROW(client->delete_topic(make_string_identifier(stream_name),
make_string_identifier(topic_name)),
@@ -368,8 +369,8 @@ TEST_F(LowLevelE2E_Topic,
DeleteTopicAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->delete_stream(make_string_identifier(stream_name)));
ForgetTrackedStream(stream_name);
@@ -386,8 +387,8 @@ TEST_F(LowLevelE2E_Topic, DeleteTopicBeforeLoginThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Client *unauthenticated_client = GetLoggedOutClient();
@@ -414,8 +415,8 @@ TEST_F(LowLevelE2E_Topic,
DeleteTopicWithInvalidStreamIdentifierThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Identifier invalid_kind_id;
invalid_kind_id.kind = "invalid";
@@ -440,8 +441,8 @@ TEST_F(LowLevelE2E_Topic,
DeleteTopicWithInvalidTopicIdentifierThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Identifier invalid_kind_id;
invalid_kind_id.kind = "invalid";
@@ -489,8 +490,8 @@ TEST_F(LowLevelE2E_Topic, GetTopicBeforeLoginThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Client *unauthenticated_client = GetLoggedOutClient();
@@ -535,8 +536,8 @@ TEST_F(LowLevelE2E_Topic, GetTopicWithWrongTopicThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_THROW(client->get_topic(make_string_identifier(stream_name),
make_string_identifier(wrong_topic_name)),
std::exception);
@@ -550,8 +551,8 @@ TEST_F(LowLevelE2E_Topic,
GetTopicAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->delete_stream(make_string_identifier(stream_name)));
ForgetTrackedStream(stream_name);
@@ -567,8 +568,8 @@ TEST_F(LowLevelE2E_Topic, GetTopicAfterTopicDeletionThrows)
{
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->delete_topic(make_string_identifier(stream_name),
make_string_identifier(topic_name)));
ASSERT_THROW(client->get_topic(make_string_identifier(stream_name),
make_string_identifier(topic_name)),
@@ -583,8 +584,8 @@ TEST_F(LowLevelE2E_Topic,
GetTopicReturnsEmptyPartitionsForZeroPartitionTopic) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 0, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 0, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW({
const auto topic_details =
@@ -694,8 +695,8 @@ TEST_F(LowLevelE2E_Topic,
GetTopicsReturnsCreatedTopicInputFields) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
first_topic_name, 2, "gzip",
- "duration", 1000, "1GiB"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
first_topic_name, 2, "gzip", "duration",
+ 1000, "1GiB"));
ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
second_topic_name, 0, "none",
"never_expire", 0, "unlimited"));
@@ -849,7 +850,8 @@ TEST_F(LowLevelE2E_Topic,
UpdateTopicDoesNotChangePartitionsCount) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
original_topic, partitions_count, "none", "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
original_topic, partitions_count, "none",
+ "server_default", 0,
"server_default"));
ASSERT_NO_THROW(client->update_topic(make_string_identifier(stream_name),
make_string_identifier(original_topic),
updated_topic_name, "gzip",
"duration", 1000, "1GiB"));
@@ -921,21 +923,20 @@ TEST_F(LowLevelE2E_Topic,
UpdateTopicWithAllOptionCombinationsUpdatesInputFields
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 2, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 2, "none", "server_default",
+ 0, "server_default"));
for (const auto &compression_algorithm : compression_algorithms) {
for (const auto &expiry_option : expiry_options) {
for (const auto &max_topic_size : max_topic_sizes) {
const std::string updated_topic_name = GetRandomName();
- SCOPED_TRACE(
- "compression=" + compression_algorithm + ", expiry_kind="
+ expiry_option.kind +
- ", expiry_value=" + std::to_string(expiry_option.value) +
", max_topic_size=" + max_topic_size);
-
-
ASSERT_NO_THROW(client->update_topic(make_string_identifier(stream_name),
-
make_string_identifier(topic_name), updated_topic_name,
- compression_algorithm,
expiry_option.kind,
- expiry_option.value,
max_topic_size));
+ SCOPED_TRACE("compression=" + compression_algorithm + ",
expiry_kind=" + expiry_option.kind +
+ ", expiry_value=" +
std::to_string(expiry_option.value) +
+ ", max_topic_size=" + max_topic_size);
+
+ ASSERT_NO_THROW(client->update_topic(
+ make_string_identifier(stream_name),
make_string_identifier(topic_name), updated_topic_name,
+ compression_algorithm, expiry_option.kind,
expiry_option.value, max_topic_size));
topic_name = updated_topic_name;
}
}
@@ -1005,8 +1006,8 @@ TEST_F(LowLevelE2E_Topic,
UpdateTopicWithInvalidNamesThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
const std::vector<std::string> invalid_topic_names = {
"",
@@ -1066,8 +1067,8 @@ TEST_F(LowLevelE2E_Topic, UpdateTopicBeforeLoginThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Client *unauthenticated_client = GetLoggedOutClient();
@@ -1126,8 +1127,8 @@ TEST_F(LowLevelE2E_Topic,
GetTopicsAfterStreamDeletionReturnsEmpty) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->delete_stream(make_string_identifier(stream_name)));
ForgetTrackedStream(stream_name);
@@ -1157,8 +1158,8 @@ TEST_F(LowLevelE2E_Topic,
PurgeTopicAfterStreamDeletionThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
ASSERT_NO_THROW(client->delete_stream(make_string_identifier(stream_name)));
ForgetTrackedStream(stream_name);
@@ -1189,8 +1190,8 @@ TEST_F(LowLevelE2E_Topic,
PurgeTopicWithInvalidStreamIdentifierThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Identifier invalid_kind_id;
invalid_kind_id.kind = "invalid";
@@ -1215,8 +1216,8 @@ TEST_F(LowLevelE2E_Topic,
PurgeTopicWithInvalidTopicIdentifierThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Identifier invalid_kind_id;
invalid_kind_id.kind = "invalid";
@@ -1370,8 +1371,8 @@ TEST_F(LowLevelE2E_Topic,
PurgeTopicAcrossMultiplePartitionsClearsAllPartitions)
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 3, "none", "server_default",
+ 0, "server_default"));
const auto created_stream =
client->get_stream(make_string_identifier(stream_name));
ASSERT_EQ(created_stream.topics.size(), 1u);
@@ -1414,8 +1415,8 @@ TEST_F(LowLevelE2E_Topic,
PurgeTopicThenSendMessagesAgainSucceeds) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
const auto created_stream =
client->get_stream(make_string_identifier(stream_name));
ASSERT_EQ(created_stream.topics.size(), 1u);
@@ -1542,8 +1543,8 @@ TEST_F(LowLevelE2E_Topic, PurgeTopicBeforeLoginThrows) {
ASSERT_NO_THROW(client->create_stream(stream_name));
TrackStream(stream_name);
- ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none",
- "server_default", 0,
"server_default"));
+ ASSERT_NO_THROW(client->create_topic(make_string_identifier(stream_name),
topic_name, 1, "none", "server_default",
+ 0, "server_default"));
iggy::ffi::Client *unauthenticated_client = GetLoggedOutClient();
diff --git
a/foreign/java/bench/src/main/java/org/apache/iggy/bench/common/provision/ResourceProvisioner.java
b/foreign/java/bench/src/main/java/org/apache/iggy/bench/common/provision/ResourceProvisioner.java
index eee42713a..8fdae2b9d 100644
---
a/foreign/java/bench/src/main/java/org/apache/iggy/bench/common/provision/ResourceProvisioner.java
+++
b/foreign/java/bench/src/main/java/org/apache/iggy/bench/common/provision/ResourceProvisioner.java
@@ -32,7 +32,6 @@ import org.slf4j.LoggerFactory;
import java.math.BigInteger;
import java.util.ArrayList;
import java.util.List;
-import java.util.Optional;
public final class ResourceProvisioner {
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/TopicsTcpClient.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/TopicsTcpClient.java
index f2abaead0..27b88d0ee 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/TopicsTcpClient.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/TopicsTcpClient.java
@@ -43,7 +43,6 @@ import java.util.concurrent.CompletableFuture;
import java.util.function.Supplier;
import static org.apache.iggy.serde.BytesSerializer.toBytes;
-import static org.apache.iggy.serde.BytesSerializer.toBytesAsU64;
/**
* Async TCP implementation of TopicsClient using Netty for non-blocking I/O.
@@ -109,8 +108,8 @@ public class TopicsTcpClient implements TopicsClient {
BigInteger maxTopicSize,
String name) {
- var payload = createTopicPayload(
- streamId, partitionsCount, compressionAlgorithm,
messageExpiry, maxTopicSize, name);
+ var payload =
+ createTopicPayload(streamId, partitionsCount,
compressionAlgorithm, messageExpiry, maxTopicSize, name);
return connection().send(CommandCode.Topic.CREATE.getValue(),
payload).thenApply(response -> {
try {
@@ -134,15 +133,13 @@ public class TopicsTcpClient implements TopicsClient {
payload.writeBytes(toBytes(streamId));
payload.writeIntLE(partitionsCount.intValue());
payload.writeBytes(BytesSerializer.toBytes(name));
- payload.writeBytes(BytesSerializer.toBytes(
- createTopicOptions(compressionAlgorithm, messageExpiry,
maxTopicSize)));
+ payload.writeBytes(
+
BytesSerializer.toBytes(createTopicOptions(compressionAlgorithm, messageExpiry,
maxTopicSize)));
return payload;
}
private static Map<HeaderKey, HeaderValue> createTopicOptions(
- CompressionAlgorithm compressionAlgorithm,
- BigInteger messageExpiry,
- BigInteger maxTopicSize) {
+ CompressionAlgorithm compressionAlgorithm, BigInteger
messageExpiry, BigInteger maxTopicSize) {
// Server-default sentinels (compression none, expiry 0, size 0) are
// omitted so the admitting server resolves them from its config.
Map<HeaderKey, HeaderValue> options = new LinkedHashMap<>();
@@ -179,8 +176,8 @@ public class TopicsTcpClient implements TopicsClient {
// Settings ride the options block. A default value means the caller
did
// not set the key, so it is omitted and the server leaves the topic's
// current value alone.
- payload.writeBytes(BytesSerializer.toBytes(
- createTopicOptions(compressionAlgorithm, messageExpiry,
maxTopicSize)));
+ payload.writeBytes(
+
BytesSerializer.toBytes(createTopicOptions(compressionAlgorithm, messageExpiry,
maxTopicSize)));
return connection()
.send(CommandCode.Topic.UPDATE.getValue(), payload)
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/TopicsClient.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/TopicsClient.java
index dbee3d16c..820b23f07 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/TopicsClient.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/TopicsClient.java
@@ -51,12 +51,7 @@ public interface TopicsClient {
BigInteger maxTopicSize,
String name) {
return createTopic(
- StreamId.of(streamId),
- partitionsCount,
- compressionAlgorithm,
- messageExpiry,
- maxTopicSize,
- name);
+ StreamId.of(streamId), partitionsCount, compressionAlgorithm,
messageExpiry, maxTopicSize, name);
}
TopicDetails createTopic(
@@ -75,12 +70,7 @@ public interface TopicsClient {
BigInteger maxTopicSize,
String name) {
updateTopic(
- StreamId.of(streamId),
- TopicId.of(topicId),
- compressionAlgorithm,
- messageExpiry,
- maxTopicSize,
- name);
+ StreamId.of(streamId), TopicId.of(topicId),
compressionAlgorithm, messageExpiry, maxTopicSize, name);
}
void updateTopic(
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/http/TopicsHttpClient.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/http/TopicsHttpClient.java
index c93a302bf..6c8693c67 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/http/TopicsHttpClient.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/http/TopicsHttpClient.java
@@ -63,8 +63,7 @@ class TopicsHttpClient implements TopicsClient {
String name) {
var request = httpClient.preparePostRequest(
STREAMS + "/" + streamId + TOPICS,
- new CreateTopic(
- partitionsCount, compressionAlgorithm, messageExpiry,
maxTopicSize, name));
+ new CreateTopic(partitionsCount, compressionAlgorithm,
messageExpiry, maxTopicSize, name));
return httpClient.execute(request, new TypeReference<>() {});
}
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/tcp/TopicsTcpClient.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/tcp/TopicsTcpClient.java
index 3e2737a3c..1e2c833ba 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/tcp/TopicsTcpClient.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/tcp/TopicsTcpClient.java
@@ -68,8 +68,8 @@ final class TopicsTcpClient implements TopicsClient {
BigInteger messageExpiry,
BigInteger maxTopicSize,
String name) {
- FutureUtil.resolve(delegate.updateTopic(
- streamId, topicId, compressionAlgorithm, messageExpiry,
maxTopicSize, name));
+ FutureUtil.resolve(
+ delegate.updateTopic(streamId, topicId, compressionAlgorithm,
messageExpiry, maxTopicSize, name));
}
@Override
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncClientIntegrationTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncClientIntegrationTest.java
index 6f7ddb503..fee52d7c0 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncClientIntegrationTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncClientIntegrationTest.java
@@ -91,12 +91,7 @@ public class AsyncClientIntegrationTest extends
BaseIntegrationTest {
client.streams().createStream(streamName).get(TIMEOUT_SECONDS,
TimeUnit.SECONDS);
client.topics()
.createTopic(
- streamId,
- 2L,
- CompressionAlgorithm.None,
- BigInteger.ZERO,
- BigInteger.ZERO,
- "test-topic")
+ streamId, 2L, CompressionAlgorithm.None,
BigInteger.ZERO, BigInteger.ZERO, "test-topic")
.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
List<Message> messages = new ArrayList<>();
@@ -145,12 +140,7 @@ public class AsyncClientIntegrationTest extends
BaseIntegrationTest {
client.streams().createStream(streamName).get(TIMEOUT_SECONDS,
TimeUnit.SECONDS);
client.topics()
.createTopic(
- streamId,
- 2L,
- CompressionAlgorithm.None,
- BigInteger.ZERO,
- BigInteger.ZERO,
- "test-topic")
+ streamId, 2L, CompressionAlgorithm.None,
BigInteger.ZERO, BigInteger.ZERO, "test-topic")
.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
// when — send messages concurrently from multiple threads
@@ -207,12 +197,7 @@ public class AsyncClientIntegrationTest extends
BaseIntegrationTest {
client.streams().createStream(streamName).get(TIMEOUT_SECONDS,
TimeUnit.SECONDS);
client.topics()
.createTopic(
- streamId,
- 2L,
- CompressionAlgorithm.None,
- BigInteger.ZERO,
- BigInteger.ZERO,
- "test-topic")
+ streamId, 2L, CompressionAlgorithm.None,
BigInteger.ZERO, BigInteger.ZERO, "test-topic")
.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
// when — send messages in concurrent batches
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncConnectionPoolAuthTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncConnectionPoolAuthTest.java
index aa8a267d8..bce4d1340 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncConnectionPoolAuthTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncConnectionPoolAuthTest.java
@@ -37,7 +37,6 @@ import org.slf4j.LoggerFactory;
import java.math.BigInteger;
import java.util.ArrayList;
import java.util.List;
-import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncConsumerGroupsTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncConsumerGroupsTest.java
index ab7116846..e7dec176b 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncConsumerGroupsTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/AsyncConsumerGroupsTest.java
@@ -82,13 +82,7 @@ public class AsyncConsumerGroupsTest extends
BaseIntegrationTest {
client.streams().createStream(TEST_STREAM).get(TIMEOUT_SECONDS,
TimeUnit.SECONDS);
client.topics()
- .createTopic(
- STREAM_ID,
- 1L,
- CompressionAlgorithm.None,
- BigInteger.ZERO,
- BigInteger.ZERO,
- TEST_TOPIC)
+ .createTopic(STREAM_ID, 1L, CompressionAlgorithm.None,
BigInteger.ZERO, BigInteger.ZERO, TEST_TOPIC)
.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
}
@@ -267,13 +261,7 @@ public class AsyncConsumerGroupsTest extends
BaseIntegrationTest {
try {
client.streams().createStream(streamName).get(TIMEOUT_SECONDS,
TimeUnit.SECONDS);
client.topics()
- .createTopic(
- streamId,
- 1L,
- CompressionAlgorithm.None,
- BigInteger.ZERO,
- BigInteger.ZERO,
- topicName)
+ .createTopic(streamId, 1L, CompressionAlgorithm.None,
BigInteger.ZERO, BigInteger.ZERO, topicName)
.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
// when
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/TopicsTcpClientPayloadTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/TopicsTcpClientPayloadTest.java
index 652e32572..2edc78140 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/TopicsTcpClientPayloadTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/TopicsTcpClientPayloadTest.java
@@ -41,12 +41,7 @@ class TopicsTcpClientPayloadTest {
@Test
void shouldWritePartitionsCountAsFixedFieldBeforeName() {
ByteBuf payload = TopicsTcpClient.createTopicPayload(
- StreamId.of(1L),
- 3L,
- CompressionAlgorithm.None,
- BigInteger.ZERO,
- BigInteger.ZERO,
- "orders");
+ StreamId.of(1L), 3L, CompressionAlgorithm.None,
BigInteger.ZERO, BigInteger.ZERO, "orders");
assertThat(payload.readUnsignedByte()).isEqualTo((short) 1); //
numeric identifier kind
assertThat(payload.readUnsignedByte()).isEqualTo((short) 4); //
identifier length
@@ -59,12 +54,7 @@ class TopicsTcpClientPayloadTest {
@Test
void shouldOmitServerDefaultOptions() {
ByteBuf payload = TopicsTcpClient.createTopicPayload(
- StreamId.of("my-stream"),
- 1L,
- CompressionAlgorithm.None,
- BigInteger.ZERO,
- BigInteger.ZERO,
- "events");
+ StreamId.of("my-stream"), 1L, CompressionAlgorithm.None,
BigInteger.ZERO, BigInteger.ZERO, "events");
payload.skipBytes(2 + "my-stream".length());
assertThat(payload.readUnsignedIntLE()).isEqualTo(1L);
@@ -87,8 +77,7 @@ class TopicsTcpClientPayloadTest {
readString(payload);
Map<String, TypedValue> options = readOptions(payload);
- assertThat(options)
- .containsOnlyKeys("compression_algorithm", "message_expiry",
"max_topic_size");
+ assertThat(options).containsOnlyKeys("compression_algorithm",
"message_expiry", "max_topic_size");
assertThat(options.get("compression_algorithm").kind()).isEqualTo(HeaderKind.String);
assertThat(new String(options.get("compression_algorithm").value(),
StandardCharsets.UTF_8))
.isEqualTo("gzip");
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/IntegrationTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/IntegrationTest.java
index 17fbbd10e..ecd2f0d81 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/IntegrationTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/IntegrationTest.java
@@ -29,7 +29,6 @@ import java.math.BigInteger;
import java.util.ArrayList;
import java.util.List;
-import static java.util.Optional.empty;
import static org.apache.iggy.TestConstants.STREAM_NAME;
import static org.apache.iggy.TestConstants.TOPIC_NAME;
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/TopicsClientBaseTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/TopicsClientBaseTest.java
index 3aa1bc6ce..7b1e25f93 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/TopicsClientBaseTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/TopicsClientBaseTest.java
@@ -26,7 +26,6 @@ import org.junit.jupiter.api.Test;
import java.math.BigInteger;
-import static java.util.Optional.empty;
import static org.apache.iggy.TestConstants.STREAM_NAME;
import static org.assertj.core.api.Assertions.assertThat;
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/tcp/TlsConnectionTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/tcp/TlsConnectionTest.java
index d555b4ac9..f284ad5c8 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/tcp/TlsConnectionTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/blocking/tcp/TlsConnectionTest.java
@@ -43,7 +43,6 @@ import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
-import static java.util.Optional.empty;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -131,13 +130,7 @@ class TlsConnectionTest extends BaseIntegrationTest {
assertThat(stream).isNotNull();
client.topics()
- .createTopic(
- streamId,
- 1L,
- CompressionAlgorithm.None,
- BigInteger.ZERO,
- BigInteger.ZERO,
- topicName);
+ .createTopic(streamId, 1L, CompressionAlgorithm.None,
BigInteger.ZERO, BigInteger.ZERO, topicName);
List<Message> messages =
List.of(Message.of("tls-message-1"),
Message.of("tls-message-2"), Message.of("tls-message-3"));