This is an automated email from the ASF dual-hosted git repository.
wwbmmm 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 cb84e83c fix rpcz root client span lifetime (#3420) (#3421)
cb84e83c is described below
commit cb84e83c59cf41a11b5e426116d663f2800de64b
Author: lh2debug-2 <[email protected]>
AuthorDate: Sun Aug 9 15:03:54 2026 +0800
fix rpcz root client span lifetime (#3420) (#3421)
Keep the current RPC span alive from Controller until the RPC finishes,
SubmitSpan runs, or the Controller is reset. This lets root client spans
without a local parent be submitted to rpcz instead of being destroyed
after the caller-side temporary shared_ptr goes out of scope.
Child client spans remain linked to their parent through weak local-parent
references and parent-owned client lists, so they are still serialized
under their parent without introducing shared_ptr cycles.
Co-authored-by: lh2debug <[email protected]>
---
src/brpc/controller.cpp | 66 +++++++++++++++++++--------------------
src/brpc/controller.h | 2 +-
src/brpc/span.cpp | 20 +++++++-----
src/brpc/span.h | 8 +++++
test/brpc_controller_unittest.cpp | 43 +++++++++++++++++++++++++
5 files changed, 97 insertions(+), 42 deletions(-)
diff --git a/src/brpc/controller.cpp b/src/brpc/controller.cpp
index 5583231c..8a8410be 100644
--- a/src/brpc/controller.cpp
+++ b/src/brpc/controller.cpp
@@ -183,8 +183,8 @@ static void CreateIgnoreAllRead() { s_ignore_all_read = new
IgnoreAllRead; }
// you don't have to set the fields to initial state after deletion since
// they'll be set uniformly after this method is called.
void Controller::ResetNonPods() {
- if (auto span = _span.lock()) {
- Span::Submit(span, butil::cpuwide_time_us());
+ if (_span) {
+ Span::Submit(_span, butil::cpuwide_time_us());
}
_span.reset();
_error_text.clear();
@@ -463,9 +463,9 @@ void Controller::SetFailed(const std::string& reason) {
AppendServerIdentiy();
}
_error_text.append(reason);
- if (auto span = _span.lock()) {
- span->set_error_code(_error_code);
- span->Annotate(reason);
+ if (_span) {
+ _span->set_error_code(_error_code);
+ _span->Annotate(reason);
}
UpdateResponseHeader(this);
}
@@ -492,9 +492,9 @@ void Controller::SetFailed(int error_code, const char*
reason_fmt, ...) {
va_start(ap, reason_fmt);
butil::string_vappendf(&_error_text, reason_fmt, ap);
va_end(ap);
- if (auto span = _span.lock()) {
- span->set_error_code(_error_code);
- span->AnnotateCStr(_error_text.c_str() + old_size, 0);
+ if (_span) {
+ _span->set_error_code(_error_code);
+ _span->AnnotateCStr(_error_text.c_str() + old_size, 0);
}
UpdateResponseHeader(this);
}
@@ -520,9 +520,9 @@ void Controller::CloseConnection(const char* reason_fmt,
...) {
va_start(ap, reason_fmt);
butil::string_vappendf(&_error_text, reason_fmt, ap);
va_end(ap);
- if (auto span = _span.lock()) {
- span->set_error_code(_error_code);
- span->AnnotateCStr(_error_text.c_str() + old_size, 0);
+ if (_span) {
+ _span->set_error_code(_error_code);
+ _span->AnnotateCStr(_error_text.c_str() + old_size, 0);
}
UpdateResponseHeader(this);
}
@@ -982,9 +982,9 @@ void Controller::EndRPC(const CompletionInfo& info) {
}
// RPC finished, now it's safe to release `LoadBalancerWithNaming'
_lb.reset();
- if (auto span = _span.lock()) {
- span->set_ending_cid(info.id);
- span->set_async(_done);
+ if (_span) {
+ _span->set_ending_cid(info.id);
+ _span->set_async(_done);
// Submit the span if we're in async RPC. For sync RPC, the span
// is submitted after Join() to get a more accurate resuming timestamp.
if (_done) {
@@ -1058,14 +1058,14 @@ void Controller::DoneInBackupThread() {
void Controller::SubmitSpan() {
const int64_t now = butil::cpuwide_time_us();
- if (auto span = _span.lock()) {
- span->set_start_callback_us(now);
- if (auto parent_span = span->local_parent().lock()) {
+ if (_span) {
+ _span->set_start_callback_us(now);
+ if (auto parent_span = _span->local_parent().lock()) {
if (parent_span->is_active()) {
parent_span->AsParent();
}
}
- Span::Submit(span, now);
+ Span::Submit(_span, now);
_span.reset();
}
}
@@ -1176,12 +1176,12 @@ void Controller::IssueRPC(int64_t start_realtime_us) {
CHECK_EQ(_remote_side, tmp_sock->remote_side());
}
- if (auto span = _span.lock()) {
+ if (_span) {
if (_current_call.nretry == 0) {
- span->set_remote_side(_remote_side);
+ _span->set_remote_side(_remote_side);
} else {
- span->Annotate("Retrying %s",
- endpoint2str(_remote_side).c_str());
+ _span->Annotate("Retrying %s",
+ endpoint2str(_remote_side).c_str());
}
}
// Handle connection type
@@ -1292,7 +1292,7 @@ void Controller::IssueRPC(int64_t start_realtime_us) {
int rc;
size_t packet_size = 0;
if (user_packet_guard) {
- if (auto span = _span.lock()) {
+ if (_span) {
packet_size = user_packet_guard->EstimatedByteSize();
}
rc = _current_call.sending_sock->Write(user_packet_guard, &wopt);
@@ -1300,13 +1300,13 @@ void Controller::IssueRPC(int64_t start_realtime_us) {
packet_size = packet.size();
rc = _current_call.sending_sock->Write(&packet, &wopt);
}
- if (auto span = _span.lock()) {
+ if (_span) {
if (_current_call.nretry == 0) {
- span->set_sent_us(butil::cpuwide_time_us());
- span->set_request_size(packet_size);
+ _span->set_sent_us(butil::cpuwide_time_us());
+ _span->set_request_size(packet_size);
} else {
- span->Annotate("Requested(%lld) [%d]",
- (long long)packet_size, _current_call.nretry + 1);
+ _span->Annotate("Requested(%lld) [%d]",
+ (long long)packet_size, _current_call.nretry + 1);
}
}
if (using_auth) {
@@ -1447,15 +1447,15 @@ const Controller* Controller::sub(int index) const {
}
uint64_t Controller::trace_id() const {
- if (auto span = _span.lock()) {
- return span->trace_id();
+ if (_span) {
+ return _span->trace_id();
}
return 0;
}
uint64_t Controller::span_id() const {
- if (auto span = _span.lock()) {
- return span->span_id();
+ if (_span) {
+ return _span->span_id();
}
return 0;
}
@@ -1802,7 +1802,7 @@ ControllerPrivateAccessor&
ControllerPrivateAccessor::set_span(Span* span) {
}
std::shared_ptr<Span> ControllerPrivateAccessor::span() const {
- return _cntl->_span.lock();
+ return _cntl->_span;
}
} // namespace brpc
diff --git a/src/brpc/controller.h b/src/brpc/controller.h
index 5d38de0b..564c0875 100644
--- a/src/brpc/controller.h
+++ b/src/brpc/controller.h
@@ -877,7 +877,7 @@ private:
private:
// NOTE: align and group fields to make Controller as compact as possible.
- std::weak_ptr<Span> _span;
+ std::shared_ptr<Span> _span;
uint32_t _flags; // all boolean fields inside Controller
int32_t _error_code;
std::string _error_text;
diff --git a/src/brpc/span.cpp b/src/brpc/span.cpp
index 1863f01a..d5807121 100644
--- a/src/brpc/span.cpp
+++ b/src/brpc/span.cpp
@@ -206,6 +206,7 @@ std::shared_ptr<Span> Span::CreateClientSpan(const
std::string& full_method_name
return nullptr;
}
std::shared_ptr<Span> span(span_raw, SpanDeleter());
+ span->_submitted.store(false, butil::memory_order_relaxed);
span->_log_id = 0;
span->_base_cid = INVALID_BTHREAD_ID;
span->_ending_cid = INVALID_BTHREAD_ID; // Client Span uses ending_cid
@@ -248,6 +249,7 @@ std::shared_ptr<Span> Span::CreateBthreadSpan(const
std::string& full_method_nam
return nullptr;
}
std::shared_ptr<Span> span(span_raw, SpanDeleter());
+ span->_submitted.store(false, butil::memory_order_relaxed);
span->_log_id = 0;
span->_base_cid = INVALID_BTHREAD_ID;
span->_ending_tid = INVALID_BTHREAD; // Bthread Span uses ending_tid
@@ -298,6 +300,7 @@ std::shared_ptr<Span> Span::CreateServerSpan(
return nullptr;
}
std::shared_ptr<Span> span(span_raw, SpanDeleter());
+ span->_submitted.store(false, butil::memory_order_relaxed);
span->_trace_id = (trace_id ? trace_id : GenerateTraceId());
span->_span_id = (span_id ? span_id : GenerateSpanId());
span->_parent_span_id = parent_span_id;
@@ -335,7 +338,9 @@ void Span::ResetServerSpanName(const std::string&
full_method_name) {
}
void Span::submit(int64_t cpuwide_us) {
- // Note: this method is not called for client-side spans.
+ // Called for server spans and root client spans (those without a local
+ // parent). Child client spans are serialized under their parent via
+ // _client_list instead.
EndAsParent();
// If memory allocation fails, the server span will not be submitted for
persistence.
// The server span will be destroyed later when its shared_ptr refcount
drops to zero
@@ -581,12 +586,10 @@ inline int GetSpanDB(butil::intrusive_ptr<SpanDB>* db) {
}
void Span::Submit(std::shared_ptr<Span> span, int64_t cpuwide_time_us) {
- // Only submit spans without a local parent (i.e., server spans).
- // Server spans hold shared_ptr references to their child spans (via
_client_list),
- // ensuring child spans remain alive until the server span is submitted
and dumped.
- // Client spans are not submitted here because their lifetime is managed
by their
- // parent server span.
- if (span->local_parent().expired()) {
+ // Submit root spans without a local parent. Server spans and root client
+ // spans are submitted independently; child client spans with a live local
+ // parent are serialized under the parent to avoid duplicate submissions.
+ if (span->local_parent().expired() && span->try_mark_submitted()) {
span->submit(cpuwide_time_us);
}
}
@@ -787,7 +790,8 @@ leveldb::Status SpanDB::Index(std::shared_ptr<const Span>
span, std::string* val
// be modified by other threads, which could lead to inconsistent data when
// serializing to database.
for (auto it = all_child_spans.rbegin(); it != all_child_spans.rend();
++it) {
- if (*it && it->get() != span.get() && !(*it)->is_active()) {
+ if (*it && it->get() != span.get() && !(*it)->is_active() &&
+ (*it)->try_mark_submitted()) {
RpczSpan* child_proto = value_proto.add_client_spans();
Span2Proto((*it).get(), child_proto);
}
diff --git a/src/brpc/span.h b/src/brpc/span.h
index efa394b5..c6b3306c 100644
--- a/src/brpc/span.h
+++ b/src/brpc/span.h
@@ -28,6 +28,7 @@
#include <ostream>
#include <memory>
#include <pthread.h>
+#include "butil/atomicops.h"
#include "butil/macros.h"
#include "butil/endpoint.h"
#include "butil/string_splitter.h"
@@ -198,6 +199,11 @@ private:
void dump_to_db();
void submit(int64_t cpuwide_us);
+ bool try_mark_submitted() const {
+ bool expected = false;
+ return _submitted.compare_exchange_strong(
+ expected, true, butil::memory_order_relaxed);
+ }
bvar::CollectorSpeedLimit* speed_limit();
bvar::CollectorPreprocessor* preprocessor();
@@ -252,6 +258,8 @@ private:
// Also protects against concurrent iteration (e.g., CountClientSpans,
SpanDB::Index)
// while the list is being modified.
mutable pthread_spinlock_t _client_list_spinlock;
+
+ mutable butil::atomic<bool> _submitted;
};
class SpanContainer : public bvar::Collected {
diff --git a/test/brpc_controller_unittest.cpp
b/test/brpc_controller_unittest.cpp
index 3f410a25..77de2099 100644
--- a/test/brpc_controller_unittest.cpp
+++ b/test/brpc_controller_unittest.cpp
@@ -28,6 +28,7 @@
#include "brpc/server.h"
#include "brpc/channel.h"
#include "brpc/controller.h"
+#include "brpc/span.h"
class ControllerTest : public ::testing::Test{
protected:
@@ -74,6 +75,48 @@ TEST_F(ControllerTest, notify_on_destruction) {
ASSERT_TRUE(cancel);
}
+TEST_F(ControllerTest, root_client_span_kept_alive_until_reset) {
+ brpc::ClearTlsParentSpan();
+
+ brpc::Controller cntl;
+ std::weak_ptr<brpc::Span> weak_span;
+ {
+ std::shared_ptr<brpc::Span> span =
+ brpc::Span::CreateClientSpan("test.RootClient/Call", 0);
+ ASSERT_TRUE(span);
+ ASSERT_TRUE(span->local_parent().expired());
+ weak_span = span;
+ cntl._span = span;
+ }
+
+ ASSERT_FALSE(weak_span.expired());
+ ASSERT_TRUE(cntl._span);
+
+ cntl.Reset();
+ ASSERT_FALSE(cntl._span);
+}
+
+TEST_F(ControllerTest, root_client_span_released_by_submit_span) {
+ brpc::ClearTlsParentSpan();
+
+ brpc::Controller cntl;
+ std::weak_ptr<brpc::Span> weak_span;
+ {
+ std::shared_ptr<brpc::Span> span =
+ brpc::Span::CreateClientSpan("test.RootClient/Call", 0);
+ ASSERT_TRUE(span);
+ ASSERT_TRUE(span->local_parent().expired());
+ weak_span = span;
+ cntl._span = span;
+ }
+
+ ASSERT_FALSE(weak_span.expired());
+ ASSERT_TRUE(cntl._span);
+
+ cntl.SubmitSpan();
+ ASSERT_FALSE(cntl._span);
+}
+
#if ! BRPC_WITH_GLOG
static bool endsWith(const std::string& s1, const butil::StringPiece& s2) {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]