From 92811c258ffe185790e9e80b4d6a7355d17ff3ab Mon Sep 17 00:00:00 2001 From: coopops Date: Fri, 5 Jun 2026 13:00:11 +0000 Subject: [PATCH] =?UTF-8?q?feat(parser):=20=D0=B0=D0=B2=D1=82=D0=BE-=D0=BF?= =?UTF-8?q?=D0=B5=D1=80=D0=B5=D0=BF=D0=BE=D0=B4=D0=BA=D0=BB=D1=8E=D1=87?= =?UTF-8?q?=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=BA=20SHiP=20=D1=87=D0=B5=D1=80?= =?UTF-8?q?=D0=B5=D0=B7=20ReconnectSupervisor=20=D1=81=20backoff?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ReconnectSupervisor был написан, протестирован и экспортнут, но НИГДЕ не подключён к циклу чтения. Parser.start() имел голый for-await без try/catch: падение ноды → ShipConnectionError наверх → темп переподключений целиком зависел от внешнего супервизора процесса. Мёртвая нода = рестарт без пауз = долбёжка ноды без ограничителей. Что сделано: - Цикл чтения вынесен в runStreamLoop (core/streamLoop.ts) с инъекцией зависимостей — чтобы покрыть reconnect/resume юнит-тестами на фейках, без реального Redis и SHiP. - runStreamLoop оборачивает (пере)подключение в ReconnectSupervisor: backoff 1→2→5→15→60с, после maxAttempts без прогресса — process.exit(1) (оркестратор перезапустит штатно). - Позиция возобновления перечитывается из Redis на КАЖДОМ подключении (sync-hash пишется после каждого блока) — продолжаем ровно с последнего записанного блока, без потерь и дублей. - ReconnectSupervisor.run теперь передаёт fn callback resetBackoff: каждый обработанный блок сбрасывает счётчик попыток, чтобы долгая стабильная сессия после редкого разрыва не копила attempt до выхода. - Штатная остановка (stop закрывает ws) ловится и не вызывает лишний backoff. - Parser получил структурный логгер (opts.logger) для onAttempt/onGiveUp. Тесты: +4 runStreamLoop (reconnect+resume, stop без reconnect, give-up на мёртвой ноде, resetBackoff на прогрессе), +4 ReconnectSupervisor (resetBackoff). Всего 229 unit. README: секция «Авто-переподключение к SHiP». Co-Authored-By: Claude Opus 4.8 --- README.md | 17 +- packages/parser2/package.json | 2 +- packages/parser2/src/core/Parser.ts | 86 +++++----- .../parser2/src/core/ReconnectSupervisor.ts | 13 +- packages/parser2/src/core/streamLoop.ts | 129 ++++++++++++++ .../test/unit/reconnectSupervisor.test.ts | 66 ++++++++ packages/parser2/test/unit/streamLoop.test.ts | 160 ++++++++++++++++++ 7 files changed, 424 insertions(+), 49 deletions(-) create mode 100644 packages/parser2/src/core/streamLoop.ts create mode 100644 packages/parser2/test/unit/streamLoop.test.ts diff --git a/README.md b/README.md index 0cec704..3f41048 100644 --- a/README.md +++ b/README.md @@ -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-теста ### Бенчмарк diff --git a/packages/parser2/package.json b/packages/parser2/package.json index 6a4faae..77384b1 100644 --- a/packages/parser2/package.json +++ b/packages/parser2/package.json @@ -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", diff --git a/packages/parser2/src/core/Parser.ts b/packages/parser2/src/core/Parser.ts index d679dda..de95eb4 100644 --- a/packages/parser2/src/core/Parser.ts +++ b/packages/parser2/src/core/Parser.ts @@ -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 { diff --git a/packages/parser2/src/core/ReconnectSupervisor.ts b/packages/parser2/src/core/ReconnectSupervisor.ts index 80a880f..1dbe035 100644 --- a/packages/parser2/src/core/ReconnectSupervisor.ts +++ b/packages/parser2/src/core/ReconnectSupervisor.ts @@ -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(fn: () => Promise): Promise { + async run(fn: (resetBackoff: () => void) => Promise): Promise { let attempt = 0 + const resetBackoff = (): void => { attempt = 0 } for (;;) { try { - return await fn() + return await fn(resetBackoff) } catch (err) { attempt++ if (attempt >= this.maxAttempts) { diff --git a/packages/parser2/src/core/streamLoop.ts b/packages/parser2/src/core/streamLoop.ts new file mode 100644 index 0000000..8a76f62 --- /dev/null +++ b/packages/parser2/src/core/streamLoop.ts @@ -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 + streamBlocks(opts: GetBlocksOptions): AsyncIterable + ack(n: number): void +} + +/** Узкий срез RedisStore, нужный циклу. */ +export interface StreamLoopRedis { + hget(key: string, field: string): Promise + hset(key: string, fields: Record): Promise + xadd(stream: string, fields: Record): Promise +} + +/** Узкий срез BlockProcessor. */ +export interface StreamLoopBlockProcessor { + process(block: ShipBlock): Promise +} + +export interface StreamLoopDeps { + chainClient: StreamLoopChainClient + redis: StreamLoopRedis + blockProcessor: StreamLoopBlockProcessor + backpressureGate: BackpressureGate | null + forkDetector: Pick + supervisor: ReconnectSupervisor + syncKey: string + eventsStream: string + irreversibleOnly: boolean + eventToFields: (event: ParserEvent) => Record + isStopped: () => boolean + log: Logger | null +} + +export async function runStreamLoop(deps: StreamLoopDeps): Promise { + const { + chainClient, redis, blockProcessor, backpressureGate, forkDetector, + supervisor, syncKey, eventsStream, irreversibleOnly, eventToFields, isStopped, + } = deps + + const runOneConnection = async (resetBackoff: () => void): Promise => { + 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) +} diff --git a/packages/parser2/test/unit/reconnectSupervisor.test.ts b/packages/parser2/test/unit/reconnectSupervisor.test.ts index 7fc6879..afe240b 100644 --- a/packages/parser2/test/unit/reconnectSupervisor.test.ts +++ b/packages/parser2/test/unit/reconnectSupervisor.test.ts @@ -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]) + }) +}) diff --git a/packages/parser2/test/unit/streamLoop.test.ts b/packages/parser2/test/unit/streamLoop.test.ts new file mode 100644 index 0000000..47ed984 --- /dev/null +++ b/packages/parser2/test/unit/streamLoop.test.ts @@ -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> } { + const hash: Record> = {} + 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[0]>): Parameters[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') + }) +})