From 69ef364ef9b874db3c88f7a3f84f6ab5182d22c7 Mon Sep 17 00:00:00 2001 From: Rossen Stoyanchev Date: Tue, 28 May 2013 12:03:04 -0400 Subject: [PATCH] Introduce messaging package org.springframework.web.stomp is now org.springframework.web.messaging.stomp Also classes in the ~.stomp.server and ~.stomp.adapter packages have been renamed. --- .../{ => messaging}/stomp/StompCommand.java | 2 +- .../{ => messaging}/stomp/StompException.java | 2 +- .../{ => messaging}/stomp/StompHeaders.java | 2 +- .../{ => messaging}/stomp/StompMessage.java | 2 +- .../{ => messaging}/stomp/StompSession.java | 11 +++-- .../stomp/adapter/StompMessageHandler.java} | 12 +++--- .../stomp/adapter/StompWebSocketHandler.java | 27 ++++++------ .../stomp/adapter/WebSocketStompSession.java | 29 ++++++++++--- .../stomp/server/RelayStompService.java} | 24 +++++------ .../server/ServerStompMessageHandler.java} | 42 +++++++++++-------- .../stomp/server/SimpleStompService.java} | 22 +++++----- .../stomp/support/StompMessageConverter.java | 10 ++--- 12 files changed, 106 insertions(+), 79 deletions(-) rename spring-websocket/src/main/java/org/springframework/web/{ => messaging}/stomp/StompCommand.java (94%) rename spring-websocket/src/main/java/org/springframework/web/{ => messaging}/stomp/StompException.java (95%) rename spring-websocket/src/main/java/org/springframework/web/{ => messaging}/stomp/StompHeaders.java (99%) rename spring-websocket/src/main/java/org/springframework/web/{ => messaging}/stomp/StompMessage.java (97%) rename spring-websocket/src/main/java/org/springframework/web/{ => messaging}/stomp/StompSession.java (81%) rename spring-websocket/src/main/java/org/springframework/web/{stomp/adapter/StompMessageProcessor.java => messaging/stomp/adapter/StompMessageHandler.java} (67%) rename spring-websocket/src/main/java/org/springframework/web/{ => messaging}/stomp/adapter/StompWebSocketHandler.java (75%) rename spring-websocket/src/main/java/org/springframework/web/{ => messaging}/stomp/adapter/WebSocketStompSession.java (72%) rename spring-websocket/src/main/java/org/springframework/web/{stomp/server/RelayStompReactorService.java => messaging/stomp/server/RelayStompService.java} (88%) rename spring-websocket/src/main/java/org/springframework/web/{stomp/server/ReactorServerStompMessageProcessor.java => messaging/stomp/server/ServerStompMessageHandler.java} (86%) rename spring-websocket/src/main/java/org/springframework/web/{stomp/server/SimpleStompReactorService.java => messaging/stomp/server/SimpleStompService.java} (84%) rename spring-websocket/src/main/java/org/springframework/web/{ => messaging}/stomp/support/StompMessageConverter.java (95%) diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/StompCommand.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompCommand.java similarity index 94% rename from spring-websocket/src/main/java/org/springframework/web/stomp/StompCommand.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompCommand.java index 752933b0bd..83edcb9011 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/StompCommand.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompCommand.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.web.stomp; +package org.springframework.web.messaging.stomp; /** diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/StompException.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompException.java similarity index 95% rename from spring-websocket/src/main/java/org/springframework/web/stomp/StompException.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompException.java index 08ba53b6b0..7f2915aa9e 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/StompException.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompException.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.web.stomp; +package org.springframework.web.messaging.stomp; import org.springframework.core.NestedRuntimeException; diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/StompHeaders.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompHeaders.java similarity index 99% rename from spring-websocket/src/main/java/org/springframework/web/stomp/StompHeaders.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompHeaders.java index 206bc2f147..798418ecad 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/StompHeaders.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompHeaders.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.web.stomp; +package org.springframework.web.messaging.stomp; import java.io.Serializable; import java.util.Collection; diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/StompMessage.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompMessage.java similarity index 97% rename from spring-websocket/src/main/java/org/springframework/web/stomp/StompMessage.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompMessage.java index 96a85e4ccc..6c8da5e3e7 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/StompMessage.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompMessage.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.web.stomp; +package org.springframework.web.messaging.stomp; import java.nio.charset.Charset; diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/StompSession.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompSession.java similarity index 81% rename from spring-websocket/src/main/java/org/springframework/web/stomp/StompSession.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompSession.java index b334a0e60b..9ddc8cf2ec 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/StompSession.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/StompSession.java @@ -13,13 +13,12 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.web.stomp; +package org.springframework.web.messaging.stomp; import java.io.IOException; /** - * * @author Rossen Stoyanchev * @since 4.0 */ @@ -28,9 +27,15 @@ public interface StompSession { String getId(); /** + * TODO... + *

* If the message is a STOMP ERROR message, the session will also be closed. - * */ void sendMessage(StompMessage message) throws IOException; + /** + * Register a task to be invoked if the underlying connection is closed. + */ + void registerConnectionClosedCallback(Runnable task); + } diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/adapter/StompMessageProcessor.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/StompMessageHandler.java similarity index 67% rename from spring-websocket/src/main/java/org/springframework/web/stomp/adapter/StompMessageProcessor.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/StompMessageHandler.java index 8d2af8e888..60f9a1a744 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/adapter/StompMessageProcessor.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/StompMessageHandler.java @@ -14,20 +14,18 @@ * limitations under the License. */ -package org.springframework.web.stomp.adapter; +package org.springframework.web.messaging.stomp.adapter; -import org.springframework.web.stomp.StompMessage; -import org.springframework.web.stomp.StompSession; +import org.springframework.web.messaging.stomp.StompMessage; +import org.springframework.web.messaging.stomp.StompSession; /** * @author Rossen Stoyanchev * @since 4.0 */ -public interface StompMessageProcessor { +public interface StompMessageHandler { - void processMessage(StompSession stompSession, StompMessage message); - - void processConnectionClosed(StompSession stompSession); + void handleMessage(StompSession stompSession, StompMessage message); } diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/adapter/StompWebSocketHandler.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/StompWebSocketHandler.java similarity index 75% rename from spring-websocket/src/main/java/org/springframework/web/stomp/adapter/StompWebSocketHandler.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/StompWebSocketHandler.java index 837ae9bf57..b12ab63d75 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/adapter/StompWebSocketHandler.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/StompWebSocketHandler.java @@ -13,39 +13,38 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.web.stomp.adapter; +package org.springframework.web.messaging.stomp.adapter; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import org.springframework.util.Assert; +import org.springframework.web.messaging.stomp.StompCommand; +import org.springframework.web.messaging.stomp.StompHeaders; +import org.springframework.web.messaging.stomp.StompMessage; +import org.springframework.web.messaging.stomp.StompSession; +import org.springframework.web.messaging.stomp.support.StompMessageConverter; import org.springframework.web.socket.CloseStatus; import org.springframework.web.socket.TextMessage; import org.springframework.web.socket.WebSocketSession; import org.springframework.web.socket.adapter.TextWebSocketHandlerAdapter; -import org.springframework.web.stomp.StompCommand; -import org.springframework.web.stomp.StompHeaders; -import org.springframework.web.stomp.StompMessage; -import org.springframework.web.stomp.StompSession; -import org.springframework.web.stomp.support.StompMessageConverter; /** - * * @author Rossen Stoyanchev * @since 4.0 */ public class StompWebSocketHandler extends TextWebSocketHandlerAdapter { - private final StompMessageProcessor messageProcessor; + private final StompMessageHandler messageHandler; private final StompMessageConverter messageConverter = new StompMessageConverter(); - private final Map sessions = new ConcurrentHashMap(); + private final Map sessions = new ConcurrentHashMap(); - public StompWebSocketHandler(StompMessageProcessor messageProcessor) { - this.messageProcessor = messageProcessor; + public StompWebSocketHandler(StompMessageHandler messageHandler) { + this.messageHandler = messageHandler; } @@ -68,7 +67,7 @@ public class StompWebSocketHandler extends TextWebSocketHandlerAdapter { // TODO: validate size limits // http://stomp.github.io/stomp-specification-1.2.html#Size_Limits - this.messageProcessor.processMessage(stompSession, stompMessage); + this.messageHandler.handleMessage(stompSession, stompMessage); // TODO: send RECEIPT message if incoming message has "receipt" header // http://stomp.github.io/stomp-specification-1.2.html#Header_receipt @@ -89,9 +88,9 @@ public class StompWebSocketHandler extends TextWebSocketHandlerAdapter { @Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { - StompSession stompSession = this.sessions.remove(session.getId()); + WebSocketStompSession stompSession = this.sessions.remove(session.getId()); if (stompSession != null) { - this.messageProcessor.processConnectionClosed(stompSession); + stompSession.handleConnectionClosed(); } } diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/adapter/WebSocketStompSession.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/WebSocketStompSession.java similarity index 72% rename from spring-websocket/src/main/java/org/springframework/web/stomp/adapter/WebSocketStompSession.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/WebSocketStompSession.java index fa71c7b772..8af38ecd6c 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/adapter/WebSocketStompSession.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/adapter/WebSocketStompSession.java @@ -14,18 +14,20 @@ * limitations under the License. */ -package org.springframework.web.stomp.adapter; +package org.springframework.web.messaging.stomp.adapter; import java.io.IOException; +import java.util.ArrayList; +import java.util.List; import org.springframework.util.Assert; +import org.springframework.web.messaging.stomp.StompCommand; +import org.springframework.web.messaging.stomp.StompMessage; +import org.springframework.web.messaging.stomp.StompSession; +import org.springframework.web.messaging.stomp.support.StompMessageConverter; import org.springframework.web.socket.CloseStatus; import org.springframework.web.socket.TextMessage; import org.springframework.web.socket.WebSocketSession; -import org.springframework.web.stomp.StompCommand; -import org.springframework.web.stomp.StompMessage; -import org.springframework.web.stomp.StompSession; -import org.springframework.web.stomp.support.StompMessageConverter; /** @@ -40,6 +42,8 @@ public class WebSocketStompSession implements StompSession { private final StompMessageConverter messageConverter; + private final List connectionClosedTasks = new ArrayList(); + public WebSocketStompSession(WebSocketSession webSocketSession, StompMessageConverter messageConverter) { Assert.notNull(webSocketSession, "webSocketSession is required"); @@ -70,4 +74,19 @@ public class WebSocketStompSession implements StompSession { } } + public void registerConnectionClosedCallback(Runnable task) { + this.connectionClosedTasks.add(task); + } + + public void handleConnectionClosed() { + for (Runnable task : this.connectionClosedTasks) { + try { + task.run(); + } + catch (Throwable t) { + // ignore + } + } + } + } \ No newline at end of file diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/server/RelayStompReactorService.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/RelayStompService.java similarity index 88% rename from spring-websocket/src/main/java/org/springframework/web/stomp/server/RelayStompReactorService.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/RelayStompService.java index 97274b59cf..697e279814 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/server/RelayStompReactorService.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/RelayStompService.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.web.stomp.server; +package org.springframework.web.messaging.stomp.server; import java.io.BufferedInputStream; import java.io.BufferedOutputStream; @@ -31,10 +31,10 @@ import javax.net.SocketFactory; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.core.task.TaskExecutor; -import org.springframework.web.stomp.StompCommand; -import org.springframework.web.stomp.StompHeaders; -import org.springframework.web.stomp.StompMessage; -import org.springframework.web.stomp.support.StompMessageConverter; +import org.springframework.web.messaging.stomp.StompCommand; +import org.springframework.web.messaging.stomp.StompHeaders; +import org.springframework.web.messaging.stomp.StompMessage; +import org.springframework.web.messaging.stomp.support.StompMessageConverter; import reactor.Fn; import reactor.core.Reactor; @@ -47,9 +47,9 @@ import reactor.util.Assert; * @author Rossen Stoyanchev * @since 4.0 */ -public class RelayStompReactorService { +public class RelayStompService { - private static final Log logger = LogFactory.getLog(RelayStompReactorService.class); + private static final Log logger = LogFactory.getLog(RelayStompService.class); private final Reactor reactor; @@ -61,7 +61,7 @@ public class RelayStompReactorService { private final TaskExecutor taskExecutor; - public RelayStompReactorService(Reactor reactor, TaskExecutor executor) { + public RelayStompService(Reactor reactor, TaskExecutor executor) { this.reactor = reactor; this.taskExecutor = executor; // For now, a naively way to manage socket reading @@ -91,7 +91,7 @@ public class RelayStompReactorService { } private RelaySession getRelaySession(String stompSessionId) { - RelaySession session = RelayStompReactorService.this.relaySessions.get(stompSessionId); + RelaySession session = RelayStompService.this.relaySessions.get(stompSessionId); Assert.notNull(session, "RelaySession not found"); return session; } @@ -188,8 +188,8 @@ public class RelayStompReactorService { } else if (b == 0x00) { byte[] bytes = out.toByteArray(); - StompMessage message = RelayStompReactorService.this.converter.toStompMessage(bytes); - RelayStompReactorService.this.reactor.notify(replyTo, Fn.event(message)); + StompMessage message = RelayStompService.this.converter.toStompMessage(bytes); + RelayStompService.this.reactor.notify(replyTo, Fn.event(message)); out.reset(); } else { @@ -209,7 +209,7 @@ public class RelayStompReactorService { StompHeaders headers = new StompHeaders(); headers.setMessage("Lost connection"); StompMessage errorMessage = new StompMessage(StompCommand.ERROR, headers); - RelayStompReactorService.this.reactor.notify(replyTo, Fn.event(errorMessage)); + RelayStompService.this.reactor.notify(replyTo, Fn.event(errorMessage)); } } diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/server/ReactorServerStompMessageProcessor.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/ServerStompMessageHandler.java similarity index 86% rename from spring-websocket/src/main/java/org/springframework/web/stomp/server/ReactorServerStompMessageProcessor.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/ServerStompMessageHandler.java index cabc8ec7f5..9a1210e45b 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/server/ReactorServerStompMessageProcessor.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/ServerStompMessageHandler.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.web.stomp.server; +package org.springframework.web.messaging.stomp.server; import java.io.IOException; import java.util.ArrayList; @@ -25,12 +25,12 @@ import java.util.concurrent.ConcurrentHashMap; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.util.CollectionUtils; -import org.springframework.web.stomp.StompCommand; -import org.springframework.web.stomp.StompException; -import org.springframework.web.stomp.StompHeaders; -import org.springframework.web.stomp.StompMessage; -import org.springframework.web.stomp.StompSession; -import org.springframework.web.stomp.adapter.StompMessageProcessor; +import org.springframework.web.messaging.stomp.StompCommand; +import org.springframework.web.messaging.stomp.StompException; +import org.springframework.web.messaging.stomp.StompHeaders; +import org.springframework.web.messaging.stomp.StompMessage; +import org.springframework.web.messaging.stomp.StompSession; +import org.springframework.web.messaging.stomp.adapter.StompMessageHandler; import reactor.Fn; import reactor.core.Reactor; @@ -43,24 +43,26 @@ import reactor.fn.Registration; * @author Rossen Stoyanchev * @since 4.0 */ -public class ReactorServerStompMessageProcessor implements StompMessageProcessor { +public class ServerStompMessageHandler implements StompMessageHandler { - private static Log logger = LogFactory.getLog(ReactorServerStompMessageProcessor.class); + private static Log logger = LogFactory.getLog(ServerStompMessageHandler.class); private final Reactor reactor; - private Map>> registrationsBySession = new ConcurrentHashMap>>(); + private Map>> registrationsBySession = + new ConcurrentHashMap>>(); - public ReactorServerStompMessageProcessor(Reactor reactor) { + public ServerStompMessageHandler(Reactor reactor) { this.reactor = reactor; } - public void processMessage(StompSession session, StompMessage message) { + public void handleMessage(StompSession session, StompMessage message) { try { StompCommand command = message.getCommand(); if (StompCommand.CONNECT.equals(command) || StompCommand.STOMP.equals(command)) { + registerConnectionClosedCallback(session); connect(session, message); } else if (StompCommand.SUBSCRIBE.equals(command)) { @@ -92,6 +94,16 @@ public class ReactorServerStompMessageProcessor implements StompMessageProcessor } } + private void registerConnectionClosedCallback(final StompSession session) { + session.registerConnectionClosedCallback(new Runnable() { + @Override + public void run() { + removeSubscriptions(session); + reactor.notify("CONNECTION_CLOSED", Fn.event(session.getId())); + } + }); + } + private void handleError(final StompSession session, Throwable t) { logger.error("Terminating STOMP session due to failure to send message: ", t); sendErrorMessage(session, t.getMessage()); @@ -233,10 +245,4 @@ public class ReactorServerStompMessageProcessor implements StompMessageProcessor return true; } - @Override - public void processConnectionClosed(StompSession session) { - removeSubscriptions(session); - this.reactor.notify("CONNECTION_CLOSED", Fn.event(session.getId())); - } - } diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/server/SimpleStompReactorService.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/SimpleStompService.java similarity index 84% rename from spring-websocket/src/main/java/org/springframework/web/stomp/server/SimpleStompReactorService.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/SimpleStompService.java index 294962d143..48caf1f87d 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/server/SimpleStompReactorService.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/server/SimpleStompService.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.web.stomp.server; +package org.springframework.web.messaging.stomp.server; import java.util.ArrayList; import java.util.List; @@ -23,9 +23,9 @@ import java.util.concurrent.ConcurrentHashMap; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.springframework.web.stomp.StompCommand; -import org.springframework.web.stomp.StompHeaders; -import org.springframework.web.stomp.StompMessage; +import org.springframework.web.messaging.stomp.StompCommand; +import org.springframework.web.messaging.stomp.StompHeaders; +import org.springframework.web.messaging.stomp.StompMessage; import reactor.Fn; import reactor.core.Reactor; @@ -38,16 +38,16 @@ import reactor.fn.Registration; * @author Rossen Stoyanchev * @since 4.0 */ -public class SimpleStompReactorService { +public class SimpleStompService { - private static final Log logger = LogFactory.getLog(SimpleStompReactorService.class); + private static final Log logger = LogFactory.getLog(SimpleStompService.class); private final Reactor reactor; private Map>> subscriptionsBySession = new ConcurrentHashMap>>(); - public SimpleStompReactorService(Reactor reactor) { + public SimpleStompService(Reactor reactor) { this.reactor = reactor; this.reactor.on(Fn.$(StompCommand.SUBSCRIBE), new SubscribeConsumer()); this.reactor.on(Fn.$(StompCommand.SEND), new SendConsumer()); @@ -85,7 +85,7 @@ public class SimpleStompReactorService { logger.debug("Subscribe " + message); } - Registration registration = SimpleStompReactorService.this.reactor.on( + Registration registration = SimpleStompService.this.reactor.on( Fn.$("destination:" + message.getHeaders().getDestination()), new Consumer>() { @Override @@ -94,7 +94,7 @@ public class SimpleStompReactorService { StompHeaders headers = new StompHeaders(); headers.setDestination(inMessage.getHeaders().getDestination()); StompMessage outMessage = new StompMessage(StompCommand.MESSAGE, headers, inMessage.getPayload()); - SimpleStompReactorService.this.reactor.notify(event.getReplyTo(), Fn.event(outMessage)); + SimpleStompService.this.reactor.notify(event.getReplyTo(), Fn.event(outMessage)); } }); @@ -110,7 +110,7 @@ public class SimpleStompReactorService { logger.debug("Message received: " + message); String destination = message.getHeaders().getDestination(); - SimpleStompReactorService.this.reactor.notify("destination:" + destination, Fn.event(message)); + SimpleStompService.this.reactor.notify("destination:" + destination, Fn.event(message)); } } @@ -119,7 +119,7 @@ public class SimpleStompReactorService { @Override public void accept(Event event) { String sessionId = event.getData(); - SimpleStompReactorService.this.removeSubscriptions(sessionId); + SimpleStompService.this.removeSubscriptions(sessionId); } } diff --git a/spring-websocket/src/main/java/org/springframework/web/stomp/support/StompMessageConverter.java b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/support/StompMessageConverter.java similarity index 95% rename from spring-websocket/src/main/java/org/springframework/web/stomp/support/StompMessageConverter.java rename to spring-websocket/src/main/java/org/springframework/web/messaging/stomp/support/StompMessageConverter.java index edccc5cb26..b120459309 100644 --- a/spring-websocket/src/main/java/org/springframework/web/stomp/support/StompMessageConverter.java +++ b/spring-websocket/src/main/java/org/springframework/web/messaging/stomp/support/StompMessageConverter.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.web.stomp.support; +package org.springframework.web.messaging.stomp.support; import java.io.ByteArrayOutputStream; import java.io.IOException; @@ -21,10 +21,10 @@ import java.util.List; import java.util.Map.Entry; import org.springframework.util.Assert; -import org.springframework.web.stomp.StompCommand; -import org.springframework.web.stomp.StompException; -import org.springframework.web.stomp.StompHeaders; -import org.springframework.web.stomp.StompMessage; +import org.springframework.web.messaging.stomp.StompCommand; +import org.springframework.web.messaging.stomp.StompException; +import org.springframework.web.messaging.stomp.StompHeaders; +import org.springframework.web.messaging.stomp.StompMessage; /** * @author Gary Russell