[989-4][@ant] fix(process-registry): end-to-end работает на живом стенде — baseline через live blockchain+parser+controller

FIXES после live-прогона через pnpm run reboot + add-test-user:

1) parser/src/config.ts: добавлены ledger2+marketplace в subscribedContracts —
   без этого parser не пускал ledger2-дельты в Redis stream и controller не
   получал wjournal/journal, т.е. весь pipeline стоял.

2) controller/src/infrastructure/blockchain/blockchain-consumer.service.ts:
   processDeltaDelayed фильтровал по `delta.value?.coopname`, но ledger2
   (wjournal/journal/wallets/accounts) и большинство кооп-scope таблиц
   хранят coopname В SCOPE, а не в value.jsonb. Добавлен fallback на scope:
   `deltaCoop = delta.value?.coopname ?? delta.scope`.

3) domain/process-registry/services/process-registry.service.ts:
   - LEDGER2_CODE изменён с `_ledger2` на `ledger2` (имя on-chain аккаунта
     реальное, без подчёркивания — verified через cleos get abi ledger2).
   - phase A + listProcesses: coopname-скоупинг по `d.scope` (вместо
     несуществующего value->>'coopname' для ledger2 journals).
   - phase B (scanEntityDeltas): принимает оба варианта — `scope = $coop
     OR value->>'coopname' = $coop`, чтобы охватить и per-coop-scope
     таблицы, и singleton-scope контракты (registrator.regs).
   - LOWER() обе стороны при сравнении process_hash: ончейн хранит
     checksum256 uppercase, а нормализация API — lowercase.
   - listProcesses возвращает processHash в lowercase через LOWER() в
     SELECT (единообразно с getProcess).
   - countDeltasByHash/countDocumentsByHash: LOWER() + scope=coop.

4) migrations/V2.1.0: expression-индексы ledger2 journals теперь на
   (process_hash, scope) и (process_type, scope) — совпадают с фактическим
   where-условием сервиса.

5) cooptypes/common/names: `_ledger2.production/testnet = "ledger2"`
   (без подчёркивания, соответствует on-chain имени).

6) boot: добавлен CLI `pnpm run cli add-test-user <username>` для smoke-
   проверки Epic 4 на живом стенде после reboot. Также установлен
   `registration_hash: generateRandomSHA256()` в:
   - boot/init/infra.ts: adduser(ant) + adduser для 4 дополнительных
     членов совета (boot:extra mode);
   - boot/init/participant.ts: addUser+addUser2.

ПРОВЕРКА (live через curl + JWT подписанный JWT_SECRET из .env):

cleos get table ledger2 voskhod wjournal — 2 записи:
  id=0 reg.minshare, process_type=reg.regist, process_hash=<sha256>
  id=1 reg.entrfee,  process_type=reg.regist, process_hash=<sha256>
  (тот же hash — мульти-операционный процесс reg.regist)

query { process(hash:"28fe0b46...",coopname:"voskhod") }
→ process_type="reg.regist"
  delta_history: 4 (2 wjournal + 2 journal)
  actions: 3 (registrator::adduser + 2×ledger2::apply)
  documents: []

query { processes(filter:{coopname:"voskhod"}, pagination:{...}) }
→ totalCount=1, processHash lowercase, actionCount=3, deltaCount=4

Redis cache: ключ `process::voskhod::<hash>`, TTL=60s ✓
Auth: без JWT → 401 Unauthorized ✓

