fix(capital): регистрация долей Благороста при создании проекта по событию (не ждать крон) #151
+65
@@ -0,0 +1,65 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { OnEvent } from '@nestjs/event-emitter';
|
||||
import { CapitalContract } from 'cooptypes';
|
||||
import config from '~/config/config';
|
||||
import type { IDelta } from '~/types/common';
|
||||
import { ProjectStatus } from '../../domain/enums/project-status.enum';
|
||||
import { ProgramShareRegistrationService } from '../services/program-share-registration.service';
|
||||
|
||||
/**
|
||||
* Слушатель дельт `capital::projects`. Как только проект (в т.ч. компонент)
|
||||
* оказывается в статусе active — СРАЗУ регистрирует доли всех активных пайщиков
|
||||
* Благороста в нём, не дожидаясь периодического scheduler'а. В pending не
|
||||
* заводим: заход только в активные проекты (решение пользователя 2026-06-16);
|
||||
* при переходе pending→active придёт новая дельта со status=active.
|
||||
*
|
||||
* Why: окно между созданием проекта и его переводом в `result` может быть
|
||||
* минутами. Контракт `regshare` принимает только pending|active, а отката
|
||||
* `result → active` нет — пайщики, не успевшие попасть до закрытия окна, теряют
|
||||
* долю в компоненте безвозвратно (инцидент voskhod, компонент 011bcd92…,
|
||||
* 2026-06-16). Событие закрывает это окно. Зеркалит
|
||||
* `ProgramShareRegistrationOnUserWalletDeltaListener` (тот реагирует на
|
||||
* изменение баланса, этот — на появление проекта).
|
||||
*
|
||||
* Неблокирующий: обработчик получает событие и неспешно обходит пайщиков;
|
||||
* на dispatch-pipeline/парсер не влияет (EventEmitter2.emit — fire-and-forget).
|
||||
* Промахи (downtime/потеря события) подбирает периодический scheduler-бэкстоп.
|
||||
*/
|
||||
@Injectable()
|
||||
export class ProgramShareRegistrationOnProjectDeltaListener {
|
||||
private readonly logger = new Logger(ProgramShareRegistrationOnProjectDeltaListener.name);
|
||||
|
||||
constructor(
|
||||
private readonly programShareRegistrationService: ProgramShareRegistrationService
|
||||
) {}
|
||||
|
||||
@OnEvent(`delta::${CapitalContract.contractName.production}::${CapitalContract.Tables.Projects.tableName}`)
|
||||
async handleProjectDelta(delta: IDelta): Promise<void> {
|
||||
if (!delta.present) return;
|
||||
if (delta.scope !== config.coopname) return;
|
||||
|
||||
const value = delta.value as CapitalContract.Tables.Projects.IProject | undefined;
|
||||
if (!value?.project_hash) return;
|
||||
|
||||
// Заводим доли только в active-проектах: pending ещё не готов к заходу
|
||||
// (решение пользователя 2026-06-16). При переходе pending→active придёт
|
||||
// новая дельта со status=active — на ней и зайдём.
|
||||
const status = String(value.status);
|
||||
if (status !== ProjectStatus.ACTIVE) return;
|
||||
|
||||
const project_hash = String(value.project_hash);
|
||||
try {
|
||||
await this.programShareRegistrationService.syncProgramSharesForProject(
|
||||
delta.scope,
|
||||
project_hash
|
||||
);
|
||||
} catch (error: unknown) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
const stack = error instanceof Error ? error.stack : undefined;
|
||||
this.logger.warn(
|
||||
`regshare event-driven (project) sync не выполнен: coop=${delta.scope} project=${project_hash}: ${message}`,
|
||||
stack
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
+43
-14
@@ -47,12 +47,12 @@ export class ProgramShareRegistrationService {
|
||||
) {}
|
||||
|
||||
/**
|
||||
* Обход участников в статусах active/import, проектов pending/active; при изменении user_shares относительно capital_contributor_shares — regshare.
|
||||
* Обход участников в статусах active/import по active-проектам; при изменении user_shares относительно capital_contributor_shares — regshare.
|
||||
*/
|
||||
async syncProgramSharesForCoop(coopname: string): Promise<void> {
|
||||
const projects = await this.findActiveProjects(coopname);
|
||||
if (projects.length === 0) {
|
||||
this.logger.debug(`Синхронизация regshare: нет проектов в статусах pending/active для ${coopname}`);
|
||||
this.logger.debug(`Синхронизация regshare: нет active-проектов для ${coopname}`);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -62,15 +62,16 @@ export class ProgramShareRegistrationService {
|
||||
(c.status === ContributorStatus.ACTIVE || c.status === ContributorStatus.IMPORT)
|
||||
);
|
||||
|
||||
const projectHashes = projects.map((p) => p.project_hash);
|
||||
for (const contributor of contributors) {
|
||||
await this.syncContributor(coopname, contributor, projects);
|
||||
await this.syncContributor(coopname, contributor, projectHashes);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Точечная синхронизация regshare для одного пайщика — вызывается из listener'а
|
||||
* на дельты `ledger2::userwallets[w.cap.blago]`. Не пишет лог, если у пайщика
|
||||
* нет ни одного pending/active проекта.
|
||||
* нет ни одного active-проекта.
|
||||
*/
|
||||
async syncProgramSharesForUser(coopname: string, username: string): Promise<void> {
|
||||
const projects = await this.findActiveProjects(coopname);
|
||||
@@ -84,21 +85,49 @@ export class ProgramShareRegistrationService {
|
||||
);
|
||||
if (!contributor) return;
|
||||
|
||||
await this.syncContributor(coopname, contributor, projects);
|
||||
await this.syncContributor(coopname, contributor, projects.map((p) => p.project_hash));
|
||||
}
|
||||
|
||||
/**
|
||||
* Точечная регистрация долей всех активных пайщиков в один проект — вызывается
|
||||
* из listener'а на дельты `capital::projects` сразу при появлении проекта.
|
||||
*
|
||||
* Why: между созданием проекта и его переводом в `result` может пройти меньше
|
||||
* минуты; контракт `regshare` принимает только статусы pending|active, а откат
|
||||
* `result → active` не предусмотрен — значит пайщики, не успевшие попасть в
|
||||
* проект до закрытия окна, теряют долю в нём безвозвратно (инцидент voskhod,
|
||||
* компонент 011bcd92…, 2026-06-16). Реакция на событие закрывает окно, не
|
||||
* дожидаясь периодического scheduler'а.
|
||||
*
|
||||
* Переиспользует тот же `syncContributor`, что и обход по расписанию.
|
||||
*/
|
||||
async syncProgramSharesForProject(coopname: string, project_hash: string): Promise<void> {
|
||||
const contributors = (await this.contributorRepository.findAll()).filter(
|
||||
(c) =>
|
||||
c.coopname === coopname &&
|
||||
(c.status === ContributorStatus.ACTIVE || c.status === ContributorStatus.IMPORT)
|
||||
);
|
||||
if (contributors.length === 0) return;
|
||||
|
||||
for (const contributor of contributors) {
|
||||
await this.syncContributor(coopname, contributor, [project_hash]);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Только active-проекты: заход долей в pending отключён (решение пользователя
|
||||
* 2026-06-16) — заводим/сверяем доли лишь в активных проектах.
|
||||
*/
|
||||
private async findActiveProjects(coopname: string): Promise<ProjectDomainEntity[]> {
|
||||
return (await this.projectRepository.findAll()).filter(
|
||||
(p) =>
|
||||
p.coopname === coopname &&
|
||||
(p.status === ProjectStatus.PENDING || p.status === ProjectStatus.ACTIVE)
|
||||
(p) => p.coopname === coopname && p.status === ProjectStatus.ACTIVE
|
||||
);
|
||||
}
|
||||
|
||||
private async syncContributor(
|
||||
coopname: string,
|
||||
contributor: ContributorDomainEntity,
|
||||
projects: ProjectDomainEntity[]
|
||||
projectHashes: string[]
|
||||
): Promise<void> {
|
||||
const programId = getProgramId(ProgramType.BLAGOROST);
|
||||
|
||||
@@ -124,10 +153,10 @@ export class ProgramShareRegistrationService {
|
||||
const targetParsed = AssetUtils.parseAsset(targetShares);
|
||||
if (!targetParsed.symbol) return;
|
||||
|
||||
for (const project of projects) {
|
||||
for (const projectHash of projectHashes) {
|
||||
const segment = await this.capitalBlockchainPort.getSegmentByProjectUser(
|
||||
coopname,
|
||||
project.project_hash,
|
||||
projectHash,
|
||||
contributor.username
|
||||
);
|
||||
|
||||
@@ -139,19 +168,19 @@ export class ProgramShareRegistrationService {
|
||||
try {
|
||||
await this.capitalBlockchainPort.registerShare({
|
||||
coopname,
|
||||
project_hash: project.project_hash,
|
||||
project_hash: projectHash,
|
||||
username: contributor.username,
|
||||
user_shares: targetShares,
|
||||
});
|
||||
this.logger.log(
|
||||
`regshare: ${contributor.username} → проект ${project.project_hash}, user_shares=${targetShares} (было ${registeredStr})`
|
||||
`regshare: ${contributor.username} → проект ${projectHash}, user_shares=${targetShares} (было ${registeredStr})`
|
||||
);
|
||||
await delay(REGSHARE_TX_GAP_MS);
|
||||
} catch (error: unknown) {
|
||||
const message = error instanceof HttpApiError ? error.message : error instanceof Error ? error.message : String(error);
|
||||
const stack = error instanceof Error ? error.stack : undefined;
|
||||
this.logger.warn(
|
||||
`regshare не выполнен: coop=${coopname} project=${project.project_hash} user=${contributor.username}: ${message}`,
|
||||
`regshare не выполнен: coop=${coopname} project=${projectHash} user=${contributor.username}: ${message}`,
|
||||
stack
|
||||
);
|
||||
}
|
||||
|
||||
@@ -311,6 +311,7 @@ import { GitHubSyncSchedulerService } from './infrastructure/services/github-syn
|
||||
import { ProgramShareRegistrationSchedulerService } from './infrastructure/services/program-share-registration-scheduler.service';
|
||||
import { ProgramShareRegistrationService } from './application/services/program-share-registration.service';
|
||||
import { ProgramShareRegistrationOnUserWalletDeltaListener } from './application/listeners/program-share-registration-on-user-wallet-delta.listener';
|
||||
import { ProgramShareRegistrationOnProjectDeltaListener } from './application/listeners/program-share-registration-on-project-delta.listener';
|
||||
import { CapitalDevelopmentRepositoryGitSyncService } from './application/services/capital-development-repository-git-sync.service';
|
||||
import { CapitalGithubExtensionLifecycleListener } from './application/listeners/capital-github-extension-lifecycle.listener';
|
||||
import { GitCommitMarkersSyncService } from './application/services/git-commit-markers-sync.service';
|
||||
@@ -875,6 +876,7 @@ IssueIdGenerationService,
|
||||
ProgramShareRegistrationService,
|
||||
ProgramShareRegistrationSchedulerService,
|
||||
ProgramShareRegistrationOnUserWalletDeltaListener,
|
||||
ProgramShareRegistrationOnProjectDeltaListener,
|
||||
GitCommitMarkersSyncService,
|
||||
{
|
||||
provide: ISSUE_LINKED_GIT_COMMIT_REPOSITORY,
|
||||
|
||||
+8
-1
@@ -3,8 +3,15 @@ import { config } from '~/config';
|
||||
import { ProgramShareRegistrationService } from '../../application/services/program-share-registration.service';
|
||||
|
||||
/**
|
||||
* Периодическая синхронизация долей участников (regshare) по балансу программы Благорост.
|
||||
* Периодическая сверка долей участников (regshare) с балансом программы Благорост.
|
||||
* Интервал задаётся в конфигурации расширения Capital (минуты); 0 — отключено.
|
||||
*
|
||||
* Роль — reconciliation-бэкстоп, не основной путь. Горячие сценарии покрыты
|
||||
* событиями: появление проекта — `ProgramShareRegistrationOnProjectDeltaListener`,
|
||||
* изменение баланса — `ProgramShareRegistrationOnUserWalletDeltaListener`. Крон
|
||||
* оставлен, потому что он ещё и ДОобновляет уже зарегистрированные доли при
|
||||
* дрейфе баланса (контракт `upsert_contributor_segment` это допускает) и
|
||||
* подбирает события, потерянные при downtime контроллера.
|
||||
*/
|
||||
@Injectable()
|
||||
export class ProgramShareRegistrationSchedulerService implements OnModuleDestroy {
|
||||
|
||||
+68
@@ -0,0 +1,68 @@
|
||||
/**
|
||||
* Unit-тесты ProgramShareRegistrationOnProjectDeltaListener.
|
||||
*
|
||||
* Фокус — гейтинг: доли регистрируются СРАЗУ при появлении проекта в статусе
|
||||
* active, и только в своём кооперативе. На pending (заводим только в active),
|
||||
* present=false, чужой scope и прочие статусы реакции нет.
|
||||
*/
|
||||
|
||||
import config from '~/config/config';
|
||||
import { ProgramShareRegistrationOnProjectDeltaListener } from '~/extensions/capital/application/listeners/program-share-registration-on-project-delta.listener';
|
||||
import { ProjectStatus } from '~/extensions/capital/domain/enums/project-status.enum';
|
||||
import type { IDelta } from '~/types/common';
|
||||
|
||||
const PROJECT_HASH = '011bcd92cc6fc6c3fbac1aa31e5ab302590993fedb666642a5bd2f88e96e6a0e';
|
||||
|
||||
function makeServiceStub() {
|
||||
return { syncProgramSharesForProject: jest.fn(async () => undefined) } as any;
|
||||
}
|
||||
|
||||
function makeListener(service: any) {
|
||||
return new ProgramShareRegistrationOnProjectDeltaListener(service);
|
||||
}
|
||||
|
||||
function makeDelta(overrides: Partial<IDelta> = {}, value: Record<string, any> = {}): IDelta {
|
||||
return {
|
||||
present: true,
|
||||
scope: config.coopname,
|
||||
value: { project_hash: PROJECT_HASH, status: ProjectStatus.ACTIVE, ...value },
|
||||
...overrides,
|
||||
} as IDelta;
|
||||
}
|
||||
|
||||
describe('ProgramShareRegistrationOnProjectDeltaListener', () => {
|
||||
it('active в своём кооперативе → регистрирует доли по project_hash', async () => {
|
||||
const service = makeServiceStub();
|
||||
await makeListener(service).handleProjectDelta(makeDelta());
|
||||
expect(service.syncProgramSharesForProject).toHaveBeenCalledWith(config.coopname, PROJECT_HASH);
|
||||
});
|
||||
|
||||
it('pending → не реагирует (заводим только в active, при переходе в active придёт новая дельта)', async () => {
|
||||
const service = makeServiceStub();
|
||||
await makeListener(service).handleProjectDelta(makeDelta({}, { status: ProjectStatus.PENDING }));
|
||||
expect(service.syncProgramSharesForProject).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('result → не реагирует (окно закрыто, откат статуса контракт не даёт)', async () => {
|
||||
const service = makeServiceStub();
|
||||
await makeListener(service).handleProjectDelta(makeDelta({}, { status: ProjectStatus.RESULT }));
|
||||
expect(service.syncProgramSharesForProject).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('present=false → не реагирует', async () => {
|
||||
const service = makeServiceStub();
|
||||
await makeListener(service).handleProjectDelta(makeDelta({ present: false }));
|
||||
expect(service.syncProgramSharesForProject).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('чужой кооператив (scope) → не реагирует', async () => {
|
||||
const service = makeServiceStub();
|
||||
await makeListener(service).handleProjectDelta(makeDelta({ scope: 'othercoop' }));
|
||||
expect(service.syncProgramSharesForProject).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('ошибка сервиса не пробрасывается (best-effort, бэкстоп — scheduler)', async () => {
|
||||
const service = { syncProgramSharesForProject: jest.fn(async () => { throw new Error('boom'); }) } as any;
|
||||
await expect(makeListener(service).handleProjectDelta(makeDelta())).resolves.toBeUndefined();
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user