feat(relayer): цена SOL из NATS + WS-раздача + деплой-артефакты
deploy-contract / deploy-contract (release) Failing after 3m34s
deploy-contract / deploy-contract (release) Failing after 3m34s
priceSource больше не рандомит: цена берётся из NATS (market.price.solusdt), конверсия строка->целое x1e8 без float. Шлюз раздаёт цену по WS (/ws, нормализация /api/... -> /...) — фронт просто выводит её. Добавлены Dockerfile и helm-чарт для деплоя в k3s.
This commit is contained in:
+157
-5
@@ -9,6 +9,7 @@
|
||||
*
|
||||
* Эндпоинты (CORS: Access-Control-Allow-Origin: *):
|
||||
* GET /state — Global: counter, min/max, multiplier, expiry, paused, vault
|
||||
* GET /price — текущая сырая цена SOL из NATS (для отладки); 503 если ещё нет
|
||||
* POST /bet {address?, side, amountUnits}
|
||||
* — address: custodial-беттор шлюза; если не задан — шлюз
|
||||
* генерирует новый и регистрирует в keystore.
|
||||
@@ -18,9 +19,17 @@
|
||||
* GET /faucet/{address} — налив SOL + фантиков на произвольный адрес
|
||||
* (только localnet, allowAirdrop=1; иначе 403). Ответ:
|
||||
* {address, sol, solLamports, tokenUnits, tokenHuman}
|
||||
* GET /ws (WebSocket) — push цены SOL из NATS; на подключение шлёт
|
||||
* снапшот (если цена уже есть), дальше — каждый тик.
|
||||
*
|
||||
* Пути нормализуются: ingress пробрасывает путь как есть, поэтому
|
||||
* /api/... и /... должны работать одинаково (strip /api).
|
||||
*
|
||||
* Запуск: npm run gateway (ts-node src/gateway.ts)
|
||||
* Env: GW_PORT (default 8895), GW_BETTOR_KEYSTORE, SMART_* (как у релейера).
|
||||
* Env: GW_PORT (default 8895), GW_BETTOR_KEYSTORE, SMART_* (как у релейера),
|
||||
* NATS_URL (default nats://192.168.88.93:4222),
|
||||
* NATS_PRICE_SUBJECT (default market.price.solusdt),
|
||||
* PRICE_WAIT_TIMEOUT_MS (default 15000).
|
||||
*/
|
||||
import http from "http";
|
||||
import * as fs from "fs";
|
||||
@@ -29,6 +38,7 @@ import { Keypair, LAMPORTS_PER_SOL, PublicKey } from "@solana/web3.js";
|
||||
import { getAssociatedTokenAddress, getMint } from "@solana/spl-token";
|
||||
import BN from "bn.js";
|
||||
import type { Program } from "@anchor-lang/core";
|
||||
import { WebSocketServer, WebSocket } from "ws";
|
||||
import { loadConfig, loadKeypair, SIDE_DOWN, SIDE_UP } from "./config";
|
||||
import {
|
||||
closeBet,
|
||||
@@ -40,7 +50,8 @@ import {
|
||||
} from "./ops";
|
||||
import { betPda } from "./pda";
|
||||
import { makeConnection, makeProgram } from "./program";
|
||||
import { getPrice } from "./priceSource";
|
||||
import { currentPrice, onPrice, connectNats } from "./natsPrice";
|
||||
import { currentPriceHuman, getPrice } from "./priceSource";
|
||||
|
||||
import type { SmartUpdown } from "../idl/smart_updown";
|
||||
|
||||
@@ -51,6 +62,22 @@ const STATUS_HOUSE_WON = 2;
|
||||
const LAMPORTS_PER_SOL_NUM = Number(LAMPORTS_PER_SOL);
|
||||
const MAX_BODY_BYTES = 100 * 1024;
|
||||
|
||||
const NATS_URL_DEFAULT = "nats://192.168.88.93:4222";
|
||||
const NATS_PRICE_SUBJECT_DEFAULT = "market.price.solusdt";
|
||||
const WS_PATH = "/ws";
|
||||
const WS_HEARTBEAT_MS = 30_000;
|
||||
const WS_MISSED_PONGS = 2;
|
||||
|
||||
/**
|
||||
* Ingress пробрасывает путь как есть, без stripPrefix. Нормализуем
|
||||
* `/api/...` к `/...`, чтобы `/api/state` и `/state` работали одинаково.
|
||||
*/
|
||||
function normalizePath(p: string): string {
|
||||
if (p === "/api" || p === "/api/") return "/";
|
||||
if (p.startsWith("/api/")) return p.slice(4);
|
||||
return p;
|
||||
}
|
||||
|
||||
/** Токен-юниты -> "N.NNNNNN" по decimals (дефолт 6, как в контракте). */
|
||||
function humanize(units: string | number | bigint, decimals: number): string {
|
||||
const n = BigInt(units);
|
||||
@@ -170,6 +197,22 @@ async function handleState(ctx: Ctx): Promise<HttpJson> {
|
||||
};
|
||||
}
|
||||
|
||||
async function handlePrice(): Promise<HttpJson> {
|
||||
const human = currentPriceHuman();
|
||||
if (human === null) {
|
||||
return { status: 503, body: { error: "price not available yet" } };
|
||||
}
|
||||
const last = currentPrice();
|
||||
return {
|
||||
status: 200,
|
||||
body: {
|
||||
symbol: "solusdt",
|
||||
price: human,
|
||||
ts: last ? last.ts : Date.now(),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
async function handleBet(ctx: Ctx, body: Record<string, unknown>): Promise<HttpJson> {
|
||||
const side = normalizeSide(body.side);
|
||||
if (side === null) {
|
||||
@@ -404,6 +447,28 @@ function readJsonBody(req: http.IncomingMessage): Promise<Record<string, unknown
|
||||
});
|
||||
}
|
||||
|
||||
/** WS-клиенты с флагом живости (2 цикла пропущенных pong -> terminate). */
|
||||
interface WsClient {
|
||||
ws: WebSocket;
|
||||
alive: boolean;
|
||||
}
|
||||
|
||||
function broadcastPrice(wsClients: Set<WsClient>): (p: { price: string; ts: number }) => void {
|
||||
return (p) => {
|
||||
const msg = JSON.stringify({
|
||||
type: "price",
|
||||
symbol: "solusdt",
|
||||
price: p.price,
|
||||
ts: p.ts,
|
||||
});
|
||||
for (const c of wsClients) {
|
||||
if (c.ws.readyState === WebSocket.OPEN) {
|
||||
c.ws.send(msg);
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
async function main(): Promise<void> {
|
||||
const cfg = loadConfig();
|
||||
const admin = loadKeypair(cfg.adminKeypairPath);
|
||||
@@ -437,6 +502,14 @@ async function main(): Promise<void> {
|
||||
console.log("gateway: contract NOT initialized — /state будет 409 до `npm run init`");
|
||||
}
|
||||
|
||||
// NATS-подписка на цену: НЕ блокируем старт сервера — если NATS недоступен,
|
||||
// шлюз всё равно поднимается, /bet будет ждать цену до таймаута.
|
||||
const natsUrl = process.env.NATS_URL ?? NATS_URL_DEFAULT;
|
||||
const natsSubject = process.env.NATS_PRICE_SUBJECT ?? NATS_PRICE_SUBJECT_DEFAULT;
|
||||
connectNats(natsUrl, natsSubject).catch((e) => {
|
||||
console.error(`gateway: NATS connect failed (${natsUrl}): ${e instanceof Error ? e.message : String(e)}`);
|
||||
});
|
||||
|
||||
const port = Number(process.env.GW_PORT ?? "8895");
|
||||
const server = http.createServer((req, res) => {
|
||||
const started = Date.now();
|
||||
@@ -455,13 +528,14 @@ async function main(): Promise<void> {
|
||||
res.writeHead(status);
|
||||
res.end();
|
||||
}
|
||||
// Логируем ИСХОДНЫЙ путь, как требует ТЗ.
|
||||
const url = req.url ?? "/";
|
||||
console.log(`gateway: ${req.method} ${url} -> ${status} (${Date.now() - started}ms)`);
|
||||
};
|
||||
|
||||
(async () => {
|
||||
const reqUrl = new URL(req.url ?? "/", "http://gateway.local");
|
||||
const p = reqUrl.pathname;
|
||||
const p = normalizePath(reqUrl.pathname);
|
||||
const m = req.method ?? "GET";
|
||||
|
||||
if (m === "OPTIONS") {
|
||||
@@ -473,6 +547,8 @@ async function main(): Promise<void> {
|
||||
try {
|
||||
if (m === "GET" && p === "/state") {
|
||||
result = await handleState(ctx);
|
||||
} else if (m === "GET" && p === "/price") {
|
||||
result = await handlePrice();
|
||||
} else if (m === "POST" && p === "/bet") {
|
||||
const body = await readJsonBody(req);
|
||||
result = await withBetLock(() => handleBet(ctx, body));
|
||||
@@ -509,15 +585,91 @@ async function main(): Promise<void> {
|
||||
});
|
||||
});
|
||||
|
||||
// WS-сервер на том же порту (noServer=true), обновляем апгрейды вручную.
|
||||
const wss = new WebSocketServer({ noServer: true });
|
||||
const wsClients = new Set<WsClient>();
|
||||
|
||||
wss.on("connection", (ws: WebSocket) => {
|
||||
const client: WsClient = { ws, alive: true };
|
||||
wsClients.add(client);
|
||||
ws.on("pong", () => {
|
||||
client.alive = true;
|
||||
});
|
||||
ws.on("close", () => {
|
||||
wsClients.delete(client);
|
||||
});
|
||||
ws.on("error", () => {
|
||||
wsClients.delete(client);
|
||||
});
|
||||
// Снапшот сразу, если цена уже пришла из NATS.
|
||||
const snap = currentPrice();
|
||||
if (snap !== null) {
|
||||
const msg = JSON.stringify({
|
||||
type: "price",
|
||||
symbol: "solusdt",
|
||||
price: snap.price,
|
||||
ts: snap.ts,
|
||||
});
|
||||
ws.send(msg);
|
||||
}
|
||||
});
|
||||
|
||||
// Рассылка каждого нового NATS-тика по всем подключённым WS.
|
||||
onPrice(broadcastPrice(wsClients));
|
||||
|
||||
// Heartbeat: пинг раз в 30с, нет pong 2 цикла — terminate.
|
||||
const heartbeat = setInterval(() => {
|
||||
for (const c of wsClients) {
|
||||
if (!c.alive) {
|
||||
try {
|
||||
c.ws.terminate();
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
wsClients.delete(c);
|
||||
continue;
|
||||
}
|
||||
c.alive = false;
|
||||
try {
|
||||
c.ws.ping();
|
||||
} catch {
|
||||
wsClients.delete(c);
|
||||
}
|
||||
}
|
||||
}, WS_HEARTBEAT_MS);
|
||||
|
||||
server.on("upgrade", (req, socket, head) => {
|
||||
const rawUrl = req.url ?? "/";
|
||||
const reqUrl = new URL(rawUrl, "http://gateway.local");
|
||||
const p = normalizePath(reqUrl.pathname);
|
||||
if (p !== WS_PATH) {
|
||||
socket.destroy();
|
||||
return;
|
||||
}
|
||||
wss.handleUpgrade(req, socket, head, (ws) => {
|
||||
wss.emit("connection", ws, req);
|
||||
});
|
||||
});
|
||||
|
||||
server.listen(port, "0.0.0.0", () => {
|
||||
console.log(`gateway: listening on 0.0.0.0:${port} (RPC ${cfg.rpcUrl}, mint ${cfg.tokenMint.toBase58()})`);
|
||||
console.log(
|
||||
`gateway: endpoints: GET /state, POST /bet, POST /close, GET /balance/{address}, GET /faucet/{address}`
|
||||
`gateway: endpoints: GET /state, GET /price, POST /bet, POST /close, GET /balance/{address}, GET /faucet/{address}, WS /ws`
|
||||
);
|
||||
});
|
||||
|
||||
// Корректная остановка heartbeat при завершении.
|
||||
const shutdown = () => {
|
||||
clearInterval(heartbeat);
|
||||
wss.close();
|
||||
server.close();
|
||||
process.exit(0);
|
||||
};
|
||||
process.on("SIGINT", shutdown);
|
||||
process.on("SIGTERM", shutdown);
|
||||
}
|
||||
|
||||
main().catch((e) => {
|
||||
console.error(e);
|
||||
process.exit(1);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,166 @@
|
||||
/**
|
||||
* NATS-источник цены SOL.
|
||||
*
|
||||
* Подписывается на subject вида `market.price.solusdt` (payload — JSON
|
||||
* `{"symbol":"solusdt","type":"price","price":"101.89000000",...}`),
|
||||
* хранит последнее валидное значение и рассылает его всем подписчикам.
|
||||
*
|
||||
* Reconnect/autoreconnect включены по умолчанию (пакет `nats`); в Core NATS
|
||||
* история не сохраняется — новые значения приходят только живой подписке.
|
||||
*
|
||||
* Публикаций наружу нет — только чтение.
|
||||
*/
|
||||
import { connect, StringCodec, NatsConnection } from "nats";
|
||||
|
||||
export interface PriceTick {
|
||||
/** Сырая строка цены из payload (например, "101.89000000"). */
|
||||
price: string;
|
||||
/** Unix ms: берётся из поля `ts` payload'а, иначе Date.now(). */
|
||||
ts: number;
|
||||
}
|
||||
|
||||
let nc: NatsConnection | null = null;
|
||||
let last: PriceTick | null = null;
|
||||
const listeners: Array<(p: PriceTick) => void> = [];
|
||||
let firstLogged = false;
|
||||
|
||||
function notify(p: PriceTick): void {
|
||||
last = p;
|
||||
for (const cb of listeners) {
|
||||
try {
|
||||
cb(p);
|
||||
} catch (e) {
|
||||
console.error(`natsPrice: listener threw: ${e instanceof Error ? e.message : String(e)}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function isValidPriceString(s: unknown): s is string {
|
||||
if (typeof s !== "string") return false;
|
||||
if (s.length === 0) return false;
|
||||
// Только цифры и одна точка. Иначе — мусор.
|
||||
if (!/^[0-9]+(?:\.[0-9]+)?$/.test(s)) return false;
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Подключиться к NATS и подписаться на subject. Возвращает управление
|
||||
* сразу после установки подписки, не дожидаясь первой цены. Сообщения с
|
||||
* невалидным JSON / невалидным `price` молча пропускаются (с логом).
|
||||
*/
|
||||
export async function connectNats(url: string, subject: string): Promise<void> {
|
||||
if (nc) {
|
||||
console.log(`natsPrice: already connected, ignoring reconnect request to ${url}`);
|
||||
return;
|
||||
}
|
||||
const conn = await connect({ servers: [url] });
|
||||
nc = conn;
|
||||
console.log(`nats: connected ${url}`);
|
||||
|
||||
const sc = StringCodec();
|
||||
const sub = conn.subscribe(subject, {
|
||||
callback: (err, msg) => {
|
||||
if (err) {
|
||||
console.error(`natsPrice: subscription error: ${err.message}`);
|
||||
return;
|
||||
}
|
||||
let parsed: unknown;
|
||||
try {
|
||||
parsed = JSON.parse(sc.decode(msg.data));
|
||||
} catch (e) {
|
||||
console.log(
|
||||
`natsPrice: malformed JSON on ${subject}: ${e instanceof Error ? e.message : String(e)}`
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (!parsed || typeof parsed !== "object") {
|
||||
console.log(`natsPrice: non-object payload on ${subject}`);
|
||||
return;
|
||||
}
|
||||
const obj = parsed as Record<string, unknown>;
|
||||
const priceRaw = obj.price;
|
||||
if (!isValidPriceString(priceRaw)) {
|
||||
console.log(`natsPrice: invalid price on ${subject}: ${String(priceRaw)}`);
|
||||
return;
|
||||
}
|
||||
const tsNum =
|
||||
typeof obj.ts === "number" && Number.isFinite(obj.ts)
|
||||
? (obj.ts as number)
|
||||
: Date.now();
|
||||
const tick: PriceTick = { price: priceRaw, ts: tsNum };
|
||||
if (!firstLogged) {
|
||||
console.log(`natsPrice: first price ${priceRaw} (ts=${tsNum})`);
|
||||
firstLogged = true;
|
||||
}
|
||||
notify(tick);
|
||||
},
|
||||
});
|
||||
// keep reference so it isn't GC'd
|
||||
void sub;
|
||||
}
|
||||
|
||||
/** Текущая последняя цена (или null, если ещё не приходила). */
|
||||
export function currentPrice(): PriceTick | null {
|
||||
return last;
|
||||
}
|
||||
|
||||
/** Регистрирует слушателя на КАЖДОЕ новое значение. */
|
||||
export function onPrice(cb: (p: PriceTick) => void): void {
|
||||
listeners.push(cb);
|
||||
}
|
||||
|
||||
/**
|
||||
* Ждёт первую цену из NATS и возвращает её в ЦЕЛЫХ ЕДИНИЦАХ
|
||||
* (price × 1e8, через строковое преобразование — без потери точности).
|
||||
* При таймауте — throw.
|
||||
*/
|
||||
export async function waitForPrice(timeoutMs: number): Promise<number> {
|
||||
if (last !== null) {
|
||||
return priceStringToInt(last.price);
|
||||
}
|
||||
return new Promise<number>((resolve, reject) => {
|
||||
const timer = setTimeout(() => {
|
||||
reject(new Error(`price timeout: no data from NATS after ${timeoutMs}ms`));
|
||||
}, timeoutMs);
|
||||
onPrice((p) => {
|
||||
clearTimeout(timer);
|
||||
resolve(priceStringToInt(p.price));
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Чистая конверсия: сырая строка из NATS → целое (price × 1e8).
|
||||
* "101.89000000" → 10189000000
|
||||
* "0.00000001" → 1
|
||||
* "1234567.12345678" → 123456712345678
|
||||
* "100" → 10000000000
|
||||
* Без float — только string/BigInt.
|
||||
*/
|
||||
export function priceStringToInt(raw: string): number {
|
||||
if (typeof raw !== "string" || raw.length === 0) {
|
||||
throw new Error(`priceStringToInt: empty/invalid input`);
|
||||
}
|
||||
if (!/^[0-9]+(?:\.[0-9]+)?$/.test(raw)) {
|
||||
throw new Error(`priceStringToInt: malformed price string: ${raw}`);
|
||||
}
|
||||
const dot = raw.indexOf(".");
|
||||
let intPart: string;
|
||||
let fracPart: string;
|
||||
if (dot === -1) {
|
||||
intPart = raw;
|
||||
fracPart = "";
|
||||
} else {
|
||||
intPart = raw.slice(0, dot);
|
||||
fracPart = raw.slice(dot + 1);
|
||||
}
|
||||
if (fracPart.length > 8) {
|
||||
fracPart = fracPart.slice(0, 8);
|
||||
} else if (fracPart.length < 8) {
|
||||
fracPart = fracPart.padEnd(8, "0");
|
||||
}
|
||||
const combined = `${intPart}${fracPart}`;
|
||||
// strip leading zeros
|
||||
const stripped = combined.replace(/^0+(?=\d)/, "");
|
||||
return Number(stripped);
|
||||
}
|
||||
+34
-13
@@ -1,22 +1,43 @@
|
||||
/**
|
||||
* STUB — SOL price source.
|
||||
*
|
||||
* TODO(пользователь): замени `getPrice` на реальный запрос к своей БД.
|
||||
* Договоримся, что функция возвращает цену SOL как целое число в «центах»
|
||||
* (то же юнит-пространство, что у веса entry/exit на контракте). Контракт
|
||||
* сравнивает тол ько «больше/меньше/равно», поэтому само юнит-пространство
|
||||
* не критично — главное, чтобы оно было согласовано между entry и exit.
|
||||
* Цена SOL. Источник — NATS (`relayer/src/natsPrice.ts`). `getPrice()`
|
||||
* возвращает цену как ЦЕЛОЕ (price × 1e8, без float — точность нужна для
|
||||
* сравнения на контракте). Override из старой заглушки сохранён для
|
||||
* тестов «equal price» / offline-сценариев.
|
||||
*/
|
||||
let simulated = 100_000_000;
|
||||
import { currentPrice, priceStringToInt, waitForPrice } from "./natsPrice";
|
||||
|
||||
/** Цена SOL сейчас (заглушка — детерминированно «гуляет»). */
|
||||
const PRICE_WAIT_TIMEOUT_MS = Number(process.env.PRICE_WAIT_TIMEOUT_MS ?? "15000");
|
||||
|
||||
let override: number | undefined = undefined;
|
||||
|
||||
/**
|
||||
* Вернуть целое число = price × 1e8. Конверсия — через строковый парсер
|
||||
* (`priceStringToInt`), БЕЗ `Number(price) * 1e8` (float теряет точность
|
||||
* на 8 знаках).
|
||||
*
|
||||
* Override выставлен → возвращаем его и НЕ ходим в NATS.
|
||||
* NATS ещё не прислал цену → ждём первую до PRICE_WAIT_TIMEOUT_MS.
|
||||
*/
|
||||
export async function getPrice(): Promise<number> {
|
||||
const drift = Math.round(Math.random() * 200_000); // ±0.2% движка
|
||||
simulated += Math.random() < 0.5 ? drift : -drift;
|
||||
return simulated;
|
||||
if (override !== undefined) return override;
|
||||
// waitForPrice сам разруливает "уже есть" vs "ещё ждём".
|
||||
return waitForPrice(PRICE_WAIT_TIMEOUT_MS);
|
||||
}
|
||||
|
||||
/** Управляемая заглушка: вернуть ровно `value` (полезна для теста «equal price»). */
|
||||
export function setPriceOverride(value: number | undefined): void {
|
||||
simulated = value ?? 100_000_000;
|
||||
override = value;
|
||||
}
|
||||
|
||||
/**
|
||||
* Последняя СЫРАЯ строка цены из NATS (например, `"101.89000000"`) —
|
||||
* для UI. `null`, если цены ещё не было. При выставленном override
|
||||
* возвращает его строковую форму с 8 знаками (UI-привычно).
|
||||
*/
|
||||
export function currentPriceHuman(): string | null {
|
||||
if (override !== undefined) {
|
||||
return override.toString().padStart(9, "0"); // минимум 1 знак слева
|
||||
}
|
||||
const p = currentPrice();
|
||||
return p ? p.price : null;
|
||||
}
|
||||
Reference in New Issue
Block a user