From 4556b5848e0d81e03609d18e4ed05a86a8fa6dfa Mon Sep 17 00:00:00 2001 From: She-ge Date: Wed, 30 Sep 2026 12:26:05 +0100 Subject: [PATCH] feat(tokens): add price feed worker and circuit breaker --- .env.example | 7 + .env.mainnet.example | 5 + .env.staging.example | 5 + .env.testnet.example | 5 + .../20260930000001_price_history/down.sql | 1 + .../migration.sql | 13 ++ prisma/schema.prisma | 14 +- src/config/env.validation.ts | 7 + src/tokens/price-feed.provider.spec.ts | 46 ++++++ src/tokens/price-feed.provider.ts | 75 +++++++++ src/tokens/price-feed.worker.spec.ts | 91 ++++++++++ src/tokens/price-feed.worker.ts | 155 ++++++++++++++++++ src/tokens/tokens.module.ts | 5 + 13 files changed, 428 insertions(+), 1 deletion(-) create mode 100644 prisma/migrations/20260930000001_price_history/down.sql create mode 100644 prisma/migrations/20260930000001_price_history/migration.sql create mode 100644 src/tokens/price-feed.provider.spec.ts create mode 100644 src/tokens/price-feed.worker.spec.ts create mode 100644 src/tokens/price-feed.worker.ts diff --git a/.env.example b/.env.example index 00b5cbe..29c851f 100644 --- a/.env.example +++ b/.env.example @@ -73,6 +73,13 @@ SOLVER_CHAINS=stellar,ethereum,base,polygon,arbitrum,optimism,avalanche INTENTS_PERSISTENCE=memory ALLOW_LEGACY_STELLAR_SIGNATURES=false SOLVERS_PERSISTENCE=memory +TOKENS_PERSISTENCE=memory + +# CoinGecko-backed token USD pricing. Add aliases as {"SYMBOL":"coin-id"}. +PRICE_FEED_COIN_IDS={} +PRICE_FEED_API_KEY= +PRICE_FEED_REFRESH_INTERVAL_MS=60000 +PRICE_FEED_CIRCUIT_BREAKER_THRESHOLD_PERCENT=50 # ─── On-chain writes ──────────────────────────────────────────────────────── # Register intents with the settlement contract on create(). diff --git a/.env.mainnet.example b/.env.mainnet.example index 1564f9f..752d94e 100644 --- a/.env.mainnet.example +++ b/.env.mainnet.example @@ -35,6 +35,11 @@ NODE_ENV=production # REQUIRED: mainnet only. STELLAR_NETWORK=mainnet INTENTS_PERSISTENCE=prisma +TOKENS_PERSISTENCE=prisma +PRICE_FEED_COIN_IDS={} +PRICE_FEED_API_KEY= +PRICE_FEED_REFRESH_INTERVAL_MS=60000 +PRICE_FEED_CIRCUIT_BREAKER_THRESHOLD_PERCENT=50 ETHEREUM_RPC_URL= ETHEREUM_ESCROW_ADDRESS= BASE_RPC_URL= diff --git a/.env.staging.example b/.env.staging.example index 293a442..b9a2262 100644 --- a/.env.staging.example +++ b/.env.staging.example @@ -12,6 +12,11 @@ NODE_ENV=staging STELLAR_NETWORK=testnet INTENTS_PERSISTENCE=prisma +TOKENS_PERSISTENCE=prisma +PRICE_FEED_COIN_IDS={} +PRICE_FEED_API_KEY= +PRICE_FEED_REFRESH_INTERVAL_MS=60000 +PRICE_FEED_CIRCUIT_BREAKER_THRESHOLD_PERCENT=50 ETHEREUM_RPC_URL= ETHEREUM_ESCROW_ADDRESS= BASE_RPC_URL= diff --git a/.env.testnet.example b/.env.testnet.example index 38d4107..93b1cb8 100644 --- a/.env.testnet.example +++ b/.env.testnet.example @@ -21,6 +21,11 @@ NODE_ENV=development # ─── Stellar / Soroban ─────────────────────────────────────────────────────── STELLAR_NETWORK=testnet INTENTS_PERSISTENCE=prisma +TOKENS_PERSISTENCE=prisma +PRICE_FEED_COIN_IDS={} +PRICE_FEED_API_KEY= +PRICE_FEED_REFRESH_INTERVAL_MS=60000 +PRICE_FEED_CIRCUIT_BREAKER_THRESHOLD_PERCENT=50 ETHEREUM_RPC_URL= ETHEREUM_ESCROW_ADDRESS= BASE_RPC_URL= diff --git a/prisma/migrations/20260930000001_price_history/down.sql b/prisma/migrations/20260930000001_price_history/down.sql new file mode 100644 index 0000000..d55cc0c --- /dev/null +++ b/prisma/migrations/20260930000001_price_history/down.sql @@ -0,0 +1 @@ +DROP TABLE "price_history"; \ No newline at end of file diff --git a/prisma/migrations/20260930000001_price_history/migration.sql b/prisma/migrations/20260930000001_price_history/migration.sql new file mode 100644 index 0000000..f95bee1 --- /dev/null +++ b/prisma/migrations/20260930000001_price_history/migration.sql @@ -0,0 +1,13 @@ +CREATE TABLE "price_history" ( + "id" BIGSERIAL NOT NULL, + "token_id" TEXT NOT NULL, + "price_usd" DOUBLE PRECISION NOT NULL, + "recorded_at" TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "price_history_pkey" PRIMARY KEY ("id"), + CONSTRAINT "price_history_token_id_fkey" FOREIGN KEY ("token_id") + REFERENCES "tokens"("id") ON DELETE CASCADE ON UPDATE CASCADE +); + +CREATE INDEX CONCURRENTLY "price_history_token_id_recorded_at_idx" + ON "price_history"("token_id", "recorded_at" DESC); \ No newline at end of file diff --git a/prisma/schema.prisma b/prisma/schema.prisma index 0607f08..0cd0e4e 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -119,7 +119,6 @@ model Intent { /// Realized protocol fee charged on this fill (destination-token base units). feeAmount String? @map("fee_amount") /// On-chain transaction hash of the Stellar fill transaction. - txHash String? @map("tx_hash") /// USD value of `srcAmount` at creation time, computed from the resolved /// source-token price. Stored as a numeric so range filters /// (`minAmountUsd` / `maxAmountUsd`) and USD sorting are indexable. @@ -233,6 +232,7 @@ model Token { priceUsd Float? @map("price_usd") /// Whether this is a destination-side Stellar token. isStellar Boolean @default(false) @map("is_stellar") + priceHistory PriceHistory[] @@unique([address, chain]) @@index([chain]) @@ -240,6 +240,18 @@ model Token { @@map("tokens") } +/// Append-only USD price samples written by the price-feed worker. +model PriceHistory { + id BigInt @id @default(autoincrement()) @map("id") + tokenId String @map("token_id") + priceUsd Float @map("price_usd") + recordedAt DateTime @default(now()) @map("recorded_at") @db.Timestamptz + token Token @relation(fields: [tokenId], references: [id], onDelete: Cascade) + + @@index([tokenId, recordedAt(sort: Desc)]) + @@map("price_history") +} + // ─── KillSwitch (issue #477) ───────────────────────────────────────────────── // Hierarchical emergency pause. One row per (scope, chain, token, operation) // tuple; activating the row closes that scope for the gated operation. diff --git a/src/config/env.validation.ts b/src/config/env.validation.ts index 300a941..0439f8b 100644 --- a/src/config/env.validation.ts +++ b/src/config/env.validation.ts @@ -101,6 +101,13 @@ export const envValidationSchema = Joi.object({ // to a live database. Intended for production / staging. INTENTS_PERSISTENCE: Joi.string().valid("memory", "prisma").default("memory"), SOLVERS_PERSISTENCE: Joi.string().valid("memory", "prisma").default("memory"), + TOKENS_PERSISTENCE: Joi.string().valid("memory", "prisma").default("memory"), + + // Map token symbols to CoinGecko IDs for assets not covered by defaults. + PRICE_FEED_COIN_IDS: Joi.string().default("{}"), + PRICE_FEED_API_KEY: Joi.string().allow("").default(""), + PRICE_FEED_REFRESH_INTERVAL_MS: Joi.number().integer().min(1000).default(60000), + PRICE_FEED_CIRCUIT_BREAKER_THRESHOLD_PERCENT: Joi.number().positive().max(10000).default(50), // ── Intent retention (in-memory store hygiene) ───────────────────────────── // How long terminal intents are kept in the in-memory adapter, and how often diff --git a/src/tokens/price-feed.provider.spec.ts b/src/tokens/price-feed.provider.spec.ts new file mode 100644 index 0000000..eadb027 --- /dev/null +++ b/src/tokens/price-feed.provider.spec.ts @@ -0,0 +1,46 @@ +import { CoinGeckoPriceFeedProvider } from "./price-feed.provider"; +import { EgressPurpose, EgressResponse, HttpEgressService } from "../common/http-egress"; + +describe("CoinGeckoPriceFeedProvider", () => { + const originalCoinIds = process.env.PRICE_FEED_COIN_IDS; + const originalApiKey = process.env.PRICE_FEED_API_KEY; + let egressFetch: jest.SpyInstance; + + afterEach(() => { + if (originalCoinIds === undefined) delete process.env.PRICE_FEED_COIN_IDS; + else process.env.PRICE_FEED_COIN_IDS = originalCoinIds; + if (originalApiKey === undefined) delete process.env.PRICE_FEED_API_KEY; + else process.env.PRICE_FEED_API_KEY = originalApiKey; + jest.restoreAllMocks(); + }); + + it("resolves a symbol to a CoinGecko ID and returns its USD price", async () => { + process.env.PRICE_FEED_COIN_IDS = JSON.stringify({ TEST: "test-coin" }); + egressFetch = jest.spyOn(HttpEgressService.prototype, "fetch").mockResolvedValue({ + statusCode: 200, + body: JSON.stringify({ "test-coin": { usd: 12.5 } }), + } as EgressResponse); + + await expect(new CoinGeckoPriceFeedProvider().getUsdPrice("test")).resolves.toBe(12.5); + const [request, options] = egressFetch.mock.calls[0]; + expect(request).toContain("ids=test-coin"); + expect(options.purpose).toBe(EgressPurpose.ORACLE); + }); + + it("rejects aliases that do not resolve to a USD price", async () => { + process.env.PRICE_FEED_COIN_IDS = JSON.stringify({ TEST: "test-coin" }); + jest.spyOn(HttpEgressService.prototype, "fetch").mockResolvedValue({ + statusCode: 200, + body: JSON.stringify({ "test-coin": { usd: "12.5" } }), + } as EgressResponse); + + await expect(new CoinGeckoPriceFeedProvider().getUsdPrice("TEST")).rejects.toThrow( + "invalid USD price", + ); + }); + + it("rejects malformed symbol mappings during construction", () => { + process.env.PRICE_FEED_COIN_IDS = "[]"; + expect(() => new CoinGeckoPriceFeedProvider()).toThrow("must be a JSON object"); + }); +}); \ No newline at end of file diff --git a/src/tokens/price-feed.provider.ts b/src/tokens/price-feed.provider.ts index af3664d..33c82bc 100644 --- a/src/tokens/price-feed.provider.ts +++ b/src/tokens/price-feed.provider.ts @@ -1,3 +1,78 @@ +import { Injectable } from "@nestjs/common"; +import { EgressPurpose, HttpEgressService } from "../common/http-egress"; + export interface PriceFeedProvider { getUsdPrice(symbol: string): Promise; } + +export const PRICE_FEED_PROVIDER = Symbol("PRICE_FEED_PROVIDER"); + +const DEFAULT_COIN_IDS: Record = { + USDC: "usd-coin", + USDT: "tether", + WETH: "ethereum", + "WETH.E": "ethereum", + WBTC: "wrapped-bitcoin", + MATIC: "matic-network", + XLM: "stellar", +}; + +@Injectable() +export class CoinGeckoPriceFeedProvider implements PriceFeedProvider { + private readonly coinIds: Record; + private readonly egress = new HttpEgressService({ + timeoutMs: 10_000, + maxRedirects: 0, + maxBodySizeBytes: 16_384, + allowlist: ["api.coingecko.com"], + blockPrivateRanges: true, + }); + + constructor() { + let configured: Record = {}; + try { + const parsed: unknown = JSON.parse(process.env.PRICE_FEED_COIN_IDS ?? "{}"); + if ( + typeof parsed !== "object" || + parsed === null || + Array.isArray(parsed) || + Object.values(parsed).some((id) => typeof id !== "string" || id.trim() === "") + ) { + throw new Error("invalid mapping"); + } + configured = parsed as Record; + } catch { + throw new Error("PRICE_FEED_COIN_IDS must be a JSON object of token symbols to CoinGecko IDs"); + } + this.coinIds = { + ...DEFAULT_COIN_IDS, + ...Object.fromEntries( + Object.entries(configured).map(([symbol, id]) => [symbol.toUpperCase(), id]), + ), + }; + } + + async getUsdPrice(symbol: string): Promise { + const coinId = this.coinIds[symbol.toUpperCase()]; + if (!coinId) throw new Error(`No CoinGecko ID configured for token symbol ${symbol}`); + + const url = new URL("https://api.coingecko.com/api/v3/simple/price"); + url.searchParams.set("ids", coinId); + url.searchParams.set("vs_currencies", "usd"); + const apiKey = process.env.PRICE_FEED_API_KEY; + const response = await this.egress.fetch(url.toString(), { + purpose: EgressPurpose.ORACLE, + ...(apiKey ? { headers: { "x-cg-demo-api-key": apiKey } } : {}), + }); + if (response.statusCode < 200 || response.statusCode >= 300) { + throw new Error(`CoinGecko returned HTTP ${response.statusCode}`); + } + + const data = JSON.parse(response.body) as Record; + const price = data[coinId]?.usd; + if (typeof price !== "number" || !Number.isFinite(price) || price <= 0) { + throw new Error(`CoinGecko returned an invalid USD price for ${symbol}`); + } + return price; + } +} diff --git a/src/tokens/price-feed.worker.spec.ts b/src/tokens/price-feed.worker.spec.ts new file mode 100644 index 0000000..956c736 --- /dev/null +++ b/src/tokens/price-feed.worker.spec.ts @@ -0,0 +1,91 @@ +import { ConfigService } from "@nestjs/config"; +import { LeaderElectionService } from "../common/leader-election"; +import { KillSwitchService } from "../killswitch/killswitch.service"; +import { PrismaService } from "../prisma/prisma.service"; +import { PriceFeedProvider } from "./price-feed.provider"; +import { PriceFeedWorker } from "./price-feed.worker"; + +describe("PriceFeedWorker", () => { + const tokenFindMany = jest.fn(); + const tokenUpdate = jest.fn(); + const historyCreate = jest.fn(); + const transaction = jest.fn(); + const getUsdPrice = jest.fn(); + const pause = jest.fn(); + const registerWorker = jest.fn(); + const configGet = jest.fn(); + let worker: PriceFeedWorker; + + beforeEach(() => { + jest.clearAllMocks(); + tokenFindMany.mockResolvedValue([]); + transaction.mockResolvedValue([]); + configGet.mockImplementation((_key: string, fallback: number) => fallback); + worker = new PriceFeedWorker( + { + token: { findMany: tokenFindMany, update: tokenUpdate }, + priceHistory: { create: historyCreate }, + $transaction: transaction, + } as unknown as PrismaService, + { getUsdPrice } as PriceFeedProvider, + { pause } as unknown as KillSwitchService, + { registerWorker } as unknown as LeaderElectionService, + { get: configGet } as unknown as ConfigService, + ); + }); + + it("persists each valid price and a matching historical sample", async () => { + tokenFindMany.mockResolvedValue([ + { id: "token-1", symbol: "XLM", priceUsd: 0.1 }, + { id: "token-2", symbol: "XLM", priceUsd: 0.1 }, + ]); + getUsdPrice.mockResolvedValue(0.11); + + const result = await worker.refreshPrices(); + + expect(result).toEqual({ updated: 2, failed: 0, circuitBroken: false }); + expect(getUsdPrice).toHaveBeenCalledTimes(1); + expect(transaction).toHaveBeenCalledTimes(2); + expect(tokenUpdate).toHaveBeenCalledWith({ + where: { id: "token-1" }, + data: { priceUsd: 0.11 }, + }); + expect(historyCreate).toHaveBeenCalledWith({ + data: expect.objectContaining({ tokenId: "token-1", priceUsd: 0.11 }), + }); + expect(pause).not.toHaveBeenCalled(); + }); + + it("trips the global kill switch when a price moves past the threshold", async () => { + tokenFindMany.mockResolvedValue([{ id: "token-1", symbol: "XLM", priceUsd: 0.1 }]); + getUsdPrice.mockResolvedValue(0.2); + + const result = await worker.refreshPrices(); + + expect(result.circuitBroken).toBe(true); + expect(pause).toHaveBeenCalledWith( + expect.objectContaining({ + scope: "global", + reasonCode: "PRICE_FEED_EXTREME_MOVE", + activatedBy: "price-feed-worker", + }), + ); + }); + + it("does not persist invalid provider prices", async () => { + tokenFindMany.mockResolvedValue([{ id: "token-1", symbol: "XLM", priceUsd: null }]); + getUsdPrice.mockResolvedValue(Number.NaN); + + const result = await worker.refreshPrices(); + + expect(result).toEqual({ updated: 0, failed: 1, circuitBroken: false }); + expect(transaction).not.toHaveBeenCalled(); + expect(pause).not.toHaveBeenCalled(); + }); + + it("registers as a leader-elected worker", () => { + worker.onModuleInit(); + expect(registerWorker).toHaveBeenCalledWith("price-feed", expect.any(Function)); + worker.onModuleDestroy(); + }); +}); \ No newline at end of file diff --git a/src/tokens/price-feed.worker.ts b/src/tokens/price-feed.worker.ts new file mode 100644 index 0000000..fb97d08 --- /dev/null +++ b/src/tokens/price-feed.worker.ts @@ -0,0 +1,155 @@ +import { + Inject, + Injectable, + Logger, + OnModuleDestroy, + OnModuleInit, +} from "@nestjs/common"; +import { ConfigService } from "@nestjs/config"; +import { LeaderElectionService, Singleton } from "../common/leader-election"; +import { KillSwitchService } from "../killswitch/killswitch.service"; +import { PrismaService } from "../prisma/prisma.service"; +import { PRICE_FEED_PROVIDER, PriceFeedProvider } from "./price-feed.provider"; + +const DEFAULT_REFRESH_INTERVAL_MS = 60_000; +const DEFAULT_CIRCUIT_BREAKER_THRESHOLD_PERCENT = 50; + +@Singleton("price-feed") +@Injectable() +export class PriceFeedWorker implements OnModuleInit, OnModuleDestroy { + private readonly logger = new Logger(PriceFeedWorker.name); + private interval?: NodeJS.Timeout; + private refreshInProgress = false; + + constructor( + private readonly prisma: PrismaService, + @Inject(PRICE_FEED_PROVIDER) private readonly provider: PriceFeedProvider, + private readonly killSwitch: KillSwitchService, + private readonly leaderElection: LeaderElectionService, + private readonly config: ConfigService, + ) {} + + onModuleInit(): void { + this.leaderElection.registerWorker("price-feed", (isLeader) => { + if (isLeader) { + if (this.interval) return; + void this.refreshPrices().catch((error: unknown) => this.logFailure(error)); + this.interval = setInterval(() => { + void this.refreshPrices().catch((error: unknown) => this.logFailure(error)); + }, this.refreshIntervalMs()); + } else { + this.stopInterval(); + } + }); + } + + onModuleDestroy(): void { + this.stopInterval(); + } + + async refreshPrices(): Promise<{ updated: number; failed: number; circuitBroken: boolean }> { + if (this.refreshInProgress) { + return { updated: 0, failed: 0, circuitBroken: false }; + } + this.refreshInProgress = true; + try { + return await this.runRefresh(); + } finally { + this.refreshInProgress = false; + } + } + + private async runRefresh(): Promise<{ updated: number; failed: number; circuitBroken: boolean }> { + const tokens = await this.prisma.token.findMany({ + select: { id: true, symbol: true, priceUsd: true }, + }); + const pricesBySymbol = new Map(); + let updated = 0; + let failed = 0; + const extremeMoves: Array<{ symbol: string; previous: number; current: number; change: number }> = []; + const recordedAt = new Date(); + const threshold = this.breakerThresholdPercent(); + + for (const token of tokens) { + try { + let price = pricesBySymbol.get(token.symbol.toUpperCase()); + if (price === undefined) { + price = await this.provider.getUsdPrice(token.symbol); + if (!Number.isFinite(price) || price <= 0) { + throw new Error(`Price provider returned an invalid price for ${token.symbol}`); + } + pricesBySymbol.set(token.symbol.toUpperCase(), price); + } + + await this.prisma.$transaction([ + this.prisma.token.update({ where: { id: token.id }, data: { priceUsd: price } }), + this.prisma.priceHistory.create({ + data: { tokenId: token.id, priceUsd: price, recordedAt }, + }), + ]); + updated++; + + if (token.priceUsd !== null && token.priceUsd > 0) { + const change = (Math.abs(price - token.priceUsd) / token.priceUsd) * 100; + if (change > threshold) { + extremeMoves.push({ + symbol: token.symbol, + previous: token.priceUsd, + current: price, + change, + }); + } + } + } catch (error) { + failed++; + this.logger.warn( + `Price refresh failed for ${token.symbol} (${token.id}): ${this.errorMessage(error)}`, + ); + } + } + + if (extremeMoves.length > 0) { + const move = extremeMoves[0]; + const reason = extremeMoves + .map(({ symbol, previous, current, change }) => + `${symbol} moved ${change.toFixed(2)}% from $${previous} to $${current}`, + ) + .join("; "); + await this.killSwitch.pause({ + scope: "global", + reasonCode: "PRICE_FEED_EXTREME_MOVE", + reason: `Price circuit breaker triggered: ${reason}`, + activatedBy: "price-feed-worker", + }); + this.logger.error( + `Price circuit breaker triggered by ${move.symbol}: ${move.change.toFixed(2)}% move`, + ); + } + + return { updated, failed, circuitBroken: extremeMoves.length > 0 }; + } + + private refreshIntervalMs(): number { + return this.config.get("PRICE_FEED_REFRESH_INTERVAL_MS", DEFAULT_REFRESH_INTERVAL_MS); + } + + private breakerThresholdPercent(): number { + return this.config.get( + "PRICE_FEED_CIRCUIT_BREAKER_THRESHOLD_PERCENT", + DEFAULT_CIRCUIT_BREAKER_THRESHOLD_PERCENT, + ); + } + + private stopInterval(): void { + if (this.interval) clearInterval(this.interval); + this.interval = undefined; + } + + private logFailure(error: unknown): void { + this.logger.error(`Price refresh cycle failed: ${this.errorMessage(error)}`); + } + + private errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); + } +} \ No newline at end of file diff --git a/src/tokens/tokens.module.ts b/src/tokens/tokens.module.ts index 908c14f..5e50bec 100644 --- a/src/tokens/tokens.module.ts +++ b/src/tokens/tokens.module.ts @@ -5,6 +5,8 @@ import { TokensService } from "./tokens.service"; import { TOKENS_REPOSITORY } from "./tokens.repository"; import { InMemoryTokensRepository } from "./in-memory-tokens.repository"; import { PrismaTokensRepository } from "./prisma-tokens.repository"; +import { PriceFeedWorker } from "./price-feed.worker"; +import { CoinGeckoPriceFeedProvider, PRICE_FEED_PROVIDER } from "./price-feed.provider"; @Module({ controllers: [TokensController], @@ -23,6 +25,9 @@ import { PrismaTokensRepository } from "./prisma-tokens.repository"; }, }, TokensService, + PriceFeedWorker, + CoinGeckoPriceFeedProvider, + { provide: PRICE_FEED_PROVIDER, useExisting: CoinGeckoPriceFeedProvider }, ], exports: [TokensService], })