Revert "Use optimistic locking where possible in ResponseBodyEmitter"
This reverts commit e67f892e44.
Closes gh-34762
This commit is contained in:
@@ -21,7 +21,6 @@ import java.util.ArrayList;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.springframework.http.MediaType;
|
||||
@@ -77,7 +76,7 @@ public class ResponseBodyEmitter {
|
||||
private final Set<DataWithMediaType> earlySendAttempts = new LinkedHashSet<>(8);
|
||||
|
||||
/** Store successful completion before the handler is initialized. */
|
||||
private final AtomicBoolean complete = new AtomicBoolean();
|
||||
private boolean complete;
|
||||
|
||||
/** Store an error before the handler is initialized. */
|
||||
@Nullable
|
||||
@@ -128,7 +127,7 @@ public class ResponseBodyEmitter {
|
||||
this.earlySendAttempts.clear();
|
||||
}
|
||||
|
||||
if (this.complete.get()) {
|
||||
if (this.complete) {
|
||||
if (this.failure != null) {
|
||||
this.handler.completeWithError(this.failure);
|
||||
}
|
||||
@@ -143,12 +142,11 @@ public class ResponseBodyEmitter {
|
||||
}
|
||||
}
|
||||
|
||||
void initializeWithError(Throwable ex) {
|
||||
if (this.complete.compareAndSet(false, true)) {
|
||||
this.failure = ex;
|
||||
this.earlySendAttempts.clear();
|
||||
this.errorCallback.accept(ex);
|
||||
}
|
||||
synchronized void initializeWithError(Throwable ex) {
|
||||
this.complete = true;
|
||||
this.failure = ex;
|
||||
this.earlySendAttempts.clear();
|
||||
this.errorCallback.accept(ex);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -186,7 +184,7 @@ public class ResponseBodyEmitter {
|
||||
* @throws java.lang.IllegalStateException wraps any other errors
|
||||
*/
|
||||
public synchronized void send(Object object, @Nullable MediaType mediaType) throws IOException {
|
||||
Assert.state(!this.complete.get(), () -> "ResponseBodyEmitter has already completed" +
|
||||
Assert.state(!this.complete, () -> "ResponseBodyEmitter has already completed" +
|
||||
(this.failure != null ? " with error: " + this.failure : ""));
|
||||
if (this.handler != null) {
|
||||
try {
|
||||
@@ -214,7 +212,7 @@ public class ResponseBodyEmitter {
|
||||
* @since 6.0.12
|
||||
*/
|
||||
public synchronized void send(Set<DataWithMediaType> items) throws IOException {
|
||||
Assert.state(!this.complete.get(), () -> "ResponseBodyEmitter has already completed" +
|
||||
Assert.state(!this.complete, () -> "ResponseBodyEmitter has already completed" +
|
||||
(this.failure != null ? " with error: " + this.failure : ""));
|
||||
sendInternal(items);
|
||||
}
|
||||
@@ -247,8 +245,9 @@ public class ResponseBodyEmitter {
|
||||
* to complete request processing. It should not be used after container
|
||||
* related events such as an error while {@link #send(Object) sending}.
|
||||
*/
|
||||
public void complete() {
|
||||
if (this.complete.compareAndSet(false, true) && this.handler != null) {
|
||||
public synchronized void complete() {
|
||||
this.complete = true;
|
||||
if (this.handler != null) {
|
||||
this.handler.complete();
|
||||
}
|
||||
}
|
||||
@@ -264,12 +263,11 @@ public class ResponseBodyEmitter {
|
||||
* container related events such as an error while
|
||||
* {@link #send(Object) sending}.
|
||||
*/
|
||||
public void completeWithError(Throwable ex) {
|
||||
if (this.complete.compareAndSet(false, true)) {
|
||||
this.failure = ex;
|
||||
if (this.handler != null) {
|
||||
this.handler.completeWithError(ex);
|
||||
}
|
||||
public synchronized void completeWithError(Throwable ex) {
|
||||
this.complete = true;
|
||||
this.failure = ex;
|
||||
if (this.handler != null) {
|
||||
this.handler.completeWithError(ex);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -278,7 +276,7 @@ public class ResponseBodyEmitter {
|
||||
* called from a container thread when an async request times out.
|
||||
* <p>As of 6.2, one can register multiple callbacks for this event.
|
||||
*/
|
||||
public void onTimeout(Runnable callback) {
|
||||
public synchronized void onTimeout(Runnable callback) {
|
||||
this.timeoutCallback.addDelegate(callback);
|
||||
}
|
||||
|
||||
@@ -289,7 +287,7 @@ public class ResponseBodyEmitter {
|
||||
* <p>As of 6.2, one can register multiple callbacks for this event.
|
||||
* @since 5.0
|
||||
*/
|
||||
public void onError(Consumer<Throwable> callback) {
|
||||
public synchronized void onError(Consumer<Throwable> callback) {
|
||||
this.errorCallback.addDelegate(callback);
|
||||
}
|
||||
|
||||
@@ -300,7 +298,7 @@ public class ResponseBodyEmitter {
|
||||
* detecting that a {@code ResponseBodyEmitter} instance is no longer usable.
|
||||
* <p>As of 6.2, one can register multiple callbacks for this event.
|
||||
*/
|
||||
public void onCompletion(Runnable callback) {
|
||||
public synchronized void onCompletion(Runnable callback) {
|
||||
this.completionCallback.addDelegate(callback);
|
||||
}
|
||||
|
||||
@@ -371,15 +369,15 @@ public class ResponseBodyEmitter {
|
||||
|
||||
private class DefaultCallback implements Runnable {
|
||||
|
||||
private final List<Runnable> delegates = new ArrayList<>(1);
|
||||
private List<Runnable> delegates = new ArrayList<>(1);
|
||||
|
||||
public synchronized void addDelegate(Runnable delegate) {
|
||||
public void addDelegate(Runnable delegate) {
|
||||
this.delegates.add(delegate);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
ResponseBodyEmitter.this.complete.compareAndSet(false, true);
|
||||
ResponseBodyEmitter.this.complete = true;
|
||||
for (Runnable delegate : this.delegates) {
|
||||
delegate.run();
|
||||
}
|
||||
@@ -389,15 +387,15 @@ public class ResponseBodyEmitter {
|
||||
|
||||
private class ErrorCallback implements Consumer<Throwable> {
|
||||
|
||||
private final List<Consumer<Throwable>> delegates = new ArrayList<>(1);
|
||||
private List<Consumer<Throwable>> delegates = new ArrayList<>(1);
|
||||
|
||||
public synchronized void addDelegate(Consumer<Throwable> callback) {
|
||||
public void addDelegate(Consumer<Throwable> callback) {
|
||||
this.delegates.add(callback);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void accept(Throwable t) {
|
||||
ResponseBodyEmitter.this.complete.compareAndSet(false, true);
|
||||
ResponseBodyEmitter.this.complete = true;
|
||||
for(Consumer<Throwable> delegate : this.delegates) {
|
||||
delegate.accept(t);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user