From b3ef834ae69c552e362fd79ae502ef19136533b9 Mon Sep 17 00:00:00 2001 From: rstoyanchev Date: Wed, 29 Apr 2026 12:10:40 +0100 Subject: [PATCH] Reliably detect broadcast messages Closes gh-36662 --- .../messaging/simp/user/UserDestinationMessageHandler.java | 6 ++++-- .../simp/stomp/StompBrokerRelayMessageHandlerTests.java | 7 ++++++- .../simp/user/UserDestinationMessageHandlerTests.java | 5 +++-- 3 files changed, 13 insertions(+), 5 deletions(-) diff --git a/spring-messaging/src/main/java/org/springframework/messaging/simp/user/UserDestinationMessageHandler.java b/spring-messaging/src/main/java/org/springframework/messaging/simp/user/UserDestinationMessageHandler.java index c220f492705..f0122d5df9a 100644 --- a/spring-messaging/src/main/java/org/springframework/messaging/simp/user/UserDestinationMessageHandler.java +++ b/spring-messaging/src/main/java/org/springframework/messaging/simp/user/UserDestinationMessageHandler.java @@ -38,6 +38,7 @@ import org.springframework.messaging.simp.SimpMessageHeaderAccessor; import org.springframework.messaging.simp.SimpMessageType; import org.springframework.messaging.simp.SimpMessagingTemplate; import org.springframework.messaging.simp.broker.OrderedMessageChannelDecorator; +import org.springframework.messaging.simp.stomp.StompBrokerRelayMessageHandler; import org.springframework.messaging.support.MessageBuilder; import org.springframework.messaging.support.MessageHeaderAccessor; import org.springframework.messaging.support.MessageHeaderInitializer; @@ -333,8 +334,9 @@ public class UserDestinationMessageHandler implements MessageHandler, SmartLifec SimpMessageHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, SimpMessageHeaderAccessor.class); Assert.state(accessor != null, "No SimpMessageHeaderAccessor"); - if (accessor.getSessionId() == null) { - // Our own broadcast + if (accessor.getSessionId() == null || + !accessor.getSessionId().equals(StompBrokerRelayMessageHandler.SYSTEM_SESSION_ID)) { + // Our own or not a broadcast return null; } destination = accessor.getFirstNativeHeader(SimpMessageHeaderAccessor.ORIGINAL_DESTINATION); diff --git a/spring-messaging/src/test/java/org/springframework/messaging/simp/stomp/StompBrokerRelayMessageHandlerTests.java b/spring-messaging/src/test/java/org/springframework/messaging/simp/stomp/StompBrokerRelayMessageHandlerTests.java index 0c5729ffe5c..eb712b655cf 100644 --- a/spring-messaging/src/test/java/org/springframework/messaging/simp/stomp/StompBrokerRelayMessageHandlerTests.java +++ b/spring-messaging/src/test/java/org/springframework/messaging/simp/stomp/StompBrokerRelayMessageHandlerTests.java @@ -245,7 +245,12 @@ class StompBrokerRelayMessageHandlerTests { ArgumentCaptor captor = ArgumentCaptor.forClass(Message.class); verify(handler).handleMessage(captor.capture()); - assertThat(captor.getValue()).isSameAs(message); + + Message actual = captor.getValue(); + assertThat(actual).isSameAs(message); + + accessor = StompHeaderAccessor.getAccessor(actual, StompHeaderAccessor.class); + assertThat(accessor.getSessionId()).isEqualTo(StompBrokerRelayMessageHandler.SYSTEM_SESSION_ID); } @Test diff --git a/spring-messaging/src/test/java/org/springframework/messaging/simp/user/UserDestinationMessageHandlerTests.java b/spring-messaging/src/test/java/org/springframework/messaging/simp/user/UserDestinationMessageHandlerTests.java index 937694fdf3d..cb6f5b03cc3 100644 --- a/spring-messaging/src/test/java/org/springframework/messaging/simp/user/UserDestinationMessageHandlerTests.java +++ b/spring-messaging/src/test/java/org/springframework/messaging/simp/user/UserDestinationMessageHandlerTests.java @@ -30,6 +30,7 @@ import org.springframework.messaging.StubMessageChannel; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.simp.SimpMessageHeaderAccessor; import org.springframework.messaging.simp.SimpMessageType; +import org.springframework.messaging.simp.stomp.StompBrokerRelayMessageHandler; import org.springframework.messaging.simp.stomp.StompCommand; import org.springframework.messaging.simp.stomp.StompHeaderAccessor; import org.springframework.messaging.support.MessageBuilder; @@ -151,7 +152,7 @@ class UserDestinationMessageHandlerTests { given(this.brokerChannel.send(Mockito.any(Message.class))).willReturn(true); StompHeaderAccessor accessor = StompHeaderAccessor.create(StompCommand.MESSAGE); - accessor.setSessionId("system123"); + accessor.setSessionId(StompBrokerRelayMessageHandler.SYSTEM_SESSION_ID); accessor.setDestination("/topic/unresolved"); accessor.setNativeHeader(ORIGINAL_DESTINATION, "/user/joe/queue/foo"); accessor.setNativeHeader("customHeader", "customHeaderValue"); @@ -175,7 +176,7 @@ class UserDestinationMessageHandlerTests { given(this.brokerChannel.send(Mockito.any(Message.class))).willReturn(true); StompHeaderAccessor accessor = StompHeaderAccessor.create(StompCommand.MESSAGE); - accessor.setSessionId("system123"); + accessor.setSessionId(StompBrokerRelayMessageHandler.SYSTEM_SESSION_ID); accessor.setDestination("/topic/unresolved"); accessor.setNativeHeader(ORIGINAL_DESTINATION, "/user/joe/queue/foo"); accessor.setLeaveMutable(true);