Subversion Repositories SmartDukaan

Rev

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

Rev Author Line No. Line
36337 amit 1
package com.spice.profitmandi.service.cron;
2
 
3
import com.spice.profitmandi.dao.entity.transaction.CronBatch;
4
import com.spice.profitmandi.dao.entity.transaction.CronBatchItem;
5
import com.spice.profitmandi.dao.enumuration.transaction.CronBatchItemStatus;
6
import com.spice.profitmandi.dao.enumuration.transaction.CronBatchStatus;
7
import com.spice.profitmandi.dao.repository.transaction.CronBatchItemRepository;
8
import com.spice.profitmandi.dao.repository.transaction.CronBatchRepository;
9
import com.spice.profitmandi.service.mail.MailOutboxService;
10
import org.apache.logging.log4j.LogManager;
11
import org.apache.logging.log4j.Logger;
12
import org.springframework.beans.factory.annotation.Autowired;
13
import org.springframework.stereotype.Service;
14
import org.springframework.transaction.annotation.Propagation;
15
import org.springframework.transaction.annotation.Transactional;
16
 
17
import java.time.LocalDateTime;
18
import java.time.format.DateTimeFormatter;
19
import java.util.List;
20
import java.util.Map;
37644 amit 21
import java.util.stream.Collectors;
36337 amit 22
 
