This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new f642c183dc fix(amber): release the reconfiguration service's engine
callback on teardown (#7823)
f642c183dc is described below
commit f642c183dc5461f31b84f6bd37dc12e917ca3fdd
Author: Xinyuan Lin <[email protected]>
AuthorDate: Mon Sep 21 01:26:58 2026 +0000
fix(amber): release the reconfiguration service's engine callback on
teardown (#7823)
### What changes were proposed in this PR?
`ExecutionReconfigurationService.registerWorkerCompletionCallback`
discarded the `Disposable` returned by `client.registerCallback`, so its
engine callback outlived the service. It now hands it to
`addSubscription(...)`, matching `registerCompletionDiffHandler`
directly below it.
It was the only `registerCallback` site in
`org.apache.texera.web.service` not doing this:
| File | sites | wrapped |
|---|---|---|
| `ExecutionStatsService.scala` | 5 | yes |
| `ExecutionResultService.scala` | 3 | yes |
| `ExecutionConsoleService.scala` | 2 | yes |
| `ExecutionRuntimeService.scala` | 1 | yes |
| `ExecutionReconfigurationService.scala` | 1 | **no** |
### Why this actually leaks, rather than being tidied up by something
else
The teardown path is live: `WorkflowExecutionService.unsubscribeAll`
(`WorkflowExecutionService.scala:185`) calls
`executionReconfigurationService.unsubscribeAll()`, which released every
other subscription but not this one.
And `client.shutdown()` is not an alternative release — it flips
`isActive` and poison-pills the `ClientActor`, but never touches
`registeredObservables`, which is written only in `registerCallback` and
never cleared. The `PublishSubject` chain therefore keeps the subscriber
closure, and through it the service, its `ExecutionStateStore` and its
`Workflow`, for as long as the `AmberClient` is reachable — which
`WorkflowExecutionService.client` guarantees.
### The fix is pinned
New test: `"the worker completion callback" should "release the engine
subscription once the service is unsubscribed"`. Verified with the
production file reverted and restored:
| | production reverted | with fix |
|---|---|---|
| `ExecutionReconfigurationServiceSpec` | **13 passed, 1 failed** | **14
passed** |
The before-state failure is the right one — `ArrayBuffer() did not
contain element class …UpdateExecutorCompleted`, i.e. the disposable was
never disposed because it was never registered.
Two additions to the existing `TestAmberClient` double were needed, and
the reasons are worth stating: an explicit `disposedCallbacks` record,
because disposal-removal from the `callbacks` map alone is
indistinguishable from "never registered" or "cleared by `reset`"; and a
`fireIfRegistered` helper, because the strict `fire` calls `fail(...)`
when nothing is registered and so cannot express "a late event is
inert".
### Two comments were rewritten, not just tests added
Both would otherwise have contradicted the code:
1. `TestAmberClient.registerCallback`'s scaladoc said production
"currently DISCARDS it" and that "no test below fires an engine event
after `unsubscribeAll`, so the suite neither depends on the leak nor
breaks when it is fixed." Both clauses are now false.
2. The routing comment in `"stop announcing completions once the service
is unsubscribed"` called the engine callback an "unrelated, currently
broken, seam". The routing decision — drive that test through
`onWorkerReconfigured` rather than the engine event — is still right and
unchanged, but re-justified: the two are separate subscriptions, so
routing through the engine event would leave an empty batch explainable
by either one being released. That test's assertion is untouched, so it
stays agnostic exactly as #7692 intended.
### Verification
- `ExecutionReconfigurationServiceSpec` 14/14.
- Dependents and neighbours: `ExecutionRuntimeServiceSpec`,
`WorkflowExecutionServiceSpec`, `ExecutionStatsServiceSpec`,
`ExecutionConsoleServiceSpec`, `WorkflowWebsocketResourceSpec`,
`TexeraWebSocketEventSpec` — **59/59 across 6 suites, 0 aborted**.
- `scalafmtCheck`, `Test/scalafmtCheck`, `scalafixAll --check` all pass.
(The one scalafix warning is a pre-existing `// scalafix:ok` in
`OutputManagerSpec`, unrelated.)
### One thing deliberately left alone
The suite still carries an assertion recorded as *observed, not
endorsed*: an N-worker operator fires N `ModifyLogicCompletedEvent`s,
which contradicts `registerCompletionDiffHandler`'s own comment claiming
the frontend is notified once **all** workers finish. That is a separate
defect with its own decision to make, so this PR does not touch it.
### Any related issues, documentation, discussions?
Closes #7822
### How was this PR tested?
```
STORAGE_ICEBERG_CATALOG_TYPE=postgres sbt
"WorkflowExecutionService/testOnly
org.apache.texera.web.service.ExecutionReconfigurationServiceSpec"
```
```
[info] Total number of tests run: 14
[info] Tests: succeeded 14, failed 0, canceled 0, ignored 0, pending 0
```
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Opus 5)
---
.../service/ExecutionReconfigurationService.scala | 8 +--
.../ExecutionReconfigurationServiceSpec.scala | 62 +++++++++++++++++-----
2 files changed, 55 insertions(+), 15 deletions(-)
diff --git
a/amber/src/main/scala/org/apache/texera/web/service/ExecutionReconfigurationService.scala
b/amber/src/main/scala/org/apache/texera/web/service/ExecutionReconfigurationService.scala
index 31611413fd..2147231abb 100644
---
a/amber/src/main/scala/org/apache/texera/web/service/ExecutionReconfigurationService.scala
+++
b/amber/src/main/scala/org/apache/texera/web/service/ExecutionReconfigurationService.scala
@@ -114,9 +114,11 @@ class ExecutionReconfigurationService(
// events into the reconfiguration store so the diff handler above can fire
// ModifyLogicCompletedEvent for the frontend.
protected def registerWorkerCompletionCallback(): Unit = {
- client.registerCallback[UpdateExecutorCompleted]((evt:
UpdateExecutorCompleted) => {
- onWorkerReconfigured(evt.id)
- })
+ addSubscription(
+ client.registerCallback[UpdateExecutorCompleted]((evt:
UpdateExecutorCompleted) => {
+ onWorkerReconfigured(evt.id)
+ })
+ )
}
// Exposed (instead of inlined in the callback) so tests can drive the
diff --git
a/amber/src/test/scala/org/apache/texera/web/service/ExecutionReconfigurationServiceSpec.scala
b/amber/src/test/scala/org/apache/texera/web/service/ExecutionReconfigurationServiceSpec.scala
index 45275ef308..dfc3569385 100644
---
a/amber/src/test/scala/org/apache/texera/web/service/ExecutionReconfigurationServiceSpec.scala
+++
b/amber/src/test/scala/org/apache/texera/web/service/ExecutionReconfigurationServiceSpec.scala
@@ -259,6 +259,12 @@ class ExecutionReconfigurationServiceSpec
/** (method name, request) of every call made through
`coordinatorInterface`. */
val coordinatorCalls: mutable.ArrayBuffer[(String, Any)] =
mutable.ArrayBuffer.empty
+ /** Event class of every Disposable handed out by `registerCallback` that
was later disposed.
+ * Recorded separately from `callbacks` because disposal removes the
entry, which on its own
+ * is indistinguishable from a callback that was never registered (or one
dropped by `reset`).
+ */
+ val disposedCallbacks: mutable.ArrayBuffer[Class[_]] =
mutable.ArrayBuffer.empty
+
override val coordinatorInterface: CoordinatorServiceFs2Grpc[Future, Unit]
=
Proxy
.newProxyInstance(
@@ -281,19 +287,18 @@ class ExecutionReconfigurationServiceSpec
.asInstanceOf[CoordinatorServiceFs2Grpc[Future, Unit]]
/**
- * Honours the Disposable it hands out. Production
- * (`ExecutionReconfigurationService.registerWorkerCompletionCallback`)
currently DISCARDS
- * it, so `unsubscribeAll()` does not release the engine callback -- the
only
- * `registerCallback` call site in this package that is not wrapped in
`addSubscription`
- * (contrast ExecutionStatsService, ExecutionConsoleService,
ExecutionResultService and
- * ExecutionRuntimeService). A fake that ignored disposal would silently
absolve that line.
- * No test below fires an engine event after `unsubscribeAll`, so the
suite neither depends
- * on the leak nor breaks when it is fixed.
+ * Honours the Disposable it hands out, the way the real client's rx
subscription does:
+ * disposing it detaches the callback, so nothing is delivered to it
afterwards. A fake that
+ * ignored disposal would make every `registerCallback` call site in this
package look
+ * correctly scoped, including one that dropped the Disposable on the
floor.
*/
override def registerCallback[T](callback: T => Unit)(implicit ct:
ClassTag[T]): Disposable = {
val clazz = ct.runtimeClass
callbacks(clazz) = callback.asInstanceOf[Any => Unit]
- Disposable.fromAction(() => callbacks.remove(clazz))
+ Disposable.fromAction(() => {
+ callbacks.remove(clazz)
+ disposedCallbacks += clazz
+ })
}
/** Delivers an engine event the way the client's observable would. */
@@ -303,10 +308,17 @@ class ExecutionReconfigurationServiceSpec
fail(s"no callback is registered for ${ct.runtimeClass.getSimpleName}")
)(event)
+ /** Same delivery as `fire`, but tolerates there being no live callback
instead of failing --
+ * the only way to check that an event arriving after teardown really is
inert.
+ */
+ def fireIfRegistered[T <: AnyRef](event: T)(implicit ct: ClassTag[T]):
Unit =
+ callbacks.get(ct.runtimeClass).foreach(deliver => deliver(event))
+
/** Per-test isolation for the one client the suite shares. */
def reset(): Unit = {
callbacks.clear()
coordinatorCalls.clear()
+ disposedCallbacks.clear()
}
def dispose(): Unit = super.shutdown()
@@ -684,9 +696,9 @@ class ExecutionReconfigurationServiceSpec
f.service.unsubscribeAll()
// Driven through the service's own entry point rather than the engine
event on purpose.
// What this test is about is the diff handler's disposable reaching
the
- // SubscriptionManager; the engine callback has an independent
lifetime (see
- // TestAmberClient.registerCallback), and routing through it would
make this test's
- // verdict depend on that unrelated, currently broken, seam.
+ // SubscriptionManager; the engine callback is a separate subscription
with its own
+ // disposable (covered by the test below), and routing through it
would leave the empty
+ // batch explainable by either one being released.
f.service.onWorkerReconfigured(workerB)
}
@@ -700,4 +712,30 @@ class ExecutionReconfigurationServiceSpec
Set(workerA, workerB)
}
}
+
+ it should "release the engine subscription once the service is unsubscribed"
in {
+ val workerA = mkWorker("opA", 0)
+ val workerB = mkWorker("opB", 0)
+ withLiveService(physicalOps = Set(mkPhysicalOp("opA"),
mkPhysicalOp("opB"))) { f =>
+ // A live registration first, so a disposed one below cannot be confused
with one that was
+ // never made.
+ f.client.fire(UpdateExecutorCompleted(workerA))
+ f.stateStore.reconfigurationStore.getState.completedReconfigurations
shouldBe Set(workerA)
+
+ f.service.unsubscribeAll()
+
+ // The Disposable `registerCallback` handed back has to reach the
SubscriptionManager, the
+ // same way the diff handler's does; this is the only thing that can
release it. The
+ // client's own `shutdown` is not an alternative: it flips `isActive`
and poison-pills the
+ // ClientActor but never touches `registeredObservables`, so the subject
keeps the
+ // subscriber -- and with it this service, its state store and its
Workflow -- attached for
+ // as long as the AmberClient is reachable, which
WorkflowExecutionService keeps it.
+ f.client.disposedCallbacks should
contain(classOf[UpdateExecutorCompleted])
+
+ // ...and the detached callback is inert: a late engine event does not
write into the store
+ // of a finished execution.
+ f.client.fireIfRegistered(UpdateExecutorCompleted(workerB))
+ f.stateStore.reconfigurationStore.getState.completedReconfigurations
shouldBe Set(workerA)
+ }
+ }
}