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
4 changes: 4 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,10 @@ OWNERSHIP_SNAPSHOT_CLEANUP_INTERVAL_MINUTES=60
DETECT_PRICE_MOVEMENTS_ENABLED=true
DETECT_PRICE_MOVEMENTS_INTERVAL_MINUTES=5

# TWAP computation job (#963) — recomputes 1h/4h/24h TWAP per active key
TWAP_COMPUTATION_ENABLED=true
TWAP_COMPUTATION_INTERVAL_MINUTES=5

# Request body size limits (see docs/body-size-limits.md)
BODY_SIZE_LIMIT_DEFAULT=10mb
# BODY_SIZE_LIMIT_AUTH=100kb
Expand Down
8 changes: 8 additions & 0 deletions src/config.schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -262,6 +262,14 @@ export const envSchema = z
.positive()
.default(5),

// TWAP computation job (#963) — recomputes 1h/4h/24h TWAP per active key.
TWAP_COMPUTATION_ENABLED: booleanCoerce.default(true),
TWAP_COMPUTATION_INTERVAL_MINUTES: z.coerce
.number()
.int()
.positive()
.default(5),

// Governance proposal sync job
GOVERNANCE_SYNC_ENABLED: booleanCoerce.default(false),
GOVERNANCE_SYNC_INTERVAL_MINUTES: z.coerce
Expand Down
29 changes: 29 additions & 0 deletions src/constants/redis.constants.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
// src/constants/redis.constants.ts
// Redis key helpers and TTLs for TWAP price caching (#963).
// Kept out of notifications.constants.ts so price-cache concerns live
// in one place alongside future Redis-backed caches.

export const TWAP_WINDOWS = ['1h', '4h', '24h'] as const;
export type TwapWindow = (typeof TWAP_WINDOWS)[number];

export const TWAP_WINDOW_MS: Record<TwapWindow, number> = {
'1h': 60 * 60 * 1000,
'4h': 4 * 60 * 60 * 1000,
'24h': 24 * 60 * 60 * 1000,
};

// TTL matches the window size so longer windows stay cached longer.
export const TWAP_CACHE_TTL_SECONDS: Record<TwapWindow, number> = {
'1h': 60 * 60,
'4h': 4 * 60 * 60,
'24h': 24 * 60 * 60,
};

export const twapRedisKey = (keyId: string, window: TwapWindow): string =>
`twap:${keyId}:${window}`;

// Stale when the computation job is behind by 2x its 5-minute interval.
export const TWAP_STALE_THRESHOLD_MS = 10 * 60 * 1000;

// Cap on snapshots scanned per TWAP computation (matches price-history cap).
export const TWAP_MAX_SNAPSHOTS = 5000;
98 changes: 98 additions & 0 deletions src/jobs/twap-computation.job.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
// src/jobs/twap-computation.job.test.ts
jest.mock('../config', () => ({
envConfig: {
TWAP_COMPUTATION_ENABLED: true,
TWAP_COMPUTATION_INTERVAL_MINUTES: 5,
},
}));

jest.mock('../utils/prisma.utils', () => ({
prisma: {
creatorProfile: { findMany: jest.fn() },
},
}));

jest.mock('../utils/logger.utils', () => ({
logger: { info: jest.fn(), warn: jest.fn(), error: jest.fn() },
}));

jest.mock('../modules/keys/key-twap.service', () => ({
computeAndCacheTwap: jest.fn(),
KeyNotFoundError: class KeyNotFoundError extends Error {},
}));

jest.mock('../modules/keys/key-registration.service', () => ({
keyEventEmitter: { on: jest.fn(), removeListener: jest.fn() },
}));

import { prisma } from '../utils/prisma.utils';
import { computeAndCacheTwap } from '../modules/keys/key-twap.service';
import {
backfillTwapForKey,
computeTwapForAllKeys,
} from './twap-computation.job';

const mockPrisma = prisma as unknown as {
creatorProfile: { findMany: jest.Mock };
};
const mockCompute = computeAndCacheTwap as jest.Mock;

