Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion ydb/functional_tests/basic/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,8 @@ project(userver-ydb-tests-basic CXX)
add_executable(
${PROJECT_NAME}
"ydb_service.cpp" "views/describe-table/post/view.cpp" "views/select-list/post/view.cpp"
"views/select-rows/post/view.cpp" "views/upsert-row/post/view.cpp" "views/upsert-row-old/post/view.cpp"
"views/select-rows/post/view.cpp" "views/insert-row/post/view.cpp" "views/upsert-row/post/view.cpp"
"views/upsert-row-old/post/view.cpp"
)
target_link_libraries(${PROJECT_NAME} userver::ydb)
target_include_directories(${PROJECT_NAME} PRIVATE ${CMAKE_CURRENT_SOURCE_DIR})
Expand Down
6 changes: 6 additions & 0 deletions ydb/functional_tests/basic/static_config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,8 @@ components_manager:
max_pool_size: 10
min_pool_size: 5
get_session_retry_limit: 1
grpc-compression-algorithm: gzip
grpc-load-balancing-policy: round_robin
# Second logical database used by the YDB_DATABASE_ROUTING test.
# The testsuite ydb recipe points every database to the same
# physical instance, so routing here is observed via per-database
Expand All @@ -100,6 +102,10 @@ components_manager:
method: POST
path: /ydb/upsert-row
task_processor: main-task-processor
handler-insert-row:
method: POST
path: /ydb/insert-row
task_processor: main-task-processor
handler-upsert-row-old:
method: POST
path: /ydb/upsert-row-old
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,8 @@ ydb.by-transaction.success: ydb_database=sampledb, ydb_transaction=trx RATE 3
ydb.by-transaction.timings: ydb_database=sampledb, ydb_transaction=trx HIST_RATE [5]=3,[10]=0,[20]=0,[35]=0,[60]=0,[100]=0,[173]=0,[300]=0,[520]=0,[1000]=0,[3200]=0,[10000]=0,[32000]=0,[100000]=0,[inf]=0
ydb.by-transaction.total: ydb_database=sampledb, ydb_transaction=trx RATE 3
ydb.by-transaction.transport-error: ydb_database=sampledb, ydb_transaction=trx RATE 0
ydb.grpc-compression-algorithm: algorithm=gzip, ydb_database=sampledb GAUGE 1
ydb.grpc-load-balancing-policy: policy=round_robin, ydb_database=sampledb GAUGE 1
ydb.native.Discovery/FailedTransportError: database=/local, ydb_database=sampledb RATE 0
ydb.native.Discovery/Regular: database=/local, ydb_database=sampledb RATE 0
ydb.native.Discovery/TooManyBadEndpoints: database=/local, ydb_database=sampledb RATE 0
Expand Down
22 changes: 22 additions & 0 deletions ydb/functional_tests/basic/tests/test_compression.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
async def test_grpc_compression_algorithm_metric(service_client, monitor_client):
response = await service_client.post(
'ydb/upsert-row',
json={
'id': 'id-compression',
'name': 'name-compression',
'service': 'srv',
'channel': 123,
},
)
assert response.status_code == 200
assert response.json() == {}

metrics = await monitor_client.metrics(prefix='ydb')
assert (
metrics.value_at(
path='ydb.grpc-compression-algorithm',
labels={'ydb_database': 'sampledb', 'algorithm': 'gzip'},
default=0,
)
== 1
)
35 changes: 35 additions & 0 deletions ydb/functional_tests/basic/tests/test_constraint_violation.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
async def _insert_row(service_client, row):
return await service_client.post('ydb/insert-row', json=row)


async def test_insert_row(service_client, ydb):
response = await _insert_row(
service_client,
{
'id': 'id-insert',
'name': 'name-insert',
'service': 'srv',
'channel': 123,
},
)
assert response.status_code == 200
assert response.json() == {}

cursor = ydb.execute('SELECT * FROM events WHERE id = "id-insert"')
assert len(cursor) == 1
assert len(cursor[0].rows) == 1


async def test_insert_row_duplicate_pk_conflict(service_client):
row = {
'id': 'id-insert-duplicate',
'name': 'name-insert-duplicate',
'service': 'srv',
'channel': 123,
}

response = await _insert_row(service_client, row)
assert response.status_code == 200

response = await _insert_row(service_client, row)
assert response.status_code == 409
20 changes: 20 additions & 0 deletions ydb/functional_tests/basic/tests/test_load_balancing.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
async def test_grpc_load_balancing_policy_metric(service_client, monitor_client):
response = await service_client.post(
'ydb/select-rows',
json={
'service': 'srv',
'channels': [1, 2, 3],
'created': '2019-10-30T11:20:00+00:00',
},
)
assert response.status_code == 200

