From ab71516ca6760fe71bedeb59108e98b75da1c41b Mon Sep 17 00:00:00 2001 From: kulikov Date: Sun, 28 Jun 2026 20:01:43 +0300 Subject: Allow custom socket streams in library/cpp/http/server Prepare to use external coroutine/fiber/etc pool with non-blocking read/writes instead of plain blocking system threads. We already can replace system thread pool with something else, allow to inject different streams around connection socket: - move common part of http connection (http streams and output buffer) directly into it; - replace connection Impl with socket streams provider; - add TClientRequest::CreateHttpConnection factory; - check that we are able to override socket streams in unittest. commit_hash:afe39ce57ee1d10673f4c36a12e01b467d9f77b0 --- library/cpp/http/server/conn.cpp | 81 ++++++++++++--------------- library/cpp/http/server/conn.h | 33 +++++++++-- library/cpp/http/server/http.cpp | 10 ++-- library/cpp/http/server/http.h | 2 + library/cpp/http/server/http_ut.cpp | 109 ++++++++++++++++++++++++++++++++++++ 5 files changed, 179 insertions(+), 56 deletions(-) (limited to 'library/cpp/http/server') diff --git a/library/cpp/http/server/conn.cpp b/library/cpp/http/server/conn.cpp index 38a76c4c309..4b72c3f77f0 100644 --- a/library/cpp/http/server/conn.cpp +++ b/library/cpp/http/server/conn.cpp @@ -1,47 +1,38 @@ #include "conn.h" #include -#include -class THttpServerConn::TImpl { -public: - inline TImpl(const TSocket& s, size_t outputBufferSize) - : S_(s) - , SI_(S_) - , SO_(S_) - , BO_(&SO_, outputBufferSize) - , HI_(&SI_) - , HO_(&BO_, &HI_) - { - } - - inline ~TImpl() { - } - - inline THttpInput* Input() noexcept { - return &HI_; - } +namespace { + class TBlockingSocketStreams: public THttpServerConn::ISocketStreams { + public: + explicit TBlockingSocketStreams(const TSocket& s) + : Socket_(s) + , Input_(Socket_) + , Output_(Socket_) + { + } - inline THttpOutput* Output() noexcept { - return &HO_; - } + IInputStream* Input() override { + return &Input_; + } - inline void Reset() { - if (S_ != INVALID_SOCKET) { - // send RST packet to client - S_.SetLinger(true, 0); - S_.Close(); + IOutputStream* Output() override { + return &Output_; } - } -private: - TSocket S_; - TSocketInput SI_; - TSocketOutput SO_; - TBufferedOutput BO_; - THttpInput HI_; - THttpOutput HO_; -}; + void Reset() override { + if (Socket_ != INVALID_SOCKET) { + // send RST packet to client + Socket_.SetLinger(true, 0); + Socket_.Close(); + } + } + private: + TSocket Socket_; + TSocketInput Input_; + TSocketOutput Output_; + }; +} THttpServerConn::THttpServerConn(const TSocket& s) : THttpServerConn(s, s.MaximumTransferUnit()) @@ -49,21 +40,21 @@ THttpServerConn::THttpServerConn(const TSocket& s) } THttpServerConn::THttpServerConn(const TSocket& s, size_t outputBufferSize) - : Impl_(new TImpl(s, outputBufferSize)) + : THttpServerConn(MakeHolder(s), outputBufferSize) { } -THttpServerConn::~THttpServerConn() { -} - -THttpInput* THttpServerConn::Input() noexcept { - return Impl_->Input(); +THttpServerConn::THttpServerConn(THolder socketStreams, size_t outputBufferSize) + : SocketStreams_(std::move(socketStreams)) + , BufferedOutput_(SocketStreams_->Output(), outputBufferSize) + , HttpInput_(SocketStreams_->Input()) + , HttpOutput_(&BufferedOutput_, &HttpInput_) +{ } -THttpOutput* THttpServerConn::Output() noexcept { - return Impl_->Output(); +THttpServerConn::~THttpServerConn() { } void THttpServerConn::Reset() { - return Impl_->Reset(); + return SocketStreams_->Reset(); } diff --git a/library/cpp/http/server/conn.h b/library/cpp/http/server/conn.h index 3aa5329af42..62a3e86ac95 100644 --- a/library/cpp/http/server/conn.h +++ b/library/cpp/http/server/conn.h @@ -2,18 +2,37 @@ #include #include +#include +class IInputStream; +class IOutputStream; class TSocket; /// Потоки ввода/вывода для получения запросов и отправки ответов HTTP-сервера. class THttpServerConn { public: - explicit THttpServerConn(const TSocket& s); + class ISocketStreams { + public: + virtual ~ISocketStreams() {} + + virtual IInputStream* Input() = 0; + virtual IOutputStream* Output() = 0; + virtual void Reset() = 0; + }; +public: + THttpServerConn(const TSocket& s); THttpServerConn(const TSocket& s, size_t outputBufferSize); + THttpServerConn(THolder socketStreams, size_t outputBufferSize); + ~THttpServerConn(); - THttpInput* Input() noexcept; - THttpOutput* Output() noexcept; + THttpInput* Input() noexcept { + return &HttpInput_; + } + + THttpOutput* Output() noexcept { + return &HttpOutput_; + } inline const THttpInput* Input() const noexcept { return const_cast(this)->Input(); @@ -30,8 +49,10 @@ public: } void Reset(); - private: - class TImpl; - THolder Impl_; + THolder SocketStreams_; + + TBufferedOutput BufferedOutput_; + THttpInput HttpInput_; + THttpOutput HttpOutput_; }; diff --git a/library/cpp/http/server/http.cpp b/library/cpp/http/server/http.cpp index 13088f5e442..79a1194c770 100644 --- a/library/cpp/http/server/http.cpp +++ b/library/cpp/http/server/http.cpp @@ -725,6 +725,10 @@ void TClientRequest::ResetConnection() { } } +THolder TClientRequest::CreateHttpConnection(const TSocket& s, size_t outputBufferSize) { + return MakeHolder(s, outputBufferSize); +} + void TClientRequest::Process(void* ThreadSpecificResource) { THolder this_(this); @@ -733,11 +737,7 @@ void TClientRequest::Process(void* ThreadSpecificResource) { try { if (!HttpConn_) { const size_t outputBufferSize = HttpServ()->Options().OutputBufferSize; - if (outputBufferSize) { - HttpConn_.Reset(new THttpServerConn(Socket(), outputBufferSize)); - } else { - HttpConn_.Reset(new THttpServerConn(Socket())); - } + HttpConn_ = CreateHttpConnection(Socket(), outputBufferSize ? outputBufferSize : Socket().MaximumTransferUnit()); auto maxRequestsPerConnection = HttpServ()->Options().MaxRequestsPerConnection; HttpConn_->Output()->EnableKeepAlive(HttpServ()->Options().KeepAliveEnabled && (!maxRequestsPerConnection || Conn_->ReceivedRequests < maxRequestsPerConnection)); diff --git a/library/cpp/http/server/http.h b/library/cpp/http/server/http.h index 62dfd5a50fa..5b23cb21a00 100644 --- a/library/cpp/http/server/http.h +++ b/library/cpp/http/server/http.h @@ -141,6 +141,8 @@ private: Y_UNUSED(ThreadSpecificResource); return true; } + virtual THolder CreateHttpConnection(const TSocket& s, size_t outputBuffer); + void Process(void* ThreadSpecificResource) override; public: diff --git a/library/cpp/http/server/http_ut.cpp b/library/cpp/http/server/http_ut.cpp index 0fc44fa5983..b3280745929 100644 --- a/library/cpp/http/server/http_ut.cpp +++ b/library/cpp/http/server/http_ut.cpp @@ -7,6 +7,7 @@ #include #include #include +#include #include #include #include @@ -1016,4 +1017,112 @@ Y_UNIT_TEST_SUITE(THttpServerTest) { t2->Join(); UNIT_ASSERT_EQUAL_C(results, (THashSet({"Zoooo", "TTL Exceed"})), "Results is {" + ToString(results) + "}"); } + + Y_UNIT_TEST(TestCustomSocketStreams) { + // Check that TClientRequest::CreateHttpConnection override works + + class TCustomSocketStreamsServer: public THttpServer::ICallBack { + public: + TStringStream InputCopy; + TStringStream OutputCopy; + private: + class TTeeInput + : public IInputStream + { + public: + TTeeInput(IInputStream* s, IOutputStream* copy) + : S_(s) + , Copy_(copy) + { + } + + size_t DoRead(void* buf, size_t len) override { + void* begin = buf; + size_t res = S_->Read(buf, len); + Copy_->Write(begin, res); + return res; + } + + private: + IInputStream* const S_; + IOutputStream* const Copy_; + }; + + class TBlockingSocketStreams: public THttpServerConn::ISocketStreams { + public: + explicit TBlockingSocketStreams(const TSocket& s, TCustomSocketStreamsServer* server) + : Socket_(s) + , Input_(Socket_) + , TI_(&Input_, &server->InputCopy) + , Output_(Socket_) + , TO_(&Output_, &server->OutputCopy) + { + } + + IInputStream* Input() override { + return &TI_; + } + + IOutputStream* Output() override { + return &TO_; + } + + void Reset() override {} + private: + TSocket Socket_; + TSocketInput Input_; + TTeeInput TI_; + TSocketOutput Output_; + TTeeOutput TO_; + }; + + class TRequest: public TClientRequest { + public: + TRequest(TCustomSocketStreamsServer* server) + : Server_(server) + { + } + + bool Reply(void* /*tsr*/) override { + Output() << "HTTP/1.1 200 Ok\r\n\r\n"; + Output().Finish(); + return true; + } + + THolder CreateHttpConnection(const TSocket& s, size_t outputBuffer) override { + return MakeHolder(MakeHolder(s, Server_), outputBuffer); + } + private: + TCustomSocketStreamsServer* const Server_; + }; + + public: + TClientRequest* CreateClient() override { + return new TRequest(this); + } + + void OnException() override { + ExceptionMessage = CurrentExceptionMessage(); + } + + TString ExceptionMessage; + }; + + TPortManager portManager; + const ui16 port = portManager.GetPort(); + TCustomSocketStreamsServer server; + THttpServer::TOptions options(port); + options.nThreads = 1; + options.MaxConnections = 2; + THttpServer srv(&server, options); + + UNIT_ASSERT(srv.Start()); + + TSocket socket(TNetworkAddress("localhost", port), TDuration::Seconds(10)); + + SendRequest(socket, port); + + UNIT_ASSERT_STRINGS_EQUAL(server.InputCopy.Str(), TStringBuilder() << "GET / HTTP/1.1\r\nHost: localhost:" << port << "\r\nConnection: Keep-Alive\r\n\r\n"); + UNIT_ASSERT_STRINGS_EQUAL(server.OutputCopy.Str(), TStringBuilder() << "HTTP/1.1 200 Ok\r\nConnection: Keep-Alive\r\nTransfer-Encoding: chunked\r\n\r\n0\r\n\r\n"); + } } -- cgit v1.3