describe('twap-computation job', () => {
beforeEach(() => {
jest.clearAllMocks();
mockCompute.mockResolvedValue({});
});

it('computes all three windows per active key', async () => {
mockPrisma.creatorProfile.findMany.mockResolvedValue([
{ id: 'key-1' },
{ id: 'key-2' },
]);

const result = await computeTwapForAllKeys();

expect(mockPrisma.creatorProfile.findMany).toHaveBeenCalledWith({
where: { deprecatedAt: null },
select: { id: true },
});
expect(result.scannedKeys).toBe(2);
// 2 keys x 3 windows
expect(mockCompute).toHaveBeenCalledTimes(6);
expect(result.computedWrites).toBe(6);
expect(result.failedWrites).toBe(0);
});

it('counts per-key failures without aborting the run', async () => {
mockPrisma.creatorProfile.findMany.mockResolvedValue([{ id: 'key-1' }]);
mockCompute
.mockResolvedValueOnce({})
.mockRejectedValueOnce(new Error('boom'))
.mockResolvedValueOnce({});

const result = await computeTwapForAllKeys();

expect(result.computedWrites).toBe(2);
expect(result.failedWrites).toBe(1);
});

it('backfills all windows for a single new key', async () => {
await backfillTwapForKey('new-key');

expect(mockCompute).toHaveBeenCalledTimes(3);
expect(mockCompute).toHaveBeenCalledWith(
'new-key',
'1h',
expect.any(Date)
);
expect(mockCompute).toHaveBeenCalledWith(
'new-key',
'4h',
expect.any(Date)
);
expect(mockCompute).toHaveBeenCalledWith(
'new-key',
'24h',
expect.any(Date)
);
});
});
160 changes: 160 additions & 0 deletions src/jobs/twap-computation.job.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,160 @@
// src/jobs/twap-computation.job.ts
// Background TWAP computation for all active creator keys (#963).
// Runs every 5 minutes per key per window (1h, 4h, 24h) and pre-warms
// the Redis cache so GET /keys/:keyId/price/twap stays a cache hit.

import { envConfig } from '../config';
import { logger } from '../utils/logger.utils';
import { prisma } from '../utils/prisma.utils';
import { TWAP_WINDOWS } from '../constants/redis.constants';
import {
computeAndCacheTwap,
KeyNotFoundError,
} from '../modules/keys/key-twap.service';
import { keyEventEmitter } from '../modules/keys/key-registration.service';

export type TwapComputationResult = {
scannedKeys: number;
computedWrites: number;
failedWrites: number;
};

/**
* Compute and cache TWAP for every active key (deprecatedAt null)
* across all windows. Missing keys are covered on each run, which is
* also the backfill path for newly created keys.
*/
export async function computeTwapForAllKeys(
now: Date = new Date()
): Promise<TwapComputationResult> {
const activeKeys = await prisma.creatorProfile.findMany({
where: { deprecatedAt: null },
select: { id: true },
});

let computedWrites = 0;
let failedWrites = 0;

for (const key of activeKeys as Array<{ id: string }>) {
for (const window of TWAP_WINDOWS) {
try {
await computeAndCacheTwap(key.id, window, now);
computedWrites += 1;
} catch (error) {
failedWrites += 1;
logger.warn(
{ error, keyId: key.id, window },
'twap-computation: failed for key/window'
);
}
}
}

logger.info(
{ scannedKeys: activeKeys.length, computedWrites, failedWrites },
'twap-computation: completed'
);

return {
scannedKeys: activeKeys.length,
computedWrites,
failedWrites,
};
}

/** Immediate backfill for a single key across all windows. */
export async function backfillTwapForKey(
keyId: string,
now: Date = new Date()
): Promise<void> {
try {
for (const window of TWAP_WINDOWS) {
await computeAndCacheTwap(keyId, window, now);
}
} catch (error) {
if (error instanceof KeyNotFoundError) {
logger.warn(
{ keyId },
'twap-computation: backfill skipped, key not found'
);
return;
}
throw error;
}
}

let twapTimer: ReturnType<typeof setInterval> | null = null;
let registrationHooked = false;

function onKeyRegistered(payload: { keyAddress?: string }): void {
// RegisteredKey.keyAddress has no direct CreatorProfile mapping yet,
// so best-effort backfill the address as a key id and always sweep
// missing keys so a new CreatorProfile is covered within seconds.
void (async () => {
try {
if (payload?.keyAddress) {
await backfillTwapForKey(payload.keyAddress).catch(() => {});
}
await computeTwapForAllKeys().catch(() => {});
} catch (error) {
logger.error(
{ err: error },
'twap-computation: registration backfill failed'
);
}
})();
}

export function startTwapComputationJob(): void {
if (!envConfig.TWAP_COMPUTATION_ENABLED) {
logger.info('twap-computation job is disabled');
return;
}

const intervalMs = envConfig.TWAP_COMPUTATION_INTERVAL_MINUTES * 60 * 1000;

const run = async () => {
try {
await computeTwapForAllKeys();
} catch (error) {
logger.error(
{ err: error },
'twap-computation failed with an unexpected error'
);
}
};

void run();
twapTimer = setInterval(() => {
void run();
}, intervalMs);

if (
typeof (twapTimer as unknown as { unref?: () => void }).unref ===
'function'
) {
(twapTimer as unknown as { unref: () => void }).unref();
}

if (!registrationHooked) {
keyEventEmitter.on('key_registered', onKeyRegistered);
registrationHooked = true;
}

logger.info(
{ intervalMinutes: envConfig.TWAP_COMPUTATION_INTERVAL_MINUTES },
'twap-computation job started'
);
}

export function stopTwapComputationJob(): void {
if (twapTimer) {
clearInterval(twapTimer);
twapTimer = null;
}
if (registrationHooked) {
keyEventEmitter.removeListener('key_registered', onKeyRegistered);
registrationHooked = false;
}
logger.info('twap-computation job stopped');
}
Loading
Loading