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

commit 950ee0373a5525d3ece4918c2d4b8ccee3f2e390
Author: Xinyuan Lin <[email protected]>
AuthorDate: Sat Aug 29 07:07:42 2026 +0000

    test(amber): cover WorkflowExecutionService state events and teardown 
(#8035)
    
    ### What changes were proposed in this PR?
    
    `WorkflowExecutionServiceSpec` goes from 3 tests to 11. The three
    existing tests stop at construction; the new ones drive the
    metadata-store handler and the teardown path.
    
    | Metric | Before | After |
    |---|---|---|
    | **Codecov (fully-covered lines)** | 26/72 = 36.1% | **36/72 = 50.0%**
    |
    | JaCoCo line-hit | 30/72 | 39/72 |
    | Branch arms | 5/18 | 10/18 |
    
    **+10 fully-covered lines and +5 branch arms.** The two metrics differ
    by one because line 179 was already line-hit and flips only by
    completing its second arm — they are not interchangeable, so both are
    given.
    
    **Of the +10, six are logic and four are not.** Lines 95, 179, 181–184
    are the recovery-banner arm and the four `unsubscribeAll` calls. Lines
    101, 107, 108 and 109 are scalac-generated public accessor pairs
    (confirmed with `javap -p -l`) that move because the teardown test
    assigns those vars from outside. Those four are a real consequence of
    driving teardown rather than coverage farming, but they are not logic
    and I would rather split them out than present all ten as equivalent.
    
    ### What the reviewers found
    
    Two independent adversarial reviewers ran against the first draft. **Six
    of eight fresh mutants survived it.** The three worth naming:
    
    - **`workflowContext.workflowSettings = request.workflowSettings` (line
    90) could be deleted outright and all 8 tests still passed.** Every test
    executes that line, so it was fully covered and entirely unconstrained.
    It is not cosmetic: `outputPortsNeedingStorage` is what
    `CostBasedScheduleGenerator` reads to decide which output ports get
    materialized, and `dataTransferBatchSize` / `executionMode` feed the
    resource allocator. Dropping it silently reverts every execution to the
    default settings while the request's are ignored.
    - **The recovering-state guard was pinned one-sidedly.** The test only
    ever drove `isRecovering` false → true, so weakening
    `newState.isRecovering != oldState.isRecovering` to a bare
    `newState.isRecovering` survived. That mutant is the mirror of the
    failure the test's own comment claimed to guard: an update that *clears*
    the flag emits no event, so the frontend's Recovering banner never comes
    down.
    - **Both the state guard and the `fatalErrors` guard could be deleted
    wholesale** and nothing failed.
    
    All are now killed by a named test. The published mutation table was
    also re-run from scratch, one mutant at a time, because one row's
    failure message had been copy-pasted from another row.
    
    ### Verification
    
    Measured with an identical suite-name filter on both sides, one fresh
    sbt JVM per measurement, `rm -rf` of the jacoco dir between runs,
    counters read per-line out of `jacoco.xml`.
    
    An independent measurer re-derived the figure a second way — over the
    **whole amber unit module**, the actual CI scope, rather than the
    six-suite filter — and got byte-identical per-line data on both sides.
    That independently confirms the six-suite list is complete and that 26 →
    36 is what Codecov will show.
    
    Two corrections that measurement forced, both worth stating:
    
    - The repair round added 3 tests **after** the original measurement, and
    the expectation was that the figure would rise above 36. It did not.
    Those 3 tests bought **zero** additional fully-covered lines — they add
    branch arms on line 94 and mutation-kill strength only. 36 is the
    number.
    - Line 185 already counted as a Codecov hit before this PR despite
    `mi=4, ci=1`, because the method's `return` instruction is attributed to
    it. It is genuinely executed only now.
    
    Full amber unit scope, both sides: **190 suites, the same 7 pre-existing
    Windows-only failures by name, zero new.** The only per-suite change
    anywhere is this spec going 3 → 11 tests, so the new `beforeAll` inserts
    do not leak — `MockTexeraDB` gives each suite its own database, and ids
    9207–9210 are unique across the repo because the Iceberg statistics URI
    is machine-global.
    
    ### Deliberately not included
    
    `executeWorkflow`'s live-runtime half — the 44-line hole — has no seam.
    Line 124 calls `ComputingUnitMaster.createAmberRuntime`, which builds an
    `AmberClient` over `AmberRuntime.actorSystem`, a JVM-global `private
    var` that nothing in unit scope initialises. Injecting a seam would be a
    production change.
    
    Lines 113/114/115 are refused for a sharper reason: JaCoCo probes the
    `try` block only at its exit, so they flip only if compile *and*
    `Workflow.fromCompilationResult` both succeed — which falls straight
    into `createAmberRuntime`. Whether that is survivable depends on
    `ClientEventSpec` having restored the global to `null` in its
    `afterAll`. `WorkflowServiceSpec`'s own header documents this
    cross-suite hazard as its reason for deliberately steering into a
    compile failure instead. Cementing an accident is worse than leaving
    three lines.
    
    Three mutants are reported as live rather than dropped: reordering
    `client.shutdown()` against the four `unsubscribeAll` calls (order is
    not a stated contract), and two whose only kill would be to assert
    current behaviour that is arguably wrong — adding `FAILED` to the
    stuck-banner guard, and pinning which field the teardown guard reads,
    where the only discriminating state currently NPEs.
    
    No production file is touched.
    
    ### Any related issues, documentation, discussions?
    
    Closes #8033
    
    ### How was this PR tested?
    
    ```
    sbt "WorkflowExecutionService/testOnly 
org.apache.texera.web.service.WorkflowExecutionServiceSpec"
    ```
    
    ```
    [info] Tests: succeeded 11, failed 0, canceled 0, ignored 0, pending 0
    [info] All tests passed.
    ```
    
    `WorkflowExecutionService/Test/scalafmtCheck` and
    `WorkflowExecutionService/Test/scalafix --check` both pass.
    
    ### 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]>
---
 .../web/service/WorkflowExecutionServiceSpec.scala | 432 ++++++++++++++++++++-
 1 file changed, 426 insertions(+), 6 deletions(-)

diff --git 
a/amber/src/test/scala/org/apache/texera/web/service/WorkflowExecutionServiceSpec.scala
 
b/amber/src/test/scala/org/apache/texera/web/service/WorkflowExecutionServiceSpec.scala
index d8af24de4b..abf2bb039c 100644
--- 
a/amber/src/test/scala/org/apache/texera/web/service/WorkflowExecutionServiceSpec.scala
+++ 
b/amber/src/test/scala/org/apache/texera/web/service/WorkflowExecutionServiceSpec.scala
@@ -20,13 +20,38 @@
 package org.apache.texera.web.service
 
 import com.google.protobuf.timestamp.Timestamp
-import org.apache.texera.amber.core.workflow.{WorkflowContext, 
WorkflowSettings}
+import io.reactivex.rxjava3.disposables.Disposable
+import org.apache.pekko.actor.ActorSystem
+import org.apache.pekko.testkit.TestKit
+import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, 
WorkflowIdentity}
+import org.apache.texera.amber.core.workflow.{PhysicalPlan, WorkflowContext, 
WorkflowSettings}
 import 
