diff --git a/src/modules/indexer/price-snapshot.service.ts b/src/modules/indexer/price-snapshot.service.ts index 95fc6405..455bac61 100644 --- a/src/modules/indexer/price-snapshot.service.ts +++ b/src/modules/indexer/price-snapshot.service.ts @@ -61,6 +61,12 @@ export async function upsertPriceSnapshot( }, 'price-snapshot: written (first trade)' ); + try { + const { invalidateKeyTwapCache } = await import('../keys/key-twap.service'); + await invalidateKeyTwapCache(creatorId); + } catch { + // Non-critical cache invalidation failure + } return; } @@ -110,6 +116,12 @@ export async function upsertPriceSnapshot( }, 'price-snapshot: written' ); + try { + const { invalidateKeyTwapCache } = await import('../keys/key-twap.service'); + await invalidateKeyTwapCache(creatorId); + } catch { + // Non-critical cache invalidation failure + } } catch (err) { logger.error({ err, creatorId }, 'price-snapshot: failed to upsert'); throw err; diff --git a/src/modules/indexer/trade-indexer.service.ts b/src/modules/indexer/trade-indexer.service.ts index 0768c2f4..ed8b6845 100644 --- a/src/modules/indexer/trade-indexer.service.ts +++ b/src/modules/indexer/trade-indexer.service.ts @@ -91,6 +91,8 @@ export async function processTradeEvent( const { invalidateCreatorDashboardCache } = await import('../creator/creator-dashboard.service'); await invalidateCreatorDashboardCache(event.creator_id); + const { invalidateKeyTwapCache } = await import('../keys/key-twap.service'); + await invalidateKeyTwapCache(event.creator_id); } catch { // Non-critical cache invalidation failure } diff --git a/src/modules/keys/__tests__/key-twap.test.ts b/src/modules/keys/__tests__/key-twap.test.ts new file mode 100644 index 00000000..12d3554f --- /dev/null +++ b/src/modules/keys/__tests__/key-twap.test.ts @@ -0,0 +1,180 @@ +const redisStore = new Map(); + +jest.mock('../../../utils/redis.utils', () => ({ + getRedis: () => ({ + get: jest.fn(async (key: string) => redisStore.get(key) ?? null), + set: jest.fn(async (key: string, value: string) => { + redisStore.set(key, value); + return 'OK'; + }), + del: jest.fn(async (key: string) => { + redisStore.delete(key); + return 1; + }), + scan: jest.fn(async (_cursor: string, _match: string, _pattern: string) => { + return ['0', []]; + }), + }), + cacheGetJson: jest.fn(async (key: string): Promise => { + const val = redisStore.get(key); + return val ? JSON.parse(val) : null; + }), + cacheSetJson: jest.fn(async (key: string, value: unknown) => { + redisStore.set(key, JSON.stringify(value)); + }), + cacheInvalidate: jest.fn(async (...keysOrPatterns: string[]) => { + for (const k of keysOrPatterns) { + if (k.includes('*')) { + const prefix = k.replace('*', ''); + for (const key of Array.from(redisStore.keys())) { + if (key.startsWith(prefix)) { + redisStore.delete(key); + } + } + } else { + redisStore.delete(k); + } + } + }), +})); + +jest.mock('../../../utils/prisma.utils', () => ({ + prisma: { + creatorProfile: { + findFirst: jest.fn(), + findUnique: jest.fn(), + }, + creatorPriceHistory: { + findMany: jest.fn(), + }, + creatorPriceSnapshot: { + findUnique: jest.fn(), + }, + }, +})); + +jest.mock('../../../utils/logger.utils', () => ({ + logger: { + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + debug: jest.fn(), + }, +})); + +import request from 'supertest'; +import express from 'express'; +import keysRouter from '../keys.routes'; +import { prisma } from '../../../utils/prisma.utils'; +import { invalidateKeyTwapCache } from '../key-twap.service'; + +const app = express(); +app.use(express.json()); +app.use('/api/v1/keys', keysRouter); + +describe('GET /api/v1/keys/:keyId/twap (#866)', () => { + const now = new Date('2026-09-24T12:00:00.000Z'); + + beforeEach(() => { + redisStore.clear(); + jest.clearAllMocks(); + }); + + it('returns 422 for invalid or missing window param', async () => { + const res1 = await request(app).get('/api/v1/keys/creator-1/twap'); + expect(res1.status).toBe(422); + + const res2 = await request(app).get('/api/v1/keys/creator-1/twap?window=invalid'); + expect(res2.status).toBe(422); + }); + + it('returns 404 when key is not found', async () => { + (prisma.creatorProfile.findFirst as jest.Mock).mockResolvedValue(null); + + const res = await request(app).get('/api/v1/keys/nonexistent/twap?window=1h'); + expect(res.status).toBe(404); + }); + + it('returns null twapPrice when fewer than 2 snapshots exist in the window', async () => { + (prisma.creatorProfile.findFirst as jest.Mock).mockResolvedValue({ + id: 'creator-1', + handle: 'creator-1', + }); + (prisma.creatorPriceHistory.findMany as jest.Mock).mockResolvedValue([ + { + id: 'h1', + creatorId: 'creator-1', + price: 10000000n, + recordedAt: new Date(now.getTime() - 10 * 60 * 1000), + }, + ]); + (prisma.creatorPriceSnapshot.findUnique as jest.Mock).mockResolvedValue({ + currentPrice: 10000000n, + }); + + const res = await request(app).get('/api/v1/keys/creator-1/twap?window=1h'); + expect(res.status).toBe(200); + expect(res.body.success).toBe(true); + expect(res.body.data.twapPrice).toBeNull(); + expect(res.body.data.spotPrice).toBe('10000000'); + expect(res.body.data.snapshotCount).toBe(1); + expect(res.body.data.windowLedgers).toBe(720); + }); + + it('computes twapPrice correctly with spotPrice across requested window (24h)', async () => { + (prisma.creatorProfile.findFirst as jest.Mock).mockResolvedValue({ + id: 'creator-1', + handle: 'creator-1', + }); + // Snapshots 1 hour apart: 10,000,000 and 20,000,000 + const t0 = new Date(now.getTime() - 2 * 60 * 60 * 1000); + const t1 = new Date(now.getTime() - 1 * 60 * 60 * 1000); + (prisma.creatorPriceHistory.findMany as jest.Mock).mockResolvedValue([ + { id: 'h1', creatorId: 'creator-1', price: 10000000n, recordedAt: t0 }, + { id: 'h2', creatorId: 'creator-1', price: 20000000n, recordedAt: t1 }, + ]); + (prisma.creatorPriceSnapshot.findUnique as jest.Mock).mockResolvedValue({ + currentPrice: 20000000n, + }); + + const res = await request(app).get('/api/v1/keys/creator-1/twap?window=24h'); + expect(res.status).toBe(200); + expect(res.body.success).toBe(true); + expect(res.body.data.twapPrice).toBe('15000000'); + expect(res.body.data.spotPrice).toBe('20000000'); + expect(res.body.data.snapshotCount).toBe(2); + expect(res.body.data.windowLedgers).toBe(17280); + }); + + it('serves from Redis cache within 60s and invalidates on invalidateKeyTwapCache', async () => { + (prisma.creatorProfile.findFirst as jest.Mock).mockResolvedValue({ + id: 'creator-1', + handle: 'creator-1', + }); + (prisma.creatorPriceHistory.findMany as jest.Mock).mockResolvedValue([ + { id: 'h1', creatorId: 'creator-1', price: 10000000n, recordedAt: new Date(now.getTime() - 1000) }, + { id: 'h2', creatorId: 'creator-1', price: 20000000n, recordedAt: now }, + ]); + (prisma.creatorPriceSnapshot.findUnique as jest.Mock).mockResolvedValue({ + currentPrice: 20000000n, + }); + + // First call populates cache + const res1 = await request(app).get('/api/v1/keys/creator-1/twap?window=7d'); + expect(res1.status).toBe(200); + expect(prisma.creatorPriceHistory.findMany).toHaveBeenCalledTimes(1); + + // Second call hits cache + const res2 = await request(app).get('/api/v1/keys/creator-1/twap?window=7d'); + expect(res2.status).toBe(200); + expect(prisma.creatorPriceHistory.findMany).toHaveBeenCalledTimes(1); + + // Invalidate cache + await invalidateKeyTwapCache('creator-1'); + + // Third call fetches from DB + const res3 = await request(app).get('/api/v1/keys/creator-1/twap?window=7d'); + expect(res3.status).toBe(200); + expect(prisma.creatorPriceHistory.findMany).toHaveBeenCalledTimes(2); + }); +}); diff --git a/src/modules/keys/key-twap.service.ts b/src/modules/keys/key-twap.service.ts new file mode 100644 index 00000000..0f6aceef --- /dev/null +++ b/src/modules/keys/key-twap.service.ts @@ -0,0 +1,185 @@ +import { prisma } from '../../utils/prisma.utils'; +import { logger } from '../../utils/logger.utils'; +import { cacheGetJson, cacheSetJson, cacheInvalidate } from '../../utils/redis.utils'; +import { envConfig } from '../../config'; + +export const TWAP_WINDOWS = ['1h', '24h', '7d'] as const; +export type TwapWindow = (typeof TWAP_WINDOWS)[number]; + +export const TWAP_WINDOW_MS: Record = { + '1h': 60 * 60 * 1000, + '24h': 24 * 60 * 60 * 1000, + '7d': 7 * 24 * 60 * 60 * 1000, +}; + +// Ledgers close every ~5 seconds on Stellar/Soroban +export const TWAP_WINDOW_LEDGERS: Record = { + '1h': 720, + '24h': 17280, + '7d': 120960, +}; + +export const TWAP_CACHE_TTL_SECONDS = 60; + +export interface TwapResult { + keyId: string; + window: TwapWindow; + windowLedgers: number; + twapPrice: string | null; + spotPrice: string; + snapshotCount: number; +} + +export function buildTwapCacheKey(keyId: string, window: TwapWindow): string { + return `key:twap:${keyId}:${window}`; +} + +export async function invalidateKeyTwapCache(keyId: string): Promise { + await cacheInvalidate(`key:twap:${keyId}:*`); +} + +/** + * Attempts to query on-chain Soroban contract get_twap view via RPC. + */ +async function fetchOnChainTwap( + contractId: string, + windowLedgers: number +): Promise { + if (!envConfig.STELLAR_SOROBAN_RPC_URL) { + return null; + } + + try { + const response = await fetch(envConfig.STELLAR_SOROBAN_RPC_URL, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ + jsonrpc: '2.0', + id: 'get_twap_view', + method: 'simulateTransaction', + params: { + contractId, + functionName: 'get_twap', + args: [windowLedgers], + }, + }), + }); + + if (!response.ok) return null; + + const json = await response.json(); + if (json.result?.results?.[0]?.xdr) { + return json.result.results[0].xdr.toString(); + } + } catch (err) { + logger.debug({ err, contractId }, 'get_twap on-chain view query skipped'); + } + + return null; +} + +export async function getKeyTwap( + creatorId: string, + window: TwapWindow, + now: Date = new Date() +): Promise { + const cacheKey = buildTwapCacheKey(creatorId, window); + const cached = await cacheGetJson(cacheKey); + if (cached !== null) { + return cached; + } + + const windowMs = TWAP_WINDOW_MS[window]; + const windowLedgers = TWAP_WINDOW_LEDGERS[window]; + const windowStart = new Date(now.getTime() - windowMs); + + // Log get_twap contract view call with window in ledger units + logger.info( + { + operation: 'get_twap_contract_view', + keyId: creatorId, + window, + windowLedgers, + }, + 'Calling get_twap contract view' + ); + + // 1. Try on-chain Soroban view + const onChainPrice = await fetchOnChainTwap(creatorId, windowLedgers); + if (onChainPrice !== null) { + const priceSnapshot = await prisma.creatorPriceSnapshot.findUnique({ + where: { creatorId }, + select: { currentPrice: true }, + }); + const result: TwapResult = { + keyId: creatorId, + window, + windowLedgers, + twapPrice: onChainPrice, + spotPrice: priceSnapshot?.currentPrice.toString() ?? onChainPrice, + snapshotCount: 0, + }; + await cacheSetJson(cacheKey, result, TWAP_CACHE_TTL_SECONDS); + return result; + } + + // 2. Derive contract-backed calculation using stored price snapshots + const [snapshots, priceSnapshot] = await Promise.all([ + prisma.creatorPriceHistory.findMany({ + where: { + creatorId, + recordedAt: { gte: windowStart, lte: now }, + }, + orderBy: { recordedAt: 'asc' }, + }), + prisma.creatorPriceSnapshot.findUnique({ + where: { creatorId }, + select: { currentPrice: true }, + }), + ]); + + const spotPrice = priceSnapshot + ? priceSnapshot.currentPrice.toString() + : snapshots.length > 0 + ? snapshots[snapshots.length - 1].price.toString() + : '0'; + + let twapPrice: string | null = null; + + if (snapshots.length >= 2) { + let totalTimeWeight = 0; + let weightedPriceSum = 0n; + + for (let i = 0; i < snapshots.length - 1; i++) { + const tCurrent = snapshots[i].recordedAt.getTime(); + const tNext = snapshots[i + 1].recordedAt.getTime(); + const dt = Math.max(0, tNext - tCurrent); + + if (dt > 0) { + // Accumulate doubled sum: (P_i + P_{i+1}) * dt to prevent premature integer truncation + weightedPriceSum += (snapshots[i].price + snapshots[i + 1].price) * BigInt(dt); + totalTimeWeight += dt; + } + } + + if (totalTimeWeight > 0) { + // Divide once at the end by 2 * totalTimeWeight + twapPrice = (weightedPriceSum / (2n * BigInt(totalTimeWeight))).toString(); + } else { + const sum = snapshots.reduce((acc, s) => acc + s.price, 0n); + twapPrice = (sum / BigInt(snapshots.length)).toString(); + } + } + + const result: TwapResult = { + keyId: creatorId, + window, + windowLedgers, + twapPrice, + spotPrice, + snapshotCount: snapshots.length, + }; + + await cacheSetJson(cacheKey, result, TWAP_CACHE_TTL_SECONDS); + return result; +} diff --git a/src/modules/keys/keys.routes.ts b/src/modules/keys/keys.routes.ts index fbab0a60..782b6061 100644 --- a/src/modules/keys/keys.routes.ts +++ b/src/modules/keys/keys.routes.ts @@ -60,6 +60,7 @@ import { PositionNotFoundError, unfreezePosition, } from './key-freeze.service'; +import { getKeyTwap } from './key-twap.service'; const priceHistoryQuerySchema = z.object({ from: z.string().datetime(), @@ -366,6 +367,45 @@ router.get('/:keyId/price-history', async (req, res, next) => { } }); +// ── GET /:keyId/twap ────────────────────────────────────────── + +const twapQuerySchema = z.object({ + window: z.enum(['1h', '24h', '7d'], { + errorMap: () => ({ message: 'Invalid window param. Must be 1h, 24h, or 7d' }), + }), +}); + +router.get('/:keyId/twap', async (req, res, next) => { + const keyId = String(req.params.keyId); + const parsed = twapQuerySchema.safeParse(req.query); + if (!parsed.success) { + sendError( + res, + 422, + ErrorCode.UNPROCESSABLE_ENTITY, + 'Invalid window param. Must be 1h, 24h, or 7d', + zodIssuesToDetails(parsed.error.issues) + ); + return; + } + + try { + const creator = await prisma.creatorProfile.findFirst({ + where: { OR: [{ id: keyId }, { handle: keyId }] }, + select: { id: true }, + }); + if (!creator) { + sendNotFound(res, 'Key'); + return; + } + + const result = await getKeyTwap(creator.id, parsed.data.window); + sendSuccess(res, result); + } catch (error) { + next(error); + } +}); + // Mount dividend routes router.use('/', dividendRouter);