From e6a4def4e0e380393612234096a0e4f2e9f7ce6d Mon Sep 17 00:00:00 2001 From: Aleksei Kobzev Date: Thu, 20 Aug 2026 16:02:08 +0300 Subject: [PATCH] Migrate logbroker to standalone YDB C++ SDK --- .../logbroker/iface/credentials_provider.cpp | 2 +- cloud/blockstore/libs/logbroker/iface/ya.make | 2 +- .../logbroker/topic_api_impl/topic_api.cpp | 54 ++- .../libs/logbroker/topic_api_impl/topic_api.h | 15 + .../logbroker/topic_api_impl/topic_api_ut.cpp | 350 ++++++++---------- .../libs/logbroker/topic_api_impl/ut/ya.make | 11 - .../libs/logbroker/topic_api_impl/ya.make | 6 +- 7 files changed, 218 insertions(+), 222 deletions(-) diff --git a/cloud/blockstore/libs/logbroker/iface/credentials_provider.cpp b/cloud/blockstore/libs/logbroker/iface/credentials_provider.cpp index 639444271be..fe631b35b70 100644 --- a/cloud/blockstore/libs/logbroker/iface/credentials_provider.cpp +++ b/cloud/blockstore/libs/logbroker/iface/credentials_provider.cpp @@ -2,7 +2,7 @@ #include -#include +#include #include diff --git a/cloud/blockstore/libs/logbroker/iface/ya.make b/cloud/blockstore/libs/logbroker/iface/ya.make index 3cb48df1ee1..09f100e5a5c 100644 --- a/cloud/blockstore/libs/logbroker/iface/ya.make +++ b/cloud/blockstore/libs/logbroker/iface/ya.make @@ -14,7 +14,7 @@ PEERDIR( library/cpp/monlib/service/pages library/cpp/threading/future - contrib/ydb/public/sdk/cpp/src/library/jwt + contrib/libs/ydb-cpp-sdk/src/client/iam ) END() diff --git a/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.cpp b/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.cpp index 605db1ea6e5..8dce397ec06 100644 --- a/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.cpp +++ b/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.cpp @@ -7,9 +7,9 @@ #include #include -#include -#include -#include +#include +#include +#include #include #include @@ -91,16 +91,19 @@ class TService final std::shared_ptr CredentialsProviderFactory; + TWriteSessionFactory WriteSessionFactory; public: TService( TLogbrokerConfigPtr config, ILoggingServicePtr logging, std::shared_ptr - credentialsProviderFactory) + credentialsProviderFactory, + TWriteSessionFactory writeSessionFactory) : Config(std::move(config)) , Logging(std::move(logging)) , CredentialsProviderFactory(std::move(credentialsProviderFactory)) + , WriteSessionFactory(std::move(writeSessionFactory)) {} TFuture Write( @@ -119,15 +122,21 @@ class TService final Batch = std::make_unique(std::move(messages), now); if (!Session) { - TTopicClient client{GetDriver()}; - - Session = client.CreateWriteSession( - TWriteSessionSettings() - .Path(Config->GetTopic()) - .ProducerId(Config->GetSourceId()) - .MessageGroupId(Config->GetSourceId()) - .RetryPolicy( - NYdb::NTopic::IRetryPolicy::GetNoRetryPolicy())); + if (WriteSessionFactory) { + Session = WriteSessionFactory(); + } else { + TTopicClient client{GetDriver()}; + + Session = client.CreateWriteSession( + TWriteSessionSettings() + .Path(Config->GetTopic()) + .ProducerId(Config->GetSourceId()) + .MessageGroupId(Config->GetSourceId()) + .RetryPolicy( + NYdb::NTopic::IRetryPolicy::GetNoRetryPolicy())); + } + + Y_ABORT_UNLESS(Session); } else if (ContinuationToken.has_value()) { TContinuationToken token{std::move(ContinuationToken.value())}; ContinuationToken.reset(); @@ -314,13 +323,26 @@ class TService final IServicePtr CreateTopicAPIService( TLogbrokerConfigPtr config, ILoggingServicePtr logging, - std::shared_ptr - credentialsProviderFactory) + std::shared_ptr credentialsProviderFactory, + TWriteSessionFactory writeSessionFactory) { return std::make_shared( std::move(config), std::move(logging), - std::move(credentialsProviderFactory)); + std::move(credentialsProviderFactory), + std::move(writeSessionFactory)); +} + +IServicePtr CreateTopicAPIService( + TLogbrokerConfigPtr config, + ILoggingServicePtr logging, + std::shared_ptr credentialsProviderFactory) +{ + return CreateTopicAPIService( + std::move(config), + std::move(logging), + std::move(credentialsProviderFactory), + {}); } IServicePtr CreateTopicAPIService( diff --git a/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.h b/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.h index 79d358f30b8..d524d0f829c 100644 --- a/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.h +++ b/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.h @@ -5,10 +5,25 @@ #include +#include + +namespace NYdb::inline V3::NTopic { +class IWriteSession; +} // namespace NYdb::inline V3::NTopic + namespace NCloud::NBlockStore::NLogbroker { //////////////////////////////////////////////////////////////////////////////// +using TWriteSessionFactory = + std::function()>; + +IServicePtr CreateTopicAPIService( + TLogbrokerConfigPtr config, + ILoggingServicePtr logging, + std::shared_ptr credentialsProviderFactory, + TWriteSessionFactory writeSessionFactory); + IServicePtr CreateTopicAPIService( TLogbrokerConfigPtr config, ILoggingServicePtr logging, diff --git a/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api_ut.cpp b/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api_ut.cpp index 92db0fddab7..c5c2cfd6357 100644 --- a/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api_ut.cpp +++ b/cloud/blockstore/libs/logbroker/topic_api_impl/topic_api_ut.cpp @@ -1,182 +1,191 @@ #include "topic_api.h" -#include - #include #include -#include #include #include -#include +#include #include -#include -#include -#include +#include -#include +#include namespace NCloud::NBlockStore::NLogbroker { using namespace NThreading; -using namespace std::chrono_literals; +using namespace NYdb::NTopic; namespace { //////////////////////////////////////////////////////////////////////////////// -const TString TestConsumer = "test-consumer"; -const TString TestTopic = "test-topic"; -const TString TestSource = "test-source"; -const TString Database = "/Root"; -const TDuration WaitTimeout = 15s; +class TContinuationTokenIssuer final + : private NYdb::NTopic::TContinuationTokenIssuer +{ +public: + static TContinuationToken Issue() + { + return IssueContinuationToken(); + } +}; //////////////////////////////////////////////////////////////////////////////// -NKikimr::Tests::TServerSettings MakeServerSettings() +class TTestWriteSession final + : public IWriteSession { - using namespace NKikimrServices; - using namespace NActors::NLog; - - auto settings = NKikimr::NPersQueueTests::PQSettings(0); - settings.SetDomainName("Root"); - settings.SetNodeCount(1); - settings.PQConfig.SetTopicsAreFirstClassCitizen(true); - settings.PQConfig.SetRoot("/Root"); - settings.PQConfig.SetDatabase("/Root"); - - settings.SetLoggerInitializer([] (NActors::TTestActorRuntime& runtime) { - runtime.SetLogPriority(PQ_READ_PROXY, PRI_DEBUG); - runtime.SetLogPriority(PQ_WRITE_PROXY, PRI_DEBUG); - runtime.SetLogPriority(PQ_MIRRORER, PRI_DEBUG); - runtime.SetLogPriority(PQ_METACACHE, PRI_DEBUG); - runtime.SetLogPriority(PERSQUEUE, PRI_DEBUG); - runtime.SetLogPriority(PERSQUEUE_CLUSTER_TRACKER, PRI_DEBUG); - }); - - return settings; -} +private: + std::vector Events; + TVector Messages; -//////////////////////////////////////////////////////////////////////////////// +public: + explicit TTestWriteSession( + std::optional initialError = std::nullopt) + { + if (initialError) { + Events.emplace_back(TSessionClosedEvent( + *initialError, + NYdb::NIssue::TIssues{})); + } else { + AddReadyToAcceptEvent(); + } + } -struct TFixture - : public NUnitTest::TBaseFixture -{ - std::optional Server; + const TVector& GetMessages() const + { + return Messages; + } - std::optional Driver; - std::optional Client; + TFuture WaitEvent() override + { + UNIT_ASSERT(!Events.empty()); + return MakeFuture(); + } - ILoggingServicePtr Logging = CreateLoggingService( - "console", - { - .FiltrationLevel = TLOG_DEBUG, - }); + std::optional GetEvent(bool) override + { + if (Events.empty()) { + return std::nullopt; + } - void SetUp(NUnitTest::TTestContext& /*context*/) override + auto event = std::move(Events.front()); + Events.erase(Events.begin()); + return event; + } + + std::vector GetEvents( + bool, + std::optional) override { - Server.emplace(MakeServerSettings()); + return std::exchange(Events, {}); + } - Driver.emplace(NYdb::TDriverConfig() - .SetEndpoint("localhost:" + ToString(Server->GrpcPort)) - .SetDatabase(Database)); - Client.emplace(*Driver); + TFuture GetInitSeqNo() override + { + return MakeFuture(0); + } - auto settings = NYdb::NTopic::TCreateTopicSettings() - .PartitioningSettings(1, 1); + void Write( + TContinuationToken&&, + TWriteMessage&&, + TTransaction*) override + { + UNIT_FAIL("unexpected TWriteMessage overload"); + } - NYdb::NTopic::TConsumerSettings consumers( - settings, - TestConsumer); + void Write( + TContinuationToken&&, + std::string_view data, + std::optional seqNo, + std::optional) override + { + UNIT_ASSERT(seqNo.has_value()); + Messages.push_back({TString{data}, *seqNo}); - settings.AppendConsumers(consumers); + TWriteSessionEvent::TAcksEvent event; + event.Acks.push_back({ + .SeqNo = *seqNo, + .State = TWriteSessionEvent::TWriteAck::EES_WRITTEN, + }); + Events.emplace_back(std::move(event)); + AddReadyToAcceptEvent(); + } - auto status = Client->CreateTopic(TestTopic, settings) - .GetValueSync(); + void WriteEncoded( + TContinuationToken&&, + TWriteMessage&&, + TTransaction*) override + { + UNIT_FAIL("unexpected WriteEncoded overload"); + } - UNIT_ASSERT_C(status.IsSuccess(), status); + void WriteEncoded( + TContinuationToken&&, + std::string_view, + ECodec, + uint32_t, + std::optional, + std::optional) override + { + UNIT_FAIL("unexpected WriteEncoded overload"); + } - Server->WaitInit(TestTopic); + bool Close(TDuration) override + { + return true; } - auto Read(size_t count) + TWriterCounters::TPtr GetCounters() override { - auto session = Client->CreateReadSession(NYdb::NTopic::TReadSessionSettings() - .ConsumerName(TestConsumer) - .AppendTopics(std::string(TestTopic)) - .Decompress(true) - .Log(Logging->CreateLog("Read"))); - - TVector messages; - - bool sessionClosed = false; - - while (!sessionClosed) { - session->WaitEvent().Wait(); - - auto events = session->GetEvents(); - for (auto& event: events) { - std::visit(TOverloaded { - [&] (NYdb::NTopic::TReadSessionEvent::TDataReceivedEvent& ev) { - UNIT_ASSERT(!ev.HasCompressedMessages()); - - for (auto& m: ev.GetMessages()) { - messages.emplace_back(TMessage{ - TString{m.GetData()}, - m.GetSeqNo() - }); - } - - count -= std::min(ev.GetMessagesCount(), count); - - ev.Commit(); - - if (!count) { - session->Close(1s); - } - }, - [&] (NYdb::NTopic::TReadSessionEvent::TStartPartitionSessionEvent& ev) { - ev.Confirm(); - }, - [&] (NYdb::NTopic::TReadSessionEvent::TStopPartitionSessionEvent& ev) { - ev.Confirm(); - }, - [&] (NYdb::NTopic::TSessionClosedEvent& ev) { - UNIT_ASSERT_C(ev.IsSuccess(), ev.DebugString()); - sessionClosed = true; - }, - [] (const auto&) {} - }, event); - } - } + return {}; + } - return messages; +private: + void AddReadyToAcceptEvent() + { + Events.emplace_back(TWriteSessionEvent::TReadyToAcceptEvent( + TContinuationTokenIssuer::Issue())); } }; +//////////////////////////////////////////////////////////////////////////////// + +TLogbrokerConfigPtr CreateConfig() +{ + NProto::TLogbrokerConfig config; + config.SetTopic("test-topic"); + config.SetSourceId("test-source"); + return std::make_shared(std::move(config)); +} + +IServicePtr CreateTestService(std::shared_ptr session) +{ + auto logging = CreateLoggingService("console", TLogSettings{}); + + return CreateTopicAPIService( + CreateConfig(), + std::move(logging), + {}, + [session = std::move(session)] + { + return session; + }); +} + } // namespace //////////////////////////////////////////////////////////////////////////////// Y_UNIT_TEST_SUITE(TLogbrokerTest) { - Y_UNIT_TEST_F(ShouldWriteData, TFixture) + Y_UNIT_TEST(ShouldWriteData) { - NProto::TLogbrokerConfig config; - - config.SetAddress("localhost"); - config.SetPort(Server->GrpcPort); - config.SetDatabase(Database); - config.SetTopic(TestTopic); - config.SetSourceId(TestSource); - - auto service = CreateTopicAPIService( - std::make_shared(config), - Logging); - + auto session = std::make_shared(); + auto service = CreateTestService(session); service->Start(); const TVector expectedData{ @@ -186,83 +195,44 @@ Y_UNIT_TEST_SUITE(TLogbrokerTest) {"bar", 1001}, }; - { - auto future = - service->Write({expectedData[0], expectedData[1]}, Now()); + auto first = service->Write( + {expectedData[0], expectedData[1]}, + Now()).GetValueSync(); + UNIT_ASSERT_C(!HasError(first), FormatError(first)); - UNIT_ASSERT(future.Wait(WaitTimeout)); + auto second = service->Write( + {expectedData[2], expectedData[3]}, + Now()).GetValueSync(); + UNIT_ASSERT_C(!HasError(second), FormatError(second)); - const auto& error = future.GetValue(); - - UNIT_ASSERT_C(!HasError(error), FormatError(error)); - } - - { - auto future = - service->Write({expectedData[2], expectedData[3]}, Now()); - - UNIT_ASSERT(future.Wait(WaitTimeout)); - - const auto& error = future.GetValue(); - - UNIT_ASSERT_C(!HasError(error), FormatError(error)); - } - - auto data = Read(expectedData.size()); - - UNIT_ASSERT_VALUES_EQUAL(expectedData.size(), data.size()); + UNIT_ASSERT_VALUES_EQUAL( + expectedData.size(), + session->GetMessages().size()); for (size_t i = 0; i != expectedData.size(); ++i) { - UNIT_ASSERT_VALUES_EQUAL(expectedData[i].Payload, data[i].Payload); - UNIT_ASSERT_VALUES_EQUAL(expectedData[i].SeqNo, data[i].SeqNo); + UNIT_ASSERT_VALUES_EQUAL( + expectedData[i].Payload, + session->GetMessages()[i].Payload); + UNIT_ASSERT_VALUES_EQUAL( + expectedData[i].SeqNo, + session->GetMessages()[i].SeqNo); } service->Stop(); } - void ShouldHandleErrorImpl(TLogbrokerConfigPtr config) + Y_UNIT_TEST(ShouldHandleSessionError) { - auto logging = CreateLoggingService("console", TLogSettings{}); - - auto service = CreateTopicAPIService(config, logging); + auto session = std::make_shared( + NYdb::EStatus::UNAVAILABLE); + auto service = CreateTestService(std::move(session)); service->Start(); - auto future = service->Write({TMessage{"hello", 42}}, Now()); - - UNIT_ASSERT(future.Wait(WaitTimeout)); - - const auto& error = future.GetValue(); - UNIT_ASSERT(HasError(error)); + auto error = service->Write({TMessage{"hello", 42}}, Now()) + .GetValueSync(); + UNIT_ASSERT_VALUES_EQUAL(E_REJECTED, error.GetCode()); service->Stop(); } - - Y_UNIT_TEST_F(ShouldHandleConnectError, TFixture) - { - NProto::TLogbrokerConfig proto; - - proto.SetDatabase(Database); - proto.SetTopic(TestTopic); - proto.SetSourceId("test"); - proto.SetAddress("unknown"); - - ShouldHandleErrorImpl(std::make_shared(proto)); - } - - Y_UNIT_TEST_F(ShouldHandleUnknownTopic, TFixture) - { - NProto::TLogbrokerConfig proto; - - proto.SetAddress("localhost"); - proto.SetPort(Server->GrpcPort); - proto.SetDatabase(Database); - proto.SetTopic("unknown-topic"); - proto.SetSourceId(Sprintf( - "test:%s:%lu", - GetFQDNHostName(), - TInstant::Now().MilliSeconds())); - - ShouldHandleErrorImpl(std::make_shared(proto)); - } } } // namespace NCloud::NBlockStore::NLogbroker diff --git a/cloud/blockstore/libs/logbroker/topic_api_impl/ut/ya.make b/cloud/blockstore/libs/logbroker/topic_api_impl/ut/ya.make index d0713afab30..851529dc3d4 100644 --- a/cloud/blockstore/libs/logbroker/topic_api_impl/ut/ya.make +++ b/cloud/blockstore/libs/logbroker/topic_api_impl/ut/ya.make @@ -1,20 +1,9 @@ UNITTEST_FOR(cloud/blockstore/libs/logbroker/topic_api_impl) -INCLUDE(${ARCADIA_ROOT}/cloud/storage/core/tests/recipes/medium.inc) - -ADDINCL ( - contrib/ydb/public/sdk/cpp -) - SRCS( topic_api_ut.cpp ) -PEERDIR( - contrib/ydb/public/sdk/cpp/src/client/topic/ut/ut_utils - contrib/ydb/public/sdk/cpp/src/client/persqueue_public/ut/ut_utils -) - YQL_LAST_ABI_VERSION() END() diff --git a/cloud/blockstore/libs/logbroker/topic_api_impl/ya.make b/cloud/blockstore/libs/logbroker/topic_api_impl/ya.make index c9808c07874..2bb7899c185 100644 --- a/cloud/blockstore/libs/logbroker/topic_api_impl/ya.make +++ b/cloud/blockstore/libs/logbroker/topic_api_impl/ya.make @@ -9,9 +9,9 @@ PEERDIR( cloud/blockstore/libs/diagnostics cloud/blockstore/libs/logbroker/iface - contrib/ydb/public/sdk/cpp/src/client/iam - contrib/ydb/public/sdk/cpp/src/client/driver - contrib/ydb/public/sdk/cpp/src/client/topic + contrib/libs/ydb-cpp-sdk/src/client/iam + contrib/libs/ydb-cpp-sdk/src/client/driver + contrib/libs/ydb-cpp-sdk/src/client/topic library/cpp/threading/future )