diff --git a/src/modules/admin/sequencer.controllers.ts b/src/modules/admin/sequencer.controllers.ts new file mode 100644 index 0000000..6f01582 --- /dev/null +++ b/src/modules/admin/sequencer.controllers.ts @@ -0,0 +1,24 @@ +import { AsyncController } from '../../types/auth.types'; +import { clearDrift } from '../../utils/supply-drift-guard.utils'; +import { sendSuccess, sendValidationError } from '../../utils/api-response.utils'; + +export const httpClearDrift: AsyncController = async (req, res, next) => { + try { + const rawParam = req.params.creatorWallet; + const creatorWallet = Array.isArray(rawParam) ? rawParam[0] : rawParam; + if (!creatorWallet) { + sendValidationError(res, 'Missing creatorWallet parameter'); + return; + } + + await clearDrift(creatorWallet); + + sendSuccess(res, { + creatorWallet, + status: 'cleared', + message: `Supply drift flag cleared for ${creatorWallet}`, + }); + } catch (err) { + next(err); + } +}; diff --git a/src/modules/admin/sequencer.routes.ts b/src/modules/admin/sequencer.routes.ts new file mode 100644 index 0000000..fa6c605 --- /dev/null +++ b/src/modules/admin/sequencer.routes.ts @@ -0,0 +1,8 @@ +import { Router } from 'express'; +import { httpClearDrift } from './sequencer.controllers'; + +const sequencerRouter = Router(); + +sequencerRouter.post('/sequencer/clear-drift/:creatorWallet', httpClearDrift); + +export default sequencerRouter; diff --git a/src/modules/index.ts b/src/modules/index.ts index 9696cc5..d7b6b04 100644 --- a/src/modules/index.ts +++ b/src/modules/index.ts @@ -14,6 +14,7 @@ import webhookRouter from './webhooks/webhook.router'; import walletsRouter from './wallets/wallets.routes'; import alertsRouter from './alerts/alert.router'; import tradingRouter from './trading/multi-buy.routes'; +import sequencerRouter from './admin/sequencer.routes'; import { BASE as CREATORS_BASE } from '../constants/creator.constants'; import { routeBodySizeLimit } from '../middlewares/body-size-limit.middleware'; @@ -38,5 +39,6 @@ router.use(CREATORS_BASE, routeBodySizeLimit('creators'), webhookRouter); router.use('/wallets', routeBodySizeLimit('default'), walletsRouter); router.use('/alerts', routeBodySizeLimit('default'), alertsRouter); router.use('/trading', routeBodySizeLimit('default'), tradingRouter); +router.use('/internal', routeBodySizeLimit('default'), sequencerRouter); export default router; diff --git a/src/utils/creator-sequencer.utils.test.ts b/src/utils/creator-sequencer.utils.test.ts new file mode 100644 index 0000000..166f18b --- /dev/null +++ b/src/utils/creator-sequencer.utils.test.ts @@ -0,0 +1,124 @@ +import { + CreatorSequencer, + SequencerTimeoutError, +} from './creator-sequencer.utils'; + +jest.mock('./logger.utils', () => ({ + logger: { + debug: jest.fn(), + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + }, +})); + +describe('CreatorSequencer (#758)', () => { + let sequencer: CreatorSequencer; + + beforeEach(() => { + sequencer = new CreatorSequencer(); + }); + + afterEach(() => { + sequencer.dispose(); + jest.restoreAllMocks(); + }); + + it('serialises operations for the same creator in FIFO order', async () => { + const executionOrder: number[] = []; + + const op1 = () => + new Promise((resolve) => { + setTimeout(() => { + executionOrder.push(1); + resolve(); + }, 50); + }); + + const op2 = () => + new Promise((resolve) => { + executionOrder.push(2); + resolve(); + }); + + const op3 = () => + new Promise((resolve) => { + executionOrder.push(3); + resolve(); + }); + + const p1 = sequencer.enqueue('creator-wallet-1', op1); + const p2 = sequencer.enqueue('creator-wallet-1', op2); + const p3 = sequencer.enqueue('creator-wallet-1', op3); + + await Promise.all([p1, p2, p3]); + + expect(executionOrder).toEqual([1, 2, 3]); + }); + + it('allows concurrent operations for different creators', async () => { + const active: string[] = []; + + const opA = () => + new Promise((resolve) => { + active.push('A-start'); + setTimeout(() => { + active.push('A-end'); + resolve(); + }, 50); + }); + + const opB = () => + new Promise((resolve) => { + active.push('B-start'); + setTimeout(() => { + active.push('B-end'); + resolve(); + }, 10); + }); + + const pA = sequencer.enqueue('creator-A', opA); + const pB = sequencer.enqueue('creator-B', opB); + + await Promise.all([pA, pB]); + + expect(active).toContain('A-start'); + expect(active).toContain('B-start'); + }); + + it('rejects queued operations with SequencerTimeoutError if waiting > 10s', async () => { + let resolveSlowOp: () => void = () => {}; + const slowOp = () => + new Promise((resolve) => { + resolveSlowOp = resolve; + }); + + const queuedOp = () => Promise.resolve('ok'); + + const p1 = sequencer.enqueue('creator-slow', slowOp); + const p2 = sequencer.enqueue('creator-slow', queuedOp); + + const realNow = Date.now; + jest.spyOn(Date, 'now').mockReturnValue(realNow() + 11_000); + + resolveSlowOp(); + + await expect(p2).rejects.toThrow(SequencerTimeoutError); + await p1; + }); + + it('cleans up inactive queues after 5 minutes', async () => { + jest.useFakeTimers(); + + const op = () => Promise.resolve('done'); + await sequencer.enqueue('creator-idle', op); + + expect(sequencer.activeQueues).toBe(1); + + jest.advanceTimersByTime(5 * 60 * 1000 + 100); + + expect(sequencer.activeQueues).toBe(0); + + jest.useRealTimers(); + }); +}); diff --git a/src/utils/creator-sequencer.utils.ts b/src/utils/creator-sequencer.utils.ts new file mode 100644 index 0000000..84d5b4e --- /dev/null +++ b/src/utils/creator-sequencer.utils.ts @@ -0,0 +1,134 @@ +import { logger } from './logger.utils'; + +type QueuedOperation = { + execute: () => Promise; + resolve: (value: T) => void; + reject: (reason: unknown) => void; + enqueuedAt: number; +}; + +const QUEUE_TIMEOUT_MS = 10_000; +const IDLE_GC_MS = 5 * 60 * 1000; + +class AsyncQueue { + private queue: QueuedOperation[] = []; + private processing = false; + + async enqueue(operation: () => Promise): Promise { + return new Promise((resolve, reject) => { + this.queue.push({ + execute: operation as () => Promise, + resolve: resolve as (value: unknown) => void, + reject, + enqueuedAt: Date.now(), + }); + this.drain(); + }); + } + + get pending(): number { + return this.queue.length; + } + + private async drain(): Promise { + if (this.processing) return; + this.processing = true; + + while (this.queue.length > 0) { + const item = this.queue.shift()!; + const waited = Date.now() - item.enqueuedAt; + + if (waited > QUEUE_TIMEOUT_MS) { + item.reject( + new SequencerTimeoutError( + `Operation waited ${waited}ms in queue, exceeding ${QUEUE_TIMEOUT_MS}ms limit` + ) + ); + continue; + } + + try { + const result = await item.execute(); + item.resolve(result); + } catch (err) { + item.reject(err); + } + } + + this.processing = false; + } +} + +export class SequencerTimeoutError extends Error { + public readonly code = 'sequencer_timeout'; + + constructor(message: string) { + super(message); + this.name = 'SequencerTimeoutError'; + } +} + +class CreatorSequencer { + private queues = new Map(); + private gcTimers = new Map>(); + + async enqueue( + creatorWallet: string, + operation: () => Promise + ): Promise { + this.resetGcTimer(creatorWallet); + + let queue = this.queues.get(creatorWallet); + if (!queue) { + queue = new AsyncQueue(); + this.queues.set(creatorWallet, queue); + logger.debug( + { creator_wallet: creatorWallet }, + 'Creator sequencer queue created' + ); + } + + return queue.enqueue(operation); + } + + private resetGcTimer(creatorWallet: string): void { + const existing = this.gcTimers.get(creatorWallet); + if (existing) { + clearTimeout(existing); + } + + const timer = setTimeout(() => { + const queue = this.queues.get(creatorWallet); + if (queue && queue.pending === 0) { + this.queues.delete(creatorWallet); + this.gcTimers.delete(creatorWallet); + logger.debug( + { creator_wallet: creatorWallet }, + 'Creator sequencer queue garbage collected after inactivity' + ); + } + }, IDLE_GC_MS); + + if (timer.unref) { + timer.unref(); + } + + this.gcTimers.set(creatorWallet, timer); + } + + get activeQueues(): number { + return this.queues.size; + } + + dispose(): void { + for (const timer of this.gcTimers.values()) { + clearTimeout(timer); + } + this.gcTimers.clear(); + this.queues.clear(); + } +} + +export const creatorSequencer = new CreatorSequencer(); + +export { CreatorSequencer }; diff --git a/src/utils/sequencer-lock.utils.ts b/src/utils/sequencer-lock.utils.ts new file mode 100644 index 0000000..04e9652 --- /dev/null +++ b/src/utils/sequencer-lock.utils.ts @@ -0,0 +1,78 @@ +import { getRedis } from './redis.utils'; +import { logger } from './logger.utils'; + +const LOCK_TTL_SECONDS = 15; +const LOCK_ACQUIRE_TIMEOUT_MS = 8_000; +const LOCK_RENEWAL_INTERVAL_MS = 5_000; +const LOCK_RETRY_DELAY_MS = 100; + +export class SequencerContentionError extends Error { + public readonly code = 'sequencer_contention'; + + constructor(message: string) { + super(message); + this.name = 'SequencerContentionError'; + } +} + +function lockKey(creatorWallet: string): string { + return `seq_lock:${creatorWallet}`; +} + +export async function acquireSequencerLock( + creatorWallet: string +): Promise<{ release: () => Promise }> { + const redis = getRedis(); + const key = lockKey(creatorWallet); + const lockValue = `${process.pid}:${Date.now()}`; + const deadline = Date.now() + LOCK_ACQUIRE_TIMEOUT_MS; + + while (Date.now() < deadline) { + const acquired = await redis.set(key, lockValue, 'EX', LOCK_TTL_SECONDS, 'NX'); + + if (acquired === 'OK') { + let renewalInterval: ReturnType | null = null; + + renewalInterval = setInterval(async () => { + try { + const current = await redis.get(key); + if (current === lockValue) { + await redis.expire(key, LOCK_TTL_SECONDS); + } + } catch { + // renewal failure is non-fatal; the TTL provides a safety net + } + }, LOCK_RENEWAL_INTERVAL_MS); + + if (renewalInterval.unref) { + renewalInterval.unref(); + } + + const release = async () => { + if (renewalInterval) { + clearInterval(renewalInterval); + renewalInterval = null; + } + try { + const current = await redis.get(key); + if (current === lockValue) { + await redis.del(key); + } + } catch (err) { + logger.warn( + { creator_wallet: creatorWallet, err }, + 'Failed to release sequencer lock' + ); + } + }; + + return { release }; + } + + await new Promise((resolve) => setTimeout(resolve, LOCK_RETRY_DELAY_MS)); + } + + throw new SequencerContentionError( + `Could not acquire sequencer lock for ${creatorWallet} within ${LOCK_ACQUIRE_TIMEOUT_MS}ms` + ); +} diff --git a/src/utils/supply-drift-guard.utils.test.ts b/src/utils/supply-drift-guard.utils.test.ts new file mode 100644 index 0000000..bce6f08 --- /dev/null +++ b/src/utils/supply-drift-guard.utils.test.ts @@ -0,0 +1,76 @@ +import { + verifySupplyAndGuard, + isDriftHalted, + assertNoSupplyDrift, + clearDrift, + SupplyDriftHaltedError, +} from './supply-drift-guard.utils'; + +jest.mock('./logger.utils', () => ({ + logger: { + debug: jest.fn(), + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + }, +})); + +jest.mock('./redis.utils', () => { + const store = new Map(); + return { + getRedis: () => ({ + exists: jest.fn().mockImplementation(async (key: string) => (store.has(key) ? 1 : 0)), + set: jest.fn().mockImplementation(async (key: string, val: string) => { + store.set(key, val); + return 'OK'; + }), + del: jest.fn().mockImplementation(async (key: string) => { + store.delete(key); + return 1; + }), + _store: store, + }), + }; +}); + +describe('Supply Drift Guard (#758)', () => { + const creatorWallet = 'GCREATOR12345'; + + beforeEach(async () => { + await clearDrift(creatorWallet); + }); + + it('returns true when expected supply matches actual supply', async () => { + const match = await verifySupplyAndGuard(creatorWallet, 10, 10); + expect(match).toBe(true); + + const halted = await isDriftHalted(creatorWallet); + expect(halted).toBe(false); + }); + + it('detects supply drift, sets Redis flag, and returns false when supplies diverge', async () => { + const match = await verifySupplyAndGuard(creatorWallet, 10, 12); + expect(match).toBe(false); + + const halted = await isDriftHalted(creatorWallet); + expect(halted).toBe(true); + }); + + it('assertNoSupplyDrift throws SupplyDriftHaltedError when drift flag is set', async () => { + await verifySupplyAndGuard(creatorWallet, 10, 15); + + await expect(assertNoSupplyDrift(creatorWallet)).rejects.toThrow( + SupplyDriftHaltedError + ); + }); + + it('clearDrift removes the drift flag and allows operations to resume', async () => { + await verifySupplyAndGuard(creatorWallet, 10, 15); + expect(await isDriftHalted(creatorWallet)).toBe(true); + + await clearDrift(creatorWallet); + + expect(await isDriftHalted(creatorWallet)).toBe(false); + await expect(assertNoSupplyDrift(creatorWallet)).resolves.not.toThrow(); + }); +}); diff --git a/src/utils/supply-drift-guard.utils.ts b/src/utils/supply-drift-guard.utils.ts new file mode 100644 index 0000000..90604e2 --- /dev/null +++ b/src/utils/supply-drift-guard.utils.ts @@ -0,0 +1,63 @@ +import { getRedis } from './redis.utils'; +import { logger } from './logger.utils'; + +export class SupplyDriftHaltedError extends Error { + public readonly code = 'supply_drift_halted'; + + constructor(message: string) { + super(message); + this.name = 'SupplyDriftHaltedError'; + } +} + +function driftKey(creatorWallet: string): string { + return `drift:${creatorWallet}`; +} + +export async function isDriftHalted(creatorWallet: string): Promise { + const redis = getRedis(); + const exists = await redis.exists(driftKey(creatorWallet)); + return exists === 1; +} + +export async function assertNoSupplyDrift(creatorWallet: string): Promise { + const halted = await isDriftHalted(creatorWallet); + if (halted) { + throw new SupplyDriftHaltedError( + `Operations halted for creator ${creatorWallet} due to detected supply drift` + ); + } +} + +export async function verifySupplyAndGuard( + creatorWallet: string, + expectedSupply: number, + actualSupply: number +): Promise { + if (expectedSupply !== actualSupply) { + logger.error( + { + event: 'supply_drift_detected', + creator_wallet: creatorWallet, + expected_supply: expectedSupply, + actual_supply: actualSupply, + }, + 'Supply drift detected! Operations halted for creator until cleared.' + ); + + const redis = getRedis(); + await redis.set(driftKey(creatorWallet), '1'); + return false; + } + + return true; +} + +export async function clearDrift(creatorWallet: string): Promise { + const redis = getRedis(); + await redis.del(driftKey(creatorWallet)); + logger.info( + { creator_wallet: creatorWallet }, + 'Supply drift flag cleared for creator' + ); +}