From c478f0833d25af3978cca07a71914a56122ca304 Mon Sep 17 00:00:00 2001 From: Pavel Misko Date: Fri, 17 Jul 2026 19:03:03 +0200 Subject: [PATCH] Journalled Devices server: handle backend exceptions --- .../journalled_device_tcp_server/server.cpp | 155 ++++++---- .../server_ut.cpp | 275 ++++++++++++++---- 2 files changed, 310 insertions(+), 120 deletions(-) diff --git a/cloud/storage/core/libs/journalled_device_tcp_server/server.cpp b/cloud/storage/core/libs/journalled_device_tcp_server/server.cpp index aedaf4ddd7a..cab35d307c2 100644 --- a/cloud/storage/core/libs/journalled_device_tcp_server/server.cpp +++ b/cloud/storage/core/libs/journalled_device_tcp_server/server.cpp @@ -21,6 +21,63 @@ namespace { //////////////////////////////////////////////////////////////////////////////// +#define STORAGE_JD_SERVER(xxx, ...) \ + xxx(AcquireDevices, __VA_ARGS__) \ + xxx(ReleaseDevices, __VA_ARGS__) \ + xxx(ReadPages, __VA_ARGS__) \ + xxx(WriteLogRecord, __VA_ARGS__) + +// STORAGE_JD_SERVER + +template < + typename TProtoRequest, + typename TProtoResponse, + TFuture (IServerBackend::*M)(TInstant, TProtoRequest), + TProtoRequest* (NProto::TDeviceProtocolRequest::*F1)(), + TProtoResponse* (NProto::TDeviceProtocolResponse::*F2)()> +struct TServerMethod +{ + using TRequest = TProtoRequest; + using TResponse = TProtoResponse; + + static auto Execute( + IServerBackend& backend, + TInstant now, + TRequest&& request) -> TFuture + { + return (backend.*M)(now, std::move(request)); + } + + static auto MutableProto(NProto::TDeviceProtocolRequest& request) + -> TRequest& + { + return *(request.*F1)(); + } + + static auto MutableProto(NProto::TDeviceProtocolResponse& response) + -> TResponse& + { + return *(response.*F2)(); + } +}; + +#define STORAGE_DECLARE_METHOD(name, ...) \ + struct T##name##Method \ + : TServerMethod< \ + NProto::T##name##Request, \ + NProto::T##name##Response, \ + &IServerBackend::name, \ + &NProto::TDeviceProtocolRequest::Mutable##name, \ + &NProto::TDeviceProtocolResponse::Mutable##name> \ + { \ + constexpr static TStringBuf Name = #name; \ + }; \ + // STORAGE_DECLARE_METHOD + +STORAGE_JD_SERVER(STORAGE_DECLARE_METHOD) + +//////////////////////////////////////////////////////////////////////////////// + struct TConnection { TSocketHolder Socket; @@ -120,6 +177,40 @@ class TServer final void HandleRequest( NProto::TDeviceProtocolRequest& request, TConnectionPtr conn); + + template + void ProcessRequest( + TConnectionPtr conn, + NProto::TDeviceProtocolRequest&& request) + { + using TResponse = typename TMethod::TResponse; + + TFuture future; + + const ui64 requestId = request.GetRequestId(); + + try { + auto& proto = TMethod::MutableProto(request); + + future = TMethod::Execute(*Backend, Now(), std::move(proto)); + } catch (...) { + STORAGE_ERROR( + TMethod::Name << " failed: " << CurrentExceptionMessage()); + + future = MakeErrorFuture(std::current_exception()); + } + + future.Subscribe( + [conn, requestId](const auto& future) + { + NProto::TDeviceProtocolResponse response; + response.SetRequestId(requestId); + TMethod::MutableProto(response).CopyFrom( + SafeExecute([&] { return future.GetValue(); })); + + conn->ResponseQueue.Enqueue(std::move(response)); + }); + } }; //////////////////////////////////////////////////////////////////////////////// @@ -314,75 +405,19 @@ void TServer::HandleRequest( switch (request.GetRequestCase()) { case ERequestCase::kAcquireDevices: { - auto future = Backend->AcquireDevices( - Now(), - std::move(*request.MutableAcquireDevices())); - - future.Subscribe( - [conn, requestId]( - const TFuture& future) - { - NProto::TDeviceProtocolResponse response; - response.SetRequestId(requestId); - response.MutableAcquireDevices()->CopyFrom( - future.GetValue()); - - conn->ResponseQueue.Enqueue(std::move(response)); - }); - + ProcessRequest(conn, std::move(request)); break; } case ERequestCase::kReleaseDevices: { - auto future = Backend->ReleaseDevices( - Now(), - std::move(*request.MutableReleaseDevices())); - - future.Subscribe( - [conn, requestId]( - const TFuture& future) - { - NProto::TDeviceProtocolResponse response; - response.SetRequestId(requestId); - response.MutableReleaseDevices()->CopyFrom( - future.GetValue()); - - conn->ResponseQueue.Enqueue(std::move(response)); - }); + ProcessRequest(conn, std::move(request)); break; } case ERequestCase::kReadPages: { - auto future = Backend->ReadPages( - Now(), - std::move(*request.MutableReadPages())); - - future.Subscribe( - [conn, - requestId](const TFuture& future) - { - NProto::TDeviceProtocolResponse response; - response.SetRequestId(requestId); - response.MutableReadPages()->CopyFrom(future.GetValue()); - - conn->ResponseQueue.Enqueue(std::move(response)); - }); + ProcessRequest(conn, std::move(request)); break; } case ERequestCase::kWriteLogRecord: { - auto future = Backend->WriteLogRecord( - Now(), - std::move(*request.MutableWriteLogRecord())); - - future.Subscribe( - [conn, requestId]( - const TFuture& future) - { - NProto::TDeviceProtocolResponse response; - response.SetRequestId(requestId); - response.MutableWriteLogRecord()->CopyFrom( - future.GetValue()); - - conn->ResponseQueue.Enqueue(std::move(response)); - }); + ProcessRequest(conn, std::move(request)); break; } default: { diff --git a/cloud/storage/core/libs/journalled_device_tcp_server/server_ut.cpp b/cloud/storage/core/libs/journalled_device_tcp_server/server_ut.cpp index ef1c3ab7563..3734b07c19c 100644 --- a/cloud/storage/core/libs/journalled_device_tcp_server/server_ut.cpp +++ b/cloud/storage/core/libs/journalled_device_tcp_server/server_ut.cpp @@ -38,13 +38,10 @@ struct TTestBackend: public IServerBackend std::function( NProto::TWriteLogRecordRequest)>; - TPromise AcquireDevicesImpl = - NewPromise(); - TPromise ReleaseDevicesImpl = - NewPromise(); - TPromise ReadPagesImpl = NewPromise(); - TPromise WriteLogRecordImpl = - NewPromise(); + TAcquireDevicesFunc AcquireDevicesImpl; + TReleaseDevicesFunc ReleaseDevicesImpl; + TReadPagesFunc ReadPagesImpl; + TWriteLogRecordFunc WriteLogRecordImpl; [[nodiscard]] auto AcquireDevices( TInstant now, @@ -53,9 +50,7 @@ struct TTestBackend: public IServerBackend { Y_UNUSED(now); - return AcquireDevicesImpl.GetFuture().Apply( - [request](const auto& future) - { return future.GetValue()(request); }); + return AcquireDevicesImpl(std::move(request)); } [[nodiscard]] auto ReleaseDevices( @@ -65,9 +60,7 @@ struct TTestBackend: public IServerBackend { Y_UNUSED(now); - return ReleaseDevicesImpl.GetFuture().Apply( - [request](const auto& future) - { return future.GetValue()(request); }); + return ReleaseDevicesImpl(std::move(request)); } [[nodiscard]] auto ReadPages( @@ -77,9 +70,7 @@ struct TTestBackend: public IServerBackend { Y_UNUSED(now); - return ReadPagesImpl.GetFuture().Apply( - [request](const auto& future) - { return future.GetValue()(request); }); + return ReadPagesImpl(std::move(request)); } [[nodiscard]] auto WriteLogRecord( @@ -89,9 +80,7 @@ struct TTestBackend: public IServerBackend { Y_UNUSED(now); - return WriteLogRecordImpl.GetFuture().Apply( - [request](const auto& future) - { return future.GetValue()(request); }); + return WriteLogRecordImpl(std::move(request)); } }; @@ -361,6 +350,40 @@ Y_UNIT_TEST_SUITE(TDeviceTCPServerTest) return proto; }(); + std::mutex mutex; + std::optional acquireDevicesRequest; + std::optional releaseDevicesRequest; + std::optional readPagesRequest; + std::optional writeLogRecordRequest; + + Backend->AcquireDevicesImpl = [&](auto request) + { + std::unique_lock lock(mutex); + acquireDevicesRequest = std::move(request); + return MakeFuture(expectedAcquireDevicesResponse); + }; + + Backend->ReleaseDevicesImpl = [&](auto request) + { + std::unique_lock lock(mutex); + releaseDevicesRequest = std::move(request); + return MakeFuture(expectedReleaseDevicesResponse); + }; + + Backend->ReadPagesImpl = [&](auto request) + { + std::unique_lock lock(mutex); + readPagesRequest = std::move(request); + return MakeFuture(expectedReadPagesResponse); + }; + + Backend->WriteLogRecordImpl = [&](auto request) + { + std::unique_lock lock(mutex); + writeLogRecordRequest = std::move(request); + return MakeFuture(expectedWriteLogRecordResponse); + }; + TTestClient client{Port}; { @@ -394,44 +417,6 @@ Y_UNIT_TEST_SUITE(TDeviceTCPServerTest) client.Send(request); } - std::mutex mutex; - std::optional acquireDevicesRequest; - std::optional releaseDevicesRequest; - std::optional readPagesRequest; - std::optional writeLogRecordRequest; - - Backend->AcquireDevicesImpl.SetValue( - [&](auto request) - { - std::unique_lock lock(mutex); - acquireDevicesRequest = std::move(request); - return MakeFuture(expectedAcquireDevicesResponse); - }); - - Backend->ReleaseDevicesImpl.SetValue( - [&](auto request) - { - std::unique_lock lock(mutex); - releaseDevicesRequest = std::move(request); - return MakeFuture(expectedReleaseDevicesResponse); - }); - - Backend->ReadPagesImpl.SetValue( - [&](auto request) - { - std::unique_lock lock(mutex); - readPagesRequest = std::move(request); - return MakeFuture(expectedReadPagesResponse); - }); - - Backend->WriteLogRecordImpl.SetValue( - [&](auto request) - { - std::unique_lock lock(mutex); - writeLogRecordRequest = std::move(request); - return MakeFuture(expectedWriteLogRecordResponse); - }); - TVector responses; for (ui32 i = 0; i != 4; ++i) { @@ -500,9 +485,10 @@ Y_UNIT_TEST_SUITE(TDeviceTCPServerTest) { const ui32 requestCount = 100; - Backend->AcquireDevicesImpl.SetValue( - [&](auto) - { return MakeFuture(NProto::TAcquireDevicesResponse()); }); + Backend->AcquireDevicesImpl = [&](auto) + { + return MakeFuture(NProto::TAcquireDevicesResponse()); + }; TTestClient client1{Port}; TTestClient client2{Port}; @@ -557,6 +543,175 @@ Y_UNIT_TEST_SUITE(TDeviceTCPServerTest) queue.Stop(); } + + Y_UNIT_TEST_F(ShouldHandleBackendExecptions, TFixture) + { + const ui64 requestId = 42; + + TTestClient client{Port}; + + auto acquire = [&] + { + NProto::TDeviceProtocolRequest request; + request.SetRequestId(requestId); + request.MutableAcquireDevices(); + client.Send(request); + + auto response = client.Receive(); + + UNIT_ASSERT_VALUES_EQUAL(requestId, response.GetRequestId()); + return response.GetAcquireDevices().GetError(); + }; + + auto release = [&] + { + NProto::TDeviceProtocolRequest request; + request.SetRequestId(requestId); + request.MutableReleaseDevices(); + client.Send(request); + + auto response = client.Receive(); + + UNIT_ASSERT_VALUES_EQUAL(requestId, response.GetRequestId()); + return response.GetReleaseDevices().GetError(); + }; + + auto readPages = [&] + { + NProto::TDeviceProtocolRequest request; + request.SetRequestId(requestId); + request.MutableReadPages(); + client.Send(request); + + auto response = client.Receive(); + + UNIT_ASSERT_VALUES_EQUAL(requestId, response.GetRequestId()); + return response.GetReadPages().GetError(); + }; + + auto writeLogRecord = [&] + { + NProto::TDeviceProtocolRequest request; + request.SetRequestId(requestId); + request.MutableWriteLogRecord(); + client.Send(request); + + auto response = client.Receive(); + + UNIT_ASSERT_VALUES_EQUAL(requestId, response.GetRequestId()); + return response.GetWriteLogRecord().GetError(); + }; + + Backend->AcquireDevicesImpl = + [](auto) -> TFuture + { + throw TServiceError(E_FAIL) << "acquire-inline-error"; + }; + + Backend->ReleaseDevicesImpl = + [](auto) -> TFuture + { + throw TServiceError(E_FAIL) << "release-inline-error"; + }; + + Backend->ReadPagesImpl = [](auto) -> TFuture + { + throw TServiceError(E_FAIL) << "readPages-inline-error"; + }; + Backend->WriteLogRecordImpl = + [](auto) -> TFuture + { + throw TServiceError(E_FAIL) << "writeLogRecord-inline-error"; + }; + + { + auto error = acquire(); + UNIT_ASSERT_VALUES_EQUAL(E_FAIL, error.GetCode()); + UNIT_ASSERT_VALUES_EQUAL( + "acquire-inline-error", + error.GetMessage()); + } + + { + auto error = release(); + UNIT_ASSERT_VALUES_EQUAL(E_FAIL, error.GetCode()); + UNIT_ASSERT_VALUES_EQUAL( + "release-inline-error", + error.GetMessage()); + } + + { + auto error = readPages(); + UNIT_ASSERT_VALUES_EQUAL(E_FAIL, error.GetCode()); + UNIT_ASSERT_VALUES_EQUAL( + "readPages-inline-error", + error.GetMessage()); + } + + { + auto error = writeLogRecord(); + UNIT_ASSERT_VALUES_EQUAL(E_FAIL, error.GetCode()); + UNIT_ASSERT_VALUES_EQUAL( + "writeLogRecord-inline-error", + error.GetMessage()); + } + + Backend->AcquireDevicesImpl = [](auto) + { + return MakeErrorFuture( + std::make_exception_ptr( + TServiceError{E_ARGUMENT} << "acquire-async-error")); + }; + + Backend->ReleaseDevicesImpl = [](auto) + { + return MakeErrorFuture( + std::make_exception_ptr( + TServiceError{E_ARGUMENT} << "release-async-error")); + }; + + Backend->ReadPagesImpl = [](auto) + { + return MakeErrorFuture( + std::make_exception_ptr( + TServiceError{E_ARGUMENT} << "readPages-async-error")); + }; + + Backend->WriteLogRecordImpl = [](auto) + { + return MakeErrorFuture( + std::make_exception_ptr( + TServiceError{E_ARGUMENT} << "writeLogRecord-async-error")); + }; + + { + auto error = acquire(); + UNIT_ASSERT_VALUES_EQUAL(E_ARGUMENT, error.GetCode()); + UNIT_ASSERT_VALUES_EQUAL("acquire-async-error", error.GetMessage()); + } + + { + auto error = release(); + UNIT_ASSERT_VALUES_EQUAL(E_ARGUMENT, error.GetCode()); + UNIT_ASSERT_VALUES_EQUAL("release-async-error", error.GetMessage()); + } + + { + auto error = readPages(); + UNIT_ASSERT_VALUES_EQUAL(E_ARGUMENT, error.GetCode()); + UNIT_ASSERT_VALUES_EQUAL( + "readPages-async-error", + error.GetMessage()); + } + + { + auto error = writeLogRecord(); + UNIT_ASSERT_VALUES_EQUAL(E_ARGUMENT, error.GetCode()); + UNIT_ASSERT_VALUES_EQUAL( + "writeLogRecord-async-error", + error.GetMessage()); + } + } } } // namespace NCloud::NJournalled