forked from wangziqi/gongxue-base
feat: add provider bill download reconciliation jobs
This commit is contained in:
@@ -10,6 +10,7 @@
|
||||
"check": "tsc -p tsconfig.json --noEmit",
|
||||
"crm:once": "tsx src/index.ts --once --job crm",
|
||||
"commerce:once": "tsx src/index.ts --once --job commerce",
|
||||
"provider-bills:once": "tsx src/index.ts --once --job provider-bills",
|
||||
"assets:once": "tsx src/index.ts --once --job assets",
|
||||
"imports:once": "tsx src/index.ts --once --job imports",
|
||||
"public-banks:once": "tsx src/index.ts --once --job public-banks",
|
||||
@@ -20,6 +21,7 @@
|
||||
"@supabase/storage-js": "^2.108.2",
|
||||
"ali-oss": "^6.23.0",
|
||||
"docx": "^9.7.1",
|
||||
"iconv-lite": "^0.6.3",
|
||||
"jszip": "^3.10.1",
|
||||
"pdfkit": "^0.19.1",
|
||||
"pg": "^8.16.3"
|
||||
|
||||
@@ -13,6 +13,9 @@ export interface WorkerConfig {
|
||||
commerceBatchSize: number;
|
||||
commerceMinAgeSeconds: number;
|
||||
commerceRequestTimeoutMs: number;
|
||||
providerBillBatchSize: number;
|
||||
providerBillWorkerId: string;
|
||||
providerBillClaimStaleSeconds: number;
|
||||
assetBatchSize: number;
|
||||
assetMinAgeSeconds: number;
|
||||
assetRecheckIntervalSeconds: number;
|
||||
@@ -69,6 +72,9 @@ export const config: WorkerConfig = {
|
||||
commerceBatchSize: envNumber('WORKER_COMMERCE_BATCH_SIZE', 20),
|
||||
commerceMinAgeSeconds: envNumber('WORKER_COMMERCE_MIN_AGE_SECONDS', 300),
|
||||
commerceRequestTimeoutMs: envNumber('WORKER_COMMERCE_REQUEST_TIMEOUT_MS', 10_000),
|
||||
providerBillBatchSize: envNumber('WORKER_PROVIDER_BILL_BATCH_SIZE', 5),
|
||||
providerBillWorkerId: envString('WORKER_PROVIDER_BILL_ID', `provider-bills-${process.pid}`),
|
||||
providerBillClaimStaleSeconds: envNumber('WORKER_PROVIDER_BILL_CLAIM_STALE_SECONDS', 15 * 60),
|
||||
assetBatchSize: envNumber('WORKER_ASSET_BATCH_SIZE', 50),
|
||||
assetMinAgeSeconds: envNumber('WORKER_ASSET_MIN_AGE_SECONDS', 300),
|
||||
assetRecheckIntervalSeconds: envNumber('WORKER_ASSET_RECHECK_INTERVAL_SECONDS', 60 * 60 * 24),
|
||||
|
||||
@@ -31,6 +31,17 @@ async function runOnce() {
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (job === 'provider-bills') {
|
||||
const { closePool: closeApiPool } = await import('../../api/src/core/db.js');
|
||||
const { processProviderBillBatch } = await import('./jobs/provider-bills.js');
|
||||
extraClosers.add(closeApiPool);
|
||||
const result = await processProviderBillBatch();
|
||||
console.log(
|
||||
`[worker] provider-bills batch processed=${result.processed}`
|
||||
+ ` completed=${result.completed} failed=${result.failed} skipped=${result.skipped}`,
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (job === 'assets') {
|
||||
const result = await processAssetBatch();
|
||||
console.log(
|
||||
|
||||
784
apps/worker/src/jobs/provider-bills.ts
Normal file
784
apps/worker/src/jobs/provider-bills.ts
Normal file
@@ -0,0 +1,784 @@
|
||||
import crypto from 'node:crypto';
|
||||
import JSZip from 'jszip';
|
||||
import iconv from 'iconv-lite';
|
||||
import type pg from 'pg';
|
||||
import { pool } from '../db.js';
|
||||
import { config } from '../config.js';
|
||||
import {
|
||||
importReconciliationBatch,
|
||||
normalizeReconciliationBillDate,
|
||||
normalizeReconciliationBillType,
|
||||
normalizeReconciliationProvider,
|
||||
type BillType,
|
||||
type ReconciliationProvider,
|
||||
} from '../../../api/src/features/commerce/reconciliation.js';
|
||||
|
||||
type ProviderBillProvider = Exclude<ReconciliationProvider, 'manual'>;
|
||||
|
||||
interface SecretRow {
|
||||
secretValue: string | null;
|
||||
secretJson: Record<string, unknown> | null;
|
||||
}
|
||||
|
||||
interface ProviderConfig {
|
||||
tenantId: string;
|
||||
provider: ProviderBillProvider;
|
||||
configPublic: Record<string, unknown>;
|
||||
secret: SecretRow | null;
|
||||
}
|
||||
|
||||
interface ProviderBillJob {
|
||||
id: string;
|
||||
tenantId: string;
|
||||
provider: ProviderBillProvider;
|
||||
billDate: string;
|
||||
billType: BillType;
|
||||
sourceName: string;
|
||||
requestedBy: string | null;
|
||||
metadata: Record<string, unknown>;
|
||||
}
|
||||
|
||||
interface DownloadedBill {
|
||||
rows: Record<string, unknown>[];
|
||||
sourceHash: string;
|
||||
downloadHashType: string | null;
|
||||
downloadHashValue: string | null;
|
||||
downloadUrlHost: string | null;
|
||||
metadata: Record<string, unknown>;
|
||||
}
|
||||
|
||||
interface ProviderBillWorkerResult {
|
||||
processed: number;
|
||||
completed: number;
|
||||
failed: number;
|
||||
skipped: number;
|
||||
}
|
||||
|
||||
class ProviderBillWorkerError extends Error {
|
||||
readonly code: string;
|
||||
|
||||
constructor(message: string, code = 'PROVIDER_BILL_WORKER_ERROR') {
|
||||
super(message);
|
||||
this.code = code;
|
||||
}
|
||||
}
|
||||
|
||||
function objectValue(value: unknown): Record<string, unknown> {
|
||||
return value && typeof value === 'object' && !Array.isArray(value) ? value as Record<string, unknown> : {};
|
||||
}
|
||||
|
||||
function safeString(value: unknown) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : '';
|
||||
}
|
||||
|
||||
function optionalPublicString(providerConfig: ProviderConfig, keys: string[]) {
|
||||
for (const key of keys) {
|
||||
const value = providerConfig.configPublic[key];
|
||||
if (typeof value === 'string' && value.trim()) return value.trim();
|
||||
}
|
||||
return '';
|
||||
}
|
||||
|
||||
function requirePublicString(providerConfig: ProviderConfig, keys: string[]) {
|
||||
const value = optionalPublicString(providerConfig, keys);
|
||||
if (value) return value;
|
||||
throw new ProviderBillWorkerError(`${providerConfig.provider} public config is missing ${keys[0]}`, 'PAYMENT_PUBLIC_CONFIG_REQUIRED');
|
||||
}
|
||||
|
||||
function secretString(providerConfig: ProviderConfig, keys: string[]) {
|
||||
const secret = providerConfig.secret;
|
||||
if (!secret) return '';
|
||||
if (typeof secret.secretValue === 'string' && secret.secretValue.trim()) return secret.secretValue.trim();
|
||||
const json = objectValue(secret.secretJson);
|
||||
for (const key of keys) {
|
||||
const value = json[key];
|
||||
if (typeof value === 'string' && value.trim()) return value.trim();
|
||||
}
|
||||
return '';
|
||||
}
|
||||
|
||||
function requireSecretString(providerConfig: ProviderConfig, keys: string[]) {
|
||||
const value = secretString(providerConfig, keys);
|
||||
if (value) return value;
|
||||
throw new ProviderBillWorkerError(`${providerConfig.provider} secret is missing ${keys[0]}`, 'PAYMENT_SECRET_REQUIRED');
|
||||
}
|
||||
|
||||
function normalizePem(value: string, label: 'PRIVATE KEY' | 'PUBLIC KEY') {
|
||||
if (value.includes('-----BEGIN')) return value;
|
||||
const wrapped = value.match(/.{1,64}/g)?.join('\n') || value;
|
||||
return `-----BEGIN ${label}-----\n${wrapped}\n-----END ${label}-----`;
|
||||
}
|
||||
|
||||
function rsaSignSha256(privateKey: string, message: string) {
|
||||
return crypto.createSign('RSA-SHA256').update(message, 'utf8').sign(privateKey, 'base64');
|
||||
}
|
||||
|
||||
function randomNonce(size = 16) {
|
||||
return crypto.randomBytes(size).toString('base64url');
|
||||
}
|
||||
|
||||
function canonicalForm(params: Record<string, string>) {
|
||||
return Object.keys(params)
|
||||
.filter(key => params[key] !== undefined && params[key] !== null && params[key] !== '')
|
||||
.sort()
|
||||
.map(key => `${key}=${params[key]}`)
|
||||
.join('&');
|
||||
}
|
||||
|
||||
function encodedForm(params: Record<string, string>) {
|
||||
return Object.keys(params)
|
||||
.filter(key => params[key] !== undefined && params[key] !== null && params[key] !== '')
|
||||
.sort()
|
||||
.map(key => `${encodeURIComponent(key)}=${encodeURIComponent(params[key])}`)
|
||||
.join('&');
|
||||
}
|
||||
|
||||
function parseSecretRef(value: unknown) {
|
||||
if (typeof value !== 'string' || !value.trim()) return null;
|
||||
const parts = value.trim().split(':');
|
||||
if (parts.length !== 3 || parts[0] !== 'app_private.tenant_secrets') return null;
|
||||
if (parts[1] !== 'payment' || !parts[2]) return null;
|
||||
return { scope: parts[1], key: parts[2] };
|
||||
}
|
||||
|
||||
function normalizeProvider(value: string | null | undefined): ProviderBillProvider | null {
|
||||
const normalized = normalizeReconciliationProvider(String(value || ''));
|
||||
return normalized === 'manual' ? null : normalized;
|
||||
}
|
||||
|
||||
function providerAliases(provider: ProviderBillProvider) {
|
||||
return provider === 'wechat_pay'
|
||||
? ['wechat_pay', 'wechat-pay', 'wechatpay', 'wxpay']
|
||||
: ['alipay', 'ali_pay'];
|
||||
}
|
||||
|
||||
function isLocalDevHost(hostname: string) {
|
||||
return ['127.0.0.1', 'localhost', '::1'].includes(hostname);
|
||||
}
|
||||
|
||||
function isAllowedHost(hostname: string, allowedHosts: string[]) {
|
||||
return allowedHosts.some(host => hostname === host || hostname.endsWith(`.${host}`));
|
||||
}
|
||||
|
||||
function providerEndpointForKeys(
|
||||
providerConfig: ProviderConfig,
|
||||
keys: string[],
|
||||
fallback: string,
|
||||
allowedHosts: string[],
|
||||
code = 'PAYMENT_ENDPOINT_NOT_ALLOWED',
|
||||
) {
|
||||
const raw = optionalPublicString(providerConfig, keys) || fallback;
|
||||
let endpoint: URL;
|
||||
try {
|
||||
endpoint = new URL(raw);
|
||||
} catch {
|
||||
throw new ProviderBillWorkerError(`${providerConfig.provider} endpoint is invalid`, 'PAYMENT_ENDPOINT_INVALID');
|
||||
}
|
||||
|
||||
const localDev = process.env.NODE_ENV !== 'production' && isLocalDevHost(endpoint.hostname);
|
||||
if (!localDev && endpoint.protocol !== 'https:') {
|
||||
throw new ProviderBillWorkerError(`${providerConfig.provider} endpoint must use HTTPS`, code);
|
||||
}
|
||||
if (!localDev && !isAllowedHost(endpoint.hostname, allowedHosts)) {
|
||||
throw new ProviderBillWorkerError(`${providerConfig.provider} endpoint host is not allowed`, code);
|
||||
}
|
||||
endpoint.username = '';
|
||||
endpoint.password = '';
|
||||
return endpoint.toString();
|
||||
}
|
||||
|
||||
function assertDownloadUrlAllowed(url: string, provider: ProviderBillProvider) {
|
||||
let parsed: URL;
|
||||
try {
|
||||
parsed = new URL(url);
|
||||
} catch {
|
||||
throw new ProviderBillWorkerError('Provider bill download URL is invalid', 'PROVIDER_BILL_DOWNLOAD_URL_INVALID');
|
||||
}
|
||||
const localDev = process.env.NODE_ENV !== 'production' && isLocalDevHost(parsed.hostname);
|
||||
if (localDev) return parsed;
|
||||
if (parsed.protocol !== 'https:') {
|
||||
throw new ProviderBillWorkerError('Provider bill download URL must use HTTPS', 'PROVIDER_BILL_DOWNLOAD_URL_NOT_ALLOWED');
|
||||
}
|
||||
const hosts = provider === 'wechat_pay'
|
||||
? ['api.mch.weixin.qq.com', 'download.mch.weixin.qq.com']
|
||||
: ['openapi.alipay.com', 'dwbillcenter.alipay.com', 'download.alipay.com'];
|
||||
if (!isAllowedHost(parsed.hostname, hosts)) {
|
||||
throw new ProviderBillWorkerError('Provider bill download URL host is not allowed', 'PROVIDER_BILL_DOWNLOAD_URL_NOT_ALLOWED');
|
||||
}
|
||||
return parsed;
|
||||
}
|
||||
|
||||
async function loadSecret(client: pg.PoolClient, tenantId: string, configPublic: Record<string, unknown>, provider: ProviderBillProvider) {
|
||||
const secretRef = parseSecretRef(configPublic.secretRef);
|
||||
const secretKey = secretRef?.key || provider;
|
||||
const result = await client.query<SecretRow>(
|
||||
`
|
||||
select secret_value as "secretValue", secret_json as "secretJson"
|
||||
from app_private.tenant_secrets
|
||||
where tenant_id = $1
|
||||
and secret_scope = 'payment'
|
||||
and secret_key = $2
|
||||
limit 1
|
||||
`,
|
||||
[tenantId, secretKey],
|
||||
);
|
||||
return result.rows[0] || null;
|
||||
}
|
||||
|
||||
async function loadPaymentProviderConfig(client: pg.PoolClient, tenantId: string, provider: ProviderBillProvider): Promise<ProviderConfig> {
|
||||
const aliases = providerAliases(provider);
|
||||
const accountResult = await client.query<{
|
||||
provider: string;
|
||||
configPublic: Record<string, unknown> | null;
|
||||
}>(
|
||||
`
|
||||
select provider, config_public as "configPublic"
|
||||
from public.tenant_payment_accounts
|
||||
where tenant_id = $1
|
||||
and provider = any($2::text[])
|
||||
and status = 'active'
|
||||
order by array_position($2::text[], provider)
|
||||
limit 1
|
||||
`,
|
||||
[tenantId, aliases],
|
||||
);
|
||||
const account = accountResult.rows[0];
|
||||
if (!account) throw new ProviderBillWorkerError(`${provider} provider is not configured`, 'PAYMENT_PROVIDER_NOT_CONFIGURED');
|
||||
const normalized = normalizeProvider(account.provider);
|
||||
if (!normalized) throw new ProviderBillWorkerError(`${account.provider} provider is unsupported`, 'PAYMENT_PROVIDER_NOT_SUPPORTED');
|
||||
const configPublic = objectValue(account.configPublic);
|
||||
const secret = await loadSecret(client, tenantId, configPublic, normalized);
|
||||
return { tenantId, provider: normalized, configPublic, secret };
|
||||
}
|
||||
|
||||
async function fetchJson(url: string, init: RequestInit, timeoutMs: number) {
|
||||
const controller = new AbortController();
|
||||
const timeout = setTimeout(() => controller.abort(), timeoutMs);
|
||||
try {
|
||||
const response = await fetch(url, { ...init, signal: controller.signal });
|
||||
const raw = objectValue(await response.json().catch(() => ({})));
|
||||
return { response, raw };
|
||||
} finally {
|
||||
clearTimeout(timeout);
|
||||
}
|
||||
}
|
||||
|
||||
async function fetchBuffer(url: string, init: RequestInit, timeoutMs: number) {
|
||||
const controller = new AbortController();
|
||||
const timeout = setTimeout(() => controller.abort(), timeoutMs);
|
||||
try {
|
||||
const response = await fetch(url, { ...init, signal: controller.signal });
|
||||
const buffer = Buffer.from(await response.arrayBuffer());
|
||||
return { response, buffer };
|
||||
} finally {
|
||||
clearTimeout(timeout);
|
||||
}
|
||||
}
|
||||
|
||||
function wechatPrivateKey(providerConfig: ProviderConfig) {
|
||||
return normalizePem(requireSecretString(providerConfig, ['privateKey', 'merchantPrivateKey']), 'PRIVATE KEY');
|
||||
}
|
||||
|
||||
function signedWechatGetHeaders(providerConfig: ProviderConfig, url: URL) {
|
||||
const mchId = requirePublicString(providerConfig, ['merchantId', 'mchId']);
|
||||
const merchantSerialNo = requirePublicString(providerConfig, ['merchantSerialNo']);
|
||||
const timestamp = Math.floor(Date.now() / 1000).toString();
|
||||
const nonce = randomNonce();
|
||||
const message = ['GET', `${url.pathname}${url.search}`, timestamp, nonce, ''].join('\n') + '\n';
|
||||
const signature = rsaSignSha256(wechatPrivateKey(providerConfig), message);
|
||||
return {
|
||||
accept: 'application/json',
|
||||
authorization: `WECHATPAY2-SHA256-RSA2048 mchid="${mchId}",nonce_str="${nonce}",signature="${signature}",timestamp="${timestamp}",serial_no="${merchantSerialNo}"`,
|
||||
};
|
||||
}
|
||||
|
||||
function wechatBillEndpoint(providerConfig: ProviderConfig, billType: BillType) {
|
||||
if (billType === 'refund') {
|
||||
return providerEndpointForKeys(
|
||||
providerConfig,
|
||||
['wechatRefundBillEndpoint', 'refundBillEndpoint', 'tradeBillEndpoint'],
|
||||
'https://api.mch.weixin.qq.com/v3/bill/tradebill',
|
||||
['api.mch.weixin.qq.com'],
|
||||
);
|
||||
}
|
||||
return providerEndpointForKeys(
|
||||
providerConfig,
|
||||
['wechatTradeBillEndpoint', 'tradeBillEndpoint'],
|
||||
'https://api.mch.weixin.qq.com/v3/bill/tradebill',
|
||||
['api.mch.weixin.qq.com'],
|
||||
);
|
||||
}
|
||||
|
||||
function mapWechatBillType(billType: BillType) {
|
||||
if (billType === 'refund') return 'REFUND';
|
||||
return 'ALL';
|
||||
}
|
||||
|
||||
async function downloadWechatBill(providerConfig: ProviderConfig, job: ProviderBillJob): Promise<DownloadedBill> {
|
||||
const endpoint = wechatBillEndpoint(providerConfig, job.billType);
|
||||
const url = new URL(endpoint);
|
||||
url.searchParams.set('bill_date', job.billDate);
|
||||
url.searchParams.set('bill_type', mapWechatBillType(job.billType));
|
||||
const { response, raw } = await fetchJson(url.toString(), {
|
||||
method: 'GET',
|
||||
headers: signedWechatGetHeaders(providerConfig, url),
|
||||
}, config.commerceRequestTimeoutMs);
|
||||
if (!response.ok) throw new ProviderBillWorkerError('WeChat Pay bill URL request failed', 'PROVIDER_BILL_URL_REQUEST_FAILED');
|
||||
|
||||
const downloadUrl = safeString(raw.download_url);
|
||||
if (!downloadUrl) throw new ProviderBillWorkerError('WeChat Pay bill response missing download_url', 'PROVIDER_BILL_RESPONSE_INVALID');
|
||||
const hashType = safeString(raw.hash_type) || null;
|
||||
const hashValue = safeString(raw.hash_value) || null;
|
||||
const parsedDownloadUrl = assertDownloadUrlAllowed(downloadUrl, 'wechat_pay');
|
||||
const { response: downloadResponse, buffer } = await fetchBuffer(downloadUrl, {
|
||||
method: 'GET',
|
||||
headers: signedWechatGetHeaders(providerConfig, parsedDownloadUrl),
|
||||
}, config.commerceRequestTimeoutMs);
|
||||
if (!downloadResponse.ok) throw new ProviderBillWorkerError('WeChat Pay bill download failed', 'PROVIDER_BILL_DOWNLOAD_FAILED');
|
||||
verifyProviderHash(buffer, hashType, hashValue);
|
||||
const rows = await parseBillRows(buffer, downloadResponse.headers.get('content-type') || '', parsedDownloadUrl.pathname, 'wechat_pay');
|
||||
return {
|
||||
rows,
|
||||
sourceHash: crypto.createHash('sha256').update(buffer).digest('hex'),
|
||||
downloadHashType: hashType,
|
||||
downloadHashValue: hashValue,
|
||||
downloadUrlHost: parsedDownloadUrl.hostname,
|
||||
metadata: {
|
||||
providerResponse: {
|
||||
hashType,
|
||||
hasHashValue: Boolean(hashValue),
|
||||
},
|
||||
downloadedBytes: buffer.length,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function alipayBillType(billType: BillType) {
|
||||
if (billType === 'refund') return 'trade';
|
||||
return 'trade';
|
||||
}
|
||||
|
||||
async function downloadAlipayBill(providerConfig: ProviderConfig, job: ProviderBillJob): Promise<DownloadedBill> {
|
||||
const appId = requirePublicString(providerConfig, ['appId']);
|
||||
const gateway = providerEndpointForKeys(
|
||||
providerConfig,
|
||||
['billEndpoint', 'endpoint'],
|
||||
'https://openapi.alipay.com/gateway.do',
|
||||
['openapi.alipay.com'],
|
||||
);
|
||||
const privateKey = normalizePem(requireSecretString(providerConfig, ['privateKey', 'appPrivateKey']), 'PRIVATE KEY');
|
||||
const params: Record<string, string> = {
|
||||
app_id: appId,
|
||||
method: 'alipay.data.dataservice.bill.downloadurl.query',
|
||||
charset: 'utf-8',
|
||||
sign_type: 'RSA2',
|
||||
timestamp: new Date().toISOString().replace('T', ' ').slice(0, 19),
|
||||
version: '1.0',
|
||||
biz_content: JSON.stringify({
|
||||
bill_type: alipayBillType(job.billType),
|
||||
bill_date: job.billDate,
|
||||
}),
|
||||
};
|
||||
params.sign = rsaSignSha256(privateKey, canonicalForm(params));
|
||||
|
||||
const { response, raw } = await fetchJson(gateway, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
accept: 'application/json',
|
||||
'content-type': 'application/x-www-form-urlencoded;charset=utf-8',
|
||||
},
|
||||
body: encodedForm(params),
|
||||
}, config.commerceRequestTimeoutMs);
|
||||
if (!response.ok) throw new ProviderBillWorkerError('Alipay bill URL request failed', 'PROVIDER_BILL_URL_REQUEST_FAILED');
|
||||
const responseBody = objectValue(raw.alipay_data_dataservice_bill_downloadurl_query_response);
|
||||
if (safeString(responseBody.code) !== '10000') {
|
||||
throw new ProviderBillWorkerError('Alipay bill URL request was rejected', 'PROVIDER_BILL_URL_REQUEST_REJECTED');
|
||||
}
|
||||
const downloadUrl = safeString(responseBody.bill_download_url);
|
||||
if (!downloadUrl) throw new ProviderBillWorkerError('Alipay bill response missing bill_download_url', 'PROVIDER_BILL_RESPONSE_INVALID');
|
||||
const parsedDownloadUrl = assertDownloadUrlAllowed(downloadUrl, 'alipay');
|
||||
const { response: downloadResponse, buffer } = await fetchBuffer(downloadUrl, {
|
||||
method: 'GET',
|
||||
headers: { accept: '*/*' },
|
||||
}, config.commerceRequestTimeoutMs);
|
||||
if (!downloadResponse.ok) throw new ProviderBillWorkerError('Alipay bill download failed', 'PROVIDER_BILL_DOWNLOAD_FAILED');
|
||||
const rows = await parseBillRows(buffer, downloadResponse.headers.get('content-type') || '', parsedDownloadUrl.pathname, 'alipay');
|
||||
return {
|
||||
rows,
|
||||
sourceHash: crypto.createHash('sha256').update(buffer).digest('hex'),
|
||||
downloadHashType: null,
|
||||
downloadHashValue: null,
|
||||
downloadUrlHost: parsedDownloadUrl.hostname,
|
||||
metadata: {
|
||||
providerResponse: {
|
||||
code: responseBody.code,
|
||||
},
|
||||
downloadedBytes: buffer.length,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function verifyProviderHash(buffer: Buffer, hashType: string | null, hashValue: string | null) {
|
||||
if (!hashType && !hashValue) return;
|
||||
if (!hashType || !hashValue) {
|
||||
throw new ProviderBillWorkerError('Provider bill hash metadata is incomplete', 'PROVIDER_BILL_HASH_INVALID');
|
||||
}
|
||||
const normalized = hashType.trim().toUpperCase();
|
||||
const algorithm = normalized === 'SHA1' ? 'sha1' : normalized === 'SHA256' ? 'sha256' : '';
|
||||
if (!algorithm) {
|
||||
throw new ProviderBillWorkerError(`Provider bill hash type ${hashType} is unsupported`, 'PROVIDER_BILL_HASH_UNSUPPORTED');
|
||||
}
|
||||
const actual = crypto.createHash(algorithm).update(buffer).digest('hex');
|
||||
if (actual.toLowerCase() !== hashValue.toLowerCase()) {
|
||||
throw new ProviderBillWorkerError('Provider bill hash mismatch', 'PROVIDER_BILL_HASH_MISMATCH');
|
||||
}
|
||||
}
|
||||
|
||||
function decodeText(buffer: Buffer) {
|
||||
if (buffer.length >= 3 && buffer[0] === 0xef && buffer[1] === 0xbb && buffer[2] === 0xbf) {
|
||||
return buffer.subarray(3).toString('utf8');
|
||||
}
|
||||
const utf8 = buffer.toString('utf8');
|
||||
if (utf8.includes('<27>')) {
|
||||
return iconv.decode(buffer, 'gb18030');
|
||||
}
|
||||
return utf8;
|
||||
}
|
||||
|
||||
async function parseBillRows(buffer: Buffer, contentType: string, pathname: string, provider: ProviderBillProvider) {
|
||||
const lowerPath = pathname.toLowerCase();
|
||||
const lowerType = contentType.toLowerCase();
|
||||
if (lowerPath.endsWith('.zip') || lowerType.includes('zip')) {
|
||||
const zip = await JSZip.loadAsync(buffer);
|
||||
const rows: Record<string, unknown>[] = [];
|
||||
for (const entry of Object.values(zip.files)) {
|
||||
if (entry.dir) continue;
|
||||
const name = entry.name.toLowerCase();
|
||||
if (!name.endsWith('.csv') && !name.endsWith('.txt') && !name.endsWith('.json')) continue;
|
||||
const entryBuffer = Buffer.from(await entry.async('uint8array'));
|
||||
rows.push(...parseBillRowsFromText(decodeText(entryBuffer), provider));
|
||||
}
|
||||
return rows;
|
||||
}
|
||||
return parseBillRowsFromText(decodeText(buffer), provider);
|
||||
}
|
||||
|
||||
function parseBillRowsFromText(text: string, provider: ProviderBillProvider) {
|
||||
const trimmed = text.trim();
|
||||
if (!trimmed) return [];
|
||||
if (trimmed.startsWith('{') || trimmed.startsWith('[')) {
|
||||
const parsed = JSON.parse(trimmed) as unknown;
|
||||
const parsedObject = objectValue(parsed);
|
||||
const rows: unknown[] = Array.isArray(parsed)
|
||||
? parsed
|
||||
: Array.isArray(parsedObject.rows)
|
||||
? parsedObject.rows
|
||||
: [];
|
||||
return rows.map(row => objectValue(row));
|
||||
}
|
||||
return parseDelimitedBillRows(text, provider);
|
||||
}
|
||||
|
||||
function parseDelimitedBillRows(text: string, provider: ProviderBillProvider) {
|
||||
const lines = text
|
||||
.split(/\r?\n/)
|
||||
.map(line => line.trim())
|
||||
.filter(line => line && !line.startsWith('#') && !line.startsWith('`----------'));
|
||||
const headerIndex = lines.findIndex(line => {
|
||||
const normalized = line.replace(/`/g, '');
|
||||
return normalized.includes(',') && (
|
||||
normalized.includes('商户订单号')
|
||||
|| normalized.includes('out_trade_no')
|
||||
|| normalized.includes('providerTradeNo')
|
||||
|| normalized.includes('交易号')
|
||||
|| normalized.includes('微信支付订单号')
|
||||
);
|
||||
});
|
||||
if (headerIndex < 0) {
|
||||
throw new ProviderBillWorkerError('Provider bill header not found', 'PROVIDER_BILL_PARSE_FAILED');
|
||||
}
|
||||
const delimiter = lines[headerIndex].includes('\t') ? '\t' : ',';
|
||||
const headers = splitDelimitedLine(lines[headerIndex], delimiter).map(cleanCell);
|
||||
const rows: Record<string, unknown>[] = [];
|
||||
for (const line of lines.slice(headerIndex + 1)) {
|
||||
const values = splitDelimitedLine(line, delimiter).map(cleanCell);
|
||||
if (values.length < 2) continue;
|
||||
const raw: Record<string, unknown> = {};
|
||||
headers.forEach((header, index) => {
|
||||
if (header) raw[header] = values[index] || '';
|
||||
});
|
||||
const normalized = provider === 'wechat_pay' ? normalizeWechatBillRow(raw) : normalizeAlipayBillRow(raw);
|
||||
if (Object.keys(normalized).length > 0) rows.push(normalized);
|
||||
}
|
||||
return rows;
|
||||
}
|
||||
|
||||
function splitDelimitedLine(line: string, delimiter: string) {
|
||||
if (delimiter === '\t') return line.split('\t');
|
||||
const cells: string[] = [];
|
||||
let current = '';
|
||||
let quoted = false;
|
||||
for (let index = 0; index < line.length; index += 1) {
|
||||
const char = line[index];
|
||||
if (char === '"') {
|
||||
if (quoted && line[index + 1] === '"') {
|
||||
current += '"';
|
||||
index += 1;
|
||||
} else {
|
||||
quoted = !quoted;
|
||||
}
|
||||
} else if (char === ',' && !quoted) {
|
||||
cells.push(current);
|
||||
current = '';
|
||||
} else {
|
||||
current += char;
|
||||
}
|
||||
}
|
||||
cells.push(current);
|
||||
return cells;
|
||||
}
|
||||
|
||||
function cleanCell(value: string) {
|
||||
return value.trim().replace(/^`/, '').replace(/^"|"$/g, '').trim();
|
||||
}
|
||||
|
||||
function firstValue(row: Record<string, unknown>, keys: string[]) {
|
||||
for (const key of keys) {
|
||||
const value = row[key];
|
||||
if (typeof value === 'string' && value.trim()) return value.trim();
|
||||
if (typeof value === 'number' && Number.isFinite(value)) return String(value);
|
||||
}
|
||||
return '';
|
||||
}
|
||||
|
||||
function normalizeAmount(value: string) {
|
||||
if (!value) return '';
|
||||
const cleaned = value.replace(/,/g, '').replace(/[¥¥]/g, '').trim();
|
||||
const parsed = Number(cleaned);
|
||||
return Number.isFinite(parsed) ? cleaned : '';
|
||||
}
|
||||
|
||||
function rowTypeFromText(value: string) {
|
||||
const normalized = value.toLowerCase();
|
||||
if (normalized.includes('refund') || value.includes('退款')) return 'refund';
|
||||
return 'payment';
|
||||
}
|
||||
|
||||
function normalizeWechatBillRow(row: Record<string, unknown>) {
|
||||
const transactionType = rowTypeFromText(firstValue(row, ['账单类型', '交易类型', 'transaction_type', 'type']));
|
||||
const amount = normalizeAmount(firstValue(row, ['订单金额(元)', '应结订单金额(元)', '总金额', 'amount', 'total_amount']));
|
||||
const refundAmount = normalizeAmount(firstValue(row, ['申请退款金额(元)', '退款金额(元)', 'refund_amount', 'refund_fee']));
|
||||
return {
|
||||
transactionType,
|
||||
providerTradeNo: firstValue(row, ['微信支付订单号', 'transaction_id', 'providerTradeNo']),
|
||||
providerRefundNo: firstValue(row, ['微信退款单号', 'refund_id', 'providerRefundNo']),
|
||||
orderNo: firstValue(row, ['商户订单号', 'out_trade_no', 'orderNo']),
|
||||
refundNo: firstValue(row, ['商户退款单号', 'out_refund_no', 'refundNo']),
|
||||
amountYuan: transactionType === 'payment' ? amount : '',
|
||||
refundAmount: transactionType === 'refund' ? (refundAmount || amount) : '',
|
||||
feeYuan: normalizeAmount(firstValue(row, ['手续费(元)', 'service_fee', 'fee'])),
|
||||
paidAt: firstValue(row, ['支付完成时间', '交易时间', 'success_time', 'paidAt']),
|
||||
refundedAt: firstValue(row, ['退款成功时间', 'refund_success_time', 'refundedAt']),
|
||||
providerStatus: firstValue(row, ['交易状态', '退款状态', 'trade_state', 'status']) || 'SUCCESS',
|
||||
rawProviderRow: row,
|
||||
};
|
||||
}
|
||||
|
||||
function normalizeAlipayBillRow(row: Record<string, unknown>) {
|
||||
const transactionType = rowTypeFromText(firstValue(row, ['业务类型', '交易类型', 'transaction_type', 'type']));
|
||||
const amount = normalizeAmount(firstValue(row, ['订单金额(元)', '订单金额(元)', '收入(+元)', '收入(+元)', 'amount', 'total_amount']));
|
||||
const refundAmount = normalizeAmount(firstValue(row, ['退款金额(元)', '退款金额(元)', '支出(-元)', '支出(-元)', 'refund_amount']));
|
||||
return {
|
||||
transactionType,
|
||||
providerTradeNo: firstValue(row, ['支付宝交易号', '交易号', 'trade_no', 'providerTradeNo']),
|
||||
providerRefundNo: firstValue(row, ['支付宝退款单号', 'refund_id', 'providerRefundNo']),
|
||||
orderNo: firstValue(row, ['商户订单号', 'out_trade_no', 'orderNo']),
|
||||
refundNo: firstValue(row, ['商户退款单号', 'out_request_no', 'refundNo']),
|
||||
amountYuan: transactionType === 'payment' ? amount : '',
|
||||
refundAmount: transactionType === 'refund' ? (refundAmount || amount) : '',
|
||||
feeYuan: normalizeAmount(firstValue(row, ['服务费(元)', '服务费(元)', 'fee', 'service_fee'])),
|
||||
paidAt: firstValue(row, ['交易创建时间', '付款时间', 'gmt_payment', 'paidAt']),
|
||||
refundedAt: firstValue(row, ['退款时间', 'gmt_refund', 'refundedAt']),
|
||||
providerStatus: firstValue(row, ['交易状态', 'trade_status', 'status']) || 'TRADE_SUCCESS',
|
||||
rawProviderRow: row,
|
||||
};
|
||||
}
|
||||
|
||||
async function claimProviderBillJobs(client: pg.PoolClient, limit = config.providerBillBatchSize) {
|
||||
const result = await client.query<ProviderBillJob>(
|
||||
`
|
||||
update public.commerce_bill_download_jobs j
|
||||
set status = 'running',
|
||||
claimed_by = $2,
|
||||
claimed_at = now(),
|
||||
updated_at = now(),
|
||||
error_code = null,
|
||||
error_message = null
|
||||
where j.id in (
|
||||
select id
|
||||
from public.commerce_bill_download_jobs
|
||||
where status = 'queued'
|
||||
or (
|
||||
status = 'running'
|
||||
and claimed_at < now() - ($3::int * interval '1 second')
|
||||
)
|
||||
order by created_at asc
|
||||
limit $1
|
||||
for update skip locked
|
||||
)
|
||||
returning id, tenant_id as "tenantId", provider,
|
||||
bill_date::text as "billDate", bill_type as "billType",
|
||||
source_name as "sourceName", requested_by as "requestedBy",
|
||||
metadata
|
||||
`,
|
||||
[limit, config.providerBillWorkerId, config.providerBillClaimStaleSeconds],
|
||||
);
|
||||
return result.rows.map(row => ({
|
||||
...row,
|
||||
provider: normalizeProvider(row.provider) || 'wechat_pay',
|
||||
billDate: normalizeReconciliationBillDate(row.billDate),
|
||||
billType: normalizeReconciliationBillType(row.billType),
|
||||
metadata: objectValue(row.metadata),
|
||||
}));
|
||||
}
|
||||
|
||||
async function processProviderBillJob(job: ProviderBillJob) {
|
||||
const client = await pool.connect();
|
||||
try {
|
||||
const providerConfig = await loadPaymentProviderConfig(client, job.tenantId, job.provider);
|
||||
const downloaded = job.provider === 'wechat_pay'
|
||||
? await downloadWechatBill(providerConfig, job)
|
||||
: await downloadAlipayBill(providerConfig, job);
|
||||
await client.query('begin');
|
||||
const imported = await importReconciliationBatch(client, {
|
||||
tenantId: job.tenantId,
|
||||
actorUserId: job.requestedBy,
|
||||
provider: job.provider,
|
||||
billDate: job.billDate,
|
||||
billType: job.billType,
|
||||
source: 'provider_download',
|
||||
sourceName: `${job.sourceName}:${job.id}`,
|
||||
rows: downloaded.rows,
|
||||
metadata: {
|
||||
...job.metadata,
|
||||
providerBillJobId: job.id,
|
||||
providerBillSourceHash: downloaded.sourceHash,
|
||||
providerBill: downloaded.metadata,
|
||||
},
|
||||
});
|
||||
await client.query(
|
||||
`
|
||||
update public.commerce_bill_download_jobs
|
||||
set status = 'completed',
|
||||
source_hash = $3,
|
||||
row_count = $4,
|
||||
download_hash_type = $5,
|
||||
download_hash_value = $6,
|
||||
download_url_host = $7,
|
||||
reconciliation_batch_id = $8::uuid,
|
||||
completed_at = now(),
|
||||
failed_at = null,
|
||||
error_code = null,
|
||||
error_message = null,
|
||||
metadata = metadata || $9::jsonb,
|
||||
updated_at = now()
|
||||
where tenant_id = $1 and id = $2
|
||||
`,
|
||||
[
|
||||
job.tenantId,
|
||||
job.id,
|
||||
downloaded.sourceHash,
|
||||
downloaded.rows.length,
|
||||
downloaded.downloadHashType,
|
||||
downloaded.downloadHashValue,
|
||||
downloaded.downloadUrlHost,
|
||||
imported.batch.id,
|
||||
JSON.stringify({
|
||||
completedBy: config.providerBillWorkerId,
|
||||
completedAt: new Date().toISOString(),
|
||||
reconciliationBatchId: imported.batch.id,
|
||||
reconciliationStatus: imported.batch.status,
|
||||
reconciliationSummary: {
|
||||
totalCount: imported.batch.totalCount,
|
||||
matchedCount: imported.batch.matchedCount,
|
||||
mismatchCount: imported.batch.mismatchCount,
|
||||
missingLocalCount: imported.batch.missingLocalCount,
|
||||
missingProviderCount: imported.batch.missingProviderCount,
|
||||
duplicateCount: imported.batch.duplicateCount,
|
||||
},
|
||||
}),
|
||||
],
|
||||
);
|
||||
await client.query('commit');
|
||||
return 'completed';
|
||||
} catch (error) {
|
||||
await client.query('rollback').catch(() => {});
|
||||
const failure = publicFailure(error);
|
||||
await client.query(
|
||||
`
|
||||
update public.commerce_bill_download_jobs
|
||||
set status = 'failed',
|
||||
failed_at = now(),
|
||||
error_code = $3,
|
||||
error_message = $4,
|
||||
metadata = metadata || $5::jsonb,
|
||||
updated_at = now()
|
||||
where tenant_id = $1 and id = $2
|
||||
`,
|
||||
[
|
||||
job.tenantId,
|
||||
job.id,
|
||||
failure.code,
|
||||
failure.message,
|
||||
JSON.stringify({
|
||||
failedBy: config.providerBillWorkerId,
|
||||
failedAt: new Date().toISOString(),
|
||||
}),
|
||||
],
|
||||
).catch(() => {});
|
||||
return 'failed';
|
||||
} finally {
|
||||
client.release();
|
||||
}
|
||||
}
|
||||
|
||||
function publicFailure(error: unknown) {
|
||||
const code = error instanceof ProviderBillWorkerError ? error.code : 'PROVIDER_BILL_WORKER_ERROR';
|
||||
const message = error instanceof Error ? error.message : 'Provider bill worker failed';
|
||||
return {
|
||||
code,
|
||||
message: message
|
||||
.replace(/-----BEGIN[\s\S]+?-----END [A-Z ]+-----/g, '[redacted-pem]')
|
||||
.slice(0, 500),
|
||||
};
|
||||
}
|
||||
|
||||
export async function processProviderBillBatch(limit = config.providerBillBatchSize): Promise<ProviderBillWorkerResult> {
|
||||
const client = await pool.connect();
|
||||
let jobs: ProviderBillJob[] = [];
|
||||
try {
|
||||
await client.query('begin');
|
||||
jobs = await claimProviderBillJobs(client, limit);
|
||||
await client.query('commit');
|
||||
} catch (error) {
|
||||
await client.query('rollback');
|
||||
throw error;
|
||||
} finally {
|
||||
client.release();
|
||||
}
|
||||
|
||||
const result: ProviderBillWorkerResult = {
|
||||
processed: jobs.length,
|
||||
completed: 0,
|
||||
failed: 0,
|
||||
skipped: 0,
|
||||
};
|
||||
|
||||
for (const job of jobs) {
|
||||
const status = await processProviderBillJob(job);
|
||||
if (status === 'completed') result.completed += 1;
|
||||
else if (status === 'failed') result.failed += 1;
|
||||
else result.skipped += 1;
|
||||
}
|
||||
return result;
|
||||
}
|
||||
Reference in New Issue
Block a user