wellkilo opened a new pull request, #1177:
URL: https://github.com/apache/flink-agents/pull/1177

   <!-- Please add the relevant components in the PR title. -->
   
   Linked issue: #1173
   
   ### Purpose of change
   
   A caller no longer receives a successful return from 
`KafkaActionStateStore.put(...)` when Kafka
   later rejects the state record. The asynchronous broker failure now fails 
the task before the
   attempt-local cache changes, preventing a checkpoint from treating an 
action-state update as
   durable when its Kafka record was never acknowledged.
   
   Runtime flow:
   
   1. `put(...)` encodes the composite state key and sends the `ProducerRecord`.
   2. It waits on the `Future<RecordMetadata>` returned by that specific 
`send()` call.
   3. After the acknowledgement succeeds, it flushes the producer and updates 
the attempt-local
      cache.
   4. Only then does `put(...)` return successfully.
   
   #### Key decisions
   
   - Wait on the send future instead of relying on `flush()`. Kafka's 
`MockProducer` demonstrates
     that `flush()` can return normally while the per-record future retains an 
asynchronous failure.
   - Keep the existing synchronous persistence contract: successful return 
still means both the
     acknowledgement and flush completed.
   - Re-throw `InterruptedException` with the thread's interrupt flag restored, 
rather than wrapping
     it as a generic Kafka persistence failure.
   - Do not add a new acknowledgement timeout or configuration. This patch only 
makes the existing
     wait observable; it does not change the configured producer behavior.
   
   #### Interaction decisions
   
   | Send / acknowledgement / flush outcome | Result |
   | --- | --- |
   | `send()` throws synchronously | `RuntimeException`; cache unchanged |
   | Send future completes exceptionally | `RuntimeException` with the broker 
failure as root cause; flush and cache update are skipped |
   | Acknowledgement succeeds, `flush()` fails | `RuntimeException`; cache 
update is skipped |
   | Acknowledgement wait is interrupted | `InterruptedException`; interrupt 
flag restored; flush and cache update are skipped |
   | Acknowledgement and flush succeed | Cache is updated, then `put(...)` 
returns normally |
   
   #### Behavioral contracts
   
   - A successful `put(...)` means the specific record was acknowledged, the 
producer was flushed,
     and the attempt-local cache contains the state.
   - An asynchronous acknowledgement failure is propagated to the caller and 
does not populate the
     attempt-local cache.
   - An interrupted acknowledgement wait preserves 
`Thread.currentThread().isInterrupted()` and
     throws `InterruptedException`.
   - The method does not update the cache on send, acknowledgement, or flush 
failure.
   
   #### Failure behavior
   
   Kafka send and flush exceptions are surfaced as
   `RuntimeException("Failed to send action state to Kafka", cause)`. 
Interruption is not wrapped:
   it is re-thrown with the interrupt status preserved. A null producer retains 
the pre-existing
   behavior of logging and returning.
   
   ### Tests
   
   | Behavioral contract | Evidence |
   | --- | --- |
   | Asynchronous broker failure fails `put(...)` and leaves the cache 
unchanged | `testPutPropagatesAsyncBrokerFailureBeforeCachingState` |
   | Interruption is propagated, the interrupt flag is restored, and the cache 
is unchanged | `testPutPreservesInterruptWhileWaitingForBrokerAcknowledgement` |
   | The successful send, flush, and cache path remains intact | 
`testPutActionState` and the existing `KafkaActionStateStoreTest` suite |
   
   Local verification completed on the branch: `KafkaActionStateStoreTest` 
(51/51), repository
   lint/Spotless, and Apache RAT all pass.
   
   Not verified: live-broker failover/recovery and a never-completing broker 
acknowledgement. The
   regression uses `MockProducer` to make the record-level failure 
deterministic; no acknowledgement
   timeout is introduced by this change.
   
   <details>
   <summary>Implementation invariants and supporting evidence</summary>
   
   - `actionStates.put(...)` is reached only after `future.get()` and 
`producer.flush()` complete
     successfully.
   - The `InterruptedException` catch restores the flag before re-throwing the 
same exception.
   - The existing `ActionStateStore.put(...)` signature already declares 
`throws Exception`, so no
     public interface change is required.
   - The Kafka record schema and the attempt-local cache key remain unchanged.
   
   </details>
   
   ### API
   
   No public API signature, event schema, checkpoint format, or configuration 
change is introduced.
   The caller-visible behavior change is limited to propagating failures that 
were previously hidden
   by ignoring the send future.
   
   ### Documentation
   
   - [ ] `doc-needed`
   - [ ] `doc-not-needed`
   - [x] `doc-included`
   
   ### Was this patch authored or co-authored using generative AI tooling?
   
   - [x] Yes
   - [ ] No
   
   Generated-by: TraeCode (GPT-5)
   


-- 
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