This is an automated email from the ASF dual-hosted git repository.

snuyanzin pushed a commit to branch release-2.2
in repository https://gitbox.apache.org/repos/asf/flink.git

commit ae9a8cc7d5662d575267f38b3a94f3de480c09fd
Author: Sergey Nuyanzin <[email protected]>
AuthorDate: Mon Aug 3 12:39:58 2026 +0200

    [FLINK-40307][table] `RestoreTestCompleteness` was never executed in CI
---
 .../exec/stream/AsyncCorrelateRestoreTest.java     |   2 +-
 .../stream/ProcessTableFunctionRestoreTests.java   |   2 +-
 .../exec/stream/WatermarkAssignerRestoreTest.java  |   2 +-
 ...eness.java => RestoreTestCompletenessTest.java} |  82 +++++++++++++++------
 .../plan/async-correlate-catalog-func.json         |   0
 .../savepoint/_metadata                            | Bin
 .../plan/async-correlate-exception.json            |   0
 .../async-correlate-exception/savepoint/_metadata  | Bin
 .../plan/async-correlate-join-filter.json          |   0
 .../savepoint/_metadata                            | Bin
 .../plan/async-correlate-left-join.json            |   0
 .../async-correlate-left-join/savepoint/_metadata  | Bin
 .../plan/async-correlate-system-func.json          |   0
 .../savepoint/_metadata                            | Bin
 14 files changed, 61 insertions(+), 27 deletions(-)

diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/AsyncCorrelateRestoreTest.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/AsyncCorrelateRestoreTest.java
index 2d52118b225..6d8214e1caa 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/AsyncCorrelateRestoreTest.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/AsyncCorrelateRestoreTest.java
@@ -28,7 +28,7 @@ import java.util.List;
 public class AsyncCorrelateRestoreTest extends RestoreTestBase {
 
     public AsyncCorrelateRestoreTest() {
-        super(StreamExecCorrelate.class);
+        super(StreamExecAsyncCorrelate.class);
     }
 
     @Override
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionRestoreTests.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionRestoreTests.java
index baa60c612c0..0738f5639ef 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionRestoreTests.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionRestoreTests.java
@@ -26,7 +26,7 @@ import java.util.List;
 /** Restore tests for {@link StreamExecProcessTableFunction}. */
 public class ProcessTableFunctionRestoreTests extends RestoreTestBase {
 
-    protected ProcessTableFunctionRestoreTests() {
+    public ProcessTableFunctionRestoreTests() {
         super(StreamExecProcessTableFunction.class);
     }
 
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/WatermarkAssignerRestoreTest.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/WatermarkAssignerRestoreTest.java
index 0a684e4c133..3d63a1523e1 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/WatermarkAssignerRestoreTest.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/WatermarkAssignerRestoreTest.java
@@ -24,7 +24,7 @@ import org.apache.flink.table.test.program.TableTestProgram;
 import java.util.List;
 
 /** Restore tests for {@link StreamExecWatermarkAssigner}. */
-class WatermarkAssignerRestoreTest extends RestoreTestBase {
+public class WatermarkAssignerRestoreTest extends RestoreTestBase {
 
     public WatermarkAssignerRestoreTest() {
         super(StreamExecWatermarkAssigner.class);
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompleteness.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompletenessTest.java
similarity index 62%
rename from 
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompleteness.java
rename to 
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompletenessTest.java
index 9e6e124449a..6044d1a8951 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompleteness.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/testutils/RestoreTestCompletenessTest.java
@@ -19,6 +19,12 @@
 package org.apache.flink.table.planner.plan.nodes.exec.testutils;
 
 import org.apache.flink.table.planner.plan.nodes.exec.ExecNode;
+import 
org.apache.flink.table.planner.plan.nodes.exec.batch.BatchExecHashAggregate;
+import 
org.apache.flink.table.planner.plan.nodes.exec.batch.BatchExecNestedLoopJoin;
+import 
org.apache.flink.table.planner.plan.nodes.exec.batch.BatchExecSortAggregate;
+import 
org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecDeltaJoin;
+import 
org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecGlobalWindowAggregate;
+import 
org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecLocalWindowAggregate;
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonCalc;
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonCorrelate;
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecPythonGroupAggregate;
@@ -30,7 +36,6 @@ import 
org.apache.flink.table.planner.plan.utils.ExecNodeMetadataUtil.ExecNodeNa
 
 import org.apache.flink.shaded.guava33.com.google.common.reflect.ClassPath;
 
-import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
 import java.io.IOException;
@@ -42,21 +47,33 @@ import java.util.Map;
 import java.util.Set;
 import java.util.stream.Collectors;
 
+import static org.assertj.core.api.Assertions.fail;
+
 /** Validate restore tests exists for Exec Nodes. */
-public class RestoreTestCompleteness {
+class RestoreTestCompletenessTest {
 
     private static final Set<Class<? extends ExecNode<?>>> SKIP_EXEC_NODES =
-            new HashSet<Class<? extends ExecNode<?>>>() {
-                {
-                    /** Ignoring python based exec nodes temporarily. */
-                    add(StreamExecPythonCalc.class);
-                    add(StreamExecPythonCorrelate.class);
-                    add(StreamExecPythonOverAggregate.class);
-                    add(StreamExecPythonGroupAggregate.class);
-                    add(StreamExecPythonGroupTableAggregate.class);
-                    add(StreamExecPythonGroupWindowAggregate.class);
-                }
-            };
+            Set.of(
+                    /* Ignoring python based exec nodes temporarily. */
+                    StreamExecPythonCalc.class,
+                    StreamExecPythonCorrelate.class,
+                    StreamExecPythonOverAggregate.class,
+                    StreamExecPythonGroupAggregate.class,
+                    StreamExecPythonGroupTableAggregate.class,
+                    StreamExecPythonGroupWindowAggregate.class,
+
+                    // Covered by tests in WindowAggregateEventTimeRestoreTest
+                    StreamExecLocalWindowAggregate.class,
+                    StreamExecGlobalWindowAggregate.class,
+                    // Test added in later releases
+                    StreamExecDeltaJoin.class,
+
+                    // There is jira for these 2 batch tests
+                    // https://issues.apache.org/jira/browse/FLINK-40306
+                    BatchExecHashAggregate.class,
+                    BatchExecNestedLoopJoin.class,
+                    // Restore test added in later release
+                    BatchExecSortAggregate.class);
 
     private Class<? extends ExecNode<?>> getExecNode(Class<?> restoreTest)
             throws NoSuchMethodException,
@@ -85,7 +102,7 @@ public class RestoreTestCompleteness {
     }
 
     @Test
-    public void testMissingRestoreTest()
+    void testMissingRestoreTest()
             throws IOException,
                     NoSuchMethodException,
                     InstantiationException,
@@ -95,12 +112,14 @@ public class RestoreTestCompleteness {
                 ExecNodeMetadataUtil.getVersionedExecNodes();
 
         Set<ClassPath.ClassInfo> classesInPackage =
-                ClassPath.from(this.getClass().getClassLoader())
-                        .getTopLevelClassesRecursive(
-                                
"org.apache.flink.table.planner.plan.nodes.exec.stream")
-                        .stream()
-                        .filter(x -> 
RestoreTestBase.class.isAssignableFrom(x.load()))
-                        .collect(Collectors.toSet());
+                new HashSet<>(
+                        gatherClasses(
+                                RestoreTestBase.class,
+                                
"org.apache.flink.table.planner.plan.nodes.exec.stream"));
+        classesInPackage.addAll(
+                gatherClasses(
+                        BatchRestoreTestBase.class,
+                        
"org.apache.flink.table.planner.plan.nodes.exec.batch"));
 
         Set<Class<? extends ExecNode<?>>> execNodesWithRestoreTests = new 
HashSet<>();
 
@@ -116,18 +135,33 @@ public class RestoreTestCompleteness {
             }
         }
 
+        Set<Class<? extends ExecNode<?>>> productionExecNodes = 
ExecNodeMetadataUtil.execNodes();
         for (Map.Entry<ExecNodeNameVersion, Class<? extends ExecNode<?>>> 
entry :
                 versionedExecNodes.entrySet()) {
             ExecNodeNameVersion execNodeNameVersion = entry.getKey();
             Class<? extends ExecNode<?>> execNode = entry.getValue();
-            if (!SKIP_EXEC_NODES.contains(execNode)) {
-                final String msg =
+            // Ignore test-only nodes that other tests leak into the shared 
LOOKUP_MAP via
+            // addTestNode().
+            if (!productionExecNodes.contains(execNode)) {
+                continue;
+            }
+            if (!SKIP_EXEC_NODES.contains(execNode)
+                    && !execNodesWithRestoreTests.contains(execNode)) {
+                fail(
                         "Missing restore test for "
                                 + execNodeNameVersion
                                 + "\nPlease add a restore test for "
-                                + execNode.toString();
-                
Assertions.assertTrue(execNodesWithRestoreTests.contains(execNode), msg);
+                                + execNode.toString());
             }
         }
     }
+
+    private Set<ClassPath.ClassInfo> gatherClasses(Class<?> clazz, String 
packageName)
+            throws IOException {
+        return ClassPath.from(this.getClass().getClassLoader())
+                .getTopLevelClassesRecursive(packageName)
+                .stream()
+                .filter(x -> clazz.isAssignableFrom(x.load()))
+                .collect(Collectors.toSet());
+    }
 }
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-catalog-func/plan/async-correlate-catalog-func.json
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-catalog-func/plan/async-correlate-catalog-func.json
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-catalog-func/plan/async-correlate-catalog-func.json
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-catalog-func/plan/async-correlate-catalog-func.json
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-catalog-func/savepoint/_metadata
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-catalog-func/savepoint/_metadata
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-catalog-func/savepoint/_metadata
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-catalog-func/savepoint/_metadata
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-exception/plan/async-correlate-exception.json
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-exception/plan/async-correlate-exception.json
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-exception/plan/async-correlate-exception.json
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-exception/plan/async-correlate-exception.json
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-exception/savepoint/_metadata
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-exception/savepoint/_metadata
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-exception/savepoint/_metadata
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-exception/savepoint/_metadata
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-join-filter/plan/async-correlate-join-filter.json
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-join-filter/plan/async-correlate-join-filter.json
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-join-filter/plan/async-correlate-join-filter.json
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-join-filter/plan/async-correlate-join-filter.json
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-join-filter/savepoint/_metadata
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-join-filter/savepoint/_metadata
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-join-filter/savepoint/_metadata
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-join-filter/savepoint/_metadata
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-left-join/plan/async-correlate-left-join.json
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-left-join/plan/async-correlate-left-join.json
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-left-join/plan/async-correlate-left-join.json
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-left-join/plan/async-correlate-left-join.json
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-left-join/savepoint/_metadata
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-left-join/savepoint/_metadata
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-left-join/savepoint/_metadata
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-left-join/savepoint/_metadata
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-system-func/plan/async-correlate-system-func.json
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-system-func/plan/async-correlate-system-func.json
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-system-func/plan/async-correlate-system-func.json
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-system-func/plan/async-correlate-system-func.json
diff --git 
a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-system-func/savepoint/_metadata
 
b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-system-func/savepoint/_metadata
similarity index 100%
rename from 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-correlate_1/async-correlate-system-func/savepoint/_metadata
rename to 
flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-async-correlate_1/async-correlate-system-func/savepoint/_metadata

Reply via email to