summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorIlnaz Nizametdinov <[email protected]>2026-05-28 09:31:02 +0300
committerGitHub <[email protected]>2026-05-28 09:31:02 +0300
commite55c625445f270656e786b2ff5d6abb3fcefcfce (patch)
tree9d185f8fa988b7a13dc6b72b372e5408618be79f
parent45347918b1621b661ce60bdad96639281f3b774e (diff)
Allow to create invalid VIEW during import (#41265)
-rw-r--r--ydb/core/kqp/host/kqp_host.cpp1
-rw-r--r--ydb/core/kqp/host/kqp_translate.cpp2
-rw-r--r--ydb/core/kqp/host/kqp_translate.h6
-rw-r--r--ydb/core/kqp/ut/view/view_ut.cpp36
-rw-r--r--ydb/core/tx/schemeshard/ut_restore/ut_restore.cpp169
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<NSQLTranslation::EYqlSelect> 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 <ydb/core/protos/schemeshard/operations.pb.h>
#include <ydb/core/tablet/resource_broker.h>
#include <ydb/core/testlib/actors/block_events.h>
+#include <ydb/core/testlib/actors/wait_events.h>
#include <ydb/core/testlib/audit_helpers/audit_helper.h>
#include <ydb/core/tx/datashard/datashard.h>
#include <ydb/core/tx/schemeshard/schemeshard_billing_helpers.h>
@@ -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<TString, TTestDataWithScheme> 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<TEvSchemeShard::TEvModifySchemeTransaction> 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<TEvSchemeShard::TEvModifySchemeTransaction> 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<TEvPrivate::TEvImportSchemeQueryResult> 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<TString>(&event->Get()->Result);
- UNIT_ASSERT(error);
- UNIT_ASSERT_STRING_CONTAINS(*error, "Cannot find table");
- return true;
+ TWaitForFirstEvent<TEvSchemeShard::TEvModifySchemeTransactionResult> 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<TEvSchemeShard::TEvModifySchemeTransaction> 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<TEvPrivate::TEvImportSchemeQueryResult>(
- [&](const TEvPrivate::TEvImportSchemeQueryResult::TPtr& event) {
- if (!schemeshardActorId
- || event->Recipient != schemeshardActorId
- || event->Get()->Status != Ydb::StatusIds::SCHEME_ERROR) {
- return;
- }
- const auto* error = std::get_if<TString>(&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");