From f3577e257c1c5ef479813427ee223f463fe33f8c Mon Sep 17 00:00:00 2001 From: Codex Date: Mon, 29 Jun 2026 21:59:42 +0800 Subject: [PATCH] feat: add provider bill download reconciliation jobs --- README.md | 21 +- apps/api/src/features/commerce/index.ts | 4 + .../src/features/commerce/reconciliation.ts | 472 ++++++++--- apps/api/src/features/tenant-admin/auth.ts | 1 + apps/worker/package.json | 2 + apps/worker/src/config.ts | 6 + apps/worker/src/index.ts | 11 + apps/worker/src/jobs/provider-bills.ts | 784 ++++++++++++++++++ docs/refactor/auth-payment-provider-plan.md | 7 +- docs/refactor/backend-capability-status.md | 3 +- docs/refactor/backend-handoff-roadmap.md | 4 +- docs/refactor/backend-progress.md | 2 +- docs/refactor/blueprint-coverage.md | 2 +- docs/refactor/implementation-status.md | 4 +- docs/refactor/legacy-feature-gap-matrix.md | 2 +- docs/refactor/next-development-todo.md | 7 +- docs/refactor/taro-frontend-integration.md | 64 +- package-lock.json | 1 + scripts/api-integration-test.js | 57 ++ scripts/commerce-worker-integration-test.js | 339 +++++++- ...2606290029_commerce_bill_download_jobs.sql | 56 ++ 21 files changed, 1691 insertions(+), 158 deletions(-) create mode 100644 apps/worker/src/jobs/provider-bills.ts create mode 100644 supabase/migrations/202606290029_commerce_bill_download_jobs.sql diff --git a/README.md b/README.md index f31ddcde..824fdadd 100644 --- a/README.md +++ b/README.md @@ -19,17 +19,17 @@ - 公共题库商业化能力:租户可采纳平台授权题库为本租户副本,并可手动或由 worker 自动同步平台新增/更新题目;同步会保护租户自改题目,返回冲突而不覆盖,后台可查询冲突明细。 - 题库导出能力:租户内容编辑可按题目集合、内容入口或分类节点导出 JSON、`paper_json`、打印 payload、PDF、Word 和每日一练图片 ZIP 素材包,后端强制租户隔离、答案/解析开关、复合题子题脱敏、导出 job 和审计;PDF/Word/ZIP 由 exports worker 生成水印文件或运营素材并发布到 `content_assets`;`daily_practice` 支持每日一练九宫格 metadata、PDF/Word 版式、9 张 PNG/SVG 卡片和拼图包。 - 销售/代理/CRM 增长链路:邀请码、扫码/分享事件、首绑客资保护、销售统计、团队关系、CRM 配置、跟进分配策略和队列。 -- `apps/worker` 后台任务进程:CRM webhook 队列消费、generic/钉钉/飞书/企微机器人发送、签名、失败重试和日志;commerce worker 可补偿查询微信/支付宝支付和退款状态;assets worker 可复检托管资源元数据、执行内置安全扫描并自动下架异常资源;imports worker 可执行大批量导入;public-banks worker 可自动同步公共题库采纳副本;exports worker 可渲染 PDF/Word 导出文件和每日一练 ZIP 图片素材包。 +- `apps/worker` 后台任务进程:CRM webhook 队列消费、generic/钉钉/飞书/企微机器人发送、签名、失败重试和日志;commerce worker 可补偿查询微信/支付宝支付和退款状态;provider-bills worker 可下载微信/支付宝官方账单并导入资金对账;assets worker 可复检托管资源元数据、执行内置安全扫描并自动下架异常资源;imports worker 可执行大批量导入;public-banks worker 可自动同步公共题库采纳副本;exports worker 可渲染 PDF/Word 导出文件和每日一练 ZIP 图片素材包。 - 销售/代理分佣结算基础闭环:租户默认比例、成员比例、激活码批次比例、订单/激活码归因、结算单生成、审核、线下打款状态、CSV/JSON 导出、打款凭证登记/复核和权限隔离。 - 订单售后基础闭环:退款请求、审核、处理状态流、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、退款金额累计、部分/全额退款订单状态、全额退款权益撤销、退款事件和审计日志。 -- 资金对账和差错工单闭环:租户财务/运营可通过 `/api/commerce/reconciliation/*` 导入或预览支付/退款账单行,后端按租户隔离比对本地订单、支付、退款记录,识别已匹配、金额不一致、状态不一致、供应商有本地无、本地有供应商无、重复行和无效行,并写入对账批次、明细和审计日志;异常明细可创建差错工单,支持分配、开始处理、升级、解决、忽略、重开和事件留痕。工单只做财务审核闭环,不直接修改订单、支付、退款或权益。 +- 资金对账和差错工单闭环:租户财务/运营可通过 `/api/commerce/reconciliation/*` 导入或预览支付/退款账单行,也可创建微信/支付宝官方账单下载任务;后端按租户隔离比对本地订单、支付、退款记录,识别已匹配、金额不一致、状态不一致、供应商有本地无、本地有供应商无、重复行和无效行,并写入对账批次、明细和审计日志;异常明细可创建差错工单,支持分配、开始处理、升级、解决、忽略、重开和事件留痕。工单只做财务审核闭环,不直接修改订单、支付、退款或权益。 - PocketBase schema/数据导入器雏形和导入后校验脚本。 - 本地 Supabase reset、烟测 seed、API 集成测试、完整重构检查命令。 还没有达到生产交付的部分: - Supabase Auth/JWT、租户角色模板、班级/教师/学生范围权限已可联调;生产前还要做真实云端 Auth/JWKS 回归和 RLS 深测。 -- 阿里云/腾讯云短信、微信小程序登录、微信网页登录、QQ 登录、手机号绑定/换绑、微信支付、支付宝主链路、微信/支付宝发起退款/查询确认/退款通知、支付/退款补偿 worker 已完成本地适配;资金对账已支持手工/API 账单导入比对和差错工单处理,微信/支付宝官方账单自动下载、异常订单运营台和真实生产账号联调还没接完。 +- 阿里云/腾讯云短信、微信小程序登录、微信网页登录、QQ 登录、手机号绑定/换绑、微信支付、支付宝主链路、微信/支付宝发起退款/查询确认/退款通知、支付/退款补偿 worker 已完成本地适配;资金对账已支持手工/API 账单导入、微信/支付宝官方账单下载任务、provider-bills worker 自动导入比对和差错工单处理;异常订单运营台、人工调整凭证复核报表和真实生产账号联调还没接完。 - OSS/COS/Supabase Storage 上传下载签名 provider 已接入;上传后校验、PDF/图片预览、资源访问事件、动态水印上下文、锁定资源 CDN 边界、资源复检 worker、内置 `metadata_rules` 安全扫描和外部 HTTP 杀毒/内容安全 scanner 接入层已完成。生产还要配置真实扫描服务 endpoint/token,并继续补转码/CDN 级水印、CDN 刷新和对象生命周期策略。 - Excel/CSV 导入解析已完成并复用 `content_import_jobs/items/issues` 管线;大批量异步导入 worker 基础已接入,支持 queued job 消费、重试和审计;导入后复检、模板下载和字段映射 API 已完成,前端 UI 待接。 - 题库导出已完成服务端结构化 payload、PDF/Word 二进制 worker、每日一练基础导出和每日一练 ZIP 图片素材包;后续还要补更精细试卷模板、多模板排版和导出操作台体验。 @@ -61,7 +61,7 @@ ```text apps/api/ Node.js 业务 API apps/taro/ Taro 4 React 跨端前端,H5 三入口,后续扩展小程序 -apps/worker/ 后台异步任务:CRM webhook、支付/退款补偿、资源复检、导入执行、公共题库同步、题库导出渲染等 +apps/worker/ 后台异步任务:CRM webhook、支付/退款补偿、官方账单下载、资源复检、导入执行、公共题库同步、题库导出渲染等 packages/config/ 共享配置 packages/db/ PostgreSQL 连接池和查询封装 packages/domain/ 领域常量和共享类型 @@ -165,6 +165,12 @@ npm --workspace @tiku-saas/worker run crm:once npm --workspace @tiku-saas/worker run commerce:once ``` +单次运行微信/支付宝官方账单下载 worker: + +```bash +npm --workspace @tiku-saas/worker run provider-bills:once +``` + 单次运行内容资源复检 worker: ```bash @@ -310,7 +316,7 @@ API 身份上下文: - 题库入口和分类使用 `content_entries/content_nodes`;题目列表和练习规则使用 `question_collections/practice_blueprints`,前端不要再把旧树字段当成唯一业务结构。 - 批量导入必须先写 `content_import_jobs/items/issues`,保留原始 payload、规范化 payload、逐行问题和审计记录。题目、单词、知识手册、分数线和视频 JSON/CSV/Excel 导入已走这套后台校验管线;大批量任务可提交 `executionMode=async`,由 imports worker 消费,前端只轮询 job 状态和展示 issues。 - 题库导出必须由后端按权限生成,不允许前端直接读取数据库拼导出文件;不开启答案/解析时,顶层题目和复合题子题都必须脱敏;PDF/Word/每日一练 ZIP 只通过 exports worker 写入 `content_assets` 后再签名下载/预览。 -- 支付 webhook 必须先设计幂等键和验签流程,再进入生产使用;生产环境还应定时运行 commerce worker 兜底供应商漏通知和处理中退款,并定期通过资金对账接口导入供应商账单核对本地订单。对账差错工单只允许记录财务处理结论和凭证,不允许前端或工单接口直接篡改订单、支付、退款或权益状态。 +- 支付 webhook 必须先设计幂等键和验签流程,再进入生产使用;生产环境还应定时运行 commerce worker 兜底供应商漏通知和处理中退款,并定时运行 provider-bills worker 下载官方账单核对本地订单。官方账单下载任务只保存下载域名、hash 和对账批次 ID,不向前端暴露下载 URL 或商户密钥。对账差错工单只允许记录财务处理结论和凭证,不允许前端或工单接口直接篡改订单、支付、退款或权益状态。 ## 最近一次验证 @@ -320,6 +326,7 @@ API 身份上下文: npx supabase db reset npm run check:api npm run check:worker +npm run test:worker:commerce npm run test:worker:assets npm run test:worker:exports npm run test:api @@ -328,7 +335,7 @@ npm run audit:runtime git diff --check ``` -结果:通过。`npm run test:api` 覆盖资源访问事件、锁定 CDN 资源拒绝、provider-managed CDN 显式放行、学生短 TTL 下载/预览、访问记录查询和安全扫描门禁。`npm run test:worker:assets` 覆盖托管资源复检、内置安全扫描、外部 HTTP scanner 通过/失败/不可用 fail-closed、扫描失败/跳过事件和异常资源自动下架。`npm run test:worker:exports` 覆盖导出 worker 生成可信资源并标记 `securityScanStatus=passed`。`npm run audit:runtime` 无 high/critical 漏洞;当前运行时依赖树仍有 `exceljs -> uuid` 的 moderate 级提示,修复需要破坏性降级 `exceljs`,后续应在导入 Excel 回归充分后单独处理。 +结果:通过。`npm run test:api` 覆盖资源访问事件、锁定 CDN 资源拒绝、provider-managed CDN 显式放行、学生短 TTL 下载/预览、访问记录查询、安全扫描门禁、官方账单下载任务权限和脱敏响应。`npm run test:worker:commerce` 覆盖支付/退款补偿、微信/支付宝官方账单下载、账单 hash 校验、导入 `provider_download` 对账批次和密钥不泄露。`npm run test:worker:assets` 覆盖托管资源复检、内置安全扫描、外部 HTTP scanner 通过/失败/不可用 fail-closed、扫描失败/跳过事件和异常资源自动下架。`npm run test:worker:exports` 覆盖导出 worker 生成可信资源并标记 `securityScanStatus=passed`。`npm run audit:runtime` 无 high/critical 漏洞;当前运行时依赖树仍有 `exceljs -> uuid` 的 moderate 级提示,修复需要破坏性降级 `exceljs`,后续应在导入 Excel 回归充分后单独处理。 注意:`apps/taro` 是静态构建工程,线上发布 `apps/taro/dist/**`,不发布 `node_modules`。Taro 4.2.0 当前构建工具链仍会触发 `npm run audit:taro:toolchain` 的上游 high/critical 提示,不能用 `npm audit fix --force` 降级到 Taro 3 破坏构建;上线验收时以 `audit:runtime`、构建产物、前端密钥检查和静态服务器配置为准,并持续跟进 Taro 官方修复。 @@ -340,4 +347,4 @@ git diff --check 2. 继续补 Taro 前端:学生端视频/反馈/模考报告/订单收银台,租户后台写入表单/导入操作台/公共题库同步/角色模板 UI,平台后台租户详情/审计/自动计费增强,小程序兼容验证。 3. 对象存储真实 AV/内容安全扫描服务联调、CDN 防盗链、转码/CDN 级水印和生命周期策略。 4. 题库导出模板精排、导出操作台、真实数据 dry-run、导入字段映射 UI 和复检结果操作台。 -5. 真实 OAuth/短信/支付生产账号联调、微信/支付宝官方账单自动下载、异常订单运营台、公共题库版本通知/冲突处理操作台、积分活动深化,以及排行榜防刷/预聚合。 +5. 真实 OAuth/短信/支付生产账号联调、异常订单运营台、人工调整凭证复核、公共题库版本通知/冲突处理操作台、积分活动深化,以及排行榜防刷/预聚合。 diff --git a/apps/api/src/features/commerce/index.ts b/apps/api/src/features/commerce/index.ts index 02e4465f..5bcdde11 100644 --- a/apps/api/src/features/commerce/index.ts +++ b/apps/api/src/features/commerce/index.ts @@ -21,11 +21,13 @@ import { createReconciliationIssueRoute, importReconciliationRoute, previewReconciliationRoute, + providerBillDownloadJobsRoute, reconciliationAnomaliesRoute, reconciliationBatchesRoute, reconciliationIssueEventsRoute, reconciliationIssuesRoute, reconciliationItemsRoute, + requestProviderBillDownloadRoute, updateReconciliationIssueStatusRoute, } from './reconciliation.js'; @@ -52,6 +54,8 @@ export const commerceRoutes: RouteDefinition[] = [ ['GET', '/api/commerce/reconciliation/batches', reconciliationBatchesRoute], ['GET', '/api/commerce/reconciliation/items', reconciliationItemsRoute], ['GET', '/api/commerce/reconciliation/anomalies', reconciliationAnomaliesRoute], + ['POST', '/api/commerce/reconciliation/provider-bills/request', requestProviderBillDownloadRoute], + ['GET', '/api/commerce/reconciliation/provider-bills/jobs', providerBillDownloadJobsRoute], ['POST', '/api/commerce/reconciliation/issues/create', createReconciliationIssueRoute], ['GET', '/api/commerce/reconciliation/issues', reconciliationIssuesRoute], ['POST', '/api/commerce/reconciliation/issues/status', updateReconciliationIssueStatusRoute], diff --git a/apps/api/src/features/commerce/reconciliation.ts b/apps/api/src/features/commerce/reconciliation.ts index 543bc019..b4806113 100644 --- a/apps/api/src/features/commerce/reconciliation.ts +++ b/apps/api/src/features/commerce/reconciliation.ts @@ -6,9 +6,9 @@ import { intParam, optionalInteger, optionalString, readJsonBody, requiredString import { query, transaction } from '../../core/db.js'; import { requireTenantAdmin, requireTenantPermission, type TenantAdminAuth } from '../tenant-admin/auth.js'; -type ReconciliationProvider = 'wechat_pay' | 'alipay' | 'manual'; -type BillType = 'payment' | 'refund' | 'combined'; -type ReconciliationSource = 'manual_upload' | 'provider_download' | 'api' | 'worker'; +export type ReconciliationProvider = 'wechat_pay' | 'alipay' | 'manual'; +export type BillType = 'payment' | 'refund' | 'combined'; +export type ReconciliationSource = 'manual_upload' | 'provider_download' | 'api' | 'worker'; type TransactionType = 'payment' | 'refund'; type MatchStatus = | 'matched' @@ -62,6 +62,7 @@ const PROVIDER_SUCCESS_STATUSES = new Set([ ]); const LOCAL_PAID_PAYMENT_STATUSES = new Set(['paid', 'partially_refunded', 'refunded']); const LOCAL_SUCCEEDED_REFUND_STATUSES = new Set(['succeeded']); +const MAX_PROVIDER_DOWNLOAD_ROWS = 50_000; interface RawBillRow { [key: string]: unknown; @@ -164,7 +165,19 @@ interface ReconciliationBuildResult { items: ReconciliationItemDraft[]; } -function normalizeProvider(value: string): ReconciliationProvider { +export interface ReconciliationBatchImportInput { + tenantId: string; + actorUserId: string | null; + provider: ReconciliationProvider; + billDate: string; + billType: BillType; + source: ReconciliationSource; + sourceName: string | null; + rows: unknown[]; + metadata?: Record; +} + +export function normalizeReconciliationProvider(value: string): ReconciliationProvider { const normalized = value.trim().toLowerCase().replace(/[-\s]/g, '_'); if (['wechat', 'wechatpay', 'wxpay', 'wx_pay', 'wechat_pay'].includes(normalized)) return 'wechat_pay'; if (['alipay', 'ali_pay'].includes(normalized)) return 'alipay'; @@ -178,7 +191,7 @@ function providerAliases(provider: ReconciliationProvider) { return ['manual']; } -function normalizeBillType(value: unknown): BillType { +export function normalizeReconciliationBillType(value: unknown): BillType { if (typeof value !== 'string' || !value.trim()) return 'combined'; const normalized = value.trim().toLowerCase(); if (BILL_TYPES.has(normalized as BillType)) return normalized as BillType; @@ -192,7 +205,7 @@ function normalizeSource(value: unknown): ReconciliationSource { throw new HttpError(400, 'source is invalid', 'RECONCILIATION_SOURCE_INVALID'); } -function normalizeBillDate(value: string) { +export function normalizeReconciliationBillDate(value: string) { const trimmed = value.trim(); if (!/^\d{4}-\d{2}-\d{2}$/.test(trimmed)) { throw new HttpError(400, 'billDate must be YYYY-MM-DD', 'RECONCILIATION_BILL_DATE_INVALID'); @@ -823,9 +836,9 @@ async function buildReconciliation( } function readReconciliationInput(body: Record) { - const provider = normalizeProvider(requiredString(body, 'provider')); - const billDate = normalizeBillDate(requiredString(body, 'billDate')); - const billType = normalizeBillType(body.billType ?? body.bill_type); + const provider = normalizeReconciliationProvider(requiredString(body, 'provider')); + const billDate = normalizeReconciliationBillDate(requiredString(body, 'billDate')); + const billType = normalizeReconciliationBillType(body.billType ?? body.bill_type); const source = normalizeSource(body.source); const sourceName = optionalString(body, 'sourceName') || optionalString(body, 'source_name') || null; const rawRows = body.rows; @@ -839,6 +852,160 @@ function readReconciliationInput(body: Record) { return { provider, billDate, billType, source, sourceName, rows }; } +function normalizeProviderDownloadedRows(provider: ReconciliationProvider, billType: BillType, rows: unknown[]) { + if (rows.length > MAX_PROVIDER_DOWNLOAD_ROWS) { + throw new HttpError(413, `Too many provider bill rows. Max ${MAX_PROVIDER_DOWNLOAD_ROWS}.`, 'RECONCILIATION_PROVIDER_ROWS_TOO_MANY'); + } + return normalizeRows(provider, billType, rows); +} + +export async function importReconciliationBatch( + client: pg.PoolClient, + input: ReconciliationBatchImportInput, +) { + const provider = input.provider; + const billDate = normalizeReconciliationBillDate(input.billDate); + const billType = normalizeReconciliationBillType(input.billType); + const source = normalizeSource(input.source); + const sourceName = input.sourceName; + const rows = normalizeProviderDownloadedRows(provider, billType, input.rows); + const built = await buildReconciliation(client, { + tenantId: input.tenantId, + provider, + billDate, + billType, + source, + sourceName, + rows, + }); + const insertedBatch = await client.query>( + ` + insert into public.commerce_reconciliation_batches ( + tenant_id, provider, bill_date, bill_type, source, source_name, source_hash, + status, total_count, matched_count, mismatch_count, missing_local_count, + missing_provider_count, duplicate_count, ignored_count, amount_cents, + refund_amount_cents, fee_cents, created_by, completed_at, metadata + ) + values ( + $1, $2, $3::date, $4, $5, $6, $7, + $8, $9, $10, $11, $12, + $13, $14, $15, $16, + $17, $18, $19::uuid, now(), $20::jsonb + ) + returning id, provider, bill_date as "billDate", bill_type as "billType", + source, source_name as "sourceName", source_hash as "sourceHash", + status, total_count as "totalCount", matched_count as "matchedCount", + mismatch_count as "mismatchCount", missing_local_count as "missingLocalCount", + missing_provider_count as "missingProviderCount", duplicate_count as "duplicateCount", + ignored_count as "ignoredCount", amount_cents as "amountCents", + refund_amount_cents as "refundAmountCents", fee_cents as "feeCents", + created_by as "createdBy", completed_at as "completedAt", error, + metadata, created_at as "createdAt", updated_at as "updatedAt" + `, + [ + input.tenantId, + built.batch.provider, + built.batch.billDate, + built.batch.billType, + built.batch.source, + built.batch.sourceName, + built.batch.sourceHash, + built.batch.status, + built.batch.totalCount, + built.batch.matchedCount, + built.batch.mismatchCount, + built.batch.missingLocalCount, + built.batch.missingProviderCount, + built.batch.duplicateCount, + built.batch.ignoredCount, + built.batch.amountCents, + built.batch.refundAmountCents, + built.batch.feeCents, + input.actorUserId, + JSON.stringify(input.metadata || {}), + ], + ); + const batch = insertedBatch.rows[0]; + const batchId = String(batch.id); + + for (const item of built.items) { + await client.query( + ` + insert into public.commerce_reconciliation_items ( + tenant_id, batch_id, row_no, provider, transaction_type, + provider_trade_no, provider_refund_no, order_no, refund_no, + amount_cents, refund_amount_cents, fee_cents, paid_at, refunded_at, + provider_status, local_status, order_id, payment_id, refund_request_id, + match_status, severity, issue_code, details + ) + values ( + $1, $2, $3, $4, $5, + $6, $7, $8, $9, + $10, $11, $12, $13::timestamptz, $14::timestamptz, + $15, $16, $17::uuid, $18::uuid, $19::uuid, + $20, $21, $22, $23::jsonb + ) + `, + [ + input.tenantId, + batchId, + item.rowNo, + item.provider, + item.transactionType, + item.providerTradeNo, + item.providerRefundNo, + item.orderNo, + item.refundNo, + item.amountCents, + item.refundAmountCents, + item.feeCents, + item.paidAt, + item.refundedAt, + item.providerStatus, + item.localStatus, + item.orderId, + item.paymentId, + item.refundRequestId, + item.matchStatus, + item.severity, + item.issueCode, + JSON.stringify(item.details), + ], + ); + } + + await client.query( + ` + insert into public.audit_logs (tenant_id, actor_user_id, action, target_type, target_id, details) + values ($1, $2::uuid, 'commerce.reconciliation.imported', 'commerce_reconciliation_batch', $3, $4::jsonb) + `, + [ + input.tenantId, + input.actorUserId, + batchId, + JSON.stringify({ + provider: built.batch.provider, + billDate: built.batch.billDate, + billType: built.batch.billType, + source: built.batch.source, + sourceName: built.batch.sourceName, + sourceHash: built.batch.sourceHash, + summary: { + totalCount: built.batch.totalCount, + matchedCount: built.batch.matchedCount, + mismatchCount: built.batch.mismatchCount, + missingLocalCount: built.batch.missingLocalCount, + missingProviderCount: built.batch.missingProviderCount, + duplicateCount: built.batch.duplicateCount, + ignoredCount: built.batch.ignoredCount, + }, + }), + ], + ); + + return { batch, built }; +} + function batchPayload(row: Record) { return { id: row.id, @@ -1065,6 +1232,158 @@ async function authorizeWrite(ctx: RequestContext) { return auth; } +async function authorizeDownload(ctx: RequestContext) { + const auth = await requireTenantAdmin(ctx); + requireTenantPermission(auth, 'tenant:reconciliation:download'); + return auth; +} + +function providerBillJobPayload(row: Record) { + return { + id: row.id, + provider: row.provider, + billDate: row.billDate, + billType: row.billType, + status: row.status, + sourceName: row.sourceName, + sourceHash: row.sourceHash, + rowCount: row.rowCount, + downloadHashType: row.downloadHashType, + downloadHashValue: row.downloadHashValue, + downloadUrlHost: row.downloadUrlHost, + reconciliationBatchId: row.reconciliationBatchId, + requestedBy: row.requestedBy, + claimedBy: row.claimedBy, + claimedAt: row.claimedAt, + completedAt: row.completedAt, + failedAt: row.failedAt, + errorCode: row.errorCode, + errorMessage: row.errorMessage, + metadata: row.metadata || {}, + createdAt: row.createdAt, + updatedAt: row.updatedAt, + }; +} + +export async function requestProviderBillDownloadRoute(ctx: RequestContext) { + const auth = await authorizeDownload(ctx); + const body = await readJsonBody(ctx); + const provider = normalizeReconciliationProvider(requiredString(body, 'provider')); + if (provider === 'manual') { + throw new HttpError(400, 'Manual provider does not support official bill download', 'PROVIDER_BILL_DOWNLOAD_UNSUPPORTED'); + } + const billDate = normalizeReconciliationBillDate(requiredString(body, 'billDate')); + const billType = normalizeReconciliationBillType(body.billType ?? body.bill_type); + const metadata = objectValue(body.metadata); + const sourceName = `provider-bill:${provider}:${billDate}:${billType}`; + + const result = await transaction(async client => { + const inserted = await client.query>( + ` + insert into public.commerce_bill_download_jobs ( + tenant_id, provider, bill_date, bill_type, status, + source_name, requested_by, metadata + ) + values ($1, $2, $3::date, $4, 'queued', $5, $6::uuid, $7::jsonb) + on conflict (tenant_id, provider, bill_date, bill_type) + where status in ('queued', 'running', 'completed') + do update set + metadata = public.commerce_bill_download_jobs.metadata || excluded.metadata, + updated_at = now() + returning id, provider, bill_date as "billDate", bill_type as "billType", + status, source_name as "sourceName", source_hash as "sourceHash", + row_count as "rowCount", download_hash_type as "downloadHashType", + download_hash_value as "downloadHashValue", download_url_host as "downloadUrlHost", + reconciliation_batch_id as "reconciliationBatchId", + requested_by as "requestedBy", claimed_by as "claimedBy", + claimed_at as "claimedAt", completed_at as "completedAt", + failed_at as "failedAt", error_code as "errorCode", + error_message as "errorMessage", metadata, + created_at as "createdAt", updated_at as "updatedAt", + (xmax = 0) as "inserted" + `, + [ + auth.tenantId, + provider, + billDate, + billType, + sourceName, + auth.userId, + JSON.stringify({ + ...metadata, + requestedFrom: 'api', + }), + ], + ); + const job = inserted.rows[0]; + await recordReconciliationAudit( + client, + auth, + job.inserted ? 'commerce.provider_bill_download.requested' : 'commerce.provider_bill_download.request_reused', + String(job.id), + { + provider, + billDate, + billType, + sourceName, + status: job.status, + }, + ); + return job; + }); + + return { + item: providerBillJobPayload(result), + idempotent: result.inserted === false, + }; +} + +export async function providerBillDownloadJobsRoute(ctx: RequestContext) { + const auth = await authorizeRead(ctx); + const limit = intParam(ctx, 'limit', 50, 200); + const provider = stringParam(ctx, 'provider'); + const status = stringParam(ctx, 'status'); + const billDate = stringParam(ctx, 'billDate'); + + const params: unknown[] = [auth.tenantId, limit]; + const where = ['tenant_id = $1']; + if (provider) { + const normalized = normalizeReconciliationProvider(provider); + if (normalized === 'manual') throw new HttpError(400, 'provider is invalid for bill download', 'PROVIDER_BILL_DOWNLOAD_UNSUPPORTED'); + params.push(normalized); + where.push(`provider = $${params.length}`); + } + if (status) { + params.push(status); + where.push(`status = $${params.length}`); + } + if (billDate) { + params.push(normalizeReconciliationBillDate(billDate)); + where.push(`bill_date = $${params.length}::date`); + } + + const items = await query>( + ` + select id, provider, bill_date as "billDate", bill_type as "billType", + status, source_name as "sourceName", source_hash as "sourceHash", + row_count as "rowCount", download_hash_type as "downloadHashType", + download_hash_value as "downloadHashValue", download_url_host as "downloadUrlHost", + reconciliation_batch_id as "reconciliationBatchId", + requested_by as "requestedBy", claimed_by as "claimedBy", + claimed_at as "claimedAt", completed_at as "completedAt", + failed_at as "failedAt", error_code as "errorCode", + error_message as "errorMessage", metadata, + created_at as "createdAt", updated_at as "updatedAt" + from public.commerce_bill_download_jobs + where ${where.join(' and ')} + order by created_at desc + limit $2 + `, + params, + ); + return { items: items.map(providerBillJobPayload) }; +} + export async function previewReconciliationRoute(ctx: RequestContext) { const auth = await authorizeRead(ctx); const body = await readJsonBody(ctx, { maxBytes: config.maxImportJsonBodyBytes }); @@ -1090,122 +1409,21 @@ export async function importReconciliationRoute(ctx: RequestContext) { const metadata = objectValue(body.metadata); const result = await transaction(async client => { - const built = await buildReconciliation(client, { tenantId: auth.tenantId, ...input }); - const insertedBatch = await client.query>( - ` - insert into public.commerce_reconciliation_batches ( - tenant_id, provider, bill_date, bill_type, source, source_name, source_hash, - status, total_count, matched_count, mismatch_count, missing_local_count, - missing_provider_count, duplicate_count, ignored_count, amount_cents, - refund_amount_cents, fee_cents, created_by, completed_at, metadata - ) - values ( - $1, $2, $3::date, $4, $5, $6, $7, - $8, $9, $10, $11, $12, - $13, $14, $15, $16, - $17, $18, $19, now(), $20::jsonb - ) - returning id, provider, bill_date as "billDate", bill_type as "billType", - source, source_name as "sourceName", source_hash as "sourceHash", - status, total_count as "totalCount", matched_count as "matchedCount", - mismatch_count as "mismatchCount", missing_local_count as "missingLocalCount", - missing_provider_count as "missingProviderCount", duplicate_count as "duplicateCount", - ignored_count as "ignoredCount", amount_cents as "amountCents", - refund_amount_cents as "refundAmountCents", fee_cents as "feeCents", - created_by as "createdBy", completed_at as "completedAt", error, - metadata, created_at as "createdAt", updated_at as "updatedAt" - `, - [ - auth.tenantId, - built.batch.provider, - built.batch.billDate, - built.batch.billType, - built.batch.source, - built.batch.sourceName, - built.batch.sourceHash, - built.batch.status, - built.batch.totalCount, - built.batch.matchedCount, - built.batch.mismatchCount, - built.batch.missingLocalCount, - built.batch.missingProviderCount, - built.batch.duplicateCount, - built.batch.ignoredCount, - built.batch.amountCents, - built.batch.refundAmountCents, - built.batch.feeCents, - auth.userId, - JSON.stringify(metadata), - ], - ); - const batch = insertedBatch.rows[0]; - const batchId = String(batch.id); - - for (const item of built.items) { - await client.query( - ` - insert into public.commerce_reconciliation_items ( - tenant_id, batch_id, row_no, provider, transaction_type, - provider_trade_no, provider_refund_no, order_no, refund_no, - amount_cents, refund_amount_cents, fee_cents, paid_at, refunded_at, - provider_status, local_status, order_id, payment_id, refund_request_id, - match_status, severity, issue_code, details - ) - values ( - $1, $2, $3, $4, $5, - $6, $7, $8, $9, - $10, $11, $12, $13::timestamptz, $14::timestamptz, - $15, $16, $17::uuid, $18::uuid, $19::uuid, - $20, $21, $22, $23::jsonb - ) - `, - [ - auth.tenantId, - batchId, - item.rowNo, - item.provider, - item.transactionType, - item.providerTradeNo, - item.providerRefundNo, - item.orderNo, - item.refundNo, - item.amountCents, - item.refundAmountCents, - item.feeCents, - item.paidAt, - item.refundedAt, - item.providerStatus, - item.localStatus, - item.orderId, - item.paymentId, - item.refundRequestId, - item.matchStatus, - item.severity, - item.issueCode, - JSON.stringify(item.details), - ], - ); - } - - await recordReconciliationAudit(client, auth, 'commerce.reconciliation.imported', batchId, { - provider: built.batch.provider, - billDate: built.batch.billDate, - billType: built.batch.billType, - source: built.batch.source, - sourceName: built.batch.sourceName, - sourceHash: built.batch.sourceHash, - summary: { - totalCount: built.batch.totalCount, - matchedCount: built.batch.matchedCount, - mismatchCount: built.batch.mismatchCount, - missingLocalCount: built.batch.missingLocalCount, - missingProviderCount: built.batch.missingProviderCount, - duplicateCount: built.batch.duplicateCount, - ignoredCount: built.batch.ignoredCount, - }, + const imported = await importReconciliationBatch(client, { + tenantId: auth.tenantId, + actorUserId: auth.userId, + provider: input.provider, + billDate: input.billDate, + billType: input.billType, + source: input.source, + sourceName: input.sourceName, + rows: body.rows as unknown[], + metadata, }); - - return { batch, previewItems: built.items.slice(0, optionalInteger(body, 'previewLimit', 50)) }; + return { + batch: imported.batch, + previewItems: imported.built.items.slice(0, optionalInteger(body, 'previewLimit', 50)), + }; }); return { @@ -1224,7 +1442,7 @@ export async function reconciliationBatchesRoute(ctx: RequestContext) { const params: unknown[] = [auth.tenantId, limit]; const where = ['tenant_id = $1']; if (provider) { - params.push(normalizeProvider(provider)); + params.push(normalizeReconciliationProvider(provider)); where.push(`provider = $${params.length}`); } if (status) { @@ -1232,7 +1450,7 @@ export async function reconciliationBatchesRoute(ctx: RequestContext) { where.push(`status = $${params.length}`); } if (billDate) { - params.push(normalizeBillDate(billDate)); + params.push(normalizeReconciliationBillDate(billDate)); where.push(`bill_date = $${params.length}::date`); } @@ -1309,7 +1527,7 @@ export async function reconciliationAnomaliesRoute(ctx: RequestContext) { const auth = await authorizeRead(ctx); const limit = intParam(ctx, 'limit', 50, 200); const providerParam = stringParam(ctx, 'provider'); - const provider = providerParam ? normalizeProvider(providerParam) : ''; + const provider = providerParam ? normalizeReconciliationProvider(providerParam) : ''; const issueParams: unknown[] = [auth.tenantId, limit]; const issueWhere = [`i.tenant_id = $1`, `i.match_status not in ('matched', 'ignored')`]; diff --git a/apps/api/src/features/tenant-admin/auth.ts b/apps/api/src/features/tenant-admin/auth.ts index 991d8b91..9a807046 100644 --- a/apps/api/src/features/tenant-admin/auth.ts +++ b/apps/api/src/features/tenant-admin/auth.ts @@ -92,6 +92,7 @@ export function tenantPermissionCatalog() { { key: 'tenant:payment:write', label: '商户配置管理' }, { key: 'tenant:reconciliation:read', label: '资金对账查看' }, { key: 'tenant:reconciliation:write', label: '资金对账导入' }, + { key: 'tenant:reconciliation:download', label: '官方账单下载' }, { key: 'tenant:refund:read', label: '退款查看' }, { key: 'tenant:refund:write', label: '退款申请/处理' }, { key: 'tenant:refund:review', label: '退款审核' }, diff --git a/apps/worker/package.json b/apps/worker/package.json index 707de4ad..ac56c4a6 100644 --- a/apps/worker/package.json +++ b/apps/worker/package.json @@ -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" diff --git a/apps/worker/src/config.ts b/apps/worker/src/config.ts index 535c61e5..5e66d177 100644 --- a/apps/worker/src/config.ts +++ b/apps/worker/src/config.ts @@ -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), diff --git a/apps/worker/src/index.ts b/apps/worker/src/index.ts index 3510c91b..e99503e7 100644 --- a/apps/worker/src/index.ts +++ b/apps/worker/src/index.ts @@ -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( diff --git a/apps/worker/src/jobs/provider-bills.ts b/apps/worker/src/jobs/provider-bills.ts new file mode 100644 index 00000000..dfa00ed4 --- /dev/null +++ b/apps/worker/src/jobs/provider-bills.ts @@ -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; + +interface SecretRow { + secretValue: string | null; + secretJson: Record | null; +} + +interface ProviderConfig { + tenantId: string; + provider: ProviderBillProvider; + configPublic: Record; + secret: SecretRow | null; +} + +interface ProviderBillJob { + id: string; + tenantId: string; + provider: ProviderBillProvider; + billDate: string; + billType: BillType; + sourceName: string; + requestedBy: string | null; + metadata: Record; +} + +interface DownloadedBill { + rows: Record[]; + sourceHash: string; + downloadHashType: string | null; + downloadHashValue: string | null; + downloadUrlHost: string | null; + metadata: Record; +} + +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 { + return value && typeof value === 'object' && !Array.isArray(value) ? value as Record : {}; +} + +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) { + 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) { + 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, provider: ProviderBillProvider) { + const secretRef = parseSecretRef(configPublic.secretRef); + const secretKey = secretRef?.key || provider; + const result = await client.query( + ` + 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 { + const aliases = providerAliases(provider); + const accountResult = await client.query<{ + provider: string; + configPublic: Record | 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 { + 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 { + 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 = { + 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('�')) { + 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[] = []; + 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[] = []; + for (const line of lines.slice(headerIndex + 1)) { + const values = splitDelimitedLine(line, delimiter).map(cleanCell); + if (values.length < 2) continue; + const raw: Record = {}; + 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, 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) { + 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) { + 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( + ` + 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 { + 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; +} diff --git a/docs/refactor/auth-payment-provider-plan.md b/docs/refactor/auth-payment-provider-plan.md index 8afba192..7ef2e10b 100644 --- a/docs/refactor/auth-payment-provider-plan.md +++ b/docs/refactor/auth-payment-provider-plan.md @@ -281,6 +281,8 @@ body: { "code": "", "redirectUri": "https://h5.example.com/auth/q - `GET /api/commerce/reconciliation/batches`:查询对账批次。 - `GET /api/commerce/reconciliation/items`:查询逐行结果。 - `GET /api/commerce/reconciliation/anomalies`:查询金额不一致、状态不一致、供应商有本地无、本地有供应商无、重复行等异常。 +- `POST /api/commerce/reconciliation/provider-bills/request`:创建微信/支付宝官方账单下载任务,body 为 `{ provider, billDate, billType }`。 +- `GET /api/commerce/reconciliation/provider-bills/jobs`:查看官方账单下载任务状态、下载域名、hash、行数和生成的对账批次 ID。 - `POST /api/commerce/reconciliation/issues/create`:从异常明细创建差错工单,重复创建同一未关闭明细会返回现有工单。 - `GET /api/commerce/reconciliation/issues`:按状态、严重级别、负责人、批次、订单号筛选差错工单。 - `POST /api/commerce/reconciliation/issues/status`:执行 `start/assign/resolve/ignore/escalate/reopen` 状态流转。 @@ -290,8 +292,11 @@ body: { "code": "", "redirectUri": "https://h5.example.com/auth/q - `tenant:reconciliation:read`:查看/预览对账。 - `tenant:reconciliation:write`:导入对账批次、创建/处理差错工单。 +- `tenant:reconciliation:download`:创建微信/支付宝官方账单下载任务。 -对账和差错工单只生成差异台账、处理记录和审计,不自动修改订单、支付、退款和权益。`resolve/ignore` 只是财务审核结论,例如 `manual_adjustment`、`provider_confirmed`、`false_positive`,最终落账仍必须走退款状态机、支付补偿、手工支付确认或后续专门的人工调整命令。真实生产中,微信/支付宝官方账单下载 adapter 应复用同一套 `commerce_reconciliation_batches/items`,下载后的 CSV/JSON 先规范化为 `rows`,再调用同一套匹配逻辑。官方账单自动下载、人工调整凭证附件、财务复核报表和异常订单运营台仍需要继续补。 +官方账单下载由 `apps/worker --job provider-bills` 执行。API 只创建 `commerce_bill_download_jobs`,不会在前端返回供应商 `download_url`、商户私钥、微信 API v3 key 或支付宝应用私钥。worker 使用租户 `tenant_payment_accounts.config_public.secretRef` 找到 `app_private.tenant_secrets(secret_scope='payment')`,在后端签名申请下载 URL,校验微信返回的 `hash_type/hash_value`,解析 JSON/CSV/ZIP 账单后复用 `importReconciliationBatch` 写入 `commerce_reconciliation_batches/items`,`source='provider_download'`。 + +对账和差错工单只生成差异台账、处理记录和审计,不自动修改订单、支付、退款和权益。`resolve/ignore` 只是财务审核结论,例如 `manual_adjustment`、`provider_confirmed`、`false_positive`,最终落账仍必须走退款状态机、支付补偿、手工支付确认或后续专门的人工调整命令。人工调整凭证附件、财务复核报表、异常订单运营台和真实生产账单格式抽样验收仍需要继续补。 B 端合作商年费、服务费、服务器资源费不走学生端 `orders`,而是走平台账务: diff --git a/docs/refactor/backend-capability-status.md b/docs/refactor/backend-capability-status.md index 557307c4..3cfa8f68 100644 --- a/docs/refactor/backend-capability-status.md +++ b/docs/refactor/backend-capability-status.md @@ -121,7 +121,8 @@ | 支付/退款补偿 worker | 可联调 | `apps/worker --job commerce` 查询微信/支付宝订单和处理中退款,补偿漏通知支付、补发权益、确认退款、全额退款撤销权益;`npm run test:worker:commerce` 覆盖幂等和密钥不泄露 | | 资金流水对账 | 可联调 | `commerce_reconciliation_batches/items` + `/api/commerce/reconciliation/preview/import/batches/items/anomalies`;租户后台需 `tenant:reconciliation:read/write`,支持支付/退款账单行手工或 API 导入、来源 hash、批次统计、逐行匹配、金额/状态差异、本地缺失、供应商缺失、重复行、无效行和审计;对账只生成差异,不自动改订单/权益 | | 对账差错工单 | 可联调 | `commerce_reconciliation_issues/events` + `/api/commerce/reconciliation/issues*`;异常明细可创建工单,支持分配、开始处理、升级、解决、忽略、重开、事件轨迹和审计;处理结论只作为财务审核记录,不直接修改订单、支付、退款或权益 | -| 官方账单下载和异常运营台 | 待补齐 | 后续补微信/支付宝账单自动下载、人工调整凭证附件/复核、异常订单运营台和财务报表 | +| 官方账单下载 | 可联调 | `commerce_bill_download_jobs` + `/api/commerce/reconciliation/provider-bills/request/jobs` + `apps/worker --job provider-bills`;租户后台需 `tenant:reconciliation:download` 创建任务,worker 后端使用租户商户密钥申请微信/支付宝官方账单下载 URL、校验 hash、解析 JSON/CSV/ZIP 账单并复用同一套 `provider_download` 对账导入;响应只暴露任务状态、下载域名、hash 和对账批次 ID,不暴露下载 URL 或密钥 | +| 异常订单运营台和财务报表 | 待补齐 | 后续补人工调整凭证附件/复核、异常订单运营台、财务复核报表和真实生产账单格式抽样验收 | ## 租户后台与平台后台 diff --git a/docs/refactor/backend-handoff-roadmap.md b/docs/refactor/backend-handoff-roadmap.md index 93545430..8a01124d 100644 --- a/docs/refactor/backend-handoff-roadmap.md +++ b/docs/refactor/backend-handoff-roadmap.md @@ -29,7 +29,7 @@ | 分数线 | 可联调 | 院校、专业、动态字段、记录、年份、趋势、后台维护、JSON/CSV/Excel 导入 | 复杂筛选、AI 择校上下文 | | 视频解析 | 可联调 | 单题视频、批量查询、后台视频绑定、JSON/CSV/Excel 导入、会员播放权限、播放次数扣减、签名 URL、播放日志和动态水印上下文 | 深度防盗链、转码级水印、播放统计 | | 资料下载 | 部分完成 | 资源台账、SVIP 权限校验、`local_dev`/阿里云 OSS/腾讯 COS/Supabase Storage 上传下载签名、上传确认、PDF/图片预览签名、访问审计、动态水印上下文、assets worker 复检、内置安全扫描、外部 HTTP scanner 接入层、题库导出 PDF/Word/每日一练 ZIP 可生成可信 `content_assets` 并走签名下载/预览 | CDN 防盗链、真实 AV/内容安全服务联调、转码级水印、资料前端操作体验 | -| 会员与订单 | 可联调 | 下单、订单详情/状态轮询、优惠券领取/抵扣、零元订单自动开通、手工确认权限保护、激活码预检查/兑换、微信支付、支付宝、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、支付/退款补偿 worker、权益发放、资金对账批次/明细/异常查询 API、对账差错工单和事件轨迹 | 微信/支付宝官方账单自动下载、财务复核报表和异常订单运营台 | +| 会员与订单 | 可联调 | 下单、订单详情/状态轮询、优惠券领取/抵扣、零元订单自动开通、手工确认权限保护、激活码预检查/兑换、微信支付、支付宝、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、支付/退款补偿 worker、权益发放、资金对账批次/明细/异常查询 API、微信/支付宝官方账单下载 worker、对账差错工单和事件轨迹 | 财务复核报表、异常订单运营台和真实生产账单格式验收 | | 登录认证 | 可联调 | 短信 mock、阿里云/腾讯云短信 adapter、迁移期 session、Supabase Auth JWT、微信小程序登录、微信网页登录、QQ 登录、手机号绑定/换绑、OAuth 配置表 | 真实生产账号和回调域名联调 | | 销售/代理/CRM | 可联调 | 邀请码、首绑保护、团队关系、销售统计、CRM 配置/队列、钉钉/飞书/企微 worker、分佣规则、成员分佣比例、订单/激活码归因、结算生成、审核和打款状态;Taro 营销中心已接第一版操作台 | 小程序码真实生成、CRM 分配策略、结算导出、真实打款、凭证、财务复核和销售转化看板 | | 内容导入 | 可联调 | 题目、单词、知识手册、分数线、视频 JSON/CSV/Excel preview/import、issue、job/detail、审计、幂等、`executionMode=async`、imports worker、导入后复检、模板下载、字段映射 API、字段映射覆盖白名单校验、PocketBase JSON dry-run 报告;Taro 租户内容页已接上传/粘贴预览、模板文件下载、字段别名编辑、同步/异步执行、异步轮询和复检详情第一版 | 真实数据 dry-run 执行验收、抽样校验和导入性能压测 | @@ -83,7 +83,7 @@ ### P1:商用收费和运营能力 -- 资金对账已支持手工/API 账单导入比对、异常查询和差错工单处理;继续补微信/支付宝官方账单下载、财务复核报表和异常订单运营台。 +- 资金对账已支持手工/API 账单导入比对、微信/支付宝官方账单下载任务、异常查询和差错工单处理;继续补真实生产账单格式验收、财务复核报表和异常订单运营台。 - XPay 或其它实际支付网关 adapter。 - 阿里云/腾讯云短信、微信小程序登录、微信网页登录、QQ 登录真实账号联调。 - 公共题库/地区题库自动同步 worker 已具备单批执行能力,租户后台已有同步通知、单条/批量冲突采纳平台或保留本地操作;继续补生产定时调度、失败告警,以及租户按 SaaS 套餐购买地区、科目和题库范围的更细计费策略。 diff --git a/docs/refactor/backend-progress.md b/docs/refactor/backend-progress.md index 3c8c29a4..c4ec26fb 100644 --- a/docs/refactor/backend-progress.md +++ b/docs/refactor/backend-progress.md @@ -269,7 +269,7 @@ GET /api/tenant-admin/audit-logs 1. 完善内容导入和文件上传:字段映射 UI、真实数据 dry-run、PDF 预览渲染、防盗链、真实 AV/内容安全服务联调和转码/CDN 级视频水印。 2. 完成真实短信 provider 联调:阿里云/腾讯云,密钥放 `app_private.tenant_secrets` 或生产 Vault。 3. 完成真实 OAuth provider 联调:微信网页、微信小程序、QQ,确认回调域名、开放平台账号和旧 PocketBase 身份映射策略。 -4. 补微信/支付宝官方账单自动下载、异常订单运营台和优惠券核销报表;支付/退款补偿、退款查询确认、退款通知、资金对账导入比对和差错工单主链路已完成。 +4. 补异常订单运营台、真实生产账单格式验收和优惠券核销报表;支付/退款补偿、退款查询确认、退款通知、微信/支付宝官方账单下载、资金对账导入比对和差错工单主链路已完成。 5. 扩展 `apps/worker`:日报统计、CRM 死信告警、公共题库同步失败告警和更完整冲突处理运营台;公共题库同步 worker 和同步通知已具备基础闭环。 6. 开始 Taro scaffold,把 `supabaseApi` 抽到跨端包或适配层。 diff --git a/docs/refactor/blueprint-coverage.md b/docs/refactor/blueprint-coverage.md index 6616adcd..c342469c 100644 --- a/docs/refactor/blueprint-coverage.md +++ b/docs/refactor/blueprint-coverage.md @@ -30,7 +30,7 @@ | CRM 系统 | 可联调 | CRM 配置、密钥私密存储、客资入队、队列查询、generic/钉钉/飞书/企微 worker、签名、重试和日志 | 定向/轮询分配、富卡片模板、失败告警、死信运营台 | | 数据看板 | 可联调 | 租户 dashboard 聚合接口,收益、注册、学习、内容、激活码、反馈、趋势、24h 活跃、套餐销量和运营动态 | 预聚合 worker、缓存、慢 SQL 监控和销售转化看板 | | 登录认证 | 可联调 | 短信 mock、阿里云/腾讯云短信 adapter、迁移期 session、Supabase Auth JWT、微信小程序登录、微信网页登录、QQ 登录、手机号绑定/换绑、OAuth 配置表 | 真实生产账号和回调域名联调 | -| 支付 | 可联调 | 订单、支付记录、手动确认权限保护、权益发放、租户商户配置、微信支付 JSAPI、支付宝 WAP/H5、webhook 幂等、退款状态机、退款通知、补偿 worker、资金对账导入比对、异常查询、差错工单和事件轨迹 | 官方账单自动下载、异常订单运营台、服务商/平台代收模式 | +| 支付 | 可联调 | 订单、支付记录、手动确认权限保护、权益发放、租户商户配置、微信支付 JSAPI、支付宝 WAP/H5、webhook 幂等、退款状态机、退款通知、补偿 worker、微信/支付宝官方账单下载 worker、资金对账导入比对、异常查询、差错工单和事件轨迹 | 异常订单运营台、真实生产账单格式验收、服务商/平台代收模式 | | AI 择校推荐 | 部分完成 | SVIP 门禁、学生输入 schema、地区/分数线上下文、`local_rules` 稳定 JSON 报告、报告台账/列表/详情、Taro 学生端基础页 | 真实 AI provider、prompt 版本管理后台、PDF 报告生成、人工复核和运营配置 | | Taro 跨端 | 未开始 | 旧 Web 新 API 适配开始 | `apps/taro`、共享 API client、H5/小程序统一构建 | diff --git a/docs/refactor/implementation-status.md b/docs/refactor/implementation-status.md index 19bfec8b..66f4c9e5 100644 --- a/docs/refactor/implementation-status.md +++ b/docs/refactor/implementation-status.md @@ -27,7 +27,7 @@ | 刷题题库 | 已建题库、题目、题目版本、内容入口、任意深度分类树、考试意向标记、题目集合、练习蓝图、导入任务台账、导出任务台账、公共题库授权/采纳表、租户内容通知表 | 已支持核心映射,JSON/CSV/Excel 导入可落到新入口/节点/集合,阅读理解/案例分析子题沿用 `subQuestions/sub_questions` | 题目列表、内容入口、分类树、集合题目、顺序/随机/全真模拟 session、答题提交、复合题 `subAnswers` 判分和报告明细、租户后台题目录入/更新、JSON/CSV/Excel 预览/导入、JSON/试卷 payload 导出、PDF/Word 异步导出 worker、每日一练九宫格 metadata、PDF/Word 运营版式、ZIP 图片素材包、异步导入 worker、平台公共题库授权、租户采纳快照、手动同步、自动同步 worker、同步通知、冲突查询和单条/批量冲突处理已实现 | 核心 API 集成测试含导航、组卷、复合题后台录入/练习/判分/报告、导入、导出权限/脱敏、每日一练导出 metadata、异步 PDF/Word/每日一练 ZIP job 创建、exports worker、公共题库授权、采纳后组卷、同步新增题、通知隔离/已读/自动 resolved、租户自改冲突保护、单条/批量冲突处理和 worker 自动同步断言 | 新题库导航和组卷基础闭环可跑,阅读理解/案例分析多小题第一版可联调,公共题库采纳/手动/自动同步、同步通知、冲突查询/处理、导入后复检、模板下载、字段映射 API、JSON/PDF/Word/每日一练 ZIP 基础导出可联调;公共题库生产调度/失败告警、更精细导出模板和更完整运营消息仍需补齐 | | 错题本 | 已建 `wrong_questions` | 已支持旧错题归一化 | 错题列表、答题自动入错题、移出错题已实现 | 仅烟测 | 基础功能已实现,复习计划和统计未完成 | | 收藏夹 | 已建 `favorite_questions` | 已支持旧收藏归一化 | 收藏/取消收藏、收藏列表已实现 | 仅烟测 | 基础功能已实现 | -| 用户订阅/题库会员/SVIP | 已建 `orders`、`payments`、`entitlements`、`svip_plans`、激活码 | 已映射旧 SVIP/会员权益 | 下单、订单详情/状态轮询、手工支付确认权限保护、微信/支付宝支付、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、激活码预检查/兑换、优惠券抵扣、零元订单自动开通、权益查询、支付/退款补偿、资金对账和差错工单已实现 | API 集成测试 | 商城主链路可联调,官方账单自动下载、异常订单运营台和生产账号联调待补 | +| 用户订阅/题库会员/SVIP | 已建 `orders`、`payments`、`entitlements`、`svip_plans`、激活码 | 已映射旧 SVIP/会员权益 | 下单、订单详情/状态轮询、手工支付确认权限保护、微信/支付宝支付、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、激活码预检查/兑换、优惠券抵扣、零元订单自动开通、权益查询、支付/退款补偿、微信/支付宝官方账单下载、资金对账和差错工单已实现 | API 集成测试、commerce worker 集成测试 | 商城主链路可联调,异常订单运营台、真实生产账单格式验收和生产账号联调待补 | | 背单词 | 已建单词单元、单词、进度、收藏表,并可绑定 `content_entries/content_nodes` | 已支持内容和部分用户状态映射 | 单元/单词只读、进度、收藏、统计、每日复习计划、租户后台单词维护 API、旧模板/新模板 JSON 预览导入、排行榜已实现 | 核心 API 集成测试含导入和排行榜断言 | 学生端学习状态、后台维护、批量 JSON 导入和基础排行榜已实现,更细复习参数和后台统计待完善 | | 知识手册 | 已建手册科目、章节、条目,并可绑定 `content_entries/content_nodes` | 已支持内容导入 | 只读 API、租户后台手册科目/章节/条目维护 API、嵌套 JSON 预览导入已实现 | 核心 API 集成测试含导入断言 | 学生端阅读、后台维护和批量 JSON 导入基础可用,富文本资源/版本管理待补 | | 分数线 | 已建院校、专业、字段、记录表 | 已支持导入映射 | 字段、院校、专业、记录、趋势、年份、租户后台维护 API、JSON 预览导入已实现 | 核心 API 集成测试含导入断言 | 查询、后台维护和批量 JSON 导入基础闭环已实现,复杂动态筛选和 AI 择校上下文待补 | @@ -300,4 +300,4 @@ platform-admin: 3. 补学习统计增强:排行榜防刷/预聚合、断点续练、专项练习策略和更细题型分析。 4. 补视频商用控制:深度防盗链、转码级水印和播放统计。 5. 补 AI 择校推荐报告、排行榜防刷/预聚合、勋章自动发放。 -6. 接真实短信/OAuth 生产账号、微信/支付宝官方账单下载和异常订单运营台,并继续推进 Taro scaffold。 +6. 接真实短信/OAuth 生产账号、真实生产账单格式验收和异常订单运营台,并继续推进 Taro scaffold。 diff --git a/docs/refactor/legacy-feature-gap-matrix.md b/docs/refactor/legacy-feature-gap-matrix.md index 0911e2cd..87ebadd6 100644 --- a/docs/refactor/legacy-feature-gap-matrix.md +++ b/docs/refactor/legacy-feature-gap-matrix.md @@ -30,7 +30,7 @@ | 背单词 | `VocabularyPage.tsx`、`VocabularyQuiz.tsx` | 部分覆盖 | 单词列表、进度、收藏、统计、每日计划和后端复习调度已覆盖;后续补收藏练习体验、发音/音频策略、排行榜和更精细的间隔算法参数 | | 知识手册 | `Handbook*.tsx` | 已覆盖 | 前端需做好 Markdown/公式/图片渲染和搜索体验 | | 分数线 | `ScorelinePage.tsx` | 已覆盖 | 动态字段/趋势、后台维护和 JSON 批量导入已有;后续补复杂筛选优化和 AI 择校数据上下文 | -| 商城/SVIP | `Store.tsx`、`SvipModal.tsx` | 部分覆盖 | 套餐、订单、订单详情/状态轮询、权益、激活码预检查/兑换、优惠券领取/下单抵扣、微信支付/支付宝 provider 主链路、内部退款状态机、微信/支付宝发起退款、退款查询确认、退款通知 webhook、支付/退款补偿 worker、全额退款权益撤销、资金对账手工/API 导入比对、异常查询、差错工单和事件轨迹已有;缺微信/支付宝官方账单自动下载、异常订单运营台和前端收银台/售后体验 | +| 商城/SVIP | `Store.tsx`、`SvipModal.tsx` | 部分覆盖 | 套餐、订单、订单详情/状态轮询、权益、激活码预检查/兑换、优惠券领取/下单抵扣、微信支付/支付宝 provider 主链路、内部退款状态机、微信/支付宝发起退款、退款查询确认、退款通知 webhook、支付/退款补偿 worker、全额退款权益撤销、资金对账手工/API 导入比对、微信/支付宝官方账单下载 worker、异常查询、差错工单和事件轨迹已有;缺异常订单运营台、人工调整凭证复核报表和前端收银台/售后体验 | | 个人中心 | `Profile.tsx` | 部分覆盖 | 基本资料、手机号绑定/换绑、权益、订单统计、练习历史、学习统计、签到积分、考试倒计时、趋势和勋章展示 API 已有;缺学习报告可视化 | | 资料下载 | `QuestionExporterPublishModal.tsx` 等 | 部分覆盖 | 资源台账、上传确认、签名下载、PDF/图片预览、动态水印上下文、worker 复检、内置安全扫描和外部 HTTP scanner 接入层已有;缺前端水印渲染、深度防盗链、真实 AV/内容安全服务联调和生命周期策略 | | AI 择校推荐 | 业务规划新增 | 部分覆盖 | 已有 SVIP 门禁、学生输入 schema、地区/分数线上下文、稳定 JSON 输出、报告台账、审计和 Taro 学生端基础页;真实 AI provider、prompt 版本管理后台、报告 PDF 渲染和更细推荐算法待补 | diff --git a/docs/refactor/next-development-todo.md b/docs/refactor/next-development-todo.md index 78f123a1..90713eb2 100644 --- a/docs/refactor/next-development-todo.md +++ b/docs/refactor/next-development-todo.md @@ -28,6 +28,7 @@ - 租户后台数据看板已完成首版聚合 API:`GET /api/tenant-admin/dashboard`,支持租户/地区维度的收益、注册、学习、内容、激活码、反馈、趋势、24h 活跃、套餐销量和运营动态,前端可直接联调。 - 支付/退款补偿 worker 已完成:`apps/worker --job commerce` 可查询微信/支付宝支付和处理中退款,补偿漏通知订单,支付成功幂等开通权益,退款成功幂等更新退款/订单/支付并在全额退款时撤销订单权益。 - 资金对账和差错工单闭环已完成:`commerce_reconciliation_batches/items` 和 `/api/commerce/reconciliation/*` 支持手工/API 导入供应商账单行、预览差异、生成批次统计、查询异常、租户隔离、权限点 `tenant:reconciliation:read/write` 和审计日志;`commerce_reconciliation_issues/events` 支持异常明细创建工单、分配、开始处理、升级、解决、忽略、重开和事件留痕,且不直接修改订单/支付/退款/权益。 +- 微信/支付宝官方账单下载地基已完成:`commerce_bill_download_jobs`、`POST /api/commerce/reconciliation/provider-bills/request`、`GET /api/commerce/reconciliation/provider-bills/jobs` 和 `apps/worker --job provider-bills` 已接入,worker 负责后端签名申请下载 URL、hash 校验、JSON/CSV/ZIP 账单解析、复用 `provider_download` 对账导入、任务状态回写和密钥脱敏。 - 内容资源复检与安全扫描 worker 已完成:`apps/worker --job assets` 可复检 `content_assets` 中的托管对象元数据,并执行内置 `metadata_rules` 和可选外部 HTTP scanner;正常资源写回复检/扫描证据,异常资源自动置为 `failed/skipped + draft` 或 `security_scan_status=failed`,外部 scanner 不可用默认 fail-closed,并写入审计、扫描事件和安全标记。 - 题库导出 worker 已完成:`apps/worker --job exports` 可抢占 `pdf/docx/daily_practice_zip` 导出任务,渲染 PDF/Word、水印或每日一练图片素材包,写入对象存储或本地开发存储,创建 `content_assets` 并回填 `assetId/hash/size`;`daily_practice` 已支持每日一练九宫格 metadata、PDF/Word 基础版式、9 张 PNG/SVG 卡片和拼图 ZIP。 - 租户后台媒体运营报表已完成:`/api/tenant-content/media-analytics/summary`、`asset-events`、`video-events` 可按 7/30/90 天、资源、视频、用户和水印 traceId 查询资料访问、视频播放、拒绝访问、Top 资源/视频和日趋势;仅开放给 owner/admin/operator 或 `content:analytics:read` 权限角色,前端不会拿到签名 URL 或播放 token。 @@ -75,7 +76,7 @@ - 已完成微信支付 JSAPI、支付宝 WAP/H5 的创建支付参数和 webhook 幂等开通权益。 - 已完成内部退款状态机、退款申请/审核/处理接口、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、部分/全额退款状态、全额退款权益撤销和审计事件。 - 已完成支付/退款补偿 worker,可兜底供应商漏通知、处理中退款和重复执行幂等。 - - 已完成资金对账手工/API 导入比对、批次/明细/异常查询、差错工单状态流和审计;继续补微信/支付宝官方账单自动下载、人工调整凭证附件、财务复核报表和异常订单运营台。 + - 已完成资金对账手工/API 导入比对、微信/支付宝官方账单下载任务、批次/明细/异常查询、差错工单状态流和审计;继续补真实生产账单格式抽样验收、人工调整凭证附件、财务复核报表和异常订单运营台。 - 租户自有商户收款和平台代收/服务商模式。 2. 国内登录和短信 @@ -114,7 +115,7 @@ 8. 订单和营销体验 - 已完成订单详情、订单状态轮询、激活码预检查、优惠券前台领取、下单抵扣计算和内部退款状态机。 - - 已完成支付/退款补偿 worker、资金对账导入比对和差错工单;继续补异常订单运营台、优惠券核销报表和复杂活动规则。 + - 已完成支付/退款补偿 worker、官方账单下载 worker、资金对账导入比对和差错工单;继续补异常订单运营台、优惠券核销报表和复杂活动规则。 9. 积分和反馈增强 - 已完成每日签到、积分流水、反馈提交、租户后台处理、奖励积分幂等。 @@ -228,5 +229,5 @@ 3. 补平台后台增强:租户详情/编辑、平台审计报表、自动计费、账单批量操作和更细平台权限点。 4. 云服务器部署 Supabase/PostgreSQL 和 API,配置对象存储生产环境变量,跑 `check:refactor` 的远程等价测试。 5. 导出现有 PocketBase 数据,做完整 dry-run 迁移。 -6. 并行补真实登录、微信/支付宝官方账单自动下载、异常订单运营台、对象存储真实 AV/内容安全服务联调、转码/CDN 级水印/生命周期、题库导出模板精排/操作台、公共题库生产定时调度和失败告警。 +6. 并行补真实登录、真实生产账单格式验收、异常订单运营台、对象存储真实 AV/内容安全服务联调、转码/CDN 级水印/生命周期、题库导出模板精排/操作台、公共题库生产定时调度和失败告警。 7. 前后端联调通过后,再做支付、权限、数据导入、资料下载、视频播放的商用验收。 diff --git a/docs/refactor/taro-frontend-integration.md b/docs/refactor/taro-frontend-integration.md index 6f8a4920..5b709e86 100644 --- a/docs/refactor/taro-frontend-integration.md +++ b/docs/refactor/taro-frontend-integration.md @@ -1853,7 +1853,7 @@ approved/processing -> failed - 退款通知地址由支付账户或 `submit_provider_refund.providerNotifyUrl` 配置,后端公开接收路径为 `POST /api/commerce/refunds/notify/wechat_pay?tenantId=`、`POST /api/commerce/refunds/notify/alipay?tenantId=`。这是支付平台回调地址,Taro 前端不要主动调用。 - 退款通知只会推进已经审核/处理中的退款申请;未审核的 `requested` 退款不能被外部通知直接落账。 - 已经 `succeeded` 的退款不能再次查询或再次标记成功,避免订单退款金额重复累加。前端应按接口返回状态展示,不要假设点击后立即到账。 -- 自动补偿 worker 已接入:支付漏通知和处理中退款会由后端定时查询供应商并幂等落账。资金对账已支持租户后台手工/API 导入供应商账单、查询差异和差错工单处理;官方账单自动下载和异常订单运营台后续继续补。生产联调时仍需保留人工确认/失败登记入口。 +- 自动补偿 worker 已接入:支付漏通知和处理中退款会由后端定时查询供应商并幂等落账。资金对账已支持租户后台手工/API 导入供应商账单、微信/支付宝官方账单下载任务、查询差异和差错工单处理;异常订单运营台和人工调整凭证复核后续继续补。生产联调时仍需保留人工确认/失败登记入口。 ### 租户后台资金对账 @@ -1955,6 +1955,66 @@ GET /api/commerce/reconciliation/issues/events?issueId= 权限:tenant:reconciliation:read ``` +官方账单下载: + +```text +POST /api/commerce/reconciliation/provider-bills/request +权限:tenant:reconciliation:download +body: { + "provider": "wechat_pay | alipay", + "billDate": "2026-06-29", + "billType": "payment | refund | combined", + "metadata": { + "remark": "财务手动申请" + } +} +``` + +返回: + +```json +{ + "item": { + "id": "", + "provider": "wechat_pay", + "billDate": "2026-06-29", + "billType": "payment", + "status": "queued", + "sourceName": "provider-bill:wechat_pay:2026-06-29:payment", + "sourceHash": null, + "rowCount": 0, + "downloadUrlHost": null, + "reconciliationBatchId": null + }, + "idempotent": false +} +``` + +轮询任务: + +```text +GET /api/commerce/reconciliation/provider-bills/jobs?provider=wechat_pay&billDate=2026-06-29 +权限:tenant:reconciliation:read +``` + +任务状态: + +```text +queued 已排队,等待 provider-bills worker +running worker 正在申请和下载官方账单 +completed 已下载、校验、导入对账批次 +failed 下载、hash 校验、解析或导入失败 +cancelled 已取消 +``` + +前端处理规则: + +- 前端只创建任务和轮询状态,不直接请求微信/支付宝账单 URL。 +- 后端响应只会返回 `downloadUrlHost`,不会返回完整下载 URL、商户私钥、API v3 key、支付宝应用私钥。 +- `completed` 后用 `reconciliationBatchId` 跳转到对账批次明细。 +- `failed` 时展示 `errorCode/errorMessage`,让财务重新发起或联系技术处理。 +- 官方账单下载由服务器定时运行 `apps/worker --job provider-bills`,前端不要自行触发供应商接口。 + 工单状态: ```text @@ -1995,7 +2055,7 @@ ignored 无效行或不符合本次 billType - `sourceHash` 可作为同一文件内容的识别线索,但当前接口不会阻止重复导入;前端应展示最近同名/同 hash 批次提醒。 - 金额统一是分,前端不要传元。 - 对账差异和差错工单只是运营判断依据,`resolve/ignore` 不会落账。最终订单修正必须走退款、补偿、人工确认或后续专门的人工调整接口。 -- 当前后端支持 JSON 行导入;CSV/Excel 可以先由前端或后续后端 parser 转成上述 `rows`。微信/支付宝官方账单自动下载仍是后续后端任务。 +- 当前后端支持 JSON 行手工导入;官方账单下载 worker 支持微信/支付宝账单 JSON/CSV/ZIP 解析,并复用同一套对账导入逻辑。真实生产接入时仍要用真实账单文件抽样验收字段映射。 ### 激活码预检查与兑换 diff --git a/package-lock.json b/package-lock.json index 3eb2f070..97246523 100644 --- a/package-lock.json +++ b/package-lock.json @@ -88,6 +88,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" diff --git a/scripts/api-integration-test.js b/scripts/api-integration-test.js index a93051f8..5ff63eff 100644 --- a/scripts/api-integration-test.js +++ b/scripts/api-integration-test.js @@ -2484,6 +2484,63 @@ async function testCommerce() { }); assert.equal(crossTenantReconciliationDenied.code, 'TENANT_ADMIN_REQUIRED', 'reconciliation items must be tenant isolated'); + const studentProviderBillRequestDenied = await request('/api/commerce/reconciliation/provider-bills/request', { + method: 'POST', + body: { + provider: 'wechat_pay', + billDate, + billType: 'payment', + }, + expectStatus: 403, + }); + assert.equal(studentProviderBillRequestDenied.code, 'TENANT_ADMIN_REQUIRED', 'students must not request official provider bill downloads'); + + const providerBillJob = await request('/api/commerce/reconciliation/provider-bills/request', { + userId: TENANT_ADMIN_USER_ID, + method: 'POST', + body: { + provider: 'wechat_pay', + billDate, + billType: 'payment', + metadata: { source: 'api-integration-test' }, + }, + }); + assert.ok(providerBillJob.item?.id, 'tenant admin should request official provider bill download job'); + assert.equal(providerBillJob.item?.status, 'queued', 'new provider bill job should be queued'); + assert.equal(providerBillJob.item?.provider, 'wechat_pay', 'provider bill job should keep provider'); + assert.ok(!JSON.stringify(providerBillJob).includes(paymentFixture.wechatApiV3Key), 'provider bill job response must not leak payment secrets'); + assert.ok(!JSON.stringify(providerBillJob).includes('PRIVATE KEY'), 'provider bill job response must not leak private keys'); + + const providerBillJobAgain = await request('/api/commerce/reconciliation/provider-bills/request', { + userId: TENANT_ADMIN_USER_ID, + method: 'POST', + body: { + provider: 'wechat_pay', + billDate, + billType: 'payment', + }, + }); + assert.equal(providerBillJobAgain.item?.id, providerBillJob.item.id, 'provider bill job request should be idempotent by provider/date/type'); + assert.equal(providerBillJobAgain.idempotent, true, 'provider bill duplicate request should return idempotent=true'); + + const providerBillJobs = await request('/api/commerce/reconciliation/provider-bills/jobs', { + userId: TENANT_ADMIN_USER_ID, + query: { provider: 'wechat_pay', billDate }, + }); + assert.ok( + providerBillJobs.items?.some(item => item.id === providerBillJob.item.id), + 'provider bill jobs endpoint should list tenant jobs', + ); + assert.ok(!JSON.stringify(providerBillJobs).includes(paymentFixture.wechatApiV3Key), 'provider bill jobs list must not leak payment secrets'); + + const crossTenantProviderBillJobsDenied = await request('/api/commerce/reconciliation/provider-bills/jobs', { + tenantId: PARTNER_TENANT_ID, + userId: TENANT_ADMIN_USER_ID, + query: { provider: 'wechat_pay', billDate }, + expectStatus: 403, + }); + assert.equal(crossTenantProviderBillJobsDenied.code, 'TENANT_ADMIN_REQUIRED', 'provider bill jobs must be tenant isolated'); + const fakeWechatPay = await startFakeWechatPayServer(); const wechatAccount = await request('/api/tenant-admin/payment-accounts', { userId: TENANT_ADMIN_USER_ID, diff --git a/scripts/commerce-worker-integration-test.js b/scripts/commerce-worker-integration-test.js index 7ac18bf3..77017574 100644 --- a/scripts/commerce-worker-integration-test.js +++ b/scripts/commerce-worker-integration-test.js @@ -17,18 +17,29 @@ const ids = { refundPayment: '10000000-0000-0000-0000-000000000704', refundRequest: '10000000-0000-0000-0000-000000000705', refundEntitlement: '10000000-0000-0000-0000-000000000706', + providerBillJobWechat: '10000000-0000-0000-0000-000000000707', + providerBillJobAlipay: '10000000-0000-0000-0000-000000000708', + providerBillOrderWechat: '10000000-0000-0000-0000-000000000709', + providerBillPaymentWechat: '10000000-0000-0000-0000-000000000710', + providerBillOrderAlipay: '10000000-0000-0000-0000-000000000711', + providerBillPaymentAlipay: '10000000-0000-0000-0000-000000000712', }; const orderNos = { payment: 'CW-WX-PAY-001', refund: 'CW-WX-REFUND-001', + providerBillWechat: 'CW-WX-BILL-001', + providerBillAlipay: 'CW-ALI-BILL-001', }; const refundNo = 'RF-CW-WX-001'; +const providerBillDate = '2026-06-29'; const paymentFixture = (() => { const wechatMerchant = crypto.generateKeyPairSync('rsa', { modulusLength: 2048 }); + const alipayApp = crypto.generateKeyPairSync('rsa', { modulusLength: 2048 }); return { wechatMerchantPrivateKey: wechatMerchant.privateKey.export({ type: 'pkcs8', format: 'pem' }).toString(), + alipayAppPrivateKey: alipayApp.privateKey.export({ type: 'pkcs8', format: 'pem' }).toString(), wechatApiV3Key: '12345678901234567890123456789012', }; })(); @@ -48,6 +59,25 @@ async function startFakeWechatPayServer() { const port = await getFreePort(); const baseUrl = `http://127.0.0.1:${port}`; const requests = []; + const billRows = [ + { + transactionType: 'payment', + orderNo: orderNos.providerBillWechat, + providerTradeNo: `wx-bill-trade-${orderNos.providerBillWechat}`, + amountCents: 990, + providerStatus: 'SUCCESS', + paidAt: '2026-06-29T10:00:00+08:00', + }, + { + transactionType: 'payment', + orderNo: 'CW-WX-PROVIDER-ONLY', + providerTradeNo: 'wx-provider-only-trade', + amountCents: 1888, + providerStatus: 'SUCCESS', + }, + ]; + const billBody = Buffer.from(JSON.stringify({ rows: billRows }), 'utf8'); + const billHash = crypto.createHash('sha256').update(billBody).digest('hex'); const server = http.createServer((req, res) => { const url = new URL(req.url || '/', baseUrl); let raw = ''; @@ -85,6 +115,22 @@ async function startFakeWechatPayServer() { return; } + if (req.method === 'GET' && url.pathname === '/v3/bill/tradebill') { + res.writeHead(200, { 'content-type': 'application/json' }); + res.end(JSON.stringify({ + hash_type: 'SHA256', + hash_value: billHash, + download_url: `${baseUrl}/download/wechat-tradebill.json`, + })); + return; + } + + if (req.method === 'GET' && url.pathname === '/download/wechat-tradebill.json') { + res.writeHead(200, { 'content-type': 'application/json' }); + res.end(billBody); + return; + } + res.writeHead(404, { 'content-type': 'application/json' }); res.end(JSON.stringify({ code: 'NOT_FOUND' })); }); @@ -96,6 +142,61 @@ async function startFakeWechatPayServer() { return { paymentQueryEndpoint: `${baseUrl}/v3/pay/transactions/out-trade-no`, refundEndpoint: `${baseUrl}/v3/refund/domestic/refunds`, + tradeBillEndpoint: `${baseUrl}/v3/bill/tradebill`, + requests, + close: () => new Promise(resolve => server.close(resolve)), + }; +} + +async function startFakeAlipayServer() { + const port = await getFreePort(); + const baseUrl = `http://127.0.0.1:${port}`; + const requests = []; + const billRows = [ + { + transactionType: 'payment', + orderNo: orderNos.providerBillAlipay, + providerTradeNo: `ali-bill-trade-${orderNos.providerBillAlipay}`, + amountCents: 1288, + providerStatus: 'TRADE_SUCCESS', + paidAt: '2026-06-29 11:00:00', + }, + ]; + const billBody = Buffer.from(JSON.stringify({ rows: billRows }), 'utf8'); + const server = http.createServer((req, res) => { + const url = new URL(req.url || '/', baseUrl); + let raw = ''; + req.on('data', chunk => { + raw += chunk.toString(); + }); + req.on('end', () => { + requests.push({ method: req.method, pathname: url.pathname, search: url.search, body: raw }); + if (req.method === 'POST' && url.pathname === '/gateway.do') { + res.writeHead(200, { 'content-type': 'application/json' }); + res.end(JSON.stringify({ + alipay_data_dataservice_bill_downloadurl_query_response: { + code: '10000', + msg: 'Success', + bill_download_url: `${baseUrl}/download/alipay-tradebill.json`, + }, + })); + return; + } + if (req.method === 'GET' && url.pathname === '/download/alipay-tradebill.json') { + res.writeHead(200, { 'content-type': 'application/json' }); + res.end(billBody); + return; + } + res.writeHead(404, { 'content-type': 'application/json' }); + res.end(JSON.stringify({ code: 'NOT_FOUND' })); + }); + }); + await new Promise((resolve, reject) => { + server.once('error', reject); + server.listen(port, '127.0.0.1', resolve); + }); + return { + endpoint: `${baseUrl}/gateway.do`, requests, close: () => new Promise(resolve => server.close(resolve)), }; @@ -129,7 +230,55 @@ async function runWorkerOnce() { return output; } +async function runProviderBillsWorkerOnce() { + const child = spawn(process.execPath, ['apps/worker/dist/apps/worker/src/index.js', '--once', '--job', 'provider-bills'], { + cwd: process.cwd(), + env: { + ...process.env, + DATABASE_URL: databaseUrl, + WORKER_PROVIDER_BILL_BATCH_SIZE: '10', + WORKER_PROVIDER_BILL_ID: 'provider-bills-integration-test', + WORKER_COMMERCE_REQUEST_TIMEOUT_MS: '5000', + }, + stdio: ['ignore', 'pipe', 'pipe'], + windowsHide: true, + }); + let output = ''; + child.stdout.on('data', chunk => { + output += chunk.toString(); + }); + child.stderr.on('data', chunk => { + output += chunk.toString(); + }); + const code = await new Promise(resolve => child.on('exit', resolve)); + assert.equal(code, 0, `provider bill worker should exit 0\n${output}`); + assert.match(output, /provider-bills batch processed=\d+/, 'worker output should include provider bill summary'); + assert.ok(!output.includes(paymentFixture.wechatApiV3Key), 'provider bill worker output must not leak WeChat API v3 key'); + assert.ok(!output.includes('PRIVATE KEY'), 'provider bill worker output must not leak private keys'); + return output; +} + async function cleanup(pool) { + await pool.query( + ` + delete from public.commerce_bill_download_jobs + where tenant_id = $1 + and ( + id in ($2::uuid, $3::uuid) + or source_name like 'provider-bill:%:2026-06-29:%' + ) + `, + [tenantId, ids.providerBillJobWechat, ids.providerBillJobAlipay], + ); + await pool.query( + ` + delete from public.commerce_reconciliation_batches + where tenant_id = $1 + and source = 'provider_download' + and source_name like 'provider-bill:%:2026-06-29:%' + `, + [tenantId], + ); await pool.query( ` delete from public.payment_events @@ -169,20 +318,20 @@ async function cleanup(pool) { await pool.query( ` delete from public.payments - where tenant_id = $1 and id in ($2::uuid, $3::uuid) + where tenant_id = $1 and id in ($2::uuid, $3::uuid, $4::uuid, $5::uuid) `, - [tenantId, ids.payment, ids.refundPayment], + [tenantId, ids.payment, ids.refundPayment, ids.providerBillPaymentWechat, ids.providerBillPaymentAlipay], ); await pool.query( ` delete from public.orders - where tenant_id = $1 and id in ($2::uuid, $3::uuid) + where tenant_id = $1 and id in ($2::uuid, $3::uuid, $4::uuid, $5::uuid) `, - [tenantId, ids.paymentOrder, ids.refundOrder], + [tenantId, ids.paymentOrder, ids.refundOrder, ids.providerBillOrderWechat, ids.providerBillOrderAlipay], ); } -async function seed(pool, fakeWechat) { +async function seed(pool, fakeWechat, fakeAlipay) { await pool.query( ` insert into app_private.tenant_secrets ( @@ -201,6 +350,23 @@ async function seed(pool, fakeWechat) { })], ); + await pool.query( + ` + insert into app_private.tenant_secrets ( + tenant_id, secret_scope, secret_key, secret_json, provider, last_rotated_at + ) + values ($1, 'payment', 'alipay', $2::jsonb, 'alipay', now()) + on conflict (tenant_id, secret_scope, secret_key) + do update set secret_json = excluded.secret_json, + provider = excluded.provider, + last_rotated_at = now(), + updated_at = now() + `, + [tenantId, JSON.stringify({ + privateKey: paymentFixture.alipayAppPrivateKey, + })], + ); + await pool.query( ` insert into public.tenant_payment_accounts ( @@ -221,10 +387,32 @@ async function seed(pool, fakeWechat) { paymentQueryEndpoint: fakeWechat.paymentQueryEndpoint, refundEndpoint: fakeWechat.refundEndpoint, refundQueryEndpoint: fakeWechat.refundEndpoint, + tradeBillEndpoint: fakeWechat.tradeBillEndpoint, secretRef: 'app_private.tenant_secrets:payment:wechat_pay', })], ); + await pool.query( + ` + insert into public.tenant_payment_accounts ( + tenant_id, provider, mode, display_name, status, config_public + ) + values ($1, 'alipay', 'tenant_collect', 'Worker Alipay', 'active', $2::jsonb) + on conflict (tenant_id, provider) + do update set mode = excluded.mode, + display_name = excluded.display_name, + status = excluded.status, + config_public = excluded.config_public, + updated_at = now() + `, + [tenantId, JSON.stringify({ + appId: 'ali-worker-appid', + endpoint: fakeAlipay.endpoint, + billEndpoint: fakeAlipay.endpoint, + secretRef: 'app_private.tenant_secrets:payment:alipay', + })], + ); + await pool.query( ` insert into public.orders ( @@ -234,9 +422,24 @@ async function seed(pool, fakeWechat) { ) values ($1, $2, $3, $4, 'pending', 'svip', 'Worker 补偿月卡', 990, 'jsapi', 'wechat_pay', $5, 30, $6, '{"source":"commerce-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes'), - ($7, $2, $3, $8, 'paid', 'svip', 'Worker 退款月卡', 500, 'jsapi', 'wechat_pay', $5, 30, $6, '{"source":"commerce-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes') + ($7, $2, $3, $8, 'paid', 'svip', 'Worker 退款月卡', 500, 'jsapi', 'wechat_pay', $5, 30, $6, '{"source":"commerce-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes'), + ($9, $2, $3, $10, 'paid', 'svip', 'Worker 微信账单月卡', 990, 'jsapi', 'wechat_pay', $5, 30, $6, '{"source":"provider-bill-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes'), + ($11, $2, $3, $12, 'paid', 'svip', 'Worker 支付宝账单月卡', 1288, 'wap', 'alipay', $5, 30, $6, '{"source":"provider-bill-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes') `, - [ids.paymentOrder, tenantId, studentUserId, orderNos.payment, planId, regionId, ids.refundOrder, orderNos.refund], + [ + ids.paymentOrder, + tenantId, + studentUserId, + orderNos.payment, + planId, + regionId, + ids.refundOrder, + orderNos.refund, + ids.providerBillOrderWechat, + orderNos.providerBillWechat, + ids.providerBillOrderAlipay, + orderNos.providerBillAlipay, + ], ); await pool.query( @@ -247,9 +450,24 @@ async function seed(pool, fakeWechat) { ) values ($1, $2, $3, 'wechat_pay', 'jsapi', 'pending', 990, null, null, '{"source":"commerce-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes'), - ($4, $2, $5, 'wechat_pay', 'jsapi', 'paid', 500, $6, now() - interval '9 minutes', '{"source":"commerce-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes') + ($4, $2, $5, 'wechat_pay', 'jsapi', 'paid', 500, $6, now() - interval '9 minutes', '{"source":"commerce-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes'), + ($7, $2, $8, 'wechat_pay', 'jsapi', 'paid', 990, $9, now() - interval '9 minutes', '{"source":"provider-bill-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes'), + ($10, $2, $11, 'alipay', 'wap', 'paid', 1288, $12, now() - interval '9 minutes', '{"source":"provider-bill-worker-test"}'::jsonb, now() - interval '10 minutes', now() - interval '10 minutes') `, - [ids.payment, tenantId, ids.paymentOrder, ids.refundPayment, ids.refundOrder, `wx-trade-${orderNos.refund}`], + [ + ids.payment, + tenantId, + ids.paymentOrder, + ids.refundPayment, + ids.refundOrder, + `wx-trade-${orderNos.refund}`, + ids.providerBillPaymentWechat, + ids.providerBillOrderWechat, + `wx-bill-trade-${orderNos.providerBillWechat}`, + ids.providerBillPaymentAlipay, + ids.providerBillOrderAlipay, + `ali-bill-trade-${orderNos.providerBillAlipay}`, + ], ); await pool.query( @@ -281,16 +499,52 @@ async function seed(pool, fakeWechat) { `, [ids.refundRequest, tenantId, ids.refundOrder, ids.refundPayment, refundNo], ); + + await pool.query( + ` + insert into public.commerce_bill_download_jobs ( + id, tenant_id, provider, bill_date, bill_type, status, + source_name, requested_by, metadata + ) + values + ($1, $3, 'wechat_pay', $4::date, 'payment', 'queued', $5, null, '{"source":"commerce-worker-test"}'::jsonb), + ($2, $3, 'alipay', $4::date, 'payment', 'queued', $6, null, '{"source":"commerce-worker-test"}'::jsonb) + on conflict (tenant_id, provider, bill_date, bill_type) + where status in ('queued', 'running', 'completed') + do update set status = 'queued', + source_name = excluded.source_name, + reconciliation_batch_id = null, + source_hash = null, + row_count = 0, + claimed_by = null, + claimed_at = null, + completed_at = null, + failed_at = null, + error_code = null, + error_message = null, + metadata = excluded.metadata, + updated_at = now() + `, + [ + ids.providerBillJobWechat, + ids.providerBillJobAlipay, + tenantId, + providerBillDate, + `provider-bill:wechat_pay:${providerBillDate}:payment`, + `provider-bill:alipay:${providerBillDate}:payment`, + ], + ); } async function main() { const fakeWechat = await startFakeWechatPayServer(); + const fakeAlipay = await startFakeAlipayServer(); const pool = new pg.Pool({ connectionString: databaseUrl }); let seeded = false; try { await pool.query('begin'); await cleanup(pool); - await seed(pool, fakeWechat); + await seed(pool, fakeWechat, fakeAlipay); await pool.query('commit'); seeded = true; @@ -378,6 +632,70 @@ async function main() { assert.ok(!JSON.stringify(events.rows).includes(paymentFixture.wechatApiV3Key), 'payment event payload must not leak API v3 key'); assert.ok(!JSON.stringify(events.rows).includes('PRIVATE KEY'), 'payment event payload must not leak private key'); + const providerBillOutput = await runProviderBillsWorkerOnce(); + assert.match(providerBillOutput, /completed=2/, 'provider bill worker should complete two provider bill jobs'); + assert.ok( + fakeWechat.requests.some(item => item.method === 'GET' && item.pathname === '/v3/bill/tradebill'), + 'provider bill worker should request WeChat trade bill URL', + ); + assert.ok( + fakeWechat.requests.some(item => item.method === 'GET' && item.pathname === '/download/wechat-tradebill.json'), + 'provider bill worker should download WeChat trade bill', + ); + assert.ok( + fakeAlipay.requests.some(item => item.method === 'POST' && item.pathname === '/gateway.do' && item.body.includes('alipay.data.dataservice.bill.downloadurl.query')), + 'provider bill worker should request Alipay bill download URL', + ); + assert.ok( + fakeAlipay.requests.some(item => item.method === 'GET' && item.pathname === '/download/alipay-tradebill.json'), + 'provider bill worker should download Alipay trade bill', + ); + + const providerJobs = await pool.query( + ` + select provider, status, source_hash, row_count, download_url_host, + reconciliation_batch_id, error_code, error_message + from public.commerce_bill_download_jobs + where tenant_id = $1 and id in ($2::uuid, $3::uuid) + order by provider + `, + [tenantId, ids.providerBillJobWechat, ids.providerBillJobAlipay], + ); + assert.equal(providerJobs.rows.length, 2, 'provider bill jobs should be persisted'); + assert.ok(providerJobs.rows.every(row => row.status === 'completed'), 'provider bill jobs should be completed'); + assert.ok(providerJobs.rows.every(row => row.source_hash && !row.error_code && !row.error_message), 'completed provider bill jobs should keep source hash and no errors'); + assert.ok(providerJobs.rows.every(row => row.download_url_host === '127.0.0.1'), 'provider bill jobs should persist only download URL host'); + + const providerBatches = await pool.query( + ` + select b.id, b.provider, b.source, b.source_name, b.status, + b.matched_count, b.missing_local_count, b.total_count, + count(i.id)::int as item_count + from public.commerce_reconciliation_batches b + join public.commerce_reconciliation_items i + on i.tenant_id = b.tenant_id and i.batch_id = b.id + where b.tenant_id = $1 + and b.source = 'provider_download' + and b.source_name like 'provider-bill:%:2026-06-29:%' + group by b.id + order by b.provider + `, + [tenantId], + ); + assert.equal(providerBatches.rows.length, 2, 'provider bill downloads should create two reconciliation batches'); + assert.ok( + providerBatches.rows.some(row => row.provider === 'wechat_pay' && row.status === 'completed_with_issues' && row.matched_count >= 1 && row.missing_local_count >= 1), + 'WeChat bill batch should include matched and provider-only anomaly rows', + ); + assert.ok( + providerBatches.rows.some(row => row.provider === 'alipay' && ['completed', 'completed_with_issues'].includes(row.status) && row.matched_count >= 1), + 'Alipay bill batch should include matched rows', + ); + + const jobPayload = JSON.stringify(providerJobs.rows); + assert.ok(!jobPayload.includes(paymentFixture.wechatApiV3Key), 'provider bill job rows must not leak API v3 key'); + assert.ok(!jobPayload.includes('PRIVATE KEY'), 'provider bill job rows must not leak private keys'); + console.log('Commerce worker integration test complete.'); } catch (error) { await pool.query('rollback').catch(() => {}); @@ -388,6 +706,7 @@ async function main() { } await pool.end(); await fakeWechat.close(); + await fakeAlipay.close(); } } diff --git a/supabase/migrations/202606290029_commerce_bill_download_jobs.sql b/supabase/migrations/202606290029_commerce_bill_download_jobs.sql new file mode 100644 index 00000000..f0937645 --- /dev/null +++ b/supabase/migrations/202606290029_commerce_bill_download_jobs.sql @@ -0,0 +1,56 @@ +create table if not exists public.commerce_bill_download_jobs ( + id uuid primary key default gen_random_uuid(), + tenant_id uuid not null references public.tenants(id) on delete cascade, + provider text not null + check (provider in ('wechat_pay', 'alipay')), + bill_date date not null, + bill_type text not null default 'payment' + check (bill_type in ('payment', 'refund', 'combined')), + status text not null default 'queued' + check (status in ('queued', 'running', 'completed', 'failed', 'cancelled')), + source_name text not null, + source_hash text, + row_count integer not null default 0 check (row_count >= 0), + download_hash_type text, + download_hash_value text, + download_url_host text, + reconciliation_batch_id uuid references public.commerce_reconciliation_batches(id) on delete set null, + requested_by uuid references public.platform_users(id) on delete set null, + claimed_by text, + claimed_at timestamptz, + completed_at timestamptz, + failed_at timestamptz, + error_code text, + error_message text, + metadata jsonb not null default '{}'::jsonb, + created_at timestamptz not null default now(), + updated_at timestamptz not null default now() +); + +create unique index if not exists idx_commerce_bill_download_jobs_idempotency + on public.commerce_bill_download_jobs(tenant_id, provider, bill_date, bill_type) + where status in ('queued', 'running', 'completed'); + +create index if not exists idx_commerce_bill_download_jobs_claim + on public.commerce_bill_download_jobs(status, created_at) + where status = 'queued'; + +create index if not exists idx_commerce_bill_download_jobs_tenant + on public.commerce_bill_download_jobs(tenant_id, provider, bill_date desc, created_at desc); + +create index if not exists idx_commerce_bill_download_jobs_batch + on public.commerce_bill_download_jobs(tenant_id, reconciliation_batch_id) + where reconciliation_batch_id is not null; + +alter table public.commerce_bill_download_jobs enable row level security; + +drop policy if exists tenant_isolation on public.commerce_bill_download_jobs; +create policy tenant_isolation on public.commerce_bill_download_jobs + for all + using (tenant_id = app.current_tenant_id() or app.is_platform_admin()) + with check (tenant_id = app.current_tenant_id() or app.is_platform_admin()); + +drop trigger if exists set_updated_at on public.commerce_bill_download_jobs; +create trigger set_updated_at + before update on public.commerce_bill_download_jobs + for each row execute function app.touch_updated_at();