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-8688-2b154f1456ad4d8c0b3729a8d99e029169590172 in repository https://gitbox.apache.org/repos/asf/texera.git
commit d102e26978cf0a377992b7d3ded03eaccd5d5c75 Author: Xinyuan Lin <[email protected]> AuthorDate: Sat Sep 26 22:59:22 2026 +0000 refactor(amber): remove the unused BackpressurePause (#8688) ### What changes were proposed in this PR? Deletes `BackpressurePause`. It is a `PauseType` that no production code has passed to `PauseManager` since flow control moved onto `ActorMessage`. There is no behaviour change: **+11/−19 lines**. ### History | | | | --- | --- | | **Introduced by** | #1636 (2022-08-20), "Introduce types of pause in Amber". Backpressure paused the worker through `PauseManager` under its own pause type | | **Usage removed by** | #2237 (2023-12-02), "Use ActorMessage for flow control". It deleted the `pauseManager.pause(BackpressurePause)` / `resume(BackpressurePause)` calls | It has been dead for nearly three years. Backpressure still works, but it now bypasses `PauseManager`: `Backpressure(enabled)` arrives as an `ActorCommand` and flips `DPThread.backpressureStatus`. #4533 removed the sibling `SchedulerTimeSlotExpiredPause` for the same reason. > Reviewer note: the specs change in two places, and both are fixture swaps, not lost coverage. > - In `PauseTypeSpec`, the singleton / identity / pattern-match / `Set` cases now cover the remaining three kinds. > - The two `WorkerManagersSpec` `PauseManager` cases used `BackpressurePause` only as "some other pause type". They now use what production actually passes: `OperatorLogicPause` for a global pause (as `DataProcessor` does) and `ECMPause` for a per-channel pause (as ECM alignment does). ### Any related issues, documentation, discussions? Closes #8686 ### How was this PR tested? No new tests. The two existing specs keep their cases with the fixtures swapped. Locally, from the repo root with Java 17: - `sbt "WorkflowExecutionService/Test/compile"`: success. - `sbt "WorkflowExecutionService/testOnly *PauseTypeSpec *WorkerManagersSpec"`: 27 tests, all pass. - `sbt "WorkflowExecutionService/scalafmtCheckAll" "WorkflowExecutionService/scalafixAll --check"`: clean. To re-check: ``` git grep -n BackpressurePause # no hits ``` ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5.5) 🤖 Generated with [Claude Code](https://claude.com/claude-code) --- .../amber/engine/architecture/worker/PauseType.scala | 2 -- .../amber/engine/architecture/worker/PauseTypeSpec.scala | 12 +----------- .../worker/managers/WorkerManagersSpec.scala | 16 ++++++++++------ 3 files changed, 11 insertions(+), 19 deletions(-) diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/PauseType.scala b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/PauseType.scala index 57fd5cee3a..b37b235cb5 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/PauseType.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/PauseType.scala @@ -25,8 +25,6 @@ sealed trait PauseType object UserPause extends PauseType -object BackpressurePause extends PauseType - object OperatorLogicPause extends PauseType case class ECMPause(id: EmbeddedControlMessageIdentity) extends PauseType diff --git a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/PauseTypeSpec.scala b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/PauseTypeSpec.scala index 68af4cec3b..a558aa83aa 100644 --- a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/PauseTypeSpec.scala +++ b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/PauseTypeSpec.scala @@ -36,19 +36,14 @@ class PauseTypeSpec extends AnyFlatSpec { // Widen to PauseType so the compiler doesn't reduce inter-singleton // comparisons to constant `false` at compile time. val u: PauseType = UserPause - val b: PauseType = BackpressurePause val o: PauseType = OperatorLogicPause assert(u == UserPause) - assert(b == BackpressurePause) assert(o == OperatorLogicPause) - assert(u != b) assert(u != o) - assert(b != o) } it should "be the same singleton instance per access (object identity)" in { assert((UserPause: AnyRef) eq UserPause) - assert((BackpressurePause: AnyRef) eq BackpressurePause) assert((OperatorLogicPause: AnyRef) eq OperatorLogicPause) } @@ -75,7 +70,6 @@ class PauseTypeSpec extends AnyFlatSpec { // an ECMPause (with any id) must not collide with any singleton kind. val p: PauseType = ECMPause(EmbeddedControlMessageIdentity("ckpt")) assert(p != UserPause) - assert(p != BackpressurePause) assert(p != OperatorLogicPause) } @@ -85,12 +79,10 @@ class PauseTypeSpec extends AnyFlatSpec { def label(p: PauseType): String = p match { case UserPause => "user" - case BackpressurePause => "backpressure" case OperatorLogicPause => "operator-logic" case ECMPause(_) => "ecm" } assert(label(UserPause) == "user") - assert(label(BackpressurePause) == "backpressure") assert(label(OperatorLogicPause) == "operator-logic") assert(label(ECMPause(EmbeddedControlMessageIdentity("x"))) == "ecm") } @@ -107,13 +99,11 @@ class PauseTypeSpec extends AnyFlatSpec { it should "coexist as distinct elements in a Set without aliasing" in { val active: Set[PauseType] = Set( UserPause, - BackpressurePause, OperatorLogicPause, ECMPause(EmbeddedControlMessageIdentity("ckpt-1")) ) - assert(active.size == 4, "all four pause kinds must be distinct Set elements") + assert(active.size == 3, "all three pause kinds must be distinct Set elements") assert(active.contains(UserPause)) - assert(active.contains(BackpressurePause)) assert(active.contains(OperatorLogicPause)) assert(active.contains(ECMPause(EmbeddedControlMessageIdentity("ckpt-1")))) } diff --git a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/managers/WorkerManagersSpec.scala b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/managers/WorkerManagersSpec.scala index 1932823f5d..4d4f489cc2 100644 --- a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/managers/WorkerManagersSpec.scala +++ b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/managers/WorkerManagersSpec.scala @@ -21,7 +21,11 @@ package org.apache.texera.amber.engine.architecture.worker.managers import org.apache.texera.amber.core.executor.OperatorExecutor import org.apache.texera.amber.core.tuple.{Tuple, TupleLike} -import org.apache.texera.amber.core.virtualidentity.{ActorVirtualIdentity, ChannelIdentity} +import org.apache.texera.amber.core.virtualidentity.{ + ActorVirtualIdentity, + ChannelIdentity, + EmbeddedControlMessageIdentity +} import org.apache.texera.amber.core.workflow.PortIdentity import org.scalatest.flatspec.AnyFlatSpec @@ -171,7 +175,7 @@ class WorkerManagersSpec extends AnyFlatSpec { import org.apache.texera.amber.engine.architecture.logreplay.OrderEnforcer import org.apache.texera.amber.engine.architecture.messaginglayer.{AmberFIFOChannel, InputGateway} import org.apache.texera.amber.engine.architecture.worker.{ - BackpressurePause, + ECMPause, OperatorLogicPause, PauseManager, UserPause @@ -241,9 +245,9 @@ class WorkerManagersSpec extends AnyFlatSpec { val (gw, a, _, _) = newGateway() val pm = new PauseManager(workerId, gw) pm.pause(UserPause) - pm.pause(BackpressurePause) + pm.pause(OperatorLogicPause) pm.resume(UserPause) - // backpressure still pausing → channels stay disabled + // operator-logic pause still active → channels stay disabled assert(pm.isPaused) assert(!a.isEnabled) } @@ -262,10 +266,10 @@ class WorkerManagersSpec extends AnyFlatSpec { val (gw, a, b, _) = newGateway() val pm = new PauseManager(workerId, gw) pm.pauseInputChannel(OperatorLogicPause, List(dataA)) - pm.pauseInputChannel(BackpressurePause, List(dataB)) + pm.pauseInputChannel(ECMPause(EmbeddedControlMessageIdentity("ecm-b")), List(dataB)) pm.resume(OperatorLogicPause) // dataA's only specific pause was OperatorLogicPause → re-enabled. - // dataB still has BackpressurePause → still disabled. + // dataB still has its ECMPause → still disabled. assert(a.isEnabled) assert(!b.isEnabled) }
