Skip to content

feat(state): control-plane store SPI + deprecate zombie code (#5301 Sub-PR A) - #5310

Merged
qqeasonchen merged 4 commits into
apache:developfrom
qqeasonchen:feat/state-stores-spi
Aug 27, 2026
Merged

qqeasonchen merged 4 commits into
apache:developfrom
qqeasonchen:feat/state-stores-spi

Conversation

@qqeasonchen

Copy link
Copy Markdown
Contributor

Sub-PR A: Control-plane store SPI + deprecate zombie code (issue #5301)

This is the first of four sub-PRs that deliver issue #5301
("Unified state control-plane stores"). Sub-PR A adds the SPI surface, adapts
two existing implementations to it, and deprecates three zombie-code classes
that the new architecture will eventually replace.

It does not change runtime behavior: nothing currently calls the new
interfaces yet (call sites move over in Sub-PRs B and C).

What's in this PR

New public interfaces (in org.apache.eventmesh.runtime.state)

Interface Backed by (target) Purpose
SubscriptionStore Meta prefix-watch + local cache Cluster-shared subscription registry
SessionStore Meta prefix-watch + local cache Cluster-shared agent / binding / session registry
DeadLetterStore Meta CAS (idempotent) Durable ledger of dead-lettered deliveries
TaskStore Meta CAS + epoch-protected status transitions A2A task state (with TTL via expireStale)

OffsetStore and DeliveryStateStore are intentionally not added in this
PR — they are introduced in Sub-PRs B and C with the actual storage backends.

Adapted implementations

  • ClusterSubscriptionStore now implements SubscriptionStore. remove(...)
    returns the boolean result from Meta.delete(...) so callers can detect
    when a subscription was already gone (idempotent teardown).
  • SessionRegistry now implements SessionStore. register/markReady/ unregister are renamed to registerAgent/markAgentReady/unregisterAgent
    to match the interface contract (no callers existed in the current tree
    for the old names).

Zombie code marked @Deprecated(forRemoval = true)

Pending issue #5309
("PARTITION_OWNED_PULL delivery topology"):

Unit tests

SubscriptionStoreTest (5), SessionStoreTest (6), DeadLetterStoreTest
(2), TaskStoreTest (3) — total 16 tests, all in-process with hand-
rolled test doubles (no Nacos/RocksDB/Testcontainers needed). They pin the
SPI contract so subsequent sub-PRs can build on it.

Out of scope (handled in follow-up sub-PRs)

  • Sub-PR B — DeliveryStateStore (RocksDB) + OffsetStore (Meta async
    flush); wire DeliveryStateStore into the new consumer path.
  • Sub-PR C — DeadLetterStore Meta-backed implementation + TaskStore
    Meta-backed implementation; wire A2A TaskRegistry over TaskStore.
  • Sub-PR D — fault-injection Testcontainers E2E covering partition
    transfer, Meta outage, DLQ replay, and TaskStore epoch conflict.

Open questions for the maintainer

(Inherited from the #5301 issue body — please review there.)

  1. Should DeadLetterStore.recordDeadLetter carry a dlqOffset (current
    design), or just the dlqTopic so the offset is derived?
  2. Should TaskStore.TaskRecord carry full input / output payloads, or
    only a payloadRef that points to a separate blob store?
  3. For SubscriptionStore.put, do we want a separate putIfAbsent variant
    to make subscription reconciliation idempotent without overwriting?
  4. Should the SessionStore interface include a Watch / Subscription
    callback hook (similar to Nacos Listener) so callers can react to
    session-state changes without polling?

Testing

This PR was developed on the uni-runtime branch (develop at
7260581). The full eventmesh-runtime:compileJava plus the four new
test classes pass locally with JDK 21 + Gradle 8.7. Upstream CI will be
the source of truth — if anything in the SPI needs to change, the four
sub-PRs are independent so a fix here will not block B/C/D.

@qqeasonchen
qqeasonchen force-pushed the feat/state-stores-spi branch 2 times, most recently from 1625149 to 0b00f83 Compare August 26, 2026 07:49
 Sub-PR A)

This is the first of four sub-PRs that deliver issue apache#5301 (unified state
control-plane stores). Sub-PR A adds the SPI surface and deprecates
zombie code, without changing runtime behavior.

New public interfaces in org.apache.eventmesh.runtime.state:
  - SubscriptionStore   : cluster-shared subscription registry
                          (Meta prefix-watch with local cache)
  - SessionStore        : cluster-shared agent/binding/session registry
                          (Meta prefix-watch with local cache)
  - DeadLetterStore     : durable ledger of dead-lettered deliveries
                          (Meta CAS, idempotent recordDeadLetter)
  - TaskStore           : A2A task state with epoch-protected status
                          transitions (Meta CAS, TTL via expireStale)

Existing classes adapted to the new SPI:
  - ClusterSubscriptionStore now implements SubscriptionStore; remove()
    returns the boolean result from Meta.delete so callers can detect
    when a subscription was already gone.
  - SessionRegistry now implements SessionStore; register/markReady/
    unregister are renamed to registerAgent/markAgentReady/
    unregisterAgent to match the interface contract.

Zombie code marked @deprecated(forRemoval = true) pending issue apache#5309:
  - cluster/PartitionOwnership   (only callable from a closed delivery
                                   mode re-introduced by PARTITION_OWNED_PULL)
  - cluster/ClusterCoordinator   (superseded by the new Meta CAS path)
  - cluster/MetaBackedOffsetStore (superseded by the OffsetStore in
                                   Sub-PR B)

Unit tests cover the new interfaces and the adapted implementations
with in-process test doubles (no external dependencies):
  - SubscriptionStoreTest  (5 tests)
  - SessionStoreTest       (6 tests)
  - DeadLetterStoreTest    (2 tests)
  - TaskStoreTest          (3 tests)

Refs: apache#5301, apache#5309
@qqeasonchen
qqeasonchen force-pushed the feat/state-stores-spi branch from 0b00f83 to 6418c55 Compare August 26, 2026 08:02
@qqeasonchen
qqeasonchen merged commit 187bb92 into apache:develop Aug 27, 2026
7 checks passed
qqeasonchen added a commit to qqeasonchen/eventmesh that referenced this pull request Aug 31, 2026
…pache#5314)

CrossStoreFaultInjectionTest covers the fault modes that span two or more
stores (or the runtime + cluster-shared Meta). The individual store contract
tests (apache#5310/apache#5311/apache#5312/apache#5313) cannot observe these, because the invariant
lives at the seam:

  1. CrashMidAckReAck  - crash after offset-write, before MQ-ACK callback;
                        recovery retires without re-invoking the channel
                        (issue apache#5291 idempotency).
  2. MetaPartitionDuringDlq - Meta unreachable while dead-letter recording;
                        the store throws MetaPartitionException rather than
                        silently no-op'ing, so the dispatcher keeps the
                        delivery in flight and retries on heal (apache#5292).
  3. A2aCancelMidStream - cancel lands between PENDING and RUNNING; the
                        taskEpoch guard rejects the late transition and the
                        task converges on one terminal state (apache#5302).
  4. SubscriptionReRegisterAfterSplit - update during a Meta partition; after
                        heal the latest write wins, nothing dropped or
                        duplicated (apache#5288, apache#5301 SubscriptionStore).
  5. OffsetStoreRaceVsDeliveryStore - cross-thread offset-advance vs retire
                        race; the probe log proves every DELIVERY_REMOVE is
                        preceded by an OFFSET_WRITE at the same offset
                        (apache#5289 at-least-once).
  6. A2aDispatchRaceVsTaskStore - two dispatchers race on one task record;
                        stale-epoch writes are rejected and createTask yields
                        exactly one winner (apache#5291).

Every scenario runs in-process and deterministically. The JvmCrashHarness
from the previous commit remains the optional cross-JVM verification path.
qqeasonchen added a commit to qqeasonchen/eventmesh that referenced this pull request Aug 31, 2026
…pache#5314)

CrossStoreFaultInjectionTest covers the fault modes that span two or more
stores (or the runtime + cluster-shared Meta). The individual store contract
tests (apache#5310/apache#5311/apache#5312/apache#5313) cannot observe these, because the invariant
lives at the seam:

  1. CrashMidAckReAck  - crash after offset-write, before MQ-ACK callback;
                        recovery retires without re-invoking the channel
                        (issue apache#5291 idempotency).
  2. MetaPartitionDuringDlq - Meta unreachable while dead-letter recording;
                        the store throws MetaPartitionException rather than
                        silently no-op'ing, so the dispatcher keeps the
                        delivery in flight and retries on heal (apache#5292).
  3. A2aCancelMidStream - cancel lands between PENDING and RUNNING; the
                        taskEpoch guard rejects stale-epoch late transitions
                        and the task converges on one terminal state (apache#5302).
  4. SubscriptionReRegisterAfterSplit - update during a Meta partition; after
                        heal the latest write wins, nothing dropped or
                        duplicated (apache#5288, apache#5301 SubscriptionStore).
  5. OffsetStoreRaceVsDeliveryStore - cross-thread offset-advance vs retire
                        race; the probe log proves every DELIVERY_REMOVE is
                        preceded by an OFFSET_WRITE at the same offset
                        (apache#5289 at-least-once).
  6. A2aDispatchRaceVsTaskStore - two dispatchers race on one task record;
                        stale-epoch writes are rejected and createTask yields
                        exactly one winner (apache#5291).

Every scenario runs in-process and deterministically. The JvmCrashHarness
from the previous commit remains the optional cross-JVM verification path.

Note on the taskEpoch contract exercised by scenarios 3 and 6: the epoch is set
at createTask and never reset, so updateStatus rejects any epoch that differs
from the record's. Same-epoch writes are last-writer-wins by design - the
Runtime dispatcher is the sole writer and the epoch guards against a restarted
instance's stale handle, not against intra-JVM ordering.
qqeasonchen added a commit that referenced this pull request Aug 31, 2026
…ntrol plane (issue #5314) (#5318)

* test(state): add cross-store fault-injection harness (issue #5314)

Four in-process primitives shared by the six #5314 scenarios:

  MetaPartitionSwitch      - MetaStore wrapper; open()/close() simulates a network
                             partition (mutating ops throw MetaPartitionException,
                             reads continue against the pre-partition snapshot).
  CrossStoreRaceProbe      - ordered log of cross-store operations (DELIVERY_PUT/
                             REMOVE, OFFSET_WRITE/READ, TASK_UPDATE) with a monotonic
                             seq so a test can assert happens-before relationships.
  JvmCrashHarness          - child JVM + sentinel-file-driven destroyForcibly()
                             (SIGKILL / TerminateProcess), then relaunch against the
                             same on-disk stores. Gated on ENABLE_JVM_CRASH_HARNESS.
  InMemorySubscriptionStore- ConcurrentHashMap-backed SubscriptionStore for the
                             split-brain scenario (two views of the world).

All four are test-only and run fully in-process: no Nacos, no Docker, no
Testcontainers. See §13.2.12 of the architecture doc (added in a follow-up commit).

* test(state): add cross-store fault-injection 6-scenario test (issue #5314)

CrossStoreFaultInjectionTest covers the fault modes that span two or more
stores (or the runtime + cluster-shared Meta). The individual store contract
tests (#5310/#5311/#5312/#5313) cannot observe these, because the invariant
lives at the seam:

  1. CrashMidAckReAck  - crash after offset-write, before MQ-ACK callback;
                        recovery retires without re-invoking the channel
                        (issue #5291 idempotency).
  2. MetaPartitionDuringDlq - Meta unreachable while dead-letter recording;
                        the store throws MetaPartitionException rather than
                        silently no-op'ing, so the dispatcher keeps the
                        delivery in flight and retries on heal (#5292).
  3. A2aCancelMidStream - cancel lands between PENDING and RUNNING; the
                        taskEpoch guard rejects stale-epoch late transitions
                        and the task converges on one terminal state (#5302).
  4. SubscriptionReRegisterAfterSplit - update during a Meta partition; after
                        heal the latest write wins, nothing dropped or
                        duplicated (#5288, #5301 SubscriptionStore).
  5. OffsetStoreRaceVsDeliveryStore - cross-thread offset-advance vs retire
                        race; the probe log proves every DELIVERY_REMOVE is
                        preceded by an OFFSET_WRITE at the same offset
                        (#5289 at-least-once).
  6. A2aDispatchRaceVsTaskStore - two dispatchers race on one task record;
                        stale-epoch writes are rejected and createTask yields
                        exactly one winner (#5291).

Every scenario runs in-process and deterministically. The JvmCrashHarness
from the previous commit remains the optional cross-JVM verification path.

Note on the taskEpoch contract exercised by scenarios 3 and 6: the epoch is set
at createTask and never reset, so updateStatus rejects any epoch that differs
from the record's. Same-epoch writes are last-writer-wins by design - the
Runtime dispatcher is the sole writer and the epoch guards against a restarted
instance's stale handle, not against intra-JVM ordering.

* docs(architecture): add §13.2.12 cross-store fault-injection (issue #5314)

Documents the 6-scenario CrossStoreFaultInjectionTest: why single-store contract
tests cannot cover cross-store invariants, the four harness primitives, the
per-scenario assertion table, why Testcontainers is not used, and the two
extension points (RocksDB-backed crash scenario, MetaBackedOffsetStore going
active).

Numbered 13.2.12 because §13.2.11 was taken by the dual-topology matrix in
PR #5317.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Architecture Review][P0] Configurable DeliveryTopology for cluster delivery

1 participant