diff --git a/CMakeLists.txt b/CMakeLists.txt index 9889f8bb..10e20533 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -78,7 +78,8 @@ set(EXTENSION_SOURCES ${LPTS_DIR}/src/lpts_ast_renderer.cpp ${LPTS_DIR}/src/lpts_ast_builder.cpp ${LPTS_DIR}/src/lpts_ast_flattener.cpp - ${LPTS_DIR}/src/dialect_function_map.cpp) + ${LPTS_DIR}/src/dialect_function_map.cpp + ${LPTS_DIR}/src/spark_scalar_functions.cpp) build_static_extension(${TARGET_NAME} ${EXTENSION_SOURCES}) build_loadable_extension(${TARGET_NAME} " " ${EXTENSION_SOURCES}) diff --git a/src/delta/operators/join.cpp b/src/delta/operators/join.cpp index 428a9cc9..ec33edda 100644 --- a/src/delta/operators/join.cpp +++ b/src/delta/operators/join.cpp @@ -834,20 +834,61 @@ void AppendMultiplicityToAncestorProjectionMaps(unique_ptr &ter if (mul_idx == DConstants::INVALID_INDEX) { continue; } + // Appending to this join's own left_projection_map grows its LEFT + // contribution width, which shifts the absolute position where its RIGHT + // contribution starts within its own combined GetColumnBindings(). A + // grandparent ancestor may already have a projection map entry referencing + // (by that now-stale absolute position) a column from this join's right + // side; left as-is, that entry would silently start pointing at the + // newly-inserted column instead, dropping the real column it used to + // select. Appending to right_projection_map never has this effect: right + // contributions are always placed last, so a new entry there only ever + // extends the combined output with a brand-new highest index. + idx_t old_width = proj_map.size(); + auto shift_stale_parent_indexes = [&](idx_t added) { + if (child_side != 0 || added == 0 || depth == 0) { + return; + } + size_t parent_side = leaf_path[depth - 1]; + auto *parent_join = dynamic_cast(ancestors[depth - 1]); + if (!parent_join || parent_side >= parent_join->children.size()) { + return; + } + auto &parent_map = + (parent_side == 0) ? parent_join->left_projection_map : parent_join->right_projection_map; + for (auto &parent_idx : parent_map) { + if (parent_idx >= old_width) { + parent_idx += added; + } + } + }; if (preserve_full_child) { idx_t projectable_count = MinValue(mul_idx + 1, child_bindings.size()); + idx_t added = 0; for (idx_t binding_idx = 0; binding_idx < projectable_count; binding_idx++) { if (std::find(proj_map.begin(), proj_map.end(), binding_idx) != proj_map.end()) { continue; } proj_map.push_back(binding_idx); + added++; OPENIVM_DEBUG_PRINT("[%s] Preserved child col %lu in immediate %s proj_map\n", context_label, (unsigned long)binding_idx, child_side == 0 ? "left" : "right"); } - } else if (std::find(proj_map.begin(), proj_map.end(), mul_idx) == proj_map.end()) { - proj_map.push_back(mul_idx); - OPENIVM_DEBUG_PRINT("[%s] Added mul col %lu to ancestor %s proj_map\n", context_label, - (unsigned long)mul_idx, child_side == 0 ? "left" : "right"); + shift_stale_parent_indexes(added); + } else { + // proj_map entries are positions into the child's *current* combined + // GetColumnBindings(). Testing raw index membership of mul_idx against + // proj_map can alias onto an unrelated pre-existing entry that now shares + // the same numeric position after a deeper level's own map grew. Compare + // by column identity against what this ancestor currently exposes instead + // of trusting the raw index. + auto exposed = join->GetColumnBindings(); + if (std::find(exposed.begin(), exposed.end(), mul_binding) == exposed.end()) { + proj_map.push_back(mul_idx); + shift_stale_parent_indexes(1); + OPENIVM_DEBUG_PRINT("[%s] Added mul col %lu to ancestor %s proj_map\n", context_label, + (unsigned long)mul_idx, child_side == 0 ? "left" : "right"); + } } join->ResolveOperatorTypes(); } @@ -1497,9 +1538,20 @@ BuildInclusionExclusionTerms(DeltaOperatorInput input, ClientContext &context, B for (size_t i = 0; i < N; i++) { if (mask & (1ULL << i)) { if (leaves[i].get) { - DeltaGetResult delta_i = CreateDeltaGetNode(context, binder, leaves[i].get, input.context.view); + // leaves[] was collected once on the ORIGINAL input.plan, before this + // mask's own renumber_and_rebind_subtree pass, so leaves[i].get is a + // stale pointer carrying the pre-renumbering table_index. The rest of + // `term` (join conditions, transitioning-key guards, etc.) was rebound + // to the FRESH per-term index, so the replacement delta node -- which + // reuses old_get->table_index verbatim -- must be built from term's own, + // already-renumbered GET at this leaf's (renumbering-invariant) path, + // not from leaves[i].get, or every reference elsewhere in `term` to this + // leaf's fresh index is left dangling. + auto &leaf_node_ref = GetNodeAtPath(term, leaves[i].path); + auto &term_local_get = leaf_node_ref->Cast(); + DeltaGetResult delta_i = CreateDeltaGetNode(context, binder, &term_local_get, input.context.view); mul_bindings.push_back(delta_i.mul_binding); - GetNodeAtPath(term, leaves[i].path) = std::move(delta_i.node); + leaf_node_ref = std::move(delta_i.node); UpdateParentProjectionMap(term, leaves[i], delta_i.mul_binding); } else { auto &subtree_ref = GetNodeAtPath(term, leaves[i].path); @@ -1616,6 +1668,22 @@ static bool HasOnlyInnerJoins(LogicalOperator *node) { return true; } +static bool HasOnlyInnerOrLeftJoins(LogicalOperator *node) { + if (node->type == LogicalOperatorType::LOGICAL_COMPARISON_JOIN || + node->type == LogicalOperatorType::LOGICAL_ANY_JOIN) { + auto *join = dynamic_cast(node); + if (!join || (join->join_type != JoinType::INNER && join->join_type != JoinType::LEFT)) { + return false; + } + } + for (auto &child : node->children) { + if (!HasOnlyInnerOrLeftJoins(child.get())) { + return false; + } + } + return true; +} + static bool SupportsRegularNtermLeaf(const JoinLeafInfo &leaf) { if (leaf.get) { return leaf.get->GetTable().get() != nullptr; @@ -1705,7 +1773,7 @@ static DeltaPlanFragment CompileRegularLeafDelta(const DeltaOperatorInput &input static vector> BuildRegularJoinTerms(DeltaOperatorInput input, ClientContext &context, Binder &binder, const vector &leaves, - uint64_t unchanged_mask) { + uint64_t unchanged_mask, bool has_left_join) { vector> terms; // Base scans see post-DML state. Term i uses current state before i, delta i, and reconstructs old state after i as // current - delta. These disjoint telescoping terms cover every non-empty delta combination exactly once. @@ -1735,6 +1803,16 @@ static vector> BuildRegularJoinTerms(DeltaOperatorIn LogicalOperator *term_root = term.get(); CollectJoinLeaves(term.get(), {}, term_leaves); D_ASSERT(term_leaves.size() == leaves.size()); + + // LEFT-JOIN telescoping: demote only the outer join(s) whose NULL-supplying + // subtree contains this term's single delta leaf, mirroring the DuckLake + // N-term path (DemoteLeftJoinsForMask). Preserved outer joins elsewhere keep + // their NULL-padded rows; the upsert layer's key-based partial recompute + // (BuildLeftJoinProjectionRefresh) fixes NULL<->match transition rows. + if (has_left_join) { + DemoteLeftJoinsForMask(term.get(), term_leaves, (1ULL << delta_leaf)); + } + vector mul_bindings; for (size_t leaf = 0; leaf < term_leaves.size(); leaf++) { @@ -1880,11 +1958,22 @@ DeltaPlanFragment CompileJoinDelta(DeltaOperatorInput input) { } auto compile_facts = openivm::CompileFactsContextSlot::Get(context); auto unchanged_mask = ComputeFactsUnchangedMask(compile_facts, leaves); - bool regular_nterm = !all_ducklake && compile_facts.compile_only && !has_left_join && - input.context.model.type == RefreshType::SIMPLE_PROJECTION && - HasOnlyInnerJoins(input.plan.get()) && - RegularNtermPreservesFKPruning(context, compile_facts, leaves, input.plan.get()) && - SqlUtils::GetBoolSetting(context, "openivm_regular_nterm", true); + bool regular_nterm_base = !all_ducklake && compile_facts.compile_only && + input.context.model.type == RefreshType::SIMPLE_PROJECTION && + SqlUtils::GetBoolSetting(context, "openivm_regular_nterm", true); + bool regular_nterm; + if (has_left_join) { + // LEFT-JOIN telescoping: the regular N-term delta extends to LEFT joins via + // per-term demotion of only the outer join whose NULL-supplying side carries + // that term's delta (see BuildRegularJoinTerms). FULL OUTER / RIGHT shapes and + // the inclusion-exclusion FK-pruning path are out of scope; NULL-padded row + // correctness is completed by BuildLeftJoinProjectionRefresh in the upsert layer. + regular_nterm = regular_nterm_base && HasOnlyInnerOrLeftJoins(input.plan.get()) && + SqlUtils::GetBoolSetting(context, "openivm_regular_nterm_left", true); + } else { + regular_nterm = regular_nterm_base && HasOnlyInnerJoins(input.plan.get()) && + RegularNtermPreservesFKPruning(context, compile_facts, leaves, input.plan.get()); + } if (regular_nterm) { for (auto &leaf : leaves) { if (!SupportsRegularNtermLeaf(leaf)) { @@ -1902,7 +1991,7 @@ DeltaPlanFragment CompileJoinDelta(DeltaOperatorInput input) { if (all_ducklake) { terms = BuildDuckLakeJoinTerms(input, context, binder, leaves, has_left_join, flattened_ducklake); } else if (regular_nterm) { - terms = BuildRegularJoinTerms(input, context, binder, leaves, unchanged_mask); + terms = BuildRegularJoinTerms(input, context, binder, leaves, unchanged_mask, has_left_join); } else { terms = BuildInclusionExclusionTerms(input, context, binder, leaves, has_left_join, transition_ctes); } diff --git a/src/include/upsert/refresh_compiler.hpp b/src/include/upsert/refresh_compiler.hpp index ee06f309..76aeaa0c 100644 --- a/src/include/upsert/refresh_compiler.hpp +++ b/src/include/upsert/refresh_compiler.hpp @@ -58,6 +58,28 @@ string CompileWindowRecompute(const string &view_name, const string &view_query_ const vector &column_names = {}, bool running_window_incremental = false); string CompileFullRecompute(const string &view_name, const string &view_query_sql, const string &catalog_prefix = ""); +/// Full recompute of `openivm_data_` that ALSO emits the exact signed +/// multiset view-delta into `openivm_delta_`. +/// +/// The partial-recompute paths (`WINDOW_PARTITION`, `GROUP_RECOMPUTE`) degrade to +/// a full recompute whenever the affected partition/group key set cannot be +/// scoped from the source deltas — an unpartitioned surrogate-key +/// `ROW_NUMBER() OVER (ORDER BY ...)`, a partition key that is a computed +/// expression absent from every source delta table, or incomplete multi-source +/// lineage. The view keeps its `WINDOW_PARTITION` / `GROUP_RECOMPUTE` +/// classification, so a caller that asked for a cascade delta +/// (`CompileFacts::force_view_delta_cascade`) would otherwise receive a program +/// that writes no `openivm_delta_` rows at all, and every downstream MV +/// would have to be demoted to a full refresh. +/// +/// Retracting the whole pre-refresh content at multiplicity -1 and adding the +/// whole post-refresh content at +1 is the exact Z-set delta of the view +/// (`new_bag - old_bag`): unchanged rows contribute cancelling -1/+1 pairs, so +/// bag semantics are preserved exactly. Statement shapes mirror the +/// `CompileWindowRecompute` / `CompileGroupRecompute` cascade branches. +string CompileFullRecomputeWithCascadeDelta(const string &view_name, const string &view_query_sql, + const string &catalog_prefix = ""); + /// Group-level partial recompute, used by `RefreshType::GROUP_RECOMPUTE` /// (inner-DISTINCT under aggregate). For each base table T_i with a non-empty /// delta, builds a "view query with T_i restricted to its delta" variant by diff --git a/src/openivm_extension.cpp b/src/openivm_extension.cpp index ba49b33d..3c347ed6 100644 --- a/src/openivm_extension.cpp +++ b/src/openivm_extension.cpp @@ -2,6 +2,7 @@ #include "core/openivm_extension.hpp" #include "compile_facts.hpp" +#include "spark_scalar_functions.hpp" #include "core/openivm_constants.hpp" #include "core/refresh_metadata.hpp" #include "core/refresh_daemon.hpp" @@ -166,6 +167,8 @@ static void LoadInternal(ExtensionLoader &loader) { // statement OpenIVM does not recognize with DuckDB's native parser. db_config.SetOption(AllowParserOverrideExtensionSetting::SettingIndex, Value("fallback")); + RegisterSparkScalarFunctions(loader); + db_config.AddExtensionOption("openivm_files_path", "path for compiled SQL reference files", LogicalType::VARCHAR); db_config.AddExtensionOption("openivm_refresh_mode", "refresh strategy: incremental, full, or auto", LogicalType::VARCHAR, Value("incremental")); @@ -191,6 +194,9 @@ static void LoadInternal(ExtensionLoader &loader) { LogicalType::BOOLEAN, Value::BOOLEAN(true)); db_config.AddExtensionOption("openivm_regular_nterm", "use N-term telescoping for compile-only regular inner joins", LogicalType::BOOLEAN, Value::BOOLEAN(true)); + db_config.AddExtensionOption("openivm_regular_nterm_left", + "extend compile-only N-term telescoping to LEFT-join projection views", + LogicalType::BOOLEAN, Value::BOOLEAN(true)); db_config.AddExtensionOption("openivm_fk_pruning", "prune inclusion-exclusion join terms using FK constraints", LogicalType::BOOLEAN, Value::BOOLEAN(true)); db_config.AddExtensionOption("openivm_scd2_range_join_accel", diff --git a/src/upsert/refresh_compiler.cpp b/src/upsert/refresh_compiler.cpp index 66ecb0db..3b2b1f37 100644 --- a/src/upsert/refresh_compiler.cpp +++ b/src/upsert/refresh_compiler.cpp @@ -1455,6 +1455,28 @@ string CompileFullRecompute(const string &view_name, const string &view_query_sq return SqlUtils::BuildFullRecomputeSQL(data_table, view_query_sql); } +string CompileFullRecomputeWithCascadeDelta(const string &view_name, const string &view_query_sql, + const string &catalog_prefix) { + string data_table = catalog_prefix + SqlUtils::QuoteIdentifier(IncrementalTableNames::DataTableName(view_name)); + string delta_table = catalog_prefix + SqlUtils::QuoteIdentifier(SqlUtils::DeltaName(view_name)); + string old_temp_table = SqlUtils::QuoteIdentifier(string(openivm::TEMP_TABLE_PREFIX) + view_name); + string new_temp_table = SqlUtils::QuoteIdentifier(string("openivm_new_") + view_name); + + string sql; + sql += "CREATE OR REPLACE TEMP TABLE " + old_temp_table + " AS\nSELECT * FROM " + data_table + " openivm_old;\n\n"; + sql += "CREATE OR REPLACE TEMP TABLE " + new_temp_table + " AS\nSELECT * FROM (" + view_query_sql + + ") openivm_recompute;\n\n"; + sql += "DELETE FROM " + data_table + ";\n"; + sql += "INSERT INTO " + data_table + "\nSELECT * FROM " + new_temp_table + ";\n"; + sql += "\n" + BuildSignedMultisetDeltaInsertSQL(delta_table, old_temp_table, new_temp_table); + sql += "DROP TABLE IF EXISTS " + old_temp_table + ";\n"; + sql += "DROP TABLE IF EXISTS " + new_temp_table + ";\n"; + OPENIVM_DEBUG_PRINT("[CompileFullRecomputeWithCascadeDelta] unscopable recompute for '%s' — emitting signed " + "whole-view cascade delta\n", + view_name.c_str()); + return sql; +} + string CompileGroupRecompute(const string &view_name, const string &view_query_sql, const vector &group_columns, const vector &delta_table_specs, const string &catalog_prefix, const string &lpts_table_prefix, bool emit_cascade_delta, @@ -1462,8 +1484,11 @@ string CompileGroupRecompute(const string &view_name, const string &view_query_s string data_table = catalog_prefix + SqlUtils::QuoteIdentifier(IncrementalTableNames::DataTableName(view_name)); // No GROUP BY columns or no source deltas registered → can't scope; fall back to full. + // A cascade delta was still requested, so emit the whole-view signed delta rather than + // silently producing a program with no `openivm_delta_` rows. if (group_columns.empty() || delta_table_specs.empty()) { - return CompileFullRecompute(view_name, view_query_sql, catalog_prefix); + return emit_cascade_delta ? CompileFullRecomputeWithCascadeDelta(view_name, view_query_sql, catalog_prefix) + : CompileFullRecompute(view_name, view_query_sql, catalog_prefix); } string group_csv = SqlUtils::JoinQuotedColumns(group_columns); diff --git a/src/upsert/refresh_compiler_aux.cpp b/src/upsert/refresh_compiler_aux.cpp index a1fe3c0d..fb79b334 100644 --- a/src/upsert/refresh_compiler_aux.cpp +++ b/src/upsert/refresh_compiler_aux.cpp @@ -1600,7 +1600,11 @@ string CompileWindowRecompute(const string &view_name, const string &view_query_ bool running_window_incremental) { bool have_affected_keys = !affected_keys_sql.empty(); if (!have_affected_keys && (partition_columns.empty() || partition_delta_specs.empty())) { - return CompileFullRecompute(view_name, view_query_sql, catalog_prefix); + // No PARTITION BY (global surrogate-key window) or no partition key resolvable in any + // source delta table → nothing to scope the recompute to. Keep the cascade delta the + // caller asked for so downstream MVs stay incremental. + return emit_cascade_delta ? CompileFullRecomputeWithCascadeDelta(view_name, view_query_sql, catalog_prefix) + : CompileFullRecompute(view_name, view_query_sql, catalog_prefix); } if (running_window_incremental) { auto suffix_sql = BuildRunningWindowSuffixRefreshSQL(view_name, view_query_sql, delta_ts_filter, catalog_prefix, diff --git a/src/upsert/refresh_sql.cpp b/src/upsert/refresh_sql.cpp index 26026b3d..43cbeb5d 100644 --- a/src/upsert/refresh_sql.cpp +++ b/src/upsert/refresh_sql.cpp @@ -532,9 +532,9 @@ static void PropagateRefreshPlanningSettings(ClientContext &from, ClientContext // session-scoped planning settings still need to be mirrored onto the fresh // planning connection. static const char *PLANNING_SETTINGS[] = { - "openivm_adaptive_refresh", "openivm_cost_decay", "openivm_skip_empty_deltas", - "openivm_fk_pruning", "openivm_ducklake_nterm", "openivm_scd2_range_join_accel", - "openivm_regular_nterm", + "openivm_adaptive_refresh", "openivm_cost_decay", "openivm_skip_empty_deltas", + "openivm_fk_pruning", "openivm_ducklake_nterm", "openivm_scd2_range_join_accel", + "openivm_regular_nterm", "openivm_regular_nterm_left", }; for (auto setting_name : PLANNING_SETTINGS) { CopyOpenIvmSetting(from, to, setting_name); diff --git a/src/upsert/refresh_window.cpp b/src/upsert/refresh_window.cpp index e77b4807..82f0f2cf 100644 --- a/src/upsert/refresh_window.cpp +++ b/src/upsert/refresh_window.cpp @@ -577,6 +577,9 @@ string BuildWindowPartitionRefresh(RefreshMetadata &metadata, Connection &con, c if (any_ducklake) { OPENIVM_DEBUG_PRINT( "[UPSERT] Compiling upsert for type: WINDOW_PARTITION (DuckLake, full recompute fallback)\n"); + if (emit_cascade_delta) { + return CompileFullRecomputeWithCascadeDelta(view_name, view_query_sql, internal_catalog_prefix); + } return "DELETE FROM " + data_table + ";\n" + "INSERT INTO " + data_table + " " + view_query_sql + ";\n"; } auto lineage_result = BuildLineageStandardAffectedKeysSQL( @@ -585,7 +588,9 @@ string BuildWindowPartitionRefresh(RefreshMetadata &metadata, Connection &con, c if (lineage_result == LineageAffectedKeysResult::UNSAFE) { OPENIVM_DEBUG_PRINT("[UPSERT] WINDOW_PARTITION lineage is unsafe for '%s' — full recompute fallback\n", view_name.c_str()); - return CompileFullRecompute(view_name, view_query_sql, internal_catalog_prefix); + return emit_cascade_delta + ? CompileFullRecomputeWithCascadeDelta(view_name, view_query_sql, internal_catalog_prefix) + : CompileFullRecompute(view_name, view_query_sql, internal_catalog_prefix); } have_lineage_affected_keys = lineage_result == LineageAffectedKeysResult::AVAILABLE; if (!have_lineage_affected_keys && delta_table_names.size() > 1 && @@ -593,7 +598,9 @@ string BuildWindowPartitionRefresh(RefreshMetadata &metadata, Connection &con, c OPENIVM_DEBUG_PRINT("[UPSERT] WINDOW_PARTITION lineage incomplete for '%s' (%zu sources) — full recompute " "fallback\n", view_name.c_str(), delta_table_names.size()); - return CompileFullRecompute(view_name, view_query_sql, internal_catalog_prefix); + return emit_cascade_delta + ? CompileFullRecomputeWithCascadeDelta(view_name, view_query_sql, internal_catalog_prefix) + : CompileFullRecompute(view_name, view_query_sql, internal_catalog_prefix); } OPENIVM_DEBUG_PRINT("[UPSERT] Compiling upsert for type: WINDOW_PARTITION (%zu partition cols, lineage keys: %s)\n", partition_cols.size(), have_lineage_affected_keys ? "yes" : "no"); diff --git a/test/sql/cascade_simple_projection_join.test b/test/sql/cascade_simple_projection_join.test index 8de65d6c..147e136b 100644 --- a/test/sql/cascade_simple_projection_join.test +++ b/test/sql/cascade_simple_projection_join.test @@ -565,3 +565,345 @@ SELECT COUNT(*) FROM ( ); ---- 0 + +# Reduced "arc_machine_status_transaction" shape: a CTE (INNER JOIN + LEFT +# JOIN) feeding an outer query with two further chained LEFT JOINs. Before +# src/delta/operators/join.cpp's AppendMultiplicityToAncestorProjectionMaps +# fix, a deeper join's own left_projection_map growing (to append a new +# per-leaf multiplicity column) silently shifted where its right-side +# contribution begins in its own combined column numbering. A grandparent +# join's pre-existing projection-map entry, fixed before that growth +# happened, could then alias onto the just-added multiplicity column +# instead of the real business column (`cust_key`) it used to select, +# permanently dropping it and surfacing as an internal arity mismatch. +statement ok +CREATE TABLE cspj_arcm_msf(id INT, coll_key INT); + +statement ok +CREATE TABLE cspj_arcm_acd(coll_key INT, sub_id INT, arm_id VARCHAR); + +statement ok +CREATE TABLE cspj_arcm_hw(arm_id VARCHAR, provider VARCHAR); + +statement ok +CREATE TABLE cspj_arcm_tpid(sub_id INT, tp_id INT); + +statement ok +CREATE TABLE cspj_arcm_cust(tp_id INT, cust_key INT); + +statement ok +INSERT INTO cspj_arcm_msf VALUES (1, 100), (2, 100); + +statement ok +INSERT INTO cspj_arcm_acd VALUES (100, 10, 'arm-1'); + +statement ok +INSERT INTO cspj_arcm_hw VALUES ('arm-1', 'azure'); + +statement ok +INSERT INTO cspj_arcm_tpid VALUES (10, 200); + +statement ok +INSERT INTO cspj_arcm_cust VALUES (200, 21); + +statement ok +CREATE MATERIALIZED VIEW cspj_arcm_mv AS + WITH bff AS ( + SELECT msf.id, acd.sub_id, + coalesce(hw.provider, 'N/A') AS provider + FROM cspj_arcm_msf msf + INNER JOIN cspj_arcm_acd acd ON msf.coll_key = acd.coll_key + LEFT JOIN cspj_arcm_hw hw ON acd.arm_id = hw.arm_id + ) + SELECT bff.id, + coalesce(cust.cust_key, 1) AS cust_key + FROM bff + LEFT JOIN cspj_arcm_tpid tpid ON bff.sub_id = tpid.sub_id + LEFT JOIN cspj_arcm_cust cust ON coalesce(tpid.tp_id, -1) = cust.tp_id; + +statement ok +SELECT COUNT(*) FROM openivm_compile_with_facts( + 'cspj_arcm_mv', + '{"target_dialect":"duckdb","compile_only":true}' +); + +query I +SELECT CASE + WHEN contains(content, 'INSERT INTO openivm_delta_cspj_arcm_mv') AND + contains(content, 'openivm_delta_cspj_arcm_msf') AND + contains(content, 'openivm_delta_cspj_arcm_acd') AND + contains(content, 'openivm_delta_cspj_arcm_hw') AND + contains(content, 'openivm_delta_cspj_arcm_tpid') AND + contains(content, 'openivm_delta_cspj_arcm_cust') AND + NOT contains(content, 'SELECT NULL::') + THEN 1 ELSE 0 + END +FROM read_text('__TEST_DIR__/openivm_upsert_queries_cspj_arcm_mv.sql'); +---- +1 + +statement ok +INSERT INTO cspj_arcm_msf VALUES (3, 300); + +statement ok +INSERT INTO cspj_arcm_acd VALUES (300, 30, 'arm-2'); + +statement ok +INSERT INTO cspj_arcm_hw VALUES ('arm-2', 'gcp'); + +statement ok +INSERT INTO cspj_arcm_tpid VALUES (30, 400); + +statement ok +DELETE FROM cspj_arcm_msf WHERE id = 2; + +statement ok +UPDATE cspj_arcm_cust SET cust_key = 99 WHERE tp_id = 200; + +statement ok +PRAGMA refresh('cspj_arcm_mv'); + +query I +SELECT COUNT(*) FROM ( + SELECT * FROM cspj_arcm_mv + EXCEPT ALL + SELECT bff.id, coalesce(cust.cust_key, 1) AS cust_key + FROM ( + SELECT msf.id, acd.sub_id, coalesce(hw.provider, 'N/A') AS provider + FROM cspj_arcm_msf msf + INNER JOIN cspj_arcm_acd acd ON msf.coll_key = acd.coll_key + LEFT JOIN cspj_arcm_hw hw ON acd.arm_id = hw.arm_id + ) bff + LEFT JOIN cspj_arcm_tpid tpid ON bff.sub_id = tpid.sub_id + LEFT JOIN cspj_arcm_cust cust ON coalesce(tpid.tp_id, -1) = cust.tp_id +); +---- +0 + +query I +SELECT COUNT(*) FROM ( + SELECT bff.id, coalesce(cust.cust_key, 1) AS cust_key + FROM ( + SELECT msf.id, acd.sub_id, coalesce(hw.provider, 'N/A') AS provider + FROM cspj_arcm_msf msf + INNER JOIN cspj_arcm_acd acd ON msf.coll_key = acd.coll_key + LEFT JOIN cspj_arcm_hw hw ON acd.arm_id = hw.arm_id + ) bff + LEFT JOIN cspj_arcm_tpid tpid ON bff.sub_id = tpid.sub_id + LEFT JOIN cspj_arcm_cust cust ON coalesce(tpid.tp_id, -1) = cust.tp_id + EXCEPT ALL + SELECT * FROM cspj_arcm_mv +); +---- +0 + +# Reduced "int_instance_status_transaction" shape: a UNION ALL of two LEFT +# JOIN branches partitioned by an IS NOT NULL / IS NULL predicate on the +# joined side. Before the `third_party/lpts` pin was advanced to commit +# 77ed5e5 (fork PR #18, "Fix set-op column binding remaps"), remapping +# duplicated UNION ALL key-projection output aliases could leave a stale +# reference to a child alias that no longer existed in the rewritten plan, +# surfacing as a runtime Binder Error ("Referenced column ... not found in +# FROM clause"). +statement ok +CREATE TABLE cspj_inti_isf(id INT, k INT); + +statement ok +CREATE TABLE cspj_inti_ml(k INT, res VARCHAR); + +statement ok +INSERT INTO cspj_inti_isf VALUES (1, 100), (2, 200); + +statement ok +INSERT INTO cspj_inti_ml VALUES (100, 'machine-1'); + +statement ok +CREATE MATERIALIZED VIEW cspj_inti_mv AS + SELECT isf.id, isf.k FROM cspj_inti_isf isf LEFT JOIN cspj_inti_ml ml ON isf.k = ml.k WHERE ml.res IS NOT NULL + UNION ALL + SELECT isf.id, isf.k FROM cspj_inti_isf isf LEFT JOIN cspj_inti_ml ml ON isf.k = ml.k WHERE ml.res IS NULL; + +statement ok +SELECT COUNT(*) FROM openivm_compile_with_facts( + 'cspj_inti_mv', + '{"target_dialect":"duckdb","compile_only":true}' +); + +query I +SELECT CASE + WHEN contains(content, 'INSERT INTO openivm_delta_cspj_inti_mv') AND + contains(content, 'openivm_delta_cspj_inti_isf') AND + contains(content, 'openivm_delta_cspj_inti_ml') AND + contains(content, 'UNION ALL') AND + NOT contains(content, 'SELECT NULL::') + THEN 1 ELSE 0 + END +FROM read_text('__TEST_DIR__/openivm_upsert_queries_cspj_inti_mv.sql'); +---- +1 + +statement ok +INSERT INTO cspj_inti_isf VALUES (3, 300), (4, 400); + +statement ok +INSERT INTO cspj_inti_ml VALUES (300, 'machine-2'), (400, NULL); + +statement ok +DELETE FROM cspj_inti_isf WHERE id = 2; + +statement ok +UPDATE cspj_inti_ml SET res = NULL WHERE k = 100; + +statement ok +PRAGMA refresh('cspj_inti_mv'); + +query I +SELECT COUNT(*) FROM ( + SELECT * FROM cspj_inti_mv + EXCEPT ALL + ( + SELECT isf.id, isf.k FROM cspj_inti_isf isf LEFT JOIN cspj_inti_ml ml ON isf.k = ml.k WHERE ml.res IS NOT NULL + UNION ALL + SELECT isf.id, isf.k FROM cspj_inti_isf isf LEFT JOIN cspj_inti_ml ml ON isf.k = ml.k WHERE ml.res IS NULL + ) +); +---- +0 + +query I +SELECT COUNT(*) FROM ( + ( + SELECT isf.id, isf.k FROM cspj_inti_isf isf LEFT JOIN cspj_inti_ml ml ON isf.k = ml.k WHERE ml.res IS NOT NULL + UNION ALL + SELECT isf.id, isf.k FROM cspj_inti_isf isf LEFT JOIN cspj_inti_ml ml ON isf.k = ml.k WHERE ml.res IS NULL + ) + EXCEPT ALL + SELECT * FROM cspj_inti_mv +); +---- +0 + +# Reduced "machine-status left-deep UNION" shape: a literal three-way +# (left-deep) UNION ALL of LEFT JOIN branches, each partitioned by a mutually +# exclusive predicate on the joined side (active / inactive / unmatched). +# Internally this compiles to a nested union whose own union is itself one +# arm of a further union (`(A UNION ALL B) UNION ALL C`), carrying duplicated +# `openivm_left_key`/`openivm_multiplicity` output columns at every level of +# the chain. This matches the shape of upstream `third_party/lpts` commit +# 754c797's own `test/sql/union.test` regression ("OpenIVM join-delta UNION +# terms carry ... a duplicated hidden left key and multiplicity. A left-deep +# UNION must preserve and bind all ... positions"); the fix retains a +# trailing rewritten UNION binding as an alias of its physical multiplicity +# output rather than losing it. Compiled with force_view_delta_cascade under +# target_dialect=spark, this must classify SIMPLE_PROJECTION -- never +# COMPILE_FAILED/FULL_REFRESH -- and a real batched multi-table +# CREATE + INSERT/DELETE/UPDATE + PRAGMA refresh must stay in bidirectional +# bag-equality with the view definition. +statement ok +CREATE TABLE cspj_ms_terms(event_time INT, machine_arm_id VARCHAR, customer_key INT); + +statement ok +CREATE TABLE cspj_ms_status(machine_arm_id VARCHAR, status_label VARCHAR); + +statement ok +INSERT INTO cspj_ms_terms VALUES (10, 'arm-1', 1), (20, 'arm-2', 6); + +statement ok +INSERT INTO cspj_ms_status VALUES ('arm-1', 'active'); + +statement ok +CREATE MATERIALIZED VIEW cspj_ms_mv AS + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label = 'active' + UNION ALL + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label = 'inactive' + UNION ALL + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label IS NULL; + +statement ok +INSERT INTO cspj_ms_terms VALUES (30, 'arm-3', 11); + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'cspj_ms_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +2 SIMPLE_PROJECTION + +query I +SELECT CASE + WHEN contains(content, 'INSERT INTO openivm_delta_cspj_ms_mv') AND + contains(content, 'openivm_delta_cspj_ms_terms') AND + contains(content, 'openivm_delta_cspj_ms_status') AND + contains(content, 'UNION ALL') AND + NOT contains(content, 'SELECT NULL::') + THEN 1 ELSE 0 + END +FROM read_text('__TEST_DIR__/openivm_upsert_queries_cspj_ms_mv.sql'); +---- +1 + +statement ok +INSERT INTO cspj_ms_terms VALUES (40, 'arm-4', 16); + +statement ok +INSERT INTO cspj_ms_status VALUES ('arm-3', 'active'), ('arm-4', NULL); + +statement ok +DELETE FROM cspj_ms_terms WHERE event_time = 20; + +statement ok +UPDATE cspj_ms_status SET status_label = 'inactive' WHERE machine_arm_id = 'arm-1'; + +statement ok +PRAGMA refresh('cspj_ms_mv'); + +query I +SELECT COUNT(*) FROM ( + SELECT * FROM cspj_ms_mv + EXCEPT ALL + ( + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label = 'active' + UNION ALL + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label = 'inactive' + UNION ALL + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label IS NULL + ) +); +---- +0 + +query I +SELECT COUNT(*) FROM ( + ( + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label = 'active' + UNION ALL + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label = 'inactive' + UNION ALL + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM cspj_ms_terms t LEFT JOIN cspj_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label IS NULL + ) + EXCEPT ALL + SELECT * FROM cspj_ms_mv +); +---- +0 diff --git a/test/sql/cascade_window_unscopable_delta.test b/test/sql/cascade_window_unscopable_delta.test new file mode 100644 index 00000000..38f93446 --- /dev/null +++ b/test/sql/cascade_window_unscopable_delta.test @@ -0,0 +1,529 @@ +# name: test/sql/cascade_window_unscopable_delta.test +# description: WINDOW_PARTITION plans whose partition keys cannot be scoped to a source delta still emit signed cascade deltas when facts.force_view_delta_cascade=true +# group: [sql] + +require openivm + +statement ok +SET openivm_files_path='__TEST_DIR__'; + +# ========================================================================== +# Case 1 - global surrogate key: ROW_NUMBER() OVER (ORDER BY ...) with no +# PARTITION BY. The classifier still selects WINDOW_PARTITION, but there are +# no partition columns to scope an affected-key set with, so the compiler +# falls back to a full recompute of the view body. That fallback must still +# honour an explicitly requested cascade view delta - otherwise every +# downstream materialized view observes "no delta" and gets demoted to a +# full refresh. +# ========================================================================== + +statement ok +CREATE TABLE uwd_src (id INT, grp VARCHAR, val INT); + +statement ok +INSERT INTO uwd_src VALUES + (1, 'a', 10), + (2, 'a', 20), + (3, 'b', 5), + (4, 'b', 15); + +statement ok +CREATE MATERIALIZED VIEW uwd_dim AS + SELECT CAST(ROW_NUMBER() OVER (ORDER BY grp, id) AS INT) AS dim_key, id, grp, val + FROM uwd_src; + +# The refresh type is unchanged by the cascade request: no relabeling, no +# demotion to FULL_REFRESH. +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'uwd_dim', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":false}' +) +LIMIT 1; +---- +5 WINDOW_PARTITION + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'uwd_dim', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +5 WINDOW_PARTITION + +# Without a cascade request the program is the plain in-place recompute and +# must not write into the view's own delta table. +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%DELETE FROM %openivm_data_uwd_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%INSERT INTO %openivm_data_uwd_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) NOT LIKE '%INSERT INTO openivm_delta_uwd_dim%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'uwd_dim', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":false}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# With a cascade request the same recompute must be bracketed by old/new +# snapshots and publish the signed whole-view delta. +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CREATE OR REPLACE TEMP TABLE openivm_old_uwd_dim AS%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CREATE OR REPLACE TEMP TABLE openivm_new_uwd_dim AS%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%DELETE FROM %openivm_data_uwd_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%INSERT INTO %openivm_data_uwd_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%INSERT INTO openivm_delta_uwd_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CAST%-1%INTEGER%openivm_old_uwd_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%UNION ALL%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CAST%1%INTEGER%openivm_new_uwd_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%DROP TABLE IF EXISTS openivm_old_uwd_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%DROP TABLE IF EXISTS openivm_new_uwd_dim%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'uwd_dim', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# ========================================================================== +# Case 2 - computed partition key: PARTITION BY lower(a) || '_' || lower(b) +# is not a column of any source delta table, so no partition delta spec can +# be built. Same requirement. +# ========================================================================== + +statement ok +CREATE TABLE uwd_os_src (os_name VARCHAR, os_sku VARCHAR); + +statement ok +INSERT INTO uwd_os_src VALUES ('linux', 'a'), ('windows', 'b'), ('linux', 'c'); + +statement ok +CREATE MATERIALIZED VIEW uwd_os_dim AS + SELECT CAST(ROW_NUMBER() OVER (ORDER BY os_full) AS INT) AS os_key, os_full, os_name, os_sku + FROM ( + SELECT lower(os_name) || '_' || lower(os_sku) AS os_full, os_name, os_sku, + ROW_NUMBER() OVER (PARTITION BY lower(os_name) || '_' || lower(os_sku) ORDER BY os_name) AS rn + FROM uwd_os_src + ) t + WHERE rn = 1; + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'uwd_os_dim', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +5 WINDOW_PARTITION + +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) NOT LIKE '%INSERT INTO openivm_delta_uwd_os_dim%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'uwd_os_dim', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":false}' +) +WHERE stmt_kind = 'data'; +---- +1 + +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CREATE OR REPLACE TEMP TABLE openivm_old_uwd_os_dim AS%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CREATE OR REPLACE TEMP TABLE openivm_new_uwd_os_dim AS%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%INSERT INTO openivm_delta_uwd_os_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CAST%-1%INTEGER%openivm_old_uwd_os_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CAST%1%INTEGER%openivm_new_uwd_os_dim%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%DROP TABLE IF EXISTS openivm_new_uwd_os_dim%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'uwd_os_dim', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# ========================================================================== +# Case 3 - multi-source join whose partition key is a computed expression: +# the window partition lineage cannot cover every source delta table. +# ========================================================================== + +statement ok +CREATE TABLE uwd_inst (instance_id INT, sub_id INT, region VARCHAR); + +statement ok +CREATE TABLE uwd_sub (sub_id INT, sub_name VARCHAR); + +statement ok +INSERT INTO uwd_inst VALUES (1, 10, 'eus'), (2, 10, 'wus'), (3, 11, 'eus'); + +statement ok +INSERT INTO uwd_sub VALUES (10, 's10'), (11, 's11'); + +statement ok +CREATE MATERIALIZED VIEW uwd_target AS + SELECT CAST(ROW_NUMBER() OVER (PARTITION BY upper(i.region) ORDER BY i.instance_id) AS INT) AS rn, + i.instance_id, i.region, s.sub_name + FROM uwd_inst i JOIN uwd_sub s ON i.sub_id = s.sub_id; + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'uwd_target', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +5 WINDOW_PARTITION + +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CREATE OR REPLACE TEMP TABLE openivm_old_uwd_target AS%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CREATE OR REPLACE TEMP TABLE openivm_new_uwd_target AS%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%INSERT INTO openivm_delta_uwd_target%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CAST%-1%INTEGER%openivm_old_uwd_target%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CAST%1%INTEGER%openivm_new_uwd_target%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'uwd_target', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# ========================================================================== +# End-to-end: run the emitted cascade program for uwd_dim after a batch of +# conflicting DML (insert + update + delete applied before a single refresh) +# and prove (a) the signed delta reconstructs the new view content from the +# old content under bag semantics, and (b) a downstream materialized view +# refreshes incrementally off that delta and stays bidirectionally equal to +# its base query. PRAGMA refresh() uses the default CompileFacts (cascade +# off), so the cascade program is executed explicitly here exactly as the +# engine driver executes it. +# ========================================================================== + +statement ok +CREATE MATERIALIZED VIEW uwd_down AS + SELECT dim_key, grp, val FROM uwd_dim WHERE val > 6; + +statement ok +INSERT INTO uwd_src VALUES (5, 'a', 15), (6, 'c', 1); + +statement ok +DELETE FROM uwd_src WHERE id = 1; + +statement ok +UPDATE uwd_src SET val = 12 WHERE id = 3; + +statement ok +CREATE TABLE uwd_old_snapshot AS SELECT * FROM openivm_data_uwd_dim; + +statement ok +CREATE OR REPLACE TEMP TABLE openivm_old_uwd_dim AS +SELECT * FROM openivm_data_uwd_dim openivm_old; + +statement ok +CREATE OR REPLACE TEMP TABLE openivm_new_uwd_dim AS +SELECT * FROM ( + SELECT CAST(ROW_NUMBER() OVER (ORDER BY grp, id) AS INT) AS dim_key, id, grp, val + FROM uwd_src +) openivm_recompute; + +statement ok +DELETE FROM openivm_data_uwd_dim; + +statement ok +INSERT INTO openivm_data_uwd_dim +SELECT * FROM openivm_new_uwd_dim; + +statement ok +INSERT INTO openivm_delta_uwd_dim +SELECT *, CAST(-1 AS INTEGER), CURRENT_TIMESTAMP FROM openivm_old_uwd_dim +UNION ALL +SELECT *, CAST(1 AS INTEGER), CURRENT_TIMESTAMP FROM openivm_new_uwd_dim; + +statement ok +DROP TABLE IF EXISTS openivm_old_uwd_dim; + +statement ok +DROP TABLE IF EXISTS openivm_new_uwd_dim; + +# The retraction leg is exactly the pre-refresh content (bag equality). +query I +SELECT COUNT(*) FROM ( + SELECT dim_key, id, grp, val FROM openivm_delta_uwd_dim WHERE openivm_multiplicity = -1 + EXCEPT ALL + SELECT dim_key, id, grp, val FROM uwd_old_snapshot +); +---- +0 + +query I +SELECT COUNT(*) FROM ( + SELECT dim_key, id, grp, val FROM uwd_old_snapshot + EXCEPT ALL + SELECT dim_key, id, grp, val FROM openivm_delta_uwd_dim WHERE openivm_multiplicity = -1 +); +---- +0 + +# Applying the signed delta to the old content reproduces the new content +# exactly, in both directions. +query I +WITH applied AS ( + ( + SELECT dim_key, id, grp, val FROM uwd_old_snapshot + UNION ALL + SELECT dim_key, id, grp, val FROM openivm_delta_uwd_dim WHERE openivm_multiplicity = 1 + ) + EXCEPT ALL + SELECT dim_key, id, grp, val FROM openivm_delta_uwd_dim WHERE openivm_multiplicity = -1 +) +SELECT COUNT(*) FROM ( + SELECT dim_key, id, grp, val FROM applied + EXCEPT ALL + SELECT dim_key, id, grp, val FROM openivm_data_uwd_dim +); +---- +0 + +query I +WITH applied AS ( + ( + SELECT dim_key, id, grp, val FROM uwd_old_snapshot + UNION ALL + SELECT dim_key, id, grp, val FROM openivm_delta_uwd_dim WHERE openivm_multiplicity = 1 + ) + EXCEPT ALL + SELECT dim_key, id, grp, val FROM openivm_delta_uwd_dim WHERE openivm_multiplicity = -1 +) +SELECT COUNT(*) FROM ( + SELECT dim_key, id, grp, val FROM openivm_data_uwd_dim + EXCEPT ALL + SELECT dim_key, id, grp, val FROM applied +); +---- +0 + +# The refreshed view itself matches its base query. +query I +SELECT COUNT(*) FROM ( + SELECT dim_key, id, grp, val FROM openivm_data_uwd_dim + EXCEPT ALL + SELECT CAST(ROW_NUMBER() OVER (ORDER BY grp, id) AS INT) AS dim_key, id, grp, val FROM uwd_src +); +---- +0 + +query I +SELECT COUNT(*) FROM ( + SELECT CAST(ROW_NUMBER() OVER (ORDER BY grp, id) AS INT) AS dim_key, id, grp, val FROM uwd_src + EXCEPT ALL + SELECT dim_key, id, grp, val FROM openivm_data_uwd_dim +); +---- +0 + +# The downstream view consumes the cascade delta incrementally and is +# bidirectionally equal to its base query. +statement ok +PRAGMA refresh('uwd_down'); + +query I +SELECT COUNT(*) FROM ( + SELECT dim_key, grp, val FROM openivm_data_uwd_down + EXCEPT ALL + SELECT dim_key, grp, val FROM ( + SELECT CAST(ROW_NUMBER() OVER (ORDER BY grp, id) AS INT) AS dim_key, id, grp, val FROM uwd_src + ) d + WHERE val > 6 +); +---- +0 + +query I +SELECT COUNT(*) FROM ( + SELECT dim_key, grp, val FROM ( + SELECT CAST(ROW_NUMBER() OVER (ORDER BY grp, id) AS INT) AS dim_key, id, grp, val FROM uwd_src + ) d + WHERE val > 6 + EXCEPT ALL + SELECT dim_key, grp, val FROM openivm_data_uwd_down +); +---- +0 + +# ========================================================================== +# Duplicate-heavy view: the signed delta must carry exact multiplicities, +# not a de-duplicated set. +# ========================================================================== + +statement ok +CREATE TABLE uwd_bag_src (grp VARCHAR, val INT); + +statement ok +INSERT INTO uwd_bag_src VALUES ('a', 1), ('a', 1), ('a', 1), ('b', 2), ('b', 2); + +statement ok +CREATE MATERIALIZED VIEW uwd_bag_mv AS + SELECT grp, val, CAST(SUM(val) OVER () AS INT) AS total FROM uwd_bag_src; + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'uwd_bag_mv', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +5 WINDOW_PARTITION + +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%INSERT INTO openivm_delta_uwd_bag_mv%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CAST%-1%INTEGER%openivm_old_uwd_bag_mv%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%CAST%1%INTEGER%openivm_new_uwd_bag_mv%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'uwd_bag_mv', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +statement ok +INSERT INTO uwd_bag_src VALUES ('a', 1), ('c', 3); + +statement ok +CREATE TABLE uwd_bag_old AS SELECT * FROM openivm_data_uwd_bag_mv; + +statement ok +CREATE OR REPLACE TEMP TABLE openivm_old_uwd_bag_mv AS +SELECT * FROM openivm_data_uwd_bag_mv openivm_old; + +statement ok +CREATE OR REPLACE TEMP TABLE openivm_new_uwd_bag_mv AS +SELECT * FROM ( + SELECT grp, val, CAST(SUM(val) OVER () AS INT) AS total FROM uwd_bag_src +) openivm_recompute; + +statement ok +DELETE FROM openivm_data_uwd_bag_mv; + +statement ok +INSERT INTO openivm_data_uwd_bag_mv +SELECT * FROM openivm_new_uwd_bag_mv; + +statement ok +INSERT INTO openivm_delta_uwd_bag_mv +SELECT *, CAST(-1 AS INTEGER), CURRENT_TIMESTAMP FROM openivm_old_uwd_bag_mv +UNION ALL +SELECT *, CAST(1 AS INTEGER), CURRENT_TIMESTAMP FROM openivm_new_uwd_bag_mv; + +statement ok +DROP TABLE IF EXISTS openivm_old_uwd_bag_mv; + +statement ok +DROP TABLE IF EXISTS openivm_new_uwd_bag_mv; + +# Three identical ('a', 1, 5) rows must be retracted three times, and four +# identical ('a', 1, 9) rows must be inserted four times. +query III +SELECT grp, val, COUNT(*) FROM openivm_delta_uwd_bag_mv +WHERE openivm_multiplicity = -1 AND grp = 'a' +GROUP BY grp, val; +---- +a 1 3 + +query III +SELECT grp, val, COUNT(*) FROM openivm_delta_uwd_bag_mv +WHERE openivm_multiplicity = 1 AND grp = 'a' +GROUP BY grp, val; +---- +a 1 4 + +query I +WITH applied AS ( + ( + SELECT grp, val, total FROM uwd_bag_old + UNION ALL + SELECT grp, val, total FROM openivm_delta_uwd_bag_mv WHERE openivm_multiplicity = 1 + ) + EXCEPT ALL + SELECT grp, val, total FROM openivm_delta_uwd_bag_mv WHERE openivm_multiplicity = -1 +) +SELECT COUNT(*) FROM ( + SELECT grp, val, total FROM applied + EXCEPT ALL + SELECT grp, val, total FROM openivm_data_uwd_bag_mv +); +---- +0 + +query I +WITH applied AS ( + ( + SELECT grp, val, total FROM uwd_bag_old + UNION ALL + SELECT grp, val, total FROM openivm_delta_uwd_bag_mv WHERE openivm_multiplicity = 1 + ) + EXCEPT ALL + SELECT grp, val, total FROM openivm_delta_uwd_bag_mv WHERE openivm_multiplicity = -1 +) +SELECT COUNT(*) FROM ( + SELECT grp, val, total FROM openivm_data_uwd_bag_mv + EXCEPT ALL + SELECT grp, val, total FROM applied +); +---- +0 + +# ========================================================================== +# Requesting a cascade delta must not relabel or rescue an unsupported plan. +# A window view whose predicate is non-deterministic (the same guard that +# demotes production predicates built on current_date()) stays FULL_REFRESH, +# and an unknown view still fails explicitly. +# ========================================================================== + +statement ok +CREATE TABLE uwd_vol_src (id INT, event_date DATE, val INT); + +statement ok +INSERT INTO uwd_vol_src VALUES (1, DATE '2024-01-01', 10); + +statement ok +CREATE MATERIALIZED VIEW uwd_vol_mv AS + SELECT CAST(ROW_NUMBER() OVER (ORDER BY id) AS INT) AS k, id, event_date, val + FROM uwd_vol_src + WHERE val > random(); + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'uwd_vol_mv', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +3 FULL_REFRESH + +statement error +SELECT * FROM openivm_compile_with_facts( + 'uwd_does_not_exist', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +); +---- +materialized view 'uwd_does_not_exist' not found diff --git a/test/sql/compile_refresh.test b/test/sql/compile_refresh.test index 531b4181..95afb325 100644 --- a/test/sql/compile_refresh.test +++ b/test/sql/compile_refresh.test @@ -387,6 +387,212 @@ WHERE stmt_kind = 'data'; ---- 1 +# ========================================== +# Test 10: reduced "arc_machine_status_transaction" shape — a CTE joining an +# INNER JOIN with a LEFT JOIN, whose result feeds an outer query with two +# further chained LEFT JOINs. Before src/delta/operators/join.cpp's +# AppendMultiplicityToAncestorProjectionMaps fix, a deeper join's own +# left_projection_map growing (to carry a new per-leaf multiplicity column up +# to the root) silently shifted where its right-side contribution begins in +# its own combined column numbering. A grandparent join's pre-existing +# projection-map entry — fixed before that growth happened — could then +# alias onto the just-added multiplicity column instead of the real column +# it used to select, permanently dropping a business column and surfacing +# as an internal arity mismatch ("union lhs column ref not in column_map"). +# ========================================== + +statement ok +CREATE TABLE arcm_msf(id INT, coll_key INT); + +statement ok +CREATE TABLE arcm_acd(coll_key INT, sub_id INT, arm_id VARCHAR); + +statement ok +CREATE TABLE arcm_hw(arm_id VARCHAR, provider VARCHAR); + +statement ok +CREATE TABLE arcm_tpid(sub_id INT, tp_id INT); + +statement ok +CREATE TABLE arcm_cust(tp_id INT, cust_key INT); + +statement ok +INSERT INTO arcm_msf VALUES (1, 100); + +statement ok +INSERT INTO arcm_acd VALUES (100, 10, 'arm-1'); + +statement ok +INSERT INTO arcm_hw VALUES ('arm-1', 'azure'); + +statement ok +INSERT INTO arcm_tpid VALUES (10, 200); + +statement ok +INSERT INTO arcm_cust VALUES (200, 21); + +statement ok +CREATE MATERIALIZED VIEW mv_arcm AS + WITH bff AS ( + SELECT msf.id, acd.sub_id, + coalesce(hw.provider, 'N/A') AS provider + FROM arcm_msf msf + INNER JOIN arcm_acd acd ON msf.coll_key = acd.coll_key + LEFT JOIN arcm_hw hw ON acd.arm_id = hw.arm_id + ) + SELECT bff.id, + coalesce(cust.cust_key, 1) AS cust_key + FROM bff + LEFT JOIN arcm_tpid tpid ON bff.sub_id = tpid.sub_id + LEFT JOIN arcm_cust cust ON coalesce(tpid.tp_id, -1) = cust.tp_id; + +statement ok +INSERT INTO arcm_msf VALUES (2, 100); + +# Must compile to a real incremental cascade delta (SIMPLE_PROJECTION), never +# FULL_REFRESH, and must not lose the `cust_key` column from the final insert. +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'mv_arcm', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +2 SIMPLE_PROJECTION + +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) + LIKE '%INSERT INTO openivm_delta_mv_arcm (id, cust_key, openivm_multiplicity)%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_msf%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_acd%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_hw%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_tpid%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_cust%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'mv_arcm', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# Batched multi-leaf delta variant of the same shape: simultaneous +# inserts/delete/update spread across ALL FIVE base tables (not just a +# single-row insert into one leaf). This is the exact reduction that exposed +# a second, deeper defect in src/delta/operators/join.cpp's +# BuildInclusionExclusionTerms: leaf substitution reused a stale pointer into +# the ORIGINAL (pre-mask-renumbering) plan instead of the mask's own +# freshly-renumbered leaf, so a leaf feeding two joins (msf join acd AND acd +# join hw) got its delta substituted at the wrong table_index, leaving the +# second join's condition dangling and surfacing as +# LPTS_UNSUPPORTED_COLUMN_REF once more than one leaf changed at once. +statement ok +INSERT INTO arcm_msf VALUES (3, 300); + +statement ok +INSERT INTO arcm_acd VALUES (300, 30, 'arm-2'); + +statement ok +INSERT INTO arcm_hw VALUES ('arm-2', 'gcp'); + +statement ok +INSERT INTO arcm_tpid VALUES (30, 400); + +statement ok +DELETE FROM arcm_msf WHERE id = 2; + +statement ok +UPDATE arcm_cust SET cust_key = 99 WHERE tp_id = 200; + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'mv_arcm', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +2 SIMPLE_PROJECTION + +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) + LIKE '%INSERT INTO openivm_delta_mv_arcm (id, cust_key, openivm_multiplicity)%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_msf%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_acd%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_hw%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_tpid%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_arcm_cust%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'mv_arcm', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# ========================================== +# Test 11: reduced "int_instance_status_transaction" shape — a UNION ALL of +# two LEFT JOIN branches partitioned by an IS NOT NULL / IS NULL predicate on +# the joined side. Before the `third_party/lpts` pin was advanced to commit +# 77ed5e5 (fork PR #18, "Fix set-op column binding remaps"), remapping +# duplicated UNION ALL key-projection output aliases could leave a stale +# reference to a child alias that no longer existed in the rewritten plan, +# surfacing as a runtime Binder Error ("Referenced column ... not found in +# FROM clause"). +# ========================================== + +statement ok +CREATE TABLE inti_isf(id INT, k INT); + +statement ok +CREATE TABLE inti_ml(k INT, res VARCHAR); + +statement ok +INSERT INTO inti_isf VALUES (1, 100), (2, 200); + +statement ok +INSERT INTO inti_ml VALUES (100, 'machine-1'); + +statement ok +CREATE MATERIALIZED VIEW mv_inti AS + SELECT isf.id, isf.k FROM inti_isf isf LEFT JOIN inti_ml ml ON isf.k = ml.k WHERE ml.res IS NOT NULL + UNION ALL + SELECT isf.id, isf.k FROM inti_isf isf LEFT JOIN inti_ml ml ON isf.k = ml.k WHERE ml.res IS NULL; + +statement ok +INSERT INTO inti_isf VALUES (3, 300); + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'mv_inti', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +2 SIMPLE_PROJECTION + +query I +SELECT CASE WHEN + string_agg(sql, ' ' ORDER BY stmt_order) + LIKE '%INSERT INTO openivm_delta_mv_inti (id, k, openivm_left_key, openivm_multiplicity)%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_inti_isf%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%openivm_delta_inti_ml%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%UNION ALL%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'mv_inti', + '{"target_dialect":"duckdb","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + # ========================================== # Test 9: openivm_compile_with_facts on a non-existent view fails cleanly # ========================================== diff --git a/test/sql/compile_spark_dialect_hardening.test b/test/sql/compile_spark_dialect_hardening.test index 246eefd3..8d4b180b 100644 --- a/test/sql/compile_spark_dialect_hardening.test +++ b/test/sql/compile_spark_dialect_hardening.test @@ -195,3 +195,260 @@ WHERE stmt_kind = 'data'; statement ok SET openivm_refresh_mode = 'incremental'; + +# ========================================== +# Spark hardening for the reduced "arc_machine_status_transaction" shape: a +# CTE (INNER JOIN + LEFT JOIN) feeding an outer query with two further +# chained LEFT JOINs must still compile to a real incremental cascade delta +# under target_dialect=spark, with a spark-portable (no `::`) cast and +# backtick-quoted identifiers, and without dropping the `cust_key` column +# from the final signed insert — the regression this shape covers is in +# src/delta/operators/join.cpp's AppendMultiplicityToAncestorProjectionMaps, +# not dialect-specific, but the original failure was first observed against +# a spark-targeted compile. +# ========================================== +statement ok +CREATE TABLE sph_arcm_msf(id INT, coll_key INT); + +statement ok +CREATE TABLE sph_arcm_acd(coll_key INT, sub_id INT, arm_id VARCHAR); + +statement ok +CREATE TABLE sph_arcm_hw(arm_id VARCHAR, provider VARCHAR); + +statement ok +CREATE TABLE sph_arcm_tpid(sub_id INT, tp_id INT); + +statement ok +CREATE TABLE sph_arcm_cust(tp_id INT, cust_key INT); + +statement ok +INSERT INTO sph_arcm_msf VALUES (1, 100); + +statement ok +INSERT INTO sph_arcm_acd VALUES (100, 10, 'arm-1'); + +statement ok +INSERT INTO sph_arcm_hw VALUES ('arm-1', 'azure'); + +statement ok +INSERT INTO sph_arcm_tpid VALUES (10, 200); + +statement ok +INSERT INTO sph_arcm_cust VALUES (200, 21); + +statement ok +CREATE MATERIALIZED VIEW sph_arcm_mv AS + WITH bff AS ( + SELECT msf.id, acd.sub_id, + coalesce(hw.provider, 'N/A') AS provider + FROM sph_arcm_msf msf + INNER JOIN sph_arcm_acd acd ON msf.coll_key = acd.coll_key + LEFT JOIN sph_arcm_hw hw ON acd.arm_id = hw.arm_id + ) + SELECT bff.id, + coalesce(cust.cust_key, 1) AS cust_key + FROM bff + LEFT JOIN sph_arcm_tpid tpid ON bff.sub_id = tpid.sub_id + LEFT JOIN sph_arcm_cust cust ON coalesce(tpid.tp_id, -1) = cust.tp_id; + +statement ok +INSERT INTO sph_arcm_msf VALUES (2, 100); + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'sph_arcm_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +2 SIMPLE_PROJECTION + +query I +SELECT CASE WHEN string_agg(sql, ' ' ORDER BY stmt_order) NOT LIKE '%::%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%`%' + AND string_agg(sql, ' ' ORDER BY stmt_order) + LIKE '%INSERT INTO openivm_delta_sph_arcm_mv (id, cust_key, openivm_multiplicity)%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'sph_arcm_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# ========================================== +# Spark hardening for the reduced "int_instance_status_transaction" shape: a +# UNION ALL of two LEFT JOIN branches partitioned by an IS NOT NULL / IS NULL +# predicate on the joined side must still compile cleanly under +# target_dialect=spark without a stale child-alias reference leaking into the +# emitted SQL (the failure this shape covers, fixed upstream in +# `third_party/lpts` at commit 77ed5e5, surfaced as a Binder Error referencing +# a column that no longer existed in the rewritten plan). +# ========================================== +statement ok +CREATE TABLE sph_inti_isf(id INT, k INT); + +statement ok +CREATE TABLE sph_inti_ml(k INT, res VARCHAR); + +statement ok +INSERT INTO sph_inti_isf VALUES (1, 100), (2, 200); + +statement ok +INSERT INTO sph_inti_ml VALUES (100, 'machine-1'); + +statement ok +CREATE MATERIALIZED VIEW sph_inti_mv AS + SELECT isf.id, isf.k FROM sph_inti_isf isf LEFT JOIN sph_inti_ml ml ON isf.k = ml.k WHERE ml.res IS NOT NULL + UNION ALL + SELECT isf.id, isf.k FROM sph_inti_isf isf LEFT JOIN sph_inti_ml ml ON isf.k = ml.k WHERE ml.res IS NULL; + +statement ok +INSERT INTO sph_inti_isf VALUES (3, 300); + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'sph_inti_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +2 SIMPLE_PROJECTION + +query I +SELECT CASE WHEN string_agg(sql, ' ' ORDER BY stmt_order) NOT LIKE '%::%' + AND string_agg(sql, ' ' ORDER BY stmt_order) + LIKE '%INSERT INTO openivm_delta_sph_inti_mv (id, k, openivm_left_key, openivm_multiplicity)%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%UNION ALL%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'sph_inti_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# ========================================== +# Spark hardening for the "bounded HUGEINT" canary: a plain (non-aggregate) +# projection that widens a BIGINT to HUGEINT via COALESCE(CAST(...), 0) must +# still compile to a real incremental cascade delta under +# target_dialect=spark. Before `third_party/lpts` was advanced to commit +# 754c797 ("Fix bounded Spark HUGEINT and UNION aliases"), rendering ANY cast +# to HUGEINT unconditionally raised LPTS_UNSUPPORTED_TYPE for Spark, even +# when the source type (BIGINT here) provably fits within Spark's +# DECIMAL(38,0); the fix maps a provably-bounded HUGEINT cast/literal to +# DECIMAL(38,0) instead of failing the compile. Note a literal +# COALESCE(SUM(...), 0) is classified GROUP_RECOMPUTE by OpenIVM (bypassing +# LPTS entirely, echoing native SQL), so this canary uses an equivalent +# non-aggregate COALESCE(CAST(... AS HUGEINT), 0) shape that is classified +# SIMPLE_PROJECTION and does go through LPTS. +# ========================================== +statement ok +CREATE TABLE sph_hugeint_src(id INT, amount BIGINT); + +statement ok +INSERT INTO sph_hugeint_src VALUES (1, 5000000000), (2, NULL); + +statement ok +CREATE MATERIALIZED VIEW sph_hugeint_mv AS + SELECT id, COALESCE(CAST(amount AS HUGEINT), 0) AS wide_amount + FROM sph_hugeint_src; + +statement ok +INSERT INTO sph_hugeint_src VALUES (3, 2000000000); + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'sph_hugeint_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +2 SIMPLE_PROJECTION + +query I +SELECT CASE WHEN string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%DECIMAL(38,0)%' + AND string_agg(sql, ' ' ORDER BY stmt_order) NOT LIKE '% HUGEINT%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'sph_hugeint_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 + +# ========================================== +# Spark hardening for the "machine-status left-deep UNION" canary: a literal +# three-way (left-deep) UNION ALL of LEFT JOIN branches, each partitioned by a +# mutually exclusive predicate on the joined side, must still compile to a +# real incremental cascade delta under target_dialect=spark, correctly +# binding every column -- including the trailing duplicated +# `openivm_left_key`/`openivm_multiplicity` output columns -- at every level +# of the left-deep chain (internally: a nested union node whose own union is +# itself one arm of a further union, e.g. `(A UNION ALL B) UNION ALL C`). +# This is upstream `third_party/lpts` commit 754c797's second fix ("retain +# trailing rewritten UNION bindings as aliases of their physical multiplicity +# output"); unlike the two-way `sph_inti_mv` shape above, this exercises an +# actual left-deep union chain, matching the shape of LPTS's own +# `test/sql/union.test` regression (a left-deep "machine status" UNION ALL +# carrying duplicated key/multiplicity output columns). +# ========================================== +statement ok +CREATE TABLE sph_ms_terms(event_time INT, machine_arm_id VARCHAR, customer_key INT); + +statement ok +CREATE TABLE sph_ms_status(machine_arm_id VARCHAR, status_label VARCHAR); + +statement ok +INSERT INTO sph_ms_terms VALUES (10, 'arm-1', 1), (20, 'arm-2', 6); + +statement ok +INSERT INTO sph_ms_status VALUES ('arm-1', 'active'); + +statement ok +CREATE MATERIALIZED VIEW sph_ms_mv AS + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM sph_ms_terms t LEFT JOIN sph_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label = 'active' + UNION ALL + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM sph_ms_terms t LEFT JOIN sph_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label = 'inactive' + UNION ALL + SELECT t.event_time, t.machine_arm_id, t.customer_key + FROM sph_ms_terms t LEFT JOIN sph_ms_status s ON t.machine_arm_id = s.machine_arm_id + WHERE s.status_label IS NULL; + +statement ok +INSERT INTO sph_ms_terms VALUES (30, 'arm-3', 11); + +query IT +SELECT refresh_type, refresh_type_name +FROM openivm_compile_with_facts( + 'sph_ms_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +LIMIT 1; +---- +2 SIMPLE_PROJECTION + +query I +SELECT CASE WHEN string_agg(sql, ' ' ORDER BY stmt_order) NOT LIKE '%::%' + AND string_agg(sql, ' ' ORDER BY stmt_order) + LIKE '%INSERT INTO openivm_delta_sph_ms_mv (event_time, machine_arm_id, customer_key, openivm_left_key, openivm_multiplicity)%' + AND string_agg(sql, ' ' ORDER BY stmt_order) LIKE '%UNION ALL%UNION ALL%' + THEN 1 ELSE 0 END +FROM openivm_compile_with_facts( + 'sph_ms_mv', + '{"target_dialect":"spark","compile_only":true,"force_view_delta_cascade":true}' +) +WHERE stmt_kind = 'data'; +---- +1 diff --git a/test/sql/left_join_regular_nterm.test b/test/sql/left_join_regular_nterm.test new file mode 100644 index 00000000..5b4177d3 --- /dev/null +++ b/test/sql/left_join_regular_nterm.test @@ -0,0 +1,161 @@ +# name: test/sql/left_join_regular_nterm.test +# description: Compile-only N-term telescoping delta for LEFT-JOIN SIMPLE_PROJECTION views (openivm_regular_nterm_left) +# group: [sql] + +# A deep LEFT-join star that would blow up under 2^N inclusion-exclusion must +# compile to a linear delta and stay bag-correct. + +require openivm + +statement ok +SET openivm_files_path='__TEST_DIR__'; + +statement ok +SET openivm_regular_nterm_left=true; + +# ========================================== +# Star schema: 1 fact + 5 LEFT-joined dimensions, projection (no aggregation). +# ========================================== + +statement ok +CREATE TABLE fact(id INTEGER, amount INTEGER, k1 INTEGER, k2 INTEGER, k3 INTEGER, k4 INTEGER, k5 INTEGER); + +statement ok +CREATE TABLE d1(k INTEGER, v VARCHAR); + +statement ok +CREATE TABLE d2(k INTEGER, v VARCHAR); + +statement ok +CREATE TABLE d3(k INTEGER, v VARCHAR); + +statement ok +CREATE TABLE d4(k INTEGER, v VARCHAR); + +statement ok +CREATE TABLE d5(k INTEGER, v VARCHAR); + +statement ok +INSERT INTO d1 VALUES (11, 'd1a'), (21, 'd1b'); + +statement ok +INSERT INTO d2 VALUES (12, 'd2a'), (22, 'd2b'); + +statement ok +INSERT INTO d3 VALUES (13, 'd3a'), (23, 'd3b'); + +statement ok +INSERT INTO d4 VALUES (14, 'd4a'), (24, 'd4b'); + +statement ok +INSERT INTO d5 VALUES (15, 'd5a'), (25, 'd5b'); + +# Row 2 has k1=91 which has NO matching d1 yet (NULL-padded until d1 gets 91). +statement ok +INSERT INTO fact VALUES + (1, 100, 11, 12, 13, 14, 15), + (2, 200, 91, 12, 13, 14, 15), + (3, 300, 21, 22, 23, 24, 25), + (4, 400, 11, 22, 13, 24, 15); + +statement ok +CREATE MATERIALIZED VIEW mv AS + SELECT f.id, f.amount, + d1.v AS v1, d2.v AS v2, d3.v AS v3, d4.v AS v4, d5.v AS v5 + FROM fact f + LEFT JOIN d1 ON f.k1 = d1.k + LEFT JOIN d2 ON f.k2 = d2.k + LEFT JOIN d3 ON f.k3 = d3.k + LEFT JOIN d4 ON f.k4 = d4.k + LEFT JOIN d5 ON f.k5 = d5.k; + +# Snapshot the MV so we can prove openivm_compile_with_facts does not mutate it. +statement ok +CREATE TABLE mv_before AS SELECT * FROM mv; + +# ========================================== +# Mixed DML batch: NULL->match (d1 gains key 91), match->NULL (d2 loses key 22), +# a dimension value update, and fact insert/update/delete. +# ========================================== + +statement ok +INSERT INTO d1 VALUES (91, 'd1_late'); + +statement ok +DELETE FROM d2 WHERE k = 22; + +statement ok +UPDATE d3 SET v = 'd3_upd' WHERE k = 13; + +statement ok +UPDATE fact SET amount = amount + 5 WHERE id = 1; + +statement ok +INSERT INTO fact VALUES (5, 500, 91, 22, 13, 14, 15); + +statement ok +DELETE FROM fact WHERE id = 3; + +# ========================================== +# Test 1: compile-only path emits a SIMPLE_PROJECTION join delta into +# openivm_delta_mv without exploding (a LEFT star that would be 2^6-1 under +# inclusion-exclusion compiles to a linear N-term telescoping delta). +# ========================================== + +query I +SELECT COUNT(*) > 0 +FROM openivm_compile_with_facts( + 'mv', + '{"target_dialect":"duckdb","compile_only":true, + "delta_shape":{"fact":"MIXED","d1":"INSERT_ONLY","d2":"MIXED","d3":"MIXED","d4":"UNCHANGED","d5":"UNCHANGED"}}' +) +WHERE stmt_kind = 'data' AND refresh_type_name = 'SIMPLE_PROJECTION' + AND sql LIKE '%openivm_delta_mv%'; +---- +true + +# Test 2: the compile call left the materialized view untouched. +query I +SELECT (SELECT count(*) FROM ((SELECT * FROM mv) EXCEPT ALL (SELECT * FROM mv_before))) = 0 AND + (SELECT count(*) FROM ((SELECT * FROM mv_before) EXCEPT ALL (SELECT * FROM mv))) = 0 AS unchanged; +---- +true + +# ========================================== +# Test 3: a real incremental refresh applies the batched deltas and the MV is +# bag-equal to the fully recomputed base query in BOTH directions (full IVM +# correctness cross-check, including the NULL<->match transition rows). +# ========================================== + +statement ok +PRAGMA refresh('mv'); + +query I +SELECT count(*) FROM ( + SELECT f.id, f.amount, d1.v AS v1, d2.v AS v2, d3.v AS v3, d4.v AS v4, d5.v AS v5 + FROM fact f + LEFT JOIN d1 ON f.k1 = d1.k + LEFT JOIN d2 ON f.k2 = d2.k + LEFT JOIN d3 ON f.k3 = d3.k + LEFT JOIN d4 ON f.k4 = d4.k + LEFT JOIN d5 ON f.k5 = d5.k + EXCEPT ALL + SELECT id, amount, v1, v2, v3, v4, v5 FROM mv +); +---- +0 + +query I +SELECT count(*) FROM ( + SELECT id, amount, v1, v2, v3, v4, v5 FROM mv + EXCEPT ALL + SELECT f.id, f.amount, d1.v AS v1, d2.v AS v2, d3.v AS v3, d4.v AS v4, d5.v AS v5 + FROM fact f + LEFT JOIN d1 ON f.k1 = d1.k + LEFT JOIN d2 ON f.k2 = d2.k + LEFT JOIN d3 ON f.k3 = d3.k + LEFT JOIN d4 ON f.k4 = d4.k + LEFT JOIN d5 ON f.k5 = d5.k +); +---- +0 diff --git a/test/sql/spark_add_months.test b/test/sql/spark_add_months.test new file mode 100644 index 00000000..236503a2 --- /dev/null +++ b/test/sql/spark_add_months.test @@ -0,0 +1,86 @@ +# name: test/sql/spark_add_months.test +# description: add_months incremental MV maintenance through the openivm compiler +# group: [sql] + +# +# Exhaustive scalar-correctness coverage for add_months lives in lpts +# (third_party/lpts test/sql/spark_add_months.test). openivm consumes the +# function from lpts; this test covers openivm's concern: that add_months +# resolves in the compiler and drives a real SIMPLE_PROJECTION delta rather +# than a silent FULL_REFRESH demotion. + +require openivm + +statement ok +SET openivm_files_path='__TEST_DIR__'; + +# add_months resolves in an openivm session (registered via lpts source) +query I +SELECT add_months(DATE '2015-01-31', 1) = DATE '2015-02-28'; +---- +true + +# --- Incremental materialized-view maintenance using add_months --- + +statement ok +CREATE TABLE am_item (id INT, d DATE, n INT); + +statement ok +INSERT INTO am_item VALUES + (1, DATE '2015-01-31', 1), + (2, DATE '2016-02-29', 12), + (3, DATE '2015-01-15', 13); + +statement ok +CREATE MATERIALIZED VIEW am_mv AS + SELECT id, d, n, add_months(d, n) AS shifted + FROM am_item; + +statement ok +SELECT COUNT(*) FROM openivm_compile_with_facts( + 'am_mv', + '{"target_dialect":"duckdb","compile_only":true}' +); + +# Real SIMPLE_PROJECTION delta emitted (no full-refresh demotion) +query I +SELECT CASE + WHEN contains(content, 'INSERT INTO openivm_delta_am_mv') AND + contains(content, 'openivm_delta_am_item') AND + NOT contains(content, 'SELECT NULL::') + THEN 1 ELSE 0 + END +FROM read_text('__TEST_DIR__/openivm_upsert_queries_am_mv.sql'); +---- +1 + +statement ok +UPDATE am_item SET d = DATE '2015-12-31', n = 1 WHERE id = 1; + +statement ok +INSERT INTO am_item VALUES (4, DATE '2015-06-30', 3); + +statement ok +DELETE FROM am_item WHERE id = 3; + +statement ok +PRAGMA refresh('am_mv'); + +# Incrementally maintained MV matches full recomputation, both directions +query I +SELECT COUNT(*) FROM ( + SELECT * FROM am_mv + EXCEPT ALL + SELECT id, d, n, add_months(d, n) AS shifted FROM am_item +); +---- +0 + +query I +SELECT COUNT(*) FROM ( + SELECT id, d, n, add_months(d, n) AS shifted FROM am_item + EXCEPT ALL + SELECT * FROM am_mv +); +---- +0 diff --git a/third_party/lpts b/third_party/lpts index 13786cbd..b3baf0bb 160000 --- a/third_party/lpts +++ b/third_party/lpts @@ -1 +1 @@ -Subproject commit 13786cbd8f32216aa5d92c99f8bf53c7cc52c9fb +Subproject commit b3baf0bbd974fc08161f1b7d9a540d60a9df96b8