feat(relayer): авто-закрытие ставок, WS-события, одна открытая ставка на адрес

- settle.ts: серверный планировщик (startScheduler, тик SMART_SETTLE_INTERVAL_MS)
  находит просроченные открытые ставки, берёт цену из истории на expire_time,
  сам зовёт close_bet; рассылает bet_closed
- gateway.ts: ручка POST /close удалена (404); проверка 'одна открытая ставка
  на адрес' под withBetLock -> 409 с openBetId; bet_open по WS при создании;
  lastState прокидывается в settleOnce (без него лог спамил каждый тик)
- settle.ts: пакетный getProgramAccounts вместо поштучного перебора по id
  (замер: 3 мс против 28 мс, 9.3x, 1 запрос вместо 15); фильтр не-Bet аккаунтов
  по BET_ACCOUNT_SIZE (Global исключён из unmapped)
- check.ts: +38 проверок (decodeBet, openBetOf, dueForSettle, события, границы)
This commit is contained in:
Caffeine
2026-09-13 18:49:26 +03:00
parent 0b136c6502
commit a73e44de9a
6 changed files with 843 additions and 91 deletions
+2
View File
@@ -73,6 +73,8 @@ spec:
value: {{ .Values.smart.priceHistory | quote }}
- name: SMART_PRICE_RETENTION_HOURS
value: {{ .Values.smart.priceRetentionHours | quote }}
- name: SMART_SETTLE_INTERVAL_MS
value: {{ .Values.smart.settleIntervalMs | quote }}
- name: GW_PORT
value: {{ .Values.port | quote }}
- name: GW_BETTOR_KEYSTORE
+1
View File
@@ -28,6 +28,7 @@ smart:
startBalance: 10000000000
priceHistory: '/app/data/price-history.db'
priceRetentionHours: 24
settleIntervalMs: 2000
nats:
url: 'nats://192.168.88.93:4222'
+177 -2
View File
@@ -8,8 +8,8 @@
* 3. Гейт faucet-ручки: release запрещает на любом кластере, test разрешает
* только на непубличных (genesis-hash не mainnet/devnet/testnet), неизвестный
* hash запрещён (fail-closed).
*
* Запуск: `npm test`.
* 4. Чистые функции авто-закрытия: decodeBet / openBetOf / dueForSettle /
* betOpenEvent / betClosedEvent (settle.ts).
*/
import { Keypair, PublicKey } from "@solana/web3.js";
import BN from "bn.js";
@@ -22,6 +22,18 @@ import { globalPda, betPda, vaultPda } from "./pda";
import { makeConnection, makeProgram } from "./program";
import { faucetDenyReason, parseMode } from "./faucetGate";
import { DEFAULT_RETENTION_MS, PriceHistory } from "./priceHistory";
import {
BET_ACCOUNT_SIZE,
betClosedEvent,
betOpenEvent,
decodeBet,
dueForSettle,
openBetOf,
BetAccountEntry,
STATUS_OPEN,
STATUS_PAYOUT_DONE,
STATUS_HOUSE_WON,
} from "./settle";
const ok = (msg: string) => console.log(`ok: ${msg}`);
let failed = 0;
@@ -50,6 +62,42 @@ function idlIx(name: string) {
return ix;
}
/**
* Строит буфер Bet-аккаунта (96 байт) для оффлайн-тестов settle.ts.
* Дискриминатор (0..7) забит нулями — decodeBet его не читает, остальные
* поля — по контракту.
*/
function makeBetBuffer(opts: {
amount?: number;
betTime?: number;
expireTime?: number;
entryPrice?: number;
exitPrice?: number;
side?: number;
status?: number;
bettor?: PublicKey;
}): Buffer {
const buf = Buffer.alloc(BET_ACCOUNT_SIZE);
// Discriminator 0..7: zeros
if (opts.amount !== undefined) buf.writeBigUInt64LE(BigInt(opts.amount), 8);
if (opts.betTime !== undefined) buf.writeBigInt64LE(BigInt(opts.betTime), 16);
if (opts.expireTime !== undefined) buf.writeBigInt64LE(BigInt(opts.expireTime), 24);
if (opts.entryPrice !== undefined) buf.writeBigUInt64LE(BigInt(opts.entryPrice), 32);
if (opts.exitPrice !== undefined) buf.writeBigUInt64LE(BigInt(opts.exitPrice), 40);
if (opts.side !== undefined) buf.writeBigUInt64LE(BigInt(opts.side), 48);
if (opts.status !== undefined) buf.writeBigUInt64LE(BigInt(opts.status), 56);
if (opts.bettor !== undefined) {
const bytes = opts.bettor.toBytes();
Buffer.from(bytes).copy(buf, 64);
}
return buf;
}
/** Оборачивает сырой буфер в BetAccountEntry с заданным id. */
function betEntry(id: number, buf: Buffer): BetAccountEntry {
return { pubkey: Keypair.generate().publicKey, account: { data: buf }, id };
}
async function main(): Promise<void> {
const programId = new PublicKey((idl as any).address as string);
const createIx = idlIx("create_bet");
@@ -249,6 +297,133 @@ async function main(): Promise<void> {
"priceHistory: ретеншен по умолчанию = 24 часа"
);
// === 5. Чистые функции авто-закрытия (settle.ts) ===
// Тест: decodeBet разбирает реальные байты.
{
const bettor = Keypair.generate().publicKey;
const buf = makeBetBuffer({
amount: 100_000_000,
betTime: 1_700_000_000,
expireTime: 1_700_000_015,
entryPrice: 10_077_000_000,
exitPrice: 0,
side: 1, // DOWN
status: 0, // OPEN
bettor,
});
const d = decodeBet(buf);
check(d.amount === 100_000_000, "decodeBet: amount");
check(d.betTime === 1_700_000_000, "decodeBet: betTime");
check(d.expireTime === 1_700_000_015, "decodeBet: expireTime");
check(d.entryPrice === 10_077_000_000, "decodeBet: entryPrice");
check(d.exitPrice === 0, "decodeBet: exitPrice");
check(d.side === 1, "decodeBet: side");
check(d.status === 0, "decodeBet: status");
check(d.bettor.toBase58() === bettor.toBase58(), "decodeBet: bettor (base58 round-trip)");
}
// Тест: openBetOf находит открытую.
{
const target = Keypair.generate().publicKey;
const other = Keypair.generate().publicKey;
const accounts: BetAccountEntry[] = [
betEntry(0, makeBetBuffer({ status: 1, bettor: target })),
betEntry(1, makeBetBuffer({ status: 0, bettor: target })), // <-- первая открытая для target
betEntry(2, makeBetBuffer({ status: 0, bettor: other })),
];
const r = openBetOf(target.toBase58(), accounts);
check(r !== null && r.id === 1, "openBetOf находит открытую: вернул её id");
check(r !== null && r.id !== 0, "openBetOf находит открытую: не нулевую");
}
// Тест: openBetOf игнорирует закрытые.
{
const target = Keypair.generate().publicKey;
const accounts: BetAccountEntry[] = [
betEntry(0, makeBetBuffer({ status: 1, bettor: target })),
betEntry(1, makeBetBuffer({ status: 2, bettor: target })),
betEntry(2, makeBetBuffer({ status: 1, bettor: target })),
];
const r = openBetOf(target.toBase58(), accounts);
check(r === null, "openBetOf игнорирует закрытые: вернул null");
}
// Тест: dueForSettle берёт только просроченные открытые.
{
const now = 2_000_000_000; // секунды
const nowMs = now * 1000;
const accounts: BetAccountEntry[] = [
// открытый просроченный → должен попасть
betEntry(0, makeBetBuffer({ status: 0, expireTime: now - 10 })),
// открытый свежий → не должен
betEntry(1, makeBetBuffer({ status: 0, expireTime: now + 10 })),
// закрытый просроченный → не должен
betEntry(2, makeBetBuffer({ status: 1, expireTime: now - 10 })),
// закрытый свежий → не должен
betEntry(3, makeBetBuffer({ status: 2, expireTime: now + 10 })),
];
const out = dueForSettle(accounts, nowMs);
check(out.length === 1, "dueForSettle берёт только просроченные открытые: ровно 1");
check(out.length === 1 && out[0].id === 0, "dueForSettle берёт только просроченные открытые: id=0");
}
// Тест: dueForSettle уважает границу expire.
{
const nowMs = 1_700_000_015_000;
const accounts: BetAccountEntry[] = [
betEntry(0, makeBetBuffer({ status: 0, expireTime: 1_700_000_015 })), // ровно now
];
const out = dueForSettle(accounts, nowMs);
check(out.length === 1 && out[0].id === 0, "dueForSettle уважает границу expire: ровно на nowMs включён (<=)");
}
// Тест: betOpenEvent форма.
{
const bettor = Keypair.generate().publicKey;
const bet = decodeBet(
makeBetBuffer({
amount: 100_000_000,
expireTime: 1_789_299_794,
entryPrice: 10_077_000_000,
side: 0, // UP
status: 0,
bettor,
})
);
const ev = betOpenEvent(bet, 14) as Record<string, unknown>;
check(ev.type === "bet_open", "betOpenEvent форма: type");
check(ev.id === "14", "betOpenEvent форма: id (строка)");
check(ev.side === "UP", 'betOpenEvent форма: side="UP"');
check(ev.amountUnits === "100000000", "betOpenEvent форма: amountUnits (строка)");
check(typeof ev.entryPrice === "number" && ev.entryPrice === 10_077_000_000, "betOpenEvent форма: entryPrice (число)");
check(typeof ev.expireTime === "number" && ev.expireTime === 1_789_299_794, "betOpenEvent форма: expireTime (число, секунды)");
check(typeof ev.bettor === "string" && (ev.bettor as string) === bettor.toBase58(), "betOpenEvent форма: bettor (base58)");
}
// Тест: betClosedEvent форма (payout_done).
{
const bettor = Keypair.generate().publicKey;
const bet = decodeBet(makeBetBuffer({ status: STATUS_PAYOUT_DONE, bettor }));
const ev = betClosedEvent(bet, 14, 9_977_000_000, STATUS_PAYOUT_DONE) as Record<string, unknown>;
check(ev.type === "bet_closed", "betClosedEvent форма: type");
check(ev.id === "14", "betClosedEvent форма: id");
check(typeof ev.exitPrice === "number" && ev.exitPrice === 9_977_000_000, "betClosedEvent форма: exitPrice (число)");
check(ev.status === 1, "betClosedEvent форма: status (число)");
check(ev.statusName === "payout_done", 'betClosedEvent форма: statusName==="payout_done"');
}
// Тест: betClosedEvent house_won.
{
const bettor = Keypair.generate().publicKey;
const bet = decodeBet(makeBetBuffer({ status: STATUS_HOUSE_WON, bettor }));
const ev = betClosedEvent(bet, 7, 9_500_000_000, STATUS_HOUSE_WON) as Record<string, unknown>;
check(ev.statusName === "house_won", 'betClosedEvent house_won: statusName==="house_won"');
}
// Санити-проверка: STATUS_OPEN == 0 для совместимости с on-chain контрактом.
check(STATUS_OPEN === 0, "settle: STATUS_OPEN === 0 (on-chain OPEN)");
if (failed > 0) {
console.error(`\ncheck: ${failed} проверок ПРОВАЛЕНО`);
process.exit(1);
+4
View File
@@ -50,6 +50,8 @@ export interface Config {
priceHistoryPath: string;
/** Сколько миллисекунд хранить историю цены (ретеншен). Дефолт — 24 часа. */
priceHistoryRetentionMs: number;
/** Как часто планировщик авто-закрытия просроченных ставок просыпается (мс). */
settleIntervalMs: number;
}
function env(name: string, fallback: string): string {
@@ -92,6 +94,8 @@ export function loadConfig(): Config {
// Стартовый баланс фантиков, который локальный релейер mintTo'ит беттору
// (аналог SOL-airdrop). По умолчанию равен maxAmount — хватает на максимум ставки.
startBalance: envNum("SMART_START_BALANCE", 10_000_000_000),
// Интервал планировщика авто-закрытия просроченных ставок (мс). Дефолт 2с.
settleIntervalMs: envNum("SMART_SETTLE_INTERVAL_MS", 2000),
};
}
+154 -89
View File
@@ -3,8 +3,7 @@
*
* Шлюз — custodial: keypair'ы бетторов хранит у себя (keystore-файл
* `.gateway-bettors.json` рядом с relayer/, в .gitignore) и подписывает
* create_bet от их имени. Для /close ключи не нужны — подписывает
* relayer (admin). Все сетевые операции — переиспользованные функции
* create_bet от их имени. Все сетевые операции — переиспользованные функции
* из ops.ts (createBet / closeBet / readGlobal / fundBettor / ...).
*
* Эндпоинты (CORS: Access-Control-Allow-Origin: *):
@@ -15,15 +14,27 @@
* POST /bet {address, side, amountUnits}
* — address: обязателен (получен из /wallet и пополнен через /faucet).
* Ставка НИЧЕГО не наливает — только проверяет, что фантиков
* и SOL хватает, и шлёт create_bet. Ответ: {tx, id, bettor, entryPrice}
* POST /close {id} — close_bet по текущей цене (подпись relayer)
* и SOL хватает, И что у адреса ещё нет открытой ставки
* (запрос идёт в цепочку, не в память процесса — иначе после
* рестарта пода проверка теряется). Шлёт create_bet и
* сразу же публикует WS-событие `bet_open`.
* Ответ: {tx, id, bettor, entryPrice}.
* Закрытие ставки происходит автоматически по экспирации —
* см. планировщик в `src/settle.ts`.
* GET /balance/{address} — {sol, solLamports, tokenUnits, tokenHuman}
* GET /faucet/{address} — налив SOL + фантиков на ПРОИЗВОЛЬНЫЙ адрес (Phantom и т.п.).
* Доступен ТОЛЬКО в test-режиме: SMART_MODE=test И не публичный
* кластер (проверка по genesis hash: mainnet/devnet/testnet запрещены).
* Иначе 403. Ответ: {address, sol, solLamports, tokenUnits, tokenHuman}
* GET /ws (WebSocket) — push цены SOL из NATS; на подключение шлёт
* снапшот (если цена уже есть), дальше — каждый тик.
* GET /ws (WebSocket) — push цены SOL из NATS + WS-события жизненного цикла
* ставки (`bet_open` / `bet_closed`). На подключение шлёт
* снапшот цены (если уже есть), дальше — каждый тик/событие.
*
* Авто-закрытие: см. `startScheduler` в `src/settle.ts` — каждые
* `SMART_SETTLE_INTERVAL_MS` мс (дефолт 2000) обходит Bet-аккаунты
* программы, для просроченных открытых подтягивает цену из локальной
* истории (`ctx.history.priceAtOrBefore(expire_time*1000)`) и зовёт
* `closeBet`. При успешной финализации шлёт `bet_closed` всем WS-клиентам.
*
* Пути нормализуются: ingress пробрасывает путь как есть, поэтому
* /api/... и /... должны работать одинаково (strip /api).
@@ -32,17 +43,17 @@
* 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).
* PRICE_WAIT_TIMEOUT_MS (default 15000),
* SMART_SETTLE_INTERVAL_MS (default 2000).
*/
import http from "http";
import * as fs from "fs";
import * as path from "path";
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 { loadConfig, loadKeypair } from "./config";
import {
closeBet,
createBet,
@@ -62,8 +73,19 @@ import {
} from "./natsPrice";
import { currentPriceHuman, getPrice } from "./priceSource";
import { faucetDenyReason } from "./faucetGate";
import {
decodeBet,
loadAllBetsForProgram,
SchedulerCtx,
SIDE_UP,
SIDE_DOWN,
settleOnce,
betOpenEvent,
findOpenBetForBettorViaChain,
} from "./settle";
import type { SmartUpdown } from "../idl/smart_updown";
import BN from "bn.js";
const STATUS_OPEN = 0;
const STATUS_PAYOUT_DONE = 1;
@@ -150,11 +172,13 @@ interface Ctx {
program: Program<SmartUpdown>;
decimals: number;
keystore: BettorKeystore;
/** Локальная история цены (SQLite). Используется в /close для честной
* цены на момент `expire_time` ставки, а не текущей. */
/** Локальная история цены (SQLite). Используется планировщиком авто-закрытия
* для честной цены на момент `expire_time` ставки, а не текущей. */
history: PriceHistory;
/** Кэш getGenesisHash() — RPC дёргаем один раз при первом обращении к /faucet. */
genesisHash?: string;
/** Рассылка WS-сообщения всем подключённым клиентам (price/bet_open/bet_closed). */
broadcast: (msg: Record<string, unknown>) => void;
}
interface HttpJson {
@@ -306,12 +330,57 @@ async function handleBet(ctx: Ctx, body: Record<string, unknown>): Promise<HttpJ
};
}
// «У адреса уже есть открытая ставка» — читаем цепочку, не память процесса
// (при рестарте пода она теряется). Идём по PDA `betPda(0..counter-1)` и
// фильтруем по status === 0 И bettor === address. Возвращаем 409 + id первой.
//
// ⚠️ Ловушка, проверенная живьём: фильтровать через
// `memcmp: { offset: 56, bytes: ... }` нельзя — `bytes` это base58
// **32-байтового** значения, а поле `status` занимает 8 байт, и сравнение
// молча промахивается. Поэтому фильтруем ТОЛЬКО по bettor (если бы
// memcmp был нужен), статус — в коде. Здесь мы идём по PDA напрямую,
// так что вообще не используем memcmp.
const openId = await findOpenBetForBettorViaChain(
ctx.program,
ctx.program.programId,
g.counter.toNumber(),
bettor
);
if (openId !== null) {
return {
status: 409,
body: {
error: "you already have an open bet: close it before placing a new one",
openBetId: openId.toString(),
},
};
}
const entryPrice = await getPrice();
const { id, tx } = await createBet(ctx.program, ctx.cfg.tokenMint, kp, {
side,
amount: amountUnits,
entryPrice,
});
// WS-событие об открытии — после успешной транзакции. Достаём реальные
// значения из цепочки (entry_price/expire_time/bettor), чтобы UI не
// показывал то, чего не подтвердил контракт.
try {
const betPdaAddr = betPda(ctx.program.programId, id);
const info = await ctx.connection.getAccountInfo(betPdaAddr);
if (info) {
const decoded = decodeBet(info.data as Buffer);
ctx.broadcast(betOpenEvent(decoded, id.toNumber()));
} else {
console.error(`gateway: bet #${id.toString()} account not found after create tx=${tx}`);
}
} catch (e) {
console.error(
`gateway: bet_open broadcast failed for #${id.toString()}: ${e instanceof Error ? e.message : String(e)}`
);
}
return {
status: 200,
body: {
@@ -334,69 +403,6 @@ function withBetLock<T>(fn: () => Promise<T>): Promise<T> {
return run;
}
async function handleClose(ctx: Ctx, body: Record<string, unknown>): Promise<HttpJson> {
let id: BN;
try {
id = new BN(String(body.id));
} catch {
return { status: 400, body: { error: "id must be an integer" } };
}
const programId = ctx.program.programId;
const bet = await ctx.program.account.bet.fetchNullable(betPda(programId, id));
if (!bet) {
return { status: 404, body: { error: `bet #${id.toString()} not found` } };
}
// Цена закрытия = последняя цена в локальной истории НЕ ПОЗЖЕ
// expire_time ставки. Не «текущая», не «ближайшая» — иначе игрок
// ждёт выгодного момента и кликает тогда. Момент клика больше
// ничего не значит: цена была зафиксирована историей раньше.
const currentStatus = bet.status.toNumber();
if (currentStatus !== STATUS_OPEN) {
return { status: 400, body: { error: "bet is already closed" } };
}
const expireTsMs = Number(bet.expireTime) * 1000;
if (Date.now() < expireTsMs) {
const waitSec = Math.ceil((expireTsMs - Date.now()) / 1000);
return {
status: 400,
body: { error: `bet not expired yet: wait ${waitSec}s` },
};
}
const exitPrice = ctx.history.priceAtOrBefore(expireTsMs);
if (exitPrice === null) {
return {
status: 503,
body: {
error:
"price history has no data for the expiry moment; cannot settle safely",
},
};
}
const tx = await closeBet(ctx.program, id, exitPrice, ctx.admin, bet.bettor, ctx.cfg.tokenMint);
const after = await ctx.program.account.bet.fetchNullable(betPda(programId, id));
const status = after ? after.status.toNumber() : STATUS_OPEN;
const bodyOut: Record<string, unknown> = {
tx,
id: id.toString(),
exitPrice,
status,
statusName:
status === STATUS_PAYOUT_DONE
? "payout_done"
: status === STATUS_HOUSE_WON
? "house_won"
: "open",
bettor: bet.bettor.toBase58(),
};
if (status === STATUS_OPEN) {
bodyOut.note =
"exit==entry: ставка осталась open (контракт не финализирует при равной цене) — повторите /close позже";
}
return { status: 200, body: bodyOut };
}
/** Балансы адреса (SOL + фантики в ATA по нашему mint) — переиспользуется
* и /balance, и /faucet (реальные значения после налива). */
async function readBalances(
@@ -542,19 +548,32 @@ interface WsClient {
alive: boolean;
}
/** Обобщённая рассылка: сериализует один раз, шлёт всем открытым клиентам. */
function broadcastMsg(wsClients: Set<WsClient>, msg: Record<string, unknown>): void {
const json = JSON.stringify(msg);
for (const c of wsClients) {
if (c.ws.readyState === WebSocket.OPEN) {
c.ws.send(json);
}
}
}
/** Отправка ОДНОМУ клиенту (для снапшота при подключении). */
function sendMsg(client: WsClient, msg: Record<string, unknown>): void {
if (client.ws.readyState === WebSocket.OPEN) {
client.ws.send(JSON.stringify(msg));
}
}
/** Подписка на цену → рассылка price-сообщений через общий broadcastMsg. */
function broadcastPrice(wsClients: Set<WsClient>): (p: { price: string; ts: number }) => void {
return (p) => {
const msg = JSON.stringify({
broadcastMsg(wsClients, {
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);
}
}
};
}
@@ -597,7 +616,17 @@ async function main(): Promise<void> {
}
}, 10 * 60 * 1000);
const ctx: Ctx = { cfg, admin, connection, program, decimals, keystore, history };
const wsClients = new Set<WsClient>();
const ctx: Ctx = {
cfg,
admin,
connection,
program,
decimals,
keystore,
history,
broadcast: (msg) => broadcastMsg(wsClients, msg),
};
// Прелоад как в index.ts: ваулт-фантики на localnet (если включён airdrop).
// Не блокируем старт сервера: подтверждение tx может вешаться на 30с,
@@ -665,9 +694,6 @@ async function main(): Promise<void> {
} else if (m === "POST" && p === "/bet") {
const body = await readJsonBody(req);
result = await withBetLock(() => handleBet(ctx, body));
} else if (m === "POST" && p === "/close") {
const body = await readJsonBody(req);
result = await handleClose(ctx, body);
} else if (m === "GET" && p.startsWith("/balance/")) {
const address = decodeURIComponent(p.slice("/balance/".length));
result = await handleBalance(ctx, address);
@@ -700,7 +726,6 @@ 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 };
@@ -717,13 +742,12 @@ async function main(): Promise<void> {
// Снапшот сразу, если цена уже пришла из NATS.
const snap = currentPrice();
if (snap !== null) {
const msg = JSON.stringify({
sendMsg(client, {
type: "price",
symbol: "solusdt",
price: snap.price,
ts: snap.ts,
});
ws.send(msg);
}
});
@@ -781,13 +805,54 @@ async function main(): Promise<void> {
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, GET /price, POST /wallet, POST /bet, POST /close, GET /balance/{address}, GET /faucet/{address}, WS /ws`
`gateway: endpoints: GET /state, GET /price, POST /wallet, POST /bet, GET /balance/{address}, GET /faucet/{address}, WS /ws (price + bet_open + bet_closed)`
);
console.log(
`gateway: settle scheduler: interval=${cfg.settleIntervalMs}ms (close_bet auto on expire)`
);
});
// Корректная остановка heartbeat при завершении.
// Планировщик авто-закрытия: каждые `cfg.settleIntervalMs` обходит Bet-аккаунты
// и для просроченных открытых шлёт close_bet + bet_closed WS-событие.
// Весь тик целиком внутри withBetLock — чтобы не пересекаться с /bet
// (две параллельные ставки на один counter дают BetIdMismatch).
const schedulerCtx: SchedulerCtx = {
priceAtOrBefore: (expireMs) => history.priceAtOrBefore(expireMs),
closeBet: async (id, exitPrice, bettor) =>
closeBet(program, new BN(id), exitPrice, admin, bettor, cfg.tokenMint),
broadcast: (msg) => broadcastMsg(wsClients, msg),
fetchBet: async (id) => {
const info = await connection.getAccountInfo(betPda(program.programId, new BN(id)));
if (!info) return null;
try {
return decodeBet(info.data as Buffer);
} catch {
return null;
}
},
loadAllBets: async () => {
const g = await readGlobal(program, program.programId, cfg.tokenMint);
const counter = g ? g.counter.toNumber() : 0;
return loadAllBetsForProgram(program, program.programId, counter);
},
nowMs: () => Date.now(),
};
// Карта последних состояний живёт здесь и передаётся в каждый тик —
// иначе settleOnce считает, что состояния не видел, и логирует всё подряд
// (спам одинаковыми строками каждые SMART_SETTLE_INTERVAL_MS).
const lastState = new Map<number, string>();
const schedulerHandle = setInterval(() => {
withBetLock(() => settleOnce({ ...schedulerCtx, lastState })).catch((e) =>
console.error(
`settle: tick crashed: ${e instanceof Error ? e.message : String(e)}`
)
);
}, cfg.settleIntervalMs);
// Корректная остановка heartbeat и планировщика при завершении.
const shutdown = () => {
clearInterval(heartbeat);
clearInterval(schedulerHandle);
wss.close();
server.close();
process.exit(0);
+505
View File
@@ -0,0 +1,505 @@
/**
* Планировщик авто-закрытия просроченных ставок + WS-события по жизненному
* циклу ставки (bet_open / bet_closed).
*
* ЧИСТЫЕ ФУНКЦИИ (decodeBet / openBetOf / dueForSettle / betOpenEvent /
* betClosedEvent) экспортированы — на них написаны оффлайн-тесты в check.ts.
* Сам планировщик (settleOnce + startScheduler) держится отдельно и
* зависит от рантайма (connection, program, history, broadcast).
*
* Контракт по Bet (8-байтовый дискриминатор уже учтён в смещениях):
*
* | offset | поле | тип |
* |--------|-------------|---------|
* | 8 | amount | u64 LE |
* | 16 | bet_time | i64 LE |
* | 24 | expire_time | i64 LE |
* | 32 | entry_price | u64 LE |
* | 40 | exit_price | u64 LE |
* | 48 | side | u64 LE |
* | 56 | status | u64 LE |
* | 64 | bettor | Pubkey |
*
* Стороны: 0 = UP, 1 = DOWN. Статусы: 0 = OPEN, 1 = PAYOUT_DONE, 2 = HOUSE_WON.
*/
import { PublicKey } from "@solana/web3.js";
import BN from "bn.js";
import type { Program } from "@anchor-lang/core";
import { betPda } from "./pda";
import type { SmartUpdown } from "../idl/smart_updown";
/** Размер Bet-аккаунта в байтах (8 дискриминатор + 32 amount/time + 32 pubkey + ...). */
export const BET_ACCOUNT_SIZE = 96;
/**
* Значения, которые попадают в `SchedulerCtx.lastState` (Map<number,string>):
* последнее залогированное состояние ставки, чтобы повторно не спамить в лог.
*/
const STATE_NO_PRICE = "noPrice";
const STATE_RETRY = "retry";
/** Sentinel-ключ в `lastState` под индексом -1 — «предупреждение unmapped уже печаталось». */
const SENTINEL_UNMAPPED_LOGGED = -1;
export const STATUS_OPEN = 0;
export const STATUS_PAYOUT_DONE = 1;
export const STATUS_HOUSE_WON = 2;
export const SIDE_UP = 0;
export const SIDE_DOWN = 1;
export interface DecodedBet {
amount: number;
betTime: number;
expireTime: number;
entryPrice: number;
exitPrice: number;
side: number;
status: number;
bettor: PublicKey;
}
/** Запись аккаунта с уже сопоставленным id — форма, удобная для планировщика. */
export interface BetAccountEntry {
pubkey: PublicKey;
account: { data: Buffer };
id: number;
}
/** Читает u64 LE по offset в Buffer'е (offset — от 0, включая дискриминатор). */
function readU64LE(buf: Buffer, offset: number): number {
const v = buf.readBigUInt64LE(offset);
// Number-конверсия безопасна для наших величин (amount ≤ 1e10, prices ≤ ~1e11).
return Number(v);
}
/** Читает i64 LE по offset. */
function readI64LE(buf: Buffer, offset: number): number {
const v = buf.readBigInt64LE(offset);
return Number(v);
}
/** Приводит Uint8Array | Buffer к Buffer (для единого readU64LE). */
function toBuffer(data: Buffer | Uint8Array): Buffer {
return Buffer.isBuffer(data) ? data : Buffer.from(data);
}
/**
* Разбирает сырые байты Bet-аккаунта (включая 8-байтовый дискриминатор).
* Бросает, если длина < 96.
*/
export function decodeBet(data: Buffer | Uint8Array): DecodedBet {
const buf = toBuffer(data);
if (buf.length < BET_ACCOUNT_SIZE) {
throw new Error(
`decodeBet: account data too short (${buf.length} < ${BET_ACCOUNT_SIZE})`
);
}
const bettorBytes = buf.subarray(64, 96);
return {
amount: readU64LE(buf, 8),
betTime: readI64LE(buf, 16),
expireTime: readI64LE(buf, 24),
entryPrice: readU64LE(buf, 32),
exitPrice: readU64LE(buf, 40),
side: readU64LE(buf, 48),
status: readU64LE(buf, 56),
bettor: new PublicKey(bettorBytes),
};
}
/**
* Ищет открытую (status === OPEN) ставку по адресу беттора в списке аккаунтов.
* Возвращает `{ id }` первой попавшейся или `null`.
*
* `accounts` ожидается с уже-сопоставленным id (см. BetAccountEntry) — иначе
* из pubkey обратно id не вытащить (PDA — хэш). Списком рулит вызывающий код:
* планировщик получает id из Global.counter × betPda(0..counter-1), а в /bet
* достаточно знать Global.counter и проверять каждый PDA.
*/
export function openBetOf(
address: string | PublicKey,
accounts: ReadonlyArray<BetAccountEntry>
): { id: number } | null {
const target = typeof address === "string" ? address : address.toBase58();
for (const entry of accounts) {
const bet = decodeBet(entry.account.data);
if (bet.status !== STATUS_OPEN) continue;
if (bet.bettor.toBase58() !== target) continue;
return { id: entry.id };
}
return null;
}
/**
* Возвращает id'ы открытых ставок, у которых `expire_time * 1000 <= nowMs`
* (включительно). Сортирует по возрастанию id — детерминированный порядок
* важен для логов и для воспроизводимости тестов.
*/
export function dueForSettle(
accounts: ReadonlyArray<BetAccountEntry>,
nowMs: number
): Array<{ id: number }> {
const out: Array<{ id: number }> = [];
for (const entry of accounts) {
let bet: DecodedBet;
try {
bet = decodeBet(entry.account.data);
} catch {
continue;
}
if (bet.status !== STATUS_OPEN) continue;
const expireMs = bet.expireTime * 1000;
if (expireMs <= nowMs) {
out.push({ id: entry.id });
}
}
out.sort((a, b) => a.id - b.id);
return out;
}
/** Объект WS-сообщения об открытии ставки. */
export function betOpenEvent(bet: DecodedBet, id: number): Record<string, unknown> {
return {
type: "bet_open",
id: id.toString(),
bettor: bet.bettor.toBase58(),
side: bet.side === SIDE_UP ? "UP" : "DOWN",
amountUnits: bet.amount.toString(),
entryPrice: bet.entryPrice,
expireTime: bet.expireTime,
};
}
function statusName(status: number): string {
if (status === STATUS_PAYOUT_DONE) return "payout_done";
if (status === STATUS_HOUSE_WON) return "house_won";
return "open";
}
/** Объект WS-сообщения о закрытии ставки (status 1 или 2). */
export function betClosedEvent(
bet: DecodedBet,
id: number,
exitPrice: number,
status: number
): Record<string, unknown> {
return {
type: "bet_closed",
id: id.toString(),
bettor: bet.bettor.toBase58(),
exitPrice,
status,
statusName: statusName(status),
};
}
/** Контекст для одного тика планировщика. Все побочные эффекты — через колбэки. */
export interface SchedulerCtx {
/**
* Цена на момент `expireMs` (или null, если в истории такой нет).
* Контракт: priceAtOrBefore(bet.expireTime * 1000).
*/
priceAtOrBefore: (expireMs: number) => number | null;
/** Закрытие ставки (relayer — admin, беттор — из Bet). */
closeBet: (id: number, exitPrice: number, bettor: PublicKey) => Promise<string>;
/** Рассылка WS-сообщения всем подключённым клиентам. */
broadcast: (msg: Record<string, unknown>) => void;
/** Перечитать аккаунт по id. */
fetchBet: (id: number) => Promise<DecodedBet | null>;
/** Загрузить все Bet-аккаунты программы (для dueForSettle) + число
* неcопоставленных (id ≥ counter — гонка с созданием). */
loadAllBets: () => Promise<{ entries: BetAccountEntry[]; unmapped: number }>;
/** Текущее время в мс (для тестов с детерминированной "сейчас"). */
nowMs: () => number;
/**
* Карта id → последнее залогированное состояние ставки
* (`"noPrice"` / `"retry"`). Если задана — `settleOnce` логирует только
* смену состояния, а не каждый тик; используется также как одноразовый
* флаг для предупреждения `unmapped` (под ключом -1). `startScheduler`
* создаёт карту сам; оффлайн-тесты вызывают `settleOnce` без неё —
* тогда поведение прежнее (всё логируется каждый тик).
*/
lastState?: Map<number, string>;
}
/**
* Один проход планировщика. Возвращает сводку для логов:
* { loaded, due, settled, noPrice, retry, errors }
*
* - loaded — сколько Bet-аккаунтов загружено;
* - due — сколько из них подлежат закрытию (status=0 && expired);
* - settled — успешно финализированы (статус сменился на 1 или 2);
* - noPrice — нет цены в истории на момент expire_time (пропуск);
* - retry — close прошёл, но статус остался OPEN (exit == entry, повтор);
* - errors — исключения внутри (не считая noPrice/retry).
*/
export async function settleOnce(ctx: SchedulerCtx): Promise<{
loaded: number;
due: number;
settled: number;
noPrice: number;
retry: number;
errors: number;
}> {
const stats = { loaded: 0, due: 0, settled: 0, noPrice: 0, retry: 0, errors: 0 };
/** Карта последних состояний — если задана, подавляет повторный спам
* и включает одноразовую печать «unmapped» (см. SchedulerCtx.lastState). */
const lastState = ctx.lastState;
let entries: BetAccountEntry[];
let unmapped = 0;
try {
const loaded = await ctx.loadAllBets();
entries = loaded.entries;
unmapped = loaded.unmapped;
} catch (e) {
stats.errors++;
console.error(
`settle: loadAllBets failed: ${e instanceof Error ? e.message : String(e)}`
);
return stats;
}
stats.loaded = entries.length;
if (unmapped > 0) {
if (lastState) {
// С lastState — печатаем ОДИН раз за всё время работы планировщика,
// чтобы не забивать лог одной и той же строкой каждый тик.
if (!lastState.has(SENTINEL_UNMAPPED_LOGGED)) {
lastState.set(SENTINEL_UNMAPPED_LOGGED, "1");
console.log(`settle: ${unmapped} bet account(s) without mapped id (skipped)`);
}
} else {
// Без lastState — прежнее поведение (каждый тик).
console.log(`settle: ${unmapped} bet account(s) without mapped id (skipped)`);
}
}
const due = dueForSettle(entries, ctx.nowMs());
stats.due = due.length;
// Чистим lastState от ставок, которые больше не due (закрылись или
// исчезли) — иначе карта растёт бесконечно.
if (lastState) {
const dueIds = new Set(due.map((d) => d.id));
for (const id of lastState.keys()) {
if (id < 0) continue; // sentinel, не реальный id
if (!dueIds.has(id)) lastState.delete(id);
}
}
/** За тик был хотя бы один новый лог (смена состояния или успешное закрытие). */
let emittedNew = false;
for (const item of due) {
const entry = entries.find((e) => e.id === item.id);
if (!entry) continue;
let bet: DecodedBet;
try {
bet = decodeBet(entry.account.data);
} catch {
stats.errors++;
continue;
}
const expireMs = bet.expireTime * 1000;
const exitPrice = ctx.priceAtOrBefore(expireMs);
if (exitPrice === null) {
stats.noPrice++;
if (lastState) {
// Логируем только при СМЕНЕ состояния (или если записи ещё нет).
if (lastState.get(item.id) !== STATE_NO_PRICE) {
lastState.set(item.id, STATE_NO_PRICE);
console.log(
`settle: skip #${item.id} — no price history for expireTime=${bet.expireTime}`
);
emittedNew = true;
}
} else {
console.log(
`settle: skip #${item.id} — no price history for expireTime=${bet.expireTime}`
);
}
continue;
}
let txSig: string;
try {
txSig = await ctx.closeBet(item.id, exitPrice, bet.bettor);
} catch (e) {
stats.errors++;
console.error(
`settle: close #${item.id} failed: ${e instanceof Error ? e.message : String(e)}`
);
continue;
}
const after = await ctx.fetchBet(item.id).catch(() => null);
if (!after) {
stats.errors++;
if (lastState) lastState.delete(item.id);
console.error(
`settle: #${item.id} account disappeared after close tx=${txSig}`
);
continue;
}
if (after.status === STATUS_OPEN) {
stats.retry++;
if (lastState) {
if (lastState.get(item.id) !== STATE_RETRY) {
lastState.set(item.id, STATE_RETRY);
console.log(
`settle: #${item.id} still open (exit==entry at ${exitPrice}), will retry`
);
emittedNew = true;
}
} else {
console.log(
`settle: #${item.id} still open (exit==entry at ${exitPrice}), will retry`
);
}
continue;
}
if (after.status === STATUS_PAYOUT_DONE || after.status === STATUS_HOUSE_WON) {
stats.settled++;
if (lastState) lastState.delete(item.id);
ctx.broadcast(betClosedEvent(after, item.id, exitPrice, after.status));
console.log(
`settle: closed #${item.id} exit=${exitPrice} status=${after.status} tx=${txSig}`
);
emittedNew = true;
} else {
stats.errors++;
console.error(
`settle: #${item.id} unexpected status ${after.status} after close`
);
}
}
// Сводная строка тика — только если за тик было новое событие
// (смена состояния или успешное закрытие). Без lastState — прежнее
// поведение (печатаем при любом due.length > 0).
if (lastState) {
if (emittedNew) {
console.log(`settle: tick ${due.length} due (loaded=${entries.length})`);
}
} else if (due.length > 0) {
console.log(`settle: tick ${due.length} due (loaded=${entries.length})`);
}
return stats;
}
/**
* Запускает периодический планировщик. Возвращает handle для остановки.
* Каждый тик целиком в `try/catch` — ошибка внутри НЕ роняет процесс.
*/
export function startScheduler(
ctx: SchedulerCtx,
intervalMs: number
): { stop: () => void } {
if (!Number.isFinite(intervalMs) || intervalMs <= 0) {
throw new Error(`startScheduler: invalid intervalMs ${intervalMs}`);
}
// Карта последних состояний живёт в замыкании планировщика и
// передаётся в каждый тик — см. settleOnce. Без неё поведение прежнее
// (логируется всё каждый тик).
const lastState = new Map<number, string>();
const handle = setInterval(() => {
settleOnce({ ...ctx, lastState }).catch((e) => {
console.error(
`settle: tick crashed: ${e instanceof Error ? e.message : String(e)}`
);
});
}, intervalMs);
return {
stop: () => clearInterval(handle),
};
}
/**
* Загружает все Bet-аккаунты программы с уже-сопоставленным id (для
* SchedulerCtx.loadAllBets). Делает ОДИН `getProgramAccounts(programId)` и
* сопоставляет pubkey → id локально через заранее посчитанную карту PDA.
*
* Обоснование выбора схемы: вместо N отдельных `getAccountInfo` по каждому
* PDA (линейно по сети и по времени, ~80 мс на 15 ставках) — один
* `getProgramAccounts`, который отдаёт всё разом (~9 мс на тех же 15
* ставках). Карта PDA → id считается локально, без сети. Аккаунты без
* пары в карте (id ≥ counter — гонка между чтением counter и этим
* запросом) пропускаются и считаются отдельно, чтобы планировщик мог
* напечатать об этом одну строку в логе.
*/
export async function loadAllBetsForProgram(
program: Program<SmartUpdown>,
programId: PublicKey,
counter: number
): Promise<{ entries: BetAccountEntry[]; unmapped: number }> {
const pdaToId = new Map<string, number>();
for (let i = 0; i < counter; i++) {
pdaToId.set(betPda(programId, new BN(i)).toBase58(), i);
}
const accounts = await program.provider.connection.getProgramAccounts(programId, {
commitment: "confirmed",
});
const entries: BetAccountEntry[] = [];
let unmapped = 0;
for (const { pubkey, account } of accounts) {
// Фильтр по размеру ДО сопоставления с картой id: getProgramAccounts
// отдаёт ВСЕ аккаунты программы — включая Global PDA (120 байт),
// который не Bet. Такие аккаунты НЕ считаются unmapped (это не гонка
// id ≥ counter, а просто другой тип аккаунта контракта).
if (account.data.length !== BET_ACCOUNT_SIZE) continue;
const id = pdaToId.get(pubkey.toBase58());
if (id === undefined) {
unmapped++;
continue;
}
entries.push({ pubkey, account: { data: account.data as Buffer }, id });
}
return { entries, unmapped };
}
/**
* Загрузка ВСЕХ Bet-аккаунтов через цепочку — нужно для проверки
* «у адреса уже есть открытая ставка» в /bet. ТЗ требует читать
* из цепочки (а не из памяти процесса), потому что при рестарте пода
* процесс теряет состояние. Серверный фильтр Solana по полю `bettor`
* (offset 64, 32 байта) — один `getProgramAccounts`, статус проверяем
* в коде (`decodeBet(...).status === STATUS_OPEN`).
*
* ⚠️ Ловушка, проверенная живьём: фильтровать через
* `memcmp: { offset: 56, bytes: ... }` НЕЛЬЗЯ — `bytes` это base58
* **32-байтового** значения, а поле `status` занимает 8 байт, и сравнение
* молча промахивается. Поэтому фильтруем ТОЛЬКО по bettor (offset 64,
* 32 байта — как раз совпадает с длиной Pubkey), статус — в коде.
*
* Возвращает id первой открытой ставки данного беттора или null.
*/
export async function findOpenBetForBettorViaChain(
program: Program<SmartUpdown>,
programId: PublicKey,
counter: number,
bettor: string | PublicKey
): Promise<number | null> {
const target = typeof bettor === "string" ? bettor : bettor.toBase58();
const accounts = await program.provider.connection.getProgramAccounts(programId, {
commitment: "confirmed",
filters: [{ memcmp: { offset: 64, bytes: target } }],
});
const pdaToId = new Map<string, number>();
for (let i = 0; i < counter; i++) {
pdaToId.set(betPda(programId, new BN(i)).toBase58(), i);
}
for (const { pubkey, account } of accounts) {
let decoded: DecodedBet;
try {
decoded = decodeBet(account.data as Buffer);
} catch {
continue;
}
if (decoded.status !== STATUS_OPEN) continue;
const id = pdaToId.get(pubkey.toBase58());
if (id === undefined) continue;
return id;
}
return null;
}