test(amber): add specs for the checkpoint and ECM promise handlers - #7118
Open
aglinxinyuan wants to merge 1 commit into
Open
test(amber): add specs for the checkpoint and ECM promise handlers#7118aglinxinyuan wants to merge 1 commit into
aglinxinyuan wants to merge 1 commit into
Conversation
These five handlers were the engine's lowest-coverage cluster (4.5-8.3%). All five specs follow the existing precedent in this layer (RetrieveWorkflowStateHandlerSpec): a real CoordinatorProcessor plus a second handler initializer, capturing what the output handler emits, so the tests drive the real handler methods rather than a re-implementation. - EmbeddedControlMessageHandlerSpec / EvaluatePythonExpressionHandlerSpec: dispatch shape, per-worker fan-out, channel scoping, the "latest execution" worker selection, and the fan-in that must not answer early. - TakeGlobalCheckpointHandlerSpec: the single-physical-op guard, the already-taken short-circuit that downgrades to an estimate, the per-checkpoint subfolder the restore path later resolves, and the propagation request's shape. Uses `ram://` in-memory VFS storage, per the LoggingSpec / DPThreadSpec precedent, with a distinct folder per case. - PrepareCheckpointHandlerSpec / FinalizeCheckpointHandlerSpec: the estimation branches, DP-state serialization under the requested id, the CheckpointSupport hand-off (and the skip branch's observable consequence), and the recorded-message fold. Assertion strength was checked by mutation across four rounds: 21 of 22 non-trivial cases were shown to fail on a realistic production change, then every mutation was reverted. No production file is touched. Also removes the two commented-out test blocks in CheckpointSpec, which are dead beyond repair. The first matches on OpExecInitInfoWithCode / OpExecInitInfoWithFunc -- names that now appear nowhere else in the repo, since OpExecInitInfo became a ScalaPB sealed oneof -- and casts to CheckpointSupport, which no production operator implements; it also has zero assertions and one branch is literally `???`. Its intent is already covered by SerializationManagerSpec and CheckpointSubsystemSpec. The second calls AmberClient with six constructor arguments where the current signature takes five, and its body starts a real two-operator workflow with 30s awaits, so it belongs in integration scope if anywhere.
Contributor
Automated Reviewer SuggestionsBased on the
|
Contributor
|
| config | throughput | MB/s | latency | max Δ latest / 7d | |
|---|---|---|---|---|---|
| ⚪ | bs=10 sw=10 sl=64 | 395 | 0.241 | 23,729/36,084/36,084 us | ⚪ within ±5% / 🔴 +128.8% |
| 🔴 | bs=100 sw=10 sl=64 | 887 | 0.541 | 107,245/157,049/157,049 us | 🔴 +20.8% / 🔴 +46.3% |
| 🔴 | bs=1000 sw=10 sl=64 | 1,096 | 0.669 | 908,807/986,619/986,619 us | 🔴 +6.6% / 🟢 -7.3% |
Baseline details
Latest main 301e3a8 from same runner
| config | metric | PR | latest main | 7d avg | Δ latest | Δ 7d |
|---|---|---|---|---|---|---|
| bs=10 sw=10 sl=64 | throughput | 395 tuples/sec | 408 tuples/sec | 786.12 tuples/sec | -3.2% | -49.8% |
| bs=10 sw=10 sl=64 | MB/s | 0.241 MB/s | 0.249 MB/s | 0.48 MB/s | -3.2% | -49.8% |
| bs=10 sw=10 sl=64 | p50 | 23,729 us | 24,086 us | 12,305 us | -1.5% | +92.8% |
| bs=10 sw=10 sl=64 | p95 | 36,084 us | 37,058 us | 15,774 us | -2.6% | +128.8% |
| bs=10 sw=10 sl=64 | p99 | 36,084 us | 37,058 us | 18,978 us | -2.6% | +90.1% |
| bs=100 sw=10 sl=64 | throughput | 887 tuples/sec | 960 tuples/sec | 999.71 tuples/sec | -7.6% | -11.3% |
| bs=100 sw=10 sl=64 | MB/s | 0.541 MB/s | 0.586 MB/s | 0.61 MB/s | -7.7% | -11.3% |
| bs=100 sw=10 sl=64 | p50 | 107,245 us | 102,487 us | 100,616 us | +4.6% | +6.6% |
| bs=100 sw=10 sl=64 | p95 | 157,049 us | 130,009 us | 107,356 us | +20.8% | +46.3% |
| bs=100 sw=10 sl=64 | p99 | 157,049 us | 130,009 us | 113,255 us | +20.8% | +38.7% |
| bs=1000 sw=10 sl=64 | throughput | 1,096 tuples/sec | 1,122 tuples/sec | 1,031 tuples/sec | -2.3% | +6.3% |
| bs=1000 sw=10 sl=64 | MB/s | 0.669 MB/s | 0.685 MB/s | 0.63 MB/s | -2.3% | +6.3% |
| bs=1000 sw=10 sl=64 | p50 | 908,807 us | 896,211 us | 980,328 us | +1.4% | -7.3% |
| bs=1000 sw=10 sl=64 | p95 | 986,619 us | 925,811 us | 1,027,528 us | +6.6% | -4.0% |
| bs=1000 sw=10 sl=64 | p99 | 986,619 us | 925,811 us | 1,054,298 us | +6.6% | -6.4% |
Raw CSV
config_idx,batch_size,schema_width,string_len,num_batches,total_ms,total_tuples,total_bytes,tuples_per_sec,mb_per_sec,lat_p50_us,lat_p95_us,lat_p99_us
0,10,10,64,20,506.57,200,128000,395,0.241,23729.07,36084.21,36084.21
1,100,10,64,20,2254.66,2000,1280000,887,0.541,107244.78,157049.18,157049.18
2,1000,10,64,20,18253.57,20000,12800000,1096,0.669,908807.14,986618.66,986618.66
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #7118 +/- ##
============================================
+ Coverage 79.57% 79.71% +0.14%
- Complexity 3837 3849 +12
============================================
Files 1159 1159
Lines 46145 46145
Branches 5127 5127
============================================
+ Hits 36718 36784 +66
+ Misses 7799 7721 -78
- Partials 1628 1640 +12
*This pull request uses carry forward flags. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What changes were proposed in this PR?
The checkpoint / ECM promise handlers were the engine's lowest-coverage cluster:
TakeGlobalCheckpointHandlerTakeGlobalCheckpointHandlerSpec(6 tests)FinalizeCheckpointHandlerFinalizeCheckpointHandlerSpec(5)PrepareCheckpointHandlerPrepareCheckpointHandlerSpec(4)EmbeddedControlMessageHandlerEmbeddedControlMessageHandlerSpec(6)EvaluatePythonExpressionHandlerEvaluatePythonExpressionHandlerSpec(6)All five follow the existing precedent in this layer (
RetrieveWorkflowStateHandlerSpec): a realCoordinatorProcessorplus a second handler initializer, capturing what the output handler emits — so the tests drive the real handler methods, not a re-implementation. Storage usesram://in-memory VFS per theLoggingSpec/DPThreadSpecprecedent, with a distinct folder per case sinceVFS.getManageris a JVM-wide singleton.Covered: the single-physical-op guard; the already-taken short-circuit that downgrades a checkpoint to an estimate; the per-checkpoint subfolder that
WorkflowActor.setupReplaylater resolves; channel scoping and the "latest execution" worker selection; the estimation branches on both worker handlers; theCheckpointSupporthand-off and the observable consequence of its skip branch; and the fan-in that must not answer before every targeted worker has replied.Assertion strength was measured, not asserted. Four rounds of production mutation, 21 of 22 non-trivial cases shown to go red, every mutation then reverted (
git diff -- amber/src/mainis empty):physicalOps.size != 1→> 2msg.targetOps.flatMap→msg.scope.flatMapif (!msg.estimationOnly)→ invertedfindLast→find(latest execution)&&→ `markAsReplayDestinationhoisted beforecollectcheckpoints.contains→!containsAlso removes the two commented-out blocks in
CheckpointSpec(−73 lines), verified dead rather than assumed: the first matchesOpExecInitInfoWithCode/OpExecInitInfoWithFunc, names that a repo-wide grep finds only inside those comment lines (OpExecInitInfois now a ScalaPB sealed oneof), and casts toCheckpointSupport, which no production operator implements — plus it has zero assertions and one branch is literally???. Its intent is already covered bySerializationManagerSpecandCheckpointSubsystemSpec. The second callsAmberClientwith six constructor arguments where the signature takes five, and starts a real two-operator workflow with 30s awaits.Two production bugs surfaced while writing these; neither is fixed here, and neither is pinned as current behaviour.
TakeGlobalCheckpointHandlerunder-reportstotalSize. It accumulates into a capturedvarvia a per-future.onSuccess { resp => totalSize += resp.size }registered beforeFuture.collectwraps those futures. Acom.twitter.util.Promiseruns continuations in reverse registration order, so for the last-satisfied future thecollectcontinuation — and the.mapthat reads the counter — runs before that future'sonSuccesshas contributed. Observed directly: two workers reporting 4000 and 56 producedTakeGlobalCheckpointResponse(4000); a single worker reporting 7 produced0. Fix direction is to sum from the collectedSeqinside the.map. The spec case is named for the fan-in it does verify and carries a comment explaining why the total is deliberately not asserted.A
CheckpointStatecan never be written to storage with the shipped config.SequentialRecordWriter.writeRecorddoesAmberRuntime.serde.serialize(obj).get, butCheckpointStateis a plainclassthat does not extendjava.io.Serializable, whilecluster.confsetsallow-java-serialization = offwith"java.io.Serializable" = kryoas the only binding — so the write throwsNotSerializableException. That kills both persist paths, which means global checkpointing cannot produce a checkpoint on disk and the checkpoint-restore branch ofWorkflowActor.setupReplayhas nothing to read. Existing specs missed it because they only exercisechkpt.save/chkpt.load(SerializedStateis a case class). Consequence for this PR: the two write-branch cases wrap the call inTryand assert only the state the handler mutates beforehand — assertions that hold both before and after a fix, so they are not inverted change-detectors. The class scaladoc says so explicitly.A third, smaller hazard is recorded in the notes but not tested:
EvaluatePythonExpressionHandlerstatically types worker results asEvaluatedValue, which is not a member of theControlReturnsealed oneof. Erasure hides it on the JVM side and the Python side sets a stray attribute, so an emptyControlReturncomes back. That is why the spec asserts dispatch and fan-in shape but never fulfils a worker promise with a fakeEvaluatedValue— doing so would cement a type confusion.Any related issues, documentation, discussions?
Closes #7117
How was this PR tested?
Five new specs, 27 tests. Run together with the three adjacent checkpoint specs they must not duplicate or break — 49 tests, Java 17:
Isolation was checked too: the new specs were re-run interleaved with
LoggingSpec,WorkflowWorkerSpec,DPThreadSpecand the sibling promise-handler specs (16 suites, concurrent in one JVM), and then a second time in the same JVM to confirm theram://folders leave no residue. No new spec mutatesAmberRuntime's private_serde/_actorSystem.Test/scalafmtCheckandTest/scalafix --checkboth[success].Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Opus 5)