metrics = await monitor_client.metrics(prefix='ydb')
assert (
metrics.value_at(
path='ydb.grpc-load-balancing-policy',
labels={'ydb_database': 'sampledb', 'policy': 'round_robin'},
default=0,
)
== 1
)
66 changes: 66 additions & 0 deletions ydb/functional_tests/basic/views/insert-row/post/view.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
#include "view.hpp"

#include <userver/formats/json.hpp>
#include <userver/server/handlers/exceptions.hpp>
#include <userver/utest/using_namespace_userver.hpp>

#include <userver/ydb/exceptions.hpp>
#include <userver/ydb/table.hpp>

namespace {

const ydb::Query kInsertQuery{
R"(
--!syntax_v1
DECLARE $id_key AS String;
DECLARE $name_key AS Utf8;
DECLARE $service_key AS String;
DECLARE $channel_key AS Int64;
DECLARE $state_key AS Json?;

INSERT INTO events (id, name, service, channel, created, state)
VALUES ($id_key, $name_key, $service_key, $channel_key, CurrentUtcTimestamp(), $state_key);
)",
ydb::Query::Name{"insert-row"},
};

} // namespace

namespace sample {

formats::json::Value
InsertRowHandler::HandleRequestJsonThrow(const server::http::HttpRequest&, const formats::json::Value& request, server::request::RequestContext&)
const {
try {
Ydb().RetryTx("trx", {.tx_mode = ydb::TransactionMode::kSerializableRW}, [&](ydb::TxActor& tx) {
auto response = tx.Execute(
kInsertQuery, //
"$id_key",
request["id"].As<std::string>(), //
"$name_key",
ydb::Utf8{request["name"].As<std::string>()}, //
"$service_key",
request["service"].As<std::string>(), //
"$channel_key",
request["channel"].As<int64_t>(), //
"$state_key",
request["state"].As<std::optional<formats::json::Value>>() //
);

if (response.GetCursorCount()) {
throw std::runtime_error("Unexpected response data");
}

return ydb::TxAction::kCommit;
});
} catch (const ydb::YdbResponseError& ex) {
if (ex.IsConstraintViolation()) {
throw server::handlers::CustomHandlerException(server::handlers::HandlerErrorCode::kConflictState);
}
throw;
}

return formats::json::MakeObject();
}

} // namespace sample
20 changes: 20 additions & 0 deletions ydb/functional_tests/basic/views/insert-row/post/view.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
#pragma once

#include <views/base_handler.hpp>

namespace sample {

class InsertRowHandler final : public BaseHandler {
public:
static constexpr std::string_view kName = "handler-insert-row";

using BaseHandler::BaseHandler;

formats::json::Value HandleRequestJsonThrow(
const server::http::HttpRequest& request,
const formats::json::Value& request_json,
server::request::RequestContext& context
) const override;
};

} // namespace sample
2 changes: 2 additions & 0 deletions ydb/functional_tests/basic/ydb_service.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
#include <userver/ydb/dist_lock/component_base.hpp>

#include <views/describe-table/post/view.hpp>
#include <views/insert-row/post/view.hpp>
#include <views/select-list/post/view.hpp>
#include <views/select-rows/post/view.hpp>
#include <views/upsert-row-old/post/view.hpp>
Expand Down Expand Up @@ -100,6 +101,7 @@ int main(int argc, char* argv[]) {
.Append<sample::SelectRowsHandler>()
.Append<sample::DescribeTableHandler>()
.Append<sample::SelectListHandler>()
.Append<sample::InsertRowHandler>()
.Append<sample::UpsertRowHandler>()
.Append<sample::UpsertRowOldHandler>()
.Append<ydb::YdbComponent>()
Expand Down
2 changes: 2 additions & 0 deletions ydb/include/userver/ydb/exceptions.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@ class YdbResponseError : public BaseError {

const NYdb::TStatus& GetStatus() const noexcept;

bool IsConstraintViolation() const noexcept;

private:
NYdb::TStatus status_;
};
Expand Down
14 changes: 14 additions & 0 deletions ydb/src/ydb/component.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,20 @@ properties:
type: boolean
description: |
NYdb::TDriverConfig::SetGRpcKeepAlivePermitWithoutCalls. If omitted, not set (SDK default).
grpc-compression-algorithm:
type: string
enum: [none, gzip, deflate]
default: none
description: |
NYdb::TDriverConfig::SetGRpcCompressionAlgorithm. If omitted, not set (SDK default).
Valid values: "gzip", "deflate", "none".
grpc-load-balancing-policy:
type: string
enum: [round_robin, pick_first]
default: round_robin
description: |
NYdb::TDriverConfig::SetGRpcLoadBalancingPolicy. If omitted, not set (SDK default).
Valid values: "round_robin", "pick_first"
prefer_local_dc:
type: boolean
default: true
Expand Down
4 changes: 4 additions & 0 deletions ydb/src/ydb/exceptions.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,10 @@ YdbResponseError::YdbResponseError(std::string_view operation_name, NYdb::TStatu

const NYdb::TStatus& YdbResponseError::GetStatus() const noexcept { return status_; }

bool YdbResponseError::IsConstraintViolation() const noexcept {
return NYdb::NStatusHelpers::StatusContainsIssueWithCode(status_, NYdb::NIssue::CONSTRAINT_VIOLATION);
}

UndefinedDatabaseError::UndefinedDatabaseError(std::string error)
: BaseError(std::move(error))
{}
Expand Down
40 changes: 40 additions & 0 deletions ydb/src/ydb/impl/config.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
#include <userver/formats/json/serialize.hpp>
#include <userver/formats/json/value.hpp>
#include <userver/formats/parse/common_containers.hpp>
#include <userver/utils/assert.hpp>
#include <userver/utils/retry_budget.hpp>
#include <userver/yaml_config/yaml_config.hpp>

Expand Down Expand Up @@ -32,6 +33,25 @@ T MergeWithSecdist(
);
}

NYdb::EGrpcCompressionAlgorithm ToCompressionAlgorithm(std::string_view alg) {
if (alg == "gzip") {
return NYdb::EGrpcCompressionAlgorithm::Gzip;
}
if (alg == "deflate") {
return NYdb::EGrpcCompressionAlgorithm::Deflate;
}
if (alg == "none") {
return NYdb::EGrpcCompressionAlgorithm::None;
}
throw yaml_config::Exception(fmt::format("Unknown grpc-compression-algorithm: {}", alg));
}

void ValidateLoadBalancingPolicy(std::string_view policy) {
if (policy != "round_robin" && policy != "pick_first") {
throw yaml_config::Exception(fmt::format("Unknown grpc-load-balancing-policy: {}", policy));
}
}

} // namespace

