Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions src/modules/admin/sequencer.controllers.ts
Original file line number Diff line number Diff line change
@@ -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);
}
};
8 changes: 8 additions & 0 deletions src/modules/admin/sequencer.routes.ts
Original file line number Diff line number Diff line change
@@ -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;
2 changes: 2 additions & 0 deletions src/modules/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand All @@ -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;
124 changes: 124 additions & 0 deletions src/utils/creator-sequencer.utils.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>((resolve) => {
setTimeout(() => {
executionOrder.push(1);
resolve();
}, 50);
});

const op2 = () =>
new Promise<void>((resolve) => {
executionOrder.push(2);
resolve();
});

const op3 = () =>
new Promise<void>((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<void>((resolve) => {
active.push('A-start');
setTimeout(() => {
active.push('A-end');
resolve();
}, 50);
});

const opB = () =>
new Promise<void>((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<void>((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();
});
});
134 changes: 134 additions & 0 deletions src/utils/creator-sequencer.utils.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
import { logger } from './logger.utils';

type QueuedOperation<T> = {
execute: () => Promise<T>;
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<unknown>[] = [];
private processing = false;

async enqueue<T>(operation: () => Promise<T>): Promise<T> {
return new Promise<T>((resolve, reject) => {
this.queue.push({
execute: operation as () => Promise<unknown>,
resolve: resolve as (value: unknown) => void,
reject,
enqueuedAt: Date.now(),
});
this.drain();
});
}

get pending(): number {
return this.queue.length;
}

private async drain(): Promise<void> {
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<string, AsyncQueue>();
private gcTimers = new Map<string, ReturnType<typeof setTimeout>>();

async enqueue<T>(
creatorWallet: string,
operation: () => Promise<T>
): Promise<T> {
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 };
78 changes: 78 additions & 0 deletions src/utils/sequencer-lock.utils.ts
Original file line number Diff line number Diff line change
@@ -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<void> }> {
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<typeof setInterval> | 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`
);
}
Loading
Loading