fix(worker): lease premium reconciliation queue entries (#2296)

This commit is contained in:
Hampus
2026-09-01 00:14:08 +02:00
committed by GitHub
parent 2ea2e79f6f
commit 7d710d881a
7 changed files with 220 additions and 5 deletions
@@ -88,6 +88,7 @@ export interface IKVProvider {
refillIntervalMs: number,
): Promise<number>;
scheduleBulkDeletion(queueKey: string, secondaryKey: string, score: number, value: string): Promise<void>;
claimBulkDeletion(queueKey: string, member: string, maxScore: number, leaseScore: number): Promise<boolean>;
removeBulkDeletion(queueKey: string, secondaryKey: string, member?: string): Promise<boolean>;
scan(pattern: string, count: number): Promise<Array<string>>;
dequeuePurgeBatch(
+24
View File
@@ -223,6 +223,17 @@ tokens = tokens - #urls
redis.call('SET', bucketKey, cjson.encode({tokens = tokens, lastRefill = lastRefill}), 'EX', 3600)
return cjson.encode({urls = urls, tokens = #urls})
`;
const CLAIM_BULK_DELETION_SCRIPT = `
local score = redis.call('ZSCORE', KEYS[1], ARGV[1])
if not score then
return 0
end
if tonumber(score) > tonumber(ARGV[2]) then
return 0
end
redis.call('ZADD', KEYS[1], ARGV[3], ARGV[1])
return 1
`;
const REMOVE_BULK_DELETION_SCRIPT = `
local member = ARGV[1]
if member ~= '' and redis.call('ZREM', KEYS[1], member) == 1 then
@@ -613,6 +624,19 @@ export class KVClient implements IKVProvider {
);
}
async claimBulkDeletion(queueKey: string, member: string, maxScore: number, leaseScore: number): Promise<boolean> {
const result = await this.executeScript(
'claimBulkDeletion',
CLAIM_BULK_DELETION_SCRIPT,
1,
queueKey,
member,
maxScore,
leaseScore,
);
return Number(result) === 1;
}
async removeBulkDeletion(queueKey: string, secondaryKey: string, member = ''): Promise<boolean> {
const result = await this.executeScript(
'removeBulkDeletion',
@@ -139,6 +139,12 @@ describe('KVClient script execution', () => {
keyCount: 2,
run: async (client) => client.scheduleBulkDeletion('queue:key', 'secondary:key', 1, 'value'),
},
{
name: 'claimBulkDeletion',
reply: 1,
keyCount: 1,
run: async (client) => client.claimBulkDeletion('queue:key', 'member', 1, 2),
},
{
name: 'removeBulkDeletion',
reply: 1,
@@ -29,6 +29,16 @@ export class PremiumStateReconciliationQueueService {
}
}
async claimUser(userId: UserID, nowMs: number, leaseUntilMs: number): Promise<boolean> {
try {
const value = this.serializeQueueValue(userId);
return await this.kvClient.claimBulkDeletion(QUEUE_KEY, value, nowMs, leaseUntilMs);
} catch (error) {
Logger.error({error, userId: userId.toString()}, 'Failed to claim user for premium state reconciliation');
throw error;
}
}
async removeUser(userId: UserID): Promise<boolean> {
try {
const secondaryKey = this.getSecondaryKey(userId);
@@ -155,6 +155,7 @@ export class MockKVProvider implements IKVProvider {
readonly checkLeakyBucketLimitSpy = vi.fn();
readonly tryConsumeTokensSpy = vi.fn();
readonly scheduleBulkDeletionSpy = vi.fn();
readonly claimBulkDeletionSpy = vi.fn();
readonly removeBulkDeletionSpy = vi.fn();
readonly scanSpy = vi.fn();
readonly dequeuePurgeBatchSpy = vi.fn();
@@ -667,6 +668,17 @@ export class MockKVProvider implements IKVProvider {
this.expiries.delete(secondaryKey);
}
async claimBulkDeletion(queueKey: string, member: string, maxScore: number, leaseScore: number): Promise<boolean> {
this.claimBulkDeletionSpy(queueKey, member, maxScore, leaseScore);
this.evictIfExpired(queueKey);
const score = this.zsetStore.get(queueKey)?.get(member);
if (score === undefined || score > maxScore) {
return false;
}
await this.zadd(queueKey, leaseScore, member);
return true;
}
async removeBulkDeletion(queueKey: string, secondaryKey: string, member = ''): Promise<boolean> {
this.removeBulkDeletionSpy(queueKey, secondaryKey, member);
this.evictIfExpired(secondaryKey);
@@ -1113,6 +1125,7 @@ export class MockKVProvider implements IKVProvider {
this.checkLeakyBucketLimitSpy.mockClear();
this.tryConsumeTokensSpy.mockClear();
this.scheduleBulkDeletionSpy.mockClear();
this.claimBulkDeletionSpy.mockClear();
this.removeBulkDeletionSpy.mockClear();
this.scanSpy.mockClear();
this.dequeuePurgeBatchSpy.mockClear();
@@ -26,6 +26,7 @@ interface ReconcileResult {
const MAX_USERS_PER_RUN = 250;
const RETRY_DELAY_MS = 5 * 60 * 1000;
const CLAIM_LEASE_MS = 10 * 60 * 1000;
function getStripeSubscriptionCustomerId(subscription: Stripe.Subscription): string | null {
if (!subscription.customer) {
@@ -305,14 +306,18 @@ const processPremiumStateReconciliationQueue: WorkerTaskHandler = async (_payloa
let strippedNoSubscriptionCount = 0;
let failedCount = 0;
let requeuedCount = 0;
let claimedElsewhereCount = 0;
for (const userId of readyUserIds) {
const claimedAt = Date.now();
let claimed = false;
try {
await premiumStateReconciliationQueueService.removeUser(userId);
claimed = await premiumStateReconciliationQueueService.claimUser(userId, claimedAt, claimedAt + CLAIM_LEASE_MS);
} catch (error) {
Logger.warn(
{error, userId: userId.toString()},
'Failed to remove user from premium reconciliation queue before processing',
);
Logger.warn({error, userId: userId.toString()}, 'Failed to claim premium reconciliation queue entry');
}
if (!claimed) {
claimedElsewhereCount += 1;
continue;
}
try {
const result = await reconcileUserPremiumStateFromStripe({
@@ -344,6 +349,7 @@ const processPremiumStateReconciliationQueue: WorkerTaskHandler = async (_payloa
} else {
skippedCount += 1;
}
await premiumStateReconciliationQueueService.removeUser(userId);
} catch (error) {
failedCount += 1;
Logger.error(
@@ -371,6 +377,7 @@ const processPremiumStateReconciliationQueue: WorkerTaskHandler = async (_payloa
strippedNoSubscriptionCount,
failedCount,
requeuedCount,
claimedElsewhereCount,
},
'Finished processing premium reconciliation queue',
);
@@ -0,0 +1,154 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {WorkerTaskHelpers} from '@pkgs/worker/src/contracts/WorkerTask';
import type Stripe from 'stripe';
import {afterEach, describe, expect, test} from 'vitest';
import {createUserID} from '../../BrandedTypes';
import {PremiumStateReconciliationQueueService} from '../../infrastructure/PremiumStateReconciliationQueueService';
import {MockKVProvider} from '../../test/mocks/MockKVProvider';
import {NoopLogger} from '../../test/mocks/NoopLogger';
import type {UserRepository} from '../../user/repositories/UserRepository';
import processPremiumStateReconciliationQueue from '../tasks/ProcessPremiumStateReconciliationQueue';
import {clearWorkerDependencies, setWorkerDependenciesForTest} from '../WorkerContext';
const USER_ID = createUserID(834271905123471361n);
const ONE_HOUR_MS = 60 * 60 * 1000;
function createHelpers(): WorkerTaskHelpers {
return {
logger: new NoopLogger(),
jobId: 1n,
addJob: async () => 0n,
reportProgress: async () => {},
shouldCancel: async () => false,
setContextLink: async () => {},
};
}
function createQueueService(): PremiumStateReconciliationQueueService {
return new PremiumStateReconciliationQueueService(new MockKVProvider());
}
describe('processPremiumStateReconciliationQueue', () => {
afterEach(() => {
clearWorkerDependencies();
});
test('keeps the user in the queue when the worker dies mid-reconciliation', async () => {
const queueService = createQueueService();
await queueService.enqueueUser(USER_ID, new Date(Date.now() - 1000));
let signalReconcileStarted: () => void = () => {};
const reconcileStarted = new Promise<void>((resolve) => {
signalReconcileStarted = resolve;
});
const userRepository = {
findUnique: async () => {
signalReconcileStarted();
return await new Promise<never>(() => {});
},
} as unknown as UserRepository;
setWorkerDependenciesForTest({
premiumStateReconciliationQueueService: queueService,
stripe: {} as Stripe,
userRepository,
});
void processPremiumStateReconciliationQueue({}, createHelpers());
await reconcileStarted;
expect(await queueService.getQueueSize()).toBe(1);
expect(await queueService.getReadyUserIds(Date.now(), 10)).toEqual([]);
expect(await queueService.getReadyUserIds(Date.now() + ONE_HOUR_MS, 10)).toEqual([USER_ID]);
});
test('removes the user from the queue once reconciliation commits', async () => {
const queueService = createQueueService();
await queueService.enqueueUser(USER_ID, new Date(Date.now() - 1000));
const userRepository = {
findUnique: async () => null,
} as unknown as UserRepository;
setWorkerDependenciesForTest({
premiumStateReconciliationQueueService: queueService,
stripe: {} as Stripe,
userRepository,
});
await processPremiumStateReconciliationQueue({}, createHelpers());
expect(await queueService.getQueueSize()).toBe(0);
expect(await queueService.getReadyUserIds(Date.now() + ONE_HOUR_MS, 10)).toEqual([]);
});
test('rejects a second claim until the lease expires', async () => {
const queueService = createQueueService();
await queueService.enqueueUser(USER_ID, new Date(Date.now() - 1000));
const now = Date.now();
expect(await queueService.claimUser(USER_ID, now, now + ONE_HOUR_MS)).toBe(true);
expect(await queueService.claimUser(USER_ID, now, now + ONE_HOUR_MS)).toBe(false);
expect(await queueService.claimUser(USER_ID, now + ONE_HOUR_MS, now + 2 * ONE_HOUR_MS)).toBe(true);
});
test('reconciles the user once when two runs claim the same entry', async () => {
const queueService = createQueueService();
await queueService.enqueueUser(USER_ID, new Date(Date.now() - 1000));
let signalFirstClaimEntered: () => void = () => {};
const firstClaimEntered = new Promise<void>((resolve) => {
signalFirstClaimEntered = resolve;
});
let releaseFirstClaim: () => void = () => {};
const firstClaimReleased = new Promise<void>((resolve) => {
releaseFirstClaim = resolve;
});
const claimUser = queueService.claimUser.bind(queueService);
let firstClaim = true;
queueService.claimUser = async (userId, nowMs, leaseUntilMs) => {
if (firstClaim) {
firstClaim = false;
signalFirstClaimEntered();
await firstClaimReleased;
}
return await claimUser(userId, nowMs, leaseUntilMs);
};
let signalReconcileStarted: () => void = () => {};
const reconcileStarted = new Promise<void>((resolve) => {
signalReconcileStarted = resolve;
});
let releaseReconcile: () => void = () => {};
const reconcileReleased = new Promise<void>((resolve) => {
releaseReconcile = resolve;
});
let findUniqueCalls = 0;
const userRepository = {
findUnique: async () => {
findUniqueCalls += 1;
signalReconcileStarted();
await reconcileReleased;
return null;
},
} as unknown as UserRepository;
setWorkerDependenciesForTest({
premiumStateReconciliationQueueService: queueService,
stripe: {} as Stripe,
userRepository,
});
const firstRun = processPremiumStateReconciliationQueue({}, createHelpers());
await firstClaimEntered;
const secondRun = processPremiumStateReconciliationQueue({}, createHelpers());
await reconcileStarted;
releaseFirstClaim();
releaseReconcile();
await Promise.all([firstRun, secondRun]);
expect(findUniqueCalls).toBe(1);
expect(await queueService.getQueueSize()).toBe(0);
});
});