Rev 37589 | Blame | Compare with Previous | Last modification | View Log | RSS feed
package com.smartdukaan.cron.scheduled;import com.spice.profitmandi.dao.entity.transaction.CronBatch;import com.spice.profitmandi.service.cron.CronBatchService;import com.spice.profitmandi.service.offers.OfferBatchService;import com.spice.profitmandi.service.pricing.PriceDropBatchService;import com.spice.profitmandi.service.transaction.PartnerLimitUpdateData;import org.apache.logging.log4j.LogManager;import org.apache.logging.log4j.Logger;import org.springframework.beans.factory.annotation.Autowired;import org.springframework.stereotype.Component;import java.util.LinkedHashMap;import java.util.List;/*** Batch-aware scheduled tasks with per-partner transaction isolation.* NO @Transactional at class level — each partner gets its own transaction* via REQUIRES_NEW in the helper beans.** Tracks every run in cron_batch / cron_batch_item tables.* Sends failure summary email on partial failures.*/@Componentpublic class BatchScheduledTasks {private static final Logger LOGGER = LogManager.getLogger(BatchScheduledTasks.class);/*** A RUNNING batch older than this was abandoned by a run that died before finalizing.* Generous on purpose: the longest real batch (offer processing over ~1,500 partners)* is minutes, so 6h can only catch a genuinely dead run.*/private static final int STALE_BATCH_HOURS = 6;@Autowiredprivate CronBatchService cronBatchService;@Autowiredprivate OfferBatchService offerBatchService;@Autowiredprivate PartnerLimitHelper partnerLimitHelper;@Autowiredprivate PriceDropBatchService priceDropBatchService;@Autowiredprivate com.spice.profitmandi.service.PartnerInvestmentSweepService partnerInvestmentSweepService;/*** Recomputes every active partner's investment snapshot, then recomputes the credit limit for* the partners whose base actually moved.** <p>Runs every 2 minutes. Measured on production: ~3.5s of batched reads for ~980 partners,* with an average of 3 partners (peak 18) showing a changed base per window — so the limit work* that follows is a handful of rows, not a full scan.** <p>Investment and utilisation are read in the same pass deliberately. Reading them at* different vintages is what let an advance payment be counted twice — once as a cached wallet* balance and once as the loan repayment it funded — and overstate a limit by the payment x the* partner's tier.*/public void sweepPartnerInvestment() throws Exception {List<Integer> changedBases;try {changedBases = partnerInvestmentSweepService.sweep();} catch (Exception e) {LOGGER.error("Investment sweep failed: {}", e.getMessage(), e);return;}if (changedBases.isEmpty()) {return;}updatePartnerLimitWithBatch(changedBases);}/** Daily refresh of the time-driven aged-Apple and live-demo haircuts. */public void refreshAgedStockDaily() {try {partnerInvestmentSweepService.refreshAgedStockDaily();} catch (Exception e) {LOGGER.error("Aged stock daily refresh failed: {}", e.getMessage(), e);}}/*** Closes out batches left RUNNING by a run that died mid-loop (the cron JVM is* kernel-killed on this box often enough that this is routine, not exotic).** Until this existed such a batch stayed RUNNING forever, which mattered twice over:* it never sent its failure email, and — for offer batches — hasUnfinishedBatch()* refused every later run, so the job could never be retried. Two were stuck in prod* for months on that path.** The loop lives HERE and not inside CronBatchService on purpose: a self-invoked* finalizeBatch() would bypass the Spring proxy, losing REQUIRES_NEW and running with* no transaction at all, so the writes would be silently discarded.*/public void reapStaleBatches() {List<Integer> staleBatchIds = cronBatchService.selectStaleRunningBatchIds(STALE_BATCH_HOURS);if (staleBatchIds.isEmpty()) {return;}LOGGER.warn("Reaping {} stale RUNNING batch(es) older than {}h: {}",staleBatchIds.size(), STALE_BATCH_HOURS, staleBatchIds);for (Integer batchId : staleBatchIds) {try {cronBatchService.finalizeBatch(batchId);} catch (Exception e) {LOGGER.error("Failed to reap stale batch {}: {}", batchId, e.getMessage(), e);}}}/*** CLI entrypoint for cron: delegates to shared OfferBatchService (also used by web/fofo controllers).*/public void processOfferWithBatch(int offerId) throws Exception {offerBatchService.processOfferWithBatch(offerId);}/*** Scheduled entrypoint: delegates to shared PriceDropBatchService.* Reprocesses (rejects) price drops for IMEIs activated before the drop date,* each drop in its own REQUIRES_NEW transaction so user_wallet locks are held* per-drop instead of for the whole run.*/public void reprocessPriceDropsWithBatch() {priceDropBatchService.reprocessPriceDropsWithBatch();}/*** Recalculates partner credit limits. Only writes to partners where values actually changed.** Flow:* 1. Read phase: calculate limits for all 1,500 partners, compare with current values* 2. Create batch with only changed partners (~50-100 typically)* 3. Per-partner: update in REQUIRES_NEW transaction* 4. Finalize: counts + failure email*/public void updatePartnerLimitWithBatch() throws Exception {updatePartnerLimitWithBatch(null);}/*** @param restrictTo when non-null, only recalculate these partners. Passed by the investment* sweep with the partners whose base moved.*/public void updatePartnerLimitWithBatch(java.util.Collection<Integer> restrictTo) throws Exception {List<PartnerLimitUpdateData> changedPartners;try {changedPartners = partnerLimitHelper.calculateChangedPartnerLimits(restrictTo);} catch (Exception e) {LOGGER.error("Failed to calculate partner limits: {}", e.getMessage());return;}if (changedPartners.isEmpty()) {LOGGER.info("No partner limits changed, skipping");return;}LinkedHashMap<Integer, String> fofoIdPartnerNameMap = new LinkedHashMap<>();for (PartnerLimitUpdateData data : changedPartners) {fofoIdPartnerNameMap.put(data.getFofoId(), "fofo-" + data.getFofoId());}CronBatch batch = cronBatchService.createBatch("updatePartnerLimit", fofoIdPartnerNameMap);for (PartnerLimitUpdateData data : changedPartners) {try {partnerLimitHelper.updateSinglePartnerLimit(batch.getId(), data);} catch (Exception e) {LOGGER.error("updatePartnerLimit failed for fofoId={}: {}", data.getFofoId(), e.getMessage());cronBatchService.markItemFailed(batch.getId(), data.getFofoId(), e.getMessage());}}cronBatchService.finalizeBatch(batch.getId());}}