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 <noreply@anthropic.com>
22 KiB
Coopenomics Parser
Универсальный индексер блокчейнов EOSIO / Antelope: читает блоки из State History Plugin (SHiP) по WebSocket, декодирует actions и дельты таблиц с учётом исторических ABI, публикует унифицированный поток событий в Redis Stream. Потребители получают события через ParserClient с single-active-consumer lock'ом, recovery после сбоев и dead-letter для poison-messages.
Зачем
SHiP отдаёт сырые бинарные блоки — чтобы превратить их в прикладной поток событий, нужно: держать кэш ABI за каждый блок (контракты обновляют ABI, старые блоки декодируются старой версией), обрабатывать форки и переподключения, распределять нагрузку между подписчиками. Этот проект закрывает весь этот слой и отдаёт вам единый стрим action / delta / native-delta / fork событий.
Пакеты монорепы
| Пакет | Описание | npm |
|---|---|---|
@coopenomics/parser2 |
Ядро индексера (Parser) + подписочный клиент (ParserClient), CLI, observability | @coopenomics/parser2 |
@coopenomics/coopos-ship-reader |
Низкоуровневый WebSocket SHiP клиент с поддержкой 24 нативных системных таблиц | @coopenomics/coopos-ship-reader |
Быстрый старт
1. Установка
pnpm add @coopenomics/parser2
# или если используете только SHiP-клиент без Redis-пайплайна:
pnpm add @coopenomics/coopos-ship-reader
Требования рантайма:
- Node ≥ 20
- Redis ≥ 7 (с persistence — AOF/RDB)
- Доступ до SHiP endpoint блокчейн-ноды (
ws://node:8080) - Доступ до Chain API ноды (
http://node:8888)
2. Запуск парсера
Создай parser.config.yaml:
ship:
url: ws://my-nodeos:8080
timeoutMs: 15000
chain:
url: http://my-nodeos:8888
redis:
url: redis://localhost:6379
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
health:
enabled: true
port: 8081
metrics:
enabled: true
port: 9090
Запусти через CLI (после глобальной установки пакета):
parser start --config parser.config.yaml
Парсер подключится к SHiP, начнёт читать блоки с последней checkpoint-позиции (или с head если запускается впервые), и публиковать события в Redis stream ce:parser2:<chainId>:events.
3. Подписка на события из приложения
import { ParserClient } from '@coopenomics/parser2'
const client = new ParserClient({
subscriptionId: 'my-app',
filters: [
{ kind: 'action', account: 'eosio.token', name: 'transfer' },
{ kind: 'action', account: 'eosio.token', name: 'issue' },
],
startFrom: 'last_known',
redis: { url: 'redis://localhost:6379' },
chain: { id: 'eb004c7dcb6e92f5ba9c98e0d86a616e79ec4e7c80bc1a66c0d6c8d6c...' },
})
for await (const event of client.stream()) {
if (event.kind === 'action') {
console.log(
`${event.block_num} | ${event.account}::${event.name}`,
event.data,
)
// Если обработчик упадёт (throw) — событие попадёт в dead-letter после 3 попыток
}
}
Несколько реплик одного subscriptionId автоматически выберут одного active-потребителя через distributed lock; остальные встанут в standby и подхватят при падении active.
4. CLI-утилиты
# Посмотреть все зарегистрированные подписки и их отставание
parser list-subscriptions
# Сбросить cursor подписки на начало стрима
parser reset-subscription --sub-id my-app --start-from 0
# Посмотреть dead-letter сообщения
parser list-dead-letters --sub-id my-app
# Replay одного dead-letter обратно в основной поток
parser replay-dead-letter --sub-id my-app --dl-id 1699999999999-0
# Удалить старые ABI-версии (GC)
parser abi-prune --keep-last 10 --all-contracts
Docker
Публичный образ: dicoop/parser2 — multi-arch (linux/amd64, linux/arm64).
docker pull dicoop/parser2:latest
# или конкретная версия
docker pull dicoop/parser2:1.0.1
Минимальный запуск
Парсер читает конфигурацию из YAML-файла, путь передаётся через --config:
docker run --rm \
-v $(pwd)/parser.config.yaml:/app/parser.config.yaml:ro \
--network host \
dicoop/parser2:latest \
start --config /app/parser.config.yaml
docker-compose.yml
Полная конфигурация с Redis:
services:
redis:
image: redis:7-alpine
command: >-
redis-server
--appendonly yes
--appendfsync everysec
volumes:
- redis-data:/data
ports:
- "6379:6379"
parser:
image: dicoop/parser2:1.0.1
depends_on:
- redis
volumes:
- ./parser.config.yaml:/app/parser.config.yaml:ro
command: ["start", "--config", "/app/parser.config.yaml"]
ports:
- "8081:8081" # /health
- "9090:9090" # /metrics (Prometheus)
restart: unless-stopped
volumes:
redis-data:
parser.config.yaml рядом с compose-файлом:
ship:
url: ws://nodeos:8080
timeoutMs: 15000
chain:
url: http://nodeos:8888
redis:
url: redis://redis:6379
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]
logger:
level: info
pretty: false # для production — JSON logs в stdout
health:
enabled: true
port: 8081
metrics:
enabled: true
port: 9090
irreversibleOnly: false # читать head-блоки (false) или ждать last_irreversible (true)
Запуск: docker compose up -d. /health отдаст 200 OK как только парсер подключился к SHiP, /metrics — экспорт для Prometheus.
Переопределение CLI-команды
Контейнер по умолчанию запускает parser start, но можно использовать любую другую команду:
# Посмотреть подписки в прод-Redis
docker run --rm --network host \
-e REDIS_URL=redis://localhost:6379 \
dicoop/parser2:latest \
list-subscriptions
Архитектура
┌──────────────┐ WS ┌──────────────┐ ┌──────────────┐
│ EOSIO node │────────▶│ Parser │──XADD────▶│ │
│ (SHiP+RPC) │ │ │ │ │
└──────────────┘ │ • AbiStore │ │ Redis │
│ • Workers │ │ │
│ • ForkDet. │ │ Streams + │
│ • XtrimSup. │ │ Hashes + │
│ • Backpr. │ │ ZSets │
└──────────────┘ │ │
│ │
┌──────────────┐ │ │
│ ParserClient │◀─XREADGR──│ │
│ #1 active │ │ │
└──────────────┘ └──────────────┘
┌──────────────┐ ▲
│ ParserClient │────────────────────┘
│ #2 standby │ (ждёт lock)
└──────────────┘
Подробности:
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 +--+----------------------------> время
пишем ждём пишем ждём пишем
Конфиг (включён по умолчанию):
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) перезапустил процесс штатно.
reconnect:
maxAttempts: 10
backoffSeconds: [1, 2, 5, 10, 30, 60, 120, 300, 600, 1800]
Известные ограничения
Схлопывание дельт create+remove внутри одного блока
Парсер не увидит жизненный цикл строки таблицы, у которой emplace и erase произошли внутри одного блока (даже в разных транзакциях). Это поведение унаследовано от upstream-плагина state_history_plugin (SHiP) ноджеса: на уровне chainbase::undo_index::on_remove строка, созданная в текущей undo session и тут же удалённая, физически уничтожается без записи в _removed_values — это by-design оптимизация для отката блока (нечего откатывать), но она же делает event невидимым для SHiP-снимка состояния. SHiP пакует только то, что есть в undo session, и не имеет шанса увидеть транзиентную пару.
Симптомы:
- На цепочке таблица пустая (получена и тут же удалена).
- ACTION-логи и
getclearance/apprvappndx-подобные действия видны нормально — экшены идут параллельным потоком и не зависят от undo. - В Redis stream / БД-индексе нет ни
present:true, ниpresent:falseсобытия для конкретногоprimary_key.
Когда это всплывает на практике:
- Контракт делает
emplace + eraseодной сущности в одной транзакции (anti-pattern, но встречается). - Несколько коротких транзакций над одной сущностью пакуются производителем блоков в один блок (массовый импорт, auto-approve, bulk-операции, миграции, batch-сценарии).
Workaround на стороне приложения:
- Разнести
emplaceиeraseпо разным блокам — в seed/bulk-сценариях добавить sleep ≥ 1 ×block_time(≈500 мс для EOSIO) между транзакциями. В production-flow ручного approve'а это не нужно — между нажатием кнопки и второй транзакцией всегда проходит более одного блока. - Если транзиентные пары — часть нормального дизайна контракта, рассмотреть переписывание на
modify(status="approved")вместоerase: SHiP корректно эмиттит дельту изменения статуса.
Системный фикс (планируется):
Drop-in патч chainbase + state_history_plugin в нашем форке ~/coopos (uENOSIO) — отдельный список _transient_removed_values в undo_index + паковка пар [present:true, present:false] в table_delta_v0 без изменения wire-формата. Wire-совместимо со всеми существующими SHiP-клиентами; после деплоя форка парсер автоматически начнёт видеть транзиентные пары без изменений в коде. Технические детали и план работ — в задаче f3-shipchainbase-patch-vidimost-tranzientnykh-delt-createremove-vnutri-bloka.md проекта parser2 в _blago.
Разработка
Структура монорепы
.
├── packages/
│ ├── parser/ # @coopenomics/parser2 — ядро + CLI + ParserClient
│ └── ship-reader/ # @coopenomics/coopos-ship-reader — SHiP WS клиент
├── docs/ # redis taxonomy, disaster recovery
├── examples/ # пример verifier-like подписчика
└── .github/workflows/ # CI + Release
Установка зависимостей
pnpm install
Сборка
pnpm build # build всех пакетов
pnpm --filter @coopenomics/parser2 build # только parser
Тесты
pnpm test # unit во всех пакетах
pnpm --filter @coopenomics/parser2 test:unit
pnpm --filter @coopenomics/parser2 test:integration # нужен Docker для блокчейн-ноды
Текущее покрытие:
@coopenomics/parser2: 229 unit-тестов + integration (xtrim/backpressure на реальном Redis)@coopenomics/coopos-ship-reader: 54 unit-теста
Бенчмарк
pnpm --filter @coopenomics/coopos-ship-reader bench
Замеряет throughput wharfkit-десериализатора — используется для контроля регрессий перформанса между версиями.
Ветки и релизы
dev— рабочая ветка разработки; push / PR сюда запускает CI (lint, typecheck, unit, integration, build)main— релизная ветка; merge изdevавтоматически публикует через Lerna
Процесс релиза через Lerna (independent versioning):
# 1. В dev сделать изменения, закоммитить через blago flow
# 2. Посмотреть какие пакеты будут опубликованы:
pnpm changed
# 3. Поднять версии интерактивно (Lerna сам предложит семвер по коммитам):
pnpm release
# → изменит packages/*/package.json, создаст тег(и), закоммитит + запушит
# (локально можно также dry-run: pnpm release:version --no-push --no-git-tag-version)
# 4. Смёрджить dev → main (обычным PR)
# 5. GitHub Actions release.yml автоматически:
# • lerna publish from-package — опубликует в npm только пакеты с новой версией
# • changelogithub — создаст GitHub Release с changelog
# • docker build-push — опубликует образ в docker.io/dicoop/parser2:<version>
Альтернативно, если lerna version не использовалась — можно вручную поднять версию в packages/*/package.json, смёрджить в main, и workflow опубликует.
Лицензия
MIT. См. LICENSE и NOTICE для атрибуции third-party компонентов.