import { and, desc, eq, sql } from "drizzle-orm"; import { db } from "../db"; import { billingPayments, billingTransactions, organizationBilling, organizations } from "@shared/schema"; import type { BillingPayment } from "@shared/schema"; import { logger } from "../utils/logger"; // Сервис платежей биллинга: создание, выборки и зачисление успешных платежей. const log = logger("billing"); export interface CreateBillingPaymentParams { organizationId: number; provider: string; externalId: string; amount: number; currency: string; idempotencyKey: string; metadata?: Record; } // Создание платежа с защитой от дублей: повторный вызов с тем же // idempotency_key возвращает существующую запись, а не создаёт новую. export async function createBillingPayment(params: CreateBillingPaymentParams): Promise { const [row] = await db .insert(billingPayments) .values({ organizationId: params.organizationId, provider: params.provider, externalId: params.externalId, amount: params.amount.toFixed(2), currency: params.currency, status: 'pending', idempotencyKey: params.idempotencyKey, metadata: params.metadata ?? null, }) .onConflictDoNothing({ target: billingPayments.idempotencyKey }) .returning(); if (row) return row; const [existing] = await db .select() .from(billingPayments) .where(eq(billingPayments.idempotencyKey, params.idempotencyKey)) .limit(1); return existing; } export async function listBillingPayments(organizationId: number, limit = 50): Promise { return db .select() .from(billingPayments) .where(eq(billingPayments.organizationId, organizationId)) .orderBy(desc(billingPayments.createdAt)) .limit(limit); } export async function getBillingPaymentById(id: number, organizationId?: number): Promise { const conditions = organizationId !== undefined ? and(eq(billingPayments.id, id), eq(billingPayments.organizationId, organizationId)) : eq(billingPayments.id, id); const [row] = await db.select().from(billingPayments).where(conditions).limit(1); return row ?? null; } export async function getBillingPaymentByExternalId(externalId: string): Promise { const [row] = await db .select() .from(billingPayments) .where(eq(billingPayments.externalId, externalId)) .limit(1); return row ?? null; } export async function markBillingPaymentCanceled(id: number): Promise { const [row] = await db .update(billingPayments) .set({ status: 'canceled', updatedAt: new Date() }) .where(and(eq(billingPayments.id, id), eq(billingPayments.status, 'pending'))) .returning(); return row ?? null; } // Зачисление успешного платежа в одной транзакции: // 1) атомарный перевод pending→succeeded (guard от гонок и повторных webhook'ов); // 2) credit-запись в billing_transactions; // 3) инкремент organization_billing.balance; // 4) авто-снятие billingBlocked при положительном балансе // (та же пара обновлений, что и storage.setBillingBlocked(orgId, false)). // Возвращает applied=false, если платёж уже обработан — повторное зачисление невозможно. export async function applySucceededPayment( paymentId: number, ): Promise<{ applied: boolean; payment: BillingPayment | null }> { return db.transaction(async (tx) => { const [payment] = await tx .update(billingPayments) .set({ status: 'succeeded', updatedAt: new Date() }) .where(and(eq(billingPayments.id, paymentId), eq(billingPayments.status, 'pending'))) .returning(); if (!payment) { const [current] = await tx .select() .from(billingPayments) .where(eq(billingPayments.id, paymentId)) .limit(1); // Повторный webhook/гонка: платёж уже не в pending — зачисление пропускаем log.warn(`Повторное зачисление платежа пропущено (paymentId=${paymentId}, status=${current?.status ?? 'not found'})`); return { applied: false, payment: current ?? null }; } const orgId = payment.organizationId; // Гарантируем наличие строки биллинга организации перед инкрементом. await tx .insert(organizationBilling) .values({ organizationId: orgId }) .onConflictDoNothing({ target: organizationBilling.organizationId }); const [billing] = await tx .update(organizationBilling) .set({ balance: sql`${organizationBilling.balance} + ${payment.amount}::numeric`, updatedAt: new Date(), }) .where(eq(organizationBilling.organizationId, orgId)) .returning(); await tx.insert(billingTransactions).values({ organizationId: orgId, amount: payment.amount, type: 'credit', description: `Пополнение баланса (платёж ${payment.provider} ${payment.externalId})`, createdBy: null, }); const newBalance = parseFloat(billing?.balance ?? '0'); if (newBalance > 0) { await tx .update(organizations) .set({ billingBlocked: false, updatedAt: new Date() }) .where(and(eq(organizations.id, orgId), eq(organizations.billingBlocked, true))); await tx .update(organizationBilling) .set({ blockedAt: null, updatedAt: new Date() }) .where(eq(organizationBilling.organizationId, orgId)); } return { applied: true, payment }; }); }