diff --git a/kafka/src/kafka/producer.cpp b/kafka/src/kafka/producer.cpp index 996766aa3a30..75de2d7a2ada 100644 --- a/kafka/src/kafka/producer.cpp +++ b/kafka/src/kafka/producer.cpp @@ -1,5 +1,6 @@ #include +#include #include #include #include @@ -106,6 +107,17 @@ void HandleDeliveryErrors(const std::vector& delivery_resu } } +void RunAsyncTask(auto name, auto& task_processor, auto code_block) { + engine::current_task::CancellationPoint(); + + { + const engine::TaskCancellationBlocker cancellation_blocker; + utils::Async(task_processor, name, code_block).Get(); + } + + engine::current_task::CancellationPoint(); +} + } // namespace Producer::Producer( @@ -136,9 +148,9 @@ void Producer::Send( std::optional partition, HeaderViews headers ) const { - utils::Async(producer_task_processor_, "producer_send", [this, topic_name, key, message, partition, &headers] { + RunAsyncTask("producer_send", producer_task_processor_, [this, topic_name, key, message, partition, &headers] { SendImpl(topic_name, key, message, partition, impl::HeadersHolder{headers}); - }).Get(); + }); } void Producer::SendWrapper( @@ -148,13 +160,13 @@ void Producer::SendWrapper( std::optional partition, HeaderViews headers ) const { - utils::Async( - producer_task_processor_, + RunAsyncTask( "producer_send_bulk", + producer_task_processor_, [this, topic_name, key, &messages, partition, &headers] { SendImpl(topic_name, key, messages, partition, BuildHeaderHolders(headers, messages.Size())); } - ).Get(); + ); } engine::TaskWithResult Producer::SendAsync( diff --git a/kafka/tests/producer_kafkatest.cpp b/kafka/tests/producer_kafkatest.cpp index 17244c195097..e9f90e8f7e70 100644 --- a/kafka/tests/producer_kafkatest.cpp +++ b/kafka/tests/producer_kafkatest.cpp @@ -135,7 +135,7 @@ UTEST_F(ProducerTest, LargeMessages) { const kafka::impl::ProducerConfiguration producer_configuration{}; - auto producer = MakeProducer("kafka-producer"); + auto producer = MakeProducer("kafka-producer", producer_configuration); const std::string topic = GenerateTopic(); const std::string big_message(producer_configuration.message_max_bytes / 2, 'm'); @@ -526,4 +526,124 @@ UTEST_F_MT(ProducerTest, OneProducerManySendAsyncMt, 1 + 4) { UEXPECT_NO_THROW(engine::WaitAllChecked(*results_lock)); } +UTEST_F(ProducerTest, NoUseAfterFreeOnTaskCancellationBeforeSendCall) { + kafka::impl::ProducerConfiguration producer_configuration{}; + producer_configuration.delivery_timeout = std::chrono::seconds{2}; + producer_configuration.queue_buffering_max = std::chrono::seconds{1}; + producer_configuration.queue_buffering_max_messages = 100; + + auto producer = MakeProducer("kafka-producer", producer_configuration); + const std::string topic = GenerateTopic(); + + bool sendThrowsTaskCancelledException = false; + + auto send_task = utils::Async("send_task", [&] { + const std::string key = "test-key"; + const std::string message = "test-msg"; + + engine::current_task::RequestCancel(); + + try { + producer.Send(topic, key, message); + } catch (...) { + sendThrowsTaskCancelledException = true; + throw; + } + }); + + UEXPECT_THROW(send_task.Get(), engine::TaskCancelledException); + EXPECT_TRUE(sendThrowsTaskCancelledException); +} + +UTEST_F(ProducerTest, NoUseAfterFreeOnTaskCancellationAfterSendCall) { + kafka::impl::ProducerConfiguration producer_configuration{}; + producer_configuration.delivery_timeout = std::chrono::seconds{2}; + producer_configuration.queue_buffering_max = std::chrono::seconds{1}; + producer_configuration.queue_buffering_max_messages = 100; + + auto producer = MakeProducer("kafka-producer", producer_configuration); + const std::string topic = GenerateTopic(); + + bool sendThrowsTaskCancelledException = false; + + auto send_task = utils::Async("send_task", [&] { + const std::string key = "test-key"; + const std::string message = "test-msg"; + + try { + producer.Send(topic, key, message); + } catch (...) { + sendThrowsTaskCancelledException = true; + throw; + } + }); + + engine::SleepFor(std::chrono::milliseconds(10)); + + send_task.RequestCancel(); + + UEXPECT_THROW(send_task.Get(), engine::TaskCancelledException); + EXPECT_TRUE(sendThrowsTaskCancelledException); +} + +UTEST_F(ProducerTest, NoUseAfterFreeOnTaskCancellationBeforeBulkSendCall) { + kafka::impl::ProducerConfiguration producer_configuration{}; + producer_configuration.delivery_timeout = std::chrono::seconds{2}; + producer_configuration.queue_buffering_max = std::chrono::seconds{1}; + producer_configuration.queue_buffering_max_messages = 100; + + auto producer = MakeProducer("kafka-producer", producer_configuration); + const std::string topic = GenerateTopic(); + + bool sendThrowsTaskCancelledException = false; + + auto send_task = utils::Async("send_task", [&] { + const std::string key = "test-key"; + const std::vector messages{"test-msg-1", "test-msg-2", "test-msg-3"}; + + engine::current_task::RequestCancel(); + + try { + producer.Send(topic, key, messages); + } catch (...) { + sendThrowsTaskCancelledException = true; + throw; + } + }); + + UEXPECT_THROW(send_task.Get(), engine::TaskCancelledException); + EXPECT_TRUE(sendThrowsTaskCancelledException); +} + +UTEST_F(ProducerTest, NoUseAfterFreeOnTaskCancellationAfterBulkSendCall) { + kafka::impl::ProducerConfiguration producer_configuration{}; + producer_configuration.delivery_timeout = std::chrono::seconds{2}; + producer_configuration.queue_buffering_max = std::chrono::seconds{1}; + producer_configuration.queue_buffering_max_messages = 100; + + auto producer = MakeProducer("kafka-producer", producer_configuration); + const std::string topic = GenerateTopic(); + + bool sendThrowsTaskCancelledException = false; + + auto send_task = utils::Async("send_task", [&] { + const std::string key = "test-key"; + const std::vector messages{"test-msg-1", "test-msg-2", "test-msg-3"}; + + try { + producer.Send(topic, key, messages); + } catch (...) { + sendThrowsTaskCancelledException = true; + throw; + } + }); + + engine::SleepFor(std::chrono::milliseconds(10)); + + send_task.RequestCancel(); + + UEXPECT_THROW(send_task.Get(), engine::TaskCancelledException); + EXPECT_TRUE(sendThrowsTaskCancelledException); +} + USERVER_NAMESPACE_END diff --git a/scripts/kafka/ubuntu_install_kafka.sh b/scripts/kafka/ubuntu_install_kafka.sh index cd45b02ef177..d4dec1eb3e5a 100755 --- a/scripts/kafka/ubuntu_install_kafka.sh +++ b/scripts/kafka/ubuntu_install_kafka.sh @@ -3,12 +3,11 @@ # Exit on any error and treat unset variables as errors, print all commands set -o errexit -o nounset -o pipefail -o posix -x -KAFKA_VERSION=4.0.1 +KAFKA_VERSION=4.3.1 KAFKA_HOME=/etc/kafka +KAFKA_URL="https://www.apache.org/dyn/closer.lua/kafka/${KAFKA_VERSION}/kafka_2.13-${KAFKA_VERSION}.tgz?action=download" DEBIAN_FRONTEND=noninteractive sudo apt install -y openjdk-17-jdk -curl "https://dlcdn.apache.org/kafka/$KAFKA_VERSION/kafka_2.13-$KAFKA_VERSION.tgz" -o kafka.tgz mkdir -p "$KAFKA_HOME" -tar -xzf kafka.tgz --directory="$KAFKA_HOME" --strip-components=1 -rm kafka.tgz +curl -fsSL "$KAFKA_URL" | tar -xzf - --directory="$KAFKA_HOME" --strip-components=1