15 Commits

Author SHA1 Message Date
claude 2218c7c1b7 Merge pull request 'feat(parser): авто-переподключение к SHiP с backoff (1.3.0)' (#12) from dev into main
Release / Release (push) Successful in 9m35s
2026-06-05 13:01:02 +00:00
coopops 92811c258f feat(parser): авто-переподключение к SHiP через ReconnectSupervisor с backoff
CI / Lint (pull_request) Successful in 9m51s
CI / Integration tests (pull_request) Successful in 9m46s
CI / Build (pull_request) Successful in 9m25s
CI / Typecheck (pull_request) Successful in 9m27s
CI / Unit tests (pull_request) Successful in 9m33s
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>
2026-06-05 13:00:11 +00:00
claude 3bf7cf9c7a Merge pull request 'docs(readme): backpressure-окно со схемой + конфиг' (#11) from dev into main
Release / Release (push) Successful in 9m21s
2026-06-04 05:44:26 +00:00
coopops be67debf51 docs(readme): раздел про backpressure-окно со схемой + конфиг в примерах
CI / Typecheck (pull_request) Successful in 9m25s
CI / Integration tests (pull_request) Successful in 9m48s
CI / Build (pull_request) Successful in 9m25s
CI / Lint (pull_request) Successful in 9m51s
CI / Unit tests (pull_request) Successful in 9m32s
Описание окна поглощения writer↔consumer (highWater/lowWater гистерезис, пауза
чтения SHiP по XLEN, антидребезг), ASCII-схема качелей backlog, рекомендация
maxmemory+noeviction. Конфиг backpressure добавлен в оба YAML-примера, строка
Backpr. в арх-диаграмму, актуализированы счётчики тестов (221/54).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-04 05:43:48 +00:00
claude 3dbd2c9116 Merge pull request 'feat(backpressure): окно поглощения writer↔consumer (защита Redis RAM на репарсинге)' (#10) from dev into main
Release / Release (push) Successful in 9m33s
2026-06-03 16:08:12 +00:00
coopops 87af7a89d2 feat(backpressure): окно поглощения writer↔consumer — пауза чтения SHiP по backlog стрима
CI / Typecheck (pull_request) Successful in 9m25s
CI / Lint (pull_request) Successful in 9m27s
CI / Unit tests (pull_request) Successful in 9m32s
CI / Build (pull_request) Successful in 9m26s
CI / Integration tests (pull_request) Successful in 13m37s
Проблема: на длинном репарсинге (genesis→head) Parser XADD'ит события в Redis
быстрее, чем consumer (controller) их вычитывает. RAM парсера ограничена
SHiP-окном (max_messages_in_flight), а Redis-стрим — нет → растёт неограниченно
в памяти Redis → OOM. XtrimSupervisor не режет непрочитанное.

Решение: BackpressureGate. Перед обработкой каждого блока Parser проверяет
backlog = XLEN стрима. Если backlog >= highWater — чтение SHiP встаёт на паузу
(блок не ack'ается, SHiP сам тормозит после in-flight окна), ждём пока consumer
сольёт backlog <= lowWater, затем продолжаем. Ничего не теряется — события
остаются в стриме. Нет consumer → стрим дорастает до highWater и writer замирает.

Метрика backlog = XLEN: ровно та память Redis, что бережём; не зависит от числа
групп и версии Redis.

Конфиг backpressure: { enabled (default true), highWater (100k), lowWater (50k),
pollMs (200) }. По умолчанию включено — защита репарсинга из коробки.

Также: добавлен del в IRedisClient (латентный type-error в xtrim-интеграционном
тесте из прошлого коммита — tsc по тестам не гонялся в release-пуше).

Проверено на реальном Redis (integration-тест в CI): пауза при backlog>=highWater,
возобновление после слива. Unit: 221 зелёных, typecheck чистый. parser2 1.1.3→1.2.0.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-03 16:07:45 +00:00
claude 19a905d54b Merge pull request 'fix(xtrim): не триммить un-acked pending + null-гард parseStreamEntries (баг #3)' (#9) from dev into main
Release / Release (push) Successful in 9m34s
2026-06-03 13:02:49 +00:00
coopops c8b7107c3a fix(xtrim): не триммить un-acked pending + гард null-записей в parseStreamEntries
CI / Typecheck (pull_request) Has been cancelled
CI / Unit tests (pull_request) Has been cancelled
CI / Integration tests (pull_request) Has been cancelled
CI / Build (pull_request) Has been cancelled
CI / Lint (pull_request) Has been cancelled
Баг #3: consumer падал навсегда с TypeError "Cannot read properties of null
(reading 'length')" и переставал читать стрим. Два корня:

1. XtrimSupervisor.trim() брал MINID = lastDeliveredId. Но un-acked pending
   записи старше lastDeliveredId (доставлены, не подтверждены) → XTRIM сносил
   собственные pending группы. Теперь нижняя граница trim — самый старый
   un-acked ID из XPENDING (новый метод RedisStore.xpendingMinId), а сравнение
   stream-ID числовое по <ms>-<seq>, не лексикографическое ('1000' < '999').

2. parseStreamEntries падал на null rawFields: при перечитывании PEL
   (XREADGROUP ... 0) Redis отдаёт [id, null] для обрезанных ID. Добавлен
   гард `if (!rawFields) continue` — обрезанные записи пропускаются.

Триггер: consumer отстаёт от писателя (бутстрап с genesis) → растёт pending →
XtrimSupervisor сносит pending → PEL отдаёт null → краш в бесконечном цикле.

Воспроизведено и проверено на реальном Redis (новый integration-тест
xtrim.integration.test.ts, гоняется в CI): trim сохраняет все un-acked pending,
PEL-перечитка не падает на null. Unit: 211 зелёных.

parser2 1.1.2→1.1.3. ship-reader без изменений.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-03 13:02:21 +00:00
claude 71f7a22436 Merge pull request 'fix(block-time): on-chain timestamp блока вместо wall-clock' (#8) from dev into main
Release / Release (push) Successful in 9m36s
2026-06-03 11:57:39 +00:00
coopops 5319209cb2 fix(block-time): брать on-chain timestamp блока из signed_block header, не wall-clock
CI / Typecheck (pull_request) Has been cancelled
CI / Unit tests (pull_request) Has been cancelled
CI / Integration tests (pull_request) Has been cancelled
CI / Build (pull_request) Has been cancelled
CI / Lint (pull_request) Has been cancelled
block_time во всех событиях равнялся времени запуска парсера (new Date()),
одинаковому для всех блоков — point-in-time запросы («был ли ключ активен в
момент T») были невозможны.

ship-reader: decodeBlocksResult теперь декодит r.block (fetch_block:true уже
запрашивал, но байты игнорились) как block_header и берёт timestamp; ShipBlock
получает блок-уровневое поле blockTime; убран new Date()-сид в BlockStream.
parser2: BlockProcessor берёт block.blockTime — общий для всех событий блока,
включая trace-less genesis; фейковый new Date()-фолбэк убран, пустое время
логируется warn'ом.

ship-reader 0.3.0→0.3.1, parser2 1.1.1→1.1.2.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-03 11:57:14 +00:00
claude 30475c055f Merge pull request 'fix(native-delta): декодировать нативные строки ABI вместо JSON.parse' (#7) from dev into main
Release / Release (push) Successful in 10m3s
2026-06-03 09:07:56 +00:00
coopops 39eabb737b fix(native-delta): декодировать нативные строки ABI, а не JSON.parse
CI / Typecheck (pull_request) Has been cancelled
CI / Unit tests (pull_request) Has been cancelled
CI / Integration tests (pull_request) Has been cancelled
CI / Lint (pull_request) Has been cancelled
CI / Build (pull_request) Has been cancelled
deserializeNativeDelta делал JSON.parse(rowRaw), но rowRaw нативных таблиц —
это ABI-сериализованные байты state_history, а не JSON. Каждая native-строка
всегда падала на десериализации, а BlockProcessor глотал ошибку в catch {} →
поток native-delta событий был всегда пустой (permission/account/resource_*).

Фикс:
- WharfkitDeserializer.deserializeNativeDelta декодит через Serializer.decode
  с ship-ABI (type = имя таблицы, это variant *_v0) и распаковывает variant
  [typeName, row] → плоскую строку.
- Deserializer/streamNativeDeltas/ShipClient прокидывают ship-ABI до метода
  (ShipClient.abi getter); адаптер parser2 передаёт this.client.abi.
- BlockProcessor больше не глушит ошибку молча — warn с block_num/table/err.
- Тесты native-rows переписаны под ABI-байты + регресс-гард «JSON теперь бросает».

Проверено на live mono-ai-1: 30/30 native-дельт декодятся (было 0/90).
ship-reader 0.2.0→0.3.0 (breaking: +abi в Deserializer/streamNativeDeltas),
parser2 1.1.0→1.1.1 (багфикс, публичный API не менялся).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-03 09:06:13 +00:00
claude d71a2d696d release: синк release.yml с main (без docker, Gitea API, NPM_TOKEN) (#6) 2026-05-28 11:57:17 +00:00
claude f6360132f0 release: убрать docker job + создавать релиз через Gitea API (#5)
Release / Release (push) Successful in 9m21s
2026-05-28 09:53:20 +00:00
claude 16bd87cca1 release: публикация npm через NODE_AUTH_TOKEN (Gitea не выдаёт OIDC) (#4)
Release / Release (push) Successful in 9m23s
Release / Build & push Docker image (push) Has been skipped
Хотфикс release.yml: NPM auth через NODE_AUTH_TOKEN=NPM_TOKEN, снят provenance и id-token: write — Gitea Actions не выдаёт OIDC id-token, поэтому trusted publishing/provenance работать не могут.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-28 05:56:17 +00:00
33 changed files with 1326 additions and 175 deletions
+32 -57
View File
@@ -5,18 +5,16 @@ name: Release
# существующим git-тегом v<version>) — тогда:
# 1. Публикуем workspace-пакеты в npm
# 2. Создаём git-тег v<version> и пушим его
# 3. Создаём GitHub Release с changelog'ом (changelogithub)
# 4. Собираем и пушим Docker-образ в docker.io (Docker Hub)
# 3. Создаём Gitea Release с changelog'ом через Gitea API
# Если версия не менялась — workflow no-op (без ошибки).
#
# Требуемые secrets в настройках репозитория:
# NPM_TOKEN — npm access token с правами publish на @coopenomics/*
# (Gitea Actions не выдаёт OIDC id-token, поэтому npm
# trusted publishing/provenance здесь недоступны —
# используем классический auth через NODE_AUTH_TOKEN)
# DOCKERHUB_USERNAME — Docker Hub логин с доступом к организации dicoop
# DOCKERHUB_TOKEN — PAT с правами Read,Write из hub.docker.com/settings/security
# GITHUB_TOKEN — выдаётся автоматически GitHub Actions
# NPM_TOKEN — npm access token с правами publish на @coopenomics/*
# (Gitea Actions не выдаёт OIDC id-token, поэтому npm
# trusted publishing/provenance здесь недоступны —
# используем классический auth через NODE_AUTH_TOKEN)
# GITHUB_TOKEN — выдаётся автоматически Gitea Actions; используется для
# создания релиза через ${GITHUB_SERVER_URL}/api/v1/...
on:
push:
@@ -27,7 +25,7 @@ concurrency:
cancel-in-progress: false
permissions:
contents: write # tag + GitHub Release
contents: write # tag + Gitea Release
jobs:
release:
@@ -39,7 +37,7 @@ jobs:
steps:
- uses: actions/checkout@v4
with:
fetch-depth: 0 # нужно для git tag проверки и changelogithub
fetch-depth: 0 # нужно для git tag проверки и git log changelog
- uses: pnpm/action-setup@v4
with: { version: '10' }
- uses: actions/setup-node@v4
@@ -90,52 +88,29 @@ jobs:
git tag "v${{ steps.ver.outputs.version }}"
git push origin "v${{ steps.ver.outputs.version }}"
- name: Create GitHub Release
# changelogithub знает только github.com и падает на git.coopenomics.world
# ("Can not parse GitHub repo from url ..."). Создаём релиз напрямую через
# Gitea REST API: GITHUB_SERVER_URL и GITHUB_REPOSITORY проставляются
# раннером, GITHUB_TOKEN имеет права contents:write на тот же репо.
# Changelog — git log между предыдущим v-тегом и текущим.
- name: Create Gitea Release
if: steps.ver.outputs.released == 'true'
run: npx changelogithub
run: |
VERSION="${{ steps.ver.outputs.version }}"
PREV=$(git tag --list 'v*' --sort=-v:refname | grep -vx "v${VERSION}" | head -1)
if [ -n "$PREV" ]; then
CHANGELOG=$(git log "${PREV}..v${VERSION}" --pretty=format:'- %s')
else
CHANGELOG=$(git log "v${VERSION}" --pretty=format:'- %s')
fi
BODY=$(jq -n \
--arg tag "v${VERSION}" \
--arg body "$CHANGELOG" \
'{tag_name:$tag, name:$tag, body:$body, draft:false, prerelease:false}')
curl -fsS -X POST \
-H "Authorization: token ${GITHUB_TOKEN}" \
-H "Content-Type: application/json" \
-d "$BODY" \
"${GITHUB_SERVER_URL}/api/v1/repos/${GITHUB_REPOSITORY}/releases"
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
docker:
name: Build & push Docker image
runs-on: ubuntu-22.04
needs: release
if: needs.release.outputs.released == 'true'
steps:
- uses: actions/checkout@v4
with:
ref: v${{ needs.release.outputs.version }}
- uses: pnpm/action-setup@v4
with: { version: '10' }
- uses: actions/setup-node@v4
with: { node-version: '20', cache: 'pnpm' }
- run: pnpm install --frozen-lockfile
- run: pnpm --filter @coopenomics/coopos-ship-reader build
- run: pnpm --filter @coopenomics/parser2 build
- uses: docker/setup-buildx-action@v3
# Docker Hub auth: требует secrets DOCKERHUB_USERNAME + DOCKERHUB_TOKEN
# (Access Token из hub.docker.com/settings/security с правами Read,Write)
- uses: docker/login-action@v3
with:
username: ${{ secrets.DOCKERHUB_USERNAME }}
password: ${{ secrets.DOCKERHUB_TOKEN }}
- name: Docker meta
id: meta
uses: docker/metadata-action@v5
with:
images: dicoop/parser2
tags: |
type=semver,pattern={{version}},value=v${{ needs.release.outputs.version }}
type=semver,pattern={{major}}.{{minor}},value=v${{ needs.release.outputs.version }}
type=raw,value=latest
- uses: docker/build-push-action@v5
with:
context: .
file: Dockerfile
push: true
tags: ${{ steps.meta.outputs.tags }}
labels: ${{ steps.meta.outputs.labels }}
cache-from: type=gha
cache-to: type=gha,mode=max
platforms: linux/amd64,linux/arm64
+65 -3
View File
@@ -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 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@coopenomics/parser2",
"version": "1.1.0",
"version": "1.3.0",
"description": "Universal EOSIO/Antelope SHiP-to-Redis blockchain indexer (parser)",
"license": "MIT",
"author": "Coopenomics contributors",
+23 -2
View File
@@ -37,10 +37,13 @@ interface IRedisClient {
count: string, countVal: number,
block: string, blockMs: number,
streams: string, stream: string, id: string,
): Promise<Array<[string, Array<[string, string[]]>]> | null>
): Promise<Array<[string, Array<[string, string[] | null]>]> | null>
// XPENDING key group (summary): [count, min-id, max-id, [[consumer, count],...]|null]
xpending(key: string, group: string): Promise<[number, string | null, string | null, Array<[string, string]> | null] | null>
xrange(key: string, start: string, end: string, count: string, countVal: number): Promise<Array<[string, string[]]>>
xrevrange(key: string, end: string, start: string, count: string, countVal: number): Promise<Array<[string, string[]]>>
xlen(key: string): Promise<number>
del(...keys: string[]): Promise<number>
xdel(key: string, ...ids: string[]): Promise<number>
xack(stream: string, group: string, id: string): Promise<number>
zadd(key: string, score: number, member: string): Promise<number>
@@ -94,9 +97,14 @@ return 0
* Redis возвращает: [[id, [field1, val1, field2, val2, ...]], ...]
* Мы конвертируем в: [{id, fields: {field1: val1, ...}}, ...]
*/
function parseStreamEntries(raw: Array<[string, string[]]>): StreamMessage[] {
function parseStreamEntries(raw: Array<[string, string[] | null]>): StreamMessage[] {
const messages: StreamMessage[] = []
for (const [msgId, rawFields] of raw) {
// Обрезанные/удалённые записи: при чтении PEL (XREADGROUP ... 0) Redis
// отдаёт [id, null] для ID, которых уже нет в стриме (XTRIM/XDEL их снёс).
// rawFields === null → .length бросает TypeError и убивает consumer навсегда.
// Пропускаем — запись недоступна, восстановить нечего.
if (!rawFields) continue
const fields: Record<string, string> = {}
// rawFields — плоский массив [key, val, key, val, ...], шагаем по 2
for (let i = 0; i + 1 < rawFields.length; i += 2) {
@@ -246,6 +254,19 @@ export class IoRedisStore implements RedisStore {
await this.client.xack(stream, group, id)
}
/**
* XPENDING stream group (summary-форма) → самый старый un-acked ID в PEL группы.
* null если pending нет. Нужен для безопасного XTRIM: триммить можно только
* НИЖЕ этого ID, иначе снесём собственные недоставленные/неподтверждённые записи.
*/
async xpendingMinId(stream: string, group: string): Promise<string | null> {
const res = await this.client.xpending(stream, group)
if (!Array.isArray(res)) return null
const [count, minId] = res
if (!count || !minId) return null
return String(minId)
}
/** ZADD key score member. */
async zadd(key: string, score: number, member: string): Promise<void> {
await this.client.zadd(key, score, member)
@@ -86,6 +86,6 @@ export class ShipReaderAdapter implements ChainClient {
* содержит hardcoded-схемы нативных таблиц (permission, account, …).
*/
deserializeNativeDelta(delta: ShipDelta): ShipNativeDeltaEvent {
return this.client.deserializer.deserializeNativeDelta(delta)
return this.client.deserializer.deserializeNativeDelta(delta, this.client.abi)
}
}
+6
View File
@@ -40,6 +40,12 @@ export interface ParserOptions {
abiFallback?: 'rpc-current' | 'fail'
/** XtrimSupervisor: интервал проверки и включение/отключение автообрезки стрима. */
xtrim?: { intervalMs?: number; enabled?: boolean }
/**
* Backpressure-окно: пауза чтения SHiP когда backlog стрима (XLEN) достигает
* highWater, возобновление при сливе до lowWater. Защита Redis RAM на
* репарсинге, когда consumer медленнее писателя. По умолчанию включено.
*/
backpressure?: { enabled?: boolean; highWater?: number; lowWater?: number; pollMs?: number }
/** ReconnectSupervisor: максимум попыток и backoff-таблица в секундах. */
reconnect?: { maxAttempts?: number; backoffSeconds?: number[] }
/** Десериализатор ABI-данных. Единственный вариант — 'wharfkit'. */
+20
View File
@@ -82,6 +82,26 @@ export const configSchema = {
},
additionalProperties: false,
},
/**
* Backpressure: окно поглощения между писателем (SHiP) и consumer'ом.
* Когда backlog стрима (XLEN) ≥ highWater — чтение блоков встаёт на паузу,
* пока consumer не сольёт backlog ≤ lowWater. Бережёт Redis RAM на длинном
* репарсинге. Ничего не теряет — события остаются в стриме.
*/
backpressure: {
type: 'object',
properties: {
/** Включить/отключить. По умолчанию true. */
enabled: { type: 'boolean', default: true },
/** Верхняя граница backlog (XLEN) — пауза чтения. По умолчанию 100000. */
highWater: { type: 'number', default: 100000 },
/** Нижняя граница — возобновление. Должна быть < highWater. По умолчанию 50000. */
lowWater: { type: 'number', default: 50000 },
/** Интервал опроса XLEN во время паузы, мс. По умолчанию 200. */
pollMs: { type: 'number', default: 200 },
},
additionalProperties: false,
},
/** ReconnectSupervisor: поведение при разрыве SHiP соединения. */
reconnect: {
type: 'object',
@@ -0,0 +1,100 @@
/**
* Backpressure-окно между чтением блокчейна (писатель) и consumer'ом (читатель).
*
* Проблема: Parser читает блоки из SHiP и XADD'ит события в Redis-стрим быстрее,
* чем ParserClient (controller) их вычитывает. На длинном репарсинге (genesis →
* head, миллионы блоков) стрим растёт в ПАМЯТИ REDIS неограниченно → Redis OOM.
* RAM самого парсера ограничена SHiP-окном (max_messages_in_flight), а вот
* Redis-стрим — нет. XtrimSupervisor не спасает: он не режет непрочитанное.
*
* Решение: перед обработкой каждого блока проверяем backlog = XLEN стрима.
* Если backlog >= highWater — СТАВИМ чтение на паузу (не ack'аем блок, SHiP
* перестаёт слать) и ждём пока consumer сольёт backlog до lowWater. Потом
* продолжаем. Ничего не теряем: события лежат в стриме, мы лишь не доливаем.
*
* Если consumer вообще не подключён — backlog растёт до highWater и писатель
* замирает, держа Redis в пределах ~highWater событий.
*
* Почему XLEN, а не lag группы: XLEN — это ровно тот объём, что занят в Redis
* RAM (единственный защищаемый ресурс). Не зависит от числа/наличия групп и
* версии Redis (lag появился в 7.0). Консервативен: считает и непрочитанное,
* и ещё не обрезанное XtrimSupervisor'ом — пауза чуть раньше, это безопасно.
*/
import type { RedisStore } from '../ports/RedisStore.js'
import { rootLogger, type Logger } from '../logger.js'
const sleep = (ms: number): Promise<void> => new Promise(resolve => setTimeout(resolve, ms))
export interface BackpressureGateOpts {
redis: RedisStore
/** Стрим, backlog которого ограничиваем (ce:parser2:<chainId>:events). */
stream: string
/** Достигли этого XLEN — пауза чтения SHiP. */
highWater: number
/** Возобновляем чтение когда backlog слился до этого XLEN. < highWater. */
lowWater: number
/** Интервал опроса XLEN во время паузы, мс. По умолчанию 200. */
pollMs?: number
/** Логгер (по умолчанию rootLogger). */
logger?: Logger
}
export class BackpressureGate {
private readonly redis: RedisStore
private readonly stream: string
private readonly highWater: number
private readonly lowWater: number
private readonly pollMs: number
private readonly log: Logger
private pausedNow = false
constructor(opts: BackpressureGateOpts) {
if (!Number.isFinite(opts.highWater) || opts.highWater <= 0) {
throw new Error(`BackpressureGate: highWater must be > 0, got ${opts.highWater}`)
}
if (!Number.isFinite(opts.lowWater) || opts.lowWater < 0 || opts.lowWater >= opts.highWater) {
throw new Error(`BackpressureGate: lowWater must be in [0, highWater), got ${opts.lowWater} (highWater=${opts.highWater})`)
}
this.redis = opts.redis
this.stream = opts.stream
this.highWater = opts.highWater
this.lowWater = opts.lowWater
this.pollMs = opts.pollMs ?? 200
this.log = (opts.logger ?? rootLogger).child({ component: 'BackpressureGate', stream: opts.stream })
}
/** true пока чтение блоков стоит на паузе (для метрик/health). */
get isPaused(): boolean {
return this.pausedNow
}
/**
* Блокирует выполнение пока backlog (XLEN стрима) >= highWater.
* Возвращается когда backlog опустился < lowWater ИЛИ shouldStop() стал true
* (graceful shutdown). Если backlog уже < highWater — возвращается сразу,
* не делая лишних RTT после первого XLEN.
*
* @param shouldStop предикат прерывания (обычно () => this.stopSignal)
*/
async waitForCapacity(shouldStop: () => boolean = () => false): Promise<void> {
let backlog = await this.redis.xlen(this.stream)
if (backlog < this.highWater) return
this.pausedNow = true
this.log.warn(
{ backlog, highWater: this.highWater, lowWater: this.lowWater },
'backpressure: чтение блоков на паузе — consumer отстаёт, ждём слива стрима',
)
try {
while (!shouldStop()) {
await sleep(this.pollMs)
backlog = await this.redis.xlen(this.stream)
if (backlog <= this.lowWater) break
}
} finally {
this.pausedNow = false
}
this.log.info({ backlog, lowWater: this.lowWater }, 'backpressure: backlog слит, чтение возобновлено')
}
}
+17 -4
View File
@@ -29,6 +29,7 @@ import { computeEventId } from '../events/eventId.js'
import type { AbiBootstrapper } from '../abi/AbiBootstrapper.js'
import type { AbiStore } from '../abi/AbiStore.js'
import type { ChainClient } from '../ports/ChainClient.js'
import { rootLogger, type Logger } from '../logger.js'
interface BlockProcessorOptions {
/** Идентификатор цепи — проставляется в каждое событие. */
@@ -68,6 +69,7 @@ export class BlockProcessor {
private abiBootstrapper: AbiBootstrapper
private abiStore: AbiStore
private chainClient: ChainClient
private log: Logger
constructor(opts: BlockProcessorOptions) {
this.chainId = opts.chainId
@@ -75,6 +77,7 @@ export class BlockProcessor {
this.abiBootstrapper = opts.abiBootstrapper
this.abiStore = opts.abiStore
this.chainClient = opts.chainClient
this.log = rootLogger.child({ component: 'BlockProcessor', chain_id: opts.chainId })
this.queue = new PQueue({ concurrency: 1 })
}
@@ -94,8 +97,13 @@ export class BlockProcessor {
const blockNum = block.thisBlock.blockNum
const blockId = block.thisBlock.blockId
// blockTime берём из первой трассировки; если трассировок нет — текущее время
const blockTime = block.traces[0]?.blockTime ?? new Date().toISOString()
// block_time — on-chain время блока (signed_block.timestamp), общее для всех событий
// блока, включая trace-less блоки (genesis). НЕ wall-clock: подставлять new Date()
// нельзя — это ломает point-in-time запросы. Пустое = block не запрошен у SHiP.
const blockTime = block.blockTime
if (!blockTime) {
this.log.warn({ block_num: blockNum }, 'block_time missing from ship block (fetch_block disabled?)')
}
// ── Фаза 1: Action traces ─────────────────────────────────────────────────
for (const trace of block.traces) {
@@ -237,8 +245,13 @@ export class BlockProcessor {
present: native.present,
}
nativeDeltaEvents.push({ ...partial, event_id: computeEventId(partial) })
} catch {
// Ошибки в отдельных нативных дельтах не должны прерывать весь блок
} catch (err) {
// Ошибки в отдельных нативных дельтах не должны прерывать весь блок,
// но молча терять их нельзя — иначе регрессии десериализации невидимы.
this.log.warn(
{ block_num: blockNum, table: delta.name, err: err instanceof Error ? err.message : String(err) },
'native delta deserialize failed',
)
}
}
}
+55 -34
View File
@@ -6,11 +6,16 @@ import { WorkerPool } from '../workers/WorkerPool.js'
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 {
@@ -20,6 +25,8 @@ export class Parser {
private workerPool: WorkerPool | null = null
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
@@ -66,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 })
@@ -81,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,
@@ -100,38 +105,54 @@ export class Parser {
this.xtrimSupervisor.start()
}
const streamOpts = {
startBlock: havePositions[0]?.blockNum ?? 0,
havePositions,
// Backpressure-окно: пауза чтения SHiP когда backlog стрима упирается в
// highWater, чтобы Redis не рос неограниченно на репарсинге. По умолчанию вкл.
if (this.opts.backpressure?.enabled !== false) {
const bp = this.opts.backpressure
this.backpressureGate = new BackpressureGate({
redis: this.redis,
stream: eventsStream,
highWater: bp?.highWater ?? 100_000,
lowWater: bp?.lowWater ?? 50_000,
...(bp?.pollMs !== undefined ? { pollMs: bp.pollMs } : {}),
})
}
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)
},
})
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) {
+34 -8
View File
@@ -7,9 +7,15 @@
* Стратегия MINID:
* Вместо хранения фиксированного числа записей (MAXLEN), мы сохраняем все
* записи, которые ещё не подтверждены (pending) хотя бы одной consumer group.
* minId = min(lastDeliveredId всех групп с pending > 0).
* minId = min(самый старый un-acked pending ID всех групп с pending > 0).
* XTRIM stream MINID minId удаляет всё с ID < minId.
*
* ВАЖНО: триммить по lastDeliveredId НЕЛЬЗЯ — un-acked pending записи старше
* lastDeliveredId (доставлены, но не подтверждены: их ID < last-delivered).
* XTRIM MINID lastDeliveredId снёс бы собственные pending группы → consumer при
* перечитывании PEL получает [id, null] и падает. Поэтому нижняя граница trim —
* самый старый un-acked ID из XPENDING, а не lastDeliveredId.
*
* Это гарантирует, что ни один consumer не потеряет сообщения при trim:
* группа с отставанием «тормозит» trim, пока не догонит.
*
@@ -18,6 +24,19 @@
import type { RedisStore } from '../ports/RedisStore.js'
/**
* Числовое сравнение Redis Stream ID формата `<ms>-<seq>`.
* Строковое сравнение (a < b) неверно: '1000-0' < '999-0' лексикографически
* true, хотя по времени 1000 > 999. Возвращает <0, 0, >0.
*/
function compareStreamIds(a: string, b: string): number {
const [aMs, aSeq = '0'] = a.split('-')
const [bMs, bSeq = '0'] = b.split('-')
const ms = Number(aMs) - Number(bMs)
if (ms !== 0) return ms
return Number(aSeq) - Number(bSeq)
}
export interface XtrimSupervisorOpts {
redis: RedisStore
/** Имя стрима для очистки (обычно ce:parser2:<chainId>:events). */
@@ -60,8 +79,8 @@ export class XtrimSupervisor {
* Один цикл очистки:
* 1. Получаем список consumer groups через XINFO GROUPS.
* 2. Фильтруем группы у которых есть pending сообщения (pending > 0).
* 3. Находим минимальный lastDeliveredId среди таких групп.
* 4. XTRIM stream MINID minId — удаляем всё старее этого ID.
* 3. Для каждой берём самый старый un-acked pending ID (XPENDING).
* 4. minId = числовой минимум этих ID; XTRIM stream MINID minId.
*
* Если pending-групп нет — trim не делается (всё подтверждено).
* Если стрим не существует или XInfo бросает — тихо игнорируем (best-effort).
@@ -75,12 +94,19 @@ export class XtrimSupervisor {
const pendingGroups = groups.filter(g => g.pending > 0)
if (pendingGroups.length === 0) return
// Наименьший lastDeliveredId = самый отстающий consumer
const minId = pendingGroups
.map(g => g.lastDeliveredId)
.reduce((a, b) => (a < b ? a : b))
// Самый старый un-acked pending ID по каждой группе. НЕ lastDeliveredId:
// он выше pending, trim по нему снёс бы собственные неподтверждённые записи.
const oldestPending: string[] = []
for (const g of pendingGroups) {
const id = await this.redis.xpendingMinId(this.stream, g.name)
if (id) oldestPending.push(id)
}
if (oldestPending.length === 0) return
if (minId) await this.redis.xtrim(this.stream, minId)
// Нижняя граница trim = самый старый un-acked ID среди всех групп
const minId = oldestPending.reduce((a, b) => (compareStreamIds(a, b) <= 0 ? a : b))
await this.redis.xtrim(this.stream, minId)
} catch {
// XTRIM — best-effort: ошибки не должны влиять на основной поток обработки
}
+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)
}
+7
View File
@@ -79,6 +79,13 @@ export interface RedisStore {
/** XACK stream group id — подтверждает обработку, убирает из PEL. */
xack(stream: string, group: string, id: string): Promise<void>
/**
* XPENDING stream group → ID самого старого un-acked сообщения в PEL группы
* (или null если pending нет). Используется XtrimSupervisor: trim допустим
* только НИЖЕ этого ID, иначе un-acked записи будут потеряны.
*/
xpendingMinId(stream: string, group: string): Promise<string | null>
// ── Sorted Set ────────────────────────────────────────────────────────────
/** ZADD key score member. */
@@ -0,0 +1,80 @@
/**
* Integration тест против реального Redis: backpressure-окно.
*
* Проверяем на живом Redis (chain/SHiP не нужны):
* - backlog < highWater → waitForCapacity не паузит, возвращается сразу
* - backlog >= highWater → пауза; после слива стрима до <= lowWater — возврат
*
* Требование: Redis на REDIS_URL (в CI всегда есть). Без Redis — тесты skip.
*/
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
import { IoRedisStore } from '../../src/adapters/IoRedisStore.js'
import { BackpressureGate } from '../../src/core/BackpressureGate.js'
import type { Logger } from '../../src/logger.js'
const REDIS_URL = process.env['REDIS_URL'] ?? 'redis://localhost:6379/14'
const STREAM = '__it_bp__:stream'
const silentLog = {
child: () => silentLog, warn: () => {}, info: () => {}, debug: () => {}, error: () => {},
} as unknown as Logger
const sleep = (ms: number): Promise<void> => new Promise(r => setTimeout(r, ms))
let store: IoRedisStore
let redisUp = false
beforeAll(async () => {
store = new IoRedisStore({ url: REDIS_URL })
try {
await store.connect()
await store.client.del(STREAM)
redisUp = true
} catch {
redisUp = false
}
})
afterAll(async () => {
if (redisUp) {
await store.client.del(STREAM)
await store.quit()
}
})
describe('backpressure — окно на реальном Redis', () => {
it('не паузит когда backlog < highWater', async () => {
if (!redisUp) return
await store.client.del(STREAM)
await store.xadd(STREAM, { n: '1' })
const gate = new BackpressureGate({ redis: store, stream: STREAM, highWater: 5, lowWater: 2, pollMs: 10, logger: silentLog })
await gate.waitForCapacity()
expect(gate.isPaused).toBe(false)
})
it('паузит при backlog >= highWater и возобновляет после слива до <= lowWater', async () => {
if (!redisUp) return
await store.client.del(STREAM)
const ids: string[] = []
for (let i = 0; i < 6; i++) ids.push(await store.xadd(STREAM, { n: String(i) }))
expect(await store.xlen(STREAM)).toBe(6)
const gate = new BackpressureGate({ redis: store, stream: STREAM, highWater: 5, lowWater: 2, pollMs: 10, logger: silentLog })
let resolved = false
const p = gate.waitForCapacity().then(() => { resolved = true })
// даём gate войти в паузу и подождать
await sleep(40)
expect(gate.isPaused).toBe(true)
expect(resolved).toBe(false)
// сливаем backlog: удаляем 4 записи → XLEN=2 <= lowWater
await store.client.xdel(STREAM, ids[0]!, ids[1]!, ids[2]!, ids[3]!)
expect(await store.xlen(STREAM)).toBe(2)
await p // должно разблокироваться
expect(resolved).toBe(true)
expect(gate.isPaused).toBe(false)
})
})
@@ -0,0 +1,103 @@
/**
* Integration тест против реального Redis: регресс-гард на баг #3.
*
* Баг: XtrimSupervisor триммил по lastDeliveredId и сносил собственные un-acked
* pending записи; при перечитывании PEL Redis отдавал [id, null], а
* parseStreamEntries падал на null.length → consumer умирал навсегда.
*
* Здесь воспроизводим точный путь на живом Redis (chain/SHiP не нужны):
* A. trim по lastDeliveredId зануляет pending → xreadGroup(PEL) НЕ падает на null.
* B. XtrimSupervisor.trim() (новая логика) сохраняет все un-acked pending.
*
* Требование: Redis на REDIS_URL (в CI всегда есть). Без Redis — тесты skip.
*/
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
import { IoRedisStore } from '../../src/adapters/IoRedisStore.js'
import { XtrimSupervisor } from '../../src/core/XtrimSupervisor.js'
const REDIS_URL = process.env['REDIS_URL'] ?? 'redis://localhost:6379/14'
const STREAM = '__it_bug3__:stream'
const GROUP = 'g'
let store: IoRedisStore
let redisUp = false
// XtrimSupervisor.trim приватный — дёргаем один цикл напрямую через каст.
function runTrimOnce(sup: XtrimSupervisor): Promise<void> {
return (sup as unknown as { trim(): Promise<void> }).trim()
}
beforeAll(async () => {
store = new IoRedisStore({ url: REDIS_URL })
try {
await store.connect()
await store.client.del(STREAM)
redisUp = true
} catch {
redisUp = false
}
})
afterAll(async () => {
if (redisUp) {
await store.client.del(STREAM)
await store.quit()
}
})
describe('bug #3 — XTRIM/PEL на реальном Redis', () => {
it('xreadGroup(PEL) не падает на обрезанных [id,null] записях', async () => {
if (!redisUp) return
await store.client.del(STREAM)
await store.xgroupCreate(STREAM, GROUP, '0')
for (let i = 0; i < 10; i++) await store.xadd(STREAM, { kind: 'action', n: String(i) })
// читаем 5 → они pending (un-acked); lastDeliveredId = id 5-й записи
const read = await store.xreadGroup(STREAM, GROUP, 'c1', 5, 0, '>')
expect(read).toHaveLength(5)
const lastDelivered = read[4]!.id
// ИМИТАЦИЯ СТАРОГО БАГА: trim по lastDeliveredId сносит первые 4 pending
await store.client.xtrim(STREAM, 'MINID', lastDelivered)
// consumer перечитывает PEL — Redis отдаёт [id,null] для снесённых.
// Не должно бросать (старый parseStreamEntries падал на null.length).
const pel = await store.xreadGroup(STREAM, GROUP, 'c1', 100, 0, '0')
// обрезанные пропущены, выжившие — с полями
expect(pel.every(m => Object.keys(m.fields).length > 0)).toBe(true)
expect(pel.some(m => m.id === lastDelivered)).toBe(true)
})
it('XtrimSupervisor.trim() сохраняет все un-acked pending записи', async () => {
if (!redisUp) return
await store.client.del(STREAM)
await store.xgroupCreate(STREAM, GROUP, '0')
for (let i = 0; i < 10; i++) await store.xadd(STREAM, { kind: 'action', n: String(i) })
const read = await store.xreadGroup(STREAM, GROUP, 'c1', 5, 0, '>')
expect(read).toHaveLength(5)
const oldestPending = read[0]!.id
const sup = new XtrimSupervisor({ redis: store, stream: STREAM, intervalMs: 999_999 })
await runTrimOnce(sup)
// consumer перечитывает PEL — все 5 pending целы, без null, без краша
const pel = await store.xreadGroup(STREAM, GROUP, 'c1', 100, 0, '0')
expect(pel).toHaveLength(5)
expect(pel.every(m => Object.keys(m.fields).length > 0)).toBe(true)
// нижняя граница trim = oldestPending → первый pending на месте
expect(pel[0]!.id).toBe(oldestPending)
})
it('xpendingMinId возвращает самый старый un-acked ID, null без pending', async () => {
if (!redisUp) return
await store.client.del(STREAM)
await store.xgroupCreate(STREAM, GROUP, '0')
expect(await store.xpendingMinId(STREAM, GROUP)).toBeNull()
for (let i = 0; i < 3; i++) await store.xadd(STREAM, { kind: 'action', n: String(i) })
const read = await store.xreadGroup(STREAM, GROUP, 'c1', 3, 0, '>')
expect(await store.xpendingMinId(STREAM, GROUP)).toBe(read[0]!.id)
})
})
@@ -0,0 +1,126 @@
/**
* Тесты BackpressureGate — окна поглощения между писателем и consumer'ом.
*
* Проверяем:
* - валидацию highWater/lowWater в конструкторе
* - мгновенный возврат когда backlog < highWater (без паузы)
* - паузу при backlog >= highWater и возобновление при сливе до <= lowWater
* - isPaused отражает состояние
* - shouldStop() прерывает паузу (graceful shutdown) даже если backlog высок
*/
import { describe, it, expect, vi } from 'vitest'
import { BackpressureGate } from '../../src/core/BackpressureGate.js'
import type { RedisStore } from '../../src/ports/RedisStore.js'
import type { Logger } from '../../src/logger.js'
// Молчаливый логгер — не засоряет вывод тестов.
const silentLog = {
child: () => silentLog,
warn: () => {},
info: () => {},
debug: () => {},
error: () => {},
} as unknown as Logger
function makeRedis(xlen: ReturnType<typeof vi.fn>): RedisStore {
return { xlen } as unknown as RedisStore
}
function makeGate(xlen: ReturnType<typeof vi.fn>, over: Partial<{ highWater: number; lowWater: number; pollMs: number }> = {}): BackpressureGate {
return new BackpressureGate({
redis: makeRedis(xlen),
stream: 's',
highWater: over.highWater ?? 100,
lowWater: over.lowWater ?? 50,
pollMs: over.pollMs ?? 1,
logger: silentLog,
})
}
describe('BackpressureGate — конструктор', () => {
const redis = makeRedis(vi.fn())
it('бросает если highWater <= 0', () => {
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 0, lowWater: 0 })).toThrow(/highWater/)
expect(() => new BackpressureGate({ redis, stream: 's', highWater: -5, lowWater: 0 })).toThrow(/highWater/)
})
it('бросает если lowWater >= highWater', () => {
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: 100 })).toThrow(/lowWater/)
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: 150 })).toThrow(/lowWater/)
})
it('бросает если lowWater < 0', () => {
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: -1 })).toThrow(/lowWater/)
})
it('принимает валидные границы', () => {
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: 50 })).not.toThrow()
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: 0 })).not.toThrow()
})
})
describe('BackpressureGate — waitForCapacity', () => {
it('возвращается сразу когда backlog < highWater (одна проверка XLEN, без паузы)', async () => {
const xlen = vi.fn().mockResolvedValue(50)
const gate = makeGate(xlen)
await gate.waitForCapacity()
expect(xlen).toHaveBeenCalledTimes(1)
expect(gate.isPaused).toBe(false)
})
it('возвращается сразу при backlog = highWater - 1', async () => {
const xlen = vi.fn().mockResolvedValue(99)
const gate = makeGate(xlen)
await gate.waitForCapacity()
expect(xlen).toHaveBeenCalledTimes(1)
})
it('ставит паузу при backlog >= highWater и возобновляет при сливе до <= lowWater', async () => {
// 150 (пауза) → опрос: 150, 70 (>50), 30 (<=50 → break)
const xlen = vi.fn()
.mockResolvedValueOnce(150)
.mockResolvedValueOnce(150)
.mockResolvedValueOnce(70)
.mockResolvedValueOnce(30)
const gate = makeGate(xlen)
await gate.waitForCapacity()
expect(xlen).toHaveBeenCalledTimes(4)
expect(gate.isPaused).toBe(false)
})
it('пауза срабатывает ровно на backlog == highWater (граница >=)', async () => {
const xlen = vi.fn().mockResolvedValueOnce(100).mockResolvedValueOnce(40)
const gate = makeGate(xlen)
await gate.waitForCapacity()
expect(xlen).toHaveBeenCalledTimes(2) // 100 → пауза, 40 → выход
})
it('isPaused = true во время паузы', async () => {
const snapshots: boolean[] = []
let gate: BackpressureGate
const xlen = vi.fn().mockImplementation(() => {
// снимок состояния паузы в момент опроса XLEN
snapshots.push(gate.isPaused)
return Promise.resolve(snapshots.length >= 3 ? 10 : 200)
})
gate = makeGate(xlen)
await gate.waitForCapacity()
// 1-й вызов — до установки паузы (false); внутри цикла — true
expect(snapshots[0]).toBe(false)
expect(snapshots.some(s => s === true)).toBe(true)
expect(gate.isPaused).toBe(false)
})
it('shouldStop() прерывает паузу даже при высоком backlog (graceful shutdown)', async () => {
const xlen = vi.fn().mockResolvedValue(9999) // никогда не сольётся
const gate = makeGate(xlen)
let stop = false
const p = gate.waitForCapacity(() => stop)
// через тик включаем стоп
setTimeout(() => { stop = true }, 5)
await p // не должно зависнуть
expect(gate.isPaused).toBe(false)
})
})
@@ -51,6 +51,7 @@ function makeBlock(blockNum = 1, numTraces = 0, numDeltas = 0): ShipBlock {
head: blockPosition,
lastIrreversible: blockPosition,
prevBlock: null,
blockTime: '2024-06-01T12:00:00.000',
traces: Array.from({ length: numTraces }, (_, i) => ({
account: 'eosio.token',
name: 'transfer',
@@ -175,6 +176,7 @@ describe('BlockProcessor — ABI updates (Story 4.3)', () => {
head: { blockNum: 500, blockId: 'c'.repeat(64) },
lastIrreversible: { blockNum: 500, blockId: 'c'.repeat(64) },
prevBlock: null,
blockTime: '2024-06-01T12:00:00.000',
traces: [{
account: 'eosio',
name: 'setabi',
@@ -214,6 +216,7 @@ describe('BlockProcessor — ABI updates (Story 4.3)', () => {
head: { blockNum: 501, blockId: 'c'.repeat(64) },
lastIrreversible: { blockNum: 501, blockId: 'c'.repeat(64) },
prevBlock: null,
blockTime: '2024-06-01T12:00:00.000',
traces: [{
account: 'eosio',
name: 'setabi',
@@ -253,6 +256,7 @@ describe('BlockProcessor — ABI updates (Story 4.3)', () => {
head: { blockNum: 600, blockId: 'e'.repeat(64) },
lastIrreversible: { blockNum: 600, blockId: 'e'.repeat(64) },
prevBlock: null,
blockTime: '2024-06-01T12:00:00.000',
traces: [],
deltas: [{
name: 'account' as never,
@@ -283,6 +287,7 @@ describe('BlockProcessor — native deltas (Epic 6)', () => {
head: { blockNum: 10, blockId: 'a'.repeat(64) },
lastIrreversible: { blockNum: 10, blockId: 'a'.repeat(64) },
prevBlock: null,
blockTime: '2024-06-01T12:00:00.000',
traces: [],
deltas: nativeTableNames.map(name => ({
name: name as never,
@@ -338,6 +343,7 @@ describe('BlockProcessor — native deltas (Epic 6)', () => {
head: { blockNum: 20, blockId: 'b'.repeat(64) },
lastIrreversible: { blockNum: 20, blockId: 'b'.repeat(64) },
prevBlock: null,
blockTime: '2024-06-01T12:00:00.000',
traces: [
{ account: 'eosio.token', name: 'transfer', authorization: [], actRaw: new Uint8Array([1]), actionOrdinal: 1, creatorActionOrdinal: 0, globalSequence: BigInt(1), receipt: null, contextFree: false, elapsed: 0, console: '', accountRamDeltas: [], blockNum: 20, blockId: 'b'.repeat(64), blockTime: '2024-01-01T00:00:00.000', transactionId: 'c'.repeat(64) },
{ account: 'eosio.token', name: 'transfer', authorization: [], actRaw: new Uint8Array([1]), actionOrdinal: 2, creatorActionOrdinal: 0, globalSequence: BigInt(2), receipt: null, contextFree: false, elapsed: 0, console: '', accountRamDeltas: [], blockNum: 20, blockId: 'b'.repeat(64), blockTime: '2024-01-01T00:00:00.000', transactionId: 'c'.repeat(64) },
@@ -370,6 +376,7 @@ describe('BlockProcessor — native deltas (Epic 6)', () => {
head: { blockNum: 5, blockId: 'd'.repeat(64) },
lastIrreversible: { blockNum: 5, blockId: 'd'.repeat(64) },
prevBlock: null,
blockTime: '2024-06-01T12:00:00.000',
traces: [{ account: 'eosio', name: 'newaccount', authorization: [], actRaw: new Uint8Array([1]), actionOrdinal: 1, creatorActionOrdinal: 0, globalSequence: BigInt(10), receipt: null, contextFree: false, elapsed: 0, console: '', accountRamDeltas: [], blockNum: 5, blockId: 'd'.repeat(64), blockTime: '2024-01-01T00:00:00.000', transactionId: 'e'.repeat(64) }],
deltas: [
{ name: 'contract_row', present: true, rowRaw: new Uint8Array([1]), code: 'eosio', scope: 'eosio', table: 'global', primaryKey: '0' },
@@ -390,6 +397,7 @@ describe('BlockProcessor — native deltas (Epic 6)', () => {
head: { blockNum: 700, blockId: 'f'.repeat(64) },
lastIrreversible: { blockNum: 700, blockId: 'f'.repeat(64) },
prevBlock: null,
blockTime: '2024-06-01T12:00:00.000',
traces: [],
deltas: [{ name: 'account' as never, present: true, rowRaw: new Uint8Array([1, 2, 3]) }],
}
@@ -5,6 +5,8 @@ vi.mock('ioredis', () => {
const mockRedis = {
xadd: vi.fn().mockResolvedValue('1700000000000-0'),
xtrim: vi.fn().mockResolvedValue(5),
xreadgroup: vi.fn().mockResolvedValue(null),
xpending: vi.fn().mockResolvedValue([0, null, null, null]),
zadd: vi.fn().mockResolvedValue(1),
zrangebyscore: vi.fn().mockResolvedValue(['{"version":"eosio::abi/1.0"}']),
zrevrangebyscore: vi.fn().mockResolvedValue(['{"version":"eosio::abi/1.0"}']),
@@ -37,6 +39,34 @@ describe('IoRedisStore', () => {
expect(store.client.xtrim).toHaveBeenCalledWith('mystream', 'MINID', '1234567890000-0')
})
it('xreadGroup skips trimmed entries with null fields (no crash — баг #3)', async () => {
// Redis отдаёт [id, null] для обрезанных/удалённых ID при чтении PEL.
// Раньше rawFields.length на null убивал consumer навсегда.
vi.mocked(store.client.xreadgroup).mockResolvedValueOnce([
['mystream', [
['100-0', ['kind', 'action']],
['101-0', null], // обрезанная запись (XTRIM/XDEL)
['102-0', ['kind', 'delta']],
]],
])
const msgs = await store.xreadGroup('mystream', 'g', 'c', 10, 0, '0')
expect(msgs).toHaveLength(2)
expect(msgs.map(m => m.id)).toEqual(['100-0', '102-0'])
expect(msgs[0]!.fields).toEqual({ kind: 'action' })
})
it('xpendingMinId returns oldest pending id from XPENDING summary', async () => {
vi.mocked(store.client.xpending).mockResolvedValueOnce([5, '90-0', '100-0', [['c', '5']]])
const id = await store.xpendingMinId('mystream', 'g')
expect(id).toBe('90-0')
})
it('xpendingMinId returns null when group has no pending', async () => {
vi.mocked(store.client.xpending).mockResolvedValueOnce([0, null, null, null])
const id = await store.xpendingMinId('mystream', 'g')
expect(id).toBeNull()
})
it('zadd calls redis.zadd with score and member', async () => {
await store.zadd('parser2:abi:eosio', 100, '{"version":"eosio::abi/1.0"}')
expect(store.client.zadd).toHaveBeenCalledWith('parser2:abi:eosio', 100, '{"version":"eosio::abi/1.0"}')
@@ -69,6 +69,7 @@ class FakeRedis implements RedisStore {
xlen = vi.fn(async (): Promise<number> => 0)
xdel = vi.fn(async (): Promise<number> => 0)
xack = vi.fn(async (): Promise<void> => {})
xpendingMinId = vi.fn(async (): Promise<string | null> => null)
zadd = vi.fn(async (key: string, score: number, member: string): Promise<void> => {
if (!this.zsets.has(key)) this.zsets.set(key, new Map())
@@ -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')
})
})
@@ -3,20 +3,28 @@
*
* Проверяем:
* - start/stop идемпотентны
* - trim использует MINID по отстающей группе с pending > 0
* - trim использует MINID по самому старому un-acked pending ID (XPENDING),
* а НЕ по lastDeliveredId (иначе сносит собственные pending — баг #3)
* - числовое сравнение ID (а не лексикографическое)
* - группы без pending не мешают trim'у
* - XTRIM не вызывается если групп нет или ни у одной нет pending
* - ошибки в xinfoGroups не прерывают supervisor
* - XTRIM не вызывается если групп нет, ни у одной нет pending, или pending-ID null
* - ошибки в xinfoGroups/xtrim не прерывают supervisor
*/
import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'
import { XtrimSupervisor } from '../../src/core/XtrimSupervisor.js'
import type { RedisStore, XGroupInfo } from '../../src/ports/RedisStore.js'
function makeRedis(groups: XGroupInfo[] = []): RedisStore {
function makeRedis(
groups: XGroupInfo[] = [],
pendingIds: Record<string, string | null> = {},
): RedisStore {
return {
xinfoGroups: vi.fn().mockResolvedValue(groups),
xtrim: vi.fn().mockResolvedValue(0),
xpendingMinId: vi.fn().mockImplementation((_stream: string, group: string) =>
Promise.resolve(pendingIds[group] ?? null),
),
// остальные методы не используются в supervisor
} as unknown as RedisStore
}
@@ -25,6 +33,7 @@ function makeRedis(groups: XGroupInfo[] = []): RedisStore {
async function flushPromises(): Promise<void> {
await Promise.resolve()
await Promise.resolve()
await Promise.resolve()
}
describe('XtrimSupervisor — lifecycle', () => {
@@ -124,36 +133,95 @@ describe('XtrimSupervisor — trim logic', () => {
sup.stop()
})
it('trims to minimum lastDeliveredId among groups with pending > 0', async () => {
// g1 pending но отстал на 100 — должен стать minId
// g2 pending и догнал до 500 — игнорируется для выбора min
// g3 нет pending — вообще не учитывается
const redis = makeRedis([
{ name: 'g1', pending: 5, lastDeliveredId: '100-0', lag: 5, consumers: 1 },
{ name: 'g2', pending: 2, lastDeliveredId: '500-0', lag: 2, consumers: 1 },
{ name: 'g3', pending: 0, lastDeliveredId: '50-0', lag: 0, consumers: 1 },
])
it('trims to the OLDEST un-acked pending ID (XPENDING), NOT lastDeliveredId', async () => {
// Регресс-гард на баг #3: pending-записи старше lastDeliveredId.
// g1: lastDelivered=100-0, но самый старый un-acked pending = 90-0.
// Триммить надо по 90-0, иначе un-acked 90..99 теряются и consumer падает.
const redis = makeRedis(
[{ name: 'g1', pending: 5, lastDeliveredId: '100-0', lag: 5, consumers: 1 }],
{ g1: '90-0' },
)
const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 })
sup.start()
await vi.advanceTimersByTimeAsync(100)
await flushPromises()
expect(redis.xtrim).toHaveBeenCalledTimes(1)
expect(redis.xtrim).toHaveBeenCalledWith('s', '100-0')
expect(redis.xtrim).toHaveBeenCalledWith('s', '90-0')
// именно НЕ lastDeliveredId
expect(redis.xtrim).not.toHaveBeenCalledWith('s', '100-0')
sup.stop()
})
it('uses lexicographic comparison for stream IDs (which is also numeric-correct for same-length)', async () => {
const redis = makeRedis([
{ name: 'g1', pending: 1, lastDeliveredId: '1000-0', lag: 1, consumers: 1 },
{ name: 'g2', pending: 1, lastDeliveredId: '999-0', lag: 1, consumers: 1 },
])
it('trims to minimum oldest-pending-id among groups with pending > 0', async () => {
// g1 отстал (oldest pending 90-0) → станет minId
// g2 догнал (oldest pending 480-0) → не минимум
// g3 нет pending → не учитывается (xpendingMinId не зовётся)
const redis = makeRedis(
[
{ name: 'g1', pending: 5, lastDeliveredId: '100-0', lag: 5, consumers: 1 },
{ name: 'g2', pending: 2, lastDeliveredId: '500-0', lag: 2, consumers: 1 },
{ name: 'g3', pending: 0, lastDeliveredId: '50-0', lag: 0, consumers: 1 },
],
{ g1: '90-0', g2: '480-0' },
)
const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 })
sup.start()
await vi.advanceTimersByTimeAsync(100)
await flushPromises()
// lexicographic: '1000-0' < '999-0' → выбран '1000-0'
expect(redis.xtrim).toHaveBeenCalledWith('s', '1000-0')
expect(redis.xtrim).toHaveBeenCalledTimes(1)
expect(redis.xtrim).toHaveBeenCalledWith('s', '90-0')
// g3 без pending — XPENDING не запрашивался
expect(redis.xpendingMinId).not.toHaveBeenCalledWith('s', 'g3')
sup.stop()
})
it('uses NUMERIC (not lexicographic) comparison for stream IDs', async () => {
// Лексикографически '1000-0' < '999-0' (неверно). Численно 999 < 1000.
const redis = makeRedis(
[
{ name: 'g1', pending: 1, lastDeliveredId: '1000-0', lag: 1, consumers: 1 },
{ name: 'g2', pending: 1, lastDeliveredId: '999-0', lag: 1, consumers: 1 },
],
{ g1: '1000-0', g2: '999-0' },
)
const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 })
sup.start()
await vi.advanceTimersByTimeAsync(100)
await flushPromises()
// численный минимум = '999-0'
expect(redis.xtrim).toHaveBeenCalledWith('s', '999-0')
sup.stop()
})
it('compares by sequence when ms part is equal', async () => {
const redis = makeRedis(
[
{ name: 'g1', pending: 1, lastDeliveredId: '5-9', lag: 1, consumers: 1 },
{ name: 'g2', pending: 1, lastDeliveredId: '5-2', lag: 1, consumers: 1 },
],
{ g1: '5-9', g2: '5-2' },
)
const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 })
sup.start()
await vi.advanceTimersByTimeAsync(100)
await flushPromises()
expect(redis.xtrim).toHaveBeenCalledWith('s', '5-2')
sup.stop()
})
it('skips trim when pending groups exist but XPENDING returns null (race)', async () => {
// pending>0 в XINFO, но к моменту XPENDING всё подтверждено → null → не триммим
const redis = makeRedis(
[{ name: 'g1', pending: 3, lastDeliveredId: '100-0', lag: 3, consumers: 1 }],
{ g1: null },
)
const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 })
sup.start()
await vi.advanceTimersByTimeAsync(100)
await flushPromises()
expect(redis.xtrim).not.toHaveBeenCalled()
sup.stop()
})
@@ -173,9 +241,10 @@ describe('XtrimSupervisor — trim logic', () => {
})
it('silently swallows errors from xtrim itself', async () => {
const redis = makeRedis([
{ name: 'g1', pending: 1, lastDeliveredId: '100-0', lag: 1, consumers: 1 },
])
const redis = makeRedis(
[{ name: 'g1', pending: 1, lastDeliveredId: '100-0', lag: 1, consumers: 1 }],
{ g1: '95-0' },
)
;(redis.xtrim as ReturnType<typeof vi.fn>).mockRejectedValue(new Error('xtrim failed'))
const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 })
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@coopenomics/coopos-ship-reader",
"version": "0.2.0",
"version": "0.3.1",
"description": "Clean-room SHiP WebSocket client for EOSIO/Antelope blockchains",
"license": "MIT",
"author": "Coopenomics contributors",
+1 -2
View File
@@ -79,7 +79,6 @@ export async function* createBlockStream(
}
let blockNum = opts.startBlock
let blockTime = new Date().toISOString()
while (true) {
const msg = await nextMessage()
@@ -91,7 +90,7 @@ export async function* createBlockStream(
const [type, raw] = decodeResult(new Uint8Array(msg), abi)
if (type !== 'get_blocks_result_v0') continue
const block = decodeBlocksResult(raw, abi, blockNum, '', blockTime)
const block = decodeBlocksResult(raw, abi, blockNum, '')
blockNum = block.thisBlock.blockNum
yield block
+3 -1
View File
@@ -1,3 +1,4 @@
import type { ABI } from '@wharfkit/antelope'
import type { ShipDelta } from './types/ship.js'
import type { NativeDeltaEvent } from './native-tables/index.js'
import { isNativeTableName } from './native-tables/index.js'
@@ -10,9 +11,10 @@ export function filterNativeDeltas(deltas: readonly ShipDelta[]): ShipDelta[] {
export function* streamNativeDeltas(
deltas: readonly ShipDelta[],
deserializer: WharfkitDeserializer,
abi: ABI,
): Generator<NativeDeltaEvent> {
for (const delta of deltas) {
if (!isNativeTableName(delta.name)) continue
yield deserializer.deserializeNativeDelta(delta)
yield deserializer.deserializeNativeDelta(delta, abi)
}
}
+9
View File
@@ -21,6 +21,15 @@ export class ShipClient {
this.deserializer = new WharfkitDeserializer()
}
/**
* SHiP-ABI (state_history), доступен после connect().
* Нужен для десериализации нативных дельт (их строки сериализованы этим ABI).
*/
get abi(): ShipAbi {
if (!this.shipAbi) throw new ShipConnectionError('Call connect() first')
return this.shipAbi
}
async connect(): Promise<void> {
if (this.ws?.readyState === WebSocket.OPEN) return
+19 -2
View File
@@ -43,6 +43,7 @@ interface RawBlocksResult {
last_irreversible: RawBlockPos
this_block: RawBlockPos | null
prev_block: RawBlockPos | null
block: Bytes | null
traces: Bytes | null
deltas: Bytes | null
}
@@ -113,7 +114,7 @@ export function decodeStatusResult(raw: unknown): { chainId: string; head: Block
}
}
export function decodeBlocksResult(raw: unknown, abi: ShipAbi, blockNum: number, blockId: string, blockTime: string): ShipBlock {
export function decodeBlocksResult(raw: unknown, abi: ShipAbi, blockNum: number, blockId: string): ShipBlock {
const r = raw as RawBlocksResult
// ABI decoder возвращает checksum256 как Checksum256 объект (не строку) —
@@ -127,6 +128,22 @@ export function decodeBlocksResult(raw: unknown, abi: ShipAbi, blockNum: number,
? { blockNum: r.prev_block.block_num, blockId: String(r.prev_block.block_id) }
: null
// Время блока берём из on-chain header (signed_block.timestamp), а НЕ из wall-clock.
// r.block — optional<bytes> с сериализованным signed_block, его база — block_header,
// первое поле которого timestamp. Декодим как block_header: читаются только
// header-поля, тело блока (transactions) не десериализуем.
let blockTime = ''
if (r.block && r.block.length > 0) {
try {
const headerDecoded = Serializer.decode({ data: r.block.array, type: 'block_header', abi })
const header = Serializer.objectify(headerDecoded) as { timestamp?: string }
blockTime = header.timestamp ? String(header.timestamp) : ''
} catch {
// Не смогли распарсить header — оставляем пустым, downstream залогирует.
blockTime = ''
}
}
const txTraces = decodeVector<[string, RawTransactionTrace]>(r.traces, 'transaction_trace', abi)
const tableDeltaVariants = decodeVector<[string, RawTableDelta]>(r.deltas, 'table_delta', abi)
@@ -228,5 +245,5 @@ export function decodeBlocksResult(raw: unknown, abi: ShipAbi, blockNum: number,
}
}
return { thisBlock, head, lastIrreversible, prevBlock, traces, deltas }
return { thisBlock, head, lastIrreversible, prevBlock, blockTime, traces, deltas }
}
@@ -5,6 +5,6 @@ import type { NativeDeltaEvent } from '../native-tables/index.js'
export interface Deserializer {
deserializeAction<T = Record<string, unknown>>(trace: ShipTrace, abi: ABI): Action<T>
deserializeContractRow<T = Record<string, unknown>>(delta: ShipDelta, abi: ABI): Delta<T>
deserializeNativeDelta<T = Record<string, unknown>>(delta: ShipDelta): NativeDeltaEvent<T>
deserializeNativeDelta<T = Record<string, unknown>>(delta: ShipDelta, abi: ABI): NativeDeltaEvent<T>
readonly name: 'wharfkit'
}
@@ -57,13 +57,18 @@ export class WharfkitDeserializer implements Deserializer {
}
}
deserializeNativeDelta<T = Record<string, unknown>>(delta: ShipDelta): NativeDeltaEvent<T> {
deserializeNativeDelta<T = Record<string, unknown>>(delta: ShipDelta, abi: ABI): NativeDeltaEvent<T> {
if (!isNativeTableName(delta.name)) {
throw new UnknownNativeTableError(delta.name)
}
const table = delta.name as NativeTableName
try {
const data = JSON.parse(Buffer.from(delta.rowRaw).toString('utf8')) as T
// rowRaw нативных таблиц — это ABI-сериализованные байты state_history,
// а НЕ JSON. Тип строки = имя таблицы: в ship-ABI это variant (*_v0).
const decoded = Serializer.decode({ data: delta.rowRaw, type: table, abi })
const objectified = Serializer.objectify(decoded as ABISerializable)
// variant оформлен как [typeName, row] — распаковываем саму строку.
const data = (Array.isArray(objectified) ? objectified[1] : objectified) as T
const lookup_key = computeLookupKey(table, data as never)
const present: boolean = delta.present
return { present, table, data, lookup_key }
+2
View File
@@ -62,6 +62,8 @@ export interface ShipBlock {
readonly head: BlockPosition
readonly lastIrreversible: BlockPosition
readonly prevBlock: BlockPosition | null
/** On-chain время блока (signed_block.timestamp), ISO-строка. '' если block не запрошен. */
readonly blockTime: string
readonly traces: readonly ShipTrace[]
readonly deltas: readonly ShipDelta[]
}
@@ -1,11 +1,50 @@
import { describe, it, expect } from 'vitest'
import { ABI, Serializer } from '@wharfkit/antelope'
import { WharfkitDeserializer } from '../../src/deserializers/WharfkitDeserializer.js'
import { computeLookupKey } from '../../src/native-tables/index.js'
import { isNativeTableName, NATIVE_TABLE_NAMES } from '../../src/native-tables/types.js'
import { UnknownNativeTableError } from '../../src/errors.js'
import { UnknownNativeTableError, DeserializationError } from '../../src/errors.js'
import type { ShipDelta } from '../../src/types/ship.js'
import type { NativePermissionRow, NativePermissionLinkRow } from '../../src/native-tables/types.js'
/**
* Мини-ABI в стиле state_history: каждая нативная таблица — это variant (*_v0).
* Достаточно полей, нужных computeLookupKey. Кодируем и декодируем одним ABI —
* схема самосогласована, повторять реальный EOSIO целиком не требуется.
*/
function makeShipAbi(): ABI {
return ABI.from({
version: 'eosio::abi/1.1',
types: [],
structs: [
{ name: 'permission_v0', base: '', fields: [
{ name: 'owner', type: 'name' },
{ name: 'name', type: 'name' },
] },
{ name: 'account_v0', base: '', fields: [
{ name: 'name', type: 'name' },
{ name: 'creation_date', type: 'string' },
{ name: 'abi', type: 'string' },
] },
],
actions: [],
tables: [],
variants: [
{ name: 'permission', types: ['permission_v0'] },
{ name: 'account', types: ['account_v0'] },
],
})
}
const SHIP_ABI = makeShipAbi()
/** Дельта с ABI-сериализованным rowRaw (как реально приходит из SHiP). */
function abiDelta(table: string, obj: Record<string, unknown>, present = true): ShipDelta {
const encoded = Serializer.encode({ object: [`${table}_v0`, obj], type: table, abi: SHIP_ABI })
return { name: table, present, rowRaw: encoded.array }
}
/** Дельта с JSON-байтами — раньше «работала» из-за бага JSON.parse, теперь должна падать. */
function jsonDelta(name: string, data: unknown, present = true): ShipDelta {
return {
name,
@@ -64,36 +103,38 @@ describe('computeLookupKey', () => {
describe('WharfkitDeserializer — deserializeNativeDelta', () => {
const deser = new WharfkitDeserializer()
it('deserializes permission delta with correct lookup_key', () => {
const permData: NativePermissionRow = {
owner: 'alice', name: 'active', parent: 'owner',
last_updated: '2024-01-01T00:00:00.000',
auth: { threshold: 1, keys: [], accounts: [], waits: [] },
}
const delta = jsonDelta('permission', permData)
const event = deser.deserializeNativeDelta<NativePermissionRow>(delta)
it('deserializes ABI-encoded permission delta with correct lookup_key', () => {
const delta = abiDelta('permission', { owner: 'alice', name: 'active' })
const event = deser.deserializeNativeDelta<NativePermissionRow>(delta, SHIP_ABI)
expect(event.table).toBe('permission')
expect(event.lookup_key).toBe('alice:active')
expect(event.present).toBe(true)
expect(event.data.owner).toBe('alice')
expect(String(event.data.owner)).toBe('alice')
})
it('unwraps the *_v0 variant into a flat row (not a [name, row] tuple)', () => {
const delta = abiDelta('permission', { owner: 'alice', name: 'active' })
const event = deser.deserializeNativeDelta<NativePermissionRow>(delta, SHIP_ABI)
expect(Array.isArray(event.data)).toBe(false)
expect(String(event.data.name)).toBe('active')
})
it('present field is boolean (not string)', () => {
const delta = jsonDelta('account', { name: 'bob', creation_date: '', abi: '' }, false)
const event = deser.deserializeNativeDelta(delta)
const delta = abiDelta('account', { name: 'bob', creation_date: '2024-01-01T00:00:00.000', abi: '' }, false)
const event = deser.deserializeNativeDelta(delta, SHIP_ABI)
expect(typeof event.present).toBe('boolean')
expect(event.present).toBe(false)
})
it('throws UnknownNativeTableError for unknown table', () => {
const delta: ShipDelta = { name: 'my_custom_table', present: true, rowRaw: new Uint8Array([]) }
expect(() => deser.deserializeNativeDelta(delta)).toThrow(UnknownNativeTableError)
expect(() => deser.deserializeNativeDelta(delta, SHIP_ABI)).toThrow(UnknownNativeTableError)
})
it('all NATIVE_TABLE_NAMES tables deserialize without throwing (smoke test)', () => {
for (const table of NATIVE_TABLE_NAMES) {
const delta = jsonDelta(table, { name: 'test' })
expect(() => deser.deserializeNativeDelta(delta)).not.toThrow()
}
it('throws DeserializationError on non-ABI (JSON) bytes — no silent JSON.parse', () => {
// Регресс-гард на исходный баг: нативные строки декодятся ABI, а не JSON.parse.
// JSON-байты — не валидная ABI-сериализация, метод обязан бросить.
const delta = jsonDelta('permission', { owner: 'alice', name: 'active' })
expect(() => deser.deserializeNativeDelta(delta, SHIP_ABI)).toThrow(DeserializationError)
})
})
@@ -123,6 +123,13 @@ function makeShipAbi(): ABI {
{ name: 'name', type: 'string' },
{ name: 'rows', type: 'row[]' },
] },
// Упрощённый block_header: реальный signed_block начинается с этих байт,
// decodeBlocksResult декодит r.block как 'block_header' и читает timestamp.
// Достаточно первых двух полей — self-consistent encode/decode тем же ABI.
{ name: 'block_header', base: '', fields: [
{ name: 'timestamp', type: 'block_timestamp_type' },
{ name: 'producer', type: 'name' },
] },
],
actions: [],
tables: [],
@@ -186,7 +193,7 @@ describe('ShipProtocol — decodeBlocksResult (no traces / no deltas)', () => {
traces: null,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 100, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 100, 'c'.repeat(64))
expect(block.traces).toHaveLength(0)
expect(block.deltas).toHaveLength(0)
expect(block.thisBlock.blockNum).toBe(100)
@@ -202,7 +209,7 @@ describe('ShipProtocol — decodeBlocksResult (no traces / no deltas)', () => {
traces: null,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 42, 'fallback-id', '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 42, 'fallback-id')
expect(block.thisBlock.blockNum).toBe(42)
expect(block.thisBlock.blockId).toBe('fallback-id')
expect(block.prevBlock).toBeNull()
@@ -219,7 +226,7 @@ describe('ShipProtocol — decodeBlocksResult (no traces / no deltas)', () => {
traces: null,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 100, 'fallback', '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 100, 'fallback')
expect(typeof block.thisBlock.blockId).toBe('string')
expect(typeof block.head.blockId).toBe('string')
expect(typeof block.lastIrreversible.blockId).toBe('string')
@@ -316,7 +323,7 @@ describe('ShipProtocol — decodeBlocksResult (real decoded traces)', () => {
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64))
expect(block.traces).toHaveLength(1)
const t = block.traces[0]!
@@ -346,7 +353,7 @@ describe('ShipProtocol — decodeBlocksResult (real decoded traces)', () => {
traces: tracesBytes,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64))
const r = block.traces[0]!.receipt!
expect(typeof r.receiver).toBe('string')
expect(typeof r.actDigest).toBe('string')
@@ -367,7 +374,7 @@ describe('ShipProtocol — decodeBlocksResult (real decoded traces)', () => {
traces: tracesBytes,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64))
expect(block.traces[0]!.globalSequence).toBe(42n)
})
@@ -382,7 +389,7 @@ describe('ShipProtocol — decodeBlocksResult (real decoded traces)', () => {
traces: tracesBytes,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64))
expect(block.traces[0]!.receipt).toBeNull()
// Без receipt и top-level global_sequence — fallback на 0n
expect(block.traces[0]!.globalSequence).toBe(0n)
@@ -399,7 +406,7 @@ describe('ShipProtocol — decodeBlocksResult (real decoded traces)', () => {
traces: tracesBytes,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64))
expect(block.traces[0]!.actRaw).toBeInstanceOf(Uint8Array)
})
})
@@ -435,7 +442,7 @@ describe('ShipProtocol — decodeBlocksResult (deltas)', () => {
traces: null,
deltas: tableDelta,
}
const block = decodeBlocksResult(raw, abi, 300, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 300, 'c'.repeat(64))
expect(block.deltas).toHaveLength(1)
const d = block.deltas[0]!
expect(d.name).toBe('contract_row')
@@ -470,7 +477,7 @@ describe('ShipProtocol — decodeBlocksResult (deltas)', () => {
traces: null,
deltas: tableDelta,
}
const block = decodeBlocksResult(raw, abi, 300, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 300, 'c'.repeat(64))
expect(block.deltas).toHaveLength(1)
const d = block.deltas[0]!
expect(d.name).toBe('permission')
@@ -499,12 +506,51 @@ describe('ShipProtocol — decodeBlocksResult (deltas)', () => {
traces: null,
deltas: tableDelta,
}
const block = decodeBlocksResult(raw, abi, 300, 'c'.repeat(64), '2024-01-01T00:00:00.000')
const block = decodeBlocksResult(raw, abi, 300, 'c'.repeat(64))
// Некорректная row пропущена без throw
expect(block.deltas).toHaveLength(0)
})
})
describe('ShipProtocol — decodeBlocksResult (block_time из header)', () => {
// Регресс-гард на баг #2: block_time брался из wall-clock (new Date()),
// одинаковый для всех блоков. Теперь — из on-chain signed_block.timestamp.
it('extracts on-chain block_time from signed_block header (not wall-clock)', () => {
const abi = makeShipAbi()
const headerBytes = Serializer.encode({
object: { timestamp: '2026-05-31T18:01:22.000', producer: 'eosio' },
type: 'block_header',
abi,
})
const raw = {
head: { block_num: 200, block_id: 'a'.repeat(64) },
last_irreversible: { block_num: 190, block_id: 'b'.repeat(64) },
this_block: { block_num: 200, block_id: 'c'.repeat(64) },
prev_block: null,
block: headerBytes,
traces: null,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 200, 'c'.repeat(64))
expect(block.blockTime).toContain('2026-05-31T18:01:22')
})
it('block_time is empty when block not fetched (fetch_block disabled)', () => {
const abi = makeShipAbi()
const raw = {
head: { block_num: 100, block_id: 'a'.repeat(64) },
last_irreversible: { block_num: 90, block_id: 'b'.repeat(64) },
this_block: { block_num: 100, block_id: 'c'.repeat(64) },
prev_block: null,
block: null,
traces: null,
deltas: null,
}
const block = decodeBlocksResult(raw, abi, 100, 'c'.repeat(64))
expect(block.blockTime).toBe('')
})
})
describe('ShipProtocol — defence against wharfkit returning typed objects', () => {
// Регрессионный тест: если кто-то удалит String() обёртки — эти тесты упадут.
it('all trace string fields survive JSON.stringify roundtrip without [object Object]', () => {