summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorAlek5andr-Kotov <[email protected]>2026-07-20 14:36:37 +0300
committerGitHub <[email protected]>2026-07-20 14:36:37 +0300
commitd204eee5192cf315f8a3fda408c2d7075bac1bd3 (patch)
treea12a84948587513ccecd90ab39a25dcd35d6490b
parent392ce6865a7b33fb4d4ede5c7070be19c499d33b (diff)
Harden multi-partition deferred publication finalize (#46491)
Co-authored-by: Alexander Kotov <[email protected]>
-rw-r--r--ydb/core/persqueue/pqtablet/pq_impl.cpp31
-rw-r--r--ydb/core/persqueue/ut/pqtablet_ut.cpp85
-rw-r--r--ydb/services/persqueue_v1/ut/topic_deferred_publish_ut.cpp122
3 files changed, 235 insertions, 3 deletions
diff --git a/ydb/core/persqueue/pqtablet/pq_impl.cpp b/ydb/core/persqueue/pqtablet/pq_impl.cpp
index f8df5032ba8..fe7d4f689ab 100644
--- a/ydb/core/persqueue/pqtablet/pq_impl.cpp
+++ b/ydb/core/persqueue/pqtablet/pq_impl.cpp
@@ -3454,6 +3454,25 @@ TMaybe<EDeferredFinalizeOp> GetSingleDeferredPublicationFinalizeOp(const NKikimr
return result;
}
+bool DeferredPublicationFinalizePartitionsMatchStaged(
+ const NKikimrPQ::TDataTransaction& txBody,
+ const THashMap<ui32, TPartitionId>& stagedPartitions)
+{
+ THashSet<ui32> opPartitions;
+ for (const auto& operation : txBody.GetOperations()) {
+ opPartitions.insert(operation.GetPartitionId());
+ }
+ if (opPartitions.size() != stagedPartitions.size()) {
+ return false;
+ }
+ for (const auto& [partitionId, _] : stagedPartitions) {
+ if (!opPartitions.contains(partitionId)) {
+ return false;
+ }
+ }
+ return true;
+}
+
} // namespace
void TPersQueue::HandleDataTransaction(TAutoPtr<TEvPersQueue::TEvProposeTransaction> ev,
@@ -3532,6 +3551,18 @@ void TPersQueue::HandleDataTransaction(TAutoPtr<TEvPersQueue::TEvProposeTransact
ctx);
return;
}
+
+ const TWriteId writeId = GetWriteId(txBody);
+ if (TxWrites.contains(writeId)
+ && !DeferredPublicationFinalizePartitionsMatchStaged(txBody, TxWrites.at(writeId).Partitions)) {
+ PQ_LOG_TX_W("TxId " << event.GetTxId() << " deferred publication finalize partition set mismatch");
+ SendProposeTransactionAbort(ActorIdFromProto(event.GetSourceActor()),
+ event.GetTxId(),
+ NKikimrPQ::TError::BAD_REQUEST,
+ "deferred publication finalize partition set mismatch",
+ ctx);
+ return;
+ }
} else if (txBody.HasWriteId() && GetWriteId(txBody).IsDeferredPublicationApiTransaction()) {
PQ_LOG_TX_W("TxId " << event.GetTxId() << " deferred publication WriteId requires finalize operation");
SendProposeTransactionAbort(ActorIdFromProto(event.GetSourceActor()),
diff --git a/ydb/core/persqueue/ut/pqtablet_ut.cpp b/ydb/core/persqueue/ut/pqtablet_ut.cpp
index 32ff61b3555..8984a3eed64 100644
--- a/ydb/core/persqueue/ut/pqtablet_ut.cpp
+++ b/ydb/core/persqueue/ut/pqtablet_ut.cpp
@@ -3360,6 +3360,91 @@ Y_UNIT_TEST_F(DeferredPublication_Several_Partitions_One_Tablet_Successful_Commi
UNIT_ASSERT_VALUES_EQUAL(messages1[0], "deferred-publish-payload");
}
+Y_UNIT_TEST_F(DeferredPublication_Several_Partitions_One_Tablet_Cancel, TPQTabletFixture) {
+ using TDeferredPublicationApi = NKikimrPQ::TPartitionOperation::TWriteOp::TDeferredPublicationApi;
+ const TWriteId writeId = NHelpers::MakeDeferredWriteId(52, "ext-52");
+ const ui64 txId = 70013;
+
+ PQTabletPrepare({.partitions=2}, {}, *Ctx);
+ EnsurePipeExist();
+
+ const TString ownerCookie0 = CreateSupportivePartitionForDeferredPublication(writeId, 0);
+ const TString ownerCookie1 = CreateSupportivePartitionForDeferredPublication(writeId, 1);
+
+ SendDeferredPublicationWriteRequest(writeId, ownerCookie0, 0);
+ SendDeferredPublicationWriteRequest(writeId, ownerCookie1, 1);
+ WaitForExactTxWritesCount(2);
+
+ CommitDeferredPublicationFinalize(writeId, txId, TDeferredPublicationApi::Cancel, {0, 1});
+
+ UNIT_ASSERT_VALUES_EQUAL(ReadMainPartitionMessages(0).size(), 0u);
+ UNIT_ASSERT_VALUES_EQUAL(ReadMainPartitionMessages(1).size(), 0u);
+}
+
+Y_UNIT_TEST_F(DeferredPublication_Finalize_MixedPublishAndCancel_Aborted, TPQTabletFixture) {
+ using TDeferredPublicationApi = NKikimrPQ::TPartitionOperation::TWriteOp::TDeferredPublicationApi;
+ const TWriteId writeId = NHelpers::MakeDeferredWriteId(53, "ext-53");
+ const ui64 txId = 70014;
+
+ PQTabletPrepare({.partitions=2}, {}, *Ctx);
+ EnsurePipeExist();
+
+ SendProposeTransactionRequest({.TxId=txId,
+ .Senders={Ctx->TabletId},
+ .Receivers={Ctx->TabletId},
+ .TxOps={
+ {.Partition=0, .Path="/topic", .DeferredPublicationOp=TDeferredPublicationApi::Publish},
+ {.Partition=1, .Path="/topic", .DeferredPublicationOp=TDeferredPublicationApi::Cancel},
+ },
+ .WriteId=writeId});
+ WaitProposeTransactionResponse({.TxId=txId,
+ .Status=NKikimrPQ::TEvProposeTransactionResult::ABORTED});
+}
+
+Y_UNIT_TEST_F(DeferredPublication_Finalize_WithReadOperation_Aborted, TPQTabletFixture) {
+ using TDeferredPublicationApi = NKikimrPQ::TPartitionOperation::TWriteOp::TDeferredPublicationApi;
+ const TWriteId writeId = NHelpers::MakeDeferredWriteId(54, "ext-54");
+ const ui64 txId = 70015;
+
+ PQTabletPrepare({.partitions=2}, {}, *Ctx);
+ EnsurePipeExist();
+
+ SendProposeTransactionRequest({.TxId=txId,
+ .Senders={Ctx->TabletId},
+ .Receivers={Ctx->TabletId},
+ .TxOps={
+ {.Partition=0, .Path="/topic", .DeferredPublicationOp=TDeferredPublicationApi::Publish},
+ {.Partition=1, .Consumer="user", .Begin=0, .End=0, .Path="/topic"},
+ },
+ .WriteId=writeId});
+ WaitProposeTransactionResponse({.TxId=txId,
+ .Status=NKikimrPQ::TEvProposeTransactionResult::ABORTED});
+}
+
+Y_UNIT_TEST_F(DeferredPublication_Finalize_PartialPartitionSet_Aborted, TPQTabletFixture) {
+ using TDeferredPublicationApi = NKikimrPQ::TPartitionOperation::TWriteOp::TDeferredPublicationApi;
+ const TWriteId writeId = NHelpers::MakeDeferredWriteId(55, "ext-55");
+ const ui64 txId = 70016;
+
+ PQTabletPrepare({.partitions=2}, {}, *Ctx);
+ EnsurePipeExist();
+
+ const TString ownerCookie0 = CreateSupportivePartitionForDeferredPublication(writeId, 0);
+ const TString ownerCookie1 = CreateSupportivePartitionForDeferredPublication(writeId, 1);
+
+ SendDeferredPublicationWriteRequest(writeId, ownerCookie0, 0);
+ SendDeferredPublicationWriteRequest(writeId, ownerCookie1, 1);
+ WaitForExactTxWritesCount(2);
+
+ SendProposeTransactionRequest({.TxId=txId,
+ .Senders={Ctx->TabletId},
+ .Receivers={Ctx->TabletId},
+ .TxOps={{.Partition=0, .Path="/topic", .DeferredPublicationOp=TDeferredPublicationApi::Publish}},
+ .WriteId=writeId});
+ WaitProposeTransactionResponse({.TxId=txId,
+ .Status=NKikimrPQ::TEvProposeTransactionResult::ABORTED});
+}
+
Y_UNIT_TEST_F(DeferredPublication_Publish_Before_Write_Ack, TPQTabletFixture) {
using TDeferredPublicationApi = NKikimrPQ::TPartitionOperation::TWriteOp::TDeferredPublicationApi;
const TWriteId writeId = NHelpers::MakeDeferredWriteId(45, "ext-45");
diff --git a/ydb/services/persqueue_v1/ut/topic_deferred_publish_ut.cpp b/ydb/services/persqueue_v1/ut/topic_deferred_publish_ut.cpp
index 10c58a12240..cfbe9e4c964 100644
--- a/ydb/services/persqueue_v1/ut/topic_deferred_publish_ut.cpp
+++ b/ydb/services/persqueue_v1/ut/topic_deferred_publish_ut.cpp
@@ -210,13 +210,14 @@ TMaybe<TString> ParseLegacyReadPayload(const TString& raw) {
TMaybe<TString> TryReadFirstTopicMessage(
NPersQueue::TTestServer& server,
const TString& topicShortName,
- TDuration timeout = TDuration::Seconds(30))
+ TDuration timeout = TDuration::Seconds(30),
+ ui32 partitionId = 0)
{
const TString topic = "rt3.dc1--" + topicShortName;
const TInstant deadline = TInstant::Now() + timeout;
while (TInstant::Now() < deadline) {
THolder<NMsgBusProxy::TBusPersQueue> request = TRequestReadPQ{
- topic, 0, 0, 100, "user", 0}.GetRequest();
+ topic, partitionId, 0, 100, "user", 0}.GetRequest();
request.Get()->Record.SetTicket("root@builtin");
const auto response = server.AnnoyingClient->CallPersQueueGRPC(request->Record);
@@ -718,8 +719,59 @@ std::unique_ptr<Ydb::Topic::V1::TopicService::Stub> MakeTopicServiceStub(
return Ydb::Topic::V1::TopicService::NewStub(channel);
}
+TString MakeLegacyStreamWriteTopicPath(const TString& topicShortName) {
+ return "/Root/PQ/rt3.dc1--" + topicShortName;
+}
+
+void AssertPartitionsOnSameTablet(
+ NPersQueue::TTestServer& server,
+ const TString& topicShortName,
+ const TVector<ui32>& partitionIds)
+{
+ UNIT_ASSERT(!partitionIds.empty());
+
+ const auto response = server.AnnoyingClient->Ls(MakeLegacyStreamWriteTopicPath(topicShortName));
+ UNIT_ASSERT(response);
+ UNIT_ASSERT_VALUES_EQUAL(response->Record.GetSchemeStatus(), NKikimrScheme::StatusSuccess);
+
+ THashMap<ui32, ui64> tabletByPartition;
+ for (const auto& partition : response->Record.GetPathDescription().GetPersQueueGroup().GetPartitions()) {
+ tabletByPartition[partition.GetPartitionId()] = partition.GetTabletId();
+ }
+
+ UNIT_ASSERT(tabletByPartition.contains(partitionIds.front()));
+ const ui64 expectedTabletId = tabletByPartition.at(partitionIds.front());
+ for (const ui32 partitionId : partitionIds) {
+ UNIT_ASSERT(tabletByPartition.contains(partitionId));
+ UNIT_ASSERT_VALUES_EQUAL(tabletByPartition.at(partitionId), expectedTabletId);
+ }
+}
+
void CreateLegacyStreamWriteTopic(NPersQueue::TTestServer& server, const TString& topicShortName, ui32 partitions = 2) {
- server.AnnoyingClient->CreateTopicNoLegacy("rt3.dc1--" + topicShortName, partitions);
+ const TString fullName = "rt3.dc1--" + topicShortName;
+ auto pqClient = NYdb::NPersQueue::TPersQueueClient(*server.AnnoyingClient->GetDriver());
+ auto settings = NYdb::NPersQueue::TCreateTopicSettings()
+ .PartitionsCount(partitions)
+ .PartitionsPerTablet(Max(partitions, 2u));
+ settings.ReadRules({NYdb::NPersQueue::TReadRuleSettings{}.ConsumerName("user")});
+
+ TString path = fullName;
+ if (!path.StartsWith("/Root")) {
+ path = TStringBuilder() << "/Root/PQ/" << fullName;
+ }
+
+ auto result = pqClient.CreateTopic(path, settings);
+ result.Wait();
+ UNIT_ASSERT_C(result.GetValue().IsSuccess(), result.GetValue().GetIssues().ToString());
+ server.AnnoyingClient->AddTopic(fullName);
+
+ if (partitions >= 2) {
+ TVector<ui32> partitionIds(Reserve(partitions));
+ for (ui32 partitionId = 0; partitionId < partitions; ++partitionId) {
+ partitionIds.push_back(partitionId);
+ }
+ AssertPartitionsOnSameTablet(server, topicShortName, partitionIds);
+ }
}
void AssertStreamWriteSuccess(const Ydb::Topic::StreamWriteMessage::FromServer& response) {
@@ -2005,6 +2057,70 @@ Y_UNIT_TEST(PublishAfterStreamWriteClearsRegistryAndMakesDataVisible) {
UNIT_ASSERT_VALUES_EQUAL(*message, TString(payload));
}
+Y_UNIT_TEST(PublishAfterStreamWriteToTwoPartitionsMakesDataVisible) {
+ auto fixture = TDeferredStreamWriteFixture::Enabled("finalize-two-partitions-topic", "ext-two-partitions");
+ GrantPublicationRegistryDelete(fixture.Server, "root@builtin");
+
+ constexpr TStringBuf payload0 = "deferred-payload-part-0";
+ constexpr TStringBuf payload1 = "deferred-payload-part-1";
+ {
+ auto session = fixture.OpenWriteStream("producer-part-0", 0);
+ WriteAndExpectWriteResponse(*session->Stream, MakeStreamWriteRequest(
+ 1,
+ TString(payload0),
+ std::make_pair(fixture.IntPublicationId, fixture.ExtPublicationId)));
+ }
+ {
+ auto session = fixture.OpenWriteStream("producer-part-1", 1);
+ WriteAndExpectWriteResponse(*session->Stream, MakeStreamWriteRequest(
+ 1,
+ TString(payload1),
+ std::make_pair(fixture.IntPublicationId, fixture.ExtPublicationId)));
+ }
+
+ const auto publishOutcome = CallPublish(*fixture.DeferredStub, "/Root", fixture.IntPublicationId);
+ UNIT_ASSERT(publishOutcome.RpcStatus.ok());
+ UNIT_ASSERT(publishOutcome.Operation.ready());
+ UNIT_ASSERT_VALUES_EQUAL(publishOutcome.Operation.status(), Ydb::StatusIds::SUCCESS);
+ UNIT_ASSERT_VALUES_EQUAL(CountPublications(fixture.Server, "root@builtin"), 0u);
+
+ const auto message0 = TryReadFirstTopicMessage(fixture.Server, fixture.TopicShortName, TDuration::Seconds(30), 0);
+ const auto message1 = TryReadFirstTopicMessage(fixture.Server, fixture.TopicShortName, TDuration::Seconds(30), 1);
+ UNIT_ASSERT(message0.Defined());
+ UNIT_ASSERT(message1.Defined());
+ UNIT_ASSERT_VALUES_EQUAL(*message0, TString(payload0));
+ UNIT_ASSERT_VALUES_EQUAL(*message1, TString(payload1));
+}
+
+Y_UNIT_TEST(CancelAfterStreamWriteToTwoPartitionsClearsRegistryWithoutData) {
+ auto fixture = TDeferredStreamWriteFixture::Enabled("finalize-cancel-two-partitions-topic", "ext-cancel-two-partitions");
+ GrantPublicationRegistryDelete(fixture.Server, "root@builtin");
+
+ {
+ auto session = fixture.OpenWriteStream("producer-cancel-0", 0);
+ WriteAndExpectWriteResponse(*session->Stream, MakeStreamWriteRequest(
+ 1,
+ "deferred-payload-cancel-0",
+ std::make_pair(fixture.IntPublicationId, fixture.ExtPublicationId)));
+ }
+ {
+ auto session = fixture.OpenWriteStream("producer-cancel-1", 1);
+ WriteAndExpectWriteResponse(*session->Stream, MakeStreamWriteRequest(
+ 1,
+ "deferred-payload-cancel-1",
+ std::make_pair(fixture.IntPublicationId, fixture.ExtPublicationId)));
+ }
+
+ const auto cancelOutcome = CallCancelPublication(*fixture.DeferredStub, "/Root", fixture.IntPublicationId);
+ UNIT_ASSERT(cancelOutcome.RpcStatus.ok());
+ UNIT_ASSERT(cancelOutcome.Operation.ready());
+ UNIT_ASSERT_VALUES_EQUAL(cancelOutcome.Operation.status(), Ydb::StatusIds::SUCCESS);
+ UNIT_ASSERT_VALUES_EQUAL(CountPublications(fixture.Server, "root@builtin"), 0u);
+
+ UNIT_ASSERT(!TryReadFirstTopicMessage(fixture.Server, fixture.TopicShortName, TDuration::Seconds(2), 0).Defined());
+ UNIT_ASSERT(!TryReadFirstTopicMessage(fixture.Server, fixture.TopicShortName, TDuration::Seconds(2), 1).Defined());
+}
+
Y_UNIT_TEST(CancelAfterStreamWriteClearsRegistryWithoutData) {
auto fixture = TDeferredStreamWriteFixture::Enabled("finalize-cancel-topic", "ext-cancel");
GrantPublicationRegistryDelete(fixture.Server, "root@builtin");