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-7996-b25375bebf6cb44378aa07e5f630eba6b87f3199 in repository https://gitbox.apache.org/repos/asf/texera.git
commit be5144325e9f61b46387faf9ebcf77090e64185e Author: anthonychengit <[email protected]> AuthorDate: Sun Aug 30 03:39:01 2026 +0000 fix(amber): enable checkpoint serialization (#7996) ### What changes were proposed in this PR? Make <code>CheckpointState</code> implement <code>Serializable</code>, so the existing Kryo binding can persist worker checkpoints. Before: finalize checkpoint → serialize <code>CheckpointState</code> → no matching binding → <code>NotSerializableException</code> After: finalize checkpoint → Kryo binding → flush/close → restore the original checkpoint payload The regression test uses a real worker write path, stores a non-empty Unicode payload, reads the record back from sequential storage, and asserts the returned size and restored value. Two existing write-branch tests now await successful completion instead of discarding failures. ### Any related issues, documentation, discussions? Closes #7907 ### How was this PR tested? Added positive persistence/round-trip coverage while retaining the existing estimate-only and checkpoint-isolation cases. ~~~powershell $env:STORAGE_JDBC_URL='jdbc:postgresql://localhost:15432/texera_codex_jooq?currentSchema=texera_db,public' $env:STORAGE_JDBC_USERNAME='postgres' $env:STORAGE_JDBC_PASSWORD='' sbt "WorkflowExecutionService/testOnly org.apache.texera.amber.engine.architecture.worker.promisehandlers.FinalizeCheckpointHandlerSpec" sbt "scalafixAll --check" "scalafmtCheckAll" ~~~ Result: 6 tests passed; Scalafix and Scalafmt checks passed. The focused suite used an isolated embedded PostgreSQL database and port. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: OpenAI Codex (GPT-5) --- .../amber/engine/common/CheckpointState.scala | 2 +- .../FinalizeCheckpointHandlerSpec.scala | 39 ++++++++++++++-------- 2 files changed, 27 insertions(+), 14 deletions(-) diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/common/CheckpointState.scala b/amber/src/main/scala/org/apache/texera/amber/engine/common/CheckpointState.scala index 0f1159a2ef..0da3eaba6f 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/common/CheckpointState.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/common/CheckpointState.scala @@ -21,7 +21,7 @@ package org.apache.texera.amber.engine.common import scala.collection.mutable -class CheckpointState { +class CheckpointState extends Serializable { private val states = new mutable.HashMap[String, SerializedState]() diff --git a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/promisehandlers/FinalizeCheckpointHandlerSpec.scala b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/promisehandlers/FinalizeCheckpointHandlerSpec.scala index 0efa8816b3..81cd6919d9 100644 --- a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/promisehandlers/FinalizeCheckpointHandlerSpec.scala +++ b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/promisehandlers/FinalizeCheckpointHandlerSpec.scala @@ -64,7 +64,6 @@ import java.net.URI import java.util.concurrent.LinkedBlockingQueue import scala.collection.mutable import scala.collection.mutable.ArrayBuffer -import scala.util.Try /** * `finalizeCheckpoint` is the second half of a worker's checkpoint. It has two disjoint jobs, @@ -82,13 +81,6 @@ import scala.util.Try * one that can. * - saving another checkpoint's recorded messages, or leaving this checkpoint's recording in * place so the worker keeps buffering input forever. - * - * The write itself is not asserted: `SequentialRecordWriter` serializes through - * `AmberRuntime.serde`, and `CheckpointState` is not `java.io.Serializable`, so the shipped Pekko - * config (kryo bound to `java.io.Serializable`, java serialization off) has no binding for it. The - * write-branch case below therefore asserts only the state the handler mutates before writing, and - * does not depend on whether the call as a whole succeeds. - * * The write branch hands a closure to the worker's main thread and blocks until it runs, so that * case uses a real `WorkflowWorker` behind a `TestActorRef` (synchronous dispatch, as in * `WorkflowWorkerSpec`) and the worker's own `DataProcessor`. The estimate branch never reaches the @@ -195,6 +187,30 @@ class FinalizeCheckpointHandlerSpec assert(!storageAt(parent).containsFolder("checkpoint-folder")) } + it should "persist and restore a non-empty worker checkpoint" in { + val destination = "ram:///finalize-persist/" + val worker = liveWorker() + val dp = worker.underlyingActor.dp + val checkpoint = new CheckpointState() + checkpoint.save("unicode-payload", "checkpoint-δ") + dp.ecmManager.checkpoints(checkpointId) = checkpoint + val handler = new DataProcessorRPCHandlerInitializer(dp) + + val response = await( + handler.finalizeCheckpoint( + FinalizeCheckpointRequest(checkpointId, destination), + rpcContext + ) + ) + + val restored = SequentialRecordStorage + .fetchAllRecords(storageAt(destination), workerId.name.replace("Worker:", "")) + .toList + assert(response.size > 0L) + assert(restored.size == 1) + assert(restored.head.load[String]("unicode-payload") == "checkpoint-δ") + } + it should "fold this checkpoint's recorded messages in and stop only its recording" in { val worker = liveWorker() val dp = worker.underlyingActor.dp @@ -208,9 +224,7 @@ class FinalizeCheckpointHandlerSpec worker.underlyingActor.recordedInputs(unrelatedCheckpointId) = ArrayBuffer(recordedMessage(99)) val handler = new DataProcessorRPCHandlerInitializer(dp) - // The outcome is intentionally ignored: everything asserted below happens on the worker's main - // thread, before the storage write this call ends with (see the note in the class comment). - Try( + await( handler.finalizeCheckpoint( FinalizeCheckpointRequest(checkpointId, "ram:///finalize-fold/"), rpcContext @@ -234,8 +248,7 @@ class FinalizeCheckpointHandlerSpec dp.ecmManager.checkpoints(checkpointId) = checkpoint val handler = new DataProcessorRPCHandlerInitializer(dp) - // Outcome ignored for the same reason as the case above. - Try( + await( handler.finalizeCheckpoint( FinalizeCheckpointRequest(checkpointId, "ram:///finalize-fold-empty/"), rpcContext
