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 e750074eb9 test(pyamber): cover the main loop's control, console and 
end-channel paths (#8042)
e750074eb9 is described below

commit e750074eb92f57a6a945b2a09d07b8e742c0b82f
Author: Xinyuan Lin <[email protected]>
AuthorDate: Fri Aug 28 02:56:20 2026 +0000

    test(pyamber): cover the main loop's control, console and end-channel paths 
(#8042)
    
    ### What changes were proposed in this PR?
    
    `test_main_loop.py` goes from 34 tests to 44, covering the
    console-message RPC path, the end-channel completion rules, the
    per-element exception backstop, and the ECM/loop-counter forwarding
    paths.
    
    | Metric | Before | After |
    |---|---|---|
    | Codecov fully-covered lines **today** | 270/285 | **284/285** |
    | Codecov fully-covered lines **after #8040** | 259/285 | **280/285** |
    | Branch arms newly taken | — | **+11** |
    
    **+14 fully-covered lines move on Codecov today; the full +21 only shows
    up once #8040 lands.** Lines 176, 397, 658, 661, 706, 714 and 837 are
    line-hit with one arm untaken, so today's branch-less report already
    scores them covered — closing their second arm is real and
    mutation-proven but invisible until branch coverage is on. I would
    rather state both figures than quote the larger one and have it not
    materialise.
    
    The `[tool.coverage.run] branch = true` setting is deliberately **not**
    added here; it belongs to #8040 and adding it would conflict.
    
    Newly fully covered: 176, 177, 397, 400, 658, 661, 662, 706, 714, 715,
    718, 719, 725, 813, 814, 817, 837, 838, 839, 849, 850. Nothing lost.
    Line 795 remains uncovered and four lines (349, 650, 740, 794) remain
    branch-partial.
    
    ### What the reviewers found
    
    Two adversarial reviewers reported eight surviving mutants against the
    first draft. Every one was reproduced before being fixed — seven gave a
    clean `43 passed` on the delivered spec and one produced no summary at
    all. **Eight of the claimed lines were covered vacuously.** The ones
    worth naming:
    
    - **The console-message test blocked on `output_queue.get()` instead of
    failing** — a hang, not an assertion failure, which is the worst way for
    a test to be wrong.
    - **The two-input-port end-channel test could not distinguish the port
    lookup from a constant**, so the rule it exists to pin was
    unconstrained.
    - **The no-output-ports test could not see an `EndChannel` broadcast at
    all.**
    - **`assert switched == [True]` claimed a statement *order* it cannot
    observe.**
    - A `emitted_counter == 0` assertion was a degenerate literal, and a
    `reset_calls == []` assertion was call-graph-guaranteed to hold.
    
    All eight now die. **33 mutations applied one at a time, 32 killed, 1
    judged equivalent** (a flush/pause order exchange, with the reasoning
    recorded).
    
    Four of the 32 kills are honest about their conditions: the line-706,
    line-661, line-837 negations and the 393–394 order swap die only when
    the run is scoped to a single test, because they make **pre-existing**
    thread-driven tests block, and pytest-timeout's Windows `thread` method
    then kills the session before the killing test runs. They are stated as
    scoped kills rather than credited as clean whole-file kills.
    
    One further mutant was discarded rather than reported as surviving:
    negating line 661 makes the pre-existing
    `test_main_loop_thread_can_process_messages` block forever, yielding no
    per-test signal. The same statement is covered by a `return`→`pass`
    mutant that dies exclusively against a new test, so nothing is left
    unproven.
    
    ### A defect found and deliberately not pinned
    
    `_process_data_element`'s per-element backstop (813–814) is `except
    Exception as err: logger.exception(err)` — it logs and continues
    **without calling `context.report_exception`**. A runtime failure on one
    element therefore never reaches the coordinator, and the workflow can
    report SUCCESS on a short result.
    
    The new test asserts only that iteration continues and that nothing
    propagates. It does **not** assert that nothing is reported, because
    that half is arguably a silent-wrong-results bug and pinning it would
    cement it. The test carries a comment saying so; anyone strengthening it
    should keep that boundary.
    
    ### Verification
    
    Measured with the full CI-shaped suite from `amber/`, `--cov-branch`
    passed explicitly on both sides (this branch is on `main` and so lacks
    #8040's config), `LOGURU_LEVEL=WARNING`, per-line `hits` and
    `condition-coverage` read out of `coverage.xml` by script rather than
    from the summary percentage. The before-state was obtained by copying
    the original spec back from a scratch snapshot — never via `git restore`
    — with the production tree verified pristine throughout.
    
    **No regression.** Full pyamber suite: `5 failed, 1169 passed, 7 errors`
    against a baseline of `5 failed, 1159 passed, 7 errors`. The
    `FAILED`/`ERROR` node-id sets were extracted from both runs, sorted and
    diffed: **identical by identity**, 12 entries (7 `test_iceberg_document`
    errors, 2 iceberg failures, `test_tuple::test_hash`, 2
    `test_expression_evaluator` repr-formatting failures). +10 passing. The
    scoped spec passes `44 passed` on five consecutive runs.
    
    `ruff check` and `ruff format --check` pass on CI's exact scope
    (`src/main/python src/test/python`). `git diff --numstat` shows `666 0`
    — 666 added, **zero deleted**.
    
    **A flake seen once and reported rather than buried:** one no-branch
    after-run came back 282/285 with lines 324–325 at zero hits; three
    consecutive re-runs all gave 284/285 with both covered. Those two lines
    are in `_check_and_process_control`'s while body, reached from a daemon
    thread, so a scheduling-dependent miss is plausible. The published
    figure is the reproducible one.
    
    ### On overlap with open PRs
    
    This file is not untouched work. **#7624** (open) and **#7388** (draft)
    both modify `main_loop.py` *and* `test_main_loop.py`. I checked both
    diffs: none of the 21 lines here falls inside either PR's edited hunks,
    and the new tests are appended at the end of `TestMainLoop`, away from
    the test-file hunks. The collision risk is a routine rebase rather than
    a semantic conflict — but it is not zero, and the earlier claim that
    this target was free was wrong.
    
    ### Deliberately not included
    
    Line 795 and the four remaining partials (349, 650, 740, 794) are left.
    Reaching 813/814 already required monkeypatching a `MainLoop` private —
    established practice in this file but white-box, and flagged as the one
    such test in the bundle.
    
    No production file is touched; `git diff -- 'amber/src/main/*'` is
    empty, checked after every mutation revert and as the final action.
    
    ### Any related issues, documentation, discussions?
    
    Closes #8041
    
    ### How was this PR tested?
    
    ```
    cd amber && python -m pytest -m "not integration" -q 
src/test/python/core/runnables/test_main_loop.py
    ```
    
    ```
    44 passed, 1 warning
    ```
    
    ```
    cd amber && ruff check src/main/python src/test/python && ruff format 
--check src/main/python src/test/python
    ```
    
    `amber/junit.xml`, `amber/coverage.xml` and `amber/.coverage` are
    regenerated by every run and are not committed.
    
    ### Was this PR authored or co-authored using generative AI tooling?
    
    Generated-by: Claude Code (Opus 5)
    
    ---------
    
    Signed-off-by: Xinyuan Lin <[email protected]>
    Co-authored-by: Copilot Autofix powered by AI 
<[email protected]>
---
 .../test/python/core/runnables/test_main_loop.py   | 667 +++++++++++++++++++++
 1 file changed, 667 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 5aa259c428..c5a566ada9 100644
--- a/amber/src/test/python/core/runnables/test_main_loop.py
+++ b/amber/src/test/python/core/runnables/test_main_loop.py
@@ -3157,3 +3157,670 @@ class TestMainLoop:
         finally:
             first.executor_manager.close()
             second.executor_manager.close()
+
+    # ------------------------------------------------------------------ #
+    # Deferred loop consume: the "nothing was stashed" shape
+    # ------------------------------------------------------------------ #
+
+    @pytest.mark.timeout(5)
+    def test_deferred_consume_is_a_noop_when_no_state_was_stashed(
+        self, main_loop, monkeypatch
+    ):
+        # A Loop End may legally reach EndChannel having never taken a matching
+        # loop state -- LoopEndOperator.eval_condition's `_loop_table is None`
+        # guard makes that a supported shape (the loop simply does not 
iterate),
+        # and _check_loop_state_arrived deliberately does not treat it as an
+        # error. The deferred consume must therefore return without doing any
+        # loop work: reading the Loop Start's input-port materialization for a
+        # loop that never ran would touch storage for no reason, and clearing
+        # the fan-in dedup flag would rewrite bookkeeping owned by a consume
+        # that did not happen.
+        executor = _FalseLoopEnd()
+        main_loop.context.executor_manager.executor = executor
+        assert main_loop._pending_loop_state is None
+
+        reads = []
+        loop_table = Table([Tuple({"v": 1})])
+
+        def _read():
+            reads.append(True)
+            return loop_table
+
+        monkeypatch.setattr(main_loop, "_read_loop_input_table", _read)
+        consumed = []
+        monkeypatch.setattr(
+            executor, "process_state", lambda st, port: consumed.append((st, 
port))
+        )
+        # Arm the dedup flag so the early return is distinguishable from the
+        # consume path, whose whole point is to clear it.
+        main_loop._loop_state_consumed = True
+
+        main_loop._consume_pending_loop_state(executor)
+
+        assert reads == [], "no stashed state means no storage read"
+        assert consumed == [], "the operator's update must not run"
+        assert executor._attached_table is None, "no table may be attached"
+        assert main_loop._loop_state_consumed is True, (
+            "the early return must not rewrite the fan-in dedup flag"
+        )
+
+    @pytest.mark.timeout(5)
+    def test_loopstart_stamps_its_own_logical_op_id_on_its_output_state(
+        self, main_loop, monkeypatch
+    ):
+        # The stamp is what lets the matching Loop End find the loop to jump
+        # back to (it rides the StateFrame envelope and is captured in
+        # _process_state_frame). A LoopStart must therefore REPLACE whatever id
+        # arrived with its own logical op id rather than forward the inbound
+        # one -- forwarding "" is exactly the lost-envelope bug class
+        # #6660/#6661 fixed, and _check_loop_state_arrived only notices that
+        # once a whole input port has drained.
+        class StubLoopStart(LoopStartOperator):
+            def process_table(self, table, port):
+                yield
+
+        main_loop.context.executor_manager.executor = StubLoopStart()
+        # get_logical_op_id parses the canonical worker actor name
+        # "Worker:WF<wf>-<opId>-<layer>-<idx>" and raises on anything else, so
+        # the fixture's "dummy_worker_id" has to be replaced here.
+        main_loop.context.worker_id = "Worker:WF7-my-loop-start-main-0"
+        emitted, _, _ = self._capture_state_emit(main_loop, monkeypatch)
+        # Both stubs append to ONE list so the ORDER is asserted rather than
+        # merely that each happened. _switch_context is what hands control to
+        # the DataProcessor thread that PRODUCES the state, so reading the
+        # output state first would emit the previous iteration's state (or
+        # None). Recording the switch in its own separate list -- what the
+        # shared _capture_state_emit helper does -- can only witness THAT the
+        # switch occurred, never that it came first.
+        order = []
+        monkeypatch.setattr(
+            main_loop, "_switch_context", lambda: order.append("switch")
+        )
+        monkeypatch.setattr(
+            main_loop.context.state_processing_manager,
+            "get_output_state",
+            lambda: (order.append("read"), State({"i": 3}))[1],
+        )
+
+        # The inbound envelope carried no id -- the shape a first-entry state
+        # and the back-edge write both have.
+        main_loop.process_input_state(output_loop_counter=0, 
output_loop_start_id="")
+
+        assert order == ["switch", "read"], (
+            f"the operator must run before its output is read; got {order}"
+        )
+        assert len(emitted) == 1
+        emitted_state, emitted_counter, emitted_id = emitted[0]
+        assert emitted_state == State({"i": 3})
+        # 0 is simultaneously the inbound value, process_input_state's own
+        # parameter default and the expectation, so this pins the emitted
+        # tuple's SHAPE only -- it cannot tell a pass-through from a hardcoded
+        # 0. The pass-through itself is pinned by
+        # test_body_operator_state_forwards_the_inbound_loop_counter_and_id,
+        # which passes values no default can supply.
+        assert emitted_counter == 0
+        assert emitted_id == "my-loop-start", (
+            "a LoopStart must stamp its own logical op id, not forward the "
+            f"inbound one; emitted id: {emitted_id!r}"
+        )
+        # `reset_calls` is deliberately NOT asserted here: reset_output_storage
+        # has exactly one call site, in _process_state_frame 
(main_loop.py:512),
+        # so nothing process_input_state reaches could fire it and the
+        # assertion could never fail. That reset is pinned where it happens, by
+        # test_loopend_passthrough_decrements_resets_output_and_skips_operator.
+
+    @pytest.mark.timeout(5)
+    def test_body_operator_state_forwards_the_inbound_loop_counter_and_id(
+        self, main_loop, monkeypatch
+    ):
+        # A plain loop-BODY operator is the shape that carries a NON-zero
+        # counter: _process_state_frame's default tail calls
+        # process_input_state(output_loop_counter=in_counter,
+        # output_loop_start_id=frame.loop_start_id), so both envelope fields
+        # have to reach _emit_and_save_state unchanged. Blanking the counter
+        # would strip the iteration number off every body operator's boundary
+        # state -- the lost-envelope class #6660/#6661 fixed -- and the values
+        # used here (7, "outer-loop") are ones no parameter default can supply,
+        # so this discriminates a real forward from a constant.
+        assert not isinstance(
+            main_loop.context.executor_manager.executor, LoopStartOperator
+        ), "a body operator is not a LoopStart, so no id stamping happens here"
+        emitted, _, _ = self._capture_state_emit(main_loop, monkeypatch)
+        monkeypatch.setattr(
+            main_loop.context.state_processing_manager,
+            "get_output_state",
+            lambda: State({"i": 3}),
+        )
+
+        main_loop.process_input_state(
+            output_loop_counter=7, output_loop_start_id="outer-loop"
+        )
+
+        assert emitted == [(State({"i": 3}), 7, "outer-loop")], (
+            f"both envelope fields must be forwarded verbatim; emitted: 
{emitted}"
+        )
+
+    # ------------------------------------------------------------------ #
+    # _process_end_channel: the two guards that hold a worker open
+    # ------------------------------------------------------------------ #
+
+    def _register_input_port(self, main_loop, schema, port_id, sender):
+        """Register one input port carrying a single data channel and return
+        that channel. Uses the public InputManager API rather than a DCM
+        round-trip so a test can build a multi-port worker directly."""
+        channel_id = ChannelIdentity(
+            ActorVirtualIdentity(sender),
+            ActorVirtualIdentity("dummy_worker_id"),
+            False,
+        )
+        main_loop.context.input_manager.add_input_port(port_id, schema, [], [])
+        main_loop.context.input_manager.register_input(channel_id, port_id)
+        return channel_id
+
+    def _capture_end_channel_effects(self, main_loop, monkeypatch):
+        """Stub out everything _process_end_channel does apart from the
+        completion decision itself, and return
+        (port_completed_calls, closed, completed)."""
+        port_completed_calls = []
+        closed = []
+        completed = []
+        monkeypatch.setattr(main_loop, "process_input_state", lambda *a, **k: 
None)
+        monkeypatch.setattr(main_loop, "process_input_tuple", lambda: None)
+        monkeypatch.setattr(main_loop, "complete", lambda: 
completed.append(True))
+        monkeypatch.setattr(
+            main_loop.context.output_manager,
+            "close_port_storage_writers",
+            lambda: closed.append(True),
+        )
+
+        class _Coordinator:
+            def port_completed(self, request):
+                port_completed_calls.append(request)
+
+        monkeypatch.setattr(
+            main_loop._async_rpc_client, "coordinator_stub", lambda: 
_Coordinator()
+        )
+        return port_completed_calls, closed, completed
+
+    @pytest.mark.timeout(5)
+    def test_end_channel_does_not_complete_while_a_second_input_port_is_open(
+        self, main_loop, monkeypatch, mock_raw_schema
+    ):
+        # A worker with several input ports gets one EndChannel per port. The
+        # first of them must report ITS port complete and stop there: closing
+        # the storage writers or calling complete() while another port is still
+        # streaming would truncate that port's results and let the coordinator
+        # mark the region done early (region completion is port-based).
+        schema = Schema(raw_schema=mock_raw_schema)
+        port_0 = PortIdentity(0, internal=False)
+        port_1 = PortIdentity(1, internal=False)
+        self._register_input_port(main_loop, schema, port_0, "sender-0")
+        channel_1 = self._register_input_port(main_loop, schema, port_1, 
"sender-1")
+        # The port that FINISHES is deliberately the NON-zero one. With port 0
+        # finishing, "report the arriving channel's port" and "always report
+        # port 0" are the same program: port 0 would be the finished port, the
+        # lowest port id and the expected value at once. A worker hardcoded to
+        # port 0 lets the coordinator close a port that is still streaming,
+        # since region completion is port-based.
+        main_loop.context.input_manager.complete_current_port(channel_1)
+        main_loop.context.current_input_channel_id = channel_1
+        assert not main_loop.context.input_manager.all_ports_completed(), (
+            "port 0 must still be open for this test to mean anything"
+        )
+        # The output port is what makes this test non-vacuous: with none, the
+        # is_missing_output_ports() guard below returns early regardless of 
what
+        # all_ports_completed() answered, and a broken port-completion check
+        # would go unnoticed.
+        main_loop.context.output_manager.add_output_port(port_0, schema)
+        assert not main_loop.context.output_manager.is_missing_output_ports()
+
+        port_completed_calls, closed, completed = 
self._capture_end_channel_effects(
+            main_loop, monkeypatch
+        )
+
+        main_loop._process_end_channel()
+
+        assert port_completed_calls == [
+            PortCompletedRequest(port_id=port_1, input=True)
+        ], (
+            "only the FINISHED input port may be reported complete; got "
+            f"{port_completed_calls}"
+        )
+        assert closed == [], "an open input port must keep the storage writers 
open"
+        assert completed == [], "the worker must not complete with a port 
still open"
+
+    @staticmethod
+    def _queued_size(output_queue, channel):
+        """How many elements are queued for `channel`.
+
+        Sub-queues are created lazily on first put, so an untouched channel is
+        simply absent -- reading it as 0 rather than raising lets a test assert
+        "nothing was sent here" without depending on whether the sub-queue was
+        ever created.
+        """
+        try:
+            return output_queue._queue.size(channel)
+        except KeyError:
+            return 0
+
+    @pytest.mark.timeout(5)
+    def test_end_channel_holds_the_worker_open_when_it_has_no_output_ports(
+        self, main_loop, monkeypatch, output_queue, mock_raw_schema
+    ):
+        # The two-phase dependee-port region shape (see
+        # OutputManager.is_missing_output_ports): this worker's only input port
+        # has finished, but it has no output port at all, which means it is
+        # running the dependee-port phase and must stay open for the
+        # non-dependee-port phase that follows. So it reports its input port
+        # complete and then stops -- no storage close, no EndChannel ECM, no
+        # complete().
+        schema = Schema(raw_schema=mock_raw_schema)
+        port_0 = PortIdentity(0, internal=False)
+        channel_0 = self._register_input_port(main_loop, schema, port_0, 
"sender-0")
+        main_loop.context.input_manager.complete_current_port(channel_0)
+        main_loop.context.current_input_channel_id = channel_0
+        assert main_loop.context.input_manager.all_ports_completed()
+        assert main_loop.context.output_manager.is_missing_output_ports()
+        # A downstream data CHANNEL with no output PORT. OutputManager keeps
+        # _ports and _channels in independent dicts (output_manager.py:85-86)
+        # and add_partitioning only writes _channels, so the hold-open guard
+        # still fires while get_output_channel_ids() has somewhere to send.
+        # Without this channel _send_ecm_to_data_channels is a no-op whatever
+        # the guard does, and a premature EndChannel broadcast injected into
+        # the guard -- which would close downstream ports before the
+        # non-dependee-port phase runs, the exact failure the two-phase scheme
+        # exists to prevent -- would go unnoticed.
+        downstream = ChannelIdentity(
+            ActorVirtualIdentity("dummy_worker_id"),
+            ActorVirtualIdentity("downstream"),
+            False,
+        )
+        main_loop.context.output_manager.add_partitioning(
+            PhysicalLink(
+                from_op_id=PhysicalOpIdentity(OperatorIdentity("from"), 
"from"),
+                from_port_id=PortIdentity(0, internal=False),
+                to_op_id=PhysicalOpIdentity(OperatorIdentity("to"), "to"),
+                to_port_id=PortIdentity(0, internal=False),
+            ),
+            set_one_of(
+                Partitioning,
+                OneToOnePartitioning(batch_size=1, channels=[downstream]),
+            ),
+        )
+        assert main_loop.context.output_manager.is_missing_output_ports(), (
+            "registering a channel must not create an output port"
+        )
+        assert list(main_loop.context.output_manager.get_output_channel_ids()) 
== [
+            downstream
+        ], "the broadcast must have a live channel to reach"
+
+        port_completed_calls, closed, completed = 
self._capture_end_channel_effects(
+            main_loop, monkeypatch
+        )
+
+        main_loop._process_end_channel()
+
+        assert port_completed_calls == [
+            PortCompletedRequest(port_id=port_0, input=True)
+        ], "the input port is still reported complete"
+        assert closed == [], (
+            "the dependee-port phase must not close the storage writers"
+        )
+        assert self._queued_size(output_queue, downstream) == 0, (
+            "the dependee-port phase must not send EndChannel downstream"
+        )
+        assert completed == [], "the worker must stay open for the next phase"
+
+    # ------------------------------------------------------------------ #
+    # _process_ecm: per-worker command dispatch and scoped forwarding
+    # ------------------------------------------------------------------ #
+
+    @staticmethod
+    def _no_op_invocation():
+        return ControlInvocation(
+            "NoOperation",
+            ControlRequest(empty_request=EmptyRequest()),
+            AsyncRpcContext(ActorVirtualIdentity(), ActorVirtualIdentity()),
+            -1,
+        )
+
+    @pytest.mark.timeout(5)
+    def test_ecm_dispatches_only_the_command_addressed_to_this_worker(
+        self, main_loop, monkeypatch, mock_data_input_channel
+    ):
+        # An ECM's command_mapping is keyed by worker id: a message travelling
+        # a scope carries commands only for the workers that must act on it,
+        # and every other worker on the path still has to align and propagate
+        # it. Handing the missing entry (None) to the RPC server would dispatch
+        # a control invocation with no method. Both directions are asserted
+        # here so the guard is pinned as a discriminator rather than as a path
+        # that merely happens to be taken.
+        received = []
+        monkeypatch.setattr(
+            main_loop._async_rpc_server,
+            "receive",
+            lambda channel_id, command: received.append((channel_id, command)),
+        )
+        main_loop.context.current_input_channel_id = mock_data_input_channel
+
+        def deliver(ecm_id, command_mapping):
+            main_loop._process_ecm(
+                ECMElement(
+                    tag=mock_data_input_channel,
+                    payload=EmbeddedControlMessage(
+                        EmbeddedControlMessageIdentity(ecm_id),
+                        EmbeddedControlMessageType.NO_ALIGNMENT,
+                        [],
+                        command_mapping,
+                    ),
+                )
+            )
+
+        deliver(
+            "ecm-for-somebody-else", {"some-other-worker": 
self._no_op_invocation()}
+        )
+        assert received == [], (
+            "an ECM carrying no command for this worker must dispatch nothing; 
"
+            f"dispatched: {received}"
+        )
+
+        mine = self._no_op_invocation()
+        deliver("ecm-for-me", {"dummy_worker_id": mine})
+        assert received == [(mock_data_input_channel, mine)], (
+            "an ECM carrying a command for this worker must dispatch it; "
+            f"dispatched: {received}"
+        )
+
+    @pytest.mark.timeout(5)
+    def test_ecm_is_forwarded_only_to_the_output_channels_in_its_scope(
+        self, main_loop, output_queue, mock_data_input_channel
+    ):
+        # An ECM's scope is the set of channels it is allowed to travel. A
+        # worker with several downstream channels must forward the message only
+        # along the ones the scope names -- sending it down a channel outside
+        # the scope would inject an alignment barrier into a region the message
+        # was never meant to reach. Two output channels are required for this
+        # to mean anything: with one, "forward to the in-scope channel" and
+        # "forward to every channel" are the same program.
+        in_scope = ChannelIdentity(
+            ActorVirtualIdentity("dummy_worker_id"),
+            ActorVirtualIdentity("downstream-in-scope"),
+            False,
+        )
+        out_of_scope = ChannelIdentity(
+            ActorVirtualIdentity("dummy_worker_id"),
+            ActorVirtualIdentity("downstream-out-of-scope"),
+            False,
+        )
+        link = PhysicalLink(
+            from_op_id=PhysicalOpIdentity(OperatorIdentity("from"), "from"),
+            from_port_id=PortIdentity(0, internal=False),
+            to_op_id=PhysicalOpIdentity(OperatorIdentity("to"), "to"),
+            to_port_id=PortIdentity(0, internal=False),
+        )
+        main_loop.context.output_manager.add_partitioning(
+            link,
+            set_one_of(
+                Partitioning,
+                OneToOnePartitioning(batch_size=1, channels=[in_scope, 
out_of_scope]),
+            ),
+        )
+        assert set(main_loop.context.output_manager.get_output_channel_ids()) 
== {
+            in_scope,
+            out_of_scope,
+        }
+
+        main_loop.context.current_input_channel_id = mock_data_input_channel
+        # The scope is expressed with a freshly built (equal, not identical)
+        # ChannelIdentity, the way a real scope arrives off the wire.
+        scoped_ecm = EmbeddedControlMessage(
+            EmbeddedControlMessageIdentity("scoped-ecm"),
+            EmbeddedControlMessageType.NO_ALIGNMENT,
+            [
+                ChannelIdentity(
+                    ActorVirtualIdentity("dummy_worker_id"),
+                    ActorVirtualIdentity("downstream-in-scope"),
+                    False,
+                )
+            ],
+            {},
+        )
+
+        main_loop._process_ecm(
+            ECMElement(tag=mock_data_input_channel, payload=scoped_ecm)
+        )
+
+        assert self._queued_size(output_queue, out_of_scope) == 0, (
+            "a channel outside the ECM's scope must receive nothing"
+        )
+        assert self._queued_size(output_queue, in_scope) == 1, (
+            "the channel named by the scope must receive the ECM"
+        )
+        element = output_queue.get()
+        assert isinstance(element, ECMElement)
+        assert element.tag == in_scope
+        assert element.payload.id == 
EmbeddedControlMessageIdentity("scoped-ecm")
+
+    # ------------------------------------------------------------------ #
+    # Console / debugger reporting
+    # ------------------------------------------------------------------ #
+
+    @pytest.mark.timeout(5)
+    def test_console_message_is_sent_to_the_coordinator_as_an_rpc(
+        self, main_loop, output_queue
+    ):
+        # Every console report -- user prints, operator errors, debugger events
+        # -- funnels through _send_console_message, and the surrounding tests
+        # all stub that method off the instance, so nothing pins the RPC it
+        # actually makes. Let the real method run and read the resulting
+        # control message off the output queue.
+        msg = ConsoleMessage(
+            worker_id="dummy_worker_id",
+            timestamp=current_time_in_local_timezone(),
+            msg_type=ConsoleMessageType.PRINT,
+            source="pytest",
+            title="hello from the operator",
+            message="",
+        )
+        main_loop.context.console_message_manager.put_message(msg)
+
+        main_loop._check_and_report_console_messages(force_flush=True)
+
+        coordinator_channel = ChannelIdentity(
+            ActorVirtualIdentity("dummy_worker_id"),
+            ActorVirtualIdentity("COORDINATOR"),
+            True,
+        )
+        # Check the queue is non-empty BEFORE reading it: output_queue.get()
+        # blocks forever, so a regression that drops the RPC entirely would
+        # hang here instead of failing (and on Windows pytest-timeout's default
+        # `thread` method then kills the whole session, yielding no per-test
+        # signal at all).
+        queued = {
+            str(key): output_queue._queue.size(key)
+            for key in output_queue._queue.sub_queues
+        }
+        assert self._queued_size(output_queue, coordinator_channel) == 1, (
+            f"the console message must be queued for the coordinator; queued: 
{queued}"
+        )
+
+        element = output_queue.get()
+        assert isinstance(element, DCMElement)
+        assert element.tag == coordinator_channel, (
+            "console messages go to the coordinator on the control channel"
+        )
+        invocation = element.payload.control_invocation
+        assert invocation.method_name == "ConsoleMessageTriggered"
+        assert (
+            
invocation.command.console_message_triggered_request.console_message == msg
+        )
+
+    @pytest.mark.timeout(5)
+    def 
test_debug_event_is_reported_as_a_debugger_message_and_pauses_the_worker(
+        self, main_loop, monkeypatch
+    ):
+        # pdb writes its output into a SingleBlockingIO that the worker polls
+        # after every context switch. An event there has to reach the frontend
+        # as a DEBUGGER console message AND pause the worker -- without the
+        # pause the debugger would report a breakpoint and then run straight
+        # past it. The buffered prints have to go out on the same beat: the
+        # worker is about to stop, so anything still sitting in the console
+        # buffer when the pause lands would stay invisible until a resume.
+        console_msgs = []
+        pauses = []
+        monkeypatch.setattr(
+            main_loop, "_send_console_message", lambda msg: 
console_msgs.append(msg)
+        )
+        # The stub RECORDS change_state instead of dropping it: that kwarg is
+        # what decides whether the worker actually reports itself PAUSED
+        # (pause_manager.py:59-62 gates transit_to(WorkerState.PAUSED) on it).
+        # With change_state=False the input queue is disabled but the state
+        # never transits, i.e. the frontend shows RUNNING while the operator
+        # sits at a breakpoint. The stub's own default is True, so an omitted
+        # kwarg records True and an explicit False records False -- which is
+        # what makes the pair discriminating rather than decorative.
+        monkeypatch.setattr(
+            main_loop.context.pause_manager,
+            "pause",
+            lambda pause_type, change_state=True: pauses.append(
+                (pause_type, change_state)
+            ),
+        )
+        buffered_print = ConsoleMessage(
+            worker_id="dummy_worker_id",
+            timestamp=current_time_in_local_timezone(),
+            msg_type=ConsoleMessageType.PRINT,
+            source="pytest",
+            title="printed just before the breakpoint",
+            message="",
+        )
+        main_loop.context.console_message_manager.put_message(buffered_print)
+        # flush() is what makes the buffered text readable; without it
+        # has_debug_event() stays False and readline() would block.
+        debug_out = main_loop.context.debug_manager.debugger.stdout
+        debug_out.write("> <string>(3)update()")
+        debug_out.flush()
+        assert main_loop.context.debug_manager.has_debug_event()
+
+        before = current_time_in_local_timezone()
+        main_loop._check_and_report_debug_event()
+        after = current_time_in_local_timezone()
+
+        assert len(console_msgs) == 2, f"sent: {console_msgs}"
+        assert console_msgs[1] == buffered_print, (
+            "the console buffer must be flushed before the worker pauses; "
+            f"sent: {console_msgs}"
+        )
+        reported = console_msgs[0]
+        assert reported.msg_type == ConsoleMessageType.DEBUGGER
+        assert reported.source == "(Pdb)"
+        assert "> <string>(3)update()" in reported.title, (
+            f"the pdb event is the message title; got title={reported.title!r} 
"
+            f"message={reported.message!r}"
+        )
+        assert reported.message == ""
+        assert reported.worker_id == "dummy_worker_id"
+        # `timestamp` is a plain betterproto datetime field, so 
`timestamp=None`
+        # constructs cleanly and a DEBUGGER message would reach the frontend
+        # with no time on it. Bracketing it against the same helper the runtime
+        # uses pins that a real clock reading was taken.
+        assert reported.timestamp is not None, "the console message needs a 
time"
+        assert before <= reported.timestamp <= after, (
+            f"timestamp {reported.timestamp!r} outside [{before!r}, {after!r}]"
+        )
+        assert pauses == [(PauseType.DEBUG_PAUSE, True)], (
+            "the debug pause must also transit the worker to PAUSED "
+            f"(change_state); recorded: {pauses}"
+        )
+        # The event was consumed, so a second poll reports nothing.
+        assert not main_loop.context.debug_manager.has_debug_event()
+        main_loop._check_and_report_debug_event()
+        assert len(console_msgs) == 2
+
+    @pytest.mark.timeout(5)
+    def test_a_failing_element_does_not_abandon_the_rest_of_the_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.
+        schema = Schema(raw_schema=mock_raw_schema)
+        port_0 = PortIdentity(0, internal=False)
+        channel = self._register_input_port(main_loop, schema, port_0, 
"sender")
+        main_loop.context.current_input_channel_id = channel
+
+        element = DataElement(
+            tag=channel,
+            payload=DataFrame(
+                frame=pyarrow.Table.from_pandas(
+                    pandas.DataFrame(
+                        [
+                            {"test-1": "a", "test-2": 0},
+                            {"test-1": "b", "test-2": 1},
+                        ]
+                    )
+                )
+            ),
+        )
+
+        attempted = []
+
+        def _boom(port_id, size):
+            attempted.append(
+                
main_loop.context.tuple_processing_manager.current_input_tuple["test-2"]
+            )
+            raise RuntimeError("statistics backend unavailable")
+
+        monkeypatch.setattr(
+            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}"

Reply via email to