92811c258f
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>
403 lines
22 KiB
Markdown
403 lines
22 KiB
Markdown
# Coopenomics Parser
|
||
|
||
Универсальный индексер блокчейнов EOSIO / Antelope: читает блоки из State History Plugin (SHiP) по WebSocket, декодирует actions и дельты таблиц с учётом исторических ABI, публикует унифицированный поток событий в Redis Stream. Потребители получают события через `ParserClient` с single-active-consumer lock'ом, recovery после сбоев и dead-letter для poison-messages.
|
||
|
||
[](https://github.com/coopenomics/parser2/actions/workflows/ci.yml)
|
||
[](https://www.npmjs.com/package/@coopenomics/parser2)
|
||
[](LICENSE)
|
||
|
||
## Зачем
|
||
|
||
SHiP отдаёт сырые бинарные блоки — чтобы превратить их в прикладной поток событий, нужно: держать кэш ABI за каждый блок (контракты обновляют ABI, старые блоки декодируются старой версией), обрабатывать форки и переподключения, распределять нагрузку между подписчиками. Этот проект закрывает весь этот слой и отдаёт вам единый стрим `action` / `delta` / `native-delta` / `fork` событий.
|
||
|
||
## Пакеты монорепы
|
||
|
||
| Пакет | Описание | npm |
|
||
|:---|:---|:---|
|
||
| [`@coopenomics/parser2`](packages/parser2) | Ядро индексера (Parser) + подписочный клиент (ParserClient), CLI, observability | [`@coopenomics/parser2`](https://www.npmjs.com/package/@coopenomics/parser2) |
|
||
| [`@coopenomics/coopos-ship-reader`](packages/ship-reader) | Низкоуровневый WebSocket SHiP клиент с поддержкой 24 нативных системных таблиц | [`@coopenomics/coopos-ship-reader`](https://www.npmjs.com/package/@coopenomics/coopos-ship-reader) |
|
||
|
||
## Быстрый старт
|
||
|
||
### 1. Установка
|
||
|
||
```bash
|
||
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`:
|
||
|
||
```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 (после глобальной установки пакета):
|
||
|
||
```bash
|
||
parser start --config parser.config.yaml
|
||
```
|
||
|
||
Парсер подключится к SHiP, начнёт читать блоки с последней checkpoint-позиции (или с head если запускается впервые), и публиковать события в Redis stream `ce:parser2:<chainId>:events`.
|
||
|
||
### 3. Подписка на события из приложения
|
||
|
||
```typescript
|
||
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-утилиты
|
||
|
||
```bash
|
||
# Посмотреть все зарегистрированные подписки и их отставание
|
||
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`](https://hub.docker.com/r/dicoop/parser2) — multi-arch (`linux/amd64`, `linux/arm64`).
|
||
|
||
```bash
|
||
docker pull dicoop/parser2:latest
|
||
# или конкретная версия
|
||
docker pull dicoop/parser2:1.0.1
|
||
```
|
||
|
||
### Минимальный запуск
|
||
|
||
Парсер читает конфигурацию из YAML-файла, путь передаётся через `--config`:
|
||
|
||
```bash
|
||
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:
|
||
|
||
```yaml
|
||
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-файлом:
|
||
|
||
```yaml
|
||
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`, но можно использовать любую другую команду:
|
||
|
||
```bash
|
||
# Посмотреть подписки в прод-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)
|
||
└──────────────┘
|
||
```
|
||
|
||
Подробности:
|
||
- [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 внутри одного блока
|
||
|
||
Парсер **не увидит** жизненный цикл строки таблицы, у которой `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
|
||
```
|
||
|
||
### Установка зависимостей
|
||
|
||
```bash
|
||
pnpm install
|
||
```
|
||
|
||
### Сборка
|
||
|
||
```bash
|
||
pnpm build # build всех пакетов
|
||
pnpm --filter @coopenomics/parser2 build # только parser
|
||
```
|
||
|
||
### Тесты
|
||
|
||
```bash
|
||
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-теста
|
||
|
||
### Бенчмарк
|
||
|
||
```bash
|
||
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):
|
||
|
||
```bash
|
||
# 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](LICENSE) и [NOTICE](packages/ship-reader/NOTICE) для атрибуции third-party компонентов.
|