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

Reply via email to