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");
        }
    }
}
