diff --git a/relayer/helm/templates/deployment.yaml b/relayer/helm/templates/deployment.yaml index f4945de..4be7193 100644 --- a/relayer/helm/templates/deployment.yaml +++ b/relayer/helm/templates/deployment.yaml @@ -69,6 +69,10 @@ spec: value: {{ .Values.smart.expirySeconds | quote }} - name: SMART_START_BALANCE value: {{ .Values.smart.startBalance | quote }} + - name: SMART_PRICE_HISTORY_PATH + value: {{ .Values.smart.priceHistory | quote }} + - name: SMART_PRICE_RETENTION_HOURS + value: {{ .Values.smart.priceRetentionHours | 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 2dc5eb5..34790e5 100644 --- a/relayer/helm/values.yaml +++ b/relayer/helm/values.yaml @@ -26,6 +26,8 @@ smart: multiplierBps: 13000 expirySeconds: 15 startBalance: 10000000000 + priceHistory: '/app/data/price-history.db' + priceRetentionHours: 24 nats: url: 'nats://192.168.88.93:4222' diff --git a/relayer/src/check.ts b/relayer/src/check.ts index 2dad642..98d5f50 100644 --- a/relayer/src/check.ts +++ b/relayer/src/check.ts @@ -13,11 +13,15 @@ */ import { Keypair, PublicKey } from "@solana/web3.js"; import BN from "bn.js"; +import * as fs from "fs"; +import * as os from "os"; +import * as path from "path"; import { getAssociatedTokenAddress, TOKEN_PROGRAM_ID } from "@solana/spl-token"; import idl from "../idl/smart_updown.json"; import { globalPda, betPda, vaultPda } from "./pda"; import { makeConnection, makeProgram } from "./program"; import { faucetDenyReason, parseMode } from "./faucetGate"; +import { DEFAULT_RETENTION_MS, PriceHistory } from "./priceHistory"; const ok = (msg: string) => console.log(`ok: ${msg}`); let failed = 0; @@ -164,6 +168,87 @@ async function main(): Promise { "parseMode: неизвестное значение = release" ); + // 4. PriceHistory (SQLite): пишется каждый NATS-тик цены, читается в + // /close для фиксированной цены на момент expire_time. Тесты гоняются + // на РЕАЛЬНОЙ БД в /tmp (файл удаляется до и после). + { + const dbPath = path.join( + os.tmpdir(), + `price-history-check-${process.pid}-${Date.now()}.db` + ); + fs.rmSync(dbPath, { force: true }); + const h = new PriceHistory(dbPath); + try { + // Тест 1: пишем точку, читаем её же. + h.record(1000, 111); + check( + h.priceAtOrBefore(1000) === 111, + "priceHistory: пишет точку и читает её" + ); + + // Тест 4: до самой первой точки — пусто. + check( + h.priceAtOrBefore(999) === null, + "priceHistory: до первой точки = null" + ); + + // Тесты 2/3/5: две точки 1000→111 и 5000→222 — разные режимы выборки. + h.record(5000, 222); + check( + h.priceAtOrBefore(3000) === 111, + "priceHistory: отдаёт предыдущую точку между записями" + ); + check( + h.priceAtOrBefore(4999) === 111, + "priceHistory: НЕ отдаёт цену позже момента" + ); + check( + h.priceAtOrBefore(5000) === 222, + "priceHistory: ровно на точке отдаёт эту точку" + ); + + // Тест 6: повтор той же цены (последняя записанная = 222) — отказ. + const before6 = h.count(); + const r6 = h.record(6000, 222); + check( + r6 === false && h.count() === before6, + "priceHistory: повтор той же цены не пишет строку" + ); + + // Тест 7: изменившаяся цена — пишется. + const before7 = h.count(); + const r7 = h.record(7000, 333); + check( + r7 === true && h.count() === before7 + 1, + "priceHistory: изменившаяся цена пишет строку" + ); + + // Тест 8: 3 точки (1000/5000/7000), prune(now=7000, retention=1000). + // cutoff = 7000 - 1000 = 6000 → удаляются 1000 и 5000, остаётся 7000 + // (ровно 2 удалено, ровно сколько старше). + const removed = h.prune(7000, 1000); + check( + removed === 2 && h.priceAtOrBefore(3000) === null, + "priceHistory: prune удаляет старше ретеншена" + ); + + // Тест 9: после prune свежая точка (7000 → 333) всё ещё читается. + check( + h.priceAtOrBefore(7000) === 333, + "priceHistory: prune не трогает свежие" + ); + } finally { + h.close(); + fs.rmSync(dbPath, { force: true }); + } + } + + // Тест 10: дефолтный ретеншен = 24 часа. + check( + DEFAULT_RETENTION_MS === 24 * 60 * 60 * 1000, + "priceHistory: ретеншен по умолчанию = 24 часа" + ); + if (failed > 0) { console.error(`\ncheck: ${failed} проверок ПРОВАЛЕНО`); process.exit(1); diff --git a/relayer/src/config.ts b/relayer/src/config.ts index 0e82bca..9930705 100644 --- a/relayer/src/config.ts +++ b/relayer/src/config.ts @@ -46,6 +46,10 @@ export interface Config { startBalance: number; multiplierBps: number; expirySeconds: number; + /** Путь к файлу локальной истории цены (SQLite). */ + priceHistoryPath: string; + /** Сколько миллисекунд хранить историю цены (ретеншен). Дефолт — 24 часа. */ + priceHistoryRetentionMs: number; } function env(name: string, fallback: string): string { @@ -80,6 +84,11 @@ export function loadConfig(): Config { maxAmount: envNum("SMART_MAX_AMOUNT", 10_000_000_000), // 10 000 фантиков multiplierBps: envNum("SMART_MULTIPLIER_BPS", 13_000), // 1.3x expirySeconds: envNum("SMART_EXPIRY_SECONDS", 15), + // Локальная история цены (SQLite): цена пишется на каждом тике при + // изменении, чтобы при /close отдавать цену на момент expire_time ставки, + // а не текущую (иначе игрок ждёт выгодного момента и кликает тогда). + priceHistoryPath: env("SMART_PRICE_HISTORY_PATH", "/app/data/price-history.db"), + priceHistoryRetentionMs: envNum("SMART_PRICE_RETENTION_HOURS", 24) * 60 * 60 * 1000, // Стартовый баланс фантиков, который локальный релейер mintTo'ит беттору // (аналог SOL-airdrop). По умолчанию равен maxAmount — хватает на максимум ставки. startBalance: envNum("SMART_START_BALANCE", 10_000_000_000), diff --git a/relayer/src/gateway.ts b/relayer/src/gateway.ts index d039c64..55f638d 100644 --- a/relayer/src/gateway.ts +++ b/relayer/src/gateway.ts @@ -52,8 +52,14 @@ import { readGlobal, } from "./ops"; import { betPda } from "./pda"; +import { PriceHistory } from "./priceHistory"; import { makeConnection, makeProgram } from "./program"; -import { currentPrice, onPrice, connectNats } from "./natsPrice"; +import { + currentPrice, + onPrice, + connectNats, + priceStringToInt, +} from "./natsPrice"; import { currentPriceHuman, getPrice } from "./priceSource"; import { faucetDenyReason } from "./faucetGate"; @@ -144,6 +150,9 @@ interface Ctx { program: Program; decimals: number; keystore: BettorKeystore; + /** Локальная история цены (SQLite). Используется в /close для честной + * цены на момент `expire_time` ставки, а не текущей. */ + history: PriceHistory; /** Кэш getGenesisHash() — RPC дёргаем один раз при первом обращении к /faucet. */ genesisHash?: string; } @@ -338,7 +347,32 @@ async function handleClose(ctx: Ctx, body: Record): Promise { process.env.GW_BETTOR_KEYSTORE ?? path.join(__dirname, "..", ".gateway-bettors.json") ); - const ctx: Ctx = { cfg, admin, connection, program, decimals, keystore }; + // Локальная история цены (SQLite): пишется каждый тик из NATS (только при + // изменении цены), читается в handleClose для фиксированной цены на момент + // expire_time ставки. Ретеншен чистится периодически ОДНИМ DELETE. + const history = new PriceHistory(cfg.priceHistoryPath); + try { + const removed = history.prune(Date.now(), cfg.priceHistoryRetentionMs); + if (removed > 0) { + console.log(`gateway: price history pruned ${removed} stale row(s) at startup`); + } + } catch (e) { + console.error( + `gateway: initial price-history prune failed: ${e instanceof Error ? e.message : String(e)}` + ); + } + setInterval(() => { + try { + history.prune(Date.now(), cfg.priceHistoryRetentionMs); + } catch { + // тихий фоновый prune — ронять процесс нельзя + } + }, 10 * 60 * 1000); + + const ctx: Ctx = { cfg, admin, connection, program, decimals, keystore, history }; // Прелоад как в index.ts: ваулт-фантики на localnet (если включён airdrop). // Не блокируем старт сервера: подтверждение tx может вешаться на 30с, @@ -674,6 +730,20 @@ async function main(): Promise { // Рассылка каждого нового NATS-тика по всем подключённым WS. onPrice(broadcastPrice(wsClients)); + // Параллельно пишем тик в локальную историю цены (только при изменении + // цены — `record` сам это разруливает и вернёт false на дубликат). + // Ошибку записи НЕ роняем в процесс — это фоновая телеметрия, + // шлюз без неё работает. + onPrice((p) => { + try { + history.record(p.eventTime, priceStringToInt(p.price)); + } catch (e) { + console.error( + `gateway: price-history record failed: ${e instanceof Error ? e.message : String(e)}` + ); + } + }); + // Heartbeat: пинг раз в 30с, нет pong 2 цикла — terminate. const heartbeat = setInterval(() => { for (const c of wsClients) { diff --git a/relayer/src/natsPrice.ts b/relayer/src/natsPrice.ts index b72fdf2..79af702 100644 --- a/relayer/src/natsPrice.ts +++ b/relayer/src/natsPrice.ts @@ -17,6 +17,14 @@ export interface PriceTick { price: string; /** Unix ms: берётся из поля `ts` payload'а, иначе Date.now(). */ ts: number; + /** + * Unix ms: время СДЕЛКИ у Binance (поле `eventTime` payload'а). + * Источник истины по времени: ставка в шлюзе закрывается по цене + * на момент `expire_time`, а не по текущей — поэтому нужно время + * сделки, а не локальная метка приёма. Fallback на `ts`, если + * `eventTime` отсутствует или не число. + */ + eventTime: number; } let nc: NatsConnection | null = null; @@ -87,7 +95,11 @@ export async function connectNats(url: string, subject: string): Promise { typeof obj.ts === "number" && Number.isFinite(obj.ts) ? (obj.ts as number) : Date.now(); - const tick: PriceTick = { price: priceRaw, ts: tsNum }; + const eventTime = + typeof obj.eventTime === "number" && Number.isFinite(obj.eventTime) + ? (obj.eventTime as number) + : tsNum; + const tick: PriceTick = { price: priceRaw, ts: tsNum, eventTime }; if (!firstLogged) { console.log(`natsPrice: first price ${priceRaw} (ts=${tsNum})`); firstLogged = true; diff --git a/relayer/src/priceHistory.ts b/relayer/src/priceHistory.ts new file mode 100644 index 0000000..22e23b2 --- /dev/null +++ b/relayer/src/priceHistory.ts @@ -0,0 +1,116 @@ +/** + * Локальная история цены SOL (SQLite, node:sqlite). + * + * Используется шлюзом, чтобы при `close_bet` подставлять в транзакцию + * **зафиксированную заранее** цену на момент `expire_time` ставки, а не + * текущую с биржи. Это закрывает окно «подождать выгодного момента и + * кликнуть позже»: цена фиксируется историей, момент клика больше + * ничего не значит. + * + * Схема (одна таблица, PK = `ts` = unix-миллисекунды сделки): + * + * ts INTEGER PRIMARY KEY -- время СДЕЛКИ (eventTime Binance), unix ms + * price INTEGER NOT NULL -- цена × 1e8, целое (без float) + * + * Запись идёт **только при изменении цены** (≈300 строк/сутки вместо + * 92k тиков), ретеншен по умолчанию — 24 часа, чистится периодически + * `DELETE FROM price_history WHERE ts < ?` одним запросом. + */ +import { DatabaseSync } from "node:sqlite"; + +/** Ретеншен по умолчанию: 24 часа (в миллисекундах). */ +export const DEFAULT_RETENTION_MS: number = 24 * 60 * 60 * 1000; + +export class PriceHistory { + private readonly db: DatabaseSync; + private readonly stmtLast: ReturnType; + private readonly stmtLastPrice: ReturnType; + private readonly stmtAtOrBefore: ReturnType; + private readonly stmtInsert: ReturnType; + private readonly stmtPrune: ReturnType; + private readonly stmtCount: ReturnType; + + constructor(filePath: string) { + this.db = new DatabaseSync(filePath); + // WAL: параллельные читатели не блокируют пишущего тик-цикл. + this.db.exec("PRAGMA journal_mode=WAL"); + // PK(ts) уже индекс для `WHERE ts <= ?` и `ORDER BY ts DESC LIMIT 1`, + // отдельный индекс не нужен. + this.db.exec(` + CREATE TABLE IF NOT EXISTS price_history ( + ts INTEGER PRIMARY KEY, + price INTEGER NOT NULL + ) + `); + this.stmtLast = this.db.prepare("SELECT ts FROM price_history ORDER BY ts DESC LIMIT 1"); + this.stmtLastPrice = this.db.prepare( + "SELECT price FROM price_history ORDER BY ts DESC LIMIT 1" + ); + this.stmtAtOrBefore = this.db.prepare( + "SELECT price FROM price_history WHERE ts <= ? ORDER BY ts DESC LIMIT 1" + ); + this.stmtInsert = this.db.prepare( + "INSERT OR REPLACE INTO price_history (ts, price) VALUES (?, ?)" + ); + this.stmtPrune = this.db.prepare("DELETE FROM price_history WHERE ts < ?"); + this.stmtCount = this.db.prepare("SELECT COUNT(*) AS c FROM price_history"); + } + + /** + * Записать точку. Пишем ТОЛЬКО если цена отличается от последней записанной + * — иначе возвращаем false и в БД ничего не добавляем (≈300 строк/сутки + * вместо 92k тиков). Первая запись всегда пишется. `tsMs` пишется как + * есть (PK; INSERT OR REPLACE на случай коллизии `ts`). + * + * Возвращает true, если строка реально добавлена. + */ + record(tsMs: number, priceInt: number): boolean { + if (!Number.isFinite(tsMs) || !Number.isFinite(priceInt)) return false; + const prev = this.stmtLastPrice.get() as { price: number } | undefined; + if (prev !== undefined && prev.price === priceInt) return false; + this.stmtInsert.run(tsMs, priceInt); + return true; + } + + /** + * Последняя цена, записанная НЕ ПОЗЖЕ `tsMs` (ts <= tsMs). + * Или null, если такой нет. + * + * Именно этот метод определяет честность расчёта: цена ПОЗЖЕ момента + * не должна возвращаться никогда. Поэтому `ORDER BY ts DESC` + `LIMIT 1` + * под условием `ts <= ?` — никаких «ближайших». + */ + priceAtOrBefore(tsMs: number): number | null { + if (!Number.isFinite(tsMs)) return null; + const row = this.stmtAtOrBefore.get(tsMs) as { price: number } | undefined; + return row !== undefined ? row.price : null; + } + + /** Unix-ms последней записи, или null если истории нет. */ + lastTs(): number | null { + const row = this.stmtLast.get() as { ts: number } | undefined; + return row !== undefined ? row.ts : null; + } + + /** + * Удалить записи старше `(nowMs - retentionMs)`. Одним запросом + * `DELETE FROM price_history WHERE ts < ?`. Возвращает число удалённых строк. + */ + prune(nowMs: number, retentionMs: number = DEFAULT_RETENTION_MS): number { + if (!Number.isFinite(nowMs) || !Number.isFinite(retentionMs)) return 0; + const cutoff = nowMs - retentionMs; + const res = this.stmtPrune.run(cutoff); + return Number(res.changes); + } + + /** Сколько строк в таблице (для диагностики/тестов). */ + count(): number { + const row = this.stmtCount.get() as { c: number } | undefined; + return row !== undefined ? Number(row.c) : 0; + } + + /** Закрыть БД. */ + close(): void { + this.db.close(); + } +}