Support receipt on DISCONNECT with simple broker
Issue: SPR-14568
This commit is contained in:
@@ -34,6 +34,7 @@ import org.springframework.context.ApplicationEventPublisher;
|
||||
import org.springframework.context.ApplicationEventPublisherAware;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.simp.SimpAttributes;
|
||||
import org.springframework.messaging.simp.SimpAttributesContextHolder;
|
||||
import org.springframework.messaging.simp.SimpMessageHeaderAccessor;
|
||||
@@ -449,8 +450,15 @@ public class StompSubProtocolHandler implements SubProtocolHandler, ApplicationE
|
||||
stompAccessor = convertConnectAcktoStompConnected(stompAccessor);
|
||||
}
|
||||
else if (SimpMessageType.DISCONNECT_ACK.equals(messageType)) {
|
||||
stompAccessor = StompHeaderAccessor.create(StompCommand.ERROR);
|
||||
stompAccessor.setMessage("Session closed.");
|
||||
String receipt = getDisconnectReceipt(stompAccessor);
|
||||
if (receipt != null) {
|
||||
stompAccessor = StompHeaderAccessor.create(StompCommand.RECEIPT);
|
||||
stompAccessor.setReceiptId(receipt);
|
||||
}
|
||||
else {
|
||||
stompAccessor = StompHeaderAccessor.create(StompCommand.ERROR);
|
||||
stompAccessor.setMessage("Session closed.");
|
||||
}
|
||||
}
|
||||
else if (SimpMessageType.HEARTBEAT.equals(messageType)) {
|
||||
stompAccessor = StompHeaderAccessor.createForHeartbeat();
|
||||
@@ -503,6 +511,16 @@ public class StompSubProtocolHandler implements SubProtocolHandler, ApplicationE
|
||||
return connectedHeaders;
|
||||
}
|
||||
|
||||
private String getDisconnectReceipt(SimpMessageHeaderAccessor simpHeaders) {
|
||||
String name = StompHeaderAccessor.DISCONNECT_MESSAGE_HEADER;
|
||||
Message<?> message = (Message<?>) simpHeaders.getHeader(name);
|
||||
if (message != null) {
|
||||
StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class);
|
||||
return accessor.getReceipt();
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
protected StompHeaderAccessor toMutableAccessor(StompHeaderAccessor headerAccessor, Message<?> message) {
|
||||
return (headerAccessor.isMutable() ? headerAccessor : StompHeaderAccessor.wrap(message));
|
||||
}
|
||||
|
||||
@@ -169,6 +169,40 @@ public class StompSubProtocolHandlerTests {
|
||||
"user-name:joe\n" + "\n" + "\u0000", actual.getPayload());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void handleMessageToClientWithSimpDisconnectAck() {
|
||||
|
||||
StompHeaderAccessor accessor = StompHeaderAccessor.create(StompCommand.DISCONNECT);
|
||||
Message<?> connectMessage = MessageBuilder.createMessage(EMPTY_PAYLOAD, accessor.getMessageHeaders());
|
||||
|
||||
SimpMessageHeaderAccessor ackAccessor = SimpMessageHeaderAccessor.create(SimpMessageType.DISCONNECT_ACK);
|
||||
ackAccessor.setHeader(SimpMessageHeaderAccessor.DISCONNECT_MESSAGE_HEADER, connectMessage);
|
||||
Message<byte[]> ackMessage = MessageBuilder.createMessage(EMPTY_PAYLOAD, ackAccessor.getMessageHeaders());
|
||||
this.protocolHandler.handleMessageToClient(this.session, ackMessage);
|
||||
|
||||
assertEquals(1, this.session.getSentMessages().size());
|
||||
TextMessage actual = (TextMessage) this.session.getSentMessages().get(0);
|
||||
assertEquals("ERROR\n" + "message:Session closed.\n" + "content-length:0\n" +
|
||||
"\n\u0000", actual.getPayload());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void handleMessageToClientWithSimpDisconnectAckAndReceipt() {
|
||||
|
||||
StompHeaderAccessor accessor = StompHeaderAccessor.create(StompCommand.DISCONNECT);
|
||||
accessor.setReceipt("message-123");
|
||||
Message<?> connectMessage = MessageBuilder.createMessage(EMPTY_PAYLOAD, accessor.getMessageHeaders());
|
||||
|
||||
SimpMessageHeaderAccessor ackAccessor = SimpMessageHeaderAccessor.create(SimpMessageType.DISCONNECT_ACK);
|
||||
ackAccessor.setHeader(SimpMessageHeaderAccessor.DISCONNECT_MESSAGE_HEADER, connectMessage);
|
||||
Message<byte[]> ackMessage = MessageBuilder.createMessage(EMPTY_PAYLOAD, ackAccessor.getMessageHeaders());
|
||||
this.protocolHandler.handleMessageToClient(this.session, ackMessage);
|
||||
|
||||
assertEquals(1, this.session.getSentMessages().size());
|
||||
TextMessage actual = (TextMessage) this.session.getSentMessages().get(0);
|
||||
assertEquals("RECEIPT\n" + "receipt-id:message-123\n" + "\n\u0000", actual.getPayload());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void handleMessageToClientWithSimpHeartbeat() {
|
||||
|
||||
|
||||
Reference in New Issue
Block a user