diff --git a/src/main/java/io/dietz/ed/companion/controllers/StatusController.java b/src/main/java/io/dietz/ed/companion/controllers/StatusController.java index 056c51e..d868dfd 100644 --- a/src/main/java/io/dietz/ed/companion/controllers/StatusController.java +++ b/src/main/java/io/dietz/ed/companion/controllers/StatusController.java @@ -13,16 +13,12 @@ import org.springframework.web.bind.annotation.RestController; import io.dietz.ed.companion.domain.State; import io.dietz.ed.companion.models.UpdateStatusResponse; import io.dietz.ed.companion.service.StateHandlerService; -import io.dietz.ed.companion.service.UpdateFeedsService; import reactor.core.publisher.Mono; @RestController @RequestMapping("/update") public class StatusController { - @Autowired - private UpdateFeedsService updateChannelsService; - @Autowired private StateHandlerService stateHandlerService; @@ -31,24 +27,18 @@ public class StatusController { @PostMapping("/status/{feedId}") private Mono updateStatus(@PathVariable String feedId, @RequestBody String rawPayload) { - System.out.println(rawPayload); return Mono.create(sink -> { - if (updateChannelsService.hasActiveFeed(feedId)) { - State payload = handlePayload(rawPayload); - if (payload != null) { - stateHandlerService.handleState(feedId, payload); - } + State state = new State(feedId); + + try { + objectMapper.readerForUpdating(state).readValue(rawPayload); + stateHandlerService.handleState(state); + } catch (JsonProcessingException e) { + sink.error(e); } + sink.success(new UpdateStatusResponse()); }); } - private State handlePayload(String payload) { - try { - return objectMapper.readValue(payload, State.class); - } catch (JsonProcessingException e) { - e.printStackTrace(); - } - return null; - } } \ No newline at end of file diff --git a/src/main/java/io/dietz/ed/companion/domain/State.java b/src/main/java/io/dietz/ed/companion/domain/State.java index e263787..5c31d6a 100644 --- a/src/main/java/io/dietz/ed/companion/domain/State.java +++ b/src/main/java/io/dietz/ed/companion/domain/State.java @@ -3,10 +3,20 @@ package io.dietz.ed.companion.domain; import com.fasterxml.jackson.annotation.JsonProperty; public class State { - + + private final String feedId; + @JsonProperty("Credits") private long credits; + public State(String feedId) { + this.feedId = feedId; + } + + public String getFeedId() { + return feedId; + } + public long getCredits() { return credits; } @@ -19,5 +29,5 @@ public class State { public String toString() { return "State [credits=" + credits + "]"; } - + } \ No newline at end of file diff --git a/src/main/java/io/dietz/ed/companion/service/StateHandlerService.java b/src/main/java/io/dietz/ed/companion/service/StateHandlerService.java index fdfc019..b75132d 100644 --- a/src/main/java/io/dietz/ed/companion/service/StateHandlerService.java +++ b/src/main/java/io/dietz/ed/companion/service/StateHandlerService.java @@ -1,7 +1,5 @@ package io.dietz.ed.companion.service; -import java.util.Map; - import org.springframework.beans.factory.annotation.Autowired; import org.springframework.integration.dsl.MessageChannels; import org.springframework.messaging.Message; @@ -12,6 +10,10 @@ import org.springframework.stereotype.Service; import io.dietz.ed.companion.domain.State; import io.dietz.ed.companion.models.UpdateFeedItem; +import reactor.core.publisher.Flux; +import reactor.core.publisher.FluxSink; +import reactor.core.publisher.FluxSink.OverflowStrategy; +import reactor.core.scheduler.Schedulers; @Service public class StateHandlerService { @@ -19,26 +21,27 @@ public class StateHandlerService { @Autowired private UpdateFeedsService updateFeedsService; - private SubscribableChannel channel = MessageChannels.publishSubscribe().get(); + private SubscribableChannel channel = MessageChannels.publishSubscribe("statePublishingChannel").get(); public StateHandlerService() { - MessageHandler stateHandler = message -> { - String feedId = message.getHeaders().get("feedId").toString(); - - System.out.println("state of '" + feedId + "' received"); - - State state = State.class.cast(message.getPayload()); - - UpdateFeedItem updateFeedItem = new UpdateFeedItem(feedId, "credits"); - updateFeedItem.setContent(Long.valueOf(state.getCredits())); - updateFeedsService.sendFeedUpdate(updateFeedItem); - }; - - channel.subscribe(stateHandler); + Flux.create((FluxSink sink) -> { + MessageHandler handler = message -> sink.next(State.class.cast(message.getPayload())); + sink.onCancel(() -> channel.unsubscribe(handler)); + channel.subscribe(handler); + }, OverflowStrategy.LATEST) + .publishOn(Schedulers.boundedElastic()) + .filter(s -> updateFeedsService.hasActiveFeed(s.getFeedId())) + .map(s -> { + UpdateFeedItem updateFeedItem = new UpdateFeedItem(s.getFeedId(), "credits"); + updateFeedItem.setContent(s.getCredits()); + return updateFeedItem; + }) + .subscribe((UpdateFeedItem u) -> updateFeedsService.sendFeedUpdate(u)); } - public void handleState(String feedId, State state) { - Message message = new GenericMessage(state, Map.of("feedId", feedId)); - channel.send(message); + public void handleState(State state) { + Message m = new GenericMessage<>(state); + channel.send(m); } + } diff --git a/src/main/java/io/dietz/ed/companion/service/UpdateFeedsService.java b/src/main/java/io/dietz/ed/companion/service/UpdateFeedsService.java index 4110835..da47d66 100644 --- a/src/main/java/io/dietz/ed/companion/service/UpdateFeedsService.java +++ b/src/main/java/io/dietz/ed/companion/service/UpdateFeedsService.java @@ -24,11 +24,11 @@ public class UpdateFeedsService { private Map> activeFeeds = new HashMap<>(); - private SubscribableChannel channel = MessageChannels.publishSubscribe().get(); + private SubscribableChannel channel = MessageChannels.publishSubscribe("feedPublishingChannel").get(); public Flux> connectToFeed(String feedId) { System.out.println("creating flux for '" + feedId + "'"); - Flux flux = Flux + Flux> flux = Flux .create((FluxSink sink) -> { MessageHandler handler = message -> sink.next(UpdateFeedItem.class.cast(message.getPayload())); sink.onCancel(() -> { @@ -42,7 +42,8 @@ public class UpdateFeedsService { System.out.println("connection '" + feedId + "' opened"); }, OverflowStrategy.LATEST) - .filter(updateFeedItem -> updateFeedItem.getFeedId().equals(feedId)); + .filter(updateFeedItem -> updateFeedItem.getFeedId().equals(feedId)) + .map(t -> ServerSentEvent.builder(t).build()); return wrap(flux); } @@ -56,8 +57,8 @@ public class UpdateFeedsService { channel.send(message); } - private Flux> wrap(Flux fluxToWrap) { - return Flux.merge(fluxToWrap.map(t -> ServerSentEvent.builder(t).build()), Flux.interval(Duration.ofSeconds(15)).map(aLong -> ServerSentEvent.builder().comment("keep alive").build())); + private Flux> wrap(Flux> fluxToWrap) { + return Flux.merge(fluxToWrap, Flux.interval(Duration.ofSeconds(15)).map(aLong -> ServerSentEvent.builder().comment("keep alive").build())); } private void storeActiveFeedHandler(String feedId, MessageHandler handler) {