Tracking sub-issue of #5354. Part of the production-HA acceptance plan (Phase 1, P0).
What this PR changes
Every offset write, broker ACK, DLQ transition, and dispatch call carries the owner's FencingToken. A new StaleOwnerException is raised when the in-memory token is below the value in Meta; the dispatcher catches it and stops processing the partition. PartitionOwnership re-tries ownership acquisition on the next tick.
Concretely
- Add
FencingToken currentToken() to ReliableDispatcher, set at construction from the PartitionOwnership scheduler.
- Pass
currentToken to every OffsetStore.commit, StoragePlugin.ack, DeadLetterSink.dispatch, and AckCallback.onAck call.
- Before the write, the storage / offset / DLQ layer calls
metaStore.get("em/assignments/" + topic + "#" + partition) and compares the persisted token to the local one. If the persisted token is higher, throw StaleOwnerException (new class in cluster package).
ReliableDispatcher.dispatch catches StaleOwnerException, removes the partition from the local owned set, increments a new fencedPartitions metric, and re-subscribes via the PartitionOwnership reaper.
- Add
boundedDelayedAckQueue cleanup on instance shutdown: pending POP ACKs are flushed to the OffsetStore before JVM exit; a unit test asserts the flush runs on graceful shutdown but NOT on kill -9 (the next instance's startup must rebuild the queue from the broker, then from the local OffsetStore).
Acceptance criteria
ReliableDispatcher has a FencingToken constructor argument and forwards it on every storage/offset/DLQ write.
- A new
StaleOwnerIntegrationTest runs two Runtime instances against an in-memory MetaStore, has instance A claim partitions, then bumps the token in Meta to simulate takeover. Instance A's next dispatch throws StaleOwnerException and is excluded from ownedPartitions on the next tick.
OffsetStore.commit is unreachable from ReliableDispatcher without a FencingToken (ArchUnit rule mirrors the runtime contract).
boundedDelayedAckQueue flush test: start instance, enqueue 10 POP ACKs, send SIGTERM, restart, assert all 10 ACKs are still pending (not lost) AND not double-applied (the broker is queried, not local state).
Verification
./gradlew :eventmesh-runtime:test --tests "*StaleOwnerIntegrationTest*"
./gradlew :eventmesh-runtime:test --tests "*ReliableDispatcherTest*"
./gradlew :eventmesh-runtime:test --tests "*BoundedDelayedAckQueueTest*"
Depends on / blocks
References
eventmesh-runtime/.../cluster/FencingToken.java
eventmesh-runtime/.../cluster/PartitionOwnership.java
eventmesh-runtime/.../cluster/MetaStore.java
eventmesh-runtime/.../delivery/ReliableDispatcher.java
eventmesh-runtime/.../offset/OffsetStore.java
Part of the production-HA topology in #5354. See also #5352 and #5353.
Tracking sub-issue of #5354. Part of the production-HA acceptance plan (Phase 1, P0).
What this PR changes
Every offset write, broker ACK, DLQ transition, and dispatch call carries the owner's
FencingToken. A newStaleOwnerExceptionis raised when the in-memory token is below the value in Meta; the dispatcher catches it and stops processing the partition.PartitionOwnershipre-tries ownership acquisition on the next tick.Concretely
FencingToken currentToken()toReliableDispatcher, set at construction from thePartitionOwnershipscheduler.currentTokento everyOffsetStore.commit,StoragePlugin.ack,DeadLetterSink.dispatch, andAckCallback.onAckcall.metaStore.get("em/assignments/" + topic + "#" + partition)and compares the persisted token to the local one. If the persisted token is higher, throwStaleOwnerException(new class inclusterpackage).ReliableDispatcher.dispatchcatchesStaleOwnerException, removes the partition from the local owned set, increments a newfencedPartitionsmetric, and re-subscribes via thePartitionOwnershipreaper.boundedDelayedAckQueuecleanup on instance shutdown: pending POP ACKs are flushed to the OffsetStore before JVM exit; a unit test asserts the flush runs on graceful shutdown but NOT onkill -9(the next instance's startup must rebuild the queue from the broker, then from the local OffsetStore).Acceptance criteria
ReliableDispatcherhas aFencingTokenconstructor argument and forwards it on every storage/offset/DLQ write.StaleOwnerIntegrationTestruns two Runtime instances against an in-memoryMetaStore, has instance A claim partitions, then bumps the token in Meta to simulate takeover. Instance A's next dispatch throwsStaleOwnerExceptionand is excluded fromownedPartitionson the next tick.OffsetStore.commitis unreachable fromReliableDispatcherwithout aFencingToken(ArchUnit rule mirrors the runtime contract).boundedDelayedAckQueueflush test: start instance, enqueue 10 POP ACKs, send SIGTERM, restart, assert all 10 ACKs are still pending (not lost) AND not double-applied (the broker is queried, not local state).Verification
Depends on / blocks
References
eventmesh-runtime/.../cluster/FencingToken.javaeventmesh-runtime/.../cluster/PartitionOwnership.javaeventmesh-runtime/.../cluster/MetaStore.javaeventmesh-runtime/.../delivery/ReliableDispatcher.javaeventmesh-runtime/.../offset/OffsetStore.javaPart of the production-HA topology in #5354. See also #5352 and #5353.