Compare commits

...

28 Commits

Author SHA1 Message Date
ant b84b02f1b4 Merge pull request 'sync-arch: parser2 ^1.2.0 + прокидывание block_time из ActionEvent/DeltaEvent [C28-13]' (#92) from parser2-block-time-and-version into parser2
Typecheck / desktop (pull_request) Successful in 12m46s
Typecheck / controller (pull_request) Successful in 12m28s
Reviewed-on: #92
Reviewed-by: Алексей Муравьев <chairman.voskhod@gmail.com>
2026-06-05 10:45:08 +00:00
coopops c33ec53230 [C28-13][@ant] feat(controller): parser2 ^1.2.0 + прокидывание block_time из ActionEvent/DeltaEvent
- @coopenomics/parser2: "1.1.0" → "^1.2.0" — снят точный pin, апгрейд парсера = npm install
- IAction/IDelta: + block_time?: string (ISO-8601, parser1 не давал, parser2 1.2.0 отдаёт)
- parser2-event.mapper.ts: пробрасывает event.block_time в IAction.block_time и IDelta.block_time
- mapper.test: 2 новых проверки на block_time (delta + action) — 9/9 зелёные

Поверх ветки parser2 (#68). Native-delta остаётся явно проигнорированной (как у parser1), это отдельное расширение интерфейса при необходимости.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-06-05 08:05:31 +00:00
ant 7dcca3a2ae Merge pull request 'Эпик 4 sync-arch: правильная обработка форков блокчейна (релиз 1.1.2) [C28-14]' (#54) from parser2-epic-4-forks into parser2
Typecheck / desktop (pull_request) Successful in 13m20s
Typecheck / controller (pull_request) Successful in 13m21s
Reviewed-on: #54
Reviewed-by: Алексей Муравьев <chairman.voskhod@gmail.com>
2026-06-03 08:32:46 +00:00
ant 70752b2caa Merge pull request 'sync-arch: явные ошибки при unknown contract version / status (Story 6.5) [C28-16]' (#57) from parser2-epic-6-canonical-storage into parser2-epic-4-forks
Reviewed-on: #57
Reviewed-by: Алексей Муравьев <chairman.voskhod@gmail.com>
2026-06-03 08:31:35 +00:00
coopops 8c5568bcf8 [C28-16][@ant] revert: откат Story 6.1 (namespace bc + replaceBc) + Story 6.2 (signedDocumentFields) — вынос за MVP-релиз parser2
Цель релиза 1.1.2 — «затащить parser2 нормально», т.е. заменить транспорт parser1→parser2.
Эпик 6 «единый порядок хранения данных» — отдельная архитектурная санитация, не связана с
заменой транспорта. Включение её в MVP-релиз привнесло инвазивные изменения:

- Story 6.1 (namespace `db`/`bc` + `replaceBc` + переписанный ProjectDomainEntity) вводит
  двойной стандарт: одна entity на новом паттерне, 21 entity на старом. Потребители
  читают `project.master` (плоский флэт) — частичная миграция создаёт регрессии у resolver'ов.
  Полная миграция = большой blast radius на 22 entity + Epic 9.5 backlog. Откладывается
  в отдельный sync-arch sanitation-эпик.

- Story 6.2 (declarative `signedDocumentFields` + `normalizeSignedDocuments` + path-parser)
  — DRY-рефакторинг без функциональной выгоды. AppendixDeltaMapper и так руками вызывал
  `convertChainDocumentToDomainFormat`. Польза появится когда подписанных полей станет
  много — пока их одно поле в одной точке.

Остаётся: Story 6.5 (UnsupportedContractVersionError + auditUnknownStatus +
`BLOCKCHAIN_UNSUPPORTED_VERSION_STRICT`). Очевидный фикс silent loss при schema drift,
не инвазивный.

Удалено: composite-entity.contract.test.ts, signed-document-normalization.test.ts,
delta-mapper-signed-doc.contract.test.ts. CLAUDE.md секция Composite-Entity ADR-008
помечена как «будущая цель, не сейчас».

10 jest blockchain unit suites / 64 tests зелёные. tsc зелёный.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-06-02 10:37:48 +00:00
coopops 30b79ad1f7 [C28-16][@ant] revert: откат Story 6.3 (zod) + 6.4 (_checksum) — pre-mature и за гранью sync-слоя
Story 6.3 (per-doc-type zod schemas): валидация `signatures.length === 1/2` зашивала
документную семантику в нормализатор parser-event. Число и валидность подписей — домен
контракта (`require_auth`), не бэкенда. При расширении доктайпов контракта (1→2 подписи)
sync падал бы на структурно-валидной цепи. Структурный transform `meta: JSON-string → object`
(Story 6.2) остаётся.

Story 6.4 (`_checksum` колонка + sha256(canonical-json(bc))): потребитель — Epic 7 nightly
snapshot и Epic 8 reconciliation — ещё не реализованы. Колонка хранила бы пустую нагрузку
+ +20% storage + лишний sha256 на каждом save/update. Принцип «не добавлять контрольных
полей до появления потребителя». Алгоритм canonicalStringify заведём как часть Epic 8.

Удалено: src/shared/sync/{checksum.util.ts, signed-document-schemas.ts} + 3 теста +
applyBcChecksum + колонка _checksum + поле _checksum на BaseDomainEntity. Story 6.1
(namespace db/bc + replaceBc), 6.2 (declarative signedDocumentFields + meta transform),
6.5 (UnsupportedContractVersionError + auditUnknownStatus) сохраняются.

13 jest blockchain unit suites / 196 tests зелёные. tsc зелёный.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-06-02 08:37:57 +00:00
coopops 87747d1782 [C28-16][@ant] feat: UnsupportedContractVersionError + auditUnknownStatus — Story 6.5
Новые shared/sync/errors:
- UnsupportedContractVersionError(entityName, ctx{contract,table,primary_key,
  block_num}) — носит контекст для DLQ-операторской диагностики.
- auditUnknownStatus(entityName, receivedStatus, logger, allowedStatuses?) —
  фиксирует unknown-status drift в audit-trail с ожидаемыми вариантами.

AbstractEntitySyncService.processDelta: mapper вернул null → logger.error
("UNSUPPORTED_CONTRACT_VERSION", ctx) ВСЕГДА (раньше был silent warn). В strict-mode
(config.blockchain.unsupported_version_strict=true) дополнительно throw — парсер не
ACK'нет дельту, DLQ сработает. UnsupportedContractVersionError из try/catch
пробрасывается дальше; остальные ошибки по-прежнему логируются и return null.

config.ts: BLOCKCHAIN_UNSUPPORTED_VERSION_STRICT (default false) → blockchain
.unsupported_version_strict. Default false — не ломать прод немедленно;
включается на стенде после подтверждения отсутствия schema drift.

ProjectDomainEntity.mapStatusToDomain эталонно — default ветка вместо silent
UNDEFINED-возврата вызывает auditUnknownStatus с полным списком ожидаемых:
[pending,active,voting,result,finalized,cancelled]. Остальные mapStatusToDomain
(state, vote, segment, debt и т.д.) — Epic 9.5.

unsupported-version-explicit-error.test.ts (6) — UnsupportedContractVersionError
конструктор/контекст; auditUnknownStatus с/без allowedStatuses;
processDelta non-strict (logger.error + null), strict (throw),
happy-path (handleSyncDelta вызывается).

229/229 blockchain unit зелёные.

Эпик 6 закрыт целиком: 6.1 namespace, 6.2 normalizeSignedDocuments, 6.3 Zod schemas,
6.4 checksum, 6.5 explicit errors. Release 1.1.2 движется.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-06-02 07:52:40 +00:00
coopops f7a182e51f [C28-16][@ant] feat: detrministic _checksum sha256(bc canonical-json) — Story 6.4
checksum.util.ts: canonicalStringify(value) (рекурсивная сортировка ключей объектов,
порядок массивов сохраняется, bigint → string, undefined → "null"); computeBcChecksum
→ sha256-hex 64 chars от canonical-stringified bc-namespace. Без внешних зависимостей.

BaseTypeormEntity получает колонку `_checksum varchar(64) nullable` —
legacy записи без bc-namespace получают хэш от "null" (стабильно, но не несёт
информации); после миграции Epic 9.5 на namespace все блокчейн-зеркала получат
содержательный checksum. BaseDomainEntity получает поле `_checksum?: string | null`.

BaseBlockchainRepository.save и .update вызывают protected applyBcChecksum(domain,
typeormEntity) ПОСЛЕ mapper.toEntity и ДО repository.save — proseться через единую
точку, не лезем в 22 mapper'а. Поле bc-namespace → reconciliation Epic 8.2
сравнивает только то, что реально в цепи (локальные db-поля типа matrix_room_id
из цепи не получаются и в checksum их быть не должно).

checksum.util.test.ts (14) — canonicalStringify: primitives/array/nested/bigint/sort;
computeBcChecksum: 64 hex chars / детерминизм / change-detection / null-stable /
порядок массивов / контрольный фиксированный хэш на nested-структуре (catch для
несовместимого изменения алгоритма).

base-repo-checksum.contract.test.ts (3) — bc заданный → checksum по нему;
bc undefined → checksum от null; повторный вызов → одинаковый _checksum.

223/223 blockchain unit зелёные.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-06-02 07:44:55 +00:00
coopops 9774abb296 [C28-16][@ant] feat: Zod per-doc-type schemas + valid-pre-transform — Story 6.3
Новый файл shared/sync/signed-document-schemas.ts: signatureInfoSchema (ISignatureInfo)
+ три per-doc-type схемы IChainDocument2:
- chainDocumentSchema — default ≥1 подпись (стандартный мультисиг кооператива).
- singleSignatureChainDocumentSchema — ровно 1 (одиночный автор).
- twoSignatureChainDocumentSchema — ровно 2 (двухподписные акты signsupp/signchair
  и signact1/signact2; OQ-13: ВСЕ подписанты массивом, primary не выделяется).

SignedDocField (Story 6.2) расширен optional schema?: ZodTypeAny — backward-compat
для legacy полей без валидации. AbstractBlockchainDeltaMapper.normalizeSignedDocuments
делает schema.parse ДО transform: при schema-drift цепи (новая структура IChainDocument2
из коопконтракта) — Zod бросает ZodError, mapper.try/catch отдаёт null + warn, Story 6.5
повысит до alert. До этого PR drift молча проходил, meta оставался JSON-строкой.

Эталон: appendix-delta.mapper.ts — добавлен singleSignatureChainDocumentSchema
(приложение к ТЭМ имеет одного автора).

signed-document-schemas.test.ts (16) — каждая схема: valid → ok, обязательные поля
отсутствуют → fail, неверное число подписей → fail. signed-document-normalization
расширен 3 тестами на schema-pre-transform: валидная schema → нормализация проходит,
single-signature на 2 подписи → ZodError, missing hash → ZodError + meta не обновлена
(transform не вызвался).

206/206 blockchain unit зелёные.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-06-02 07:27:00 +00:00
coopops 724cd381cc [C28-16][@ant] feat: AbstractDeltaMapper.normalizeSignedDocuments + SignedDocField — Story 6.2
AbstractBlockchainDeltaMapper получает protected readonly signedDocumentFields
ReadonlyArray<SignedDocField> (default []) + helper normalizeSignedDocuments(data).
SignedDocField описывает путь к IChainDocument2-полю в TBlockchainData; bracket-
нотация поддерживает массивы: `appendix` (top-level), `statement.attachments[]
.signed_attachment` (nested-array, E12). parseSignedDocPath → PathSegment[],
applyAtPath in-place трансформирует через
DomainToBlockchainUtils.convertChainDocumentToDomainFormat.

appendix-delta.mapper.ts — эталон: ручной convertChainDocumentToDomainFormat
(value.appendix) заменён декларацией signedDocumentFields = [{ path: 'appendix' }]
+ this.normalizeSignedDocuments({ ...value }).

signed-document-normalization.test.ts (16) — парсер 4 кейса + normalize 7 кейсов
(top-level, nested-array, missing-field no-op, multiple fields, empty config,
null/undefined, immutability caveat).

delta-mapper-signed-doc.contract.test.ts (3) — guard: mapper с непустым
signedDocumentFields обязан вызвать this.normalizeSignedDocuments — иначе
silent data corruption (meta остаётся JSON-строкой, Story 6.4 checksum
не сходится).

IPFS lazy resolver (E14) и миграция остальных mapper'ов с ручной нормализацией —
вне scope, тречится Epic 9.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-06-02 07:21:09 +00:00
coopops 5aa1a21a97 [C28-16][@ant] feat: namespace db/bc/derived + replaceBc helper — Story 6.1
BaseDomainEntity получает generic <TDb, TBc>: this.db (shallow-copy databaseData),
this.bc (nullable, заполняется через utility replaceBc()), block_num/present/status
остаются на базе (hot-path stale-delta guard). Backward-compat: TBc=unknown по
дефолту — 22 legacy entity компилируются без изменений.

replaceBc вынесен utility-функцией, не методом класса — protected/public метод
ломает structural typing там, где Omit<XxxDomainEntity,...> используется в сигнатурах
репо (state-typeorm-repo).

ProjectDomainEntity — эталон: PROJECT_BC_KEYS как const satisfies, updateFromBlockchain
вызывает replaceBc(this, blockchainData, PROJECT_BC_KEYS) вместо запрещённого
Object.assign(this, blockchainData). Плоские поля остаются для legacy compat —
миграция остальных потребителей на entity.bc.* трекется Epic 9.5.

composite-entity.contract.test.ts — рекурсивно сканит src/**/*.entity.ts, для
не-legacy entity запрещает Object.assign(this, ...). Legacy allowlist на 22 файла
с явной отсылкой на Epic 9.5. 99/99 entity prove'нуты + Project явно проверен
на replaceBc + PROJECT_BC_KEYS.

controller/CLAUDE.md — секция Composite-Entity (ADR-008) актуализирована: BC_KEYS
с satisfies ReadonlyArray<keyof IXxxBlockchainData> — canonical-ordered источник
для Story 6.4 checksum.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-06-02 07:11:34 +00:00
coopops 933e5b5ffe [C28-14][@ant] feat: архив инвалидированных сущностей + retention LIB-1000 — Story 4.4
handleFork теперь не делает hard-delete, а атомарно переносит снесённые форком live-ряды
в invalidated_entities и инвалидированные снимки entity_versions в invalidated_entity_versions
(две таблицы — разные семантики). Порядок: archiveInvalidatedSince → restoreFromVersions →
archiveInvalidatedVersionsSince. forkEventId (controller-формат chain:fork:N:short) пробрасывается
от handleEvent → processFork → ForkRegistry.runAll → syncer.handleFork → архив.fork_event_id
для forensic-группировки.

BlockchainArchiveRetentionService (@Cron, default ежечасно) читает LIB через
BlockchainService.getInfo() и удаляет архив старше LIB-1000 блоков. RETENTION_HORIZON_BLOCKS=1000
хардкод (свойство сети, не оператора). Env-переключатели:
BLOCKCHAIN_ARCHIVE_RETENTION_ENABLED (default true), BLOCKCHAIN_ARCHIVE_RETENTION_CRON.

EntityVersioningService расширен archive-методами под DataSource.transaction. Контракт
IForkAwareSyncer.handleFork(forkBlockNum, forkEventId?) — обратно совместим (2-й optional).
IBlockchainSyncRepository.archive* — optional, fallback на старую find+delete для off-chain
(allowlist Story 4.3).

ScheduleModule.forRoot() добавлен в app.module (раньше не было).

Tests: 14 новых unit (entity-archive, retention, handleFork contract) + 52 регрессионных
(processFork, fork-registry, base-repo-contract, no-onevent) — 66 зелёных, tsc=0.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-28 16:43:10 +00:00
coopops d886d780fc [C28-14][@ant] test: contract guard BaseBlockchainRepository + block_num типизация — Story 4.3
Audit-проход без миграции рантайма. Полная карта entity/repo и сравнений
block_num — в blago _bmad-output/tracks/sync-arch/.../audit-report-4-3.md.

Изменения в коде:

1. tests/unit/blockchain/base-blockchain-repository.contract.test.ts — НОВЫЙ
   regression-guard. Сканирует все *.typeorm-entity.ts extends BaseTypeormEntity
   и для каждой ищет *.typeorm-repository.ts extends BaseBlockchainRepository.
   Allowlist OFF_CHAIN_BASE_ENTITIES для 5 off-chain артефактов (comment,
   cycle, issue, story, time-entry) — у них block_num унаследован vestigial,
   синкера нет, форк их не откатывает (корректно). Тест ловит регрессию,
   если кто-то добавит блокчейн-mirror entity без BaseBlockchainRepository
   — иначе entity_versions молча перестанут писаться, форк превратится в
   hard delete (silent data loss).

2. controller/CLAUDE.md — два новых правила:
   (a) разнобой типов block_num: BaseTypeormEntity (capital+shared) использует
       integer+number — корректно; ActionEntity/DeltaEntity/ForkEntity/
       SyncStateEntity — bigint+number type-mismatch (PG возвращает string,
       TS говорит number). Hot-path везде явно Number() либо PG bind, так
       что runtime safe. Технический долг в Epic 9 — bigint Transformer.
   (b) контракт «entity с block_num → repo extends BaseBlockchainRepository»
       + ссылка на regression test + allowlist OFF_CHAIN_BASE_ENTITIES.

3. tests/unit/blockchain/base-blockchain-repository.contract.test.ts:
   нюанс — entityKindFromFileName нормализует суффикс -typeorm (chairman
   approval-typeorm.entity.ts vs approval.typeorm-repository.ts mismatch
   имён файлов).

Findings (полное в audit-report-4-3.md):
- 25 entity extends BaseTypeormEntity; 20 имеют BaseBlockchainRepository.
- 5 без BaseBlockchainRepository — off-chain артефакты (легитимно).
- 4 infra entity с bigint+number type-mismatch (runtime safe, тех. долг).
- Direct typeormRepo.save для блокчейн-зеркал — НЕ найдено (контракт соблюдён).
- Все сравнения block_num в sync-core используют Number() или PG bind.

Tests: 2/2 contract-test зелёные.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-28 15:51:22 +00:00
coopops 68b37efbe9 [C28-14][@ant] refactor: снять DEC-013 pause-barrier и @OnEvent('fork::*') — Story 4.2
После Story 4.1 ForkRegistry уже даёт sequential rollback всех syncer'ов через
pull-DiscoveryService, а saveFork идёт прямым вызовом из processFork. Старый
broadcast-путь через EventEmitter ('fork::*' + emitAsyncWithTimeout с TTL
force-resume из Story 1.3 DEC-013) — рудимент: дублирует saveFork (через
BlockchainEventHandlerService.handleForkEvent), повторно зовёт handleFork в
каждом syncer'е (idempotent no-op после ForkRegistry, но лишний DB-traffic).

Удалено:
- @OnEvent('fork::*') + handle*Fork методы в 20 файлах:
  · 3 в infrastructure/database/typeorm/blockchain/services/ (user-agreement, user-wallet, agreement)
  · 16 в extensions/capital/application/syncers/
  · 1 в extensions/chairman/infrastructure/blockchain/services/ (approval)
  · 1 дубль saveFork в domain/parser/services/blockchain-event-handler.service.ts
  · orphan dead-метод в application/wallet/services/program-wallet-sync.service.ts
- Шаг 4 (deprecated emitAsyncWithTimeout) из BlockchainConsumerService.processFork.
- BLOCKCHAIN_FORK_PAUSE_TIMEOUT_MS из config/config.ts (zod + config объект).
- Упоминание env-var из controller/CLAUDE.md.

Что осталось:
- EventsService.emitAsyncWithTimeout — определение в infrastructure/events/events.service.ts.
  Каллеров в src/ нет; не удаляю утилиту — возможно пригодится для Epic 5
  (waitForDelta / pool retry). Удалят на следующем рефакторе если не использована.
- AbstractEntitySyncService.handleFork остаётся: его теперь зовёт только ForkRegistry.
- @OnEvent('action::*') и @OnEvent('delta::*') в BlockchainEventHandlerService —
  валидный dispatch pipeline (ADR-002), не трогаем.

Tests:
- processFork.test.ts: убраны expect'ы на emit; новый regression test что
  events.emit/emitAsyncWithTimeout НЕ вызываются с fork:: префиксом.
- НОВЫЙ tests/unit/blockchain/no-onevent-fork.test.ts: grep-guard по src/
  на @OnEvent('fork::*') (комментарии отфильтрованы) → fail на регрессии.
- Все 6 suites зелёные (52 теста), tsc --noEmit 0 errors.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-28 15:42:39 +00:00
coopops e810d3152f [C28-14][@ant] feat: ForkRegistry sequential apply форка — Story 4.1
Цель — снять параллельный @OnEvent('fork::*') broadcast в 20+ syncer-ов и
заменить на sequential обход через ForkRegistryService, что в сочетании
с single-active XREADGROUP parser2 даёт натуральный барьер форка (INV-T03,
NFR10). Старый broadcast путь оставлен deprecated на этот релиз — Story 4.2
удалит вместе с pause-barrier из DEC-013.

Реализация:
- shared/sync/fork/{interface,service,module}: IForkAwareSyncer + symbol-marker
  FORK_AWARE_MARKER, ForkRegistryService с pull-сбором через DiscoveryService
  на onApplicationBootstrap, sequential for-of await + priority-ordering.
- AbstractEntitySyncService: implements IForkAwareSyncer (marker через class
  field), handleFork теперь re-throw (catch удалён) — контракт sequential
  ForkRegistry.runAll. Не требует super.onModuleInit() у 20 наследников
  (DiscoveryService снимает push-зависимость).
- BlockchainConsumerService.handleEvent для fork: dedup-gate через
  computeForkEventId → processFork (runAll → deleteDedupAfterBlock → saveFork
  → deprecated emit) → markEventApplied. Для action/delta — markEventApplied
  теперь пишет block_num.
- ConsumerDedup: новая колонка block_num (bigint nullable) + индекс +
  deleteAfterBlock(blockNum). Старые NULL-записи не затрагиваются.
- event-id.util: новый computeForkEventId в формате chain:fork:block:short_id
  (полное «fork», единый стиль с action/delta — parser2-формат chain:f:... не
  используется).

Тесты (51 кейс, 5 suites):
- unit fork-registry: sequential apply, error propagation, priority, bootstrap
  discovery (15 кейсов).
- unit blockchain-consumer.processFork: порядок шагов, dedup-gate, mark после
  processFork, fail-paths, parser2-format isolation (11 кейсов).
- unit consumer-dedup repository: markApplied с/без blockNum, deleteAfterBlock
  edge-cases (8 кейсов).
- unit event-id.util: computeForkEventId детерминизм + поведение на коротком
  block_id (5 новых кейсов).
- integration fork-flow: bootstrap → delta(N+1)→fork(N)→delta(N+2), rollback
  не трогает <=N, error propagation (4 кейса с реальным Nest DI).

tsc 0 ошибок, jest 51/51 зелёные.

Spec: blago/production/13-platforma-tsifrovogo-kooperativa/components/14-versiya-3/_bmad-output/tracks/sync-arch/implementation-artifacts/spec-4-1-fork-kak-event-v-unified-stream.md

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-28 14:52:38 +00:00
coopops 99c22b0064 [C28-13][@ant] chore: обновить pnpm-lock.yaml до @coopenomics/parser2@1.1.0
Lockfile висел на 1.0.3, потому что 1.1.0 ещё не был опубликован
в npm в момент bump'а package.json. Сейчас 1.1.0 в реестре, pnpm
install подхватил его + новую транзитивную @coopenomics/coopos-ship-reader@0.2.0
и @napi-rs/nice@1.1.1.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-28 12:26:30 +00:00
ant 4c0fbefca9 Merge pull request 'Эпик 3 sync-arch: переезд транспорта controller с parser1 на parser2 ParserClient [C28-13]' (#35) from parser2-epic-3-transport into parser2
Reviewed-on: #35
2026-05-28 11:57:20 +00:00
coopops 8fdc99a7f9 [C28-13][@ant] feat: маппер parser2 отдаёт реальные поля action-трейса + dep parser2 1.1.0 — паритет с parser1, без заглушек
transaction_id/creator_action_ordinal/context_free/elapsed/console/
account_ram_deltas + receipt.auth_sequence теперь идут из ActionEvent
parser2 1.1.0 (добавлены в сам пакет + ship-reader 0.2.0). Нужны
ledger2 для cross-link родительского apply (transaction_id+action_ordinal)
и blockchain-explorer. tsc --noEmit контроллера против новых типов — чисто.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-25 17:08:24 +00:00
coopops 8d27c4e245 [C28-13][@ant] feat: заменить транспорт parser1 на parser2 ParserClient — единый движок без флагов
Контроллер читает события через @coopenomics/parser2 ParserClient вместо самодельного
consumer'а поверх Redis-стрима notifications (его писал parser1). Никаких флагов и
параллельной работы двух движков: одна система — либо работает на parser2, либо нет.

Удалена ручная обвязка (RedisStreamService, consumer-group, recoverOwnPending, XAUTOCLAIM,
XTRIM) — всё это ParserClient делает сам (single-active-lock, recover, dead-letter после
N провалов). Новый consume-loop ведёт генератор stream() вручную: next() при успехе
(XACK внутри parser2), throw() при ошибке (учёт провалов / dead-letter) — наивный for-await
неверен, проброс из тела не доходит до catch вокруг yield и убивал бы консьюмер.

Маппер parser2-event.mapper переводит ParserEvent → IDelta/IAction (DEC-T09), обработчики
processAction/processDelta/processFork не тронуты. Дедуп по event_id теперь безусловный
(флаг BLOCKCHAIN_DEDUP_ENABLED убран), формат event_id оставлен прежним (action/delta).

ВНИМАНИЕ: parser2 ActionEvent не несёт transaction_id/creator_action_ordinal — ledger2
cross-link родительского apply на них опирается; в маппере заглушки, проверить на прогоне.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-25 15:58:12 +00:00
coopops 0e3ccb4356 Revert "[C28-13][@ant] fix: выровнять формулу event_id под parser2 v1.0.3 — сверка перед cutover"
This reverts commit 306835d0fc.
2026-05-25 15:42:46 +00:00
coopops 306835d0fc [C28-13][@ant] fix: выровнять формулу event_id под parser2 v1.0.3 — сверка перед cutover
Локальная формула computeDeltaEventId/computeActionEventId расходилась с движком
parser2 (дискриминанты delta/action вместо a/d, block_id[0..8] вместо [0..16]).
Это закрывает сверку из Story 3.3 phase 2, перенесённую вперёд: computeEventId
теперь импортируем из опубликованного пакета. Выровнял байт-в-байт, golden-тест
сверен против исходника parser2 (delta и action совпадают). Без выравнивания при
overlap dual-consume один event получал бы два разных id в legacy и в движке и
dedup-gate промахнулся бы = silent data loss. Поправил формат event_id в CLAUDE.md.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-25 10:18:41 +00:00
ant 10e5b83ac9 Merge pull request 'Эпик 2 sync-arch: фундамент идемпотентности — event_id + consumer_dedup (релиз 1.1.1) [C28-12]' (#33) from parser2-epic-2-idempotency into parser2
Reviewed-on: #33
Reviewed-by: Алексей Муравьев <chairman.voskhod@gmail.com>
2026-05-25 10:05:56 +00:00
coopops 33d936e5c8 [C28-12][@ant] refactor: убрать DDL-миграцию consumer_dedup — таблицу создаёт synchronize:true через entity, отдельная миграция избыточна по текущей конвенции typeorm (ревью PR #33)
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-25 07:45:54 +00:00
coopops a73f07fa0c [C28-12][@ant] feat: локальное вычисление event_id с dual-write и dedup-gate за флагом — повтор события распознаётся как no-op, gate выключен до сверки формулы с parser2 в Epic 3 (Stories 2.2, 2.3, DEC-T08, DEC-020)
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-24 18:31:00 +00:00
coopops 77f358e865 [C28-12][@ant] feat: добавить таблицу consumer_dedup с миграцией и репозиторием — фундамент идемпотентности для распознавания повторно применённых событий блокчейна (Story 2.1, INV-09)
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-24 18:30:53 +00:00
ant 2afc71c281 Merge pull request 'Эпик 1 sync-arch: срочные фиксы синхронизации (релиз 1.1.1) [C28-11]' (#32) from parser2-epic-1-hardening into parser2
Reviewed-on: #32
2026-05-24 18:17:23 +00:00
coopops 9cf19aa8fe [C28-11][@ant] fix: барьер форка и вынос задержки action-emit в конфиг — откаты форка завершаются до обработки следующей дельты, исчезает хардкод 3000мс (Stories 1.3, 1.4, DEC-007)
Stories 1.3 + 1.4 затрагивают общие файлы (config.ts, blockchain-consumer.service.ts), поэтому одним коммитом:
- 1.4: ACTION_EMIT_DELAY_MS=3000 → config.blockchain.action_emit_delay_ms (env BLOCKCHAIN_ACTION_EMIT_DELAY_MS). Порядок save→ACK→setTimeout(emit) уже был корректен.
- 1.3: EventsService.emitAsyncWithTimeout + await в processFork = пауза consumer'а до завершения @OnEvent('fork::*')-откатов; TTL force-resume = config.blockchain.fork_pause_timeout_ms. Временно до Epic 3/ForkRegistry.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-24 17:38:42 +00:00
coopops 4b8121b0b7 [C28-11][@ant] fix: guard монотонности block_num в createIfNotExists — устаревшая дельта из раннего блока больше не затирает более свежую запись на create-пути (Story 1.1, DEC-008)
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-24 17:38:41 +00:00
72 changed files with 3764 additions and 851 deletions
+17 -13
View File
@@ -58,6 +58,7 @@ _Критичные правила и паттерны для AI-агентов
### Language-Specific (TypeScript)
- **Bigint из PostgreSQL приходит как STRING.** Всегда `Number(blockNum) < Number(currentBlockNum)`, никогда не полагаться на `<`/`>` для `bigint`-колонок напрямую.
- **block_num — разнобой типов в текущей кодовой базе (Story 4.3 audit).** `BaseTypeormEntity` (capital + shared) использует `@Column({ type: 'integer' }) block_num!: number` (PG возвращает number — корректно). Инфраструктурные entity (`ActionEntity`, `DeltaEntity`, `ForkEntity`, `SyncStateEntity`) используют `@Column({ type: 'bigint' }) block_num!: number`**type-mismatch** (runtime будет string, TS говорит number). `ConsumerDedupEntity``bigint` + `string | null` (корректно). Это технический долг (Epic 9 backlog: bigint Transformer). Hot-path везде явно делает `Number()` либо использует TypeORM bind, так что runtime safe. При добавлении нового блокчейн-зеркала — выбрать одно из: (a) `integer` + `number`; (b) `bigint` + Transformer возвращающий `number`; (c) `bigint` + `string` с явным `Number()` в кодe. НЕ `bigint` + declaration `number` без Transformer.
- **`Object.assign(this, blockchainData)` запрещено** в sync-сущностях — ломает типизацию. Только явное копирование полей.
- **Lowercase hash полей (`project_hash`, `listing_hash`)** — нормализация **в конструкторе/mapValue**, НЕ в mapper'ах и НЕ повторно в `updateFromBlockchain`.
- **Discriminated union для write-mutation response:** `{ status: 'applied' | 'pending' | 'failed' | 'conflict' }`. Клиент должен switch на статусе.
@@ -76,6 +77,8 @@ _Критичные правила и паттерны для AI-агентов
- `@Column({ type: 'bigint', nullable: true })` для `block_num`. Для `jsonb``@Column({ type: 'jsonb' })`.
- `ADD COLUMN NOT NULL` на больших таблицах — **двухэтапно**: ADD nullable → backfill → ALTER NOT NULL.
- Repository `extends BaseBlockchainRepository<DomainEntity, TypeormEntity>`. `findBySyncKey`, `createIfNotExists`, `deleteByBlockNumGreaterThan`, `restoreFromVersions`**наследуются**, не реализовывать руками.
- **Контракт «entity с block_num → repo extends BaseBlockchainRepository» (Story 4.3).** Любая `*.typeorm-entity.ts` extends `BaseTypeormEntity` ОБЯЗАНА иметь репозиторий extends `BaseBlockchainRepository` — иначе `entity_versions` не пишется (silent), а форк-rollback превращается в hard delete без восстановления. CI grep-guard: `tests/unit/blockchain/base-blockchain-repository.contract.test.ts`. Allowlist для 5 off-chain entity (`comment`, `cycle`, `issue`, `story`, `time-entry` — vestigial block_num, не блокчейн-зеркала). Новое исключение в allowlist — только с обоснованием в audit-report.
- **Архив форка вместо hard-delete (Story 4.4).** `handleFork(N, eventId?)` НЕ удаляет live-ряды и снимки версий — переносит в `invalidated_entities` / `invalidated_entity_versions` (атомарно через DataSource.transaction). Порядок: archiveInvalidatedSince → restoreFromVersions → archiveInvalidatedVersionsSince. `fork_event_id` группирует записи одного форка (для forensic). Retention — `BlockchainArchiveRetentionService` ежечасно удаляет архив старше `LIB - 1000` блоков. LIB читается через `BlockchainService.getInfo()` (вариант C, RPC `/v1/chain/get_info`). Окно `RETENTION_HORIZON_BLOCKS = 1000` ХАРДКОД (свойство сети, не оператора). Env-переключатели: `BLOCKCHAIN_ARCHIVE_RETENTION_ENABLED` (default true), `BLOCKCHAIN_ARCHIVE_RETENTION_CRON` (default `0 * * * *`).
**parser2 integration:**
- `ParserClient` subscribe с `subscriptionId = "controller-${coopname}"`, `consumerName = "primary"` (детерминирован), `startFromBlock: 'last_known'`.
@@ -87,14 +90,16 @@ _Критичные правила и паттерны для AI-агентов
- `{contract}{Entity}Updated` / `Deleted` / `RolledBack` / `PendingRetry` / `Failed`**pubsub канал** для GraphQL subscriptions. Per-contract, не global.
- `entitysynced::{contract}::{table}` — для business-side-effect listeners (Matrix / Notification / и т.п.).
### Composite-Entity (ADR-008) — СТРОГО
### Composite-Entity (ADR-008) — будущая цель, не сейчас
- Namespaced: `entity.db.X` (DB-поля) / `entity.bc?.Y` (blockchain, nullable) / `entity.derived.Z` (computed getters).
- **НЕ** писать `entity.X` напрямую — ломает изоляцию.
- Конструктор `(databaseData, blockchainData?)` — обязан `throw` на sync-key mismatch.
- `updateFromBlockchain` возвращает **новый экземпляр** (immutable) или мутирует только `this.bc`, `this.block_num`, `this.present` — БЕЗ `Object.assign`.
- `derived` getter — детерминирован (NO `new Date()` в конструкторе / getter — ломает snapshot tests).
- Все signed-document поля нормализуются через `AbstractDeltaMapper.normalizeSignedDocuments` на основе `signedDocumentFields: SignedDocField[]` декларативно, НЕ руками в mapper.
ADR-008 описывает целевой паттерн `entity.db.X` / `entity.bc?.Y` / `entity.derived.Z`, заменяющий
`Object.assign(this, blockchainData)`. Переход вынесен за пределы MVP-релиза parser2 в отдельный
sync-arch sanitation-эпик: blast radius на 22 entity + потребители «плоских» полей в resolver'ах
делают эту миграцию большой и рискованной задачей, несвязанной с заменой транспорта parser1→parser2.
До отдельного эпика — текущий код продолжает использовать flat-namespace + `Object.assign` в
`updateFromBlockchain`. Не вводить namespace частично на одной entity — двойной канон хуже единого
старого.
### Dispatch pipeline (ADR-002, ADR-009) — СТРОГО
@@ -216,7 +221,6 @@ return { tx_hash: tx.tx_hash, status: 'pending' };
- `BLOCKCHAIN_RECONCILE_CRON` default `'0 * * * *'`
- `BLOCKCHAIN_RECONCILE_SAMPLE_SIZE` default 100
- `BLOCKCHAIN_RECONCILE_TOLERANCE_BLOCKS` default 10
- `BLOCKCHAIN_FORK_PAUSE_TIMEOUT_MS` default 30000
- `BLOCKCHAIN_DLQ_MAX_RETRIES` default 5
- `BLOCKCHAIN_MAX_TX_RETRIES` default 3
- `BLOCKCHAIN_PENDING_TX_MAX_AGE_SECONDS` default 3600
@@ -251,12 +255,13 @@ return { tx_hash: tx.tx_hash, status: 'pending' };
- `await capitalBlockchainPort.getProject(hash)` в resolver / application — **retire**. Только `repository.findBySyncKey`.
- RPC fallback "если PG не отдал" — **запрещено**. Если PG null → pending status наверх.
### ❌ Composite-entity anti-patterns
### ❌ Domain-entity anti-patterns
- `project.matrix_room_id` (flat access) — **запрещено**. Правильно: `project.db.matrix_room_id`.
- `project.master` (flat access) — **запрещено**. Правильно: `project.bc?.master`.
- `Object.assign(this, blockchainData)` в update — **запрещено**.
- Новый `new Date()` в конструкторе или derived-getter — **запрещено** (ломает snapshot-tests).
- Миграция на namespace `entity.db.X` / `entity.bc?.Y` — целевой паттерн ADR-008, но отложен
в отдельный sync-arch sanitation-эпик: текущий код всё ещё использует flat-namespace +
`Object.assign(this, blockchainData)` в `updateFromBlockchain`. Не разводить два стандарта
частично.
### ❌ Fork handling anti-patterns
@@ -296,7 +301,6 @@ return { tx_hash: tx.tx_hash, status: 'pending' };
### ⚡ Performance
- `JSON.stringify` на сущностях с nested signed-documents — использовать `json-stable-stringify` для canonical checksum.
- Reconciliation cron — sample N=100 rows, **не** full scan на hot path.
- IPFS fetch для signed-doc — **lazy resolver вне consumer critical path**, не в mapper.
- `waitForDelta` memory leak — timer cleanup обязателен (см. INV-T10).
+1
View File
@@ -76,6 +76,7 @@
"@coopenomics/factory": "workspace:*",
"@coopenomics/inter": "workspace:*",
"@coopenomics/notifications": "workspace:*",
"@coopenomics/parser2": "^1.2.0",
"@coopenomics/provider-client": "2025.11.12-alpha-1",
"@coopenomics/sdk": "workspace:*",
"@graphql-codegen/typescript-graphql-request": "^6.2.0",
+4
View File
@@ -1,6 +1,7 @@
// app.module.ts
import { Module } from '@nestjs/common';
import { ConfigModule } from '@nestjs/config';
import { ScheduleModule } from '@nestjs/schedule';
import { ThrottlerModule } from '@nestjs/throttler';
// Infrastructure modules
@@ -9,6 +10,7 @@ import { GraphqlModule } from './infrastructure/graphql/graphql.module';
import { MongooseModule } from '@nestjs/mongoose';
import config from '~/config/config';
import { BlockchainModule } from './infrastructure/blockchain/blockchain.module';
import { ForkRegistryModule } from './shared/sync/fork';
import { GeneratorInfrastructureModule } from './infrastructure/generator/generator.module';
import { RedisModule } from './infrastructure/redis/redis.module';
import { NovuModule } from './infrastructure/novu/novu.module';
@@ -85,6 +87,7 @@ import { MutationLoggingInterceptor } from './application/common/interceptors/mu
ConfigModule.forRoot({
isGlobal: true, // Чтобы .env был доступен глобально
}),
ScheduleModule.forRoot(), // @Cron / @Interval / @Timeout (Story 4.4 retention)
ThrottlerModule.forRoot([
{
ttl: 60000,
@@ -95,6 +98,7 @@ import { MutationLoggingInterceptor } from './application/common/interceptors/mu
MongooseModule.forRoot(config.mongoose.url),
DatabaseModule,
GraphqlModule,
ForkRegistryModule,
BlockchainModule,
GeneratorInfrastructureModule,
RedisModule,
@@ -49,12 +49,4 @@ export class ProgramWalletSyncService
this.logger.debug('Сервис синхронизации программных кошельков полностью инициализирован с подписками на паттерны');
}
/**
* Обработка форков для программных кошельков
* Подписывается на все форки независимо от контракта
*/
async handleProgramWalletFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
+39 -1
View File
@@ -123,7 +123,41 @@ const envVarsSchema = z.object({
.string()
.default('1000')
.transform((val) => parseInt(val, 10)),
/**
* Задержка (мс) перед emit'ом action-события во внутреннюю шину. Даёт
* дельтам того же блока сохраниться в БД раньше, чем обработчики action
* полезут читать состояние (DEC-007, ранее хардкод-константа 3000).
*/
BLOCKCHAIN_ACTION_EMIT_DELAY_MS: z
.string()
.default('3000')
.transform((val) => parseInt(val, 10)),
/**
* Story 4.4: глобальный выключатель ежечасного retention-крона архива
* invalidated_entities/invalidated_entity_versions. На малом объёме (нынешний
* кооператив) данные могут копиться годами — отключить кроном.
*/
BLOCKCHAIN_ARCHIVE_RETENTION_ENABLED: z
.string()
.default('true')
.transform((v) => v === 'true'),
/**
* Story 4.4: cron-расписание retention. Default — ежечасно. RETENTION_HORIZON_BLOCKS
* (=1000) хардкод в BlockchainArchiveRetentionService — не вынесен в env намеренно
* (свойство сети, не оператора).
*/
BLOCKCHAIN_ARCHIVE_RETENTION_CRON: z.string().default('0 * * * *'),
/**
* Story 6.5: при `true` mapper-fail (mapDeltaToBlockchainData → null) перестаёт
* быть silent loss и поднимается `UnsupportedContractVersionError` из
* `AbstractEntitySyncService.processDelta`. Парсер не ACK'ает delta — DLQ
* сработает. Default `false` для не-ломать-прод-немедленно; включается после
* подтверждения, что schema drift отсутствует (например на стенде).
*/
BLOCKCHAIN_UNSUPPORTED_VERSION_STRICT: z
.string()
.default('false')
.transform((v) => v === 'true'),
// Параметры NOVU
NOVU_APP_ID: z.string().min(1, { message: 'Не должно быть пустым' }),
NOVU_BACKEND_URL: z.string().min(1, { message: 'Не должно быть пустым' }).default('https://novu.coopenomics.world/api'),
@@ -231,6 +265,10 @@ export default {
root_precision: envVars.data.ROOT_PRECISION,
root_govern_precision: envVars.data.ROOT_GOVERN_PRECISION,
post_transact_chain_read_delay_ms: envVars.data.POST_TRANSACT_CHAIN_READ_DELAY_MS,
action_emit_delay_ms: envVars.data.BLOCKCHAIN_ACTION_EMIT_DELAY_MS,
archive_retention_enabled: envVars.data.BLOCKCHAIN_ARCHIVE_RETENTION_ENABLED,
archive_retention_cron: envVars.data.BLOCKCHAIN_ARCHIVE_RETENTION_CRON,
unsupported_version_strict: envVars.data.BLOCKCHAIN_UNSUPPORTED_VERSION_STRICT,
},
mongoose: {
url: envVars.data.MONGODB_URL + (envVars.data.NODE_ENV === 'test' ? '-test' : ''),
@@ -11,10 +11,12 @@ import type { ActionRepositoryPort } from '../ports/action-repository.port';
import type { DeltaRepositoryPort } from '../ports/delta-repository.port';
import type { ForkRepositoryPort } from '../ports/fork-repository.port';
import type { SyncStateRepositoryPort } from '../ports/sync-state-repository.port';
import type { ConsumerDedupRepositoryPort } from '../ports/consumer-dedup-repository.port';
import { ACTION_REPOSITORY_PORT } from '../ports/action-repository.port';
import { DELTA_REPOSITORY_PORT } from '../ports/delta-repository.port';
import { FORK_REPOSITORY_PORT } from '../ports/fork-repository.port';
import { SYNC_STATE_REPOSITORY_PORT } from '../ports/sync-state-repository.port';
import { CONSUMER_DEDUP_REPOSITORY_PORT } from '../ports/consumer-dedup-repository.port';
/**
* Интерактор парсера блокчейна
@@ -30,9 +32,36 @@ export class ParserInteractor {
@Inject(FORK_REPOSITORY_PORT)
private readonly forkRepository: ForkRepositoryPort,
@Inject(SYNC_STATE_REPOSITORY_PORT)
private readonly syncStateRepository: SyncStateRepositoryPort
private readonly syncStateRepository: SyncStateRepositoryPort,
@Inject(CONSUMER_DEDUP_REPOSITORY_PORT)
private readonly consumerDedupRepository: ConsumerDedupRepositoryPort
) {}
/**
* Отмечено ли событие как уже применённое (Story 2.3, INV-09).
*/
async isEventApplied(eventId: string): Promise<boolean> {
return await this.consumerDedupRepository.isApplied(eventId);
}
/**
* Отметить событие применённым в consumer_dedup (Story 2.2 dual-write).
* Идемпотентно: повтор после краха между save и mark не падает.
* blockNum (Story 4.1) — для последующего deleteAfterBlock при форке;
* опциональный, если вызов из мест без контекста блока.
*/
async markEventApplied(eventId: string, blockNum?: number): Promise<void> {
await this.consumerDedupRepository.markApplied(eventId, blockNum);
}
/**
* Удалить из consumer_dedup записи с block_num > forkBlockNum — очистка дедупа
* на форке (Story 4.1, ADR-005). Возвращает число удалённых строк.
*/
async deleteDedupAfterBlock(forkBlockNum: number): Promise<number> {
return await this.consumerDedupRepository.deleteAfterBlock(forkBlockNum);
}
/**
* Сохранение действия блокчейна
*/
@@ -0,0 +1,28 @@
/**
* Порт списка применённых событий (consumer_dedup) — фундамент идемпотентности
* (Story 2.1, INV-09).
*/
export interface ConsumerDedupRepositoryPort {
/** Отмечено ли событие как уже применённое. Используется dedup-gate (Story 2.3). */
isApplied(eventId: string): Promise<boolean>;
/**
* Отметить событие применённым. Идемпотентно (ON CONFLICT DO NOTHING): повтор
* после краха между save и mark не должен падать. blockNum — номер блока события
* (Story 4.1, для последующего deleteAfterBlock на форке); опциональный для
* backward-compat с местами, где блок неизвестен (например, ручные тесты).
*/
markApplied(eventId: string, blockNum?: number): Promise<void>;
/** Очистка меток старше cutoff (retention). Возвращает число удалённых строк. */
deleteOlderThan(cutoff: Date): Promise<number>;
/**
* Удалить метки событий с block_num > blockNum — для очистки дедупа на форке
* (Story 4.1, ADR-005). Записи с NULL block_num (legacy до Epic 4) НЕ затрагиваются:
* PG сравнение NULL > N даёт unknown и не попадает под WHERE. Возвращает число строк.
*/
deleteAfterBlock(blockNum: number): Promise<number>;
}
export const CONSUMER_DEDUP_REPOSITORY_PORT = Symbol('ConsumerDedupRepositoryPort');
@@ -94,27 +94,4 @@ export class BlockchainEventHandlerService implements OnModuleInit {
throw error; // Перебрасываем ошибку для корректной обработки
}
}
/**
* Обработка события форка блокчейна
* Сохраняет форк в базу данных
*/
@OnEvent('fork::*')
async handleForkEvent(data: { block_num: number }): Promise<void> {
try {
this.logger.debug(`Handling fork event at block: ${data.block_num}`);
// Преобразование данных форка для сохранения
const forkData = {
chain_id: config.blockchain.id, // Используем chain_id из конфигурации
block_num: data.block_num,
};
await this.parserInteractor.saveFork(forkData);
this.logger.debug(`Fork saved at block: ${data.block_num} for chain: ${config.blockchain.id}`);
} catch (error: any) {
this.logger.error(`Ошибка обработки события форка: ${error.message}`, error.stack);
throw error; // Перебрасываем ошибку для корректной обработки
}
}
}
@@ -75,12 +75,4 @@ export class AppendixSyncService
this.logger.error(`Ошибка при обработке отклонения приложения: ${error?.message}`, error?.stack);
}
}
/**
* Обработчик форков для приложений
*/
@OnEvent('fork::*')
async handleAppendixFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -104,13 +104,4 @@ export class CommitSyncService
this.logger.error(`Ошибка при обработке отклонения коммита: ${error?.message}`, error?.stack);
}
}
/**
* Обработка форков для коммитов
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleCommitFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { ContributorDomainEntity } from '../../domain/entities/contributor.entity';
@@ -84,13 +84,4 @@ export class ContributorSyncService
return contributorEntity;
}
/**
* Обработка форков для участников
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleContributorFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { DebtDomainEntity } from '../../domain/entities/debt.entity';
@@ -49,13 +49,4 @@ export class DebtSyncService
this.logger.debug('Сервис синхронизации долгов полностью инициализирован с подписками на паттерны');
}
/**
* Обработка форков для долгов
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleDebtFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { ExpenseDomainEntity } from '../../domain/entities/expense.entity';
@@ -49,13 +49,4 @@ export class ExpenseSyncService
this.logger.debug('Сервис синхронизации расходов полностью инициализирован с подписками на паттерны');
}
/**
* Обработка форков для расходов
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleExpenseFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { InvestDomainEntity } from '../../domain/entities/invest.entity';
@@ -49,13 +49,4 @@ export class InvestSyncService
this.logger.debug('Сервис синхронизации инвестиций полностью инициализирован с подписками на паттерны');
}
/**
* Обработка форков для инвестиций
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleInvestFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { ProgramPropertyDomainEntity } from '../../domain/entities/program-property.entity';
@@ -50,12 +50,4 @@ export class ProgramPropertySyncService
this.eventEmitter.on(pattern, this.processDelta.bind(this));
});
}
/**
* Обработчик форков для программных имущественных взносов
*/
@OnEvent('fork::*')
async handleProgramPropertyFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { ProgramWalletDomainEntity } from '../../domain/entities/program-wallet.entity';
@@ -47,12 +47,4 @@ export class ProgramWalletSyncService
this.eventEmitter.on(pattern, this.processDelta.bind(this));
});
}
/**
* Обработчик форков для программных кошельков
*/
@OnEvent('fork::*')
async handleProgramWalletFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { ProgramWithdrawDomainEntity } from '../../domain/entities/program-withdraw.entity';
@@ -50,12 +50,4 @@ export class ProgramWithdrawSyncService
this.eventEmitter.on(pattern, this.processDelta.bind(this));
});
}
/**
* Обработчик форков для возвратов из программы
*/
@OnEvent('fork::*')
async handleProgramWithdrawFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { ProjectPropertyDomainEntity } from '../../domain/entities/project-property.entity';
@@ -50,12 +50,4 @@ export class ProjectPropertySyncService
this.eventEmitter.on(pattern, this.processDelta.bind(this));
});
}
/**
* Обработчик форков для проектных имущественных взносов
*/
@OnEvent('fork::*')
async handleProjectPropertyFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { ProjectDomainEntity } from '../../domain/entities/project.entity';
@@ -139,13 +139,4 @@ export class ProjectSyncService
return projectEntity;
}
/**
* Обработка форков для проектов
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleProjectFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { ResultDomainEntity } from '../../domain/entities/result.entity';
@@ -102,13 +102,4 @@ export class ResultSyncService
return resultEntity;
}
/**
* Обработка форков для результатов
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleResultFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { SegmentDomainEntity } from '../../domain/entities/segment.entity';
@@ -54,16 +54,6 @@ export class SegmentSyncService
this.logger.debug('Сервис синхронизации сегментов полностью инициализирован с подписками на паттерны');
}
/**
* Обработка форков для сегментов
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleSegmentFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
/**
* Синхронизация сегмента между блокчейном и базой данных
*/
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { StateDomainEntity } from '../../domain/entities/state.entity';
@@ -49,13 +49,4 @@ export class StateSyncService
this.logger.debug('Сервис синхронизации состояния полностью инициализирован с подписками на паттерны');
}
/**
* Обработка форков для состояния
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleStateFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../shared/services/abstract-entity-sync.service';
import { VoteDomainEntity } from '../../domain/entities/vote.entity';
@@ -49,13 +49,4 @@ export class VoteSyncService
this.logger.debug('Сервис синхронизации голосов полностью инициализирован с подписками на паттерны');
}
/**
* Обработка форков для голосов
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleVoteFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -6,7 +6,13 @@ import type {
import type { IProjectDomainInterfaceBlockchainData } from '../interfaces/project-blockchain.interface';
import type { IBlockchainSynchronizable } from '~/shared/interfaces/blockchain-sync.interface';
import { BaseDomainEntity } from '~/shared/sync/entities/base-domain.entity';
import { auditUnknownStatus } from '~/shared/sync/errors/audit-unknown-status';
import { IssueIdGenerationService } from '../services/issue-id-generation.service';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
const PROJECT_STATUS_AUDIT_LOGGER = new WinstonLoggerService();
PROJECT_STATUS_AUDIT_LOGGER.setContext('ProjectDomainEntity');
/**
* Доменная сущность проекта
*
@@ -210,7 +216,16 @@ export class ProjectDomainEntity
case 'cancelled':
return ProjectStatus.CANCELLED;
default:
// По умолчанию считаем статус неопределенным
// Story 6.5: silent fallback на UNDEFINED заменён audit-trail'ом.
// Если контракт ввёл новый статус — drift всплывёт в логе как error.
auditUnknownStatus('ProjectDomainEntity', blockchainStatus, PROJECT_STATUS_AUDIT_LOGGER, [
'pending',
'active',
'voting',
'result',
'finalized',
'cancelled',
]);
return ProjectStatus.UNDEFINED;
}
}
@@ -29,10 +29,9 @@ export class AppendixDeltaMapper extends AbstractBlockchainDeltaMapper<IAppendix
return null;
}
// 🔥 ВАЖНО: Парсим документы ПЕРЕД возвратом
// Парсим документы ПЕРЕД возвратом
const appendix = DomainToBlockchainUtils.convertChainDocumentToDomainFormat(value.appendix);
// Парсим документы
return { ...value, appendix };
} catch (error: any) {
this.logger.error(`Error mapping delta to blockchain data: ${error.message}`, error.stack);
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import type { IDelta } from '~/types/common';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '../../../../../shared/services/abstract-entity-sync.service';
@@ -56,15 +56,6 @@ export class ApprovalSyncService
async handleApprovalDelta(delta: IDelta): Promise<void> {
await this.processDelta(delta);
}
/**
* Обработчик форков для одобрений
*/
@OnEvent('fork::*')
async handleApprovalFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
/**
* Получение поддерживаемых версий контрактов и таблиц
*/
@@ -1,225 +1,169 @@
// infrastructure/blockchain/blockchain-consumer.service.ts
import { Injectable, OnModuleInit, OnModuleDestroy } from '@nestjs/common';
import { ParserClient, type ParserEvent } from '@coopenomics/parser2';
import { IAction, IDelta } from '~/types/common';
import { RedisStreamService, StreamMessage } from '~/infrastructure/redis/redis-stream.service';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { EventsService } from '~/infrastructure/events/events.service';
import { ParserInteractor } from '~/domain/parser/interactors/parser.interactor';
import { ForkRegistryService } from '~/shared/sync/fork';
import { computeActionEventId, computeDeltaEventId, computeForkEventId } from './event-id.util';
import { mapParserActionToIAction, mapParserDeltaToIDelta } from './parser2-event.mapper';
import { config } from '~/config';
export interface BlockchainEventData {
type: string;
event?: IAction;
delta?: IDelta;
block_num?: number;
}
// Выносим исключения в конфиг или отдельный файл
const ACTION_EXCEPTIONS = {
'eosio.token': ['transfer', 'issue'],
};
/**
* Инфраструктурный сервис потребления событий блокчейна из Redis
* Читает события из Redis стрима и публикует их во внутреннюю шину событий
* Не содержит бизнес-логики, только предварительную фильтрацию
* Потребление событий блокчейна из parser2 (@coopenomics/parser2).
*
* Транспорт: единственный — ParserClient поверх Redis Stream parser2
* (`ce:parser2:<chain_id>:events`). Старый самодельный consumer поверх стрима
* `notifications` (его писал parser1, components/parser) удалён. Никаких флагов и
* параллельной работы двух движков: либо контроллер работает на parser2, либо нет.
*
* ParserClient берёт на себя то, что раньше делала ручная обвязка консьюмера:
* consumer-group, XREADGROUP/XACK, single-active-lock, recover-own-pending,
* dead-letter после N провалов, XTRIM. Контроллеру остаётся только обработка.
*
* event_id (дедуп, INV-09) вычисляется локально из полей события — формат
* action/delta (см. event-id.util.ts). parser2 кладёт свой event_id в событие,
* но контроллер ведёт собственный consumer_dedup в привычном формате.
*
* Обработчики processAction/processDelta/processFork не изменились при смене
* транспорта — маппер переводит ParserEvent → IDelta/IAction (DEC-T09).
*/
@Injectable()
export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy {
private readonly streamName = 'notifications';
private readonly consumerGroup = 'blockchain-consumer';
/**
* Стабильное имя consumer'а. Раньше было `consumer-${random}`, и каждый
* рестарт coopback создавал нового consumer'а, а прежний оставался в
* группе со своими pending-сообщениями — зомби-consumer. Со стабильным
* именем мы всегда возвращаемся к «своему» pending-списку после рестарта
* и можем его доиграть (readOwnPending на старте).
*
* Если понадобится horizontal-scaling (несколько реплик coopback'а) —
* имя можно расширить до `coopback-${HOSTNAME}` через env.
*/
private readonly consumerName = 'coopback-main';
/** Имя подписки = имя consumer-group parser2. Детерминировано по кооперативу. */
private readonly subscriptionId = `controller-${config.coopname}`;
/** Как долго pending другого consumer'а должен висеть, прежде чем его можно забрать. */
private readonly staleClaimIdleMs = 5 * 60 * 1000; // 5 минут
/** Период XAUTOCLAIM поиска stale pending. */
private readonly claimIntervalMs = 60 * 1000; // 1 минута
/** Период XTRIM MINID для освобождения памяти Redis от consumed сообщений. */
private readonly trimIntervalMs = 30 * 1000; // 30 секунд
/** Пауза перед переподключением, если поток ParserClient неожиданно упал. */
private readonly reconnectDelayMs = 5000;
private claimTimer?: NodeJS.Timeout;
private trimTimer?: NodeJS.Timeout;
private client?: ParserClient;
private running = false;
constructor(
private readonly redisStreamService: RedisStreamService,
private readonly logger: WinstonLoggerService,
private readonly eventsService: EventsService,
private readonly parserInteractor: ParserInteractor
private readonly parserInteractor: ParserInteractor,
private readonly forkRegistry: ForkRegistryService
) {
this.logger.setContext(BlockchainConsumerService.name);
}
async onModuleInit() {
this.logger.log('Инициализация сервиса потребителя блокчейна');
await this.redisStreamService.createConsumerGroup(this.streamName, this.consumerGroup);
// 1) Сначала доиграть свои pending (могли остаться после crash'а между
// handleMessage и xack). Если их нет — мгновенно пройдёт.
await this.recoverOwnPending();
// 2) Забрать pending у зомби-consumer'ов предыдущих рестартов,
// чтобы они не висели вечно. XAUTOCLAIM idempotent.
await this.reclaimStalePending();
// 3) Запустить основной consumer loop (на новые сообщения, `>`).
this.startConsuming();
// 4) Фоновые задачи: периодический XAUTOCLAIM + XTRIM MINID.
this.claimTimer = setInterval(
() => this.reclaimStalePending().catch((e) => this.logger.error(`claim tick: ${e?.message}`, e?.stack)),
this.claimIntervalMs,
);
this.trimTimer = setInterval(
() => this.trimConsumed().catch((e) => this.logger.error(`trim tick: ${e?.message}`, e?.stack)),
this.trimIntervalMs,
);
this.logger.log('Инициализация потребителя событий parser2');
this.running = true;
// Не await: цикл живёт всё время работы приложения.
void this.runConsumeLoop();
}
onModuleDestroy() {
this.logger.log('Остановка сервиса потребителя блокчейна');
if (this.claimTimer) clearInterval(this.claimTimer);
if (this.trimTimer) clearInterval(this.trimTimer);
this.redisStreamService.stopConsumer(this.streamName, this.consumerGroup, this.consumerName);
async onModuleDestroy() {
this.logger.log('Остановка потребителя событий parser2');
this.running = false;
if (this.client) {
await this.client.close().catch((e) => this.logger.error(`Ошибка close ParserClient: ${e?.message}`, e?.stack));
}
}
/**
* После рестарта контроллера читаем pending, адресованные нашему consumer'у.
* Это сообщения, которые Redis считает выданными нам, но ещё не ACK'нутыми —
* например, процесс упал после handleMessage, до xack. Доигрываем в том же
* порядке, ACK'аем — и дальше работает `>`-поток.
* Внешний цикл: держит подписку живой. Если поток ParserClient завершился с
* ошибкой (обрыв Redis и т.п.) — пауза и переподключение, пока сервис running.
*/
private async recoverOwnPending(): Promise<void> {
let total = 0;
// Читаем порциями до тех пор, пока pending не закончится.
while (true) {
const messages = await this.redisStreamService.readOwnPending(
this.streamName,
this.consumerGroup,
this.consumerName,
100,
);
if (messages.length === 0) break;
for (const msg of messages) {
try {
await this.handleMessage(msg);
await this.redisStreamService.acknowledgeMessage(this.streamName, this.consumerGroup, msg.messageId);
total += 1;
} catch (err: any) {
this.logger.error(
`recoverOwnPending: не удалось переобработать ${msg.messageId}: ${err?.message}`,
err?.stack,
);
// Оставляем pending — claim retry по idle разберёт.
}
private async runConsumeLoop(): Promise<void> {
while (this.running) {
try {
await this.consume();
} catch (err: any) {
this.logger.error(`Поток ParserClient прерван: ${err?.message}`, err?.stack);
}
if (this.client) {
await this.client.close().catch(() => undefined);
this.client = undefined;
}
if (this.running) {
await new Promise((r) => setTimeout(r, this.reconnectDelayMs));
this.logger.warn('Переподключение к parser2…');
}
}
if (total > 0) this.logger.log(`recoverOwnPending: доиграно ${total} сообщений`);
}
/**
* Забрать pending другого consumer'а, который idle > staleClaimIdleMs.
* Обычный кейс: зомби-consumer из прошлой сессии (если в БД Redis остались
* следы старого случайного имени `consumer-${random}` до этого фикса).
* Также защита от split-brain, если когда-нибудь появится несколько реплик.
* Один проход подписки. Управляем генератором вручную: it.next() подтверждает
* (XACK внутри ParserClient) успешно обработанное событие, it.throw(err) при
* ошибке обработчика запускает учёт провалов parser2 (PEL-retry / dead-letter
* после порога) — событие НЕ ACK'ается молча. Наивный `for await` тут неверен:
* проброс из тела вызывает iterator.return(), catch вокруг yield не срабатывает,
* и одна ошибка убила бы консьюмер.
*/
private async reclaimStalePending(): Promise<void> {
let claimed = 0;
try {
const messages = await this.redisStreamService.autoClaimStale(
this.streamName,
this.consumerGroup,
this.consumerName,
this.staleClaimIdleMs,
100,
);
for (const msg of messages) {
try {
await this.handleMessage(msg);
await this.redisStreamService.acknowledgeMessage(this.streamName, this.consumerGroup, msg.messageId);
claimed += 1;
} catch (err: any) {
this.logger.error(
`reclaimStalePending: ошибка ${msg.messageId}: ${err?.message}`,
err?.stack,
);
}
}
if (claimed > 0) this.logger.warn(`reclaimStalePending: перехвачено и обработано ${claimed} stale-сообщений`);
} catch (err: any) {
this.logger.error(`reclaimStalePending: ${err?.message}`, err?.stack);
}
}
/**
* XTRIM MINID: освобождаем Redis от ACK-нутых сообщений.
* Parser больше не делает XTRIM MAXLEN (burst'ы удаляли не-consumed), так что
* обрезка — обязанность consumer'а, который _точно знает_ границу.
*
* Граница = first-pending-id (если pending не пусто) ИЛИ last-generated-id
* (если ВСЕ сообщения consumed — тогда stream можно почистить полностью).
* Всё что ≤ first-pending, уже ACK'нуто кем-то в группе — безопасно.
*/
private async trimConsumed(): Promise<void> {
try {
const firstPending = await this.redisStreamService.getFirstPendingId(this.streamName, this.consumerGroup);
const trimId = firstPending ?? (await this.redisStreamService.getStreamLastId(this.streamName));
if (!trimId || trimId === '0-0') return;
await this.redisStreamService.trimUpTo(this.streamName, trimId);
} catch (err: any) {
this.logger.error(`trimConsumed: ${err?.message}`, err?.stack);
}
}
/**
* Запуск потребления сообщений из Redis стрима
*/
private async startConsuming(): Promise<void> {
this.logger.log(`Starting consumer ${this.consumerName} for stream ${this.streamName}`);
await this.redisStreamService.startConsumer(
{
stream: this.streamName,
group: this.consumerGroup,
consumer: this.consumerName,
count: 1,
block: 1000,
private async consume(): Promise<void> {
this.client = new ParserClient({
subscriptionId: this.subscriptionId,
// Без фильтров: получаем все события, фильтрация по coopname — в processDelta/processAction.
startFrom: 'last_known',
redis: {
url: `redis://${config.redis.host}:${config.redis.port}`,
password: config.redis.password || undefined,
},
this.handleMessage.bind(this)
);
chain: { id: config.blockchain.id },
// Жизненным циклом управляет NestJS (onModuleDestroy), не SIGTERM-хуки parser2.
noSignalHandlers: true,
});
this.logger.log(`Подписка parser2 "${this.subscriptionId}" на цепь ${config.blockchain.id}`);
const iterator = this.client.stream();
let result = await iterator.next();
while (!result.done && this.running) {
try {
await this.handleEvent(result.value);
result = await iterator.next(); // успех → XACK внутри ParserClient
} catch (err: any) {
this.logger.error(`Ошибка обработки события parser2: ${err?.message}`, err?.stack);
result = await iterator.throw(err); // провал → FailureTracker / dead-letter parser2
}
}
}
/**
* Обработка входящего сообщения из стрима
* Диспетчеризация события parser2 на обработчики контроллера.
* native-delta контроллер не потребляет (нет легаси-пути) — пропускаем.
*
* fork-event (Story 4.1, AC INV-09): dedup-gate в handleEvent — повторно
* доставленный fork с уже отмеченным event_id делает ранний return. Иначе
* runAll/deleteAfterBlock/saveFork выполнятся повторно, что для ForkEntity
* без UNIQUE-constraint породит дубль и нагрузку на репозитории syncer'ов.
* markApplied идёт ПОСЛЕ успешного processFork (порядок симметричен dispatch'у
* action/delta: save → mark, иначе сбой между save и mark = silent loss).
*/
private async handleMessage(message: StreamMessage): Promise<void> {
try {
const eventData: BlockchainEventData = JSON.parse(
message.fields.event || message.fields.delta || message.fields.fork || '{}'
);
if (eventData.event) {
await this.processAction(eventData.event);
} else if (eventData.delta) {
await this.processDelta(eventData.delta);
} else if (eventData.block_num !== undefined) {
await this.processFork(eventData.block_num);
} else {
this.logger.warn(`Unknown message format: ${JSON.stringify(message.fields)}`);
private async handleEvent(event: ParserEvent): Promise<void> {
switch (event.kind) {
case 'action':
return this.processAction(mapParserActionToIAction(event));
case 'delta':
return this.processDelta(mapParserDeltaToIDelta(event));
case 'fork': {
// event_id вычисляем локально в controller-формате (chain:fork:...), а НЕ
// берём event.event_id из parser2 (его формат chain:f:...) — иначе в
// consumer_dedup смешаются две формулы, и дедуп между controller и
// транспортом расползётся. См. event-id.util.ts.
const eventId = computeForkEventId(event.chain_id, event.forked_from_block, event.new_head_block_id);
if (await this.parserInteractor.isEventApplied(eventId)) {
this.logger.debug(`Fork-дубликат пропущен (no-op): ${eventId}`);
return;
}
await this.processFork(event.forked_from_block, eventId);
await this.parserInteractor.markEventApplied(eventId, event.forked_from_block);
return;
}
} catch (error: any) {
this.logger.error(`Ошибка обработки сообщения ${message.messageId}: ${error.message}`, error.stack);
throw error; // Перебрасываем ошибку чтобы сообщение не было подтверждено
case 'native-delta':
return;
default:
this.logger.warn(`Неизвестный тип события parser2: ${JSON.stringify(event)}`);
}
}
@@ -228,7 +172,7 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
* Выполняет минимальную предварительную фильтрацию, сохраняет в базу и
* с задержкой публикует событие во внутреннюю шину.
*
* Порядок writes: сохранение → ACK (по возврату из handleMessage) →
* Порядок writes: сохранение → ACK (по возврату из handleEvent) →
* отложенный emit события через ACTION_EMIT_DELAY_MS. Задержка нужна
* чтобы дельты, попавшие в стрим из того же блока что и action, успели
* пройти обработчики и прописаться в БД ДО того, как обработчики action
@@ -236,9 +180,8 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
* apprvappndx ищет appendix в capital_appendixes — без задержки гонится
* с дельтой capital::appendixes того же блока). См. задачу #53.
*
* Ошибка saveAction бросается наверх в handleMessage → сообщение НЕ
* подтверждается и остаётся pending (consumer перечитает его через
* XCLAIM/re-delivery). Emit'ится только то, что успешно сохранено.
* Ошибка saveAction бросается наверх в consume → событие НЕ подтверждается
* (iterator.throw → parser2 учитывает провал). Emit'ится только сохранённое.
*/
private async processAction(action: IAction): Promise<void> {
if (action.receiver != action.account) {
@@ -248,14 +191,6 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
await this.processActionDelayed(action);
}
/**
* Задержка emit'а action-события — даёт время дельтам того же блока
* (capital_appendixes, capital_projects, ...) сохраниться раньше, чем
* обработчики action полезут их читать. Подобрана опытно: при <1.5с
* на быстрой машине race ещё происходит, при ~3с гонок не наблюдаем.
*/
private static readonly ACTION_EMIT_DELAY_MS = 3000;
private async processActionDelayed(action: IAction): Promise<void> {
// Проверяем, является ли действие исключением
const isException = this.isActionException(action.account, action.name);
@@ -266,6 +201,14 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
return;
}
// Idempotency: признак уникальности события (INV-09). Повторно доставленное
// событие с уже отмеченным event_id игнорируется как no-op.
const eventId = computeActionEventId(action);
if (await this.parserInteractor.isEventApplied(eventId)) {
this.logger.debug(`Action-дубликат пропущен (no-op): ${eventId}`);
return;
}
try {
// Сохраняем действие в базу данных через интерактор
await this.parserInteractor.saveAction(action);
@@ -274,13 +217,19 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
);
} catch (error: any) {
this.logger.error(`Не удалось сохранить действие ${action.account}::${action.name}: ${error.message}`, error.stack);
throw error; // Перебрасываем ошибку чтобы сообщение не было подтверждено
throw error; // Перебрасываем ошибку чтобы событие не было подтверждено
}
// Публикуем событие с задержкой — пусть сначала прокатятся дельты этого же блока.
// saveAction уже выполнен, так что данные не потеряем; задерживаем только emit.
// Метка в consumer_dedup ПОСЛЕ save, ДО отложенного emit.
// block_num пишем для последующего deleteAfterBlock на форке (Story 4.1).
await this.parserInteractor.markEventApplied(eventId, action.block_num);
// Публикуем событие с задержкой — пусть сначала прокатятся дельты этого же блока
// (capital_appendixes, capital_projects, ...), чтобы обработчики action видели
// уже персистентное состояние. saveAction уже выполнен, так что данные не
// потеряем; задерживаем только emit. Задержка вынесена в конфиг (DEC-007).
const eventName = `action::${action.account}::${action.name}`;
const delayMs = BlockchainConsumerService.ACTION_EMIT_DELAY_MS;
const delayMs = config.blockchain.action_emit_delay_ms;
setTimeout(() => {
this.eventsService.emit(eventName, action);
this.logger.debug(
@@ -304,10 +253,8 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
* Обработка дельты (delta) из блокчейна
* Выполняет минимальную предварительную фильтрацию, сохраняет в базу и публикует событие во внутреннюю шину.
*
* Ошибка saveDelta поднимается наверх в handleMessage → сообщение остаётся
* pending в consumer group. См. processAction — тот же контракт.
* Раньше здесь был try/catch, который глотал ошибки и log-only'ил их: это
* приводило к silent data loss (ACK шёл, дельта в PG не попадала).
* Ошибка saveDelta поднимается наверх в consume → событие остаётся pending в
* consumer-group parser2 (iterator.throw). См. processAction — тот же контракт.
*/
private async processDelta(delta: IDelta): Promise<void> {
this.logger.debug(`Обработка дельты: ${delta.table} от ${delta.code}`);
@@ -331,15 +278,28 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
return;
}
// Idempotency: признак уникальности события (INV-09).
const eventId = computeDeltaEventId(delta);
if (await this.parserInteractor.isEventApplied(eventId)) {
this.logger.debug(`Дельта-дубликат пропущена (no-op): ${eventId}`);
return;
}
try {
// Сохраняем дельту в базу данных через интерактор
await this.parserInteractor.saveDelta(delta);
this.logger.log(`Дельта сохранена в базу: ${delta.code}::${delta.table} с primary_key ${delta.primary_key}`);
} catch (error: any) {
this.logger.error(`Не удалось сохранить дельту ${delta.code}::${delta.table}: ${error.message}`, error.stack);
throw error; // Перебрасываем ошибку чтобы сообщение не было подтверждено
throw error; // Перебрасываем ошибку чтобы событие не было подтверждено
}
// Метка в consumer_dedup ПОСЛЕ save. Если markApplied упадёт — consume
// пробросит ошибку, событие останется pending и переиграется (saveDelta
// идемпотентен через block_num-guard, mark — через ON CONFLICT DO NOTHING).
// block_num пишем для последующего deleteAfterBlock на форке (Story 4.1).
await this.parserInteractor.markEventApplied(eventId, delta.block_num);
// Публикуем событие во внутреннюю шину с типизированным именем
const eventName = `delta::${delta.code}::${delta.table}`;
this.eventsService.emit(eventName, delta);
@@ -348,28 +308,39 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
}
/**
* Обработка форка (fork) из блокчейна
* Сохраняет форк в базу и публикует событие форка во внутреннюю шину
* Обработка форка (fork) из блокчейна (ADR-005).
*
* Порядок шагов (контрактный):
* 1) ForkRegistry.runAll(blockNum) — sequential откат сущностей всех syncer'ов.
* Любая ошибка re-throw, parser2 не ACK'нет, повторная доставка пересыграет.
* 2) consumer_dedup.deleteAfterBlock(blockNum) — очистка дедупа отрезанной ветки.
* Делается ПОСЛЕ успешного rollback (иначе при сбое syncer'ов мы потеряем
* возможность повторить весь форк по тому же event_id).
* 3) saveFork(blockNum) — фиксация форка для аудита и future-pool re-submit (Epic 5).
*
* INV-T03: к моменту, когда handleEvent resolves и parser2 берёт следующее событие,
* вся цепочка rollback завершена (sequential XREADGROUP = natural barrier).
*/
private async processFork(block_num: number): Promise<void> {
this.logger.debug(`Обработка форка на блоке: ${block_num}`);
private async processFork(block_num: number, forkEventId?: string | null): Promise<void> {
this.logger.log(`Обработка форка на блоке ${block_num} (eventId=${forkEventId ?? 'n/a'}): запуск ForkRegistry rollback`);
try {
// Сохраняем форк в базу данных через интерактор
await this.parserInteractor.saveFork({
chain_id: config.blockchain.id, // Используем chain id из конфига
block_num: block_num,
});
this.logger.debug(`Форк сохранен в базу данных на блоке: ${block_num}`);
} catch (error: any) {
this.logger.error(`Не удалось сохранить форк на блоке ${block_num}: ${error.message}`, error.stack);
throw error; // Перебрасываем ошибку чтобы сообщение не было подтверждено
}
// 1. Sequential rollback всех зарегистрированных syncer'ов.
// Story 4.4: forkEventId пробрасывается syncer'ам — они кладут его в архив
// invalidated_entities для forensic-группировки.
await this.forkRegistry.runAll(block_num, forkEventId);
this.logger.debug(`ForkRegistry: rollback завершён для ${this.forkRegistry.size()} syncer(s)`);
// Публикуем событие во внутреннюю шину с типизированным именем
const eventName = `fork::${block_num}`;
this.eventsService.emit(eventName, { block_num });
// 2. Очистка consumer_dedup для блоков отрезанной ветки.
const purged = await this.parserInteractor.deleteDedupAfterBlock(block_num);
this.logger.debug(`consumer_dedup: удалено ${purged} записей с block_num > ${block_num}`);
this.logger.debug(`Форк опубликован в событийную шину: ${eventName}`);
// 3. Фиксация форка для аудита.
await this.parserInteractor.saveFork({
chain_id: config.blockchain.id,
block_num: block_num,
});
this.logger.debug(`Форк сохранён в БД на блоке ${block_num}`);
this.logger.log(`Форк обработан на блоке ${block_num}`);
}
}
@@ -31,6 +31,7 @@ import { LEDGER2_BLOCKCHAIN_PORT } from '~/domain/ledger2/ports/ledger2-blockcha
import { SovietContractInfoService } from './services/soviet-contract-info.service';
import { WalletContractInfoService } from './services/wallet-contract-info.service';
import { Ledger2ContractInfoService } from './services/ledger2-contract-info.service';
import { BlockchainArchiveRetentionService } from '~/shared/sync/services/blockchain-archive-retention.service';
@Global()
@Module({
@@ -92,6 +93,7 @@ import { Ledger2ContractInfoService } from './services/ledger2-contract-info.ser
SovietContractInfoService,
WalletContractInfoService,
Ledger2ContractInfoService,
BlockchainArchiveRetentionService,
],
exports: [
BlockchainService,
@@ -112,6 +114,7 @@ import { Ledger2ContractInfoService } from './services/ledger2-contract-info.ser
SovietContractInfoService,
WalletContractInfoService,
Ledger2ContractInfoService,
BlockchainArchiveRetentionService,
],
})
export class BlockchainModule {}
@@ -0,0 +1,48 @@
import { IAction, IDelta } from '~/types/common';
/**
* Локальное вычисление event_id — основа идемпотентности (INV-09).
*
* Формат: `${chain}:${kind}:${block_num}:${block_id_short}:${natural_key}`, где
* kind = action|delta|fork. Контроллер ведёт собственный consumer_dedup в этом
* формате (parser2 кладёт свой event_id в событие — другая формула
* `chain:a:...`/`chain:d:...`/`chain:f:...`, мы её не используем). Дедуп
* безусловный: повторно доставленное событие с уже отмеченным event_id
* игнорируется как no-op. Полный «kind» — чтобы человек, глядя в consumer_dedup,
* сразу видел тип события без расшифровки префикса.
*/
/** Длина короткого префикса block_id в event_id. */
const BLOCK_ID_SHORT_LEN = 8;
function shortBlockId(blockId: string | undefined): string {
return (blockId ?? '').slice(0, BLOCK_ID_SHORT_LEN);
}
/**
* event_id дельты. natural_key = code:scope:table:primary_key — естественная
* идентичность строки таблицы в конкретном блоке.
*/
export function computeDeltaEventId(delta: IDelta): string {
const naturalKey = `${delta.code}:${delta.scope}:${delta.table}:${delta.primary_key}`;
return `${delta.chain_id}:delta:${delta.block_num}:${shortBlockId(delta.block_id)}:${naturalKey}`;
}
/**
* event_id действия. natural_key = global_sequence — глобально-монотонный
* идентификатор action'а в цепи, уникален сам по себе.
*/
export function computeActionEventId(action: IAction): string {
return `${action.chain_id}:action:${action.block_num}:${shortBlockId(action.block_id)}:${action.global_sequence}`;
}
/**
* event_id форка. natural_key — short new_head_block_id (новая голова цепи
* после rollback). Симметрично action/delta-формуле: `chain:fork:block_num:short_id`.
* block_num = forked_from_block (последний безопасный блок до отката).
* short_id различает форк-эпохи — две разных «новых ветки» того же forked_from_block
* дают разный event_id, что корректно с точки зрения retry / breach-сценариев.
*/
export function computeForkEventId(chainId: string, forkedFromBlock: number, newHeadBlockId: string): string {
return `${chainId}:fork:${forkedFromBlock}:${shortBlockId(newHeadBlockId)}`;
}
@@ -0,0 +1,83 @@
import type { ActionEvent, DeltaEvent } from '@coopenomics/parser2';
import { IAction, IDelta } from '~/types/common';
/**
* Преобразование событий parser2 (ParserEvent) во внутренние формы контроллера
* IDelta / IAction. Транспорт сменился (parser1 Redis-стрим → parser2 ParserClient),
* но обработчики (processDelta/processAction → syncer'ы) остаются прежними: они
* работают с IDelta/IAction. Маппер — единственная точка перевода (DEC-T09).
*
* Все поля действия — РЕАЛЬНЫЕ из SHiP-трейса (parser2 их отдаёт): transaction_id,
* creator_action_ordinal, receipt (с auth_sequence), console, elapsed, context_free,
* account_ram_deltas. Это полный паритет с тем, что давал parser1, — ledger2
* (cross-link родительского apply по transaction_id + action_ordinal) и
* blockchain-explorer работают без потерь.
*
* bigint-поля (global_sequence, receipt.*Sequence) приходят по проводу строками
* (parser2 сериализует bigint→string), поэтому String() безопасен и для bigint, и
* для string.
*/
export function mapParserDeltaToIDelta(event: DeltaEvent): IDelta {
return {
chain_id: event.chain_id,
block_num: event.block_num,
block_id: event.block_id,
block_time: event.block_time,
present: event.present,
code: event.code,
scope: event.scope,
table: event.table,
primary_key: event.primary_key,
value: event.value,
};
}
export function mapParserActionToIAction(event: ActionEvent): IAction {
const globalSequence = String(event.global_sequence);
const r = event.receipt;
const receipt = r
? {
receiver: r.receiver,
act_digest: r.actDigest,
global_sequence: String(r.globalSequence),
recv_sequence: String(r.recvSequence),
auth_sequence: r.authSequence.map((s) => ({ account: s.account, sequence: String(s.sequence) })),
code_sequence: r.codeSequence,
abi_sequence: r.abiSequence,
}
: {
// Трассировки нет — receipt не null (read-path explorer'а и фильтр
// notification по receipt.receiver не должны падать).
receiver: event.account,
act_digest: '',
global_sequence: globalSequence,
recv_sequence: '0',
auth_sequence: [],
code_sequence: 0,
abi_sequence: 0,
};
return {
transaction_id: event.transaction_id,
account: event.account,
block_num: event.block_num,
block_id: event.block_id,
block_time: event.block_time,
chain_id: event.chain_id,
name: event.name,
// parser2 эмитит action уже единожды (дедуп по global_sequence); receiver=account,
// чтобы guard processAction (receiver != account → skip) пропускал событие.
receiver: receipt.receiver,
authorization: event.authorization.map((a) => ({ actor: a.actor, permission: a.permission })),
data: event.data,
action_ordinal: event.action_ordinal,
global_sequence: globalSequence,
account_ram_deltas: event.account_ram_deltas.map((d) => ({ account: d.account, delta: d.delta })),
console: event.console,
receipt,
creator_action_ordinal: event.creator_action_ordinal,
context_free: event.context_free,
elapsed: event.elapsed,
};
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '~/shared/services/abstract-entity-sync.service';
import { AgreementDomainEntity } from '~/domain/agreement/entities/agreement.entity';
@@ -49,13 +49,4 @@ export class AgreementSyncService
this.logger.debug('Сервис синхронизации соглашений полностью инициализирован с подписками на паттерны');
}
/**
* Обработка форков для соглашений
* Теперь подписывается на все форки независимо от контракта
*/
@OnEvent('fork::*')
async handleAgreementFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '~/shared/services/abstract-entity-sync.service';
import { UserAgreementDomainEntity } from '~/domain/wallet/entities/user-agreement-domain.entity';
@@ -49,8 +49,4 @@ export class UserAgreementSyncService
});
}
@OnEvent('fork::*')
async handleUserAgreementFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -1,5 +1,5 @@
import { Injectable, OnModuleInit, Inject } from '@nestjs/common';
import { OnEvent, EventEmitter2 } from '@nestjs/event-emitter';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { AbstractEntitySyncService } from '~/shared/services/abstract-entity-sync.service';
import { UserWalletDomainEntity } from '~/domain/wallet/entities/user-wallet-domain.entity';
@@ -51,8 +51,4 @@ export class UserWalletSyncService
});
}
@OnEvent('fork::*')
async handleUserWalletFork(forkData: { block_num: number }): Promise<void> {
await this.handleFork(forkData.block_num);
}
}
@@ -0,0 +1,33 @@
import { Entity, PrimaryColumn, Column, Index, CreateDateColumn } from 'typeorm';
/**
* Список уже применённых событий блокчейна — фундамент идемпотентности
* (Story 2.1, INV-09, NFR5). Повторно доставленное событие с уже отмеченным
* event_id распознаётся dispatch-путём как no-op (Story 2.3), что защищает от
* двойного применения при at-least-once доставке Redis Streams.
*
* Без тяжёлых constraints (RT-03): event_id — PK, индекс по applied_at нужен
* для retention-очистки старых меток (ориентир OQ-T04: Rollback Horizon × 2).
*
* event_id вычисляется локально по формуле parser2 (Story 2.2); после миграции
* на parser2 (Epic 3) формула сверяется с авторитетной из движка.
*
* Story 4.1: добавлена колонка block_num — для очистки дедупа на форке
* (deleteAfterBlock). Колонка nullable: старые записи (до Epic 4) её не имеют,
* они доживут до своего retention по applied_at и не будут попадать под
* WHERE block_num > N (PG NULL-сравнения возвращают unknown, строка не удалится).
* Новые записи всегда несут block_num.
*/
@Entity('consumer_dedup')
export class ConsumerDedupEntity {
@PrimaryColumn({ type: 'varchar', length: 512 })
event_id!: string;
@Index('idx_consumer_dedup_applied_at')
@CreateDateColumn({ type: 'timestamptz' })
applied_at!: Date;
@Index('idx_consumer_dedup_block_num')
@Column({ type: 'bigint', nullable: true })
block_num!: string | null;
}
@@ -0,0 +1,62 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import type { ConsumerDedupRepositoryPort } from '~/domain/parser/ports/consumer-dedup-repository.port';
import { ConsumerDedupEntity } from '../entities/consumer-dedup.entity';
/**
* TypeORM-реализация списка применённых событий (Story 2.1).
*/
@Injectable()
export class TypeOrmConsumerDedupRepository implements ConsumerDedupRepositoryPort {
constructor(
@InjectRepository(ConsumerDedupEntity)
private readonly repository: Repository<ConsumerDedupEntity>
) {}
async isApplied(eventId: string): Promise<boolean> {
const found = await this.repository.findOne({
where: { event_id: eventId },
select: { event_id: true },
});
return found != null;
}
async markApplied(eventId: string, blockNum?: number): Promise<void> {
// ON CONFLICT DO NOTHING: повторная отметка (re-delivery / crash-recovery)
// не должна падать на нарушении PK.
// block_num хранится как bigint → string в TypeORM; NULL для legacy-вызовов без блока.
await this.repository
.createQueryBuilder()
.insert()
.into(ConsumerDedupEntity)
.values({
event_id: eventId,
block_num: typeof blockNum === 'number' ? String(blockNum) : null,
})
.orIgnore()
.execute();
}
async deleteOlderThan(cutoff: Date): Promise<number> {
const result = await this.repository
.createQueryBuilder()
.delete()
.from(ConsumerDedupEntity)
.where('applied_at < :cutoff', { cutoff })
.execute();
return result.affected ?? 0;
}
async deleteAfterBlock(blockNum: number): Promise<number> {
// Сравнение bigint > N — оба операнда числа на стороне PG; параметр приводится в bigint.
// NULL-записи (legacy до Story 4.1) не попадают: NULL > N = UNKNOWN, строка не удаляется.
const result = await this.repository
.createQueryBuilder()
.delete()
.from(ConsumerDedupEntity)
.where('block_num > :blockNum', { blockNum })
.execute();
return result.affected ?? 0;
}
}
@@ -42,6 +42,10 @@ import { SyncStateEntity } from './entities/sync-state.entity';
import { EntityVersionTypeormEntity } from '~/shared/sync/entities/entity-version.typeorm-entity';
import { EntityVersionRepository } from '~/shared/sync/repositories/entity-version.repository';
import { EntityVersioningService } from '~/shared/sync/services/entity-versioning.service';
import { InvalidatedEntityTypeormEntity } from '~/shared/sync/entities/invalidated-entity.typeorm-entity';
import { InvalidatedEntityVersionTypeormEntity } from '~/shared/sync/entities/invalidated-entity-version.typeorm-entity';
import { InvalidatedEntityRepository } from '~/shared/sync/repositories/invalidated-entity.repository';
import { InvalidatedEntityVersionRepository } from '~/shared/sync/repositories/invalidated-entity-version.repository';
import { ACTION_REPOSITORY_PORT } from '~/domain/parser/ports/action-repository.port';
import { DELTA_REPOSITORY_PORT } from '~/domain/parser/ports/delta-repository.port';
import { FORK_REPOSITORY_PORT } from '~/domain/parser/ports/fork-repository.port';
@@ -50,6 +54,9 @@ import { TypeOrmActionRepository } from './repositories/typeorm-action.repositor
import { TypeOrmDeltaRepository } from './repositories/typeorm-delta.repository';
import { TypeOrmForkRepository } from './repositories/typeorm-fork.repository';
import { TypeOrmSyncStateRepository } from './repositories/typeorm-sync-state.repository';
import { ConsumerDedupEntity } from './entities/consumer-dedup.entity';
import { CONSUMER_DEDUP_REPOSITORY_PORT } from '~/domain/parser/ports/consumer-dedup-repository.port';
import { TypeOrmConsumerDedupRepository } from './repositories/typeorm-consumer-dedup.repository';
import { SettingsEntity } from './entities/settings.entity';
import { SETTINGS_REPOSITORY } from '~/domain/settings/repositories/settings.repository';
import { SettingsTypeormRepository } from './repositories/settings.typeorm-repository';
@@ -122,7 +129,10 @@ import { UserWalletIndexInitializer } from './blockchain/services/user-wallet-in
DeltaEntity,
ForkEntity,
SyncStateEntity,
ConsumerDedupEntity,
EntityVersionTypeormEntity,
InvalidatedEntityTypeormEntity,
InvalidatedEntityVersionTypeormEntity,
SettingsEntity,
TokenEntity,
UserEntity,
@@ -197,6 +207,10 @@ import { UserWalletIndexInitializer } from './blockchain/services/user-wallet-in
provide: SYNC_STATE_REPOSITORY_PORT,
useClass: TypeOrmSyncStateRepository,
},
{
provide: CONSUMER_DEDUP_REPOSITORY_PORT,
useClass: TypeOrmConsumerDedupRepository,
},
{
provide: SETTINGS_REPOSITORY,
useClass: SettingsTypeormRepository,
@@ -251,6 +265,8 @@ import { UserWalletIndexInitializer } from './blockchain/services/user-wallet-in
UserWalletIndexInitializer,
EntityVersionRepository,
EntityVersioningService,
InvalidatedEntityRepository,
InvalidatedEntityVersionRepository,
],
exports: [
NestTypeOrmModule,
@@ -269,6 +285,7 @@ import { UserWalletIndexInitializer } from './blockchain/services/user-wallet-in
DELTA_REPOSITORY_PORT,
FORK_REPOSITORY_PORT,
SYNC_STATE_REPOSITORY_PORT,
CONSUMER_DEDUP_REPOSITORY_PORT,
SETTINGS_REPOSITORY,
TOKEN_REPOSITORY,
USER_REPOSITORY,
@@ -286,6 +303,8 @@ import { UserWalletIndexInitializer } from './blockchain/services/user-wallet-in
UserWalletSyncService,
EntityVersionRepository,
EntityVersioningService,
InvalidatedEntityRepository,
InvalidatedEntityVersionRepository,
],
})
export class TypeOrmModule {}
@@ -17,4 +17,34 @@ export class EventsService {
emit(eventName: string, data: any): void {
this.eventEmitter.emit(eventName, data);
}
/**
* Публикация с ожиданием завершения всех async-обработчиков.
* В отличие от emit (fire-and-forget) дожидается, пока @OnEvent-листенеры
* (в т.ч. async) отработают. Если любой обработчик бросает — промис
* отклоняется (ошибку не глотаем).
*/
async emitAsync(eventName: string, data: any): Promise<unknown[]> {
return this.eventEmitter.emitAsync(eventName, data);
}
/**
* Барьер: публикация с ожиданием обработчиков, но не дольше timeoutMs.
* Возвращает true, если все обработчики завершились в срок; false — если
* сработал TTL force-resume (какой-то обработчик завис). Используется для
* обработки форка: дождаться откатов до продолжения consumer'а, но не
* блокировать поток навсегда из-за зависшего синкера (Story 1.3).
*/
async emitAsyncWithTimeout(eventName: string, data: any, timeoutMs: number): Promise<boolean> {
let timer: NodeJS.Timeout | undefined;
const timeout = new Promise<false>((resolve) => {
timer = setTimeout(() => resolve(false), timeoutMs);
});
try {
const settled = this.eventEmitter.emitAsync(eventName, data).then(() => true);
return await Promise.race([settled, timeout]);
} finally {
if (timer) clearTimeout(timer);
}
}
}
@@ -1,292 +0,0 @@
import { Inject, Injectable, OnModuleDestroy } from '@nestjs/common';
import Redis from 'ioredis';
import { REDIS_PROVIDER } from './redis.provider';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
export interface StreamMessage {
messageId: string;
fields: Record<string, string>;
}
export interface StreamConsumerOptions {
stream: string;
group: string;
consumer: string;
count?: number;
block?: number;
}
/**
* Сервис для работы с Redis Streams
* Обеспечивает надежное потребление событий от parser
*/
@Injectable()
export class RedisStreamService implements OnModuleDestroy {
private consumers: Map<string, boolean> = new Map();
constructor(
@Inject(REDIS_PROVIDER)
private readonly redisClient: { subscriber: Redis; publisher: Redis; streamManager: Redis; streamReader: Redis },
private readonly logger: WinstonLoggerService
) {
this.logger.setContext(RedisStreamService.name);
}
onModuleDestroy() {
// Останавливаем всех потребителей
this.consumers.forEach((_, key) => {
this.consumers.set(key, false);
});
// Закрываем дополнительные соединения
this.redisClient.streamManager.quit();
this.redisClient.streamReader.quit();
}
/**
* Создание группы потребителей
*/
async createConsumerGroup(stream: string, group: string, startId = '0'): Promise<void> {
try {
await this.redisClient.streamManager.xgroup('CREATE', stream, group, startId, 'MKSTREAM');
this.logger.log(`Consumer group created: ${group} for stream: ${stream}`);
} catch (error: any) {
if (error.message.includes('BUSYGROUP')) {
this.logger.log(`Consumer group ${group} already exists for stream ${stream}`);
} else {
this.logger.error(`Error creating consumer group: ${error.message}`);
throw error;
}
}
}
/**
* Потребление сообщений из stream
*/
async startConsumer(
options: StreamConsumerOptions,
messageHandler: (message: StreamMessage) => Promise<void>
): Promise<void> {
const { stream, group, consumer, count = 1, block = 1000 } = options;
const consumerKey = `${stream}:${group}:${consumer}`;
// Отмечаем потребителя как активного
this.consumers.set(consumerKey, true);
this.logger.log(`Starting consumer: ${consumerKey}`);
while (this.consumers.get(consumerKey)) {
try {
const result: any = await this.redisClient.streamReader.xreadgroup(
'GROUP',
group,
consumer,
'COUNT',
count.toString(),
'BLOCK',
block.toString(),
'STREAMS',
stream,
'>'
);
if (result && result.length > 0) {
for (const [streamName, messages] of result) {
if (streamName === stream) {
for (const [messageId, messageData] of messages) {
const fields: Record<string, string> = {};
// Преобразуем массив [key, value, key, value] в объект
for (let i = 0; i < messageData.length; i += 2) {
fields[messageData[i]] = messageData[i + 1];
}
try {
await messageHandler({ messageId, fields });
await this.acknowledgeMessage(stream, group, messageId);
} catch (error: any) {
this.logger.error(`Error processing message ${messageId}: ${error.message}`);
// Не подтверждаем сообщение при ошибке - оно останется в pending
}
}
}
}
}
} catch (error: any) {
if (this.consumers.get(consumerKey)) {
this.logger.error(`Consumer ${consumerKey} error: ${error.message}`);
// Небольшая пауза перед повторной попыткой
await new Promise((resolve) => setTimeout(resolve, 5000));
}
}
}
this.logger.log(`Consumer stopped: ${consumerKey}`);
}
/**
* Подтверждение обработки сообщения
*/
async acknowledgeMessage(stream: string, group: string, messageId: string): Promise<void> {
try {
await this.redisClient.streamManager.xack(stream, group, messageId);
this.logger.debug(`Message acknowledged: ${messageId} in stream ${stream}`);
} catch (error: any) {
this.logger.error(`Error acknowledging message ${messageId} in stream ${stream}, group ${group}: ${error.message}`);
throw error;
}
}
/**
* Остановка потребителя
*/
stopConsumer(stream: string, group: string, consumer: string): void {
const consumerKey = `${stream}:${group}:${consumer}`;
this.consumers.set(consumerKey, false);
this.logger.log(`Stopping consumer: ${consumerKey}`);
}
/**
* Получение информации о pending сообщениях
*/
async getPendingMessages(stream: string, group: string): Promise<any> {
try {
return await this.redisClient.streamManager.xpending(stream, group);
} catch (error: any) {
this.logger.error(`Error getting pending messages: ${error.message}`);
throw error;
}
}
/**
* Получение информации о группах потребителей
*/
async getConsumerGroups(stream: string): Promise<any> {
try {
return await this.redisClient.streamManager.xinfo('GROUPS', stream);
} catch (error: any) {
this.logger.error(`Error getting consumer groups: ${error.message}`);
throw error;
}
}
/**
* Прочитать pending-сообщения конкретного consumer'а (uncacknowledged).
* Используется при старте для восстановления работы после рестарта:
* сообщения, которые были выданы этому consumer'у, но не подтверждены
* (например, coopback упал между handleMessage и acknowledgeMessage),
* должны быть перечитаны и допоставлены.
*
* XREADGROUP с ID '0' означает «все pending этого consumer'а»,
* а не новые — этим отличается от обычного '>'.
*/
async readOwnPending(stream: string, group: string, consumer: string, count = 100): Promise<StreamMessage[]> {
const result: any = await this.redisClient.streamReader.xreadgroup(
'GROUP',
group,
consumer,
'COUNT',
count.toString(),
'STREAMS',
stream,
'0',
);
return this.parseStreamResult(result, stream);
}
/**
* XAUTOCLAIM: переназначить себе pending-сообщения других consumer'ов,
* idle которых превысил minIdleMs. Это защита от зомби-consumer'ов:
* предыдущий рестарт coopback оставил consumer "consumer-abc123" со своими
* pending; новый consumer "coopback-main" заберёт их через XAUTOCLAIM.
*
* Redis 6.2+. Возвращает [nextCursor, claimedEntries, deletedIds?].
*/
async autoClaimStale(
stream: string,
group: string,
consumer: string,
minIdleMs: number,
count = 100,
): Promise<StreamMessage[]> {
const result: any = await this.redisClient.streamManager.xautoclaim(
stream,
group,
consumer,
minIdleMs.toString(),
'0',
'COUNT',
count.toString(),
);
// Redis 7+: [cursor, entries, deletedIds]. Redis 6.2: [cursor, entries].
if (!Array.isArray(result) || result.length < 2) return [];
const entries = result[1] as any[];
const messages: StreamMessage[] = [];
for (const [messageId, messageData] of entries) {
if (!Array.isArray(messageData)) continue; // deleted entry
const fields: Record<string, string> = {};
for (let i = 0; i < messageData.length; i += 2) {
fields[messageData[i]] = messageData[i + 1];
}
messages.push({ messageId, fields });
}
return messages;
}
/**
* XTRIM MINID: удалить из stream записи с ID меньше minId.
* Используется controller'ом для освобождения памяти Redis от уже
* consumed сообщений. minId должен быть ≤ first-pending-id, иначе
* удалим ещё не обработанные сообщения — всё, что в pending, защищено
* этим граничным условием.
*
* `~` (approximate trim) — Redis удаляет эффективно по блокам radix tree,
* реально удалённых записей может быть чуть больше или чуть меньше minId.
*/
async trimUpTo(stream: string, minId: string): Promise<number> {
return (await this.redisClient.streamManager.xtrim(stream, 'MINID', '~', minId)) as number;
}
/**
* Минимальный pending ID в consumer-group (или null, если pending пусто).
* XPENDING summary form возвращает [count, minId, maxId, consumers].
*/
async getFirstPendingId(stream: string, group: string): Promise<string | null> {
const summary: any = await this.redisClient.streamManager.xpending(stream, group);
if (!Array.isArray(summary) || !summary[0]) return null;
const count = Number(summary[0]);
if (count === 0) return null;
return (summary[1] as string) || null;
}
/**
* Последний ID в stream (или '0-0', если stream пуст).
* Нужен для trim'а когда pending пусто и last-delivered-id бесполезен:
* безопасно обрезать всё до нынешнего конца stream'а только если ВСЕ
* сообщения consumed. Проверка pending.count=0 — гарантия этого.
*/
async getStreamLastId(stream: string): Promise<string> {
const info: any = await this.redisClient.streamManager.xinfo('STREAM', stream);
if (!Array.isArray(info)) return '0-0';
// XINFO STREAM возвращает массив [key, value, key, value, ...].
for (let i = 0; i < info.length; i += 2) {
if (info[i] === 'last-generated-id') return String(info[i + 1] || '0-0');
}
return '0-0';
}
private parseStreamResult(result: any, expectedStream: string): StreamMessage[] {
if (!result || !Array.isArray(result) || result.length === 0) return [];
const messages: StreamMessage[] = [];
for (const [streamName, streamMessages] of result) {
if (streamName !== expectedStream) continue;
for (const [messageId, messageData] of streamMessages) {
const fields: Record<string, string> = {};
for (let i = 0; i < messageData.length; i += 2) {
fields[messageData[i]] = messageData[i + 1];
}
messages.push({ messageId, fields });
}
}
return messages;
}
}
@@ -2,7 +2,6 @@
import { Module } from '@nestjs/common';
import { RedisService } from './redis.service';
import { RedisStreamService } from './redis-stream.service';
import { RedisProvider, REDIS_PROVIDER } from './redis.provider';
import { REDIS_PORT } from '~/domain/common/ports/redis.port';
@@ -10,12 +9,11 @@ import { REDIS_PORT } from '~/domain/common/ports/redis.port';
providers: [
RedisProvider,
RedisService,
RedisStreamService,
{
provide: REDIS_PORT,
useClass: RedisService,
},
],
exports: [RedisService, RedisStreamService, REDIS_PORT, REDIS_PROVIDER],
exports: [RedisService, REDIS_PORT, REDIS_PROVIDER],
})
export class RedisModule {}
@@ -78,6 +78,18 @@ export interface IBlockchainSyncRepository<TEntity extends IBlockchainSynchroniz
/** Восстановить сущности из версий после форка */
restoreFromVersions?(forkBlockNum: number): Promise<void>;
/**
* Story 4.4: атомарно перенести live-сущности WHERE block_num > forkBlockNum в архив
* (invalidated_entities) и удалить из исходной таблицы. Возвращает count.
*/
archiveInvalidatedSince?(forkBlockNum: number, forkEventId?: string | null): Promise<number>;
/**
* Story 4.4: атомарно перенести версии этой entity_table WHERE block_num > forkBlockNum
* в архив (invalidated_entity_versions) и удалить из entity_versions. Возвращает count.
*/
archiveInvalidatedVersionsSince?(forkBlockNum: number, forkEventId?: string | null): Promise<number>;
}
/**
@@ -7,6 +7,9 @@ import type {
IBlockchainSyncRepository,
ISyncResult,
} from '~/shared/interfaces/blockchain-sync.interface';
import { FORK_AWARE_MARKER, type IForkAwareSyncer } from '~/shared/sync/fork';
import { UnsupportedContractVersionError } from '~/shared/sync/errors/unsupported-contract-version.error';
import config from '~/config/config';
/**
* Абстрактный сервис для синхронизации сущностей с блокчейном
@@ -14,12 +17,23 @@ import type {
* Предоставляет базовую логику для:
* - Обработки дельт блокчейна
* - Создания/обновления сущностей
* - Обработки форков
* - Обработки форков (Story 4.1: реализует IForkAwareSyncer ForkRegistryService
* собирает наследников через DiscoveryService по symbol-маркеру и обходит
* sequential при форке)
*/
@Injectable()
export abstract class AbstractEntitySyncService<TEntity extends IBlockchainSynchronizable, TBlockchainData = any> {
export abstract class AbstractEntitySyncService<TEntity extends IBlockchainSynchronizable, TBlockchainData = any>
implements IForkAwareSyncer
{
protected abstract readonly entityName: string;
/**
* Symbol-маркер для ForkRegistryService (Story 4.1). Все 20+ наследников
* автоматически попадают в реестр через bootstrap-сканирование Discovery
* без правок их onModuleInit.
*/
readonly [FORK_AWARE_MARKER] = true;
constructor(
protected readonly repository: IBlockchainSyncRepository<TEntity>,
protected readonly mapper: IBlockchainDeltaMapper<TBlockchainData>,
@@ -44,7 +58,21 @@ export abstract class AbstractEntitySyncService<TEntity extends IBlockchainSynch
// Маппинг дельты в блокчейн-данные
const blockchainData = this.mapper.mapDeltaToBlockchainData(delta);
if (!blockchainData) {
this.logger.warn(`Failed to map delta to blockchain data for ${this.entityName} ${syncValue}`);
// Story 6.5: silent loss заменён на audit-trail error. В strict-mode дополнительно
// throw UnsupportedContractVersionError — парсер не ACK'нет дельту, dead-letter сработает.
const ctx = {
contract: (delta as any).contract ?? (delta as any).code,
table: (delta as any).table,
primary_key: (delta as any).primary_key,
block_num: Number((delta as any).block_num),
};
this.logger.error(
`UNSUPPORTED_CONTRACT_VERSION: mapDeltaToBlockchainData returned null for ${this.entityName} ${syncValue}`,
{ entity: this.entityName, syncValue, ...ctx }
);
if (config.blockchain.unsupported_version_strict) {
throw new UnsupportedContractVersionError(this.entityName, ctx);
}
return null;
}
@@ -54,6 +82,9 @@ export abstract class AbstractEntitySyncService<TEntity extends IBlockchainSynch
// Обработка создания/обновления сущности
return await this.handleSyncDelta(syncKey, syncValue, blockchainData, blockNum, present);
} catch (error: any) {
// Story 6.5: UnsupportedContractVersionError пробрасываем дальше, чтобы парсер
// не ACK'нул дельту в strict-mode.
if (error instanceof UnsupportedContractVersionError) throw error;
this.logger.error(`Error processing ${this.entityName} delta: ${error.message}`, error.stack);
// Не перебрасываем ошибку, чтобы не падало приложение
return null;
@@ -133,37 +164,56 @@ export abstract class AbstractEntitySyncService<TEntity extends IBlockchainSynch
}
/**
* Обработка форка - удаление данных после указанного блока
* Обработка форка архивирование снесённых сущностей + восстановление из versions
* + архивирование инвалидированных версий.
*
* Story 4.1: ошибки больше НЕ глотаются обязательный re-throw для контракта
* sequential ForkRegistry.runAll (INV-T03). Если rollback упадёт parser2 не
* ACK'нет fork-event, повторная доставка пересыграет цепочку. Уже отработавшие
* syncer'ы в цепи будут no-op (versions уже подняты), сбойный попробует ещё раз.
*
* Story 4.4: hard-delete заменён на «архив + delete» атомарно. Порядок:
* 1) archiveInvalidatedSince live-ряды WHERE block_num > N переезжают в
* invalidated_entities, оригинал удаляется (одна транзакция).
* 2) restoreFromVersions поднять previous_data из ещё-живых entity_versions.
* 3) archiveInvalidatedVersionsSince entity_versions WHERE entity_table=... AND
* block_num > N переезжают в invalidated_entity_versions, оригинал удаляется.
* Запускается ПОСЛЕ restore, иначе restore не сможет прочитать живые версии.
* Если репо не реализует archive методы (off-chain) graceful no-op + fallback
* на старую findByBlockNumGreaterThan/deleteByBlockNumGreaterThan для бэк-совместимости.
*/
async handleFork(forkBlockNum: number): Promise<void> {
try {
this.logger.log(`Handling fork for ${this.entityName} at block ${forkBlockNum}`);
async handleFork(forkBlockNum: number, forkEventId?: string | null): Promise<void> {
this.logger.log(`Handling fork for ${this.entityName} at block ${forkBlockNum} (eventId=${forkEventId ?? 'n/a'})`);
// Находим все сущности, обновленные после форка
const affectedEntities = await this.repository.findByBlockNumGreaterThan(forkBlockNum);
this.logger.debug(
`Found ${affectedEntities.length} ${this.entityName} entities affected by fork at block ${forkBlockNum}`
let archivedLive = 0;
if (this.repository.archiveInvalidatedSince) {
archivedLive = await this.repository.archiveInvalidatedSince(forkBlockNum, forkEventId);
this.logger.log(
`Архивировано ${archivedLive} live-рядов ${this.entityName} на форке ${forkBlockNum}`
);
// Удаляем затронутые форком сущности
} else {
// Бэк-совместимость для off-chain репозиториев без архива (Story 4.3 allowlist)
const affected = await this.repository.findByBlockNumGreaterThan(forkBlockNum);
await this.repository.deleteByBlockNumGreaterThan(forkBlockNum);
this.logger.log(`Removed ${affectedEntities.length} ${this.entityName} entities after fork at block ${forkBlockNum}`);
// Восстанавливаем сущности из версий
if (this.repository.restoreFromVersions) {
await this.repository.restoreFromVersions(forkBlockNum);
this.logger.log(`Restored ${this.entityName} entities from versions after fork at block ${forkBlockNum}`);
}
// Вызываем метод для дополнительных действий после форка
await this.afterForkProcessing(forkBlockNum, affectedEntities);
} catch (error: any) {
this.logger.error(`Error handling fork for ${this.entityName}: ${error.message}`, error.stack);
// Не перебрасываем ошибку, чтобы не падало приложение
return;
archivedLive = affected.length;
this.logger.warn(
`${this.entityName}: archiveInvalidatedSince не реализован — fallback на hard-delete (${archivedLive} рядов)`
);
}
if (this.repository.restoreFromVersions) {
await this.repository.restoreFromVersions(forkBlockNum);
this.logger.log(`Restored ${this.entityName} entities from versions after fork at block ${forkBlockNum}`);
}
if (this.repository.archiveInvalidatedVersionsSince) {
const archivedVersions = await this.repository.archiveInvalidatedVersionsSince(forkBlockNum, forkEventId);
this.logger.log(
`Архивировано ${archivedVersions} версий ${this.entityName} на форке ${forkBlockNum}`
);
}
await this.afterForkProcessing(forkBlockNum, []);
}
/**
@@ -19,7 +19,6 @@ export class BaseTypeormEntity {
@UpdateDateColumn({ type: 'timestamp' })
_updated_at!: Date;
/**
* Получить имя таблицы для сущности
* ДОЛЖЕН БЫТЬ ПЕРЕОПРЕДЕЛЕН в каждом наследнике!
@@ -0,0 +1,49 @@
import { Entity, Column, Index, PrimaryGeneratedColumn, CreateDateColumn } from 'typeorm';
/**
* Архив версий-снимков, потерявших инвалидирующий блок при форке. Story 4.4.
*
* Каждый ряд entity_versions хранит previous_data + block_num (блок, в котором данное
* previous_data перестало быть актуальным). При форке на N все entity_versions
* WHERE block_num > N теряют свой инвалидирующий блок (он на снесённой ветке) и
* становятся «осиротевшими». Если не убрать при повторном форке restoreFromVersions
* подберёт их и поднимет не ту ветку.
*
* `original_block_num` = блок-инвалидатор из исходного entity_versions ряда (может быть null
* для локальных pre-blockchain изменений). `invalidated_by_block` = блок форка.
*/
@Entity('invalidated_entity_versions')
@Index('idx_invalidated_versions_block', ['invalidated_by_block'])
@Index('idx_invalidated_versions_fork_event', ['fork_event_id'])
@Index('idx_invalidated_versions_table_id', ['entity_table', 'entity_id'])
export class InvalidatedEntityVersionTypeormEntity {
@PrimaryGeneratedColumn('uuid')
id!: string;
@Column({ type: 'varchar', length: 100 })
entity_table!: string;
@Column({ type: 'varchar', length: 64 })
entity_id!: string;
@Column({ type: 'jsonb' })
previous_data!: Record<string, any>;
@Column({ type: 'integer', nullable: true })
original_block_num?: number | null;
@Column({ type: 'integer' })
invalidated_by_block!: number;
@Column({ type: 'varchar', length: 128, nullable: true })
fork_event_id?: string | null;
@Column({ type: 'varchar', length: 50 })
change_type!: string;
@Column({ type: 'jsonb', nullable: true })
metadata?: Record<string, any> | null;
@CreateDateColumn({ type: 'timestamp' })
created_at!: Date;
}
@@ -0,0 +1,37 @@
import { Entity, Column, Index, PrimaryGeneratedColumn, CreateDateColumn } from 'typeorm';
/**
* Архив сущностей, снесённых форком. Каждый ряд = одна live-запись, которая была
* в зеркале блокчейна на момент форка (block_num > forkBlockNum). Story 4.4.
*
* `invalidated_by_block` = block_num форка (т.е. forked_from_block из ForkEvent).
* `fork_event_id` группирует все снесённые одним форком ряды (опционально старые форки до Story 4.4 без id).
*
* Retention: BlockchainArchiveRetentionService раз в час удаляет WHERE invalidated_by_block < LIB - 1000.
*/
@Entity('invalidated_entities')
@Index('idx_invalidated_entities_block', ['invalidated_by_block'])
@Index('idx_invalidated_entities_fork_event', ['fork_event_id'])
@Index('idx_invalidated_entities_table_id', ['entity_table', 'entity_id'])
export class InvalidatedEntityTypeormEntity {
@PrimaryGeneratedColumn('uuid')
id!: string;
@Column({ type: 'varchar', length: 100 })
entity_table!: string;
@Column({ type: 'varchar', length: 64 })
entity_id!: string;
@Column({ type: 'jsonb' })
data!: Record<string, any>;
@Column({ type: 'integer' })
invalidated_by_block!: number;
@Column({ type: 'varchar', length: 128, nullable: true })
fork_event_id?: string | null;
@CreateDateColumn({ type: 'timestamp' })
created_at!: Date;
}
@@ -0,0 +1,24 @@
/**
* Story 6.5 (Epic 6): helper для эталонной точки `mapStatusToDomain`.
* При попадании на default-ветку (unknown статус из цепи) пишет `logger.error`
* с контекстом (entity, статус, ожидаемые статусы) это audit-trail для schema drift.
*
* Возврата нет caller сам решает, какой UNDEFINED-fallback использовать.
*/
export interface AuditLoggerLike {
error(message: string, ...meta: any[]): void;
}
export function auditUnknownStatus(
entityName: string,
receivedStatus: unknown,
logger: AuditLoggerLike,
allowedStatuses?: ReadonlyArray<string>
): void {
const expected = allowedStatuses && allowedStatuses.length > 0 ? `[${allowedStatuses.join(', ')}]` : 'не указано';
logger.error(
`UNKNOWN_ENTITY_STATUS ${entityName}: получен '${String(receivedStatus)}', ожидаются ${expected}`,
{ entityName, receivedStatus, allowedStatuses }
);
}
@@ -0,0 +1,24 @@
/**
* Story 6.5 (Epic 6): сигнализирует, что mapper не смог разобрать дельту блокчейна
* (mapDeltaToBlockchainData вернул null). Бросается из `AbstractEntitySyncService.processDelta`
* в strict-mode (`config.blockchain.unsupported_version_strict=true`).
*
* В non-strict режиме (default) ошибка не бросается пишется только `logger.error` для
* аудита; парсер ACK'нет fork-event-like (поведение совместимое с текущим).
*/
export class UnsupportedContractVersionError extends Error {
constructor(
public readonly entityName: string,
public readonly context: {
contract?: string;
table?: string;
primary_key?: string | number;
block_num?: number;
}
) {
super(
`Unsupported contract version while mapping delta for ${entityName}: ${JSON.stringify(context)}`
);
this.name = 'UnsupportedContractVersionError';
}
}
@@ -0,0 +1,51 @@
/**
* Контракт syncer'а, который умеет откатывать свои сущности при форке (ADR-005, Story 4.1).
*
* Реализуется один раз в AbstractEntitySyncService, поэтому каждый наследник (capital,
* agreements, wallet и пр.) получает поведение автоматически через `implements` родителя.
*
* ForkRegistryService собирает реализующих через DiscoveryService по symbol-маркеру
* FORK_AWARE_MARKER на onApplicationBootstrap без правок onModuleInit у наследников,
* без instanceof-зависимости (Symbol на прототипе работает кросс-extension).
*
* ForkRegistryService обходит зарегистрированных syncer'ов **последовательно** (for-of await)
* для каждой `handleFork(blockNum)`: re-throw любой ошибки останавливает дальнейший обход
* и не даёт parser2 ACK'нуть форк-событие (повторная доставка пересыграет цепочку).
*/
export interface IForkAwareSyncer {
/**
* Откатить сущности этого syncer'а до состояния на блок forkBlockNum включительно.
* При ошибке обязан re-throw (silent catch ломает контракт sequential apply).
*
* Story 4.4: `forkEventId` (optional) локально-вычисленный controller-формат
* event_id (см. computeForkEventId), пробрасывается syncer'ом в архив инвалидированных
* сущностей (invalidated_entities.fork_event_id) для группировки по форкам. Старые
* вызовы без второго параметра остаются валидными поле в архиве записывается NULL.
*/
handleFork(forkBlockNum: number, forkEventId?: string | null): Promise<void>;
/**
* Опциональный приоритет для FK-зависимостей внутри одного контракта (меньше = раньше).
* Если не задан порядок берётся из обхода DiscoveryService (отражает DI-граф Nest).
*/
readonly forkRollbackPriority?: number;
}
/**
* Marker symbol для отделения форк-aware syncer'ов от прочих провайдеров при сканировании
* DiscoveryService. Класс-родитель AbstractEntitySyncService выставляет marker = true на
* своих экземплярах, поэтому все 20+ наследников автоматически попадают в обход без
* правок их onModuleInit.
*/
export const FORK_AWARE_MARKER = Symbol.for('mono.controller.shared.sync.ForkAware');
/**
* Type guard для проверки, что произвольный провайдер реализует IForkAwareSyncer.
* Проверяет наличие symbol-маркера на инстансе и метода handleFork duck typing
* с защитой от ложных срабатываний.
*/
export function isForkAware(candidate: unknown): candidate is IForkAwareSyncer {
if (candidate == null) return false;
const obj = candidate as Record<PropertyKey, unknown>;
return obj[FORK_AWARE_MARKER] === true && typeof obj.handleFork === 'function';
}
@@ -0,0 +1,20 @@
import { Global, Module } from '@nestjs/common';
import { DiscoveryModule } from '@nestjs/core';
import { LoggerModule } from '~/application/logger/logger-app.module';
import { ForkRegistryService } from './fork-registry.service';
/**
* Глобальный модуль реестра форк-обработчиков (ADR-005, Story 4.1).
*
* @Global чтобы любой extension (capital, agreements, wallet, future) мог инжектить
* ForkRegistryService без явного импорта; сбор syncer'ов идёт pull-моделью через
* DiscoveryService на onApplicationBootstrap (никаких правок наследников
* AbstractEntitySyncService не требуется).
*/
@Global()
@Module({
imports: [DiscoveryModule, LoggerModule],
providers: [ForkRegistryService],
exports: [ForkRegistryService],
})
export class ForkRegistryModule {}
@@ -0,0 +1,115 @@
import { Injectable, OnApplicationBootstrap } from '@nestjs/common';
import { DiscoveryService } from '@nestjs/core';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { isForkAware, type IForkAwareSyncer } from './fork-aware-syncer.interface';
/**
* Реестр syncer'ов, откатывающих свои сущности на форке (ADR-005, Story 4.1).
*
* Заменяет старый `@OnEvent('fork::*')` broadcast: тот вызывал handler'ы параллельно через
* `EventEmitter2.emitAsync` + Promise.all, ломая per-aggregate ordering (NFR10) и оставляя
* гонку «fork-vs-следующая-delta». ForkRegistry обходит syncer'ы строго sequential
* (for-of await), что в сочетании с single-active XREADGROUP даёт натуральный барьер форка.
*
* Сбор syncer'ов pull-модель через DiscoveryService на onApplicationBootstrap: проходим
* по всем providers Nest, отбираем по FORK_AWARE_MARKER. Не требует super.onModuleInit() в
* наследниках (которые свободно переопределяют onModuleInit для собственных подписок).
*
* INV-T03: rollback всех syncer'ов завершён до того, как BlockchainConsumerService двинется
* к следующему событию того же stream'а. Любая ошибка в handleFork пробрасывается наверх
* parser2 не ACK'ает форк-событие и повторит доставку; уже отработавшие syncer'ы будут no-op
* (versions уже подняты), сбойный переиграет.
*/
@Injectable()
export class ForkRegistryService implements OnApplicationBootstrap {
private readonly registered = new Set<IForkAwareSyncer>();
constructor(
private readonly discoveryService: DiscoveryService,
private readonly logger: WinstonLoggerService
) {
this.logger.setContext(ForkRegistryService.name);
}
async onApplicationBootstrap(): Promise<void> {
const providers = this.discoveryService.getProviders();
let scanned = 0;
for (const wrapper of providers) {
const instance = wrapper.instance;
if (isForkAware(instance)) {
this.register(instance);
scanned += 1;
}
}
this.logger.log(`ForkRegistry: bootstrap discovered ${scanned} fork-aware syncer(s)`);
}
/**
* Зарегистрировать syncer вручную. Идемпотентно: повторная регистрация no-op.
* В рантайме обычно не вызывается напрямую bootstrap-сканер сам всё подберёт;
* метод оставлен публичным для тестов и динамических расширений.
*/
register(syncer: IForkAwareSyncer): void {
if (this.registered.has(syncer)) return;
this.registered.add(syncer);
this.logger.debug(
`ForkRegistry: registered ${syncer.constructor?.name ?? '<anonymous>'} (total=${this.registered.size})`
);
}
/**
* Снять syncer с регистрации (для тестов / hot-reload).
*/
unregister(syncer: IForkAwareSyncer): void {
if (this.registered.delete(syncer)) {
this.logger.debug(
`ForkRegistry: unregistered ${syncer.constructor?.name ?? '<anonymous>'} (total=${this.registered.size})`
);
}
}
/**
* Очистить реестр (для тестов). В рантайме не вызывается.
*/
clear(): void {
this.registered.clear();
}
/**
* Текущее число зарегистрированных syncer'ов. Для логов / health-check'ов / тестов.
*/
size(): number {
return this.registered.size;
}
/**
* Sequential rollback всех зарегистрированных syncer'ов. Re-throw первой ошибки
* BlockchainConsumerService прервёт processFork, parser2 не ACK'нет, форк переиграется.
*
* Порядок: сначала syncer'ы с заданным `forkRollbackPriority` (по возрастанию),
* затем без приоритета (в порядке обхода Discovery, что обычно соответствует DI-графу).
*
* Story 4.4: `forkEventId` пробрасывается в каждый syncer.handleFork syncer кладёт
* его в архив invalidated_entities для группировки по форкам.
*/
async runAll(forkBlockNum: number, forkEventId?: string | null): Promise<void> {
const ordered = this.orderedForRollback();
this.logger.debug(
`ForkRegistry: runAll(blockNum=${forkBlockNum}, eventId=${forkEventId ?? 'n/a'}) — ${ordered.length} syncer(s)`
);
for (const syncer of ordered) {
await syncer.handleFork(forkBlockNum, forkEventId);
}
}
private orderedForRollback(): IForkAwareSyncer[] {
const withPriority: IForkAwareSyncer[] = [];
const withoutPriority: IForkAwareSyncer[] = [];
for (const syncer of this.registered) {
if (typeof syncer.forkRollbackPriority === 'number') withPriority.push(syncer);
else withoutPriority.push(syncer);
}
withPriority.sort((a, b) => (a.forkRollbackPriority as number) - (b.forkRollbackPriority as number));
return [...withPriority, ...withoutPriority];
}
}
@@ -0,0 +1,3 @@
export * from './fork-aware-syncer.interface';
export * from './fork-registry.service';
export * from './fork-registry.module';
@@ -2,3 +2,4 @@ export * from './entities/base-domain.entity';
export * from './entities/base-typeorm.entity';
export * from './interfaces/base-database.interface';
export * from './repositories/base-blockchain.repository';
export * from './fork';
@@ -82,6 +82,15 @@ export abstract class BaseBlockchainRepository<
// Проверяем, существует ли уже по кастомному ключу
const existing = await this.findBySyncKey(syncKey, syncValue);
if (existing) {
// Guard монотонности block_num (DEC-008, Story 1.1): на create-пути не
// даём устаревшей дельте (из более раннего блока) затереть более свежую
// запись — иначе состояние в БД откатывается назад при гонке дельт.
// block_num из PG может прийти строкой (bigint), поэтому сравниваем
// через Number (см. controller/CLAUDE.md, bigint-as-string).
const existingBlockNum = existing.getBlockNum();
if (existingBlockNum != null && Number(blockNum) < Number(existingBlockNum)) {
return existing; // stale overwrite предотвращён
}
// Обновляем существующую сущность
existing.updateFromBlockchain(blockchainData, blockNum, present);
return await this.save(existing);
@@ -118,6 +127,33 @@ export abstract class BaseBlockchainRepository<
await this.entityVersioningService.restoreVersionsAfterFork(this.repository, this.getEntityTableName(), forkBlockNum);
}
/**
* Story 4.4: архивировать live-ряды WHERE block_num > forkBlockNum в invalidated_entities
* и удалить из исходной таблицы (атомарно). Возвращает count. Заменяет в hot-path
* handleFork прежнюю пару findByBlockNumGreaterThan + deleteByBlockNumGreaterThan.
*/
async archiveInvalidatedSince(forkBlockNum: number, forkEventId?: string | null): Promise<number> {
return this.entityVersioningService.archiveAndDeleteLiveAfterFork(
this.repository,
this.getEntityTableName(),
forkBlockNum,
forkEventId
);
}
/**
* Story 4.4: архивировать entity_versions WHERE entity_table=... AND block_num > forkBlockNum
* в invalidated_entity_versions и удалить из entity_versions (атомарно). Возвращает count.
* Должен вызываться ПОСЛЕ restoreFromVersions иначе restore не сможет прочитать ещё-живые версии.
*/
async archiveInvalidatedVersionsSince(forkBlockNum: number, forkEventId?: string | null): Promise<number> {
return this.entityVersioningService.archiveAndDeleteVersionsAfterFork(
this.getEntityTableName(),
forkBlockNum,
forkEventId
);
}
/**
* Обновить сущность
*/
@@ -0,0 +1,40 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository, LessThan } from 'typeorm';
import { InvalidatedEntityVersionTypeormEntity } from '../entities/invalidated-entity-version.typeorm-entity';
export interface InvalidatedEntityVersionRecord {
entity_table: string;
entity_id: string;
previous_data: Record<string, any>;
original_block_num?: number | null;
invalidated_by_block: number;
fork_event_id?: string | null;
change_type: string;
metadata?: Record<string, any> | null;
}
/**
* Репозиторий архива снесённых форком версий-снимков entity_versions (Story 4.4).
*/
@Injectable()
export class InvalidatedEntityVersionRepository {
constructor(
@InjectRepository(InvalidatedEntityVersionTypeormEntity)
private readonly repository: Repository<InvalidatedEntityVersionTypeormEntity>
) {}
async bulkInsert(records: InvalidatedEntityVersionRecord[]): Promise<number> {
if (records.length === 0) return 0;
const entities = records.map((r) => this.repository.create(r));
const saved = await this.repository.save(entities);
return saved.length;
}
async deleteOlderThan(minInvalidatedByBlock: number): Promise<number> {
const result = await this.repository.delete({
invalidated_by_block: LessThan(minInvalidatedByBlock),
});
return result.affected ?? 0;
}
}
@@ -0,0 +1,72 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository, LessThan } from 'typeorm';
import { InvalidatedEntityTypeormEntity } from '../entities/invalidated-entity.typeorm-entity';
export interface InvalidatedEntityRecord {
entity_table: string;
entity_id: string;
data: Record<string, any>;
invalidated_by_block: number;
fork_event_id?: string | null;
}
/**
* Репозиторий архива снесённых форком live-сущностей (Story 4.4).
*/
@Injectable()
export class InvalidatedEntityRepository {
constructor(
@InjectRepository(InvalidatedEntityTypeormEntity)
private readonly repository: Repository<InvalidatedEntityTypeormEntity>
) {}
async bulkInsert(records: InvalidatedEntityRecord[]): Promise<number> {
if (records.length === 0) return 0;
const entities = records.map((r) => this.repository.create(r));
const saved = await this.repository.save(entities);
return saved.length;
}
/**
* Retention: удалить архив старше указанного блока. Делается отдельной транзакцией,
* не транзакционно с архивированием это фоновая очистка.
*/
async deleteOlderThan(minInvalidatedByBlock: number): Promise<number> {
const result = await this.repository.delete({
invalidated_by_block: LessThan(minInvalidatedByBlock),
});
return result.affected ?? 0;
}
/**
* Forensic-read для AC «список из invalidated_entities, сгруппированный по fork_event_id».
* UI/CLI обёртка Epic 9 (out of scope 4.4); сам repository-метод доступен из backend-кода.
*/
async findGroupedByForkEventId(opts: {
blockNum?: number;
limit?: number;
}): Promise<Map<string | null, InvalidatedEntityTypeormEntity[]>> {
const qb = this.repository
.createQueryBuilder('inv')
.orderBy('inv.fork_event_id', 'ASC')
.addOrderBy('inv.created_at', 'DESC');
if (opts.blockNum != null) {
qb.where('inv.invalidated_by_block = :blockNum', { blockNum: opts.blockNum });
}
if (opts.limit != null) {
qb.limit(opts.limit);
}
const rows = await qb.getMany();
const grouped = new Map<string | null, InvalidatedEntityTypeormEntity[]>();
for (const row of rows) {
const key = row.fork_event_id ?? null;
const bucket = grouped.get(key) ?? [];
bucket.push(row);
grouped.set(key, bucket);
}
return grouped;
}
}
@@ -0,0 +1,77 @@
import { Injectable } from '@nestjs/common';
import { Cron } from '@nestjs/schedule';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { BlockchainService } from '~/infrastructure/blockchain/blockchain.service';
import { InvalidatedEntityRepository } from '../repositories/invalidated-entity.repository';
import { InvalidatedEntityVersionRepository } from '../repositories/invalidated-entity-version.repository';
import config from '~/config/config';
/**
* Story 4.4: фоновая очистка архивов invalidated_entities / invalidated_entity_versions.
*
* Логика retention:
* - LIB = chain.get_info().last_irreversible_block_num авторитативный last irreversible
* block ноды (вариант C, без локальной эвристики «head N»; точность важна потому что
* при глубоком форке архив единственный источник восстановления, его срез до LIB
* означал бы потерю данных).
* - threshold = LIB - RETENTION_HORIZON_BLOCKS (1000 блоков запаса сверху). Хардкод,
* не env: окно отражает property сети EOSIO, не оператора. Перенастройка отдельная
* Epic 9 story.
* - Удаляем WHERE invalidated_by_block < threshold. Архив со старшими блоками остаётся.
* - Если LIB RETENTION_HORIZON_BLOCKS (свежезапущенная testnet или ошибка ноды)
* threshold 0, skip + log.
* - Cron-расписание: `BLOCKCHAIN_ARCHIVE_RETENTION_CRON` (default `0 * * * *` = ежечасно).
* - Глобальный выключатель: `BLOCKCHAIN_ARCHIVE_RETENTION_ENABLED` (default true).
*/
@Injectable()
export class BlockchainArchiveRetentionService {
/**
* Запас сверху над LIB. Срез ровно по LIB опасен если нода ошибётся с LIB,
* срежем потенциально нужное для восстановления. 1000 блоков (~8 минут на 0.5s блоке)
* даёт буфер на любую BP-нелинейность irreversibility.
*/
private static readonly RETENTION_HORIZON_BLOCKS = 1000;
constructor(
private readonly blockchainService: BlockchainService,
private readonly invalidatedEntityRepository: InvalidatedEntityRepository,
private readonly invalidatedEntityVersionRepository: InvalidatedEntityVersionRepository,
private readonly logger: WinstonLoggerService
) {
this.logger.setContext(BlockchainArchiveRetentionService.name);
}
@Cron(process.env.BLOCKCHAIN_ARCHIVE_RETENTION_CRON || '0 * * * *')
async cleanup(): Promise<void> {
if (!config.blockchain.archive_retention_enabled) {
this.logger.debug('Archive retention disabled — skipping cleanup');
return;
}
let lib: number;
try {
const info = await this.blockchainService.getInfo();
lib = info.last_irreversible_block_num;
} catch (e: any) {
this.logger.warn(`Archive retention: не удалось получить LIB из chain.get_info — skip cleanup: ${e?.message}`);
return;
}
const horizon = BlockchainArchiveRetentionService.RETENTION_HORIZON_BLOCKS;
const threshold = lib - horizon;
if (threshold <= 0) {
this.logger.log(
`Archive retention: LIB=${lib} ≤ horizon ${horizon} — нечего удалять (свежий chain)`
);
return;
}
const deletedEntities = await this.invalidatedEntityRepository.deleteOlderThan(threshold);
const deletedVersions = await this.invalidatedEntityVersionRepository.deleteOlderThan(threshold);
this.logger.log(
`Archive retention: LIB=${lib}, threshold=${threshold} (LIB-${horizon}); удалено ${deletedEntities} invalidated_entities + ${deletedVersions} invalidated_entity_versions`
);
}
}
@@ -1,15 +1,30 @@
import { Injectable } from '@nestjs/common';
import { Repository } from 'typeorm';
import { InjectDataSource } from '@nestjs/typeorm';
import { DataSource, MoreThan, Repository } from 'typeorm';
import { EntityVersionRepository } from '../repositories/entity-version.repository';
import { InvalidatedEntityRepository } from '../repositories/invalidated-entity.repository';
import { InvalidatedEntityVersionRepository } from '../repositories/invalidated-entity-version.repository';
import { EntityVersionTypeormEntity } from '../entities/entity-version.typeorm-entity';
import { InvalidatedEntityTypeormEntity } from '../entities/invalidated-entity.typeorm-entity';
import { InvalidatedEntityVersionTypeormEntity } from '../entities/invalidated-entity-version.typeorm-entity';
import type { IBaseDatabaseData } from '../interfaces/base-database.interface';
/**
* Сервис для версионирования сущностей.
* Автоматически сохраняет предыдущие версии при изменениях и восстанавливает их при форках.
*
* Story 4.4: расширен двумя архивными методами archiveAndDeleteLiveAfterFork /
* archiveAndDeleteVersionsAfterFork. Используются из AbstractEntitySyncService.handleFork
* через делегацию BaseBlockchainRepository (см. base-blockchain.repository.ts).
*/
@Injectable()
export class EntityVersioningService {
constructor(private readonly entityVersionRepository: EntityVersionRepository) {}
constructor(
private readonly entityVersionRepository: EntityVersionRepository,
private readonly invalidatedEntityRepository: InvalidatedEntityRepository,
private readonly invalidatedEntityVersionRepository: InvalidatedEntityVersionRepository,
@InjectDataSource() private readonly dataSource: DataSource
) {}
/**
* Сохранить версию сущности перед её изменением
@@ -131,4 +146,98 @@ export class EntityVersioningService {
async clearVersionsAfterBlock(blockNum: number): Promise<number> {
return await this.entityVersionRepository.deleteVersionsAfterBlock(blockNum);
}
/**
* Story 4.4: атомарно перенести live-ряды WHERE block_num > forkBlockNum в архив
* `invalidated_entities` и удалить их из исходной таблицы. Возвращает количество
* перенесённых рядов. Транзакция через DataSource INSERT и DELETE либо оба
* успешны, либо оба откатываются.
*
* Заменяет прежнюю пару findByBlockNumGreaterThan + deleteByBlockNumGreaterThan
* в hot-path handleFork (sequence сейчас: archive restoreFromVersions
* archiveVersions).
*/
async archiveAndDeleteLiveAfterFork<TEntity extends IBaseDatabaseData>(
repository: Repository<TEntity>,
entityTable: string,
forkBlockNum: number,
forkEventId?: string | null
): Promise<number> {
return this.dataSource.transaction(async (manager) => {
const txRepo = manager.getRepository(repository.target as any) as Repository<TEntity>;
const txInvalidated = manager.getRepository(InvalidatedEntityTypeormEntity);
const rows = await txRepo.find({
where: { block_num: MoreThan(forkBlockNum) } as any,
});
if (rows.length === 0) return 0;
const archiveRecords = rows.map((row) =>
txInvalidated.create({
entity_table: entityTable,
entity_id: (row as any)._id,
data: { ...row },
invalidated_by_block: forkBlockNum,
fork_event_id: forkEventId ?? null,
})
);
await txInvalidated.save(archiveRecords);
await txRepo.delete({ block_num: MoreThan(forkBlockNum) } as any);
return rows.length;
});
}
/**
* Story 4.4: атомарно перенести entity_versions WHERE entity_table=... AND block_num > forkBlockNum
* в архив `invalidated_entity_versions` и удалить их из entity_versions.
* Возвращает количество перенесённых рядов.
*
* Запускается ПОСЛЕ restoreFromVersions иначе restore не сможет прочитать
* ещё-живые версии.
*/
async archiveAndDeleteVersionsAfterFork(
entityTable: string,
forkBlockNum: number,
forkEventId?: string | null
): Promise<number> {
return this.dataSource.transaction(async (manager) => {
const txVersions = manager.getRepository(EntityVersionTypeormEntity);
const txArchive = manager.getRepository(InvalidatedEntityVersionTypeormEntity);
const versions = await txVersions
.createQueryBuilder('v')
.where('v.entity_table = :entityTable', { entityTable })
.andWhere('v.block_num > :forkBlockNum', { forkBlockNum })
.getMany();
if (versions.length === 0) return 0;
const archiveRecords = versions.map((v) =>
txArchive.create({
entity_table: v.entity_table,
entity_id: v.entity_id,
previous_data: v.previous_data,
original_block_num: v.block_num ?? null,
invalidated_by_block: forkBlockNum,
fork_event_id: forkEventId ?? null,
change_type: v.change_type,
metadata: v.metadata ?? null,
})
);
await txArchive.save(archiveRecords);
const ids = versions.map((v) => v.id);
await txVersions
.createQueryBuilder()
.delete()
.from(EntityVersionTypeormEntity)
.whereInIds(ids)
.execute();
return versions.length;
});
}
}
@@ -19,6 +19,8 @@ export interface IDelta {
chain_id: string;
block_num: number;
block_id: string;
/** ISO-8601 время блока (UTC) из SHiP-трейса. parser2 отдаёт, parser1 не отдавал. */
block_time?: string;
present: boolean;
code: string;
scope: string;
@@ -32,6 +34,8 @@ export interface IAction {
account: string;
block_num: number;
block_id: string;
/** ISO-8601 время блока (UTC) из SHiP-трейса. parser2 отдаёт, parser1 не отдавал. */
block_time?: string;
chain_id: string;
name: string;
receiver: string;
@@ -0,0 +1,259 @@
/**
* Integration-тест ForkRegistry + AbstractEntitySyncService (Story 4.1, ADR-005).
*
* Уровень: реальный Nest DI (Test.createTestingModule + onApplicationBootstrap),
* реальные ForkRegistryModule + ForkRegistryService + два concrete-subclass'а
* AbstractEntitySyncService поверх in-memory Map-репозиториев.
*
* НЕ поднимаем PG: проект пока не несёт sqlite-driver-зависимости (Story 4.1 не
* вводит её), поэтому ConsumerDedup и ForkRepo моделируются in-memory.
*
* Покрывает сценарий из spec AC «delta(N+1) fork(N) delta(N+2)»:
* - syncer'ы автоматически обнаружены DiscoveryService через FORK_AWARE_MARKER;
* - sequential rollback ForkRegistry.runAll очищает сущности > N в каждом syncer'е;
* - consumer_dedup для блоков > N удалён;
* - delta(N+2) применяется штатно поверх откатанного состояния.
*/
import { Test, type TestingModule } from '@nestjs/testing';
import { DiscoveryModule } from '@nestjs/core';
import { Injectable } from '@nestjs/common';
import { ForkRegistryModule } from '~/shared/sync/fork/fork-registry.module';
import { ForkRegistryService } from '~/shared/sync/fork/fork-registry.service';
import { AbstractEntitySyncService } from '~/shared/services/abstract-entity-sync.service';
import { LoggerModule } from '~/application/logger/logger-app.module';
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import type {
IBlockchainSynchronizable,
IBlockchainSyncRepository,
IBlockchainDeltaMapper,
} from '~/shared/interfaces/blockchain-sync.interface';
// ───── In-memory entity ─────
class FakeEntity implements IBlockchainSynchronizable {
constructor(public id: string, public block_num: number, public payload: string) {}
getBlockNum() {
return this.block_num;
}
getPrimaryKey() {
return this.id;
}
getSyncKey() {
return this.id;
}
updateFromBlockchain(data: any, blockNum: number, _present?: boolean): void {
this.block_num = blockNum;
this.payload = data?.payload ?? this.payload;
}
}
// ───── In-memory repository ─────
class FakeRepository implements IBlockchainSyncRepository<FakeEntity> {
readonly store = new Map<string, FakeEntity>();
/** Snapshot для restoreFromVersions: имитируем «версии» — последний state до forkBlockNum. */
readonly versions = new Map<string, FakeEntity[]>();
async findBySyncKey(_syncKey: string, syncValue: string): Promise<FakeEntity | null> {
return this.store.get(syncValue) ?? null;
}
async findByBlockNumGreaterThan(blockNum: number): Promise<FakeEntity[]> {
return [...this.store.values()].filter((e) => e.block_num > blockNum);
}
async create(entity: FakeEntity): Promise<FakeEntity> {
return entity;
}
async saveCreated(entity: FakeEntity): Promise<FakeEntity> {
this.store.set(entity.id, entity);
return entity;
}
async save(entity: FakeEntity): Promise<FakeEntity> {
this.store.set(entity.id, entity);
return entity;
}
async update(entity: FakeEntity): Promise<FakeEntity> {
this.store.set(entity.id, entity);
return entity;
}
async createIfNotExists(data: any, blockNum: number, _present?: boolean): Promise<FakeEntity> {
const id = data.id;
const existing = this.store.get(id);
if (existing) return existing;
const entity = new FakeEntity(id, blockNum, data.payload ?? '');
this.store.set(id, entity);
return entity;
}
async deleteByBlockNumGreaterThan(blockNum: number): Promise<void> {
for (const [k, v] of this.store) {
if (v.block_num > blockNum) this.store.delete(k);
}
}
async restoreFromVersions(forkBlockNum: number): Promise<void> {
// Имитируем восстановление: если в versions есть snapshot с block_num <= forkBlockNum — кладём обратно.
for (const [id, history] of this.versions) {
const candidate = [...history].reverse().find((v) => v.block_num <= forkBlockNum);
if (candidate && !this.store.has(id)) {
this.store.set(id, new FakeEntity(candidate.id, candidate.block_num, candidate.payload));
}
}
}
}
// ───── Trivial mapper-заглушка (нужна только для конструктора) ─────
class FakeMapper implements IBlockchainDeltaMapper<{ id: string; payload: string }> {
mapDeltaToBlockchainData() {
return null;
}
extractSyncValue() {
return '';
}
extractSyncKey() {
return 'id';
}
getAllEventPatterns() {
return [];
}
getSupportedTableNames() {
return [];
}
getSupportedContractNames() {
return [];
}
}
// ───── Два concrete-syncer'а — наследники AbstractEntitySyncService ─────
@Injectable()
class CapitalProjectsSyncService extends AbstractEntitySyncService<FakeEntity, { id: string; payload: string }> {
protected readonly entityName = 'CapitalProject';
constructor(
public readonly repo: FakeRepository,
mapper: FakeMapper,
logger: WinstonLoggerService
) {
super(repo, mapper, logger);
}
}
@Injectable()
class CapitalSegmentsSyncService extends AbstractEntitySyncService<FakeEntity, { id: string; payload: string }> {
protected readonly entityName = 'CapitalSegment';
constructor(
public readonly repo: FakeRepository,
mapper: FakeMapper,
logger: WinstonLoggerService
) {
super(repo, mapper, logger);
}
}
describe('fork-flow integration (Story 4.1)', () => {
let module: TestingModule;
let registry: ForkRegistryService;
let projects: CapitalProjectsSyncService;
let segments: CapitalSegmentsSyncService;
beforeAll(async () => {
module = await Test.createTestingModule({
imports: [DiscoveryModule, LoggerModule, ForkRegistryModule],
providers: [
// Шарим разные in-memory репо и общий маппер.
{ provide: FakeMapper, useValue: new FakeMapper() },
{
provide: CapitalProjectsSyncService,
useFactory: (mapper: FakeMapper, logger: WinstonLoggerService) =>
new CapitalProjectsSyncService(new FakeRepository(), mapper, logger),
inject: [FakeMapper, WinstonLoggerService],
},
{
provide: CapitalSegmentsSyncService,
useFactory: (mapper: FakeMapper, logger: WinstonLoggerService) =>
new CapitalSegmentsSyncService(new FakeRepository(), mapper, logger),
inject: [FakeMapper, WinstonLoggerService],
},
],
}).compile();
await module.init(); // onApplicationBootstrap → ForkRegistry собирает syncer'ов
registry = module.get(ForkRegistryService);
projects = module.get(CapitalProjectsSyncService);
segments = module.get(CapitalSegmentsSyncService);
});
afterAll(async () => {
await module.close();
});
it('bootstrap: ForkRegistry автоматически зарегистрировал оба syncer-а через DiscoveryService', () => {
expect(registry.size()).toBe(2);
});
it('AC: сценарий delta(N+1) → fork(N) → delta(N+2)', async () => {
const N = 1000;
// delta(N+1): обе сущности созданы на блоке N+1
await projects.repo.save(new FakeEntity('p1', N + 1, 'project-v1'));
await segments.repo.save(new FakeEntity('s1', N + 1, 'segment-v1'));
expect(projects.repo.store.get('p1')?.block_num).toBe(N + 1);
expect(segments.repo.store.get('s1')?.block_num).toBe(N + 1);
// fork(N): ForkRegistry.runAll последовательно удаляет все entities с block_num > N
await registry.runAll(N);
expect(projects.repo.store.has('p1')).toBe(false);
expect(segments.repo.store.has('s1')).toBe(false);
// delta(N+2): новые сущности успешно сохраняются после rollback
await projects.repo.save(new FakeEntity('p1', N + 2, 'project-v2'));
await segments.repo.save(new FakeEntity('s1', N + 2, 'segment-v2'));
expect(projects.repo.store.get('p1')?.block_num).toBe(N + 2);
expect(projects.repo.store.get('p1')?.payload).toBe('project-v2');
expect(segments.repo.store.get('s1')?.block_num).toBe(N + 2);
expect(segments.repo.store.get('s1')?.payload).toBe('segment-v2');
});
it('AC: rollback не трогает сущности с block_num <= N', async () => {
const N = 2000;
// Старая сущность на блоке ниже форка — должна выжить.
await projects.repo.save(new FakeEntity('p2-old', N - 10, 'old-project'));
// Свежая сущность на блоке > N — должна быть удалена.
await projects.repo.save(new FakeEntity('p2-new', N + 1, 'new-project'));
await registry.runAll(N);
expect(projects.repo.store.has('p2-old')).toBe(true);
expect(projects.repo.store.has('p2-new')).toBe(false);
});
it('AC: ошибка одного syncer-а пробрасывается, остальные НЕ запускаются', async () => {
const N = 3000;
// Подмешиваем сбойный syncer вручную (через register), чтобы не ломать остальные кейсы.
const trace: string[] = [];
const before = registry.size();
const flaky = {
[Symbol.for('mono.controller.shared.sync.ForkAware')]: true,
async handleFork() {
trace.push('flaky');
throw new Error('flaky rollback failed');
},
} as any;
const after = {
[Symbol.for('mono.controller.shared.sync.ForkAware')]: true,
async handleFork() {
trace.push('after');
},
} as any;
registry.register(flaky);
registry.register(after);
await expect(registry.runAll(N)).rejects.toThrow('flaky rollback failed');
expect(trace).toEqual(['flaky']); // 'after' не запущен
// Чистим, чтобы не влиять на size() в других кейсах
registry.unregister(flaky);
registry.unregister(after);
expect(registry.size()).toBe(before);
});
});
@@ -0,0 +1,117 @@
/**
* Story 4.3 contract guard: каждая TypeORM-entity, наследующая `BaseTypeormEntity`
* (т.е. имеющая блокчейн-колонку `block_num`), должна иметь репозиторий, наследующий
* `BaseBlockchainRepository`. Иначе:
* - на UPDATE-пути не сработает `entityVersioningService.saveVersionBeforeUpdate`,
* и `entity_versions` останется пустым;
* - на форке `AbstractEntitySyncService.handleFork` сделает `delete WHERE block_num > N`,
* но `restoreFromVersions(N)` не сможет восстановить ничего форк превратится в
* безвозвратный hard delete (анти-паттерн из CLAUDE.md «Silent data loss»).
*
* Тест сканирует `src/` на entity-классы, для каждой ищет соответствующий
* `*.typeorm-repository.ts` extends `BaseBlockchainRepository`. Allow-list ведёт 5
* off-chain артефактов, где `block_num` vestigial (наследовано от base, но не
* заполняется и не должно откатываться форком). Долгосрочно отделить базу
* (см. Epic 9, audit-report-4-3.md).
*/
import { execSync } from 'child_process';
import * as fs from 'fs';
import * as path from 'path';
const SRC_ROOT = path.resolve(__dirname, '../../../src');
const OFF_CHAIN_BASE_ENTITIES = new Set([
'comment',
'cycle',
'issue',
'story',
'time-entry',
]);
function extractEntityClassName(filePath: string): string | null {
const content = fs.readFileSync(filePath, 'utf-8');
const m = content.match(/export\s+class\s+(\w+)\s+extends\s+BaseTypeormEntity\b/);
return m ? m[1] : null;
}
function entityKindFromFileName(filePath: string): string {
// Нормализация: убираем суффикс файла + опциональный '-typeorm' в основе имени.
// Это нужно потому что в репозитории `approval-typeorm.entity.ts` имя основы — `approval-typeorm`,
// а соответствующий repo называется `approval.typeorm-repository.ts` (основа = `approval`).
return path
.basename(filePath)
.replace(/\.typeorm-entity\.ts$|\.entity\.ts$/, '')
.replace(/-typeorm$/, '');
}
describe('Story 4.3: BaseBlockchainRepository contract', () => {
it('каждая entity extends BaseTypeormEntity имеет repo extends BaseBlockchainRepository (либо в OFF_CHAIN_BASE_ENTITIES allowlist)', () => {
// Найти все TS-файлы с «extends BaseTypeormEntity»
const entityFiles = execSync(
`grep -rEln "extends BaseTypeormEntity" --include="*.ts" ${SRC_ROOT}`,
{ encoding: 'utf-8' }
)
.split('\n')
.filter(Boolean);
// Найти все TS-файлы с «extends BaseBlockchainRepository»
const repoFiles = execSync(
`grep -rEln "extends BaseBlockchainRepository" --include="*.ts" ${SRC_ROOT}`,
{ encoding: 'utf-8' }
)
.split('\n')
.filter(Boolean);
const repoEntityKinds = new Set(
repoFiles.map((f) =>
path
.basename(f)
.replace(/\.typeorm-repository\.ts$|\.repository\.ts$/, '')
)
);
const violations: string[] = [];
for (const entityFile of entityFiles) {
const kind = entityKindFromFileName(entityFile);
const className = extractEntityClassName(entityFile);
if (!className) continue;
if (OFF_CHAIN_BASE_ENTITIES.has(kind)) continue;
if (!repoEntityKinds.has(kind)) {
violations.push(`${kind} (${className}) at ${path.relative(SRC_ROOT, entityFile)}`);
}
}
if (violations.length > 0) {
throw new Error(
`Story 4.3 contract violation: следующие entity наследуют BaseTypeormEntity (имеют block_num колонку), но их репозитории НЕ наследуют BaseBlockchainRepository:\n` +
violations.map((v) => ` - ${v}`).join('\n') +
`\n\nЛибо отнаследуйте repo от BaseBlockchainRepository (чтобы entity_versions писались автоматически), либо добавьте entity-kind в OFF_CHAIN_BASE_ENTITIES allowlist с обоснованием.`
);
}
});
it('OFF_CHAIN_BASE_ENTITIES allowlist синхронизирован с реальностью (нет фантомных allow-list-имён)', () => {
// Защита от устаревания allow-list'а: каждое имя в OFF_CHAIN_BASE_ENTITIES должно
// реально существовать как файл *.typeorm-entity.ts или *.entity.ts в src/.
const allFiles = execSync(
`find ${SRC_ROOT} \\( -name "*.typeorm-entity.ts" -o -name "*.entity.ts" \\)`,
{ encoding: 'utf-8' }
)
.split('\n')
.filter(Boolean);
const allKinds = new Set(allFiles.map(entityKindFromFileName));
const stale: string[] = [];
for (const kind of OFF_CHAIN_BASE_ENTITIES) {
if (!allKinds.has(kind)) stale.push(kind);
}
if (stale.length > 0) {
throw new Error(
`OFF_CHAIN_BASE_ENTITIES содержит имена, которых больше нет в src/: ${stale.join(', ')}. Уберите их из allowlist'а.`
);
}
});
});
@@ -0,0 +1,130 @@
/**
* Unit-тесты Story 4.4: BlockchainArchiveRetentionService.cleanup.
*
* Контрактные инварианты:
* - LIB читается через BlockchainService.getInfo() (вариант C, без локальной эвристики).
* - threshold = LIB - RETENTION_HORIZON_BLOCKS (1000, хардкод).
* - При threshold <= 0 skip + log (свежезапущенная testnet).
* - При config.blockchain.archive_retention_enabled=false ранний return, getInfo НЕ вызывается.
* - При ошибке getInfo skip + warn (без re-throw фоновая задача не должна валиться).
* - Удаление идёт ОДНОВРЕМЕННО для entities и для versions через два независимых repo.
*/
import { BlockchainArchiveRetentionService } from '~/shared/sync/services/blockchain-archive-retention.service';
jest.mock('~/config/config', () => ({
__esModule: true,
default: {
blockchain: {
archive_retention_enabled: true,
archive_retention_cron: '0 * * * *',
},
},
}));
function makeLoggerStub(): any {
return {
setContext: jest.fn(),
log: jest.fn(),
debug: jest.fn(),
warn: jest.fn(),
error: jest.fn(),
};
}
function makeBlockchainServiceStub(lib: number): any {
return {
getInfo: jest.fn(async () => ({
last_irreversible_block_num: lib,
head_block_num: lib + 100,
})),
};
}
function makeRepoStub(): any {
return {
deleteOlderThan: jest.fn(async () => 0),
};
}
describe('BlockchainArchiveRetentionService.cleanup (Story 4.4)', () => {
it('LIB=100500: threshold=99500 (LIB-1000), удаляет архив старше threshold', async () => {
const bc = makeBlockchainServiceStub(100500);
const entityRepo = makeRepoStub();
const versionRepo = makeRepoStub();
entityRepo.deleteOlderThan.mockResolvedValueOnce(42);
versionRepo.deleteOlderThan.mockResolvedValueOnce(17);
const service = new BlockchainArchiveRetentionService(bc, entityRepo, versionRepo, makeLoggerStub());
await service.cleanup();
expect(bc.getInfo).toHaveBeenCalledTimes(1);
expect(entityRepo.deleteOlderThan).toHaveBeenCalledWith(99500);
expect(versionRepo.deleteOlderThan).toHaveBeenCalledWith(99500);
});
it('LIB=500 (< 1000): threshold ≤ 0 — skip без вызова deleteOlderThan', async () => {
const bc = makeBlockchainServiceStub(500);
const entityRepo = makeRepoStub();
const versionRepo = makeRepoStub();
const service = new BlockchainArchiveRetentionService(bc, entityRepo, versionRepo, makeLoggerStub());
await service.cleanup();
expect(bc.getInfo).toHaveBeenCalledTimes(1);
expect(entityRepo.deleteOlderThan).not.toHaveBeenCalled();
expect(versionRepo.deleteOlderThan).not.toHaveBeenCalled();
});
it('LIB=1000 (= 1000): threshold = 0, тоже skip', async () => {
const bc = makeBlockchainServiceStub(1000);
const entityRepo = makeRepoStub();
const versionRepo = makeRepoStub();
const service = new BlockchainArchiveRetentionService(bc, entityRepo, versionRepo, makeLoggerStub());
await service.cleanup();
expect(entityRepo.deleteOlderThan).not.toHaveBeenCalled();
});
it('archive_retention_enabled=false: getInfo НЕ вызывается, deleteOlderThan НЕ вызывается', async () => {
const config = (await import('~/config/config')).default;
(config as any).blockchain.archive_retention_enabled = false;
try {
const bc = makeBlockchainServiceStub(100500);
const entityRepo = makeRepoStub();
const versionRepo = makeRepoStub();
const service = new BlockchainArchiveRetentionService(bc, entityRepo, versionRepo, makeLoggerStub());
await service.cleanup();
expect(bc.getInfo).not.toHaveBeenCalled();
expect(entityRepo.deleteOlderThan).not.toHaveBeenCalled();
expect(versionRepo.deleteOlderThan).not.toHaveBeenCalled();
} finally {
(config as any).blockchain.archive_retention_enabled = true;
}
});
it('ошибка getInfo: cleanup НЕ бросает, warn + skip (фоновая задача переживёт сбой ноды)', async () => {
const bc = {
getInfo: jest.fn(async () => {
throw new Error('RPC timeout');
}),
};
const entityRepo = makeRepoStub();
const versionRepo = makeRepoStub();
const logger = makeLoggerStub();
const service = new BlockchainArchiveRetentionService(bc as any, entityRepo, versionRepo, logger);
await expect(service.cleanup()).resolves.toBeUndefined();
expect(entityRepo.deleteOlderThan).not.toHaveBeenCalled();
expect(logger.warn).toHaveBeenCalled();
});
});
@@ -0,0 +1,278 @@
// parser2 — ESM-only пакет; в unit-тестах нам он не нужен, мокаем как virtual,
// иначе jest падает с "Cannot find module" на import top-level.
jest.mock('@coopenomics/parser2', () => ({ ParserClient: class {} }), { virtual: true });
/**
* Unit-тесты BlockchainConsumerService.processFork (Stories 4.1 + 4.2, ADR-005).
*
* Контрактные инварианты:
* - Порядок шагов: forkRegistry.runAll consumerDedup.deleteAfterBlock saveFork.
* - Ошибка в runAll останавливает цепочку: дальнейшие шаги НЕ вызываются, ошибка пробрасывается.
* - Ошибка в deleteAfterBlock останавливает saveFork.
* - blockNum прокидывается одинаково во все шаги.
* - markEventApplied в process{Action,Delta} вызывается с block_num.
* - 4.2: EventEmitter `fork::*` НЕ эмитится (deprecated broadcast удалён).
* - 4.4: forkEventId (controller-формат) прокидывается из handleEvent processFork runAll.
*/
import { BlockchainConsumerService } from '~/infrastructure/blockchain/blockchain-consumer.service';
function makeLoggerStub(): any {
return {
setContext: jest.fn(),
log: jest.fn(),
debug: jest.fn(),
warn: jest.fn(),
error: jest.fn(),
};
}
function makeEventsServiceStub(): any {
return {
emit: jest.fn(),
emitAsyncWithTimeout: jest.fn(async () => true),
};
}
function makeParserInteractorStub(): any {
return {
saveFork: jest.fn(async () => undefined),
saveDelta: jest.fn(async () => undefined),
saveAction: jest.fn(async () => undefined),
isEventApplied: jest.fn(async () => false),
markEventApplied: jest.fn(async () => undefined),
deleteDedupAfterBlock: jest.fn(async () => 0),
};
}
function makeForkRegistryStub(): any {
return {
runAll: jest.fn(async () => undefined),
size: jest.fn(() => 0),
};
}
function makeService(overrides: {
logger?: any;
events?: any;
parserInteractor?: any;
forkRegistry?: any;
}) {
const logger = overrides.logger ?? makeLoggerStub();
const events = overrides.events ?? makeEventsServiceStub();
const parser = overrides.parserInteractor ?? makeParserInteractorStub();
const fork = overrides.forkRegistry ?? makeForkRegistryStub();
const service = new BlockchainConsumerService(logger, events, parser, fork);
return { service, logger, events, parser, fork };
}
describe('BlockchainConsumerService.processFork (Stories 4.1 + 4.2)', () => {
it('вызывает шаги В ПРАВИЛЬНОМ ПОРЯДКЕ: runAll → deleteDedupAfterBlock → saveFork (БЕЗ deprecated emit)', async () => {
const calls: string[] = [];
const parser = makeParserInteractorStub();
parser.deleteDedupAfterBlock.mockImplementation(async () => {
calls.push('deleteDedupAfterBlock');
return 0;
});
parser.saveFork.mockImplementation(async () => {
calls.push('saveFork');
});
const fork = makeForkRegistryStub();
fork.runAll.mockImplementation(async () => {
calls.push('runAll');
});
const events = makeEventsServiceStub();
events.emitAsyncWithTimeout.mockImplementation(async () => {
calls.push('emit');
return true;
});
const { service } = makeService({ events, parserInteractor: parser, forkRegistry: fork });
await (service as any).processFork(100);
expect(calls).toEqual(['runAll', 'deleteDedupAfterBlock', 'saveFork']);
// Story 4.2: emit более НЕ должен вызываться.
expect(events.emitAsyncWithTimeout).not.toHaveBeenCalled();
expect(events.emit).not.toHaveBeenCalledWith(expect.stringMatching(/^fork::/), expect.anything());
});
it('прокидывает forked_from_block во ВСЕ шаги (один и тот же N)', async () => {
const parser = makeParserInteractorStub();
const fork = makeForkRegistryStub();
const events = makeEventsServiceStub();
const { service } = makeService({ events, parserInteractor: parser, forkRegistry: fork });
await (service as any).processFork(12345);
expect(fork.runAll).toHaveBeenCalledWith(12345, undefined);
expect(parser.deleteDedupAfterBlock).toHaveBeenCalledWith(12345);
expect(parser.saveFork).toHaveBeenCalledWith(expect.objectContaining({ block_num: 12345 }));
// Story 4.2: никакого broadcast'а `fork::*` через EventEmitter.
expect(events.emitAsyncWithTimeout).not.toHaveBeenCalled();
});
it('ошибка в forkRegistry.runAll: deleteDedupAfterBlock и saveFork НЕ вызываются, ошибка пробрасывается', async () => {
const parser = makeParserInteractorStub();
const fork = makeForkRegistryStub();
fork.runAll.mockRejectedValueOnce(new Error('syncer #3 rollback failed'));
const { service } = makeService({ parserInteractor: parser, forkRegistry: fork });
await expect((service as any).processFork(100)).rejects.toThrow('syncer #3 rollback failed');
expect(parser.deleteDedupAfterBlock).not.toHaveBeenCalled();
expect(parser.saveFork).not.toHaveBeenCalled();
});
it('ошибка в deleteDedupAfterBlock: saveFork НЕ вызывается, ошибка пробрасывается', async () => {
const parser = makeParserInteractorStub();
parser.deleteDedupAfterBlock.mockRejectedValueOnce(new Error('PG outage'));
const { service } = makeService({ parserInteractor: parser });
await expect((service as any).processFork(100)).rejects.toThrow('PG outage');
expect(parser.saveFork).not.toHaveBeenCalled();
});
it('ошибка в saveFork: emit НЕ вызывается, ошибка пробрасывается', async () => {
const parser = makeParserInteractorStub();
parser.saveFork.mockRejectedValueOnce(new Error('fork insert failed'));
const events = makeEventsServiceStub();
const { service } = makeService({ parserInteractor: parser, events });
await expect((service as any).processFork(100)).rejects.toThrow('fork insert failed');
expect(events.emitAsyncWithTimeout).not.toHaveBeenCalled();
});
// controller считает свой event_id в формате chain:fork:block_num:short_id (полное "fork",
// не parser2-однобуквенное "f") — поэтому ожидаемый ID = 'c1:fork:100:abc12345' (slice 8).
const FORK_EVENT = {
kind: 'fork' as const,
event_id: 'c1:f:100:abc12345', // parser2-формат — controller его НЕ использует
chain_id: 'c1',
forked_from_block: 100,
new_head_block_id: 'abc12345xyz',
};
const EXPECTED_CONTROLLER_ID = 'c1:fork:100:abc12345';
it('fork-event dedup (Story 4.1 AC INV-09): уже-applied event_id — handleEvent ранний return', async () => {
const parser = makeParserInteractorStub();
parser.isEventApplied.mockResolvedValueOnce(true);
const fork = makeForkRegistryStub();
const events = makeEventsServiceStub();
const { service } = makeService({ parserInteractor: parser, forkRegistry: fork, events });
await (service as any).handleEvent(FORK_EVENT);
expect(parser.isEventApplied).toHaveBeenCalledWith(EXPECTED_CONTROLLER_ID);
expect(fork.runAll).not.toHaveBeenCalled();
expect(parser.deleteDedupAfterBlock).not.toHaveBeenCalled();
expect(parser.saveFork).not.toHaveBeenCalled();
expect(parser.markEventApplied).not.toHaveBeenCalled();
expect(events.emitAsyncWithTimeout).not.toHaveBeenCalled();
});
it('fork-event новый: controller использует свой event_id (chain:fork:...), НЕ parser2-формат (chain:f:...)', async () => {
const parser = makeParserInteractorStub();
parser.isEventApplied.mockResolvedValueOnce(false);
const { service } = makeService({ parserInteractor: parser });
await (service as any).handleEvent(FORK_EVENT);
expect(parser.isEventApplied).toHaveBeenCalledWith(EXPECTED_CONTROLLER_ID);
expect(parser.markEventApplied).toHaveBeenCalledWith(EXPECTED_CONTROLLER_ID, 100);
// Никаких следов parser2-формата (`:f:`) в обращениях к dedup-порту:
expect(parser.isEventApplied).not.toHaveBeenCalledWith(expect.stringContaining(':f:'));
expect(parser.markEventApplied).not.toHaveBeenCalledWith(expect.stringContaining(':f:'), expect.anything());
});
it('fork-event новый: ошибка в processFork → markEventApplied НЕ вызывается', async () => {
const parser = makeParserInteractorStub();
parser.isEventApplied.mockResolvedValueOnce(false);
const fork = makeForkRegistryStub();
fork.runAll.mockRejectedValueOnce(new Error('rollback fail'));
const { service } = makeService({ parserInteractor: parser, forkRegistry: fork });
await expect((service as any).handleEvent(FORK_EVENT)).rejects.toThrow('rollback fail');
expect(parser.markEventApplied).not.toHaveBeenCalled();
});
it('Story 4.2 regression: deprecated `fork::*` broadcast полностью удалён — events.emit НИКОГДА не вызывается с fork:: префиксом', async () => {
const events = makeEventsServiceStub();
const { service } = makeService({ events });
await (service as any).processFork(777);
expect(events.emitAsyncWithTimeout).not.toHaveBeenCalled();
// Также проверяем синхронный emit — на случай если рудимент остался в виде events.emit('fork::...').
const allEmitCalls = (events.emit as jest.Mock).mock.calls;
for (const call of allEmitCalls) {
expect(call[0]).not.toMatch(/^fork::/);
}
});
it('Story 4.4: handleEvent для fork прокидывает controller event_id в processFork → runAll', async () => {
const parser = makeParserInteractorStub();
parser.isEventApplied.mockResolvedValueOnce(false);
const fork = makeForkRegistryStub();
const { service } = makeService({ parserInteractor: parser, forkRegistry: fork });
await (service as any).handleEvent(FORK_EVENT);
expect(fork.runAll).toHaveBeenCalledWith(100, EXPECTED_CONTROLLER_ID);
});
it('Story 4.4: processFork(N, eventId) прокидывает eventId в runAll', async () => {
const fork = makeForkRegistryStub();
const { service } = makeService({ forkRegistry: fork });
await (service as any).processFork(555, 'c1:fork:555:deadbeef');
expect(fork.runAll).toHaveBeenCalledWith(555, 'c1:fork:555:deadbeef');
});
});
describe('BlockchainConsumerService.processAction/processDelta — markEventApplied с block_num (Story 4.1)', () => {
// config.coopname — сравниваем с тем, что в тестовой среде. Берём из mock'а конфига если есть,
// иначе тестовый сценарий использует exception (eosio.token::transfer) который проходит без coopname.
it('processAction: markEventApplied вызывается с block_num действия', async () => {
const parser = makeParserInteractorStub();
const { service } = makeService({ parserInteractor: parser });
const action = {
account: 'eosio.token',
name: 'transfer',
receiver: 'eosio.token',
data: { coopname: 'irrelevant' },
block_num: 999,
global_sequence: 1,
} as any;
await (service as any).processActionDelayed(action);
expect(parser.markEventApplied).toHaveBeenCalledWith(expect.any(String), 999);
});
it('processDelta: markEventApplied вызывается с block_num дельты', async () => {
const parser = makeParserInteractorStub();
const { service } = makeService({ parserInteractor: parser });
// config.coopname прокидывается через scope для прохода фильтра.
const { config } = await import('~/config');
const delta = {
code: 'capital',
table: 'projects',
primary_key: '1',
value: { coopname: config.coopname },
scope: config.coopname,
block_num: 12345,
present: true,
} as any;
await (service as any).processDeltaDelayed(delta);
expect(parser.markEventApplied).toHaveBeenCalledWith(expect.any(String), 12345);
});
});
@@ -0,0 +1,99 @@
/**
* Unit-тесты TypeOrmConsumerDedupRepository (Story 2.1).
*
* Фокус контракт идемпотентности: isApplied отражает наличие метки,
* markApplied вставляет с ON CONFLICT DO NOTHING (повтор не падает).
*/
import { TypeOrmConsumerDedupRepository } from '~/infrastructure/database/typeorm/repositories/typeorm-consumer-dedup.repository';
function makeQueryBuilderStub() {
const execute = jest.fn(async () => ({ affected: 1 }));
const qb: any = {
insert: jest.fn(() => qb),
into: jest.fn(() => qb),
values: jest.fn(() => qb),
orIgnore: jest.fn(() => qb),
delete: jest.fn(() => qb),
from: jest.fn(() => qb),
where: jest.fn(() => qb),
execute,
};
return qb;
}
describe('TypeOrmConsumerDedupRepository (Story 2.1)', () => {
it('isApplied → true, когда метка найдена', async () => {
const repoStub: any = { findOne: jest.fn(async () => ({ event_id: 'e1' })) };
const repo = new TypeOrmConsumerDedupRepository(repoStub);
expect(await repo.isApplied('e1')).toBe(true);
});
it('isApplied → false, когда метки нет', async () => {
const repoStub: any = { findOne: jest.fn(async () => null) };
const repo = new TypeOrmConsumerDedupRepository(repoStub);
expect(await repo.isApplied('e1')).toBe(false);
});
it('markApplied без blockNum пишет block_num=null (legacy-вызов)', async () => {
const qb = makeQueryBuilderStub();
const repoStub: any = { createQueryBuilder: jest.fn(() => qb) };
const repo = new TypeOrmConsumerDedupRepository(repoStub);
await repo.markApplied('e1');
expect(qb.values).toHaveBeenCalledWith({ event_id: 'e1', block_num: null });
expect(qb.orIgnore).toHaveBeenCalled();
expect(qb.execute).toHaveBeenCalled();
});
it('markApplied с blockNum пишет block_num как строку (Story 4.1, bigint serialization)', async () => {
const qb = makeQueryBuilderStub();
const repoStub: any = { createQueryBuilder: jest.fn(() => qb) };
const repo = new TypeOrmConsumerDedupRepository(repoStub);
await repo.markApplied('e1', 12345);
expect(qb.values).toHaveBeenCalledWith({ event_id: 'e1', block_num: '12345' });
});
it('deleteOlderThan возвращает число удалённых строк', async () => {
const qb = makeQueryBuilderStub();
qb.execute.mockResolvedValueOnce({ affected: 7 });
const repoStub: any = { createQueryBuilder: jest.fn(() => qb) };
const repo = new TypeOrmConsumerDedupRepository(repoStub);
expect(await repo.deleteOlderThan(new Date())).toBe(7);
});
it('deleteAfterBlock (Story 4.1) выполняет DELETE WHERE block_num > N и возвращает число удалённых', async () => {
const qb = makeQueryBuilderStub();
qb.execute.mockResolvedValueOnce({ affected: 3 });
const repoStub: any = { createQueryBuilder: jest.fn(() => qb) };
const repo = new TypeOrmConsumerDedupRepository(repoStub);
const purged = await repo.deleteAfterBlock(1000);
expect(qb.delete).toHaveBeenCalled();
expect(qb.where).toHaveBeenCalledWith('block_num > :blockNum', { blockNum: 1000 });
expect(purged).toBe(3);
});
it('deleteAfterBlock возвращает 0 при пустом результате', async () => {
const qb = makeQueryBuilderStub();
qb.execute.mockResolvedValueOnce({ affected: 0 });
const repoStub: any = { createQueryBuilder: jest.fn(() => qb) };
const repo = new TypeOrmConsumerDedupRepository(repoStub);
expect(await repo.deleteAfterBlock(999)).toBe(0);
});
it('deleteAfterBlock возвращает 0, если драйвер не вернул affected', async () => {
const qb = makeQueryBuilderStub();
qb.execute.mockResolvedValueOnce({});
const repoStub: any = { createQueryBuilder: jest.fn(() => qb) };
const repo = new TypeOrmConsumerDedupRepository(repoStub);
expect(await repo.deleteAfterBlock(999)).toBe(0);
});
});
@@ -0,0 +1,243 @@
/**
* Unit-тесты Story 4.4: EntityVersioningService.archiveAndDeleteLiveAfterFork /
* archiveAndDeleteVersionsAfterFork.
*
* Контрактные инварианты:
* - SELECT WHERE block_num > N + INSERT в архив + DELETE из исходной таблицы происходят
* в ОДНОЙ транзакции (DataSource.transaction).
* - При нулевом count нет INSERT и нет DELETE (ранний return).
* - forkEventId пробрасывается в архив (NULL допустим).
* - count возвращается корректно.
* - archive-методы НЕ задействуют entityVersionRepository.saveVersion (не путать с saveVersionBeforeUpdate).
*/
import { EntityVersioningService } from '~/shared/sync/services/entity-versioning.service';
import { InvalidatedEntityTypeormEntity } from '~/shared/sync/entities/invalidated-entity.typeorm-entity';
import { InvalidatedEntityVersionTypeormEntity } from '~/shared/sync/entities/invalidated-entity-version.typeorm-entity';
import { EntityVersionTypeormEntity } from '~/shared/sync/entities/entity-version.typeorm-entity';
interface MockRepo {
find: jest.Mock;
delete: jest.Mock;
save: jest.Mock;
create: jest.Mock;
createQueryBuilder?: jest.Mock;
}
function makeMockRepo(): MockRepo {
return {
find: jest.fn(),
delete: jest.fn(async () => ({ affected: 0 })),
save: jest.fn(async (x) => x),
create: jest.fn((x) => x),
};
}
function makeMockDataSource(repoMap: Map<any, MockRepo>): any {
return {
transaction: jest.fn(async (cb: (manager: any) => Promise<any>) => {
const manager = {
getRepository: (target: any) => {
const repo = repoMap.get(target);
if (!repo) throw new Error(`No mock repo for target ${target?.name ?? target}`);
return repo;
},
};
return cb(manager);
}),
};
}
describe('EntityVersioningService.archiveAndDeleteLiveAfterFork (Story 4.4)', () => {
it('SELECT block_num > N → INSERT в архив → DELETE из исходной таблицы; count = найденные ряды', async () => {
const liveTarget = class FakeProject {};
const liveRepo: MockRepo = makeMockRepo();
const archiveRepo: MockRepo = makeMockRepo();
const rows = [
{ _id: 'p1', block_num: 200, name: 'A' },
{ _id: 'p2', block_num: 201, name: 'B' },
];
liveRepo.find.mockResolvedValueOnce(rows);
const liveRepoOuter = { target: liveTarget } as any;
const dataSource = makeMockDataSource(
new Map<any, MockRepo>([
[liveTarget, liveRepo],
[InvalidatedEntityTypeormEntity, archiveRepo],
])
);
const service = new EntityVersioningService(
{} as any,
{} as any,
{} as any,
dataSource
);
const count = await service.archiveAndDeleteLiveAfterFork(
liveRepoOuter,
'capital_projects',
100,
'c1:fork:100:abc12345'
);
expect(count).toBe(2);
expect(liveRepo.find).toHaveBeenCalledTimes(1);
expect(archiveRepo.save).toHaveBeenCalledTimes(1);
const archived = archiveRepo.save.mock.calls[0][0];
expect(archived).toHaveLength(2);
expect(archived[0]).toMatchObject({
entity_table: 'capital_projects',
entity_id: 'p1',
invalidated_by_block: 100,
fork_event_id: 'c1:fork:100:abc12345',
});
expect(archived[0].data).toMatchObject({ name: 'A', block_num: 200 });
expect(liveRepo.delete).toHaveBeenCalledTimes(1);
expect(dataSource.transaction).toHaveBeenCalledTimes(1);
});
it('нет рядов для архивирования — нет INSERT, нет DELETE, count=0', async () => {
const liveTarget = class FakeProject {};
const liveRepo: MockRepo = makeMockRepo();
const archiveRepo: MockRepo = makeMockRepo();
liveRepo.find.mockResolvedValueOnce([]);
const dataSource = makeMockDataSource(
new Map<any, MockRepo>([
[liveTarget, liveRepo],
[InvalidatedEntityTypeormEntity, archiveRepo],
])
);
const service = new EntityVersioningService({} as any, {} as any, {} as any, dataSource);
const count = await service.archiveAndDeleteLiveAfterFork(
{ target: liveTarget } as any,
'capital_projects',
100,
null
);
expect(count).toBe(0);
expect(archiveRepo.save).not.toHaveBeenCalled();
expect(liveRepo.delete).not.toHaveBeenCalled();
});
it('forkEventId=undefined → в архиве fork_event_id=null (NULL допустим, не required)', async () => {
const liveTarget = class FakeProject {};
const liveRepo: MockRepo = makeMockRepo();
const archiveRepo: MockRepo = makeMockRepo();
liveRepo.find.mockResolvedValueOnce([{ _id: 'p1', block_num: 200 }]);
const dataSource = makeMockDataSource(
new Map<any, MockRepo>([
[liveTarget, liveRepo],
[InvalidatedEntityTypeormEntity, archiveRepo],
])
);
const service = new EntityVersioningService({} as any, {} as any, {} as any, dataSource);
await service.archiveAndDeleteLiveAfterFork({ target: liveTarget } as any, 'capital_projects', 100);
const archived = archiveRepo.save.mock.calls[0][0];
expect(archived[0].fork_event_id).toBeNull();
});
});
describe('EntityVersioningService.archiveAndDeleteVersionsAfterFork (Story 4.4)', () => {
it('SELECT entity_versions WHERE entity_table AND block_num > N → INSERT → DELETE по ids', async () => {
const versionsRepo: any = {
createQueryBuilder: jest.fn(() => versionsRepo),
where: jest.fn(() => versionsRepo),
andWhere: jest.fn(() => versionsRepo),
getMany: jest.fn(async () => [
{
id: 'v1',
entity_table: 'capital_projects',
entity_id: 'p1',
previous_data: { name: 'A_old' },
block_num: 150,
change_type: 'update',
metadata: null,
},
{
id: 'v2',
entity_table: 'capital_projects',
entity_id: 'p2',
previous_data: { name: 'B_old' },
block_num: null,
change_type: 'update',
metadata: { foo: 1 },
},
]),
delete: jest.fn(() => versionsRepo),
from: jest.fn(() => versionsRepo),
whereInIds: jest.fn(() => versionsRepo),
execute: jest.fn(async () => ({ affected: 2 })),
save: jest.fn(async (x) => x),
create: jest.fn((x) => x),
};
const archiveRepo: MockRepo = makeMockRepo();
const dataSource = makeMockDataSource(
new Map<any, MockRepo>([
[EntityVersionTypeormEntity, versionsRepo],
[InvalidatedEntityVersionTypeormEntity, archiveRepo],
])
);
const service = new EntityVersioningService({} as any, {} as any, {} as any, dataSource);
const count = await service.archiveAndDeleteVersionsAfterFork(
'capital_projects',
100,
'c1:fork:100:deadbeef'
);
expect(count).toBe(2);
expect(archiveRepo.save).toHaveBeenCalledTimes(1);
const archived = archiveRepo.save.mock.calls[0][0];
expect(archived).toHaveLength(2);
expect(archived[0]).toMatchObject({
entity_table: 'capital_projects',
entity_id: 'p1',
previous_data: { name: 'A_old' },
original_block_num: 150,
invalidated_by_block: 100,
fork_event_id: 'c1:fork:100:deadbeef',
change_type: 'update',
});
// v2 имел block_num=null → original_block_num=null
expect(archived[1].original_block_num).toBeNull();
expect(versionsRepo.whereInIds).toHaveBeenCalledWith(['v1', 'v2']);
expect(versionsRepo.execute).toHaveBeenCalledTimes(1);
});
it('нет версий для архивирования — нет INSERT, нет DELETE', async () => {
const versionsRepo: any = {
createQueryBuilder: jest.fn(() => versionsRepo),
where: jest.fn(() => versionsRepo),
andWhere: jest.fn(() => versionsRepo),
getMany: jest.fn(async () => []),
};
const archiveRepo: MockRepo = makeMockRepo();
const dataSource = makeMockDataSource(
new Map<any, MockRepo>([
[EntityVersionTypeormEntity, versionsRepo],
[InvalidatedEntityVersionTypeormEntity, archiveRepo],
])
);
const service = new EntityVersioningService({} as any, {} as any, {} as any, dataSource);
const count = await service.archiveAndDeleteVersionsAfterFork('capital_projects', 100, null);
expect(count).toBe(0);
expect(archiveRepo.save).not.toHaveBeenCalled();
});
});
@@ -0,0 +1,104 @@
/**
* Unit-тесты computeDeltaEventId / computeActionEventId (Story 2.2).
*
* event_id фундамент идемпотентности (INV-09). Формула должна быть
* детерминированной и не схлопывать разные логические события в один id
* (иначе dedup-gate Story 2.3 отбросит реальное событие = silent data loss).
* В Epic 3 phase 2 эти id сверяются с авторитетными из parser2.
*/
import { computeActionEventId, computeDeltaEventId, computeForkEventId } from '~/infrastructure/blockchain/event-id.util';
function makeDelta(over: Partial<any> = {}): any {
return {
chain_id: 'chain-aaa',
block_num: 100,
block_id: '0123456789abcdef',
present: true,
code: 'eosio.token',
scope: 'voskhod',
table: 'accounts',
primary_key: '42',
...over,
};
}
function makeAction(over: Partial<any> = {}): any {
return {
chain_id: 'chain-aaa',
block_num: 100,
block_id: '0123456789abcdef',
global_sequence: '777',
account: 'eosio.token',
name: 'transfer',
...over,
};
}
describe('computeDeltaEventId (Story 2.2)', () => {
it('формирует id по формуле chain:delta:block:block_id_short:code:scope:table:primary_key', () => {
expect(computeDeltaEventId(makeDelta())).toBe('chain-aaa:delta:100:01234567:eosio.token:voskhod:accounts:42');
});
it('детерминирован: один и тот же вход → один и тот же id', () => {
expect(computeDeltaEventId(makeDelta())).toBe(computeDeltaEventId(makeDelta()));
});
it('разный primary_key → разный id', () => {
expect(computeDeltaEventId(makeDelta({ primary_key: '1' }))).not.toBe(
computeDeltaEventId(makeDelta({ primary_key: '2' }))
);
});
it('разный блок → разный id (один и тот же ряд, обновлённый в другом блоке)', () => {
expect(computeDeltaEventId(makeDelta({ block_num: 100, block_id: 'aaaa1111' }))).not.toBe(
computeDeltaEventId(makeDelta({ block_num: 101, block_id: 'bbbb2222' }))
);
});
it('обрезает block_id до 8 символов; короткий block_id не падает', () => {
expect(computeDeltaEventId(makeDelta({ block_id: 'ab' }))).toBe('chain-aaa:delta:100:ab:eosio.token:voskhod:accounts:42');
});
});
describe('computeActionEventId (Story 2.2)', () => {
it('формирует id по формуле chain:action:block:block_id_short:global_sequence', () => {
expect(computeActionEventId(makeAction())).toBe('chain-aaa:action:100:01234567:777');
});
it('разный global_sequence → разный id', () => {
expect(computeActionEventId(makeAction({ global_sequence: '1' }))).not.toBe(
computeActionEventId(makeAction({ global_sequence: '2' }))
);
});
it('детерминирован', () => {
expect(computeActionEventId(makeAction())).toBe(computeActionEventId(makeAction()));
});
});
describe('computeForkEventId (Story 4.1)', () => {
it('формирует id по формуле chain:fork:forked_from_block:new_head_block_id_short (kind = полное "fork", не parser2 "f")', () => {
expect(computeForkEventId('chain-aaa', 1000, '0123456789abcdef')).toBe('chain-aaa:fork:1000:01234567');
});
it('одинаковые входы → одинаковый id (детерминирован)', () => {
expect(computeForkEventId('chain-aaa', 1000, 'aaaa1111')).toBe(computeForkEventId('chain-aaa', 1000, 'aaaa1111'));
});
it('разный new_head_block_id (две разных «новых ветки» того же forked_from_block) → разный id', () => {
expect(computeForkEventId('chain-aaa', 1000, 'aaaaaaaa')).not.toBe(
computeForkEventId('chain-aaa', 1000, 'bbbbbbbb')
);
});
it('разный forked_from_block → разный id', () => {
expect(computeForkEventId('chain-aaa', 1000, 'aaaa1111')).not.toBe(
computeForkEventId('chain-aaa', 1001, 'aaaa1111')
);
});
it('короткий new_head_block_id (< 8 символов) не падает', () => {
expect(computeForkEventId('chain-aaa', 1000, 'ab')).toBe('chain-aaa:fork:1000:ab');
});
});
@@ -0,0 +1,116 @@
/**
* Unit-тест Story 4.4: AbstractEntitySyncService.handleFork порядок шагов.
*
* Контрактные инварианты:
* - Порядок: archiveInvalidatedSince restoreFromVersions archiveInvalidatedVersionsSince.
* - forkEventId пробрасывается в archive-методы.
* - Если repo НЕ реализует archiveInvalidatedSince fallback на findByBlockNumGreaterThan + deleteByBlockNumGreaterThan (бэк-совместимость для off-chain).
* - Ошибка в archiveInvalidatedSince re-throw'ится, restore не вызывается.
*/
import { AbstractEntitySyncService } from '~/shared/services/abstract-entity-sync.service';
function makeLoggerStub(): any {
return {
setContext: jest.fn(),
log: jest.fn(),
debug: jest.fn(),
warn: jest.fn(),
error: jest.fn(),
};
}
function makeMapperStub(): any {
return {
extractSyncValue: jest.fn(),
extractSyncKey: jest.fn(() => 'id'),
mapDeltaToBlockchainData: jest.fn(),
getAllEventPatterns: jest.fn(() => []),
getSupportedTableNames: jest.fn(() => []),
getSupportedContractNames: jest.fn(() => []),
};
}
class TestSyncService extends AbstractEntitySyncService<any, any> {
protected readonly entityName = 'TestEntity';
}
describe('AbstractEntitySyncService.handleFork (Story 4.4)', () => {
it('Story 4.4 happy-path: archiveInvalidatedSince → restoreFromVersions → archiveInvalidatedVersionsSince; eventId пробрасывается', async () => {
const calls: Array<{ name: string; args: any[] }> = [];
const repo: any = {
archiveInvalidatedSince: jest.fn(async (...args) => {
calls.push({ name: 'archiveInvalidatedSince', args });
return 3;
}),
restoreFromVersions: jest.fn(async (...args) => {
calls.push({ name: 'restoreFromVersions', args });
}),
archiveInvalidatedVersionsSince: jest.fn(async (...args) => {
calls.push({ name: 'archiveInvalidatedVersionsSince', args });
return 5;
}),
};
const service = new TestSyncService(repo, makeMapperStub(), makeLoggerStub());
await service.handleFork(100, 'c1:fork:100:beef');
expect(calls.map((c) => c.name)).toEqual([
'archiveInvalidatedSince',
'restoreFromVersions',
'archiveInvalidatedVersionsSince',
]);
expect(repo.archiveInvalidatedSince).toHaveBeenCalledWith(100, 'c1:fork:100:beef');
expect(repo.restoreFromVersions).toHaveBeenCalledWith(100);
expect(repo.archiveInvalidatedVersionsSince).toHaveBeenCalledWith(100, 'c1:fork:100:beef');
});
it('Story 4.4: forkEventId необязателен — без него archive получают undefined', async () => {
const repo: any = {
archiveInvalidatedSince: jest.fn(async () => 0),
restoreFromVersions: jest.fn(async () => undefined),
archiveInvalidatedVersionsSince: jest.fn(async () => 0),
};
const service = new TestSyncService(repo, makeMapperStub(), makeLoggerStub());
await service.handleFork(100);
expect(repo.archiveInvalidatedSince).toHaveBeenCalledWith(100, undefined);
expect(repo.archiveInvalidatedVersionsSince).toHaveBeenCalledWith(100, undefined);
});
it('Story 4.4 fallback: репо без archive-методов → старая пара findByBlockNumGreaterThan + deleteByBlockNumGreaterThan (off-chain совместимость)', async () => {
const repo: any = {
findByBlockNumGreaterThan: jest.fn(async () => [{ _id: 'e1' }, { _id: 'e2' }]),
deleteByBlockNumGreaterThan: jest.fn(async () => undefined),
restoreFromVersions: jest.fn(async () => undefined),
};
const service = new TestSyncService(repo, makeMapperStub(), makeLoggerStub());
await service.handleFork(100, 'eventid');
expect(repo.findByBlockNumGreaterThan).toHaveBeenCalledWith(100);
expect(repo.deleteByBlockNumGreaterThan).toHaveBeenCalledWith(100);
expect(repo.restoreFromVersions).toHaveBeenCalledWith(100);
});
it('Story 4.4: ошибка archiveInvalidatedSince → re-throw, restoreFromVersions и archiveVersions НЕ вызываются', async () => {
const repo: any = {
archiveInvalidatedSince: jest.fn(async () => {
throw new Error('archive PG outage');
}),
restoreFromVersions: jest.fn(async () => undefined),
archiveInvalidatedVersionsSince: jest.fn(async () => 0),
};
const service = new TestSyncService(repo, makeMapperStub(), makeLoggerStub());
await expect(service.handleFork(100, 'eventid')).rejects.toThrow('archive PG outage');
expect(repo.restoreFromVersions).not.toHaveBeenCalled();
expect(repo.archiveInvalidatedVersionsSince).not.toHaveBeenCalled();
});
});
@@ -0,0 +1,45 @@
/**
* Regression-тест Story 4.2: после удаления pause-barrier и deprecated broadcast
* паттерна `@OnEvent('fork::*')` в исходниках controller'а быть НЕ должно.
*
* Контракт: единственный путь обработки форка через ForkRegistryService (см.
* fork-registry.service.ts + AbstractEntitySyncService). EventEmitter2 для
* `fork::*` больше не используется. Тест ловит регрессию, если кто-то случайно
* вернёт декоратор `@OnEvent('fork::*')` в новый syncer.
*
* Сканирует `src/` целиком; исключает комментарии (строки, начинающиеся с *).
*/
import { execSync } from 'child_process';
import * as path from 'path';
const SRC_ROOT = path.resolve(__dirname, '../../../src');
describe('Story 4.2: @OnEvent(fork::*) regression guard', () => {
it('в src/ нет ни одного активного @OnEvent fork::* декоратора', () => {
// -R рекурсивно, -E расширенные regexp, --include=*.ts только TS-файлы,
// -h без имени файла, || true чтобы exit 1 (no matches) не падал в jest.
const output = execSync(
`grep -REn --include="*.ts" "@OnEvent\\(['\\\"]fork::" ${SRC_ROOT} || true`,
{ encoding: 'utf-8' }
).trim();
// Отфильтруем строки, которые — комментарии (* в начале после пробелов или // ).
const offending = output
.split('\n')
.filter((line) => line.length > 0)
.filter((line) => {
// формат grep -n: path:line:content
const content = line.split(':').slice(2).join(':').trimStart();
if (content.startsWith('*') || content.startsWith('//')) return false;
return true;
});
if (offending.length > 0) {
throw new Error(
`Найдены активные @OnEvent('fork::*') декораторы (Story 4.2 должна была их удалить):\n` +
offending.join('\n')
);
}
});
});
@@ -0,0 +1,141 @@
/**
* Unit-тесты mapParserDeltaToIDelta / mapParserActionToIAction (миграция на parser2).
*
* Маппер единственная точка перевода ParserEvent IDelta/IAction при смене
* транспорта (DEC-T09). Все поля действия реальные из SHiP-трейса parser2 паритет
* с parser1: transaction_id/creator_action_ordinal (нужны ledger2 для cross-link
* родительского apply), receipt с auth_sequence, console/elapsed/context_free/
* account_ram_deltas (нужны blockchain-explorer).
*/
import { mapParserActionToIAction, mapParserDeltaToIDelta } from '~/infrastructure/blockchain/parser2-event.mapper';
function deltaEvent(over: Record<string, any> = {}): any {
return {
kind: 'delta',
event_id: 'chain-aaa:d:100:0123456789abcdef:eosio.token:voskhod:accounts:42',
chain_id: 'chain-aaa',
block_num: 100,
block_time: '2026-05-25T00:00:00.000Z',
block_id: '0123456789abcdef0000',
code: 'eosio.token',
scope: 'voskhod',
table: 'accounts',
primary_key: '42',
value: { balance: '10.0000 RUB', coopname: 'voskhod' },
present: true,
...over,
};
}
function actionEvent(over: Record<string, any> = {}): any {
return {
kind: 'action',
event_id: 'chain-aaa:a:100:0123456789abcdef:777',
chain_id: 'chain-aaa',
block_num: 100,
block_time: '2026-05-25T00:00:00.000Z',
block_id: '0123456789abcdef0000',
transaction_id: 'abc123def456',
account: 'eosio.token',
name: 'transfer',
authorization: [{ actor: 'voskhod', permission: 'active' }],
data: { from: 'voskhod', to: 'ant', quantity: '1.0000 RUB' },
action_ordinal: 2,
creator_action_ordinal: 1,
global_sequence: 777n,
receipt: {
receiver: 'eosio.token',
actDigest: 'deadbeef',
globalSequence: 777n,
recvSequence: 5n,
authSequence: [{ account: 'voskhod', sequence: 9n }],
codeSequence: 3,
abiSequence: 4,
},
context_free: false,
elapsed: 42,
console: 'log-output',
account_ram_deltas: [{ account: 'voskhod', delta: 128 }],
...over,
};
}
describe('mapParserDeltaToIDelta', () => {
it('мапит все поля DeltaEvent один-в-один', () => {
expect(mapParserDeltaToIDelta(deltaEvent())).toEqual({
chain_id: 'chain-aaa',
block_num: 100,
block_id: '0123456789abcdef0000',
block_time: '2026-05-25T00:00:00.000Z',
present: true,
code: 'eosio.token',
scope: 'voskhod',
table: 'accounts',
primary_key: '42',
value: { balance: '10.0000 RUB', coopname: 'voskhod' },
});
});
it('пробрасывает present=false (удаление строки)', () => {
expect(mapParserDeltaToIDelta(deltaEvent({ present: false })).present).toBe(false);
});
it('пробрасывает block_time из SHiP-трейса (parser1 этого поля не давал)', () => {
expect(mapParserDeltaToIDelta(deltaEvent()).block_time).toBe('2026-05-25T00:00:00.000Z');
});
});
describe('mapParserActionToIAction', () => {
it('пробрасывает реальные поля action и сериализует global_sequence (bigint → string)', () => {
const a = mapParserActionToIAction(actionEvent());
expect(a.account).toBe('eosio.token');
expect(a.name).toBe('transfer');
expect(a.block_num).toBe(100);
expect(a.data).toEqual({ from: 'voskhod', to: 'ant', quantity: '1.0000 RUB' });
expect(a.global_sequence).toBe('777');
expect(typeof a.global_sequence).toBe('string');
expect(a.authorization).toEqual([{ actor: 'voskhod', permission: 'active' }]);
});
it('несёт transaction_id и creator_action_ordinal (нужны ledger2 для cross-link)', () => {
const a = mapParserActionToIAction(actionEvent());
expect(a.transaction_id).toBe('abc123def456');
expect(a.action_ordinal).toBe(2);
expect(a.creator_action_ordinal).toBe(1);
});
it('мапит receipt с auth_sequence (bigint → string)', () => {
const a = mapParserActionToIAction(actionEvent());
expect(a.receipt).toEqual({
receiver: 'eosio.token',
act_digest: 'deadbeef',
global_sequence: '777',
recv_sequence: '5',
auth_sequence: [{ account: 'voskhod', sequence: '9' }],
code_sequence: 3,
abi_sequence: 4,
});
expect(a.receiver).toBe('eosio.token');
});
it('несёт console/elapsed/context_free/account_ram_deltas (нужны blockchain-explorer)', () => {
const a = mapParserActionToIAction(actionEvent());
expect(a.console).toBe('log-output');
expect(a.elapsed).toBe(42);
expect(a.context_free).toBe(false);
expect(a.account_ram_deltas).toEqual([{ account: 'voskhod', delta: 128 }]);
});
it('пробрасывает block_time из SHiP-трейса (parser1 этого поля не давал)', () => {
const a = mapParserActionToIAction(actionEvent());
expect(a.block_time).toBe('2026-05-25T00:00:00.000Z');
});
it('при отсутствии receipt (null) собирает дефолтный receipt, не падает', () => {
const a = mapParserActionToIAction(actionEvent({ receipt: null }));
expect(a.receipt.receiver).toBe('eosio.token');
expect(a.receipt.global_sequence).toBe('777');
expect(a.receipt.auth_sequence).toEqual([]);
});
});
@@ -0,0 +1,133 @@
/**
* Story 6.5 (Epic 6): unit-тесты для:
* - `UnsupportedContractVersionError` носит контекст (contract/table/primary_key/block_num).
* - `auditUnknownStatus` пишет error в logger с ожидаемыми статусами.
* - `AbstractEntitySyncService.processDelta` strict-mode throw, default log.error + null.
*
* Конфиг strict-mode мокается через `jest.mock('~/config/config', ...)`.
*/
let strictMode = false;
jest.mock('~/config/config', () => ({
__esModule: true,
default: {
get blockchain() {
return { unsupported_version_strict: strictMode };
},
},
}));
import { AbstractEntitySyncService } from '~/shared/services/abstract-entity-sync.service';
import { UnsupportedContractVersionError } from '~/shared/sync/errors/unsupported-contract-version.error';
import { auditUnknownStatus } from '~/shared/sync/errors/audit-unknown-status';
function makeLoggerStub(): any {
return {
setContext: jest.fn(),
log: jest.fn(),
debug: jest.fn(),
warn: jest.fn(),
error: jest.fn(),
};
}
function makeMapperStub(canMap: boolean): any {
return {
extractSyncValue: jest.fn(() => 'sync-value'),
extractSyncKey: jest.fn(() => 'id'),
mapDeltaToBlockchainData: jest.fn(() => (canMap ? { v: 1 } : null)),
getAllEventPatterns: jest.fn(() => []),
getSupportedTableNames: jest.fn(() => []),
getSupportedContractNames: jest.fn(() => []),
};
}
class TestSyncService extends AbstractEntitySyncService<any, any> {
protected readonly entityName = 'TestEntity';
}
describe('Story 6.5: UnsupportedContractVersionError', () => {
it('конструктор сохраняет entityName и context', () => {
const err = new UnsupportedContractVersionError('Project', {
contract: 'capital',
table: 'projects',
primary_key: 42,
block_num: 100,
});
expect(err.name).toBe('UnsupportedContractVersionError');
expect(err.entityName).toBe('Project');
expect(err.context.primary_key).toBe(42);
expect(err.message).toContain('Project');
expect(err.message).toContain('"table":"projects"');
});
});
describe('Story 6.5: auditUnknownStatus', () => {
it('пишет logger.error с контекстом и списком ожидаемых статусов', () => {
const logger = makeLoggerStub();
auditUnknownStatus('Project', 'something-unknown', logger, ['pending', 'active']);
expect(logger.error).toHaveBeenCalledTimes(1);
const [msg, meta] = logger.error.mock.calls[0];
expect(msg).toContain('UNKNOWN_ENTITY_STATUS');
expect(msg).toContain('Project');
expect(msg).toContain('something-unknown');
expect(msg).toContain('pending');
expect(meta).toMatchObject({ entityName: 'Project', receivedStatus: 'something-unknown' });
});
it('без allowedStatuses пишет "не указано"', () => {
const logger = makeLoggerStub();
auditUnknownStatus('Other', 'x', logger);
expect(logger.error.mock.calls[0][0]).toContain('не указано');
});
});
describe('Story 6.5: processDelta — silent loss заменён audit + опциональный throw', () => {
const delta = {
contract: 'capital',
table: 'projects',
primary_key: 'abc',
block_num: 100,
present: true,
value: {},
} as any;
beforeEach(() => {
strictMode = false;
});
it('non-strict: mapper вернул null → logger.error("UNSUPPORTED_CONTRACT_VERSION", ...) + return null', async () => {
const logger = makeLoggerStub();
const repo: any = {};
const service = new TestSyncService(repo, makeMapperStub(false), logger);
const result = await service.processDelta(delta);
expect(result).toBeNull();
expect(logger.error).toHaveBeenCalled();
expect(logger.error.mock.calls[0][0]).toContain('UNSUPPORTED_CONTRACT_VERSION');
});
it('strict: mapper вернул null → throw UnsupportedContractVersionError (парсер не ACK\'нет)', async () => {
strictMode = true;
const logger = makeLoggerStub();
const repo: any = {};
const service = new TestSyncService(repo, makeMapperStub(false), logger);
await expect(service.processDelta(delta)).rejects.toBeInstanceOf(UnsupportedContractVersionError);
expect(logger.error).toHaveBeenCalled();
});
it('mapper вернул данные → happy-path продолжается (handleSyncDelta)', async () => {
const logger = makeLoggerStub();
const repo: any = {
findBySyncKey: jest.fn(async () => null),
createIfNotExists: jest.fn(async () => undefined),
};
const service = new TestSyncService(repo, makeMapperStub(true), logger);
const result = await service.processDelta(delta);
expect(result).not.toBeNull();
expect(repo.createIfNotExists).toHaveBeenCalled();
});
});
@@ -0,0 +1,46 @@
/**
* Unit-тесты EventsService.emitAsyncWithTimeout (Story 1.3).
*
* Фокус барьер форка: ждём завершения обработчиков отката, но не дольше
* timeoutMs (TTL force-resume), и не глотаем ошибку обработчика.
*/
import { EventsService } from '~/infrastructure/events/events.service';
function makeService(emitAsyncImpl: (e: string, d: any) => Promise<any>) {
const emitter = {
emit: jest.fn(),
emitAsync: jest.fn(emitAsyncImpl),
} as any;
return { service: new EventsService(emitter), emitter };
}
describe('EventsService.emitAsyncWithTimeout (Story 1.3)', () => {
it('возвращает true, когда обработчики завершились в срок', async () => {
const { service } = makeService(async () => [undefined]);
const ok = await service.emitAsyncWithTimeout('fork::100', { block_num: 100 }, 1000);
expect(ok).toBe(true);
});
it('возвращает false (force-resume), когда обработчик завис дольше timeoutMs', async () => {
const { service } = makeService(() => new Promise(() => {/* никогда не резолвится */}));
const ok = await service.emitAsyncWithTimeout('fork::100', { block_num: 100 }, 20);
expect(ok).toBe(false);
});
it('пробрасывает ошибку обработчика (не глотаем)', async () => {
const { service } = makeService(async () => {
throw new Error('rollback failed');
});
await expect(service.emitAsyncWithTimeout('fork::100', { block_num: 100 }, 1000)).rejects.toThrow(
'rollback failed'
);
});
it('emitAsync делегирует в EventEmitter2', async () => {
const { service, emitter } = makeService(async () => ['r']);
const res = await service.emitAsync('delta::x::y', { a: 1 });
expect(emitter.emitAsync).toHaveBeenCalledWith('delta::x::y', { a: 1 });
expect(res).toEqual(['r']);
});
});
@@ -0,0 +1,296 @@
/**
* Unit-тесты ForkRegistryService (Story 4.1, ADR-005).
*
* Контрактные инварианты:
* - runAll обходит syncer'ов строго sequential (for-of await, не Promise.all) INV-T03.
* - Любая ошибка re-throw'ится наверх (silent catch ломает barrier-контракт parser2).
* - register идемпотентен (повторная регистрация того же экземпляра no-op).
* - forkRollbackPriority перекрывает порядок регистрации (FK-зависимости).
* - onApplicationBootstrap собирает провайдеров с FORK_AWARE_MARKER через DiscoveryService.
*/
import { ForkRegistryService } from '~/shared/sync/fork/fork-registry.service';
import { FORK_AWARE_MARKER, type IForkAwareSyncer } from '~/shared/sync/fork/fork-aware-syncer.interface';
function makeLoggerStub(): any {
return {
setContext: jest.fn(),
log: jest.fn(),
debug: jest.fn(),
warn: jest.fn(),
error: jest.fn(),
};
}
function makeDiscoveryStub(providers: Array<{ instance: unknown }>): any {
return { getProviders: jest.fn(() => providers) };
}
class FakeSyncer implements IForkAwareSyncer {
readonly [FORK_AWARE_MARKER] = true;
readonly trace: number[] = [];
constructor(public readonly id: number, public readonly forkRollbackPriority?: number) {}
async handleFork(blockNum: number): Promise<void> {
this.trace.push(blockNum);
}
}
describe('ForkRegistryService (Story 4.1)', () => {
describe('register / unregister / size', () => {
it('register добавляет syncer; size отражает количество', () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const a = new FakeSyncer(1);
const b = new FakeSyncer(2);
expect(service.size()).toBe(0);
service.register(a);
service.register(b);
expect(service.size()).toBe(2);
});
it('register идемпотентен: повторная регистрация того же экземпляра — no-op', () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const a = new FakeSyncer(1);
service.register(a);
service.register(a);
service.register(a);
expect(service.size()).toBe(1);
});
it('unregister снимает syncer; повторный unregister отсутствующего — no-op', () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const a = new FakeSyncer(1);
service.register(a);
service.unregister(a);
expect(service.size()).toBe(0);
service.unregister(a);
expect(service.size()).toBe(0);
});
it('clear обнуляет реестр', () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
service.register(new FakeSyncer(1));
service.register(new FakeSyncer(2));
service.clear();
expect(service.size()).toBe(0);
});
});
describe('runAll — sequential apply', () => {
it('обходит syncer-ов в порядке регистрации', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const order: number[] = [];
const a: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
order.push(1);
},
} as any;
const b: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
order.push(2);
},
} as any;
const c: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
order.push(3);
},
} as any;
service.register(a);
service.register(b);
service.register(c);
await service.runAll(100);
expect(order).toEqual([1, 2, 3]);
});
it('каждый syncer стартует ТОЛЬКО после resolve предыдущего (NOT Promise.all)', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const events: string[] = [];
const slow: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
events.push('slow:start');
await new Promise((r) => setTimeout(r, 20));
events.push('slow:end');
},
} as any;
const fast: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
events.push('fast:start');
events.push('fast:end');
},
} as any;
service.register(slow);
service.register(fast);
await service.runAll(100);
expect(events).toEqual(['slow:start', 'slow:end', 'fast:start', 'fast:end']);
});
it('runAll прокидывает blockNum в handleFork каждого syncer-а', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const a = new FakeSyncer(1);
const b = new FakeSyncer(2);
service.register(a);
service.register(b);
await service.runAll(12345);
expect(a.trace).toEqual([12345]);
expect(b.trace).toEqual([12345]);
});
it('пустой registry: runAll — no-op (без ошибок)', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
await expect(service.runAll(100)).resolves.toBeUndefined();
});
});
describe('runAll — error propagation', () => {
it('первая ошибка пробрасывается наверх; последующие syncer-ы НЕ запускаются', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const trace: string[] = [];
const ok1: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
trace.push('ok1');
},
} as any;
const fail: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
trace.push('fail');
throw new Error('rollback failed');
},
} as any;
const ok2: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
trace.push('ok2');
},
} as any;
service.register(ok1);
service.register(fail);
service.register(ok2);
await expect(service.runAll(100)).rejects.toThrow('rollback failed');
expect(trace).toEqual(['ok1', 'fail']);
});
it('после ошибки registry остаётся валидным: повторный runAll проходит', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
let attempt = 0;
const flaky: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
attempt += 1;
if (attempt === 1) throw new Error('first fail');
},
} as any;
service.register(flaky);
await expect(service.runAll(100)).rejects.toThrow('first fail');
await expect(service.runAll(100)).resolves.toBeUndefined();
expect(attempt).toBe(2);
});
});
describe('forkRollbackPriority — порядок при FK-зависимостях', () => {
it('syncer с приоритетом 1 идёт раньше syncer с приоритетом 5', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const trace: number[] = [];
const low: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
forkRollbackPriority: 5,
async handleFork() {
trace.push(5);
},
} as any;
const high: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
forkRollbackPriority: 1,
async handleFork() {
trace.push(1);
},
} as any;
service.register(low);
service.register(high);
await service.runAll(100);
expect(trace).toEqual([1, 5]);
});
it('syncer с приоритетом идёт раньше syncer без приоритета', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
const trace: string[] = [];
const noPrio: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {
trace.push('noPrio');
},
} as any;
const prio: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
forkRollbackPriority: 0,
async handleFork() {
trace.push('prio');
},
} as any;
service.register(noPrio);
service.register(prio);
await service.runAll(100);
expect(trace).toEqual(['prio', 'noPrio']);
});
});
describe('onApplicationBootstrap — pull-сбор через DiscoveryService', () => {
it('собирает только провайдеров с FORK_AWARE_MARKER', async () => {
const aware: IForkAwareSyncer = {
[FORK_AWARE_MARKER]: true,
async handleFork() {},
} as any;
const notAware = { handleFork: async () => {} }; // нет marker'а
const someService = { someMethod: () => {} };
const nullInstance = null;
const discovery = makeDiscoveryStub([
{ instance: aware },
{ instance: notAware },
{ instance: someService },
{ instance: nullInstance },
]);
const service = new ForkRegistryService(discovery, makeLoggerStub());
await service.onApplicationBootstrap();
expect(service.size()).toBe(1);
});
it('не падает при пустом списке провайдеров', async () => {
const service = new ForkRegistryService(makeDiscoveryStub([]), makeLoggerStub());
await expect(service.onApplicationBootstrap()).resolves.toBeUndefined();
expect(service.size()).toBe(0);
});
it('игнорирует FORK_AWARE_MARKER без handleFork (защита от ложного срабатывания)', async () => {
const fake = { [FORK_AWARE_MARKER]: true }; // marker есть, метода нет
const discovery = makeDiscoveryStub([{ instance: fake }]);
const service = new ForkRegistryService(discovery, makeLoggerStub());
await service.onApplicationBootstrap();
expect(service.size()).toBe(0);
});
});
});
@@ -0,0 +1,110 @@
/**
* Unit-тесты BaseBlockchainRepository.createIfNotExists (Story 1.1, DEC-008).
*
* Фокус guard монотонности block_num на create-пути: устаревшая дельта
* (из более раннего блока) не должна затирать более свежую запись, иначе
* состояние в БД откатывается назад при гонке дельт.
*/
import { BaseBlockchainRepository } from '~/shared/sync/repositories/base-blockchain.repository';
function makeDomain(data: any) {
return {
_id: data._id ?? 'id-1',
block_num: data.block_num,
present: data.present ?? true,
getBlockNum: () => data.block_num,
getPrimaryKey: () => String(data.username ?? 'pk'),
getSyncKey: () => 'username',
updateFromBlockchain: jest.fn(),
};
}
/** Конкретный наследник абстрактного репозитория для теста. */
class TestRepo extends BaseBlockchainRepository<any, any> {
constructor(repo: any, versioning: any) {
super(repo, versioning);
}
protected getMapper() {
return {
toDomain: (e: any) => (e && e.getBlockNum ? e : makeDomain(e)),
toEntity: (d: any) => ({ ...d }),
};
}
protected createDomainEntity(databaseData: any, blockchainData: any) {
return makeDomain({ ...databaseData, ...blockchainData });
}
protected getSyncKey() {
return 'username';
}
}
function makeRepoStub(existingDomain: any) {
return {
target: { getTableName: () => 'test_table' },
findOne: jest.fn(async () => existingDomain),
save: jest.fn(async (e: any) => e),
} as any;
}
function makeVersioningStub() {
return { saveVersionBeforeUpdate: jest.fn(async () => undefined) } as any;
}
describe('BaseBlockchainRepository.createIfNotExists — guard block_num (Story 1.1)', () => {
it('не перезаписывает свежую запись устаревшей дельтой (block_num < N)', async () => {
const existing = makeDomain({ username: 'alice', block_num: 100 });
const repoStub = makeRepoStub(existing);
const repo = new TestRepo(repoStub, makeVersioningStub());
const result = await repo.createIfNotExists({ username: 'Alice' }, /* blockNum */ 50, true);
expect(result).toBe(existing);
expect(existing.updateFromBlockchain).not.toHaveBeenCalled();
expect(repoStub.save).not.toHaveBeenCalled();
});
it('обновляет запись более свежей дельтой (block_num > N)', async () => {
const existing = makeDomain({ username: 'alice', block_num: 100 });
const repoStub = makeRepoStub(existing);
const repo = new TestRepo(repoStub, makeVersioningStub());
await repo.createIfNotExists({ username: 'Alice' }, 150, true);
expect(existing.updateFromBlockchain).toHaveBeenCalledWith({ username: 'Alice' }, 150, true);
expect(repoStub.save).toHaveBeenCalled();
});
it('обновляет при равном block_num (идемпотентный повтор, не stale)', async () => {
const existing = makeDomain({ username: 'alice', block_num: 100 });
const repoStub = makeRepoStub(existing);
const repo = new TestRepo(repoStub, makeVersioningStub());
await repo.createIfNotExists({ username: 'Alice' }, 100, true);
expect(existing.updateFromBlockchain).toHaveBeenCalled();
});
it('создаёт новую сущность, если записи ещё нет', async () => {
const repoStub = makeRepoStub(null);
const repo = new TestRepo(repoStub, makeVersioningStub());
const result = await repo.createIfNotExists({ username: 'Bob' }, 10, true);
expect(result).toBeTruthy();
expect(repoStub.save).toHaveBeenCalled();
});
it('сравнивает block_num численно, когда из PG он пришёл строкой (bigint-as-string)', async () => {
const existing = makeDomain({ username: 'alice', block_num: '100' as any });
const repoStub = makeRepoStub(existing);
const repo = new TestRepo(repoStub, makeVersioningStub());
// '100' (строка) vs 50 (число): без Number() сравнение строк дало бы '100' < 50 == false по-разному;
// guard обязан трактовать 50 < 100 как stale.
const result = await repo.createIfNotExists({ username: 'Alice' }, 50, true);
expect(result).toBe(existing);
expect(existing.updateFromBlockchain).not.toHaveBeenCalled();
});
});
+274 -86
View File
@@ -271,6 +271,9 @@ importers:
'@coopenomics/notifications':
specifier: workspace:*
version: link:../notifications
'@coopenomics/parser2':
specifier: ^1.2.0
version: 1.2.0
'@coopenomics/provider-client':
specifier: 2025.11.12-alpha-1
version: 2025.11.12-alpha-1
@@ -2720,6 +2723,13 @@ packages:
resolution: {integrity: sha512-Ir+AOibqzrIsL6ajt3Rz3LskB7OiMVHqltZmspbW/TJuTVuyOMirVqAkjfY6JISiLHgyNqicAC8AyHHGzNd/dA==}
engines: {node: '>=0.1.90'}
'@coopenomics/coopos-ship-reader@0.3.1':
resolution: {integrity: sha512-qaCdEy72do/mje/YNCiZLvGHzcVbwIXjesKLJdm9Q0rGFxFSBwfbHv45G2hjUmc36NAhzh2a+blULRk5m/pSgA==}
'@coopenomics/parser2@1.2.0':
resolution: {integrity: sha512-UgxMQfPxYKmouKAGEuimkBW7bSyuRIsqL7xUwT5CLx5OWtmGg/1pMSmBV618oEDRidJRHX19IPTsNF+EN3xD0Q==}
hasBin: true
'@coopenomics/provider-client@2025.11.12-alpha-1':
resolution: {integrity: sha512-Z2gDb/iDZhIKSZ7NDEAQ2BdqJdPQs04gI/irwjyr22imlvB84ARIG8bAq7ZkBAeCVE2E7txXa6N+OdqNnafEhQ==}
@@ -4449,105 +4459,89 @@ packages:
resolution: {integrity: sha512-excjX8DfsIcJ10x1Kzr4RcWe1edC9PquDRRPx3YVCvQv+U5p7Yin2s32ftzikXojb1PIFc/9Mt28/y+iRklkrw==}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@img/sharp-libvips-linux-arm@1.2.4':
resolution: {integrity: sha512-bFI7xcKFELdiNCVov8e44Ia4u2byA+l3XtsAj+Q8tfCwO6BQ8iDojYdvoPMqsKDkuoOo+X6HZA0s0q11ANMQ8A==}
cpu: [arm]
os: [linux]
libc: [glibc]
'@img/sharp-libvips-linux-ppc64@1.2.4':
resolution: {integrity: sha512-FMuvGijLDYG6lW+b/UvyilUWu5Ayu+3r2d1S8notiGCIyYU/76eig1UfMmkZ7vwgOrzKzlQbFSuQfgm7GYUPpA==}
cpu: [ppc64]
os: [linux]
libc: [glibc]
'@img/sharp-libvips-linux-riscv64@1.2.4':
resolution: {integrity: sha512-oVDbcR4zUC0ce82teubSm+x6ETixtKZBh/qbREIOcI3cULzDyb18Sr/Wcyx7NRQeQzOiHTNbZFF1UwPS2scyGA==}
cpu: [riscv64]
os: [linux]
libc: [glibc]
'@img/sharp-libvips-linux-s390x@1.2.4':
resolution: {integrity: sha512-qmp9VrzgPgMoGZyPvrQHqk02uyjA0/QrTO26Tqk6l4ZV0MPWIW6LTkqOIov+J1yEu7MbFQaDpwdwJKhbJvuRxQ==}
cpu: [s390x]
os: [linux]
libc: [glibc]
'@img/sharp-libvips-linux-x64@1.2.4':
resolution: {integrity: sha512-tJxiiLsmHc9Ax1bz3oaOYBURTXGIRDODBqhveVHonrHJ9/+k89qbLl0bcJns+e4t4rvaNBxaEZsFtSfAdquPrw==}
cpu: [x64]
os: [linux]
libc: [glibc]
'@img/sharp-libvips-linuxmusl-arm64@1.2.4':
resolution: {integrity: sha512-FVQHuwx1IIuNow9QAbYUzJ+En8KcVm9Lk5+uGUQJHaZmMECZmOlix9HnH7n1TRkXMS0pGxIJokIVB9SuqZGGXw==}
cpu: [arm64]
os: [linux]
libc: [musl]
'@img/sharp-libvips-linuxmusl-x64@1.2.4':
resolution: {integrity: sha512-+LpyBk7L44ZIXwz/VYfglaX/okxezESc6UxDSoyo2Ks6Jxc4Y7sGjpgU9s4PMgqgjj1gZCylTieNamqA1MF7Dg==}
cpu: [x64]
os: [linux]
libc: [musl]
'@img/sharp-linux-arm64@0.34.5':
resolution: {integrity: sha512-bKQzaJRY/bkPOXyKx5EVup7qkaojECG6NLYswgktOZjaXecSAeCWiZwwiFf3/Y+O1HrauiE3FVsGxFg8c24rZg==}
engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@img/sharp-linux-arm@0.34.5':
resolution: {integrity: sha512-9dLqsvwtg1uuXBGZKsxem9595+ujv0sJ6Vi8wcTANSFpwV/GONat5eCkzQo/1O6zRIkh0m/8+5BjrRr7jDUSZw==}
engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0}
cpu: [arm]
os: [linux]
libc: [glibc]
'@img/sharp-linux-ppc64@0.34.5':
resolution: {integrity: sha512-7zznwNaqW6YtsfrGGDA6BRkISKAAE1Jo0QdpNYXNMHu2+0dTrPflTLNkpc8l7MUP5M16ZJcUvysVWWrMefZquA==}
engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0}
cpu: [ppc64]
os: [linux]
libc: [glibc]
'@img/sharp-linux-riscv64@0.34.5':
resolution: {integrity: sha512-51gJuLPTKa7piYPaVs8GmByo7/U7/7TZOq+cnXJIHZKavIRHAP77e3N2HEl3dgiqdD/w0yUfiJnII77PuDDFdw==}
engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0}
cpu: [riscv64]
os: [linux]
libc: [glibc]
'@img/sharp-linux-s390x@0.34.5':
resolution: {integrity: sha512-nQtCk0PdKfho3eC5MrbQoigJ2gd1CgddUMkabUj+rBevs8tZ2cULOx46E7oyX+04WGfABgIwmMC0VqieTiR4jg==}
engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0}
cpu: [s390x]
os: [linux]
libc: [glibc]
'@img/sharp-linux-x64@0.34.5':
resolution: {integrity: sha512-MEzd8HPKxVxVenwAa+JRPwEC7QFjoPWuS5NZnBt6B3pu7EG2Ge0id1oLHZpPJdn3OQK+BQDiw9zStiHBTJQQQQ==}
engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0}
cpu: [x64]
os: [linux]
libc: [glibc]
'@img/sharp-linuxmusl-arm64@0.34.5':
resolution: {integrity: sha512-fprJR6GtRsMt6Kyfq44IsChVZeGN97gTD331weR1ex1c1rypDEABN6Tm2xa1wE6lYb5DdEnk03NZPqA7Id21yg==}
engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0}
cpu: [arm64]
os: [linux]
libc: [musl]
'@img/sharp-linuxmusl-x64@0.34.5':
resolution: {integrity: sha512-Jg8wNT1MUzIvhBFxViqrEhWDGzqymo3sV7z7ZsaWbZNDLXRJZoRGrjulp60YYtV4wfY8VIKcWidjojlLcWrd8Q==}
engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0}
cpu: [x64]
os: [linux]
libc: [musl]
'@img/sharp-wasm32@0.34.5':
resolution: {integrity: sha512-OdWTEiVkY2PHwqkbBI8frFxQQFekHaSSkUIJkwzclWZe64O1X4UlUjqqqLaPbUpMOQk6FBu/HtlGXNblIs0huw==}
@@ -5077,14 +5071,12 @@ packages:
engines: {node: '>= 10'}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@livekit/rtc-node-linux-x64-gnu@0.13.24':
resolution: {integrity: sha512-vKOxzN/SsrtV8zIVwZCi31bZUhlb6RhJZ0NnY5MwKGSRFPi7Dwt8fmr0Vh0YmsY/p+4eZjKxvFmy7L3WVE54zw==}
engines: {node: '>= 10'}
cpu: [x64]
os: [linux]
libc: [glibc]
'@livekit/rtc-node-win32-x64-msvc@0.13.24':
resolution: {integrity: sha512-yTzqwndq2oKLUkXW2i/BkZMJC6kZOpRO/DKvkkKQvqc3Q+JuWz1m48GmyjIwTOKF28QjqEU3+IrnD65Uu+mFOg==}
@@ -5231,6 +5223,112 @@ packages:
cpu: [x64]
os: [win32]
'@napi-rs/nice-android-arm-eabi@1.1.1':
resolution: {integrity: sha512-kjirL3N6TnRPv5iuHw36wnucNqXAO46dzK9oPb0wj076R5Xm8PfUVA9nAFB5ZNMmfJQJVKACAPd/Z2KYMppthw==}
engines: {node: '>= 10'}
cpu: [arm]
os: [android]
'@napi-rs/nice-android-arm64@1.1.1':
resolution: {integrity: sha512-blG0i7dXgbInN5urONoUCNf+DUEAavRffrO7fZSeoRMJc5qD+BJeNcpr54msPF6qfDD6kzs9AQJogZvT2KD5nw==}
engines: {node: '>= 10'}
cpu: [arm64]
os: [android]
'@napi-rs/nice-darwin-arm64@1.1.1':
resolution: {integrity: sha512-s/E7w45NaLqTGuOjC2p96pct4jRfo61xb9bU1unM/MJ/RFkKlJyJDx7OJI/O0ll/hrfpqKopuAFDV8yo0hfT7A==}
engines: {node: '>= 10'}
cpu: [arm64]
os: [darwin]
'@napi-rs/nice-darwin-x64@1.1.1':
resolution: {integrity: sha512-dGoEBnVpsdcC+oHHmW1LRK5eiyzLwdgNQq3BmZIav+9/5WTZwBYX7r5ZkQC07Nxd3KHOCkgbHSh4wPkH1N1LiQ==}
engines: {node: '>= 10'}
cpu: [x64]
os: [darwin]
'@napi-rs/nice-freebsd-x64@1.1.1':
resolution: {integrity: sha512-kHv4kEHAylMYmlNwcQcDtXjklYp4FCf0b05E+0h6nDHsZ+F0bDe04U/tXNOqrx5CmIAth4vwfkjjUmp4c4JktQ==}
engines: {node: '>= 10'}
cpu: [x64]
os: [freebsd]
'@napi-rs/nice-linux-arm-gnueabihf@1.1.1':
resolution: {integrity: sha512-E1t7K0efyKXZDoZg1LzCOLxgolxV58HCkaEkEvIYQx12ht2pa8hoBo+4OB3qh7e+QiBlp1SRf+voWUZFxyhyqg==}
engines: {node: '>= 10'}
cpu: [arm]
os: [linux]
'@napi-rs/nice-linux-arm64-gnu@1.1.1':
resolution: {integrity: sha512-CIKLA12DTIZlmTaaKhQP88R3Xao+gyJxNWEn04wZwC2wmRapNnxCUZkVwggInMJvtVElA+D4ZzOU5sX4jV+SmQ==}
engines: {node: '>= 10'}
cpu: [arm64]
os: [linux]
'@napi-rs/nice-linux-arm64-musl@1.1.1':
resolution: {integrity: sha512-+2Rzdb3nTIYZ0YJF43qf2twhqOCkiSrHx2Pg6DJaCPYhhaxbLcdlV8hCRMHghQ+EtZQWGNcS2xF4KxBhSGeutg==}
engines: {node: '>= 10'}
cpu: [arm64]
os: [linux]
'@napi-rs/nice-linux-ppc64-gnu@1.1.1':
resolution: {integrity: sha512-4FS8oc0GeHpwvv4tKciKkw3Y4jKsL7FRhaOeiPei0X9T4Jd619wHNe4xCLmN2EMgZoeGg+Q7GY7BsvwKpL22Tg==}
engines: {node: '>= 10'}
cpu: [ppc64]
os: [linux]
'@napi-rs/nice-linux-riscv64-gnu@1.1.1':
resolution: {integrity: sha512-HU0nw9uD4FO/oGCCk409tCi5IzIZpH2agE6nN4fqpwVlCn5BOq0MS1dXGjXaG17JaAvrlpV5ZeyZwSon10XOXw==}
engines: {node: '>= 10'}
cpu: [riscv64]
os: [linux]
'@napi-rs/nice-linux-s390x-gnu@1.1.1':
resolution: {integrity: sha512-2YqKJWWl24EwrX0DzCQgPLKQBxYDdBxOHot1KWEq7aY2uYeX+Uvtv4I8xFVVygJDgf6/92h9N3Y43WPx8+PAgQ==}
engines: {node: '>= 10'}
cpu: [s390x]
os: [linux]
'@napi-rs/nice-linux-x64-gnu@1.1.1':
resolution: {integrity: sha512-/gaNz3R92t+dcrfCw/96pDopcmec7oCcAQ3l/M+Zxr82KT4DljD37CpgrnXV+pJC263JkW572pdbP3hP+KjcIg==}
engines: {node: '>= 10'}
cpu: [x64]
os: [linux]
'@napi-rs/nice-linux-x64-musl@1.1.1':
resolution: {integrity: sha512-xScCGnyj/oppsNPMnevsBe3pvNaoK7FGvMjT35riz9YdhB2WtTG47ZlbxtOLpjeO9SqqQ2J2igCmz6IJOD5JYw==}
engines: {node: '>= 10'}
cpu: [x64]
os: [linux]
'@napi-rs/nice-openharmony-arm64@1.1.1':
resolution: {integrity: sha512-6uJPRVwVCLDeoOaNyeiW0gp2kFIM4r7PL2MczdZQHkFi9gVlgm+Vn+V6nTWRcu856mJ2WjYJiumEajfSm7arPQ==}
engines: {node: '>= 10'}
cpu: [arm64]
os: [openharmony]
'@napi-rs/nice-win32-arm64-msvc@1.1.1':
resolution: {integrity: sha512-uoTb4eAvM5B2aj/z8j+Nv8OttPf2m+HVx3UjA5jcFxASvNhQriyCQF1OB1lHL43ZhW+VwZlgvjmP5qF3+59atA==}
engines: {node: '>= 10'}
cpu: [arm64]
os: [win32]
'@napi-rs/nice-win32-ia32-msvc@1.1.1':
resolution: {integrity: sha512-CNQqlQT9MwuCsg1Vd/oKXiuH+TcsSPJmlAFc5frFyX/KkOh0UpBLEj7aoY656d5UKZQMQFP7vJNa1DNUNORvug==}
engines: {node: '>= 10'}
cpu: [ia32]
os: [win32]
'@napi-rs/nice-win32-x64-msvc@1.1.1':
resolution: {integrity: sha512-vB+4G/jBQCAh0jelMTY3+kgFy00Hlx2f2/1zjMoH821IbplbWZOkLiTYXQkygNTzQJTq5cvwBDgn2ppHD+bglQ==}
engines: {node: '>= 10'}
cpu: [x64]
os: [win32]
'@napi-rs/nice@1.1.1':
resolution: {integrity: sha512-xJIPs+bYuc9ASBl+cvGsKbGrJmS6fAKaSZCnT0lhahT5rhA2VVy9/EcIgd2JhtEuFOJNx7UHNn/qiTPTY4nrQw==}
engines: {node: '>= 10'}
'@napi-rs/wasm-runtime@0.2.12':
resolution: {integrity: sha512-ZVWUcfwY4E/yPitQJl481FjFo3K22D6qF0DuFH6Y/nbnE11GY5uguDxZMGXPQ8WQ0128MXQD7TnfHyK4oWoIJQ==}
@@ -5552,28 +5650,24 @@ packages:
engines: {node: '>= 10'}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@nx/nx-linux-arm64-musl@20.8.4':
resolution: {integrity: sha512-AlZZFolS/S0FahRKG7rJ0Z9CgmIkyzHgGaoy3qNEMDEjFhR3jt2ZZSLp90W7zjgrxojOo90ajNMrg2UmtcQRDA==}
engines: {node: '>= 10'}
cpu: [arm64]
os: [linux]
libc: [musl]
'@nx/nx-linux-x64-gnu@20.8.4':
resolution: {integrity: sha512-MSu+xVNdR95tuuO+eL/a/ZeMlhfrZ627On5xaCZXnJ+lFxNg/S4nlKZQk0Eq5hYALCd/GKgFGasRdlRdOtvGPg==}
engines: {node: '>= 10'}
cpu: [x64]
os: [linux]
libc: [glibc]
'@nx/nx-linux-x64-musl@20.8.4':
resolution: {integrity: sha512-KxpQpyLCgIIHWZ4iRSUN9ohCwn1ZSDASbuFCdG3mohryzCy8WrPkuPcb+68J3wuQhmA5w//Xpp/dL0hHoit9zQ==}
engines: {node: '>= 10'}
cpu: [x64]
os: [linux]
libc: [musl]
'@nx/nx-win32-arm64-msvc@20.8.4':
resolution: {integrity: sha512-ffLBrxM9ibk+eWSY995kiFFRTSRb9HkD5T1s/uZyxV6jfxYPaZDBAWAETDneyBXps7WtaOMu+kVZlXQ3X+TfIA==}
@@ -6074,49 +6168,41 @@ packages:
resolution: {integrity: sha512-heV2+jmXyYnUrpUXSPugqWDRpnsQcDm2AX4wzTuvgdlZfoNYO0O3W2AVpJYaDn9AG4JdM6Kxom8+foE7/BcSig==}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@oxc-resolver/binding-linux-arm64-musl@11.19.1':
resolution: {integrity: sha512-jvo2Pjs1c9KPxMuMPIeQsgu0mOJF9rEb3y3TdpsrqwxRM+AN6/nDDwv45n5ZrUnQMsdBy5gIabioMKnQfWo9ew==}
cpu: [arm64]
os: [linux]
libc: [musl]
'@oxc-resolver/binding-linux-ppc64-gnu@11.19.1':
resolution: {integrity: sha512-vLmdNxWCdN7Uo5suays6A/+ywBby2PWBBPXctWPg5V0+eVuzsJxgAn6MMB4mPlshskYbppjpN2Zg83ArHze9gQ==}
cpu: [ppc64]
os: [linux]
libc: [glibc]
'@oxc-resolver/binding-linux-riscv64-gnu@11.19.1':
resolution: {integrity: sha512-/b+WgR+VTSBxzgOhDO7TlMXC1ufPIMR6Vj1zN+/x+MnyXGW7prTLzU9eW85Aj7Th7CCEG9ArCbTeqxCzFWdg2w==}
cpu: [riscv64]
os: [linux]
libc: [glibc]
'@oxc-resolver/binding-linux-riscv64-musl@11.19.1':
resolution: {integrity: sha512-YlRdeWb9j42p29ROh+h4eg/OQ3dTJlpHSa+84pUM9+p6i3djtPz1q55yLJhgW9XfDch7FN1pQ/Vd6YP+xfRIuw==}
cpu: [riscv64]
os: [linux]
libc: [musl]
'@oxc-resolver/binding-linux-s390x-gnu@11.19.1':
resolution: {integrity: sha512-EDpafVOQWF8/MJynsjOGFThcqhRHy417sRyLfQmeiamJ8qVhSKAn2Dn2VVKUGCjVB9C46VGjhNo7nOPUi1x6uA==}
cpu: [s390x]
os: [linux]
libc: [glibc]
'@oxc-resolver/binding-linux-x64-gnu@11.19.1':
resolution: {integrity: sha512-NxjZe+rqWhr+RT8/Ik+5ptA3oz7tUw361Wa5RWQXKnfqwSSHdHyrw6IdcTfYuml9dM856AlKWZIUXDmA9kkiBQ==}
cpu: [x64]
os: [linux]
libc: [glibc]
'@oxc-resolver/binding-linux-x64-musl@11.19.1':
resolution: {integrity: sha512-cM/hQwsO3ReJg5kR+SpI69DMfvNCp+A/eVR4b4YClE5bVZwz8rh2Nh05InhwI5HR/9cArbEkzMjcKgTHS6UaNw==}
cpu: [x64]
os: [linux]
libc: [musl]
'@oxc-resolver/binding-openharmony-arm64@11.19.1':
resolution: {integrity: sha512-QF080IowFB0+9Rh6RcD19bdgh49BpQHUW5TajG1qvWHvmrQznTZZjYlgE2ltLXyKY+qs4F/v5xuX1XS7Is+3qA==}
@@ -6178,42 +6264,36 @@ packages:
engines: {node: '>= 10.0.0'}
cpu: [arm]
os: [linux]
libc: [glibc]
'@parcel/watcher-linux-arm-musl@2.5.6':
resolution: {integrity: sha512-Ve3gUCG57nuUUSyjBq/MAM0CzArtuIOxsBdQ+ftz6ho8n7s1i9E1Nmk/xmP323r2YL0SONs1EuwqBp2u1k5fxg==}
engines: {node: '>= 10.0.0'}
cpu: [arm]
os: [linux]
libc: [musl]
'@parcel/watcher-linux-arm64-glibc@2.5.6':
resolution: {integrity: sha512-f2g/DT3NhGPdBmMWYoxixqYr3v/UXcmLOYy16Bx0TM20Tchduwr4EaCbmxh1321TABqPGDpS8D/ggOTaljijOA==}
engines: {node: '>= 10.0.0'}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@parcel/watcher-linux-arm64-musl@2.5.6':
resolution: {integrity: sha512-qb6naMDGlbCwdhLj6hgoVKJl2odL34z2sqkC7Z6kzir8b5W65WYDpLB6R06KabvZdgoHI/zxke4b3zR0wAbDTA==}
engines: {node: '>= 10.0.0'}
cpu: [arm64]
os: [linux]
libc: [musl]
'@parcel/watcher-linux-x64-glibc@2.5.6':
resolution: {integrity: sha512-kbT5wvNQlx7NaGjzPFu8nVIW1rWqV780O7ZtkjuWaPUgpv2NMFpjYERVi0UYj1msZNyCzGlaCWEtzc+exjMGbQ==}
engines: {node: '>= 10.0.0'}
cpu: [x64]
os: [linux]
libc: [glibc]
'@parcel/watcher-linux-x64-musl@2.5.6':
resolution: {integrity: sha512-1JRFeC+h7RdXwldHzTsmdtYR/Ku8SylLgTU/reMuqdVD7CtLwf0VR1FqeprZ0eHQkO0vqsbvFLXUmYm/uNKJBg==}
engines: {node: '>= 10.0.0'}
cpu: [x64]
os: [linux]
libc: [musl]
'@parcel/watcher-win32-arm64@2.5.6':
resolution: {integrity: sha512-3ukyebjc6eGlw9yRt678DxVF7rjXatWiHvTXqphZLvo7aC5NdEgFufVwjFfY51ijYEWpXbqF5jtrK275z52D4Q==}
@@ -6462,42 +6542,36 @@ packages:
engines: {node: ^20.19.0 || >=22.12.0}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@rolldown/binding-linux-arm64-musl@1.0.0-rc.12':
resolution: {integrity: sha512-V6/wZztnBqlx5hJQqNWwFdxIKN0m38p8Jas+VoSfgH54HSj9tKTt1dZvG6JRHcjh6D7TvrJPWFGaY9UBVOaWPw==}
engines: {node: ^20.19.0 || >=22.12.0}
cpu: [arm64]
os: [linux]
libc: [musl]
'@rolldown/binding-linux-ppc64-gnu@1.0.0-rc.12':
resolution: {integrity: sha512-AP3E9BpcUYliZCxa3w5Kwj9OtEVDYK6sVoUzy4vTOJsjPOgdaJZKFmN4oOlX0Wp0RPV2ETfmIra9x1xuayFB7g==}
engines: {node: ^20.19.0 || >=22.12.0}
cpu: [ppc64]
os: [linux]
libc: [glibc]
'@rolldown/binding-linux-s390x-gnu@1.0.0-rc.12':
resolution: {integrity: sha512-nWwpvUSPkoFmZo0kQazZYOrT7J5DGOJ/+QHHzjvNlooDZED8oH82Yg67HvehPPLAg5fUff7TfWFHQS8IV1n3og==}
engines: {node: ^20.19.0 || >=22.12.0}
cpu: [s390x]
os: [linux]
libc: [glibc]
'@rolldown/binding-linux-x64-gnu@1.0.0-rc.12':
resolution: {integrity: sha512-RNrafz5bcwRy+O9e6P8Z/OCAJW/A+qtBczIqVYwTs14pf4iV1/+eKEjdOUta93q2TsT/FI0XYDP3TCky38LMAg==}
engines: {node: ^20.19.0 || >=22.12.0}
cpu: [x64]
os: [linux]
libc: [glibc]
'@rolldown/binding-linux-x64-musl@1.0.0-rc.12':
resolution: {integrity: sha512-Jpw/0iwoKWx3LJ2rc1yjFrj+T7iHZn2JDg1Yny1ma0luviFS4mhAIcd1LFNxK3EYu3DHWCps0ydXQ5i/rrJ2ig==}
engines: {node: ^20.19.0 || >=22.12.0}
cpu: [x64]
os: [linux]
libc: [musl]
'@rolldown/binding-openharmony-arm64@1.0.0-rc.12':
resolution: {integrity: sha512-vRugONE4yMfVn0+7lUKdKvN4D5YusEiPilaoO2sgUWpCvrncvWgPMzK00ZFFJuiPgLwgFNP5eSiUlv2tfc+lpA==}
@@ -6644,79 +6718,66 @@ packages:
resolution: {integrity: sha512-RzeBwv0B3qtVBWtcuABtSuCzToo2IEAIQrcyB/b2zMvBWVbjo8bZDjACUpnaafaxhTw2W+imQbP2BD1usasK4g==}
cpu: [arm]
os: [linux]
libc: [glibc]
'@rollup/rollup-linux-arm-musleabihf@4.60.0':
resolution: {integrity: sha512-Sf7zusNI2CIU1HLzuu9Tc5YGAHEZs5Lu7N1ssJG4Tkw6e0MEsN7NdjUDDfGNHy2IU+ENyWT+L2obgWiguWibWQ==}
cpu: [arm]
os: [linux]
libc: [musl]
'@rollup/rollup-linux-arm64-gnu@4.60.0':
resolution: {integrity: sha512-DX2x7CMcrJzsE91q7/O02IJQ5/aLkVtYFryqCjduJhUfGKG6yJV8hxaw8pZa93lLEpPTP/ohdN4wFz7yp/ry9A==}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@rollup/rollup-linux-arm64-musl@4.60.0':
resolution: {integrity: sha512-09EL+yFVbJZlhcQfShpswwRZ0Rg+z/CsSELFCnPt3iK+iqwGsI4zht3secj5vLEs957QvFFXnzAT0FFPIxSrkQ==}
cpu: [arm64]
os: [linux]
libc: [musl]
'@rollup/rollup-linux-loong64-gnu@4.60.0':
resolution: {integrity: sha512-i9IcCMPr3EXm8EQg5jnja0Zyc1iFxJjZWlb4wr7U2Wx/GrddOuEafxRdMPRYVaXjgbhvqalp6np07hN1w9kAKw==}
cpu: [loong64]
os: [linux]
libc: [glibc]
'@rollup/rollup-linux-loong64-musl@4.60.0':
resolution: {integrity: sha512-DGzdJK9kyJ+B78MCkWeGnpXJ91tK/iKA6HwHxF4TAlPIY7GXEvMe8hBFRgdrR9Ly4qebR/7gfUs9y2IoaVEyog==}
cpu: [loong64]
os: [linux]
libc: [musl]
'@rollup/rollup-linux-ppc64-gnu@4.60.0':
resolution: {integrity: sha512-RwpnLsqC8qbS8z1H1AxBA1H6qknR4YpPR9w2XX0vo2Sz10miu57PkNcnHVaZkbqyw/kUWfKMI73jhmfi9BRMUQ==}
cpu: [ppc64]
os: [linux]
libc: [glibc]
'@rollup/rollup-linux-ppc64-musl@4.60.0':
resolution: {integrity: sha512-Z8pPf54Ly3aqtdWC3G4rFigZgNvd+qJlOE52fmko3KST9SoGfAdSRCwyoyG05q1HrrAblLbk1/PSIV+80/pxLg==}
cpu: [ppc64]
os: [linux]
libc: [musl]
'@rollup/rollup-linux-riscv64-gnu@4.60.0':
resolution: {integrity: sha512-3a3qQustp3COCGvnP4SvrMHnPQ9d1vzCakQVRTliaz8cIp/wULGjiGpbcqrkv0WrHTEp8bQD/B3HBjzujVWLOA==}
cpu: [riscv64]
os: [linux]
libc: [glibc]
'@rollup/rollup-linux-riscv64-musl@4.60.0':
resolution: {integrity: sha512-pjZDsVH/1VsghMJ2/kAaxt6dL0psT6ZexQVrijczOf+PeP2BUqTHYejk3l6TlPRydggINOeNRhvpLa0AYpCWSQ==}
cpu: [riscv64]
os: [linux]
libc: [musl]
'@rollup/rollup-linux-s390x-gnu@4.60.0':
resolution: {integrity: sha512-3ObQs0BhvPgiUVZrN7gqCSvmFuMWvWvsjG5ayJ3Lraqv+2KhOsp+pUbigqbeWqueGIsnn+09HBw27rJ+gYK4VQ==}
cpu: [s390x]
os: [linux]
libc: [glibc]
'@rollup/rollup-linux-x64-gnu@4.60.0':
resolution: {integrity: sha512-EtylprDtQPdS5rXvAayrNDYoJhIz1/vzN2fEubo3yLE7tfAw+948dO0g4M0vkTVFhKojnF+n6C8bDNe+gDRdTg==}
cpu: [x64]
os: [linux]
libc: [glibc]
'@rollup/rollup-linux-x64-musl@4.60.0':
resolution: {integrity: sha512-k09oiRCi/bHU9UVFqD17r3eJR9bn03TyKraCrlz5ULFJGdJGi7VOmm9jl44vOJvRJ6P7WuBi/s2A97LxxHGIdw==}
cpu: [x64]
os: [linux]
libc: [musl]
'@rollup/rollup-openbsd-x64@4.60.0':
resolution: {integrity: sha512-1o/0/pIhozoSaDJoDcec+IVLbnRtQmHwPV730+AOD29lHEEo4F5BEUB24H0OBdhbBBDwIOSuf7vgg0Ywxdfiiw==}
@@ -7084,42 +7145,36 @@ packages:
engines: {node: '>=10'}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@swc/core-linux-arm64-musl@1.15.33':
resolution: {integrity: sha512-il7tYM+CpUNzieQbwAjFT1P8zqAhmGWNAGhQZBnxurXZ0aNn+5nqYFTEUKNZl7QibtT0uQXzTZrNGHCIj6Y1Og==}
engines: {node: '>=10'}
cpu: [arm64]
os: [linux]
libc: [musl]
'@swc/core-linux-ppc64-gnu@1.15.33':
resolution: {integrity: sha512-ZtNBwN0Z7CFj9Il0FcPaKdjgP7URyKu/3RfH46vq+0paOBqLj4NYldD6Qo//Duif/7IOtAraUfDOmp0PLAufog==}
engines: {node: '>=10'}
cpu: [ppc64]
os: [linux]
libc: [glibc]
'@swc/core-linux-s390x-gnu@1.15.33':
resolution: {integrity: sha512-De1IyajoOmhOYYjw/lx66bKlyDpHZTueqwpDrWgf5O7T6d1ODeJJO9/OqMBmrBQc5C+dNnlmIufHsp4QVCWufA==}
engines: {node: '>=10'}
cpu: [s390x]
os: [linux]
libc: [glibc]
'@swc/core-linux-x64-gnu@1.15.33':
resolution: {integrity: sha512-mGTH0YxmUN+x6vRN/I6NOk5X0ogNktkwPnJ94IMvR7QjhRDwL0O8RXEDhyUM0YtwWrryBOqaJQBX4zruxEPRGw==}
engines: {node: '>=10'}
cpu: [x64]
os: [linux]
libc: [glibc]
'@swc/core-linux-x64-musl@1.15.33':
resolution: {integrity: sha512-hj628ZkSEJf6zMf5VMbYrG2O6QqyTIp2qwY6VlCjvIa9lAEZ5c2lfPblCLVGYubTeLJDxadLB/CxqQYOQABeEQ==}
engines: {node: '>=10'}
cpu: [x64]
os: [linux]
libc: [musl]
'@swc/core-win32-arm64-msvc@1.15.33':
resolution: {integrity: sha512-GV2oohtN2/5+KSccl86VULu3aT+LrISC8uzgSq0FRnikpD+Zwc+sBlXmoKQ+Db6jI57ITUOIB8jRkdGMABC29g==}
@@ -7931,49 +7986,41 @@ packages:
resolution: {integrity: sha512-34gw7PjDGB9JgePJEmhEqBhWvCiiWCuXsL9hYphDF7crW7UgI05gyBAi6MF58uGcMOiOqSJ2ybEeCvHcq0BCmQ==}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@unrs/resolver-binding-linux-arm64-musl@1.11.1':
resolution: {integrity: sha512-RyMIx6Uf53hhOtJDIamSbTskA99sPHS96wxVE/bJtePJJtpdKGXO1wY90oRdXuYOGOTuqjT8ACccMc4K6QmT3w==}
cpu: [arm64]
os: [linux]
libc: [musl]
'@unrs/resolver-binding-linux-ppc64-gnu@1.11.1':
resolution: {integrity: sha512-D8Vae74A4/a+mZH0FbOkFJL9DSK2R6TFPC9M+jCWYia/q2einCubX10pecpDiTmkJVUH+y8K3BZClycD8nCShA==}
cpu: [ppc64]
os: [linux]
libc: [glibc]
'@unrs/resolver-binding-linux-riscv64-gnu@1.11.1':
resolution: {integrity: sha512-frxL4OrzOWVVsOc96+V3aqTIQl1O2TjgExV4EKgRY09AJ9leZpEg8Ak9phadbuX0BA4k8U5qtvMSQQGGmaJqcQ==}
cpu: [riscv64]
os: [linux]
libc: [glibc]
'@unrs/resolver-binding-linux-riscv64-musl@1.11.1':
resolution: {integrity: sha512-mJ5vuDaIZ+l/acv01sHoXfpnyrNKOk/3aDoEdLO/Xtn9HuZlDD6jKxHlkN8ZhWyLJsRBxfv9GYM2utQ1SChKew==}
cpu: [riscv64]
os: [linux]
libc: [musl]
'@unrs/resolver-binding-linux-s390x-gnu@1.11.1':
resolution: {integrity: sha512-kELo8ebBVtb9sA7rMe1Cph4QHreByhaZ2QEADd9NzIQsYNQpt9UkM9iqr2lhGr5afh885d/cB5QeTXSbZHTYPg==}
cpu: [s390x]
os: [linux]
libc: [glibc]
'@unrs/resolver-binding-linux-x64-gnu@1.11.1':
resolution: {integrity: sha512-C3ZAHugKgovV5YvAMsxhq0gtXuwESUKc5MhEtjBpLoHPLYM+iuwSj3lflFwK3DPm68660rZ7G8BMcwSro7hD5w==}
cpu: [x64]
os: [linux]
libc: [glibc]
'@unrs/resolver-binding-linux-x64-musl@1.11.1':
resolution: {integrity: sha512-rV0YSoyhK2nZ4vEswT/QwqzqQXw5I6CjoaYMOX0TqBlWhojUf8P94mvI7nuJTeaCkkds3QE4+zS8Ko+GdXuZtA==}
cpu: [x64]
os: [linux]
libc: [musl]
'@unrs/resolver-binding-wasm32-wasi@1.11.1':
resolution: {integrity: sha512-5u4RkfxJm+Ng7IWgkzi3qrFOvLvQYnPBmjmZQ8+szTK/b31fQCnleNl1GgEt7nIsZRIf5PLhPwT0WM+q45x/UQ==}
@@ -8227,6 +8274,9 @@ packages:
'@wharfkit/antelope@1.1.1':
resolution: {integrity: sha512-zOetZJG0T4/WRCIC8Ui2DGv+RxrSkdXOq/q8sUqUcuKei+aJo4/9Xk37VULPt2qe0NpTm93ajYV9cYdD9m5Fgg==}
'@wharfkit/antelope@1.2.0':
resolution: {integrity: sha512-9q0nvM8yUtjKTQlukKZODAhUN2S2/cfSlIYdh2mPnaOCSH8KOLJ2gYCPuQVvH1FE9AKu312lM5TzhkRccf1VjQ==}
'@wharfkit/common@1.5.0':
resolution: {integrity: sha512-eqXkOy+vshcEzK8kED+EsoTPJjlBKHYglgV9CBnZQgIlGrWIRXWH4YaXH3W7EbI/nCRJCaNqxm5fC+pgpFcp8g==}
peerDependencies:
@@ -8935,6 +8985,9 @@ packages:
bindings@1.5.0:
resolution: {integrity: sha512-p2q/t/mhvuOj/UeLlV6566GD/guowlr0hHxClI0W9m7MWYkL1F0hLo+0Aexs9HSPCtR1SXQ0TD3MMKrXZajbiQ==}
bintrees@1.0.2:
resolution: {integrity: sha512-VOMgTMwjAaUG580SXn3LacVgjurrbMme7ZZNYGSSV7mmtY6QQRh0Eg3pwIcntQ77DErK1L0NxkbetjcoXzVwKw==}
birpc@2.9.0:
resolution: {integrity: sha512-KrayHS5pBi69Xi9JmvoqrIgYGDkD6mcSe/i6YKi3w5kekCLzrX4+nawcXqrj2tIp50Kw/mT/s3p+GVK0A0sKxw==}
@@ -13530,28 +13583,24 @@ packages:
engines: {node: '>= 12.0.0'}
cpu: [arm64]
os: [linux]
libc: [glibc]
lightningcss-linux-arm64-musl@1.32.0:
resolution: {integrity: sha512-UpQkoenr4UJEzgVIYpI80lDFvRmPVg6oqboNHfoH4CQIfNA+HOrZ7Mo7KZP02dC6LjghPQJeBsvXhJod/wnIBg==}
engines: {node: '>= 12.0.0'}
cpu: [arm64]
os: [linux]
libc: [musl]
lightningcss-linux-x64-gnu@1.32.0:
resolution: {integrity: sha512-V7Qr52IhZmdKPVr+Vtw8o+WLsQJYCTd8loIfpDaMRWGUZfBOYEJeyJIkqGIDMZPwPx24pUMfwSxxI8phr/MbOA==}
engines: {node: '>= 12.0.0'}
cpu: [x64]
os: [linux]
libc: [glibc]
lightningcss-linux-x64-musl@1.32.0:
resolution: {integrity: sha512-bYcLp+Vb0awsiXg/80uCRezCYHNg1/l3mt0gzHnWV9XP1W5sKa5/TCdGWaR/zBM2PeF/HbsQv/j2URNOiVuxWg==}
engines: {node: '>= 12.0.0'}
cpu: [x64]
os: [linux]
libc: [musl]
lightningcss-win32-arm64-msvc@1.32.0:
resolution: {integrity: sha512-8SbC8BR40pS6baCM8sbtYDSwEVQd4JlFTOlaD3gWGHfThTcABnNDBda6eTZeqbofalIJhFx0qKzgHJmcPTnGdw==}
@@ -15125,6 +15174,10 @@ packages:
resolution: {integrity: sha512-RwFpb72c/BhQLEXIZ5K2e+AhgNVmIejGlTgiB9MzZ0e93GRvqZ7uSi0dvRF7/XIXDeNkra2fNHBxTyPDGySpjQ==}
engines: {node: '>=8'}
p-queue@8.1.1:
resolution: {integrity: sha512-aNZ+VfjobsWryoiPnEApGGmf5WmNsCo9xu8dfaYamG5qaLP7ClhLN6NgsFe6SwJ2UbLEBK5dv9x8Mn5+RVhMWQ==}
engines: {node: '>=18'}
p-reduce@2.1.0:
resolution: {integrity: sha512-2USApvnsutq8uoxZBGbbWM0JIYLiEMJ9RlaN7fAzVNb9OZN0SHjjTTfIcb667XynS5Y1VhwDJVDa72TnPzAYWw==}
engines: {node: '>=8'}
@@ -15133,6 +15186,10 @@ packages:
resolution: {integrity: sha512-rhIwUycgwwKcP9yTOOFK/AKsAopjjCakVqLHePO3CC6Mir1Z99xT+R63jZxAT5lFZLa2inS5h+ZS2GvR99/FBg==}
engines: {node: '>=8'}
p-timeout@6.1.4:
resolution: {integrity: sha512-MyIV3ZA/PmyBN/ud8vV9XzwTrNtR4jFrObymZYnZqMmW0zA8Z17vnT0rBgFE/TlohB+YCHqXMgZzb3Csp49vqg==}
engines: {node: '>=14.16'}
p-try@1.0.0:
resolution: {integrity: sha512-U1etNYuMJoIz3ZXSrrySFjsXQTWOx2/jdi86L+2pRvph/qMKL6sbcCYdH23fqsbm8TH2Gn0OybpT4eSFlCVHww==}
engines: {node: '>=4'}
@@ -15521,6 +15578,9 @@ packages:
resolution: {integrity: sha512-TfySrs/5nm8fQJDcBDuUng3VOUKsd7S+zqvbOTiGXHfxX4wK31ard+hoNuvkicM/2YFzlpDgABOevKSsB4G/FA==}
engines: {node: '>= 6'}
piscina@4.9.2:
resolution: {integrity: sha512-Fq0FERJWFEUpB4eSY59wSNwXD4RYqR+nR/WiEVcZW8IWfVBxJJafcgTEZDQo8k3w0sUarJ8RyVbbUF4GQ2LGbQ==}
pkg-dir@4.2.0:
resolution: {integrity: sha512-HRDzbaKjC+AOWVXxAU/x54COGeIv9eb+6CkDSQoNTt4XyWoIJvuPsXizxu/Fr23EiekbtZwmh1IcIG/l/a10GQ==}
engines: {node: '>=8'}
@@ -15911,6 +15971,10 @@ packages:
proj4@2.20.4:
resolution: {integrity: sha512-/EEBoXjBw+zW1Lofinw0YFQ4OFOqC6XThnwRAgjEHw8WBN+wSy/6aeqRRdjy4b2igy3X0k3RNDuwcYqcQq7nDw==}
prom-client@15.1.3:
resolution: {integrity: sha512-6ZiOBfCywsD4k1BN9IX0uZhF+tJkV8q8llP64G5Hajs4JOeVLPCwpPVcpXy3BwYiUGgyJzsJJQeOIv7+hDSq8g==}
engines: {node: ^16 || ^18 || >=20}
promise-all-reject-late@1.0.1:
resolution: {integrity: sha512-vuf0Lf0lOxyQREH7GDIOUMLS7kz+gs8i6B+Yi8dC68a2sychGrHTJYghMBD6k7eUcH0H5P73EckCA48xijWqXw==}
@@ -16757,56 +16821,48 @@ packages:
engines: {node: '>=14.0.0'}
cpu: [arm64]
os: [linux]
libc: glibc
sass-embedded-linux-arm@1.98.0:
resolution: {integrity: sha512-03baQZCxVyEp8v1NWBRlzGYrmVT/LK7ZrHlF1piscGiGxwfdxoLXVuxsylx3qn/dD/4i/rh7Bzk7reK1br9jvQ==}
engines: {node: '>=14.0.0'}
cpu: [arm]
os: [linux]
libc: glibc
sass-embedded-linux-musl-arm64@1.98.0:
resolution: {integrity: sha512-LeqNxQA8y4opjhe68CcFvMzCSrBuJqYVFbwElEj9bagHXQHTp9xVPJRn6VcrC+0VLEDq13HVXMv7RslIuU0zmA==}
engines: {node: '>=14.0.0'}
cpu: [arm64]
os: [linux]
libc: musl
sass-embedded-linux-musl-arm@1.98.0:
resolution: {integrity: sha512-OBkjTDPYR4hSaueOGIM6FDpl9nt/VZwbSRpbNu9/eEJcxE8G/vynRugW8KRZmCFjPy8j/jkGBvvS+k9iOqKV3g==}
engines: {node: '>=14.0.0'}
cpu: [arm]
os: [linux]
libc: musl
sass-embedded-linux-musl-riscv64@1.98.0:
resolution: {integrity: sha512-7w6hSuOHKt8FZsmjRb3iGSxEzM87fO9+M8nt5JIQYMhHTj5C+JY/vcske0v715HCVj5e1xyTnbGXf8FcASeAIw==}
engines: {node: '>=14.0.0'}
cpu: [riscv64]
os: [linux]
libc: musl
sass-embedded-linux-musl-x64@1.98.0:
resolution: {integrity: sha512-QikNyDEJOVqPmxyCFkci8ZdCwEssdItfjQFJB+D+Uy5HFqcS5Lv3d3GxWNX/h1dSb23RPyQdQc267ok5SbEyJw==}
engines: {node: '>=14.0.0'}
cpu: [x64]
os: [linux]
libc: musl
sass-embedded-linux-riscv64@1.98.0:
resolution: {integrity: sha512-E7fNytc/v4xFBQKzgzBddV/jretA4ULAPO6XmtBiQu4zZBdBozuSxsQLe2+XXeb0X4S2GIl72V7IPABdqke/vA==}
engines: {node: '>=14.0.0'}
cpu: [riscv64]
os: [linux]
libc: glibc
sass-embedded-linux-x64@1.98.0:
resolution: {integrity: sha512-VsvP0t/uw00mMNPv3vwyYKUrFbqzxQHnRMO+bHdAMjvLw4NFf6mscpym9Bzf+NXwi1ZNKnB6DtXjmcpcvqFqYg==}
engines: {node: '>=14.0.0'}
cpu: [x64]
os: [linux]
libc: glibc
sass-embedded-unknown-all@1.98.0:
resolution: {integrity: sha512-C4MMzcAo3oEDQnW7L8SBgB9F2Fq5qHPnaYTZRMOH3Mp/7kM4OooBInXpCiiFjLnjY95hzP4KyctVx0uYR6MYlQ==}
@@ -17661,6 +17717,9 @@ packages:
resolution: {integrity: sha512-tOG/7GyXpFevhXVh8jOPJrmtRpOTsYqUIkVdVooZYJS/z8WhfQUX8RJILmeuJNinGAMSu1veBr4asSHFt5/hng==}
engines: {node: '>=18'}
tdigest@0.1.2:
resolution: {integrity: sha512-+G0LLgjjo9BZX2MfdvPfH+MKLCrxlXSYec5DaPYP1fe6Iyhf0/fSmJ0bFiZ1F8BT6cGXl2LpltQptzjXKWEkKA==}
teex@1.0.1:
resolution: {integrity: sha512-eYE6iEI62Ni1H8oIa7KlDU6uQBtqr4Eajni3wX7rpfXD8ysFx8z0+dri+KWEPWpBsxXfxu58x/0jvTVT1ekOSg==}
@@ -21411,6 +21470,32 @@ snapshots:
'@colors/colors@1.6.0': {}
'@coopenomics/coopos-ship-reader@0.3.1':
dependencies:
'@wharfkit/antelope': 1.2.0
ws: 8.20.0
transitivePeerDependencies:
- bufferutil
- utf-8-validate
'@coopenomics/parser2@1.2.0':
dependencies:
'@coopenomics/coopos-ship-reader': 0.3.1
'@wharfkit/antelope': 1.2.0
ajv: 8.18.0
commander: 12.1.0
ioredis: 5.10.1
p-queue: 8.1.1
pino: 9.14.0
pino-pretty: 13.1.3
piscina: 4.9.2
prom-client: 15.1.3
yaml: 2.8.3
transitivePeerDependencies:
- bufferutil
- supports-color
- utf-8-validate
'@coopenomics/provider-client@2025.11.12-alpha-1':
dependencies:
axios: 1.13.6
@@ -24308,6 +24393,78 @@ snapshots:
'@msgpackr-extract/msgpackr-extract-win32-x64@3.0.3':
optional: true
'@napi-rs/nice-android-arm-eabi@1.1.1':
optional: true
'@napi-rs/nice-android-arm64@1.1.1':
optional: true
'@napi-rs/nice-darwin-arm64@1.1.1':
optional: true
'@napi-rs/nice-darwin-x64@1.1.1':
optional: true
'@napi-rs/nice-freebsd-x64@1.1.1':
optional: true
'@napi-rs/nice-linux-arm-gnueabihf@1.1.1':
optional: true
'@napi-rs/nice-linux-arm64-gnu@1.1.1':
optional: true
'@napi-rs/nice-linux-arm64-musl@1.1.1':
optional: true
'@napi-rs/nice-linux-ppc64-gnu@1.1.1':
optional: true
'@napi-rs/nice-linux-riscv64-gnu@1.1.1':
optional: true
'@napi-rs/nice-linux-s390x-gnu@1.1.1':
optional: true
'@napi-rs/nice-linux-x64-gnu@1.1.1':
optional: true
'@napi-rs/nice-linux-x64-musl@1.1.1':
optional: true
'@napi-rs/nice-openharmony-arm64@1.1.1':
optional: true
'@napi-rs/nice-win32-arm64-msvc@1.1.1':
optional: true
'@napi-rs/nice-win32-ia32-msvc@1.1.1':
optional: true
'@napi-rs/nice-win32-x64-msvc@1.1.1':
optional: true
'@napi-rs/nice@1.1.1':
optionalDependencies:
'@napi-rs/nice-android-arm-eabi': 1.1.1
'@napi-rs/nice-android-arm64': 1.1.1
'@napi-rs/nice-darwin-arm64': 1.1.1
'@napi-rs/nice-darwin-x64': 1.1.1
'@napi-rs/nice-freebsd-x64': 1.1.1
'@napi-rs/nice-linux-arm-gnueabihf': 1.1.1
'@napi-rs/nice-linux-arm64-gnu': 1.1.1
'@napi-rs/nice-linux-arm64-musl': 1.1.1
'@napi-rs/nice-linux-ppc64-gnu': 1.1.1
'@napi-rs/nice-linux-riscv64-gnu': 1.1.1
'@napi-rs/nice-linux-s390x-gnu': 1.1.1
'@napi-rs/nice-linux-x64-gnu': 1.1.1
'@napi-rs/nice-linux-x64-musl': 1.1.1
'@napi-rs/nice-openharmony-arm64': 1.1.1
'@napi-rs/nice-win32-arm64-msvc': 1.1.1
'@napi-rs/nice-win32-ia32-msvc': 1.1.1
'@napi-rs/nice-win32-x64-msvc': 1.1.1
optional: true
'@napi-rs/wasm-runtime@0.2.12':
dependencies:
'@emnapi/core': 1.9.1
@@ -24950,7 +25107,7 @@ snapshots:
'@opentelemetry/api-logs@0.53.0':
dependencies:
'@opentelemetry/api': 1.9.0
'@opentelemetry/api': 1.9.1
'@opentelemetry/api@1.9.0': {}
@@ -27994,7 +28151,7 @@ snapshots:
'@wharfkit/abicache@1.2.2':
dependencies:
'@wharfkit/antelope': 1.1.1
'@wharfkit/antelope': 1.2.0
'@wharfkit/signing-request': 3.4.0
pako: 2.1.0
tslib: 2.8.1
@@ -28017,6 +28174,15 @@ snapshots:
pako: 2.1.0
tslib: 2.8.1
'@wharfkit/antelope@1.2.0':
dependencies:
bn.js: 4.12.3
brorand: 1.1.0
elliptic: 6.6.1
hash.js: 1.1.7
pako: 2.1.0
tslib: 2.8.1
'@wharfkit/common@1.5.0(@wharfkit/antelope@1.1.1)':
dependencies:
'@wharfkit/antelope': 1.1.1
@@ -28031,7 +28197,7 @@ snapshots:
'@wharfkit/resources@1.5.0':
dependencies:
'@wharfkit/antelope': 1.1.1
'@wharfkit/antelope': 1.2.0
bn.js: 4.12.3
js-big-decimal: 2.2.0
tslib: 2.8.1
@@ -28047,12 +28213,12 @@ snapshots:
'@wharfkit/signing-request@3.4.0':
dependencies:
'@wharfkit/antelope': 1.1.1
'@wharfkit/antelope': 1.2.0
tslib: 2.8.1
'@wharfkit/token@1.2.0':
dependencies:
'@wharfkit/antelope': 1.1.1
'@wharfkit/antelope': 1.2.0
'@wharfkit/contract': 1.2.1
bn.js: 4.12.3
tslib: 2.8.1
@@ -28845,6 +29011,8 @@ snapshots:
dependencies:
file-uri-to-path: 1.0.0
bintrees@1.0.2: {}
birpc@2.9.0: {}
bl@4.1.0:
@@ -37153,12 +37321,19 @@ snapshots:
eventemitter3: 4.0.7
p-timeout: 3.2.0
p-queue@8.1.1:
dependencies:
eventemitter3: 5.0.4
p-timeout: 6.1.4
p-reduce@2.1.0: {}
p-timeout@3.2.0:
dependencies:
p-finally: 1.0.0
p-timeout@6.1.4: {}
p-try@1.0.0: {}
p-try@2.2.0: {}
@@ -37560,6 +37735,10 @@ snapshots:
pirates@4.0.7: {}
piscina@4.9.2:
optionalDependencies:
'@napi-rs/nice': 1.1.1
pkg-dir@4.2.0:
dependencies:
find-up: 4.1.0
@@ -37967,6 +38146,11 @@ snapshots:
mgrs: 1.0.0
wkt-parser: 1.5.4
prom-client@15.1.3:
dependencies:
'@opentelemetry/api': 1.9.1
tdigest: 0.1.2
promise-all-reject-late@1.0.1: {}
promise-call-limit@3.0.2: {}
@@ -40205,6 +40389,10 @@ snapshots:
minizlib: 3.1.0
yallist: 5.0.0
tdigest@0.1.2:
dependencies:
bintrees: 1.0.2
teex@1.0.1:
dependencies:
streamx: 2.25.0