1fanwang commented on code in PR #1161:
URL: https://github.com/apache/flink-agents/pull/1161#discussion_r4126049883


##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -161,9 +166,11 @@ public void put(Object key, long seqNum, Action action, 
Event event, ActionState
         try {
             ProducerRecord<String, ActionState> kafkaRecord =
                     new ProducerRecord<>(topic, stateKey, state);
-            producer.send(kafkaRecord);
+            RecordMetadata metadata = producer.send(kafkaRecord).get();

Review Comment:
   Done in `f109972f`. `put()` now catches `InterruptedException` before the 
generic catch, restores the flag with `Thread.currentThread().interrupt()`, and 
rethrows the original exception instead of wrapping it. 
`testPutRestoresInterruptFlagWhenSendIsInterrupted` covers it: the send fails 
like a cancelled `FutureTask` (clears the flag, throws), and the test asserts 
the flag is set again after `put()` returns.



##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -299,6 +310,11 @@ public void setOwnershipFilter(IntPredicate 
ownershipFilter) {
         this.ownershipFilter = ownershipFilter;
     }
 
+    @Override
+    public void markCheckpointedSequence(Object key, long seqNum) {
+        latestKeySeqNum.merge(keyEncoder.generateBusinessKeyIdentity(key), 
seqNum, Math::max);

Review Comment:
   Done in `f109972f`. `pruneState()` now drops the key from `latestKeySeqNum` 
once the identity has no live state left in the cache, so entries stop 
accumulating per key for the job lifetime. Losing the boundary is harmless: any 
later record for that identity is newer than the last completed sequence and 
must be replayed either way. 
`testPruneDropsCheckpointedBoundaryWhenNoLiveStateRemains` covers it.



##########
runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStoreTest.java:
##########
@@ -233,6 +233,87 @@ void testRecoveryMarker() throws Exception {
         assertThat((Map<Integer, Long>) secondMarker).containsEntry(1, 3L);
     }
 
+    @Test
+    void testRebuildStateKeepsPendingResultBeforeRecoveryMarker() throws 
Exception {
+        ActionState pendingState = new ActionState(testEvent);
+        actionStateStore.put(TEST_KEY, 1L, testAction, testEvent, 
pendingState);

Review Comment:
   Replaced by `KafkaActionStateStoreRecoveryTest` in `b07f9bd5`, which now 
executes a real durable call and takes a real checkpoint in the window the 
issue describes:
   
   - The action runs `context.durableExecute(...)` with a real supplier, then 
snapshots the operator from inside the action body while the durable result is 
still pending, then fails without completing.
   - Recovery takes the checkpoint into a fresh `KafkaActionStateStore` via 
`initializeState`, so it exercises the real `handleRecovery` -> `rebuildState` 
path against the checkpoint's marker.
   - The assertions check the pending result survived: the recovered 
`ActionState` is incomplete, carries the one successful `CallResult`, and the 
supplier ran exactly once.
   
   A constraint: under `mvn test` the operator always runs the synchronous 
executor, because the JDK 21 continuation executor lives only in the packaged 
multi-release JAR and is never on the surefire classpath, so a public-API 
checkpoint only observes post-invocation state. The checkpoint is therefore 
taken from inside the action body to land in the pre-completion window. I kept 
the original unit test because it still exercises the same marker rewind and 
rebuild path at unit scope, cheaply.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to