package com.sonymusic.hierdispatcher.service;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;

import java.util.Collections;

import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;

import com.fasterxml.jackson.databind.ObjectMapper;
import com.sonymusic.hierdispatcher.model.DraLambdaAuditLog;
import com.sonymusic.hierdispatcher.model.EventType;
import com.sonymusic.hierdispatcher.model.GroupProcessStatus;
import com.sonymusic.hierdispatcher.model.KafkaMessage;
import com.sonymusic.hierdispatcher.model.PerfHierarchyProcessStatus;
import com.sonymusic.hierdispatcher.model.ProcessType;
import com.sonymusic.hierdispatcher.repository.DraLambdaAuditLogRepository;
import com.sonymusic.hierdispatcher.repository.GroupProcessStatusRepository;
import com.sonymusic.hierdispatcher.repository.PerfHierarchyProcessStatusRepository;

@ExtendWith(MockitoExtension.class)
class KafkaConsumerServiceTest {

    @Mock
    private GroupProcessStatusRepository groupProcessStatusRepository;

    @Mock
    private PerfHierarchyProcessStatusRepository perfHierarchyProcessStatusRepository;

    @Mock
    private DraLambdaAuditLogRepository draLambdaAuditLogRepository;

    @Mock
    private ObjectMapper objectMapper;

    @Captor
    private ArgumentCaptor<DraLambdaAuditLog> draLambdaAuditLogCaptor;

    @Captor
    private ArgumentCaptor<GroupProcessStatus> groupProcessStatusCaptor;

    @Captor
    private ArgumentCaptor<PerfHierarchyProcessStatus> perfHierarchyProcessStatusCaptor;

    @InjectMocks
    private KafkaConsumerService kafkaConsumerService;

    private KafkaMessage kafkaMessage;
    private String payloadJson = "{\"trackNo\":\"123\",\"trackExt\":\"ext\",\"messageId\":\"msg1\",\"recordType\":\"type1\",\"products\":[]}";

    @BeforeEach
    void setUp() {
        kafkaMessage = new KafkaMessage();
        kafkaMessage.setTrackNo("123");
        kafkaMessage.setTrackExt("ext");
        kafkaMessage.setMessageId("msg1");
        kafkaMessage.setRecordType("type1");
        kafkaMessage.setProducts(Collections.emptyList());
    }

    @Test
    void testProcessMessage_NewRecord_Success() throws Exception {
        when(groupProcessStatusRepository.existsByTrackNoAndTrackExtAndStatus("123", "ext", "R")).thenReturn(false);

        kafkaConsumerService.processMessage(kafkaMessage, payloadJson);

        verify(groupProcessStatusRepository).save(any(GroupProcessStatus.class));
        verify(perfHierarchyProcessStatusRepository).save(any(PerfHierarchyProcessStatus.class));
        verify(draLambdaAuditLogRepository).save(any(DraLambdaAuditLog.class));
    }

    @Test
    void testProcessMessage_DuplicateRecord_SkipGroupInsert() throws Exception {
        when(groupProcessStatusRepository.existsByTrackNoAndTrackExtAndStatus("123", "ext", "R")).thenReturn(true);

        kafkaConsumerService.processMessage(kafkaMessage, payloadJson);

        verify(groupProcessStatusRepository, never()).save(any(GroupProcessStatus.class));
        verify(perfHierarchyProcessStatusRepository).save(any(PerfHierarchyProcessStatus.class));
        verify(draLambdaAuditLogRepository).save(any(DraLambdaAuditLog.class));
    }

    @Test
    void testProcessMessage_NewRecord_Success_AuditAndStatusValues() throws Exception {
        when(groupProcessStatusRepository.existsByTrackNoAndTrackExtAndStatus("123", "ext", "R")).thenReturn(false);

        kafkaConsumerService.processMessage(kafkaMessage, payloadJson);

        verify(groupProcessStatusRepository).save(groupProcessStatusCaptor.capture());
        verify(perfHierarchyProcessStatusRepository).save(perfHierarchyProcessStatusCaptor.capture());
        verify(draLambdaAuditLogRepository).save(draLambdaAuditLogCaptor.capture());

        GroupProcessStatus savedGroupProcessStatus = groupProcessStatusCaptor.getValue();
        assertEquals("R", savedGroupProcessStatus.getStatus());
        assertEquals(ProcessType.HIER, savedGroupProcessStatus.getType());
        assertEquals("123", savedGroupProcessStatus.getTrackNo());
        assertEquals("ext", savedGroupProcessStatus.getTrackExt());
        assertEquals("type1", savedGroupProcessStatus.getRecordType());

        PerfHierarchyProcessStatus savedPerfProcessStatus = perfHierarchyProcessStatusCaptor.getValue();
        assertEquals("R", savedPerfProcessStatus.getStatus());
        assertEquals("123", savedPerfProcessStatus.getTrackNo());
        assertEquals("ext", savedPerfProcessStatus.getTrackExt());
        assertEquals("type1", savedPerfProcessStatus.getRecordType());

        DraLambdaAuditLog savedAuditLog = draLambdaAuditLogCaptor.getValue();
        assertEquals("SUCCESS", savedAuditLog.getStatus());
        assertNull(savedAuditLog.getErrorMsg());
        assertEquals(payloadJson, savedAuditLog.getPayload());
        assertEquals(EventType.HIER_DISPATCHER, savedAuditLog.getEventType());
        assertEquals("123", savedAuditLog.getTrackNo());
        assertEquals("ext", savedAuditLog.getTrackExt());
    }

    @Test
    void testProcessMessage_BlankTrackExt_ProcessesSuccessfully() throws Exception {
        kafkaMessage.setTrackExt("   ");
        when(groupProcessStatusRepository.existsByTrackNoAndTrackExtAndStatus("123", "   ", "R")).thenReturn(false);

        kafkaConsumerService.processMessage(kafkaMessage, payloadJson);

        verify(draLambdaAuditLogRepository).save(draLambdaAuditLogCaptor.capture());
        verify(groupProcessStatusRepository).save(any(GroupProcessStatus.class));
        verify(perfHierarchyProcessStatusRepository).save(any(PerfHierarchyProcessStatus.class));

        DraLambdaAuditLog savedAuditLog = draLambdaAuditLogCaptor.getValue();
        assertEquals("SUCCESS", savedAuditLog.getStatus());
        assertNull(savedAuditLog.getErrorMsg());
    }

    @Test
    void testProcessMessage_DBFailure_FailedAuditAndException() throws Exception {
        when(groupProcessStatusRepository.existsByTrackNoAndTrackExtAndStatus("123", "ext", "R")).thenReturn(false);
        when(groupProcessStatusRepository.save(any(GroupProcessStatus.class))).thenThrow(new RuntimeException("DB error"));

        RuntimeException thrown = assertThrows(RuntimeException.class, () ->
                kafkaConsumerService.processMessage(kafkaMessage, payloadJson));
        assertEquals("DB error", thrown.getMessage());

        verify(draLambdaAuditLogRepository).save(draLambdaAuditLogCaptor.capture());
        DraLambdaAuditLog savedAuditLog = draLambdaAuditLogCaptor.getValue();
        assertEquals("FAILED", savedAuditLog.getStatus());
        assertEquals("DB error", savedAuditLog.getErrorMsg());
    }

}
