Copilot commented on code in PR #3484:
URL: https://github.com/apache/brpc/pull/3484#discussion_r3985801895


##########
src/brpc/input_messenger.h:
##########
@@ -92,6 +95,26 @@ class InputMessageClosure {
     InputMessageBase* _msg;
 };
 
+class InputMessageBatch {
+public:
+    InputMessageBatch() {}
+    explicit InputMessageBatch(size_t capacity);
+    ~InputMessageBatch() noexcept;

Review Comment:
   This type owns the raw message pointers and its destructor 
processes/destroys them, but the user-declared destructor does not disable the 
implicitly generated copy constructor and assignment operator. Copying an 
`InputMessageBatch` would make two owners and cause the same messages to be 
processed or destroyed twice. Make the batch non-copyable (or define explicit 
move-only ownership).



##########
test/brpc_input_messenger_unittest.cpp:
##########
@@ -60,6 +206,214 @@ class MessengerTest : public ::testing::Test{
     };
 };
 
+TEST_F(MessengerTest, input_message_batch_runs_in_order_once) {
+    BatchRecorder recorder;
+    brpc::InputMessageBatch batch(8);
+    batch.add(NewBatchTestMessage(1, &recorder));
+    batch.add(nullptr);
+    batch.add(NewBatchTestMessage(2, &recorder));
+    batch.add(NewBatchTestMessage(3, &recorder));
+    ASSERT_EQ(3u, batch.size());
+
+    batch.Run();
+    EXPECT_TRUE(batch.empty());
+    EXPECT_EQ((std::vector<int>{1, 2, 3}), recorder.Snapshot());
+    EXPECT_EQ(3, recorder.destroyed.load());
+
+    batch.Run();
+    EXPECT_EQ((std::vector<int>{1, 2, 3}), recorder.Snapshot());
+    EXPECT_EQ(3, recorder.destroyed.load());
+}
+
+TEST_F(MessengerTest, input_message_batch_destructor_contains_exceptions) {
+    static_assert(
+        std::is_nothrow_destructible<brpc::InputMessageBatch>::value,
+        "InputMessageBatch must be nothrow destructible");
+
+    BatchRecorder recorder;
+    EXPECT_NO_THROW({
+        brpc::InputMessageBatch batch(2);
+        batch.add(NewBatchTestMessage(1, &recorder, true));
+        batch.add(NewBatchTestMessage(2, &recorder));
+    });
+    EXPECT_EQ((std::vector<int>{1}), recorder.Snapshot());
+    EXPECT_EQ(2, recorder.destroyed.load());
+
+    BatchRecorder worker_recorder;
+    brpc::InputMessageBatch* batch = new brpc::InputMessageBatch(2);
+    batch->add(NewBatchTestMessage(1, &worker_recorder, true));
+    batch->add(NewBatchTestMessage(2, &worker_recorder));
+    EXPECT_NO_THROW(brpc::ProcessInputMessageBatch(batch));
+    EXPECT_EQ((std::vector<int>{1, 2}), worker_recorder.Snapshot());
+    EXPECT_EQ(2, worker_recorder.destroyed.load());
+}
+
+TEST_F(MessengerTest, input_message_batch_flag_validation) {
+    InputBatchFlagGuard flag_guard;
+    EXPECT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "-1").empty());
+    EXPECT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "0").empty());
+    EXPECT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "1").empty());
+    EXPECT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "8").empty());
+    EXPECT_TRUE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "-2").empty());
+}
+
+TEST_F(MessengerTest, adaptive_input_message_batch_rises_falls_and_caps) {
+    uint32_t ema_q8 = 256;
+    uint32_t batch_size = 1;
+    std::vector<uint32_t> rising_levels(1, batch_size);
+    for (int i = 0; i < 32 && batch_size < 16; ++i) {
+        const uint32_t old_batch_size = batch_size;
+        batch_size = brpc::InputMessengerProcessor::UpdateAdaptiveBatchSize(
+            &ema_q8, batch_size, std::numeric_limits<size_t>::max());

Review Comment:
   This new test cannot compile: `UpdateAdaptiveBatchSize` is declared under 
`private:` in `InputMessengerProcessor`, and `MessengerTest` is not a friend. 
Expose a deliberate public/test-only seam or drive the adaptive behavior 
through `ProcessNewMessage` instead.



##########
src/brpc/input_messenger_processor.cpp:
##########
@@ -167,8 +176,78 @@ size_t InputMessengerProcessor::OnceReadSize() const {
 
 void InputMessengerProcessor::Reset() {
     _read_buf.clear();
-    _last_msg_size = 0;
-    _avg_msg_size = 0;
+    ResetMsgSizeStats();
+}
+
+void InputMessengerProcessor::QueueInputMessageBatch(
+        std::unique_ptr<InputMessageBatch>* batch,
+        int* num_bthread_created) {
+    if (!batch->get() || (*batch)->empty()) {
+        return;
+    }
+    _socket->_transport->QueueMessages(batch->release(), num_bthread_created);
+}
+
+void InputMessengerProcessor::QueueLastMessageOrBatch(
+        InputMessageClosure& last_msg,
+        std::unique_ptr<InputMessageBatch>* batch,
+        int* num_bthread_created, size_t batch_size) {
+    InputMessageBase* msg = last_msg.release();
+    if (msg == nullptr) {
+        return;
+    }
+    if (!batch->get()) {
+        batch->reset(new (std::nothrow) InputMessageBatch(batch_size));
+    }
+    if (!batch->get()) {
+        last_msg.reset(msg);
+        _socket->_transport->QueueMessage(
+            last_msg, num_bthread_created, false);
+        return;
+    }
+    (*batch)->add(msg);

Review Comment:
   `add` can throw `std::bad_alloc` while growing `_msgs` for a fixed batch 
size larger than the capped initial reservation. At this point `msg` has 
already been released from `last_msg`, so the exception leaves it unowned and 
`ProcessNewMessage` has no synchronous fallback; retain ownership and handle 
the failure while preserving ordering with any messages already in the batch.



##########
src/brpc/input_messenger_processor.cpp:
##########
@@ -167,8 +176,78 @@ size_t InputMessengerProcessor::OnceReadSize() const {
 
 void InputMessengerProcessor::Reset() {
     _read_buf.clear();
-    _last_msg_size = 0;
-    _avg_msg_size = 0;
+    ResetMsgSizeStats();
+}
+
+void InputMessengerProcessor::QueueInputMessageBatch(
+        std::unique_ptr<InputMessageBatch>* batch,
+        int* num_bthread_created) {
+    if (!batch->get() || (*batch)->empty()) {
+        return;
+    }
+    _socket->_transport->QueueMessages(batch->release(), num_bthread_created);
+}
+
+void InputMessengerProcessor::QueueLastMessageOrBatch(
+        InputMessageClosure& last_msg,
+        std::unique_ptr<InputMessageBatch>* batch,
+        int* num_bthread_created, size_t batch_size) {
+    InputMessageBase* msg = last_msg.release();
+    if (msg == nullptr) {
+        return;
+    }
+    if (!batch->get()) {
+        batch->reset(new (std::nothrow) InputMessageBatch(batch_size));
+    }

Review Comment:
   `new (std::nothrow)` only converts allocation failure from `operator new` 
into `nullptr`; exceptions from the `InputMessageBatch` constructor still 
propagate. Since the constructor calls `vector::reserve`, a reserve failure can 
unwind here after `msg` was released from `last_msg`, bypassing the documented 
synchronous fallback and losing ownership of the message. Catch the constructor 
failure before taking the fallback path.



##########
test/brpc_input_messenger_unittest.cpp:
##########
@@ -60,6 +206,214 @@ class MessengerTest : public ::testing::Test{
     };
 };
 
+TEST_F(MessengerTest, input_message_batch_runs_in_order_once) {
+    BatchRecorder recorder;
+    brpc::InputMessageBatch batch(8);
+    batch.add(NewBatchTestMessage(1, &recorder));
+    batch.add(nullptr);
+    batch.add(NewBatchTestMessage(2, &recorder));
+    batch.add(NewBatchTestMessage(3, &recorder));
+    ASSERT_EQ(3u, batch.size());
+
+    batch.Run();
+    EXPECT_TRUE(batch.empty());
+    EXPECT_EQ((std::vector<int>{1, 2, 3}), recorder.Snapshot());
+    EXPECT_EQ(3, recorder.destroyed.load());
+
+    batch.Run();
+    EXPECT_EQ((std::vector<int>{1, 2, 3}), recorder.Snapshot());
+    EXPECT_EQ(3, recorder.destroyed.load());
+}
+
+TEST_F(MessengerTest, input_message_batch_destructor_contains_exceptions) {
+    static_assert(
+        std::is_nothrow_destructible<brpc::InputMessageBatch>::value,
+        "InputMessageBatch must be nothrow destructible");
+
+    BatchRecorder recorder;
+    EXPECT_NO_THROW({
+        brpc::InputMessageBatch batch(2);
+        batch.add(NewBatchTestMessage(1, &recorder, true));
+        batch.add(NewBatchTestMessage(2, &recorder));
+    });
+    EXPECT_EQ((std::vector<int>{1}), recorder.Snapshot());
+    EXPECT_EQ(2, recorder.destroyed.load());
+
+    BatchRecorder worker_recorder;
+    brpc::InputMessageBatch* batch = new brpc::InputMessageBatch(2);
+    batch->add(NewBatchTestMessage(1, &worker_recorder, true));
+    batch->add(NewBatchTestMessage(2, &worker_recorder));
+    EXPECT_NO_THROW(brpc::ProcessInputMessageBatch(batch));
+    EXPECT_EQ((std::vector<int>{1, 2}), worker_recorder.Snapshot());
+    EXPECT_EQ(2, worker_recorder.destroyed.load());
+}
+
+TEST_F(MessengerTest, input_message_batch_flag_validation) {
+    InputBatchFlagGuard flag_guard;
+    EXPECT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "-1").empty());
+    EXPECT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "0").empty());
+    EXPECT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "1").empty());
+    EXPECT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "8").empty());
+    EXPECT_TRUE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "input_message_batch_process_size", "-2").empty());
+}
+
+TEST_F(MessengerTest, adaptive_input_message_batch_rises_falls_and_caps) {
+    uint32_t ema_q8 = 256;
+    uint32_t batch_size = 1;
+    std::vector<uint32_t> rising_levels(1, batch_size);
+    for (int i = 0; i < 32 && batch_size < 16; ++i) {
+        const uint32_t old_batch_size = batch_size;
+        batch_size = brpc::InputMessengerProcessor::UpdateAdaptiveBatchSize(
+            &ema_q8, batch_size, std::numeric_limits<size_t>::max());
+        if (batch_size != old_batch_size) {
+            rising_levels.push_back(batch_size);
+        }
+    }
+    EXPECT_EQ((std::vector<uint32_t>{1, 2, 4, 8, 16}), rising_levels);
+    EXPECT_LE(ema_q8, 32u * 256);
+
+    std::vector<uint32_t> falling_levels(1, batch_size);
+    for (int i = 0; i < 32 && batch_size > 1; ++i) {
+        const uint32_t old_batch_size = batch_size;
+        batch_size = brpc::InputMessengerProcessor::UpdateAdaptiveBatchSize(
+            &ema_q8, batch_size, 1);
+        if (batch_size != old_batch_size) {
+            falling_levels.push_back(batch_size);
+        }
+    }
+    EXPECT_EQ((std::vector<uint32_t>{16, 8, 4, 2, 1}), falling_levels);
+
+    brpc::SocketId id;
+    ASSERT_EQ(0, brpc::Socket::Create(brpc::SocketOptions(), &id));
+    brpc::SocketUniquePtr socket;
+    ASSERT_EQ(0, brpc::Socket::Address(id, &socket));
+    brpc::InputMessengerProcessor& processor = socket->fd_input_processor();
+    processor._input_messages_per_read_ema_q8 = ema_q8;
+    processor._adaptive_input_message_batch_size = batch_size;
+    ASSERT_EQ(0, socket->ResetFileDescriptor(-1));
+    EXPECT_EQ(0u, processor._input_messages_per_read_ema_q8);
+    EXPECT_EQ(0u, processor._adaptive_input_message_batch_size);
+}
+
+TEST_F(MessengerTest, batching_consumes_the_last_message) {
+    InputBatchFlagGuard flag_guard;
+    brpc::FLAGS_input_message_batch_process_size = 8;
+    brpc::FLAGS_usercode_in_coroutine = false;
+
+    BatchRecorder recorder;
+    brpc::InputMessenger messenger(4);
+    ASSERT_EQ(0, messenger.AddNonProtocolHandler(
+        MakeBatchTestHandler(&recorder)));
+    brpc::SocketId id;
+    brpc::SocketUniquePtr socket = CreateBatchTestSocket(&messenger, &id);
+    ASSERT_TRUE(socket);
+    brpc::InputMessengerProcessor& processor = socket->fd_input_processor();
+    processor.read_buf().append("abcd", 4);
+
+    brpc::InputMessageClosure last_msg;
+    ASSERT_EQ(0, processor.ProcessNewMessage(
+        4, false, 123, 456, last_msg));
+    brpc::DestroyingPtr<brpc::InputMessageBase> leftover(last_msg.release());
+    EXPECT_EQ(nullptr, leftover.get());
+    ASSERT_TRUE(recorder.WaitForSize(4));
+    EXPECT_EQ((std::vector<int>{'a', 'b', 'c', 'd'}), recorder.Snapshot());
+
+    socket->SetFailed();
+    ASSERT_TRUE(recorder.WaitForDestroyed(4));
+}
+
+TEST_F(MessengerTest, disabled_and_progressive_paths_remain_individual) {
+    InputBatchFlagGuard flag_guard;
+    brpc::FLAGS_usercode_in_coroutine = false;
+
+    {
+        brpc::FLAGS_input_message_batch_process_size = 0;
+        BatchRecorder recorder;
+        brpc::InputMessenger messenger(4);
+        ASSERT_EQ(0, messenger.AddNonProtocolHandler(
+            MakeBatchTestHandler(&recorder)));
+        brpc::SocketId id;
+        brpc::SocketUniquePtr socket = CreateBatchTestSocket(&messenger, &id);
+        ASSERT_TRUE(socket);
+        brpc::InputMessengerProcessor& processor =
+            socket->fd_input_processor();
+        processor.read_buf().append("ab", 2);
+        brpc::InputMessageClosure last_msg;
+        ASSERT_EQ(0, processor.ProcessNewMessage(
+            2, false, 123, 456, last_msg));
+        brpc::DestroyingPtr<brpc::InputMessageBase> 
leftover(last_msg.release());
+        EXPECT_NE(nullptr, leftover.get());
+        leftover.reset();
+        ASSERT_TRUE(recorder.WaitForDestroyed(2));
+        socket->SetFailed();
+    }
+
+    {
+        brpc::FLAGS_input_message_batch_process_size = 8;
+        BatchRecorder recorder;
+        brpc::InputMessenger messenger(4);
+        ASSERT_EQ(0, messenger.AddNonProtocolHandler(
+            MakeBatchTestHandler(&recorder)));
+        brpc::SocketId id;
+        brpc::SocketUniquePtr socket = CreateBatchTestSocket(&messenger, &id);
+        ASSERT_TRUE(socket);
+        socket->read_will_be_progressive(brpc::CONNECTION_TYPE_SINGLE);
+        brpc::InputMessengerProcessor& processor =
+            socket->fd_input_processor();
+        processor.read_buf().append("abc", 3);
+        brpc::InputMessageClosure last_msg;
+        ASSERT_EQ(0, processor.ProcessNewMessage(
+            3, false, 123, 456, last_msg));
+        EXPECT_EQ(nullptr, last_msg.release());
+        ASSERT_TRUE(recorder.WaitForSize(3));
+        std::vector<int> values = recorder.Snapshot();
+        std::sort(values.begin(), values.end());
+        EXPECT_EQ((std::vector<int>{'a', 'b', 'c'}), values);
+        socket->SetFailed();
+        ASSERT_TRUE(recorder.WaitForDestroyed(3));
+    }
+}
+
+TEST_F(MessengerTest, transport_batch_helper_handles_inline_and_empty_batches) 
{
+    InputBatchFlagGuard flag_guard;
+    brpc::FLAGS_usercode_in_coroutine = true;
+
+    brpc::InputMessenger messenger;
+    brpc::SocketId id;
+    brpc::SocketUniquePtr socket = CreateBatchTestSocket(&messenger, &id);
+    ASSERT_TRUE(socket);
+
+    BatchRecorder recorder;
+    int num_bthread_created = 0;
+    brpc::InputMessageBatch* batch = new brpc::InputMessageBatch(2);
+    batch->add(NewBatchTestMessage(1, &recorder));
+    batch->add(NewBatchTestMessage(2, &recorder));
+    socket->_transport->QueueMessages(batch, &num_bthread_created);
+    EXPECT_EQ((std::vector<int>{1, 2}), recorder.Snapshot());
+    EXPECT_EQ(0, num_bthread_created);
+
+    socket->_transport->QueueMessages(

Review Comment:
   `Socket::_transport` is private (socket.h:975), and this test has no 
friendship or private-to-public build workaround, so these calls cannot compile 
(the same access is repeated at line 408). Exercise the helper through a public 
input path or add a deliberate test seam.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to