Skip to content

Commit fa69ce0

Browse files
committed
Continue cascades past empty refresh nodes
1 parent 344ee4e commit fa69ce0

4 files changed

Lines changed: 340 additions & 25 deletions

File tree

docs/codebase-audit-2026-07-22.md

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,9 +25,9 @@ maintenance path. A bug must not be hidden by weakening a test or silently chang
2525
DuckDB query-pragma preprocessing ends a transaction before its returned program completes; `refresh.cpp` records the
2626
native-operator follow-up. DuckLake replacement materializes data and auxiliary state under unpublished staging names
2727
before publishing them.
28-
- [ ] **Continue cascade traversal when the current node has no source deltas.** The root empty-delta fast path returns
29-
before visiting downstream nodes. A child can have independent source deltas or pending parent-delta rows after an earlier
30-
failed cascade. Model `node skipped` separately from `graph traversal complete` and checkpoint each node independently.
28+
- [x] **Continue cascade traversal when the current node has no source deltas.** Empty-delta detection now skips only the
29+
current node; it cannot terminate traversal of the downstream DAG. Regression coverage exercises both a pending parent
30+
delta left by a cascade-off refresh and independent child-source DML, including the hook-aware empty-node path.
3131
- [ ] **Use NULL-safe affected-window partition matching.** `IN` does not select a NULL partition even though SQL window
3232
partitioning groups NULL values. Join the affected-key relation with `IS NOT DISTINCT FROM` for every partition key.
3333
- [ ] **Preserve `SUM` NULL semantics.** Weighted SUM alone cannot distinguish numeric zero from no non-NULL inputs. Persist

src/upsert/refresh.cpp

