Skip to content

concurrent-api: preserve demand across overlapping subscription switches - #3677

Open
tuannx wants to merge 2 commits into
apple:mainfrom
tuannx:codex/sequential-subscription-upstream-phase1
Open

tuannx wants to merge 2 commits into
apple:mainfrom
tuannx:codex/sequential-subscription-upstream-phase1

Conversation

@tuannx

@tuannx tuannx commented Oct 6, 2026 •

Copy link
Copy Markdown

Motivation

Overlapping SequentialSubscription.switchTo calls can lose or duplicate outstanding demand while an old request callback unwinds. Related to #1672; this does not yet prove the reported pipeline timeout's cause.

Modifications

  • Publish replacement subscriptions through an atomic mailbox; serialize demand and cancellation with the existing tryAcquireLock/releaseLock drain pattern.
  • Record transferred demand before callbacks; drain pending work even when a callback throws.
  • Preserve ordinary replacement without cancellation and invalid-request forwarding; cover zero demand, overflow, pending cancellation, and callback failures.

The design follows the mailbox/drain approach of RxJava's SubscriptionArbiter, adapted to ServiceTalk's demand accounting and callback serialization.

Result

Two commits preserve local red→green evidence: the test-only commit observed 4 lost and 37 duplicated transfers in 200,000 iterations; the fix passes all 200,000 and the latch-controlled regression.

Validation: JDK 17 concurrent-api suite 3,319 tests, 0 failures, 9 skipped, quality passed; all 31 SequentialSubscription tests pass on JDK 8. Upstream CI is pending.

A separate JDI experiment reproduces lost demand deterministically; it is omitted because upstream has no existing JDI test infrastructure I could find. Retain the 200k loop until equivalent deterministic coverage permits reducing/replacing it. Existing precedents: 500 race iterations, 10k TTL-race iterations, and 100k stack-depth iterations, the largest explicit single-loop bound found, not a project limit.

tuannx added 2 commits October 6, 2026 00:56
**Phase 1: test only. Expected red; production fix follows review.**

#### Motivation

Overlapping switches can lose outstanding demand or request it twice. Lost demand is the hypothesis investigated in apple#1672.

#### Modifications

- Stress 200,000 resubscriptions after source termination; require exactly 10 demand in every iteration.
- Retain the latch-controlled duplicate-demand test and its serialized control.

#### Result

Local stress: **6 lost-demand cases (0), 10 duplicate-demand cases (20)**; the other 199,984 received 10. Counts vary with scheduling; this is not deterministic or an end-to-end proof of apple#1672.

The latch test also reproduces **expected 1, actual 2**. Production code is unchanged.

Validation: full concurrent-api suite has only these two regression failures; quality passes. A local serialized-order control passes all 200,000 iterations.
#### Motivation

A replacement arriving while a request callback unwinds can lose or duplicate
outstanding demand.

#### Modifications

Publish replacements through an atomic mailbox and use the existing drain-lock
helpers to serialize demand accounting and callbacks. Keep pending work visible
through reentry and callback failures. Add terminal-action and demand-boundary
coverage.

#### Result

The lost/duplicate-demand regressions pass, including 200,000 stress iterations.
The concurrent-api suite and quality pass on JDK 17; focused tests pass on JDK 8.
@tuannx tuannx changed the title concurrent-api: reproduce lost and duplicated demand on resubscribe concurrent-api: preserve demand across overlapping subscription switches Oct 6, 2026
@tuannx
tuannx marked this pull request as ready for review October 6, 2026 18:45

@bryce-anderson bryce-anderson left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This looks remarkably better. I just have a few refinements to suggest but I think this is great work.

