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