diff options
| author | Nikita Vasilev <[email protected]> | 2026-07-22 11:56:22 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2026-07-22 11:56:22 +0300 |
| commit | 75900fddb7cfff27d66271c4a9c419ddd852554b (patch) | |
| tree | a98d115c3e1e5ae4ab8b9083f91e3f639ad4a055 | |
| parent | d8cb30f4bd872ed7d6e7f845053a24ea60298c6f (diff) | |
StrictSerializable: return commit timestamp (#46795)
27 files changed, 610 insertions, 80 deletions
diff --git a/ydb/core/grpc_services/query/rpc_execute_query.cpp b/ydb/core/grpc_services/query/rpc_execute_query.cpp index 406b2a6154c..ac448ec5135 100644 --- a/ydb/core/grpc_services/query/rpc_execute_query.cpp +++ b/ydb/core/grpc_services/query/rpc_execute_query.cpp @@ -473,6 +473,12 @@ private: hasTrailingMessage = true; response.mutable_tx_meta()->set_id(kqpResponse.GetTxMeta().id()); } + + if (kqpResponse.HasCommitTimestamp()) { + hasTrailingMessage = true; + response.mutable_commit_timestamp()->set_plan_step(kqpResponse.GetCommitTimestamp().plan_step()); + response.mutable_commit_timestamp()->set_tx_id(kqpResponse.GetCommitTimestamp().tx_id()); + } } if (hasTrailingMessage) { diff --git a/ydb/core/grpc_services/query/rpc_kqp_tx.cpp b/ydb/core/grpc_services/query/rpc_kqp_tx.cpp index 2914aa83aaa..24f30b1154b 100644 --- a/ydb/core/grpc_services/query/rpc_kqp_tx.cpp +++ b/ydb/core/grpc_services/query/rpc_kqp_tx.cpp @@ -181,7 +181,7 @@ public: private: virtual std::pair<TString, TString> GetReqData() const = 0; virtual void Fill(NKikimrKqp::TQueryRequest* req) const = 0; - virtual NProtoBuf::Message* CreateResult(Ydb::StatusIds::StatusCode status, const NYql::TIssues& issues) const = 0; + virtual NProtoBuf::Message* CreateResult(Ydb::StatusIds::StatusCode status, const NYql::TIssues& issues, const NKikimrKqp::TQueryResponse& kqpResponse) const = 0; void StateWork(TAutoPtr<IEventHandle>& ev) { try { @@ -233,14 +233,15 @@ private: FillCommonKqpRespFields(record, Request.get()); NYql::TIssues issues; + const NKikimrKqp::TQueryResponse* kqpResponse = nullptr; if (record.HasResponse()) { - const auto& kqpResponse = record.GetResponse(); - const auto& issueMessage = kqpResponse.GetQueryIssues(); + kqpResponse = &record.GetResponse(); + const auto& issueMessage = kqpResponse->GetQueryIssues(); NYql::IssuesFromMessage(issueMessage, issues); Request->RaiseIssues(issues); } - Reply(record.GetYdbStatus(), CreateResult(record.GetYdbStatus(), issues)); + Reply(record.GetYdbStatus(), CreateResult(record.GetYdbStatus(), issues, kqpResponse ? *kqpResponse : NKikimrKqp::TQueryResponse::default_instance())); } void InternalError(const TString& message) { @@ -285,10 +286,14 @@ private: req->MutableTxControl()->set_commit_tx(true); } - NProtoBuf::Message* CreateResult(Ydb::StatusIds::StatusCode status, const NYql::TIssues& issues) const override { + NProtoBuf::Message* CreateResult(Ydb::StatusIds::StatusCode status, const NYql::TIssues& issues, const NKikimrKqp::TQueryResponse& kqpResponse) const override { auto result = TEvCommitTransactionRequest::AllocateResult<Ydb::Query::CommitTransactionResponse>(Request); result->set_status(status); NYql::IssuesToMessage(issues, result->mutable_issues()); + if (kqpResponse.HasCommitTimestamp()) { + result->mutable_commit_timestamp()->set_plan_step(kqpResponse.GetCommitTimestamp().plan_step()); + result->mutable_commit_timestamp()->set_tx_id(kqpResponse.GetCommitTimestamp().tx_id()); + } return result; } }; @@ -308,7 +313,8 @@ private: req->SetAction(NKikimrKqp::QUERY_ACTION_ROLLBACK_TX); } - NProtoBuf::Message* CreateResult(Ydb::StatusIds::StatusCode status, const NYql::TIssues& issues) const override { + NProtoBuf::Message* CreateResult(Ydb::StatusIds::StatusCode status, const NYql::TIssues& issues, const NKikimrKqp::TQueryResponse& kqpResponse) const override { + Y_UNUSED(kqpResponse); auto result = TEvRollbackTransactionRequest::AllocateResult<Ydb::Query::RollbackTransactionResponse>(Request); result->set_status(status); NYql::IssuesToMessage(issues, result->mutable_issues()); diff --git a/ydb/core/kqp/common/buffer/events.h b/ydb/core/kqp/common/buffer/events.h index 922f32ea981..a40580e8539 100644 --- a/ydb/core/kqp/common/buffer/events.h +++ b/ydb/core/kqp/common/buffer/events.h @@ -9,6 +9,11 @@ namespace NKikimr { namespace NKqp { +struct TCommitTimestamp { + ui64 PlanStep = 0; + ui64 TxId = 0; +}; + struct TEvKqpBuffer { // To BufferActor @@ -34,8 +39,12 @@ struct TEvTerminate : public TEventLocal<TEvTerminate, TKqpBufferWriterEvents::E struct TEvResult : public TEventLocal<TEvResult, TKqpBufferWriterEvents::EvResult> { TEvResult() = default; TEvResult(NYql::NDqProto::TDqTaskStats&& stats) : Stats(std::move(stats)) {} + TEvResult(NYql::NDqProto::TDqTaskStats&& stats, std::optional<TCommitTimestamp>&& commitTimestamp) + : Stats(std::move(stats)) + , CommitTimestamp(std::move(commitTimestamp)) {} std::optional<NYql::NDqProto::TDqTaskStats> Stats; + std::optional<TCommitTimestamp> CommitTimestamp; }; struct TEvError : public TEventLocal<TEvError, TKqpBufferWriterEvents::EvError> { diff --git a/ydb/core/kqp/executer_actor/kqp_data_executer.cpp b/ydb/core/kqp/executer_actor/kqp_data_executer.cpp index 9e81c83f4ac..bd2267a62a3 100644 --- a/ydb/core/kqp/executer_actor/kqp_data_executer.cpp +++ b/ydb/core/kqp/executer_actor/kqp_data_executer.cpp @@ -301,6 +301,7 @@ public: if (ev->Get()->Stats && Stats) { Stats->AddBufferStats(std::move(*ev->Get()->Stats)); } + ResponseEv->CommitTimestamp = std::move(ev->Get()->CommitTimestamp); MakeResponseAndPassAway(); } diff --git a/ydb/core/kqp/executer_actor/kqp_executer.h b/ydb/core/kqp/executer_actor/kqp_executer.h index 6450d9fb5cf..3ead2fe51ff 100644 --- a/ydb/core/kqp/executer_actor/kqp_executer.h +++ b/ydb/core/kqp/executer_actor/kqp_executer.h @@ -4,6 +4,7 @@ #include <ydb/core/kqp/common/kqp_batch_operations.h> #include <ydb/core/kqp/common/kqp_tx.h> #include <ydb/core/kqp/common/kqp_event_ids.h> +#include <ydb/core/kqp/common/buffer/events.h> #include <ydb/core/kqp/common/kqp_user_request_context.h> #include <ydb/core/kqp/executer_actor/kqp_partition_helper.h> #include <ydb/core/kqp/executer_actor/shards_resolver/kqp_shards_resolver_events.h> @@ -32,6 +33,7 @@ struct TEvKqpExecuter { NLWTrace::TOrbit Orbit; IKqpGateway::TKqpSnapshot Snapshot; + std::optional<TCommitTimestamp> CommitTimestamp; std::optional<NYql::TKikimrPathId> BrokenLockPathId; std::optional<ui64> BrokenLockShardId; std::optional<ui64> BrokenLockQuerySpanId; diff --git a/ydb/core/kqp/runtime/kqp_write_actor.cpp b/ydb/core/kqp/runtime/kqp_write_actor.cpp index 508bcc72d84..dd2e526a1f1 100644 --- a/ydb/core/kqp/runtime/kqp_write_actor.cpp +++ b/ydb/core/kqp/runtime/kqp_write_actor.cpp @@ -4792,6 +4792,12 @@ public: case TEvTxProxy::TEvProposeTransactionStatus::EStatus::StatusPlanned: TxProxyMon->ClientTxStatusPlanned->Inc(); TxPlanned = true; + if (TxManager->GetIsolationLevel() == NKqpProto::ISOLATION_LEVEL_STRICT_SERIALIZABLE) { + AFL_ENSURE(res->Record.HasStepId()); + AFL_ENSURE(res->Record.HasTxId()); + AFL_ENSURE(TxId && *TxId == res->Record.GetTxId()); + CommitTimestamp = TCommitTimestamp{res->Record.GetStepId(), res->Record.GetTxId()}; + } break; case TEvTxProxy::TEvProposeTransactionStatus::EStatus::StatusOutdated: @@ -5553,7 +5559,8 @@ public: CA_LOG_D("Committed TxId=" << TxId.value_or(0)); OnOperationFinished(Counters->BufferActorCommitLatencyHistogram); Send<ESendingType::Tail>(ExecuterActorId, new TEvKqpBuffer::TEvResult{ - BuildStats() + BuildStats(), + std::move(CommitTimestamp) }); ExecuterActorId = {}; AFL_ENSURE(GetTotalMemory() == 0); @@ -5945,6 +5952,7 @@ private: bool IsImmediateCommit = false; bool TxPlanned = false; std::optional<ui64> Coordinator; + std::optional<TCommitTimestamp> CommitTimestamp; ui64 LocksBrokenAsBreaker = 0; ui64 LocksBrokenAsVictim = 0; diff --git a/ydb/core/kqp/session_actor/kqp_query_state.h b/ydb/core/kqp/session_actor/kqp_query_state.h index 46ac709df1f..0c9ebb1c74d 100644 --- a/ydb/core/kqp/session_actor/kqp_query_state.h +++ b/ydb/core/kqp/session_actor/kqp_query_state.h @@ -12,6 +12,7 @@ #include <ydb/core/kqp/common/kqp_resolve.h> #include <ydb/core/kqp/common/kqp_timeouts.h> #include <ydb/core/kqp/common/kqp_tx.h> +#include <ydb/core/kqp/common/buffer/events.h> #include <ydb/core/kqp/common/kqp_user_request_context.h> #include <ydb/core/kqp/common/kqp.h> #include <ydb/core/kqp/common/simple/temp_tables.h> @@ -192,6 +193,8 @@ public: bool Commit = false; bool Commited = false; + std::optional<TCommitTimestamp> CommitTimestamp; + NTopic::TTopicOperations TopicOperations; TDuration CpuTime; std::optional<NCpuTime::TCpuTimer> CurrentTimer; diff --git a/ydb/core/kqp/session_actor/kqp_session_actor.cpp b/ydb/core/kqp/session_actor/kqp_session_actor.cpp index 1b13e6dea22..68d9515531e 100644 --- a/ydb/core/kqp/session_actor/kqp_session_actor.cpp +++ b/ydb/core/kqp/session_actor/kqp_session_actor.cpp @@ -2730,6 +2730,10 @@ public: QueryState->TxCtx->AcceptIncomingSnapshot(ev->Snapshot); + if (ev->CommitTimestamp) { + QueryState->CommitTimestamp = std::move(ev->CommitTimestamp); + } + if (ev->LockHandle) { QueryState->TxCtx->LockHandle = std::move(ev->LockHandle); } @@ -3040,6 +3044,12 @@ public: FillTxInfo(response); FillPoolId(response); + if (QueryState->CommitTimestamp) { + auto* ts = response->MutableCommitTimestamp(); + ts->set_plan_step(QueryState->CommitTimestamp->PlanStep); + ts->set_tx_id(QueryState->CommitTimestamp->TxId); + } + UpdateQueryExecutionCounters(); bool replyQueryId = false; diff --git a/ydb/core/kqp/ut/tx/kqp_tx_commit_timestamp_cdc_ut.cpp b/ydb/core/kqp/ut/tx/kqp_tx_commit_timestamp_cdc_ut.cpp new file mode 100644 index 00000000000..84a59bb6b0b --- /dev/null +++ b/ydb/core/kqp/ut/tx/kqp_tx_commit_timestamp_cdc_ut.cpp @@ -0,0 +1,226 @@ +#include <ydb/core/kqp/ut/common/kqp_ut_common.h> +#include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/topic/client.h> + +namespace NKikimr::NKqp { + +using namespace NYdb; +using namespace NYdb::NQuery; +using namespace NYdb::NTopic; + +namespace { + +TVector<NJson::TJsonValue> ReadCdcMessages(TKikimrRunner& kikimr, const TString& topicPath, size_t expectedCount) { + TTopicClient topicClient(kikimr.GetDriver()); + + TReadSessionSettings readSettings; + readSettings.ConsumerName("test_consumer"); + TTopicReadSettings topicReadSettings; + topicReadSettings.Path(topicPath); + readSettings.AppendTopics(topicReadSettings); + + auto readSession = topicClient.CreateReadSession(readSettings); + + TVector<NJson::TJsonValue> result; + const TInstant deadline = TInstant::Now() + TDuration::Seconds(30); + + while (result.size() < expectedCount && TInstant::Now() < deadline) { + const TDuration remain = deadline - TInstant::Now(); + if (!readSession->WaitEvent().Wait(remain)) { + break; + } + + for (auto& ev : readSession->GetEvents(false)) { + if (auto* startEvent = std::get_if<NTopic::TReadSessionEvent::TStartPartitionSessionEvent>(&ev)) { + startEvent->Confirm(); + } else if (auto* stopEvent = std::get_if<NTopic::TReadSessionEvent::TStopPartitionSessionEvent>(&ev)) { + stopEvent->Confirm(); + } else if (auto* data = std::get_if<NTopic::TReadSessionEvent::TDataReceivedEvent>(&ev)) { + for (const auto& m : data->GetMessages()) { + NJson::TJsonValue json; + if (NJson::ReadJsonTree(m.GetData(), &json)) { + result.push_back(json); + } + } + data->Commit(); + } else if (std::get_if<NTopic::TSessionClosedEvent>(&ev)) { + return result; + } + } + } + + return result; +} + +void ExecDdl(NQuery::TSession& session, const std::string& query, const char* what) { + auto result = session.ExecuteQuery(query, TTxControl::NoTx()).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::SUCCESS, + what << ": " << result.GetIssues().ToString()); +} + +} // namespace + +Y_UNIT_TEST_SUITE(KqpTxCommitTimestampCdc) { + Y_UNIT_TEST(SingleShardWrite) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + settings.SetWithSampleTables(false); + TKikimrRunner kikimr(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + ExecDdl(session, R"( + CREATE TABLE `/Root/CdcTest` ( + Key Uint64, + Value Text, + PRIMARY KEY (Key) + ); + )", "CREATE TABLE"); + + ExecDdl(session, R"( + ALTER TABLE `/Root/CdcTest` ADD CHANGEFEED `feed` WITH ( + MODE = 'UPDATES', FORMAT = 'JSON', VIRTUAL_TIMESTAMPS = TRUE + ); + )", "ADD CHANGEFEED"); + + ExecDdl(session, R"( + ALTER TOPIC `/Root/CdcTest/feed` ADD CONSUMER `test_consumer`; + )", "ADD CONSUMER"); + + auto result = session.ExecuteQuery(R"( + UPSERT INTO `/Root/CdcTest` (Key, Value) VALUES (1u, "hello"); + )", TTxControl::BeginTx(TTxSettings::StrictSerializableRW()).CommitTx()).ExtractValueSync(); + + UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::SUCCESS, result.GetIssues().ToString()); + UNIT_ASSERT_C(result.GetCommitTimestamp().has_value(), "Commit timestamp should be present"); + const auto& commitTs = *result.GetCommitTimestamp(); + UNIT_ASSERT_C(commitTs.PlanStep > 0, "PlanStep should be nonzero"); + UNIT_ASSERT_C(commitTs.TxId > 0, "TxId should be nonzero"); + + auto records = ReadCdcMessages(kikimr, "/Root/CdcTest/feed", 1); + UNIT_ASSERT_VALUES_EQUAL(records.size(), 1); + UNIT_ASSERT_C(records[0].Has("ts"), "CDC record should have 'ts' field"); + UNIT_ASSERT_VALUES_EQUAL(records[0]["ts"][0].GetUInteger(), commitTs.PlanStep); + UNIT_ASSERT_VALUES_EQUAL(records[0]["ts"][1].GetUInteger(), commitTs.TxId); + } + + Y_UNIT_TEST(ThreeInsertsOneShard) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + settings.SetWithSampleTables(false); + TKikimrRunner kikimr(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + ExecDdl(session, R"( + CREATE TABLE `/Root/CdcMultiInsert` ( + Key Uint64, + Value Text, + PRIMARY KEY (Key) + ) WITH (UNIFORM_PARTITIONS = 1); + )", "CREATE TABLE"); + + ExecDdl(session, R"( + ALTER TABLE `/Root/CdcMultiInsert` ADD CHANGEFEED `feed` WITH ( + MODE = 'UPDATES', FORMAT = 'JSON', VIRTUAL_TIMESTAMPS = TRUE + ); + )", "ADD CHANGEFEED"); + + ExecDdl(session, R"( + ALTER TOPIC `/Root/CdcMultiInsert/feed` ADD CONSUMER `test_consumer`; + )", "ADD CONSUMER"); + + auto beginResult = session.BeginTransaction(TTxSettings::StrictSerializableRW()).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(beginResult.GetStatus(), EStatus::SUCCESS, beginResult.GetIssues().ToString()); + auto tx = beginResult.GetTransaction(); + + for (ui64 i = 1; i <= 3; ++i) { + TString query = TStringBuilder() + << "INSERT INTO `/Root/CdcMultiInsert` (Key, Value) VALUES (" + << i << "u, \"val" << i << "\");"; + auto execResult = session.ExecuteQuery(query, TTxControl::Tx(tx)).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(execResult.GetStatus(), EStatus::SUCCESS, execResult.GetIssues().ToString()); + } + + auto commitResult = tx.Commit().ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(commitResult.GetStatus(), EStatus::SUCCESS, commitResult.GetIssues().ToString()); + UNIT_ASSERT_C(commitResult.GetCommitTimestamp().has_value(), "Commit timestamp should be present"); + const auto& commitTs = *commitResult.GetCommitTimestamp(); + UNIT_ASSERT_C(commitTs.PlanStep > 0, "PlanStep should be nonzero"); + UNIT_ASSERT_C(commitTs.TxId > 0, "TxId should be nonzero"); + + auto records = ReadCdcMessages(kikimr, "/Root/CdcMultiInsert/feed", 3); + UNIT_ASSERT_VALUES_EQUAL(records.size(), 3); + + for (size_t i = 0; i < records.size(); ++i) { + UNIT_ASSERT_C(records[i].Has("ts"), TStringBuilder() << "CDC record " << i << " should have 'ts' field"); + UNIT_ASSERT_VALUES_EQUAL_C( + records[i]["ts"][0].GetUInteger(), commitTs.PlanStep, + TStringBuilder() << "CDC record " << i << " PlanStep mismatch"); + UNIT_ASSERT_VALUES_EQUAL_C( + records[i]["ts"][1].GetUInteger(), commitTs.TxId, + TStringBuilder() << "CDC record " << i << " TxId mismatch"); + } + } + + Y_UNIT_TEST(ThreeUpsertsMultiShard) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + settings.SetWithSampleTables(false); + TKikimrRunner kikimr(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + ExecDdl(session, R"( + CREATE TABLE `/Root/CdcMultiShard` ( + Key Uint64, + Value Text, + PRIMARY KEY (Key) + ) WITH (UNIFORM_PARTITIONS = 4); + )", "CREATE TABLE"); + + ExecDdl(session, R"( + ALTER TABLE `/Root/CdcMultiShard` ADD CHANGEFEED `feed` WITH ( + MODE = 'UPDATES', FORMAT = 'JSON', VIRTUAL_TIMESTAMPS = TRUE + ); + )", "ADD CHANGEFEED"); + + ExecDdl(session, R"( + ALTER TOPIC `/Root/CdcMultiShard/feed` ADD CONSUMER `test_consumer`; + )", "ADD CONSUMER"); + + auto beginResult = session.BeginTransaction(TTxSettings::StrictSerializableRW()).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(beginResult.GetStatus(), EStatus::SUCCESS, beginResult.GetIssues().ToString()); + auto tx = beginResult.GetTransaction(); + + const ui64 keys[] = {1, 1000000, 2000000}; + for (ui64 i = 0; i < 3; ++i) { + TString query = TStringBuilder() + << "UPSERT INTO `/Root/CdcMultiShard` (Key, Value) VALUES (" + << keys[i] << "u, \"val" << (i + 1) << "\");"; + auto execResult = session.ExecuteQuery(query, TTxControl::Tx(tx)).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(execResult.GetStatus(), EStatus::SUCCESS, execResult.GetIssues().ToString()); + } + + auto commitResult = tx.Commit().ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(commitResult.GetStatus(), EStatus::SUCCESS, commitResult.GetIssues().ToString()); + UNIT_ASSERT_C(commitResult.GetCommitTimestamp().has_value(), "Commit timestamp should be present"); + const auto& commitTs = *commitResult.GetCommitTimestamp(); + UNIT_ASSERT_C(commitTs.PlanStep > 0, "PlanStep should be nonzero"); + UNIT_ASSERT_C(commitTs.TxId > 0, "TxId should be nonzero"); + + auto records = ReadCdcMessages(kikimr, "/Root/CdcMultiShard/feed", 3); + UNIT_ASSERT_VALUES_EQUAL(records.size(), 3); + + for (size_t i = 0; i < records.size(); ++i) { + UNIT_ASSERT_C(records[i].Has("ts"), TStringBuilder() << "CDC record " << i << " should have 'ts' field"); + UNIT_ASSERT_VALUES_EQUAL_C( + records[i]["ts"][0].GetUInteger(), commitTs.PlanStep, + TStringBuilder() << "CDC record " << i << " PlanStep mismatch"); + UNIT_ASSERT_VALUES_EQUAL_C( + records[i]["ts"][1].GetUInteger(), commitTs.TxId, + TStringBuilder() << "CDC record " << i << " TxId mismatch"); + } + } +} + +} // namespace NKikimr::NKqp diff --git a/ydb/core/kqp/ut/tx/kqp_tx_ut.cpp b/ydb/core/kqp/ut/tx/kqp_tx_ut.cpp index 5d9d8af2225..4f487a74863 100644 --- a/ydb/core/kqp/ut/tx/kqp_tx_ut.cpp +++ b/ydb/core/kqp/ut/tx/kqp_tx_ut.cpp @@ -868,6 +868,137 @@ Y_UNIT_TEST_SUITE(KqpTx) { UNIT_ASSERT(foundInsertWithImmediate); UNIT_ASSERT(foundPrepare); } + + Y_UNIT_TEST(StrictSerializable_CommitTimestamp_ExecuteQuery) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + auto kikimr = TKikimrRunner(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + auto result = session.ExecuteQuery(R"( + UPSERT INTO `/Root/KeyValue` (Key, Value) VALUES (200u, "Commit1"); + )", NYdb::NQuery::TTxControl::BeginTx(NYdb::NQuery::TTxSettings::StrictSerializableRW()).CommitTx()).ExtractValueSync(); + + UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::SUCCESS, result.GetIssues().ToString()); + UNIT_ASSERT_C(result.GetCommitTimestamp().has_value(), "Commit timestamp should be present for StrictSerializableRW write commit"); + const auto& ts1 = *result.GetCommitTimestamp(); + UNIT_ASSERT_C(ts1.PlanStep > 0, "PlanStep should be nonzero"); + UNIT_ASSERT_C(ts1.TxId > 0, "TxId should be nonzero"); + + auto result2 = session.ExecuteQuery(R"( + UPSERT INTO `/Root/KeyValue` (Key, Value) VALUES (201u, "Commit2"); + )", NYdb::NQuery::TTxControl::BeginTx(NYdb::NQuery::TTxSettings::StrictSerializableRW()).CommitTx()).ExtractValueSync(); + + UNIT_ASSERT_VALUES_EQUAL_C(result2.GetStatus(), EStatus::SUCCESS, result2.GetIssues().ToString()); + UNIT_ASSERT_C(result2.GetCommitTimestamp().has_value(), "Commit timestamp should be present for StrictSerializableRW write commit"); + const auto& ts2 = *result2.GetCommitTimestamp(); + UNIT_ASSERT_C(ts2.PlanStep > 0, "PlanStep should be nonzero"); + UNIT_ASSERT_C(ts2.TxId > 0, "TxId should be nonzero"); + + UNIT_ASSERT_C(ts1 < ts2, "Second commit timestamp should be greater than first (lexicographic order)"); + } + + Y_UNIT_TEST(StrictSerializable_CommitTimestamp_ExplicitCommit) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + auto kikimr = TKikimrRunner(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + auto beginResult = session.BeginTransaction(NYdb::NQuery::TTxSettings::StrictSerializableRW()).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(beginResult.GetStatus(), EStatus::SUCCESS, beginResult.GetIssues().ToString()); + auto tx = beginResult.GetTransaction(); + + auto execResult = session.ExecuteQuery(R"( + UPSERT INTO `/Root/KeyValue` (Key, Value) VALUES (300u, "Explicit"); + )", NYdb::NQuery::TTxControl::Tx(tx)).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(execResult.GetStatus(), EStatus::SUCCESS, execResult.GetIssues().ToString()); + + auto commitResult = tx.Commit().ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(commitResult.GetStatus(), EStatus::SUCCESS, commitResult.GetIssues().ToString()); + UNIT_ASSERT_C(commitResult.GetCommitTimestamp().has_value(), "Commit timestamp should be present for StrictSerializableRW write commit"); + const auto& ts = *commitResult.GetCommitTimestamp(); + UNIT_ASSERT_C(ts.PlanStep > 0, "PlanStep should be nonzero"); + UNIT_ASSERT_C(ts.TxId > 0, "TxId should be nonzero"); + } + + Y_UNIT_TEST(StrictSerializable_CommitTimestamp_ReadOnly) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + auto kikimr = TKikimrRunner(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + auto result = session.ExecuteQuery(R"( + SELECT * FROM `/Root/KeyValue` WHERE Key = 100u; + )", NYdb::NQuery::TTxControl::BeginTx(NYdb::NQuery::TTxSettings::StrictSerializableRW()).CommitTx()).ExtractValueSync(); + + UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::SUCCESS, result.GetIssues().ToString()); + UNIT_ASSERT_C(!result.GetCommitTimestamp().has_value(), "Commit timestamp should not be present for read-only query"); + } + + Y_UNIT_TEST(StrictSerializable_CommitTimestamp_ReadOnly_ExplicitCommit) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + auto kikimr = TKikimrRunner(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + auto beginResult = session.BeginTransaction(NYdb::NQuery::TTxSettings::StrictSerializableRW()).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(beginResult.GetStatus(), EStatus::SUCCESS, beginResult.GetIssues().ToString()); + auto tx = beginResult.GetTransaction(); + + auto execResult = session.ExecuteQuery(R"( + SELECT * FROM `/Root/KeyValue` WHERE Key = 100u; + )", NYdb::NQuery::TTxControl::Tx(tx)).ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(execResult.GetStatus(), EStatus::SUCCESS, execResult.GetIssues().ToString()); + + auto commitResult = tx.Commit().ExtractValueSync(); + UNIT_ASSERT_VALUES_EQUAL_C(commitResult.GetStatus(), EStatus::SUCCESS, commitResult.GetIssues().ToString()); + UNIT_ASSERT_C(!commitResult.GetCommitTimestamp().has_value(), "Commit timestamp should not be present for read-only explicit commit"); + } + + Y_UNIT_TEST(StrictSerializable_CommitTimestamp_SerializableRW) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + auto kikimr = TKikimrRunner(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + auto result = session.ExecuteQuery(R"( + UPSERT INTO `/Root/KeyValue` (Key, Value) VALUES (400u, "Serializable"); + )", NYdb::NQuery::TTxControl::BeginTx(NYdb::NQuery::TTxSettings::SerializableRW()).CommitTx()).ExtractValueSync(); + + UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::SUCCESS, result.GetIssues().ToString()); + UNIT_ASSERT_C(!result.GetCommitTimestamp().has_value(), "Commit timestamp should not be present for non-StrictSerializableRW isolation"); + } + + Y_UNIT_TEST(StrictSerializable_CommitTimestamp_Order) { + TKikimrSettings settings; + settings.SetEnableStrictSerializableIsolation(true); + auto kikimr = TKikimrRunner(settings); + auto db = kikimr.GetQueryClient(); + auto session = db.GetSession().GetValueSync().GetSession(); + + auto result1 = session.ExecuteQuery(R"( + UPSERT INTO `/Root/KeyValue` (Key, Value) VALUES (500u, "Distinct1"); + )", NYdb::NQuery::TTxControl::BeginTx(NYdb::NQuery::TTxSettings::StrictSerializableRW()).CommitTx()).ExtractValueSync(); + + UNIT_ASSERT_VALUES_EQUAL_C(result1.GetStatus(), EStatus::SUCCESS, result1.GetIssues().ToString()); + UNIT_ASSERT_C(result1.GetCommitTimestamp().has_value(), "Commit timestamp should be present"); + const auto& ts1 = *result1.GetCommitTimestamp(); + + auto result2 = session.ExecuteQuery(R"( + UPSERT INTO `/Root/KeyValue` (Key, Value) VALUES (501u, "Distinct2"); + )", NYdb::NQuery::TTxControl::BeginTx(NYdb::NQuery::TTxSettings::StrictSerializableRW()).CommitTx()).ExtractValueSync(); + + UNIT_ASSERT_VALUES_EQUAL_C(result2.GetStatus(), EStatus::SUCCESS, result2.GetIssues().ToString()); + UNIT_ASSERT_C(result2.GetCommitTimestamp().has_value(), "Commit timestamp should be present"); + const auto& ts2 = *result2.GetCommitTimestamp(); + + UNIT_ASSERT_C(ts1 < ts2, "Commit timestamps must be ordered and distinct"); + } } } // namespace NKqp diff --git a/ydb/core/kqp/ut/tx/ya.make b/ydb/core/kqp/ut/tx/ya.make index 3dc1e7152b8..900ffd20057 100644 --- a/ydb/core/kqp/ut/tx/ya.make +++ b/ydb/core/kqp/ut/tx/ya.make @@ -17,6 +17,7 @@ SRCS( kqp_sink_tx_ut.cpp kqp_snapshot_isolation_ut.cpp kqp_tx_ut.cpp + kqp_tx_commit_timestamp_cdc_ut.cpp kqp_rollback.cpp kqp_online_ro_ut.cpp ) diff --git a/ydb/core/protos/kqp.proto b/ydb/core/protos/kqp.proto index 97c0af935ed..992347697f0 100644 --- a/ydb/core/protos/kqp.proto +++ b/ydb/core/protos/kqp.proto @@ -10,6 +10,7 @@ import "ydb/core/protos/kqp_tablemetadata.proto"; import "ydb/core/protos/kqp_physical.proto"; import "ydb/core/protos/kqp_stats.proto"; import "ydb/core/protos/data_events.proto"; +import "ydb/public/api/protos/ydb_common.proto"; import "ydb/public/api/protos/ydb_formats.proto"; import "ydb/public/api/protos/ydb_status_codes.proto"; import "ydb/public/api/protos/ydb_table.proto"; @@ -322,6 +323,7 @@ message TQueryResponse { optional string QueryDiagnostics = 15; optional TQueryResponseExtraInfo ExtraInfo = 16; optional string EffectivePoolId = 17; + optional Ydb.VirtualTimestamp CommitTimestamp = 18; } message TEvQueryResponse { diff --git a/ydb/public/api/protos/ydb_query.proto b/ydb/public/api/protos/ydb_query.proto index b4dc75df786..e48f4fcb056 100644 --- a/ydb/public/api/protos/ydb_query.proto +++ b/ydb/public/api/protos/ydb_query.proto @@ -142,6 +142,10 @@ message CommitTransactionRequest { message CommitTransactionResponse { StatusIds.StatusCode status = 1; repeated Ydb.Issue.IssueMessage issues = 2; + + // Commit timestamp (PlanStep, TxId) for StrictSerializableRW write transactions. + // Present only on SUCCESS and when the transaction had write effects. + Ydb.VirtualTimestamp commit_timestamp = 3; } message RollbackTransactionRequest { @@ -255,6 +259,10 @@ message ExecuteQueryResponsePart { TransactionMeta tx_meta = 6; VirtualTimestamp snapshot_timestamp = 7; + + // Commit timestamp (PlanStep, TxId) for StrictSerializableRW write transactions. + // Present only in the final (trailing) part on SUCCESS and when the transaction had write effects. + VirtualTimestamp commit_timestamp = 8; } message ExecuteScriptRequest { diff --git a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/client.h b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/client.h index 62464535561..91ed6635b04 100644 --- a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/client.h +++ b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/client.h @@ -8,6 +8,7 @@ #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/driver/driver.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/params/params.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/retry/retry.h> +#include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/tx/tx.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/request_settings.h> @@ -293,19 +294,25 @@ public: const std::optional<TTransaction>& GetTransaction() const { return Transaction_; } - TExecuteQueryPart(TStatus&& status, std::optional<TExecStats>&& queryStats, std::optional<TTransaction>&& tx) + const std::optional<NScheme::TVirtualTimestamp>& GetCommitTimestamp() const { return CommitTimestamp_; } + + TExecuteQueryPart(TStatus&& status, std::optional<TExecStats>&& queryStats, std::optional<TTransaction>&& tx, + std::optional<NScheme::TVirtualTimestamp>&& commitTimestamp = {}) : TStreamPartStatus(std::move(status)) , Stats_(std::move(queryStats)) , Transaction_(std::move(tx)) + , CommitTimestamp_(std::move(commitTimestamp)) {} TExecuteQueryPart(TStatus&& status, TResultSet&& resultSet, int64_t resultSetIndex, - std::optional<TExecStats>&& queryStats, std::optional<TTransaction>&& tx) + std::optional<TExecStats>&& queryStats, std::optional<TTransaction>&& tx, + std::optional<NScheme::TVirtualTimestamp>&& commitTimestamp = {}) : TStreamPartStatus(std::move(status)) , ResultSet_(std::move(resultSet)) , ResultSetIndex_(resultSetIndex) , Stats_(std::move(queryStats)) , Transaction_(std::move(tx)) + , CommitTimestamp_(std::move(commitTimestamp)) {} private: @@ -313,6 +320,7 @@ private: int64_t ResultSetIndex_ = 0; std::optional<TExecStats> Stats_; std::optional<TTransaction> Transaction_; + std::optional<NScheme::TVirtualTimestamp> CommitTimestamp_; }; class TExecuteQueryResult : public TStatus { @@ -325,22 +333,27 @@ public: std::optional<TTransaction> GetTransaction() const {return Transaction_; } + const std::optional<NScheme::TVirtualTimestamp>& GetCommitTimestamp() const { return CommitTimestamp_; } + TExecuteQueryResult(TStatus&& status) : TStatus(std::move(status)) {} TExecuteQueryResult(TStatus&& status, std::vector<TResultSet>&& resultSets, - std::optional<TExecStats>&& stats, std::optional<TTransaction>&& tx) + std::optional<TExecStats>&& stats, std::optional<TTransaction>&& tx, + std::optional<NScheme::TVirtualTimestamp>&& commitTimestamp = {}) : TStatus(std::move(status)) , ResultSets_(std::move(resultSets)) , Stats_(std::move(stats)) , Transaction_(std::move(tx)) + , CommitTimestamp_(std::move(commitTimestamp)) {} private: std::vector<TResultSet> ResultSets_; std::optional<TExecStats> Stats_; std::optional<TTransaction> Transaction_; + std::optional<NScheme::TVirtualTimestamp> CommitTimestamp_; }; } // namespace NYdb::NQuery diff --git a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/query.h b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/query.h index 5f0fe3462b3..91d9c54b56e 100644 --- a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/query.h +++ b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/query.h @@ -6,6 +6,7 @@ #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/result/result.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/retry/retry.h> +#include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/fluent_settings_helpers.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/operation/operation.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/request_settings.h> @@ -131,6 +132,12 @@ struct TDeleteSessionSettings : public TRequestSettings<TDeleteSessionSettings> class TCommitTransactionResult : public TStatus { public: TCommitTransactionResult(TStatus&& status); + TCommitTransactionResult(TStatus&& status, std::optional<NScheme::TVirtualTimestamp>&& commitTimestamp); + + const std::optional<NScheme::TVirtualTimestamp>& GetCommitTimestamp() const { return CommitTimestamp_; } + +private: + std::optional<NScheme::TVirtualTimestamp> CommitTimestamp_; }; using TAsyncBeginTransactionResult = NThreading::TFuture<TBeginTransactionResult>; diff --git a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/scheme/scheme.h b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/scheme/scheme.h index b638b824bb6..99e921a9105 100644 --- a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/scheme/scheme.h +++ b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/scheme/scheme.h @@ -1,6 +1,7 @@ #pragma once #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/driver/driver.h> +#include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h> namespace Ydb { class VirtualTimestamp; @@ -57,25 +58,6 @@ enum class ESchemeEntryType : i32 { Secret = 26, }; -struct TVirtualTimestamp { - uint64_t PlanStep = 0; - uint64_t TxId = 0; - - TVirtualTimestamp() = default; - TVirtualTimestamp(uint64_t planStep, uint64_t txId); - TVirtualTimestamp(const ::Ydb::VirtualTimestamp& proto); - - std::string ToString() const; - void Out(IOutputStream& out) const; - - bool operator<(const TVirtualTimestamp& rhs) const; - bool operator<=(const TVirtualTimestamp& rhs) const; - bool operator>(const TVirtualTimestamp& rhs) const; - bool operator>=(const TVirtualTimestamp& rhs) const; - bool operator==(const TVirtualTimestamp& rhs) const; - bool operator!=(const TVirtualTimestamp& rhs) const; -}; - struct TSchemeEntry { std::string Name; std::string Owner; diff --git a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h new file mode 100644 index 00000000000..ce71b9010f3 --- /dev/null +++ b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h @@ -0,0 +1,32 @@ +#pragma once + +#include <ydb/public/api/protos/ydb_common.pb.h> + +#include <util/stream/output.h> + +#include <tuple> + +namespace NYdb::inline Dev { +namespace NScheme { + +struct TVirtualTimestamp { + uint64_t PlanStep = 0; + uint64_t TxId = 0; + + TVirtualTimestamp() = default; + TVirtualTimestamp(uint64_t planStep, uint64_t txId); + TVirtualTimestamp(const ::Ydb::VirtualTimestamp& proto); + + std::string ToString() const; + void Out(IOutputStream& out) const; + + bool operator<(const TVirtualTimestamp& rhs) const; + bool operator<=(const TVirtualTimestamp& rhs) const; + bool operator>(const TVirtualTimestamp& rhs) const; + bool operator>=(const TVirtualTimestamp& rhs) const; + bool operator==(const TVirtualTimestamp& rhs) const; + bool operator!=(const TVirtualTimestamp& rhs) const; +}; + +} // namespace NScheme +} // namespace NYdb diff --git a/ydb/public/sdk/cpp/src/client/query/client.cpp b/ydb/public/sdk/cpp/src/client/query/client.cpp index 61af1e38afb..e8bdda91f2e 100644 --- a/ydb/public/sdk/cpp/src/client/query/client.cpp +++ b/ydb/public/sdk/cpp/src/client/query/client.cpp @@ -260,7 +260,12 @@ public: obs->End(commitTxStatus.GetStatus(), commitTxStatus.GetEndpoint()); - TCommitTransactionResult commitTxResult(std::move(commitTxStatus)); + std::optional<NScheme::TVirtualTimestamp> commitTimestamp; + if (response->has_commit_timestamp()) { + commitTimestamp = NScheme::TVirtualTimestamp(response->commit_timestamp()); + } + + TCommitTransactionResult commitTxResult(std::move(commitTxStatus), std::move(commitTimestamp)); promise.SetValue(std::move(commitTxResult)); } else { obs->End(status.Status, status.Endpoint); diff --git a/ydb/public/sdk/cpp/src/client/query/impl/exec_query.cpp b/ydb/public/sdk/cpp/src/client/query/impl/exec_query.cpp index 6911525c777..79589c3b923 100644 --- a/ydb/public/sdk/cpp/src/client/query/impl/exec_query.cpp +++ b/ydb/public/sdk/cpp/src/client/query/impl/exec_query.cpp @@ -88,6 +88,7 @@ public: std::optional<TExecStats> stats; std::optional<TTransaction> tx; + std::optional<NScheme::TVirtualTimestamp> commitTimestamp; if (self->Response_.has_exec_stats()) { stats = TExecStats(std::move(*self->Response_.mutable_exec_stats())); } @@ -96,16 +97,21 @@ public: tx = TTransaction(self->Session_.value(), self->Response_.tx_meta().id()); } + if (self->Response_.has_commit_timestamp()) { + commitTimestamp = NScheme::TVirtualTimestamp(self->Response_.commit_timestamp()); + } + if (self->Response_.has_result_set()) { promise.SetValue({ std::move(status), TResultSet(std::move(*self->Response_.mutable_result_set())), self->Response_.result_set_index(), std::move(stats), - std::move(tx) + std::move(tx), + std::move(commitTimestamp) }); } else { - promise.SetValue({std::move(status), std::move(stats), std::move(tx)}); + promise.SetValue({std::move(status), std::move(stats), std::move(tx), std::move(commitTimestamp)}); } } }; @@ -160,6 +166,7 @@ struct TExecuteQueryBuffer : public TThrRefBase, TNonCopyable { std::vector<Ydb::ResultSet> ResultSets_; std::optional<TExecStats> Stats_; std::optional<TTransaction> Tx_; + std::optional<NScheme::TVirtualTimestamp> CommitTimestamp_; std::vector<std::string> ArrowSchemas_; std::vector<std::vector<std::string>> BytesData_; std::vector<bool> ResultSetSeen_; @@ -174,6 +181,10 @@ struct TExecuteQueryBuffer : public TThrRefBase, TNonCopyable { self->Stats_ = st; } + if (const auto& ct = part.GetCommitTimestamp()) { + self->CommitTimestamp_ = ct; + } + if (!part.IsSuccess()) { std::optional<TExecStats> stats; std::swap(self->Stats_, stats); @@ -182,12 +193,14 @@ struct TExecuteQueryBuffer : public TThrRefBase, TNonCopyable { std::vector<NYdb::NIssue::TIssue> issues; std::vector<Ydb::ResultSet> resultProtos; std::optional<TTransaction> tx; + std::optional<NScheme::TVirtualTimestamp> commitTimestamp; std::vector<std::string> arrowSchemas; std::vector<std::vector<std::string>> bytesData; std::swap(self->Issues_, issues); std::swap(self->ResultSets_, resultProtos); std::swap(self->Tx_, tx); + std::swap(self->CommitTimestamp_, commitTimestamp); std::swap(self->ArrowSchemas_, arrowSchemas); std::swap(self->BytesData_, bytesData); @@ -204,10 +217,11 @@ struct TExecuteQueryBuffer : public TThrRefBase, TNonCopyable { TStatus(EStatus::SUCCESS, NYdb::NIssue::TIssues(std::move(issues))), std::move(resultSets), std::move(stats), - std::move(tx) + std::move(tx), + std::move(commitTimestamp) )); } else { - self->Promise_.SetValue(TExecuteQueryResult(std::move(part), {}, std::move(stats), {})); + self->Promise_.SetValue(TExecuteQueryResult(std::move(part), {}, std::move(stats), {})); // No commit timestamp on error } return; diff --git a/ydb/public/sdk/cpp/src/client/query/impl/ya.make b/ydb/public/sdk/cpp/src/client/query/impl/ya.make index 3cd591433bf..0f14ffdd428 100644 --- a/ydb/public/sdk/cpp/src/client/query/impl/ya.make +++ b/ydb/public/sdk/cpp/src/client/query/impl/ya.make @@ -14,6 +14,7 @@ PEERDIR( ydb/public/sdk/cpp/src/client/impl/session ydb/public/sdk/cpp/src/client/impl/observability ydb/public/sdk/cpp/src/client/proto + ydb/public/sdk/cpp/src/client/types ) END() diff --git a/ydb/public/sdk/cpp/src/client/query/query.cpp b/ydb/public/sdk/cpp/src/client/query/query.cpp index f4b8cddceab..d0b547fdd2f 100644 --- a/ydb/public/sdk/cpp/src/client/query/query.cpp +++ b/ydb/public/sdk/cpp/src/client/query/query.cpp @@ -66,4 +66,9 @@ TCommitTransactionResult::TCommitTransactionResult(TStatus&& status) : TStatus(std::move(status)) {} +TCommitTransactionResult::TCommitTransactionResult(TStatus&& status, std::optional<NScheme::TVirtualTimestamp>&& commitTimestamp) + : TStatus(std::move(status)) + , CommitTimestamp_(std::move(commitTimestamp)) +{} + } // namespace NYdb::NQuery diff --git a/ydb/public/sdk/cpp/src/client/query/ya.make b/ydb/public/sdk/cpp/src/client/query/ya.make index 69507e40c1a..27ccdc7d704 100644 --- a/ydb/public/sdk/cpp/src/client/query/ya.make +++ b/ydb/public/sdk/cpp/src/client/query/ya.make @@ -16,6 +16,7 @@ PEERDIR( ydb/public/sdk/cpp/src/client/metrics ydb/public/sdk/cpp/src/client/query/impl ydb/public/sdk/cpp/src/client/result + ydb/public/sdk/cpp/src/client/types ydb/public/sdk/cpp/src/client/types/operation ) diff --git a/ydb/public/sdk/cpp/src/client/scheme/scheme.cpp b/ydb/public/sdk/cpp/src/client/scheme/scheme.cpp index 479929ef268..d3b342cfec9 100644 --- a/ydb/public/sdk/cpp/src/client/scheme/scheme.cpp +++ b/ydb/public/sdk/cpp/src/client/scheme/scheme.cpp @@ -29,52 +29,6 @@ void TPermissions::SerializeTo(::Ydb::Scheme::Permissions& proto) const { } } -TVirtualTimestamp::TVirtualTimestamp(uint64_t planStep, uint64_t txId) - : PlanStep(planStep) - , TxId(txId) -{} - -TVirtualTimestamp::TVirtualTimestamp(const ::Ydb::VirtualTimestamp& proto) - : TVirtualTimestamp(proto.plan_step(), proto.tx_id()) -{} - -std::string TVirtualTimestamp::ToString() const { - TString result; - TStringOutput out(result); - Out(out); - return result; -} - -void TVirtualTimestamp::Out(IOutputStream& out) const { - out << "{ plan_step: " << PlanStep - << ", tx_id: " << TxId - << " }"; -} - -bool TVirtualTimestamp::operator<(const TVirtualTimestamp& rhs) const { - return PlanStep < rhs.PlanStep && TxId < rhs.TxId; -} - -bool TVirtualTimestamp::operator<=(const TVirtualTimestamp& rhs) const { - return PlanStep <= rhs.PlanStep && TxId <= rhs.TxId; -} - -bool TVirtualTimestamp::operator>(const TVirtualTimestamp& rhs) const { - return PlanStep > rhs.PlanStep && TxId > rhs.TxId; -} - -bool TVirtualTimestamp::operator>=(const TVirtualTimestamp& rhs) const { - return PlanStep >= rhs.PlanStep && TxId >= rhs.TxId; -} - -bool TVirtualTimestamp::operator==(const TVirtualTimestamp& rhs) const { - return PlanStep == rhs.PlanStep && TxId == rhs.TxId; -} - -bool TVirtualTimestamp::operator!=(const TVirtualTimestamp& rhs) const { - return !(*this == rhs); -} - static ESchemeEntryType ConvertProtoEntryType(::Ydb::Scheme::Entry::Type entry) { switch (entry) { case ::Ydb::Scheme::Entry::DIRECTORY: diff --git a/ydb/public/sdk/cpp/src/client/types/virtual_timestamp.cpp b/ydb/public/sdk/cpp/src/client/types/virtual_timestamp.cpp new file mode 100644 index 00000000000..8a80a50a2f3 --- /dev/null +++ b/ydb/public/sdk/cpp/src/client/types/virtual_timestamp.cpp @@ -0,0 +1,58 @@ +#include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h> + +#include <util/stream/output.h> +#include <util/string/printf.h> + +#include <tuple> + +namespace NYdb::inline Dev { +namespace NScheme { + +TVirtualTimestamp::TVirtualTimestamp(uint64_t planStep, uint64_t txId) + : PlanStep(planStep) + , TxId(txId) +{} + +TVirtualTimestamp::TVirtualTimestamp(const ::Ydb::VirtualTimestamp& proto) + : TVirtualTimestamp(proto.plan_step(), proto.tx_id()) +{} + +std::string TVirtualTimestamp::ToString() const { + TString result; + TStringOutput out(result); + Out(out); + return result; +} + +void TVirtualTimestamp::Out(IOutputStream& out) const { + out << "{ plan_step: " << PlanStep + << ", tx_id: " << TxId + << " }"; +} + +bool TVirtualTimestamp::operator<(const TVirtualTimestamp& rhs) const { + return std::tie(PlanStep, TxId) < std::tie(rhs.PlanStep, rhs.TxId); +} + +bool TVirtualTimestamp::operator<=(const TVirtualTimestamp& rhs) const { + return std::tie(PlanStep, TxId) <= std::tie(rhs.PlanStep, rhs.TxId); +} + +bool TVirtualTimestamp::operator>(const TVirtualTimestamp& rhs) const { + return std::tie(PlanStep, TxId) > std::tie(rhs.PlanStep, rhs.TxId); +} + +bool TVirtualTimestamp::operator>=(const TVirtualTimestamp& rhs) const { + return std::tie(PlanStep, TxId) >= std::tie(rhs.PlanStep, rhs.TxId); +} + +bool TVirtualTimestamp::operator==(const TVirtualTimestamp& rhs) const { + return PlanStep == rhs.PlanStep && TxId == rhs.TxId; +} + +bool TVirtualTimestamp::operator!=(const TVirtualTimestamp& rhs) const { + return !(*this == rhs); +} + +} // namespace NScheme +} // namespace NYdb diff --git a/ydb/public/sdk/cpp/src/client/types/ya.make b/ydb/public/sdk/cpp/src/client/types/ya.make index 98b9db6ea2a..e6a7e901e49 100644 --- a/ydb/public/sdk/cpp/src/client/types/ya.make +++ b/ydb/public/sdk/cpp/src/client/types/ya.make @@ -1,6 +1,7 @@ LIBRARY() SRCS( + virtual_timestamp.cpp ydb.cpp ) diff --git a/ydb/public/sdk/cpp/tests/unit/client/query/virtual_timestamp_ut.cpp b/ydb/public/sdk/cpp/tests/unit/client/query/virtual_timestamp_ut.cpp new file mode 100644 index 00000000000..34589da9439 --- /dev/null +++ b/ydb/public/sdk/cpp/tests/unit/client/query/virtual_timestamp_ut.cpp @@ -0,0 +1,42 @@ +#include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h> + +#include <library/cpp/testing/unittest/registar.h> + +using NYdb::NScheme::TVirtualTimestamp; + +Y_UNIT_TEST_SUITE(VirtualTimestamp) { + Y_UNIT_TEST(LexicographicLess) { + UNIT_ASSERT(TVirtualTimestamp(1, 100) < TVirtualTimestamp(2, 1)); + UNIT_ASSERT(TVirtualTimestamp(1, 50) < TVirtualTimestamp(1, 100)); + UNIT_ASSERT(!(TVirtualTimestamp(2, 1) < TVirtualTimestamp(1, 100))); + UNIT_ASSERT(!(TVirtualTimestamp(1, 100) < TVirtualTimestamp(1, 100))); + } + + Y_UNIT_TEST(LexicographicGreater) { + UNIT_ASSERT(TVirtualTimestamp(2, 1) > TVirtualTimestamp(1, 100)); + UNIT_ASSERT(TVirtualTimestamp(1, 100) > TVirtualTimestamp(1, 50)); + UNIT_ASSERT(!(TVirtualTimestamp(1, 100) > TVirtualTimestamp(2, 1))); + UNIT_ASSERT(!(TVirtualTimestamp(1, 100) > TVirtualTimestamp(1, 100))); + } + + Y_UNIT_TEST(LessOrEqual) { + UNIT_ASSERT(TVirtualTimestamp(1, 100) <= TVirtualTimestamp(1, 100)); + UNIT_ASSERT(TVirtualTimestamp(1, 50) <= TVirtualTimestamp(1, 100)); + UNIT_ASSERT(TVirtualTimestamp(1, 100) <= TVirtualTimestamp(2, 1)); + UNIT_ASSERT(!(TVirtualTimestamp(2, 1) <= TVirtualTimestamp(1, 100))); + } + + Y_UNIT_TEST(GreaterOrEqual) { + UNIT_ASSERT(TVirtualTimestamp(1, 100) >= TVirtualTimestamp(1, 100)); + UNIT_ASSERT(TVirtualTimestamp(1, 100) >= TVirtualTimestamp(1, 50)); + UNIT_ASSERT(TVirtualTimestamp(2, 1) >= TVirtualTimestamp(1, 100)); + UNIT_ASSERT(!(TVirtualTimestamp(1, 50) >= TVirtualTimestamp(1, 100))); + } + + Y_UNIT_TEST(Equality) { + UNIT_ASSERT(TVirtualTimestamp(1, 100) == TVirtualTimestamp(1, 100)); + UNIT_ASSERT(!(TVirtualTimestamp(1, 100) == TVirtualTimestamp(1, 101))); + UNIT_ASSERT(!(TVirtualTimestamp(1, 100) == TVirtualTimestamp(2, 100))); + UNIT_ASSERT(TVirtualTimestamp(1, 100) != TVirtualTimestamp(2, 100)); + } +} diff --git a/ydb/public/sdk/cpp/tests/unit/client/query/ya.make b/ydb/public/sdk/cpp/tests/unit/client/query/ya.make index c36998d57d6..f0572e36a5a 100644 --- a/ydb/public/sdk/cpp/tests/unit/client/query/ya.make +++ b/ydb/public/sdk/cpp/tests/unit/client/query/ya.make @@ -13,6 +13,7 @@ SRCS( client_session_ut.cpp deferred_session_creation_ut.cpp query_stats_ut.cpp + virtual_timestamp_ut.cpp ) PEERDIR( @@ -23,6 +24,7 @@ PEERDIR( ydb/public/sdk/cpp/src/client/impl/session ydb/public/sdk/cpp/src/client/query/impl ydb/public/sdk/cpp/src/client/query + ydb/public/sdk/cpp/src/client/types ) END() |
