This is an automated email from the ASF dual-hosted git repository.
wenjin272 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-agents.git
The following commit(s) were added to refs/heads/main by this push:
new b551f7a2 [test][python] Use the registered checkpoint interval key in
the Python e2e tests (#978)
b551f7a2 is described below
commit b551f7a234c1744b133f579bce1c0b7fb24d2891
Author: Weiqing Yang <[email protected]>
AuthorDate: Sat Aug 8 00:45:53 2026 -0700
[test][python] Use the registered checkpoint interval key in the Python e2e
tests (#978)
Generated-by: Claude Code 2.1.224
---
.../checkpointing_config_test.py | 58 ++++++++++++++++++++++
.../e2e_tests_mcp/mcp_test.py | 2 +-
.../e2e_tests_integration/execute_test.py | 12 ++---
.../flink_intergration_test.py | 2 +-
.../e2e_tests_integration/mock_chat_model_test.py | 4 +-
.../workflow_memory_remote_test.py | 2 +-
6 files changed, 69 insertions(+), 11 deletions(-)
diff --git
a/python/flink_agents/e2e_tests/e2e_tests_integration/checkpointing_config_test.py
b/python/flink_agents/e2e_tests/e2e_tests_integration/checkpointing_config_test.py
new file mode 100644
index 00000000..7270649f
--- /dev/null
+++
b/python/flink_agents/e2e_tests/e2e_tests_integration/checkpointing_config_test.py
@@ -0,0 +1,58 @@
+################################################################################
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#################################################################################
+"""Guard for the checkpointing configuration the Flink e2e tests set."""
+
+import re
+from pathlib import Path
+
+from pyflink.java_gateway import get_gateway
+
+_E2E_TESTS_PACKAGE = Path(__file__).parents[1]
+
+# Deliberately loose on setter name, stem, and quote style: a near miss Flink
+# would silently ignore has to be caught whatever spelling it arrives in.
+_CHECKPOINT_INTERVAL_KEY = re.compile(
+ r"set_\w+\(\s*(?P<quote>['\"])"
+ r"(?P<key>[^'\"]*checkpoint[^'\"]*interval[^'\"]*)(?P=quote)"
+)
+
+
+def test_checkpoint_interval_key_matches_the_flink_option() -> None:
+ """Every checkpoint interval the e2e tests set uses the key Flink
registers.
+
+ Flink ignores an unregistered configuration key silently, with no exception
+ and no log line, so a near miss such as ``checkpointing.interval`` leaves
+ checkpointing disabled while the test still passes. Reading the expected
key
+ off ``CheckpointingOptions#CHECKPOINTING_INTERVAL`` takes it from Flink
+ itself rather than from another hand-written copy of the string.
+ """
+ flink_config = get_gateway().jvm.org.apache.flink.configuration
+ expected = flink_config.CheckpointingOptions.CHECKPOINTING_INTERVAL.key()
+
+ found: dict[str, set[str]] = {}
+ for path in sorted(_E2E_TESTS_PACKAGE.rglob("*_test.py")):
+ for match in _CHECKPOINT_INTERVAL_KEY.finditer(path.read_text()):
+ found.setdefault(match["key"], set()).add(path.name)
+
+ assert found, "no checkpoint interval configured; the pattern above is
stale"
+ assert set(found) == {expected}, (
+ f"checkpointing must be configured with {expected!r}, found "
+ + "; ".join(
+ f"{key!r} in {sorted(files)}" for key, files in
sorted(found.items())
+ )
+ )
diff --git
a/python/flink_agents/e2e_tests/e2e_tests_integration/e2e_tests_mcp/mcp_test.py
b/python/flink_agents/e2e_tests/e2e_tests_integration/e2e_tests_mcp/mcp_test.py
index 7a8e8a72..2c625c56 100644
---
a/python/flink_agents/e2e_tests/e2e_tests_integration/e2e_tests_mcp/mcp_test.py
+++
b/python/flink_agents/e2e_tests/e2e_tests_integration/e2e_tests_mcp/mcp_test.py
@@ -211,7 +211,7 @@ def test_mcp(
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
diff --git
a/python/flink_agents/e2e_tests/e2e_tests_integration/execute_test.py
b/python/flink_agents/e2e_tests/e2e_tests_integration/execute_test.py
index edd2d4b9..08405f14 100644
--- a/python/flink_agents/e2e_tests/e2e_tests_integration/execute_test.py
+++ b/python/flink_agents/e2e_tests/e2e_tests_integration/execute_test.py
@@ -56,7 +56,7 @@ def test_durable_execute_basic_flink(tmp_path: Path) -> None:
"""Test basic synchronous durable_execute() functionality in Flink
environment."""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
@@ -108,7 +108,7 @@ def test_durable_execute_multiple_calls_flink(tmp_path:
Path) -> None:
"""Test multiple durable_execute() calls in Flink environment."""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
@@ -160,7 +160,7 @@ def test_durable_execute_with_async_flink(tmp_path: Path)
-> None:
"""Test durable_execute() and durable_execute_async() in Flink
environment."""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
@@ -212,7 +212,7 @@ def test_durable_execute_async_exception_flink(tmp_path:
Path) -> None:
"""Test durable_execute_async() exception handling in Flink environment."""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
@@ -264,7 +264,7 @@ def test_durable_execute_sync_exception_flink(tmp_path:
Path) -> None:
"""Test synchronous durable_execute() exception handling in Flink
environment."""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
@@ -316,7 +316,7 @@ def test_durable_execute_with_kwargs_flink(tmp_path: Path)
-> None:
"""Test durable_execute() with keyword arguments in Flink environment."""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
diff --git
a/python/flink_agents/e2e_tests/e2e_tests_integration/flink_intergration_test.py
b/python/flink_agents/e2e_tests/e2e_tests_integration/flink_intergration_test.py
index 4d8302ce..2bc1887c 100644
---
a/python/flink_agents/e2e_tests/e2e_tests_integration/flink_intergration_test.py
+++
b/python/flink_agents/e2e_tests/e2e_tests_integration/flink_intergration_test.py
@@ -51,7 +51,7 @@ os.environ["PYTHONPATH"] = sysconfig.get_paths()["purelib"]
def test_from_datastream_to_datastream(tmp_path: Path) -> None:
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
diff --git
a/python/flink_agents/e2e_tests/e2e_tests_integration/mock_chat_model_test.py
b/python/flink_agents/e2e_tests/e2e_tests_integration/mock_chat_model_test.py
index dd7bc3c7..cc2aad26 100644
---
a/python/flink_agents/e2e_tests/e2e_tests_integration/mock_chat_model_test.py
+++
b/python/flink_agents/e2e_tests/e2e_tests_integration/mock_chat_model_test.py
@@ -57,7 +57,7 @@ def test_built_in_chat_tool_action_content(tmp_path: Path) ->
None:
"""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
@@ -115,7 +115,7 @@ def test_chat_model_get_resource_in_action(tmp_path: Path)
-> None:
"""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
diff --git
a/python/flink_agents/e2e_tests/e2e_tests_integration/workflow_memory_remote_test.py
b/python/flink_agents/e2e_tests/e2e_tests_integration/workflow_memory_remote_test.py
index 93f5acbf..6fbd3f30 100644
---
a/python/flink_agents/e2e_tests/e2e_tests_integration/workflow_memory_remote_test.py
+++
b/python/flink_agents/e2e_tests/e2e_tests_integration/workflow_memory_remote_test.py
@@ -55,7 +55,7 @@ def test_short_term_memory_same_key_accumulation(tmp_path:
Path) -> None:
"""
config = Configuration()
config.set_string("state.backend.type", "rocksdb")
- config.set_string("checkpointing.interval", "1s")
+ config.set_string("execution.checkpointing.interval", "1s")
config.set_string("restart-strategy.type", "disable")
env = StreamExecutionEnvironment.get_execution_environment(config)
env.set_runtime_mode(RuntimeExecutionMode.STREAMING)