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

fhueske pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git


The following commit(s) were added to refs/heads/master by this push:
     new f2370b9b3f4 [FLINK-40079][table] Reject PTF calls with sys-args if 
they are disabled (#28675)
f2370b9b3f4 is described below

commit f2370b9b3f4366081b8533dbed9805a26590c1ce
Author: Fabian Hueske <[email protected]>
AuthorDate: Mon Jul 20 18:49:24 2026 +0200

    [FLINK-40079][table] Reject PTF calls with sys-args if they are disabled 
(#28675)
    
    * [FLINK-40079][table] Reject PTF calls with sys-args if they are disabled
    
    * Add a check in SqlValidator to reject PTF calls with system-args 
(on_time, uid) in SQL querys if the function disabled them.
    * Add a check in ResolveCallByArgumentsRule to reject system-args in 
functions that disabled them from Table API.
    
    Generated-By: Claude Opus 4.8 (1M context)
---
 .../resolver/rules/ResolveCallByArgumentsRule.java |  3 ++
 .../table/types/inference/SystemTypeInference.java | 28 +++++++++++
 .../planner/calcite/FlinkCalciteSqlValidator.java  | 30 +++++++++++
 .../stream/sql/MLPredictTableFunctionTest.java     | 30 ++++++++++-
 .../plan/stream/sql/ProcessTableFunctionTest.java  | 58 ++++++++++++++++++++++
 .../plan/stream/sql/SnapshotTableFunctionTest.java |  9 ++--
 .../plan/stream/sql/MLPredictTableFunctionTest.xml |  2 +-
 7 files changed, 152 insertions(+), 8 deletions(-)

diff --git 
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/expressions/resolver/rules/ResolveCallByArgumentsRule.java
 
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/expressions/resolver/rules/ResolveCallByArgumentsRule.java
index 5e4a02cb9fb..7364536099a 100644
--- 
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/expressions/resolver/rules/ResolveCallByArgumentsRule.java
+++ 
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/expressions/resolver/rules/ResolveCallByArgumentsRule.java
@@ -299,6 +299,9 @@ final class ResolveCallByArgumentsRule implements 
ResolverRule {
                                 functionName));
             }
 
+            SystemTypeInference.checkNoSystemArguments(
+                    inference.disableSystemArguments(), namedArgs.keySet(), 
functionName);
+
             fillInDefaultNamedArguments(declaredArgs, namedArgs);
             fillInPtfSpecificNamedArguments(
                     functionName, definition, declaredArgs, namedArgs, 
actualArgs);
diff --git 
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java
 
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java
index 971831d6ac4..a1f9fbc5a9f 100644
--- 
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java
+++ 
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java
@@ -46,6 +46,7 @@ import javax.annotation.Nullable;
 
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collection;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
@@ -126,6 +127,33 @@ public class SystemTypeInference {
         return !UID_FORMAT.test(uid);
     }
 
+    /**
+     * Rejects the implicit system arguments ({@code on_time}, {@code uid}) 
for a function that
+     * disables them via {@link TypeInference#disableSystemArguments()}.
+     *
+     * <p>The system arguments are not part of such a function's signature. 
Enforcing this from
+     * every translation path (SQL operand checking and Table API call 
resolution) rejects them
+     * consistently, regardless of whether the function is processed by the 
generic PTF rule or a
+     * dedicated optimizer rule (e.g. ML_PREDICT, LATERAL SNAPSHOT).
+     */
+    public static void checkNoSystemArguments(
+            boolean sysArgsDisabled,
+            Collection<String> suppliedArgumentNames,
+            String functionName) {
+        if (!sysArgsDisabled) {
+            return;
+        }
+        for (StaticArgument systemArg : PROCESS_TABLE_FUNCTION_SYSTEM_ARGS) {
+            if (suppliedArgumentNames.contains(systemArg.getName())) {
+                throw new ValidationException(
+                        String.format(
+                                "Invalid function call. The '%s' argument is 
not supported "
+                                        + "because function '%s' does not use 
system arguments.",
+                                systemArg.getName(), functionName));
+            }
+        }
+    }
+
     // 
