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-7954-d7dc698c54ca18e99cb89c98d2e0d3ccc537a875
in repository https://gitbox.apache.org/repos/asf/texera.git

commit 5fc901727c4582778c14035bfe6e4be6ccd02770
Author: Xinyuan Lin <[email protected]>
AuthorDate: Tue Aug 25 21:35:25 2026 +0000

    test(amber): cover the amber client and the async RPC server (#7954)
    
    ### What changes were proposed in this PR?
    
    Two new specs for two files that had none: `AmberClientSpec` and
    `AsyncRPCServerSpec`. 13 tests.
    
    Measured with `WorkflowExecutionService/jacoco` scoped to these two
    suites, one fresh sbt JVM per run.
    
    | File | Codecov (scoped) | JaCoCo line-hit |
    |---|---|---|
    | `AsyncRPCServer.scala` | 21/32 = 65.6% → **27/32 = 84.4%** | 28/32 →
    **32/32 = 100%** |
    | `AmberClient.scala` | 21/47 = 44.7% → **40/47 = 85.1%** | 22/47 →
    **43/47 = 91.5%** |
    
    **Those scoped figures overstate what CI will show, and the difference
    is the whole story of this file.** Codecov reports `AmberClient` at
    83.0%, not the 44.7% a scoped run sees, because the e2e specs
    (`DataProcessingSpec`, `PauseSpec`) are **not** tagged
    `@IntegrationTest` and therefore run in the coverage job, already
    driving a real `AmberClient`. So the CI-visible gain is bounded by the 8
    lines Codecov actually reports as missing — not +25.
    
    I am stating it that way because a sibling file, `ClientActor.scala`,
    was assessed earlier in exactly this trap and turned out **saturated
    rather than undertested**: ~190 lines of new test bought 1 line. These
    two files are not that — `notifyNodeFailure` is genuinely unexercised,
    and `ClusterListenerSpec`'s own header says so in writing — but the
    scoped number is not the number to quote.
    
    ### Verification
    
    17 mutations, **all 17 killed, no survivors.** Each applied one at a
    time against a hash-verified pristine tree.
    
    The first draft claimed no survivors on a 10-mutant table; **six further
    mutants survived it**, and all six now die.
    
    Four of its claims were overstated and are corrected here:
    
    - A dispatch test was described as pinning that the handler saw "the
    request and the context — not the two swapped". It did not.
    - One kill was credited to a `notifyNodeFailure` test, but it only
    pinned *which branch ran*: `ClientActor` replies `Ack` to any message,
    so the assertion could not distinguish the call from a no-op.
    - The error-handler test was described as pinning the `catch` at line
    144. It pins only that `errorHandler` is *called*; the catch's other job
    — swallowing so the stream survives — was unasserted.
    
    ### Traps recorded in the specs, because each would produce a false pass
    
    - **`AmberClient`'s constructor is not side-effect-free.** It does
    `system.actorOf(Props(new ClientActor))` and then blocks on an ask whose
    handler spawns a real `Coordinator` child. It is safe only with an empty
    `PhysicalPlan` and an all-`None` `CoordinatorConfig`, and every client
    must be `shutdown()` in a `finally` or the shared serialized amber JVM
    accumulates live actors.
    - **`AsyncRPCServer.methodsByName` is a memoized `@transient lazy val`**
    built from `getClass.getMethods` on the *first* `receive()`. Reassigning
    `server.handler` afterwards is silently ignored, so a shared server
    across tests would produce a false pass. Every test builds a fresh one.
    - **The ask inside `notifyNodeFailure` is governed by a 1-minute
    implicit timeout.** A non-replying actor would stall a full minute and
    then report "1 TEST FAILED" without naming it, so the spec awaits with
    an explicit 5-second bound instead.
    - **`registerCallback` asserts
    `clientActor.path.address.hasLocalScope`**, so the spec's `ActorSystem`
    must be plain local — not `AmberRuntime.pekkoConfig`, which is
    clustered.
    - The debug-level test restores the previous logback level in a
    `finally`. amber suites are strictly serial, so a leaked level would
    deterministically pollute later suites rather than occasionally.
    
    ### Deliberately not included
    
    Eight of the sixteen gap lines are structurally unwinnable: both
    lazy-val bitmaps, a `MatchError` arm, an `$outer` null guard, an
    `instanceof Object` arm, and three scala-logging `isEnabled` arms.
    
    `AsyncRPCServer:104`'s error-disabled arm has no test-only seam —
    `invokeMethod`, `returnResult` and `noReplyNeeded` are all private, and
    everything here is driven through the public `receive` plus the
    `handler` var. Accepted limitation; no production seam was requested.
    
    Both new specs carry the Apache licence header. No production file is
    touched.
    
    ### Any related issues, documentation, discussions?
    
    Closes #7953
    
    ### How was this PR tested?
    
    ```
    sbt "WorkflowExecutionService/testOnly 
org.apache.texera.amber.engine.common.client.AmberClientSpec 
org.apache.texera.amber.engine.common.rpc.AsyncRPCServerSpec"
    ```
    
    ```
    [info] Total number of tests run: 13
    [info] Tests: succeeded 13, failed 0, canceled 0, ignored 0, pending 0
    ```
    
    `Test/scalafmtCheck` passes.
    
    ### Was this PR authored or co-authored using generative AI tooling?
    
    Generated-by: Claude Code (Opus 5)
    
    ---------
    
    Signed-off-by: Xinyuan Lin <[email protected]>
    Co-authored-by: Copilot Autofix powered by AI 
<[email protected]>
---
 .../engine/common/client/AmberClientSpec.scala     | 323 +++++++++++++++++++++
 .../engine/common/rpc/AsyncRPCServerSpec.scala     | 282 ++++++++++++++++++
 2 files changed, 605 insertions(+)

diff --git 
a/amber/src/test/scala/org/apache/texera/amber/engine/common/client/AmberClientSpec.scala
 
b/amber/src/test/scala/org/apache/texera/amber/engine/common/client/AmberClientSpec.scala
new file mode 100644
index 0000000000..92063009a5
--- /dev/null
+++ 
b/amber/src/test/scala/org/apache/texera/amber/engine/common/client/AmberClientSpec.scala
@@ -0,0 +1,323 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.texera.amber.engine.common.client
+
+import com.twitter.util.{Await => TwitterAwait, Duration => TwitterDuration}
+import org.apache.pekko.actor.{ActorRef, ActorSystem, Address, 
UnhandledMessage}
+import org.apache.pekko.pattern.StatusReply.Ack
+import org.apache.pekko.testkit.{TestKit, TestProbe}
+import org.apache.texera.amber.core.virtualidentity.ChannelIdentity
+import org.apache.texera.amber.core.workflow.{PhysicalPlan, WorkflowContext}
+import 
org.apache.texera.amber.engine.architecture.common.WorkflowActor.{NetworkAck, 
NetworkMessage}
+import org.apache.texera.amber.engine.architecture.coordinator.{
+  CoordinatorConfig,
+  ExecutionStateUpdate,
+  WorkflowRecoveryStatus
+}
+import 
org.apache.texera.amber.engine.architecture.rpc.controlreturns.WorkflowAggregatedState
+import org.apache.texera.amber.engine.common.ambermessage.{
+  NotifyFailedNode,
+  WorkflowFIFOMessage,
+  WorkflowFIFOMessagePayload,
+  WorkflowRecoveryMessage
+}
+import org.apache.texera.amber.engine.common.virtualidentity.util.{CLIENT, 
COORDINATOR}
+import org.scalatest.BeforeAndAfterAll
+import org.scalatest.flatspec.AnyFlatSpecLike
+import org.scalatest.matchers.should.Matchers
+
+import java.util.concurrent.atomic.AtomicInteger
+import java.util.concurrent.{ConcurrentLinkedQueue, CountDownLatch, TimeUnit}
+import scala.concurrent.Await
+import scala.concurrent.duration.DurationInt
+import scala.jdk.CollectionConverters._
+
+/**
+  * Unit test for [[AmberClient]].
+  *
+  * This spec lives in `...engine.common.client` on purpose: `ClientActor` and 
its
+  * companion messages are `private[client]`, so the event-delivery tests 
below could
+  * not be written from any other package.
+  *
+  * Every client is built over the empty plan
+  * (`PhysicalPlan(Set.empty, Set.empty)` + `CoordinatorConfig(None, None, 
None, None)`),
+  * the recipe three `org.apache.texera.web.service` specs already use: the 
constructor
+  * blocks on an `InitializeRequest`, which spawns a real `Coordinator` child, 
and only
+  * an empty plan makes that complete without an engine. Clients are always 
shut down
+  * in a `finally` -- amber's suites share one serially-run JVM, so a leaked 
client
+  * would leave live actors behind for every later suite.
+  */
+class AmberClientSpec
+    extends TestKit(ActorSystem("AmberClientSpec"))
+    with AnyFlatSpecLike
+    with Matchers
+    with BeforeAndAfterAll {
+
+  override def afterAll(): Unit = {
+    try TestKit.shutdownActorSystem(system)
+    finally super.afterAll()
+  }
+
+  private val failedNode = Address("pekko", "AmberClientSpec", "127.0.0.1", 
2552)
+
+  /**
+    * `AmberClient` declares `implicit val timeout = Timeout(1.minute)` for 
the ask in
+    * `notifyNodeFailure`, so a reply that never arrives would stall the suite 
for a
+    * full minute before ScalaTest reports an unnamed failure. Every await 
below uses
+    * this short explicit bound instead.
+    */
+  private val awaitTimeout: TwitterDuration = TwitterDuration.fromSeconds(5)
+
+  private def newClient(
+      actorSystem: ActorSystem = system,
+      errorHandler: Throwable => Unit = _ => ()
+  ): AmberClient =
+    new AmberClient(
+      actorSystem,
+      new WorkflowContext(),
+      PhysicalPlan(Set.empty, Set.empty),
+      CoordinatorConfig(None, None, None, None),
+      errorHandler
+    )
+
+  private def withClient(body: AmberClient => Unit): Unit = {
+    val client = newClient()
+    try body(client)
+    finally client.shutdown()
+  }
+
+  // 
---------------------------------------------------------------------------
+  // notifyNodeFailure
+  // 
---------------------------------------------------------------------------
+
+  "notifyNodeFailure" should "forward the failed node's address to the 
coordinator and complete with Ack" in {
+    // `ClientActor` replies Ack to ANY WorkflowRecoveryMessage, so the Ack on 
its own
+    // says nothing about what was sent. The `Coordinator` it forwards to has 
no arm
+    // for WorkflowRecoveryMessage (`WorkflowActor.receive` is a fixed 
orElse-chain
+    // with no catch-all), so pekko republishes the forwarded message verbatim 
on the
+    // event stream -- which is how the address can be observed from outside a 
client
+    // whose `clientActor` field is class-private.
+    val unhandled = TestProbe()
+    system.eventStream.subscribe(unhandled.ref, classOf[UnhandledMessage])
+    try {
+      withClient { client =>
+        TwitterAwait.result(client.notifyNodeFailure(failedNode), 
awaitTimeout) shouldBe Ack
+
+        val forwarded = 
unhandled.fishForSpecificMessage[UnhandledMessage](10.seconds) {
+          case msg @ UnhandledMessage(_: WorkflowRecoveryMessage, _, _) => msg
+        }
+        forwarded.message shouldBe WorkflowRecoveryMessage(CLIENT, 
NotifyFailedNode(failedNode))
+      }
+    } finally system.eventStream.unsubscribe(unhandled.ref)
+  }
+
+  it should "return an already-satisfied unit future once the client has been 
shut down" in {
+    val client = newClient()
+    client.shutdown()
+
+    val result = client.notifyNodeFailure(failedNode)
+
+    // Short-circuited, not asked: the actor took a PoisonPill in 
`shutdown()`, so an
+    // ask would sit unanswered until the one-minute timeout instead of 
resolving now.
+    TwitterAwait.result(result, awaitTimeout) shouldBe ((): Unit)
+  }
+
+  // 
---------------------------------------------------------------------------
+  // shutdown
+  // 
---------------------------------------------------------------------------
+
+  "shutdown" should "stop the client actor" in {
+    val ownSystem = ActorSystem("AmberClientSpecShutdown")
+    try {
+      val client = newClient(ownSystem)
+      // Same addressing trick as `withIsolatedClient` below.
+      val clientActor: ActorRef = Await.result(
+        ownSystem.actorSelection("/user/$a").resolveOne(5.seconds),
+        5.seconds
+      )
+      val watcher = TestProbe()(ownSystem)
+      watcher.watch(clientActor)
+
+      client.shutdown()
+
+      // `shutdown()` has to actually stop the actor and not merely flip 
`isActive`:
+      // the ClientActor owns a live Coordinator child tree, and amber's 
suites share
+      // one JVM, so a client that only flipped the flag would leak that whole 
tree
+      // into every later suite. The two "once the client has been shut down" 
tests
+      // observe the flag; this one observes the stop.
+      watcher.expectTerminated(clientActor, 10.seconds)
+    } finally TestKit.shutdownActorSystem(ownSystem)
+  }
+
+  // 
---------------------------------------------------------------------------
+  // registerCallback
+  // 
---------------------------------------------------------------------------
+
+  "registerCallback" should "throw once the client has been shut down" in {
+    val client = newClient()
+    client.shutdown()
+
+    val thrown = intercept[RuntimeException] {
+      client.registerCallback[ExecutionStateUpdate](_ => ())
+    }
+    thrown.getMessage shouldBe "amber runtime environment is not active"
+  }
+
+  it should "return the subscription handle, so disposing it silences that 
callback alone" in {
+    withIsolatedClient("AmberClientSpecDispose") { (client, deliver) =>
+      val disposedFired = new AtomicInteger(0)
+      val liveFired = new CountDownLatch(1)
+
+      val disposed = client.registerCallback[ExecutionStateUpdate] { _ =>
+        disposedFired.incrementAndGet()
+        ()
+      }
+      val live = client.registerCallback[ExecutionStateUpdate](_ => 
liveFired.countDown())
+
+      disposed.isDisposed shouldBe false
+      disposed.dispose()
+      disposed.isDisposed shouldBe true
+      live.isDisposed shouldBe false
+
+      deliver(ExecutionStateUpdate(WorkflowAggregatedState.RUNNING))
+
+      // The returned value has to BE the subscription: any freshly built 
Disposable
+      // satisfies the isDisposed assertions above, but only the real handle
+      // unsubscribes the callback. The subject emits to its observers in 
subscription
+      // order, so the still-live callback firing means the disposed one has 
already
+      // had its chance.
+      assert(liveFired.await(10, TimeUnit.SECONDS), "the still-subscribed 
callback never fired")
+      disposedFired.get() shouldBe 0
+    }
+  }
+
+  it should "deliver each client event to every callback registered for that 
type" in {
+    withIsolatedClient("AmberClientSpecFanOut") { (client, deliver) =>
+      val bothFired = new CountDownLatch(2)
+      val seen = new ConcurrentLinkedQueue[String]
+
+      // Two callbacks for the SAME event class. The second registration must 
reuse
+      // the observable the first one created; if it built a fresh subject and
+      // registered a second partial function on the actor instead, the actor's
+      // `pf orElse handlers` chain would route the event to the newer 
function only
+      // and the first callback would never fire.
+      client.registerCallback[ExecutionStateUpdate] { evt =>
+        seen.add(s"first:${evt.state}")
+        bothFired.countDown()
+      }
+      client.registerCallback[ExecutionStateUpdate] { evt =>
+        seen.add(s"second:${evt.state}")
+        bothFired.countDown()
+      }
+      // A callback for a different class must not see this event.
+      client.registerCallback[WorkflowRecoveryStatus] { _ =>
+        seen.add("other-type")
+        ()
+      }
+
+      deliver(ExecutionStateUpdate(WorkflowAggregatedState.RUNNING))
+
+      assert(
+        bothFired.await(10, TimeUnit.SECONDS),
+        s"both callbacks should have fired, saw: ${seen.asScala.toList}"
+      )
+      seen.asScala.toSet shouldBe Set("first:RUNNING", "second:RUNNING")
+    }
+  }
+
+  it should "route an exception thrown by a callback to the error handler and 
keep the subscription alive" in {
+    val handled = new ConcurrentLinkedQueue[Throwable]
+    val seen = new ConcurrentLinkedQueue[String]
+    val handlerCalledTwice = new CountDownLatch(2)
+
+    withIsolatedClient(
+      "AmberClientSpecErrorHandler",
+      errorHandler = t => { handled.add(t); handlerCalledTwice.countDown() }
+    ) { (client, deliver) =>
+      client.registerCallback[ExecutionStateUpdate] { evt =>
+        seen.add(evt.state.toString)
+        throw new IllegalStateException(s"callback blew up on ${evt.state}")
+      }
+
+      // Two events on purpose. The catch has two jobs: route the failure to 
the
+      // errorHandler AND swallow it. If it rethrew, RxJava's LambdaObserver 
would
+      // dispose the subscription on the first escaping onNext and the 
callback would
+      // be silently unsubscribed forever -- so the second event is what 
proves the
+      // swallow, and a single-event test cannot see that regression at all.
+      deliver(ExecutionStateUpdate(WorkflowAggregatedState.RUNNING))
+      deliver(ExecutionStateUpdate(WorkflowAggregatedState.FAILED))
+
+      assert(
+        handlerCalledTwice.await(10, TimeUnit.SECONDS),
+        s"error handler saw only: ${handled.asScala.toList}"
+      )
+      seen.asScala.toList shouldBe List("RUNNING", "FAILED")
+      handled.asScala.toList.map(_.getMessage) shouldBe List(
+        "callback blew up on RUNNING",
+        "callback blew up on FAILED"
+      )
+    }
+  }
+
+  // 
---------------------------------------------------------------------------
+  // fixtures
+  // 
---------------------------------------------------------------------------
+
+  /**
+    * Runs `body` against a client that owns a private [[ActorSystem]], 
together with a
+    * function that delivers a payload to that client's `ClientActor`.
+    *
+    * `AmberClient.clientActor` is class-private, so the actor has to be 
addressed by
+    * path. It is created with an unnamed `system.actorOf`, which makes it 
`/user/$a` in
+    * a system where nothing else has been created under `/user` (pekko's 
TestKit puts
+    * its own probes under `/system`). A private system per test keeps that 
name
+    * deterministic, and `resolveOne` fails the test loudly if the assumption 
ever
+    * breaks rather than silently delivering nothing.
+    */
+  private def withIsolatedClient(
+      systemName: String,
+      errorHandler: Throwable => Unit = _ => ()
+  )(body: (AmberClient, WorkflowFIFOMessagePayload => Unit) => Unit): Unit = {
+    val ownSystem = ActorSystem(systemName)
+    try {
+      val client = newClient(ownSystem, errorHandler)
+      try {
+        val clientActor: ActorRef = Await.result(
+          ownSystem.actorSelection("/user/$a").resolveOne(5.seconds),
+          5.seconds
+        )
+        val probe = TestProbe()(ownSystem)
+        val channelId = ChannelIdentity(COORDINATOR, CLIENT, isControl = true)
+        var messageId = 0L
+        val deliver: WorkflowFIFOMessagePayload => Unit = payload => {
+          messageId += 1
+          clientActor.tell(
+            NetworkMessage(messageId, WorkflowFIFOMessage(channelId, 
messageId, payload)),
+            probe.ref
+          )
+          // The actor acks before it runs the callbacks, so this only proves 
the
+          // message was consumed; the tests still wait on their own latches.
+          probe.expectMsgType[NetworkAck](5.seconds)
+        }
+        body(client, deliver)
+      } finally client.shutdown()
+    } finally TestKit.shutdownActorSystem(ownSystem)
+  }
+}
diff --git 
a/amber/src/test/scala/org/apache/texera/amber/engine/common/rpc/AsyncRPCServerSpec.scala
 
