Skip to content

Commit e486ab9

Browse files
committed
Migrate logbroker to standalone YDB C++ SDK
1 parent febc362 commit e486ab9

418 files changed

Lines changed: 69505 additions & 222 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

cloud/blockstore/libs/logbroker/iface/credentials_provider.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22

33
#include <cloud/blockstore/libs/logbroker/iface/config.h>
44

5-
#include <contrib/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/iam/iam.h>
5+
#include <ydb-cpp-sdk/client/iam/iam.h>
66

77
#include <util/string/cast.h>
88

cloud/blockstore/libs/logbroker/iface/ya.make

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ PEERDIR(
1414
library/cpp/monlib/service/pages
1515
library/cpp/threading/future
1616

17-
contrib/ydb/public/sdk/cpp/src/library/jwt
17+
contrib/libs/ydb-cpp-sdk/src/client/iam
1818
)
1919

2020
END()

cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.cpp

Lines changed: 38 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -7,9 +7,9 @@
77
#include <cloud/storage/core/libs/common/error.h>
88
#include <cloud/storage/core/libs/diagnostics/logging.h>
99

10-
#include <contrib/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/iam/iam.h>
11-
#include <contrib/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/topic/client.h>
12-
#include <contrib/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/status_codes.h>
10+
#include <ydb-cpp-sdk/client/iam/iam.h>
11+
#include <ydb-cpp-sdk/client/topic/client.h>
12+
#include <ydb-cpp-sdk/client/types/status_codes.h>
1313

1414
#include <util/generic/overloaded.h>
1515
#include <util/stream/file.h>
@@ -91,16 +91,19 @@ class TService final
9191

9292
std::shared_ptr<NYdbICredentialsProviderFactory>
9393
CredentialsProviderFactory;
94+
TWriteSessionFactory WriteSessionFactory;
9495

9596
public:
9697
TService(
9798
TLogbrokerConfigPtr config,
9899
ILoggingServicePtr logging,
99100
std::shared_ptr<NYdbICredentialsProviderFactory>
100-
credentialsProviderFactory)
101+
credentialsProviderFactory,
102+
TWriteSessionFactory writeSessionFactory)
101103
: Config(std::move(config))
102104
, Logging(std::move(logging))
103105
, CredentialsProviderFactory(std::move(credentialsProviderFactory))
106+
, WriteSessionFactory(std::move(writeSessionFactory))
104107
{}
105108

106109
TFuture<NProto::TError> Write(
@@ -119,15 +122,21 @@ class TService final
119122
Batch = std::make_unique<TBatch>(std::move(messages), now);
120123

121124
if (!Session) {
122-
TTopicClient client{GetDriver()};
123-
124-
Session = client.CreateWriteSession(
125-
TWriteSessionSettings()
126-
.Path(Config->GetTopic())
127-
.ProducerId(Config->GetSourceId())
128-
.MessageGroupId(Config->GetSourceId())
129-
.RetryPolicy(
130-
NYdb::NTopic::IRetryPolicy::GetNoRetryPolicy()));
125+
if (WriteSessionFactory) {
126+
Session = WriteSessionFactory();
127+
} else {
128+
TTopicClient client{GetDriver()};
129+
130+
Session = client.CreateWriteSession(
131+
TWriteSessionSettings()
132+
.Path(Config->GetTopic())
133+
.ProducerId(Config->GetSourceId())
134+
.MessageGroupId(Config->GetSourceId())
135+
.RetryPolicy(
136+
NYdb::NTopic::IRetryPolicy::GetNoRetryPolicy()));
137+
}
138+
139+
Y_ABORT_UNLESS(Session);
131140
} else if (ContinuationToken.has_value()) {
132141
TContinuationToken token{std::move(ContinuationToken.value())};
133142
ContinuationToken.reset();
@@ -314,13 +323,26 @@ class TService final
314323
IServicePtr CreateTopicAPIService(
315324
TLogbrokerConfigPtr config,
316325
ILoggingServicePtr logging,
317-
std::shared_ptr<NYdbICredentialsProviderFactory>
318-
credentialsProviderFactory)
326+
std::shared_ptr<NYdbICredentialsProviderFactory> credentialsProviderFactory,
327+
TWriteSessionFactory writeSessionFactory)
319328
{
320329
return std::make_shared<TService>(
321330
std::move(config),
322331
std::move(logging),
323-
std::move(credentialsProviderFactory));
332+
std::move(credentialsProviderFactory),
333+
std::move(writeSessionFactory));
334+
}
335+
336+
IServicePtr CreateTopicAPIService(
337+
TLogbrokerConfigPtr config,
338+
ILoggingServicePtr logging,
339+
std::shared_ptr<NYdbICredentialsProviderFactory> credentialsProviderFactory)
340+
{
341+
return CreateTopicAPIService(
342+
std::move(config),
343+
std::move(logging),
344+
std::move(credentialsProviderFactory),
345+
{});
324346
}
325347

326348
IServicePtr CreateTopicAPIService(

cloud/blockstore/libs/logbroker/topic_api_impl/topic_api.h

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,10 +5,25 @@
55

66
#include <cloud/storage/core/libs/diagnostics/public.h>
77

8+
#include <functional>
9+
10+
namespace NYdb::inline V3::NTopic {
11+
class IWriteSession;
12+
} // namespace NYdb::inline V3::NTopic
13+
814
namespace NCloud::NBlockStore::NLogbroker {
915

1016
////////////////////////////////////////////////////////////////////////////////
1117

18+
using TWriteSessionFactory =
19+
std::function<std::shared_ptr<NYdb::NTopic::IWriteSession>()>;
20+
21+
IServicePtr CreateTopicAPIService(
22+
TLogbrokerConfigPtr config,
23+
ILoggingServicePtr logging,
24+
std::shared_ptr<NYdbICredentialsProviderFactory> credentialsProviderFactory,
25+
TWriteSessionFactory writeSessionFactory);
26+
1227
IServicePtr CreateTopicAPIService(
1328
TLogbrokerConfigPtr config,
1429
ILoggingServicePtr logging,

0 commit comments

Comments
 (0)