diff --git a/ydb/functional_tests/basic/CMakeLists.txt b/ydb/functional_tests/basic/CMakeLists.txt index 4d0c33dca057..acd8203d4b1e 100644 --- a/ydb/functional_tests/basic/CMakeLists.txt +++ b/ydb/functional_tests/basic/CMakeLists.txt @@ -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}) diff --git a/ydb/functional_tests/basic/static_config.yaml b/ydb/functional_tests/basic/static_config.yaml index 1f0e074b0c3d..d9fa21fb9946 100644 --- a/ydb/functional_tests/basic/static_config.yaml +++ b/ydb/functional_tests/basic/static_config.yaml @@ -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 @@ -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 diff --git a/ydb/functional_tests/basic/tests-metrics/static/metrics_values.txt b/ydb/functional_tests/basic/tests-metrics/static/metrics_values.txt index 5427330f4fdb..577c91dc4ab9 100644 --- a/ydb/functional_tests/basic/tests-metrics/static/metrics_values.txt +++ b/ydb/functional_tests/basic/tests-metrics/static/metrics_values.txt @@ -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 diff --git a/ydb/functional_tests/basic/tests/test_compression.py b/ydb/functional_tests/basic/tests/test_compression.py new file mode 100644 index 000000000000..c39dac4d8e5c --- /dev/null +++ b/ydb/functional_tests/basic/tests/test_compression.py @@ -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 + ) diff --git a/ydb/functional_tests/basic/tests/test_constraint_violation.py b/ydb/functional_tests/basic/tests/test_constraint_violation.py new file mode 100644 index 000000000000..2519b5ff9041 --- /dev/null +++ b/ydb/functional_tests/basic/tests/test_constraint_violation.py @@ -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 diff --git a/ydb/functional_tests/basic/tests/test_load_balancing.py b/ydb/functional_tests/basic/tests/test_load_balancing.py new file mode 100644 index 000000000000..e2697e6e251f --- /dev/null +++ b/ydb/functional_tests/basic/tests/test_load_balancing.py @@ -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 + ) diff --git a/ydb/functional_tests/basic/views/insert-row/post/view.cpp b/ydb/functional_tests/basic/views/insert-row/post/view.cpp new file mode 100644 index 000000000000..84e60121b6e2 --- /dev/null +++ b/ydb/functional_tests/basic/views/insert-row/post/view.cpp @@ -0,0 +1,66 @@ +#include "view.hpp" + +#include +#include +#include + +#include +#include + +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(), // + "$name_key", + ydb::Utf8{request["name"].As()}, // + "$service_key", + request["service"].As(), // + "$channel_key", + request["channel"].As(), // + "$state_key", + request["state"].As>() // + ); + + 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 diff --git a/ydb/functional_tests/basic/views/insert-row/post/view.hpp b/ydb/functional_tests/basic/views/insert-row/post/view.hpp new file mode 100644 index 000000000000..f12f6060f18a --- /dev/null +++ b/ydb/functional_tests/basic/views/insert-row/post/view.hpp @@ -0,0 +1,20 @@ +#pragma once + +#include + +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 diff --git a/ydb/functional_tests/basic/ydb_service.cpp b/ydb/functional_tests/basic/ydb_service.cpp index 0fa5c078d89a..4a54a7a697d6 100644 --- a/ydb/functional_tests/basic/ydb_service.cpp +++ b/ydb/functional_tests/basic/ydb_service.cpp @@ -18,6 +18,7 @@ #include #include +#include #include #include #include @@ -100,6 +101,7 @@ int main(int argc, char* argv[]) { .Append() .Append() .Append() + .Append() .Append() .Append() .Append() diff --git a/ydb/include/userver/ydb/exceptions.hpp b/ydb/include/userver/ydb/exceptions.hpp index 92bd46e8dd7a..4e9c5659ea8c 100644 --- a/ydb/include/userver/ydb/exceptions.hpp +++ b/ydb/include/userver/ydb/exceptions.hpp @@ -23,6 +23,8 @@ class YdbResponseError : public BaseError { const NYdb::TStatus& GetStatus() const noexcept; + bool IsConstraintViolation() const noexcept; + private: NYdb::TStatus status_; }; diff --git a/ydb/src/ydb/component.yaml b/ydb/src/ydb/component.yaml index c92dbf563ad5..1194f511cc65 100644 --- a/ydb/src/ydb/component.yaml +++ b/ydb/src/ydb/component.yaml @@ -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 diff --git a/ydb/src/ydb/exceptions.cpp b/ydb/src/ydb/exceptions.cpp index 4494616ddace..48879577c9b1 100644 --- a/ydb/src/ydb/exceptions.cpp +++ b/ydb/src/ydb/exceptions.cpp @@ -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)) {} diff --git a/ydb/src/ydb/impl/config.cpp b/ydb/src/ydb/impl/config.cpp index b9065e8205b3..0b276cf0e8f1 100644 --- a/ydb/src/ydb/impl/config.cpp +++ b/ydb/src/ydb/impl/config.cpp @@ -3,6 +3,7 @@ #include #include #include +#include #include #include @@ -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) { @@ -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, @@ -85,6 +117,14 @@ DriverSettings ParseDriverSettings( result.grpc_keepalive_permit_without_calls = dbconfig["grpc-keepalive-permit-without-calls"].As>(); + if (auto alg = dbconfig["grpc-compression-algorithm"].As>()) { + result.grpc_compression_algorithm = ToCompressionAlgorithm(*alg); + } + result.grpc_load_balancing_policy = dbconfig["grpc-load-balancing-policy"].As>(); + 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; diff --git a/ydb/src/ydb/impl/config.hpp b/ydb/src/ydb/impl/config.hpp index dc8285da8b2e..21bc6fcb306f 100644 --- a/ydb/src/ydb/impl/config.hpp +++ b/ydb/src/ydb/impl/config.hpp @@ -15,6 +15,7 @@ #include #include +#include USERVER_NAMESPACE_BEGIN @@ -59,6 +60,9 @@ struct DriverSettings { std::optional grpc_keepalive_timeout{}; std::optional grpc_keepalive_permit_without_calls{}; + std::optional grpc_compression_algorithm{}; + std::optional grpc_load_balancing_policy{}; + bool prefer_local_dc{false}; std::optional oauth_token; std::optional iam_jwt_params; @@ -68,6 +72,8 @@ struct DriverSettings { std::shared_ptr credentials_provider_factory; }; +std::string_view ToString(NYdb::EGrpcCompressionAlgorithm algorithm); + TableSettings ParseTableSettings(const yaml_config::YamlConfig& dbconfig, const secdist::DatabaseSettings& dbsecdist); DriverSettings ParseDriverSettings( diff --git a/ydb/src/ydb/impl/driver.cpp b/ydb/src/ydb/impl/driver.cpp index f01920e3f342..bad9e727f052 100644 --- a/ydb/src/ydb/impl/driver.cpp +++ b/ydb/src/ydb/impl/driver.cpp @@ -18,6 +18,8 @@ namespace ydb::impl { Driver::Driver(std::string dbname, impl::DriverSettings settings) : dbname_(std::move(dbname)), dbpath_(settings.database), + grpc_compression_algorithm_(settings.grpc_compression_algorithm), + grpc_load_balancing_policy_(settings.grpc_load_balancing_policy), native_metrics_(std::make_unique(NMonitoring::TLabels{})), retry_budget_(utils::RetryBudgetSettings{}) { @@ -70,6 +72,13 @@ Driver::Driver(std::string dbname, impl::DriverSettings settings) driver_config.SetGRpcKeepAlivePermitWithoutCalls(*settings.grpc_keepalive_permit_without_calls); } + if (settings.grpc_compression_algorithm.has_value()) { + driver_config.SetGRpcCompressionAlgorithm(*settings.grpc_compression_algorithm); + } + if (settings.grpc_load_balancing_policy.has_value()) { + driver_config.SetGRpcLoadBalancingPolicy(*settings.grpc_load_balancing_policy); + } + AppendUserverYdbBuildInfo(driver_config); driver_ = std::make_unique(driver_config); @@ -86,7 +95,17 @@ const std::string& Driver::GetDbPath() const { return dbpath_; } utils::RetryBudget& Driver::GetRetryBudget() { return retry_budget_; } -void DumpMetric(utils::statistics::Writer& writer, const Driver& driver) { writer["native"] = *driver.native_metrics_; } +void DumpMetric(utils::statistics::Writer& writer, const Driver& driver) { + writer["native"] = *driver.native_metrics_; + if (driver.grpc_compression_algorithm_.has_value()) { + writer["grpc-compression-algorithm"].ValueWithLabels( + 1, {{"algorithm", ToString(*driver.grpc_compression_algorithm_)}} + ); + } + if (driver.grpc_load_balancing_policy_.has_value()) { + writer["grpc-load-balancing-policy"].ValueWithLabels(1, {{"policy", *driver.grpc_load_balancing_policy_}}); + } +} std::string JoinPath(std::string_view database_path, std::string_view path) { UASSERT(!database_path.ends_with("/")); diff --git a/ydb/src/ydb/impl/driver.hpp b/ydb/src/ydb/impl/driver.hpp index cd4eff0276da..ccda07f4673c 100644 --- a/ydb/src/ydb/impl/driver.hpp +++ b/ydb/src/ydb/impl/driver.hpp @@ -1,12 +1,14 @@ #pragma once #include +#include #include #include #include #include +#include namespace NMonitoring { class TMetricRegistry; @@ -44,6 +46,9 @@ class Driver final { const std::string dbname_; const std::string dbpath_; + const std::optional grpc_compression_algorithm_; + const std::optional grpc_load_balancing_policy_; + std::unique_ptr native_metrics_; // The retry_budget_ is used in driver_ threads, so it must be before the // driver_ diff --git a/ydb/tests/driver_config_test.cpp b/ydb/tests/driver_config_test.cpp new file mode 100644 index 000000000000..e0b3362d0807 --- /dev/null +++ b/ydb/tests/driver_config_test.cpp @@ -0,0 +1,66 @@ +#include +#include +#include + +#include +#include + +USERVER_NAMESPACE_BEGIN + +namespace { + +ydb::impl::DriverSettings DriverSettingsFromYaml(std::string_view yaml) { + const yaml_config::YamlConfig config{formats::yaml::FromString(std::string{yaml}), {}}; + ydb::impl::secdist::DatabaseSettings secdist; + secdist.endpoint = "localhost:2136"; + secdist.database = "local"; + return ydb::impl::ParseDriverSettings(config, secdist, nullptr); +} + +} // namespace + +UTEST(YdbDriverConfig, GrpcCompressionAlgorithmGzip) { + const auto settings = DriverSettingsFromYaml("grpc-compression-algorithm: gzip"); + ASSERT_TRUE(settings.grpc_compression_algorithm.has_value()); + EXPECT_EQ(*settings.grpc_compression_algorithm, NYdb::EGrpcCompressionAlgorithm::Gzip); +} + +UTEST(YdbDriverConfig, GrpcCompressionAlgorithmDeflate) { + const auto settings = DriverSettingsFromYaml("grpc-compression-algorithm: deflate"); + ASSERT_TRUE(settings.grpc_compression_algorithm.has_value()); + EXPECT_EQ(*settings.grpc_compression_algorithm, NYdb::EGrpcCompressionAlgorithm::Deflate); +} + +UTEST(YdbDriverConfig, GrpcCompressionAlgorithmNone) { + const auto settings = DriverSettingsFromYaml("grpc-compression-algorithm: none"); + ASSERT_TRUE(settings.grpc_compression_algorithm.has_value()); + EXPECT_EQ(*settings.grpc_compression_algorithm, NYdb::EGrpcCompressionAlgorithm::None); +} + +UTEST(YdbDriverConfig, GrpcCompressionAlgorithmUnknownThrows) { + EXPECT_THROW(DriverSettingsFromYaml("grpc-compression-algorithm: brotli"), yaml_config::Exception); +} + +UTEST(YdbDriverConfig, GrpcLoadBalancingPolicyRoundRobin) { + const auto settings = DriverSettingsFromYaml("grpc-load-balancing-policy: round_robin"); + ASSERT_TRUE(settings.grpc_load_balancing_policy.has_value()); + EXPECT_EQ(*settings.grpc_load_balancing_policy, "round_robin"); +} + +UTEST(YdbDriverConfig, GrpcLoadBalancingPolicyPickFirst) { + const auto settings = DriverSettingsFromYaml("grpc-load-balancing-policy: pick_first"); + ASSERT_TRUE(settings.grpc_load_balancing_policy.has_value()); + EXPECT_EQ(*settings.grpc_load_balancing_policy, "pick_first"); +} + +UTEST(YdbDriverConfig, GrpcLoadBalancingPolicyUnknownThrows) { + EXPECT_THROW(DriverSettingsFromYaml("grpc-load-balancing-policy: random"), yaml_config::Exception); +} + +UTEST(YdbDriverConfig, MissingKeysLeaveSettingsUnset) { + const auto settings = DriverSettingsFromYaml("max_pool_size: 10"); + EXPECT_FALSE(settings.grpc_compression_algorithm.has_value()); + EXPECT_FALSE(settings.grpc_load_balancing_policy.has_value()); +} + +USERVER_NAMESPACE_END