From e55c625445f270656e786b2ff5d6abb3fcefcfce Mon Sep 17 00:00:00 2001 From: Ilnaz Nizametdinov Date: Thu, 28 May 2026 09:31:02 +0300 Subject: Allow to create invalid VIEW during import (#41265) --- ydb/core/kqp/host/kqp_host.cpp | 1 + ydb/core/kqp/host/kqp_translate.cpp | 2 + ydb/core/kqp/host/kqp_translate.h | 6 + ydb/core/kqp/ut/view/view_ut.cpp | 36 +++++ ydb/core/tx/schemeshard/ut_restore/ut_restore.cpp | 169 ++++++++++++---------- 5 files changed, 137 insertions(+), 77 deletions(-) diff --git a/ydb/core/kqp/host/kqp_host.cpp b/ydb/core/kqp/host/kqp_host.cpp index 9b14c9b8eeb..98d6adf0e35 100644 --- a/ydb/core/kqp/host/kqp_host.cpp +++ b/ydb/core/kqp/host/kqp_host.cpp @@ -1684,6 +1684,7 @@ private: settingsBuilder.SetSqlAutoCommit(false) .SetUsePgParser(settings.UsePgParser) .SetFromConfig(SessionCtx->Config()) + .SetValidateViewStatement(!settings.IsInternalCall.GetOrElse(false)) .SetYqlSelect(settings.YqlSelect); auto compileResult = CompileYqlQuery(query, /* isSql */ true, ctx, sqlVersion, settingsBuilder); diff --git a/ydb/core/kqp/host/kqp_translate.cpp b/ydb/core/kqp/host/kqp_translate.cpp index c8a1914b32f..e2849d2896f 100644 --- a/ydb/core/kqp/host/kqp_translate.cpp +++ b/ydb/core/kqp/host/kqp_translate.cpp @@ -276,6 +276,8 @@ NSQLTranslation::TTranslationSettings TKqpTranslationSettingsBuilder::Build(NYql settings.YqlSelect = *YqlSelect; } + settings.ValidateViewStatement = ValidateViewStatement; + return settings; } diff --git a/ydb/core/kqp/host/kqp_translate.h b/ydb/core/kqp/host/kqp_translate.h index f7c35c14785..e88ecbef9aa 100644 --- a/ydb/core/kqp/host/kqp_translate.h +++ b/ydb/core/kqp/host/kqp_translate.h @@ -169,6 +169,11 @@ public: return *this; } + TKqpTranslationSettingsBuilder& SetValidateViewStatement(bool value) { + ValidateViewStatement = value; + return *this; + } + private: const NYql::EKikimrQueryType QueryType; ui16 KqpYqlSyntaxVersion = 1; @@ -190,6 +195,7 @@ private: NYql::EBackportCompatibleFeaturesMode BackportMode = NYql::EBackportCompatibleFeaturesMode::Released; bool IsAmbiguityError = false; TMaybe YqlSelect = {}; + bool ValidateViewStatement = true; }; NYql::EKikimrQueryType ConvertType(NKikimrKqp::EQueryType type); diff --git a/ydb/core/kqp/ut/view/view_ut.cpp b/ydb/core/kqp/ut/view/view_ut.cpp index fa378963364..1ee9c807329 100644 --- a/ydb/core/kqp/ut/view/view_ut.cpp +++ b/ydb/core/kqp/ut/view/view_ut.cpp @@ -204,6 +204,42 @@ Y_UNIT_TEST_SUITE(TCreateAndDropViewTest) { UNIT_ASSERT_STRING_CONTAINS(creationResult.GetIssues().ToString(), "Error: Cannot divide type String and String"); } + void InvalidRef(const char* createTableQuery, const char* selectTableQuery) { + TKikimrRunner kikimr(TKikimrSettings().SetWithSampleTables(false)); + auto session = kikimr.GetQueryClient().GetSession().ExtractValueSync().GetSession(); + + { + auto result = session.ExecuteQuery(createTableQuery, NQuery::TTxControl::NoTx()).ExtractValueSync(); + UNIT_ASSERT_C(result.IsSuccess(), "Failed to CREATE TABLE: " << result.GetIssues().ToOneLineString()); + } + { + auto result = session.ExecuteQuery(std::format("CREATE VIEW `View` WITH (security_invoker = true) AS {}", selectTableQuery), + NQuery::TTxControl::NoTx()).ExtractValueSync(); + UNIT_ASSERT_C(!result.IsSuccess(), "CREATE VIEW executed successfully"); + } + } + + Y_UNIT_TEST(InvalidTableRef) { + InvalidRef( + "CREATE TABLE `Table` (key Int32, PRIMARY KEY (key))", + "SELECT `key` FROM `MissingTable`" + ); + } + + Y_UNIT_TEST(InvalidColumnRef) { + InvalidRef( + "CREATE TABLE `Table` (key Int32, PRIMARY KEY (key))", + "SELECT `MissingColumn` FROM `Table`" + ); + } + + Y_UNIT_TEST(InvalidIndexRef) { + InvalidRef( + "CREATE TABLE `Table` (key Int32, value Int32, PRIMARY KEY (key), INDEX `Index` GLOBAL ON (value))", + "SELECT `value` FROM `Table` VIEW `MissingIndex`" + ); + } + Y_UNIT_TEST(ParsingSecurityInvoker) { TKikimrRunner kikimr(TKikimrSettings().SetWithSampleTables(false)); auto session = kikimr.GetQueryClient().GetSession().ExtractValueSync().GetSession(); diff --git a/ydb/core/tx/schemeshard/ut_restore/ut_restore.cpp b/ydb/core/tx/schemeshard/ut_restore/ut_restore.cpp index 6115a556bd0..ce9e84eb9b0 100644 --- a/ydb/core/tx/schemeshard/ut_restore/ut_restore.cpp +++ b/ydb/core/tx/schemeshard/ut_restore/ut_restore.cpp @@ -11,6 +11,7 @@ #include #include #include +#include #include #include #include @@ -5622,7 +5623,64 @@ Y_UNIT_TEST_SUITE(TImportTests) { TestGetImport(runtime, txId, "/MyRoot"); } - Y_UNIT_TEST(ViewCreationRetry) { + Y_UNIT_TEST(ShouldImportInvalidView) { + TTestBasicRuntime runtime; + auto options = TTestEnvOptions() + .RunFakeConfigDispatcher(true) + .SetupKqpProxy(true); + TTestEnv env(runtime, options); + ui64 txId = 100; + + THashMap bucketContent(2); + bucketContent.emplace("/table", GenerateTestData(R"( + columns { + name: "key" + type { optional_type { item { type_id: UTF8 } } } + } + primary_key: "key" + )")); + bucketContent.emplace("/view", GenerateTestData( + { + EPathTypeView, + R"( + -- backup root: "/MyRoot" + CREATE VIEW IF NOT EXISTS `view` WITH security_invoker = TRUE AS SELECT missing FROM `table`; + )" + } + )); + + TPortManager portManager; + const ui16 port = portManager.GetPort(); + + TS3Mock s3Mock(ConvertTestData(bucketContent), TS3Mock::TSettings(port)); + UNIT_ASSERT(s3Mock.Start()); + + const ui64 importId = ++txId; + TestImport(runtime, importId, "/MyRoot", Sprintf(R"( + ImportFromS3Settings { + endpoint: "localhost:%d" + scheme: HTTP + items { + source_prefix: "table" + destination_path: "/MyRoot/table" + } + items { + source_prefix: "view" + destination_path: "/MyRoot/view" + } + } + )", port)); + + env.TestWaitNotification(runtime, importId); + TestGetImport(runtime, importId, "/MyRoot"); + + TestDescribeResult(DescribePath(runtime, "/MyRoot/view"), { + NLs::Finished, + NLs::IsView + }); + } + + Y_UNIT_TEST(ShouldNotRetryViewCreation) { TTestBasicRuntime runtime; auto options = TTestEnvOptions() .RunFakeConfigDispatcher(true) @@ -5660,33 +5718,30 @@ Y_UNIT_TEST_SUITE(TImportTests) { UNIT_ASSERT(s3Mock.Start()); ui64 tableCreationTxId = 0; - TActorId schemeshardActorId; - TBlockEvents tableCreationBlocker(runtime, - [&](const TEvSchemeShard::TEvModifySchemeTransaction::TPtr& event) { - const auto& record = event->Get()->Record; - if (record.GetTransaction(0).GetOperationType() == ESchemeOpCreateIndexedTable) { - tableCreationTxId = record.GetTxId(); - schemeshardActorId = event->Recipient; - return true; - } + ui64 viewCreationTxId = 0; + TBlockEvents tableCreationBlocker(runtime, [&](auto& ev) { + const auto& record = ev->Get()->Record; + switch (record.GetTransaction(0).GetOperationType()) { + case ESchemeOpCreateIndexedTable: + tableCreationTxId = record.GetTxId(); + return true; + case ESchemeOpCreateView: + viewCreationTxId = record.GetTxId(); + return false; + default: return false; } - ); + }); - TBlockEvents queryResultBlocker(runtime, - [&](const TEvPrivate::TEvImportSchemeQueryResult::TPtr& event) { - // The test expects the SchemeShard actor ID to be already initialized when we receive the first query result message. - // This expectation is valid because we import items in order of their appearance on the import items list. - if (!schemeshardActorId || event->Recipient != schemeshardActorId) { - return false; - } - UNIT_ASSERT_VALUES_EQUAL(event->Get()->Status, Ydb::StatusIds::SCHEME_ERROR); - const auto* error = std::get_if(&event->Get()->Result); - UNIT_ASSERT(error); - UNIT_ASSERT_STRING_CONTAINS(*error, "Cannot find table"); - return true; + TWaitForFirstEvent viewCreationAccepter(runtime, [&](auto& ev) { + const auto& record = ev->Get()->Record; + if (!viewCreationTxId || viewCreationTxId != record.GetTxId()) { + return false; } - ); + + UNIT_ASSERT_VALUES_EQUAL(record.GetStatus(), NKikimrScheme::StatusAccepted); + return true; + }); const ui64 importId = ++txId; TestImport(runtime, importId, "/MyRoot", Sprintf(R"( @@ -5705,21 +5760,25 @@ Y_UNIT_TEST_SUITE(TImportTests) { )", port)); runtime.WaitFor("table creation attempt", [&]{ return !tableCreationBlocker.empty(); }); - runtime.WaitFor("query result", [&]{ return !queryResultBlocker.empty(); }); + viewCreationAccepter.Wait(); + env.TestWaitNotification(runtime, viewCreationTxId); + TestDescribeResult(DescribePath(runtime, "/MyRoot/view"), { + NLs::Finished, + NLs::IsView, + }); + tableCreationBlocker.Unblock().Stop(); - queryResultBlocker.Unblock().Stop(); env.TestWaitNotification(runtime, tableCreationTxId); + TestDescribeResult(DescribePath(runtime, "/MyRoot/table"), { + NLs::Finished, + NLs::IsTable, + }); env.TestWaitNotification(runtime, importId); TestGetImport(runtime, importId, "/MyRoot"); - - TestDescribeResult(DescribePath(runtime, "/MyRoot/view"), { - NLs::Finished, - NLs::IsView - }); } - Y_UNIT_TEST(MultipleViewCreationRetries) { + Y_UNIT_TEST(ShouldImportMultipleViews) { TTestBasicRuntime runtime; auto options = TTestEnvOptions() .RunFakeConfigDispatcher(true) @@ -5751,40 +5810,13 @@ Y_UNIT_TEST_SUITE(TImportTests) { TS3Mock s3Mock(ConvertTestData(bucketContent), TS3Mock::TSettings(port)); UNIT_ASSERT(s3Mock.Start()); - TActorId schemeshardActorId; - TBlockEvents viewCreationBlocker(runtime, - [&](const TEvSchemeShard::TEvModifySchemeTransaction::TPtr& event) { - const auto& record = event->Get()->Record; - if (record.GetTransaction(0).GetOperationType() == ESchemeOpCreateView) { - schemeshardActorId = event->Recipient; - return true; - } - return false; - } - ); - - int missingDependencyFails = 0; - auto missingDependencyObserver = runtime.AddObserver( - [&](const TEvPrivate::TEvImportSchemeQueryResult::TPtr& event) { - if (!schemeshardActorId - || event->Recipient != schemeshardActorId - || event->Get()->Status != Ydb::StatusIds::SCHEME_ERROR) { - return; - } - const auto* error = std::get_if(&event->Get()->Result); - if (error && error->Contains("Cannot find table")) { - ++missingDependencyFails; - } - } - ); - auto importSettings = TStringBuilder() << std::format(R"( ImportFromS3Settings {{ endpoint: "localhost:{}" scheme: HTTP )", port ); - for (int i = 0; i < ViewLayers; ++i) { + for (int i = ViewLayers - 1; i >= 0; --i) { importSettings << std::format(R"( items {{ source_prefix: "view{}" @@ -5798,23 +5830,6 @@ Y_UNIT_TEST_SUITE(TImportTests) { const ui64 importId = ++txId; TestImport(runtime, importId, "/MyRoot", importSettings); - int expectedFails = 0; - for (int iteration = 1; iteration <= ViewLayers; ++iteration) { - runtime.WaitFor("blocked view creation", [&]{ return !viewCreationBlocker.empty(); }); - - expectedFails += ViewLayers - iteration; - if (iteration > 1) { - runtime.WaitFor("query results", [&]{ return missingDependencyFails >= expectedFails; }); - } else { - // the first iteration might miss some query results due to the initially unset schemeshardActorId - missingDependencyFails = expectedFails; - } - - viewCreationBlocker.Unblock(1); - } - UNIT_ASSERT(viewCreationBlocker.empty()); - viewCreationBlocker.Stop(); - env.TestWaitNotification(runtime, importId); TestGetImport(runtime, importId, "/MyRoot"); -- cgit v1.3