The GitHub Actions job "Required Checks" on texera.git/backport/6724-disable-data-sub-queues-registered-after-v1.2 has succeeded. Run started by GitHub user xuang7 (triggered by xuang7).
Head commit for run: 896ef130c3f164d36c77265df2eb5d7018d26125 / Eugene Gu <[email protected]> fix(pyamber): disable data sub-queues registered after disable_data (#6724) This PR fixes two independent crash bugs in the Python worker's `InternalQueue`, both of which silently kill a worker thread while the heartbeat thread stays alive, hanging the execution with no error reported. It also fixes the reconfiguration test harness, whose incorrect await order masqueraded as an engine deadlock. **Fix 1: data leaking out of a paused or backpressured worker.** Pause and backpressure work by disabling the data sub-queues of `InternalQueue`. However, `disable_data` only disables the sub-queues that exist when it is called, while sub-queues are created lazily on a channel's first `put` and start enabled. A data channel whose first message arrives during a disable window therefore came up enabled and was not covered by the disable at all. This is a bug because the worker can then keep processing data while reporting PAUSED, and if the leaked `DataElement` is dequeued inside the pause wait-loop in `main_loop.py`, the control-only pampy `match` raises an uncaught `MatchError` that silently kills the DP thread while the heartbeat thread stays alive, so the execution hangs forever with no error reported. Under backpressure, the leaked channel keeps feeding the congested downstream, defeating flow control. The fix in `InternalQueue.put`: when it registers a new data channel while `_queue_state` is non-empty, the new sub-queue now starts disabled. The registration is done under the same lock used by `disable_data`/`enable_data` (with a double-check to avoid duplicate registration), and the sub-queue is disabled before its first element is enqueued, so the element never becomes dequeuable during the disable window. Control channels are never disabled, and `enable_data` needs no change because it iterates the channel set at release time, which by then includes channels registered mid-pause. ```mermaid sequenceDiagram participant R as Reader thread participant Q as InternalQueue participant DP as DP thread DP->>Q: disable_data(PAUSE) R->>Q: put(): first message on a new channel Note right of Q: _queue_state is not empty, so it starts DISABLED Note right of Q: messages wait in the queue DP->>Q: enable_data(PAUSE) on resume Q-->>DP: messages delivered in order, nothing lost ``` An important consequence, pinned by a regression test: ECMs ride data channels, so an ECM that is such a channel's first message (e.g. a reconfiguration marker reaching a paused worker) is **delayed until resume, not dropped** — `enable_data` re-enables the channel and the ECM is delivered in order. This matches the engine's intended semantics: reconfigurations submitted while paused only take effect on resume (`ExecutionReconfigurationService` dispatches them without awaiting), and the JVM `DPThread` likewise refuses all data-channel traffic, ECMs included, while paused. **Harness fix: `TestUtils.shouldReconfigure` awaited the reconfiguration ack while still paused.** The three `ReconfigurationIntegrationSpec` failures that earlier looked like an engine deadlock were the harness deadlocking itself: it called `Await.result(reconfigureWorkflow(...))` before resuming, but the ack can only be produced after resume delivers the ECM — a circular wait that expired at the 30s command timeout. Production never uses this order. The harness now matches production: dispatch the reconfiguration without awaiting, await the resume ack (which `ResumeHandler` completes only once every worker has acknowledged), then await the reconfiguration ack. This also removes the tests' dependence on source timing. (There is no RUNNING state event to wait on instead: the engine only pushes `ExecutionStateUpdate` for PAUSED and terminal states.) An intermediate revision of this PR instead moved the leak fix to the dequeue side so ECMs could pass while paused; it was reverted (`87939215d`) after review established that delivering ECMs during a pause diverges from the engine's semantics and that the harness order was the actual culprit. **Fix 2: the category query methods iterated `_queue_ids` unguarded.** The six per-category query methods (`is_control_empty`, `is_data_empty`, `size_control`, `size_data`, `in_mem_size`, `is_data_enabled`) iterated the live `_queue_ids` set, which `put()` grows on a channel's first message from a network thread. CPython raises `RuntimeError: Set changed size during iteration` when a set is mutated mid-iteration, so a query racing a first-time registration kills the calling thread — e.g. the DP thread polling `is_data_enabled()` in the main loop — producing the same silent-hang failure mode as Fix 1 through an unrelated mechanism. These methods now go through two private helpers, `_control_queue_ids()` and `_data_queue_ids()`, which take a `tuple()` snapshot of the set and filter it, so the snapshot is part of the API rather than a convention each query method has to remember. This keeps the hot-path queries lock-free: the snapshot copy is a single uninterruptible C-level operation, and a snapshot at most misses a channel registered mid-call, which the next poll observes. Cleanup in the same file: the unreachable `InternalMarker` entry in the isinstance tuple is removed (`InternalMarker` does not subclass `InternalQueueElement`, so it can never reach that branch) and the two identical dispatch branches are merged. Fixes #6723. The regression tests extend the `InternalQueue` spec added by #6444. Two pre-existing bugs found while tracing the reconfiguration failure will be filed separately, as neither is needed for this suite to pass: `_check_and_process_control` processes ECMs under a stale `current_input_channel_id`, and `PauseManager.resume` re-enables ECM_PAUSE channels before checking `_global_pauses` (the Scala `PauseManager` has the correct order). `amber/src/test/python/core/models/test_internal_queue.py` passes with 34 passed and 1 xfailed (the xfail documents a pre-existing `LinkedBlockingMultiQueue` priority bug from #6444, unrelated to this PR). For Fix 1: the late-channel leak test (fails without the fix), release via `enable_data`, stacked pause+backpressure disables, control channels staying unblocked mid-pause, pre-registered and normally-registered baselines, two threaded stress tests racing first-time registrations against disable/enable toggles with exact element counts, and a new test pinning the delayed-ECM semantics (an ECM as a mid-disable channel's first message is not dequeuable while a reason is active and is delivered after `enable_data`, under both PAUSE and BACKPRESSURE). Deleting the born-disabled block turns 5 tests red, including both parametrized cases of the new ECM test. For Fix 2: 6 deterministic parametrized tests (a sentinel key interleaves a real first-time `put()` into each query's iteration) plus a threaded stress test racing 800 first-time registrations against all six queries; all fail with the `RuntimeError` when the snapshot is reverted. For the harness fix, the flip is the evidence: with the engine byte-identical, `ReconfigurationIntegrationSpec` fails 3 of 3 under the old await order and passes 3 of 3 on repeated clean local runs under the production order; `ReconfigurationSpec`, which shares `shouldReconfigure`, passes 2 of 2. The wider `amber/src/test/python/core` tree was run before and after with identical pass/fail sets apart from the new tests (the only failures need a live Iceberg catalog). `scalafmtCheck` and `scalafix --check` pass on the touched Scala scope. Co-authored by: Claude Code (Claude Fable 5) --------- Co-authored-by: Yicong Huang <[email protected]> (cherry picked from commit 2173ec57fcc237b9716caf80d4990ba3df479d99) Report URL: https://github.com/apache/texera/actions/runs/35305260766 With regards, GitHub Actions via GitBox