Lines changed: 22 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -112,7 +112,7 @@ static bool TrySkipEmptyRefresh(ClientContext &context, RefreshMetadata &metadat
112112
// Generate and execute refresh SQL for a single view under its per-view lock.
113113
// When openivm_adaptive_refresh is on, also computes a cost estimate before execution
114114
// and records execution history for the learned cost model.
115-
static bool RefreshViewLocked(ClientContext &context, const string &view_catalog_name, const string &view_schema_name,
115+
static void RefreshViewLocked(ClientContext &context, const string &view_catalog_name, const string &view_schema_name,
116116
const string &vn, bool cross_system, const string &attached_db_catalog_name,
117117
const string &attached_db_schema_name, bool skip_empty_refresh) {
118118
RefreshProfiler profiler(context, vn);
@@ -171,7 +171,7 @@ static bool RefreshViewLocked(ClientContext &context, const string &view_catalog
171171
attached_db_catalog_name, attached_db_schema_name, &delta_activity)) {
172172
profiler.AddTotal();
173173
profiler.Flush(*context.db.get());
174-
return true;
174+
return;
175175
}
176176
if (!delta_activity.active_delta_table_names.empty() || delta_activity.requires_full_refresh) {
177177
precomputed_delta_activity = &delta_activity;
@@ -363,7 +363,7 @@ static bool RefreshViewLocked(ClientContext &context, const string &view_catalog
363363
}
364364
profiler.AddTotal();
365365
profiler.Flush(*context.db.get());
366-
return false;
366+
return;
367367
} catch (...) {
368368
// Ensure the transaction is rolled back before we propagate the exception.
369369
// This covers the case where Query() itself threw (vs returning HasError) —
@@ -544,30 +544,30 @@ void UpsertDeltaQueriesLocked(ClientContext &context, const FunctionParameters &
544544

545545
// Hook-bearing refreshes keep the old pre-hook empty skip semantics. Hook-free refreshes
546546
// compute the same delta activity under the view lock and reuse it during SQL generation.
547-
if (has_refresh_hook && TrySkipEmptyRefresh(context, metadata, con, view_catalog_name, view_schema_name, view_name,
548-
attached_db_catalog_name, attached_db_schema_name, nullptr)) {
549-
return;
550-
}
551-
552-
if (!hook_sql.empty() && hook_mode == "before") {
553-
auto hr = con.Query(hook_sql);
554-
if (hr->HasError()) {
555-
Printer::Print("Warning: before-hook for '" + view_name + "' failed: " + hr->GetError());
547+
bool skip_current_node =
548+
has_refresh_hook && TrySkipEmptyRefresh(context, metadata, con, view_catalog_name, view_schema_name, view_name,
549+
attached_db_catalog_name, attached_db_schema_name, nullptr);
550+
if (!skip_current_node) {
551+
if (!hook_sql.empty() && hook_mode == "before") {
552+
auto hr = con.Query(hook_sql);
553+
if (hr->HasError()) {
554+
Printer::Print("Warning: before-hook for '" + view_name + "' failed: " + hr->GetError());
555+
}
556556
}
557-
}
558557

559-
if (hook_mode != "replace") {
560-
if (RefreshViewLocked(context, view_catalog_name, view_schema_name, view_name, cross_system,
561-
attached_db_catalog_name, attached_db_schema_name, !has_refresh_hook)) {
562-
return;
558+
if (hook_mode != "replace") {
559+
RefreshViewLocked(context, view_catalog_name, view_schema_name, view_name, cross_system,
560+
attached_db_catalog_name, attached_db_schema_name, !has_refresh_hook);
563561
}
564-
}
565562

566-
if (!hook_sql.empty() && (hook_mode == "after" || hook_mode == "replace")) {
567-
auto hr = con.Query(hook_sql);
568-
if (hr->HasError()) {
569-
Printer::Print("Warning: " + hook_mode + "-hook for '" + view_name + "' failed: " + hr->GetError());
563+
if (!hook_sql.empty() && (hook_mode == "after" || hook_mode == "replace")) {
564+
auto hr = con.Query(hook_sql);
565+
if (hr->HasError()) {
566+
Printer::Print("Warning: " + hook_mode + "-hook for '" + view_name + "' failed: " + hr->GetError());
567+
}
570568
}
569+
} else {
570+
OPENIVM_DEBUG_PRINT("[UPSERT] Skipped refresh node '%s'; continuing cascade traversal\n", view_name.c_str());
571571
}
572572

573573
// Downstream cascade: refresh dependents after

test/sql/pipeline.test

Lines changed: 203 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -752,6 +752,209 @@ SELECT COUNT(*) FROM (
752752
----
753753
0
754754

755+
# ============================================================
756+
# Empty current node must not stop downstream cascade traversal
757+
# ============================================================
758+
759+
statement ok
760+
SET openivm_cascade_refresh = 'downstream';
761+
762+
statement ok
763+
CREATE TABLE cascade_skip_root_src (id INT, amount INT);
764+
765+
statement ok
766+
CREATE TABLE cascade_skip_side_src (id INT, bonus INT);
767+
768+
statement ok
769+
INSERT INTO cascade_skip_root_src VALUES (1, 10), (2, 20);
770+
771+
statement ok
772+
INSERT INTO cascade_skip_side_src VALUES (1, 1), (2, 2);
773+
774+
statement ok
775+
CREATE MATERIALIZED VIEW cascade_skip_parent AS
776+
SELECT id, SUM(amount) AS total
777+
FROM cascade_skip_root_src
778+
GROUP BY id;
779+
780+
statement ok
781+
CREATE MATERIALIZED VIEW cascade_skip_child AS
782+
SELECT p.id, SUM(p.total + s.bonus) AS score, COUNT(*) AS cnt
783+
FROM cascade_skip_parent p
784+
JOIN cascade_skip_side_src s USING (id)
785+
GROUP BY p.id;
786+
787+
# Leave a parent delta pending for the child, then invoke a downstream cascade
788+
# when the parent's own source delta is empty.
789+
statement ok
790+
SET openivm_cascade_refresh = 'off';
791+
792+
statement ok
793+
INSERT INTO cascade_skip_root_src VALUES (3, 30), (4, 40);
794+
795+
statement ok
796+
UPDATE cascade_skip_root_src SET amount = amount + 5 WHERE id IN (1, 3);
797+
798+
statement ok
799+
DELETE FROM cascade_skip_root_src WHERE id IN (2, 4);
800+
801+
statement ok
802+
PRAGMA refresh('cascade_skip_parent');
803+
804+
query I
805+
SELECT COUNT(*) FROM (
806+
SELECT id, SUM(amount) AS total
807+
FROM cascade_skip_root_src
808+
GROUP BY id
809+
EXCEPT ALL
810+
SELECT id, total FROM cascade_skip_parent
811+
);
812+
----
813+
0
814+
815+
query I
816+
SELECT COUNT(*) FROM (
817+
SELECT id, total FROM cascade_skip_parent
818+
EXCEPT ALL
819+
SELECT id, SUM(amount) AS total
820+
FROM cascade_skip_root_src
821+
GROUP BY id
822+
);
823+
----
824+
0
825+
826+
statement ok
827+
SET openivm_cascade_refresh = 'downstream';
828+
829+
statement ok
830+
PRAGMA refresh('cascade_skip_parent');
831+
832+
query I
833+
SELECT COUNT(*) FROM (
834+
SELECT id, SUM(amount) AS total
835+
FROM cascade_skip_root_src
836+
GROUP BY id
837+
EXCEPT ALL
838+
SELECT id, total FROM cascade_skip_parent
839+
);
840+
----
841+
0
842+
843+
query I
844+
SELECT COUNT(*) FROM (
845+
SELECT id, total FROM cascade_skip_parent
846+
EXCEPT ALL
847+
SELECT id, SUM(amount) AS total
848+
FROM cascade_skip_root_src
849+
GROUP BY id
850+
);
851+
----
852+
0
853+
854+
query I
855+
SELECT COUNT(*) FROM (
856+
SELECT p.id, SUM(p.total + s.bonus) AS score, COUNT(*) AS cnt
857+
FROM (
858+
SELECT id, SUM(amount) AS total
859+
FROM cascade_skip_root_src
860+
GROUP BY id
861+
) p
862+
JOIN cascade_skip_side_src s USING (id)
863+
GROUP BY p.id
864+
EXCEPT ALL
865+
SELECT id, score, cnt FROM cascade_skip_child
866+
);
867+
----
868+
0
869+
870+
query I
871+
SELECT COUNT(*) FROM (
872+
SELECT id, score, cnt FROM cascade_skip_child
873+
EXCEPT ALL
874+
SELECT p.id, SUM(p.total + s.bonus) AS score, COUNT(*) AS cnt
875+
FROM (
876+
SELECT id, SUM(amount) AS total
877+
FROM cascade_skip_root_src
878+
GROUP BY id
879+
) p
880+
JOIN cascade_skip_side_src s USING (id)
881+
GROUP BY p.id
882+
);
883+
----
884+
0
885+
886+
# The child can also have independent source changes while the parent has none.
887+
# Batch conflicting DML before one refresh: inserted rows are updated, and one
888+
# inserted row is deleted before the cascade.
889+
statement ok
890+
INSERT INTO cascade_skip_side_src VALUES (3, 3), (4, 4);
891+
892+
statement ok
893+
UPDATE cascade_skip_side_src SET bonus = bonus + 10 WHERE id IN (1, 3);
894+
895+
statement ok
896+
DELETE FROM cascade_skip_side_src WHERE id IN (2, 4);
897+
898+
statement ok
899+
PRAGMA refresh('cascade_skip_parent');
900+
901+
query I
902+
SELECT COUNT(*) FROM (
903+
SELECT id, SUM(amount) AS total
904+
FROM cascade_skip_root_src
905+
GROUP BY id
906+
EXCEPT ALL
907+
SELECT id, total FROM cascade_skip_parent
908+
);
909+
----
910+
0
911+
912+
query I
913+
SELECT COUNT(*) FROM (
914+
SELECT id, total FROM cascade_skip_parent
915+
EXCEPT ALL
916+
SELECT id, SUM(amount) AS total
917+
FROM cascade_skip_root_src
918+
GROUP BY id
919+
);
920+
----
921+
0
922+
923+
query I
924+
SELECT COUNT(*) FROM (
925+
SELECT p.id, SUM(p.total + s.bonus) AS score, COUNT(*) AS cnt
926+
FROM (
927+
SELECT id, SUM(amount) AS total
928+
FROM cascade_skip_root_src
929+
GROUP BY id
930+
) p
931+
JOIN cascade_skip_side_src s USING (id)
932+
GROUP BY p.id
933+
EXCEPT ALL
934+
SELECT id, score, cnt FROM cascade_skip_child
935+
);
936+
----
937+
0
938+
939+
query I
940+
SELECT COUNT(*) FROM (
941+
SELECT id, score, cnt FROM cascade_skip_child
942+
EXCEPT ALL
943+
SELECT p.id, SUM(p.total + s.bonus) AS score, COUNT(*) AS cnt
944+
FROM (
945+
SELECT id, SUM(amount) AS total
946+
FROM cascade_skip_root_src
947+
GROUP BY id
948+
) p
949+
JOIN cascade_skip_side_src s USING (id)
950+
GROUP BY p.id
951+
);
952+
----
953+
0
954+
955+
statement ok
956+
SET openivm_cascade_refresh = 'upstream';
957+
755958
# Insert more data
756959
statement ok
757960
INSERT INTO pipe_up VALUES ('a', 5), ('b', 20);

0 commit comments

Comments
 (0)