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

commit 512aebde1c6750cc554705aeb783673641981fc8
Author: Xinyuan Lin <[email protected]>
AuthorDate: Sat Aug 29 07:05:59 2026 +0000

    test(amber): cover the end-channel and retry handlers and the range 
partitioner (#8088)
    
    ### What changes were proposed in this PR?
    
    Two new specs and one extended, taking the three files from 4 tests to
    20.
    
    | File | Codecov lines | Branch arms |
    |---|---|---|
    | `RetryWorkflowHandler.scala` | 1/5 → **5/5** | no branches |
    | `EndChannelHandler.scala` | 14/19 → **18/19** | 4/10 → **8/10** |
    | `RangeBasedShufflePartitioner.scala` | 17/21 → **20/21** | 9/30 →
    10/30 |
    
    **+11 fully-covered lines and +5 branch arms.** JaCoCo line-hit moves +8
    (the two metrics differ because three of the gained lines were already
    line-hit and flip only by completing a second arm).
    
    ### Two files from the original scope contribute nothing, and are named
    rather than quietly dropped
    
    - **`DataProcessorRPCHandlerInitializer`** — the nominated target of
    this bundle — measured **8/15 before and 8/15 after**. Its only real
    logic is already covered, and its five remaining lines are `???` stubs
    that the existing spec deliberately declined to pin.
    - **`OutputManager`** stays at 98/108. Its storage path looked like the
    single biggest opportunity here (~18 lines) and it is already covered by
    the `DataProcessingSpec` end-to-end workflows, which provision real
    result and state documents.
    
    The assessment put this bundle at 25 lines and 15 arms. **It delivered
    11 and 5.** The gap was almost entirely one file:
    `RangeBasedShufflePartitioner` was assessed at +16 on a claimed 4/21
    baseline, but the real baseline is 17/21 because **a spec already
    exists** — `RangeBasedShuffleSpec.scala`, under `messaginglayer` rather
    than the mirrored `sendsemantics/partitioners` path and without
    "Partitioner" in its name, which is why a class-name search missed it.
    That spec was extended in place rather than a duplicate created, and the
    real win there is +3.
    
    ### `retryWorkflow` has no caller in main
    
    `RetryWorkflowHandler.retryWorkflow` is invoked from nowhere in `main`.
    It exists as the declared `rpc RetryWorkflow(RetryWorkflowRequest)`
    endpoint in `coordinatorservice.proto`, so it is a live API surface
    rather than dead internal code — but the new spec pins the **RPC
    contract**, not observed production behaviour, and that distinction is
    worth stating before someone reads 1/5 → 5/5 as covering a hot path.
    
    ### Verification
    
    **One survivor, reported rather than killed or dressed up as
    equivalent.** Exchanging `getAllPorts.forall(portId =>
    isPortCompleted(portId))` for `getAllPorts.forall(p =>
    getPort(p).completed)` at `EndChannelHandler:64` **survives** — measured
    twice. These are genuinely different predicates:
    `InputManager.isPortCompleted` branches, and for a port fed from
    materialization it returns the reader *thread's* `finished` flag and
    ignores `.completed` entirely. The fixture registers both input ports
    with empty URI lists, so only the non-materialization branch is ever
    reachable. Closing it needs storage-backed input ports or a widened
    private — a production change, refused. **This is an out-of-reach
    mutant, not an equivalent one**, and the spec's own scaladoc says so, so
    the next reader does not mistake the two-port fixture for full coverage
    of the completion predicate.
    
    **`EndChannelHandler:54` can never become a fully-covered line** and is
    not counted as one. It moves 0/4 → 2/4 arms. The two dead arms are the
    `catch pf` isDefinedAt-false rethrow — dead because `ErrorUtils.safely`
    returns a *total* PartialFunction — and the `$outer` null guard in the
    generated anonfun's constructor.
    
    Two adversarial reviewers returned ten findings; all repaired. One
    correction worth naming: two assertions of the form `command ==
    EmptyRequest()` were **guaranteed by the generated type** —
    `controlcommands.proto:75` is literally `message EmptyRequest{}` with no
    fields, so every instance equals every other. They were removed rather
    than left looking like constraints.
    
    Measured with one fresh sbt JVM per side, the jacoco directory removed
    between runs, scoped by **suite name** via a byte-identical throwaway
    `.sbt` filter that was deleted afterwards. The before-state restored the
    extended spec with `git show HEAD:<path>` and deleted the two new files
    — never `git checkout --` — with the final specs sha1-snapshotted first
    and sha1-verified on restore. Before: 59 suites / 492 tests, 0 failures.
    After: 61 suites / 508 tests, 0 failures.
    
    No production file is touched; `git diff -- '*/src/main/*'` is empty.
    
    ### Any related issues, documentation, discussions?
    
    Closes #8087
    
    ### How was this PR tested?
    
    ```
    sbt "WorkflowExecutionService/testOnly 
org.apache.texera.amber.engine.architecture.worker.promisehandlers.EndChannelHandlerSpec
 
org.apache.texera.amber.engine.architecture.coordinator.promisehandlers.RetryWorkflowHandlerSpec
 
org.apache.texera.amber.engine.architecture.messaginglayer.RangeBasedShuffleSpec"
    ```
    
    ```
    [info] Total number of tests run: 20
    [info] Tests: succeeded 20, failed 0, canceled 0, ignored 0, pending 0
    [info] All tests passed.
    ```
    
    `WorkflowExecutionService/Test/scalafmtCheck` and
    `WorkflowExecutionService/Test/scalafix --check` both pass. Re-run after
    rebasing onto current `main`.
    
    ### 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]>
---
 .../promisehandlers/RetryWorkflowHandlerSpec.scala | 156 +++++++
 .../messaginglayer/RangeBasedShuffleSpec.scala     |  93 ++++-
 .../promisehandlers/EndChannelHandlerSpec.scala    | 446 +++++++++++++++++++++
 3 files changed, 687 insertions(+), 8 deletions(-)

diff --git 
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/promisehandlers/RetryWorkflowHandlerSpec.scala
 
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/promisehandlers/RetryWorkflowHandlerSpec.scala
new file mode 100644
index 0000000000..8d27648ed8
--- /dev/null
+++ 
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/promisehandlers/RetryWorkflowHandlerSpec.scala
@@ -0,0 +1,156 @@
+/*
+ * 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.architecture.coordinator.promisehandlers
+
+import org.apache.texera.amber.core.virtualidentity.ActorVirtualIdentity
+import org.apache.texera.amber.core.workflow.WorkflowContext
+import org.apache.texera.amber.engine.architecture.coordinator.{
+  CoordinatorAsyncRPCHandlerInitializer,
+  CoordinatorConfig,
+  CoordinatorProcessor
+}
+import org.apache.texera.amber.engine.architecture.rpc.controlcommands.{
+  AsyncRPCContext,
+  ControlInvocation,
+  RetryWorkflowRequest
+}
+import 
org.apache.texera.amber.engine.architecture.worker.WorkflowWorker.MainThreadDelegateMessage
+import org.apache.texera.amber.engine.common.ambermessage.WorkflowFIFOMessage
+import org.apache.texera.amber.engine.common.virtualidentity.util.{CLIENT, 
COORDINATOR}
+import org.scalatest.flatspec.AnyFlatSpec
+
+import scala.collection.mutable.ArrayBuffer
+
+/**
+  * `retryWorkflow` is the coordinator-side entry point behind "retry" on a 
paused, failed
+  * execution. It is a two-step fan-out: tell each worker the request names to 
discard and re-run
+  * the tuple it died on (`retryCurrentTuple`), then resume the whole workflow 
by calling
+  * `resumeWorkflow` on itself.
+  *
+  * The order is the contract. A resume that went out first would let an 
untouched worker step past
+  * the failing tuple before it was told to retry it, so the per-worker 
preparation is dispatched
+  * before the resume even though it is not awaited: the retries are 
fire-and-forget
+  * (`Future.collect(...).unit` is discarded), and the reply the caller sees 
is the resume's.
+  *
+  * There is no call site for `retryWorkflow` inside the engine — it is 
reached only from a client
+  * over the coordinator RPC surface — so these tests characterize a declared 
RPC contract rather
+  * than an internally exercised path.
+  *
+  * The harness mirrors `DebugCommandHandlerSpec`: a real 
`CoordinatorProcessor` whose output
+  * handler collects the dispatched control messages, so what is asserted is 
the wire message each
+  * worker would receive. No ActorSystem and no live workers are needed — 
which is also why the
+  * returned `Future` is never awaited, only inspected for pendency: 
`resumeWorkflow` only
+  * completes once a coordinator loop that does not exist here replies to it, 
and that pendency is
+  * itself what proves the caller was handed the resume's own reply rather 
than a fabricated one.
+  */
+class RetryWorkflowHandlerSpec extends AnyFlatSpec {
+
+  /** Two distinct workers, so "reached every worker" is distinguishable from 
"reached one". */
+  private val firstWorkerId = ActorVirtualIdentity("Worker:WF1-udf-main-0")
+  private val secondWorkerId = ActorVirtualIdentity("Worker:WF1-udf-main-1")
+
+  /** The retry comes from the client, so a handler that reused `ctx` would be 
visible. */
+  private val rpcContext = AsyncRPCContext(CLIENT, COORDINATOR)
+
+  private def newFixture()
+      : (CoordinatorAsyncRPCHandlerInitializer, 
ArrayBuffer[WorkflowFIFOMessage]) = {
+    val sent = ArrayBuffer[WorkflowFIFOMessage]()
+    val outputHandler: Either[MainThreadDelegateMessage, WorkflowFIFOMessage] 
=> Unit = {
+      case Right(m) => sent += m
+      case _        => ()
+    }
+    val cp = new CoordinatorProcessor(
+      new WorkflowContext(),
+      CoordinatorConfig(None, None, None, None),
+      COORDINATOR,
+      outputHandler
+    )
+    (new CoordinatorAsyncRPCHandlerInitializer(cp), sent)
+  }
+
+  private def dispatched(sent: ArrayBuffer[WorkflowFIFOMessage]): 
Seq[ControlInvocation] =
+    sent.toSeq.collect {
+      case WorkflowFIFOMessage(_, _, invocation: ControlInvocation) => 
invocation
+    }
+
+  behavior of "RetryWorkflowHandler"
+
+  it should "send a retry to every worker the request names" in {
+    val (init, sent) = newFixture()
+
+    init.retryWorkflow(
+      RetryWorkflowRequest(Seq(firstWorkerId, secondWorkerId)),
+      rpcContext
+    )
+
+    val retries = dispatched(sent).filter(_.methodName == "retryCurrentTuple")
+    assert(retries.size == 2)
+    // Addressed to the workers named in the request body, not to the client 
that asked.
+    assert(retries.map(_.context.receiver) == Seq(firstWorkerId, 
secondWorkerId))
+    assert(retries.forall(_.context.sender == COORDINATOR))
+    // Deliberately no assertion on `command`: the generated stub types it as 
`EmptyRequest`, a
+    // protobuf message with no fields, so every instance equals every other 
and no edit to the
+    // handler could make such an assertion fail without first failing to 
compile.
+  }
+
+  it should "resume the workflow through the coordinator itself" in {
+    val (init, sent) = newFixture()
+
+    val response = 
init.retryWorkflow(RetryWorkflowRequest(Seq(firstWorkerId)), rpcContext)
+
+    val resumes = dispatched(sent).filter(_.methodName == "resumeWorkflow")
+    assert(resumes.size == 1)
+    // The coordinator resumes the workflow by calling its own handler, so 
both ends of the
+    // invocation are the coordinator; the client is not in the loop.
+    assert(resumes.head.context == AsyncRPCContext(COORDINATOR, COORDINATOR))
+    // And the reply the caller gets IS the resume's: `resumeWorkflow` hands 
back a promise that
+    // completes only when the coordinator loop answers, and no such loop runs 
here, so the
+    // response must still be pending. A handler that dispatched the resume 
and then returned its
+    // own `Future.value(EmptyReturn())` would tell the client the workflow 
had resumed before it
+    // had.
+    assert(!response.isDefined)
+  }
+
+  it should "prepare every worker before resuming the workflow" in {
+    val (init, sent) = newFixture()
+
+    init.retryWorkflow(
+      RetryWorkflowRequest(Seq(firstWorkerId, secondWorkerId)),
+      rpcContext
+    )
+
+    // A resume dispatched ahead of a retry would let a worker run past the 
tuple it failed on
+    // before it was told to re-run it.
+    assert(
+      dispatched(sent).map(_.methodName) ==
+        Seq("retryCurrentTuple", "retryCurrentTuple", "resumeWorkflow")
+    )
+  }
+
+  it should "dispatch only the resume when the request names no workers" in {
+    val (init, sent) = newFixture()
+
+    init.retryWorkflow(RetryWorkflowRequest(Seq.empty), rpcContext)
+
+    // The worker list drives the fan-out; an empty one still resumes, because 
"retry" on an
+    // execution whose workers all completed is just a resume.
+    assert(dispatched(sent).map(_.methodName) == Seq("resumeWorkflow"))
+  }
+}
diff --git 
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/RangeBasedShuffleSpec.scala
 
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/RangeBasedShuffleSpec.scala
index 2f906f5973..668b4fda99 100644
--- 
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/RangeBasedShuffleSpec.scala
+++ 
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/RangeBasedShuffleSpec.scala
@@ -56,33 +56,43 @@ class RangeBasedShuffleSpec extends AnyFlatSpec with 
MockFactory {
   "RangeBasedShuffleSpec" should "return 0 when value is less than rangeMin" 
in {
     val tuple = Tuple.builder(schema).add(attr, -600).build()
     val idx = partitioner.getBucketIndex(tuple)
-    assert(idx.next() == 0)
+    // `toList`, not `next()`: range partitioning routes each tuple to exactly 
one bucket. The
+    // `Iterator[Int]` return type is shared with BroadcastPartitioner, which 
legitimately yields
+    // many, so "exactly one" is this partitioner's contract and not the 
type's.
+    assert(idx.toList == List(0))
+
+    // -600 sits in the one window where the clamp and the raw arithmetic 
agree: Long division
+    // truncates toward zero, so (-600 - -400) / 201 is 0 with or without the 
guard. -1000 is
+    // outside it -- unclamped the arithmetic yields -2, and 
`OutputManager.passTupleToDownstream`
+    // feeds the bucket index straight into `partitioner.allReceivers(...)`.
+    val farBelow = Tuple.builder(schema).add(attr, -1000).build()
+    assert(partitioner.getBucketIndex(farBelow).toList == List(0))
   }
 
   "RangeBasedShuffleSpec" should "return last receiver when value is more than 
rangeMax" in {
     val tuple = Tuple.builder(schema).add(attr, 800).build()
     val idx = partitioner.getBucketIndex(tuple)
-    assert(idx.next() == 4)
+    assert(idx.toList == List(4))
   }
 
   "RangeBasedShuffleSpec" should "find index correctly" in {
     var tuple = Tuple.builder(schema).add(attr, -400).build()
     var idx = partitioner.getBucketIndex(tuple)
-    assert(idx.next() == 0)
+    assert(idx.toList == List(0))
 
     tuple = Tuple.builder(schema).add(attr, -200).build()
     idx = partitioner.getBucketIndex(tuple)
-    assert(idx.next() == 0)
+    assert(idx.toList == List(0))
 
     tuple = Tuple.builder(schema).add(attr, -199).build()
     idx = partitioner.getBucketIndex(tuple)
-    assert(idx.next() == 1)
+    assert(idx.toList == List(1))
   }
 
   "RangeBasedShuffleSpec" should "handle different data types correctly" in {
     var tuple = Tuple.builder(schema).add(attr, -90).build()
     var idx = partitioner.getBucketIndex(tuple)
-    assert(idx.next() == 1)
+    assert(idx.toList == List(1))
 
     val partitioning2: RangeBasedShufflePartitioning =
       RangeBasedShufflePartitioning(
@@ -104,7 +114,7 @@ class RangeBasedShuffleSpec extends AnyFlatSpec with 
MockFactory {
     val doubleSchema: Schema = Schema().add(doubleAttr)
     tuple = Tuple.builder(doubleSchema).add(doubleAttr, -90.5).build()
     idx = partitioner2.getBucketIndex(tuple)
-    assert(idx.next() == 1)
+    assert(idx.toList == List(1))
 
     val partitioning3: RangeBasedShufflePartitioning =
       RangeBasedShufflePartitioning(
@@ -126,7 +136,74 @@ class RangeBasedShuffleSpec extends AnyFlatSpec with 
MockFactory {
     val longSchema: Schema = Schema().add(longAttr)
     tuple = Tuple.builder(longSchema).add(longAttr, -90L).build()
     idx = partitioner3.getBucketIndex(tuple)
-    assert(idx.next() == 1)
+    assert(idx.toList == List(1))
+  }
+
+  "RangeBasedShuffleSpec" should "refuse a range attribute it cannot widen to 
a Long" in {
+    // Only LONG / INTEGER / DOUBLE reach the bucket arithmetic. Anything else 
has to fail loudly:
+    // `fieldVal` starts at -1, so a silently skipped type would route every 
tuple of that column
+    // to bucket 0 (below rangeMin) instead of spreading it, and the sort a 
RangePartition exists
+    // to enable would be wrong rather than absent.
+    val stringAttr: Attribute = new Attribute("Attr4", AttributeType.STRING)
+    // Two range attributes, the unsupported one FIRST. Only 
`rangeAttributeNames.head` is
+    // consulted, so a partitioner that reached for any other name would find 
the INTEGER column
+    // and bucket the tuple instead of refusing it -- which for the one 
production declarer
+    // (SortPartitionsOpDesc) means routing by an unintended column rather 
than failing.
+    val stringSchema: Schema = Schema().add(stringAttr).add(attr)
+    val partitioning4: RangeBasedShufflePartitioning =
+      RangeBasedShufflePartitioning(
+        400,
+        List(
+          ChannelIdentity(identifier, fakeID1, isControl = false),
+          ChannelIdentity(identifier, fakeID2, isControl = false)
+        ),
+        Seq("Attr4", "Attr1"),
+        -400,
+        600
+      )
+    val partitioner4: RangeBasedShufflePartitioner = 
RangeBasedShufflePartitioner(partitioning4)
+    val tuple = Tuple.builder(stringSchema).add(stringAttr, "100").add(attr, 
5).build()
+
+    val failure = 
intercept[RuntimeException](partitioner4.getBucketIndex(tuple))
+
+    // The message names the offending *type*, not the attribute: "Attr4" 
would not tell the user
+    // what about the column is unsupported.
+    assert(failure.getMessage == "unsupported attribute type: string")
+  }
+
+  "RangeBasedShuffleSpec" should "collapse a receiver when its channels name 
it twice" in {
+    // Two channels can land on the same worker. `allReceivers` is what the 
bucket index is an
+    // index *into*, and its size is also the divisor behind the bucket width, 
so keeping the
+    // duplicate would both mis-address a bucket and shrink every bucket.
+    val duplicated: RangeBasedShufflePartitioning =
+      RangeBasedShufflePartitioning(
+        400,
+        List(
+          ChannelIdentity(identifier, fakeID1, isControl = false),
+          ChannelIdentity(identifier, fakeID2, isControl = false),
+          ChannelIdentity(identifier, fakeID1, isControl = false),
+          ChannelIdentity(identifier, fakeID3, isControl = false)
+        ),
+        Seq("Attr1"),
+        -400,
+        600
+      )
+    val deduplicating: RangeBasedShufflePartitioner = 
RangeBasedShufflePartitioner(duplicated)
+
+    // First occurrence wins and channel order is preserved.
+    assert(deduplicating.allReceivers == Seq(fakeID1, fakeID2, fakeID3))
+    // ...and the bucket width follows the deduplicated count: 3 receivers 
over [-400, 600] give
+    // 334 keys each, so -100 is still in the first bucket. With the duplicate 
counted the width
+    // would be 251 and -100 would land in the second.
+    assert(
+      deduplicating.getBucketIndex(Tuple.builder(schema).add(attr, 
-100).build()).toList == List(0)
+    )
+    // ...and so does the one place the receiver *count* addresses a bucket 
rather than sizing one:
+    // the above-rangeMax clamp returns the last index, which is 2 for the 
three deduplicated
+    // receivers and would be 3 -- past the end of `allReceivers` -- for the 
four raw channels.
+    assert(
+      deduplicating.getBucketIndex(Tuple.builder(schema).add(attr, 
800).build()).toList == List(2)
+    )
   }
 
 }
diff --git 
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/promisehandlers/EndChannelHandlerSpec.scala
 
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/promisehandlers/EndChannelHandlerSpec.scala
new file mode 100644
index 0000000000..e9ba636963
--- /dev/null
+++ 
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/promisehandlers/EndChannelHandlerSpec.scala
@@ -0,0 +1,446 @@
+/*
+ * 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.architecture.worker.promisehandlers
+
+import com.twitter.util.{Await, Duration, Future}
+import org.apache.texera.amber.core.executor.OperatorExecutor
+import org.apache.texera.amber.core.state.State
+import org.apache.texera.amber.core.tuple.{
+  AttributeType,
+  FinalizeExecutor,
+  FinalizePort,
+  Schema,
+  Tuple,
+  TupleLike
+}
+import org.apache.texera.amber.core.virtualidentity.{
+  ActorVirtualIdentity,
+  ChannelIdentity,
+  OperatorIdentity,
+  PhysicalOpIdentity
+}
+import org.apache.texera.amber.core.workflow.{PhysicalLink, PortIdentity}
+import org.apache.texera.amber.engine.architecture.rpc.controlcommands.{
+  AsyncRPCContext,
+  ConsoleMessageTriggeredRequest,
+  ControlInvocation,
+  EmptyRequest
+}
+import 
org.apache.texera.amber.engine.architecture.rpc.controlreturns.EmptyReturn
+import 
org.apache.texera.amber.engine.architecture.sendsemantics.partitionings.OneToOnePartitioning
+import org.apache.texera.amber.engine.architecture.worker.WorkflowWorker.{
+  DPInputQueueElement,
+  MainThreadDelegateMessage
+}
+import org.apache.texera.amber.engine.architecture.worker.{
+  DataProcessor,
+  DataProcessorRPCHandlerInitializer,
+  OperatorLogicPause,
+  UserPause
+}
+import org.apache.texera.amber.engine.common.ambermessage.{StateFrame, 
WorkflowFIFOMessage}
+import org.apache.texera.amber.engine.common.virtualidentity.util.COORDINATOR
+import org.scalatest.flatspec.AnyFlatSpec
+
+import java.util.concurrent.LinkedBlockingQueue
+import scala.collection.mutable.ArrayBuffer
+import scala.util.control.ControlThrowable
+
+/**
+  * `endChannel` runs on a worker when the END_CHANNEL embedded control 
message reaches it, i.e.
+  * when one input channel has delivered everything it will ever deliver. It 
is the mirror of
+  * [[StartChannelHandlerSpec]]'s handler and does four things:
+  *
+  *   1. resolves the input port from the channel the ECM arrived on and marks 
that port completed;
+  *   2. resets the channel's input batch (`initBatch(channelId, 
Array.empty)`), so the DP loop's
+  *      `hasUnfinishedInput` stops reporting leftover tuples on a channel 
that has just ended;
+  *   3. asks the operator for a boundary state (`produceStateOnFinish`) and 
emits it if there is
+  *      one, then hands the operator's finish output (`onFinishMultiPort`) to 
the output iterator,
+  *      in that order;
+  *   4. appends the input port's `FinalizePort` marker to the output stream;
+  *   5. finalizes the *output* ports, but only once every input port is 
completed.
+  *
+  * Steps 1, 2, 4 and 5 sit outside the `try`, so an operator that throws in 
step 3 — at either of
+  * its two call sites — still gets its port marked completed and its markers 
appended. What the
+  * operator throws is absorbed by `ErrorUtils.safely` and reported through
+  * `handleExecutorException` (console message + in-place pause), leaving the 
RPC itself
+  * successful; the one exception is `scala.util.control.ControlThrowable`, 
which `safely`
+  * rethrows. Both directions are pinned below.
+  *
+  * The `forall` in step 5 is why the fixture always registers two input 
ports: a worker whose
+  * other input port is still open must NOT finalize its output, or the 
downstream region would
+  * see the executor finish while data is still arriving on the second port. 
The fixture pins that
+  * `forall`, but only over link-fed ports: `InputManager.isPortCompleted` 
branches, and for a port
+  * fed from materialization it reports the reader THREAD's `finished` flag 
and ignores
+  * `completed` entirely. That arm needs storage-backed input ports, which 
this fixture cannot
+  * build without real Iceberg documents, so substituting 
`getPort(p).completed` for
+  * `isPortCompleted(p)` is not distinguishable here — a known, stated gap 
rather than coverage.
+  *
+  * `produceStateOnFinish` has no override in main today — `OperatorExecutor` 
returns `None` for
+  * every port — so the "operator produced a state" arm is reachable only 
through a test executor.
+  * It is a declared extension point of the same shape as 
`produceStateOnStart`, and the emission
+  * it drives (no loop envelope) is what a Loop End downstream depends on, so 
it is pinned here
+  * rather than left unspecified.
+  *
+  * These tests drive a real [[DataProcessor]] with no ActorSystem: the 
emitted state is read off
+  * the worker's outgoing wire messages, and the port markers are read off the 
DP output iterator
+  * the DP thread would drain.
+  */
+class EndChannelHandlerSpec extends AnyFlatSpec {
+
+  import EndChannelHandlerSpec._
+
+  private val workerId = ActorVirtualIdentity("Worker:WF1-end-channel-main-0")
+  private val upstreamWorkerId = 
ActorVirtualIdentity("Worker:WF1-upstream-main-0")
+  private val otherUpstreamWorkerId = 
ActorVirtualIdentity("Worker:WF1-other-upstream-main-0")
+  private val downstreamWorkerId = 
ActorVirtualIdentity("Worker:WF1-downstream-main-0")
+  private val rpcContext = AsyncRPCContext(COORDINATOR, workerId)
+  private val awaitTimeout = Duration.fromSeconds(5)
+
+  /** The port of the channel the ECM arrives on. Non-zero, so a hard-coded 0 
cannot pass. */
+  private val currentPortId = PortIdentity(2)
+
+  /** A second input port, wired to a channel that is *not* the one the ECM 
arrives on. */
+  private val otherPortId = PortIdentity(7)
+
+  /** The single output port. Distinct from both input port ids. */
+  private val outputPortId = PortIdentity(5)
+
+  private val currentChannelId = ChannelIdentity(upstreamWorkerId, workerId, 
isControl = false)
+  private val otherChannelId = ChannelIdentity(otherUpstreamWorkerId, 
workerId, isControl = false)
+  private val downstreamChannelId =
+    ChannelIdentity(workerId, downstreamWorkerId, isControl = false)
+  private val coordinatorChannelId = ChannelIdentity(workerId, COORDINATOR, 
isControl = true)
+
+  /** Larger than the number of tuples any test produces, so nothing flushes 
by accident. */
+  private val batchSize = 10
+
+  private val schema: Schema = Schema().add("value", AttributeType.INTEGER)
+
+  /**
+    * @param executor                  the operator under the handler.
+    * @param otherInputPortCompleted   whether the *other* input port has 
already finished. The
+    *                                  common single-input case behaves like 
`true`: the channel
+    *                                  being ended is the last one open.
+    */
+  private class Fixture(
+      executor: OperatorExecutor,
+      otherInputPortCompleted: Boolean = true
+  ) {
+    val sent: ArrayBuffer[WorkflowFIFOMessage] = ArrayBuffer()
+
+    private val outputHandler: Either[MainThreadDelegateMessage, 
WorkflowFIFOMessage] => Unit = {
+      case Right(msg) => sent += msg
+      case Left(_)    => ()
+    }
+
+    val dp: DataProcessor =
+      new DataProcessor(workerId, outputHandler, new 
LinkedBlockingQueue[DPInputQueueElement]())
+    dp.executor = executor
+
+    // Two input ports on two different channels; only one of them is the 
channel the ECM came in
+    // on, so a handler that picked an arbitrary channel would read the wrong 
port.
+    dp.inputManager.addPort(currentPortId, schema, List.empty, List.empty)
+    dp.inputManager.addPort(otherPortId, schema, List.empty, List.empty)
+    dp.inputGateway.getChannel(currentChannelId).setPortId(currentPortId)
+    dp.inputGateway.getChannel(otherChannelId).setPortId(otherPortId)
+    // Seed an unconsumed input batch on the channel about to be ended. 
`endChannel` resets it
+    // (`initBatch(channelId, Array.empty)`), and without a seed that reset is 
invisible: the
+    // batch would already be null and `hasUnfinishedInput` already false.
+    dp.inputManager.initBatch(
+      currentChannelId,
+      Array(TupleLike(9).enforceSchema(schema))
+    )
+    dp.inputManager.getPort(otherPortId).completed = otherInputPortCompleted
+
+    // One downstream data channel, which is also what creates the output 
buffer `emitState`
+    // writes to, plus a control channel that carries the console message on 
the failure paths.
+    dp.outputManager.addPort(outputPortId, schema, None)
+    dp.outputManager.addPartitionerWithPartitioning(
+      PhysicalLink(
+        PhysicalOpIdentity(OperatorIdentity("end-channel-spec-up"), "main"),
+        outputPortId,
+        PhysicalOpIdentity(OperatorIdentity("end-channel-spec-down"), "main"),
+        PortIdentity()
+      ),
+      OneToOnePartitioning(batchSize, Seq(downstreamChannelId))
+    )
+    dp.outputGateway.addOutputChannel(coordinatorChannelId)
+
+    val handler: DataProcessorRPCHandlerInitializer = new 
DataProcessorRPCHandlerInitializer(dp)
+
+    def endChannel(): EmptyReturn = await(handler.endChannel(EmptyRequest(), 
rpcContext))
+
+    /** Everything the DP thread would pull out of the output iterator, in 
order. */
+    def drainOutput(): List[TupleLike] =
+      dp.outputManager.outputIterator.map(_._1).toList
+
+    def emittedStates: Seq[StateFrame] =
+      sent.toSeq.collect { case WorkflowFIFOMessage(_, _, frame: StateFrame) 
=> frame }
+
+    def consoleMessages: Seq[ConsoleMessageTriggeredRequest] =
+      sent.toSeq
+        .collect {
+          case WorkflowFIFOMessage(_, _, invocation: ControlInvocation) =>
+            invocation.command
+        }
+        .collect { case request: ConsoleMessageTriggeredRequest => request }
+  }
+
+  private def await[T](future: Future[T]): T = Await.result(future, 
awaitTimeout)
+
+  /**
+    * `handleExecutorException` pauses with `OperatorLogicPause` specifically. 
`PauseManager` keeps
+    * its pause set private, so the type is probed through `resume`: resuming 
some other type
+    * leaves the worker paused, and resuming `OperatorLogicPause` clears it.
+    */
+  private def assertPausedByOperatorLogic(dp: DataProcessor): Unit = {
+    assert(dp.pauseManager.isPaused)
+    dp.pauseManager.resume(UserPause)
+    assert(dp.pauseManager.isPaused)
+    dp.pauseManager.resume(OperatorLogicPause)
+    assert(!dp.pauseManager.isPaused)
+  }
+
+  behavior of "EndChannelHandler"
+
+  it should "complete the input port of the channel the ECM arrived on, and 
only that one" in {
+    val executor = new RecordingExecutor()
+    val fixture = new Fixture(executor, otherInputPortCompleted = false)
+
+    assert(fixture.endChannel() == EmptyReturn())
+
+    assert(fixture.dp.inputManager.isPortCompleted(currentPortId))
+    assert(!fixture.dp.inputManager.isPortCompleted(otherPortId))
+    // `currentPortId.id`, not the other input port's and not the PortIdentity 
itself — and for
+    // BOTH finish callbacks, which take the port independently.
+    assert(executor.statePorts.toList == List(currentPortId.id))
+    assert(executor.finishOutputPorts.toList == List(currentPortId.id))
+  }
+
+  it should "emit the boundary state the operator produced on finish, with no 
loop envelope" in {
+    val state = State(Map("boundary" -> "end"))
+    val fixture = new Fixture(new RecordingExecutor(onFinishState = _ => 
Some(state)))
+
+    assert(fixture.endChannel() == EmptyReturn())
+
+    // Exactly the state the operator returned. A state a Scala handler 
originates carries the
+    // "no loop" defaults (loopCounter 0, loopStartId ""); only the JVM hop in
+    // `DataProcessor.processState` forwards an incoming envelope, and it 
never creates one.
+    assert(fixture.emittedStates == Seq(StateFrame(state, 0L, "")))
+  }
+
+  it should "emit an empty state, because presence and not content gates the 
emission" in {
+    // `isDefined` is the gate, so a State with no fields is still a state: 
empty is not absent.
+    val emptyState = State(Map.empty[String, Any])
+    val fixture = new Fixture(new RecordingExecutor(onFinishState = _ => 
Some(emptyState)))
+
+    assert(fixture.endChannel() == EmptyReturn())
+
+    assert(fixture.emittedStates == Seq(StateFrame(emptyState, 0L, "")))
+  }
+
+  it should "emit no state when the operator produces none, but still finalize 
the input port" in {
+    val fixture = new Fixture(new RecordingExecutor())
+
+    // The channel still has a seeded, unconsumed tuple when the ECM arrives...
+    assert(fixture.dp.inputManager.hasUnfinishedInput)
+
+    assert(fixture.endChannel() == EmptyReturn())
+
+    // ...and does not after: the channel has delivered everything it ever 
will, so `endChannel`
+    // resets the input batch. `hasUnfinishedInput` is what the DP loop 
consults before pulling
+    // another tuple, so leaving it set would keep feeding a channel that has 
just been declared
+    // ended and whose port was marked completed one statement earlier.
+    assert(!fixture.dp.inputManager.hasUnfinishedInput)
+
+    assert(fixture.emittedStates.isEmpty)
+    // "No state" is the ordinary case, not a failure: reaching into the empty 
Option would be
+    // absorbed by `safely` and surface as a console message and a pause 
rather than as a crash,
+    // so the absence of both is what proves the guard is doing the work.
+    assert(fixture.consoleMessages.isEmpty)
+    assert(!fixture.dp.pauseManager.isPaused)
+    assert(fixture.drainOutput().contains(FinalizePort(currentPortId, input = 
true)))
+  }
+
+  it should "put the operator's finish output ahead of the input port marker" 
in {
+    val fixture = new Fixture(
+      new RecordingExecutor(
+        onFinishTuples = _ =>
+          Iterator(
+            TupleLike(1).enforceSchema(schema),
+            TupleLike(2).enforceSchema(schema)
+          )
+      )
+    )
+
+    fixture.endChannel()
+
+    val output = fixture.drainOutput()
+    // What is pinned here is that the operator's finish output reaches the 
output iterator at
+    // all, and in the order the operator produced it.
+    assert(output.collect { case tuple: Tuple => tuple.getField[Int]("value") 
} == List(1, 2))
+    // The marker follows those tuples, but that is `DPOutputIterator`'s doing 
and not this
+    // handler's: `next()` drains `outputIter` to exhaustion before it touches 
`queue`
+    // (OutputManager.DPOutputIterator), so no statement ordering inside 
`endChannel` could put an
+    // appended marker ahead of the finish tuples. Read this as a 
characterization of the drain
+    // order, not as a pin on where `appendSpecialTupleToEnd` is called; that 
placement is pinned
+    // by the two exception tests below, which assert the marker survives a 
throwing operator.
+    assert(output.indexOf(FinalizePort(currentPortId, input = true)) == 2)
+  }
+
+  it should "swallow an operator exception raised at finish, report it, and 
still reply successfully" in {
+    val failure = new RuntimeException("onFinish blew up")
+    val state = State(Map("boundary" -> "end"))
+    // An operator that produces a boundary state AND then fails while 
flushing its finish output.
+    // The state is asked for first, so it is emitted before the failure 
rather than lost with it.
+    val fixture = new Fixture(
+      new RecordingExecutor(onFinishState = _ => Some(state), onFinishTuples = 
_ => throw failure)
+    )
+
+    // The RPC succeeds: the failure is reported out-of-band, not as a 
control-message failure.
+    assert(fixture.endChannel() == EmptyReturn())
+
+    // The boundary state made it out. This is what pins the ORDER of the two 
finish callbacks:
+    // ask for the state, then take the finish output. Taking the output first 
would mean an
+    // operator that throws there never has `produceStateOnFinish` called at 
all, and the state a
+    // downstream Loop End consumes would be silently dropped instead of 
emitted.
+    assert(fixture.emittedStates == Seq(StateFrame(state, 0L, "")))
+
+    // `handleExecutorException` sends a console message carrying the 
throwable...
+    val consoleMessages = fixture.consoleMessages
+    assert(consoleMessages.size == 1)
+    assert(consoleMessages.head.consoleMessage.title == failure.toString)
+    assert(consoleMessages.head.consoleMessage.workerId == workerId.name)
+    // ...and pauses the worker in place.
+    assertPausedByOperatorLogic(fixture.dp)
+    // Steps 4 AND 5 sit outside the try, so a failed operator still gets its 
input-port marker
+    // appended and — this fixture's other input port having already completed 
— its output ports
+    // finalized. A handler that skipped finalization on the failure path 
would leave the
+    // downstream region waiting on a worker that will never produce anything 
again, so the whole
+    // drained stream is asserted here rather than just the presence of the 
input marker.
+    assert(
+      fixture.drainOutput() == List(
+        FinalizePort(currentPortId, input = true),
+        FinalizePort(outputPortId, input = false),
+        FinalizeExecutor()
+      )
+    )
+  }
+
+  it should "swallow an operator Error as well" in {
+    // `ErrorUtils.safely` guards only against ControlThrowable, so an Error 
is handled like any
+    // other throwable rather than escaping.
+    val failure = new Error("produceStateOnFinish failed hard")
+    val fixture = new Fixture(new RecordingExecutor(onFinishState = _ => throw 
failure))
+
+    assert(fixture.endChannel() == EmptyReturn())
+
+    assert(fixture.consoleMessages.map(_.consoleMessage.title) == 
Seq(failure.toString))
+    assertPausedByOperatorLogic(fixture.dp)
+    // The marker is appended outside the try for the FIRST throw site too, 
not only the second:
+    // this operator dies in `produceStateOnFinish`, before anything else in 
the block runs, and
+    // the downstream still learns the input port is finalized.
+    assert(fixture.drainOutput().contains(FinalizePort(currentPortId, input = 
true)))
+    // Step 2 is outside the try as well, and this is the throw site that 
proves it: the batch
+    // reset happens before the operator is ever consulted, so a channel that 
has just ended
+    // cannot keep feeding the DP loop even when the operator dies at its very 
first callback.
+    assert(!fixture.dp.inputManager.hasUnfinishedInput)
+  }
+
+  it should "let a ControlThrowable escape the finish handler" in {
+    // The one throwable `safely` refuses to handle: it is rethrown, so it 
propagates out of the
+    // handler instead of being reported and paused on.
+    val fixture =
+      new Fixture(new RecordingExecutor(onFinishState = _ => throw new 
ControlThrowable {}))
+
+    intercept[ControlThrowable] {
+      fixture.endChannel()
+    }
+
+    assert(fixture.consoleMessages.isEmpty)
+    assert(!fixture.dp.pauseManager.isPaused)
+    // Characterization of the asymmetry, not an endorsement: the port was 
already marked
+    // completed before the try block, yet nothing downstream is told, because 
the markers are
+    // appended after it.
+    assert(fixture.dp.inputManager.isPortCompleted(currentPortId))
+    assert(fixture.drainOutput().isEmpty)
+  }
+
+  it should "leave the output unfinalized while another input port is still 
open" in {
+    val fixture = new Fixture(new RecordingExecutor(), otherInputPortCompleted 
= false)
+
+    fixture.endChannel()
+
+    // Only the input port's own marker. Finalizing the output here would tell 
the downstream
+    // region the executor is done while the second input port is still 
delivering.
+    assert(fixture.drainOutput() == List(FinalizePort(currentPortId, input = 
true)))
+  }
+
+  it should "finalize the output once the last input port completes" in {
+    val fixture = new Fixture(new RecordingExecutor(), otherInputPortCompleted 
= true)
+
+    fixture.endChannel()
+
+    assert(
+      fixture.drainOutput() == List(
+        FinalizePort(currentPortId, input = true),
+        FinalizePort(outputPortId, input = false),
+        FinalizeExecutor()
+      )
+    )
+  }
+}
+
+object EndChannelHandlerSpec {
+
+  /**
+    * Records the port each finish callback was called with, and produces 
whatever the test wants.
+    *
+    * The two callbacks record separately on purpose. `onFinishMultiPort` is a 
thin forward to
+    * `onFinish(port)` (see `OperatorExecutor`), and a multi-input operator — 
HashJoinProbe,
+    * Difference, Aggregate all override `onFinish(port)` — flushes the buffer 
of exactly the port
+    * it is handed, so the operand of each call is observable behaviour rather 
than bookkeeping.
+    * Recording only one of them would leave the other call site's port 
unpinned.
+    */
+  class RecordingExecutor(
+      onFinishState: Int => Option[State] = _ => None,
+      onFinishTuples: Int => Iterator[TupleLike] = _ => Iterator.empty
+  ) extends OperatorExecutor {
+
+    /** Ports handed to `produceStateOnFinish`. */
+    val statePorts: ArrayBuffer[Int] = ArrayBuffer()
+
+    /** Ports handed to `onFinish`, i.e. what `onFinishMultiPort` forwarded. */
+    val finishOutputPorts: ArrayBuffer[Int] = ArrayBuffer()
+
+    override def produceStateOnFinish(port: Int): Option[State] = {
+      statePorts += port
+      onFinishState(port)
+    }
+
+    override def onFinish(port: Int): Iterator[TupleLike] = {
+      finishOutputPorts += port
+      onFinishTuples(port)
+    }
+
+    override def processTuple(tuple: Tuple, port: Int): Iterator[TupleLike] = 
Iterator.empty
+  }
+}

Reply via email to