diff --git a/laf-cloud/docs/pass-check-settlement.txt b/laf-cloud/docs/pass-check-settlement.txt index 9f3ffa0..350ee1e 100644 --- a/laf-cloud/docs/pass-check-settlement.txt +++ b/laf-cloud/docs/pass-check-settlement.txt @@ -1,4 +1,4 @@ -战令结束结算接入说明(2026-10-08) +战令结束结算接入说明(2026-10-09) 时间契约 沿用现有 passCheck/read 返回的 time:毫秒级周期起点,结束点为 time + 30 天。 @@ -10,10 +10,31 @@ passCheckTime、wx/checkIos。passCheckUpgrade 补齐了 POST YAML。 2. 按 functions/pass-check.trigger.json 创建/更新每分钟触发器,目标 passCheckJobs。 passCheckJobs 不开放 HTTP 方法;需要现有 PASSCHECKTIME_ID 环境配置。 + 每分钟仅用于检查周期和继续未完成的批次;本期扫描完成后不再读取用户表。 3. 再发布配套客户端。旧客户端的过期进度写入将被拒绝,不能继续按旧方式领上期奖励。 4. 不清空用户历史。快照存入 users/usersAd 的 passSettlements, - passSettlementRevision 用于 CAS。后台每分钟每个用户集合扫描100条,游标保存在 - passCheckJobCursors。请按用户规模评估扫描延迟;玩家请求时同样惰性补建快照。 + passSettlementRevision 用于 CAS。后台每期每个用户集合只完成一轮分页扫描, + 每次触发各处理最多100条;进度保存在 passCheckJobCursors。 + 请按用户规模评估扫描延迟;玩家请求时同样惰性补建快照,不必等待后台扫到自己。 + +按期扫描与旧任务升级 +- 当前全局周期起点同时是上一期结束点,以该毫秒时间戳标识扫描期。 + 即使在线接口先推进了全局周期,后台仍会对比自己的 periodEnd,启动新一期扫描。 +- 新进度键为 season:users、season:usersAd;保存 periodEnd、after、completed、 + completedAt、revision、leaseUntil。仅使用已有集合的 _id 索引,无需手动初始化。 + 保留旧 users/usersAd 游标,但不再读取它们;升级后的首次运行从头补扫一次历史存档。 + 此后只有全局周期推进才开启新一轮;配置的起点还在未来时不扫描。 +- 每批成功后推进游标,不足100条时标记完成;恰好整批时下一次空页确认完成。 + 两个集合分别完成。全部完成后每次触发仅读取周期配置和两个进度文档,不扫描用户, + 不反复更新游标;云函数每分钟的触发次数本身未减少。 +- 每批取得5分钟租约,并用 revision 比较更新防止重叠执行或旧执行覆盖新进度。 + 批次失败保留起始游标、释放租约,下次重试;进程直接退出时租约到期后恢复。 + 失败重试或租约超时接管可能重复读取该批,但已有快照不重复生成,也不自动发奖。 +- 跨期仍未扫完时先续完旧扫描,再开启最新一期。停机错过多期时无需为每个中间期 + 重扫全服:freezeUser 一次补齐该玩家存档中所有可验证的到期记录。 + 扫描完成后新增或导入的历史玩家记录,由其在线请求补建,或在下一期扫描时补齐。 +- 已部署旧版时,本次只需重新发布 passCheckJobs 并更新触发器描述;cron 仍为 + * * * * *。不需要改客户端、奖励领取接口,不删除用户快照或旧游标。 协议 POST passCheckSettlement,使用现有 uid/token/gameName 鉴权。 @@ -41,7 +62,9 @@ claimFlags 仅更新实际领取的 free/passCheck 索引为0,不修改冻结 需人工退款处理,不自动退款。 验证 -Node 24:node --test laf-cloud/tests/pass-check-settlement.test.mjs(6项通过)。 +Node 24:node --test laf-cloud/tests/pass-check-settlement.test.mjs(15项通过)。 +覆盖完成后零用户扫描/零写入、分页、新周期、空集合/整批边界、失败重试、旧游标升级、 +跨期续跑、租约超时接管与旧执行防回退、在线补建、未来周期配置;保留原发奖幂等测试。 配套客户端 tools/test-battle-pass-settlement.cjs 覆盖奖励表、UI关闭=领取、 失败重试、确认弹窗、优先队列、多期顺序和支付幂等。 现有 rookie-gift.test.mjs 在导入阶段失败:goldMiner/config 缺少 paymentProductId diff --git a/laf-cloud/functions/pass-check.trigger.json b/laf-cloud/functions/pass-check.trigger.json index c54d00a..08621f5 100644 --- a/laf-cloud/functions/pass-check.trigger.json +++ b/laf-cloud/functions/pass-check.trigger.json @@ -1,5 +1,5 @@ { - "desc": "战令到期快照与离线玩家补扫,每分钟,不自动入账", + "desc": "每分钟检查战令周期,每期补扫一次,完成停止,不自动入账", "target": "passCheckJobs", "cron": "* * * * *" } diff --git a/laf-cloud/functions/passCheckJobs.ts b/laf-cloud/functions/passCheckJobs.ts index ea36ea5..a13e672 100644 --- a/laf-cloud/functions/passCheckJobs.ts +++ b/laf-cloud/functions/passCheckJobs.ts @@ -1,18 +1,58 @@ import cloud from '@lafjs/cloud'; import { currentPassStart, freezeUser } from '@/passCheckSettlement'; -/** Bounded durable scans also backfill players who were offline when this job was deployed. */ +const BATCH_SIZE = 100; +const LEASE_MS = 5 * 60000; + +/** One resumable sweep per period; completed periods never restart their user scan. */ export default async function () { - await currentPassStart(); + // The current period's start is the previous period's end, even if an online + // request has already advanced the shared clock before this job runs. + const periodEnd = await currentPassStart(); + if (periodEnd > Date.now()) return { count: 0 }; const db: any = cloud.mongo.db; let count = 0; for (const table of ['users', 'usersAd']) { const cursors = db.collection('passCheckJobCursors'); - const cursor = await cursors.findOne({ _id: table }); - const rows = await db.collection(table).find(cursor?.after ? { _id: { $gt: cursor.after } } : {}).sort({ _id: 1 }).limit(100).toArray(); - // A failed row leaves the cursor unchanged, so the next run retries the batch. - for (const user of rows) { await freezeUser(user, table); count++; } - await cursors.updateOne({ _id: table }, { $set: { after: rows.length === 100 ? rows[rows.length - 1]._id : null } }, { upsert: true }); + // Separate keys intentionally ignore legacy circular cursors: the first + // run backfills from the beginning once, without deleting old job state. + const key = { _id: 'season:' + table }; + let cursor = await cursors.findOne(key); + if (!cursor) { + try { + await cursors.updateOne(key, { $setOnInsert: { + periodEnd, after: null, completed: false, revision: 0, leaseUntil: 0 + } }, { upsert: true }); + } catch (error: any) { if (error.code !== 11000) throw error; } + cursor = await cursors.findOne(key); + } + if (cursor.completed && cursor.periodEnd >= periodEnd) continue; + const now = Date.now(); + if (cursor.leaseUntil > now) continue; + const next = cursor.completed + ? { periodEnd, after: null, completed: false, completedAt: null } + : { periodEnd: cursor.periodEnd, after: cursor.after, completed: false }; + const locked = await cursors.updateOne({ ...key, revision: cursor.revision, leaseUntil: { $lte: now } }, { + $set: { ...next, leaseUntil: now + LEASE_MS }, $inc: { revision: 1 } + }); + if (!locked.modifiedCount) continue; + const owned = { ...key, revision: cursor.revision + 1 }; + try { + const rows = await db.collection(table).find(next.after != null ? { _id: { $gt: next.after } } : {}) + .sort({ _id: 1 }).limit(BATCH_SIZE).toArray(); + for (const user of rows) { await freezeUser(user, table); count++; } + const completed = rows.length < BATCH_SIZE; + await cursors.updateOne(owned, { + $set: { after: rows.length ? rows[rows.length - 1]._id : next.after, completed, + completedAt: completed ? Date.now() : null, leaseUntil: 0 }, + $inc: { revision: 1 } + }); + } catch (error) { + // Keep this batch's starting cursor. Freezing is idempotent on retry; + // revision fencing prevents an expired worker overwriting its successor. + await cursors.updateOne(owned, { $set: { leaseUntil: 0 }, $inc: { revision: 1 } }); + throw error; + } } return { count }; } diff --git a/laf-cloud/functions/passCheckJobs.yaml b/laf-cloud/functions/passCheckJobs.yaml index d12d317..cb95374 100644 --- a/laf-cloud/functions/passCheckJobs.yaml +++ b/laf-cloud/functions/passCheckJobs.yaml @@ -1,4 +1,4 @@ name: passCheckJobs -desc: 战令到期批量冻结与离线玩家补扫 +desc: 战令每期一次分批冻结,完成后停止扫描 methods: [] tags: [] diff --git a/laf-cloud/tests/pass-check-settlement.test.mjs b/laf-cloud/tests/pass-check-settlement.test.mjs index b23e5f9..8ff2731 100644 --- a/laf-cloud/tests/pass-check-settlement.test.mjs +++ b/laf-cloud/tests/pass-check-settlement.test.mjs @@ -3,22 +3,48 @@ import assert from 'node:assert/strict'; import { registerHooks } from 'node:module'; const NOW = 1800000000000, END = NOW - 1000, PERIOD = 30 * 86400000; -let now = NOW, drop = false; -const tables = { users: [], usersAd: [], order: [], idcount: [] }; +let now = NOW, drop = false, failFreeze = '', beforeFind = null; +const scans = [], writes = []; +const tables = { users: [], usersAd: [], order: [], idcount: [], passCheckJobCursors: [] }; const clone = v => v == null ? v : structuredClone(v); const get = (r, k) => k.split('.').reduce((v, key) => v?.[key], r); -const match = (r, q) => Object.entries(q).every(([k, v]) => k === '$expr' - ? now < new Date(v.$lt[1]).getTime() : v === null ? get(r, k) == null : JSON.stringify(get(r, k)) === JSON.stringify(v)); +const match = (r, q) => Object.entries(q).every(([k, v]) => { + if (k === '$expr') return now < new Date(v.$lt[1]).getTime(); + if (v === null) return get(r, k) == null; + if (v && typeof v === 'object' && '$gt' in v) return get(r, k) > v.$gt; + if (v && typeof v === 'object' && '$lte' in v) return get(r, k) <= v.$lte; + return JSON.stringify(get(r, k)) === JSON.stringify(v); +}); const collection = name => ({ async findOne(q) { return clone(tables[name].find(r => match(r, q))); }, - async updateOne(q, update) { - const row = tables[name].find(r => match(r, q)); + find(q) { + let limit = Infinity; + const query = { + sort() { return query; }, + limit(value) { limit = value; return query; }, + async toArray() { + scans.push({ name, q: clone(q) }); + if (beforeFind) await beforeFind(name); + return clone(tables[name].filter(r => match(r, q)).sort((a, b) => a._id < b._id ? -1 : a._id > b._id ? 1 : 0).slice(0, limit)); + } + }; + return query; + }, + async updateOne(q, update, options = {}) { + writes.push({ name, q: clone(q), update: clone(update) }); + if (q._id === failFreeze && update.$set?.passSettlements) throw Error('freeze failed'); + let row = tables[name].find(r => match(r, q)); + if (!row && options.upsert) { + row = { ...clone(q), ...clone(update.$setOnInsert || {}) }; + tables[name].push(row); + } if (!row) return { matchedCount: 0, modifiedCount: 0 }; for (const [k, v] of Object.entries(update.$set || {})) { const keys = k.split('.'); let target = row; for (const part of keys.slice(0, -1)) target = target[part] ||= {}; target[keys.at(-1)] = clone(v); } + for (const [k, v] of Object.entries(update.$inc || {})) row[k] = (row[k] || 0) + v; if (drop && update.$set.timestamp) { drop = false; throw Error('lost response'); } return { matchedCount: 1, modifiedCount: 1 }; } @@ -39,11 +65,13 @@ registerHooks({ resolve(specifier, context, next) { } }); const { default: settlement, currentPassStart } = await import('../functions/passCheckSettlement.ts'); const { default: upgrade, passOrderFields } = await import('../functions/passCheckUpgrade.ts'); +const { default: runJob } = await import('../functions/passCheckJobs.ts'); const realNow = Date.now; Date.now = () => now; test.after(() => { Date.now = realNow; }); function reset() { - now = NOW; drop = false; + now = NOW; drop = false; failFreeze = ''; beforeFind = null; + scans.length = 0; writes.length = 0; for (const key of Object.keys(tables)) tables[key] = []; process.env.PASSCHECKTIME_ID = 'clock'; tables.idcount.push({ _id: 'clock', passcheckTime: END - PERIOD }); @@ -109,3 +137,141 @@ test('unauthenticated requests and expired ordinary writes cannot alter the snap assert.equal(result.code, 410); assert.equal(u.passSettlements[END].stage.experience, 4); }); + +const jobCursor = table => tables.passCheckJobCursors.find(r => r._id === 'season:' + table); +function populate(count) { + const original = reset(); + tables.users = Array.from({ length: count }, (_, i) => ({ ...clone(original), _id: 'u' + String(i).padStart(4, '0') })); +} + +test('job completes each collection once and idle ticks never scan users or write cursors', async () => { + const u = reset(); + tables.usersAd.push({ ...clone(u), _id: 'ad' }); + assert.equal((await runJob()).count, 2); + for (const table of ['users', 'usersAd']) { + assert.equal(jobCursor(table).completed, true); + assert.equal(jobCursor(table).periodEnd, END); + assert.ok(tables[table][0].passSettlements[END]); + assert.equal(tables[table][0].coinAmount, 10); + } + scans.length = 0; writes.length = 0; + for (let i = 0; i < 5; i++) { now += 60000; assert.equal((await runJob()).count, 0); } + assert.deepEqual(scans, []); + assert.deepEqual(writes, []); +}); + +test('job resumes batches without rewinding and only opens a new sweep at the next boundary', async () => { + populate(201); + assert.equal((await runJob()).count, 100); + assert.equal(jobCursor('users').after, 'u0099'); + assert.equal((await runJob()).count, 100); + assert.equal(jobCursor('users').after, 'u0199'); + assert.equal((await runJob()).count, 1); + assert.equal(jobCursor('users').completed, true); + now = END + PERIOD - 1; + assert.equal((await runJob()).count, 0); + now++; + // Online reads can advance the global clock first; this must not hide a new season. + await currentPassStart(); + const u = tables.users[0]; + u.passCheck = JSON.stringify({ 2: { time: String(now), experience: 8, free: [1] } }); + assert.equal((await runJob()).count, 100); + assert.equal(jobCursor('users').periodEnd, now); + assert.equal(jobCursor('users').after, 'u0099'); + assert.equal(u.passSettlements[now].stage.experience, 8); + assert.equal(u.coinAmount, 10); +}); + +test('empty collections and exact full batches terminate instead of restarting', async () => { + populate(100); + assert.equal((await runJob()).count, 100); + assert.equal(jobCursor('users').completed, false); + assert.equal(jobCursor('usersAd').completed, true); + assert.equal((await runJob()).count, 0); + assert.equal(jobCursor('users').completed, true); + scans.length = 0; + await runJob(); + assert.deepEqual(scans, []); +}); + +test('failed batch preserves its cursor, releases the lease and retries idempotently', async () => { + populate(102); + await runJob(); + failFreeze = 'u0101'; + await assert.rejects(runJob(), /freeze failed/); + assert.equal(jobCursor('users').after, 'u0099'); + assert.equal(jobCursor('users').completed, false); + assert.equal(jobCursor('users').leaseUntil, 0); + assert.ok(tables.users[100].passSettlements[END]); + const revision = tables.users[100].passSettlementRevision; + failFreeze = ''; + assert.equal((await runJob()).count, 2); + assert.equal(jobCursor('users').completed, true); + assert.equal(tables.users[100].passSettlementRevision, revision); +}); + +test('legacy circular cursors are preserved but never used to skip the initial backfill', async () => { + reset(); + tables.passCheckJobCursors.push({ _id: 'users', after: 'z' }); + await runJob(); + assert.ok(tables.users[0].passSettlements[END]); + assert.deepEqual(tables.passCheckJobCursors.find(r => r._id === 'users'), { _id: 'users', after: 'z' }); +}); + +test('cross-period recovery finishes the old sweep before starting the latest period', async () => { + populate(101); + await runJob(); + now = END + 3 * PERIOD; + assert.equal((await runJob()).count, 1); + assert.equal(jobCursor('users').periodEnd, END); + assert.equal(jobCursor('users').completed, true); + assert.equal((await runJob()).count, 100); + assert.equal(jobCursor('users').periodEnd, now); +}); + +test('overlapping workers skip leased batches and cannot overwrite newer progress after expiry', async () => { + populate(101); + let resume, started; + const paused = new Promise(resolve => { started = resolve; }); + beforeFind = async name => { + if (name !== 'users') return; + beforeFind = null; + started(); + await new Promise(resolve => { resume = resolve; }); + }; + const stale = runJob(); + await paused; + assert.equal((await runJob()).count, 0); + now += 5 * 60000 + 1; + assert.equal((await runJob()).count, 100); + assert.equal((await runJob()).count, 1); + assert.equal(jobCursor('users').completed, true); + const finished = clone(jobCursor('users')); + resume(); + await stale; + assert.deepEqual(jobCursor('users'), finished); + scans.length = 0; + await runJob(); + assert.deepEqual(scans, []); +}); + +test('online reads still freeze late-arriving records after the period scan is complete', async () => { + reset(); + tables.users = []; + await runJob(); + tables.users.push({ _id: 'u', token: 't', coinAmount: 10, + passCheck: JSON.stringify({ 2: { time: String(END), experience: 4 } }) }); + scans.length = 0; + await runJob(); + assert.deepEqual(scans, []); + assert.equal((await call({ action: 'read' })).data.rows[0].end, END); + assert.equal(tables.users[0].coinAmount, 10); +}); + +test('a future configured start never begins a settlement scan', async () => { + reset(); + tables.idcount[0].passcheckTime = NOW + PERIOD; + assert.equal((await runJob()).count, 0); + assert.deepEqual(scans, []); + assert.deepEqual(tables.passCheckJobCursors, []); +});