import cloud from '@lafjs/cloud'; const DAY = 24 * 60 * 60 * 1000; const LEGACY_TIME_OFFSET = 8 * 60 * 60 * 1000; const REPORT_TIMEZONE_OFFSET = 8 * 60 * 60 * 1000; // Asia/Shanghai const BATCH_SIZE = 100; function emptyStats(asOf: number) { return { version: 4, currency: 'CNY', unit: 'fen', asOf, amount15d: 0, amount30d: 0, amountTotal: 0, orderCount: 0, fallbackTimeOrderCount: 0, missingTimeOrderCount: 0, invalidOrderCount: 0, duplicateOrderCount: 0, missingOpenid: false, }; } function positiveInteger(value: any): number | null { if (typeof value !== 'number' && typeof value !== 'string') return null; const number = Number(value); return Number.isSafeInteger(number) && number > 0 ? number : null; } function timestamp(value: any): number | null { if (value instanceof Date) return positiveInteger(value.getTime()); const numeric = positiveInteger(value); if (numeric !== null) return numeric; // Date strings must include a timezone; do not depend on the server's timezone. if (typeof value === 'string' && /T.*(?:Z|[+-]\d{2}:?\d{2})$/i.test(value)) { return positiveInteger(Date.parse(value)); } return null; } // Laf timer entry point. HTTP methods are disabled in rechargeStats.yaml. export default async function () { const mongo = cloud.mongo.db; const users = mongo.collection('users'); const rechargeStatsCollection = mongo.collection('userRechargeStats'); const updatedAt = Date.now(); // End of yesterday in Beijing time, independent of execution time and server timezone. const asOf = Math.floor((updatedAt + REPORT_TIMEZONE_OFFSET) / DAY) * DAY - REPORT_TIMEZONE_OFFSET - 1; let afterId: any; let processedUsers = 0; let updatedUsers = 0; while (true) { const batch = await users.find({ pay_user: true, ...(afterId === undefined ? {} : { _id: { $gt: afterId } }), }, { projection: { _id: 1, openid: 1 } }).sort({ _id: 1 }).limit(BATCH_SIZE).toArray(); if (!batch.length) break; const summaries = new Map; seen: Set; fen15d: number; fen30d: number; fenTotal: number; }>(); for (const user of batch) { if (typeof user.openid === 'string' && user.openid.trim() && !summaries.has(user.openid)) { summaries.set(user.openid, { stats: emptyStats(asOf), seen: new Set(), fen15d: 0, fen30d: 0, fenTotal: 0 }); } } if (summaries.size) { const orders = mongo.collection('order').find({ openid: { $in: [...summaries.keys()] }, state: { $in: [1, 2] }, paymentAppEnv: { $ne: 'test' }, outTradeNo: { $not: /^wct_/ }, }, { projection: { openid: 1, outTradeNo: 1, goodsPrice: 1, itemCount: 1, chargeTime: 1, time: 1 } }) .sort({ openid: 1, outTradeNo: 1, _id: 1 }); try { for await (const order of orders) { const summary = summaries.get(order.openid)!; const price = positiveInteger(order.goodsPrice); const count = positiveInteger(order.itemCount); const amount = price === null || count === null ? NaN : price * count; if (typeof order.outTradeNo !== 'string' || !order.outTradeNo.trim() || !Number.isSafeInteger(amount)) { summary.stats.invalidOrderCount++; continue; } // Existing payment writers store new Date(Date.now() + 8h). // order.time is the original, unshifted Unix timestamp. const chargeTime = timestamp(order.chargeTime); const paidAt = chargeTime === null ? timestamp(order.time) : chargeTime - LEGACY_TIME_OFFSET; if (paidAt !== null && paidAt > asOf) continue; if (summary.seen.has(order.outTradeNo)) { summary.stats.duplicateOrderCount++; continue; } summary.seen.add(order.outTradeNo); summary.fenTotal += amount; if (!Number.isSafeInteger(summary.fenTotal)) throw new Error('Recharge total exceeds safe integer range'); summary.stats.orderCount++; if (paidAt === null) summary.stats.missingTimeOrderCount++; else { if (chargeTime === null) summary.stats.fallbackTimeOrderCount++; if (paidAt > asOf - 15 * DAY) summary.fen15d += amount; if (paidAt > asOf - 30 * DAY) summary.fen30d += amount; } } } finally { await orders.close(); } } for (const summary of summaries.values()) { summary.stats.amount15d = summary.fen15d; summary.stats.amount30d = summary.fen30d; summary.stats.amountTotal = summary.fenTotal; } const result = await rechargeStatsCollection.bulkWrite(batch.map(user => { const rechargeStats = summaries.get(user.openid)?.stats ?? { ...emptyStats(asOf), missingOpenid: true }; const record = { _id: user._id, openid: user.openid ?? null, ...rechargeStats, updatedAt }; return { updateOne: { // Match only the unique user ID, so a newer snapshot cannot cause an upsert ID conflict. // Apply the freshness check atomically inside the update, including same-day reruns. filter: { _id: user._id }, update: [{ $replaceWith: { $cond: [ { $and: [ { $lte: [{ $ifNull: ['$asOf', 0] }, asOf] }, { $lte: [{ $ifNull: ['$updatedAt', 0] }, updatedAt] }, ] }, { $literal: record }, '$$ROOT', ] } }], upsert: true, } }; })); processedUsers += batch.length; updatedUsers += result.modifiedCount + result.upsertedCount; afterId = batch[batch.length - 1]._id; } const result = { asOf, updatedAt, processedUsers, updatedUsers }; console.log('rechargeStats completed', result); return { code: 1, data: result, msg: '充值统计完成' }; }