Skip to content

A quiet ReactorMongoSubscriptionModel subscription kept a position the oplog could drop - #1181

Draft
johanhaleby wants to merge 10 commits into
mainfrom
johan/reactor-quiet-subscription-position
Draft

johanhaleby wants to merge 10 commits into
mainfrom
johan/reactor-quiet-subscription-position

Conversation

@johanhaleby

@johanhaleby johanhaleby commented Oct 1, 2026 •

Copy link
Copy Markdown
Owner

Closes #1169

A ReactorMongoSubscriptionModel subscription that matched no events kept the position of the last event that did match, so a pause and a resume, a restart of the change stream or a lease handover after a long quiet period opened the change stream at a position the oplog had dropped. The reactive driver hands over the documents of a change stream and nothing for an empty batch, so the model never saw the newer position MongoDB sends with each batch.

A subscription with an id now reads the driver's own change stream cursor one batch at a time, and moves its position to the postBatchResumeToken of a reply that had no event for it. A listener can get that position through the new reactor QuietPositionReportingSubscriptions, the counterpart of the blocking capability from #1174.

I opened this as a draft. ReactorDurableSubscriptionModel doesn't save the quiet position yet, since #1179 changes that class and is still open. Until it does, a process restart after a long quiet period can still end in lost history for a reactor durable subscription. The list of what it needs is under Still open.

How the model gets the token

ChangeStreamPublisher has no method that returns the postBatchResumeToken. The model opens the change stream with the public BatchCursorPublisher.batchCursor(int), which is the method the driver calls itself when something subscribes to the publisher. It then reads the driver's change stream cursor, an AsyncChangeStreamBatchCursor, from the private field BatchCursor.wrapped, and calls getPostBatchResumeToken() on it. It also reads the private field AsyncChangeStreamBatchCursor.wrapped, the AtomicReference that holds the cursor the driver reads with, to tell when the driver is opening the change stream again. The model never writes either field. Every batch and the close go through the public BatchCursor.next() and BatchCursor.close(), so the commands, the server they go to, the session and the driver's resume after a failover are the driver's own, as with ReactiveMongoTemplate.changeStream(..).

The model checks the route in three places:

  • When the model is made, DriverChangeStreamCursor looks up BatchCursorPublisher.batchCursor(int), AsyncAggregateResponseBatchCursor.getPostBatchResumeToken(), and getters for BatchCursor.wrapped and AsyncChangeStreamBatchCursor.wrapped with MethodHandles.privateLookupIn(..). It also checks that BatchCursor.wrapped, AsyncChangeStreamBatchCursor.wrapped and CommandCursorResult.postBatchResumeToken are final. Any failure becomes the reason in one WARN, and the model reads every subscription through ReactiveMongoTemplate.changeStream(..) as on main.
  • For every cursor, the driver's change stream cursor has to be an AsyncChangeStreamBatchCursor, and the class of the cursor inside it has to declare a volatile commandCursorResult. If not, the model closes that cursor, logs the WARN once and reads through ReactiveMongoTemplate.changeStream(..) from then on.
  • For every look, a token read that fails ends the same way, unless the driver is opening the change stream again. The driver empties AsyncChangeStreamBatchCursor.wrapped for that and then puts a new cursor in, so a look that finds it empty, or that fails and then finds it empty or holding another cursor, gives no token. The check doesn't read the exception's message.

Either way the events delivered are the same, and only the quiet position is lost. Two tests in ReactorMongoSubscriptionModelQuietPositionTest, the_change_stream_cursor_of_the_driver_can_be_read_with_the_driver_of_this_build and a_subscription_that_has_reported_a_quiet_position_still_reads_through_the_cursor_of_the_driver, fail the build when the route is off on the driver version the build uses.

The driver jars have no module-info, only an Automatic-Module-Name, and an automatic module opens every package. I checked it with a named module on the module path on Temurin 21.0.12.1, where privateLookupIn(..) and setAccessible(true) both work on BatchCursor.wrapped, and privateLookupIn(..) works on AsyncChangeStreamBatchCursor.wrapped.

Why the model doesn't send aggregate and getMore itself anymore

