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-7061-3579fc97140b985f4b6b2c02420ee60e981497c5 in repository https://gitbox.apache.org/repos/asf/texera.git
commit f28d8ec35c6fc5828a46157ffcec8c696e0060ae Author: Eugene Gu <[email protected]> AuthorDate: Wed Jul 29 19:02:21 2026 -0700 test(pyamber): extend the StartWorkerHandler unit tests (#7061) ### What changes were proposed in this PR? Rewrites `amber/src/test/python/core/architecture/handlers/control/test_start_worker_handler.py`, added by #6988 while this PR was open, to cover all three branches of `StartWorkerHandler` and the shared tail. Four behaviors the previous file did not assert: | Behavior | Why it matters | | --- | --- | | Branch precedence | `is_source` wins even when materialization reader threads exist, so the `elif` predicate must go unqueried | | Empty marker scope | A self-addressed marker propagates nowhere | | `EmptyRequest` unwrapped via `get_one_of` | `AsyncRpcServer` unwraps the payload before dispatching; an unset oneof delivers `None`, and betterproto's `__eq__` cannot see the difference | | Markers land on the data sub-queue | `ECMElement`s are data traffic, so a source worker's own start/end markers are gated by `disable_data` — pause and backpressure — unlike control messages | It also pins the ordering of `add_input_port` before `register_input` (the reverse raises `KeyError` on a real `InputManager`, since `register_input` ends with `self._ports[port_id].add_channel(...)`). Six tests run against real collaborators — a `StateManager` built on the production transition graph and a real `InputManager` sharing the context's `InternalQueue` — rather than mocks, because they pin behavior mocks cannot express: a rejected transition from `UNINITIALIZED` leaves the queue, the state and the `InputManager` untouched; the two markers are routed to the data sub-queue; and the response version tracks the real state clock on both the source and the non-source paths. The two non-source cases (`READY`/0 with reader threads, `PAUSED`/1 echoed untouched) carry over from #6988; everything else that file asserted is subsumed by a stricter assertion here. `test_double_start_enqueues_a_second_pair_of_markers` deliberately pins observed current behavior rather than correctness. `StateManager.transit_to` returns early when the target equals the current state, so a second `StartWorker` is a silent no-op on the state and appends another marker pair — including a second `PORT_ALIGNMENT` `EndChannel` on the same port. The test comment says so and asks for a rewrite if an idempotency guard is ever added. Test-only change; no production code is touched. ### Any related issues, documentation, discussions? Closes #7058. The real-collaborator tests double as a Python-side guard against the stale-state-snapshot family of bugs described in #6010 (`WorkerExecution.scala:56`); a mocked `get_state_with_version` returns a constant and cannot catch a hoisted snapshot. ### How was this PR tested? New unit tests only; no production behavior changes. ``` cd amber pytest src/test/python/core/architecture/handlers/ -q ``` - 23 cases in this file; the handlers package goes from 46 to 69 passing tests. - `black --check` and `ruff check` clean. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Opus 5) --- .../handlers/control/test_start_worker_handler.py | 513 +++++++++++++++------ 1 file changed, 373 insertions(+), 140 deletions(-) diff --git a/amber/src/test/python/core/architecture/handlers/control/test_start_worker_handler.py b/amber/src/test/python/core/architecture/handlers/control/test_start_worker_handler.py index 4d1afb30ef..24bba621bb 100644 --- a/amber/src/test/python/core/architecture/handlers/control/test_start_worker_handler.py +++ b/amber/src/test/python/core/architecture/handlers/control/test_start_worker_handler.py @@ -17,6 +17,7 @@ import asyncio from types import SimpleNamespace +from typing import Optional from unittest.mock import MagicMock import pytest @@ -28,11 +29,12 @@ from core.architecture.managers.state_manager import ( StateManager, ) from core.architecture.packaging.input_manager import InputManager -from core.models.internal_queue import ECMElement, InternalQueue +from core.models import Schema +from core.models.internal_queue import InternalQueue +from core.util.proto import get_one_of from proto.org.apache.texera.amber.core import ( ActorVirtualIdentity, ChannelIdentity, - EmbeddedControlMessageIdentity, PortIdentity, ) from proto.org.apache.texera.amber.engine.architecture.rpc import ( @@ -44,175 +46,406 @@ from proto.org.apache.texera.amber.engine.architecture.worker import WorkerState WORKER_ID = "worker-1" -SOURCE_CHANNEL = ChannelIdentity( +# The synthetic channel a source worker starts itself over: it has no upstream, +# so the handler fabricates a SOURCE_STARTER -> self channel to carry the +# start/end markers. +SOURCE_STARTER_CHANNEL = ChannelIdentity( InputManager.SOURCE_STARTER, ActorVirtualIdentity(WORKER_ID), False ) -SOURCE_PORT = PortIdentity(0, False) - - -def make_context( - is_source: bool, - initial_state: WorkerState = WorkerState.READY, - input_manager=None, -): - """Wire the slice of Context that StartWorkerHandler touches. - - The state manager and the input queue are the real ones so state - transitions and the self-fed source markers are asserted against real - behavior; only the executor is stubbed, since the handler reads nothing - from it but `is_source`. - """ - input_queue = InternalQueue() - return SimpleNamespace( + +# The no-reply marker: these invocations are fire-and-forget, so the controller +# must not be left awaiting a promise for them. +NO_REPLY_COMMAND_ID = -1 + +# A non-default version so a regression that drops state_version — or hardcodes +# it to 0 — cannot pass. +STATE_VERSION = 17 + + +def _build_handler( + is_source: bool, mat_reader_threads: Optional[dict] = None +) -> StartWorkerHandler: + """Wire a handler with a SimpleNamespace context exposing everything + start_worker touches, a real InternalQueue so enqueue ordering is observed + through the production queue, and mocked managers.""" + state_manager = MagicMock() + state_manager.get_state_with_version.return_value = ( + WorkerState.RUNNING, + STATE_VERSION, + ) + input_manager = MagicMock() + input_manager.get_input_port_mat_reader_threads.return_value = ( + {} if mat_reader_threads is None else mat_reader_threads + ) + context = SimpleNamespace( + executor_manager=SimpleNamespace(executor=SimpleNamespace(is_source=is_source)), + state_manager=state_manager, + input_manager=input_manager, + input_queue=InternalQueue(), worker_id=WORKER_ID, - input_queue=input_queue, + current_input_channel_id=None, + ) + return StartWorkerHandler(context) + + +def _build_real_handler( + is_source: bool, initial_state: WorkerState, input_manager=None +) -> StartWorkerHandler: + """Wire a handler against the real StateManager and InputManager rather than + mocks, so the collaborators' own invariants take part in the test. The + InputManager and the context share one InternalQueue, exactly as Context + constructs them. The non-source branches pass an InputManager double, since + there is no public way to make a real one report reader threads.""" + queue = InternalQueue() + context = SimpleNamespace( executor_manager=SimpleNamespace(executor=SimpleNamespace(is_source=is_source)), + state_manager=StateManager(WORKER_STATE_TRANSITIONS, initial_state), input_manager=( - input_manager - if input_manager is not None - else InputManager(WORKER_ID, input_queue) + InputManager(WORKER_ID, queue) if input_manager is None else input_manager ), - state_manager=StateManager(WORKER_STATE_TRANSITIONS, initial_state), + input_queue=queue, + worker_id=WORKER_ID, current_input_channel_id=None, ) + return StartWorkerHandler(context) -def start(context) -> WorkerStateResponse: - return asyncio.run(StartWorkerHandler(context).start_worker(EmptyRequest())) - - -def drain(input_queue: InternalQueue): - return [input_queue.get() for _ in range(input_queue.size())] - +def _drain(queue: InternalQueue) -> list: + """Pop everything currently queued, preserving enqueue order.""" + return [queue.get() for _ in range(queue.size())] -class TestStartWorkerOnSource: - def test_transits_to_running_and_reports_the_bumped_version(self): - context = make_context(is_source=True) - response = start(context) - - assert isinstance(response, WorkerStateResponse) - assert response.state == WorkerState.RUNNING - # A source worker really transitions here, so the state manager's - # logical clock must advance and be reported with the new state. - assert response.state_version == 1 - assert context.state_manager.get_current_state() == WorkerState.RUNNING - - def test_registers_the_source_starter_channel_on_port_zero(self): - context = make_context(is_source=True) - - start(context) - - # A source has no upstream, so the handler fabricates port 0 and a - # channel from the synthetic SOURCE_STARTER actor to itself. - assert context.current_input_channel_id == SOURCE_CHANNEL - assert context.input_manager.get_port_id(SOURCE_CHANNEL) == SOURCE_PORT - port = context.input_manager.get_port(SOURCE_PORT) - assert port.get_channels() == {SOURCE_CHANNEL} - # The synthetic port carries no attributes and no materialized readers. - assert port.get_schema().get_attr_names() == [] - assert context.input_manager.get_input_port_mat_reader_threads() == { - SOURCE_PORT: [] - } - - def test_self_feeds_start_then_end_channel_markers(self): - context = make_context(is_source=True) - - start(context) +def _summarize(element) -> tuple: + """Reduce an ECMElement to the fields the start protocol is defined by.""" + (command,) = element.payload.command_mapping.values() + return ( + element.tag, + element.payload.id.id, + element.payload.ecm_type, + set(element.payload.command_mapping), + command.method_name, + command.command_id, + ) - elements = drain(context.input_queue) - assert len(elements) == 2 - assert all(isinstance(element, ECMElement) for element in elements) - assert [element.tag for element in elements] == [SOURCE_CHANNEL] * 2 - start_ecm, end_ecm = (element.payload for element in elements) - # The start marker opens the channel without alignment; the end marker - # must be port-aligned so downstream completion is ordered correctly. - assert start_ecm.id == EmbeddedControlMessageIdentity("StartChannel") - assert start_ecm.ecm_type == EmbeddedControlMessageType.NO_ALIGNMENT - assert end_ecm.id == EmbeddedControlMessageIdentity("EndChannel") - assert end_ecm.ecm_type == EmbeddedControlMessageType.PORT_ALIGNMENT +class TestStartWorkerHandler: + @pytest.fixture + def source_handler(self): + return _build_handler(is_source=True) + + @pytest.fixture + def mat_reader_handler(self): + return _build_handler( + is_source=False, + mat_reader_threads={PortIdentity(0, False): [MagicMock()]}, + ) + + @pytest.fixture + def fall_through_handler(self): + return _build_handler(is_source=False) + + # --- source branch --------------------------------------------------- + + def test_source_transits_to_running(self, source_handler): + asyncio.run(source_handler.start_worker(EmptyRequest())) + source_handler.context.state_manager.transit_to.assert_called_once_with( + WorkerState.RUNNING + ) + + def test_source_adds_the_synthetic_input_port(self, source_handler): + asyncio.run(source_handler.start_worker(EmptyRequest())) + # A source has no declared input port, so port 0 is created empty just + # to give the starter channel somewhere to attach. + source_handler.context.input_manager.add_input_port.assert_called_once_with( + port_id=PortIdentity(0, False), + schema=Schema(), + storage_uris=[], + partitionings=[], + ) + + def test_source_registers_the_starter_channel_on_port_zero(self, source_handler): + asyncio.run(source_handler.start_worker(EmptyRequest())) + source_handler.context.input_manager.register_input.assert_called_once_with( + SOURCE_STARTER_CHANNEL, PortIdentity(0, False) + ) + + def test_source_creates_the_port_before_registering_the_channel( + self, source_handler + ): + asyncio.run(source_handler.start_worker(EmptyRequest())) + # register_input ends with `self._ports[port_id].add_channel(...)`, so + # the port has to exist first; doing these two in the other order + # raises KeyError on a real InputManager. + calls = [name for name, _, _ in source_handler.context.input_manager.mock_calls] + assert calls.index("add_input_port") < calls.index("register_input") + + def test_source_sets_current_input_channel_id(self, source_handler): + asyncio.run(source_handler.start_worker(EmptyRequest())) + # The main loop reads current_input_channel_id while processing the + # markers below, so it must already point at the starter channel. + assert source_handler.context.current_input_channel_id == SOURCE_STARTER_CHANNEL + + @pytest.mark.timeout(2) + def test_source_enqueues_start_channel_then_end_channel(self, source_handler): + asyncio.run(source_handler.start_worker(EmptyRequest())) + elements = _drain(source_handler.context.input_queue) + + # The whole point of the source branch: a StartChannel marker that is + # not alignment-gated, immediately followed by a port-aligned + # EndChannel that closes the synthetic port and drives the source to + # completion. Order, alignment types and the no-reply command id are + # all part of the contract. + assert [_summarize(element) for element in elements] == [ + ( + SOURCE_STARTER_CHANNEL, + "StartChannel", + EmbeddedControlMessageType.NO_ALIGNMENT, + {WORKER_ID}, + "StartChannel", + NO_REPLY_COMMAND_ID, + ), + ( + SOURCE_STARTER_CHANNEL, + "EndChannel", + EmbeddedControlMessageType.PORT_ALIGNMENT, + {WORKER_ID}, + "EndChannel", + NO_REPLY_COMMAND_ID, + ), + ] + + @pytest.mark.timeout(2) + def test_source_markers_carry_an_empty_scope(self, source_handler): + asyncio.run(source_handler.start_worker(EmptyRequest())) + # A self-addressed marker propagates nowhere, so the scope stays empty. + assert [ + e.payload.scope for e in _drain(source_handler.context.input_queue) + ] == [ + [], + [], + ] + + @pytest.mark.timeout(2) + def test_source_markers_carry_an_unwrappable_empty_request(self, source_handler): + asyncio.run(source_handler.start_worker(EmptyRequest())) + # AsyncRpcServer unwraps the payload with get_one_of before dispatching. + # Leaving the oneof unset delivers None to the StartChannel/EndChannel + # handlers, and betterproto's __eq__ cannot see the difference — so the + # request has to be checked through the same unwrapping the server does. + for element in _drain(source_handler.context.input_queue): + (command,) = element.payload.command_mapping.values() + assert get_one_of(command.command) == EmptyRequest() + + def test_source_does_not_start_mat_reader_threads(self, source_handler): + asyncio.run(source_handler.start_worker(EmptyRequest())) + input_manager = source_handler.context.input_manager + input_manager.start_input_port_mat_reader_threads.assert_not_called() - def test_markers_carry_a_command_addressed_to_this_worker(self): - context = make_context(is_source=True) + def test_source_branch_is_checked_before_the_mat_reader_branch(self): + # is_source wins even when reader threads exist: the elif is never + # evaluated, so the predicate itself must go unqueried. + handler = _build_handler( + is_source=True, + mat_reader_threads={PortIdentity(0, False): [MagicMock()]}, + ) + asyncio.run(handler.start_worker(EmptyRequest())) + input_manager = handler.context.input_manager + input_manager.get_input_port_mat_reader_threads.assert_not_called() + input_manager.start_input_port_mat_reader_threads.assert_not_called() + assert handler.context.input_queue.size() == 2 - start(context) + # --- materialization-reader branch ----------------------------------- - start_ecm, end_ecm = (element.payload for element in drain(context.input_queue)) - for ecm, method in ((start_ecm, "StartChannel"), (end_ecm, "EndChannel")): - assert ecm.scope == [] - # The mapping is keyed by the receiving worker, which is this - # worker itself since the source feeds its own input queue. - assert list(ecm.command_mapping) == [WORKER_ID] - invocation = ecm.command_mapping[WORKER_ID] - assert invocation.method_name == method - assert invocation.command_id == -1 + def test_mat_reader_branch_starts_the_reader_threads(self, mat_reader_handler): + asyncio.run(mat_reader_handler.start_worker(EmptyRequest())) + input_manager = mat_reader_handler.context.input_manager + input_manager.start_input_port_mat_reader_threads.assert_called_once_with() - def test_a_started_source_is_left_running_on_a_second_start(self): - # transit_to() is a no-op when already in the target state, so a - # repeated start must not raise nor bump the reported version. - context = make_context(is_source=True) - start(context) + def test_mat_reader_branch_skips_all_source_effects(self, mat_reader_handler): + asyncio.run(mat_reader_handler.start_worker(EmptyRequest())) + context = mat_reader_handler.context + # A non-source worker is driven to RUNNING by its upstream's channel + # markers, not by this handler; fabricating a port or a starter channel + # here would inject a second, bogus input. + context.state_manager.transit_to.assert_not_called() + context.input_manager.add_input_port.assert_not_called() + context.input_manager.register_input.assert_not_called() + assert context.current_input_channel_id is None + assert context.input_queue.size() == 0 - response = start(context) + # --- fall-through ---------------------------------------------------- + + def test_fall_through_does_nothing(self, fall_through_handler): + asyncio.run(fall_through_handler.start_worker(EmptyRequest())) + context = fall_through_handler.context + input_manager = context.input_manager + # get_input_port_mat_reader_threads returns a mapping, and the branch + # is guarded by its truthiness: an empty mapping means there is nothing + # to read, so the predicate is consulted but every effect is skipped. + input_manager.get_input_port_mat_reader_threads.assert_called_once_with() + context.state_manager.transit_to.assert_not_called() + input_manager.add_input_port.assert_not_called() + input_manager.register_input.assert_not_called() + input_manager.start_input_port_mat_reader_threads.assert_not_called() + assert context.current_input_channel_id is None + assert context.input_queue.size() == 0 - assert response.state == WorkerState.RUNNING - assert response.state_version == 1 + # --- response -------------------------------------------------------- - def test_rejects_a_start_before_the_worker_is_ready(self): - # UNINITIALIZED -> RUNNING is not a legal edge; the handler must fail - # loudly instead of half-starting the source. - context = make_context(is_source=True, initial_state=WorkerState.UNINITIALIZED) + @pytest.mark.parametrize( + "handler_fixture", + ["source_handler", "mat_reader_handler", "fall_through_handler"], + ) + def test_returns_the_state_paired_with_its_version(self, handler_fixture, request): + handler = request.getfixturevalue(handler_fixture) + result = asyncio.run(handler.start_worker(EmptyRequest())) + # Every branch reports through the same atomic read, so the version + # must travel with the state rather than be dropped or defaulted. + assert result == WorkerStateResponse( + WorkerState.RUNNING, state_version=STATE_VERSION + ) + handler.context.state_manager.get_state_with_version.assert_called_once_with() + + def test_reports_whatever_state_the_state_manager_holds(self, fall_through_handler): + # A non-source worker with no materialized inputs is still READY here; + # the handler must not substitute RUNNING for it. + state_manager = fall_through_handler.context.state_manager + state_manager.get_state_with_version.return_value = (WorkerState.READY, 3) + result = asyncio.run(fall_through_handler.start_worker(EmptyRequest())) + assert result == WorkerStateResponse(WorkerState.READY, state_version=3) + + # --- real collaborators ---------------------------------------------- + + def test_source_branch_wires_up_a_real_input_manager(self): + handler = _build_real_handler(is_source=True, initial_state=WorkerState.READY) + result = asyncio.run(handler.start_worker(EmptyRequest())) + input_manager = handler.context.input_manager + + # READY -> RUNNING is a real transition, so the version advances off 0. + assert result == WorkerStateResponse(WorkerState.RUNNING, state_version=1) + assert handler.context.input_queue.size() == 2 + + # The synthetic port really exists and really owns the starter channel. + # InputManager has no public reader for the whole port map, so only the + # exhaustive port check reaches for _ports; the rest goes through the + # public accessors. + assert list(input_manager._ports) == [PortIdentity(0, False)] + assert list(input_manager.get_all_channel_ids()) == [SOURCE_STARTER_CHANNEL] + assert input_manager.get_port_id(SOURCE_STARTER_CHANNEL) == PortIdentity( + 0, False + ) + port = input_manager.get_port(PortIdentity(0, False)) + assert port.get_channels() == {SOURCE_STARTER_CHANNEL} + assert port.get_schema() == Schema() + + # add_input_port records a reader-runnable list for every port it + # creates, even with no materialization URIs, so the real mapping is + # {port: []} and not {}. The mock-based tests above model the guard as + # "are there reader threads?"; on a real InputManager it is really + # "has any input port been added?". + assert input_manager.get_input_port_mat_reader_threads() == { + PortIdentity(0, False): [] + } + def test_rejected_transition_leaves_nothing_enqueued(self): + # UNINITIALIZED only permits READY, so a source told to start before it + # finished initializing cannot reach RUNNING. + handler = _build_real_handler( + is_source=True, initial_state=WorkerState.UNINITIALIZED + ) with pytest.raises(InvalidTransitionException): - start(context) + asyncio.run(handler.start_worker(EmptyRequest())) - assert context.state_manager.get_current_state() == WorkerState.UNINITIALIZED + context = handler.context + # transit_to runs ahead of both put calls, so the failure is clean. A + # future reordering that enqueued first would strand two markers in the + # queue with the state never having advanced. assert context.input_queue.size() == 0 + assert context.state_manager.get_state_with_version() == ( + WorkerState.UNINITIALIZED, + 0, + ) assert context.current_input_channel_id is None - - -class TestStartWorkerOnNonSource: - def test_starts_materialization_reader_threads_when_present(self): + # "Clean" has to cover the InputManager too, not just the queue: doing + # the port and channel wiring ahead of the transition would leave a + # worker that is still UNINITIALIZED owning a phantom port 0 and a + # phantom starter channel. + assert context.input_manager.get_input_port_mat_reader_threads() == {} + assert list(context.input_manager.get_all_channel_ids()) == [] + + def test_double_start_enqueues_a_second_pair_of_markers(self): + # THIS PINS OBSERVED CURRENT BEHAVIOR — it is not an assertion that the + # behavior is correct. Starting an already-running source does not + # raise, because StateManager.transit_to returns early when the target + # equals the current state: RUNNING -> RUNNING is a silent no-op that + # does not even bump the version. The handler carries no idempotency + # guard of its own, so a second pair of markers is appended — including + # a second PORT_ALIGNMENT EndChannel on the same port, which may or may + # not be intended. Revisit and rewrite this test if the handler ever + # gains such a guard. + handler = _build_real_handler(is_source=True, initial_state=WorkerState.READY) + first = asyncio.run(handler.start_worker(EmptyRequest())) + assert handler.context.input_queue.size() == 2 + + second = asyncio.run(handler.start_worker(EmptyRequest())) + assert handler.context.input_queue.size() == 4 + # The no-op transition leaves the logical clock untouched, so the two + # responses are indistinguishable. + assert second == first + assert second.state_version == first.state_version == 1 + # Port and channel registration, unlike the queue, is idempotent. + assert list(handler.context.input_manager.get_all_channel_ids()) == [ + SOURCE_STARTER_CHANNEL + ] + + def test_mat_reader_branch_leaves_a_real_state_clock_untouched(self): + # The mocked counterpart above pins this as "transit_to is not called"; + # against a real StateManager the same guarantee is visible from the + # outside, as a worker still READY at version 0. Reading a materialized + # port does not start the worker -- the main loop transitions on the + # first arriving tuple instead. input_manager = MagicMock(spec=InputManager) input_manager.get_input_port_mat_reader_threads.return_value = { - SOURCE_PORT: [object()] + PortIdentity(0, False): [MagicMock()] } - context = make_context(is_source=False, input_manager=input_manager) + handler = _build_real_handler( + is_source=False, + initial_state=WorkerState.READY, + input_manager=input_manager, + ) - response = start(context) + result = asyncio.run(handler.start_worker(EmptyRequest())) input_manager.start_input_port_mat_reader_threads.assert_called_once_with() - # Reading from a materialized port does not start the worker; the - # main loop transitions on the first arriving tuple instead. - assert response.state == WorkerState.READY - assert response.state_version == 0 - assert context.input_queue.size() == 0 - assert context.current_input_channel_id is None - - def test_does_nothing_without_materialization_reader_threads(self): - input_manager = MagicMock(spec=InputManager) - input_manager.get_input_port_mat_reader_threads.return_value = {} - context = make_context(is_source=False, input_manager=input_manager) - - response = start(context) - - input_manager.start_input_port_mat_reader_threads.assert_not_called() - input_manager.add_input_port.assert_not_called() - input_manager.register_input.assert_not_called() - assert response.state == WorkerState.READY - assert response.state_version == 0 - assert context.input_queue.size() == 0 - - def test_reports_a_paused_state_untouched(self): - # The handler is a pure reporter for a non-source worker: whatever - # state the worker is in is echoed back with its version. + assert result == WorkerStateResponse(WorkerState.READY, state_version=0) + assert handler.context.input_queue.size() == 0 + assert handler.context.current_input_channel_id is None + + def test_reports_a_real_paused_state_untouched(self): + # For a non-source worker the handler is a pure reporter, and the state + # it reports is read after the branch rather than snapshotted before + # it. A real StateManager -- unlike a mock returning a constant -- makes + # a hoisted or stale read observable, since PAUSED at version 1 can only + # come from the live clock. input_manager = MagicMock(spec=InputManager) input_manager.get_input_port_mat_reader_threads.return_value = {} - context = make_context(is_source=False, input_manager=input_manager) - context.state_manager.transit_to(WorkerState.PAUSED) - - response = start(context) - - assert response.state == WorkerState.PAUSED - assert response.state_version == 1 + handler = _build_real_handler( + is_source=False, + initial_state=WorkerState.READY, + input_manager=input_manager, + ) + handler.context.state_manager.transit_to(WorkerState.PAUSED) + + result = asyncio.run(handler.start_worker(EmptyRequest())) + + assert result == WorkerStateResponse(WorkerState.PAUSED, state_version=1) + + def test_markers_are_routed_to_the_data_sub_queue(self): + handler = _build_real_handler(is_source=True, initial_state=WorkerState.READY) + asyncio.run(handler.start_worker(EmptyRequest())) + queue = handler.context.input_queue + # ECMElements are data-queue traffic, not control traffic, so a source + # worker's own start/end markers are gated by disable_data — pause and + # backpressure — unlike direct control messages, which keep flowing on + # the control sub-queue. + assert (queue.size_control(), queue.size_data()) == (0, 2)
