initial import
This commit is contained in:
@@ -0,0 +1,77 @@
|
||||
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);
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
public boolean hasActiveFeed(String feedId) {
|
||||
return activeFeeds.containsKey(feedId);
|
||||
}
|
||||
|
||||
public void sendFeedUpdate(UpdateFeedItem updateFeedItem) {
|
||||
Message<UpdateFeedItem> message = new GenericMessage<>(updateFeedItem);
|
||||
channel.send(message);
|
||||
}
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user