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

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

commit 91a604795c70ce24b7a95949555d2412289550d1
Author: Sergey Nuyanzin <[email protected]>
AuthorDate: Tue Aug 4 09:06:33 2026 +0200

    [FLINK-40307][table]  was never executed in CI
---
 ...eness.java => RestoreTestCompletenessTest.java} | 84 +++++++++++++++-------
 1 file changed, 60 insertions(+), 24 deletions(-)

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 60%
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..99ac5b09182 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,13 @@
 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.StreamExecProcessTableFunction;
 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 +37,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 +48,34 @@ 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,
+                    // restore tests for delta join and PTF were added in 
later releases
+                    StreamExecDeltaJoin.class,
+                    StreamExecProcessTableFunction.class,
+
+                    // There is jira for these 2 batch tests
+                    // https://issues.apache.org/jira/browse/FLINK-40306
+                    BatchExecHashAggregate.class,
+                    BatchExecNestedLoopJoin.class,
+                    // restore tests for batch sort aggregate was added in 
later releases
+                    BatchExecSortAggregate.class);
 
     private Class<? extends ExecNode<?>> getExecNode(Class<?> restoreTest)
             throws NoSuchMethodException,
@@ -85,7 +104,7 @@ public class RestoreTestCompleteness {
     }
 
     @Test
-    public void testMissingRestoreTest()
+    void testMissingRestoreTest()
             throws IOException,
                     NoSuchMethodException,
                     InstantiationException,
@@ -95,12 +114,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 +137,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());
+    }
 }

Reply via email to