Subversion Repositories SmartDukaan

Rev

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

Rev Author Line No. Line
36306 amit 1
package com.smartdukaan.cron.scheduled;
2
 
3
import com.spice.profitmandi.dao.entity.transaction.CronBatch;
36338 amit 4
import com.spice.profitmandi.service.cron.CronBatchService;
36343 amit 5
import com.spice.profitmandi.service.offers.OfferBatchService;
37083 amit 6
import com.spice.profitmandi.service.pricing.PriceDropBatchService;
36306 amit 7
import com.spice.profitmandi.service.transaction.PartnerLimitUpdateData;
8
import org.apache.logging.log4j.LogManager;
9
import org.apache.logging.log4j.Logger;
10
import org.springframework.beans.factory.annotation.Autowired;
11
import org.springframework.stereotype.Component;
12
 
13
import java.util.LinkedHashMap;
14
import java.util.List;
15
 
16
/**
17
 * Batch-aware scheduled tasks with per-partner transaction isolation.
18
 * NO @Transactional at class level — each partner gets its own transaction
19
 * via REQUIRES_NEW in the helper beans.
20
 *
21
 * Tracks every run in cron_batch / cron_batch_item tables.
36343 amit 22
 * Sends failure summary email on partial failures.
36306 amit 23
 */
