Skip to content

timeout drops an element that races cancellation instead of discarding it #4395

Description

@duaneg-tp

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:

No activity

Activity on this issue will appear here.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions