made state ingress asynchronous

This commit is contained in:
2021-05-11 22:20:48 +00:00
parent 56bd0d4df0
commit fdbd14a2b2
4 changed files with 48 additions and 44 deletions
@@ -13,16 +13,12 @@ import org.springframework.web.bind.annotation.RestController;
import io.dietz.ed.companion.domain.State; import io.dietz.ed.companion.domain.State;
import io.dietz.ed.companion.models.UpdateStatusResponse; import io.dietz.ed.companion.models.UpdateStatusResponse;
import io.dietz.ed.companion.service.StateHandlerService; import io.dietz.ed.companion.service.StateHandlerService;
import io.dietz.ed.companion.service.UpdateFeedsService;
import reactor.core.publisher.Mono; import reactor.core.publisher.Mono;
@RestController @RestController
@RequestMapping("/update") @RequestMapping("/update")
public class StatusController { public class StatusController {
@Autowired
private UpdateFeedsService updateChannelsService;
@Autowired @Autowired
private StateHandlerService stateHandlerService; private StateHandlerService stateHandlerService;
@@ -31,24 +27,18 @@ public class StatusController {
@PostMapping("/status/{feedId}") @PostMapping("/status/{feedId}")
private Mono<UpdateStatusResponse> updateStatus(@PathVariable String feedId, @RequestBody String rawPayload) { private Mono<UpdateStatusResponse> updateStatus(@PathVariable String feedId, @RequestBody String rawPayload) {
System.out.println(rawPayload);
return Mono.create(sink -> { return Mono.create(sink -> {
if (updateChannelsService.hasActiveFeed(feedId)) { State state = new State(feedId);
State payload = handlePayload(rawPayload);
if (payload != null) { try {
stateHandlerService.handleState(feedId, payload); objectMapper.readerForUpdating(state).readValue(rawPayload);
} stateHandlerService.handleState(state);
} catch (JsonProcessingException e) {
sink.error(e);
} }
sink.success(new UpdateStatusResponse()); sink.success(new UpdateStatusResponse());
}); });
} }
private State handlePayload(String payload) {
try {
return objectMapper.readValue(payload, State.class);
} catch (JsonProcessingException e) {
e.printStackTrace();
}
return null;
}
} }
@@ -4,9 +4,19 @@ import com.fasterxml.jackson.annotation.JsonProperty;
public class State { public class State {
private final String feedId;
@JsonProperty("Credits") @JsonProperty("Credits")
private long credits; private long credits;
public State(String feedId) {
this.feedId = feedId;
}
public String getFeedId() {
return feedId;
}
public long getCredits() { public long getCredits() {
return credits; return credits;
} }
@@ -1,7 +1,5 @@
package io.dietz.ed.companion.service; package io.dietz.ed.companion.service;
import java.util.Map;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.integration.dsl.MessageChannels; import org.springframework.integration.dsl.MessageChannels;
import org.springframework.messaging.Message; 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.domain.State;
import io.dietz.ed.companion.models.UpdateFeedItem; 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 @Service
public class StateHandlerService { public class StateHandlerService {
@@ -19,26 +21,27 @@ public class StateHandlerService {
@Autowired @Autowired
private UpdateFeedsService updateFeedsService; private UpdateFeedsService updateFeedsService;
private SubscribableChannel channel = MessageChannels.publishSubscribe().get(); private SubscribableChannel channel = MessageChannels.publishSubscribe("statePublishingChannel").get();
public StateHandlerService() { public StateHandlerService() {
MessageHandler stateHandler = message -> { Flux.create((FluxSink<State> sink) -> {
String feedId = message.getHeaders().get("feedId").toString(); MessageHandler handler = message -> sink.next(State.class.cast(message.getPayload()));
sink.onCancel(() -> channel.unsubscribe(handler));
System.out.println("state of '" + feedId + "' received"); channel.subscribe(handler);
}, OverflowStrategy.LATEST)
State state = State.class.cast(message.getPayload()); .publishOn(Schedulers.boundedElastic())
.filter(s -> updateFeedsService.hasActiveFeed(s.getFeedId()))
UpdateFeedItem updateFeedItem = new UpdateFeedItem(feedId, "credits"); .map(s -> {
updateFeedItem.setContent(Long.valueOf(state.getCredits())); UpdateFeedItem updateFeedItem = new UpdateFeedItem(s.getFeedId(), "credits");
updateFeedsService.sendFeedUpdate(updateFeedItem); updateFeedItem.setContent(s.getCredits());
}; return updateFeedItem;
})
channel.subscribe(stateHandler); .subscribe((UpdateFeedItem u) -> updateFeedsService.sendFeedUpdate(u));
} }
public void handleState(String feedId, State state) { public void handleState(State state) {
Message<State> message = new GenericMessage<State>(state, Map.of("feedId", feedId)); Message<State> m = new GenericMessage<>(state);
channel.send(message); channel.send(m);
} }
} }
@@ -24,11 +24,11 @@ public class UpdateFeedsService {
private Map<String, Set<MessageHandler>> activeFeeds = new HashMap<>(); private Map<String, Set<MessageHandler>> activeFeeds = new HashMap<>();
private SubscribableChannel channel = MessageChannels.publishSubscribe().get(); private SubscribableChannel channel = MessageChannels.publishSubscribe("feedPublishingChannel").get();
public Flux<ServerSentEvent<UpdateFeedItem>> connectToFeed(String feedId) { public Flux<ServerSentEvent<UpdateFeedItem>> connectToFeed(String feedId) {
System.out.println("creating flux for '" + feedId + "'"); System.out.println("creating flux for '" + feedId + "'");
Flux<UpdateFeedItem> flux = Flux Flux<ServerSentEvent<UpdateFeedItem>> flux = Flux
.create((FluxSink<UpdateFeedItem> sink) -> { .create((FluxSink<UpdateFeedItem> sink) -> {
MessageHandler handler = message -> sink.next(UpdateFeedItem.class.cast(message.getPayload())); MessageHandler handler = message -> sink.next(UpdateFeedItem.class.cast(message.getPayload()));
sink.onCancel(() -> { sink.onCancel(() -> {
@@ -42,7 +42,8 @@ public class UpdateFeedsService {
System.out.println("connection '" + feedId + "' opened"); System.out.println("connection '" + feedId + "' opened");
}, OverflowStrategy.LATEST) }, OverflowStrategy.LATEST)
.filter(updateFeedItem -> updateFeedItem.getFeedId().equals(feedId)); .filter(updateFeedItem -> updateFeedItem.getFeedId().equals(feedId))
.map(t -> ServerSentEvent.builder(t).build());
return wrap(flux); return wrap(flux);
} }
@@ -56,8 +57,8 @@ public class UpdateFeedsService {
channel.send(message); channel.send(message);
} }
private <T> Flux<ServerSentEvent<T>> wrap(Flux<T> fluxToWrap) { private <T> Flux<ServerSentEvent<T>> wrap(Flux<ServerSentEvent<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())); return Flux.merge(fluxToWrap, Flux.interval(Duration.ofSeconds(15)).map(aLong -> ServerSentEvent.<T>builder().comment("keep alive").build()));
} }
private void storeActiveFeedHandler(String feedId, MessageHandler handler) { private void storeActiveFeedHandler(String feedId, MessageHandler handler) {