diff --git a/server/laf-cloud/functions/rechargeStats.ts b/server/laf-cloud/functions/rechargeStats.ts new file mode 100644 index 0000000..25e5cdf --- /dev/null +++ b/server/laf-cloud/functions/rechargeStats.ts @@ -0,0 +1,140 @@ +import cloud from '@lafjs/cloud'; + +const DAY = 24 * 60 * 60 * 1000; +const LEGACY_TIME_OFFSET = 8 * 60 * 60 * 1000; +const BATCH_SIZE = 100; + +function emptyStats(asOf: number) { + return { + version: 1, + currency: 'CNY', + 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 asOf = Date.now(); + 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 / 100; + summary.stats.amount30d = summary.fen30d / 100; + summary.stats.amountTotal = summary.fenTotal / 100; + } + const result = await users.bulkWrite(batch.map(user => { + const rechargeStats = summaries.get(user.openid)?.stats ?? { ...emptyStats(asOf), missingOpenid: true }; + return { updateOne: { + // Recheck eligibility and identity, and prevent an older run overwriting a newer snapshot. + filter: { + _id: user._id, + openid: user.openid ?? null, + pay_user: true, + $or: [{ 'rechargeStats.asOf': { $exists: false } }, { 'rechargeStats.asOf': { $lte: asOf } }], + }, + update: { $set: { rechargeStats } }, + } }; + })); + processedUsers += batch.length; + updatedUsers += result.modifiedCount; + afterId = batch[batch.length - 1]._id; + } + + const result = { asOf, processedUsers, updatedUsers }; + console.log('rechargeStats completed', result); + return { code: 1, data: result, msg: '充值统计完成' }; +} diff --git a/server/laf-cloud/functions/rechargeStats.yaml b/server/laf-cloud/functions/rechargeStats.yaml new file mode 100644 index 0000000..db989b9 --- /dev/null +++ b/server/laf-cloud/functions/rechargeStats.yaml @@ -0,0 +1,5 @@ +name: rechargeStats +desc: "付费用户近15天、近30天及累计充值统计;每小时由定时触发器执行" +methods: [] +tags: + - statistics diff --git a/server/laf-cloud/recharge-stats.trigger.json b/server/laf-cloud/recharge-stats.trigger.json new file mode 100644 index 0000000..f17c9d4 --- /dev/null +++ b/server/laf-cloud/recharge-stats.trigger.json @@ -0,0 +1,5 @@ +{ + "desc": "付费用户充值统计(每小时)", + "target": "rechargeStats", + "cron": "0 * * * *" +} diff --git a/server/laf-cloud/rechargeStats.README.md b/server/laf-cloud/rechargeStats.README.md new file mode 100644 index 0000000..6cbf150 --- /dev/null +++ b/server/laf-cloud/rechargeStats.README.md @@ -0,0 +1,73 @@ +# 付费用户充值统计定时任务 + +云函数 `rechargeStats` 每次全量重算当前 `users.pay_user === true` 的玩家,将结果覆盖到 `users.rechargeStats`。默认每小时整点执行;不修改订单和玩家付费标记。 + +## 保存字段 + +| 字段 | 含义 | +| --- | --- | +| `amount15d` | 从本轮运行时刻倒推 15 × 24 小时的充值金额,元 | +| `amount30d` | 从本轮运行时刻倒推 30 × 24 小时的充值金额,元 | +| `amountTotal` | 截至本轮运行时刻、当前库中可统计的累计充值金额,元 | +| `currency` | 固定 `CNY` | +| `asOf` | 本轮统计截止时间,原始 Unix 毫秒时间戳 | +| `orderCount` | 纳入累计金额的去重订单数 | +| `fallbackTimeOrderCount` | 无支付确认时间,使用下单时间的订单数 | +| `missingTimeOrderCount` | 支付和下单时间均缺失,仅计入累计的订单数 | +| `invalidOrderCount` | 缺失订单号、金额或数量无效而跳过的记录数 | +| `duplicateOrderCount` | 跳过的重复订单记录数 | +| `missingOpenid` | 用户是否缺少可关联订单的 openid | +| `version` | 统计规则版本,当前为 1 | + +## 统计口径 + +- 通过 `users.openid = order.openid` 关联,仅纳入数值型 `state: 1/2` 的记录。未确认支付的 `state: 0` 和临时集合 `iosOrder` 不参与。 +- 排除 `paymentAppEnv: "test"` 或 `outTradeNo` 以 `wct_` 开头的测试订单;保留没有环境字段的历史订单。 +- 每个 openid 内按 `outTradeNo` 去重,同号多条按 `_id` 升序取第一条金额有效、时间不在未来的记录。若同号金额不同,应人工核查重复数据。 +- 金额为 `goodsPrice × itemCount`,先用整数分累加,再除以 100 保存为元。金额和数量兼容数字字符串,要求正安全整数;缺失数量不默认当作 1,异常记录计入 `invalidOrderCount`。累计金额超出安全整数范围则报错,不保存失真的数值。 +- 每轮只获取一次当前时间,窗口为 `[asOf - N × 86400000, asOf]`,不是自然日,也不是从午夜倒推。 +- 优先使用 `chargeTime`。当前支付写入代码使用 `new Date(Date.now() + 8小时)`,因此还原时减去 8 小时;此处针对现有存储约定,若以后修正支付时间存储方式,必须同步调整统计规则。 +- `chargeTime` 缺失或无效时回退到原始 `order.time`,不再减 8 小时;两个时间都不可用时只计入累计,并记录异常数量。支持原始毫秒数、Date 和带时区的 ISO 字符串;不解析无时区日期字符串。 +- 已知时间在未来的订单暂不纳入任何金额。没有匹配订单的付费用户也保存零值结果;这不代表迁移前从未充值,结合 `orderCount` 核查历史完整性。 +- 当前项目没有退款同步逻辑,因此这是当前已确认订单的金额统计,不扣除退款,不还原渠道优惠后的实付。订单状态异常仍可能造成偏差。 + +## 执行与一致性 + +用户按 `_id` 游标每批 100 条读取,订单使用 MongoDB 游标遍历,避免默认查询上限截断。每批只更新 `rechargeStats`,写入时重新检查 `pay_user` 和 `openid`,且较早轮次不会覆盖较新 `asOf` 的快照。重复执行不会重复累加;没有新增充值时,旧订单也会按新截止时间退出滚动窗口。 + +整轮不是数据库事务:运行期间新支付或用户变更可能下一轮才体现。失败时抛出错误,已完成批次保留,剩余用户下次重算;每个用户的 `asOf` 可判断新旧结果,不能把混合轮次当作同一时刻的全库快照。 + +建议上线前在 Laf 数据库控制台建立以下非唯一索引(若已有等价索引则复用): + +```javascript +db.collection('users').createIndex({ pay_user: 1, _id: 1 }); +db.collection('order').createIndex({ openid: 1, outTradeNo: 1, _id: 1 }); +``` + +## 发布和启用 + +代码与触发器配置是独立资源,提交本地文件不会自动启动线上定时任务。`recharge-stats.trigger.json` 使用 Laf 创建触发器 API 的 `desc / target / cron` 字段。 + +1. 在目标 Laf 应用发布 `functions/rechargeStats.ts` 和同名 YAML,保持 `methods: []`,不开放公共 HTTP 入口。 +2. 在云函数控制台手动执行一次 `rechargeStats`(无参数),检查返回的 `processedUsers`、`updatedUsers` 和部分玩家的结果,确认运行耗时。 +3. 在触发器面板绑定函数 `rechargeStats`,使用 `recharge-stats.trigger.json` 的配置:每小时整点执行。 +4. 检查已有触发器,避免重复创建。若使用已登录且已绑定正确应用的 Laf CLI,可执行: + +```shell +laf trigger list +laf trigger create "付费用户充值统计(每小时)" rechargeStats "0 * * * *" +``` + +5. 首次定时触发后确认日志出现 `rechargeStats completed`,抽查 `users.rechargeStats.asOf` 已更新。大规模历史数据上线前应确认全量耗时在云函数执行时限内。 + +参考:[Laf 定时任务文档](https://doc.laf.run/zh/cloud-function/cron.html)、[官方 CLI 触发器命令](https://github.com/labring/laf/blob/main/cli/src/command/trigger/index.ts)、[创建触发器字段](https://github.com/labring/laf/blob/main/server/src/trigger/dto/create-trigger.dto.ts)。 + +## 本地测试 + +Node.js 24.11 或兼容 `registerHooks` / `stripTypeScriptTypes` 的版本: + +```shell +node --test laf-cloud/tests/recharge-stats.test.mjs +``` + +覆盖滚动时间边界、8 小时修正、支付时间优先、时间回退、无效金额、数量、去重、测试订单排除、超过 1000 条订单、多页用户、重跑、并发写保护、异常中断和关闭 HTTP 入口。使用内存 MongoDB 接口替身,不连接生产数据库。 diff --git a/server/laf-cloud/tests/recharge-stats.test.mjs b/server/laf-cloud/tests/recharge-stats.test.mjs new file mode 100644 index 0000000..68ad47a --- /dev/null +++ b/server/laf-cloud/tests/recharge-stats.test.mjs @@ -0,0 +1,222 @@ +import test from 'node:test'; +import assert from 'node:assert/strict'; +import { registerHooks, stripTypeScriptTypes } from 'node:module'; +import { readFileSync } from 'node:fs'; + +const DAY = 86400000; +const OFFSET = 8 * 3600000; +const NOW = Date.UTC(2026, 8, 15, 12); +const state = { users: [], order: [], writes: 0, closed: 0, beforeWrite: null, failRead: false }; +const realNow = Date.now; +Date.now = () => NOW; +test.after(() => { Date.now = realNow; }); + +function matches(row, query) { + return Object.entries(query).every(([key, expected]) => { + if (key === '$or') return expected.some(part => matches(row, part)); + const actual = key.split('.').reduce((value, part) => value?.[part], row); + if (expected === null) return actual == null; + if (typeof expected !== 'object') return actual === expected; + return Object.entries(expected).every(([op, value]) => { + if (op === '$in') return value.includes(actual); + if (op === '$ne') return actual !== value; + if (op === '$gt') return actual > value; + if (op === '$lte') return actual <= value; + if (op === '$exists') return (actual !== undefined) === value; + if (op === '$not') return !value.test(actual); + throw new Error(`Unsupported query operator: ${op}`); + }); + }); +} + +globalThis.__rechargeStatsCloud = { mongo: { db: { + collection(name) { + assert.ok(['users', 'order'].includes(name), 'must not count iosOrder temporary records'); + return { + find(query) { + let rows = state[name].filter(row => matches(row, query)); + return { + sort(sort) { + rows.sort((a, b) => { + for (const key of Object.keys(sort)) { + if (a[key] < b[key]) return -1; + if (a[key] > b[key]) return 1; + } + return 0; + }); + return this; + }, + limit(size) { rows = rows.slice(0, size); return this; }, + async toArray() { return rows.map(row => ({ ...row })); }, + async *[Symbol.asyncIterator]() { + if (state.failRead) throw new Error('database read failed'); + yield* rows; + }, + async close() { state.closed++; }, + }; + }, + async bulkWrite(operations) { + state.writes++; + state.beforeWrite?.(); + let modifiedCount = 0; + for (const { updateOne } of operations) { + const user = state.users.find(row => matches(row, updateOne.filter)); + if (user) { + Object.assign(user, structuredClone(updateOne.update.$set)); + modifiedCount++; + } + } + return { modifiedCount }; + }, + }; + }, +} } }; +registerHooks({ resolve(specifier, context, next) { + if (specifier === '@lafjs/cloud') return { + url: 'data:text/javascript,export default globalThis.__rechargeStatsCloud', shortCircuit: true, + }; + return next(specifier, context); +}, load(url, context, next) { + if (url === new URL('../functions/rechargeStats.ts', import.meta.url).href) { + return { format: 'module', source: stripTypeScriptTypes(readFileSync(new URL(url), 'utf8')), shortCircuit: true }; + } + return next(url, context); +} }); +const { default: run } = await import('../functions/rechargeStats.ts'); + +function reset() { + Object.assign(state, { + users: [{ _id: 'u1', openid: 'o1', pay_user: true }], order: [], + writes: 0, closed: 0, beforeWrite: null, failRead: false, + }); +} +function order(id, age, extra = {}) { + return { + _id: id, outTradeNo: id, openid: 'o1', state: 2, goodsPrice: 100, itemCount: 1, + chargeTime: new Date(NOW - age + OFFSET), time: NOW - age - DAY, ...extra, + }; +} + +test('rolling windows include exact boundaries, use payment time, restore +8h and exclude future payments', async () => { + reset(); + state.order = [ + order('now', 0), order('15d', 15 * DAY), order('15d-old', 15 * DAY + 1), + order('30d', 30 * DAY), order('30d-old', 30 * DAY + 1), order('future', -1), + ]; + await run(); + const stats = state.users[0].rechargeStats; + assert.equal(stats.amount15d, 2); + assert.equal(stats.amount30d, 4); + assert.equal(stats.amountTotal, 5); + assert.equal(stats.orderCount, 5); + assert.equal(stats.asOf, NOW); +}); + +test('counts quantities, numeric strings, paid states and deduplicates; excludes unpaid and test orders', async () => { + reset(); + state.order = [ + order('a', DAY, { goodsPrice: '101', itemCount: '3', state: 1 }), + order('b', DAY, { goodsPrice: 7, itemCount: 1 }), + order('duplicate', DAY, { outTradeNo: 'a', goodsPrice: 101, itemCount: 3 }), + order('pending', DAY, { state: 0 }), order('string-state', DAY, { state: '2' }), + order('test-env', DAY, { paymentAppEnv: 'test' }), order('wct_test', DAY), + ]; + await run(); + assert.equal(state.users[0].rechargeStats.amountTotal, 3.10); + assert.equal(state.users[0].rechargeStats.orderCount, 2); + assert.equal(state.users[0].rechargeStats.duplicateOrderCount, 1); +}); + +test('falls back to unshifted creation time, tracks unknown times, and rejects malformed amounts', async () => { + reset(); + state.order = [ + order('fallback', 0, { chargeTime: 0, time: NOW - 15 * DAY }), + order('unknown', 0, { chargeTime: 0, time: null }), + order('iso', 0, { chargeTime: new Date(NOW - 30 * DAY + OFFSET).toISOString() }), + order('bad-price', 0, { goodsPrice: 'bad' }), order('bad-count', 0, { itemCount: null }), + order('negative', 0, { goodsPrice: -1 }), order('missing-id', 0, { outTradeNo: '' }), + order('overflow', 0, { goodsPrice: Number.MAX_SAFE_INTEGER, itemCount: 2 }), + ]; + await run(); + const stats = state.users[0].rechargeStats; + assert.equal(stats.amount15d, 1); + assert.equal(stats.amount30d, 2); + assert.equal(stats.amountTotal, 3); + assert.equal(stats.fallbackTimeOrderCount, 1); + assert.equal(stats.missingTimeOrderCount, 1); + assert.equal(stats.invalidOrderCount, 5); +}); + +test('processes every page, restricts to boolean true, resets old totals, and flags missing openid', async () => { + reset(); + state.users = Array.from({ length: 251 }, (_, i) => ({ + _id: `u${String(i).padStart(4, '0')}`, openid: `o${i}`, pay_user: true, + rechargeStats: { asOf: NOW - DAY, amountTotal: 999 }, + })); + state.users.push({ _id: 'no-openid', pay_user: true }); + state.users.push({ _id: 'free', openid: 'free', pay_user: false }); + state.users.push({ _id: 'string', openid: 'string', pay_user: 'true' }); + state.order = [order('last', 0, { openid: 'o250' }), order('free', 0, { openid: 'free' })]; + const result = await run(); + assert.equal(result.data.processedUsers, 252); + assert.equal(state.writes, 3); + assert.equal(state.users.find(user => user.openid === 'o250').rechargeStats.amountTotal, 1); + assert.equal(state.users.find(user => user.openid === 'o0').rechargeStats.amountTotal, 0); + assert.equal(state.users.find(user => user._id === 'no-openid').rechargeStats.missingOpenid, true); + assert.equal(state.users.find(user => user._id === 'free').rechargeStats, undefined); + assert.equal(state.users.find(user => user._id === 'string').rechargeStats, undefined); +}); + +test('reruns overwrite instead of accumulating; older runs cannot overwrite newer snapshots', async () => { + reset(); + state.order = [order('paid', 0)]; + await run(); + await run(); + assert.equal(state.users[0].rechargeStats.amountTotal, 1); + state.users[0].rechargeStats = { asOf: NOW + 1, amountTotal: 2 }; + await run(); + assert.equal(state.users[0].rechargeStats.amountTotal, 2); +}); + +test('streams more than 1000 orders and expired windows clear on the next run', async () => { + reset(); + state.order = Array.from({ length: 1005 }, (_, i) => order(`paid-${i}`, 30 * DAY, { goodsPrice: 1 })); + await run(); + assert.equal(state.users[0].rechargeStats.amountTotal, 10.05); + assert.equal(state.users[0].rechargeStats.amount30d, 10.05); + try { + Date.now = () => NOW + 1; + await run(); + assert.equal(state.users[0].rechargeStats.amount30d, 0); + assert.equal(state.users[0].rechargeStats.amountTotal, 10.05); + } finally { Date.now = () => NOW; } +}); + +test('rechecks user eligibility and openid when writing', async () => { + reset(); + state.beforeWrite = () => { state.users[0].pay_user = false; }; + await run(); + assert.equal(state.users[0].rechargeStats, undefined); + reset(); + state.beforeWrite = () => { state.users[0].openid = 'changed'; }; + await run(); + assert.equal(state.users[0].rechargeStats, undefined); +}); + +test('read failure propagates, closes cursor and never replaces existing totals with zero', async () => { + reset(); + state.users[0].rechargeStats = { asOf: NOW - 1, amountTotal: 9 }; + state.failRead = true; + await assert.rejects(run(), /database read failed/); + assert.equal(state.closed, 1); + assert.equal(state.writes, 0); + assert.equal(state.users[0].rechargeStats.amountTotal, 9); +}); + +test('timer function does not expose a public HTTP endpoint', () => { + const config = readFileSync(new URL('../functions/rechargeStats.yaml', import.meta.url), 'utf8'); + assert.match(config, /methods: \[\]/); + const trigger = JSON.parse(readFileSync(new URL('../recharge-stats.trigger.json', import.meta.url), 'utf8')); + assert.equal(trigger.target, 'rechargeStats'); + assert.equal(trigger.cron, '0 * * * *'); +});