The previous version of this PR read its own cursor with commands. On a sharded cluster with two mongos routers, its getMore commands went to whichever router the driver picked, and over 20 seconds 8 of 21 failed with CursorNotFound and it opened the change stream 8 times. The driver's change stream on the same cluster opened it once, and none of its 20 getMore commands failed.

ReactorMongoSubscriptionModelTwoMongosTest starts that cluster in one mongo:8.0 container, which took 4 to 5 seconds locally and 4.6 seconds on CI with both JDK 21 and JDK 25. The whole class took 31 and 28 seconds on CI. Against the previous version of this PR, a_subscription_opens_its_change_stream_once failed with [change streams opened] Expected size: 1 but was: 12, and a_subscription_sends_no_getMore_that_fails failed on the getMore commands that got CursorNotFound. The other two tests of the class, that a quiet position is reported and that a matching event is delivered, pass on the previous version too. All four pass on this version.

Invariants

  1. At most one piece of work per id. A new run for an id reads nothing until every earlier run for that id has closed and its work, an action's Mono or a listener's, has completed or been cancelled.
  2. Every reported or kept position is at or before the next event to deliver. The model asks for the next batch only after the action's Mono has completed for every event of the batch before, so no batch is fetched ahead. Within one call to next() the driver sends another getMore only after a reply with no document, and the reply that ends the call comes last. A token that a later look in the same call finds replaced therefore came with a reply that had no document. The model never moves to the token it finds at the latest look, since the driver stores a reply before it hands over its documents.
  3. Every topology where main works still works. The model uses the same cursor object as ReactiveMongoTemplate.changeStream(..), and only the batch size of the first batch and when next() is called differ.

No lock is held across the caller's code. Each wait is a Mono.

Invariant 2 also needs a later look to read a reply at least as new as an earlier look did, although the driver stores it on another thread. On driver 5.8.0 every field on the way is final or volatile. BatchCursor.wrapped is final, AsyncChangeStreamBatchCursor.wrapped is a final AtomicReference, AsyncCommandCursor.commandCursorResult is volatile and CommandCursorResult.postBatchResumeToken is final. On 5.5.2 the cursor inside is an AsyncCommandBatchCursor, whose commandCursorResult is volatile too. The looks of one wait come one after the other. Those modifiers are private to the driver, which is why the model checks them and falls back when one is missing.

Paths

Path What happens 1 2 3
Start subscribe(..) validates the filter and the start position, registers the subscription, and subscribes the run outside the monitor. The run waits for every earlier run for the id. waitUntilStarted() completes when the change stream is subscribed to, as on main. waits opens at the start position driver cursor
Wait for a batch The model looks at the token once a second and moves to a token a later look finds replaced, then hands it to the listeners inside a step. one step at a time only a replaced token
Batch with events Each event goes to the action, and the position moves past it once the action's Mono has completed. one step at a time moves after completion
Stop, pause, cancel, shutdown The run closes, which cancels the step under way and closes the cursor, so the driver sends killCursors. A move after the close counts only while a step that started before it is still under way. next run waits for the cancel the position stays before the event in hand
Resume A new run opens at the subscription's position, which is a quiet position when the last look found one. waits for every earlier run
Failover or network error The driver resumes the cursor by itself, as on main. the driver resumes from its own token driver resume
Any other read error The run reads again from the subscription's position after the backoff. the failed step has ended
Lost history (286) Restarts at now when restartSubscriptionsOnChangeStreamHistoryLost(true). Otherwise the run ends and the subscription is forgotten, as on main. The driver's cursor hands the MongoCommandException over unwrapped, so main's one-level check finds it.
Action error Retried with the backoff without opening the change stream again, as on main.
Quiet position handler error Retried forever like an action, a CheckpointWriteConditionNotFulfilledException included, by reading again from the subscription's position after the backoff rather than in place. That loses nothing, since the position only moves to a token a later look found replaced. only a replaced token
Route off, or a field neither final nor volatile One WARN, then main's read through ReactiveMongoTemplate.changeStream(..). position of the last event main's path
Token read fails Outside the driver's reopen, the model logs the WARN once and the subscription reads through ReactiveMongoTemplate.changeStream(..) from its position. During the reopen the look gives no token. the failed step has ended position of the last event main's path

