diff options
| author | achains <[email protected]> | 2025-09-03 11:49:59 +0300 |
|---|---|---|
| committer | achains <[email protected]> | 2025-09-03 12:09:24 +0300 |
| commit | 6165c60a95c3c621ca21112b432df698412abf86 (patch) | |
| tree | 9233f38995755027bf7c5bac5a387c6083d18541 /yt/cpp/mapreduce/rpc_client/client_impl.cpp | |
| parent | 17e0c8876359d20baa5bec0d9e7860b83dd4472b (diff) | |
YT-23616: RPC proxy tests for C++ client
* Changelog entry
Type: fix
Component: cpp-sdk
Enable misc tests for RPC proxy, various fixes
commit_hash:c3c716503a2e106731ad99b66ec57ea00baf0304
Diffstat (limited to 'yt/cpp/mapreduce/rpc_client/client_impl.cpp')
| -rw-r--r-- | yt/cpp/mapreduce/rpc_client/client_impl.cpp | 95 |
1 files changed, 95 insertions, 0 deletions
diff --git a/yt/cpp/mapreduce/rpc_client/client_impl.cpp b/yt/cpp/mapreduce/rpc_client/client_impl.cpp new file mode 100644 index 00000000000..01e702231f6 --- /dev/null +++ b/yt/cpp/mapreduce/rpc_client/client_impl.cpp @@ -0,0 +1,95 @@ +#include "raw_client.h" + +#include <yt/cpp/mapreduce/client/client.h> +#include <yt/cpp/mapreduce/client/init.h> + +#include <yt/cpp/mapreduce/common/retry_lib.h> + +#include <yt/cpp/mapreduce/interface/client_method_options.h> +#include <yt/cpp/mapreduce/interface/tvm.h> + +#include <yt/yt/client/api/rpc_proxy/config.h> +#include <yt/yt/client/api/rpc_proxy/connection.h> + +namespace NYT { + +//////////////////////////////////////////////////////////////////////////////// + +namespace NDetail { + +//////////////////////////////////////////////////////////////////////////////// + +NYT::NApi::IClientPtr CreateApiClient(const TClientContext& context) +{ + auto connectionConfig = New<NApi::NRpcProxy::TConnectionConfig>(); + connectionConfig->SetDefaults(); + if (context.JobProxySocketPath) { + connectionConfig->ProxyUnixDomainSocket = *context.JobProxySocketPath; + } else { + connectionConfig->ClusterUrl = context.ServerName; + } + if (context.RpcProxyRole) { + connectionConfig->ProxyRole = *context.RpcProxyRole; + } + if (context.ProxyAddress) { + connectionConfig->ProxyAddresses = {*context.ProxyAddress}; + } + + THashMap<std::string, std::string> proxyUrlAliasingRules; + for (const auto& [clusterName, url] : context.Config->ProxyUrlAliasingRules) { + proxyUrlAliasingRules.emplace(clusterName, url); + } + + connectionConfig->ProxyUrlAliasingRules = std::move(proxyUrlAliasingRules); + + NApi::TClientOptions clientOptions; + clientOptions.Token = context.Token; + if (context.ServiceTicketAuth) { + clientOptions.ServiceTicketAuth = context.ServiceTicketAuth->Ptr; + } + if (context.ImpersonationUser) { + clientOptions.User = *context.ImpersonationUser; + } + if (context.JobProxySocketPath) { + clientOptions.MultiproxyTargetCluster = context.MultiproxyTargetCluster; + } + + auto connection = NApi::NRpcProxy::CreateConnection(connectionConfig); + return connection->CreateClient(clientOptions); +} + +//////////////////////////////////////////////////////////////////////////////// + +} // namespace NDetail + +//////////////////////////////////////////////////////////////////////////////// + +IClientPtr CreateRpcClient( + const TString& serverName, + const TCreateClientOptions& options) +{ + auto context = NDetail::CreateClientContext(serverName, options); + + auto globalTxId = GetGuid(context.Config->GlobalTxId); + + auto retryConfigProvider = options.RetryConfigProvider_; + if (!retryConfigProvider) { + retryConfigProvider = CreateDefaultRetryConfigProvider(); + } + + auto rawClient = MakeIntrusive<NDetail::TRpcRawClient>( + NDetail::CreateApiClient(context), + context.Config); + + NDetail::EnsureInitialized(); + + return new NDetail::TClient( + std::move(rawClient), + context, + globalTxId, + CreateDefaultClientRetryPolicy(retryConfigProvider, context.Config)); +} + +//////////////////////////////////////////////////////////////////////////////// + +} // namespace NYT |
