Files
qipai/backend/src/payments/payment-repository.ts
T

374 lines
14 KiB
TypeScript

import { randomBytes } from 'node:crypto';
import type { PoolConnection, ResultSetHeader, RowDataPacket } from 'mysql2/promise';
import type { MySqlPool } from '../db/mysql.js';
import type { MarketingBenefitService } from '../wallets/marketing-benefit-service.js';
import type { WalletLedgerService } from '../wallets/wallet-ledger-service.js';
export type PaymentProvider = 'WECHAT' | 'BALANCE' | 'PACKAGE' | 'GROUP_BUY' | 'TEST';
interface OrderRow extends RowDataPacket {
id: string;
storeId: string;
status: string;
totalAmountCents: number;
paidAmountCents: number;
}
interface PaymentRow extends RowDataPacket {
id: string;
orderId: string;
paymentNo: string;
provider: PaymentProvider;
status: string;
amountCents: number;
}
interface ConfigRow extends RowDataPacket {
id: string;
credentialRef: string;
settings: string | object;
scopeKey: string;
}
export class PaymentError extends Error {
constructor(public readonly code: string) { super(code); }
}
export class PaymentRepository {
constructor(
private readonly pool: MySqlPool,
private readonly wallet?: Pick<WalletLedgerService, 'debitInTransaction'>,
private readonly benefits?: Pick<MarketingBenefitService, 'confirmReservedInTransaction'>
) {}
async createPayment(input: {
tenantId: string;
platformAppId: string;
userId: string;
orderId: string;
provider: PaymentProvider;
clientRequestId: string;
testAdapterEnabled: boolean;
}) {
return this.transaction(async (connection) => {
const order = await this.loadOwnedOrder(
connection, input.tenantId, input.userId, input.orderId, true
);
if (!['PENDING_PAYMENT', 'PAID', 'RESERVED', 'IN_PROGRESS'].includes(order.status)) {
throw new PaymentError('PAYMENT_ORDER_STATUS_INVALID');
}
const amountCents = Number(order.totalAmountCents) - Number(order.paidAmountCents);
if (amountCents <= 0) throw new PaymentError('PAYMENT_NOT_REQUIRED');
if (input.provider === 'TEST' && !input.testAdapterEnabled) {
throw new PaymentError('TEST_PAYMENT_DISABLED');
}
if (!['TEST', 'BALANCE'].includes(input.provider)) {
await this.resolveConfig(
connection, input.tenantId, input.platformAppId, order.storeId, input.provider
);
}
const [existing] = await connection.execute<PaymentRow[]>(
`SELECT id, order_id AS orderId, payment_no AS paymentNo, provider,
status, amount_cents AS amountCents
FROM qipai_payments
WHERE tenant_id = ? AND client_request_id = ? LIMIT 1`,
[input.tenantId, input.clientRequestId]
);
if (existing[0]) {
if (String(existing[0].orderId) !== input.orderId
|| existing[0].provider !== input.provider
|| Number(existing[0].amountCents) !== amountCents) {
throw new PaymentError('PAYMENT_IDEMPOTENCY_CONFLICT');
}
return this.paymentResponse(existing[0], true, input.testAdapterEnabled);
}
const paymentNo = `PAY${Date.now()}${randomBytes(5).toString('hex').toUpperCase()}`;
const [result] = await connection.execute<ResultSetHeader>(
`INSERT INTO qipai_payments
(tenant_id, platform_app_id, order_id, store_id, payment_no,
channel, provider, client_request_id, status, amount_cents)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'PENDING', ?)`,
[input.tenantId, input.platformAppId, input.orderId, order.storeId,
paymentNo, input.provider, input.provider, input.clientRequestId, amountCents]
);
const paymentId = String(result.insertId);
await connection.execute(
`INSERT INTO qipai_payment_attempts
(tenant_id, payment_id, attempt_no, status, request_payload, response_payload,
completed_at)
VALUES (?, ?, 1, 'CREATED',
JSON_OBJECT('provider', ?, 'amountCents', ?),
JSON_OBJECT('adapter', ?, 'credentialExposed', FALSE), UTC_TIMESTAMP(3))`,
[input.tenantId, paymentId, input.provider, amountCents,
input.provider === 'TEST' ? 'test' : 'configured']
);
if (input.provider === 'BALANCE') {
if (!this.wallet) throw new PaymentError('WALLET_SETTLEMENT_NOT_CONFIGURED');
await this.wallet.debitInTransaction(connection, {
tenantId: input.tenantId,
userId: input.userId,
scopeType: 'STORE',
storeId: order.storeId,
businessType: 'ORDER_PAYMENT',
businessId: input.orderId,
entryType: 'CONSUME',
amountCents,
traceId: input.clientRequestId,
note: 'Customer balance payment',
metadata: { paymentId, provider: input.provider }
});
await this.applyPaymentSuccess(connection, {
paymentId,
orderId: input.orderId,
tenantId: input.tenantId,
userId: input.userId,
amountCents,
providerPaymentId: `BAL-${paymentNo}`,
traceId: input.clientRequestId,
reason: 'Balance payment completed'
});
return this.paymentResponse({
id: paymentId, orderId: input.orderId, paymentNo, provider: input.provider,
status: 'SUCCEEDED', amountCents
} as PaymentRow, false, input.testAdapterEnabled);
}
return this.paymentResponse({
id: paymentId, orderId: input.orderId, paymentNo, provider: input.provider,
status: 'PENDING', amountCents
} as PaymentRow, false, input.testAdapterEnabled);
});
}
async processTestCallback(input: {
tenantId: string;
userId: string;
paymentId: string;
callbackId: string;
amountCents: number;
testAdapterEnabled: boolean;
traceId: string;
}) {
if (!input.testAdapterEnabled) throw new PaymentError('TEST_PAYMENT_DISABLED');
return this.transaction(async (connection) => {
const payment = await this.loadPayment(connection, input.tenantId, input.paymentId, true);
await this.assertOrderOwner(connection, input.tenantId, payment.orderId, input.userId);
if (payment.provider !== 'TEST') throw new PaymentError('PAYMENT_PROVIDER_INVALID');
const [callbackResult] = await connection.execute<ResultSetHeader>(
`INSERT IGNORE INTO qipai_payment_callbacks
(tenant_id, payment_id, provider, callback_id, callback_type,
verified, payload)
VALUES (?, ?, 'TEST', ?, 'PAYMENT_SUCCEEDED', 1,
JSON_OBJECT('amountCents', ?))`,
[input.tenantId, input.paymentId, input.callbackId, input.amountCents]
);
if (callbackResult.affectedRows === 0) {
return { paymentId: input.paymentId, status: payment.status, idempotent: true };
}
if (Number(payment.amountCents) !== input.amountCents) {
await connection.execute(
`UPDATE qipai_payment_callbacks
SET processing_status = 'REJECTED', error_code = 'PAYMENT_AMOUNT_MISMATCH',
processed_at = UTC_TIMESTAMP(3)
WHERE tenant_id = ? AND provider = 'TEST' AND callback_id = ?`,
[input.tenantId, input.callbackId]
);
return {
paymentId: input.paymentId,
status: 'REJECTED',
code: 'PAYMENT_AMOUNT_MISMATCH',
idempotent: false
};
}
if (payment.status === 'SUCCEEDED') {
await connection.execute(
`UPDATE qipai_payment_callbacks
SET processing_status = 'DUPLICATE', processed_at = UTC_TIMESTAMP(3)
WHERE tenant_id = ? AND provider = 'TEST' AND callback_id = ?`,
[input.tenantId, input.callbackId]
);
return { paymentId: input.paymentId, status: 'SUCCEEDED', idempotent: true };
}
await this.applyPaymentSuccess(connection, {
paymentId: input.paymentId,
orderId: payment.orderId,
tenantId: input.tenantId,
userId: input.userId,
amountCents: Number(payment.amountCents),
providerPaymentId: `TEST-${input.callbackId}`,
traceId: input.traceId,
reason: 'Verified payment callback'
});
await connection.execute(
`UPDATE qipai_payment_callbacks
SET processing_status = 'PROCESSED', processed_at = UTC_TIMESTAMP(3)
WHERE tenant_id = ? AND provider = 'TEST' AND callback_id = ?`,
[input.tenantId, input.callbackId]
);
return { paymentId: input.paymentId, status: 'SUCCEEDED', idempotent: false };
});
}
async resolveConfig(
connection: Pick<MySqlPool, 'execute'>,
tenantId: string,
platformAppId: string,
storeId: string,
provider: PaymentProvider
) {
const [rows] = await connection.execute<ConfigRow[]>(
`SELECT id, credential_ref AS credentialRef, settings, scope_key AS scopeKey
FROM qipai_payment_configs
WHERE provider = ? AND enabled = 1
AND (tenant_id IS NULL OR tenant_id = ?)
AND (platform_app_id IS NULL OR platform_app_id = ?)
AND (store_id IS NULL OR store_id = ?)
ORDER BY (store_id IS NOT NULL) DESC,
(tenant_id IS NOT NULL) DESC,
(platform_app_id IS NOT NULL) DESC, id DESC LIMIT 1`,
[provider, tenantId, platformAppId, storeId]
);
if (!rows[0]) throw new PaymentError('PAYMENT_CONFIG_NOT_FOUND');
return {
id: String(rows[0].id),
credentialRef: rows[0].credentialRef,
scopeKey: rows[0].scopeKey,
settings: typeof rows[0].settings === 'string'
? JSON.parse(rows[0].settings) : rows[0].settings
};
}
private paymentResponse(row: PaymentRow, idempotent: boolean, testEnabled: boolean) {
return {
paymentId: String(row.id),
orderId: String(row.orderId),
paymentNo: row.paymentNo,
provider: row.provider,
status: row.status,
amountCents: Number(row.amountCents),
idempotent,
testCompletionAvailable: row.provider === 'TEST' && testEnabled
};
}
private async loadOwnedOrder(
connection: PoolConnection, tenantId: string, userId: string,
orderId: string, lock: boolean
) {
await this.assertOrderOwner(connection, tenantId, orderId, userId);
return this.loadOrder(connection, tenantId, orderId, lock);
}
private async loadOrder(
connection: PoolConnection, tenantId: string, orderId: string, lock: boolean
) {
const [rows] = await connection.execute<Array<OrderRow & { statusVersion: number }>>(
`SELECT id, store_id AS storeId, status, status_version AS statusVersion,
total_amount_cents AS totalAmountCents, paid_amount_cents AS paidAmountCents
FROM qipai_orders WHERE tenant_id = ? AND id = ? AND deleted_at IS NULL
${lock ? 'FOR UPDATE' : ''}`,
[tenantId, orderId]
);
if (!rows[0]) throw new PaymentError('ORDER_NOT_FOUND');
return rows[0];
}
private async loadPayment(
connection: PoolConnection, tenantId: string, paymentId: string, lock: boolean
) {
const [rows] = await connection.execute<PaymentRow[]>(
`SELECT id, order_id AS orderId, payment_no AS paymentNo, provider,
status, amount_cents AS amountCents
FROM qipai_payments WHERE tenant_id = ? AND id = ? AND deleted_at IS NULL
${lock ? 'FOR UPDATE' : ''}`,
[tenantId, paymentId]
);
if (!rows[0]) throw new PaymentError('PAYMENT_NOT_FOUND');
return rows[0];
}
private async assertOrderOwner(
connection: PoolConnection, tenantId: string, orderId: string, userId: string
) {
const [rows] = await connection.execute<RowDataPacket[]>(
`SELECT 1 FROM qipai_order_user_access
WHERE tenant_id = ? AND order_id = ? AND user_id = ?`,
[tenantId, orderId, userId]
);
if (!rows[0]) throw new PaymentError('ORDER_ACCESS_FORBIDDEN');
}
async applyPaymentSuccess(
connection: PoolConnection,
input: {
paymentId: string;
orderId: string;
tenantId: string;
userId?: string;
amountCents: number;
providerPaymentId: string;
traceId: string;
reason: string;
}
) {
await connection.execute(
`UPDATE qipai_payments
SET status = 'SUCCEEDED', provider_payment_id = ?,
paid_at = UTC_TIMESTAMP(3), raw_notify = JSON_OBJECT('verified', TRUE)
WHERE tenant_id = ? AND id = ? AND status = 'PENDING'`,
[input.providerPaymentId, input.tenantId, input.paymentId]
);
await connection.execute(
`UPDATE qipai_orders
SET paid_amount_cents = paid_amount_cents + ?
WHERE tenant_id = ? AND id = ?`,
[input.amountCents, input.tenantId, input.orderId]
);
const order = await this.loadOrder(connection, input.tenantId, input.orderId, true);
if (Number(order.paidAmountCents) >= Number(order.totalAmountCents)
&& order.status === 'PENDING_PAYMENT') {
if (this.benefits && input.userId) {
await this.benefits.confirmReservedInTransaction(connection, {
tenantId: input.tenantId,
userId: input.userId,
orderId: input.orderId,
traceId: input.traceId
});
}
const nextVersion = Number((order as OrderRow & { statusVersion?: number }).statusVersion ?? 1) + 1;
await connection.execute(
`UPDATE qipai_orders SET status = 'PAID', status_version = ?,
status_updated_at = UTC_TIMESTAMP(3)
WHERE tenant_id = ? AND id = ?`,
[nextVersion, input.tenantId, input.orderId]
);
await connection.execute(
`UPDATE qipai_room_reservations
SET status = 'CONSUMED', expires_at = GREATEST(expires_at, ends_at)
WHERE tenant_id = ? AND order_id = ? AND status = 'HELD'`,
[input.tenantId, input.orderId]
);
await connection.execute(
`INSERT INTO qipai_order_status_history
(tenant_id, order_id, from_status, to_status, action, actor_type,
actor_id, source, reason, trace_id, metadata)
VALUES (?, ?, 'PENDING_PAYMENT', 'PAID', 'CONFIRM_PAYMENT', 'SYSTEM',
NULL, 'PAYMENT', ?, ?, JSON_OBJECT('statusVersion', ?, 'paymentId', ?))`,
[input.tenantId, input.orderId, input.reason, input.traceId, nextVersion, input.paymentId]
);
}
}
private async transaction<T>(work: (connection: PoolConnection) => Promise<T>) {
const connection = await this.pool.getConnection();
try {
await connection.beginTransaction();
const result = await work(connection);
await connection.commit();
return result;
} catch (error) {
await connection.rollback();
throw error;
} finally {
connection.release();
}
}
}