TableSettings ParseTableSettings(const yaml_config::YamlConfig& dbconfig, const secdist::DatabaseSettings& dbsecdist) {
Expand All @@ -58,6 +78,18 @@ TableSettings ParseTableSettings(const yaml_config::YamlConfig& dbconfig, const
return result;
}

std::string_view ToString(NYdb::EGrpcCompressionAlgorithm algorithm) {
switch (algorithm) {
case NYdb::EGrpcCompressionAlgorithm::None:
return "none";
case NYdb::EGrpcCompressionAlgorithm::Deflate:
return "deflate";
case NYdb::EGrpcCompressionAlgorithm::Gzip:
return "gzip";
}
UINVARIANT(false, "Unhandled EGrpcCompressionAlgorithm value");
}

DriverSettings ParseDriverSettings(
const yaml_config::YamlConfig& dbconfig,
const secdist::DatabaseSettings& dbsecdist,
Expand Down Expand Up @@ -85,6 +117,14 @@ DriverSettings ParseDriverSettings(
result.grpc_keepalive_permit_without_calls =
dbconfig["grpc-keepalive-permit-without-calls"].As<std::optional<bool>>();

if (auto alg = dbconfig["grpc-compression-algorithm"].As<std::optional<std::string>>()) {
result.grpc_compression_algorithm = ToCompressionAlgorithm(*alg);
}
result.grpc_load_balancing_policy = dbconfig["grpc-load-balancing-policy"].As<std::optional<std::string>>();
if (result.grpc_load_balancing_policy.has_value()) {
ValidateLoadBalancingPolicy(*result.grpc_load_balancing_policy);
}

result.endpoint = MergeWithSecdist(dbsecdist.endpoint, std::move(config_endpoint), dbconfig, "endpoint");
result.database = MergeWithSecdist(dbsecdist.database, std::move(config_database), dbconfig, "database");
result.oauth_token = dbsecdist.oauth_token;
Expand Down
6 changes: 6 additions & 0 deletions ydb/src/ydb/impl/config.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
#include <userver/yaml_config/fwd.hpp>

#include <userver/ydb/settings.hpp>
#include <ydb-cpp-sdk/client/types/ydb.h>

USERVER_NAMESPACE_BEGIN

Expand Down Expand Up @@ -59,6 +60,9 @@ struct DriverSettings {
std::optional<std::chrono::milliseconds> grpc_keepalive_timeout{};
std::optional<bool> grpc_keepalive_permit_without_calls{};

std::optional<NYdb::EGrpcCompressionAlgorithm> grpc_compression_algorithm{};
std::optional<std::string> grpc_load_balancing_policy{};

bool prefer_local_dc{false};
std::optional<std::string> oauth_token;
std::optional<std::string> iam_jwt_params;
Expand All @@ -68,6 +72,8 @@ struct DriverSettings {
std::shared_ptr<NYdb::ICredentialsProviderFactory> credentials_provider_factory;
};

std::string_view ToString(NYdb::EGrpcCompressionAlgorithm algorithm);

TableSettings ParseTableSettings(const yaml_config::YamlConfig& dbconfig, const secdist::DatabaseSettings& dbsecdist);

DriverSettings ParseDriverSettings(
Expand Down
Loading
Loading