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-8332-5bc7f49a75bcaee257fe24bf7dfd91a83224f9d0 in repository https://gitbox.apache.org/repos/asf/texera.git
commit b928f82b90286b92f15efd8514134578fcb19577 Author: Xinyuan Lin <[email protected]> AuthorDate: Sun Sep 27 22:08:44 2026 +0000 chore(amber): remove the unused UpdateRecoveryStatus and ResendOutputTo payloads (#8332) ### What changes were proposed in this PR? Deletes two of the three `RecoveryPayload` subtypes — `UpdateRecoveryStatus` and `ResendOutputTo` — which have no sender and no handler. Pure deletion, no behaviour change: **−39 lines**. `NotifyFailedNode` is the only subtype still in use (sent from `AmberClient` when a cluster node fails), and it stays, as does the sealed trait. ### History | | | | --- | --- | | **Introduced by** | #1677 (2022-11-07) — "Amber Fault Tolerance: Global Recovery and Detection" | | **Usage removed by** | #2208 (2023-11-15) — "Refactoring of amber engine" removed the `case UpdateRecoveryStatus(...)` and `case ResendOutputTo(...)` handler arms along with the old recovery path | Dead for about three years. They picked up unit tests in #4829 (2026-05-03) during the coverage work, which is why they look live. Removing them also frees the `ActorRef` and `ActorVirtualIdentity` imports, which only `ResendOutputTo` used. > Reviewer note — the spec is edited, not just trimmed. `AmberMessageEnvelopesSpec` used `UpdateRecoveryStatus` as the payload fixture in its "WorkflowRecoveryMessage carries sender and payload" test; that now uses `NotifyFailedNode`, so the envelope contract stays covered. The suite-local `ActorSystem` and its `afterAll` existed **only** to give the `ResendOutputTo` test a real `ActorRef`, so they go with it — which also drops the `BeforeAndAfterAll` mixin and the pekko `ActorSystem`/`TestKit` imports. The remaining 10 tests pass. ### Any related issues, documentation, discussions? Closes #8329 ### How was this PR tested? Existing tests only — this PR adds none; it removes two types and rewrites the assertions that referenced them. Locally, from the repo root with Java 17: - `sbt "WorkflowExecutionService/Test/compile"` — success. - `sbt "WorkflowExecutionService/testOnly *AmberMessageEnvelopesSpec"` — 10 tests, all pass. - `sbt scalafmtCheckAll "scalafixAll --check"` — clean. Verification, re-runnable by a reviewer: ``` git grep -n "UpdateRecoveryStatus\|ResendOutputTo" # only the deleted declarations and their assertions git grep -n "NotifyFailedNode" -- amber/src/main # the live subtype, untouched ``` ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5) --- .../common/ambermessage/RecoveryPayload.scala | 9 +----- .../ambermessage/AmberMessageEnvelopesSpec.scala | 34 ++-------------------- 2 files changed, 4 insertions(+), 39 deletions(-) diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/common/ambermessage/RecoveryPayload.scala b/amber/src/main/scala/org/apache/texera/amber/engine/common/ambermessage/RecoveryPayload.scala index 5f9302d21a..b4b55ad1c8 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/common/ambermessage/RecoveryPayload.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/common/ambermessage/RecoveryPayload.scala @@ -19,16 +19,9 @@ package org.apache.texera.amber.engine.common.ambermessage -import org.apache.pekko.actor.{ActorRef, Address} -import org.apache.texera.amber.core.virtualidentity.ActorVirtualIdentity +import org.apache.pekko.actor.Address sealed trait RecoveryPayload extends Serializable {} -// Notify coordinator on worker recovery starts/ends -final case class UpdateRecoveryStatus(isRecovering: Boolean) extends RecoveryPayload - -// Notify upstream worker to resend output to another worker for recovery -final case class ResendOutputTo(vid: ActorVirtualIdentity, ref: ActorRef) extends RecoveryPayload - // Notify coordinator when the machine fails and triggers recovery final case class NotifyFailedNode(addr: Address) extends RecoveryPayload diff --git a/amber/src/test/scala/org/apache/texera/amber/engine/common/ambermessage/AmberMessageEnvelopesSpec.scala b/amber/src/test/scala/org/apache/texera/amber/engine/common/ambermessage/AmberMessageEnvelopesSpec.scala index bcf32f1b2a..c86823b06b 100644 --- a/amber/src/test/scala/org/apache/texera/amber/engine/common/ambermessage/AmberMessageEnvelopesSpec.scala +++ b/amber/src/test/scala/org/apache/texera/amber/engine/common/ambermessage/AmberMessageEnvelopesSpec.scala @@ -19,25 +19,12 @@ package org.apache.texera.amber.engine.common.ambermessage -import org.apache.pekko.actor.{Address, ActorSystem} -import org.apache.pekko.testkit.TestKit +import org.apache.pekko.actor.Address import org.apache.texera.amber.core.tuple.{Attribute, AttributeType, Schema, Tuple} import org.apache.texera.amber.core.virtualidentity.{ActorVirtualIdentity, ChannelIdentity} -import org.scalatest.BeforeAndAfterAll import org.scalatest.flatspec.AnyFlatSpec -class AmberMessageEnvelopesSpec extends AnyFlatSpec with BeforeAndAfterAll { - - // Suite-local actor system used only by the ResendOutputTo test below; - // shut down via TestKit.shutdownActorSystem in afterAll so threads do not - // outlive the test, matching the cleanup pattern in CoordinatorSpec / - // WorkerSpec. - private val pekkoSystem: ActorSystem = ActorSystem("amber-message-envelopes-test") - - override protected def afterAll(): Unit = { - TestKit.shutdownActorSystem(pekkoSystem) - super.afterAll() - } +class AmberMessageEnvelopesSpec extends AnyFlatSpec { private val channel = ChannelIdentity(ActorVirtualIdentity("from"), ActorVirtualIdentity("to"), isControl = false) @@ -69,7 +56,7 @@ class AmberMessageEnvelopesSpec extends AnyFlatSpec with BeforeAndAfterAll { "WorkflowRecoveryMessage" should "carry the sender and payload as constructed" in { val from = ActorVirtualIdentity("worker-1") - val payload = UpdateRecoveryStatus(isRecovering = true) + val payload = NotifyFailedNode(Address("pekko", "test")) val msg = WorkflowRecoveryMessage(from, payload) assert(msg.from == from) assert(msg.payload == payload) @@ -80,27 +67,12 @@ class AmberMessageEnvelopesSpec extends AnyFlatSpec with BeforeAndAfterAll { // --------------------------------------------------------------------------- "RecoveryPayload subtypes" should "carry their constructor arguments" in { - val update = UpdateRecoveryStatus(isRecovering = true) - assert(update.isRecovering) - - val updateOff = UpdateRecoveryStatus(isRecovering = false) - assert(!updateOff.isRecovering) - val nodeFailure = NotifyFailedNode(Address("pekko", "test")) assert(nodeFailure.addr == Address("pekko", "test")) } - it should "exercise ResendOutputTo via a real ActorRef so the case class wires correctly" in { - val deadRef = pekkoSystem.deadLetters - val vid = ActorVirtualIdentity("downstream") - val payload = ResendOutputTo(vid, deadRef) - assert(payload.vid == vid) - assert(payload.ref == deadRef) - } - it should "be Serializable on every subtype" in { val payloads: Seq[RecoveryPayload] = Seq( - UpdateRecoveryStatus(isRecovering = true), NotifyFailedNode(Address("pekko", "n")) ) payloads.foreach(p => assert(p.isInstanceOf[Serializable]))
