From 7d710d881a202c2c96a41fd488b165d0feb4749c Mon Sep 17 00:00:00 2001 From: Hampus Date: Tue, 1 Sep 2026 00:14:08 +0200 Subject: [PATCH] fix(worker): lease premium reconciliation queue entries (#2296) --- fluxer_api/pkgs/kv_client/src/IKVProvider.ts | 1 + fluxer_api/pkgs/kv_client/src/KVClient.ts | 24 +++ .../kv_client/src/__tests__/KVClient.test.ts | 6 + .../PremiumStateReconciliationQueueService.ts | 10 ++ .../src/api/test/mocks/MockKVProvider.ts | 13 ++ .../ProcessPremiumStateReconciliationQueue.ts | 17 +- ...essPremiumStateReconciliationQueue.test.ts | 154 ++++++++++++++++++ 7 files changed, 220 insertions(+), 5 deletions(-) create mode 100644 fluxer_api/src/api/worker/tests/ProcessPremiumStateReconciliationQueue.test.ts diff --git a/fluxer_api/pkgs/kv_client/src/IKVProvider.ts b/fluxer_api/pkgs/kv_client/src/IKVProvider.ts index a3fe5c437..130b5c993 100644 --- a/fluxer_api/pkgs/kv_client/src/IKVProvider.ts +++ b/fluxer_api/pkgs/kv_client/src/IKVProvider.ts @@ -88,6 +88,7 @@ export interface IKVProvider { refillIntervalMs: number, ): Promise; scheduleBulkDeletion(queueKey: string, secondaryKey: string, score: number, value: string): Promise; + claimBulkDeletion(queueKey: string, member: string, maxScore: number, leaseScore: number): Promise; removeBulkDeletion(queueKey: string, secondaryKey: string, member?: string): Promise; scan(pattern: string, count: number): Promise>; dequeuePurgeBatch( diff --git a/fluxer_api/pkgs/kv_client/src/KVClient.ts b/fluxer_api/pkgs/kv_client/src/KVClient.ts index 292a09a35..46ba39be0 100644 --- a/fluxer_api/pkgs/kv_client/src/KVClient.ts +++ b/fluxer_api/pkgs/kv_client/src/KVClient.ts @@ -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 { + 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 { const result = await this.executeScript( 'removeBulkDeletion', diff --git a/fluxer_api/pkgs/kv_client/src/__tests__/KVClient.test.ts b/fluxer_api/pkgs/kv_client/src/__tests__/KVClient.test.ts index ec1f8c97b..8ddbb46c7 100644 --- a/fluxer_api/pkgs/kv_client/src/__tests__/KVClient.test.ts +++ b/fluxer_api/pkgs/kv_client/src/__tests__/KVClient.test.ts @@ -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, diff --git a/fluxer_api/src/api/infrastructure/PremiumStateReconciliationQueueService.ts b/fluxer_api/src/api/infrastructure/PremiumStateReconciliationQueueService.ts index 8876ea314..5dfad00f4 100644 --- a/fluxer_api/src/api/infrastructure/PremiumStateReconciliationQueueService.ts +++ b/fluxer_api/src/api/infrastructure/PremiumStateReconciliationQueueService.ts @@ -29,6 +29,16 @@ export class PremiumStateReconciliationQueueService { } } + async claimUser(userId: UserID, nowMs: number, leaseUntilMs: number): Promise { + 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 { try { const secondaryKey = this.getSecondaryKey(userId); diff --git a/fluxer_api/src/api/test/mocks/MockKVProvider.ts b/fluxer_api/src/api/test/mocks/MockKVProvider.ts index 5d448c0b1..d1ee83c0a 100644 --- a/fluxer_api/src/api/test/mocks/MockKVProvider.ts +++ b/fluxer_api/src/api/test/mocks/MockKVProvider.ts @@ -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 { + 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 { 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(); diff --git a/fluxer_api/src/api/worker/tasks/ProcessPremiumStateReconciliationQueue.ts b/fluxer_api/src/api/worker/tasks/ProcessPremiumStateReconciliationQueue.ts index 6a07be309..88e857dc7 100644 --- a/fluxer_api/src/api/worker/tasks/ProcessPremiumStateReconciliationQueue.ts +++ b/fluxer_api/src/api/worker/tasks/ProcessPremiumStateReconciliationQueue.ts @@ -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', ); diff --git a/fluxer_api/src/api/worker/tests/ProcessPremiumStateReconciliationQueue.test.ts b/fluxer_api/src/api/worker/tests/ProcessPremiumStateReconciliationQueue.test.ts new file mode 100644 index 000000000..c69f0c07b --- /dev/null +++ b/fluxer_api/src/api/worker/tests/ProcessPremiumStateReconciliationQueue.test.ts @@ -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((resolve) => { + signalReconcileStarted = resolve; + }); + const userRepository = { + findUnique: async () => { + signalReconcileStarted(); + return await new Promise(() => {}); + }, + } 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((resolve) => { + signalFirstClaimEntered = resolve; + }); + let releaseFirstClaim: () => void = () => {}; + const firstClaimReleased = new Promise((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((resolve) => { + signalReconcileStarted = resolve; + }); + let releaseReconcile: () => void = () => {}; + const reconcileReleased = new Promise((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); + }); +});