Remove remaining Reactor event wrapping

master
Rossen Stoyanchev 11 years ago
parent 84c55f90db
commit 547167e8b4
  1. 4
      spring-websocket/src/main/java/org/springframework/web/messaging/service/PubSubMessageService.java
  2. 3
      spring-websocket/src/main/java/org/springframework/web/messaging/stomp/support/RelayStompService.java

@ -31,8 +31,6 @@ import org.springframework.web.messaging.event.EventBus;
import org.springframework.web.messaging.event.EventConsumer;
import org.springframework.web.messaging.event.EventRegistration;
import reactor.fn.Event;
/**
* @author Rossen Stoyanchev
@ -74,7 +72,7 @@ public class PubSubMessageService extends AbstractMessageService {
byte[] payload = payloadConverter.convertToPayload(message.getPayload(), contentType);
message = new GenericMessage<byte[]>(payload, headers);
getEventBus().send(getPublishKey(message), Event.wrap(message));
getEventBus().send(getPublishKey(message), message);
}
catch (Exception ex) {
logger.error("Failed to publish " + message, ex);

@ -40,7 +40,6 @@ import org.springframework.web.messaging.stomp.StompCommand;
import org.springframework.web.messaging.stomp.StompHeaders;
import org.springframework.web.messaging.stomp.StompMessage;
import reactor.fn.Event;
import reactor.util.Assert;
@ -243,7 +242,7 @@ public class RelayStompService extends AbstractMessageService {
StompHeaders headers = new StompHeaders();
headers.setMessage(message);
StompMessage errorMessage = new StompMessage(StompCommand.ERROR, headers);
getEventBus().send(this.replyTo, Event.wrap(errorMessage));
getEventBus().send(this.replyTo, errorMessage);
}
}

Loading…
Cancel
Save