Implement Eclipse Jetty core HTTP handler adapter
This provides an implementation of an HTTP Handler Adapter that is coded directly to the Eclipse Jetty core API, bypassing any servlet implementation. This includes a Jetty implementation of the spring `WebSocketClient` interface, `JettyWebSocketClient`, using an explicit dependency to the jetty-websocket-api. Closes gh-32097 Co-authored-by: Lachlan Roberts <lachlan@webtide.com> Co-authored-by: Arjen Poutsma <arjen.poutsma@broadcom.com>
This commit is contained in:
@@ -22,16 +22,9 @@ import java.nio.charset.StandardCharsets;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.IntPredicate;
|
||||
|
||||
import org.eclipse.jetty.util.BufferUtil;
|
||||
import org.eclipse.jetty.websocket.api.Callback;
|
||||
import org.eclipse.jetty.websocket.api.Frame;
|
||||
import org.eclipse.jetty.websocket.api.Session;
|
||||
import org.eclipse.jetty.websocket.api.annotations.OnWebSocketClose;
|
||||
import org.eclipse.jetty.websocket.api.annotations.OnWebSocketError;
|
||||
import org.eclipse.jetty.websocket.api.annotations.OnWebSocketFrame;
|
||||
import org.eclipse.jetty.websocket.api.annotations.OnWebSocketMessage;
|
||||
import org.eclipse.jetty.websocket.api.annotations.OnWebSocketOpen;
|
||||
import org.eclipse.jetty.websocket.api.annotations.WebSocket;
|
||||
import org.eclipse.jetty.websocket.core.OpCode;
|
||||
|
||||
import org.springframework.core.io.buffer.CloseableDataBuffer;
|
||||
import org.springframework.core.io.buffer.DataBuffer;
|
||||
@@ -44,18 +37,14 @@ import org.springframework.web.reactive.socket.WebSocketMessage;
|
||||
import org.springframework.web.reactive.socket.WebSocketMessage.Type;
|
||||
|
||||
/**
|
||||
* Jetty {@link WebSocket @WebSocket} handler that delegates events to a
|
||||
* Jetty {@link org.eclipse.jetty.websocket.api.Session.Listener} handler that delegates events to a
|
||||
* reactive {@link WebSocketHandler} and its session.
|
||||
*
|
||||
* @author Violeta Georgieva
|
||||
* @author Rossen Stoyanchev
|
||||
* @since 5.0
|
||||
*/
|
||||
@WebSocket
|
||||
public class JettyWebSocketHandlerAdapter {
|
||||
|
||||
private static final ByteBuffer EMPTY_PAYLOAD = ByteBuffer.wrap(new byte[0]);
|
||||
|
||||
public class JettyWebSocketHandlerAdapter implements Session.Listener {
|
||||
|
||||
private final WebSocketHandler delegateHandler;
|
||||
|
||||
@@ -74,70 +63,62 @@ public class JettyWebSocketHandlerAdapter {
|
||||
this.sessionFactory = sessionFactory;
|
||||
}
|
||||
|
||||
|
||||
@OnWebSocketOpen
|
||||
@Override
|
||||
public void onWebSocketOpen(Session session) {
|
||||
this.delegateSession = this.sessionFactory.apply(session);
|
||||
this.delegateHandler.handle(this.delegateSession)
|
||||
JettyWebSocketSession delegateSession = this.sessionFactory.apply(session);
|
||||
this.delegateSession = delegateSession;
|
||||
this.delegateHandler.handle(delegateSession)
|
||||
.checkpoint(session.getUpgradeRequest().getRequestURI() + " [JettyWebSocketHandlerAdapter]")
|
||||
.subscribe(this.delegateSession);
|
||||
.subscribe(unused -> {}, delegateSession::onHandlerError, delegateSession::onHandleComplete);
|
||||
}
|
||||
|
||||
@OnWebSocketMessage
|
||||
@Override
|
||||
public void onWebSocketText(String message) {
|
||||
if (this.delegateSession != null) {
|
||||
byte[] bytes = message.getBytes(StandardCharsets.UTF_8);
|
||||
DataBuffer buffer = this.delegateSession.bufferFactory().wrap(bytes);
|
||||
WebSocketMessage webSocketMessage = new WebSocketMessage(Type.TEXT, buffer);
|
||||
this.delegateSession.handleMessage(webSocketMessage.getType(), webSocketMessage);
|
||||
}
|
||||
Assert.state(this.delegateSession != null, "No delegate session available");
|
||||
byte[] bytes = message.getBytes(StandardCharsets.UTF_8);
|
||||
DataBuffer buffer = this.delegateSession.bufferFactory().wrap(bytes);
|
||||
WebSocketMessage webSocketMessage = new WebSocketMessage(Type.TEXT, buffer);
|
||||
this.delegateSession.handleMessage(webSocketMessage);
|
||||
}
|
||||
|
||||
@OnWebSocketMessage
|
||||
@Override
|
||||
public void onWebSocketBinary(ByteBuffer byteBuffer, Callback callback) {
|
||||
if (this.delegateSession != null) {
|
||||
DataBuffer buffer = this.delegateSession.bufferFactory().wrap(byteBuffer);
|
||||
buffer = new JettyDataBuffer(buffer, callback);
|
||||
WebSocketMessage webSocketMessage = new WebSocketMessage(Type.BINARY, buffer);
|
||||
this.delegateSession.handleMessage(webSocketMessage.getType(), webSocketMessage);
|
||||
}
|
||||
Assert.state(this.delegateSession != null, "No delegate session available");
|
||||
DataBuffer buffer = this.delegateSession.bufferFactory().wrap(byteBuffer);
|
||||
buffer = new JettyCallbackDataBuffer(buffer, callback);
|
||||
WebSocketMessage webSocketMessage = new WebSocketMessage(Type.BINARY, buffer);
|
||||
this.delegateSession.handleMessage(webSocketMessage);
|
||||
}
|
||||
|
||||
@OnWebSocketFrame
|
||||
public void onWebSocketFrame(Frame frame, Callback callback) {
|
||||
if (this.delegateSession != null) {
|
||||
if (OpCode.PONG == frame.getOpCode()) {
|
||||
ByteBuffer byteBuffer = (frame.getPayload() != null ? frame.getPayload() : EMPTY_PAYLOAD);
|
||||
DataBuffer buffer = this.delegateSession.bufferFactory().wrap(byteBuffer);
|
||||
buffer = new JettyDataBuffer(buffer, callback);
|
||||
WebSocketMessage webSocketMessage = new WebSocketMessage(Type.PONG, buffer);
|
||||
this.delegateSession.handleMessage(webSocketMessage.getType(), webSocketMessage);
|
||||
}
|
||||
}
|
||||
@Override
|
||||
public void onWebSocketPong(ByteBuffer payload) {
|
||||
Assert.state(this.delegateSession != null, "No delegate session available");
|
||||
DataBuffer buffer = this.delegateSession.bufferFactory().wrap(BufferUtil.copy(payload));
|
||||
WebSocketMessage webSocketMessage = new WebSocketMessage(Type.PONG, buffer);
|
||||
this.delegateSession.handleMessage(webSocketMessage);
|
||||
}
|
||||
|
||||
@OnWebSocketClose
|
||||
@Override
|
||||
public void onWebSocketClose(int statusCode, String reason) {
|
||||
if (this.delegateSession != null) {
|
||||
this.delegateSession.handleClose(CloseStatus.create(statusCode, reason));
|
||||
}
|
||||
Assert.state(this.delegateSession != null, "No delegate session available");
|
||||
this.delegateSession.handleClose(CloseStatus.create(statusCode, reason));
|
||||
}
|
||||
|
||||
@OnWebSocketError
|
||||
@Override
|
||||
public void onWebSocketError(Throwable cause) {
|
||||
if (this.delegateSession != null) {
|
||||
this.delegateSession.handleError(cause);
|
||||
}
|
||||
Assert.state(this.delegateSession != null, "No delegate session available");
|
||||
this.delegateSession.handleError(cause);
|
||||
}
|
||||
|
||||
|
||||
private static final class JettyDataBuffer implements CloseableDataBuffer {
|
||||
private static final class JettyCallbackDataBuffer implements CloseableDataBuffer {
|
||||
|
||||
private final DataBuffer delegate;
|
||||
|
||||
private final Callback callback;
|
||||
|
||||
public JettyDataBuffer(DataBuffer delegate, Callback callback) {
|
||||
|
||||
public JettyCallbackDataBuffer(DataBuffer delegate, Callback callback) {
|
||||
Assert.notNull(delegate, "'delegate` must not be null");
|
||||
Assert.notNull(callback, "Callback must not be null");
|
||||
this.delegate = delegate;
|
||||
@@ -272,13 +253,13 @@ public class JettyWebSocketHandlerAdapter {
|
||||
@Deprecated
|
||||
public DataBuffer slice(int index, int length) {
|
||||
DataBuffer delegateSlice = this.delegate.slice(index, length);
|
||||
return new JettyDataBuffer(delegateSlice, this.callback);
|
||||
return new JettyCallbackDataBuffer(delegateSlice, this.callback);
|
||||
}
|
||||
|
||||
@Override
|
||||
public DataBuffer split(int index) {
|
||||
DataBuffer delegateSplit = this.delegate.split(index);
|
||||
return new JettyDataBuffer(delegateSplit, this.callback);
|
||||
return new JettyCallbackDataBuffer(delegateSplit, this.callback);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -16,18 +16,26 @@
|
||||
|
||||
package org.springframework.web.reactive.socket.adapter;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.ByteBuffer;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.concurrent.locks.Lock;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
|
||||
import org.eclipse.jetty.util.BufferUtil;
|
||||
import org.eclipse.jetty.util.IteratingCallback;
|
||||
import org.eclipse.jetty.websocket.api.Callback;
|
||||
import org.eclipse.jetty.websocket.api.Session;
|
||||
import org.eclipse.jetty.websocket.api.StatusCode;
|
||||
import org.reactivestreams.Publisher;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.FluxSink;
|
||||
import reactor.core.publisher.Mono;
|
||||
import reactor.core.publisher.Sinks;
|
||||
|
||||
import org.springframework.core.io.buffer.DataBuffer;
|
||||
import org.springframework.core.io.buffer.DataBufferFactory;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.ObjectUtils;
|
||||
import org.springframework.web.reactive.socket.CloseStatus;
|
||||
import org.springframework.web.reactive.socket.HandshakeInfo;
|
||||
@@ -36,13 +44,29 @@ import org.springframework.web.reactive.socket.WebSocketSession;
|
||||
|
||||
/**
|
||||
* Spring {@link WebSocketSession} implementation that adapts to a Jetty
|
||||
* WebSocket {@link org.eclipse.jetty.websocket.api.Session}.
|
||||
* WebSocket {@link Session}.
|
||||
*
|
||||
* @author Violeta Georgieva
|
||||
* @author Rossen Stoyanchev
|
||||
* @since 5.0
|
||||
*/
|
||||
public class JettyWebSocketSession extends AbstractListenerWebSocketSession<Session> {
|
||||
public class JettyWebSocketSession extends AbstractWebSocketSession<Session> {
|
||||
|
||||
private final Flux<WebSocketMessage> flux;
|
||||
|
||||
private final Sinks.One<CloseStatus> closeStatusSink = Sinks.one();
|
||||
|
||||
private final Lock lock = new ReentrantLock();
|
||||
|
||||
private long requested = 0;
|
||||
|
||||
private boolean awaitingMessage = false;
|
||||
|
||||
@Nullable
|
||||
private FluxSink<WebSocketMessage> sink;
|
||||
|
||||
@Nullable
|
||||
private final Sinks.Empty<Void> handlerCompletionSink;
|
||||
|
||||
public JettyWebSocketSession(Session session, HandshakeInfo info, DataBufferFactory factory) {
|
||||
this(session, info, factory, null);
|
||||
@@ -51,52 +75,88 @@ public class JettyWebSocketSession extends AbstractListenerWebSocketSession<Sess
|
||||
public JettyWebSocketSession(Session session, HandshakeInfo info, DataBufferFactory factory,
|
||||
@Nullable Sinks.Empty<Void> completionSink) {
|
||||
|
||||
super(session, ObjectUtils.getIdentityHexString(session), info, factory, completionSink);
|
||||
// TODO: suspend causes failures if invoked at this stage
|
||||
// suspendReceiving();
|
||||
}
|
||||
super(session, ObjectUtils.getIdentityHexString(session), info, factory);
|
||||
this.handlerCompletionSink = completionSink;
|
||||
this.flux = Flux.create(emitter -> {
|
||||
this.sink = emitter;
|
||||
emitter.onRequest(n -> {
|
||||
boolean demand = false;
|
||||
this.lock.lock();
|
||||
try {
|
||||
this.requested = Math.addExact(this.requested, n);
|
||||
if (this.requested < 0L) {
|
||||
this.requested = Long.MAX_VALUE;
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
protected boolean canSuspendReceiving() {
|
||||
// Jetty 12 TODO: research suspend functionality in Jetty 12
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void suspendReceiving() {
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void resumeReceiving() {
|
||||
}
|
||||
|
||||
@Override
|
||||
protected boolean sendMessage(WebSocketMessage message) throws IOException {
|
||||
DataBuffer dataBuffer = message.getPayload();
|
||||
Session session = getDelegate();
|
||||
if (WebSocketMessage.Type.TEXT.equals(message.getType())) {
|
||||
getSendProcessor().setReadyToSend(false);
|
||||
String text = dataBuffer.toString(StandardCharsets.UTF_8);
|
||||
session.sendText(text, new SendProcessorCallback());
|
||||
}
|
||||
else {
|
||||
if (WebSocketMessage.Type.BINARY.equals(message.getType())) {
|
||||
getSendProcessor().setReadyToSend(false);
|
||||
}
|
||||
try (DataBuffer.ByteBufferIterator iterator = dataBuffer.readableByteBuffers()) {
|
||||
while (iterator.hasNext()) {
|
||||
ByteBuffer byteBuffer = iterator.next();
|
||||
switch (message.getType()) {
|
||||
case BINARY -> session.sendBinary(byteBuffer, new SendProcessorCallback());
|
||||
case PING -> session.sendPing(byteBuffer, new SendProcessorCallback());
|
||||
case PONG -> session.sendPong(byteBuffer, new SendProcessorCallback());
|
||||
default -> throw new IllegalArgumentException("Unexpected message type: " + message.getType());
|
||||
if (!this.awaitingMessage && this.requested > 0) {
|
||||
if (this.requested != Long.MAX_VALUE) {
|
||||
this.requested--;
|
||||
}
|
||||
this.awaitingMessage = true;
|
||||
demand = true;
|
||||
}
|
||||
}
|
||||
finally {
|
||||
this.lock.unlock();
|
||||
}
|
||||
|
||||
if (demand) {
|
||||
getDelegate().demand();
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
void handleMessage(WebSocketMessage message) {
|
||||
Assert.state(this.sink != null, "No sink available");
|
||||
this.sink.next(message);
|
||||
|
||||
boolean demand = false;
|
||||
this.lock.lock();
|
||||
try {
|
||||
if (!this.awaitingMessage) {
|
||||
throw new IllegalStateException();
|
||||
}
|
||||
this.awaitingMessage = false;
|
||||
if (this.requested > 0) {
|
||||
if (this.requested != Long.MAX_VALUE) {
|
||||
this.requested--;
|
||||
}
|
||||
this.awaitingMessage = true;
|
||||
demand = true;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
finally {
|
||||
this.lock.unlock();
|
||||
}
|
||||
|
||||
if (demand) {
|
||||
getDelegate().demand();
|
||||
}
|
||||
}
|
||||
|
||||
void handleError(Throwable ex) {
|
||||
}
|
||||
|
||||
void handleClose(CloseStatus closeStatus) {
|
||||
this.closeStatusSink.tryEmitValue(closeStatus);
|
||||
if (this.sink != null) {
|
||||
this.sink.complete();
|
||||
}
|
||||
}
|
||||
|
||||
void onHandlerError(Throwable error) {
|
||||
if (JettyWebSocketSession.this.handlerCompletionSink != null) {
|
||||
JettyWebSocketSession.this.handlerCompletionSink.tryEmitError(error);
|
||||
}
|
||||
getDelegate().close(StatusCode.SERVER_ERROR, error.getMessage(), Callback.NOOP);
|
||||
}
|
||||
|
||||
void onHandleComplete() {
|
||||
if (JettyWebSocketSession.this.handlerCompletionSink != null) {
|
||||
JettyWebSocketSession.this.handlerCompletionSink.tryEmitEmpty();
|
||||
}
|
||||
getDelegate().close(StatusCode.NORMAL, null, Callback.NOOP);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -108,25 +168,81 @@ public class JettyWebSocketSession extends AbstractListenerWebSocketSession<Sess
|
||||
public Mono<Void> close(CloseStatus status) {
|
||||
Callback.Completable callback = new Callback.Completable();
|
||||
getDelegate().close(status.getCode(), status.getReason(), callback);
|
||||
|
||||
return Mono.fromFuture(callback);
|
||||
}
|
||||
|
||||
|
||||
private final class SendProcessorCallback implements Callback {
|
||||
|
||||
@Override
|
||||
public void fail(Throwable x) {
|
||||
getSendProcessor().cancel();
|
||||
getSendProcessor().onError(x);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void succeed() {
|
||||
getSendProcessor().setReadyToSend(true);
|
||||
getSendProcessor().onWritePossible();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Mono<CloseStatus> closeStatus() {
|
||||
return this.closeStatusSink.asMono();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Flux<WebSocketMessage> receive() {
|
||||
return this.flux;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Mono<Void> send(Publisher<WebSocketMessage> messages) {
|
||||
return Flux.from(messages)
|
||||
.flatMap(this::sendMessage, 1)
|
||||
.then();
|
||||
}
|
||||
|
||||
protected Mono<Void> sendMessage(WebSocketMessage message) {
|
||||
|
||||
Callback.Completable completable = new Callback.Completable();
|
||||
DataBuffer dataBuffer = message.getPayload();
|
||||
Session session = getDelegate();
|
||||
if (WebSocketMessage.Type.TEXT.equals(message.getType())) {
|
||||
String text = dataBuffer.toString(StandardCharsets.UTF_8);
|
||||
session.sendText(text, completable);
|
||||
}
|
||||
else {
|
||||
switch (message.getType()) {
|
||||
case BINARY -> {
|
||||
@SuppressWarnings("resource")
|
||||
DataBuffer.ByteBufferIterator iterator = dataBuffer.readableByteBuffers();
|
||||
new IteratingCallback() {
|
||||
@Override
|
||||
protected Action process() {
|
||||
if (!iterator.hasNext()) {
|
||||
return Action.SUCCEEDED;
|
||||
}
|
||||
|
||||
ByteBuffer buffer = iterator.next();
|
||||
boolean last = iterator.hasNext();
|
||||
session.sendPartialBinary(buffer, last, Callback.from(this::succeeded, this::failed));
|
||||
return Action.SCHEDULED;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void onCompleteSuccess() {
|
||||
iterator.close();
|
||||
completable.succeed();
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void onCompleteFailure(Throwable cause) {
|
||||
iterator.close();
|
||||
completable.fail(cause);
|
||||
}
|
||||
}.iterate();
|
||||
}
|
||||
case PING -> {
|
||||
// Maximum size of Control frame payload is 125, per RFC 6455.
|
||||
ByteBuffer buffer = BufferUtil.allocate(125);
|
||||
dataBuffer.toByteBuffer(buffer);
|
||||
session.sendPing(buffer, completable);
|
||||
}
|
||||
case PONG -> {
|
||||
// Maximum size of Control frame payload is 125, per RFC 6455.
|
||||
ByteBuffer buffer = BufferUtil.allocate(125);
|
||||
dataBuffer.toByteBuffer(buffer);
|
||||
session.sendPong(buffer, completable);
|
||||
}
|
||||
default -> throw new IllegalArgumentException("Unexpected message type: " + message.getType());
|
||||
}
|
||||
}
|
||||
return Mono.fromFuture(completable);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,111 @@
|
||||
/*
|
||||
* Copyright 2002-2022 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* https://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.web.reactive.socket.client;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.net.URI;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import org.eclipse.jetty.client.Request;
|
||||
import org.eclipse.jetty.client.Response;
|
||||
import org.eclipse.jetty.http.HttpHeader;
|
||||
import org.eclipse.jetty.util.component.LifeCycle;
|
||||
import org.eclipse.jetty.websocket.client.ClientUpgradeRequest;
|
||||
import org.eclipse.jetty.websocket.client.JettyUpgradeListener;
|
||||
import reactor.core.publisher.Mono;
|
||||
import reactor.core.publisher.Sinks;
|
||||
|
||||
import org.springframework.context.Lifecycle;
|
||||
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
|
||||
import org.springframework.http.HttpHeaders;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.web.reactive.socket.HandshakeInfo;
|
||||
import org.springframework.web.reactive.socket.WebSocketHandler;
|
||||
import org.springframework.web.reactive.socket.adapter.JettyWebSocketHandlerAdapter;
|
||||
import org.springframework.web.reactive.socket.adapter.JettyWebSocketSession;
|
||||
|
||||
public class JettyWebSocketClient implements WebSocketClient, Lifecycle {
|
||||
|
||||
private final org.eclipse.jetty.websocket.client.WebSocketClient client;
|
||||
|
||||
public JettyWebSocketClient() {
|
||||
this(new org.eclipse.jetty.websocket.client.WebSocketClient());
|
||||
}
|
||||
|
||||
public JettyWebSocketClient(org.eclipse.jetty.websocket.client.WebSocketClient client) {
|
||||
this.client = client;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
LifeCycle.start(this.client);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop() {
|
||||
LifeCycle.stop(this.client);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isRunning() {
|
||||
return this.client.isRunning();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Mono<Void> execute(URI url, WebSocketHandler handler) {
|
||||
return execute(url, null, handler);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Mono<Void> execute(URI url, @Nullable HttpHeaders headers, WebSocketHandler handler) {
|
||||
|
||||
ClientUpgradeRequest upgradeRequest = new ClientUpgradeRequest();
|
||||
upgradeRequest.setSubProtocols(handler.getSubProtocols());
|
||||
if (headers != null) {
|
||||
headers.keySet().forEach(header -> upgradeRequest.setHeader(header, headers.getValuesAsList(header)));
|
||||
}
|
||||
|
||||
final AtomicReference<HandshakeInfo> handshakeInfo = new AtomicReference<>();
|
||||
JettyUpgradeListener jettyUpgradeListener = new JettyUpgradeListener() {
|
||||
@Override
|
||||
public void onHandshakeResponse(Request request, Response response) {
|
||||
String protocol = response.getHeaders().get(HttpHeader.SEC_WEBSOCKET_SUBPROTOCOL);
|
||||
HttpHeaders responseHeaders = new HttpHeaders();
|
||||
response.getHeaders().forEach(header -> responseHeaders.add(header.getName(), header.getValue()));
|
||||
handshakeInfo.set(new HandshakeInfo(url, responseHeaders, Mono.empty(), protocol));
|
||||
}
|
||||
};
|
||||
|
||||
Sinks.Empty<Void> completion = Sinks.empty();
|
||||
JettyWebSocketHandlerAdapter handlerAdapter = new JettyWebSocketHandlerAdapter(handler, session ->
|
||||
new JettyWebSocketSession(session, Objects.requireNonNull(handshakeInfo.get()), DefaultDataBufferFactory.sharedInstance, completion));
|
||||
try {
|
||||
this.client.connect(handlerAdapter, url, upgradeRequest, jettyUpgradeListener)
|
||||
.exceptionally(throwable -> {
|
||||
// Only fail the completion if we have an error
|
||||
// as the JettyWebSocketSession will never be opened.
|
||||
completion.tryEmitError(throwable);
|
||||
return null;
|
||||
});
|
||||
return completion.asMono();
|
||||
}
|
||||
catch (IOException ex) {
|
||||
return Mono.error(ex);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -43,6 +43,7 @@ import org.springframework.web.reactive.socket.HandshakeInfo;
|
||||
import org.springframework.web.reactive.socket.WebSocketHandler;
|
||||
import org.springframework.web.reactive.socket.server.RequestUpgradeStrategy;
|
||||
import org.springframework.web.reactive.socket.server.WebSocketService;
|
||||
import org.springframework.web.reactive.socket.server.upgrade.JettyCoreRequestUpgradeStrategy;
|
||||
import org.springframework.web.reactive.socket.server.upgrade.JettyRequestUpgradeStrategy;
|
||||
import org.springframework.web.reactive.socket.server.upgrade.ReactorNetty2RequestUpgradeStrategy;
|
||||
import org.springframework.web.reactive.socket.server.upgrade.ReactorNettyRequestUpgradeStrategy;
|
||||
@@ -76,6 +77,8 @@ public class HandshakeWebSocketService implements WebSocketService, Lifecycle {
|
||||
|
||||
private static final boolean jettyWsPresent;
|
||||
|
||||
private static final boolean jettyCoreWsPresent;
|
||||
|
||||
private static final boolean undertowWsPresent;
|
||||
|
||||
private static final boolean reactorNettyPresent;
|
||||
@@ -88,6 +91,8 @@ public class HandshakeWebSocketService implements WebSocketService, Lifecycle {
|
||||
"org.apache.tomcat.websocket.server.WsHttpUpgradeHandler", classLoader);
|
||||
jettyWsPresent = ClassUtils.isPresent(
|
||||
"org.eclipse.jetty.ee10.websocket.server.JettyWebSocketServerContainer", classLoader);
|
||||
jettyCoreWsPresent = ClassUtils.isPresent(
|
||||
"org.eclipse.jetty.websocket.server.ServerWebSocketContainer", classLoader);
|
||||
undertowWsPresent = ClassUtils.isPresent(
|
||||
"io.undertow.websockets.WebSocketProtocolHandshakeHandler", classLoader);
|
||||
reactorNettyPresent = ClassUtils.isPresent(
|
||||
@@ -278,6 +283,9 @@ public class HandshakeWebSocketService implements WebSocketService, Lifecycle {
|
||||
else if (jettyWsPresent) {
|
||||
return new JettyRequestUpgradeStrategy();
|
||||
}
|
||||
else if (jettyCoreWsPresent) {
|
||||
return new JettyCoreRequestUpgradeStrategy();
|
||||
}
|
||||
else if (undertowWsPresent) {
|
||||
return new UndertowRequestUpgradeStrategy();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,127 @@
|
||||
/*
|
||||
* Copyright 2002-2023 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* https://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.web.reactive.socket.server.upgrade;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import org.eclipse.jetty.ee10.websocket.server.JettyWebSocketServerContainer;
|
||||
import org.eclipse.jetty.server.Request;
|
||||
import org.eclipse.jetty.server.Response;
|
||||
import org.eclipse.jetty.server.Server;
|
||||
import org.eclipse.jetty.util.Callback;
|
||||
import org.eclipse.jetty.websocket.api.Configurable;
|
||||
import org.eclipse.jetty.websocket.api.exceptions.WebSocketException;
|
||||
import org.eclipse.jetty.websocket.server.ServerWebSocketContainer;
|
||||
import org.eclipse.jetty.websocket.server.WebSocketCreator;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.springframework.core.io.buffer.DataBufferFactory;
|
||||
import org.springframework.http.server.reactive.ServerHttpRequest;
|
||||
import org.springframework.http.server.reactive.ServerHttpRequestDecorator;
|
||||
import org.springframework.http.server.reactive.ServerHttpResponse;
|
||||
import org.springframework.http.server.reactive.ServerHttpResponseDecorator;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.web.reactive.socket.HandshakeInfo;
|
||||
import org.springframework.web.reactive.socket.WebSocketHandler;
|
||||
import org.springframework.web.reactive.socket.adapter.ContextWebSocketHandler;
|
||||
import org.springframework.web.reactive.socket.adapter.JettyWebSocketHandlerAdapter;
|
||||
import org.springframework.web.reactive.socket.adapter.JettyWebSocketSession;
|
||||
import org.springframework.web.reactive.socket.server.RequestUpgradeStrategy;
|
||||
import org.springframework.web.server.ServerWebExchange;
|
||||
|
||||
/**
|
||||
* A WebSocket {@code RequestUpgradeStrategy} for Jetty 12 Core.
|
||||
*
|
||||
* @author Rossen Stoyanchev
|
||||
* @since 5.3.4
|
||||
*/
|
||||
public class JettyCoreRequestUpgradeStrategy implements RequestUpgradeStrategy {
|
||||
|
||||
@Nullable
|
||||
private Consumer<Configurable> webSocketConfigurer;
|
||||
|
||||
@Nullable
|
||||
private ServerWebSocketContainer serverContainer;
|
||||
|
||||
/**
|
||||
* Add a callback to configure WebSocket server parameters on
|
||||
* {@link JettyWebSocketServerContainer}.
|
||||
* @since 6.1
|
||||
*/
|
||||
public void addWebSocketConfigurer(Consumer<Configurable> webSocketConfigurer) {
|
||||
this.webSocketConfigurer = (this.webSocketConfigurer != null ?
|
||||
this.webSocketConfigurer.andThen(webSocketConfigurer) : webSocketConfigurer);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Mono<Void> upgrade(
|
||||
ServerWebExchange exchange, WebSocketHandler handler,
|
||||
@Nullable String subProtocol, Supplier<HandshakeInfo> handshakeInfoFactory) {
|
||||
|
||||
ServerHttpRequest request = exchange.getRequest();
|
||||
ServerHttpResponse response = exchange.getResponse();
|
||||
|
||||
Request jettyRequest = ServerHttpRequestDecorator.getNativeRequest(request);
|
||||
Response jettyResponse = ServerHttpResponseDecorator.getNativeResponse(response);
|
||||
|
||||
HandshakeInfo handshakeInfo = handshakeInfoFactory.get();
|
||||
DataBufferFactory factory = response.bufferFactory();
|
||||
|
||||
// Trigger WebFlux preCommit actions before upgrade
|
||||
return exchange.getResponse().setComplete()
|
||||
.then(Mono.deferContextual(contextView -> {
|
||||
JettyWebSocketHandlerAdapter adapter = new JettyWebSocketHandlerAdapter(
|
||||
ContextWebSocketHandler.decorate(handler, contextView),
|
||||
session -> new JettyWebSocketSession(session, handshakeInfo, factory));
|
||||
|
||||
WebSocketCreator webSocketCreator = (upgradeRequest, upgradeResponse, callback) -> {
|
||||
if (subProtocol != null) {
|
||||
upgradeResponse.setAcceptedSubProtocol(subProtocol);
|
||||
}
|
||||
return adapter;
|
||||
};
|
||||
|
||||
Callback.Completable callback = new Callback.Completable();
|
||||
Mono<Void> mono = Mono.fromFuture(callback);
|
||||
ServerWebSocketContainer container = getWebSocketServerContainer(jettyRequest);
|
||||
try {
|
||||
if (!container.upgrade(webSocketCreator, jettyRequest, jettyResponse, callback)) {
|
||||
throw new WebSocketException("request could not be upgraded to websocket");
|
||||
}
|
||||
}
|
||||
catch (WebSocketException ex) {
|
||||
callback.failed(ex);
|
||||
}
|
||||
|
||||
return mono;
|
||||
}));
|
||||
}
|
||||
|
||||
private ServerWebSocketContainer getWebSocketServerContainer(Request jettyRequest) {
|
||||
if (this.serverContainer == null) {
|
||||
Server server = jettyRequest.getConnectionMetaData().getConnector().getServer();
|
||||
ServerWebSocketContainer container = ServerWebSocketContainer.get(server.getContext());
|
||||
if (this.webSocketConfigurer != null) {
|
||||
this.webSocketConfigurer.accept(container);
|
||||
}
|
||||
this.serverContainer = container;
|
||||
}
|
||||
return this.serverContainer;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -42,7 +42,7 @@ import org.springframework.web.reactive.socket.server.RequestUpgradeStrategy;
|
||||
import org.springframework.web.server.ServerWebExchange;
|
||||
|
||||
/**
|
||||
* A WebSocket {@code RequestUpgradeStrategy} for Jetty 11.
|
||||
* A WebSocket {@code RequestUpgradeStrategy} for Jetty 12 EE10.
|
||||
*
|
||||
* @author Rossen Stoyanchev
|
||||
* @since 5.3.4
|
||||
|
||||
@@ -16,7 +16,12 @@
|
||||
|
||||
package org.springframework.web.reactive.result.method.annotation;
|
||||
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import org.junit.jupiter.api.Named;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.params.ParameterizedTest;
|
||||
import org.junit.jupiter.params.provider.MethodSource;
|
||||
|
||||
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
@@ -28,10 +33,16 @@ import org.springframework.web.bind.annotation.RestController;
|
||||
import org.springframework.web.client.RestTemplate;
|
||||
import org.springframework.web.reactive.config.EnableWebFlux;
|
||||
import org.springframework.web.server.adapter.WebHttpHandlerBuilder;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.AbstractHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.HttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.JettyCoreHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.JettyHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.ReactorHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.TomcatHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.UndertowHttpServer;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.junit.jupiter.api.Named.named;
|
||||
|
||||
/**
|
||||
* Integration tests related to the use of context paths.
|
||||
@@ -40,15 +51,25 @@ import static org.assertj.core.api.Assertions.assertThat;
|
||||
*/
|
||||
class ContextPathIntegrationTests {
|
||||
|
||||
@Test
|
||||
void multipleWebFluxApps() throws Exception {
|
||||
static Stream<Named<HttpServer>> httpServers() {
|
||||
return Stream.of(
|
||||
named("Jetty", new JettyHttpServer()),
|
||||
named("Jetty Core", new JettyCoreHttpServer()),
|
||||
named("Reactor Netty", new ReactorHttpServer()),
|
||||
named("Tomcat", new TomcatHttpServer()),
|
||||
named("Undertow", new UndertowHttpServer())
|
||||
);
|
||||
}
|
||||
|
||||
@ParameterizedTest(name = "[{index}] {0}")
|
||||
@MethodSource("httpServers")
|
||||
void multipleWebFluxApps(AbstractHttpServer server) throws Exception {
|
||||
AnnotationConfigApplicationContext context1 = new AnnotationConfigApplicationContext(WebAppConfig.class);
|
||||
AnnotationConfigApplicationContext context2 = new AnnotationConfigApplicationContext(WebAppConfig.class);
|
||||
|
||||
HttpHandler webApp1Handler = WebHttpHandlerBuilder.applicationContext(context1).build();
|
||||
HttpHandler webApp2Handler = WebHttpHandlerBuilder.applicationContext(context2).build();
|
||||
|
||||
ReactorHttpServer server = new ReactorHttpServer();
|
||||
server.registerHttpHandler("/webApp1", webApp1Handler);
|
||||
server.registerHttpHandler("/webApp2", webApp2Handler);
|
||||
server.afterPropertiesSet();
|
||||
|
||||
@@ -53,6 +53,7 @@ import org.springframework.web.reactive.function.client.WebClient;
|
||||
import org.springframework.web.server.adapter.WebHttpHandlerBuilder;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.AbstractHttpHandlerIntegrationTests;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.HttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.JettyCoreHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.JettyHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.ReactorHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.TomcatHttpServer;
|
||||
@@ -127,7 +128,7 @@ class SseIntegrationTests extends AbstractHttpHandlerIntegrationTests {
|
||||
|
||||
@ParameterizedSseTest
|
||||
void sseAsEvent(HttpServer httpServer, ClientHttpConnector connector) throws Exception {
|
||||
assumeTrue(httpServer instanceof JettyHttpServer);
|
||||
assumeTrue(httpServer instanceof JettyHttpServer || httpServer instanceof JettyCoreHttpServer);
|
||||
|
||||
startServer(httpServer, connector);
|
||||
|
||||
@@ -302,18 +303,21 @@ class SseIntegrationTests extends AbstractHttpHandlerIntegrationTests {
|
||||
|
||||
static Stream<Arguments> arguments() {
|
||||
return Stream.of(
|
||||
args(new JettyHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new JettyHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new JettyHttpServer(), new HttpComponentsClientHttpConnector()),
|
||||
args(new ReactorHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new ReactorHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new ReactorHttpServer(), new HttpComponentsClientHttpConnector()),
|
||||
args(new TomcatHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new TomcatHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new TomcatHttpServer(), new HttpComponentsClientHttpConnector()),
|
||||
args(new UndertowHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new UndertowHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new UndertowHttpServer(), new HttpComponentsClientHttpConnector())
|
||||
args(new JettyHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new JettyHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new JettyHttpServer(), new HttpComponentsClientHttpConnector()),
|
||||
args(new JettyCoreHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new JettyCoreHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new JettyCoreHttpServer(), new HttpComponentsClientHttpConnector()),
|
||||
args(new ReactorHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new ReactorHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new ReactorHttpServer(), new HttpComponentsClientHttpConnector()),
|
||||
args(new TomcatHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new TomcatHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new TomcatHttpServer(), new HttpComponentsClientHttpConnector()),
|
||||
args(new UndertowHttpServer(), new ReactorClientHttpConnector()),
|
||||
args(new UndertowHttpServer(), new JettyClientHttpConnector()),
|
||||
args(new UndertowHttpServer(), new HttpComponentsClientHttpConnector())
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -45,6 +45,7 @@ import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.http.server.reactive.HttpHandler;
|
||||
import org.springframework.web.filter.reactive.ServerWebExchangeContextFilter;
|
||||
import org.springframework.web.reactive.DispatcherHandler;
|
||||
import org.springframework.web.reactive.socket.client.JettyWebSocketClient;
|
||||
import org.springframework.web.reactive.socket.client.ReactorNettyWebSocketClient;
|
||||
import org.springframework.web.reactive.socket.client.TomcatWebSocketClient;
|
||||
import org.springframework.web.reactive.socket.client.UndertowWebSocketClient;
|
||||
@@ -53,6 +54,7 @@ import org.springframework.web.reactive.socket.server.RequestUpgradeStrategy;
|
||||
import org.springframework.web.reactive.socket.server.WebSocketService;
|
||||
import org.springframework.web.reactive.socket.server.support.HandshakeWebSocketService;
|
||||
import org.springframework.web.reactive.socket.server.support.WebSocketHandlerAdapter;
|
||||
import org.springframework.web.reactive.socket.server.upgrade.JettyCoreRequestUpgradeStrategy;
|
||||
import org.springframework.web.reactive.socket.server.upgrade.JettyRequestUpgradeStrategy;
|
||||
import org.springframework.web.reactive.socket.server.upgrade.ReactorNetty2RequestUpgradeStrategy;
|
||||
import org.springframework.web.reactive.socket.server.upgrade.ReactorNettyRequestUpgradeStrategy;
|
||||
@@ -61,6 +63,7 @@ import org.springframework.web.reactive.socket.server.upgrade.UndertowRequestUpg
|
||||
import org.springframework.web.server.WebFilter;
|
||||
import org.springframework.web.server.adapter.WebHttpHandlerBuilder;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.HttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.JettyCoreHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.JettyHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.ReactorHttpServer;
|
||||
import org.springframework.web.testfixture.http.server.reactive.bootstrap.TomcatHttpServer;
|
||||
@@ -90,6 +93,7 @@ abstract class AbstractReactiveWebSocketIntegrationTests {
|
||||
|
||||
WebSocketClient[] clients = new WebSocketClient[] {
|
||||
new TomcatWebSocketClient(),
|
||||
new JettyWebSocketClient(),
|
||||
new ReactorNettyWebSocketClient(),
|
||||
new UndertowWebSocketClient(Xnio.getInstance().createWorker(OptionMap.EMPTY))
|
||||
};
|
||||
@@ -97,6 +101,7 @@ abstract class AbstractReactiveWebSocketIntegrationTests {
|
||||
Map<HttpServer, Class<?>> servers = new LinkedHashMap<>();
|
||||
servers.put(new TomcatHttpServer(TMP_DIR.getAbsolutePath(), WsContextListener.class), TomcatConfig.class);
|
||||
servers.put(new JettyHttpServer(), JettyConfig.class);
|
||||
servers.put(new JettyCoreHttpServer(), JettyCoreConfig.class);
|
||||
servers.put(new ReactorHttpServer(), ReactorNettyConfig.class);
|
||||
servers.put(new UndertowHttpServer(), UndertowConfig.class);
|
||||
|
||||
@@ -241,4 +246,12 @@ abstract class AbstractReactiveWebSocketIntegrationTests {
|
||||
}
|
||||
}
|
||||
|
||||
@Configuration
|
||||
static class JettyCoreConfig extends AbstractHandlerAdapterConfig {
|
||||
|
||||
@Override
|
||||
protected RequestUpgradeStrategy getUpgradeStrategy() {
|
||||
return new JettyCoreRequestUpgradeStrategy();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user