X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2Fchat%2Fbackend%2Fdomain%2FChatRoom.java;h=35d0c3d42afe5263466c112c2ed42ecb795fa8c8;hb=aa0efd1151673c5f0f1576c3026f6fdd0dfad691;hp=0fdea33e5375efb3ed14ab16124a97fa32f57af6;hpb=cfda873368d7b3fdb4869fbce98a0d6e8ca69ab7;p=demos%2Fkafka%2Fchat diff --git a/src/main/java/de/juplo/kafka/chat/backend/domain/ChatRoom.java b/src/main/java/de/juplo/kafka/chat/backend/domain/ChatRoom.java index 0fdea33e..35d0c3d4 100644 --- a/src/main/java/de/juplo/kafka/chat/backend/domain/ChatRoom.java +++ b/src/main/java/de/juplo/kafka/chat/backend/domain/ChatRoom.java @@ -1,22 +1,28 @@ package de.juplo.kafka.chat.backend.domain; +import lombok.EqualsAndHashCode; import lombok.Getter; +import lombok.ToString; import lombok.extern.slf4j.Slf4j; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.Sinks; +import java.time.Clock; import java.time.LocalDateTime; import java.util.*; @Slf4j +@EqualsAndHashCode(of = { "id" }) +@ToString(of = { "id", "name" }) public class ChatRoom { @Getter private final UUID id; @Getter private final String name; + private final Clock clock; private final ChatRoomService service; private final int bufferSize; private Sinks.Many sink; @@ -24,11 +30,13 @@ public class ChatRoom public ChatRoom( UUID id, String name, + Clock clock, ChatRoomService service, int bufferSize) { this.id = id; this.name = name; + this.clock = clock; this.service = service; this.bufferSize = bufferSize; this.sink = createSink(); @@ -37,20 +45,26 @@ public class ChatRoom synchronized public Mono addMessage( Long id, - LocalDateTime timestamp, String user, String text) { + Message.MessageKey key = Message.MessageKey.of(user, id); return service - .persistMessage(Message.MessageKey.of(user, id), timestamp, text) - .doOnNext(message -> - { - Sinks.EmitResult result = sink.tryEmitNext(message); - if (result.isFailure()) - { - log.warn("Emitting of message failed with {} for {}", result.name(), message); - } - }); + .getMessage(key) + .flatMap(existing -> text.equals(existing.getMessageText()) + ? Mono.just(existing) + : Mono.error(() -> new MessageMutationException(existing, text))) + .switchIfEmpty( + Mono + .fromSupplier(() ->service.persistMessage(key, LocalDateTime.now(clock), text)) + .doOnNext(m -> + { + Sinks.EmitResult result = sink.tryEmitNext(m); + if (result.isFailure()) + { + log.warn("Emitting of message failed with {} for {}", result.name(), m); + } + })); }