23
@Service
24
public class CronBatchService {
25
 
26
    private static final Logger LOGGER = LogManager.getLogger(CronBatchService.class);
27
    private static final String[] TECHNOLOGY_EMAIL = {"amit.gupta@smartdukaan.com"};
28
    private static final DateTimeFormatter DTF = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
37644 amit 29
    /** Cap on how many never-attempted items the alert lists by name. */
30
    private static final int PENDING_LIST_CAP = 50;
36337 amit 31
 
32
    @Autowired
33
    private CronBatchRepository cronBatchRepository;
34
 
35
    @Autowired
36
    private CronBatchItemRepository cronBatchItemRepository;
37
 
38
    @Autowired
39
    private MailOutboxService mailOutboxService;
40
 
36357 amit 41
    /**
42
     * Read-only guard used by OfferBatchService before creating a new batch.
43
     * Returns any RUNNING batch for this offer (sellin or activation job name),
44
     * or null if none. Wrapped in a read-only tx so callers running outside an
45
     * existing transaction (e.g. background worker thread) can query safely.
46
     */
47
    @Transactional(readOnly = true)
48
    public CronBatch findRunningForOffer(int offerId) {
49
        List<CronBatch> running = cronBatchRepository.selectRunningForOffer(offerId);
50
        return running.isEmpty() ? null : running.get(0);
51
    }
52
 
36337 amit 53
    @Transactional(propagation = Propagation.REQUIRES_NEW)
54
    public CronBatch createBatch(String jobName, Map<Integer, String> fofoIdPartnerNameMap) {
55
        CronBatch batch = new CronBatch(jobName);
56
        batch.setTotalCount(fofoIdPartnerNameMap.size());
57
        cronBatchRepository.persist(batch);
58
 
59
        for (Map.Entry<Integer, String> entry : fofoIdPartnerNameMap.entrySet()) {
60
            CronBatchItem item = new CronBatchItem(batch.getId(), entry.getKey(), entry.getValue());
61
            cronBatchItemRepository.persist(item);
62
        }
63
 
64
        LOGGER.info("Created batch {} for job {} with {} items", batch.getId(), jobName, fofoIdPartnerNameMap.size());
65
        return batch;
66
    }
67
 
37644 amit 68
    /**
69
     * REQUIRED — deliberately NOT REQUIRES_NEW.
70
     *
71
     * Every caller is the last statement of a REQUIRES_NEW per-item work method, so this
72
     * joins that work transaction and the SUCCESS marker commits (or rolls back) with the
73
     * work it describes. Under REQUIRES_NEW the marker committed on its own, so work that
74
     * later failed at flush/commit — a constraint violation, a wallet deadlock — was left
75
     * recorded as a successful payout.
76
     */
77
    @Transactional(propagation = Propagation.REQUIRED)
36337 amit 78
    public void markItemSuccess(int batchId, int fofoId) {
37644 amit 79
        CronBatchItem item = cronBatchItemRepository.selectByBatchIdAndFofoId(batchId, fofoId);
80
        if (item != null && item.getStatus() == CronBatchItemStatus.PENDING) {
81
            item.markSuccess();
36337 amit 82
        }
83
    }
84
 
37644 amit 85
    /**
86
     * REQUIRES_NEW — the opposite of markItemSuccess, and intentionally so: this is called
87
     * from the orchestrator's catch block after the work transaction has already rolled
88
     * back, so the failure record must commit independently of it.
89
     */
36337 amit 90
    @Transactional(propagation = Propagation.REQUIRES_NEW)
91
    public void markItemFailed(int batchId, int fofoId, String errorMessage) {
37644 amit 92
        CronBatchItem item = cronBatchItemRepository.selectByBatchIdAndFofoId(batchId, fofoId);
93
        if (item != null && item.getStatus() == CronBatchItemStatus.PENDING) {
94
            item.markFailed(errorMessage);
36337 amit 95
        }
96
    }
97
 
37644 amit 98
    /**
99
     * RUNNING batches older than {@code staleHours} — a run that died before finalizing.
100
     * Returns ids only; the caller must finalize each through the proxy (see
101
     * BatchScheduledTasks.reapStaleBatches) so REQUIRES_NEW actually applies.
102
     */
103
    @Transactional(readOnly = true)
104
    public List<Integer> selectStaleRunningBatchIds(int staleHours) {
105
        return cronBatchRepository.selectStaleRunning(LocalDateTime.now().minusHours(staleHours))
106
                .stream().map(CronBatch::getId).collect(Collectors.toList());
107
    }
108
 
109
    /**
110
     * Closes a batch using the item rows as the source of truth.
111
     *
112
     * successCount is COUNTED, never derived as (total - failures). Deriving it reported
113
     * every item that never reported at all — the JVM killed mid-loop, a lost marker write —
114
     * as a successful payout, and suppressed the alert, because failureCount stayed 0.
115
     *
116
     * Items still PENDING were never attempted. That leaves the batch INCOMPLETE (not
117
     * PARTIAL_FAILURE, which means "attempted and failed") and always raises an email.
118
     * PENDING rows are left as they are: "never tried" and "tried and failed" are
119
     * different facts and the retry decision depends on which one it was.
120
     */
36337 amit 121
    @Transactional(propagation = Propagation.REQUIRES_NEW)
37644 amit 122
    public BatchOutcome finalizeBatch(int batchId) {
36337 amit 123
        CronBatch batch = cronBatchRepository.selectById(batchId);
37644 amit 124
        if (batch == null) {
125
            LOGGER.warn("finalizeBatch called for unknown batch {}", batchId);
126
            return null;
127
        }
36337 amit 128
        List<CronBatchItem> failedItems = cronBatchItemRepository.selectFailedByBatchId(batchId);
129
 
130
        int failureCount = failedItems.size();
37644 amit 131
        int successCount = (int) cronBatchItemRepository
132
                .countByBatchIdAndStatus(batchId, CronBatchItemStatus.SUCCESS);
133
        int pendingCount = batch.getTotalCount() - successCount - failureCount;
36337 amit 134
 
135
        batch.setSuccessCount(successCount);
136
        batch.setFailureCount(failureCount);
137
        batch.setCompletedAt(LocalDateTime.now());
37644 amit 138
        if (pendingCount > 0) {
139
            batch.setStatus(CronBatchStatus.INCOMPLETE);
140
        } else if (failureCount > 0) {
141
            batch.setStatus(CronBatchStatus.PARTIAL_FAILURE);
142
        } else {
143
            batch.setStatus(CronBatchStatus.COMPLETED);
144
        }
36337 amit 145
 
37644 amit 146
        LOGGER.info("Batch {} finalized as {}: {} success, {} failed, {} never attempted",
147
                batchId, batch.getStatus(), successCount, failureCount, pendingCount);
36337 amit 148
 
37644 amit 149
        if (failureCount > 0 || pendingCount > 0) {
150
            sendFailureEmail(batch, failedItems, pendingCount);
36337 amit 151
        }
37644 amit 152
        return new BatchOutcome(batch.getStatus(), batch.getTotalCount(),
153
                successCount, failureCount, pendingCount);
36337 amit 154
    }
155
 
37644 amit 156
    /**
157
     * Result of a finalize, returned rather than re-read by the caller: finalizeBatch
158
     * commits in its own session, so a caller that re-reads inside its own transaction
159
     * gets the stale pre-finalize entity back out of Hibernate's first-level cache.
160
     */
161
    public static class BatchOutcome {
162
        private final CronBatchStatus status;
163
        private final int totalCount;
164
        private final int successCount;
165
        private final int failureCount;
166
        private final int neverAttemptedCount;
167
 
168
        BatchOutcome(CronBatchStatus status, int totalCount, int successCount,
169
                     int failureCount, int neverAttemptedCount) {
170
            this.status = status;
171
            this.totalCount = totalCount;
172
            this.successCount = successCount;
173
            this.failureCount = failureCount;
174
            this.neverAttemptedCount = neverAttemptedCount;
175
        }
176
 
177
        public CronBatchStatus getStatus() { return status; }
178
        public int getTotalCount() { return totalCount; }
179
        public int getSuccessCount() { return successCount; }
180
        public int getFailureCount() { return failureCount; }
181
        public int getNeverAttemptedCount() { return neverAttemptedCount; }
182
    }
183
 
184
    private void sendFailureEmail(CronBatch batch, List<CronBatchItem> failedItems, int pendingCount) {
36337 amit 185
        StringBuilder body = new StringBuilder();
186
        body.append("Cron job: ").append(batch.getJobName()).append("\n");
37644 amit 187
        body.append("Batch id: ").append(batch.getId()).append("\n");
36337 amit 188
        body.append("Run time: ").append(batch.getStartedAt().format(DTF)).append("\n");
37644 amit 189
        body.append("Outcome: ").append(batch.getStatus()).append("\n");
190
        body.append("Total: ").append(batch.getTotalCount()).append("\n");
36337 amit 191
        body.append("Success: ").append(batch.getSuccessCount()).append("\n");
37644 amit 192
        body.append("Failed: ").append(batch.getFailureCount()).append("\n");
193
        body.append("Never attempted: ").append(pendingCount).append("\n\n");
36337 amit 194
 
37644 amit 195
        if (!failedItems.isEmpty()) {
196
            body.append("Partner failures:\n");
197
            for (CronBatchItem item : failedItems) {
198
                body.append("- ").append(item.getPartnerName())
199
                        .append(" (fofoId: ").append(item.getFofoId()).append(")")
200
                        .append(" — ").append(item.getErrorMessage())
201
                        .append("\n");
202
            }
203
            body.append("\n");
36337 amit 204
        }
205
 
37644 amit 206
        if (pendingCount > 0) {
207
            body.append("The run did not reach these ").append(pendingCount)
208
                    .append(" item(s) — they were never attempted and no work was done for them:\n");
209
            List<CronBatchItem> pendingItems = cronBatchItemRepository
210
                    .selectByBatchIdAndStatus(batch.getId(), CronBatchItemStatus.PENDING);
211
            for (CronBatchItem item : pendingItems.subList(0, Math.min(pendingItems.size(), PENDING_LIST_CAP))) {
212
                body.append("- ").append(item.getPartnerName())
213
                        .append(" (fofoId: ").append(item.getFofoId()).append(")\n");
214
            }
215
            if (pendingItems.size() > PENDING_LIST_CAP) {
216
                body.append("... and ").append(pendingItems.size() - PENDING_LIST_CAP).append(" more\n");
217
            }
218
            body.append("\nRe-run the job to process them.\n");
219
        }
36337 amit 220
 
37644 amit 221
        String subject = pendingCount > 0
222
                ? String.format("[CRON ALERT] %s — INCOMPLETE, %d of %d never attempted",
223
                        batch.getJobName(), pendingCount, batch.getTotalCount())
224
                : String.format("[CRON ALERT] %s — %d of %d partners failed",
225
                        batch.getJobName(), batch.getFailureCount(), batch.getTotalCount());
226
 
36337 amit 227
        try {
228
            mailOutboxService.queueMailViaGoogle(TECHNOLOGY_EMAIL, null, subject, body.toString(),
229
                    "CronBatchService." + batch.getJobName());
230
        } catch (Exception e) {
231
            LOGGER.error("Failed to send batch failure email for batch {}: {}", batch.getId(), e.getMessage());
232
        }
233
    }
234
}