package com.xfestudio.xfeservermanager.infra.http;
import static org.assertj.core.api.Assertions.assertThat;
import java.time.Duration;
import org.junit.jupiter.api.Test;
class SseBrokerTest {
@Test
void replaysEventsAfterKnownIdAndHonorsCapacity() throws Exception {
var broker = new SseBroker(1);
var first = broker.publish("status", "{\"tps\":20}");
broker.publish("players", "{\"online\":2}");
var subscription = broker.subscribe(first.id()).orElseThrow();
assertThat(subscription.poll(Duration.ofMillis(10)).orElseThrow().type()).isEqualTo("players");
assertThat(broker.subscribe(0)).isEmpty();
subscription.close();
assertThat(broker.subscriberCount()).isZero();
assertThat(broker.subscribe(0)).isPresent();
}
@Test
void encodedEventCannotInjectFieldsThroughNewlines() {
var broker = new SseBroker(1);
String encoded = broker.publish("audit", "{\n\"ok\":true\n}").encode();
assertThat(encoded).contains("\ndata: ");
assertThat(encoded).doesNotContain("\nevent: ok");
}
@Test
void requestsResyncInsteadOfSilentlyTruncatingAnOversizedReplay() throws Exception {
var broker = new SseBroker(1);
for (int index = 0; index < 200; index++) {
broker.publish("operation", "{\"index\":" + index + "}");
}
try (var subscription = broker.subscribe(0).orElseThrow()) {
assertThat(subscription.poll(java.time.Duration.ZERO))
.get()
.extracting(SseBroker.Event::type)
.isEqualTo("resync-required");
}
}
}
package com.xfestudio.xfeservermanager.infra.http;
import static org.assertj.core.api.Assertions.assertThat;
import java.time.Duration;
import org.junit.jupiter.api.Test;
class SseBrokerTest {
@Test
void replaysEventsAfterKnownIdAndHonorsCapacity() throws Exception {
var broker = new SseBroker(1);
var first = broker.publish("status", "{\"tps\":20}");
broker.publish("players", "{\"online\":2}");
var subscription = broker.subscribe(first.id()).orElseThrow();
assertThat(subscription.poll(Duration.ofMillis(10)).orElseThrow().type()).isEqualTo("players");
assertThat(broker.subscribe(0)).isEmpty();
subscription.close();
assertThat(broker.subscriberCount()).isZero();
assertThat(broker.subscribe(0)).isPresent();
}
@Test
void encodedEventCannotInjectFieldsThroughNewlines() {
var broker = new SseBroker(1);
String encoded = broker.publish("audit", "{\n\"ok\":true\n}").encode();
assertThat(encoded).contains("\ndata: ");
assertThat(encoded).doesNotContain("\nevent: ok");
}
@Test
void requestsResyncInsteadOfSilentlyTruncatingAnOversizedReplay() throws Exception {
var broker = new SseBroker(1);
for (int index = 0; index < 200; index++) {
broker.publish("operation", "{\"index\":" + index + "}");
}
try (var subscription = broker.subscribe(0).orElseThrow()) {
assertThat(subscription.poll(java.time.Duration.ZERO))
.get()
.extracting(SseBroker.Event::type)
.isEqualTo("resync-required");
}
}
}