[B2D-15][@ant] fix(coopback billing): цепочка invoice→pay→confirm + реактивный pay-листенер — чтобы оплата реально продлевала подписки и переживала падение парсера

- BillingCronService: перед on-chain pay создаёт PENDING-invoice у провайдера
  (POST /billing/invoice — раньше invoice никто не создавал и подтверждение
  падало на 'Invoice не найден'), подтверждает на POST /billing/payment-confirmed
  (прежний /subscriptions/confirm-batch-payment на провайдере выпилен в v5).
  Ошибка anti-replay контракта = оплата уже в чейне → доносим подтверждение
  без повторного списания.
- BillingPaymentListener: второй (реактивный) путь подтверждения по событию
  парсера action::billing::pay — переживает падение backend'а между transact
  и callback; провайдер идемпотентен по payment_hash.
- BillingBlockchainAdapter.pay: подписывает оператор (config.coopname узла-хаба,
  = _provider контрактов), не кооператив-спица.
- BILLING_CRON_PAYER удалён: пайщик-плательщик = сам кооператив (w.wal.bill его),
  молчаливый пропуск тика по пустому конфигу больше невозможен.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
coopops
2026-06-11 06:16:10 +00:00
parent 38ecf7eff5
commit 9463ab7583
6 changed files with 218 additions and 63 deletions
@@ -3,6 +3,7 @@ import { BillingService } from './services/billing.service';
import { BillingResolver } from './resolvers/billing.resolver';
import { BillingProviderClient } from '~/infrastructure/billing/billing-provider.client';
import { BillingConversionListener } from '~/infrastructure/billing/billing-conversion.listener';
import { BillingPaymentListener } from '~/infrastructure/billing/billing-payment.listener';
import { BillingCronService } from '~/domain/billing/services/billing-cron.service';
import { ProviderModule } from '~/application/provider/provider.module';
import { DocumentDomainModule } from '~/domain/document/document.module';
@@ -20,6 +21,8 @@ import { DocumentDomainModule } from '~/domain/document/document.module';
*
* BillingConversionListener ловит on-chain `billing::converttoaxn` с шины
* `action::` (от парсера блокчейна) и реактивно уведомляет провайдера.
* BillingPaymentListener аналогично ловит `billing::pay` — второй (реактивный)
* путь подтверждения time-оплат; первый — синхронный confirm в cron'е.
*/
@Module({
imports: [ProviderModule, DocumentDomainModule],
@@ -28,6 +31,7 @@ import { DocumentDomainModule } from '~/domain/document/document.module';
BillingResolver,
BillingProviderClient,
BillingConversionListener,
BillingPaymentListener,
BillingCronService,
],
exports: [BillingService],
+3 -6
View File
@@ -100,11 +100,6 @@ const envVarsSchema = z.object({
.string()
.default('0 * * * *')
.describe('cron-выражение тика биллинга (по умолчанию ежечасно)'),
BILLING_CRON_PAYER: z
.string()
.default('')
.describe('пайщик-плательщик, чей w.wal.bill дебетуется (пусто — списание пропускается)'),
// Параметры союза кооперативов
UNION_LINK: z
.string()
@@ -285,7 +280,9 @@ export default {
cron_expression: envVars.data.BILLING_CRON_EXPRESSION,
// Список коопов — из on-chain `registrator.coops` через ProviderService
// (см. BillingCronService.activeCoopnames). Env-override отсутствует.
payer: envVars.data.BILLING_CRON_PAYER,
// Подписант списаний — оператор платформы (аккаунт узла, config.coopname);
// отдельного BILLING_CRON_PAYER больше нет: пайщик-плательщик = сам
// кооператив, владелец w.wal.bill.
},
union: {
link: envVars.data.UNION_LINK,
@@ -9,23 +9,30 @@ import { BillingProviderClient } from '~/infrastructure/billing/billing-provider
import { ProviderService } from '~/application/provider/services/provider.service';
/**
* Периодическое списание подписок (Epic 12, Story 12.6) — oracle-паттерн.
* Периодическое списание подписок (Epic 12, Single-Hub v5) — oracle-паттерн.
*
* Antelope не поддерживает deferred_trx, поэтому рекуррентность инициирует
* backend: на каждый тик узел запрашивает у провайдера «сумму к оплате»
* (источник истины по составу/ценам — provider, on-chain их нет), и если есть
* что списывать и срок подошёл — проводит on-chain `billing::pay`, затем
* подтверждает платёж провайдеру (тот продлевает подписки по `payment_hash`).
* backend Восхода. Полный цикл на каждый кооператив:
*
* 1. GET /subscriptions/billing-summary — сумма к оплате и срок (источник
* истины по составу/ценам — provider, on-chain их нет);
* 2. POST /billing/invoice — провайдер фиксирует PENDING-invoice и отдаёт
* детерминированный `payment_hash`. БЕЗ этого шага подтверждение оплаты
* не найдёт invoice и подписки не продлятся;
* 3. on-chain `billing::pay` — подписывает ОПЕРАТОР (`_provider`, аккаунт
* узла-хаба). Контракт отклоняет повтор `payment_hash` (anti-replay
* таблица `paidpayments`) — повторное списание тех же средств невозможно
* даже при потере подтверждения;
* 4. POST /billing/payment-confirmed — провайдер переводит invoice в PAID и
* продлевает подписки. Идемпотентно; этот же callback реактивно шлёт
* BillingPaymentListener по событию парсера, так что зависший парсер или
* упавший между шагами 3-4 backend не теряют оплату: при следующем тике
* контракт ответит «уже проведён», и тик отправит подтверждение повторно.
*
* Инварианты:
* - сумма = 0 (все подписки free) → on-chain списание НЕ выполняется;
* - идемпотентность по `payment_hash` (контракт no-op + провайдер идемпотентен);
* - падение `pay` не зацикливает узел: ошибка логируется, тик продолжает
* следующий кооператив (grace/уведомления — на стороне провайдера/Epic 9).
*
* Плательщик (`config.billing.payer`) — пайщик, чей USER_SHARED-кошелёк
* `w.wal.bill` дебетуется. Списание per-пайщик, поэтому без указанного
* плательщика тик пропускается (см. конфиг `BILLING_CRON_PAYER`).
* следующий кооператив (grace/уведомления — на стороне провайдера).
*/
@Injectable()
export class BillingCronService implements OnModuleInit, OnModuleDestroy {
@@ -81,23 +88,18 @@ export class BillingCronService implements OnModuleInit, OnModuleDestroy {
this.logger.warn('BillingCronService: предыдущий тик ещё выполняется — пропуск');
return;
}
const payer = config.billing.payer;
if (!payer) {
this.logger.warn('BillingCronService: BILLING_CRON_PAYER не задан — нечего дебетовать, пропуск тика');
return;
}
this.running = true;
try {
const coopnames = await this.activeCoopnames();
for (const coopname of coopnames) {
await this.processCoop(coopname, payer);
await this.processCoop(coopname);
}
} finally {
this.running = false;
}
}
private async processCoop(coopname: string, payer: string): Promise<void> {
private async processCoop(coopname: string): Promise<void> {
try {
const summary = await this.providerClient.getBillingSummary(coopname);
@@ -108,38 +110,75 @@ export class BillingCronService implements OnModuleInit, OnModuleDestroy {
return; // срок ещё не подошёл
}
const quantity = `${summary.total_amount.toFixed(config.blockchain.root_govern_precision)} ${summary.currency || config.blockchain.root_govern_symbol}`;
const payableItems = (summary.items ?? [])
.filter((item) => !item.is_free && item.amount > 0)
.map((item) => ({ subscription_id: item.subscription_id, period_days: summary.period_days }));
if (!payableItems.length) {
return;
}
const result = await this.blockchainPort.pay({
coopname,
username: payer,
quantity,
paymentHash: summary.payment_hash,
memo: `Оплата подписок за ${summary.period_days} дн.`,
});
// Шаг 2: PENDING-invoice у провайдера (идемпотентно).
const invoice = await this.providerClient.createInvoice(coopname, payableItems);
if (invoice.status === 'PAID') {
this.logger.log(`Invoice ${invoice.payment_hash} (${coopname}) уже PAID — пропуск`);
return;
}
const transactionId =
result && typeof result === 'object' && 'transaction_id' in result
? String((result as { transaction_id?: unknown }).transaction_id ?? '')
: '';
const quantity = `${invoice.total_amount.toFixed(config.blockchain.root_govern_precision)} ${summary.currency || config.blockchain.root_govern_symbol}`;
// Шаг 3: on-chain списание. Подписывает оператор (аккаунт узла-хаба);
// username — пайщик-кооператив, владелец биллинг-кошелька.
let transactionId = '';
try {
const result = await this.blockchainPort.pay({
coopname,
username: coopname,
quantity,
paymentHash: invoice.payment_hash,
memo: `Оплата подписок за ${summary.period_days} дн.`,
});
transactionId =
result && typeof result === 'object' && 'transaction_id' in result
? String((result as { transaction_id?: unknown }).transaction_id ?? '')
: '';
} catch (error: any) {
if (this.isAlreadyPaidOnChain(error)) {
// Списание уже в чейне (прошлый тик не дошёл до подтверждения) —
// деньги второй раз НЕ списаны (anti-replay), доносим подтверждение.
this.logger.warn(
`billing::pay ${coopname}: payment_hash=${invoice.payment_hash} уже проведён on-chain — отправляю подтверждение провайдеру`,
);
await this.providerClient.confirmPayment({
paymentHash: invoice.payment_hash,
blockchainTransactionId: '',
});
return;
}
throw error;
}
// Шаг 4: синхронное подтверждение (реактивный BillingPaymentListener
// продублирует — провайдер идемпотентен по payment_hash).
await this.providerClient.confirmPayment({
coopname,
paymentHash: summary.payment_hash,
amount: summary.total_amount,
paymentHash: invoice.payment_hash,
blockchainTransactionId: transactionId,
periodDays: summary.period_days,
});
this.logger.log(`Списано ${quantity} за подписки ${coopname} (payment_hash=${summary.payment_hash})`);
this.logger.log(`Списано ${quantity} за подписки ${coopname} (payment_hash=${invoice.payment_hash})`);
} catch (error: any) {
// Падение списания (например, недостаток средств на w.wal.bill) не должно
// зацикливать узел: фиксируем и продолжаем. Перевод подписки в past_due/grace
// и уведомления — на стороне провайдера (Epic 4/9).
// и уведомления — на стороне провайдера (Epic 4/14).
this.logger.error(`BillingCronService: списание для ${coopname} не выполнено: ${error?.message ?? error}`);
}
}
/** Ошибка anti-replay контракта billing: этот payment_hash уже проведён. */
private isAlreadyPaidOnChain(error: any): boolean {
const message = String(error?.message ?? error ?? '');
return message.includes('уже проведён') || message.includes('anti-replay');
}
/**
* Срок оплаты подошёл, если дата следующего платежа не задана (первое списание)
* либо она в прошлом/сегодня.
@@ -0,0 +1,69 @@
import { Injectable, Logger } from '@nestjs/common';
import { OnEvent } from '@nestjs/event-emitter';
import { BillingContract } from 'cooptypes';
import config from '~/config/config';
import type { ActionDomainInterface } from '~/domain/parser/interfaces/action-domain.interface';
import { BillingProviderClient } from './billing-provider.client';
/**
* Single-Hub v5 (Story 12.11, проект «Облачный провайдер») — реактивный мост
* on-chain → провайдер для time-оплат подписок.
*
* Ловит событие шины `action::billing::pay` (списание членских взносов с
* биллинг-кошелька за подписки) и пересылает подтверждение провайдеру через
* {@link BillingProviderClient.confirmPayment} — тот переводит invoice в PAID
* и продлевает подписки.
*
* Это ВТОРОЙ путь подтверждения: первый — синхронный confirm в
* BillingCronService сразу после transact. Дублирование намеренное (провайдер
* идемпотентен по `payment_hash`): если backend упал между transact и
* подтверждением, факт оплаты доносит парсер; если парсер завис — синхронный
* confirm уже прошёл. Повторного списания средств при любом сценарии не
* происходит — контракт billing отклоняет повтор `payment_hash` (anti-replay).
*
* Включается только на хабе (Восход, BILLING_HUB_MODE=true): на спицах
* BillingModule не подключается вовсе.
*
* ⚠️ Ограничение шины (как у BillingConversionListener): эмит `action::` идёт
* с задержкой ПОСЛЕ ACK сообщения, при падении хаба в этом окне реактивное
* уведомление теряется. Для pay это не критично: при следующем тике cron
* контракт ответит «уже проведён», и тик отправит подтверждение сам.
*/
@Injectable()
export class BillingPaymentListener {
private readonly logger = new Logger(BillingPaymentListener.name);
constructor(private readonly providerClient: BillingProviderClient) {}
@OnEvent(
`action::${BillingContract.contractName.production}::${BillingContract.Actions.Pay.actionName}`,
)
async onPay(action: ActionDomainInterface): Promise<void> {
// Defense-in-depth: BillingModule и так грузится лишь на хабе.
if (!config.billing.hub_mode) return;
if (!this.providerClient.isConfigured()) {
this.logger.warn('billing::pay: provider_base_url не сконфигурирован — пропуск');
return;
}
try {
const data: any = action.data ?? {};
const paymentHash = String(data.payment_hash ?? '');
const txId = String(action.transaction_id ?? '');
if (!paymentHash) {
this.logger.warn('billing::pay: пустой payment_hash в событии — пропуск');
return;
}
await this.providerClient.confirmPayment({
paymentHash,
blockchainTransactionId: txId,
});
} catch (err: any) {
// Не пробрасываем: on-chain состояние уже консистентно, подтверждение
// продублирует синхронный путь cron'а (см. docstring класса).
this.logger.error(`billing::pay → provider: ${err?.message}`, err?.stack);
}
}
}
@@ -27,6 +27,20 @@ export interface ProviderBillingSummary {
next_payment_due: string | null;
}
/**
* Invoice на батч-оплату подписок (Single-Hub v5, POST /billing/invoice).
* `payment_hash` провайдер считает детерминированно; повторный запрос с теми же
* позициями возвращает существующий PENDING-invoice. `status === 'PAID'` —
* платить нечего (оплата уже зафиксирована ранее).
*/
export interface ProviderBillingInvoice {
payment_hash: string;
coopname: string;
total_amount: number;
status: string;
expires_at: string;
}
/**
* HTTP-клиент к provider backend (Восход) для биллинга подписок
* (Epic 12/13, проект «Облачный провайдер»).
@@ -70,31 +84,46 @@ export class BillingProviderClient {
}
/**
* Подтверждение проведённого on-chain платежа: провайдер фиксирует факт и
* продлевает (`extend`) все подписки, входившие в `payment_hash`. Идемпотентно
* по `transaction_id = payment_hash` (Epic 3 / Story 12.5).
* Выписать invoice на батч-оплату подписок (Single-Hub v5).
*
* Вызывается ПЕРЕД on-chain `billing::pay`: провайдер фиксирует PENDING-invoice
* с TTL и отдаёт детерминированный `payment_hash`, который уходит в чейн.
* Без этого шага `POST /billing/payment-confirmed` не найдёт invoice и оплата
* не продлит подписки. Идемпотентно: повтор с теми же позициями возвращает
* существующий invoice.
*/
async confirmPayment(input: {
coopname: string;
paymentHash: string;
amount: number;
blockchainTransactionId: string;
periodDays?: number;
}): Promise<void> {
const url = `${this.baseUrl}/subscriptions/confirm-batch-payment`;
async createInvoice(
coopname: string,
items: Array<{ subscription_id: number; period_days: number }>,
): Promise<ProviderBillingInvoice> {
const url = `${this.baseUrl}/billing/invoice`;
const { data } = await axios.post<ProviderBillingInvoice>(
url,
{ coopname, items },
{ headers: this.headers(), timeout: 10_000 },
);
this.logger.log(`createInvoice ${coopname} payment_hash=${data.payment_hash} status=${data.status}`);
return data;
}
/**
* Подтверждение проведённого on-chain платежа (Single-Hub v5): провайдер
* переводит invoice в PAID и продлевает (`extend`) все входившие в него
* подписки. Идемпотентно по `payment_hash` (повтор → no-op), поэтому зовётся
* с двух сторон: синхронно из BillingCronService сразу после transact и
* реактивно из BillingPaymentListener по событию парсера.
*/
async confirmPayment(input: { paymentHash: string; blockchainTransactionId: string }): Promise<void> {
const url = `${this.baseUrl}/billing/payment-confirmed`;
await axios.post(
url,
{
coopname: input.coopname,
transaction_id: input.paymentHash,
payment_hash: input.paymentHash,
amount: input.amount,
blockchain_transaction_id: input.blockchainTransactionId,
period_days: input.periodDays ?? 30,
tx_id: input.blockchainTransactionId,
},
{ headers: this.headers() },
{ headers: this.headers(), timeout: 10_000 },
);
this.logger.log(`confirmPayment ${input.coopname} payment_hash=${input.paymentHash}`);
this.logger.log(`confirmPayment payment_hash=${input.paymentHash} tx=${input.blockchainTransactionId}`);
}
/**
@@ -1,6 +1,7 @@
import { Inject, Injectable, Logger } from '@nestjs/common';
import httpStatus from 'http-status';
import { BillingContract } from 'cooptypes';
import config from '~/config/config';
import { TransactResult } from '@wharfkit/session';
import { BlockchainService } from '../blockchain.service';
import { VAULT_DOMAIN_SERVICE, VaultDomainService } from '~/domain/vault/services/vault-domain.service';
@@ -16,7 +17,12 @@ import { DomainToBlockchainUtils } from '~/shared/utils/domain-to-blockchain.uti
/**
* Блокчейн-адаптер billing (Epic 12) — оплата подписок членскими взносами.
*
* Подпись `coopname@active` через `BlockchainService.transact` с WIF из vault.
* Подпись через `BlockchainService.transact` с WIF из vault:
* - `convert` — релей подписи кооператива (`coopname@active`) после JWT пайщика;
* - `pay` — ОПЕРАТОР платформы (аккаунт узла-хаба `config.coopname`, на Восходе
* = `_provider` контрактов): рекуррентные списания авторизует он, ключей
* кооперативов-спиц в vault хаба нет.
*
* Имена действий и payload — из cooptypes (`BillingContract.Actions.{Convert,Pay}`),
* без сырых строк. Состав/цены подписок on-chain не передаются — только сумма,
* payment_hash и memo.
@@ -39,6 +45,17 @@ export class BillingBlockchainAdapter implements BillingBlockchainPort {
this.blockchainService.initialize(coopname, wif);
}
/** Подпись оператором платформы — аккаунтом узла-хаба (см. docstring класса). */
private async initForOperator(): Promise<string> {
const operator = config.coopname;
const wif = await this.vaultDomainService.getWif(operator);
if (!wif) {
throw new HttpApiError(httpStatus.BAD_GATEWAY, 'Не найден приватный ключ оператора для подписания биллинг-операции');
}
this.blockchainService.initialize(operator, wif);
return operator;
}
async convert(data: BillingConvertBlockchainDomainInterface): Promise<TransactionResult> {
await this.initForCoop(data.coopname);
const formattedQuantity = this.domainToBlockchainUtils.formatQuantityWithPrecision(data.quantity);
@@ -62,7 +79,7 @@ export class BillingBlockchainAdapter implements BillingBlockchainPort {
}
async pay(data: BillingPayBlockchainDomainInterface): Promise<TransactionResult> {
await this.initForCoop(data.coopname);
const operator = await this.initForOperator();
const formattedQuantity = this.domainToBlockchainUtils.formatQuantityWithPrecision(data.quantity);
const payload: BillingContract.Actions.Pay.IPay = {
@@ -76,7 +93,7 @@ export class BillingBlockchainAdapter implements BillingBlockchainPort {
const result = (await this.blockchainService.transact({
account: BillingContract.contractName.production,
name: BillingContract.Actions.Pay.actionName,
authorization: [{ actor: data.coopname, permission: 'active' }],
authorization: [{ actor: operator, permission: 'active' }],
data: payload,
})) as TransactResult;