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

yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new a7e9c5ffeda branch-4.1: [fix](rpc) Fix AutoReleaseClosure data race 
with callback reuse (#61782) (#67340)
a7e9c5ffeda is described below

commit a7e9c5ffeda09ece0c72594892aba4b67d9e349a
Author: linrrarity <[email protected]>
AuthorDate: Thu Sep 10 11:22:11 2026 +0800

    branch-4.1: [fix](rpc) Fix AutoReleaseClosure data race with callback reuse 
(#61782) (#67340)
    
    pick: https://github.com/apache/doris/pull/61782
    
    Co-authored-by: zclllyybb <[email protected]>
---
 be/src/exec/exchange/vdata_stream_sender.h         |   1 +
 be/src/exec/operator/exchange_sink_buffer.cpp      |   3 +-
 be/src/exec/runtime_filter/runtime_filter.cpp      |  23 +-
 be/src/exec/runtime_filter/runtime_filter.h        |   7 +
 be/src/exec/runtime_filter/runtime_filter_mgr.cpp  |  73 +++--
 be/src/exec/runtime_filter/runtime_filter_mgr.h    |   6 +
 .../runtime_filter/runtime_filter_producer.cpp     |  63 +---
 .../exec/runtime_filter/runtime_filter_producer.h  |  66 +++-
 be/src/exec/sink/writer/vtablet_writer.cpp         |   1 +
 be/src/util/brpc_closure.h                         | 114 +++----
 .../runtime_filter/sync_size_callback_test.cpp     | 351 +++++++++++++++++++++
 11 files changed, 551 insertions(+), 157 deletions(-)

diff --git a/be/src/exec/exchange/vdata_stream_sender.h 
b/be/src/exec/exchange/vdata_stream_sender.h
index 29b6f400235..9636c752d3d 100644
--- a/be/src/exec/exchange/vdata_stream_sender.h
+++ b/be/src/exec/exchange/vdata_stream_sender.h
@@ -164,6 +164,7 @@ public:
 
     std::shared_ptr<ExchangeSendCallback<PTransmitDataResult>> 
get_send_callback(RpcInstance* ins,
                                                                                
  bool eos) {
+        // here we reuse the callback because it's re-construction may be 
expensive due to many parameters' capture
         if (!_send_callback) {
             _send_callback = 
ExchangeSendCallback<PTransmitDataResult>::create_shared();
         } else {
diff --git a/be/src/exec/operator/exchange_sink_buffer.cpp 
b/be/src/exec/operator/exchange_sink_buffer.cpp
index 34a1c6de666..7f7b8981554 100644
--- a/be/src/exec/operator/exchange_sink_buffer.cpp
+++ b/be/src/exec/operator/exchange_sink_buffer.cpp
@@ -346,6 +346,7 @@ Status ExchangeSinkBuffer::_send_rpc(RpcInstance& 
instance_data) {
             }
             // The eos here only indicates that the current exchange sink has 
reached eos.
             // However, the queue still contains data from other exchange 
sinks, so RPCs need to continue being sent.
+            // `_send_rpc` must be the LAST operation in this function, 
because it may reuse the callback!
             s = _send_rpc(ins);
             if (!s) {
                 _failed(ins.id,
@@ -472,9 +473,9 @@ Status ExchangeSinkBuffer::_send_rpc(RpcInstance& 
instance_data) {
             } else if (eos) {
                 _ended(ins);
             }
-
             // The eos here only indicates that the current exchange sink has 
reached eos.
             // However, the queue still contains data from other exchange 
sinks, so RPCs need to continue being sent.
+            // `_send_rpc` must be the LAST operation in this function, 
because it may reuse the callback!
             s = _send_rpc(ins);
             if (!s) {
                 _failed(ins.id,
diff --git a/be/src/exec/runtime_filter/runtime_filter.cpp 
b/be/src/exec/runtime_filter/runtime_filter.cpp
index ca1ef805251..0fd26fa0c69 100644
--- a/be/src/exec/runtime_filter/runtime_filter.cpp
+++ b/be/src/exec/runtime_filter/runtime_filter.cpp
@@ -36,13 +36,12 @@ Status RuntimeFilter::_push_to_remote(RuntimeState* state, 
const TNetworkAddress
 
     auto merge_filter_request = std::make_shared<PMergeFilterRequest>();
     merge_filter_request->set_stage(_stage);
-    auto merge_filter_callback = 
DummyBrpcCallback<PMergeFilterResponse>::create_shared();
+    _merge_filter_callback = 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared(
+            state->query_options().ignore_runtime_filter_error ? 
std::weak_ptr<QueryContext> {}
+                                                               : 
state->get_query_ctx_weak());
     auto merge_filter_closure =
-            AutoReleaseClosure<PMergeFilterRequest, 
DummyBrpcCallback<PMergeFilterResponse>>::
-                    create_unique(merge_filter_request, merge_filter_callback,
-                                  
state->query_options().ignore_runtime_filter_error
-                                          ? std::weak_ptr<QueryContext> {}
-                                          : state->get_query_ctx_weak());
+            AutoReleaseClosure<PMergeFilterRequest, 
HandleErrorBrpcCallback<PMergeFilterResponse>>::
+                    create_unique(merge_filter_request, 
_merge_filter_callback);
     void* data = nullptr;
     int len = 0;
 
@@ -54,19 +53,21 @@ Status RuntimeFilter::_push_to_remote(RuntimeState* state, 
const TNetworkAddress
     pfragment_instance_id->set_hi(BackendOptions::get_local_backend().id);
     pfragment_instance_id->set_lo((int64_t)this);
 
-    merge_filter_callback->cntl_->set_timeout_ms(
+    _merge_filter_callback->cntl_->set_timeout_ms(
             
get_execution_rpc_timeout_ms(state->get_query_ctx()->execution_timeout()));
     if (config::execution_ignore_eovercrowded) {
-        merge_filter_callback->cntl_->ignore_eovercrowded();
+        _merge_filter_callback->cntl_->ignore_eovercrowded();
     }
 
     RETURN_IF_ERROR(serialize(merge_filter_request.get(), &data, &len));
 
     if (len > 0) {
-        DCHECK(data != nullptr);
-        merge_filter_callback->cntl_->request_attachment().append(data, len);
+        if (data == nullptr) {
+            return Status::InternalError(
+                    "data is nullptr after serialization with len > 0, filter: 
{}", debug_string());
+        }
+        _merge_filter_callback->cntl_->request_attachment().append(data, len);
     }
-
     stub->merge_filter(merge_filter_closure->cntl_.get(), 
merge_filter_closure->request_.get(),
                        merge_filter_closure->response_.get(), 
merge_filter_closure.get());
     // the closure will be released by brpc during closure->Run.
diff --git a/be/src/exec/runtime_filter/runtime_filter.h 
b/be/src/exec/runtime_filter/runtime_filter.h
index ed4357dd772..c1e80dd2b88 100644
--- a/be/src/exec/runtime_filter/runtime_filter.h
+++ b/be/src/exec/runtime_filter/runtime_filter.h
@@ -19,6 +19,8 @@
 
 #include <gen_cpp/PaloInternalService_types.h>
 
+#include <vector>
+
 #include "common/exception.h"
 #include "common/status.h"
 #include "exec/runtime_filter/runtime_filter_definitions.h"
@@ -28,8 +30,11 @@
 
 namespace doris {
 #include "common/compile_check_begin.h"
+class PMergeFilterResponse;
 class RuntimeFilterWrapper;
 class RuntimeProfile;
+template <typename Response>
+class HandleErrorBrpcCallback;
 
 /// The runtimefilter is built in the join node.
 /// The main purpose is to reduce the scanning amount of the
@@ -124,6 +129,8 @@ protected:
     // runtime filter type
     RuntimeFilterType _runtime_filter_type = RuntimeFilterType::UNKNOWN_FILTER;
 
+    std::shared_ptr<HandleErrorBrpcCallback<PMergeFilterResponse>> 
_merge_filter_callback;
+
     friend class RuntimeFilterProducer;
     friend class RuntimeFilterConsumer;
     friend class RuntimeFilterMerger;
diff --git a/be/src/exec/runtime_filter/runtime_filter_mgr.cpp 
b/be/src/exec/runtime_filter/runtime_filter_mgr.cpp
index 29b7aaa3e2f..f69efa855a9 100644
--- a/be/src/exec/runtime_filter/runtime_filter_mgr.cpp
+++ b/be/src/exec/runtime_filter/runtime_filter_mgr.cpp
@@ -43,6 +43,7 @@
 #include "runtime/exec_env.h"
 #include "runtime/memory/mem_tracker.h"
 #include "runtime/query_context.h"
+#include "runtime/runtime_profile.h"
 #include "runtime/runtime_state.h"
 #include "runtime/thread_context.h"
 #include "util/brpc_client_cache.h"
@@ -76,6 +77,28 @@ std::vector<RuntimeFilterPublishTarget> 
build_runtime_filter_publish_targets(
     return publish_targets;
 }
 
+class RuntimeFilterRelayRpcClosure final : public google::protobuf::Closure {
+public:
+    RuntimeFilterRelayRpcClosure(std::shared_ptr<PPublishFilterRequestV2> 
request,
+                                 std::weak_ptr<QueryContext> query_ctx)
+            : _request(std::move(request)),
+              
_callback(HandleErrorBrpcCallback<PPublishFilterResponse>::create_shared(
+                      std::move(query_ctx))) {}
+
+    void Run() override {
+        std::unique_ptr<RuntimeFilterRelayRpcClosure> self(this);
+        _callback->call();
+    }
+
+    brpc::Controller* cntl() { return _callback->cntl_.get(); }
+    PPublishFilterRequestV2* request() { return _request.get(); }
+    PPublishFilterResponse* response() { return _callback->response_.get(); }
+
+private:
+    std::shared_ptr<PPublishFilterRequestV2> _request;
+    std::shared_ptr<HandleErrorBrpcCallback<PPublishFilterResponse>> _callback;
+};
+
 Status send_runtime_filter_relay_rpc(const RuntimeFilterPublishTask& task,
                                      const butil::IOBuf& request_attachment, 
int timeout_ms,
                                      std::weak_ptr<QueryContext> query_ctx) {
@@ -88,21 +111,17 @@ Status send_runtime_filter_relay_rpc(const 
RuntimeFilterPublishTask& task,
                                      task.receiver.addr.hostname(), 
task.receiver.addr.port());
     }
 
-    auto closure =
-            AutoReleaseClosure<PPublishFilterRequestV2, 
DummyBrpcCallback<PPublishFilterResponse>>::
-                    
create_unique(std::make_shared<PPublishFilterRequestV2>(task.request),
-                                  
DummyBrpcCallback<PPublishFilterResponse>::create_shared(),
-                                  query_ctx);
+    // brpc calls Run() exactly once; RuntimeFilterRelayRpcClosure deletes 
itself there.
+    auto* closure = new RuntimeFilterRelayRpcClosure(
+            std::make_shared<PPublishFilterRequestV2>(task.request), 
std::move(query_ctx));
     if (!request_attachment.empty()) {
-        closure->cntl_->request_attachment().append(request_attachment);
+        closure->cntl()->request_attachment().append(request_attachment);
     }
-    closure->cntl_->set_timeout_ms(timeout_ms);
+    closure->cntl()->set_timeout_ms(timeout_ms);
     if (config::execution_ignore_eovercrowded) {
-        closure->cntl_->ignore_eovercrowded();
+        closure->cntl()->ignore_eovercrowded();
     }
-    stub->apply_filterv2(closure->cntl_.get(), closure->request_.get(), 
closure->response_.get(),
-                         closure.get());
-    closure.release();
+    stub->apply_filterv2(closure->cntl(), closure->request(), 
closure->response(), closure);
     return Status::OK();
 }
 
@@ -456,9 +475,9 @@ Status 
RuntimeFilterMergeControllerEntity::send_filter_size(std::shared_ptr<Quer
     Status st = Status::OK();
     // After all runtime filters' size are collected, we should send response 
to all producers.
     if (cnt_val.merger->add_rf_size(request->filter_size())) {
-        auto ctx = query_ctx->ignore_runtime_filter_error() ? 
std::weak_ptr<QueryContext> {}
-                                                            : query_ctx;
-        for (auto addr : cnt_val.source_addrs) {
+        cnt_val.sync_size_callbacks.resize(cnt_val.source_addrs.size());
+        for (size_t i = 0; i < cnt_val.source_addrs.size(); ++i) {
+            auto& addr = cnt_val.source_addrs[i];
             std::shared_ptr<PBackendService_Stub> stub(
                     
ExecEnv::GetInstance()->brpc_internal_client_cache()->get_client(addr));
             if (stub == nullptr) {
@@ -471,10 +490,14 @@ Status 
RuntimeFilterMergeControllerEntity::send_filter_size(std::shared_ptr<Quer
             auto sync_request = std::make_shared<PSyncFilterSizeRequest>();
             sync_request->set_stage(cnt_val.stage);
 
-            auto closure = AutoReleaseClosure<PSyncFilterSizeRequest,
-                                              
DummyBrpcCallback<PSyncFilterSizeResponse>>::
-                    create_unique(sync_request,
-                                  
DummyBrpcCallback<PSyncFilterSizeResponse>::create_shared(), ctx);
+            auto callback = 
HandleErrorBrpcCallback<PSyncFilterSizeResponse>::create_shared(
+                    query_ctx->ignore_runtime_filter_error() ? 
std::weak_ptr<QueryContext> {}
+                                                             : 
query_ctx->weak_from_this());
+            cnt_val.sync_size_callbacks[i] = callback;
+            auto closure = AutoReleaseClosure<
+                    PSyncFilterSizeRequest,
+                    
HandleErrorBrpcCallback<PSyncFilterSizeResponse>>::create_unique(sync_request,
+                                                                               
      callback);
 
             auto* pquery_id = closure->request_->mutable_query_id();
             pquery_id->set_hi(query_ctx->query_id().hi);
@@ -487,7 +510,6 @@ Status 
RuntimeFilterMergeControllerEntity::send_filter_size(std::shared_ptr<Quer
 
             closure->request_->set_filter_id(filter_id);
             
closure->request_->set_filter_size(cnt_val.merger->get_received_sum_size());
-
             stub->sync_filter_size(closure->cntl_.get(), 
closure->request_.get(),
                                    closure->response_.get(), closure.get());
             closure.release();
@@ -669,11 +691,14 @@ Status 
RuntimeFilterMergeControllerEntity::_send_rf_to_target(
     }
 
     auto st = Status::OK();
-    for (auto& target : targets) {
+    cnt_val.publish_callbacks.resize(targets.size());
+    for (size_t i = 0; i < targets.size(); ++i) {
+        auto& target = targets[i];
+        auto callback = 
HandleErrorBrpcCallback<PPublishFilterResponse>::create_shared(ctx);
+        cnt_val.publish_callbacks[i] = callback;
         auto closure = AutoReleaseClosure<PPublishFilterRequestV2,
-                                          
DummyBrpcCallback<PPublishFilterResponse>>::
-                
create_unique(std::make_shared<PPublishFilterRequestV2>(apply_request),
-                              
DummyBrpcCallback<PPublishFilterResponse>::create_shared(), ctx);
+                                          
HandleErrorBrpcCallback<PPublishFilterResponse>>::
+                
create_unique(std::make_shared<PPublishFilterRequestV2>(apply_request), 
callback);
 
         if (has_attachment) {
             closure->cntl_->request_attachment().append(request_attachment);
@@ -717,6 +742,8 @@ Status GlobalMergeContext::reset(QueryContext* query_ctx) {
     merger->increase_expected_producer_num(producer_size);
     arrive_id.clear();
     source_addrs.clear();
+    sync_size_callbacks.clear();
+    publish_callbacks.clear();
     done = false;
     stage++;
     // Keep the Merger's own stage in sync for consistent debug output.
diff --git a/be/src/exec/runtime_filter/runtime_filter_mgr.h 
b/be/src/exec/runtime_filter/runtime_filter_mgr.h
index 94058c2c547..98a1bb3d006 100644
--- a/be/src/exec/runtime_filter/runtime_filter_mgr.h
+++ b/be/src/exec/runtime_filter/runtime_filter_mgr.h
@@ -53,6 +53,9 @@ class RuntimeState;
 class RuntimeFilterWrapper;
 class QueryContext;
 class ExecEnv;
+template <typename Response>
+class HandleErrorBrpcCallback;
+class SyncSizeCallback;
 
 struct RuntimeFilterPublishTarget {
     PNetworkAddress addr;
@@ -94,6 +97,9 @@ struct GlobalMergeContext {
     std::vector<TRuntimeFilterTargetParamsV2> targetv2_info;
     std::unordered_set<UniqueId> arrive_id;
     std::vector<PNetworkAddress> source_addrs;
+    
std::vector<std::shared_ptr<HandleErrorBrpcCallback<PSyncFilterSizeResponse>>>
+            sync_size_callbacks;
+    
std::vector<std::shared_ptr<HandleErrorBrpcCallback<PPublishFilterResponse>>> 
publish_callbacks;
     std::atomic<bool> done = false;
 
     // for represent the round number of recursive cte
diff --git a/be/src/exec/runtime_filter/runtime_filter_producer.cpp 
b/be/src/exec/runtime_filter/runtime_filter_producer.cpp
index ede2fb9e756..b3684a774ab 100644
--- a/be/src/exec/runtime_filter/runtime_filter_producer.cpp
+++ b/be/src/exec/runtime_filter/runtime_filter_producer.cpp
@@ -94,53 +94,6 @@ Status RuntimeFilterProducer::publish(RuntimeState* state, 
bool build_hash_table
     return Status::OK();
 }
 
-class SyncSizeClosure : public AutoReleaseClosure<PSendFilterSizeRequest,
-                                                  
DummyBrpcCallback<PSendFilterSizeResponse>> {
-    std::shared_ptr<Dependency> _dependency;
-    // Should use weak ptr here, because when query context deconstructs, 
should also delete runtime filter
-    // context, it not the memory is not released. And rpc is in another 
thread, it will hold rf context
-    // after query context because the rpc is not returned.
-    std::weak_ptr<RuntimeFilterWrapper> _wrapper;
-    using Base =
-            AutoReleaseClosure<PSendFilterSizeRequest, 
DummyBrpcCallback<PSendFilterSizeResponse>>;
-    friend class RuntimeFilterProducer;
-    ENABLE_FACTORY_CREATOR(SyncSizeClosure);
-
-    void _process_if_rpc_failed() override {
-        Defer defer {[&]() {
-            Base::_process_if_rpc_failed();
-            ((CountedFinishDependency*)_dependency.get())->sub();
-        }};
-        auto wrapper = _wrapper.lock();
-        if (!wrapper) {
-            return;
-        }
-
-        wrapper->set_state(RuntimeFilterWrapper::State::DISABLED, 
cntl_->ErrorText());
-    }
-
-    void _process_if_meet_error_status(const Status& status) override {
-        Defer defer {[&]() {
-            Base::_process_if_meet_error_status(status);
-            ((CountedFinishDependency*)_dependency.get())->sub();
-        }};
-        auto wrapper = _wrapper.lock();
-        if (!wrapper) {
-            return;
-        }
-
-        wrapper->set_state(RuntimeFilterWrapper::State::DISABLED, 
status.to_string());
-    }
-
-public:
-    SyncSizeClosure(std::shared_ptr<PSendFilterSizeRequest> req,
-                    
std::shared_ptr<DummyBrpcCallback<PSendFilterSizeResponse>> callback,
-                    std::shared_ptr<Dependency> dependency,
-                    std::shared_ptr<RuntimeFilterWrapper> wrapper,
-                    std::weak_ptr<QueryContext> context)
-            : Base(req, callback, context), 
_dependency(std::move(dependency)), _wrapper(wrapper) {}
-};
-
 void RuntimeFilterProducer::latch_dependency(
         const std::shared_ptr<CountedFinishDependency>& dependency) {
     std::unique_lock<std::recursive_mutex> l(_rmtx);
@@ -198,14 +151,13 @@ Status RuntimeFilterProducer::send_size(RuntimeState* 
state, uint64_t local_filt
 
     auto request = std::make_shared<PSendFilterSizeRequest>();
     request->set_stage(_stage);
-
-    auto callback = 
DummyBrpcCallback<PSendFilterSizeResponse>::create_shared();
+    // when failed, will check `ignore_runtime_filter_error` in callback to 
decide cancel or not
+    _sync_size_callback = SyncSizeCallback::create_shared(_dependency, 
_wrapper,
+                                                          
state->get_query_ctx()->weak_from_this());
     // RuntimeFilter maybe deconstructed before the rpc finished, so that 
could not use
     // a raw pointer in closure. Has to use the context's shared ptr.
-    auto closure = SyncSizeClosure::create_unique(request, callback, 
_dependency, _wrapper,
-                                                  
state->query_options().ignore_runtime_filter_error
-                                                          ? 
std::weak_ptr<QueryContext> {}
-                                                          : 
state->get_query_ctx_weak());
+    auto closure = AutoReleaseClosure<PSendFilterSizeRequest, 
SyncSizeCallback>::create_unique(
+            request, _sync_size_callback);
     auto* pquery_id = request->mutable_query_id();
     pquery_id->set_hi(state->get_query_ctx()->query_id().hi);
     pquery_id->set_lo(state->get_query_ctx()->query_id().lo);
@@ -217,9 +169,10 @@ Status RuntimeFilterProducer::send_size(RuntimeState* 
state, uint64_t local_filt
     request->set_filter_size(local_filter_size);
     request->set_filter_id(_wrapper->filter_id());
 
-    
callback->cntl_->set_timeout_ms(get_execution_rpc_timeout_ms(state->execution_timeout()));
+    _sync_size_callback->cntl_->set_timeout_ms(
+            get_execution_rpc_timeout_ms(state->execution_timeout()));
     if (config::execution_ignore_eovercrowded) {
-        callback->cntl_->ignore_eovercrowded();
+        _sync_size_callback->cntl_->ignore_eovercrowded();
     }
 
     if (config::enable_debug_points &&
diff --git a/be/src/exec/runtime_filter/runtime_filter_producer.h 
b/be/src/exec/runtime_filter/runtime_filter_producer.h
index 433396144a5..4308030d340 100644
--- a/be/src/exec/runtime_filter/runtime_filter_producer.h
+++ b/be/src/exec/runtime_filter/runtime_filter_producer.h
@@ -17,15 +17,78 @@
 
 #pragma once
 
+#include <glog/logging.h>
+
+#include <memory>
 #include <mutex>
 
 #include "exec/pipeline/dependency.h"
 #include "exec/runtime_filter/runtime_filter.h"
 #include "runtime/query_context.h"
-#include "runtime/runtime_profile.h"
+#include "util/brpc_closure.h"
 
 namespace doris {
 #include "common/compile_check_begin.h"
+
+// Callback for sync-size RPCs. Handles errors (disable wrapper + sub 
dependency) in call().
+class SyncSizeCallback : public DummyBrpcCallback<PSendFilterSizeResponse> {
+    ENABLE_FACTORY_CREATOR(SyncSizeCallback);
+
+public:
+    SyncSizeCallback(std::shared_ptr<Dependency> dependency,
+                     std::shared_ptr<RuntimeFilterWrapper> wrapper,
+                     std::weak_ptr<QueryContext> context)
+            : _dependency(std::move(dependency)), _wrapper(wrapper), 
_context(std::move(context)) {}
+
+    void call() override {
+        // On error: disable the wrapper and sub the dependency here, because 
set_synced_size()
+        // will never be called (the merge node won't respond with a sync).
+        // On success: do NOT sub here. The merge node will respond with 
sync_filter_size,
+        // which calls set_synced_size() -> _dependency->sub().
+        if (this->cntl_->Failed()) {
+            LOG(WARNING) << fmt::format("RPC meet failed: {}", 
this->cntl_->ErrorText());
+            if (auto w = _wrapper.lock()) {
+                w->set_state(RuntimeFilterWrapper::State::DISABLED, 
this->cntl_->ErrorText());
+            }
+            if (auto ctx = _context.lock()) {
+                if (!ctx->ignore_runtime_filter_error()) {
+                    ctx->cancel(
+                            Status::NetworkError("RPC meet failed: {}", 
this->cntl_->ErrorText()));
+                }
+            }
+            auto p = 
std::dynamic_pointer_cast<CountedFinishDependency>(_dependency);
+            DORIS_CHECK(p);
+            p->sub();
+            return;
+        }
+
+        Status status = Status::create(this->response_->status());
+        if (!status.ok()) {
+            LOG(WARNING) << "RPC meet error status: " << status;
+            if (auto w = _wrapper.lock()) {
+                w->set_state(RuntimeFilterWrapper::State::DISABLED, 
status.to_string());
+            }
+            if (auto ctx = _context.lock()) {
+                if (!ctx->ignore_runtime_filter_error()) {
+                    ctx->cancel(status);
+                }
+            }
+            auto p = 
std::dynamic_pointer_cast<CountedFinishDependency>(_dependency);
+            DORIS_CHECK(p);
+            p->sub();
+        }
+    }
+
+private:
+    std::shared_ptr<Dependency> _dependency;
+    // Should use weak ptr here, because when query context deconstructs, 
should also delete runtime filter
+    // context, it not the memory is not released. And rpc is in another 
thread, it will hold rf context
+    // after query context because the rpc is not returned.
+    std::weak_ptr<RuntimeFilterWrapper> _wrapper;
+    // input context is always meaningful. will check 
`ignore_runtime_filter_error` to decide whether to cancel query
+    std::weak_ptr<QueryContext> _context;
+};
+
 // Work on (hash/corss) join build sink node, RuntimeFilterProducerHelper will 
manage all RuntimeFilterProducer
 // Used to generate specific predicate and publish it to consumer/merger
 /**
@@ -181,6 +244,7 @@ private:
 
     int64_t _synced_size = -1;
     std::shared_ptr<CountedFinishDependency> _dependency;
+    std::shared_ptr<SyncSizeCallback> _sync_size_callback;
 
     std::atomic<State> _rf_state;
 };
diff --git a/be/src/exec/sink/writer/vtablet_writer.cpp 
b/be/src/exec/sink/writer/vtablet_writer.cpp
index c5ff0b0d191..a04d3e2c66e 100644
--- a/be/src/exec/sink/writer/vtablet_writer.cpp
+++ b/be/src/exec/sink/writer/vtablet_writer.cpp
@@ -1388,6 +1388,7 @@ void VNodeChannel::cancel(const std::string& cancel_msg) {
     request->set_sender_id(_parent->_sender_id);
     request->set_cancel_reason(cancel_msg);
 
+    // cancel is already in post-processing, so error status could be ignored. 
so not keeping cancel_callback is acceptable.
     auto cancel_callback = 
DummyBrpcCallback<PTabletWriterCancelResult>::create_shared();
     auto closure = AutoReleaseClosure<
             PTabletWriterCancelRequest,
diff --git a/be/src/util/brpc_closure.h b/be/src/util/brpc_closure.h
index e981b185fc0..c517b066a6d 100644
--- a/be/src/util/brpc_closure.h
+++ b/be/src/util/brpc_closure.h
@@ -19,15 +19,18 @@
 
 #include <google/protobuf/stubs/common.h>
 
-#include <atomic>
+#include <type_traits>
 #include <utility>
 
 #include "runtime/query_context.h"
 #include "runtime/thread_context.h"
-#include "service/brpc.h"
+#include "service/brpc.h" // IWYU pragma: keep
 
 namespace doris {
 
+template <typename T>
+concept HasStatus = requires(T* response) { response->status(); };
+
 template <typename Response>
 class DummyBrpcCallback {
     ENABLE_FACTORY_CREATOR(DummyBrpcCallback);
@@ -58,54 +61,69 @@ public:
     std::shared_ptr<Response> response_;
 };
 
+template <typename Response>
+class HandleErrorBrpcCallback : public DummyBrpcCallback<Response> {
+    ENABLE_FACTORY_CREATOR(HandleErrorBrpcCallback);
+
+public:
+    using ResponseType = Response;
+    // input context must be held by caller
+    HandleErrorBrpcCallback(std::weak_ptr<QueryContext> context = {})
+            : _context(std::move(context)) {}
+
+    ~HandleErrorBrpcCallback() override = default;
+
+    void call() override {
+        if (this->cntl_->Failed()) {
+            LOG(WARNING) << fmt::format("RPC meet failed: {}", 
this->cntl_->ErrorText());
+            if (auto ctx = _context.lock()) {
+                ctx->cancel(Status::NetworkError("RPC meet failed: {}", 
this->cntl_->ErrorText()));
+            }
+            return;
+        }
+        if constexpr (HasStatus<Response>) {
+            if (Status status = Status::create(this->response_->status()); 
!status.ok()) {
+                if (!status.is<ErrorCode::END_OF_FILE>()) {
+                    LOG(WARNING) << "RPC meet error status: " << status;
+                    if (auto ctx = _context.lock()) {
+                        ctx->cancel(std::move(status));
+                    }
+                }
+            }
+        }
+    }
+
+private:
+    std::weak_ptr<QueryContext> _context;
+};
+
 // The closure will be deleted after callback.
 // It could only be created by using shared ptr or unique ptr.
-// It will hold a weak ptr of T and call run of T
-// Callback() {
-//  xxxx;
-//  public
-//  void run() {
-//      logxxx
-//  }
-//  }
-//
-//  std::shared_ptr<Callback> b;
-//
+// Example:
 //  std::unique_ptr<AutoReleaseClosure> a(b);
 //  brpc_call(a.release());
-
-template <typename T>
-concept HasStatus = requires(T* response) { response->status(); };
-
+// the closure doesn't own the callback, so the callback MUST be kept alive 
outside.
+// closure only keep a weak ref. if outside owner destroyed (like query 
finish), the callback will be ignored.
 template <typename Request, typename Callback>
 class AutoReleaseClosure : public google::protobuf::Closure {
-    using Weak = typename std::shared_ptr<Callback>::weak_type;
     using ResponseType = typename Callback::ResponseType;
     ENABLE_FACTORY_CREATOR(AutoReleaseClosure);
 
 public:
-    AutoReleaseClosure(std::shared_ptr<Request> req, std::shared_ptr<Callback> 
callback,
-                       std::weak_ptr<QueryContext> context = {}, 
std::string_view error_msg = {})
-            : request_(req), callback_(callback), context_(std::move(context)) 
{
+    AutoReleaseClosure(std::shared_ptr<Request> req, std::shared_ptr<Callback> 
callback)
+            : request_(std::move(req)), callback_(callback) {
         this->cntl_ = callback->cntl_;
         this->response_ = callback->response_;
     }
 
     ~AutoReleaseClosure() override = default;
 
-    //  Will delete itself
+    // Will delete itself. all operations should be done in callback's call(). 
Run() only do one thing.
     void Run() override {
         Defer defer {[&]() { delete this; }};
-        // If lock failed, it means the callback object is deconstructed, then 
no need
-        // to deal with the callback any more.
         if (auto tmp = callback_.lock()) {
             tmp->call();
         }
-        if (cntl_->Failed()) {
-            _process_if_rpc_failed();
-        } else {
-            _process_status<ResponseType>(response_.get());
-        }
     }
 
     // controller has to be the same lifecycle with the closure, because brpc 
may use
@@ -116,45 +134,9 @@ public:
     // at any stage.
     std::shared_ptr<Request> request_;
     std::shared_ptr<ResponseType> response_;
-    std::string error_msg_;
-
-protected:
-    virtual void _process_if_rpc_failed() {
-        std::string error_msg =
-                fmt::format("RPC meet failed: {} {}", cntl_->ErrorText(), 
error_msg_);
-        if (auto ctx = context_.lock(); ctx) {
-            ctx->cancel(Status::NetworkError(error_msg));
-        } else {
-            LOG(WARNING) << error_msg;
-        }
-    }
-
-    virtual void _process_if_meet_error_status(const Status& status) {
-        if (status.is<ErrorCode::END_OF_FILE>()) {
-            // no need to log END_OF_FILE, reduce the unlessful log
-            return;
-        }
-        if (auto ctx = context_.lock(); ctx) {
-            ctx->cancel(status);
-        } else {
-            LOG(WARNING) << "RPC meet error status: " << status;
-        }
-    }
 
 private:
-    template <typename Response>
-    void _process_status(Response* response) {}
-
-    template <HasStatus Response>
-    void _process_status(Response* response) {
-        if (Status status = Status::create(response->status()); !status.ok()) {
-            _process_if_meet_error_status(status);
-        }
-    }
-    // Use a weak ptr to keep the callback, so that the callback can be 
deleted if the main
-    // thread is freed.
-    Weak callback_;
-    std::weak_ptr<QueryContext> context_;
+    std::weak_ptr<Callback> callback_;
 };
 
 } // namespace doris
diff --git a/be/test/exec/runtime_filter/sync_size_callback_test.cpp 
b/be/test/exec/runtime_filter/sync_size_callback_test.cpp
new file mode 100644
index 00000000000..ea1e163d542
--- /dev/null
+++ b/be/test/exec/runtime_filter/sync_size_callback_test.cpp
@@ -0,0 +1,351 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include <gtest/gtest.h>
+
+#include "exec/runtime_filter/runtime_filter_producer.h"
+#include "exec/runtime_filter/runtime_filter_test_utils.h"
+
+namespace doris {
+
+class SyncSizeCallbackTest : public RuntimeFilterTest {
+protected:
+    void SetUp() override {
+        RuntimeFilterTest::SetUp();
+        _dependency = std::make_shared<CountedFinishDependency>(0, 0, 
"TEST_DEP");
+        _dependency->add();
+        _wrapper = std::make_shared<RuntimeFilterWrapper>(
+                PrimitiveType::TYPE_INT, RuntimeFilterType::BLOOM_FILTER, 
/*filter_id=*/0,
+                RuntimeFilterWrapper::State::UNINITED);
+    }
+
+    std::shared_ptr<QueryContext> make_query_ctx(bool ignore_rf_error) {
+        auto opts = TQueryOptionsBuilder().build();
+        opts.__set_ignore_runtime_filter_error(ignore_rf_error);
+        auto fe_address = TNetworkAddress();
+        fe_address.hostname = LOCALHOST;
+        fe_address.port = DUMMY_PORT;
+        return QueryContext::create(TUniqueId(), ExecEnv::GetInstance(), opts, 
fe_address, true,
+                                    fe_address, 
QuerySource::INTERNAL_FRONTEND);
+    }
+
+    std::shared_ptr<SyncSizeCallback> make_callback(
+            std::weak_ptr<QueryContext> ctx = {}, std::shared_ptr<Dependency> 
dep = nullptr,
+            std::shared_ptr<RuntimeFilterWrapper> wrapper = nullptr) {
+        return SyncSizeCallback::create_shared(dep ? dep : _dependency,
+                                               wrapper ? wrapper : _wrapper, 
std::move(ctx));
+    }
+
+    std::shared_ptr<CountedFinishDependency> _dependency;
+    std::shared_ptr<RuntimeFilterWrapper> _wrapper;
+};
+
+// ==================== cntl_->Failed() path ====================
+
+TEST_F(SyncSizeCallbackTest, 
rpc_fail_cancels_query_when_ignore_rf_error_false) {
+    auto ctx = make_query_ctx(false);
+    auto callback = make_callback(ctx);
+    callback->cntl_->SetFailed("injected failure");
+
+    callback->call();
+
+    // Query should be cancelled
+    ASSERT_TRUE(ctx->is_cancelled());
+    // Wrapper should be disabled
+    ASSERT_EQ(_wrapper->get_state(), RuntimeFilterWrapper::State::DISABLED);
+    // Dependency should have been subbed (counter back to 0 -> ready)
+    ASSERT_TRUE(_dependency->ready());
+}
+
+TEST_F(SyncSizeCallbackTest, 
rpc_fail_does_not_cancel_query_when_ignore_rf_error_true) {
+    auto ctx = make_query_ctx(true);
+    auto callback = make_callback(ctx);
+    callback->cntl_->SetFailed("injected failure");
+
+    callback->call();
+
+    // Query should NOT be cancelled
+    ASSERT_FALSE(ctx->is_cancelled());
+    // Wrapper should still be disabled
+    ASSERT_EQ(_wrapper->get_state(), RuntimeFilterWrapper::State::DISABLED);
+    // Dependency should have been subbed
+    ASSERT_TRUE(_dependency->ready());
+}
+
+// ==================== response error status path ====================
+
+TEST_F(SyncSizeCallbackTest, 
response_error_cancels_query_when_ignore_rf_error_false) {
+    auto ctx = make_query_ctx(false);
+    auto callback = make_callback(ctx);
+    // Set a non-OK status in the response
+    auto* status_pb = callback->response_->mutable_status();
+    status_pb->set_status_code(TStatusCode::INTERNAL_ERROR);
+    status_pb->add_error_msgs("injected response error");
+
+    callback->call();
+
+    // Query should be cancelled
+    ASSERT_TRUE(ctx->is_cancelled());
+    // Wrapper should be disabled
+    ASSERT_EQ(_wrapper->get_state(), RuntimeFilterWrapper::State::DISABLED);
+    // Dependency should have been subbed
+    ASSERT_TRUE(_dependency->ready());
+}
+
+TEST_F(SyncSizeCallbackTest, 
response_error_does_not_cancel_query_when_ignore_rf_error_true) {
+    auto ctx = make_query_ctx(true);
+    auto callback = make_callback(ctx);
+    auto* status_pb = callback->response_->mutable_status();
+    status_pb->set_status_code(TStatusCode::INTERNAL_ERROR);
+    status_pb->add_error_msgs("injected response error");
+
+    callback->call();
+
+    // Query should NOT be cancelled
+    ASSERT_FALSE(ctx->is_cancelled());
+    // Wrapper should still be disabled
+    ASSERT_EQ(_wrapper->get_state(), RuntimeFilterWrapper::State::DISABLED);
+    // Dependency should have been subbed
+    ASSERT_TRUE(_dependency->ready());
+}
+
+// ==================== success path ====================
+
+TEST_F(SyncSizeCallbackTest, success_path_no_cancel_no_sub) {
+    auto ctx = make_query_ctx(false);
+    auto callback = make_callback(ctx);
+    // Default: cntl_ not failed, response status OK (status_code defaults to 
0 = OK)
+
+    callback->call();
+
+    // Query should NOT be cancelled
+    ASSERT_FALSE(ctx->is_cancelled());
+    // Wrapper should NOT be disabled
+    ASSERT_NE(_wrapper->get_state(), RuntimeFilterWrapper::State::DISABLED);
+    // Dependency should NOT have been subbed (still blocked)
+    ASSERT_FALSE(_dependency->ready());
+}
+
+// ==================== expired weak_ptr paths ====================
+
+TEST_F(SyncSizeCallbackTest, rpc_fail_with_expired_query_context_no_crash) {
+    // Create a temporary QueryContext that will be destroyed
+    std::weak_ptr<QueryContext> expired_ctx;
+    {
+        auto fe_address = TNetworkAddress();
+        fe_address.hostname = LOCALHOST;
+        fe_address.port = DUMMY_PORT;
+        auto tmp_ctx =
+                QueryContext::create(TUniqueId(), ExecEnv::GetInstance(), 
_query_options,
+                                     fe_address, true, fe_address, 
QuerySource::INTERNAL_FRONTEND);
+        expired_ctx = tmp_ctx;
+    }
+    // expired_ctx is now expired
+
+    auto callback = make_callback(expired_ctx);
+    callback->cntl_->SetFailed("injected failure");
+
+    // Should not crash
+    callback->call();
+
+    // Wrapper should be disabled
+    ASSERT_EQ(_wrapper->get_state(), RuntimeFilterWrapper::State::DISABLED);
+    // Dependency should have been subbed
+    ASSERT_TRUE(_dependency->ready());
+}
+
+TEST_F(SyncSizeCallbackTest, rpc_fail_with_expired_wrapper_no_crash) {
+    auto ctx = make_query_ctx(false);
+    auto tmp_wrapper = std::make_shared<RuntimeFilterWrapper>(
+            PrimitiveType::TYPE_INT, RuntimeFilterType::BLOOM_FILTER, 
/*filter_id=*/1,
+            RuntimeFilterWrapper::State::UNINITED);
+    auto callback = make_callback(ctx, _dependency, tmp_wrapper);
+    // Destroy the wrapper before calling
+    tmp_wrapper.reset();
+
+    callback->cntl_->SetFailed("injected failure");
+
+    // Should not crash even though wrapper is expired
+    callback->call();
+
+    // Query should be cancelled (ignore_runtime_filter_error defaults to 
false)
+    ASSERT_TRUE(ctx->is_cancelled());
+    // Dependency should have been subbed
+    ASSERT_TRUE(_dependency->ready());
+}
+
+TEST_F(SyncSizeCallbackTest, 
auto_release_closure_requires_external_callback_owner) {
+    auto ctx = make_query_ctx(false);
+    auto request = std::make_shared<PSendFilterSizeRequest>();
+    auto* closure = new AutoReleaseClosure<PSendFilterSizeRequest, 
SyncSizeCallback>(
+            request, SyncSizeCallback::create_shared(_dependency, _wrapper, 
ctx));
+    closure->cntl_->SetFailed("injected failure");
+
+    closure->Run();
+
+    ASSERT_FALSE(ctx->is_cancelled());
+    ASSERT_NE(_wrapper->get_state(), RuntimeFilterWrapper::State::DISABLED);
+    ASSERT_FALSE(_dependency->ready());
+}
+
+// =====================================================================
+// HandleErrorBrpcCallback error handling tests
+// =====================================================================
+
+class HandleErrorBrpcCallbackTest : public RuntimeFilterTest {
+protected:
+    std::shared_ptr<QueryContext> make_query_ctx() {
+        auto opts = TQueryOptionsBuilder().build();
+        auto fe_address = TNetworkAddress();
+        fe_address.hostname = LOCALHOST;
+        fe_address.port = DUMMY_PORT;
+        return QueryContext::create(TUniqueId(), ExecEnv::GetInstance(), opts, 
fe_address, true,
+                                    fe_address, 
QuerySource::INTERNAL_FRONTEND);
+    }
+};
+
+class WeakMergeFilterCallback : public DummyBrpcCallback<PMergeFilterResponse> 
{
+public:
+    explicit WeakMergeFilterCallback(std::weak_ptr<QueryContext> context = {})
+            : _context(std::move(context)) {}
+
+    void call() override {
+        if (this->cntl_->Failed()) {
+            if (auto ctx = _context.lock()) {
+                ctx->cancel(Status::NetworkError("RPC meet failed: {}", 
this->cntl_->ErrorText()));
+            }
+        }
+    }
+
+private:
+    std::weak_ptr<QueryContext> _context;
+};
+
+TEST_F(HandleErrorBrpcCallbackTest, 
rpc_fail_cancels_query_when_context_provided) {
+    auto ctx = make_query_ctx();
+    auto callback = 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared(ctx);
+    callback->cntl_->SetFailed("injected failure");
+
+    callback->call();
+
+    ASSERT_TRUE(ctx->is_cancelled());
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, rpc_fail_no_cancel_when_no_context) {
+    // Simulates ignore_runtime_filter_error=true pattern: pass empty weak_ptr
+    auto callback = 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared();
+    callback->cntl_->SetFailed("injected failure");
+
+    callback->call();
+    // No crash, no cancel (no context to cancel)
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, 
response_error_cancels_query_when_context_provided) {
+    auto ctx = make_query_ctx();
+    auto callback = 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared(ctx);
+    auto* status_pb = callback->response_->mutable_status();
+    status_pb->set_status_code(TStatusCode::INTERNAL_ERROR);
+    status_pb->add_error_msgs("injected error");
+
+    callback->call();
+
+    ASSERT_TRUE(ctx->is_cancelled());
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, response_error_no_cancel_when_no_context) {
+    auto callback = 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared();
+    auto* status_pb = callback->response_->mutable_status();
+    status_pb->set_status_code(TStatusCode::INTERNAL_ERROR);
+    status_pb->add_error_msgs("injected error");
+
+    callback->call();
+    // No crash, no cancel
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, success_path_no_cancel) {
+    auto ctx = make_query_ctx();
+    auto callback = 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared(ctx);
+
+    callback->call();
+
+    ASSERT_FALSE(ctx->is_cancelled());
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, expired_query_context_no_crash) {
+    std::weak_ptr<QueryContext> expired_ctx;
+    {
+        auto tmp_ctx = make_query_ctx();
+        expired_ctx = tmp_ctx;
+    }
+
+    auto callback = 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared(expired_ctx);
+    callback->cntl_->SetFailed("injected failure");
+
+    callback->call();
+    // Should not crash; expired context
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, end_of_file_status_does_not_cancel) {
+    auto ctx = make_query_ctx();
+    auto callback = 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared(ctx);
+    auto* status_pb = callback->response_->mutable_status();
+    status_pb->set_status_code(TStatusCode::END_OF_FILE);
+
+    callback->call();
+
+    ASSERT_FALSE(ctx->is_cancelled());
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, 
auto_release_closure_requires_external_callback_owner) {
+    auto ctx = make_query_ctx();
+    auto request = std::make_shared<PMergeFilterRequest>();
+    auto* closure = new AutoReleaseClosure<PMergeFilterRequest,
+                                           
HandleErrorBrpcCallback<PMergeFilterResponse>>(
+            request, 
HandleErrorBrpcCallback<PMergeFilterResponse>::create_shared(ctx));
+    closure->cntl_->SetFailed("injected failure");
+
+    closure->Run();
+
+    ASSERT_FALSE(ctx->is_cancelled());
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, 
weak_auto_release_closure_requires_external_callback_owner) {
+    auto ctx = make_query_ctx();
+    auto request = std::make_shared<PMergeFilterRequest>();
+    auto* closure = new AutoReleaseClosure<PMergeFilterRequest, 
WeakMergeFilterCallback>(
+            request, std::make_shared<WeakMergeFilterCallback>(ctx));
+    closure->cntl_->SetFailed("injected failure");
+
+    closure->Run();
+
+    ASSERT_FALSE(ctx->is_cancelled());
+}
+
+TEST_F(HandleErrorBrpcCallbackTest, 
weak_auto_release_closure_runs_with_external_callback_owner) {
+    auto ctx = make_query_ctx();
+    auto request = std::make_shared<PMergeFilterRequest>();
+    auto callback = std::make_shared<WeakMergeFilterCallback>(ctx);
+    auto* closure =
+            new AutoReleaseClosure<PMergeFilterRequest, 
WeakMergeFilterCallback>(request, callback);
+    closure->cntl_->SetFailed("injected failure");
+
+    closure->Run();
+
+    ASSERT_TRUE(ctx->is_cancelled());
+}
+
+} // namespace doris


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

Reply via email to