Back to main's behaviour

The previous version changed more than the quiet position needed. These are main's behaviour again:

  • The driver resumes after a failover or a network error by itself.
  • The driver sends killCursors when the run closes the cursor.
  • waitUntilStarted() completes when the change stream is subscribed to.
  • The change stream follows the read preference of the database, as ReactiveMongoTemplate.changeStream(..) does, instead of always the primary.
  • The Flux of subscribe(filter, startAt) is main's code, with no quiet position. Nothing listens for it.
  • Lost history is found with main's one-level check.

One change stays. The model asks for the next batch only after the actions of the batch before have completed. Invariant 2 needs it, since a batch fetched ahead moves the token past events the action hasn't had. It costs a round trip per batch, and the changelog and the upgrade guide say so.

What changes for callers

  • The model asks MongoDB for the next batch of a subscription with an id only once the action's Mono has completed for every event of the batch before it.
  • A test that stubs changeStream(..) on a mocked ReactiveMongoOperations no longer reaches a subscription with an id, since the model opens it from getCollection(..). Two tests of ReactorMongoSubscriptionModelResilienceTest on main stub getCollection("events") too, and their assertions are unchanged.

The changelog has these under Changes and Breaking changes, and section 19 of the upgrade guide has a subsection for them. ADR 142 has a section for the reactor model.

ApplyFilterToChangeStreamOptionsBuilder.changeStreamPipeline(..) in the common module now builds the filter stages both SpringMongoSubscriptionModel and the reactor model put after $changeStream, which the Spring model built itself before. Its output for the Spring model is the same.

Tests

mvn -pl <modules> -am test with the tests of the modules this PR touches, on Temurin 21 against colima. The reactor module ran on 396d02d. The other rows are from b1e7f34, before the token read could fall back and before a refused quiet position write was retried, and I didn't run them again:

Module Tests Failures
subscription-mongodb-spring-reactor 132 0
subscription-mongodb-spring-reactor-checkpoint-storage 32 0
subscription-mongodb-spring-blocking 109 0
subscription-mongodb-spring-blocking-checkpoint-storage 40 0
subscription-mongodb-spring-blocking-competing-consumer-strategy 13 0
subscription-durable-reactor 146 0
subscription-catchup-reactor 26 0
subscription-stream-catchup-reactor 35 0
dcb-dsl-reactor 40 0
projection-dsl-reactor 117 0
mongodb-reactive-spring-boot-starter 119 0

New test classes in the reactor module:

  • ReactorMongoSubscriptionModelQuietPositionTest, 16 tests of the quiet position itself, pause and resume, the listener, the runs for one id, killCursors and the session of every getMore, the two route tests above, and a test that the driver of this build declares every field on the way to the token final or volatile.
  • ReactorMongoSubscriptionModelTokenWatchTest, 7 tests of the look rule.
  • ReactorMongoSubscriptionModelRunTest, 8 tests of when a run moves the position and runs a step.
  • ReactorMongoSubscriptionModelDriverCursorTest, 10 tests of lost history with and without a restart, a CheckpointWriteConditionNotFulfilledException from a quiet position handler once and every time, the fallback to ReactiveMongoTemplate.changeStream(..) with its single warning, and two token reads that fail. A ByteBuddy agent in the test sources makes getPostBatchResumeToken() on the driver's change stream cursor throw. In one test every read throws after the cursor opened, and the model warns once and delivers the next event. In the other the test empties AsyncChangeStreamBatchCursor.wrapped until the model has looked four times, so at least three token reads find it empty, and the model gives no warning and reports quiet positions after it.
  • ReactorMongoSubscriptionModelDriverResumeTest, 3 tests that the driver opens the change stream again after its connection is closed and the model doesn't restart the subscription.
  • ReactorMongoSubscriptionModelTwoMongosTest, 4 tests on the two-router cluster.

