Subversion Repositories SmartDukaan

Rev

Rev 36357 | Blame | Compare with Previous | Last modification | View Log | RSS feed

package com.spice.profitmandi.service.cron;

import com.spice.profitmandi.dao.entity.transaction.CronBatch;
import com.spice.profitmandi.dao.entity.transaction.CronBatchItem;
import com.spice.profitmandi.dao.enumuration.transaction.CronBatchItemStatus;
import com.spice.profitmandi.dao.enumuration.transaction.CronBatchStatus;
import com.spice.profitmandi.dao.repository.transaction.CronBatchItemRepository;
import com.spice.profitmandi.dao.repository.transaction.CronBatchRepository;
import com.spice.profitmandi.service.mail.MailOutboxService;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Propagation;
import org.springframework.transaction.annotation.Transactional;

import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;

@Service
public class CronBatchService {

    private static final Logger LOGGER = LogManager.getLogger(CronBatchService.class);
    private static final String[] TECHNOLOGY_EMAIL = {"amit.gupta@smartdukaan.com"};
    private static final DateTimeFormatter DTF = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
    /** Cap on how many never-attempted items the alert lists by name. */
    private static final int PENDING_LIST_CAP = 50;

    @Autowired
    private CronBatchRepository cronBatchRepository;

    @Autowired
    private CronBatchItemRepository cronBatchItemRepository;

    @Autowired
    private MailOutboxService mailOutboxService;

    /**
     * Read-only guard used by OfferBatchService before creating a new batch.
     * Returns any RUNNING batch for this offer (sellin or activation job name),
     * or null if none. Wrapped in a read-only tx so callers running outside an
     * existing transaction (e.g. background worker thread) can query safely.
     */
    @Transactional(readOnly = true)
    public CronBatch findRunningForOffer(int offerId) {
        List<CronBatch> running = cronBatchRepository.selectRunningForOffer(offerId);
        return running.isEmpty() ? null : running.get(0);
    }

    @Transactional(propagation = Propagation.REQUIRES_NEW)
    public CronBatch createBatch(String jobName, Map<Integer, String> fofoIdPartnerNameMap) {
        CronBatch batch = new CronBatch(jobName);
        batch.setTotalCount(fofoIdPartnerNameMap.size());
        cronBatchRepository.persist(batch);

        for (Map.Entry<Integer, String> entry : fofoIdPartnerNameMap.entrySet()) {
            CronBatchItem item = new CronBatchItem(batch.getId(), entry.getKey(), entry.getValue());
            cronBatchItemRepository.persist(item);
        }

        LOGGER.info("Created batch {} for job {} with {} items", batch.getId(), jobName, fofoIdPartnerNameMap.size());
        return batch;
    }

    /**
     * REQUIRED — deliberately NOT REQUIRES_NEW.
     *
     * Every caller is the last statement of a REQUIRES_NEW per-item work method, so this
     * joins that work transaction and the SUCCESS marker commits (or rolls back) with the
     * work it describes. Under REQUIRES_NEW the marker committed on its own, so work that
     * later failed at flush/commit — a constraint violation, a wallet deadlock — was left
     * recorded as a successful payout.
     */
    @Transactional(propagation = Propagation.REQUIRED)
    public void markItemSuccess(int batchId, int fofoId) {
        CronBatchItem item = cronBatchItemRepository.selectByBatchIdAndFofoId(batchId, fofoId);
        if (item != null && item.getStatus() == CronBatchItemStatus.PENDING) {
            item.markSuccess();
        }
    }

    /**
     * REQUIRES_NEW — the opposite of markItemSuccess, and intentionally so: this is called
     * from the orchestrator's catch block after the work transaction has already rolled
     * back, so the failure record must commit independently of it.
     */
    @Transactional(propagation = Propagation.REQUIRES_NEW)
    public void markItemFailed(int batchId, int fofoId, String errorMessage) {
        CronBatchItem item = cronBatchItemRepository.selectByBatchIdAndFofoId(batchId, fofoId);
        if (item != null && item.getStatus() == CronBatchItemStatus.PENDING) {
            item.markFailed(errorMessage);
        }
    }

    /**
     * RUNNING batches older than {@code staleHours} — a run that died before finalizing.
     * Returns ids only; the caller must finalize each through the proxy (see
     * BatchScheduledTasks.reapStaleBatches) so REQUIRES_NEW actually applies.
     */
    @Transactional(readOnly = true)
    public List<Integer> selectStaleRunningBatchIds(int staleHours) {
        return cronBatchRepository.selectStaleRunning(LocalDateTime.now().minusHours(staleHours))
                .stream().map(CronBatch::getId).collect(Collectors.toList());
    }

