This is an automated email from the ASF dual-hosted git repository.

chenBright pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/brpc.git


The following commit(s) were added to refs/heads/master by this push:
     new 4047c4e0 Limit reassembled stream message size (#3489)
4047c4e0 is described below

commit 4047c4e0a33e49b3f144b138c902eab54057252c
Author: Xiaofeng Wang <[email protected]>
AuthorDate: Wed Aug 26 10:34:01 2026 +0800

    Limit reassembled stream message size (#3489)
    
    Reject a fragmented stream message when its accumulated payload exceeds
    max_body_size. Release buffered fragments and close the stream with
    EMSGSIZE.
    
    Cover the boundary and rejection paths through a complete client/server RPC
    stream.
---
 src/brpc/stream.cpp                  | 14 ++++++
 test/brpc_streaming_rpc_unittest.cpp | 86 ++++++++++++++++++++++++++++++++++++
 2 files changed, 100 insertions(+)

diff --git a/src/brpc/stream.cpp b/src/brpc/stream.cpp
index 8db93523..7e032b69 100644
--- a/src/brpc/stream.cpp
+++ b/src/brpc/stream.cpp
@@ -35,6 +35,7 @@
 namespace brpc {
 
 DECLARE_bool(usercode_in_pthread);
+DECLARE_uint64(max_body_size);
 DECLARE_int64(socket_max_streams_unconsumed_bytes);
 DEFINE_uint64(stream_write_max_segment_size, 512 * 1024 * 1024,
               "Stream message exceeding this size will be automatically split 
into smaller segments");
@@ -610,6 +611,19 @@ int Stream::OnReceived(const StreamFrameMeta& fm, 
butil::IOBuf *buf, Socket* soc
         CHECK(buf->empty());
         break;
     case FRAME_TYPE_DATA:
+        if (buf->length() > FLAGS_max_body_size ||
+            (_pending_buf != nullptr &&
+             _pending_buf->length() > FLAGS_max_body_size - buf->length())) {
+            LOG(WARNING) << "Close stream=" << id()
+                         << " whose pending message size="
+                         << (_pending_buf != nullptr ? _pending_buf->length() 
: 0)
+                         << " plus frame size=" << buf->length()
+                         << " exceeds max_body_size=" << FLAGS_max_body_size;
+            delete _pending_buf;
+            _pending_buf = nullptr;
+            Close(EMSGSIZE, "Reassembled stream message is too large");
+            return -1;
+        }
         if (_pending_buf != nullptr) {
             _pending_buf->append(*buf);
             buf->clear();
diff --git a/test/brpc_streaming_rpc_unittest.cpp 
b/test/brpc_streaming_rpc_unittest.cpp
index 7ae820aa..12c2ff3c 100644
--- a/test/brpc_streaming_rpc_unittest.cpp
+++ b/test/brpc_streaming_rpc_unittest.cpp
@@ -368,6 +368,92 @@ static bool WaitForTrue(const std::atomic<bool>& f, int 
timeout_ms) {
     return WaitForTrue([&f]() { return f.load(std::memory_order_acquire); }, 
timeout_ms);
 }
 
+class ReassemblyLimitHandler : public brpc::StreamInputHandler {
+public:
+    int on_received_messages(brpc::StreamId,
+                             butil::IOBuf* const messages[],
+                             size_t size) override {
+        for (size_t i = 0; i < size; ++i) {
+            received_bytes.fetch_add(messages[i]->length(),
+                                     std::memory_order_relaxed);
+        }
+        received_messages.fetch_add(size, std::memory_order_release);
+        return 0;
+    }
+
+    void on_idle_timeout(brpc::StreamId) override {}
+    void on_closed(brpc::StreamId) override {}
+    void on_failed(brpc::StreamId, int error_code,
+                   const std::string&) override {
+        failure_code.store(error_code, std::memory_order_release);
+    }
+
+    std::atomic<size_t> received_bytes{0};
+    std::atomic<size_t> received_messages{0};
+    std::atomic<int> failure_code{0};
+};
+
+TEST_F(StreamingRpcTest, limit_reassembled_message_size) {
+    std::string old_max_body_size;
+    std::string old_segment_size;
+    ASSERT_TRUE(GFLAGS_NAMESPACE::GetCommandLineOption(
+        "max_body_size", &old_max_body_size));
+    ASSERT_TRUE(GFLAGS_NAMESPACE::GetCommandLineOption(
+        "stream_write_max_segment_size", &old_segment_size));
+
+    ReassemblyLimitHandler handler;
+    brpc::StreamOptions server_stream_options;
+    server_stream_options.handler = &handler;
+    brpc::Server server;
+    MyServiceWithStream service(server_stream_options);
+    ASSERT_EQ(0, server.AddService(
+        &service, brpc::SERVER_DOESNT_OWN_SERVICE));
+    ASSERT_EQ(0, server.Start(0, nullptr));
+
+    brpc::Channel channel;
+    ASSERT_EQ(0, channel.Init(server.listen_address(), nullptr));
+    brpc::Controller cntl;
+    brpc::StreamId request_stream;
+    ASSERT_EQ(0, brpc::StreamCreate(&request_stream, cntl, nullptr));
+    brpc::ScopedStream stream_guard(request_stream);
+    test::EchoService_Stub stub(&channel);
+    stub.Echo(&cntl, &request, &response, nullptr);
+    ASSERT_FALSE(cntl.Failed()) << cntl.ErrorText();
+
+    ASSERT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "max_body_size", "64").empty());
+    ASSERT_FALSE(GFLAGS_NAMESPACE::SetCommandLineOption(
+        "stream_write_max_segment_size", "16").empty());
+    BRPC_SCOPE_EXIT {
+        GFLAGS_NAMESPACE::SetCommandLineOption(
+            "max_body_size", old_max_body_size.c_str());
+        GFLAGS_NAMESPACE::SetCommandLineOption(
+            "stream_write_max_segment_size", old_segment_size.c_str());
+    };
+
+    butil::IOBuf at_limit;
+    at_limit.append(std::string(64, 'a'));
+    ASSERT_EQ(0, brpc::StreamWrite(request_stream, at_limit));
+    ASSERT_TRUE(WaitForTrue([&handler]() {
+        return handler.received_messages.load(std::memory_order_acquire) == 1;
+    }, 2000));
+    ASSERT_EQ(64u, handler.received_bytes.load(std::memory_order_relaxed));
+
+    butil::IOBuf over_limit;
+    over_limit.append(std::string(65, 'b'));
+    ASSERT_EQ(0, brpc::StreamWrite(request_stream, over_limit));
+    ASSERT_TRUE(WaitForTrue([&handler]() {
+        return handler.failure_code.load(std::memory_order_acquire) != 0;
+    }, 2000));
+    ASSERT_EQ(EMSGSIZE,
+              handler.failure_code.load(std::memory_order_relaxed));
+    ASSERT_EQ(1u,
+              handler.received_messages.load(std::memory_order_relaxed));
+
+    server.Stop(0);
+    server.Join();
+}
+
 TEST_F(StreamingRpcTest, sanity) {
     brpc::Server server;
     MyServiceWithStream service;


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

Reply via email to