From 2a65f984938ca04a96bd29e89c2b1be56afb5c5d Mon Sep 17 00:00:00 2001 From: akastijn Date: Sat, 22 Aug 2026 18:22:12 +0200 Subject: [PATCH] Add `ServerMessageService` for broadcasting player state updates, implement `PlayerListState` model, and update `ChatService` to synchronize active player states across servers. --- .../chat/ServerMessageController.java | 2 +- .../services/chat/ChatService.java | 26 ++++++++-- .../chat/event_publisher/EventPublisher.java | 5 ++ .../chat/to_server/ServerMessageService.java | 35 ++++++++++--- .../chat/to_server/data/PlayerListState.java | 17 ++++++ .../mappers/PlayerListStateMapper.java | 19 +++++++ .../services/chat/ChatServiceTest.java | 23 +++++++- .../to_server/ServerMessageServiceTest.java | 52 +++++++++++++++++++ 8 files changed, 164 insertions(+), 15 deletions(-) create mode 100644 backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/data/PlayerListState.java create mode 100644 backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/mappers/PlayerListStateMapper.java create mode 100644 backend/src/test/java/com/alttd/altitudeweb/services/chat/to_server/ServerMessageServiceTest.java diff --git a/backend/src/main/java/com/alttd/altitudeweb/controllers/chat/ServerMessageController.java b/backend/src/main/java/com/alttd/altitudeweb/controllers/chat/ServerMessageController.java index 5617b35..4de3c8b 100644 --- a/backend/src/main/java/com/alttd/altitudeweb/controllers/chat/ServerMessageController.java +++ b/backend/src/main/java/com/alttd/altitudeweb/controllers/chat/ServerMessageController.java @@ -52,7 +52,7 @@ public class ServerMessageController implements ServerMessageApi { @Override public ResponseEntity sendAdminChat(String server, ServerMessageRequestDto serverMessageRequestDto) { validateUser(serverMessageRequestDto.getUuid()); - return send(server, "web_admin_chat", ChatFromWebMapper.fromDto(serverMessageRequestDto)); + return send(server, "web_ac_chat", ChatFromWebMapper.fromDto(serverMessageRequestDto)); } @Override diff --git a/backend/src/main/java/com/alttd/altitudeweb/services/chat/ChatService.java b/backend/src/main/java/com/alttd/altitudeweb/services/chat/ChatService.java index 1f24e11..e872c63 100644 --- a/backend/src/main/java/com/alttd/altitudeweb/services/chat/ChatService.java +++ b/backend/src/main/java/com/alttd/altitudeweb/services/chat/ChatService.java @@ -13,6 +13,8 @@ import com.alttd.altitudeweb.model.ServerStateDto; import com.alttd.altitudeweb.services.chat.event_publisher.EventPublisher; import com.alttd.altitudeweb.services.chat.event_publisher.EventUser; import com.alttd.altitudeweb.services.chat.event_publisher.MessageForUser; +import com.alttd.altitudeweb.services.chat.to_server.ServerMessageService; +import com.alttd.altitudeweb.services.chat.to_server.mappers.PlayerListStateMapper; import com.alttd.altitudeweb.setup.Connection; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -26,6 +28,7 @@ import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.time.Duration; import java.time.Instant; import java.util.*; +import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; @Slf4j @@ -35,8 +38,9 @@ public class ChatService { private static final Duration MAX_AGE = Duration.ofHours(1); private final EventPublisher eventPublisher; + private final ServerMessageService serverMessageService; private final NavigableMap chatMessages = new TreeMap<>(); - private final Map eventUserMap = new HashMap<>(); + private final Map eventUserMap = new ConcurrentHashMap<>(); private final Map serverStateCache = new HashMap<>(); @Value("${chat.allowed-servers}") @@ -71,15 +75,18 @@ public class ChatService { }); } - public SseEmitter subscribe(EventUser eventUser, String json) { + public synchronized SseEmitter subscribe(EventUser eventUser, String json) { String key = eventUser.uuid().toString(); - eventUserMap.put(key, eventUser); + boolean newUser = eventUserMap.put(key, eventUser) == null; SseEmitter emitter = eventPublisher.subscribe(key, json, this::handleSessionEnd); + if (newUser) { + sendPlayerStateToServers(); + } sendServerStateToUser(key, eventUser); return emitter; } - private void handleSessionEnd(String key, Instant sessionStart, Instant sessionEnd) { + private synchronized void handleSessionEnd(String key, Instant sessionStart, Instant sessionEnd) { EventUser eventUser = eventUserMap.get(key); if (eventUser == null) { log.error("Failed to find event user for key {}", key); @@ -91,6 +98,17 @@ public class ChatService { .session_end(sessionEnd) .build(); saveSession(chatSession); + + if (!eventPublisher.hasSubscribers(key) && eventUserMap.remove(key, eventUser)) { + sendPlayerStateToServers(); + } + } + + private void sendPlayerStateToServers() { + ArrayList activePlayers = eventUserMap.values().stream() + .map(EventUser::uuid) + .collect(Collectors.toCollection(ArrayList::new)); + serverMessageService.sendMessageToAll("player_state", PlayerListStateMapper.fromList(activePlayers)); } private void saveSession(ChatSession chatSession) { diff --git a/backend/src/main/java/com/alttd/altitudeweb/services/chat/event_publisher/EventPublisher.java b/backend/src/main/java/com/alttd/altitudeweb/services/chat/event_publisher/EventPublisher.java index fd9c8ec..b3ae288 100644 --- a/backend/src/main/java/com/alttd/altitudeweb/services/chat/event_publisher/EventPublisher.java +++ b/backend/src/main/java/com/alttd/altitudeweb/services/chat/event_publisher/EventPublisher.java @@ -104,6 +104,11 @@ public class EventPublisher { userEmitters.removeAll(dead); } + public boolean hasSubscribers(String uniqueKey) { + List emitters = emitterMap.get(uniqueKey); + return emitters != null && !emitters.isEmpty(); + } + private long countActive() { return emitterMap.values().stream().mapToLong(List::size).sum(); } diff --git a/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/ServerMessageService.java b/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/ServerMessageService.java index 88a11cf..3e7799c 100644 --- a/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/ServerMessageService.java +++ b/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/ServerMessageService.java @@ -1,8 +1,6 @@ package com.alttd.altitudeweb.services.chat.to_server; import com.alttd.altitudeweb.services.chat.event_publisher.EventPublisher; -import com.alttd.altitudeweb.services.chat.event_publisher.EventUser; -import com.alttd.altitudeweb.services.chat.to_server.data.ChatFromWeb; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import lombok.RequiredArgsConstructor; @@ -13,27 +11,48 @@ import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.time.Duration; import java.time.Instant; -import java.util.HashSet; import java.util.Set; -import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; @Slf4j @Service @RequiredArgsConstructor public class ServerMessageService { - private final Set servers = new HashSet<>(); + private static final String PLAYER_STATE_CHANNEL = "player_state"; + private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + private final Set servers = ConcurrentHashMap.newKeySet(); private final EventPublisher eventPublisher; + private volatile String latestPlayerState; - public SseEmitter subscribe(String server) { + public synchronized SseEmitter subscribe(String server) { servers.add(server); - return eventPublisher.subscribe(server, this::handleSessionEnd); + SseEmitter emitter = eventPublisher.subscribe(server, this::handleSessionEnd); + if (latestPlayerState != null) { + eventPublisher.sendToUser(server, PLAYER_STATE_CHANNEL, latestPlayerState); + } + return emitter; } - private void handleSessionEnd(String key, Instant sessionStart, Instant sessionEnd) { + private synchronized void handleSessionEnd(String key, Instant sessionStart, Instant sessionEnd) { log.info("Server session ended for {}. Active from {} for {}", key, sessionStart, Duration.between(sessionStart, sessionEnd)); } + public void sendMessageToAll(String channel, Object data) { + String json; + try { + json = OBJECT_MAPPER.writeValueAsString(data); + } catch (JsonProcessingException e) { + throw new IllegalStateException("Failed to serialize server message", e); + } + + if (PLAYER_STATE_CHANNEL.equals(channel)) { + latestPlayerState = json; + } + + servers.forEach(server -> eventPublisher.sendToUser(server, channel, json)); + } + public boolean sendMessage(String server, String channel, String json) { if (!servers.contains(server)) { log.warn("Server {} is not connected, not sending message", server); diff --git a/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/data/PlayerListState.java b/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/data/PlayerListState.java new file mode 100644 index 0000000..746f434 --- /dev/null +++ b/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/data/PlayerListState.java @@ -0,0 +1,17 @@ +package com.alttd.altitudeweb.services.chat.to_server.data; + +import lombok.Builder; +import lombok.Data; +import lombok.Getter; +import lombok.NoArgsConstructor; + +import java.util.List; +import java.util.UUID; + +@Builder +@Getter +public class PlayerListState { + + private List activePlayers; + +} diff --git a/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/mappers/PlayerListStateMapper.java b/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/mappers/PlayerListStateMapper.java new file mode 100644 index 0000000..078f202 --- /dev/null +++ b/backend/src/main/java/com/alttd/altitudeweb/services/chat/to_server/mappers/PlayerListStateMapper.java @@ -0,0 +1,19 @@ +package com.alttd.altitudeweb.services.chat.to_server.mappers; + +import com.alttd.altitudeweb.services.chat.to_server.data.PlayerListState; +import lombok.experimental.UtilityClass; + +import java.util.ArrayList; +import java.util.UUID; + +@UtilityClass +public class PlayerListStateMapper { + + public static PlayerListState fromList(ArrayList activeViewers) { + return PlayerListState.builder() + .activePlayers(activeViewers) + .build(); + } + + +} diff --git a/backend/src/test/java/com/alttd/altitudeweb/services/chat/ChatServiceTest.java b/backend/src/test/java/com/alttd/altitudeweb/services/chat/ChatServiceTest.java index 4440916..7ed2715 100644 --- a/backend/src/test/java/com/alttd/altitudeweb/services/chat/ChatServiceTest.java +++ b/backend/src/test/java/com/alttd/altitudeweb/services/chat/ChatServiceTest.java @@ -8,6 +8,8 @@ import com.alttd.altitudeweb.model.ServerStateDto; import com.alttd.altitudeweb.services.chat.event_publisher.EventPublisher; import com.alttd.altitudeweb.services.chat.event_publisher.EventUser; import com.alttd.altitudeweb.services.chat.event_publisher.MessageForUser; +import com.alttd.altitudeweb.services.chat.to_server.ServerMessageService; +import com.alttd.altitudeweb.services.chat.to_server.data.PlayerListState; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; @@ -18,6 +20,7 @@ import java.util.List; import java.util.UUID; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; @@ -26,11 +29,14 @@ import static org.mockito.Mockito.*; class ChatServiceTest { private ChatService chatService; + private EventPublisher eventPublisher; + private ServerMessageService serverMessageService; @BeforeEach void setUp() { - EventPublisher eventPublisher = mock(EventPublisher.class); - chatService = spy(new ChatService(eventPublisher)); + eventPublisher = mock(EventPublisher.class); + serverMessageService = mock(ServerMessageService.class); + chatService = spy(new ChatService(eventPublisher, serverMessageService)); ReflectionTestUtils.setField(chatService, "allowedServers", new String[]{"server1"}); } @@ -179,4 +185,17 @@ class ChatServiceTest { assertTrue(headModResult.contains("server1"), "HEAD_MOD should see allowed server"); assertTrue(headModResult.contains("server2"), "HEAD_MOD should see non-allowed server"); } + + @Test + void subscribingUserBroadcastsActivePlayerStateOnce() { + UUID uuid = UUID.randomUUID(); + EventUser eventUser = new EventUser(uuid, List.of(), null); + + chatService.subscribe(eventUser, "[]"); + chatService.subscribe(eventUser, "[]"); + + ArgumentCaptor stateCaptor = ArgumentCaptor.forClass(PlayerListState.class); + verify(serverMessageService, times(1)).sendMessageToAll(eq("player_state"), stateCaptor.capture()); + assertEquals(List.of(uuid), stateCaptor.getValue().getActivePlayers()); + } } diff --git a/backend/src/test/java/com/alttd/altitudeweb/services/chat/to_server/ServerMessageServiceTest.java b/backend/src/test/java/com/alttd/altitudeweb/services/chat/to_server/ServerMessageServiceTest.java new file mode 100644 index 0000000..0502886 --- /dev/null +++ b/backend/src/test/java/com/alttd/altitudeweb/services/chat/to_server/ServerMessageServiceTest.java @@ -0,0 +1,52 @@ +package com.alttd.altitudeweb.services.chat.to_server; + +import com.alttd.altitudeweb.services.chat.event_publisher.EventPublisher; +import com.alttd.altitudeweb.services.chat.to_server.data.PlayerListState; +import org.junit.jupiter.api.Test; +import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; + +import java.util.List; +import java.util.UUID; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class ServerMessageServiceTest { + + @Test + void broadcastsSerializedPlayerStateToEveryServer() { + EventPublisher eventPublisher = mock(EventPublisher.class); + when(eventPublisher.subscribe(any(), any(EventPublisher.SessionCleanupCallback.class))) + .thenReturn(new SseEmitter()); + ServerMessageService service = new ServerMessageService(eventPublisher); + UUID uuid = UUID.randomUUID(); + + service.subscribe("server-one"); + service.subscribe("server-two"); + service.sendMessageToAll("player_state", PlayerListState.builder() + .activePlayers(List.of(uuid)) + .build()); + + String expectedJson = "{\"activePlayers\":[\"" + uuid + "\"]}"; + verify(eventPublisher).sendToUser("server-one", "player_state", expectedJson); + verify(eventPublisher).sendToUser("server-two", "player_state", expectedJson); + } + + @Test + void newlySubscribedServerReceivesLatestPlayerState() { + EventPublisher eventPublisher = mock(EventPublisher.class); + when(eventPublisher.subscribe(eq("server-one"), any(EventPublisher.SessionCleanupCallback.class))) + .thenReturn(new SseEmitter()); + ServerMessageService service = new ServerMessageService(eventPublisher); + + service.sendMessageToAll("player_state", PlayerListState.builder() + .activePlayers(List.of()) + .build()); + service.subscribe("server-one"); + + verify(eventPublisher).sendToUser("server-one", "player_state", "{\"activePlayers\":[]}"); + } +}