Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2218c7c1b7 | |||
| 92811c258f | |||
| 3bf7cf9c7a | |||
| be67debf51 |
@@ -49,6 +49,11 @@ abiFallback: rpc-current
|
||||
xtrim:
|
||||
enabled: true
|
||||
intervalMs: 60000
|
||||
backpressure:
|
||||
enabled: true # пауза чтения SHiP когда стрим переполнен (защита Redis RAM)
|
||||
highWater: 100000 # порог паузы — число событий в стриме
|
||||
lowWater: 50000 # порог возобновления (должен быть < highWater)
|
||||
pollMs: 200 # период опроса XLEN во время паузы, мс
|
||||
logger:
|
||||
level: info
|
||||
pretty: false
|
||||
@@ -185,6 +190,11 @@ abiFallback: rpc-current
|
||||
xtrim:
|
||||
enabled: true
|
||||
intervalMs: 60000
|
||||
backpressure:
|
||||
enabled: true # пауза чтения SHiP когда стрим переполнен (защита Redis RAM)
|
||||
highWater: 100000 # порог паузы — число событий в стриме
|
||||
lowWater: 50000 # порог возобновления (должен быть < highWater)
|
||||
pollMs: 200 # период опроса XLEN во время паузы, мс
|
||||
reconnect:
|
||||
maxAttempts: 10
|
||||
backoffSeconds: [1, 2, 5, 10, 30, 60, 120, 300, 600, 1800]
|
||||
@@ -224,7 +234,8 @@ docker run --rm --network host \
|
||||
│ • Workers │ │ │
|
||||
│ • ForkDet. │ │ Streams + │
|
||||
│ • XtrimSup. │ │ Hashes + │
|
||||
└──────────────┘ │ ZSets │
|
||||
│ • Backpr. │ │ ZSets │
|
||||
└──────────────┘ │ │
|
||||
│ │
|
||||
┌──────────────┐ │ │
|
||||
│ ParserClient │◀─XREADGR──│ │
|
||||
@@ -240,6 +251,57 @@ docker run --rm --network host \
|
||||
- [Redis key taxonomy](docs/redis-key-taxonomy.md)
|
||||
- [Disaster recovery / fork scenarios](docs/disaster-recovery.md)
|
||||
|
||||
## Backpressure: защита памяти Redis на репарсинге
|
||||
|
||||
Parser читает блоки из SHiP и пишет события в Redis-стрим. На длинном репарсинге (genesis → head, миллионы блоков) писатель легко обгоняет потребителя: события копятся в **памяти Redis** быстрее, чем `ParserClient` их вычитывает. RAM самого парсера ограничена окном SHiP (`max_messages_in_flight` + pull-based чтение с ack после обработки), а вот стрим в Redis может расти неограниченно — вплоть до OOM. `XtrimSupervisor` тут не спасает: он по дизайну не режет непрочитанные события.
|
||||
|
||||
`BackpressureGate` вводит **окно поглощения с гистерезисом**. Перед обработкой каждого блока парсер смотрит backlog = `XLEN` стрима:
|
||||
|
||||
- backlog ≥ **highWater** → чтение SHiP встаёт на **паузу** (блок не подтверждается, SHiP сам перестаёт слать после своего in-flight окна);
|
||||
- во время паузы парсер опрашивает `XLEN` каждые `pollMs` и ждёт, пока потребитель сольёт backlog ≤ **lowWater**;
|
||||
- слилось → чтение **возобновляется**.
|
||||
|
||||
Ничего не теряется — события остаются в стриме, парсер лишь перестаёт доливать. Если потребитель вообще не подключён, backlog дорастает до `highWater` и писатель замирает, удерживая Redis в пределах ~`highWater` событий.
|
||||
|
||||
**Почему не дёргается.** Разрыв high/low (по умолчанию 100k / 50k) — это антидребезг (гистерезис). Один цикл «пауза ↔ работа» перемалывает 50k событий, поэтому переключения происходят редко, а не на каждом блоке. В steady-state (потребитель у головы цепи) backlog почти нулевой → пауза не включается вообще, оверхед = один `XLEN` на блок (O(1)).
|
||||
|
||||
```
|
||||
backlog
|
||||
100k | /| /| /| <- highWater -> ПАУЗА
|
||||
| / | / | / |
|
||||
50k | / |______/ |______/ |___ <- lowWater -> ПУСК
|
||||
| / (пауза) (пауза)
|
||||
0 +--+----------------------------> время
|
||||
пишем ждём пишем ждём пишем
|
||||
```
|
||||
|
||||
Конфиг (включён по умолчанию):
|
||||
|
||||
```yaml
|
||||
backpressure:
|
||||
enabled: true # выключить можно enabled: false
|
||||
highWater: 100000 # порог паузы — число событий в стриме
|
||||
lowWater: 50000 # порог возобновления (< highWater)
|
||||
pollMs: 200 # период опроса XLEN во время паузы, мс
|
||||
```
|
||||
|
||||
Рекомендация по инфраструктуре: на Redis выставить `maxmemory` + `maxmemory-policy noeviction` — тогда при достижении лимита `XADD` вернёт ошибку (парсер её обработает), а не сработает OOM-killer.
|
||||
|
||||
## Авто-переподключение к SHiP
|
||||
|
||||
Если SHiP-нода падает или рвёт соединение, парсер **не пытается переподключиться мгновенно в цикле** — это превратило бы упавшую ноду в мишень для долбёжки без пауз. Вместо этого `ReconnectSupervisor` переподключается с **экспоненциальным backoff** (по умолчанию `1 → 2 → 5 → 15 → 60` секунд, дальше держит последнее значение):
|
||||
|
||||
- разрыв соединения → пауза по таблице `backoffSeconds`, затем новая попытка;
|
||||
- успешная попытка перечитывает позицию возобновления из Redis (`sync`-hash пишется после каждого блока) и продолжает **ровно с последнего записанного блока** — без потерь и дублей;
|
||||
- каждый обработанный блок сбрасывает счётчик попыток: долгая стабильная сессия после редкого разрыва не копит попытки;
|
||||
- `maxAttempts` неудач подряд (нода действительно мертва, прогресса нет) → парсер завершается с кодом 1, чтобы оркестратор (Docker/systemd/k8s) перезапустил процесс штатно.
|
||||
|
||||
```yaml
|
||||
reconnect:
|
||||
maxAttempts: 10
|
||||
backoffSeconds: [1, 2, 5, 10, 30, 60, 120, 300, 600, 1800]
|
||||
```
|
||||
|
||||
## Известные ограничения
|
||||
|
||||
### Схлопывание дельт create+remove внутри одного блока
|
||||
@@ -298,8 +360,8 @@ pnpm --filter @coopenomics/parser2 test:integration # нужен Docker для
|
||||
```
|
||||
|
||||
**Текущее покрытие:**
|
||||
- `@coopenomics/parser2`: 81% statements / 74% functions (205 unit-тестов)
|
||||
- `@coopenomics/coopos-ship-reader`: 81% statements / 86% functions (51 unit-тест)
|
||||
- `@coopenomics/parser2`: 229 unit-тестов + integration (xtrim/backpressure на реальном Redis)
|
||||
- `@coopenomics/coopos-ship-reader`: 54 unit-теста
|
||||
|
||||
### Бенчмарк
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@coopenomics/parser2",
|
||||
"version": "1.2.0",
|
||||
"version": "1.3.0",
|
||||
"description": "Universal EOSIO/Antelope SHiP-to-Redis blockchain indexer (parser)",
|
||||
"license": "MIT",
|
||||
"author": "Coopenomics contributors",
|
||||
|
||||
@@ -7,11 +7,15 @@ import { BlockProcessor } from './BlockProcessor.js'
|
||||
import type { XtrimSupervisorOpts } from './XtrimSupervisor.js'
|
||||
import { XtrimSupervisor } from './XtrimSupervisor.js'
|
||||
import { BackpressureGate } from './BackpressureGate.js'
|
||||
import { ReconnectSupervisor } from './ReconnectSupervisor.js'
|
||||
import { runStreamLoop } from './streamLoop.js'
|
||||
import { ForkDetector } from './ForkDetector.js'
|
||||
import { RedisKeys } from '../redis/keys.js'
|
||||
import { ChainIdMismatchError } from '../errors.js'
|
||||
import { AbiStore } from '../abi/AbiStore.js'
|
||||
import { AbiBootstrapper } from '../abi/AbiBootstrapper.js'
|
||||
import { createLogger } from '../logger.js'
|
||||
import type { Logger } from '../logger.js'
|
||||
import type { ParserEvent } from '../types.js'
|
||||
|
||||
export class Parser {
|
||||
@@ -22,6 +26,7 @@ export class Parser {
|
||||
private blockProcessor: BlockProcessor | null = null
|
||||
private xtrimSupervisor: XtrimSupervisor | null = null
|
||||
private backpressureGate: BackpressureGate | null = null
|
||||
private log: Logger | null = null
|
||||
private running = false
|
||||
private stopSignal = false
|
||||
|
||||
@@ -68,6 +73,12 @@ export class Parser {
|
||||
throw new ChainIdMismatchError(this.opts.chain.id, chainId)
|
||||
}
|
||||
|
||||
this.log = createLogger({
|
||||
...(this.opts.logger?.level !== undefined ? { level: this.opts.logger.level } : {}),
|
||||
...(this.opts.logger?.pretty !== undefined ? { pretty: this.opts.logger.pretty } : {}),
|
||||
chain_id: chainId,
|
||||
}).child({ component: 'Parser' })
|
||||
|
||||
const abiFallback = this.opts.abiFallback ?? 'rpc-current'
|
||||
const abiStore = new AbiStore(this.redis)
|
||||
const abiBootstrapper = new AbiBootstrapper(this.chainClient, abiStore, { abiFallback })
|
||||
@@ -83,14 +94,6 @@ export class Parser {
|
||||
const syncKey = RedisKeys.syncHash(chainId)
|
||||
const eventsStream = RedisKeys.eventsStream(chainId)
|
||||
|
||||
const lastBlockNum = await this.redis.hget(syncKey, 'block_num')
|
||||
const lastBlockId = await this.redis.hget(syncKey, 'block_id')
|
||||
|
||||
const havePositions =
|
||||
lastBlockNum && lastBlockId
|
||||
? [{ blockNum: Number(lastBlockNum), blockId: lastBlockId }]
|
||||
: []
|
||||
|
||||
const xtrimOpts: XtrimSupervisorOpts = {
|
||||
redis: this.redis,
|
||||
stream: eventsStream,
|
||||
@@ -115,46 +118,41 @@ export class Parser {
|
||||
})
|
||||
}
|
||||
|
||||
const streamOpts = {
|
||||
startBlock: havePositions[0]?.blockNum ?? 0,
|
||||
havePositions,
|
||||
}
|
||||
|
||||
const irreversibleOnly = this.opts.irreversibleOnly ?? false
|
||||
const forkDetector = new ForkDetector(chainId)
|
||||
|
||||
for await (const block of this.chainClient.streamBlocks(streamOpts)) {
|
||||
if (this.stopSignal) break
|
||||
// ReconnectSupervisor оборачивает (пере)подключение к SHiP экспоненциальным
|
||||
// backoff. Без него разрыв ноды бросал ShipConnectionError наверх из start(),
|
||||
// и темп переподключений целиком зависел от внешнего супервизора процесса —
|
||||
// мёртвая нода = рестарт без пауз = долбёжка. Теперь паузы 1→2→5→15→60с,
|
||||
// а после maxAttempts подряд неудач — чистый выход (процесс рестартит оркестратор).
|
||||
const reconnectCfg = this.opts.reconnect ?? {}
|
||||
const supervisor = new ReconnectSupervisor({
|
||||
...(reconnectCfg.maxAttempts !== undefined ? { maxAttempts: reconnectCfg.maxAttempts } : {}),
|
||||
...(reconnectCfg.backoffSeconds !== undefined ? { backoffSeconds: reconnectCfg.backoffSeconds } : {}),
|
||||
onAttempt: (attempt, delayMs) => {
|
||||
this.log?.warn({ attempt, delayMs }, 'SHiP-соединение потеряно, переподключение через backoff')
|
||||
},
|
||||
onGiveUp: (attempts) => {
|
||||
this.log?.error({ attempts }, 'SHiP-переподключение исчерпало все попытки, остановка парсера')
|
||||
if (!this.opts.noSignalHandlers) process.exit(1)
|
||||
},
|
||||
})
|
||||
|
||||
// Backpressure: пока backlog стрима выше окна — ждём, не доливаем и не
|
||||
// ack'аем (SHiP сам притормозит после max_messages_in_flight). Блок уже в
|
||||
// RAM (≤ окна SHiP), обработаем его как только consumer освободит место.
|
||||
if (this.backpressureGate) {
|
||||
await this.backpressureGate.waitForCapacity(() => this.stopSignal)
|
||||
if (this.stopSignal) break
|
||||
}
|
||||
|
||||
if (irreversibleOnly && block.thisBlock.blockNum > block.lastIrreversible.blockNum) {
|
||||
this.chainClient.ack(1)
|
||||
continue
|
||||
}
|
||||
|
||||
const forkEvent = forkDetector.check(block.thisBlock.blockNum, block.thisBlock.blockId)
|
||||
const events: ParserEvent[] = await this.blockProcessor.process(block)
|
||||
const toPublish: ParserEvent[] = forkEvent ? [forkEvent, ...events] : events
|
||||
|
||||
for (const event of toPublish) {
|
||||
await this.redis.xadd(eventsStream, this.eventToFields(event))
|
||||
}
|
||||
|
||||
await this.redis.hset(syncKey, {
|
||||
block_num: String(block.thisBlock.blockNum),
|
||||
block_id: block.thisBlock.blockId,
|
||||
last_updated: new Date().toISOString(),
|
||||
})
|
||||
|
||||
this.chainClient.ack(1)
|
||||
}
|
||||
await runStreamLoop({
|
||||
chainClient: this.chainClient!,
|
||||
redis: this.redis!,
|
||||
blockProcessor: this.blockProcessor!,
|
||||
backpressureGate: this.backpressureGate,
|
||||
forkDetector,
|
||||
supervisor,
|
||||
syncKey,
|
||||
eventsStream,
|
||||
irreversibleOnly,
|
||||
eventToFields: (event) => this.eventToFields(event),
|
||||
isStopped: () => this.stopSignal,
|
||||
log: this.log,
|
||||
})
|
||||
}
|
||||
|
||||
private eventToFields(event: ParserEvent): Record<string, string> {
|
||||
|
||||
@@ -51,9 +51,15 @@ export class ReconnectSupervisor {
|
||||
/**
|
||||
* Запускает fn и повторяет при исключении с паузами.
|
||||
*
|
||||
* fn получает callback resetBackoff: вызови его, когда fn сделала реальный
|
||||
* прогресс (например успешно обработала блок на свежем соединении) — счётчик
|
||||
* попыток сбросится в 0. Без сброса долгоживущее соединение, которое изредка
|
||||
* рвётся, копило бы attempt до maxAttempts и завершало процесс, хотя между
|
||||
* разрывами оно часами работало штатно.
|
||||
*
|
||||
* Псевдокод:
|
||||
* loop:
|
||||
* try: return await fn()
|
||||
* try: return await fn(resetBackoff)
|
||||
* catch: attempt++
|
||||
* if attempt >= maxAttempts: onGiveUp(); throw
|
||||
* delay = backoffSeconds[min(attempt-1, len-1)] * 1000
|
||||
@@ -61,11 +67,12 @@ export class ReconnectSupervisor {
|
||||
*
|
||||
* @returns Результат первого успешного вызова fn.
|
||||
*/
|
||||
async run<T>(fn: () => Promise<T>): Promise<T> {
|
||||
async run<T>(fn: (resetBackoff: () => void) => Promise<T>): Promise<T> {
|
||||
let attempt = 0
|
||||
const resetBackoff = (): void => { attempt = 0 }
|
||||
for (;;) {
|
||||
try {
|
||||
return await fn()
|
||||
return await fn(resetBackoff)
|
||||
} catch (err) {
|
||||
attempt++
|
||||
if (attempt >= this.maxAttempts) {
|
||||
|
||||
@@ -0,0 +1,129 @@
|
||||
/**
|
||||
* Цикл чтения блоков SHiP с авто-переподключением.
|
||||
*
|
||||
* Вынесен из Parser.start() отдельной функцией с инъекцией зависимостей —
|
||||
* чтобы покрыть reconnect/resume-логику юнит-тестами на фейках, без реального
|
||||
* Redis и SHiP-ноды.
|
||||
*
|
||||
* Контракт:
|
||||
* - supervisor оборачивает одну попытку соединения экспоненциальным backoff;
|
||||
* - разрыв ноды → streamBlocks бросает → supervisor ждёт паузу и зовёт заново;
|
||||
* - позиция возобновления перечитывается из Redis на КАЖДОМ подключении
|
||||
* (sync-hash пишется после каждого блока) — продолжаем ровно с последнего
|
||||
* записанного блока, без потерь и дублей;
|
||||
* - resetBackoff() на каждом обработанном блоке: реальный прогресс возвращает
|
||||
* лесенку backoff в начало, чтобы долгая стабильная сессия после редкого
|
||||
* разрыва не копила попытки до выхода;
|
||||
* - штатная остановка (isStopped()) закрывает сокет извне → for-await бросает;
|
||||
* ловим и выходим без backoff и без retry.
|
||||
*/
|
||||
|
||||
import type { ShipBlock, GetBlocksOptions } from '@coopenomics/coopos-ship-reader'
|
||||
import type { ReconnectSupervisor } from './ReconnectSupervisor.js'
|
||||
import type { BackpressureGate } from './BackpressureGate.js'
|
||||
import type { ForkDetector } from './ForkDetector.js'
|
||||
import type { Logger } from '../logger.js'
|
||||
import type { ParserEvent } from '../types.js'
|
||||
|
||||
/** Узкий срез ChainClient, нужный циклу (упрощает фейки в тестах). */
|
||||
export interface StreamLoopChainClient {
|
||||
connect(): Promise<unknown>
|
||||
streamBlocks(opts: GetBlocksOptions): AsyncIterable<ShipBlock>
|
||||
ack(n: number): void
|
||||
}
|
||||
|
||||
/** Узкий срез RedisStore, нужный циклу. */
|
||||
export interface StreamLoopRedis {
|
||||
hget(key: string, field: string): Promise<string | null>
|
||||
hset(key: string, fields: Record<string, string>): Promise<unknown>
|
||||
xadd(stream: string, fields: Record<string, string>): Promise<unknown>
|
||||
}
|
||||
|
||||
/** Узкий срез BlockProcessor. */
|
||||
export interface StreamLoopBlockProcessor {
|
||||
process(block: ShipBlock): Promise<ParserEvent[]>
|
||||
}
|
||||
|
||||
export interface StreamLoopDeps {
|
||||
chainClient: StreamLoopChainClient
|
||||
redis: StreamLoopRedis
|
||||
blockProcessor: StreamLoopBlockProcessor
|
||||
backpressureGate: BackpressureGate | null
|
||||
forkDetector: Pick<ForkDetector, 'check'>
|
||||
supervisor: ReconnectSupervisor
|
||||
syncKey: string
|
||||
eventsStream: string
|
||||
irreversibleOnly: boolean
|
||||
eventToFields: (event: ParserEvent) => Record<string, string>
|
||||
isStopped: () => boolean
|
||||
log: Logger | null
|
||||
}
|
||||
|
||||
export async function runStreamLoop(deps: StreamLoopDeps): Promise<void> {
|
||||
const {
|
||||
chainClient, redis, blockProcessor, backpressureGate, forkDetector,
|
||||
supervisor, syncKey, eventsStream, irreversibleOnly, eventToFields, isStopped,
|
||||
} = deps
|
||||
|
||||
const runOneConnection = async (resetBackoff: () => void): Promise<void> => {
|
||||
if (isStopped()) return
|
||||
|
||||
// connect() идемпотентен при OPEN ws (первый заход — соединение уже из
|
||||
// start-connect выше); после разрыва создаёт свежий WebSocket. handshake
|
||||
// внутри отдаёт кешированный chainId, повторный get_status не шлётся.
|
||||
await chainClient.connect()
|
||||
|
||||
const resumeNum = await redis.hget(syncKey, 'block_num')
|
||||
const resumeId = await redis.hget(syncKey, 'block_id')
|
||||
const havePositions =
|
||||
resumeNum && resumeId ? [{ blockNum: Number(resumeNum), blockId: resumeId }] : []
|
||||
const streamOpts = {
|
||||
startBlock: havePositions[0]?.blockNum ?? 0,
|
||||
havePositions,
|
||||
}
|
||||
|
||||
try {
|
||||
for await (const block of chainClient.streamBlocks(streamOpts)) {
|
||||
if (isStopped()) return
|
||||
|
||||
// Backpressure: пока backlog стрима выше окна — ждём, не доливаем и не
|
||||
// ack'аем (SHiP сам притормозит после max_messages_in_flight). Блок уже
|
||||
// в RAM (≤ окна SHiP), обработаем его как только consumer освободит место.
|
||||
if (backpressureGate) {
|
||||
await backpressureGate.waitForCapacity(() => isStopped())
|
||||
if (isStopped()) return
|
||||
}
|
||||
|
||||
if (irreversibleOnly && block.thisBlock.blockNum > block.lastIrreversible.blockNum) {
|
||||
chainClient.ack(1)
|
||||
resetBackoff()
|
||||
continue
|
||||
}
|
||||
|
||||
const forkEvent = forkDetector.check(block.thisBlock.blockNum, block.thisBlock.blockId)
|
||||
const events = await blockProcessor.process(block)
|
||||
const toPublish: ParserEvent[] = forkEvent ? [forkEvent, ...events] : events
|
||||
|
||||
for (const event of toPublish) {
|
||||
await redis.xadd(eventsStream, eventToFields(event))
|
||||
}
|
||||
|
||||
await redis.hset(syncKey, {
|
||||
block_num: String(block.thisBlock.blockNum),
|
||||
block_id: block.thisBlock.blockId,
|
||||
last_updated: new Date().toISOString(),
|
||||
})
|
||||
|
||||
chainClient.ack(1)
|
||||
resetBackoff()
|
||||
}
|
||||
} catch (err) {
|
||||
// Штатная остановка закрывает WebSocket → for-await бросает
|
||||
// ShipConnectionError. Это не сбой — выходим без backoff и без retry.
|
||||
if (isStopped()) return
|
||||
throw err
|
||||
}
|
||||
}
|
||||
|
||||
await supervisor.run(runOneConnection)
|
||||
}
|
||||
@@ -71,3 +71,69 @@ describe('ReconnectSupervisor — exhaustion', () => {
|
||||
delays.forEach(d => expect(d).toBe(0))
|
||||
})
|
||||
})
|
||||
|
||||
describe('ReconnectSupervisor — resetBackoff (прогресс сбрасывает счётчик)', () => {
|
||||
it('передаёт resetBackoff в fn', async () => {
|
||||
const sup = new ReconnectSupervisor({ backoffSeconds: [0] })
|
||||
let gotCallback = false
|
||||
await sup.run(async (reset) => {
|
||||
gotCallback = typeof reset === 'function'
|
||||
return 'ok'
|
||||
})
|
||||
expect(gotCallback).toBe(true)
|
||||
})
|
||||
|
||||
it('не вызывает onGiveUp пока fn делает прогресс между разрывами', async () => {
|
||||
// Сценарий: соединение каждый раз отдаёт «блок» (прогресс), затем рвётся.
|
||||
// Без сброса счётчик дошёл бы до maxAttempts=3 и завершил процесс. С resetBackoff
|
||||
// на каждом прогрессе счётчик возвращается в 0 — даём 10 циклов разрыва, всё живёт.
|
||||
let gaveUp = false
|
||||
let drops = 0
|
||||
const sup = new ReconnectSupervisor({
|
||||
maxAttempts: 3,
|
||||
backoffSeconds: [0],
|
||||
onGiveUp: () => { gaveUp = true; throw new Error('gave up') },
|
||||
})
|
||||
await sup.run(async (reset) => {
|
||||
reset() // обработали блок — прогресс есть
|
||||
drops++
|
||||
if (drops < 10) throw new Error('node dropped') // разрыв
|
||||
return 'done'
|
||||
})
|
||||
expect(gaveUp).toBe(false)
|
||||
expect(drops).toBe(10)
|
||||
})
|
||||
|
||||
it('БЕЗ прогресса исчерпывает попытки и вызывает onGiveUp', async () => {
|
||||
// Контроль к предыдущему: мёртвая нода (resetBackoff не зовётся) → exit.
|
||||
let gaveUp = false
|
||||
const sup = new ReconnectSupervisor({
|
||||
maxAttempts: 3,
|
||||
backoffSeconds: [0],
|
||||
onGiveUp: () => { gaveUp = true; throw new Error('gave up') },
|
||||
})
|
||||
await expect(
|
||||
sup.run(async () => { throw new Error('dead node') }),
|
||||
).rejects.toThrow()
|
||||
expect(gaveUp).toBe(true)
|
||||
})
|
||||
|
||||
it('resetBackoff возвращает нумерацию onAttempt в начало', async () => {
|
||||
const attempts: number[] = []
|
||||
const sup = new ReconnectSupervisor({
|
||||
maxAttempts: 10,
|
||||
backoffSeconds: [0],
|
||||
onAttempt: (attempt) => attempts.push(attempt),
|
||||
})
|
||||
let calls = 0
|
||||
await sup.run(async (reset) => {
|
||||
calls++
|
||||
// calls: 1 throw → attempt 1; 2 throw → attempt 2; 3 reset+throw → attempt 1 снова; 4 ok
|
||||
if (calls === 3) reset()
|
||||
if (calls < 4) throw new Error('fail')
|
||||
return 'ok'
|
||||
})
|
||||
// attempt-числа: [1, 2, 1] — третий разрыв после reset снова стартует с 1
|
||||
expect(attempts).toEqual([1, 2, 1])
|
||||
})
|
||||
})
|
||||
|
||||
@@ -0,0 +1,160 @@
|
||||
/**
|
||||
* Юнит-тесты цикла чтения с авто-переподключением (runStreamLoop) на фейках —
|
||||
* без реального Redis и SHiP-ноды.
|
||||
*
|
||||
* Покрываем то, ради чего цикл вынесен из Parser:
|
||||
* 1. разрыв ноды → backoff → переподключение, чтение продолжается;
|
||||
* 2. после разрыва позиция возобновления берётся из Redis (последний блок) —
|
||||
* без потерь и дублей;
|
||||
* 3. штатная остановка не вызывает reconnect;
|
||||
* 4. полностью мёртвая нода исчерпывает попытки и завершается (onGiveUp),
|
||||
* а НЕ долбит без пауз.
|
||||
*/
|
||||
|
||||
import { describe, it, expect } from 'vitest'
|
||||
import type { GetBlocksOptions, ShipBlock } from '@coopenomics/coopos-ship-reader'
|
||||
import { runStreamLoop } from '../../src/core/streamLoop.js'
|
||||
import type { StreamLoopChainClient, StreamLoopRedis } from '../../src/core/streamLoop.js'
|
||||
import { ReconnectSupervisor } from '../../src/core/ReconnectSupervisor.js'
|
||||
|
||||
function makeBlock(n: number): ShipBlock {
|
||||
return {
|
||||
thisBlock: { blockNum: n, blockId: `id${n}` },
|
||||
lastIrreversible: { blockNum: n, blockId: `id${n}` },
|
||||
head: { blockNum: n, blockId: `id${n}` },
|
||||
prevBlock: null,
|
||||
blockTime: '2024-06-01T00:00:00.000',
|
||||
traces: [],
|
||||
deltas: [],
|
||||
} as unknown as ShipBlock
|
||||
}
|
||||
|
||||
function makeRedis(): StreamLoopRedis & { hash: Record<string, Record<string, string>> } {
|
||||
const hash: Record<string, Record<string, string>> = {}
|
||||
return {
|
||||
hash,
|
||||
async hget(key, field) { return hash[key]?.[field] ?? null },
|
||||
async hset(key, fields) { hash[key] = { ...(hash[key] ?? {}), ...fields }; return 'OK' },
|
||||
async xadd() { return '1-0' },
|
||||
}
|
||||
}
|
||||
|
||||
/** Скриптованный SHiP-клиент: массив сессий, каждая — список блоков + рвётся/нет. */
|
||||
function makeChainClient(sessions: Array<{ blocks: number[]; drop: boolean }>) {
|
||||
const captured: GetBlocksOptions[] = []
|
||||
let session = 0
|
||||
let connectCalls = 0
|
||||
const client: StreamLoopChainClient = {
|
||||
async connect() { connectCalls++ },
|
||||
ack() {},
|
||||
async *streamBlocks(opts: GetBlocksOptions) {
|
||||
captured.push(opts)
|
||||
const s = sessions[session++]
|
||||
if (!s) return
|
||||
for (const n of s.blocks) yield makeBlock(n)
|
||||
if (s.drop) throw new Error('node dropped')
|
||||
},
|
||||
}
|
||||
return {
|
||||
client,
|
||||
captured,
|
||||
get connectCalls() { return connectCalls },
|
||||
}
|
||||
}
|
||||
|
||||
const baseDeps = (over: Partial<Parameters<typeof runStreamLoop>[0]>): Parameters<typeof runStreamLoop>[0] => ({
|
||||
chainClient: over.chainClient!,
|
||||
redis: over.redis!,
|
||||
blockProcessor: over.blockProcessor ?? { process: async () => [] },
|
||||
backpressureGate: null,
|
||||
forkDetector: { check: () => null },
|
||||
supervisor: over.supervisor ?? new ReconnectSupervisor({ maxAttempts: 5, backoffSeconds: [0] }),
|
||||
syncKey: 'sync',
|
||||
eventsStream: 'events',
|
||||
irreversibleOnly: false,
|
||||
eventToFields: () => ({ data: '{}' }),
|
||||
isStopped: over.isStopped ?? (() => false),
|
||||
log: null,
|
||||
...over,
|
||||
})
|
||||
|
||||
describe('runStreamLoop — переподключение после разрыва', () => {
|
||||
it('после разрыва переподключается и продолжает с последнего записанного блока', async () => {
|
||||
const redis = makeRedis()
|
||||
const cc = makeChainClient([
|
||||
{ blocks: [1, 2, 3], drop: true }, // упала после блока 3
|
||||
{ blocks: [4, 5], drop: false }, // реконнект → дочитали и штатно завершились
|
||||
])
|
||||
|
||||
await runStreamLoop(baseDeps({ chainClient: cc.client, redis }))
|
||||
|
||||
// Подключались дважды: первичное + 1 реконнект
|
||||
expect(cc.connectCalls).toBe(2)
|
||||
expect(cc.captured).toHaveLength(2)
|
||||
// Первая сессия — с нуля (позиции в Redis ещё нет)
|
||||
expect(cc.captured[0]?.startBlock).toBe(0)
|
||||
// Вторая — возобновление ровно с блока 3 (последнего записанного), без потерь/дублей
|
||||
expect(cc.captured[1]?.startBlock).toBe(3)
|
||||
expect(cc.captured[1]?.havePositions).toEqual([{ blockNum: 3, blockId: 'id3' }])
|
||||
// Финальная позиция — блок 5
|
||||
expect(redis.hash['sync']?.['block_num']).toBe('5')
|
||||
})
|
||||
|
||||
it('штатная остановка не вызывает reconnect', async () => {
|
||||
const redis = makeRedis()
|
||||
const cc = makeChainClient([{ blocks: [1, 2, 3, 4], drop: false }])
|
||||
let processed = 0
|
||||
const deps = baseDeps({
|
||||
chainClient: cc.client,
|
||||
redis,
|
||||
blockProcessor: { process: async () => { processed++; return [] } },
|
||||
isStopped: () => processed >= 2, // останавливаемся после 2 блоков
|
||||
})
|
||||
|
||||
await runStreamLoop(deps)
|
||||
|
||||
expect(cc.connectCalls).toBe(1) // без переподключений
|
||||
expect(redis.hash['sync']?.['block_num']).toBe('2')
|
||||
})
|
||||
|
||||
it('мёртвая нода исчерпывает попытки и вызывает onGiveUp (не долбит без пауз)', async () => {
|
||||
const redis = makeRedis()
|
||||
const cc = makeChainClient([
|
||||
{ blocks: [], drop: true },
|
||||
{ blocks: [], drop: true },
|
||||
{ blocks: [], drop: true },
|
||||
{ blocks: [], drop: true },
|
||||
])
|
||||
let gaveUp = false
|
||||
const supervisor = new ReconnectSupervisor({
|
||||
maxAttempts: 3,
|
||||
backoffSeconds: [0],
|
||||
onGiveUp: () => { gaveUp = true; throw new Error('gave up') },
|
||||
})
|
||||
|
||||
await expect(
|
||||
runStreamLoop(baseDeps({ chainClient: cc.client, redis, supervisor })),
|
||||
).rejects.toThrow()
|
||||
expect(gaveUp).toBe(true)
|
||||
})
|
||||
|
||||
it('долгая стабильная сессия после редкого разрыва не копит попытки до выхода', async () => {
|
||||
// 6 разрывов подряд, но каждая сессия отдаёт блок (прогресс) → resetBackoff
|
||||
// держит счётчик у нуля. maxAttempts=3 НЕ срабатывает, т.к. между разрывами прогресс.
|
||||
const redis = makeRedis()
|
||||
const sessions = Array.from({ length: 6 }, (_, i) => ({ blocks: [i + 1], drop: i < 5 }))
|
||||
const cc = makeChainClient(sessions)
|
||||
let gaveUp = false
|
||||
const supervisor = new ReconnectSupervisor({
|
||||
maxAttempts: 3,
|
||||
backoffSeconds: [0],
|
||||
onGiveUp: () => { gaveUp = true; throw new Error('gave up') },
|
||||
})
|
||||
|
||||
await runStreamLoop(baseDeps({ chainClient: cc.client, redis, supervisor }))
|
||||
|
||||
expect(gaveUp).toBe(false)
|
||||
expect(cc.connectCalls).toBe(6)
|
||||
expect(redis.hash['sync']?.['block_num']).toBe('6')
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user