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 137c1ca3 Revert "Progressive timeout dev (#3409)" (#3453)
137c1ca3 is described below

commit 137c1ca30487811149a5d36da4c096c69a6d5c73
Author: Weibing Wang <[email protected]>
AuthorDate: Sun Aug 16 17:31:23 2026 +0800

    Revert "Progressive timeout dev (#3409)" (#3453)
    
    This reverts commit e0abb1001ef76d400dddc2b2c3a1ff42f5d3afac.
---
 example/http_c++/http_client.cpp         |  35 -----
 example/http_c++/http_server.cpp         |   7 -
 src/brpc/controller.cpp                  | 259 +------------------------------
 src/brpc/controller.h                    |   8 +-
 src/brpc/errno.proto                     |   1 -
 src/brpc/policy/http_rpc_protocol.cpp    |   1 -
 src/brpc/policy/http_rpc_protocol.h      |  12 +-
 src/brpc/progressive_reader.h            |   2 -
 test/brpc_http_rpc_protocol_unittest.cpp | 237 +---------------------------
 9 files changed, 5 insertions(+), 557 deletions(-)

diff --git a/example/http_c++/http_client.cpp b/example/http_c++/http_client.cpp
index 3a09186f..23222dee 100644
--- a/example/http_c++/http_client.cpp
+++ b/example/http_c++/http_client.cpp
@@ -22,15 +22,11 @@
 // - Access www.foo.com
 //   ./http_client www.foo.com
 
-#include <string>
 #include <gflags/gflags.h>
 #include <butil/logging.h>
 #include <brpc/channel.h>
-#include "bthread/countdown_event.h"
 
 DEFINE_string(d, "", "POST this data to the http server");
-DEFINE_bool(progressive, false, "whether or not progressive read data from 
server");
-DEFINE_int32(progressive_read_timeout_ms, 5000, "progressive read data idle 
timeout in milliseconds");
 DEFINE_string(load_balancer, "", "The algorithm for load balancing");
 DEFINE_int32(timeout_ms, 2000, "RPC timeout in milliseconds");
 DEFINE_int32(max_retry, 3, "Max retries(not including the first RPC)"); 
@@ -40,25 +36,6 @@ namespace brpc {
 DECLARE_bool(http_verbose);
 }
 
-class PartDataReader: public brpc::ProgressiveReader {
-public:
-    explicit PartDataReader(bthread::CountdownEvent* done): _done(done){}
-
-    butil::Status OnReadOnePart(const void* data, size_t length) {
-        const std::string part(static_cast<const char*>(data), length);
-        LOG(INFO) << "data: " << part << " size: " << length;
-        return butil::Status::OK();
-    }
-
-    void OnEndOfMessage(const butil::Status& status) {
-        LOG(INFO) << "progressive read data final status : " << status;
-        _done->signal();
-        delete this;
-    }
-private:
-    bthread::CountdownEvent* _done;
-};
-
 int main(int argc, char* argv[]) {
     // Parse gflags. We recommend you to use gflags as well.
     GFLAGS_NAMESPACE::ParseCommandLineFlags(&argc, &argv, true);
@@ -94,11 +71,6 @@ int main(int argc, char* argv[]) {
         cntl.request_attachment().append(FLAGS_d);
     }
 
-    if (FLAGS_progressive) {
-        
cntl.set_progressive_read_timeout_ms(FLAGS_progressive_read_timeout_ms);
-        cntl.response_will_be_read_progressively();
-    }
-
     // Because `done'(last parameter) is NULL, this function waits until
     // the response comes back or error occurs(including timedout).
     channel.CallMethod(NULL, &cntl, NULL, NULL, NULL);
@@ -106,13 +78,6 @@ int main(int argc, char* argv[]) {
         std::cerr << cntl.ErrorText() << std::endl;
         return -1;
     }
-
-    if (FLAGS_progressive) {
-        bthread::CountdownEvent done(1);
-        cntl.ReadProgressiveAttachmentBy(new PartDataReader(&done));
-        done.wait();
-        LOG(INFO) << "wait client progressive read done safely";
-    }
     // If -http_verbose is on, brpc already prints the response to stderr.
     if (!brpc::FLAGS_http_verbose) {
         std::cout << cntl.response_attachment() << std::endl;
diff --git a/example/http_c++/http_server.cpp b/example/http_c++/http_server.cpp
index 4c3c8722..05c9a0ee 100644
--- a/example/http_c++/http_server.cpp
+++ b/example/http_c++/http_server.cpp
@@ -31,7 +31,6 @@ DEFINE_int32(idle_timeout_s, -1, "Connection will be closed 
if there is no "
 DEFINE_string(certificate, "cert.pem", "Certificate file path to enable SSL");
 DEFINE_string(private_key, "key.pem", "Private key file path to enable SSL");
 DEFINE_string(ciphers, "", "Cipher suite used for SSL connections");
-DEFINE_bool(enable_progressive_timeout, false, "whether or not trigger 
progressive write attachment data timeout");
 
 namespace example {
 
@@ -105,9 +104,6 @@ public:
 
             // sleep a while to send another part.
             bthread_usleep(10000);
-            if (FLAGS_enable_progressive_timeout && i > 50) {
-                bthread_usleep(100000000UL);
-            }
         }
         return NULL;
     }
@@ -198,9 +194,6 @@ public:
 
             // sleep a while to send another part.
             bthread_usleep(10000 * 10);
-            if (FLAGS_enable_progressive_timeout && i > 50) {
-                bthread_usleep(100000000UL);
-            }
         }
         return NULL;
     }
diff --git a/src/brpc/controller.cpp b/src/brpc/controller.cpp
index 66a9a3a5..8a8410be 100644
--- a/src/brpc/controller.cpp
+++ b/src/brpc/controller.cpp
@@ -73,7 +73,6 @@ BAIDU_REGISTER_ERRNO(brpc::EEOF, "Got EOF");
 BAIDU_REGISTER_ERRNO(brpc::EUNUSED, "The socket was not needed");
 BAIDU_REGISTER_ERRNO(brpc::ESSL, "SSL related operation failed");
 BAIDU_REGISTER_ERRNO(brpc::EH2RUNOUTSTREAMS, "The H2 socket was run out of 
streams");
-BAIDU_REGISTER_ERRNO(brpc::EPROGREADTIMEOUT, "Progressive read timed out");
 
 BAIDU_REGISTER_ERRNO(brpc::EINTERNAL, "General internal error");
 BAIDU_REGISTER_ERRNO(brpc::ERESPONSE, "Bad response");
@@ -95,9 +94,8 @@ namespace brpc {
 DEFINE_bool(graceful_quit_on_sigterm, false,
             "Register SIGTERM handle func to quit graceful");
 DEFINE_bool(graceful_quit_on_sighup, false,
-            "Register SIGHUP handle func to quit graceful");
-DEFINE_bool(log_idle_progressive_read_close, false,
-            "Print log when an idle progressive read is closed");
+            "Register SIGHUP handle func to quit graceful");            
+
 const IdlNames idl_single_req_single_res = { "req", "res" };
 const IdlNames idl_single_req_multi_res = { "req", "" };
 const IdlNames idl_multi_req_single_res = { "", "res" };
@@ -176,226 +174,6 @@ public:
     void OnEndOfMessage(const butil::Status&) {}
 };
 
-struct ProgressiveReadTimeoutTask;
-
-struct ProgressiveReadTimeoutState {
-    ProgressiveReadTimeoutState(SocketId id, int32_t timeout_ms)
-        : socket_id(id)
-        , read_timeout_ms(timeout_ms)
-        , deadline_us(butil::cpuwide_time_us() + timeout_ms * 1000L)
-        , timer_id(0)
-        , timer_task(NULL)
-        , user_callback_running(false)
-        , reader_failed(false)
-        , timeout_triggered(false)
-        , end_delivered(false) {}
-
-    butil::Mutex mutex;
-    const SocketId socket_id;
-    const int32_t read_timeout_ms;
-    int64_t deadline_us;
-    bthread_timer_t timer_id;
-    ProgressiveReadTimeoutTask* timer_task;
-    bool user_callback_running;
-    bool reader_failed;
-    bool timeout_triggered;
-    bool end_delivered;
-    butil::Status timer_error;
-};
-
-struct ProgressiveReadTimeoutTask {
-    explicit ProgressiveReadTimeoutTask(
-        const std::shared_ptr<ProgressiveReadTimeoutState>& state_in)
-        : state(state_in) {}
-
-    std::shared_ptr<ProgressiveReadTimeoutState> state;
-};
-
-class ProgressiveTimeoutReader : public ProgressiveReader {
-public:
-    ProgressiveTimeoutReader(SocketId id, int32_t read_timeout_ms,
-                             ProgressiveReader* reader)
-        : _reader(reader)
-        , _state(new ProgressiveReadTimeoutState(id, read_timeout_ms)) {}
-
-    int Start() {
-        std::unique_lock<butil::Mutex> mu(_state->mutex);
-        return AddWatchdogLocked(_state, _state->read_timeout_ms * 1000L);
-    }
-
-    butil::Status OnReadOnePart(const void* data, size_t length) override {
-        {
-            std::unique_lock<butil::Mutex> mu(_state->mutex);
-            if (_state->timeout_triggered) {
-                return MakeTimeoutStatus(_state->read_timeout_ms);
-            }
-            if (!_state->timer_error.ok()) {
-                return _state->timer_error;
-            }
-            _state->user_callback_running = true;
-        }
-
-        butil::Status status = _reader->OnReadOnePart(data, length);
-        {
-            std::unique_lock<butil::Mutex> mu(_state->mutex);
-            _state->user_callback_running = false;
-            if (_state->timeout_triggered) {
-                status = MakeTimeoutStatus(_state->read_timeout_ms);
-            } else if (!_state->timer_error.ok()) {
-                status = _state->timer_error;
-            } else if (status.ok() && !_state->end_delivered) {
-                _state->deadline_us = butil::cpuwide_time_us() +
-                    _state->read_timeout_ms * 1000L;
-            } else if (!status.ok()) {
-                _state->reader_failed = true;
-            }
-        }
-        return status;
-    }
-
-    void OnEndOfMessage(const butil::Status& status) override {
-        bthread_timer_t timer_id = 0;
-        ProgressiveReadTimeoutTask* timer_task = NULL;
-        butil::Status final_status = status;
-        ProgressiveReader* reader = NULL;
-        {
-            std::unique_lock<butil::Mutex> mu(_state->mutex);
-            if (_state->end_delivered) {
-                LOG(ERROR) << "ProgressiveReader::OnEndOfMessage was called 
more than once";
-                return;
-            }
-            _state->end_delivered = true;
-            timer_id = _state->timer_id;
-            timer_task = _state->timer_task;
-            _state->timer_id = 0;
-            _state->timer_task = NULL;
-            if (_state->timeout_triggered) {
-                final_status = MakeTimeoutStatus(_state->read_timeout_ms);
-            } else if (!_state->timer_error.ok()) {
-                final_status = _state->timer_error;
-            }
-            reader = _reader;
-            _reader = NULL;
-        }
-
-        CancelWatchdog(timer_id, timer_task);
-        reader->OnEndOfMessage(final_status);
-        delete this;
-    }
-
-private:
-    ~ProgressiveTimeoutReader() override {}
-
-    static butil::Status MakeTimeoutStatus(int32_t timeout_ms) {
-        return butil::Status(
-            EPROGREADTIMEOUT,
-            "Progressive read timed out after %d ms", timeout_ms);
-    }
-
-    static butil::Status MakeTimerErrorStatus(int error_code) {
-        return butil::Status(
-            error_code, "Fail to add progressive read timeout timer: %s",
-            berror(error_code));
-    }
-
-    static void CancelWatchdog(
-        bthread_timer_t timer_id, ProgressiveReadTimeoutTask* timer_task) {
-        if (timer_id == 0) {
-            return;
-        }
-        const int rc = bthread_timer_del(timer_id);
-        if (rc == 0) {
-            delete timer_task;
-        } else if (rc == 1 || rc == EINVAL) {
-            // The callback owns timer_task once it starts running. EINVAL 
means
-            // that the callback has already finished and released the task.
-        } else {
-            LOG(ERROR) << "Unexpected bthread_timer_del error=" << rc;
-        }
-    }
-
-    static int AddWatchdogLocked(
-        const std::shared_ptr<ProgressiveReadTimeoutState>& state,
-        int64_t delay_us) {
-        if (state->end_delivered || state->reader_failed) {
-            return ECANCELED;
-        }
-        if (delay_us <= 0) {
-            delay_us = 1;
-        }
-        ProgressiveReadTimeoutTask* task =
-            new (std::nothrow) ProgressiveReadTimeoutTask(state);
-        if (task == NULL) {
-            return ENOMEM;
-        }
-        bthread_timer_t timer_id = 0;
-        const int rc = bthread_timer_add(
-            &timer_id, butil::microseconds_from_now(delay_us),
-            HandleIdleProgressiveReader, task);
-        if (rc != 0) {
-            delete task;
-            return rc;
-        }
-        state->timer_id = timer_id;
-        state->timer_task = task;
-        return 0;
-    }
-
-    static void HandleIdleProgressiveReader(void* arg) {
-        std::unique_ptr<ProgressiveReadTimeoutTask> task(
-            static_cast<ProgressiveReadTimeoutTask*>(arg));
-        const std::shared_ptr<ProgressiveReadTimeoutState> state = task->state;
-        bool fail_socket = false;
-        int error_code = 0;
-        std::string error_text;
-        {
-            std::unique_lock<butil::Mutex> mu(state->mutex);
-            if (state->timer_task == task.get()) {
-                state->timer_id = 0;
-                state->timer_task = NULL;
-            }
-            if (state->end_delivered || state->reader_failed) {
-                return;
-            }
-
-            const int64_t now_us = butil::cpuwide_time_us();
-            if (state->user_callback_running || now_us < state->deadline_us) {
-                const int64_t delay_us = state->user_callback_running
-                    ? state->read_timeout_ms * 1000L
-                    : state->deadline_us - now_us;
-                const int rc = AddWatchdogLocked(state, delay_us);
-                if (rc != 0) {
-                    state->timer_error = MakeTimerErrorStatus(rc);
-                    fail_socket = true;
-                    error_code = rc;
-                    error_text = state->timer_error.error_str();
-                }
-            } else {
-                state->timeout_triggered = true;
-                fail_socket = true;
-                error_code = EPROGREADTIMEOUT;
-                error_text = 
MakeTimeoutStatus(state->read_timeout_ms).error_str();
-            }
-        }
-
-        if (!fail_socket) {
-            return;
-        }
-        SocketUniquePtr socket;
-        if (Socket::Address(state->socket_id, &socket) != 0) {
-            LOG(ERROR) << "Fail to address socket_id=" << state->socket_id
-                       << " after progressive read timeout";
-        } else {
-            LOG_IF(INFO, FLAGS_log_idle_progressive_read_close)
-                << error_text << ", socket_id=" << state->socket_id;
-            socket->SetFailed(error_code, "%s", error_text.c_str());
-        }
-    }
-
-    ProgressiveReader* _reader;
-    const std::shared_ptr<ProgressiveReadTimeoutState> _state;
-};
-
 static IgnoreAllRead* s_ignore_all_read = NULL;
 static pthread_once_t s_ignore_all_read_once = PTHREAD_ONCE_INIT;
 static void CreateIgnoreAllRead() { s_ignore_all_read = new IgnoreAllRead; }
@@ -483,7 +261,6 @@ void Controller::ResetPods() {
     _backup_request_ms = UNSET_MAGIC_NUM;
     _backup_request_policy = NULL;
     _connect_timeout_ms = UNSET_MAGIC_NUM;
-    _progressive_read_timeout_ms = UNSET_MAGIC_NUM;
     _real_timeout_ms = UNSET_MAGIC_NUM;
     _deadline_us = -1;
     _timeout_id = 0;
@@ -559,11 +336,6 @@ void Controller::Call::Reset() {
     stream_user_data = NULL;
 }
 
-void Controller::set_progressive_read_timeout_ms(
-    int32_t progressive_read_timeout_ms) {
-    _progressive_read_timeout_ms = progressive_read_timeout_ms;
-}
-
 void Controller::set_timeout_ms(int64_t timeout_ms) {
     if (timeout_ms <= 0x7fffffff) {
         _timeout_ms = timeout_ms;
@@ -1841,33 +1613,6 @@ void 
Controller::ReadProgressiveAttachmentBy(ProgressiveReader* r) {
                          __FUNCTION__));
     }
     add_flag(FLAGS_PROGRESSIVE_READER);
-    if (progressive_read_timeout_ms() > 0) {
-        const SocketId socket_id = _rpa->GetSocketId();
-        if (socket_id == INVALID_SOCKET_ID) {
-            pthread_once(&s_ignore_all_read_once, CreateIgnoreAllRead);
-            _rpa->ReadProgressiveAttachmentBy(s_ignore_all_read);
-            return r->OnEndOfMessage(butil::Status(
-                ENOTSUP,
-                "Progressive read timeout is only supported for HTTP/1.x"));
-        }
-        ProgressiveTimeoutReader* reader = new (std::nothrow)
-            ProgressiveTimeoutReader(
-                socket_id, _progressive_read_timeout_ms, r);
-        if (reader == NULL) {
-            pthread_once(&s_ignore_all_read_once, CreateIgnoreAllRead);
-            _rpa->ReadProgressiveAttachmentBy(s_ignore_all_read);
-            return r->OnEndOfMessage(
-                butil::Status(ENOMEM, "Fail to create progressive timeout 
reader"));
-        }
-        const int rc = reader->Start();
-        if (rc != 0) {
-            pthread_once(&s_ignore_all_read_once, CreateIgnoreAllRead);
-            _rpa->ReadProgressiveAttachmentBy(s_ignore_all_read);
-            return reader->OnEndOfMessage(butil::Status(
-                rc, "Fail to add progressive read timeout timer: %s", 
berror(rc)));
-        }
-        return _rpa->ReadProgressiveAttachmentBy(reader);
-    }
     return _rpa->ReadProgressiveAttachmentBy(r);
 }
 
diff --git a/src/brpc/controller.h b/src/brpc/controller.h
index 41d42ae7..564c0875 100644
--- a/src/brpc/controller.h
+++ b/src/brpc/controller.h
@@ -199,9 +199,6 @@ public:
 
     // Set/get timeout in milliseconds for the RPC call. Use
     // ChannelOptions.timeout_ms on unset.
-    void set_progressive_read_timeout_ms(int32_t progressive_read_timeout_ms);
-    int32_t progressive_read_timeout_ms() const { return 
_progressive_read_timeout_ms; }
-
     void set_timeout_ms(int64_t timeout_ms);
     int64_t timeout_ms() const { return _timeout_ms; }
 
@@ -364,9 +361,7 @@ public:
 
     // Make the RPC end when the HTTP response has complete headers and let
     // user read the remaining body by using ReadProgressiveAttachmentBy().
-    void response_will_be_read_progressively() {
-        add_flag(FLAGS_READ_PROGRESSIVELY);
-    }
+    void response_will_be_read_progressively() { 
add_flag(FLAGS_READ_PROGRESSIVELY); }
     // Make the RPC end when the HTTP request has complete headers and let
     // user read the remaining body by using ReadProgressiveAttachmentBy().
     void request_will_be_read_progressively() { 
add_flag(FLAGS_READ_PROGRESSIVELY); }
@@ -916,7 +911,6 @@ private:
     int32_t _timeout_ms;
     int32_t _connect_timeout_ms;
     int32_t _backup_request_ms;
-    int32_t _progressive_read_timeout_ms;
     // Priority: `_backup_request_policy' > `_backup_request_ms'.
     BackupRequestPolicy* _backup_request_policy;
     // If this rpc call has retry/backup request,this var save the real 
timeout for current call
diff --git a/src/brpc/errno.proto b/src/brpc/errno.proto
index 166d82dc..26ffadc2 100644
--- a/src/brpc/errno.proto
+++ b/src/brpc/errno.proto
@@ -41,7 +41,6 @@ enum Errno {
     ESSL                    = 1016;  // SSL related error
     EH2RUNOUTSTREAMS        = 1017;  // The H2 socket was run out of streams
     EREJECT                 = 1018;  // The Request is rejected
-    EPROGREADTIMEOUT        = 1019;  // The Progressive read timeout
 
     // Errno caused by server
     EINTERNAL               = 2001;  // Internal Server Error
diff --git a/src/brpc/policy/http_rpc_protocol.cpp 
b/src/brpc/policy/http_rpc_protocol.cpp
index 3fb94088..8cbe0698 100644
--- a/src/brpc/policy/http_rpc_protocol.cpp
+++ b/src/brpc/policy/http_rpc_protocol.cpp
@@ -1201,7 +1201,6 @@ ParseResult ParseHttpMessage(butil::IOBuf *source, Socket 
*socket,
             LOG(FATAL) << "Fail to new HttpContext";
             return MakeParseError(PARSE_ERROR_NO_RESOURCE);
         }
-        http_imsg->SetSocketId(socket->id());
         // Parsing http is costly, parsing an incomplete http message from the
         // beginning repeatedly should be avoided, otherwise the cost may reach
         // O(n^2) in the worst case. Save incomplete http messages in sockets
diff --git a/src/brpc/policy/http_rpc_protocol.h 
b/src/brpc/policy/http_rpc_protocol.h
index cd41798e..bc8bd065 100644
--- a/src/brpc/policy/http_rpc_protocol.h
+++ b/src/brpc/policy/http_rpc_protocol.h
@@ -87,20 +87,11 @@ class HttpContext : public ReadableProgressiveAttachment
                   , public InputMessageBase
                   , public HttpMessage {
 public:
-    SocketId GetSocketId() override {
-        return _socket_id;
-    }
-
-    void SetSocketId(SocketId id) {
-        _socket_id = id;
-    }
-
     explicit HttpContext(bool read_body_progressively,
                          HttpMethod request_method = HTTP_METHOD_GET)
         : InputMessageBase()
         , HttpMessage(read_body_progressively, request_method)
-        , _is_stage2(false)
-        , _socket_id(INVALID_SOCKET_ID) {
+        , _is_stage2(false) {
         // add one ref for Destroy
         butil::intrusive_ptr<HttpContext>(this).detach();
     }
@@ -131,7 +122,6 @@ public:
 
 private:
     bool _is_stage2;
-    SocketId _socket_id;
 };
 
 // Implement functions required in protocol.h
diff --git a/src/brpc/progressive_reader.h b/src/brpc/progressive_reader.h
index c84be8b7..6f54ae68 100644
--- a/src/brpc/progressive_reader.h
+++ b/src/brpc/progressive_reader.h
@@ -20,7 +20,6 @@
 #define BRPC_PROGRESSIVE_READER_H
 
 #include "brpc/shared_object.h"
-#include "brpc/socket_id.h"
 
 
 namespace brpc {
@@ -85,7 +84,6 @@ public:
     // Any error occurred should destroy the reader by calling r->Destroy().
     // r->Destroy() should be guaranteed to be called once and only once.
     virtual void ReadProgressiveAttachmentBy(ProgressiveReader* r) = 0;
-    virtual SocketId GetSocketId() = 0;
 };
 
 } // namespace brpc
diff --git a/test/brpc_http_rpc_protocol_unittest.cpp 
b/test/brpc_http_rpc_protocol_unittest.cpp
index 70691314..97de6995 100644
--- a/test/brpc_http_rpc_protocol_unittest.cpp
+++ b/test/brpc_http_rpc_protocol_unittest.cpp
@@ -19,7 +19,6 @@
 
 // Date: Sun Jul 13 15:04:18 CST 2014
 
-#include <atomic>
 #include <cstddef>
 #include <string>
 #include <sys/ioctl.h>
@@ -737,13 +736,9 @@ static void CopyPAPrefixedWithSeqNo(char* buf, uint64_t 
seq_no) {
 class DownloadServiceImpl : public ::test::DownloadService {
 public:
     DownloadServiceImpl(DonePlace done_place = DONE_BEFORE_CREATE_PA,
-                        size_t num_repeat = 1,
-                        int write_interval_us = 0,
-                        int initial_write_delay_us = 0)
+                        size_t num_repeat = 1)
         : _done_place(done_place)
         , _nrep(num_repeat)
-        , _write_interval_us(write_interval_us)
-        , _initial_write_delay_us(initial_write_delay_us)
         , _nwritten(0)
         , _ever_full(false)
         , _last_errno(0) {}
@@ -767,9 +762,6 @@ public:
         if (_done_place == DONE_BEFORE_CREATE_PA) {
             done_guard.reset(NULL);
         }
-        if (_initial_write_delay_us > 0) {
-            bthread_usleep(_initial_write_delay_us);
-        }
         ASSERT_GT(PA_DATA_LEN, 8u);  // long enough to hold a 64-bit decimal.
         char buf[PA_DATA_LEN];
         for (size_t c = 0; c < _nrep;) {
@@ -786,9 +778,6 @@ public:
                 }
             } else {
                 _nwritten += PA_DATA_LEN;
-                if (_write_interval_us > 0) {
-                    bthread_usleep(_write_interval_us);
-                }
             }
             ++c;
         }
@@ -851,8 +840,6 @@ public:
 private:
     DonePlace _done_place;
     size_t _nrep;
-    int _write_interval_us;
-    int _initial_write_delay_us;
     size_t _nwritten;
     bool _ever_full;
     int _last_errno;
@@ -954,47 +941,6 @@ private:
     butil::Status _destroying_st;
 };
 
-class TimeoutReadBody : public brpc::ProgressiveReader,
-                        public brpc::SharedObject {
-public:
-    explicit TimeoutReadBody(int read_delay_us = 0, int read_error = 0)
-        : _read_delay_us(read_delay_us)
-        , _read_error(read_error)
-        , _nread(0)
-        , _nend(0)
-        , _end_error(0) {
-        butil::intrusive_ptr<TimeoutReadBody>(this).detach();
-    }
-
-    butil::Status OnReadOnePart(const void*, size_t length) override {
-        if (_read_delay_us > 0) {
-            bthread_usleep(_read_delay_us);
-        }
-        _nread.fetch_add(length);
-        if (_read_error != 0) {
-            return butil::Status(_read_error, "intended progressive read 
failure");
-        }
-        return butil::Status::OK();
-    }
-
-    void OnEndOfMessage(const butil::Status& status) override {
-        _end_error.store(status.error_code());
-        _nend.fetch_add(1);
-        butil::intrusive_ptr<TimeoutReadBody>(this, false);
-    }
-
-    size_t read_bytes() const { return _nread.load(); }
-    int end_count() const { return _nend.load(); }
-    int end_error() const { return _end_error.load(); }
-
-private:
-    const int _read_delay_us;
-    const int _read_error;
-    std::atomic<size_t> _nread;
-    std::atomic<int> _nend;
-    std::atomic<int> _end_error;
-};
-
 #ifdef BUTIL_USE_ASAN
 static const int GENERAL_DELAY_US = 1000000; // 1s
 #else
@@ -1088,187 +1034,6 @@ TEST_F(HttpTest, read_short_body_progressively) {
     }
 }
 
-TEST_F(HttpTest, progressive_read_timeout_keeps_active_reader_alive) {
-    const int port = 8923;
-    brpc::Server server;
-    DownloadServiceImpl svc(DONE_BEFORE_CREATE_PA, 8, 100000);
-    ASSERT_EQ(0, server.AddService(&svc, brpc::SERVER_DOESNT_OWN_SERVICE));
-    ASSERT_EQ(0, server.Start(port, NULL));
-
-    brpc::Channel channel;
-    brpc::ChannelOptions options;
-    options.protocol = brpc::PROTOCOL_HTTP;
-    ASSERT_EQ(0, channel.Init(butil::EndPoint(butil::my_ip(), port), 
&options));
-
-    brpc::Controller cntl;
-    cntl.response_will_be_read_progressively();
-    cntl.set_progressive_read_timeout_ms(500);
-    cntl.http_request().uri() = "/DownloadService/Download";
-    channel.CallMethod(NULL, &cntl, NULL, NULL, NULL);
-    ASSERT_FALSE(cntl.Failed()) << cntl.ErrorText();
-
-    butil::intrusive_ptr<TimeoutReadBody> reader(new TimeoutReadBody);
-    cntl.ReadProgressiveAttachmentBy(reader.get());
-    for (int i = 0; i < 200 && reader->end_count() == 0; ++i) {
-        bthread_usleep(10000);
-    }
-    ASSERT_EQ(1, reader->end_count());
-    EXPECT_EQ(0, reader->end_error());
-    EXPECT_EQ(8 * PA_DATA_LEN, reader->read_bytes());
-}
-
-TEST_F(HttpTest, progressive_read_timeout_closes_idle_http1_reader_once) {
-    const int port = 8923;
-    brpc::Server server;
-    DownloadServiceImpl svc(DONE_BEFORE_CREATE_PA, 2, 300000);
-    ASSERT_EQ(0, server.AddService(&svc, brpc::SERVER_DOESNT_OWN_SERVICE));
-    ASSERT_EQ(0, server.Start(port, NULL));
-
-    butil::intrusive_ptr<TimeoutReadBody> reader(new TimeoutReadBody);
-    {
-        brpc::Channel channel;
-        brpc::ChannelOptions options;
-        options.protocol = brpc::PROTOCOL_HTTP;
-        ASSERT_EQ(0, channel.Init(butil::EndPoint(butil::my_ip(), port), 
&options));
-        {
-            brpc::Controller cntl;
-            cntl.response_will_be_read_progressively();
-            cntl.set_progressive_read_timeout_ms(50);
-            cntl.http_request().uri() = "/DownloadService/Download";
-            channel.CallMethod(NULL, &cntl, NULL, NULL, NULL);
-            ASSERT_FALSE(cntl.Failed()) << cntl.ErrorText();
-            cntl.ReadProgressiveAttachmentBy(reader.get());
-            bthread_usleep(400000);
-            ASSERT_NE(0, svc.last_errno());
-            EXPECT_EQ(0, reader->end_count());
-        }
-    }
-    for (int i = 0; i < 100 && reader->end_count() == 0; ++i) {
-        bthread_usleep(10000);
-    }
-    ASSERT_EQ(1, reader->end_count());
-    EXPECT_EQ(brpc::EPROGREADTIMEOUT, reader->end_error());
-    bthread_usleep(400000);
-    EXPECT_EQ(1, reader->end_count());
-}
-
-TEST_F(HttpTest, progressive_read_timeout_before_first_body_part) {
-    const int port = 8923;
-    brpc::Server server;
-    DownloadServiceImpl svc(DONE_BEFORE_CREATE_PA, 1, 0, 300000);
-    ASSERT_EQ(0, server.AddService(&svc, brpc::SERVER_DOESNT_OWN_SERVICE));
-    ASSERT_EQ(0, server.Start(port, NULL));
-
-    butil::intrusive_ptr<TimeoutReadBody> reader(new TimeoutReadBody);
-    {
-        brpc::Channel channel;
-        brpc::ChannelOptions options;
-        options.protocol = brpc::PROTOCOL_HTTP;
-        ASSERT_EQ(0, channel.Init(butil::EndPoint(butil::my_ip(), port), 
&options));
-        {
-            brpc::Controller cntl;
-            cntl.response_will_be_read_progressively();
-            cntl.set_progressive_read_timeout_ms(50);
-            cntl.http_request().uri() = "/DownloadService/Download";
-            channel.CallMethod(NULL, &cntl, NULL, NULL, NULL);
-            ASSERT_FALSE(cntl.Failed()) << cntl.ErrorText();
-            cntl.ReadProgressiveAttachmentBy(reader.get());
-            bthread_usleep(400000);
-            ASSERT_NE(0, svc.last_errno());
-            EXPECT_EQ(size_t(0), reader->read_bytes());
-            EXPECT_EQ(0, reader->end_count());
-        }
-    }
-    for (int i = 0; i < 100 && reader->end_count() == 0; ++i) {
-        bthread_usleep(10000);
-    }
-    ASSERT_EQ(1, reader->end_count());
-    EXPECT_EQ(brpc::EPROGREADTIMEOUT, reader->end_error());
-}
-
-TEST_F(HttpTest, progressive_read_timeout_ignores_slow_user_callback) {
-    const int port = 8923;
-    brpc::Server server;
-    DownloadServiceImpl svc(DONE_BEFORE_CREATE_PA, 3, 50000);
-    ASSERT_EQ(0, server.AddService(&svc, brpc::SERVER_DOESNT_OWN_SERVICE));
-    ASSERT_EQ(0, server.Start(port, NULL));
-
-    brpc::Channel channel;
-    brpc::ChannelOptions options;
-    options.protocol = brpc::PROTOCOL_HTTP;
-    ASSERT_EQ(0, channel.Init(butil::EndPoint(butil::my_ip(), port), 
&options));
-
-    brpc::Controller cntl;
-    cntl.response_will_be_read_progressively();
-    cntl.set_progressive_read_timeout_ms(50);
-    cntl.http_request().uri() = "/DownloadService/Download";
-    channel.CallMethod(NULL, &cntl, NULL, NULL, NULL);
-    ASSERT_FALSE(cntl.Failed()) << cntl.ErrorText();
-
-    butil::intrusive_ptr<TimeoutReadBody> reader(
-        new TimeoutReadBody(200000));
-    cntl.ReadProgressiveAttachmentBy(reader.get());
-    for (int i = 0; i < 100 && reader->end_count() == 0; ++i) {
-        bthread_usleep(10000);
-    }
-    ASSERT_EQ(1, reader->end_count());
-    EXPECT_EQ(0, reader->end_error());
-    EXPECT_EQ(3 * PA_DATA_LEN, reader->read_bytes());
-}
-
-TEST_F(HttpTest, progressive_read_timeout_preserves_reader_error) {
-    const int port = 8923;
-    brpc::Server server;
-    DownloadServiceImpl svc(DONE_BEFORE_CREATE_PA, 10);
-    ASSERT_EQ(0, server.AddService(&svc, brpc::SERVER_DOESNT_OWN_SERVICE));
-    ASSERT_EQ(0, server.Start(port, NULL));
-
-    brpc::Channel channel;
-    brpc::ChannelOptions options;
-    options.protocol = brpc::PROTOCOL_HTTP;
-    ASSERT_EQ(0, channel.Init(butil::EndPoint(butil::my_ip(), port), 
&options));
-
-    brpc::Controller cntl;
-    cntl.response_will_be_read_progressively();
-    cntl.set_progressive_read_timeout_ms(1000);
-    cntl.http_request().uri() = "/DownloadService/Download";
-    channel.CallMethod(NULL, &cntl, NULL, NULL, NULL);
-    ASSERT_FALSE(cntl.Failed()) << cntl.ErrorText();
-
-    butil::intrusive_ptr<TimeoutReadBody> reader(
-        new TimeoutReadBody(0, EIO));
-    cntl.ReadProgressiveAttachmentBy(reader.get());
-    ASSERT_EQ(1, reader->end_count());
-    EXPECT_EQ(EIO, reader->end_error());
-}
-
-TEST_F(HttpTest, progressive_read_timeout_rejects_http2) {
-    const int port = 8923;
-    brpc::Server server;
-    ASSERT_EQ(0, server.AddService(&_svc, brpc::SERVER_DOESNT_OWN_SERVICE));
-    ASSERT_EQ(0, server.Start(port, NULL));
-
-    brpc::Channel channel;
-    brpc::ChannelOptions options;
-    options.protocol = brpc::PROTOCOL_H2;
-    ASSERT_EQ(0, channel.Init(butil::EndPoint(butil::my_ip(), port), 
&options));
-
-    brpc::Controller cntl;
-    cntl.response_will_be_read_progressively();
-    cntl.set_progressive_read_timeout_ms(1000);
-    cntl.http_request().uri() = "/EchoService/Echo";
-    test::EchoRequest req;
-    req.set_message(EXP_REQUEST);
-    channel.CallMethod(NULL, &cntl, &req, NULL, NULL);
-    ASSERT_FALSE(cntl.Failed()) << cntl.ErrorText();
-
-    butil::intrusive_ptr<TimeoutReadBody> reader(new TimeoutReadBody);
-    cntl.ReadProgressiveAttachmentBy(reader.get());
-    ASSERT_EQ(1, reader->end_count());
-    EXPECT_EQ(ENOTSUP, reader->end_error());
-    EXPECT_EQ(size_t(0), reader->read_bytes());
-}
-
 TEST_F(HttpTest, read_progressively_after_cntl_destroys) {
     DownloadServiceImpl svc(DONE_BEFORE_CREATE_PA,
                             std::numeric_limits<size_t>::max());


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

Reply via email to