@ -133,6 +133,13 @@ public class DefaultStompSession implements ConnectionHandlingStompSession {
@@ -133,6 +133,13 @@ public class DefaultStompSession implements ConnectionHandlingStompSession {
@ -141,7 +140,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -141,7 +140,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@ -451,7 +450,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -451,7 +450,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@ -495,7 +494,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -495,7 +494,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
thrownewMessageDeliveryException("Message broker not active. Consider subscribing to "+
"receive BrokerAvailabilityEvent's from an ApplicationListener Spring bean.");
handler.sendStompErrorFrameToClient("Broker not available.");
handler.clearConnection();
@ -562,14 +561,14 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -562,14 +561,14 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
logger.debug("Ignoring DISCONNECT in session "+sessionId+". Connection already cleaned up.");
@ -580,7 +579,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -580,7 +579,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
logger.debug("No TCP connection for session "+sessionId+" in "+message);
@ -611,7 +610,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -611,7 +610,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@ -634,11 +633,11 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -634,11 +633,11 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
Assert.notNull(sessionId,"'sessionId' must not be null");
Assert.notNull(connectHeaders,"'connectHeaders' must not be null");
this.sessionId=sessionId;
@ -662,6 +661,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -662,6 +661,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
returnthis.sessionId;
}
@Override
publicStompHeaderAccessorgetConnectHeaders(){
returnthis.connectHeaders;
}
@ -968,9 +968,9 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -968,9 +968,9 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@ -1099,7 +1099,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@@ -1099,7 +1099,7 @@ public class StompBrokerRelayMessageHandler extends AbstractBrokerMessageHandler
@ -186,7 +186,7 @@ public class ReactorNettyTcpClient<P> implements TcpOperations<P> {
@@ -186,7 +186,7 @@ public class ReactorNettyTcpClient<P> implements TcpOperations<P> {
@ -195,6 +195,19 @@ public class ReactorNettyTcpClient<P> implements TcpOperations<P> {
@@ -195,6 +195,19 @@ public class ReactorNettyTcpClient<P> implements TcpOperations<P> {
Assert.notNull(handler,"TcpConnectionHandler is required");
@ -207,7 +220,7 @@ public class ReactorNettyTcpClient<P> implements TcpOperations<P> {
@@ -207,7 +220,7 @@ public class ReactorNettyTcpClient<P> implements TcpOperations<P> {