Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 33d936e5c8 | |||
| a73f07fa0c | |||
| 77f358e865 | |||
| 2afc71c281 | |||
| 9cf19aa8fe | |||
| 4b8121b0b7 |
@@ -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');
|
||||
+54
-14
@@ -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}`;
|
||||
}
|
||||
+23
@@ -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;
|
||||
}
|
||||
+46
@@ -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();
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user