summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorNikita Vasilev <[email protected]>2026-07-22 11:56:22 +0300
committerGitHub <[email protected]>2026-07-22 11:56:22 +0300
commit75900fddb7cfff27d66271c4a9c419ddd852554b (patch)
treea98d115c3e1e5ae4ab8b9083f91e3f639ad4a055
parentd8cb30f4bd872ed7d6e7f845053a24ea60298c6f (diff)
StrictSerializable: return commit timestamp (#46795)
-rw-r--r--ydb/core/grpc_services/query/rpc_execute_query.cpp6
-rw-r--r--ydb/core/grpc_services/query/rpc_kqp_tx.cpp18
-rw-r--r--ydb/core/kqp/common/buffer/events.h9
-rw-r--r--ydb/core/kqp/executer_actor/kqp_data_executer.cpp1
-rw-r--r--ydb/core/kqp/executer_actor/kqp_executer.h2
-rw-r--r--ydb/core/kqp/runtime/kqp_write_actor.cpp10
-rw-r--r--ydb/core/kqp/session_actor/kqp_query_state.h3
-rw-r--r--ydb/core/kqp/session_actor/kqp_session_actor.cpp10
-rw-r--r--ydb/core/kqp/ut/tx/kqp_tx_commit_timestamp_cdc_ut.cpp226
-rw-r--r--ydb/core/kqp/ut/tx/kqp_tx_ut.cpp131
-rw-r--r--ydb/core/kqp/ut/tx/ya.make1
-rw-r--r--ydb/core/protos/kqp.proto2
-rw-r--r--ydb/public/api/protos/ydb_query.proto8
-rw-r--r--ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/client.h19
-rw-r--r--ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/query/query.h7
-rw-r--r--ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/scheme/scheme.h20
-rw-r--r--ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/virtual_timestamp.h32
-rw-r--r--ydb/public/sdk/cpp/src/client/query/client.cpp7
-rw-r--r--ydb/public/sdk/cpp/src/client/query/impl/exec_query.cpp22
-rw-r--r--ydb/public/sdk/cpp/src/client/query/impl/ya.make1
-rw-r--r--ydb/public/sdk/cpp/src/client/query/query.cpp5
-rw-r--r--ydb/public/sdk/cpp/src/client/query/ya.make1
-rw-r--r--ydb/public/sdk/cpp/src/client/scheme/scheme.cpp46
-rw-r--r--ydb/public/sdk/cpp/src/client/types/virtual_timestamp.cpp58
-rw-r--r--ydb/public/sdk/cpp/src/client/types/ya.make1
-rw-r--r--ydb/public/sdk/cpp/tests/unit/client/query/virtual_timestamp_ut.cpp42
-rw-r--r--ydb/public/sdk/cpp/tests/unit/client/query/ya.make2
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()