From e0ae707e5a95c80ec86fc57e91f85f3e94d18606 Mon Sep 17 00:00:00 2001 From: shogun444 Date: Sun, 27 Sep 2026 18:03:56 +0530 Subject: [PATCH] feat: add TWAP price calculation and caching service --- .env.example | 4 + src/config.schema.ts | 8 + src/constants/redis.constants.ts | 29 +++ src/jobs/twap-computation.job.test.ts | 98 +++++++++ src/jobs/twap-computation.job.ts | 160 ++++++++++++++ src/modules/keys/key-twap.service.test.ts | 210 +++++++++++++++++++ src/modules/keys/key-twap.service.ts | 242 ++++++++++++++++++++++ src/modules/keys/keys.routes.ts | 41 ++++ src/server.ts | 6 + 9 files changed, 798 insertions(+) create mode 100644 src/constants/redis.constants.ts create mode 100644 src/jobs/twap-computation.job.test.ts create mode 100644 src/jobs/twap-computation.job.ts create mode 100644 src/modules/keys/key-twap.service.test.ts create mode 100644 src/modules/keys/key-twap.service.ts diff --git a/.env.example b/.env.example index d2e60c89..fd8463b7 100644 --- a/.env.example +++ b/.env.example @@ -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 diff --git a/src/config.schema.ts b/src/config.schema.ts index 53aa50fd..aff00a88 100644 --- a/src/config.schema.ts +++ b/src/config.schema.ts @@ -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 diff --git a/src/constants/redis.constants.ts b/src/constants/redis.constants.ts new file mode 100644 index 00000000..69dbc62a --- /dev/null +++ b/src/constants/redis.constants.ts @@ -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 = { + '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 = { + '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; diff --git a/src/jobs/twap-computation.job.test.ts b/src/jobs/twap-computation.job.test.ts new file mode 100644 index 00000000..e1b59294 --- /dev/null +++ b/src/jobs/twap-computation.job.test.ts @@ -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) + ); + }); +}); diff --git a/src/jobs/twap-computation.job.ts b/src/jobs/twap-computation.job.ts new file mode 100644 index 00000000..d7419d87 --- /dev/null +++ b/src/jobs/twap-computation.job.ts @@ -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 { + 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 { + 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 | 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'); +} diff --git a/src/modules/keys/key-twap.service.test.ts b/src/modules/keys/key-twap.service.test.ts new file mode 100644 index 00000000..afc275b9 --- /dev/null +++ b/src/modules/keys/key-twap.service.test.ts @@ -0,0 +1,210 @@ +// src/modules/keys/key-twap.service.test.ts +const redisStore = new Map(); + +jest.mock('../../utils/redis.utils', () => ({ + cacheGetJson: jest.fn(async (key: string) => + redisStore.has(key) ? (redisStore.get(key) as unknown) : null + ), + cacheSetJson: jest.fn(async (key: string, value: unknown) => { + redisStore.set(key, value); + }), +})); + +jest.mock('../../utils/prisma.utils', () => ({ + prisma: { + creatorProfile: { findFirst: jest.fn() }, + creatorPriceHistory: { findFirst: jest.fn(), findMany: jest.fn() }, + }, +})); + +jest.mock('../../utils/logger.utils', () => ({ + logger: { + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + debug: jest.fn(), + }, +})); + +import { prisma } from '../../utils/prisma.utils'; +import { cacheSetJson } from '../../utils/redis.utils'; +import { + computeAndCacheTwap, + computeDeltaPct, + computeTwapFromSnapshots, + getTwapPrice, + KeyNotFoundError, +} from './key-twap.service'; +import { + TWAP_CACHE_TTL_SECONDS, + TWAP_STALE_THRESHOLD_MS, + twapRedisKey, +} from '../../constants/redis.constants'; + +const mockPrisma = prisma as unknown as { + creatorProfile: { findFirst: jest.Mock }; + creatorPriceHistory: { findFirst: jest.Mock; findMany: jest.Mock }; +}; + +function mockCreator(overrides: Record = {}) { + mockPrisma.creatorProfile.findFirst.mockResolvedValue({ + id: 'key-1', + circulatingSupply: '10', + creatorRoyaltyBuyBps: 0, + ...overrides, + }); +} + +describe('key-twap.service', () => { + beforeEach(() => { + redisStore.clear(); + jest.clearAllMocks(); + }); + + describe('computeTwapFromSnapshots', () => { + it('weights prices by time with a prior seed', () => { + const now = new Date('2026-09-27T12:00:00.000Z').getTime(); + const start = now - 60 * 60 * 1000; + // Seed 100 for first 30m, then 200 for last 30m => TWAP 150. + const snapshots = [ + { timestamp: new Date(start + 30 * 60 * 1000), price: 200n }, + ]; + expect(computeTwapFromSnapshots(100n, snapshots, start, now)).toBe( + 150n + ); + }); + + it('fills the leading gap with the first price when no prior exists', () => { + const now = new Date('2026-09-27T12:00:00.000Z').getTime(); + const start = now - 60 * 60 * 1000; + const snapshots = [ + { timestamp: new Date(start + 30 * 60 * 1000), price: 200n }, + ]; + // First price fills [start, t1), then 200 to now => 200. + expect(computeTwapFromSnapshots(null, snapshots, start, now)).toBe( + 200n + ); + }); + + it('returns the seed price when no in-window snapshots exist', () => { + const now = Date.now(); + const start = now - 60 * 60 * 1000; + expect(computeTwapFromSnapshots(123n, [], start, now)).toBe(123n); + }); + + it('returns null when there is no price information', () => { + const now = Date.now(); + const start = now - 60 * 60 * 1000; + expect(computeTwapFromSnapshots(null, [], start, now)).toBeNull(); + }); + }); + + describe('computeDeltaPct', () => { + it('computes ((spot - twap) / twap) * 100', () => { + expect(computeDeltaPct(110n, 100n)).toBe(10); + expect(computeDeltaPct(90n, 100n)).toBe(-10); + }); + + it('returns null when TWAP is zero (division-by-zero guard)', () => { + expect(computeDeltaPct(100n, 0n)).toBeNull(); + }); + }); + + describe('computeAndCacheTwap', () => { + it('falls back to spot with delta 0 when the key has 0 snapshots', async () => { + mockCreator(); + mockPrisma.creatorPriceHistory.findFirst.mockResolvedValue(null); + mockPrisma.creatorPriceHistory.findMany.mockResolvedValue([]); + + const result = await computeAndCacheTwap('key-1', '1h'); + + expect(result.twap).toBe(result.spotPrice); + expect(result.deltaPct).toBe(0); + expect(result.stale).toBe(false); + expect(cacheSetJson).toHaveBeenCalledWith( + twapRedisKey('key-1', '1h'), + expect.objectContaining({ keyId: 'key-1', window: '1h' }), + TWAP_CACHE_TTL_SECONDS['1h'] + ); + }); + + it('seeds carry-forward from the snapshot prior to the window start', async () => { + mockCreator(); + const now = new Date('2026-09-27T12:00:00.000Z'); + const windowStart = new Date(now.getTime() - 60 * 60 * 1000); + mockPrisma.creatorPriceHistory.findFirst.mockResolvedValue({ + price: 100n, + recordedAt: new Date(windowStart.getTime() - 1000), + }); + mockPrisma.creatorPriceHistory.findMany.mockResolvedValue([ + { + price: 200n, + recordedAt: new Date(windowStart.getTime() + 30 * 60 * 1000), + }, + ]); + + const result = await computeAndCacheTwap('key-1', '1h', now); + + expect(result.twap).toBe('150'); + }); + + it('throws KeyNotFoundError for an unknown key', async () => { + mockPrisma.creatorProfile.findFirst.mockResolvedValue(null); + await expect(computeAndCacheTwap('missing', '1h')).rejects.toThrow( + KeyNotFoundError + ); + }); + }); + + describe('getTwapPrice read-through', () => { + it('computes on a cold cache and caches with window TTL', async () => { + mockCreator(); + mockPrisma.creatorPriceHistory.findFirst.mockResolvedValue(null); + mockPrisma.creatorPriceHistory.findMany.mockResolvedValue([]); + + const result = await getTwapPrice('key-1', '4h'); + + expect(result.window).toBe('4h'); + expect(redisStore.has(twapRedisKey('key-1', '4h'))).toBe(true); + }); + + it('serves cache hits without hitting the history table', async () => { + mockCreator(); + const cached = { + keyId: 'key-1', + window: '1h' as const, + twap: '100', + spotPrice: '110', + deltaPct: 10, + computedAt: new Date().toISOString(), + stale: false, + }; + redisStore.set(twapRedisKey('key-1', '1h'), cached); + + const result = await getTwapPrice('key-1', '1h'); + + expect(result.twap).toBe('100'); + expect(result.stale).toBe(false); + expect(mockPrisma.creatorPriceHistory.findMany).not.toHaveBeenCalled(); + }); + + it('marks stale true when computedAt is older than 10 minutes', async () => { + mockCreator(); + expect(TWAP_STALE_THRESHOLD_MS).toBe(10 * 60 * 1000); + const cached = { + keyId: 'key-1', + window: '24h' as const, + twap: '100', + spotPrice: '100', + deltaPct: 0, + computedAt: new Date(Date.now() - 11 * 60 * 1000).toISOString(), + stale: false, + }; + redisStore.set(twapRedisKey('key-1', '24h'), cached); + + const result = await getTwapPrice('key-1', '24h'); + + expect(result.stale).toBe(true); + }); + }); +}); diff --git a/src/modules/keys/key-twap.service.ts b/src/modules/keys/key-twap.service.ts new file mode 100644 index 00000000..c29467b0 --- /dev/null +++ b/src/modules/keys/key-twap.service.ts @@ -0,0 +1,242 @@ +// src/modules/keys/key-twap.service.ts +// TWAP price calculation and caching service for creator keys (#963). +// +// - Time-weighted average over CreatorPriceHistory snapshots per window. +// - Redis cache with TTL matching the window size. +// - Read-through on cache miss; stale flag when the job is behind. + +import { prisma } from '../../utils/prisma.utils'; +import { cacheGetJson, cacheSetJson } from '../../utils/redis.utils'; +import { getBuyUnitPrice } from '../../utils/pricing.utils'; +import { + TWAP_CACHE_TTL_SECONDS, + TWAP_MAX_SNAPSHOTS, + TWAP_STALE_THRESHOLD_MS, + TWAP_WINDOW_MS, + twapRedisKey, + type TwapWindow, +} from '../../constants/redis.constants'; + +export class KeyNotFoundError extends Error { + constructor(keyId: string) { + super(`Key not found: ${keyId}`); + this.name = 'KeyNotFoundError'; + } +} + +export interface TwapPriceResult { + keyId: string; + window: TwapWindow; + /** TWAP in stroops (string to avoid JS precision loss). */ + twap: string; + /** Bonding-curve spot price for the next buy unit, in stroops. */ + spotPrice: string; + /** + * Signed percentage difference from TWAP to spot: + * ((spotPrice - twap) / twap) * 100 + * Positive = spot above TWAP; negative = spot below TWAP. + * 0 when falling back to spot with no history; null when TWAP is 0. + */ + deltaPct: number | null; + /** ISO-8601 timestamp of this computation. */ + computedAt: string; + /** True when computedAt is older than TWAP_STALE_THRESHOLD_MS. */ + stale: boolean; +} + +export interface TwapSnapshotPoint { + timestamp: Date; + price: bigint; +} + +/** + * Pure time-weighted average over carry-forward prices. + * + * priorPrice seeds [windowStart, firstSnapshot); when null, the first + * in-window price fills the leading gap. Returns null when there is + * no price information at all. + */ +export function computeTwapFromSnapshots( + priorPrice: bigint | null, + snapshots: TwapSnapshotPoint[], + windowStartMs: number, + nowMs: number +): bigint | null { + const totalMs = nowMs - windowStartMs; + if (totalMs <= 0) { + return null; + } + if (snapshots.length === 0) { + return priorPrice; + } + + const ordered = [...snapshots].sort( + (a, b) => a.timestamp.getTime() - b.timestamp.getTime() + ); + + let currentPrice: bigint | null = priorPrice ?? ordered[0].price; + let prevTimeMs = windowStartMs; + let weightedSum = 0n; + + for (const snapshot of ordered) { + const snapshotMs = snapshot.timestamp.getTime(); + const durationMs = Math.max(0, snapshotMs - prevTimeMs); + if (durationMs > 0 && currentPrice !== null) { + weightedSum += currentPrice * BigInt(durationMs); + } + currentPrice = snapshot.price; + prevTimeMs = Math.max(prevTimeMs, snapshotMs); + } + + const tailMs = Math.max(0, nowMs - prevTimeMs); + if (tailMs > 0 && currentPrice !== null) { + weightedSum += currentPrice * BigInt(tailMs); + } + + if (currentPrice === null) { + return null; + } + return weightedSum / BigInt(totalMs); +} + +/** Guarded delta: ((spot - twap) / twap) * 100, null when TWAP is 0. */ +export function computeDeltaPct( + spotPrice: bigint, + twap: bigint +): number | null { + if (twap === 0n) { + return null; + } + const spot = Number(spotPrice); + const base = Number(twap); + if (!Number.isFinite(spot) || !Number.isFinite(base)) { + return null; + } + return parseFloat((((spot - base) / base) * 100).toFixed(4)); +} + +type CreatorForTwap = { + id: string; + circulatingSupply: unknown; + creatorRoyaltyBuyBps: number; +}; + +async function fetchCreatorForTwap(keyId: string): Promise { + const creator = await prisma.creatorProfile.findFirst({ + where: { OR: [{ id: keyId }, { handle: keyId }] }, + select: { + id: true, + circulatingSupply: true, + creatorRoyaltyBuyBps: true, + }, + }); + if (!creator) { + throw new KeyNotFoundError(keyId); + } + return creator as CreatorForTwap; +} + +function getSpotForCreator(creator: CreatorForTwap): bigint { + const supply = Number(creator.circulatingSupply?.toString() ?? '0'); + const safeSupply = Number.isFinite(supply) + ? Math.max(0, Math.floor(supply)) + : 0; + return getBuyUnitPrice(safeSupply, creator.creatorRoyaltyBuyBps); +} + +/** + * Compute TWAP for a key/window and refresh the Redis cache. + * Falls back to spot (deltaPct = 0) when the key has 0 snapshots. + */ +export async function computeAndCacheTwap( + keyId: string, + window: TwapWindow, + now: Date = new Date() +): Promise { + const creator = await fetchCreatorForTwap(keyId); + const nowMs = now.getTime(); + const windowStartMs = nowMs - TWAP_WINDOW_MS[window]; + const windowStart = new Date(windowStartMs); + + const prior = await prisma.creatorPriceHistory.findFirst({ + where: { creatorId: creator.id, recordedAt: { lt: windowStart } }, + orderBy: { recordedAt: 'desc' }, + select: { price: true, recordedAt: true }, + }); + + const rows = await prisma.creatorPriceHistory.findMany({ + where: { + creatorId: creator.id, + recordedAt: { gte: windowStart, lte: now }, + }, + orderBy: { recordedAt: 'asc' }, + take: TWAP_MAX_SNAPSHOTS, + select: { price: true, recordedAt: true }, + }); + + const spotPrice = getSpotForCreator(creator); + const hasHistory = + (prior !== null && prior !== undefined) || rows.length > 0; + + let twap: bigint; + let deltaPct: number | null; + if (!hasHistory) { + twap = spotPrice; + deltaPct = 0; + } else { + const computed = computeTwapFromSnapshots( + prior ? (prior.price as bigint) : null, + rows.map((row: { price: bigint; recordedAt: Date }) => ({ + timestamp: row.recordedAt as Date, + price: row.price as bigint, + })), + windowStartMs, + nowMs + ); + if (computed === null) { + twap = spotPrice; + deltaPct = 0; + } else { + twap = computed; + deltaPct = computeDeltaPct(spotPrice, twap); + } + } + + const result: TwapPriceResult = { + keyId: creator.id, + window, + twap: twap.toString(), + spotPrice: spotPrice.toString(), + deltaPct, + computedAt: now.toISOString(), + stale: false, + }; + + await cacheSetJson( + twapRedisKey(creator.id, window), + result, + TWAP_CACHE_TTL_SECONDS[window] + ); + return result; +} + +/** + * Read-through TWAP fetch: serve the cached value when present + * (flagging stale when the job is behind), otherwise compute on demand. + */ +export async function getTwapPrice( + keyId: string, + window: TwapWindow, + now: Date = new Date() +): Promise { + const creator = await fetchCreatorForTwap(keyId); + const cached = await cacheGetJson( + twapRedisKey(creator.id, window) + ); + if (cached !== null) { + const ageMs = Date.now() - new Date(cached.computedAt).getTime(); + const stale = !Number.isFinite(ageMs) || ageMs > TWAP_STALE_THRESHOLD_MS; + return { ...cached, keyId: creator.id, window, stale }; + } + return computeAndCacheTwap(creator.id, window, now); +} diff --git a/src/modules/keys/keys.routes.ts b/src/modules/keys/keys.routes.ts index 69b51ae7..f2cb3b0b 100644 --- a/src/modules/keys/keys.routes.ts +++ b/src/modules/keys/keys.routes.ts @@ -23,6 +23,11 @@ import { KeyNotFoundError as OracleKeyNotFoundError, OraclePriceNotFoundError, } from './oracle-price.service'; +import { + getTwapPrice, + KeyNotFoundError as TwapKeyNotFoundError, +} from './key-twap.service'; +import { TWAP_WINDOWS } from '../../constants/redis.constants'; import { cacheControl } from '../../middlewares/cache-control.middleware'; import { envConfig } from '../../config'; import { getKeyProposals, getProposalForVoting } from './key-proposals.service'; @@ -132,6 +137,10 @@ const priceImpactQuerySchema = z.object({ direction: z.enum(['buy', 'sell']), }); +const twapQuerySchema = z.object({ + window: z.enum(TWAP_WINDOWS).optional(), +}); + const buybackPoolHistoryQuerySchema = z.object({ limit: z .string() @@ -373,6 +382,38 @@ router.get( } ); +/** + * GET /api/v1/keys/:keyId/price/twap?window=1h|4h|24h + * + * Returns the cached TWAP for the requested window, the bonding-curve + * spot price, the spot-vs-TWAP delta percentage, and a stale flag when + * the computation job is behind (>10 minutes since computedAt). + * Read-through: a cold cache is computed on demand and cached with + * a TTL matching the window size. + */ +router.get('/:keyId/price/twap', async (req, res, next) => { + const parsed = twapQuerySchema.safeParse(req.query); + if (!parsed.success) { + sendValidationError( + res, + 'Invalid twap query', + zodIssuesToDetails(parsed.error.issues) + ); + return; + } + try { + const keyId = String(req.params.keyId); + const window = parsed.data.window ?? '1h'; + sendSuccess(res, await getTwapPrice(keyId, window)); + } catch (error) { + if (error instanceof TwapKeyNotFoundError) { + sendNotFound(res, 'Key'); + return; + } + next(error); + } +}); + /** * GET /api/v1/keys/:keyId * Public key detail response includes supply milestone metadata. diff --git a/src/server.ts b/src/server.ts index 60a1d039..63deae22 100644 --- a/src/server.ts +++ b/src/server.ts @@ -24,6 +24,10 @@ import { startPriceHistoryCleanupJob, stopPriceHistoryCleanupJob, } from './jobs/price-history-cleanup.job'; +import { + startTwapComputationJob, + stopTwapComputationJob, +} from './jobs/twap-computation.job'; import { startFlashLoanViolationCleanupJob, stopFlashLoanViolationCleanupJob, @@ -81,6 +85,7 @@ async function startServer() { startDetectPriceMovementsJob(); startGovernanceSyncJob(); startPriceHistoryCleanupJob(); + startTwapComputationJob(); startFlashLoanViolationCleanupJob(); const server = app.listen(envConfig.PORT, () => { @@ -118,6 +123,7 @@ function createGracefulShutdownHandler(server: ReturnType) { stopDetectPriceMovementsJob(); stopGovernanceSyncJob(); stopPriceHistoryCleanupJob(); + stopTwapComputationJob(); stopFlashLoanViolationCleanupJob(); await prisma.$disconnect(); logger.info('Database connection closed');