Expected Behavior
An element that timeout receives but will not emit, e.g. due to cancellation, should be passed to the discard hook (as per "Dealing with Objects that Need Cleanup").
Actual Behavior
FluxTimeout.TimeoutMainSubscriber.onNext passes the element to Operators.onNextDropped in both of its early returns: when index is Long.MIN_VALUE (cancelled, or timed out, or the source already terminated), and when the CAS loses to a timeout firing concurrently.
The SerializedSubscriber that this operator wraps its own subscriber in already makes the distinction correctly: onDiscard when cancelled, onNextDropped when done (#2077). The drop here dates from #744 (2017), before the discard hook existed (#999, 3.2).
Steps to Reproduce
DEFER_CANCELLATION models an element that loses the race with a cancel deterministically:
ts.cancel();
source.next("late");
assertThat(discarded).containsExactly("late"); // fails: discarded is empty
Same with the timeout firing instead of the cancel. Both fail on 3.8.7; the third test below is a guard for the case where onNextDropped stays correct, and passes on 3.8.7.
Complete runnable test class (reactor-core, reactor-test, JUnit 5, AssertJ)
package demo;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.TimeoutException;
import org.junit.jupiter.api.Test;
import reactor.core.publisher.Hooks;
import reactor.test.StepVerifier;
import reactor.test.publisher.TestPublisher;
import reactor.test.subscriber.TestSubscriber;
import static org.assertj.core.api.Assertions.assertThat;
public class TimeoutDiscardTest {
@Test
public void discardsElementArrivingAfterCancel() {
List<Object> discarded = new ArrayList<>();
TestPublisher<String> source =
TestPublisher.createNoncompliant(TestPublisher.Violation.DEFER_CANCELLATION);
TestSubscriber<String> ts = TestSubscriber.builder().initialRequest(1).build();
source.flux()
.timeout(Duration.ofSeconds(30))
.doOnDiscard(Object.class, discarded::add)
.subscribe(ts);
ts.cancel();
source.next("late");
assertThat(discarded).as("discarded").containsExactly("late");
}
@Test
public void discardsElementArrivingAfterTimeout() {
List<Object> discarded = new ArrayList<>();
TestPublisher<String> source =
TestPublisher.createNoncompliant(TestPublisher.Violation.DEFER_CANCELLATION);
StepVerifier.withVirtualTime(() -> source.flux()
.timeout(Duration.ofMillis(100))
.doOnDiscard(Object.class, discarded::add))
.thenAwait(Duration.ofMillis(200))
.expectError(TimeoutException.class)
.verify(Duration.ofSeconds(5));
source.next("late");
assertThat(discarded).as("discarded").containsExactly("late");
}
@Test // guard: passes on 3.8.7, and must keep passing
public void dropsElementArrivingAfterSourceTerminated() {
List<Object> discarded = new ArrayList<>();
List<Object> dropped = new ArrayList<>();
TestPublisher<String> source =
TestPublisher.createNoncompliant(TestPublisher.Violation.CLEANUP_ON_TERMINATE);
TestSubscriber<String> ts = TestSubscriber.builder().initialRequest(1).build();
source.flux()
.timeout(Duration.ofSeconds(30))
.doOnDiscard(Object.class, discarded::add)
.subscribe(ts);
Hooks.onNextDropped(dropped::add);
try {
source.complete();
source.next("malformed");
}
finally {
Hooks.resetOnNextDropped();
}
assertThat(dropped).as("dropped").containsExactly("malformed");
assertThat(discarded).as("discarded").isEmpty();
}
}
Possible Solution
Route the element to Operators.onDiscard unless the source has signalled a terminal event, in which case onNextDropped stays correct.
Context
This and the sibling issues were discovered during coverage-guided fuzz testing of an internal application using JDK 21.0.12 on RHEL 8 on a 96 core host. We are happy to provide more information or update the proposed fixes as required.
Note that the proposed fixes are against main but all these issues are present in every 3.x version. We would be happy to port to 3.7.x if desired.
See also various other related, but separate, open issues:
- Discard contract violations in other operators:
- Other cancellation races:
Expected Behavior
An element that
timeoutreceives but will not emit, e.g. due to cancellation, should be passed to the discard hook (as per "Dealing with Objects that Need Cleanup").Actual Behavior
FluxTimeout.TimeoutMainSubscriber.onNextpasses the element toOperators.onNextDroppedin both of its early returns: whenindexisLong.MIN_VALUE(cancelled, or timed out, or the source already terminated), and when the CAS loses to a timeout firing concurrently.The
SerializedSubscriberthat this operator wraps its own subscriber in already makes the distinction correctly:onDiscardwhen cancelled,onNextDroppedwhen done (#2077). The drop here dates from #744 (2017), before the discard hook existed (#999, 3.2).Steps to Reproduce
DEFER_CANCELLATIONmodels an element that loses the race with a cancel deterministically:Same with the timeout firing instead of the cancel. Both fail on 3.8.7; the third test below is a guard for the case where
onNextDroppedstays correct, and passes on 3.8.7.Complete runnable test class (reactor-core, reactor-test, JUnit 5, AssertJ)
Possible Solution
Route the element to
Operators.onDiscardunless the source has signalled a terminal event, in which caseonNextDroppedstays correct.Context
This and the sibling issues were discovered during coverage-guided fuzz testing of an internal application using JDK 21.0.12 on RHEL 8 on a 96 core host. We are happy to provide more information or update the proposed fixes as required.
Note that the proposed fixes are against main but all these issues are present in every 3.x version. We would be happy to port to 3.7.x if desired.
See also various other related, but separate, open issues:
MonoSinglepropagates empty completion after cancellation #4330