From 074216082166d155ad39a1c32f7033ca77181f30 Mon Sep 17 00:00:00 2001 From: Violeta Georgieva Date: Thu, 13 Dec 2018 23:02:30 +0200 Subject: [PATCH] Propagate the cancel signal to the downstream --- .../http/server/reactive/ChannelSendOperator.java | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/spring-web/src/main/java/org/springframework/http/server/reactive/ChannelSendOperator.java b/spring-web/src/main/java/org/springframework/http/server/reactive/ChannelSendOperator.java index 9ab8145175..540979f13c 100644 --- a/spring-web/src/main/java/org/springframework/http/server/reactive/ChannelSendOperator.java +++ b/spring-web/src/main/java/org/springframework/http/server/reactive/ChannelSendOperator.java @@ -337,6 +337,8 @@ public class ChannelSendOperator extends Mono implements Scannable { private final WriteBarrier writeBarrier; + private Subscription subscription; + public WriteCompletionBarrier(CoreSubscriber subscriber, WriteBarrier writeBarrier) { this.completionSubscriber = subscriber; @@ -356,6 +358,7 @@ public class ChannelSendOperator extends Mono implements Scannable { @Override public void onSubscribe(Subscription subscription) { + this.subscription = subscription; subscription.request(Long.MAX_VALUE); } @@ -387,6 +390,10 @@ public class ChannelSendOperator extends Mono implements Scannable { @Override public void cancel() { this.writeBarrier.cancel(); + Subscription subscription = this.subscription; + if (subscription != null) { + subscription.cancel(); + } } }