package com.sonymusic.hierdispatcher.service;

import java.time.LocalDateTime;

import org.springframework.stereotype.Service;

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;

import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;

@Slf4j
@Service
@RequiredArgsConstructor
public class KafkaConsumerService {

    /** Repository for group process status records. */
    private final GroupProcessStatusRepository groupProcessStatusRepository;

    /** Repository for performance hierarchy process status records. */
    private final PerfHierarchyProcessStatusRepository
            perfHierarchyProcessStatusRepository;

    /** Repository for DRA Lambda audit log records. */
    private final DraLambdaAuditLogRepository draLambdaAuditLogRepository;

    /** JSON mapper used when auditing Kafka payloads. */
    private final ObjectMapper objectMapper;

    /**
     * Processes a Kafka payload and persists the appropriate records.
     *
     * @param kafkaMessage inbound message payload
     * @param payloadJson serialized payload JSON
     */
    public void processMessage(final KafkaMessage kafkaMessage,
                               final String payloadJson) {
        String trackNo = kafkaMessage.getTrackNo();
        String trackExt = kafkaMessage.getTrackExt();

        log.info(
                "Processing message with trackNo: {}, trackExt: {}, thread: {}",
                trackNo,
                trackExt,
                Thread.currentThread().getName());

        // Validation
        if (isInvalid(trackNo)) {
            log.error(
                    "Invalid input: trackNo is null/blank for messageId: {}",
                    kafkaMessage.getMessageId());
            insertFailedAuditLog(
                    kafkaMessage,
                    payloadJson,
                    "Invalid trackNo");
            return; // Do not acknowledge
        }

        try {
            boolean exists = groupProcessStatusRepository
                    .existsByTrackNoAndTrackExtAndStatus(trackNo,
                            trackExt, "R");

            if (!exists) {
                insertGroupProcessStatus(kafkaMessage);
            } else {
                log.info(
                        "Skipping group_process_status insert for existing trackNo: {}, trackExt: {}",
                        trackNo, trackExt);
            }

            insertPerfHierarchyProcessStatus(kafkaMessage);
            insertSuccessAuditLog(
                    kafkaMessage,
                    payloadJson);

            log.info(
                    "Successfully processed message with trackNo: {}, trackExt: {}",
                    trackNo,
                    trackExt);
        } catch (Exception e) {
            log.error(
                    "Error processing message with trackNo: {}, trackExt: {}",
                    trackNo,
                    trackExt,
                    e);
            insertFailedAuditLog(
                    kafkaMessage,
                    payloadJson,
                    e.getMessage());
            throw e; // Re-throw to prevent acknowledgment
        }
    }

    private boolean isInvalid(final String trackNo) {
        return trackNo == null
                || trackNo.trim().isEmpty();
    }

    private void insertGroupProcessStatus(final KafkaMessage kafkaMessage) {
        GroupProcessStatus gps = new GroupProcessStatus();
        gps.setAuditLogDocumentId(kafkaMessage.getMessageId());
        gps.setProductId(null);
        gps.setRecordType(kafkaMessage.getRecordType());
        gps.setStatus("R");
        gps.setTrackNo(kafkaMessage.getTrackNo());
        gps.setTrackExt(kafkaMessage.getTrackExt());
        gps.setType(ProcessType.HIER);
        gps.setUpdatedDt(LocalDateTime.now());

        groupProcessStatusRepository.save(gps);
        log.debug(
                "Inserted into group_process_status for trackNo: {}, trackExt: {}",
                kafkaMessage.getTrackNo(),
                kafkaMessage.getTrackExt());
    }

    private void insertPerfHierarchyProcessStatus(
            final KafkaMessage kafkaMessage) {
        PerfHierarchyProcessStatus phps = new PerfHierarchyProcessStatus();
        phps.setAwsRequestId(null);
        phps.setProductId(null);
        phps.setRecordType(kafkaMessage.getRecordType());
        phps.setStatus("R");
        phps.setTrackNo(kafkaMessage.getTrackNo());
        phps.setTrackExt(kafkaMessage.getTrackExt());
        phps.setUpdatedDt(LocalDateTime.now());

        perfHierarchyProcessStatusRepository.save(phps);
        log.debug(
                "Inserted into perf_hierarchy_process_status for trackNo: {}, trackExt: {}",
                kafkaMessage.getTrackNo(),
                kafkaMessage.getTrackExt());
    }

    private void insertSuccessAuditLog(
            final KafkaMessage kafkaMessage,
            final String payloadJson) {
        DraLambdaAuditLog auditLog = createBaseAuditLog(kafkaMessage, payloadJson);
        auditLog.setStatus("SUCCESS");
        auditLog.setErrorMsg(null);

        draLambdaAuditLogRepository.save(auditLog);
        log.debug(
                "Inserted SUCCESS audit log for trackNo: {}, trackExt: {}",
                kafkaMessage.getTrackNo(),
                kafkaMessage.getTrackExt());
    }

    private void insertFailedAuditLog(
            final KafkaMessage kafkaMessage,
            final String payloadJson,
            final String errorMsg) {
        DraLambdaAuditLog auditLog = createBaseAuditLog(kafkaMessage, payloadJson);
        auditLog.setStatus("FAILED");
        auditLog.setErrorMsg(errorMsg);

        draLambdaAuditLogRepository.save(auditLog);
        log.debug(
                "Inserted FAILED audit log for trackNo: {}, trackExt: {}",
                kafkaMessage.getTrackNo(),
                kafkaMessage.getTrackExt());
    }

    private DraLambdaAuditLog createBaseAuditLog(
            final KafkaMessage kafkaMessage,
            final String payloadJson) {
        DraLambdaAuditLog auditLog = new DraLambdaAuditLog();
        auditLog.setCreatedDate(LocalDateTime.now());
        auditLog.setUpdatedDate(LocalDateTime.now());
        auditLog.setEventType(EventType.HIER_DISPATCHER);
        auditLog.setPayload(payloadJson);
        auditLog.setProductId(null);
        auditLog.setRetryCount(0);
        auditLog.setTrackNo(kafkaMessage.getTrackNo());
        auditLog.setTrackExt(kafkaMessage.getTrackExt());
        return auditLog;
    }
}
