fix(capital): регистрация долей Благороста при создании проекта по событию (не ждать крон) #151

Merged
ant merged 2 commits from fix/capital-share-onproject-event into dev 2026-06-16 08:15:30 +00:00
5 changed files with 186 additions and 15 deletions
@@ -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
);
}
}
}
@@ -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,
@@ -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 {
@@ -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();
});
});