This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new 511bc27695 test(amber): cover zero-tuple-channel worker completion in
pyamber (#6544)
511bc27695 is described below
commit 511bc276953a251f668329b9cac1b121c7651c69
Author: Xinyuan Lin <[email protected]>
AuthorDate: Mon Jul 20 12:57:47 2026 -0700
test(amber): cover zero-tuple-channel worker completion in pyamber (#6544)
### What changes were proposed in this PR?
Adds two regression tests to
`amber/src/test/python/core/runnables/test_main_loop.py`
covering a Python-UDF worker whose input port receives **zero tuples**
across a
`StartChannel` → `EndChannel` bracket — e.g. the untaken branch of an
`If` operator read
through an `InputPortMaterializationReaderRunnable`, or a filter that
drops everything on a
materialized edge:
- **`test_zero_tuple_channel_completes_worker`** — `StartChannel` then
`EndChannel` with no
data tuples must still complete both ports and report
`PortCompleted(input)`,
`PortCompleted(output)`, and `WorkerExecutionCompleted`, forwarding both
channel ECMs
downstream.
- **`test_empty_on_finish_after_tuples_completes_worker`** — an operator
whose `on_finish`
yields nothing (like `BatchOperator` when the input size is an exact
multiple of
`BATCH_SIZE`) must still complete cleanly after processing tuples, since
the empty
`on_finish` is exhausted within a single MainLoop↔DataProcessor
hand-off.
Both run the MainLoop on a **daemon thread** and detect completion by
**deadline-polling
`WorkerState.COMPLETED`**, so a regression that deadlocks the handshake
fails the test
cleanly instead of hanging the whole pytest process.
This is **test-only** — no production changes. It pins behavior the
MainLoop↔DataProcessor
handshake already handles on `main` (the internal marker is dispatched
in `_process_ecm`
without an intervening context switch; the worker state graph permits
`READY → COMPLETED`,
which a zero-tuple worker needs because it never enters `RUNNING`) so it
can't silently
regress.
### Any related issues, documentation, discussions?
Closes #6543
### How was this PR tested?
```
cd amber && python -m pytest
src/test/python/core/runnables/test_main_loop.py \
-k "zero_tuple_channel or empty_on_finish"
# 2 passed, 21 deselected
cd amber && python -m pytest
src/test/python/core/runnables/test_main_loop.py
# 23 passed
```
- Both new tests pass on `main`; the whole `test_main_loop.py` stays
green (23 passed).
- `ruff check` → all checks passed; `ruff format --check` → already
formatted.
- Guard-verified that each test actually catches a regression (mutation
reverted after each):
- reintroducing the `_process_ecm` pre-marker `self._switch_context()` →
`test_zero_tuple_channel_completes_worker` deadlocks (deadline-fail);
- removing `COMPLETED` from the `READY` set in the worker state graph →
`test_zero_tuple_channel_completes_worker` fails with
`InvalidTransitionException`;
- gating `finished_current.set()` behind an "`on_finish` produced ≥1
output" flag →
`test_empty_on_finish_after_tuples_completes_worker` hangs
(deadline-fail), while
Echo-based tests stay green.
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Opus 4.8 [1M context])
---
.../test/python/core/runnables/test_main_loop.py | 375 +++++++++++++++++++++
1 file changed, 375 insertions(+)
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 79d372e2c0..5884edf27b 100644
--- a/amber/src/test/python/core/runnables/test_main_loop.py
+++ b/amber/src/test/python/core/runnables/test_main_loop.py
@@ -22,6 +22,7 @@ import pyarrow
import pytest
import sys
import time
+import uuid
from threading import Thread
from core.models import (
@@ -86,6 +87,7 @@ from proto.org.apache.texera.amber.engine.architecture.worker
import (
from proto.org.apache.texera.amber.engine.common import
DirectControlMessagePayloadV2
from pytexera.udf.examples.count_batch_operator import CountBatchOperator
from pytexera.udf.examples.echo_operator import EchoOperator
+from pytexera.udf.udf_operator import UDFOperatorV2
class _FalseLoopEnd(LoopEndOperator):
@@ -93,6 +95,22 @@ class _FalseLoopEnd(LoopEndOperator):
return False
+class EmptyOnFinishOperator(UDFOperatorV2):
+ # Echoes each input tuple, but its on_finish is a zero-yield generator
+ # (`return` before `yield` makes the body unreachable while still marking
+ # the function as a generator). This is the BatchOperator-fed-an-exact-
+ # multiple-of-BATCH_SIZE shape: the EndChannel on_finish produces NOTHING,
+ # so DataProcessor._set_output_tuple exhausts the iterator in a single
+ # hand-off (no per-output switch dance) and sets finished_current straight
+ # away. MainLoop must not lose that completion signal.
+ def process_tuple(self, tuple_, port):
+ yield tuple_
+
+ def on_finish(self, port):
+ return
+ yield
+
+
class TestMainLoop:
@pytest.fixture
def command_sequence(self):
@@ -268,6 +286,58 @@ class TestMainLoop:
),
)
+ @pytest.fixture
+ def mock_start_channel(self, mock_data_input_channel):
+ # Mirror of mock_end_of_upstream but a StartChannel ECM with
+ # NO_ALIGNMENT (the alignment a real StartChannel bracket uses).
+ return ECMElement(
+ tag=mock_data_input_channel,
+ payload=EmbeddedControlMessage(
+ EmbeddedControlMessageIdentity("StartChannel"),
+ EmbeddedControlMessageType.NO_ALIGNMENT,
+ [],
+ {
+ mock_data_input_channel.to_worker_id.name:
ControlInvocation(
+ "StartChannel",
+ ControlRequest(empty_request=EmptyRequest()),
+ AsyncRpcContext(ActorVirtualIdentity(),
ActorVirtualIdentity()),
+ -1,
+ )
+ },
+ ),
+ )
+
+ @pytest.fixture
+ def mock_initialize_empty_on_finish_executor(
+ self,
+ mock_control_input_channel,
+ mock_sender_actor,
+ mock_link,
+ command_sequence,
+ mock_raw_schema,
+ ):
+ operator_code = "from pytexera import *\n" + inspect.getsource(
+ EmptyOnFinishOperator
+ )
+ command = set_one_of(
+ ControlRequest,
+ InitializeExecutorRequest(
+ op_exec_init_info=set_one_of(
+ OpExecInitInfo, OpExecWithCode(operator_code, "python")
+ ),
+ is_source=False,
+ ),
+ )
+ payload = set_one_of(
+ DirectControlMessagePayloadV2,
+ ControlInvocation(
+ method_name="InitializeExecutor",
+ command_id=command_sequence,
+ command=command,
+ ),
+ )
+ return DCMElement(tag=mock_control_input_channel, payload=payload)
+
@pytest.fixture
def input_queue(self):
return InternalQueue()
@@ -1684,6 +1754,311 @@ class TestMainLoop:
reraise()
+ @staticmethod
+ def _expected_port_completed_dcm(
+ mock_control_output_channel, command_id, port_id, is_input
+ ):
+ return DCMElement(
+ tag=mock_control_output_channel,
+ payload=DirectControlMessagePayloadV2(
+ control_invocation=ControlInvocation(
+ method_name="PortCompleted",
+ command_id=command_id,
+ context=AsyncRpcContext(
+ sender=ActorVirtualIdentity(name="dummy_worker_id"),
+ receiver=ActorVirtualIdentity(name="COORDINATOR"),
+ ),
+ command=ControlRequest(
+ port_completed_request=PortCompletedRequest(
+ port_id=port_id, input=is_input
+ )
+ ),
+ )
+ ),
+ )
+
+ @staticmethod
+ def _expected_worker_completed_dcm(mock_control_output_channel):
+ return DCMElement(
+ tag=mock_control_output_channel,
+ payload=DirectControlMessagePayloadV2(
+ control_invocation=ControlInvocation(
+ method_name="WorkerExecutionCompleted",
+ command_id=2,
+ context=AsyncRpcContext(
+ sender=ActorVirtualIdentity(name="dummy_worker_id"),
+ receiver=ActorVirtualIdentity(name="COORDINATOR"),
+ ),
+ command=ControlRequest(empty_request=EmptyRequest()),
+ )
+ ),
+ )
+
+ @staticmethod
+ def _forwarded_ecm(mock_data_output_channel, method_name, alignment):
+ return ECMElement(
+ tag=mock_data_output_channel,
+ payload=EmbeddedControlMessage(
+ EmbeddedControlMessageIdentity(method_name),
+ alignment,
+ [],
+ {
+ mock_data_output_channel.to_worker_id.name:
ControlInvocation(
+ method_name,
+ ControlRequest(empty_request=EmptyRequest()),
+ AsyncRpcContext(ActorVirtualIdentity(),
ActorVirtualIdentity()),
+ -1,
+ )
+ },
+ ),
+ )
+
+ @staticmethod
+ def _drain_until(output_queue, done, timeout=15.0):
+ # Non-blocking drain of the output queue against a deadline. A
+ # regression that deadlocks the MainLoop/DataProcessor handshake never
+ # satisfies `done`, so we return the partial batch at the deadline and
+ # let the caller pytest.fail() -- the whole pytest process is never
+ # hung because the worker runs on a daemon thread.
+ deadline = time.time() + timeout
+ collected = []
+ while time.time() < deadline:
+ while output_queue.size() > 0:
+ collected.append(output_queue.get())
+ if done(collected):
+ return collected
+ time.sleep(0.005)
+ return collected
+
+ @pytest.mark.timeout(30)
+ def test_zero_tuple_channel_completes_worker(
+ self,
+ mock_link,
+ mock_data_output_channel,
+ mock_control_output_channel,
+ input_queue,
+ output_queue,
+ main_loop,
+ main_loop_thread,
+ mock_assign_input_port,
+ mock_assign_output_port,
+ mock_add_input_channel,
+ mock_add_partitioning,
+ mock_initialize_executor,
+ mock_start_channel,
+ mock_end_of_upstream,
+ command_sequence,
+ reraise,
+ ):
+ # A worker whose input port receives a StartChannel->EndChannel bracket
+ # with ZERO DataElements (the untaken branch of an If read through an
+ # InputPortMaterializationReaderRunnable, or a filter that drops
+ # everything on a materialized edge) must still process both ECMs and
+ # reach COMPLETED. Two sub-bugs made this hang/crash on the old
+ # loop-feb branch:
+ # (1) DEADLOCK: MainLoop._process_ecm re-read current_internal_marker
+ # AFTER a _switch_context(), so the DataProcessor could pop the
+ # marker mid-switch and MainLoop skipped _process_end_channel ->
+ # both threads park forever.
+ # (2) STATE GRAPH: a zero-tuple worker never enters RUNNING (only
+ # _process_data_element does that), so completion is a direct
+ # READY -> COMPLETED transition, which the transition graph must
+ # permit.
+ # Run on a daemon thread and detect completion by deadline-polling the
+ # state manager so a regression deadlock fails cleanly instead of
+ # hanging the whole pytest process.
+ main_loop_thread.daemon = True
+ main_loop_thread.start()
+
+ for setup_msg in [
+ mock_assign_input_port,
+ mock_assign_output_port,
+ mock_add_input_channel,
+ mock_add_partitioning,
+ mock_initialize_executor,
+ ]:
+ input_queue.put(setup_msg)
+ assert output_queue.get() == DCMElement(
+ tag=mock_control_output_channel,
+ payload=DirectControlMessagePayloadV2(
+ return_invocation=ReturnInvocation(
+ command_id=command_sequence,
+ return_value=ControlReturn(empty_return=EmptyReturn()),
+ )
+ ),
+ )
+
+ # The worker is READY here and never enters RUNNING (no data element).
+ assert main_loop.context.state_manager.confirm_state(WorkerState.READY)
+
+ # Zero tuples between StartChannel and EndChannel.
+ input_queue.put(mock_start_channel)
+ input_queue.put(mock_end_of_upstream)
+
+ expected_worker_completed = self._expected_worker_completed_dcm(
+ mock_control_output_channel
+ )
+ collected = self._drain_until(
+ output_queue,
+ lambda items: expected_worker_completed in items,
+ )
+
+ if not
main_loop.context.state_manager.confirm_state(WorkerState.COMPLETED):
+ pytest.fail(
+ "zero-tuple worker did not reach COMPLETED within the deadline
"
+ "-- likely the _process_ecm marker-after-switch deadlock or a "
+ "missing READY->COMPLETED transition. "
+ f"state={main_loop.context.state_manager.get_current_state()},
"
+ f"collected={collected}"
+ )
+
+ # Both the input and output ports complete, and the worker signals
+ # WorkerExecutionCompleted -- all on the coordinator control channel.
+ expected_input_port_completed = self._expected_port_completed_dcm(
+ mock_control_output_channel, 0, mock_link.to_port_id, True
+ )
+ expected_output_port_completed = self._expected_port_completed_dcm(
+ mock_control_output_channel, 1, PortIdentity(id=0), False
+ )
+ assert expected_input_port_completed in collected
+ assert expected_output_port_completed in collected
+ assert expected_worker_completed in collected
+
+ # Both ECMs are forwarded downstream on the data output channel.
+ assert (
+ self._forwarded_ecm(
+ mock_data_output_channel,
+ "StartChannel",
+ EmbeddedControlMessageType.NO_ALIGNMENT,
+ )
+ in collected
+ )
+ assert (
+ self._forwarded_ecm(
+ mock_data_output_channel,
+ "EndChannel",
+ EmbeddedControlMessageType.PORT_ALIGNMENT,
+ )
+ in collected
+ )
+
+ reraise()
+
+ @pytest.mark.timeout(30)
+ def test_empty_on_finish_after_tuples_completes_worker(
+ self,
+ mock_link,
+ mock_tuple,
+ mock_data_output_channel,
+ mock_control_output_channel,
+ input_queue,
+ output_queue,
+ main_loop,
+ main_loop_thread,
+ mock_assign_input_port,
+ mock_assign_output_port,
+ mock_add_input_channel,
+ mock_add_partitioning,
+ mock_initialize_empty_on_finish_executor,
+ mock_data_element,
+ mock_end_of_upstream,
+ command_sequence,
+ monkeypatch,
+ reraise,
+ ):
+ # Sibling case: after processing real tuples, an EndChannel whose
+ # on_finish yields NOTHING must also complete cleanly. The empty
+ # on_finish is exhausted inside a single hand-off (DataProcessor
+ # ._set_output_tuple runs no per-output switch dance, it just sets
+ # finished_current), and MainLoop must not lose the completion signal.
+
+ # Guard the udf-v1 executor-module-contamination landmine: force a
+ # unique module name so cross-test importlib caching can't hand us a
+ # stale operator class. (main's ExecutorManager already uses a
+ # process-wide unique counter, so this is belt-and-suspenders.)
+ unique_name = f"udf_empty_on_finish_{uuid.uuid4().hex}"
+ monkeypatch.setattr(
+ main_loop.context.executor_manager,
+ "gen_module_file_name",
+ lambda: (unique_name, f"{unique_name}.py"),
+ )
+
+ main_loop_thread.daemon = True
+ main_loop_thread.start()
+
+ for setup_msg in [
+ mock_assign_input_port,
+ mock_assign_output_port,
+ mock_add_input_channel,
+ mock_add_partitioning,
+ mock_initialize_empty_on_finish_executor,
+ ]:
+ input_queue.put(setup_msg)
+ assert output_queue.get() == DCMElement(
+ tag=mock_control_output_channel,
+ payload=DirectControlMessagePayloadV2(
+ return_invocation=ReturnInvocation(
+ command_id=command_sequence,
+ return_value=ControlReturn(empty_return=EmptyReturn()),
+ )
+ ),
+ )
+
+ # The loaded executor must be our zero-yield-on_finish operator, not a
+ # stale cached class from another test.
+ assert (
+ type(main_loop.context.executor_manager.executor).__name__
+ == "EmptyOnFinishOperator"
+ )
+
+ # One real tuple: the operator echoes it and the worker enters RUNNING.
+ input_queue.put(mock_data_element)
+ echoed: DataElement = output_queue.get()
+ assert echoed.tag == mock_data_output_channel
+ assert isinstance(echoed.payload, DataFrame)
+ assert Tuple(echoed.payload.frame.to_pylist()[0]) == mock_tuple
+
+ # EndChannel with an empty on_finish must still complete the worker.
+ input_queue.put(mock_end_of_upstream)
+
+ expected_worker_completed = self._expected_worker_completed_dcm(
+ mock_control_output_channel
+ )
+ collected = self._drain_until(
+ output_queue,
+ lambda items: expected_worker_completed in items,
+ )
+
+ if not
main_loop.context.state_manager.confirm_state(WorkerState.COMPLETED):
+ pytest.fail(
+ "worker with an empty on_finish did not reach COMPLETED within
"
+ "the deadline -- the single-hand-off completion signal was
lost. "
+ f"state={main_loop.context.state_manager.get_current_state()},
"
+ f"collected={collected}"
+ )
+
+ expected_input_port_completed = self._expected_port_completed_dcm(
+ mock_control_output_channel, 0, mock_link.to_port_id, True
+ )
+ expected_output_port_completed = self._expected_port_completed_dcm(
+ mock_control_output_channel, 1, PortIdentity(id=0), False
+ )
+ assert expected_input_port_completed in collected
+ assert expected_output_port_completed in collected
+ assert expected_worker_completed in collected
+
+ # The EndChannel ECM is forwarded downstream on the data output
channel.
+ assert (
+ self._forwarded_ecm(
+ mock_data_output_channel,
+ "EndChannel",
+ EmbeddedControlMessageType.PORT_ALIGNMENT,
+ )
+ in collected
+ )
+
+ reraise()
+
@pytest.mark.timeout(2)
def test_console_message_rpc_fires_before_exception_pause(
self, main_loop, monkeypatch