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 d32395cd52 fix(amber): reject duplicate worker initialization (#8083)
d32395cd52 is described below
commit d32395cd52da6550577a1487a757d77b02b7cf61
Author: carloea2 <[email protected]>
AuthorDate: Mon Aug 31 18:44:44 2026 +0000
fix(amber): reject duplicate worker initialization (#8083)
### What changes were proposed in this PR?
`OperatorExecution.initWorkerExecution` now checks whether the worker ID
is already a key in the execution map. This prevents a repeated
initialization from silently replacing the existing execution.
The old characterization and pending tests are replaced with one active
regression test. It verifies both that the repeated initialization fails
and that the original execution remains registered.
Before: A repeated worker ID replaced the original execution.
After: A repeated worker ID raises `AssertionError` and preserves the
original execution.
This change removes 19 lines overall.
### Any related issues, documentation, discussions?
Closes #8082
### How was this PR tested?
The regression test failed before the production change with 13 tests
passing and 1 test failing. It passed after the change with all 14 tests
passing.
```shell
sbt "set DAO / Compile / sourceGenerators := Seq.empty"
"WorkflowExecutionService / Test / testOnly
org.apache.texera.amber.engine.architecture.coordinator.execution.OperatorExecutionSpec"
sbt "set DAO / Compile / sourceGenerators := Seq.empty" "scalafixAll
--check"
sbt scalafmtCheckAll
```
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: OpenAI Codex, GPT-5
---
.../coordinator/execution/OperatorExecution.scala | 2 +-
.../execution/OperatorExecutionSpec.scala | 23 ++--------------------
2 files changed, 3 insertions(+), 22 deletions(-)
diff --git
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/OperatorExecution.scala
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/OperatorExecution.scala
index 6315fb06fe..ee7f36d488 100644
---
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/OperatorExecution.scala
+++
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/OperatorExecution.scala
@@ -52,7 +52,7 @@ case class OperatorExecution() {
*/
def initWorkerExecution(workerId: ActorVirtualIdentity): WorkerExecution = {
assert(
- !workerExecutions.contains(workerId),
+ !workerExecutions.containsKey(workerId),
s"WorkerExecution already exists for workerId: $workerId"
)
workerExecutions.put(workerId, WorkerExecution())
diff --git
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/OperatorExecutionSpec.scala
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/OperatorExecutionSpec.scala
index fa16873fbc..3a78f76f12 100644
---
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/OperatorExecutionSpec.scala
+++
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/OperatorExecutionSpec.scala
@@ -78,34 +78,15 @@ class OperatorExecutionSpec extends AnyFlatSpec {
assert(opExec.getWorkerExecution(w) eq workerExec)
}
- // The class docstring claims `initWorkerExecution` throws
- // `AssertionError` on a duplicate worker id, but the implementation's
- // `workerExecutions.contains(workerId)` call resolves to Java
- // `ConcurrentHashMap.contains(Object)`, which checks VALUES rather than
- // KEYS — so the assertion never fires and the second call silently
- // overwrites the prior WorkerExecution. We pin the CURRENT (broken)
- // behavior here so a future fix is noticed in CI, and document the
- // intended contract with `pendingUntilFixed` so the failure surfaces
- // the day the implementation is corrected.
-
it should
- "currently overwrite the previous WorkerExecution on a second init for the
same id " +
- "(characterization of the contains-by-value bug)" in {
+ "reject a second init for the same id without replacing the existing
execution" in {
val opExec = OperatorExecution()
val w = workerId("w-1")
val firstExec = opExec.initWorkerExecution(w)
- val secondExec = opExec.initWorkerExecution(w)
- assert(firstExec ne secondExec, "current impl replaces the prior
WorkerExecution instance")
- assert(opExec.getWorkerExecution(w) eq secondExec)
- }
-
- it should "(desired) throw AssertionError when initWorkerExecution is called
twice for the same id" in pendingUntilFixed {
- val opExec = OperatorExecution()
- val w = workerId("w-1")
- opExec.initWorkerExecution(w)
assertThrows[AssertionError] {
opExec.initWorkerExecution(w)
}
+ assert(opExec.getWorkerExecution(w) eq firstExec)
}
"OperatorExecution.getWorkerIds" should "be empty on a freshly constructed
operator" in {