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
7 changes: 7 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,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().
Expand Down
5 changes: 5 additions & 0 deletions .env.mainnet.example
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
5 changes: 5 additions & 0 deletions .env.staging.example
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
1 change: 1 addition & 0 deletions prisma/migrations/20260930000001_price_history/down.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
DROP TABLE "price_history";
13 changes: 13 additions & 0 deletions prisma/migrations/20260930000001_price_history/migration.sql
Original file line number Diff line number Diff line change
@@ -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);
46 changes: 46 additions & 0 deletions src/tokens/price-feed.provider.spec.ts
Original file line number Diff line number Diff line change
@@ -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");
});
});
75 changes: 75 additions & 0 deletions src/tokens/price-feed.provider.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,6 @@
import { Injectable } from "@nestjs/common";
import { EgressPurpose, HttpEgressService } from "../common/http-egress";

export interface PriceFeedProvider {
/**
* Return a positive USD quote for a token symbol. Consumers making
Expand All @@ -6,3 +9,75 @@ export interface PriceFeedProvider {
*/
getUsdPrice(symbol: string): Promise<number>;
}

export const PRICE_FEED_PROVIDER = Symbol("PRICE_FEED_PROVIDER");

const DEFAULT_COIN_IDS: Record<string, string> = {
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<string, string>;
private readonly egress = new HttpEgressService({
timeoutMs: 10_000,
maxRedirects: 0,
maxBodySizeBytes: 16_384,
allowlist: ["api.coingecko.com"],
blockPrivateRanges: true,
});

constructor() {
let configured: Record<string, string> = {};
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<string, string>;
} 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<number> {
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<string, { usd?: unknown }>;
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;
}
}
91 changes: 91 additions & 0 deletions src/tokens/price-feed.worker.spec.ts
Original file line number Diff line number Diff line change
@@ -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();
});
});
Loading
Loading