org.apache.texera.amber.core.workflowruntimestate.FatalErrorType.EXECUTION_FAILURE
 import org.apache.texera.amber.core.workflowruntimestate.WorkflowFatalError
+import 
org.apache.texera.amber.engine.architecture.coordinator.CoordinatorConfig
 import 
org.apache.texera.amber.engine.architecture.rpc.controlreturns.WorkflowAggregatedState.{
+  COMPLETED,
   FAILED,
   RUNNING
 }
+import org.apache.texera.amber.engine.common.client.AmberClient
+import org.apache.texera.amber.operator.source.scan.csv.CSVScanSourceOpDesc
+import org.apache.texera.dao.MockTexeraDB
+import org.apache.texera.dao.jooq.generated.enums.WorkflowComputingUnitTypeEnum
+import org.apache.texera.dao.jooq.generated.tables.daos.{
+  UserDao,
+  WorkflowComputingUnitDao,
+  WorkflowDao,
+  WorkflowExecutionsDao,
+  WorkflowVersionDao
+}
+import org.apache.texera.dao.jooq.generated.tables.pojos.{
+  User,
+  Workflow,
+  WorkflowComputingUnit,
+  WorkflowExecutions,
+  WorkflowVersion
+}
+import org.apache.texera.web.WebsocketInput
 import org.apache.texera.web.model.websocket.event.{
   TexeraWebSocketEvent,
   WorkflowErrorEvent,
@@ -36,12 +61,16 @@ import 
org.apache.texera.common.compiler.model.LogicalPlanPojo
 import org.apache.texera.web.model.websocket.request.WorkflowExecuteRequest
 import org.apache.texera.web.storage.ExecutionStateStore
 import org.apache.texera.web.storage.ExecutionStateStore.updateWorkflowState
-import org.scalatest.flatspec.AnyFlatSpec
+import org.scalatest.BeforeAndAfterAll
+import org.scalatest.flatspec.AnyFlatSpecLike
 import org.scalatest.matchers.should.Matchers
 
 import java.net.URI
+import java.sql.{Timestamp => SqlTimestamp}
 import java.time.Instant
 import scala.collection.mutable
+import scala.collection.mutable.ListBuffer
+import scala.reflect.ClassTag
 
 /**
   * Regression guard for the consolidated init-error reporting path (#5921):
@@ -54,19 +83,161 @@ import scala.collection.mutable
   * purpose: construction must stay side-effect-free (all throwing work is in
   * `executeWorkflow`), so a future change that dereferences them during
   * construction would fail here.
+  *
+  * The suite also owns two paths that need more scaffolding than that:
+  *
+  *   - `executeWorkflow`'s compilation-failure arm, driven for real (not 
simulated
+  *     through `errorHandler`) with a CSV scan that has no file selected -- 
the
+  *     recipe `SyncExecutionResourceSpec` establishes and 
`WorkflowServiceSpec`
+  *     reuses. Compilation rejects it inside 
`LogicalPlan.resolveScanSourceOpFileName`
+  *     before `FileResolver` is reached, so no storage is involved, and the 
early
+  *     `return` stops the method before 
`ComputingUnitMaster.createAmberRuntime`,
+  *     which would build an `AmberClient` on `AmberRuntime.actorSystem` -- 
null
+  *     outside a started coordinator, and whatever another suite left behind 
inside
+  *     the shared serial JVM.
+  *
+  *   - `unsubscribeAll`, the session-teardown contract. This is what makes the
+  *     suite carry a TestKit `ActorSystem` and `MockTexeraDB`: the method 
walks the
+  *     four runtime services, and `ExecutionStatsService`'s constructor 
creates its
+  *     iceberg runtime-statistics table and stamps the URI onto a 
`workflow_executions`
+  *     row, so a non-null value for that field cannot be had more cheaply. 
There is
+  *     no second home for the test -- one spec file per source class -- and 
stubbing
+  *     the services out would leave the four teardown calls unobserved, which 
is the
+  *     entire content of the method.
+  *
+  * The runtime half of `executeWorkflow` (everything from 
`createAmberRuntime` on) is
+  * out of reach here for the reason above and is exercised by the integration 
suites.
+  *
+  * Known, pre-existing: `ExecutionStatsService` never shuts down its private
+  * `metricsPersistThread` (it does not override `unsubscribeAll`), so the 
teardown
+  * test below leaks one single-thread executor into the shared JVM, as every
+  * construction of that service already does.
+  *
+  * Three further gaps are deliberately described and NOT pinned, because a 
test that
+  * froze the current behaviour would cement it:
+  *
+  *   - Partial construction. `executeWorkflow` assigns `client` first and the 
four
+  *     services after it, and `WorkflowService` catches whatever they throw 
while
+  *     leaving the half-built service published. `unsubscribeAll` then passes 
its
+  *     `client != null` guard and NPEs on the first null service, aborting 
teardown.
+  *
+  *   - A recovery that ends in failure. `createStateEvent` carves out 
COMPLETED only,
+  *     while the flag itself is lowered only by a 
`WorkflowRecoveryStatus(false)` from
+  *     the engine, which an execution that dies mid-recovery need not ever 
send. Such an
+  *     execution reaches FAILED with the flag still raised and is announced as
+  *     "Recovering" from then on.
+  *
+  *   - `ExecutionReconfigurationService.registerWorkerCompletionCallback` 
discards the
+  *     `Disposable` its `registerCallback` returns instead of handing it to
+  *     `addSubscription`, so that one callback outlives the teardown below. 
Not this
+  *     class's defect, and not this suite's to pin.
   */
-class WorkflowExecutionServiceSpec extends AnyFlatSpec with Matchers {
+class WorkflowExecutionServiceSpec
+    extends TestKit(ActorSystem("WorkflowExecutionServiceSpec"))
+    with AnyFlatSpecLike
+    with Matchers
+    with BeforeAndAfterAll
+    with MockTexeraDB {
+
+  // Distinct from every other suite's ids: the statistics URI is derived from 
wid/eid and
+  // `createDocument` truncates whatever table already sits at it 
(ExecutionStatsServiceSpec
+  // owns 9107/9108).
+  private val testUid: Integer = 9207
+  private val testWid: Integer = 9207
+  private val testEid: Integer = 9208
+
+  /**
+    * The computing unit the fixture execution ran on. A different literal 
from `testWid` on
+    * purpose: `getLatestExecutionID(wid, cuid)` binds two bare `Integer`s 
into one predicate, so a
+    * fixture that reused one number for both would make a transposition of 
the two arguments
+    * produce byte-identical SQL.
+    */
+  private val testCuid: Integer = 9209
+
+  /** A computing unit with no executions, so the cuid leg of the predicate is 
not vacuous. */
+  private val otherCuid: Integer = 9210
+
+  override protected def beforeAll(): Unit = {
+    super.beforeAll()
+    initializeDBAndReplaceDSLContext()
+
+    // Set explicitly rather than left to the SERIAL default: the 
computing-unit and execution
+    // rows below bind `testUid` into their `uid` foreign keys, so a generated 
value would leave
+    // them pointing at a user that does not exist.
+    val user = new User
+    user.setUid(testUid)
+    user.setName("workflow-execution-test-user")
+    user.setEmail(s"[email protected]")
+    new UserDao(getDSLContext.configuration()).insert(user)
+
+    val workflow = new Workflow
+    workflow.setWid(testWid)
+    workflow.setName(s"workflow-execution-test-$testWid")
+    workflow.setContent("{}")
+    workflow.setDescription("")
+    workflow.setCreationTime(new SqlTimestamp(System.currentTimeMillis()))
+    workflow.setLastModifiedTime(new SqlTimestamp(System.currentTimeMillis()))
+    new WorkflowDao(getDSLContext.configuration()).insert(workflow)
+
+    val version = new WorkflowVersion
+    version.setWid(testWid)
+    version.setContent("{}")
+    version.setCreationTime(new SqlTimestamp(System.currentTimeMillis()))
+    new WorkflowVersionDao(getDSLContext.configuration()).insert(version)
+
+    val computingUnitDao = new 
WorkflowComputingUnitDao(getDSLContext.configuration())
+    List(testCuid -> "workflow-execution-test-unit", otherCuid -> 
"workflow-execution-idle-unit")
+      .foreach {
+        case (cuid, name) =>
+          val unit = new WorkflowComputingUnit
+          unit.setCuid(cuid)
+          unit.setUid(testUid)
+          unit.setName(name)
+          unit.setCreationTime(new SqlTimestamp(System.currentTimeMillis()))
+          unit.setType(WorkflowComputingUnitTypeEnum.local)
+          unit.setUri("local://test")
+          unit.setResource("{}")
+          computingUnitDao.insert(unit)
+      }
+
+    // `ExecutionStatsService`'s constructor stamps the runtime-statistics URI 
onto this row,
+    // reaching it from the workflow through workflow_version. `cuid` is set 
because
+    // `getLatestExecutionID` matches on it and a NULL cuid matches no value 
at all.
+    val execution = new WorkflowExecutions
+    execution.setEid(testEid)
+    execution.setVid(version.getVid)
+    execution.setUid(testUid)
+    execution.setCuid(testCuid)
+    execution.setStatus(0.toByte)
+    execution.setStartingTime(new SqlTimestamp(System.currentTimeMillis()))
+    execution.setBookmarked(false)
+    execution.setName("workflow-execution-test-execution")
+    execution.setEnvironmentVersion("test-env")
+    new WorkflowExecutionsDao(getDSLContext.configuration()).insert(execution)
+  }
+
+  override protected def afterAll(): Unit = {
+    try {
+      TestKit.shutdownActorSystem(system)
+    } finally {
+      closeConnectionPool()
+      super.afterAll()
+    }
+  }
 
   private def buildService(
       store: ExecutionStateStore,
-      errorHandler: Throwable => Unit = (_: Throwable) => ()
+      errorHandler: Throwable => Unit = (_: Throwable) => (),
+      logicalPlan: LogicalPlanPojo =
+        LogicalPlanPojo(List.empty, List.empty, List.empty, List.empty),
+      settings: WorkflowSettings = WorkflowSettings()
   ): WorkflowExecutionService = {
     val request = WorkflowExecuteRequest(
       executionName = "test",
       engineVersion = "test",
-      logicalPlan = LogicalPlanPojo(List.empty, List.empty, List.empty, 
List.empty),
+      logicalPlan = logicalPlan,
       replayFromExecution = None,
-      workflowSettings = WorkflowSettings(),
+      workflowSettings = settings,
       emailNotificationEnabled = false,
       computingUnitId = 0,
       warehouseId = None
@@ -94,6 +265,36 @@ class WorkflowExecutionServiceSpec extends AnyFlatSpec with 
Matchers {
     events
   }
 
+  /** A subscription whose only job is to report whether the manager holding 
it let it go. */
+  private final class Tracker {
+    var disposed = false
+    val disposable: Disposable = Disposable.fromAction(() => disposed = true)
+  }
+
+  /**
+    * Empty-plan client, the shape both `ExecutionRuntimeServiceSpec` and
+    * `ExecutionStatsServiceSpec` use: real enough for the services' 
constructors, with the
+    * engine-facing callbacks stubbed so nothing is registered against a live 
coordinator.
+    */
+  private final class TestAmberClient
+      extends AmberClient(
+        system,
+        new WorkflowContext(),
+        PhysicalPlan(Set.empty, Set.empty),
+        CoordinatorConfig(None, None, None, None),
+        _ => ()
+      ) {
+    var shutdownCount = 0
+
+    override def shutdown(): Unit = {
+      shutdownCount += 1
+      super.shutdown()
+    }
+
+    override def registerCallback[T](callback: T => Unit)(implicit ct: 
ClassTag[T]): Disposable =
+      Disposable.empty()
+  }
+
   "WorkflowExecutionService" should
     "surface a recorded fatal error as a WorkflowErrorEvent via the 
metadata-store handler" in {
     val store = new ExecutionStateStore()
@@ -107,6 +308,24 @@ class WorkflowExecutionServiceSpec extends AnyFlatSpec 
with Matchers {
     val errorEvents = events.collect { case e: WorkflowErrorEvent => e }
     errorEvents should have size 1
     errorEvents.head.fatalErrors should contain(err)
+    // The other half of the handler's contract: an update that moves only 
`fatalErrors` must not
+    // also announce a state. `StateStore` filters out updates that change 
nothing at all, so the
+    // guard's whole job is this case -- a write that moves some other field 
while `state` and
+    // `isRecovering` stand still. Without it, recording an error also 
republishes a redundant
+    // WorkflowStateEvent to the session.
+    events.collect { case e: WorkflowStateEvent => e } shouldBe empty
+  }
+
+  it should "apply the request's workflow settings to the shared workflow 
context" in {
+    // The one thing construction does besides wiring the handler, and it is 
not inert: the
+    // context is what the compiler and the schedule generator read later, so 
a dropped assignment
+    // silently runs every execution on WorkflowContext's defaults while the 
request's settings --
+    // batch size, execution mode, and the output ports that need materialized 
storage -- are
+    // ignored. The value has to be a non-default one, or the context's own 
default would answer.
+    val settings = WorkflowSettings(dataTransferBatchSize = 137)
+    val service = buildService(new ExecutionStateStore(), settings = settings)
+
+    service.workflowContext.workflowSettings shouldBe settings
   }
 
   it should "report fatal errors recorded at successive phases through the 
same handler" in {
@@ -143,5 +362,206 @@ class WorkflowExecutionServiceSpec extends AnyFlatSpec 
with Matchers {
     store.metadataStore.updateState(_.withState(RUNNING))
 
     events.collect { case e: WorkflowStateEvent => e } should not be empty
+    // The mirror of the assertion in the fatal-error test: a state-only 
update must not push an
+    // empty WorkflowErrorEvent at the frontend. This is the suite's only 
update that leaves
+    // `fatalErrors` alone, so it is the only place the errors guard can be 
observed suppressing.
+    events.collect { case e: WorkflowErrorEvent => e } shouldBe empty
+  }
+
+  it should "report a recovering execution as Recovering instead of as its 
aggregated state" in {
+    // Recovery is a banner the frontend keeps up over whatever the engine is 
otherwise doing:
+    // the execution really is RUNNING throughout, and the three updates are 
kept apart so that the
+    // recovery flag has to carry the second and third events on its own -- 
the aggregated state
+    // does not move with it, and a handler that only watched `state` would go 
quiet exactly when
+    // the banner needs to go up.
+    //
+    // The flag is driven in BOTH directions on purpose. Raising it is what a 
guard that merely
+    // read `newState.isRecovering` would also do; only clearing it separates 
that from the
+    // `!=` the guard actually needs, and a missing third event is the banner 
never coming down.
+    val store = new ExecutionStateStore()
+    buildService(store)
+    val events = collectEvents(store)
+
+    store.metadataStore.updateState(_.withState(RUNNING))
+    store.metadataStore.updateState(_.withIsRecovering(true))
+    store.metadataStore.updateState(_.withIsRecovering(false))
+
+    events.collect { case e: WorkflowStateEvent => e } shouldBe
+      Seq(
+        WorkflowStateEvent("Running"),
+        WorkflowStateEvent("Recovering"),
+        WorkflowStateEvent("Running")
+      )
+  }
+
+  it should "report the aggregated state of a completed execution still marked 
as recovering" in {
+    // The recovery flag is not cleared on the way to COMPLETED, so without 
the COMPLETED
+    // carve-out a finished execution would be announced as still recovering 
and the frontend
+    // would never take the banner down.
+    //
+    // COMPLETED is the only sentinel this pins, and deliberately so. 
`createStateEvent`'s
+    // carve-out lists exactly that one terminal state, so an execution that 
dies mid-recovery
+    // (FAILED / KILLED with the flag still raised) is still announced as 
"Recovering" -- see the
+    // note in the class header. Adding FAILED to the expectation here would 
cement that; adding
+    // it to the guard is a production question, not this suite's to answer.
+    val store = new ExecutionStateStore()
+    buildService(store)
+    val events = collectEvents(store)
+
+    
store.metadataStore.updateState(_.withState(COMPLETED).withIsRecovering(true))
+
+    events.collect { case e: WorkflowStateEvent => e } shouldBe
+      Seq(WorkflowStateEvent("Completed"))
+  }
+
+  it should "report a compilation failure and leave the workflow and the 
client unset" in {
+    val store = new ExecutionStateStore()
+    val errors = ListBuffer.empty[Throwable]
+    val scan = new CSVScanSourceOpDesc()
+    scan.setOperatorId("scan-op")
+    val service = buildService(
+      store,
+      errors += _,
+      LogicalPlanPojo(List(scan), List.empty, List.empty, List.empty)
+    )
+
+    // This is the assertion that pins the early `return`, and it has to wrap 
the call itself:
+    // without the `return` the run falls through to `createAmberRuntime`, 
dereferences the still
+    // null `workflow` on the way in, and throws out of `executeWorkflow` -- 
which is also what
+    // keeps a "unit" test out of the live-runtime half on whatever actor 
system the shared JVM
+    // holds.
+    noException should be thrownBy service.executeWorkflow()
+
+    errors should have size 1
+    errors.head.getMessage should include("No file selected")
+    // Documentation, not a pin: both fields are still at their `_` defaults 
on this path, because
+    // the only assignments to them sit after the throw and after the 
`return`. Neither assertion
+    // can fail while the two above hold -- they record the post-condition the 
early return leaves
+    // behind, and they are the suite's only read of `workflow`.
+    service.workflow shouldBe null
+    service.client shouldBe null
+  }
+
+  it should "shut the client down and unsubscribe every runtime service on 
teardown" in {
+    // The end of a websocket session. Everything the execution owns has to be 
let go here:
+    // the engine client, and the four services that hold callbacks on it and 
diff handlers on
+    // the shared state store. A survivor keeps publishing into a store nobody 
reads.
+    //
+    // What is asserted is the SET of teardown effects, not their order. 
`AmberClient.shutdown`
+    // only flips a flag and posts a PoisonPill, so whether it runs before or 
after the four
+    // unsubscribes changes only a race against the client actor's mailbox -- 
nothing a
+    // deterministic test can observe, and not an order any contract here 
states.
+    val store = new ExecutionStateStore()
+    val service = buildService(store)
+    val events = collectEvents(store)
+
+    val client = new TestAmberClient
+    val wsInput = new WebsocketInput(_ => ())
+    // `workflow` is only dereferenced by the reconfiguration diff handler 
when a
+    // reconfiguration actually completes, which is not what this test drives.
+    val reconfigurationService = new ExecutionReconfigurationService(client, 
store, workflow = null)
+    val statsService = new ExecutionStatsService(
+      client,
+      store,
+      new WorkflowContext(
+        workflowId = WorkflowIdentity(testWid.longValue()),
+        executionId = ExecutionIdentity(testEid.longValue())
+      )
+    )
+    val runtimeService = new ExecutionRuntimeService(
+      client,
+      store,
+      wsInput,
+      reconfigurationService,
+      logConf = None,
+      workflowId = testWid.longValue(),
+      emailNotificationEnabled = false,
+      userEmailOpt = None,
+      sessionUri = new URI("https://texera.example/session";)
+    )
+    val consoleService = new ExecutionConsoleService(client, store, wsInput, 
new WorkflowContext())
+
+    service.client = client
+    service.executionReconfigurationService = reconfigurationService
+    service.executionStatsService = statsService
+    service.executionRuntimeService = runtimeService
+    service.executionConsoleService = consoleService
+
+    // One tracker per manager, so a teardown call that goes missing names its 
own service.
+    val ownTracker = new Tracker
+    val reconfigurationTracker = new Tracker
+    val statsTracker = new Tracker
+    val runtimeTracker = new Tracker
+    val consoleTracker = new Tracker
+    service.addSubscription(ownTracker.disposable)
+    reconfigurationService.addSubscription(reconfigurationTracker.disposable)
+    statsService.addSubscription(statsTracker.disposable)
+    runtimeService.addSubscription(runtimeTracker.disposable)
+    consoleService.addSubscription(consoleTracker.disposable)
+
+    service.unsubscribeAll()
+
+    withClue("the engine client: ") { client.shutdownCount shouldBe 1 }
+    withClue("the execution's own subscriptions: ") { ownTracker.disposed 
shouldBe true }
+    withClue("executionRuntimeService: ") { runtimeTracker.disposed shouldBe 
true }
+    withClue("executionConsoleService: ") { consoleTracker.disposed shouldBe 
true }
+    withClue("executionStatsService: ") { statsTracker.disposed shouldBe true }
+    withClue("executionReconfigurationService: ") {
+      reconfigurationTracker.disposed shouldBe true
+    }
+
+    // The constructor-registered diff handler is part of "its own 
subscriptions": after teardown
+    // the store is inert, which is what stops a closed session from still 
producing websocket
+    // events for its (now disconnected) client.
+    events.clear()
+    store.metadataStore.updateState(_.withState(RUNNING))
+    events shouldBe empty
+  }
+
+  it should "leave a never-started execution alone on teardown" in {
+    // `unsubscribeAll` also runs for an execution that failed to compile, 
where the client and all
+    // four service fields are still null: teardown of that execution must not 
throw, and its own
+    // subscriptions must still be let go.
+    //
+    // Which field the guard reads is NOT observed here, and cannot honestly 
be. The two states
+    // that would separate the five candidates are (client set, services null) 
-- whose current
+    // behaviour is an NPE, so pinning it would cement the 
partial-construction defect noted in the
+    // class header -- and (client null, services set), which the assignment 
order in
+    // `executeWorkflow` makes unreachable.
+    val store = new ExecutionStateStore()
+    val service = buildService(store)
+    val events = collectEvents(store)
+    val ownTracker = new Tracker
+    service.addSubscription(ownTracker.disposable)
+
+    noException should be thrownBy service.unsubscribeAll()
+
+    ownTracker.disposed shouldBe true
+    events.clear()
+    store.metadataStore.updateState(_.withState(RUNNING))
+    events shouldBe empty
+  }
+
+  // The companion's one expression, exercised directly rather than only 
through the two result
+  // services that call it. Transposing its two arguments -- both bare 
`Integer`s bound into one
+  // predicate -- does reach those siblings, and hard (30 of their tests fail 
on it), so this is not
+  // a hole in the module's coverage. It is the class that owns the method 
stating its own contract,
+  // and failing first and by name when that contract breaks.
+  "WorkflowExecutionService.getLatestExecutionId" should
+    "resolve the newest execution of a workflow on a given computing unit" in {
+    WorkflowExecutionService.getLatestExecutionId(
+      WorkflowIdentity(testWid.longValue()),
+      testCuid.intValue()
+    ) shouldBe Some(ExecutionIdentity(testEid.longValue()))
+  }
+
+  it should "find no execution for a computing unit the workflow has not run 
on" in {
+    // Same workflow row as the case above, so only the cuid leg can produce 
this None -- which is
+    // what keeps that leg from being vacuous. The transposition of the two 
arguments dies in the
+    // case above, where the swapped pair selects no row at all.
+    WorkflowExecutionService.getLatestExecutionId(
+      WorkflowIdentity(testWid.longValue()),
+      otherCuid.intValue()
+    ) shouldBe None
   }
 }

Reply via email to