diff --git a/relayer/helm/templates/deployment.yaml b/relayer/helm/templates/deployment.yaml index 4be7193..4c802ca 100644 --- a/relayer/helm/templates/deployment.yaml +++ b/relayer/helm/templates/deployment.yaml @@ -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 diff --git a/relayer/helm/values.yaml b/relayer/helm/values.yaml index 34790e5..1638315 100644 --- a/relayer/helm/values.yaml +++ b/relayer/helm/values.yaml @@ -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' diff --git a/relayer/src/check.ts b/relayer/src/check.ts index 98d5f50..c9db4dc 100644 --- a/relayer/src/check.ts +++ b/relayer/src/check.ts @@ -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 { const programId = new PublicKey((idl as any).address as string); const createIx = idlIx("create_bet"); @@ -249,6 +297,133 @@ async function main(): Promise { "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; + 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; + 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; + 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); diff --git a/relayer/src/config.ts b/relayer/src/config.ts index 9930705..5691440 100644 --- a/relayer/src/config.ts +++ b/relayer/src/config.ts @@ -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), }; } diff --git a/relayer/src/gateway.ts b/relayer/src/gateway.ts index 55f638d..167184d 100644 --- a/relayer/src/gateway.ts +++ b/relayer/src/gateway.ts @@ -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; 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) => void; } interface HttpJson { @@ -306,12 +330,57 @@ async function handleBet(ctx: Ctx, body: Record): Promise(fn: () => Promise): Promise { return run; } -async function handleClose(ctx: Ctx, body: Record): Promise { - 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 = { - 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, msg: Record): 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): void { + if (client.ws.readyState === WebSocket.OPEN) { + client.ws.send(JSON.stringify(msg)); + } +} + +/** Подписка на цену → рассылка price-сообщений через общий broadcastMsg. */ function broadcastPrice(wsClients: Set): (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 { } }, 10 * 60 * 1000); - const ctx: Ctx = { cfg, admin, connection, program, decimals, keystore, history }; + const wsClients = new Set(); + 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 { } 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 { // WS-сервер на том же порту (noServer=true), обновляем апгрейды вручную. const wss = new WebSocketServer({ noServer: true }); - const wsClients = new Set(); wss.on("connection", (ws: WebSocket) => { const client: WsClient = { ws, alive: true }; @@ -717,13 +742,12 @@ async function main(): Promise { // Снапшот сразу, если цена уже пришла из 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 { 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(); + 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); diff --git a/relayer/src/settle.ts b/relayer/src/settle.ts new file mode 100644 index 0000000..3c772eb --- /dev/null +++ b/relayer/src/settle.ts @@ -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): + * последнее залогированное состояние ставки, чтобы повторно не спамить в лог. + */ +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 +): { 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, + 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 { + 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 { + 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; + /** Рассылка WS-сообщения всем подключённым клиентам. */ + broadcast: (msg: Record) => void; + /** Перечитать аккаунт по id. */ + fetchBet: (id: number) => Promise; + /** Загрузить все Bet-аккаунты программы (для dueForSettle) + число + * неcопоставленных (id ≥ counter — гонка с созданием). */ + loadAllBets: () => Promise<{ entries: BetAccountEntry[]; unmapped: number }>; + /** Текущее время в мс (для тестов с детерминированной "сейчас"). */ + nowMs: () => number; + /** + * Карта id → последнее залогированное состояние ставки + * (`"noPrice"` / `"retry"`). Если задана — `settleOnce` логирует только + * смену состояния, а не каждый тик; используется также как одноразовый + * флаг для предупреждения `unmapped` (под ключом -1). `startScheduler` + * создаёт карту сам; оффлайн-тесты вызывают `settleOnce` без неё — + * тогда поведение прежнее (всё логируется каждый тик). + */ + lastState?: Map; +} + +/** + * Один проход планировщика. Возвращает сводку для логов: + * { 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(); + 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, + programId: PublicKey, + counter: number +): Promise<{ entries: BetAccountEntry[]; unmapped: number }> { + const pdaToId = new Map(); + 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, + programId: PublicKey, + counter: number, + bettor: string | PublicKey +): Promise { + 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(); + 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; +}