import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.WakeupException;
+import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.time.*;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import java.util.UUID;
+import java.util.stream.IntStream;
@RequiredArgsConstructor
public class ChatRoomChannel implements Runnable
{
private final String topic;
- private final Producer<Integer, CreateChatRoomRequestTo> producer;
- private final Consumer<Integer, CreateChatRoomRequestTo> consumer;
- private final ShardingStrategy shardingStrategy;
- private final ChatMessageChannel chatMessageChannel;
- private final Clock clock;
- private final int bufferSize;
+ private final Consumer<String, AbstractTo> consumer;
+ private final Map<UUID, ChatRoomInfo> chatrooms = new HashMap<>();
private boolean running;
- Mono<ChatRoomInfo> sendCreateChatRoomRequest(
- UUID chatRoomId,
- String name)
- {
- int shard = this.shardingStrategy.selectShard(chatRoomId);
- CreateChatRoomRequestTo createChatRoomRequestTo = CreateChatRoomRequestTo.of(chatRoomId.toString(), name, shard);
- return Mono.create(sink ->
- {
- ProducerRecord<Integer, CreateChatRoomRequestTo> record =
- new ProducerRecord<>(
- topic,
- shard,
- createChatRoomRequestTo);
-
- producer.send(record, ((metadata, exception) ->
- {
- if (metadata != null)
- {
- log.info("Successfully send chreate-request for chat room: {}", createChatRoomRequestTo);
- sink.success(createChatRoomRequestTo.toChatRoomInfo());
- }
- else
- {
- // On send-failure
- log.error(
- "Could not send create-request for chat room (id={}, name={}): {}",
- chatRoomId,
- name,
- exception);
- sink.error(exception);
- }
- }));
- });
- }
-
@Override
public void run()
{
{
try
{
- ConsumerRecords<Integer, CreateChatRoomRequestTo> records = consumer.poll(Duration.ofMinutes(5));
+ ConsumerRecords<String, AbstractTo> records = consumer.poll(Duration.ofMinutes(5));
log.info("Fetched {} messages", records.count());
- for (ConsumerRecord<Integer, CreateChatRoomRequestTo> record : records)
+ for (ConsumerRecord<String, AbstractTo> record : records)
{
- createChatRoom(record.value().toChatRoomInfo());
+ switch (record.value().getType())
+ {
+ case CHATROOM_INFO:
+ createChatRoom((ChatRoomInfoTo) record.value());
+ break;
+
+ default:
+ log.debug(
+ "Ignoring message for key {} with offset {}: {}",
+ record.key(),
+ record.offset(),
+ record.value());
+ }
}
}
catch (WakeupException e)
}
- void createChatRoom(ChatRoomInfo chatRoomInfo)
+ void createChatRoom(ChatRoomInfoTo chatRoomInfoTo)
+ {
+ ChatRoomInfo chatRoomInfo = chatRoomInfoTo.toChatRoomInfo();
+ chatrooms.put(chatRoomInfo.getId(), chatRoomInfo);
+ }
+
+ Flux<ChatRoomInfo> getChatRooms()
{
- UUID id = chatRoomInfo.getId();
- String name = chatRoomInfo.getName();
- int shard = chatRoomInfo.getShard();
- log.info("Creating ChatRoom {} with buffer-size {}", id, bufferSize);
- KafkaChatRoomService service = new KafkaChatRoomService(chatMessageChannel, id);
- ChatRoom chatRoom = new ChatRoom(id, name, shard, clock, service, bufferSize);
- chatMessageChannel.putChatRoom(chatRoom);
+ return Flux.fromIterable(chatrooms.values());
}
}