Compare commits

...

6 Commits

Author SHA1 Message Date
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
14 changed files with 575 additions and 15 deletions
@@ -123,6 +123,38 @@ 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)),
/**
* Потолок ожидания (мс) завершения обработчиков отката форка перед тем,
* как consumer продолжит обработку (TTL force-resume). Защищает от
* зависшего синкера, который иначе заблокировал бы поток навсегда.
*/
BLOCKCHAIN_FORK_PAUSE_TIMEOUT_MS: z
.string()
.default('30000')
.transform((val) => parseInt(val, 10)),
/**
* Активирует idempotency-gate (Story 2.3, INV-09): повторно доставленное
* событие с уже отмеченным event_id игнорируется как no-op. По умолчанию
* OFF в релизе 1.1.1 — локальная формула event_id ещё не сверена с
* авторитетной из parser2 (валидация в Epic 3 phase 2), а ложный дубль =
* silent data loss (тот же класс багов, что чинил Эпик 1). Dual-write меток
* (Story 2.2) наполняет consumer_dedup независимо от флага; пока флаг false
* первичной защитой остаётся block_num-guard (DEC-020). Флаг переводится в
* true после сверки формулы в Epic 3.
*/
BLOCKCHAIN_DEDUP_ENABLED: z
.string()
.default('false')
.transform((val) => val === 'true'),
// Параметры NOVU
NOVU_APP_ID: z.string().min(1, { message: 'Не должно быть пустым' }),
@@ -231,6 +263,9 @@ 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,
fork_pause_timeout_ms: envVars.data.BLOCKCHAIN_FORK_PAUSE_TIMEOUT_MS,
dedup_enabled: envVars.data.BLOCKCHAIN_DEDUP_ENABLED,
},
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,26 @@ 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 не падает.
*/
async markEventApplied(eventId: string): Promise<void> {
await this.consumerDedupRepository.markApplied(eventId);
}
/**
* Сохранение действия блокчейна
*/
@@ -0,0 +1,19 @@
/**
* Порт списка применённых событий (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 не должен падать.
*/
markApplied(eventId: string): Promise<void>;
/** Очистка меток старше cutoff (retention). Возвращает число удалённых строк. */
deleteOlderThan(cutoff: Date): Promise<number>;
}
export const CONSUMER_DEDUP_REPOSITORY_PORT = Symbol('ConsumerDedupRepositoryPort');
@@ -6,6 +6,7 @@ import { RedisStreamService, StreamMessage } from '~/infrastructure/redis/redis-
import { WinstonLoggerService } from '~/application/logger/logger-app.service';
import { EventsService } from '~/infrastructure/events/events.service';
import { ParserInteractor } from '~/domain/parser/interactors/parser.interactor';
import { computeActionEventId, computeDeltaEventId } from './event-id.util';
import { config } from '~/config';
export interface BlockchainEventData {
@@ -248,14 +249,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 +259,15 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
return;
}
// Idempotency: признак уникальности события (Story 2.2/2.3, INV-09).
const eventId = computeActionEventId(action);
// Dedup-gate (Story 2.3, за флагом — см. processDeltaDelayed).
if (config.blockchain.dedup_enabled && (await this.parserInteractor.isEventApplied(eventId))) {
this.logger.debug(`Action-дубликат пропущен (no-op): ${eventId}`);
return;
}
try {
// Сохраняем действие в базу данных через интерактор
await this.parserInteractor.saveAction(action);
@@ -277,10 +279,15 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
throw error; // Перебрасываем ошибку чтобы сообщение не было подтверждено
}
// Публикуем событие с задержкой — пусть сначала прокатятся дельты этого же блока.
// saveAction уже выполнен, так что данные не потеряем; задерживаем только emit.
// Dual-write метки в consumer_dedup ПОСЛЕ save, ДО отложенного emit (Story 2.2).
await this.parserInteractor.markEventApplied(eventId);
// Публикуем событие с задержкой — пусть сначала прокатятся дельты этого же блока
// (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(
@@ -331,6 +338,17 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
return;
}
// Idempotency: признак уникальности события (Story 2.2/2.3, INV-09).
const eventId = computeDeltaEventId(delta);
// Dedup-gate (Story 2.3): повторно доставленная дельта = no-op. За флагом —
// формула event_id ещё не сверена с parser2 (Epic 3), а ложный дубль =
// silent data loss; пока флаг false первичной защитой остаётся block_num-guard.
if (config.blockchain.dedup_enabled && (await this.parserInteractor.isEventApplied(eventId))) {
this.logger.debug(`Дельта-дубликат пропущена (no-op): ${eventId}`);
return;
}
try {
// Сохраняем дельту в базу данных через интерактор
await this.parserInteractor.saveDelta(delta);
@@ -340,6 +358,12 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
throw error; // Перебрасываем ошибку чтобы сообщение не было подтверждено
}
// Dual-write метки в consumer_dedup ПОСЛЕ save (Story 2.2). Если markApplied
// упадёт — handleMessage пробросит ошибку, сообщение останется pending и
// переиграется (saveDelta идемпотентен через block_num-guard, mark — через
// ON CONFLICT DO NOTHING).
await this.parserInteractor.markEventApplied(eventId);
// Публикуем событие во внутреннюю шину с типизированным именем
const eventName = `delta::${delta.code}::${delta.table}`;
this.eventsService.emit(eventName, delta);
@@ -366,10 +390,26 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
throw error; // Перебрасываем ошибку чтобы сообщение не было подтверждено
}
// Публикуем событие во внутреннюю шину с типизированным именем
// Барьер форка (Story 1.3, временное решение до Epic 3 / ForkRegistry).
// Раньше fork эмитился fire-and-forget: processFork возвращался и сообщение
// ACK'алось ДО завершения откатов в @OnEvent('fork::*')-синкерах, из-за чего
// следующая дельта (того же/соседнего блока) обрабатывалась параллельно с
// откатом — гонка. Consumer-loop последователен (COUNT 1, await handleMessage),
// поэтому ожидание здесь = пауза обработки: ждём завершения всех откатов до
// продолжения потока. TTL force-resume: если синкер завис, не блокируем
// consumer навсегда — продолжаем после fork_pause_timeout_ms с предупреждением.
const eventName = `fork::${block_num}`;
this.eventsService.emit(eventName, { block_num });
const completed = await this.eventsService.emitAsyncWithTimeout(
eventName,
{ block_num },
config.blockchain.fork_pause_timeout_ms
);
if (!completed) {
this.logger.warn(
`Барьер форка ${eventName}: откаты не завершились за ${config.blockchain.fork_pause_timeout_ms}ms — продолжаем (force-resume)`
);
}
this.logger.debug(`Форк опубликован в событийную шину: ${eventName}`);
this.logger.debug(`Форк обработан (откаты завершены/таймаут): ${eventName}`);
}
}
@@ -0,0 +1,38 @@
import { IAction, IDelta } from '~/types/common';
/**
* Локальное вычисление event_id по формуле parser2 (Story 2.2, DEC-T08 phase 1).
*
* Формат (см. controller/CLAUDE.md): `${chain}:${kind}:${block_num}:${block_id_short}:${natural_key}`.
* Признак уникальности события — основа идемпотентности (INV-09). До миграции
* на parser2 (Epic 3) считается здесь параллельно текущему flow; в Epic 3 phase 2
* сверяется с авторитетным event_id, который отдаёт сам движок.
*
* ВАЖНО: формула ещё не верифицирована против parser2, поэтому dedup-gate
* (Story 2.3) по умолчанию выключен (BLOCKCHAIN_DEDUP_ENABLED=false). Ложный
* дубль = silent data loss, поэтому активация — только после сверки.
*/
/** Длина короткого префикса 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}`;
}
@@ -0,0 +1,23 @@
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) формула сверяется с авторитетной из движка.
*/
@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;
}
@@ -0,0 +1,46 @@
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): Promise<void> {
// ON CONFLICT DO NOTHING: повторная отметка (re-delivery / crash-recovery)
// не должна падать на нарушении PK.
await this.repository
.createQueryBuilder()
.insert()
.into(ConsumerDedupEntity)
.values({ event_id: eventId })
.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;
}
}
@@ -50,6 +50,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,6 +125,7 @@ import { UserWalletIndexInitializer } from './blockchain/services/user-wallet-in
DeltaEntity,
ForkEntity,
SyncStateEntity,
ConsumerDedupEntity,
EntityVersionTypeormEntity,
SettingsEntity,
TokenEntity,
@@ -197,6 +201,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,
@@ -269,6 +277,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,
@@ -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);
}
}
}
@@ -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);
@@ -0,0 +1,58 @@
/**
* 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 вставляет с orIgnore (ON CONFLICT DO NOTHING)', 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' });
expect(qb.orIgnore).toHaveBeenCalled();
expect(qb.execute).toHaveBeenCalled();
});
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);
});
});
@@ -0,0 +1,78 @@
/**
* 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 } 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()));
});
});
@@ -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,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();
});
});