Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}

Expand All @@ -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)
}

Expand All @@ -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")
}
Expand All @@ -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"))))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
}
Expand All @@ -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)
}
Expand Down
Loading