Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,14 @@ public class RealtimeProperties {
@Min(5)
private long heartbeatTimeoutSeconds = 60;

/**
* Interval between unsolicited server -> client heartbeats. Must stay comfortably below the
* client-side inbound liveness timeout (the Minecraft plugin force-reconnects after 75s of
* silence), so the default gives roughly three heartbeats per client window.
*/
@Min(1000)
private long serverHeartbeatIntervalMs = 25_000;

@Min(1)
private long handshakeTimeoutSeconds = 10;

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
package gg.modl.backend.realtime.schedule;

import gg.modl.backend.realtime.config.RealtimeProperties;
import gg.modl.backend.realtime.state.RealtimeConnectionRegistry;
import gg.modl.backend.realtime.state.RealtimeConnectionState;
import gg.modl.backend.realtime.transport.RealtimeCodec;
import gg.modl.backend.realtime.transport.RealtimeSessionOperations;
import lombok.RequiredArgsConstructor;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.web.socket.WebSocketSession;

/**
* Emits unsolicited heartbeats to authenticated connections so idle sessions keep seeing inbound
* traffic.
*
* <p>Clients treat prolonged inbound silence as a dead connection. Without this, a connection that
* is healthy and fully subscribed but simply has no domain events to carry looks dead to the
* client: the Minecraft plugin tears it down and reconnects roughly every 100 seconds, re-running a
* full baseline fetch each cycle. {@link RealtimeHeartbeatSweeper} is the inbound counterpart — it
* closes connections whose <em>client</em> heartbeats have stopped.</p>
*/
@Component
@RequiredArgsConstructor
public class RealtimeServerHeartbeatEmitter {
private final RealtimeProperties properties;
private final RealtimeConnectionRegistry connectionRegistry;
private final RealtimeCodec codec;
private final RealtimeSessionOperations sessionOperations;

@Scheduled(fixedDelayString = "${modl.realtime.ws.server-heartbeat-interval-ms:25000}")
public void emitHeartbeats() {
if (!properties.isEnabled()) {
return;
}

for (RealtimeConnectionRegistry.RealtimeConnectionSnapshot snapshot : connectionRegistry.snapshot()) {
RealtimeConnectionState state = snapshot.state();
// Unauthenticated sessions are the handshake sweeper's business; sending to a closing or
// terminal session would only race its close frame.
if (!state.isAuthenticated() || state.isClosing() || state.getTerminalSince() != null) {
continue;
}
WebSocketSession session = snapshot.session();
if (!session.isOpen()) {
continue;
}
// Best effort: a failed keepalive is not itself grounds for tearing down the connection.
// A genuinely dead peer stops sending client heartbeats and the sweeper closes it.
sessionOperations.trySend(session, state, codec.heartbeat(state.nextOutboundHeartbeatSequence()));
Comment thread
greptile-apps[bot] marked this conversation as resolved.
Outdated
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ public class RealtimeConnectionState {
private final Object sendLock = new Object();
private final AtomicLong deliveryAttempts = new AtomicLong();
private final AtomicLong deliveryFailures = new AtomicLong();
private final AtomicLong outboundHeartbeatSequence = new AtomicLong();
private volatile Instant lastHeartbeat = Instant.now();
private volatile RealtimePrincipal principal;
private volatile int protocolVersion;
Expand Down Expand Up @@ -110,6 +111,11 @@ public void recordHeartbeat() {
lastHeartbeat = Instant.now();
}

/** Sequence for the next unsolicited server -> client heartbeat on this connection. */
public long nextOutboundHeartbeatSequence() {
return outboundHeartbeatSequence.incrementAndGet();
}

public void setLastAcknowledgedEventId(@Nullable String lastAcknowledgedEventId) {
this.lastAcknowledgedEventId = lastAcknowledgedEventId;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import com.google.protobuf.Timestamp;
import gg.modl.backend.realtime.config.RealtimeProperties;
import gg.modl.proto.modl.v1.ErrorCode;
import gg.modl.proto.modl.v1.Heartbeat;
import gg.modl.proto.modl.v1.RealtimeEnvelope;
import gg.modl.proto.modl.v1.ReconnectAction;
import gg.modl.proto.modl.v1.ReconnectAdvice;
Expand Down Expand Up @@ -44,6 +45,26 @@ public BinaryMessage serverHello(String connectionId, Collection<Topic> accepted
return toMessage(baseEnvelope().setServerHello(hello).build());
}

/**
* Unsolicited server -> client keepalive. Clients treat prolonged inbound silence as a dead
* connection (the Minecraft plugin force-reconnects after 75s without a frame), so idle
* connections need this even when no domain events are flowing.
*
* <p>Built without an {@code event_id} on purpose: both the plugin and the panel transport-ACK
* any frame carrying one, which would turn every keepalive into a request/response pair.</p>
*/
public BinaryMessage heartbeat(long sequence) {
Instant now = Instant.now();
return toMessage(RealtimeEnvelope.newBuilder()
.setProtocolVersion(properties.getProtocolVersion())
.setTimestamp(Timestamp.newBuilder()
.setSeconds(now.getEpochSecond())
.setNanos(now.getNano())
.build())
.setHeartbeat(Heartbeat.newBuilder().setSequence(sequence))
.build());
}

public BinaryMessage error(ErrorCode code, String message) {
gg.modl.proto.modl.v1.Error error = gg.modl.proto.modl.v1.Error.newBuilder()
.setCode(code)
Expand Down
2 changes: 2 additions & 0 deletions src/main/resources/application.properties
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,8 @@ modl.realtime.ws.protocol-version=${MODL_REALTIME_WS_PROTOCOL_VERSION:1}
modl.realtime.ws.heartbeat-timeout-seconds=${MODL_REALTIME_WS_HEARTBEAT_TIMEOUT_SECONDS:60}
modl.realtime.ws.handshake-timeout-seconds=${MODL_REALTIME_WS_HANDSHAKE_TIMEOUT_SECONDS:10}
modl.realtime.ws.heartbeat-sweep-interval-ms=${MODL_REALTIME_WS_HEARTBEAT_SWEEP_INTERVAL_MS:15000}
# Must stay well under the client inbound-liveness timeout (plugin force-reconnects after 75s of silence)
modl.realtime.ws.server-heartbeat-interval-ms=${MODL_REALTIME_WS_SERVER_HEARTBEAT_INTERVAL_MS:25000}
modl.realtime.ws.inbound-rate-limit-messages=${MODL_REALTIME_WS_INBOUND_RATE_LIMIT_MESSAGES:120}
modl.realtime.ws.inbound-rate-limit-window-seconds=${MODL_REALTIME_WS_INBOUND_RATE_LIMIT_WINDOW_SECONDS:10}
modl.realtime.ws.deploy-drain-close-code=${MODL_REALTIME_WS_DEPLOY_DRAIN_CLOSE_CODE:1012}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,167 @@
package gg.modl.backend.realtime.schedule;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;

import gg.modl.backend.realtime.auth.RealtimePrincipal;
import gg.modl.backend.realtime.config.RealtimeProperties;
import gg.modl.backend.realtime.lifecycle.RealtimeConnectionCleanup;
import gg.modl.backend.realtime.metrics.RealtimeMetrics;
import gg.modl.backend.realtime.rate.RealtimeMessageRateLimiter;
import gg.modl.backend.realtime.state.RealtimeConnectionRegistry;
import gg.modl.backend.realtime.state.RealtimeConnectionState;
import gg.modl.backend.realtime.transport.RealtimeCodec;
import gg.modl.backend.realtime.transport.RealtimeSessionOperations;
import gg.modl.backend.server.data.Server;
import gg.modl.backend.server.data.ServerPlan;
import gg.modl.proto.modl.v1.RealtimeEnvelope;
import io.micrometer.core.instrument.simple.SimpleMeterRegistry;
import java.util.concurrent.ConcurrentHashMap;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import org.springframework.web.socket.BinaryMessage;
import org.springframework.web.socket.WebSocketSession;

/**
* Guards the server -> client liveness contract. The Minecraft plugin force-reconnects when it sees
* no inbound frame for 75s, so the backend must emit unsolicited heartbeats on idle connections;
* without them a healthy, fully subscribed connection is torn down roughly every 100 seconds.
*/
class RealtimeServerHeartbeatEmitterTest {

@Test
void sendsHeartbeatToAuthenticatedSession() throws Exception {
Fixture fixture = new Fixture();
WebSocketSession session = fixture.openSession("authenticated");
RealtimeConnectionState state = fixture.registry.register(session);
state.authenticate(RealtimePrincipal.minecraft(server()), 1);

fixture.emitter.emitHeartbeats();

ArgumentCaptor<BinaryMessage> captor = ArgumentCaptor.forClass(BinaryMessage.class);
verify(session).sendMessage(captor.capture());
RealtimeEnvelope envelope = RealtimeEnvelope.parseFrom(payload(captor.getValue()));
assertEquals(RealtimeEnvelope.PayloadCase.HEARTBEAT, envelope.getPayloadCase());
}

/**
* Both the plugin and the panel transport-ACK any frame carrying a non-empty event_id. An
* event_id on a heartbeat would therefore double every keepalive into a request/response pair.
*/
@Test
void heartbeatCarriesNoEventIdSoClientsDoNotAckIt() throws Exception {
Fixture fixture = new Fixture();
WebSocketSession session = fixture.openSession("no-ack");
RealtimeConnectionState state = fixture.registry.register(session);
state.authenticate(RealtimePrincipal.minecraft(server()), 1);

fixture.emitter.emitHeartbeats();

ArgumentCaptor<BinaryMessage> captor = ArgumentCaptor.forClass(BinaryMessage.class);
verify(session).sendMessage(captor.capture());
RealtimeEnvelope envelope = RealtimeEnvelope.parseFrom(payload(captor.getValue()));
assertTrue(envelope.getEventId().isEmpty(), "heartbeat must not carry an event_id");
assertEquals(1, envelope.getProtocolVersion());
}

@Test
void heartbeatSequenceAdvancesPerConnection() throws Exception {
Fixture fixture = new Fixture();
WebSocketSession session = fixture.openSession("sequenced");
RealtimeConnectionState state = fixture.registry.register(session);
state.authenticate(RealtimePrincipal.minecraft(server()), 1);

fixture.emitter.emitHeartbeats();
fixture.emitter.emitHeartbeats();

ArgumentCaptor<BinaryMessage> captor = ArgumentCaptor.forClass(BinaryMessage.class);
verify(session, org.mockito.Mockito.times(2)).sendMessage(captor.capture());
long first = RealtimeEnvelope.parseFrom(payload(captor.getAllValues().get(0))).getHeartbeat().getSequence();
long second = RealtimeEnvelope.parseFrom(payload(captor.getAllValues().get(1))).getHeartbeat().getSequence();
assertEquals(1L, first);
assertEquals(2L, second);
}

@Test
void doesNotHeartbeatUnauthenticatedSession() throws Exception {
Fixture fixture = new Fixture();
WebSocketSession session = fixture.openSession("unauthenticated");
fixture.registry.register(session);

fixture.emitter.emitHeartbeats();

verify(session, never()).sendMessage(any(BinaryMessage.class));
}

@Test
void doesNotHeartbeatTerminalSession() throws Exception {
Fixture fixture = new Fixture();
WebSocketSession session = fixture.openSession("terminal");
RealtimeConnectionState state = fixture.registry.register(session);
state.authenticate(RealtimePrincipal.minecraft(server()), 1);
state.markClosing();
state.markTerminal();

fixture.emitter.emitHeartbeats();

verify(session, never()).sendMessage(any(BinaryMessage.class));
}

@Test
void doesNothingWhenRealtimeDisabled() throws Exception {
Fixture fixture = new Fixture();
fixture.properties.setEnabled(false);
WebSocketSession session = fixture.openSession("disabled");
RealtimeConnectionState state = fixture.registry.register(session);
state.authenticate(RealtimePrincipal.minecraft(server()), 1);

fixture.emitter.emitHeartbeats();

verify(session, never()).sendMessage(any(BinaryMessage.class));
}

private static final class Fixture {
private final RealtimeProperties properties = new RealtimeProperties();
private final RealtimeConnectionRegistry registry;
private final RealtimeServerHeartbeatEmitter emitter;

private Fixture() {
properties.setEnabled(true);
registry = new RealtimeConnectionRegistry(properties);
RealtimeMetrics metrics = new RealtimeMetrics(new SimpleMeterRegistry());
RealtimeConnectionCleanup cleanup =
new RealtimeConnectionCleanup(registry, new RealtimeMessageRateLimiter(properties), metrics);
emitter = new RealtimeServerHeartbeatEmitter(
properties,
registry,
new RealtimeCodec(properties),
new RealtimeSessionOperations(registry, cleanup, metrics)
);
}

private WebSocketSession openSession(String id) {
WebSocketSession session = mock(WebSocketSession.class);
when(session.getId()).thenReturn(id);
when(session.isOpen()).thenReturn(true);
when(session.getAttributes()).thenReturn(new ConcurrentHashMap<>());
return session;
}
}

private static byte[] payload(BinaryMessage message) {
byte[] payload = new byte[message.getPayloadLength()];
message.getPayload().get(payload);
return payload;
}

private static Server server() {
Server server = new Server("server", "server", "server_db", "admin@example.com", true, ServerPlan.FREE);
server.setId("server-id");
return server;
}
}
Loading