From 56bd0d4df0b78c260b72c35eba0b6237907605b9 Mon Sep 17 00:00:00 2001 From: Steffen Dietz Date: Sat, 8 May 2021 12:32:08 +0000 Subject: [PATCH] updated spring boot removed HomeController started to decouple status ingress and feed egress --- pom.xml | 4 +- .../companion/controllers/HomeController.java | 14 ------ .../controllers/StatusController.java | 25 +++++++---- .../service/StateHandlerService.java | 44 +++++++++++++++++++ .../companion/service/UpdateFeedsService.java | 18 ++++---- 5 files changed, 71 insertions(+), 34 deletions(-) delete mode 100644 src/main/java/io/dietz/ed/companion/controllers/HomeController.java create mode 100644 src/main/java/io/dietz/ed/companion/service/StateHandlerService.java diff --git a/pom.xml b/pom.xml index ac4b6d0..25e5d16 100644 --- a/pom.xml +++ b/pom.xml @@ -5,7 +5,7 @@ org.springframework.boot spring-boot-starter-parent - 2.3.1.RELEASE + 2.4.5 @@ -34,7 +34,7 @@ org.springframework.integration - spring-integration-core + spring-integration-webflux org.postgresql diff --git a/src/main/java/io/dietz/ed/companion/controllers/HomeController.java b/src/main/java/io/dietz/ed/companion/controllers/HomeController.java deleted file mode 100644 index b057c20..0000000 --- a/src/main/java/io/dietz/ed/companion/controllers/HomeController.java +++ /dev/null @@ -1,14 +0,0 @@ -package io.dietz.ed.companion.controllers; - -import org.springframework.stereotype.Controller; -import org.springframework.ui.Model; -import org.springframework.web.bind.annotation.GetMapping; - -@Controller -public class HomeController { - - @GetMapping("/") - public String index(final Model model) { - return "index"; - } -} \ No newline at end of file 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 1bbfd06..056c51e 100644 --- a/src/main/java/io/dietz/ed/companion/controllers/StatusController.java +++ b/src/main/java/io/dietz/ed/companion/controllers/StatusController.java @@ -11,8 +11,8 @@ import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import io.dietz.ed.companion.domain.State; -import io.dietz.ed.companion.models.UpdateFeedItem; 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; @@ -23,6 +23,9 @@ public class StatusController { @Autowired private UpdateFeedsService updateChannelsService; + @Autowired + private StateHandlerService stateHandlerService; + @Autowired private ObjectMapper objectMapper; @@ -31,17 +34,21 @@ public class StatusController { System.out.println(rawPayload); return Mono.create(sink -> { if (updateChannelsService.hasActiveFeed(feedId)) { - UpdateFeedItem updateFeedItem = new UpdateFeedItem(feedId, "credits"); - State payload; - try { - payload = objectMapper.readValue(rawPayload, State.class); - updateFeedItem.setContent(Long.valueOf(payload.getCredits())); - updateChannelsService.sendFeedUpdate(updateFeedItem); - } catch (JsonProcessingException e) { - e.printStackTrace(); + State payload = handlePayload(rawPayload); + if (payload != null) { + stateHandlerService.handleState(feedId, payload); } } 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/service/StateHandlerService.java b/src/main/java/io/dietz/ed/companion/service/StateHandlerService.java new file mode 100644 index 0000000..fdfc019 --- /dev/null +++ b/src/main/java/io/dietz/ed/companion/service/StateHandlerService.java @@ -0,0 +1,44 @@ +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; +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.domain.State; +import io.dietz.ed.companion.models.UpdateFeedItem; + +@Service +public class StateHandlerService { + + @Autowired + private UpdateFeedsService updateFeedsService; + + private SubscribableChannel channel = MessageChannels.publishSubscribe().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); + } + + public void handleState(String feedId, State state) { + Message message = new GenericMessage(state, Map.of("feedId", feedId)); + channel.send(message); + } +} 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 4e14df0..4110835 100644 --- a/src/main/java/io/dietz/ed/companion/service/UpdateFeedsService.java +++ b/src/main/java/io/dietz/ed/companion/service/UpdateFeedsService.java @@ -47,6 +47,15 @@ public class UpdateFeedsService { return wrap(flux); } + public boolean hasActiveFeed(String feedId) { + return activeFeeds.containsKey(feedId); + } + + public void sendFeedUpdate(UpdateFeedItem updateFeedItem) { + Message message = new GenericMessage<>(updateFeedItem); + 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())); } @@ -65,13 +74,4 @@ public class UpdateFeedsService { } } - public boolean hasActiveFeed(String feedId) { - return activeFeeds.containsKey(feedId); - } - - public void sendFeedUpdate(UpdateFeedItem updateFeedItem) { - Message message = new GenericMessage<>(updateFeedItem); - channel.send(message); - } - } \ No newline at end of file