Merge pull request 'feat(parser): авто-переподключение к SHiP с backoff (1.3.0)' (#12) from dev into main
Release / Release (push) Successful in 9m35s

This commit was merged in pull request #12.
This commit is contained in:
2026-06-05 13:01:02 +00:00
7 changed files with 424 additions and 49 deletions
+16 -1
View File
@@ -287,6 +287,21 @@ backpressure:
Рекомендация по инфраструктуре: на 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 внутри одного блока
@@ -345,7 +360,7 @@ pnpm --filter @coopenomics/parser2 test:integration # нужен Docker для
```
**Текущее покрытие:**
- `@coopenomics/parser2`: 221 unit-тестов + integration (xtrim/backpressure на реальном Redis)
- `@coopenomics/parser2`: 229 unit-тестов + integration (xtrim/backpressure на реальном Redis)
- `@coopenomics/coopos-ship-reader`: 54 unit-теста
### Бенчмарк
+1 -1
View File
@@ -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",
+41 -43
View File
@@ -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
// 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(),
// 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)
},
})
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) {
+129
View File
@@ -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')
})
})