pnpm test — 67/67 passed (все unit-тесты по-прежнему зелёные).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
coopops
2026-04-18 08:55:56 +00:00
parent ede721deb3
commit eeb4d50f98
9 changed files with 105 additions and 27 deletions
+15
View File
@@ -12,6 +12,7 @@ import { sleep } from './utils'
import { checkHealth } from './docker/health'
import { clearDB, clearDirectory, deleteFile } from './docker/purge'
import { deployCommand } from './docker/deploy'
import { addTestUser } from './scripts/add-test-user'
config()
@@ -30,6 +31,20 @@ const program = new Command()
program.version('0.1.0')
// Epic 4: smoke-скрипт для проверки ProcessRegistry end-to-end.
program
.command('add-test-user <username>')
.description('Добавить пайщика через registrator::adduser — для live-тестов ProcessRegistry')
.action(async (username: string) => {
try {
await addTestUser(username)
process.exit(0)
} catch (e) {
console.error('Failed:', e)
process.exit(1)
}
})
// Команда для запуска команды в контейнере
program
.command('cleos [cmd...]')
+3
View File
@@ -8,6 +8,7 @@ import type { Account, Contract } from '../types'
import config from '../configs'
import Blockchain from '../blockchain'
import { sleep } from '../utils'
import { generateRandomSHA256 } from '../utils/randomHash'
import { initUsersInPostgres, initVaultInPostgres } from '../postgres-init'
import { CooperativeClass } from './cooperative'
@@ -390,6 +391,7 @@ export async function installInitialData(blockchain: Blockchain, isExtended = fa
minimum: '200.0000 RUB',
spread_initial: true,
meta: 'Основатель кооператива ВОСХОД',
registration_hash: generateRandomSHA256(),
})
console.log('Устанавливаем дефолтный публичный ключ для ant')
@@ -481,6 +483,7 @@ export async function installInitialData(blockchain: Blockchain, isExtended = fa
minimum: '300.0000 RUB',
spread_initial: true,
meta: `Член совета кооператива ВОСХОД - ${user.first_name} ${user.middle_name} ${user.last_name}`,
registration_hash: generateRandomSHA256(),
})
console.log(`Устанавливаем дефолтный публичный ключ для ${user.username}`)
+3
View File
@@ -6,6 +6,7 @@ import type { Account, Contract, Keys } from '../types'
import config, { GOVERN_SYMBOL, SYMBOL } from '../configs'
import Blockchain from '../blockchain'
import { sendPostToCoopbackWithSecret, sleep } from '../utils'
import { generateRandomSHA256 } from '../utils/randomHash'
import { fakeDocument } from '../tests/shared/fakeDocument'
export class ParticipantsClass {
@@ -33,6 +34,7 @@ export class ParticipantsClass {
minimum: '100.0000 RUB',
spread_initial: false,
meta: '',
registration_hash: generateRandomSHA256(),
}
await this.blockchain.api.transact({
@@ -113,6 +115,7 @@ export class ParticipantsClass {
minimum: '100.0000 RUB',
spread_initial: false,
meta: '',
registration_hash: generateRandomSHA256(),
}
await this.blockchain.api.transact({
@@ -0,0 +1,38 @@
/**
* Одноразовый скрипт: добавляет тестового пайщика через registrator::adduser,
* чтобы parser уловил live ledger2-дельты и controller получил их в
* blockchain_deltas (для проверки ProcessRegistryService end-to-end).
*
* Запуск: pnpm run cli add-test-user <username>
*/
import Blockchain from '../blockchain'
import { generateRandomSHA256 } from '../utils/randomHash'
import config from '../configs'
import { RegistratorContract } from 'cooptypes'
export async function addTestUser(username: string) {
const blockchain = new Blockchain(config.network, config.private_keys)
await blockchain.update_pass_instance()
const registration_hash = generateRandomSHA256()
console.log(`\nДобавляем пайщика ${username} с registration_hash=${registration_hash}`)
const data: RegistratorContract.Actions.AddUser.IAddUser = {
coopname: 'voskhod',
referer: '',
username,
type: 'individual',
created_at: '2025-01-15T10:00:00',
initial: '100.0000 RUB',
minimum: '200.0000 RUB',
spread_initial: true,
meta: 'Тестовый пайщик для проверки ProcessRegistry',
registration_hash,
}
await blockchain.addUser(data)
console.log(`\nУчастник ${username} добавлен. process_hash=${registration_hash}`)
console.log(`Проверить ledger2: cleos get table ledger2 voskhod wjournal`)
console.log(`Проверить GraphQL: query { process(hash: "${registration_hash}", coopname: "voskhod") { ... } }`)
}
@@ -25,20 +25,23 @@ export default {
async up({ dataSource, logger }: { dataSource: DataSource; logger: MigrationLogger }): Promise<boolean> {
try {
// ----- phase A: ledger2 journals -----
// Важно: ledger2 wjournal/journal хранят coopname как scope, а не в
// value.jsonb — индексы покрывают (process_hash, scope) и аналогично
// для process_type/username.
await dataSource.query(`
CREATE INDEX IF NOT EXISTS "idx_deltas_ledger2_process_hash"
ON "blockchain_deltas" (("value"->>'process_hash'), ("value"->>'coopname'))
WHERE code = '_ledger2' AND "table" IN ('wjournal','journal')
ON "blockchain_deltas" (("value"->>'process_hash'), scope)
WHERE code = 'ledger2' AND "table" IN ('wjournal','journal')
`);
await dataSource.query(`
CREATE INDEX IF NOT EXISTS "idx_deltas_ledger2_process_type"
ON "blockchain_deltas" (("value"->>'process_type'), ("value"->>'coopname'))
WHERE code = '_ledger2' AND "table" IN ('wjournal','journal')
ON "blockchain_deltas" (("value"->>'process_type'), scope)
WHERE code = 'ledger2' AND "table" IN ('wjournal','journal')
`);
await dataSource.query(`
CREATE INDEX IF NOT EXISTS "idx_deltas_ledger2_wjournal_username"
ON "blockchain_deltas" (("value"->>'username'), ("value"->>'coopname'))
WHERE code = '_ledger2' AND "table" = 'wjournal'
ON "blockchain_deltas" (("value"->>'username'), scope)
WHERE code = 'ledger2' AND "table" = 'wjournal'
`);
// ----- phase B: entity tables из PROCESS_HASH_LOCATOR -----
@@ -26,7 +26,7 @@ import {
PaginationResult,
} from '~/application/common/dto/pagination.dto';
const LEDGER2_CODE = '_ledger2';
const LEDGER2_CODE = 'ledger2';
const LEDGER2_JOURNALS = ['wjournal', 'journal'] as const;
const HARD_LIMIT = 200;
const CACHE_TTL_SECONDS = 60;
@@ -77,12 +77,15 @@ export class ProcessRegistryService {
if (cached) return cached;
// ---------- Phase A: anchor scan в ledger2-журналах ----------
// ledger2 хранит coopname в scope (не в value.jsonb), поэтому фильтр —
// по scope вместо value->>'coopname'. checksum256 в блокчейн-дельтах
// хранится uppercase, поэтому сравниваем через LOWER() обе стороны.
const anchors = await this.deltaRepository
.createQueryBuilder('d')
.where('d.code = :code', { code: LEDGER2_CODE })
.andWhere('d.table IN (:...tables)', { tables: LEDGER2_JOURNALS })
.andWhere("d.value ->> 'process_hash' = :hash", { hash: normHash })
.andWhere("d.value ->> 'coopname' = :coop", { coop: coopname })
.andWhere("LOWER(d.value ->> 'process_hash') = :hash", { hash: normHash })
.andWhere('d.scope = :coop', { coop: coopname })
.orderBy('d.block_num', 'ASC')
.getMany();
@@ -155,13 +158,14 @@ export class ProcessRegistryService {
const limit = Math.max(1, Math.min(100, pagination.limit ?? 10));
const offset = (page - 1) * limit;
// Подсчёт total (distinct по process_hash)
// Coopname-scoping идёт по d.scope, т.к. ledger2 wjournal/journal хранят
// coopname как scope, а не в value.jsonb.
const countRow = await this.deltaRepository.manager.query(
`SELECT COUNT(DISTINCT d.value ->> 'process_hash') AS cnt
FROM blockchain_deltas d
WHERE d.code = $1
AND d.table = 'wjournal'
AND d.value ->> 'coopname' = $2
AND d.scope = $2
AND (d.value ->> 'process_hash') IS NOT NULL
${filter.processType ? "AND d.value ->> 'process_type' = $3" : ''}
${filter.username ? `AND d.value ->> 'username' = $${filter.processType ? 4 : 3}` : ''}
@@ -173,27 +177,28 @@ export class ProcessRegistryService {
const totalCount = parseInt(countRow[0]?.cnt ?? '0', 10);
const totalPages = Math.max(1, Math.ceil(totalCount / limit));
// Выборка страницы
// Выборка страницы. processHash нормализуется к lowercase для единообразия
// с getProcess (ончейн хранит uppercase, а наружу отдаём lowercase hex).
const rows = await this.deltaRepository.manager.query(
`SELECT
d.value ->> 'process_type' AS "processType",
d.value ->> 'process_hash' AS "processHash",
d.value ->> 'coopname' AS "coopname",
d.value ->> 'process_type' AS "processType",
LOWER(d.value ->> 'process_hash') AS "processHash",
d.scope AS "coopname",
MIN(d.value ->> 'username') AS "username",
MIN(d.created_at) AS "firstSeenAt",
MAX(d.created_at) AS "lastSeenAt"
FROM blockchain_deltas d
WHERE d.code = $1
AND d.table = 'wjournal'
AND d.value ->> 'coopname' = $2
AND d.scope = $2
AND (d.value ->> 'process_hash') IS NOT NULL
${filter.processType ? "AND d.value ->> 'process_type' = $3" : ''}
${filter.username ? `AND d.value ->> 'username' = $${filter.processType ? 4 : 3}` : ''}
${filter.fromBlock ? `AND d.block_num >= $${this.nextParamIdx(filter, 'from')}` : ''}
${filter.toBlock ? `AND d.block_num <= $${this.nextParamIdx(filter, 'to')}` : ''}
GROUP BY d.value ->> 'process_type',
d.value ->> 'process_hash',
d.value ->> 'coopname'
LOWER(d.value ->> 'process_hash'),
d.scope
ORDER BY MAX(d.created_at) DESC
LIMIT ${limit} OFFSET ${offset}
`,
@@ -250,12 +255,15 @@ export class ProcessRegistryService {
if (locations.length === 0) return [];
const all: DeltaEntity[] = [];
for (const loc of locations) {
// Coopname-скоупинг: часть таблиц хранит coopname в scope (ledger2,
// большинство кооп-scope таблиц), часть — в value.jsonb (singleton-scope
// контракты типа registrator). Поддерживаем оба варианта.
const rows = await this.deltaRepository
.createQueryBuilder('d')
.where('d.code = :code', { code: loc.code })
.andWhere('d.table = :table', { table: loc.table })
.andWhere(`d.value ->> :field = :hash`, { field: loc.field, hash })
.andWhere("d.value ->> 'coopname' = :coop", { coop: coopname })
.andWhere(`LOWER(d.value ->> :field) = :hash`, { field: loc.field, hash })
.andWhere("(d.scope = :coop OR d.value ->> 'coopname' = :coop)", { coop: coopname })
.orderBy('d.block_num', 'ASC')
.getMany();
all.push(...rows);
@@ -387,6 +395,7 @@ export class ProcessRegistryService {
}
private async countActionsByHash(hash: string, coopname: string): Promise<number> {
// ILIKE уже case-insensitive, так что проверим по lowercase-подстроке.
const row = await this.actionRepository.manager.query(
`SELECT COUNT(*) AS cnt FROM blockchain_actions a
WHERE (a.data::text ILIKE $1) AND (a.data::text ILIKE $2)`,
@@ -399,8 +408,8 @@ export class ProcessRegistryService {
const row = await this.deltaRepository.manager.query(
`SELECT COUNT(*) AS cnt FROM blockchain_deltas d
WHERE d.code = $1 AND d.table IN ('wjournal','journal')
AND d.value ->> 'process_hash' = $2
AND d.value ->> 'coopname' = $3`,
AND LOWER(d.value ->> 'process_hash') = $2
AND d.scope = $3`,
[LEDGER2_CODE, hash, coopname]
);
return parseInt(row[0]?.cnt ?? '0', 10);
@@ -413,8 +422,8 @@ export class ProcessRegistryService {
const row = await this.deltaRepository.manager.query(
`SELECT COUNT(*) AS cnt FROM blockchain_deltas d
WHERE d.code = $1 AND d.table IN ('wjournal','journal')
AND d.value ->> 'process_hash' = $2
AND d.value ->> 'coopname' = $3`,
AND LOWER(d.value ->> 'process_hash') = $2
AND d.scope = $3`,
[LEDGER2_CODE, hash, coopname]
);
// actions → документ считаем как "1 если есть хоть одна дельта".
@@ -180,7 +180,12 @@ export class BlockchainConsumerService implements OnModuleInit, OnModuleDestroy
}
private async processDeltaDelayed(delta: IDelta): Promise<void> {
if (delta.value?.coopname != config.coopname) {
// Пропускаем дельту, если она не относится к нашему кооперативу.
// coopname может быть в value (большинство таблиц контроллируемых таблиц)
// ИЛИ в scope (ончейн-таблицы с scope=coopname: ledger2 wjournal/journal/
// wallets/accounts, wallet deposits/withdraws, capital debts/results/...).
const deltaCoop = (delta.value as any)?.coopname ?? delta.scope;
if (deltaCoop != config.coopname) {
return;
}
@@ -74,6 +74,6 @@ export const _ledger = {
} as const
export const _ledger2 = {
production: '_ledger2',
testnet: '_ledger2',
production: 'ledger2',
testnet: 'ledger2',
} as const
+2
View File
@@ -30,6 +30,8 @@ export const subscribedContracts: string[] = [
'capital',
'wallet',
'ledger',
'ledger2',
'marketplace',
]
// Автоматически генерируем действия для всех контрактов из списка