77 lines
2.8 KiB
Java
77 lines
2.8 KiB
Java
package io.dietz.ed.companion.service;
|
|
|
|
import java.time.Duration;
|
|
import java.util.HashMap;
|
|
import java.util.HashSet;
|
|
import java.util.Map;
|
|
import java.util.Set;
|
|
|
|
import org.springframework.http.codec.ServerSentEvent;
|
|
import org.springframework.integration.dsl.MessageChannels;
|
|
import org.springframework.messaging.Message;
|
|
import org.springframework.messaging.MessageHandler;
|
|
import org.springframework.messaging.SubscribableChannel;
|
|
import org.springframework.messaging.support.GenericMessage;
|
|
import org.springframework.stereotype.Service;
|
|
|
|
import io.dietz.ed.companion.models.UpdateFeedItem;
|
|
import reactor.core.publisher.Flux;
|
|
import reactor.core.publisher.FluxSink;
|
|
import reactor.core.publisher.FluxSink.OverflowStrategy;
|
|
|
|
@Service
|
|
public class UpdateFeedsService {
|
|
|
|
private Map<String, Set<MessageHandler>> activeFeeds = new HashMap<>();
|
|
|
|
private SubscribableChannel channel = MessageChannels.publishSubscribe().get();
|
|
|
|
public Flux<ServerSentEvent<UpdateFeedItem>> connectToFeed(String feedId) {
|
|
System.out.println("creating flux for '" + feedId + "'");
|
|
Flux<UpdateFeedItem> flux = Flux
|
|
.create((FluxSink<UpdateFeedItem> sink) -> {
|
|
MessageHandler handler = message -> sink.next(UpdateFeedItem.class.cast(message.getPayload()));
|
|
sink.onCancel(() -> {
|
|
channel.unsubscribe(handler);
|
|
removeActiveFeedHandler(feedId, handler);
|
|
|
|
System.out.println("connection '" + feedId + "' closed");
|
|
});
|
|
channel.subscribe(handler);
|
|
storeActiveFeedHandler(feedId, handler);
|
|
|
|
System.out.println("connection '" + feedId + "' opened");
|
|
}, OverflowStrategy.LATEST)
|
|
.filter(updateFeedItem -> updateFeedItem.getFeedId().equals(feedId));
|
|
|
|
return wrap(flux);
|
|
}
|
|
|
|
public boolean hasActiveFeed(String feedId) {
|
|
return activeFeeds.containsKey(feedId);
|
|
}
|
|
|
|
public void sendFeedUpdate(UpdateFeedItem updateFeedItem) {
|
|
Message<UpdateFeedItem> message = new GenericMessage<>(updateFeedItem);
|
|
channel.send(message);
|
|
}
|
|
|
|
private <T> Flux<ServerSentEvent<T>> wrap(Flux<T> fluxToWrap) {
|
|
return Flux.merge(fluxToWrap.map(t -> ServerSentEvent.builder(t).build()), Flux.interval(Duration.ofSeconds(15)).map(aLong -> ServerSentEvent.<T>builder().comment("keep alive").build()));
|
|
}
|
|
|
|
private void storeActiveFeedHandler(String feedId, MessageHandler handler) {
|
|
if(!activeFeeds.containsKey(feedId)) {
|
|
activeFeeds.put(feedId, new HashSet<>());
|
|
}
|
|
activeFeeds.get(feedId).add(handler);
|
|
}
|
|
|
|
private void removeActiveFeedHandler(String feedId, MessageHandler handler) {
|
|
activeFeeds.get(feedId).remove(handler);
|
|
if(activeFeeds.get(feedId).isEmpty()) {
|
|
activeFeeds.remove(feedId);
|
|
}
|
|
}
|
|
|
|
} |