    /**
     * Closes a batch using the item rows as the source of truth.
     *
     * successCount is COUNTED, never derived as (total - failures). Deriving it reported
     * every item that never reported at all — the JVM killed mid-loop, a lost marker write —
     * as a successful payout, and suppressed the alert, because failureCount stayed 0.
     *
     * Items still PENDING were never attempted. That leaves the batch INCOMPLETE (not
     * PARTIAL_FAILURE, which means "attempted and failed") and always raises an email.
     * PENDING rows are left as they are: "never tried" and "tried and failed" are
     * different facts and the retry decision depends on which one it was.
     */
    @Transactional(propagation = Propagation.REQUIRES_NEW)
    public BatchOutcome finalizeBatch(int batchId) {
        CronBatch batch = cronBatchRepository.selectById(batchId);
        if (batch == null) {
            LOGGER.warn("finalizeBatch called for unknown batch {}", batchId);
            return null;
        }
        List<CronBatchItem> failedItems = cronBatchItemRepository.selectFailedByBatchId(batchId);

        int failureCount = failedItems.size();
        int successCount = (int) cronBatchItemRepository
                .countByBatchIdAndStatus(batchId, CronBatchItemStatus.SUCCESS);
        int pendingCount = batch.getTotalCount() - successCount - failureCount;

        batch.setSuccessCount(successCount);
        batch.setFailureCount(failureCount);
        batch.setCompletedAt(LocalDateTime.now());
        if (pendingCount > 0) {
            batch.setStatus(CronBatchStatus.INCOMPLETE);
        } else if (failureCount > 0) {
            batch.setStatus(CronBatchStatus.PARTIAL_FAILURE);
        } else {
            batch.setStatus(CronBatchStatus.COMPLETED);
        }

        LOGGER.info("Batch {} finalized as {}: {} success, {} failed, {} never attempted",
                batchId, batch.getStatus(), successCount, failureCount, pendingCount);

        if (failureCount > 0 || pendingCount > 0) {
            sendFailureEmail(batch, failedItems, pendingCount);
        }
        return new BatchOutcome(batch.getStatus(), batch.getTotalCount(),
                successCount, failureCount, pendingCount);
    }

    /**
     * Result of a finalize, returned rather than re-read by the caller: finalizeBatch
     * commits in its own session, so a caller that re-reads inside its own transaction
     * gets the stale pre-finalize entity back out of Hibernate's first-level cache.
     */
    public static class BatchOutcome {
        private final CronBatchStatus status;
        private final int totalCount;
        private final int successCount;
        private final int failureCount;
        private final int neverAttemptedCount;

        BatchOutcome(CronBatchStatus status, int totalCount, int successCount,
                     int failureCount, int neverAttemptedCount) {
            this.status = status;
            this.totalCount = totalCount;
            this.successCount = successCount;
            this.failureCount = failureCount;
            this.neverAttemptedCount = neverAttemptedCount;
        }

        public CronBatchStatus getStatus() { return status; }
        public int getTotalCount() { return totalCount; }
        public int getSuccessCount() { return successCount; }
        public int getFailureCount() { return failureCount; }
        public int getNeverAttemptedCount() { return neverAttemptedCount; }
    }

    private void sendFailureEmail(CronBatch batch, List<CronBatchItem> failedItems, int pendingCount) {
        StringBuilder body = new StringBuilder();
        body.append("Cron job: ").append(batch.getJobName()).append("\n");
        body.append("Batch id: ").append(batch.getId()).append("\n");
        body.append("Run time: ").append(batch.getStartedAt().format(DTF)).append("\n");
        body.append("Outcome: ").append(batch.getStatus()).append("\n");
        body.append("Total: ").append(batch.getTotalCount()).append("\n");
        body.append("Success: ").append(batch.getSuccessCount()).append("\n");
        body.append("Failed: ").append(batch.getFailureCount()).append("\n");
        body.append("Never attempted: ").append(pendingCount).append("\n\n");

        if (!failedItems.isEmpty()) {
            body.append("Partner failures:\n");
            for (CronBatchItem item : failedItems) {
                body.append("- ").append(item.getPartnerName())
                        .append(" (fofoId: ").append(item.getFofoId()).append(")")
                        .append(" — ").append(item.getErrorMessage())
                        .append("\n");
            }
            body.append("\n");
        }

        if (pendingCount > 0) {
            body.append("The run did not reach these ").append(pendingCount)
                    .append(" item(s) — they were never attempted and no work was done for them:\n");
            List<CronBatchItem> pendingItems = cronBatchItemRepository
                    .selectByBatchIdAndStatus(batch.getId(), CronBatchItemStatus.PENDING);
            for (CronBatchItem item : pendingItems.subList(0, Math.min(pendingItems.size(), PENDING_LIST_CAP))) {
                body.append("- ").append(item.getPartnerName())
                        .append(" (fofoId: ").append(item.getFofoId()).append(")\n");
            }
            if (pendingItems.size() > PENDING_LIST_CAP) {
                body.append("... and ").append(pendingItems.size() - PENDING_LIST_CAP).append(" more\n");
            }
            body.append("\nRe-run the job to process them.\n");
        }

        String subject = pendingCount > 0
                ? String.format("[CRON ALERT] %s — INCOMPLETE, %d of %d never attempted",
                        batch.getJobName(), pendingCount, batch.getTotalCount())
                : String.format("[CRON ALERT] %s — %d of %d partners failed",
                        batch.getJobName(), batch.getFailureCount(), batch.getTotalCount());

        try {
            mailOutboxService.queueMailViaGoogle(TECHNOLOGY_EMAIL, null, subject, body.toString(),
                    "CronBatchService." + batch.getJobName());
        } catch (Exception e) {
            LOGGER.error("Failed to send batch failure email for batch {}: {}", batch.getId(), e.getMessage());
        }
    }
}