The only change to a test from main is the two getCollection("events") stubs in ReactorMongoSubscriptionModelResilienceTest. The previous version of this PR changed ReactorMongoSubscriptionLifecycleTest and ReactorMongoSubscriptionModelResilienceTest more, and both are main's again apart from those stubs. I deleted two tests the previous version had added, a_plain_subscription_that_restarts... and waiting_until_a_subscription_has_started..., since they asserted behaviour this version takes back to main's.

A few of the failures I got when I broke the code the new tests check:

  • With the field name changed to wrappedX, the route test failed with expected: null but was: "java.lang.NoSuchFieldException: no such field: …BatchCursor.wrappedX/…".
  • With the next batch read while the action runs, the prefetch test failed with [getMore sent while the action runs] Expected size: 1 but was: 9.
  • With a new run waiting only for the last run of its id, the chain test failed with [the third run started while the first run is being cancelled] Expecting value to be false but was true.
  • With the look rule comparing tokens with equals(..) instead of by instance, a TokenWatch test failed with Expecting actual: [null, null] to contain exactly …[null, {"_data"=BsonString{value='a'}}].
  • With the model of the previous version, the resume test failed with [the change stream opens at the resume token the driver holds] Expecting value to be true but was false.
  • With the modifiers test pointed at AsyncCommandCursor.batchSize, it failed with private int com.mongodb.internal.operation.AsyncCommandCursor.batchSize isn't volatile. Pointed at AsyncChangeStreamBatchCursor.resumeToken as a field that has to be final, it failed with private volatile org.bson.BsonDocument …AsyncChangeStreamBatchCursor.resumeToken isn't final.
  • On b1e7f34, the test whose token reads start failing failed with [warnings that the resume token can't be read] Expected size: 1 but was: 0.
  • On 887e9ec, the test whose quiet position handler refuses the write once failed with [events delivered] Expecting actual: [] to contain exactly (and in same order): ["64cdf14a-…"], and the test whose handler keeps refusing failed with [quiet positions handed over] Expecting size of: [MongoResumeTokenCheckpoint[…]] to be greater than 1 but was 1.
  • With every failed token read taken as the fallback, the reopen test failed with [warnings that the resume token can't be read] Expecting empty but was: [[WARN] … (java.lang.AssertionError)…].

Three tests pass on the previous version too, so they catch a later regression and prove nothing about this change: an_event_written_after_the_driver_opened_the_change_stream_again_is_delivered_exactly_once, and the two tests of ReactorMongoSubscriptionModelTwoMongosTest named above. The reopen test passes on b1e7f34 as well, since b1e7f34 never fell back on a failed token read. It fails only on the change that falls back on every failed read.

Still open

ReactorDurableSubscriptionModel, after #1179 is merged:

  • Add a quiet position listener to the wrapped model through QuietPositionReportingSubscriptions.findIn(..) when the model is made, and remove it on shutdown(). It needs an interval in ReactorDurableSubscriptionModelConfig, like saveQuietPositionEvery(Duration) and neverSaveQuietPosition() on the blocking config, and the starter needs to pass it on.
  • Keep per subscription when a checkpoint was last stored, whether a delivery is under way, and whether the persist predicate declined the most recent event, and save a quiet position only by the blocking rule from ADR 142. A cancel deletes the checkpoint, so a quiet save that is under way when the cancel comes must not write it back.
  • ReactorCatchupSubscriptionModel, ReactorDcbCatchupSubscriptionModel and ReactorStreamCatchupSubscriptionModel have to answer capability(QuietPositionReportingSubscriptions.class) with the model they wrap. On the reactor side capability(..) is a plain instanceof and does not look through a wrapper.
  • The reactor durable model saves an event's checkpoint with storage.save(..) and no write condition, so a quiet save that is late can replace a newer checkpoint. The blocking model has the same window for a store without a condition, and ADR 142 accepts it there, since the stored position then only moves back.
  • A reactor cancel returns a Mono that completes once the stored state is gone #1179 changes cancelSubscription(..) to return Mono<Void>, which the quiet save has to wait for or check.
  • Save the quiet position with the write condition the save for an event uses, any(), so the save can't be refused. The PositionWriter A reactor cancel returns a Mono that completes once the stored state is gone #1179 gives each generation of a subscription is what keeps a quiet save from following the delete in cancelSubscription(..), since a cancel retires it and a retired writer begins no write.

A CheckpointWriteConditionNotFulfilledException from a quiet position handler is retried like any other error from it, and one from an action is retried with the backoff indefinitely, as on main. The blocking models end delivery on that exception, since there the quiet save is conditional on the version of the lease. The reactor stack has no lease, so only a listener of your own can raise it.

…e oplog could drop

The reactive driver hands over the documents of a change stream and nothing
for an empty batch, so a subscription that matched no events kept the
position of the last one that did. A pause and a resume or a restart after a
long quiet period then opened at a position MongoDB no longer had.

The model now sends aggregate and getMore itself, on a session it opens for
each change stream, since MongoDB refuses a getMore on another session than
the cursor's. A reply with no document moves the position to its post batch
resume token, and the reactor QuietPositionReportingSubscriptions hands that
position to a listener.

A named subscription reads in runs. A new run for an id reads nothing until
the previous run has closed and its action or listener Mono has completed or
been cancelled, and the next batch is asked for only once the actions for the
batch before it have completed.
The changelog, the upgrade guide and ADR 142 now cover the reactor model.
They say it reads from the primary whatever the read preference, that
waitUntilStarted() waits for MongoDB to open the change stream, and that
ReactorDurableSubscriptionModel doesn't save the quiet position yet.
Sending aggregate and getMore from the model failed with CursorNotFound
behind two mongos routers. A subscription with an id now reads the
driver's change stream cursor one batch at a time, and takes the
postBatchResumeToken from the private wrapped field of BatchCursor. When
that field can't be read, the model logs a warning and reads through
ReactiveMongoTemplate.changeStream(..) as before.

waitUntilStarted(), the read preference, the driver's own resume, the
plain Flux and the lost-history check are back to what main does.
ADR 142, the changelog and the upgrade guide described the cursor the
model read itself, with its own session, always from the primary.
A token read that failed came back as no token, so a driver whose getter
kept failing left a subscription without a quiet position and logged
nothing. Now only a read while the driver opens the change stream again
gives no token, which the model tells apart by the cursor the driver
reads with. Any other failure makes the model read through Spring's
changeStream with the one warning.

The look at the token is only safe because every field it reads through
is final or volatile in the driver. The model now checks those modifiers
when it loads and for every cursor, and falls back the same way when one
is missing.
ADR 142 said the model read the token from a plain field and relied on
HotSpot for the order of two reads. On driver 5.8.0 every field on the
way is final or volatile, so a later look reads a reply at least as new
as an earlier one. The ADR, the changelog and the upgrade guide now say
so, and say that a missing modifier or a failing token read ends in the
warning too.
One test fails when the driver of this build declares a field on the
way to the token as neither final nor volatile. Two more break the
driver's token read with a ByteBuddy agent. A read that starts failing
after the cursor opened gives one warning and events through Spring's
change stream. A read while the driver's cursor reference is empty
gives no warning and quiet positions keep coming.
A CheckpointWriteConditionNotFulfilledException from a quiet position
handler ended delivery for the subscription, which still answered
isRunning(id) with true. Nothing in the library raises it from a
reactor quiet handler, since the reactor stack has no lease, so only a
listener of your own could stop a subscription that way. The model now
reads again from the subscription's position after its backoff, as for
any other error from that handler, and an action's refused write is
retried as on main.

The two tests that asserted the old outcome now assert that a handler
refusing once lets the subscription deliver the next event, and that
one that keeps refusing makes the model read again from the quiet
position. The test of the driver's reopen waits for four looks instead
of four seconds, so a slow machine still gets three token reads in.
The reactor QuietPositionListener javadoc and ADR 142 said a
CheckpointWriteConditionNotFulfilledException from a quiet position
handler ended delivery on the node. The model now retries it forever,
like an error from an action, by reading again from the subscription's
position. Both now say so, and the ADR says why the blocking models
still end delivery on it.
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.

ReactorMongoSubscriptionModel can fail with lost history when a subscription matches no events

1 participant