Repository navigation
Conversation
**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.
bryce-anderson
left a comment
There was a problem hiding this comment.
This looks remarkably better. I just have a few refinements to suggest but I think this is great work.
| 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(); | ||
| } |
There was a problem hiding this comment.
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().
There was a problem hiding this comment.
+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:
- 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));
}- 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();
}
}There was a problem hiding this comment.
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?
There was a problem hiding this comment.
Good point. To make incremental progress, would be nice to preserve the current behavior of main and then fix later after other changes
| 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); | ||
| } | ||
| } | ||
| } |
There was a problem hiding this comment.
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.
| 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); | ||
| } | ||
| } | ||
| } | ||
| } |
There was a problem hiding this comment.
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.
| 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); | |
| } | |
| } | |
| } | |
| } |
| 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(); | ||
| } |
There was a problem hiding this comment.
+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:
- 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));
}- 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; |
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
nit: because switches are rare, consider checking pendingSubscription != null before the getAndSet to skip an atomic write on every request(n)
95a5399 to
506ce73
Compare
Motivation
Overlapping
SequentialSubscription.switchTocalls 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
tryAcquireLock/releaseLockdrain pattern.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.