57 lines
2.7 KiB
TypeScript
57 lines
2.7 KiB
TypeScript
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;
|
|
const cursors = db.collection('passCheckJobCursors');
|
|
// This key intentionally ignores the legacy circular cursor: the first
|
|
// run backfills from the beginning once, without deleting old job state.
|
|
const key = { _id: 'season:users' };
|
|
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) return { count };
|
|
const now = Date.now();
|
|
if (cursor.leaseUntil > now) return { count };
|
|
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) return { count };
|
|
const owned = { ...key, revision: cursor.revision + 1 };
|
|
try {
|
|
const rows = await db.collection('users').find(next.after != null ? { _id: { $gt: next.after } } : {})
|
|
.sort({ _id: 1 }).limit(BATCH_SIZE).toArray();
|
|
for (const user of rows) { await freezeUser(user, 'users'); 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 };
|
|
}
|