Comment on lines 96 to 99
public void cancel() {
final Subscription currSubscription = subscription;
final long currSourceRequested = sourceRequestedUpdater.getAndSet(this, CANCELLED);
// To avoid concurrent invocation with the switch thread we defer to that thread to cancel.
if (currSourceRequested >= 0) {
currSubscription.cancel();
}
cancelled = true;
drain();
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I worry a little bit about this not being able to trigger cancel on a reentrant call. Eg, we're draining and that triggers a cancel call synchronously. Could we detect if this thread is the one draining, and if so, fire the cancel signal? In my minds eye, this can be done by setting a private Thread drainingThread field on this class that gets set when we successfully enter the drain() loop, and if cancel() fails to enter the drain loop (drain can return a boolean signaling whether we were denied by contention or not) we can check if this thread is the current drainingThread, and if so it's a reentrant call and safe to call subscription.cancel().

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1, I simple reproducer would be fromIterable(it).concat(empty()) with a subscriber that cancels in onNext, it keeps draining the iterator (passes on main). Note that drainingThread proposal would only fix the same-thread case. A cancel() from another thread still waits behind the drain owner. If request() blocks, e.g. fromBlockingIterable in hasNext(timeout), the cancel that would unblock it never reaches it. Could cancel() bypass the lock and call subscription.cancel() directly, like the old implementation?

Here are examples of tests that fail on this PR and pass on main:

  1. ConcatPublisherTest (end-to-end, re-entrant cancel)
@Test
void cancelFromOnNextStopsSynchronousFirstSource() {
    final AtomicInteger pulled = new AtomicInteger();
    // Bounded so a regression fails instead of hanging: an infinite source would never stop emitting.
    final Publisher<Integer> p = fromIterable(() -> new Iterator<Integer>() {
        @Override
        public boolean hasNext() {
            return pulled.get() < 1000;
        }

        @Override
        public Integer next() {
            return pulled.incrementAndGet();
        }
    }).concat(empty());
    final List<Integer> received = new ArrayList<>();
    toSource(p).subscribe(new Subscriber<Integer>() {
        @Nullable
        private Subscription subscription;

        @Override
        public void onSubscribe(final Subscription s) {
            subscription = s;
            s.request(Long.MAX_VALUE);
        }

        @Override
        public void onNext(@Nullable final Integer item) {
            received.add(item);
            assert subscription != null;
            subscription.cancel();
        }

        @Override
        public void onError(final Throwable t) {
        }

        @Override
        public void onComplete() {
        }
    });
    assertThat(received, contains(1));
    assertThat("Items pulled from the source after cancel", pulled.get(), is(1));
}
  1. SequentialSubscriptionTest (re-entrant and cross-thread cancel)
@Test
void reentrantCancelDuringRequestCancelsActiveSubscription() {
    final AtomicBoolean cancelled = new AtomicBoolean();
    doAnswer(invocation -> {
        cancelled.set(true);
        return null;
    }).when(s1).cancel();
    doAnswer(invocation -> {
        s.cancel();
        // A synchronous source, e.g. Publisher.fromIterable, stops its emission loop only when it observes cancel.
        assertThat("Cancel did not reach the active subscription", cancelled.get(), is(true));
        return null;
    }).when(s1).request(anyLong());
    s.request(MAX_VALUE);
}

@Test
void cancelWhileRequestBlocksCancelsActiveSubscription() throws Exception {
    final CountDownLatch requested = new CountDownLatch(1);
    final CountDownLatch cancelled = new CountDownLatch(1);
    doAnswer(invocation -> {
        cancelled.countDown();
        return null;
    }).when(s1).cancel();
    doAnswer(invocation -> {
        requested.countDown();
        // A blocking source, e.g. Publisher.fromBlockingIterable, is unblocked only by cancel.
        cancelled.await(DEFAULT_TIMEOUT_SECONDS, SECONDS);
        return null;
    }).when(s1).request(anyLong());
    final Future<?> requesting = executor.submit(() -> s.request(1));
    try {
        assertThat("The request was not reached", requested.await(DEFAULT_TIMEOUT_SECONDS, SECONDS), is(true));
        s.cancel();
        assertThat("Cancel did not reach the blocked subscription", cancelled.getCount(), is(0L));
    } finally {
        cancelled.countDown();
        requesting.get();
    }
}

@bryce-anderson bryce-anderson Oct 10, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think technically allowing cancel() and request(n) calls to happen concurrently from different threads is a RS spec violation, part 2.7. That said, it looks like we already break these rules on main. It feels like the right thing to do is fix the blocking iterable (and potentially others) to not block in those calls.

Unfortunately, that puts us deeper down a rabbit hole. 😞
What do you think @idelpivnitskiy, should we allow it now and try to fix it later after we can fix the BlockingIterable (and maybe others, we'd need to do an audit), or should we try to cull any blocking in request(n) calls first?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. To make incremental progress, would be nice to preserve the current behavior of main and then fix later after other changes

Comment on lines +109 to +121
try {
// A source may terminate without demand before the drain reaches it. Only the latest source needs demand,
// but a displaced subscription must still observe a pending terminal action.
if (previous != null) {
if (cancelled) {
previous.cancel();
} else {
// Make the subscription visible before restoring the state of sourceRequested. If the Subscription
// thread observes the sourceRequested change it will also observe the subscription change. The
// Subscription thread also uses sourceRequested to make sure there is no concurrent invocation of
// the switched Subscription.
subscription = next;
final long n = requested;
if (n < 0) {
previous.request(n);
}
}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm not sure what this is for: if we're getting a switchTo call it must only be when the previous source has completed so calls to request(n) and cancel should be unnecessary.

Comment on lines +132 to 162
final Subscription next = pendingSubscriptionUpdater.getAndSet(this, null);
if (cancelled) {
final Subscription current = subscription;
subscription = EMPTY_SUBSCRIPTION_NO_THROW;
try {
current.cancel();
} finally {
if (next != null) {
next.cancel();
}
}
} else {
if (next != null) {
subscription = next;
sourceRequested = sourceEmitted;
}
final long n = requested;
if (sourceRequested >= 0) {
if (n < 0) {
sourceRequested = n;
subscription.request(n);
} else {
final long delta = n - sourceRequested;
if (delta != 0) {
// Commit before the callback: reentrant request/switch calls only enqueue more work.
sourceRequested = n;
subscription.request(delta);
}
}
}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think this can be a little simpler - if we have a non-null next I think we can completely ignore the previous subscription because we should only get a switchTo call if the previous source completed.

Suggested change
final Subscription next = pendingSubscriptionUpdater.getAndSet(this, null);
if (cancelled) {
final Subscription current = subscription;
subscription = EMPTY_SUBSCRIPTION_NO_THROW;
try {
current.cancel();
} finally {
if (next != null) {
next.cancel();
}
}
} else {
if (next != null) {
subscription = next;
sourceRequested = sourceEmitted;
}
final long n = requested;
if (sourceRequested >= 0) {
if (n < 0) {
sourceRequested = n;
subscription.request(n);
} else {
final long delta = n - sourceRequested;
if (delta != 0) {
// Commit before the callback: reentrant request/switch calls only enqueue more work.
sourceRequested = n;
subscription.request(delta);
}
}
}
}
final Subscription next = pendingSubscriptionUpdater.getAndSet(this, null);
if (next != null) {
subscription = next;
sourceRequested = sourceEmitted;
}
if (cancelled) {
final Subscription current = subscription;
subscription = EMPTY_SUBSCRIPTION_NO_THROW;
current.cancel();
} else {
final long n = requested;
if (sourceRequested >= 0) {
if (n < 0) {
sourceRequested = n;
subscription.request(n);
} else {
final long delta = n - sourceRequested;
if (delta != 0) {
// Commit before the callback: reentrant request/switch calls only enqueue more work.
sourceRequested = n;
subscription.request(delta);
}
}
}
}

Comment on lines 96 to 99
public void cancel() {
final Subscription currSubscription = subscription;
final long currSourceRequested = sourceRequestedUpdater.getAndSet(this, CANCELLED);
// To avoid concurrent invocation with the switch thread we defer to that thread to cancel.
if (currSourceRequested >= 0) {
currSubscription.cancel();
}
cancelled = true;
drain();
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1, I simple reproducer would be fromIterable(it).concat(empty()) with a subscriber that cancels in onNext, it keeps draining the iterator (passes on main). Note that drainingThread proposal would only fix the same-thread case. A cancel() from another thread still waits behind the drain owner. If request() blocks, e.g. fromBlockingIterable in hasNext(timeout), the cancel that would unblock it never reaches it. Could cancel() bypass the lock and call subscription.cancel() directly, like the old implementation?

Here are examples of tests that fail on this PR and pass on main:

  1. ConcatPublisherTest (end-to-end, re-entrant cancel)
@Test
void cancelFromOnNextStopsSynchronousFirstSource() {
    final AtomicInteger pulled = new AtomicInteger();
    // Bounded so a regression fails instead of hanging: an infinite source would never stop emitting.
    final Publisher<Integer> p = fromIterable(() -> new Iterator<Integer>() {
        @Override
        public boolean hasNext() {
            return pulled.get() < 1000;
        }

        @Override
        public Integer next() {
            return pulled.incrementAndGet();
        }
    }).concat(empty());
    final List<Integer> received = new ArrayList<>();
    toSource(p).subscribe(new Subscriber<Integer>() {
        @Nullable
        private Subscription subscription;

        @Override
        public void onSubscribe(final Subscription s) {
            subscription = s;
            s.request(Long.MAX_VALUE);
        }

        @Override
        public void onNext(@Nullable final Integer item) {
            received.add(item);
            assert subscription != null;
            subscription.cancel();
        }

        @Override
        public void onError(final Throwable t) {
        }

        @Override
        public void onComplete() {
        }
    });
    assertThat(received, contains(1));
    assertThat("Items pulled from the source after cancel", pulled.get(), is(1));
}
  1. SequentialSubscriptionTest (re-entrant and cross-thread cancel)
@Test
void reentrantCancelDuringRequestCancelsActiveSubscription() {
    final AtomicBoolean cancelled = new AtomicBoolean();
    doAnswer(invocation -> {
        cancelled.set(true);
        return null;
    }).when(s1).cancel();
    doAnswer(invocation -> {
        s.cancel();
        // A synchronous source, e.g. Publisher.fromIterable, stops its emission loop only when it observes cancel.
        assertThat("Cancel did not reach the active subscription", cancelled.get(), is(true));
        return null;
    }).when(s1).request(anyLong());
    s.request(MAX_VALUE);
}

@Test
void cancelWhileRequestBlocksCancelsActiveSubscription() throws Exception {
    final CountDownLatch requested = new CountDownLatch(1);
    final CountDownLatch cancelled = new CountDownLatch(1);
    doAnswer(invocation -> {
        cancelled.countDown();
        return null;
    }).when(s1).cancel();
    doAnswer(invocation -> {
        requested.countDown();
        // A blocking source, e.g. Publisher.fromBlockingIterable, is unblocked only by cancel.
        cancelled.await(DEFAULT_TIMEOUT_SECONDS, SECONDS);
        return null;
    }).when(s1).request(anyLong());
    final Future<?> requesting = executor.submit(() -> s.request(1));
    try {
        assertThat("The request was not reached", requested.await(DEFAULT_TIMEOUT_SECONDS, SECONDS), is(true));
        s.cancel();
        assertThat("Cancel did not reach the blocked subscription", cancelled.getCount(), is(0L));
    } finally {
        cancelled.countDown();
        requesting.get();
    }
}

sourceRequested = n;
subscription.request(n);
} else {
final long delta = n - sourceRequested;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Consider preserving the old assert delta >= 0 here. In case source over-emits, sourceEmitted > requested after a switch, and the next source gets request(<negative>)


@Test
@Timeout(60) // Deliberate stress loop; the latch does not force the handoff window.
void switchToWhileAnotherSwitchUnwindsRequestsOutstandingDemandFromNewSubscription() throws Exception {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are there any ways to use fewer iterations? It likely takes long time and may be flaky on CI. Maybe there are some opportunities to make it more deterministic latch-based test like switchToAfterReentrantRequestTransfersOutstandingDemandOnce?

boolean tryAcquire = true;
while (tryAcquire && tryAcquireLock(emittingUpdater, this)) {
try {
final Subscription next = pendingSubscriptionUpdater.getAndSet(this, null);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: because switches are rare, consider checking pendingSubscription != null before the getAndSet to skip an atomic write on every request(n)

@tuannx
tuannx force-pushed the codex/sequential-subscription-upstream-phase1 branch from 95a5399 to 506ce73 Compare October 10, 2026 08:08

This branch has not been deployed

No deployments
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.

3 participants