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

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/main/pr-8298-1b25a0ab536f0ff5d7cee521488b3d356a1f3092
in repository https://gitbox.apache.org/repos/asf/texera.git

commit 577bc408d66802f9ad8905458728ed1526068ca5
Author: carloea2 <[email protected]>
AuthorDate: Thu Sep 24 03:08:12 2026 +0000

    feat(pyamber): report batch processing errors (#8298)
    
    ### What changes were proposed in this PR?
    
    Report per-row MainLoop failures through the worker exception path,
    enter exception pause, and stop consuming the affected batch.
    
    Before: the worker logged the failure, continued later rows, and could
    finish with missing output.
    
    After: the coordinator receives the error and the batch stops at the
    failed row.
    
    ### Any related issues, documentation, discussions?
    
    Closes #8297
    
    ### How was this PR tested?
    
    Added a regression covering the failure report, exception pause, and
    stopped batch consumption.
    
    ```powershell
    $env:PYTHONDONTWRITEBYTECODE='1'; 
C:\Users\carlo\texera\texera\.venv312\Scripts\python.exe -c "import sys,pytest; 
sys.path[:0]=[r'C:\Users\carlo\texera\texera-worktrees\investigate-bug78\amber\src\main\python',r'C:\Users\carlo\texera\texera\amber\src\main\python'];
 raise 
SystemExit(pytest.main([r'amber/src/test/python/core/runnables/test_main_loop.py::TestMainLoop::test_failing_element_reports_exception_and_stops_batch','-p','no:cacheprovider','-q']))"
    C:\Users\carlo\texera\texera\.venv312\Scripts\ruff.exe check 
amber/src/main/python/core/runnables/main_loop.py 
amber/src/test/python/core/runnables/test_main_loop.py
    C:\Users\carlo\texera\texera\.venv312\Scripts\ruff.exe format --check 
amber/src/main/python/core/runnables/main_loop.py 
amber/src/test/python/core/runnables/test_main_loop.py
    ```
    
    The focused regression passed. The full MainLoop test file was also run:
    36 tests passed and 8 pre-existing integration cases failed during
    executor initialization because the shared generated protocol files do
    not yet contain `loop_start_port_uris`.
    
    ### Was this PR authored or co-authored using generative AI tooling?
    
    Generated-by: Codex
---
 amber/src/main/python/core/runnables/main_loop.py  |  4 +-
 .../test/python/core/runnables/test_main_loop.py   | 68 ++++++----------------
 2 files changed, 20 insertions(+), 52 deletions(-)

diff --git a/amber/src/main/python/core/runnables/main_loop.py 
b/amber/src/main/python/core/runnables/main_loop.py
index 6fcae0c26b..2664925fb8 100644
--- a/amber/src/main/python/core/runnables/main_loop.py
+++ b/amber/src/main/python/core/runnables/main_loop.py
@@ -812,7 +812,9 @@ class MainLoop(StoppableQueueBlockingRunnable):
                     self._process_state_frame,
                 )
             except Exception as err:
-                logger.exception(err)
+                self.context.report_exception(err)
+                self._check_exception()
+                return
 
     def _send_console_message(self, console_message: ConsoleMessage):
         self._async_rpc_client.coordinator_stub().console_message_triggered(
diff --git a/amber/src/test/python/core/runnables/test_main_loop.py 
b/amber/src/test/python/core/runnables/test_main_loop.py
index 5ece58166f..a7ad83cf27 100644
--- a/amber/src/test/python/core/runnables/test_main_loop.py
+++ b/amber/src/test/python/core/runnables/test_main_loop.py
@@ -3738,20 +3738,19 @@ class TestMainLoop:
         assert len(console_msgs) == 2
 
     @pytest.mark.timeout(5)
-    def test_a_failing_element_does_not_abandon_the_rest_of_the_batch(
+    def test_failing_element_reports_exception_and_stops_batch(
         self, main_loop, monkeypatch, mock_raw_schema
     ):
-        # The batch iterator is the only handle on the elements still to come,
-        # so letting a failure on one element escape the loop would silently
-        # drop every element behind it. _process_data_element's per-element
-        # backstop keeps iterating instead.
-        #
-        # Deliberately NOT asserted: that nothing is reported to the
-        # coordinator. The backstop only logs, so a runtime-level per-element
-        # failure never reaches Context.report_exception and the workflow can
-        # report success on a short result -- arguably a silent-wrong-result
-        # defect. Pinning that half would cement it, so this test asserts only
-        # that iteration continues and that nothing propagates.
+        console_msgs = []
+        pauses = []
+        monkeypatch.setattr(
+            main_loop, "_send_console_message", lambda msg: 
console_msgs.append(msg)
+        )
+        monkeypatch.setattr(
+            main_loop.context.pause_manager,
+            "pause",
+            lambda reason: pauses.append(reason),
+        )
         schema = Schema(raw_schema=mock_raw_schema)
         port_0 = PortIdentity(0, internal=False)
         channel = self._register_input_port(main_loop, schema, port_0, 
"sender")
@@ -3783,44 +3782,11 @@ class TestMainLoop:
             main_loop.context.statistics_manager, "increase_input_statistics", 
_boom
         )
 
-        # Because the coordinator report is deliberately not asserted (above),
-        # the log line is the swallow's ONLY remaining trace -- so pin it.
-        # `except Exception: pass` is a strictly worse regression than the
-        # defect described above (a short result with no evidence anywhere
-        # rather than a stack trace in the worker log) and is not a defensible
-        # production change, so this is a gap rather than a bug to cement.
-        # The proxy delegates every other level to the real logger so the
-        # module's debug/info calls keep working.
-        from core.runnables import main_loop as main_loop_module
-
-        class _RecordingLogger:
-            def __init__(self, delegate):
-                self.exceptions = []
-                self._delegate = delegate
-
-            def exception(self, err):
-                self.exceptions.append(err)
-
-            def __getattr__(self, name):
-                return getattr(self._delegate, name)
-
-        recorder = _RecordingLogger(main_loop_module.logger)
-        monkeypatch.setattr(main_loop_module, "logger", recorder)
-
-        # Must not raise.
         main_loop._process_data_element(element)
 
-        assert attempted == [0, 1], (
-            "a failure on one element must not abandon the rest of the batch; "
-            f"attempted: {attempted}"
-        )
-        assert [type(err) for err in recorder.exceptions] == [
-            RuntimeError,
-            RuntimeError,
-        ], (
-            "every swallowed per-element failure must leave a log trace; "
-            f"logged: {recorder.exceptions}"
-        )
-        assert all(
-            "statistics backend unavailable" in str(err) for err in 
recorder.exceptions
-        ), f"the logged trace must carry the real error; logged: 
{recorder.exceptions}"
+        assert attempted == [0]
+        assert main_loop.context.exception_manager.has_exception()
+        assert pauses == [PauseType.EXCEPTION_PAUSE]
+        error_msgs = [m for m in console_msgs if m.msg_type == 
ConsoleMessageType.ERROR]
+        assert len(error_msgs) == 1
+        assert "statistics backend unavailable" in error_msgs[0].title

Reply via email to