b/amber/src/test/scala/org/apache/texera/amber/engine/common/rpc/AsyncRPCServerSpec.scala
new file mode 100644
index 0000000000..c0c9cb7940
--- /dev/null
+++ 
b/amber/src/test/scala/org/apache/texera/amber/engine/common/rpc/AsyncRPCServerSpec.scala
@@ -0,0 +1,282 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.texera.amber.engine.common.rpc
+
+import ch.qos.logback.classic.spi.ILoggingEvent
+import ch.qos.logback.classic.{Level, Logger => LogbackLogger}
+import ch.qos.logback.core.AppenderBase
+import com.twitter.util.Future
+import org.apache.texera.amber.core.virtualidentity.{ActorVirtualIdentity, 
ChannelIdentity}
+import 
org.apache.texera.amber.engine.architecture.messaginglayer.NetworkOutputGateway
+import org.apache.texera.amber.engine.architecture.rpc.controlcommands.{
+  AsyncRPCContext,
+  ControlInvocation,
+  ControlRequest,
+  DebugCommandRequest,
+  EmptyRequest
+}
+import org.apache.texera.amber.engine.architecture.rpc.controlreturns.{
+  ControlError,
+  ControlReturn,
+  IntResponse,
+  ReturnInvocation
+}
+import org.apache.texera.amber.engine.common.ambermessage.WorkflowFIFOMessage
+import org.apache.texera.amber.util.VirtualIdentityUtils
+import org.scalatest.flatspec.AnyFlatSpec
+import org.scalatest.matchers.should.Matchers
+import org.slf4j.LoggerFactory
+
+import java.util.concurrent.ConcurrentLinkedQueue
+import scala.collection.mutable.ArrayBuffer
+import scala.jdk.CollectionConverters._
+
+/**
+  * Plain-JVM unit test for [[AsyncRPCServer]]: no actor system, no DB, no 
network.
+  * The server is driven through its public surface only -- the `handler` var 
plus
+  * `receive` -- with a [[NetworkOutputGateway]] whose send handler captures 
the
+  * emitted messages into an in-memory buffer (the same fixture shape
+  * `AsyncRPCClientSpec` uses).
+  *
+  * Every test builds a FRESH server: `AsyncRPCServer.methodsByName` is a 
memoized
+  * `@transient lazy val` built from `handler.getClass.getMethods` on the first
+  * `receive`, so re-assigning `handler` on a server that has already 
dispatched
+  * once is silently ignored and would make later tests vacuous.
+  */
+class AsyncRPCServerSpec extends AnyFlatSpec with Matchers {
+
+  private val serverId = ActorVirtualIdentity("rpc-server")
+  private val senderId = ActorVirtualIdentity("rpc-sender")
+
+  // Four DISTINCT identities on purpose. The context that travels inside the
+  // invocation carries its own two, so neither of them can stand in for the 
wire
+  // sender or for the server's own id: production is free to route the reply 
to
+  // `context.sender` instead of to the actor the message actually came from, 
and a
+  // fixture that reused one literal for both could not see the difference.
+  private val contextOrigin = ActorVirtualIdentity("rpc-context-origin")
+  private val contextTarget = ActorVirtualIdentity("rpc-context-target")
+  private val context = AsyncRPCContext(contextOrigin, contextTarget)
+
+  /** Fresh server, its handler stub, and the buffer its output gateway writes 
into. */
+  private def newFixture(
+      actorId: ActorVirtualIdentity = serverId
+  ): (AsyncRPCServer, AsyncRPCServerSpecHandler, 
ArrayBuffer[WorkflowFIFOMessage]) = {
+    val sent = ArrayBuffer[WorkflowFIFOMessage]()
+    val out = new NetworkOutputGateway(actorId, msg => { sent += msg; () })
+    val server = new AsyncRPCServer(out, actorId)
+    val handler = new AsyncRPCServerSpecHandler
+    server.handler = handler
+    (server, handler, sent)
+  }
+
+  private def invocation(
+      methodName: String,
+      commandId: Long,
+      request: ControlRequest = EmptyRequest()
+  ): ControlInvocation = ControlInvocation(methodName, request, context, 
commandId)
+
+  /** The single control message the server replied with; fails loudly if it 
sent 0 or 2+. */
+  private def soleReturn(sent: ArrayBuffer[WorkflowFIFOMessage]): 
ReturnInvocation = {
+    sent should have size 1
+    val msg = sent.head
+    msg.channelId shouldBe ChannelIdentity(serverId, senderId, isControl = 
true)
+    msg.payload shouldBe a[ReturnInvocation]
+    msg.payload.asInstanceOf[ReturnInvocation]
+  }
+
+  // 
---------------------------------------------------------------------------
+  // dispatch
+  // 
---------------------------------------------------------------------------
+
+  "receive" should "dispatch a known method case-insensitively and return its 
value under the request's commandId" in {
+    val (server, handler, sent) = newFixture()
+    // A request with FIELDS, not the default EmptyRequest: two EmptyRequests 
are
+    // equal, so an EmptyRequest here would be satisfied by a production that 
threw
+    // the wire request away and handed the handler one it invented.
+    val request = DebugCommandRequest("worker-9", "print stats")
+
+    // The lookup map is keyed on the lower-cased name, so the mixed-case
+    // methodName on the wire must still resolve.
+    server.receive(invocation("EchoNumber", 7L, request), senderId)
+
+    // The handler saw the request and the context from the invocation -- not, 
say,
+    // the two swapped, or the context rebuilt from the sender id.
+    handler.lastRequest shouldBe request
+    handler.lastContext shouldBe context
+    soleReturn(sent) shouldBe ReturnInvocation(7L, IntResponse(42))
+  }
+
+  it should "answer a failed handler future with a ControlError naming the 
failure" in {
+    val (server, _, sent) = newFixture()
+
+    server.receive(invocation("failingFuture", 11L), senderId)
+
+    val ret = soleReturn(sent)
+    ret.commandId shouldBe 11L
+    ret.returnValue shouldBe a[ControlError]
+    ret.returnValue.asInstanceOf[ControlError].errorMessage should 
include("future rejected")
+  }
+
+  it should "unwrap the cause when the handler throws synchronously, rather 
than reporting the reflective wrapper" in {
+    val (server, _, sent) = newFixture()
+
+    // A synchronous throw inside the handler reaches the server wrapped in an
+    // InvocationTargetException; production rethrows `e.getCause` so the 
sender
+    // sees the real failure.
+    server.receive(invocation("throwingCall", 12L), senderId)
+
+    val ret = soleReturn(sent)
+    ret.commandId shouldBe 12L
+    val error = ret.returnValue.asInstanceOf[ControlError]
+    error.errorMessage should include("handler exploded")
+    error.errorMessage should not include "InvocationTargetException"
+  }
+
+  it should "log an error naming the unknown method and dispatch nothing when 
no handler method matches" in {
+    val actorId = ActorVirtualIdentity("rpc-unknown-method")
+    withCapturedLogs(actorId, Level.ERROR) { events =>
+      val sent = ArrayBuffer[WorkflowFIFOMessage]()
+      val out = new NetworkOutputGateway(actorId, msg => { sent += msg; () })
+      val server = new AsyncRPCServer(out, actorId)
+      val handler = new AsyncRPCServerSpecHandler
+      server.handler = handler
+
+      server.receive(invocation("noSuchMethod", 13L), senderId)
+
+      // Nothing ran: an unresolvable name must not fall through onto some 
other
+      // entry of the lookup map.
+      handler.invocationCount shouldBe 0
+
+      // DELIBERATELY NOT asserting on `sent`. Production sends nothing at all 
from
+      // this arm -- it only logs -- so the caller's promise for this 
commandId is
+      // never resolved, and a mistyped or version-skewed method name surfaces 
as a
+      // hang on the sender's own timeout instead of as a ControlError. That 
is a
+      // defect in AsyncRPCServer, reported rather than pinned here: an 
assertion
+      // that the server stays silent would fail the day somebody fixes it.
+      val errors = events().filter(_.getLevel == 
Level.ERROR).map(_.getFormattedMessage)
+      errors should contain("No methods found with name nosuchmethod")
+    }
+  }
+
+  it should "emit the command trace naming the method and the sender when 
debug logging is enabled" in {
+    val actorId = ActorVirtualIdentity("rpc-debug-trace")
+    withCapturedLogs(actorId, Level.DEBUG) { events =>
+      val (server, _, sent) = newFixture(actorId)
+
+      server.receive(invocation("echoNumber", 14L), senderId)
+
+      sent should have size 1
+      val debugs = events().filter(_.getLevel == 
Level.DEBUG).map(_.getFormattedMessage)
+      debugs.filter(_.startsWith("receive command: echonumber ")) match {
+        case Seq(only) =>
+          only should include(s"from $senderId")
+          only should include("(controlID: 14)")
+        case other =>
+          fail(s"expected exactly one command trace for echonumber, got: 
$other")
+      }
+    }
+  }
+
+  it should "dispatch but send no reply when the commandId is negative" in {
+    val (server, handler, sent) = newFixture()
+
+    // Negative ids are the fire-and-forget convention (AsyncRPCClient's
+    // IgnoreReplyAndDoNotLog is -2): the call still runs, the result is 
dropped.
+    server.receive(invocation("echoNumber", -2L), senderId)
+
+    handler.invocationCount shouldBe 1
+    sent shouldBe empty
+  }
+
+  // 
---------------------------------------------------------------------------
+  // logging fixture
+  // 
---------------------------------------------------------------------------
+
+  /** Buffers the events a logger emits; appends may arrive from any thread. */
+  private final class CollectingAppender extends AppenderBase[ILoggingEvent] {
+    private val events = new ConcurrentLinkedQueue[ILoggingEvent]
+
+    override def append(event: ILoggingEvent): Unit = events.add(event)
+
+    def snapshot: Seq[ILoggingEvent] = events.asScala.toSeq
+  }
+
+  /**
+    * Runs `body` with the [[AsyncRPCServer]] logger for `actorId` pinned to 
`level`
+    * and a capturing appender attached, handing it a live view of what was 
logged.
+    *
+    * `AmberLogging` names the logger after the actor id, so a distinctive id 
keeps
+    * the level change isolated from every other suite -- which matters because
+    * amber's suites share one JVM and run strictly serially. The previous 
level and
+    * additivity are restored in a `finally`; additivity is switched off 
meanwhile so
+    * the captured lines do not also reach the console and rolling-file 
appenders.
+    */
+  private def withCapturedLogs[T](actorId: ActorVirtualIdentity, level: Level)(
+      body: (() => Seq[ILoggingEvent]) => T
+  ): T = {
+    val loggerName = s"${VirtualIdentityUtils.toShorterString(actorId)}] 
[AsyncRPCServer"
+    val logger = 
LoggerFactory.getLogger(loggerName).asInstanceOf[LogbackLogger]
+    val appender = new CollectingAppender
+    val previousLevel = logger.getLevel
+    val previousAdditive = logger.isAdditive
+    appender.setContext(logger.getLoggerContext)
+    appender.setName(s"async-rpc-server-spec-appender-${actorId.name}")
+    appender.start()
+    logger.addAppender(appender)
+    logger.setLevel(level)
+    logger.setAdditive(false)
+    try body(() => appender.snapshot)
+    finally {
+      logger.detachAppender(appender)
+      logger.setAdditive(previousAdditive)
+      logger.setLevel(previousLevel)
+      appender.stop()
+    }
+  }
+}
+
+/**
+  * Handler stub for [[AsyncRPCServerSpec]].
+  *
+  * Deliberately a top-level, public class. `AsyncRPCServer.methodsByName` is 
built
+  * from `handler.getClass.getMethods` and dispatch goes through 
`Method.invoke`, so
+  * a handler emitted with package-private access would fail every dispatch 
with an
+  * IllegalAccessException and route the happy-path test into the same error 
branch
+  * as the failure tests, making both vacuous.
+  */
+class AsyncRPCServerSpecHandler {
+
+  var lastRequest: ControlRequest = _
+  var lastContext: AsyncRPCContext = _
+  var invocationCount: Int = 0
+
+  def echoNumber(request: ControlRequest, ctx: AsyncRPCContext): 
Future[ControlReturn] = {
+    lastRequest = request
+    lastContext = ctx
+    invocationCount += 1
+    Future.value(IntResponse(42))
+  }
+
+  def failingFuture(request: ControlRequest, ctx: AsyncRPCContext): 
Future[ControlReturn] =
+    Future.exception(new IllegalStateException("future rejected"))
+
+  def throwingCall(request: ControlRequest, ctx: AsyncRPCContext): 
Future[ControlReturn] =
+    throw new IllegalStateException("handler exploded")
+}

Reply via email to