Make AbstractResponseBodySubscriber.onSubscribe thread-safe
When there are simultaneous invocations of onSubscribe, only the first one should succeed, the rest should cancel the provided subscriptions
This commit is contained in:
@@ -18,6 +18,7 @@ package org.springframework.http.server.reactive;
|
|||||||
|
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
import java.nio.channels.Channel;
|
import java.nio.channels.Channel;
|
||||||
|
import java.util.Objects;
|
||||||
import java.util.concurrent.atomic.AtomicReference;
|
import java.util.concurrent.atomic.AtomicReference;
|
||||||
import javax.servlet.WriteListener;
|
import javax.servlet.WriteListener;
|
||||||
|
|
||||||
@@ -169,11 +170,13 @@ abstract class AbstractResponseBodySubscriber implements Subscriber<DataBuffer>
|
|||||||
@Override
|
@Override
|
||||||
void onSubscribe(AbstractResponseBodySubscriber subscriber,
|
void onSubscribe(AbstractResponseBodySubscriber subscriber,
|
||||||
Subscription subscription) {
|
Subscription subscription) {
|
||||||
if (BackpressureUtils.validate(subscriber.subscription, subscription)) {
|
Objects.requireNonNull(subscription, "Subscription cannot be null");
|
||||||
|
if (subscriber.changeState(this, REQUESTED)) {
|
||||||
subscriber.subscription = subscription;
|
subscriber.subscription = subscription;
|
||||||
if (subscriber.changeState(this, REQUESTED)) {
|
subscription.request(1);
|
||||||
subscription.request(1);
|
}
|
||||||
}
|
else {
|
||||||
|
super.onSubscribe(subscriber, subscription);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|||||||
Reference in New Issue
Block a user