From 70b295ca00035e532eccb36b357424086f9f7876 Mon Sep 17 00:00:00 2001 From: tianrking <10758833+tianrking@users.noreply.github.com> Date: Mon, 5 Oct 2026 06:21:27 +0800 Subject: [PATCH] Release cancelled replay consumers from their buffers Related #7912 --- .../operators/flowable/FlowableReplay.java | 22 +- .../observable/ObservableReplay.java | 12 +- .../rxjava4/processors/ReplayProcessor.java | 18 +- .../rxjava4/subjects/ReplaySubject.java | 11 +- .../rxjava4/core/ReplayCancellationTest.java | 219 ++++++++++++++++++ 5 files changed, 264 insertions(+), 18 deletions(-) create mode 100644 src/test/java/io/reactivex/rxjava4/core/ReplayCancellationTest.java diff --git a/src/main/java/io/reactivex/rxjava4/internal/operators/flowable/FlowableReplay.java b/src/main/java/io/reactivex/rxjava4/internal/operators/flowable/FlowableReplay.java index 74baae5b3a..c681fafe64 100644 --- a/src/main/java/io/reactivex/rxjava4/internal/operators/flowable/FlowableReplay.java +++ b/src/main/java/io/reactivex/rxjava4/internal/operators/flowable/FlowableReplay.java @@ -465,7 +465,7 @@ static final class InnerSubscription extends AtomicLong implements Subscripti * The parent subscriber-to-source used to allow removing the child in case of * child cancellation. */ - final ReplaySubscriber parent; + volatile ReplaySubscriber parent; /** The actual child subscriber. */ final Subscriber child; /** @@ -502,11 +502,14 @@ public void request(long n) { if (BackpressureHelper.addCancel(this, n) != CANCELLED) { // increment the total request counter BackpressureHelper.add(totalRequested, n); + ReplaySubscriber p = parent; // if successful, notify the parent dispatcher this child can receive more // elements - parent.manageRequests(); - // try replaying any cached content - parent.buffer.replay(this); + if (p != null) { + p.manageRequests(); + // try replaying any cached content + p.buffer.replay(this); + } } } } @@ -533,15 +536,20 @@ public void cancel() { @Override public void dispose() { if (getAndSet(CANCELLED) != CANCELLED) { + ReplaySubscriber p = parent; + parent = null; // remove this from the parent - parent.remove(this); + p.remove(this); // After removal, we might have unblocked the other child subscribers: // let's assume this child had 0 requested before the cancellation while // the others had non-zero. By removing this 'blocking' child, the others // are now free to receive events - parent.manageRequests(); + p.manageRequests(); // make sure the last known node is not retained - index = null; + synchronized (this) { + index = null; + missed = true; + } } } /** diff --git a/src/main/java/io/reactivex/rxjava4/internal/operators/observable/ObservableReplay.java b/src/main/java/io/reactivex/rxjava4/internal/operators/observable/ObservableReplay.java index 381cfeccaf..f6eafcc46f 100644 --- a/src/main/java/io/reactivex/rxjava4/internal/operators/observable/ObservableReplay.java +++ b/src/main/java/io/reactivex/rxjava4/internal/operators/observable/ObservableReplay.java @@ -423,7 +423,7 @@ static final class InnerDisposable * The parent subscriber-to-source used to allow removing the child in case of * child dispose() call. */ - final ReplayObserver parent; + volatile ReplayObserver parent; /** The actual child subscriber. */ final Observer child; /** @@ -448,10 +448,16 @@ public boolean isDisposed() { public void dispose() { if (!cancelled) { cancelled = true; + ReplayObserver p = parent; + parent = null; // remove this from the parent - parent.remove(this); + if (p != null) { + p.remove(this); + } // make sure the last known node is not retained - index = null; + if (getAndIncrement() == 0) { + index = null; + } } } /** diff --git a/src/main/java/io/reactivex/rxjava4/processors/ReplayProcessor.java b/src/main/java/io/reactivex/rxjava4/processors/ReplayProcessor.java index 3b088e3162..92ff0c5baf 100644 --- a/src/main/java/io/reactivex/rxjava4/processors/ReplayProcessor.java +++ b/src/main/java/io/reactivex/rxjava4/processors/ReplayProcessor.java @@ -618,7 +618,7 @@ static final class ReplaySubscription<@NonNull T> extends AtomicInteger implemen @Serial private static final long serialVersionUID = 466549804534799122L; final Subscriber downstream; - final ReplayProcessor state; + volatile ReplayProcessor state; Object index; @@ -637,16 +637,24 @@ static final class ReplaySubscription<@NonNull T> extends AtomicInteger implemen @Override public void request(long n) { if (SubscriptionHelper.validate(n)) { - BackpressureHelper.add(requested, n); - state.buffer.replay(this); + ReplayProcessor s = state; + if (s != null) { + BackpressureHelper.add(requested, n); + s.buffer.replay(this); + } } } @Override public void cancel() { - if (!cancelled) { + ReplayProcessor s = state; + if (s != null) { cancelled = true; - state.remove(this); + state = null; + s.remove(this); + if (getAndIncrement() == 0) { + index = null; + } } } } diff --git a/src/main/java/io/reactivex/rxjava4/subjects/ReplaySubject.java b/src/main/java/io/reactivex/rxjava4/subjects/ReplaySubject.java index 695d00ca19..cc8525ef96 100644 --- a/src/main/java/io/reactivex/rxjava4/subjects/ReplaySubject.java +++ b/src/main/java/io/reactivex/rxjava4/subjects/ReplaySubject.java @@ -621,7 +621,7 @@ static final class ReplayDisposable extends AtomicInteger implements Disposab @Serial private static final long serialVersionUID = 466549804534799122L; final Observer downstream; - final ReplaySubject state; + volatile ReplaySubject state; Object index; @@ -634,9 +634,14 @@ static final class ReplayDisposable extends AtomicInteger implements Disposab @Override public void dispose() { - if (!cancelled) { + ReplaySubject s = state; + if (s != null) { cancelled = true; - state.remove(this); + state = null; + s.remove(this); + if (getAndIncrement() == 0) { + index = null; + } } } diff --git a/src/test/java/io/reactivex/rxjava4/core/ReplayCancellationTest.java b/src/test/java/io/reactivex/rxjava4/core/ReplayCancellationTest.java new file mode 100644 index 0000000000..e1823fd9a0 --- /dev/null +++ b/src/test/java/io/reactivex/rxjava4/core/ReplayCancellationTest.java @@ -0,0 +1,219 @@ +/* + * Copyright (c) 2016-present, RxJava Contributors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in + * compliance with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is + * distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See + * the License for the specific language governing permissions and limitations under the License. + */ + +package io.reactivex.rxjava4.core; + +import static org.junit.jupiter.api.Assertions.*; + +import java.lang.ref.*; +import java.util.concurrent.Flow.Subscription; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.Test; + +import io.reactivex.rxjava4.disposables.Disposable; +import io.reactivex.rxjava4.processors.*; +import io.reactivex.rxjava4.schedulers.TestScheduler; +import io.reactivex.rxjava4.subjects.*; +import io.reactivex.rxjava4.testsupport.TestHelper; + +public class ReplayCancellationTest extends RxJavaTest { + + @Test + public void cancelledObservableDoesNotRetainBuffer() throws Exception { + assertAll( + () -> assertReclaimed(cancelObservable(0)), + () -> assertReclaimed(cancelObservable(1)), + () -> assertReclaimed(cancelObservable(2)) + ); + } + + @Test + public void cancelledFlowableDoesNotRetainBuffer() throws Exception { + assertAll( + () -> assertReclaimed(cancelFlowable(0)), + () -> assertReclaimed(cancelFlowable(1)), + () -> assertReclaimed(cancelFlowable(2)) + ); + } + + @Test + public void cancelledSubjectDoesNotRetainBuffer() throws Exception { + assertAll( + () -> assertReclaimed(cancelSubject(0)), + () -> assertReclaimed(cancelSubject(1)), + () -> assertReclaimed(cancelSubject(2)) + ); + } + + @Test + public void cancelledProcessorDoesNotRetainBuffer() throws Exception { + assertAll( + () -> assertReclaimed(cancelProcessor(0)), + () -> assertReclaimed(cancelProcessor(1)), + () -> assertReclaimed(cancelProcessor(2)) + ); + } + + @Test + public void flowableRequestCancelRace() { + for (int i = 0; i < TestHelper.RACE_DEFAULT_LOOPS; i++) { + PublishProcessor source = PublishProcessor.create(); + ConnectableFlowable replay = source.replay(1); + ReplayConsumer consumer = new ReplayConsumer(new Object()); + replay.subscribe(consumer); + Disposable connection = replay.connect(); + try { + TestHelper.race(() -> consumer.upstream.request(1), consumer.upstream::cancel); + source.onNext(1); + assertEquals(0, consumer.values); + } finally { + connection.dispose(); + } + } + } + + @Test + public void processorRequestCancelRace() { + for (int i = 0; i < TestHelper.RACE_DEFAULT_LOOPS; i++) { + ReplayProcessor source = ReplayProcessor.createWithSize(1); + ReplayConsumer consumer = new ReplayConsumer(new Object()); + source.subscribe(consumer); + TestHelper.race(() -> consumer.upstream.request(1), consumer.upstream::cancel); + source.onNext(1); + assertEquals(0, consumer.values); + assertFalse(source.hasSubscribers()); + } + } + + static ReplayConsumer cancelObservable(int kind) { + Object value = new Object(); + ReplayConsumer consumer = new ReplayConsumer(value); + PublishSubject source = PublishSubject.create(); + ConnectableObservable replay = switch (kind) { + case 0 -> source.replay(); + case 1 -> source.replay(1); + default -> source.replay(1, 1, TimeUnit.DAYS, new TestScheduler()); + }; + replay.subscribe(consumer); + replay.connect(); + source.onNext(value); + consumer.disposable.dispose(); + consumer.disposable.dispose(); + assertTrue(consumer.disposable.isDisposed()); + return consumer; + } + + static ReplayConsumer cancelFlowable(int kind) { + Object value = new Object(); + ReplayConsumer consumer = new ReplayConsumer(value); + PublishProcessor source = PublishProcessor.create(); + ConnectableFlowable replay = switch (kind) { + case 0 -> source.replay(); + case 1 -> source.replay(1); + default -> source.replay(1, 1, TimeUnit.DAYS, new TestScheduler()); + }; + replay.subscribe(consumer); + replay.connect(); + source.onNext(value); + consumer.upstream.cancel(); + consumer.upstream.cancel(); + consumer.upstream.request(1); + return consumer; + } + + static ReplayConsumer cancelSubject(int kind) { + Object value = new Object(); + ReplayConsumer consumer = new ReplayConsumer(value); + ReplaySubject source = switch (kind) { + case 0 -> ReplaySubject.create(); + case 1 -> ReplaySubject.createWithSize(1); + default -> ReplaySubject.createWithTimeAndSize(1, TimeUnit.DAYS, new TestScheduler(), 1); + }; + source.subscribe(consumer); + source.onNext(value); + consumer.disposable.dispose(); + consumer.disposable.dispose(); + assertTrue(consumer.disposable.isDisposed()); + assertFalse(source.hasObservers()); + return consumer; + } + + static ReplayConsumer cancelProcessor(int kind) { + Object value = new Object(); + ReplayConsumer consumer = new ReplayConsumer(value); + ReplayProcessor source = switch (kind) { + case 0 -> ReplayProcessor.create(); + case 1 -> ReplayProcessor.createWithSize(1); + default -> ReplayProcessor.createWithTimeAndSize(1, TimeUnit.DAYS, new TestScheduler(), 1); + }; + source.subscribe(consumer); + source.onNext(value); + consumer.upstream.cancel(); + consumer.upstream.cancel(); + consumer.upstream.request(1); + assertFalse(source.hasSubscribers()); + return consumer; + } + + static void assertReclaimed(ReplayConsumer consumer) throws Exception { + assertEquals(1, consumer.values); + // Keep the original cancelled handle reachable, as a custom consumer may do. + try { + for (int i = 0; i < 20 && consumer.value.get() != null; i++) { + System.gc(); + Thread.sleep(50); + } + assertNull(consumer.value.get(), "A cancelled replay consumer retained its cached value"); + } finally { + Reference.reachabilityFence(consumer); + } + } + + static final class ReplayConsumer implements Observer, FlowableSubscriber { + final WeakReference value; + Disposable disposable; + Subscription upstream; + int values; + + ReplayConsumer(Object value) { + this.value = new WeakReference<>(value); + } + + @Override + public void onSubscribe(Disposable d) { + disposable = d; + } + + @Override + public void onSubscribe(Subscription s) { + upstream = s; + s.request(Long.MAX_VALUE); + } + + @Override + public void onNext(Object item) { + values++; + } + + @Override + public void onError(Throwable error) { + throw new AssertionError(error); + } + + @Override + public void onComplete() { + // The cancellation tests use nonterminating sources. + } + } +}