import cloud from '@lafjs/cloud'; import { currentPassStart, freezeUser } from '@/passCheckSettlement'; 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 () { // 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'); // 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 }; }