//package org.egl_cepgl.pm.service;
//
//import com.mongodb.client.model.changestream.FullDocument;
//import io.jsonwebtoken.Jwt;
//import jakarta.annotation.PostConstruct;
//import lombok.extern.slf4j.Slf4j;
//import org.egl_cepgl.pm.dto.*;
//import org.egl_cepgl.pm.model.ChatContact;
//import org.egl_cepgl.pm.model.ChatInbox;
//import org.egl_cepgl.pm.model.ChatMessage;
//import org.egl_cepgl.pm.model.RecentChat;
//import org.egl_cepgl.pm.repository.ChatContactRepository;
//import org.egl_cepgl.pm.repository.ChatInboxRepository;
//import org.egl_cepgl.pm.repository.ChatMessageRepository;
//import org.egl_cepgl.pm.repository.RecentChatRepository;
//import org.egl_cepgl.pm.service.keycloak.KeycloakUserService;
//import org.springframework.beans.factory.annotation.Autowired;
//import org.springframework.data.mongodb.core.ChangeStreamEvent;
//import org.springframework.data.mongodb.core.ChangeStreamOptions;
//import org.springframework.data.mongodb.core.ReactiveMongoTemplate;
//import org.springframework.data.mongodb.core.aggregation.Aggregation;
//import org.springframework.data.mongodb.core.aggregation.MatchOperation;
//import org.springframework.data.mongodb.core.messaging.DefaultMessageListenerContainer;
//import org.springframework.data.mongodb.core.query.Criteria;
//import org.springframework.security.core.Authentication;
//import org.springframework.security.core.context.SecurityContextHolder;
//import org.springframework.stereotype.Service;
//import reactor.core.publisher.Flux;
//import reactor.core.publisher.Mono;
//import reactor.core.publisher.Sinks;
//
//import java.time.Duration;
//import java.time.Instant;
//import java.util.*;
//import java.util.concurrent.ConcurrentHashMap;
//import java.util.concurrent.atomic.AtomicReference;
//
//import static org.springframework.data.mongodb.core.aggregation.Aggregation.match;
//import static org.springframework.data.mongodb.core.query.Criteria.where;
//
//@Slf4j
//@Service
//public class ChatService
//{
//    private final ReactiveMongoTemplate mongoTemplate;
//    private DefaultMessageListenerContainer container;
//
//    private ChatMessageRepository chatMessageRepository;
//    private ChatInboxRepository chatInboxRepository;
//    private ChatContactRepository ccontactRepo;
//    private RecentChatRepository recentChatRepo;
//    private KeycloakUserService kcService;
//
//    private final Sinks.Many<ChatMessage> messageSink4 = Sinks.many().replay().all();
//    private final Sinks.Many<ChatInbox> recentChatSink4 = Sinks.many().replay().all();
//    //private final Sinks.Many<ContactRecentMsg> contacts = Sinks.many().replay().all();
//    private final Set<String> connectedUsers = ConcurrentHashMap.newKeySet();
//    private final Sinks.Many<String> notificationUserSink = Sinks.many().multicast().onBackpressureBuffer();
//
//    @Autowired
//    public ChatService(
//            ChatMessageRepository chatMessageRepository,
//            ChatInboxRepository chatInboxRepository,
//            ReactiveMongoTemplate mongoTemplate,
//            ChatContactRepository ccontactRepo,
//            RecentChatRepository recentChatRepo,
//            KeycloakUserService kcService
//    ){
//        this.chatMessageRepository = chatMessageRepository;
//        this.chatInboxRepository = chatInboxRepository;
//        this.mongoTemplate = mongoTemplate;
//        this.ccontactRepo= ccontactRepo;
//        this.recentChatRepo= recentChatRepo;
//        this.kcService= kcService;
//    }
//
//    @PostConstruct
//    public void initMongoDBChangeStream()
//    {
//        //Criteria criteria = Criteria.where("fieldName").is("value").and("anotherField").gt(10);
////        Criteria criteria = Criteria.where("operationType").is(OperationType.INSERT.name().toLowerCase());
////        MatchOperation matchOperation = Aggregation.match(criteria);
////        Aggregation aggregation = Aggregation.newAggregation(matchOperation);
////        Flux<ChangeStreamEvent<ChatMessage>> changeStream =
////                mongoTemplate
////                    .changeStream(ChatMessage.class)
////                    .watchCollection("chatMessage")
////                    .filter(aggregation)
////                    .listen();
////
////        changeStream
////                .doFirst(() -> log.info("==LISTENER=="))
////            .map(ChangeStreamEvent::getBody)
////            .subscribe(messageSink::tryEmitNext); // Emits the new or updated product
//        //chatMessageRepository.findAllBy().subscribe(messageSink::tryEmitNext);
//    }
//
////    @PostConstruct
////    public void startChangeStream() {
////        container = new DefaultMessageListenerContainer(mongoTemplate);
////        container.start();
////
////        // Define your listener logic
////        MessageListener<ChangeStreamDocument<Document>, YourDomainObject> listener = changeStreamDocument -> {
////            System.out.println("Received change: " + changeStreamDocument);
////            // Process the changeStreamDocument here
////            // You can access the full document, operation type, etc.
////            // Example:
////            // if (changeStreamDocument.getOperationType() == OperationType.INSERT) {
////            //     YourDomainObject newObject = changeStreamDocument.getFullDocument();
////            //     // ... handle new object
////            // }
////        };
////
////        // Configure ChangeStreamRequestOptions
////        ChangeStreamRequestOptions options = new ChangeStreamRequestOptions(
////                "yourDatabaseName", // Database name
////                "yourCollectionName", // Collection name
////                ChangeStreamOptions.empty() // You can add options like fullDocument, filter, etc.
////                // .fullDocument(FullDocument.UPDATE_LOOKUP) // Example: to get full document on update
////        );
////
////        // Register the change stream request
////        container.register(new ChangeStreamRequest<>(listener, options), YourDomainObject.class);
////    }
//
//    public void getChatInboxChangeStream(String userEmail)
//    {
//
////        log.info("RCHT===");
////        //MatchOperation matchUpdateType = match(where("operationType").is("update"));
////        //MatchOperation matchFieldChange = match(where("updateDescription.updatedFields.messageIds").exists(true));
////        Criteria criteria = new Criteria().orOperator(
////                Criteria.where("fullDocument.username1").is(userEmail),
////                Criteria.where("fullDocument.username2").is(userEmail)
////        );
////        MatchOperation matchField = match(criteria);
////        Aggregation aggregation = Aggregation.newAggregation(matchField);
////
////        ChangeStreamOptions options = ChangeStreamOptions.builder()
////                //.fullDocumentLookup(FullDocument.UPDATE_LOOKUP)
////                //.filter(aggregation)
////                .build();
////
////        return mongoTemplate.changeStream(
////                "chatInbox",
////                options,
////                ChatInbox.class
////        );
//    }
//
//    public void findRecentChatOld(String userEmail)
//    {
////        log.info("RRRRRR");
////        Flux<RecentChatDto> recentflux= Flux.empty();
////        getChatInboxChangeStream(userEmail).flatMap((event) -> {
////            ChatInbox currIbx= event.getBody();
////            Integer asize= currIbx.getMessageIds().size();
////            AtomicReference<ChatMessage> msg= new AtomicReference<>(new ChatMessage());
////            AtomicReference<String> userTo= new AtomicReference<>("");
////            AtomicReference<Integer> unreadMsg= new AtomicReference<>(0);
////            chatMessageRepository.findById(currIbx.getMessageIds().get(asize-1))
////                    .doOnNext((mgg) -> { if(mgg.getRead()) unreadMsg.getAndSet(unreadMsg.get()+1); })
////                    .doOnNext((mg) -> msg.set(mg))
////                    .flatMap((m) -> {
////                        String msgToEmail= currIbx.getUsername1() == userEmail ? currIbx.getUsername2() : currIbx.getUsername2();
////                        return ccontactRepo.findByEmail(msgToEmail).doOnNext((cc) -> userTo.set(cc.getFirst_name()+' '+cc.getLast_name()));
////                    }).subscribe();
////            RecentChatDto rc= RecentChatDto.builder()
////                    .chatInbox(currIbx)
////                    .lastMsg(msg.get())
////                    .userTo(userTo.get())
////                    .unreadMsg(unreadMsg.get()).build();
////            return recentflux.concatMap(rf -> Mono.just(rc));
////        }).subscribe();
////        //recentflux.map(recentChatSink::tryEmitNext);
////        this.recentChatSink.asFlux().delayElements(Duration.ofMillis(100));
//    }
//
//    @PostConstruct
//    public void initChatMessages(){
//        chatMessageRepository.findAllBy().subscribe(messageSink4::tryEmitNext);
//    }
//
//    @PostConstruct
//    public void initChatInbox(){
//        chatInboxRepository.findAllBy().subscribe(recentChatSink4::tryEmitNext);
//    }
//
//    public Flux<ChatInbox> getRecentChats(String userEmail)
//    {
//        return this.recentChatSink4.asFlux()
//                .filter(rc -> rc.getUsername1().equals(userEmail) || rc.getUsername2().equals(userEmail))
//                .delayElements(Duration.ofMillis(100));
//    }
//
//    public Flux<ChatMessage> getMsgsByInbox(String inbox)
//    {
//        return this.messageSink4.asFlux()
//                   .filter(m -> m.getIndBoxId().equals(inbox))
//                   .delayElements(Duration.ofMillis(100));
//    }
//
//    public void saveMsg(String inbox, MessageText msg)
//    {
//        chatMessageRepository
//            .save(new ChatMessage(msg.getMsg(), Instant.now(), msg.getSendBy(), false, inbox))
//                  .flatMap(mg ->
//                      chatInboxRepository.findById(inbox)
//                      .flatMap(inb -> {
//                          inb.setLastMsg(mg);
//                          inb.getMessageIds().add(mg.getId());
//                          return chatInboxRepository.save(inb);
//                      })
//                      .thenReturn(mg)
//                  ).subscribe();
//    }
//
//    public Flux<ChatContact> findAllContacts()
//    {
//        return this.ccontactRepo.findAllByStatus(true);
//    }
//
//    public Mono<ChatInbox> addInboxIfNotExists(String userEmail, InboxUsersDto nbx)
//    {
//        return this.chatInboxRepository.findByNamepOrNamep(
//                nbx.getUsername1()+'_'+nbx.getUsername2(),
//                nbx.getUsername2()+'_'+nbx.getUsername1()
//              ).switchIfEmpty(Mono.defer(() -> {
//                    ChatInbox inbox= ChatInbox.builder()
//                        .namep(nbx.getUsername1()+'_'+nbx.getUsername2())
//                        .sentAt(Instant.now())
//                        .id(UUID.randomUUID().toString())
//                        .username1(nbx.getUsername1())
//                        .username2(nbx.getUsername2())
//                        .user1fullname(nbx.getUser1fullname())
//                        .user2fullname(nbx.getUser2fullname())
//                        .lastMsg(null)
//                        .unreadMsg(0)
//                        .messageIds(new ArrayList<>()).build();
//                    return this.chatInboxRepository.save(inbox);
//              }));
//    }
//}
//
//
//
//
//
//
