Add ServerMessageService for broadcasting player state updates, implement PlayerListState model, and update ChatService to synchronize active player states across servers.

This commit is contained in:
2026-08-22 18:22:12 +02:00
parent 423d9c9043
commit 2a65f98493
8 changed files with 164 additions and 15 deletions
@@ -52,7 +52,7 @@ public class ServerMessageController implements ServerMessageApi {
@Override @Override
public ResponseEntity<Void> sendAdminChat(String server, ServerMessageRequestDto serverMessageRequestDto) { public ResponseEntity<Void> sendAdminChat(String server, ServerMessageRequestDto serverMessageRequestDto) {
validateUser(serverMessageRequestDto.getUuid()); validateUser(serverMessageRequestDto.getUuid());
return send(server, "web_admin_chat", ChatFromWebMapper.fromDto(serverMessageRequestDto)); return send(server, "web_ac_chat", ChatFromWebMapper.fromDto(serverMessageRequestDto));
} }
@Override @Override
@@ -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.EventPublisher;
import com.alttd.altitudeweb.services.chat.event_publisher.EventUser; import com.alttd.altitudeweb.services.chat.event_publisher.EventUser;
import com.alttd.altitudeweb.services.chat.event_publisher.MessageForUser; 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 com.alttd.altitudeweb.setup.Connection;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; 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.Duration;
import java.time.Instant; import java.time.Instant;
import java.util.*; import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@Slf4j @Slf4j
@@ -35,8 +38,9 @@ public class ChatService {
private static final Duration MAX_AGE = Duration.ofHours(1); private static final Duration MAX_AGE = Duration.ofHours(1);
private final EventPublisher eventPublisher; private final EventPublisher eventPublisher;
private final ServerMessageService serverMessageService;
private final NavigableMap<Instant, ChatMessage> chatMessages = new TreeMap<>(); private final NavigableMap<Instant, ChatMessage> chatMessages = new TreeMap<>();
private final Map<String, EventUser> eventUserMap = new HashMap<>(); private final Map<String, EventUser> eventUserMap = new ConcurrentHashMap<>();
private final Map<String, ServerDto> serverStateCache = new HashMap<>(); private final Map<String, ServerDto> serverStateCache = new HashMap<>();
@Value("${chat.allowed-servers}") @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(); String key = eventUser.uuid().toString();
eventUserMap.put(key, eventUser); boolean newUser = eventUserMap.put(key, eventUser) == null;
SseEmitter emitter = eventPublisher.subscribe(key, json, this::handleSessionEnd); SseEmitter emitter = eventPublisher.subscribe(key, json, this::handleSessionEnd);
if (newUser) {
sendPlayerStateToServers();
}
sendServerStateToUser(key, eventUser); sendServerStateToUser(key, eventUser);
return emitter; 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); EventUser eventUser = eventUserMap.get(key);
if (eventUser == null) { if (eventUser == null) {
log.error("Failed to find event user for key {}", key); log.error("Failed to find event user for key {}", key);
@@ -91,6 +98,17 @@ public class ChatService {
.session_end(sessionEnd) .session_end(sessionEnd)
.build(); .build();
saveSession(chatSession); saveSession(chatSession);
if (!eventPublisher.hasSubscribers(key) && eventUserMap.remove(key, eventUser)) {
sendPlayerStateToServers();
}
}
private void sendPlayerStateToServers() {
ArrayList<UUID> activePlayers = eventUserMap.values().stream()
.map(EventUser::uuid)
.collect(Collectors.toCollection(ArrayList::new));
serverMessageService.sendMessageToAll("player_state", PlayerListStateMapper.fromList(activePlayers));
} }
private void saveSession(ChatSession chatSession) { private void saveSession(ChatSession chatSession) {
@@ -104,6 +104,11 @@ public class EventPublisher {
userEmitters.removeAll(dead); userEmitters.removeAll(dead);
} }
public boolean hasSubscribers(String uniqueKey) {
List<SseEmitter> emitters = emitterMap.get(uniqueKey);
return emitters != null && !emitters.isEmpty();
}
private long countActive() { private long countActive() {
return emitterMap.values().stream().mapToLong(List::size).sum(); return emitterMap.values().stream().mapToLong(List::size).sum();
} }
@@ -1,8 +1,6 @@
package com.alttd.altitudeweb.services.chat.to_server; 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.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.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
@@ -13,27 +11,48 @@ import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.time.Duration; import java.time.Duration;
import java.time.Instant; import java.time.Instant;
import java.util.HashSet;
import java.util.Set; import java.util.Set;
import java.util.UUID; import java.util.concurrent.ConcurrentHashMap;
@Slf4j @Slf4j
@Service @Service
@RequiredArgsConstructor @RequiredArgsConstructor
public class ServerMessageService { public class ServerMessageService {
private final Set<String> servers = new HashSet<>(); private static final String PLAYER_STATE_CHANNEL = "player_state";
private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
private final Set<String> servers = ConcurrentHashMap.newKeySet();
private final EventPublisher eventPublisher; private final EventPublisher eventPublisher;
private volatile String latestPlayerState;
public SseEmitter subscribe(String server) { public synchronized SseEmitter subscribe(String server) {
servers.add(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)); 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) { public boolean sendMessage(String server, String channel, String json) {
if (!servers.contains(server)) { if (!servers.contains(server)) {
log.warn("Server {} is not connected, not sending message", server); log.warn("Server {} is not connected, not sending message", server);
@@ -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<UUID> activePlayers;
}
@@ -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<UUID> activeViewers) {
return PlayerListState.builder()
.activePlayers(activeViewers)
.build();
}
}
@@ -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.EventPublisher;
import com.alttd.altitudeweb.services.chat.event_publisher.EventUser; import com.alttd.altitudeweb.services.chat.event_publisher.EventUser;
import com.alttd.altitudeweb.services.chat.event_publisher.MessageForUser; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
@@ -18,6 +20,7 @@ import java.util.List;
import java.util.UUID; import java.util.UUID;
import static org.junit.jupiter.api.Assertions.assertFalse; 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.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
@@ -26,11 +29,14 @@ import static org.mockito.Mockito.*;
class ChatServiceTest { class ChatServiceTest {
private ChatService chatService; private ChatService chatService;
private EventPublisher eventPublisher;
private ServerMessageService serverMessageService;
@BeforeEach @BeforeEach
void setUp() { void setUp() {
EventPublisher eventPublisher = mock(EventPublisher.class); eventPublisher = mock(EventPublisher.class);
chatService = spy(new ChatService(eventPublisher)); serverMessageService = mock(ServerMessageService.class);
chatService = spy(new ChatService(eventPublisher, serverMessageService));
ReflectionTestUtils.setField(chatService, "allowedServers", new String[]{"server1"}); 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("server1"), "HEAD_MOD should see allowed server");
assertTrue(headModResult.contains("server2"), "HEAD_MOD should see non-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<PlayerListState> stateCaptor = ArgumentCaptor.forClass(PlayerListState.class);
verify(serverMessageService, times(1)).sendMessageToAll(eq("player_state"), stateCaptor.capture());
assertEquals(List.of(uuid), stateCaptor.getValue().getActivePlayers());
}
} }
@@ -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\":[]}");
}
}