24
@Component
25
public class BatchScheduledTasks {
26
 
27
    private static final Logger LOGGER = LogManager.getLogger(BatchScheduledTasks.class);
28
 
37644 amit 29
    /**
30
     * A RUNNING batch older than this was abandoned by a run that died before finalizing.
31
     * Generous on purpose: the longest real batch (offer processing over ~1,500 partners)
32
     * is minutes, so 6h can only catch a genuinely dead run.
33
     */
34
    private static final int STALE_BATCH_HOURS = 6;
35
 
36306 amit 36
    @Autowired
37
    private CronBatchService cronBatchService;
38
 
39
    @Autowired
36343 amit 40
    private OfferBatchService offerBatchService;
36306 amit 41
 
42
    @Autowired
43
    private PartnerLimitHelper partnerLimitHelper;
44
 
37083 amit 45
    @Autowired
46
    private PriceDropBatchService priceDropBatchService;
47
 
37589 amit 48
    @Autowired
49
    private com.spice.profitmandi.service.PartnerInvestmentSweepService partnerInvestmentSweepService;
50
 
36306 amit 51
    /**
37589 amit 52
     * Recomputes every active partner's investment snapshot, then recomputes the credit limit for
53
     * the partners whose base actually moved.
54
     *
55
     * <p>Runs every 2 minutes. Measured on production: ~3.5s of batched reads for ~980 partners,
56
     * with an average of 3 partners (peak 18) showing a changed base per window — so the limit work
57
     * that follows is a handful of rows, not a full scan.
58
     *
59
     * <p>Investment and utilisation are read in the same pass deliberately. Reading them at
60
     * different vintages is what let an advance payment be counted twice — once as a cached wallet
61
     * balance and once as the loan repayment it funded — and overstate a limit by the payment x the
62
     * partner's tier.
63
     */
64
    public void sweepPartnerInvestment() throws Exception {
65
        List<Integer> changedBases;
66
        try {
67
            changedBases = partnerInvestmentSweepService.sweep();
68
        } catch (Exception e) {
69
            LOGGER.error("Investment sweep failed: {}", e.getMessage(), e);
70
            return;
71
        }
72
        if (changedBases.isEmpty()) {
73
            return;
74
        }
75
        updatePartnerLimitWithBatch(changedBases);
76
    }
77
 
78
    /** Daily refresh of the time-driven aged-Apple and live-demo haircuts. */
79
    public void refreshAgedStockDaily() {
80
        try {
81
            partnerInvestmentSweepService.refreshAgedStockDaily();
82
        } catch (Exception e) {
83
            LOGGER.error("Aged stock daily refresh failed: {}", e.getMessage(), e);
84
        }
85
    }
86
 
87
    /**
37644 amit 88
     * Closes out batches left RUNNING by a run that died mid-loop (the cron JVM is
89
     * kernel-killed on this box often enough that this is routine, not exotic).
90
     *
91
     * Until this existed such a batch stayed RUNNING forever, which mattered twice over:
92
     * it never sent its failure email, and — for offer batches — hasUnfinishedBatch()
93
     * refused every later run, so the job could never be retried. Two were stuck in prod
94
     * for months on that path.
95
     *
96
     * The loop lives HERE and not inside CronBatchService on purpose: a self-invoked
97
     * finalizeBatch() would bypass the Spring proxy, losing REQUIRES_NEW and running with
98
     * no transaction at all, so the writes would be silently discarded.
99
     */
100
    public void reapStaleBatches() {
101
        List<Integer> staleBatchIds = cronBatchService.selectStaleRunningBatchIds(STALE_BATCH_HOURS);
102
        if (staleBatchIds.isEmpty()) {
103
            return;
104
        }
105
        LOGGER.warn("Reaping {} stale RUNNING batch(es) older than {}h: {}",
106
                staleBatchIds.size(), STALE_BATCH_HOURS, staleBatchIds);
107
        for (Integer batchId : staleBatchIds) {
108
            try {
109
                cronBatchService.finalizeBatch(batchId);
110
            } catch (Exception e) {
111
                LOGGER.error("Failed to reap stale batch {}: {}", batchId, e.getMessage(), e);
112
            }
113
        }
114
    }
115
 
116
    /**
36343 amit 117
     * CLI entrypoint for cron: delegates to shared OfferBatchService (also used by web/fofo controllers).
36306 amit 118
     */
119
    public void processOfferWithBatch(int offerId) throws Exception {
36343 amit 120
        offerBatchService.processOfferWithBatch(offerId);
36306 amit 121
    }
122
 
123
    /**
37083 amit 124
     * Scheduled entrypoint: delegates to shared PriceDropBatchService.
125
     * Reprocesses (rejects) price drops for IMEIs activated before the drop date,
126
     * each drop in its own REQUIRES_NEW transaction so user_wallet locks are held
127
     * per-drop instead of for the whole run.
128
     */
129
    public void reprocessPriceDropsWithBatch() {
130
        priceDropBatchService.reprocessPriceDropsWithBatch();
131
    }
132
 
133
    /**
36306 amit 134
     * Recalculates partner credit limits. Only writes to partners where values actually changed.
135
     *
136
     * Flow:
137
     * 1. Read phase: calculate limits for all 1,500 partners, compare with current values
138
     * 2. Create batch with only changed partners (~50-100 typically)
139
     * 3. Per-partner: update in REQUIRES_NEW transaction
140
     * 4. Finalize: counts + failure email
141
     */
142
    public void updatePartnerLimitWithBatch() throws Exception {
37589 amit 143
        updatePartnerLimitWithBatch(null);
144
    }
145
 
146
    /**
147
     * @param restrictTo when non-null, only recalculate these partners. Passed by the investment
148
     * sweep with the partners whose base moved.
149
     */
150
    public void updatePartnerLimitWithBatch(java.util.Collection<Integer> restrictTo) throws Exception {
36306 amit 151
        List<PartnerLimitUpdateData> changedPartners;
152
        try {
37589 amit 153
            changedPartners = partnerLimitHelper.calculateChangedPartnerLimits(restrictTo);
36306 amit 154
        } catch (Exception e) {
155
            LOGGER.error("Failed to calculate partner limits: {}", e.getMessage());
156
            return;
157
        }
158
 
159
        if (changedPartners.isEmpty()) {
160
            LOGGER.info("No partner limits changed, skipping");
161
            return;
162
        }
163
 
164
        LinkedHashMap<Integer, String> fofoIdPartnerNameMap = new LinkedHashMap<>();
165
        for (PartnerLimitUpdateData data : changedPartners) {
166
            fofoIdPartnerNameMap.put(data.getFofoId(), "fofo-" + data.getFofoId());
167
        }
168
 
169
        CronBatch batch = cronBatchService.createBatch("updatePartnerLimit", fofoIdPartnerNameMap);
170
 
171
        for (PartnerLimitUpdateData data : changedPartners) {
172
            try {
173
                partnerLimitHelper.updateSinglePartnerLimit(batch.getId(), data);
174
            } catch (Exception e) {
175
                LOGGER.error("updatePartnerLimit failed for fofoId={}: {}", data.getFofoId(), e.getMessage());
176
                cronBatchService.markItemFailed(batch.getId(), data.getFofoId(), e.getMessage());
177
            }
178
        }
179
 
180
        cronBatchService.finalizeBatch(batch.getId());
181
    }
182
}