--------------------------------------------------------------------------------------------
 
     private static void checkScalarArgsOnly(List<StaticArgument> defaultArgs) {
diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkCalciteSqlValidator.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkCalciteSqlValidator.java
index 651a86925c3..0eb925352f0 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkCalciteSqlValidator.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkCalciteSqlValidator.java
@@ -30,9 +30,11 @@ import org.apache.flink.table.data.TimestampData;
 import org.apache.flink.table.functions.FunctionKind;
 import org.apache.flink.table.planner.catalog.CatalogSchemaModel;
 import org.apache.flink.table.planner.catalog.CatalogSchemaTable;
+import org.apache.flink.table.planner.functions.bridging.BridgingSqlFunction;
 import org.apache.flink.table.planner.plan.FlinkCalciteCatalogReader;
 import org.apache.flink.table.planner.plan.utils.FlinkRexUtil;
 import org.apache.flink.table.planner.utils.ShortcutUtils;
+import org.apache.flink.table.types.inference.SystemTypeInference;
 import org.apache.flink.table.types.logical.DecimalType;
 
 import org.apache.calcite.plan.RelOptCluster;
@@ -89,6 +91,7 @@ import java.math.BigDecimal;
 import java.time.ZoneId;
 import java.util.ArrayList;
 import java.util.Collections;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
@@ -383,6 +386,7 @@ public final class FlinkCalciteSqlValidator extends 
FlinkSqlParsingValidator {
 
         final SqlBasicCall call = (SqlBasicCall) node;
         checkNoNamedAndPositionalMixedArgs(call);
+        checkDisabledSystemArgs(call);
 
         // Special case for MODEL
         if (node instanceof SqlExplicitModelCall) {
@@ -456,6 +460,32 @@ public final class FlinkCalciteSqlValidator extends 
FlinkSqlParsingValidator {
         }
     }
 
+    /**
+     * Rejects the implicit PTF system arguments (on_time, uid) for functions 
that disable them.
+     *
+     * <p>This must happen before Calcite permutes named arguments, because 
unknown named arguments
+     * are silently dropped during permutation and would otherwise be lost. 
The actual rule and
+     * error message live in {@link 
SystemTypeInference#checkNoSystemArguments} so that the Table
+     * API path (which resolves calls without this validator) enforces it 
identically.
+     */
+    private static void checkDisabledSystemArgs(SqlBasicCall call) {
+        final SqlOperator operator = call.getOperator();
+        if (!(operator instanceof BridgingSqlFunction)
+                || !((BridgingSqlFunction) 
operator).getTypeInference().disableSystemArguments()) {
+            return;
+        }
+        final Set<String> suppliedArgNames = new HashSet<>();
+        for (SqlNode operand : call.getOperandList()) {
+            if (operand != null && operand.getKind() == 
SqlKind.ARGUMENT_ASSIGNMENT) {
+                final SqlNode nameNode = ((SqlCall) operand).operand(1);
+                if (nameNode instanceof SqlIdentifier) {
+                    suppliedArgNames.add(((SqlIdentifier) 
nameNode).getSimple());
+                }
+            }
+        }
+        SystemTypeInference.checkNoSystemArguments(true, suppliedArgNames, 
operator.getName());
+    }
+
     @Override
     public SqlNode maybeCast(SqlNode node, RelDataType currentType, 
RelDataType desiredType) {
         return super.maybeCast(node, currentType, desiredType);
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.java
index b409e02cf9d..93a9c5b937d 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.java
@@ -30,6 +30,7 @@ import org.junit.jupiter.api.Test;
 
 import java.util.Collections;
 
+import static org.apache.flink.core.testutils.FlinkAssertions.anyCauseMatches;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
 /** Tests for model table value function in stream mode. */
@@ -76,8 +77,7 @@ public class MLPredictTableFunctionTest extends 
MLPredictTableFunctionTestBase {
     @Test
     public void testInputTableIsInsertOnlyStream() {
         String sql =
-                "SELECT *\n"
-                        + "FROM TABLE(ML_PREDICT(TABLE MyTable, MODEL MyModel, 
DESCRIPTOR(a, b)))";
+                "SELECT *\n" + "FROM ML_PREDICT(TABLE MyTable, MODEL MyModel, 
DESCRIPTOR(a, b))";
         util.verifyRelPlan(
                 sql,
                 JavaScalaConversionUtil.toScala(
@@ -103,4 +103,30 @@ public class MLPredictTableFunctionTest extends 
MLPredictTableFunctionTestBase {
                 .hasMessageContaining(
                         "StreamPhysicalMLPredictTableFunction doesn't support 
consuming update and delete changes which is produced by node 
TableSourceScan(table=[[default_catalog, default_database, CdcTable]], 
fields=[a, b])");
     }
+
+    @Test
+    void testOnTimeArgumentNotAllowed() {
+        // ML_PREDICT disables the implicit system arguments. Supplying 
`on_time` must be rejected
+        // at the SQL level even though ML_PREDICT is handled by a dedicated 
optimizer rule.
+        String sql =
+                "SELECT * FROM ML_PREDICT(INPUT => TABLE MyTable, MODEL => 
MODEL MyModel, "
+                        + "ARGS => DESCRIPTOR(a, b), on_time => 
DESCRIPTOR(rowtime))";
+        assertThatThrownBy(() -> util.verifyRelPlan(sql))
+                .satisfies(
+                        anyCauseMatches(
+                                "The 'on_time' argument is not supported 
because function "
+                                        + "'ML_PREDICT' does not use system 
arguments."));
+    }
+
+    @Test
+    void testUidArgumentNotAllowed() {
+        String sql =
+                "SELECT * FROM ML_PREDICT(INPUT => TABLE MyTable, MODEL => 
MODEL MyModel, "
+                        + "ARGS => DESCRIPTOR(a, b), uid => 'my-uid')";
+        assertThatThrownBy(() -> util.verifyRelPlan(sql))
+                .satisfies(
+                        anyCauseMatches(
+                                "The 'uid' argument is not supported because 
function "
+                                        + "'ML_PREDICT' does not use system 
arguments."));
+    }
 }
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java
index d34b3f3bc8d..7500bef9c98 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java
@@ -76,6 +76,7 @@ import static 
org.apache.flink.table.annotation.ArgumentTrait.ROW_SEMANTIC_TABLE
 import static 
org.apache.flink.table.annotation.ArgumentTrait.SET_SEMANTIC_TABLE;
 import static org.apache.flink.table.annotation.ArgumentTrait.SUPPORT_UPDATES;
 import static org.apache.flink.table.api.Expressions.$;
+import static org.apache.flink.table.api.Expressions.lit;
 import static org.apache.flink.table.api.Expressions.row;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
@@ -283,6 +284,63 @@ class ProcessTableFunctionTest extends TableTestBase {
                                 "Disabling system arguments is not supported 
for user-defined PTF."));
     }
 
+    @Test
+    void testOnTimeArgRejectedForDisabledPtf() {
+        util.addTemporarySystemFunction("f", NoSystemArgsTableFunction.class);
+        assertThatThrownBy(
+                        () ->
+                                util.verifyRelPlan(
+                                        "SELECT * FROM f(r => TABLE 
t_watermarked, i => 1, "
+                                                + "on_time => 
DESCRIPTOR(ts));"))
+                .satisfies(
+                        anyCauseMatches(
+                                "The 'on_time' argument is not supported 
because function "
+                                        + "'f' does not use system 
arguments."));
+    }
+
+    @Test
+    void testUidArgRejectedForDisabledPtf() {
+        util.addTemporarySystemFunction("f", NoSystemArgsScalarFunction.class);
+        assertThatThrownBy(() -> util.verifyRelPlan("SELECT * FROM f(i => 1, 
uid => 'my-uid');"))
+                .satisfies(
+                        anyCauseMatches(
+                                "The 'uid' argument is not supported because 
function "
+                                        + "'f' does not use system 
arguments."));
+    }
+
+    @Test
+    void testSystemArgRejectedByNameBeforeTypeCheck() {
+        // System arguments are rejected by name, rather than a type mismatch.
+        util.addTemporarySystemFunction("f", NoSystemArgsTableFunction.class);
+        assertThatThrownBy(
+                        () ->
+                                util.verifyRelPlan(
+                                        "SELECT * FROM f(r => TABLE t, i => 1, 
on_time => 1);"))
+                .satisfies(
+                        anyCauseMatches(
+                                "The 'on_time' argument is not supported 
because function "
+                                        + "'f' does not use system 
arguments."));
+    }
+
+    @Test
+    void testSystemArgRejectedForDisabledPtfViaTableApi() {
+        // The same enforcement applies to the Table API path, which resolves 
calls via
+        // ResolveCallByArgumentsRule instead of the SQL validator.
+        util.addTemporarySystemFunction("f", NoSystemArgsTableFunction.class);
+        assertThatThrownBy(
+                        () ->
+                                util.tableEnv()
+                                        .fromCall(
+                                                "f",
+                                                
util.tableEnv().from("t").asArgument("r"),
+                                                lit(1).asArgument("i"),
+                                                
lit("my-uid").asArgument("uid")))
+                .satisfies(
+                        anyCauseMatches(
+                                "The 'uid' argument is not supported because 
function "
+                                        + "'f' does not use system 
arguments."));
+    }
+
     @Test
     void testUidPipelineSplitIntoTwoFunctions() {
         util.addTemporarySystemFunction("f", SetSemanticTableFunction.class);
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/SnapshotTableFunctionTest.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/SnapshotTableFunctionTest.java
index 6d254b220b8..eef09609aeb 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/SnapshotTableFunctionTest.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/SnapshotTableFunctionTest.java
@@ -24,7 +24,6 @@ import org.apache.flink.table.planner.utils.TableTestBase;
 import org.apache.flink.table.planner.utils.TableTestUtil;
 
 import org.junit.jupiter.api.BeforeEach;
-import org.junit.jupiter.api.Disabled;
 import org.junit.jupiter.api.Test;
 
 import static org.apache.flink.core.testutils.FlinkAssertions.anyCauseMatches;
@@ -121,9 +120,6 @@ public class SnapshotTableFunctionTest extends 
TableTestBase {
     }
 
     @Test
-    @Disabled(
-            "SNAPSHOT sets disableSystemArguments(true), but that flag is 
currently not enforced. "
-                    + "Re-enable once FLINK-40079 is fixed.")
     void testSystemArgumentsNotAllowed() {
         // SNAPSHOT disables the implicit system arguments (e.g. `on_time`). 
Passing one in a
         // LATERAL context must be rejected because the argument is not part 
of the function
@@ -136,7 +132,10 @@ public class SnapshotTableFunctionTest extends 
TableTestBase {
                                                 + "input => TABLE Rates, "
                                                 + "on_time => 
DESCRIPTOR(rate_time))) AS r "
                                                 + "WHERE o.currency = 
r.currency"))
-                .satisfies(anyCauseMatches("on_time"));
+                .satisfies(
+                        anyCauseMatches(
+                                "The 'on_time' argument is not supported 
because function "
+                                        + "'SNAPSHOT' does not use system 
arguments."));
     }
 
     @Test
diff --git 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.xml
 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.xml
index b1d51479529..778451f067e 100644
--- 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.xml
+++ 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.xml
@@ -869,7 +869,7 @@ Calc(select=[a, b, c, d, rowtime, 
PROCTIME_MATERIALIZE(proctime) AS proctime, c0
   <TestCase name="testInputTableIsInsertOnlyStream">
     <Resource name="sql">
       <![CDATA[SELECT *
-FROM TABLE(ML_PREDICT(TABLE MyTable, MODEL MyModel, DESCRIPTOR(a, b)))]]>
+FROM ML_PREDICT(TABLE MyTable, MODEL MyModel, DESCRIPTOR(a, b))]]>
     </Resource>
     <Resource name="ast">
       <![CDATA[

Reply via email to