From 9891fe9ef3eef9f42eaa9d7f6561402f9a76560a Mon Sep 17 00:00:00 2001 From: Codex Date: Tue, 30 Jun 2026 23:07:43 +0800 Subject: [PATCH] feat: add crm dead-letter operations and benchmark summary --- README.md | 21 +- apps/api/src/features/referral/index.ts | 6 + apps/api/src/features/referral/routes.ts | 439 +++++++++++++++++- apps/taro/src/pages/tenant-admin/admin.css | 35 ++ .../pages/tenant-admin/marketing/index.tsx | 97 +++- apps/taro/src/services/tenantAdmin.ts | 69 +++ apps/worker/src/jobs/crm.ts | 7 +- docs/refactor/backend-capability-status.md | 4 +- docs/refactor/backend-handoff-roadmap.md | 8 +- docs/refactor/backend-progress.md | 9 +- docs/refactor/blueprint-coverage.md | 4 +- docs/refactor/crm-worker.md | 33 +- docs/refactor/frontend-handoff-index.md | 4 +- docs/refactor/implementation-status.md | 5 +- docs/refactor/legacy-feature-gap-matrix.md | 4 +- docs/refactor/next-development-todo.md | 5 +- .../refactor/performance-benchmark-runbook.md | 12 +- .../performance-benchmark-summary-20260630.md | 39 ++ docs/refactor/taro-frontend-integration.md | 12 +- scripts/api-integration-test.js | 264 +++++++++-- scripts/api-performance-benchmark.js | 275 +++++++++-- ...02606300018_crm_dead_letter_operations.sql | 42 ++ 22 files changed, 1284 insertions(+), 110 deletions(-) create mode 100644 docs/refactor/performance-benchmark-summary-20260630.md create mode 100644 supabase/migrations/202606300018_crm_dead_letter_operations.sql diff --git a/README.md b/README.md index 363f8c3e..91be0c77 100644 --- a/README.md +++ b/README.md @@ -34,8 +34,8 @@ - Excel/CSV 导入解析已完成并复用 `content_import_jobs/items/issues` 管线;大批量异步导入 worker 基础已接入,支持 queued job 消费、重试和审计;导入后复检、模板下载、字段映射 API 和 Taro 租户内容页第一版导入操作台已完成。 - 题库导出已完成服务端结构化 payload、PDF/Word 二进制 worker、每日一练基础导出和每日一练 ZIP 图片素材包;后续还要补更精细试卷模板、多模板排版和导出操作台体验。 - 优惠券复杂规则和核销报表已可联调,包含状态启停、活动分组、最低订单金额、优惠封顶、单用户限次、首单限制、适用套餐/地区、核销明细和活动报表;Taro 租户营销中心已接优惠券规则表单、筛选、核销明细和报表第一版。 -- 勋章管理、手动发放、签到连续天数、积分阈值、反馈解决、积分活动任务、练习次数、单词掌握和模考成绩系统触发勋章已可联调;积分活动任务、积分兑换商品、兑换订单、优惠券兑换履约、租户后台配置和用户站内通知第一版已完成,Taro 学生个人中心已接积分任务/兑换/积分明细和消息中心第一版,租户营销中心已接积分任务/兑换操作台和用户通知查看第一版。学生激励默认以后台配置勋章自动发放为主,排行榜默认不开启也不在学生端默认请求。后续还要补更细活动效果看板、外部微信订阅消息/短信推送、分佣真实打款 provider、发票、批量凭证上传、CRM 富卡片模板、失败告警、死信运营台、销售转化看板、公共题库版本通知和冲突处理操作台。 -- `apps/taro` 已建立 Taro 4 React 跨端前端地基,包含 H5 学生端、租户后台、平台后台三套构建入口、租户解析、统一 API client 和 Supabase Auth client 初始化;学生端第一批页面已接入登录、首页、题库、练习、背单词、知识手册、分数线、AI 择校推荐、资料、独立消息中心和个人中心,已新增 `RichContent` 安全渲染组件用于题干、选项、解析、知识手册和逐题复盘,H5 端已用 KaTeX 渲染 `$...$`、`$$...$$`、`\(...\)`、`\[...\]` 公式,私有题图可用 `asset:`/`content_asset:` 资源引用走短期预览签名,已升级背单词为今日计划/单元学习/收藏练习、学习概览、掌握率、收藏数、计划拆分、卡片翻转、发音、美/英音切换和本地位置恢复第一版,知识手册已接章节内搜索、安全文本摘要高亮和目录定位第一版,分数线已接目标地区默认筛选、院校/专业/年份 chip、租户动态字段筛选、结果字段 chip 和趋势摘要第一版,AI 择校已接报告生成、历史报告和 Markdown/HTML 导出第一版,资料页已补齐预览/下载的短签名、水印 traceId 和强制水印容器第一版,个人中心已接学习报告、14 天趋势、题型表现、最近练习、男女预设头像选择、积分任务/兑换/积分明细和消息中心摘要第一版,独立消息中心已接状态/类型筛选、批量已读、归档/忽略和站内安全跳转第一版;学生端不默认请求排行榜,仅在租户显式开启 `enableLeaderboard` 并完成压测后进入独立排行榜页或活动页;租户后台第一批页面已接入工作台、数据看板、学生/班级、题库内容、营销中心、财务运营和租户设置,学生运营页已接跟进看板、学习督导自动化、督导规则保存和批量 CRM 推送第一版,营销中心已接 CRM、分佣结算、优惠券规则/核销报表、积分任务/兑换操作台和用户通知查看第一版,财务运营已接退款状态机、官方账单任务、对账异常、差错工单和调整凭证第一版,设置页已接主题模板、草稿预览/发布、角色模板和成员绑定第一版;平台后台已接入工作台、租户管理、账务中心、公共题库授权、平台员工管理,以及创建租户、租户详情、状态变更、账务资料维护、平台员工创建/编辑/禁用恢复、权限点勾选、平台审计查询/CSV 导出、开放审计告警确认/解决、审计告警外部通知渠道/事件状态摘要、订阅、订阅账单候选/dry-run/批量生成、自动计费 worker 生成结果查看、收款、逾期预览/催缴记录、催缴外部通知渠道/事件摘要、用量和题库授权第一版写操作。 +- 勋章管理、手动发放、签到连续天数、积分阈值、反馈解决、积分活动任务、练习次数、单词掌握和模考成绩系统触发勋章已可联调;积分活动任务、积分兑换商品、兑换订单、优惠券兑换履约、租户后台配置和用户站内通知第一版已完成,Taro 学生个人中心已接积分任务/兑换/积分明细和消息中心第一版,租户营销中心已接积分任务/兑换操作台和用户通知查看第一版。学生激励默认以后台配置勋章自动发放为主,排行榜默认不开启也不在学生端默认请求。CRM 死信运营第一版已完成失败池、日志脱敏、重试/忽略和审计闭环。后续还要补更细活动效果看板、外部微信订阅消息/短信推送、分佣真实打款 provider、发票、批量凭证上传、CRM 富卡片模板、外部失败告警升级、销售转化看板、公共题库版本通知和冲突处理操作台。 +- `apps/taro` 已建立 Taro 4 React 跨端前端地基,包含 H5 学生端、租户后台、平台后台三套构建入口、租户解析、统一 API client 和 Supabase Auth client 初始化;学生端第一批页面已接入登录、首页、题库、练习、背单词、知识手册、分数线、AI 择校推荐、资料、独立消息中心和个人中心,已新增 `RichContent` 安全渲染组件用于题干、选项、解析、知识手册和逐题复盘,H5 端已用 KaTeX 渲染 `$...$`、`$$...$$`、`\(...\)`、`\[...\]` 公式,私有题图可用 `asset:`/`content_asset:` 资源引用走短期预览签名,已升级背单词为今日计划/单元学习/收藏练习、学习概览、掌握率、收藏数、计划拆分、卡片翻转、发音、美/英音切换和本地位置恢复第一版,知识手册已接章节内搜索、安全文本摘要高亮和目录定位第一版,分数线已接目标地区默认筛选、院校/专业/年份 chip、租户动态字段筛选、结果字段 chip 和趋势摘要第一版,AI 择校已接报告生成、历史报告和 Markdown/HTML 导出第一版,资料页已补齐预览/下载的短签名、水印 traceId 和强制水印容器第一版,个人中心已接学习报告、14 天趋势、题型表现、最近练习、男女预设头像选择、积分任务/兑换/积分明细和消息中心摘要第一版,独立消息中心已接状态/类型筛选、批量已读、归档/忽略和站内安全跳转第一版;学生端不默认请求排行榜,仅在租户显式开启 `enableLeaderboard` 并完成压测后进入独立排行榜页或活动页;租户后台第一批页面已接入工作台、数据看板、学生/班级、题库内容、营销中心、财务运营和租户设置,学生运营页已接跟进看板、学习督导自动化、督导规则保存和批量 CRM 推送第一版,营销中心已接 CRM 配置、队列筛选、死信失败池、日志查看、重试/忽略、分佣结算、优惠券规则/核销报表、积分任务/兑换操作台和用户通知查看第一版,财务运营已接退款状态机、官方账单任务、对账异常、差错工单和调整凭证第一版,设置页已接主题模板、草稿预览/发布、角色模板和成员绑定第一版;平台后台已接入工作台、租户管理、账务中心、公共题库授权、平台员工管理,以及创建租户、租户详情、状态变更、账务资料维护、平台员工创建/编辑/禁用恢复、权限点勾选、平台审计查询/CSV 导出、开放审计告警确认/解决、审计告警外部通知渠道/事件状态摘要、订阅、订阅账单候选/dry-run/批量生成、自动计费 worker 生成结果查看、收款、逾期预览/催缴记录、催缴外部通知渠道/事件摘要、用量和题库授权第一版写操作。 - 根目录已清理为新 Supabase SaaS monorepo 编排层;旧 PocketBase/React 项目和旧构建产物仅保留在 `参考/` 目录作为迁移参考,不进入 Git 提交。 ## 商用功能完成度总览 @@ -53,7 +53,7 @@ | 对象存储/资料安全 | √ 可联调,待生产 AV/CDN | OSS/COS/Supabase Storage 签名、上传确认、短签名预览下载、水印 traceId、复检和安全扫描地基已完成 | | PocketBase 真实数据迁移 | √ 本地跑通,待人工复核 blocker | SQLite 导出、标准化导入、校验和抽样脚本已跑通;正式切换前处理缺用户订单和缺归属手册章节 | | Taro H5 三端前端 | √ 第一版可构建 | 学生端、租户后台、平台后台均有真实 API 页面;后续继续补小程序兼容、视觉精修、状态管理和端到端测试 | -| 生产安全/压测交付 | △ 待专项阶段 | 需执行 `@codex-security`、生产 readiness、远程 Auth/RLS、4c16g 压测、PostgreSQL 调优和上线证据门禁 | +| 生产安全/压测交付 | △ 本地真实数据压测已跑,云端待复测 | 本地 Docker/Supabase 已完成真实迁移数据 30/50/100/150 并发混合读写压测;上云后仍需执行 `@codex-security`、生产 readiness、远程 Auth/RLS、4c16g 压测、PostgreSQL 调优和上线证据门禁 | 更完整的进度看这些文档: @@ -538,6 +538,21 @@ docs/refactor/postgresql-4c16g-tuning.md docs/refactor/performance-benchmark-runbook.md ``` +最近一次本地真实迁移库已开启刷题写入闭环压测,数据规模约为 74,102 题、1,597 个题目合集、3,102 个练习蓝图、3,500 个单词、2,676 条知识手册和 3,670 个用户。压测 worker 是无停顿请求流,不能直接等同于真实在线学生数;前端完成后需要用真实页面埋点估算单个学生平均 RPS,再折算在线容量。 + +| 并发 worker | 时长 | 刷题写入比例 | 请求数 | 错误率 | 吞吐 | P95 | P99 | +| ---: | ---: | ---: | ---: | ---: | ---: | ---: | ---: | +| 30 | 120s | 10% | 112,896 | 0.00% | 934.74 req/s | 64.26 ms | 81.14 ms | +| 50 | 120s | 10% | 92,737 | 0.00% | 767.24 req/s | 121.51 ms | 159.12 ms | +| 100 | 120s | 8% | 84,608 | 0.00% | 699.27 req/s | 254.69 ms | 331.84 ms | +| 150 | 120s | 6% | 82,693 | 0.00% | 682.40 req/s | 369.90 ms | 493.69 ms | + +本地结论:100 个无停顿 worker 内 P95 仍低于 300ms;150 worker 零错误但 P95 已明显上升,可作为本机 Docker 环境的压力拐点参考。正式对外容量承诺必须在目标 4 核 16G 云服务器、生产对象存储/CDN 和真实前端请求节奏下复跑。脱敏摘要见: + +```text +docs/refactor/performance-benchmark-summary-20260630.md +``` + 正式切换前建议使用 production 严格模式: ```bash diff --git a/apps/api/src/features/referral/index.ts b/apps/api/src/features/referral/index.ts index 6060bfca..07594c56 100644 --- a/apps/api/src/features/referral/index.ts +++ b/apps/api/src/features/referral/index.ts @@ -15,6 +15,9 @@ import { } from './commission.js'; import { crmConfigRoute, + crmDeadLettersRoute, + crmQueueActionRoute, + crmQueueLogsRoute, crmQueueRoute, referralBindRoute, referralInviteCodeRoute, @@ -45,6 +48,9 @@ export const referralRoutes: RouteDefinition[] = [ ['GET', '/api/crm/config', crmConfigRoute], ['PUT', '/api/crm/config', upsertCrmConfigRoute], ['GET', '/api/crm/queue', crmQueueRoute], + ['GET', '/api/crm/dead-letters', crmDeadLettersRoute], + ['GET', '/api/crm/queue/logs', crmQueueLogsRoute], + ['POST', '/api/crm/queue/action', crmQueueActionRoute], ['GET', '/api/commission/settings', commissionSettingsRoute], ['PUT', '/api/commission/settings', updateCommissionSettingsRoute], ['PUT', '/api/commission/member-rate', updateMemberCommissionRateRoute], diff --git a/apps/api/src/features/referral/routes.ts b/apps/api/src/features/referral/routes.ts index 93728047..12c2c1bb 100644 --- a/apps/api/src/features/referral/routes.ts +++ b/apps/api/src/features/referral/routes.ts @@ -13,6 +13,10 @@ const EVENT_TYPES = ['enter', 'register', 'purchase', 'share', 'scan', 'manual_b const TRACK_SOURCES = ['share', 'qrcode', 'timeline', 'miniapp', 'h5', 'manual', 'unknown']; const CRM_ASSIGNMENT_MODES = ['none', 'direct', 'round_robin', 'referrer']; const CRM_ASSIGNABLE_ROLES = ['tenant_owner', 'tenant_admin', 'tenant_operator', 'teacher', 'sales', 'agent']; +const CRM_QUEUE_STATUSES = ['pending', 'processing', 'retrying', 'sent', 'failed', 'discarded']; +const CRM_DEAD_LETTER_STATUSES = ['failed', 'discarded']; +const CRM_QUEUE_ACTIONS = ['retry', 'ignore']; +const UUID_RE = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i; function nullableString(value: unknown) { return typeof value === 'string' && value.trim() ? value.trim() : null; @@ -30,6 +34,97 @@ function optionalChoice(value: unknown, allowed: string[], fallback: string) { return candidate; } +function optionalUuidString(value: unknown, key: string) { + const text = nullableString(value); + if (!text) return ''; + if (!UUID_RE.test(text)) { + throw new HttpError(400, `${key} must be a UUID`, 'INVALID_UUID'); + } + return text; +} + +function requiredUuidString(body: JsonBody, key: string) { + const text = optionalUuidString(body[key], key); + if (!text) throw new HttpError(400, `${key} is required`, 'REQUIRED_FIELD'); + return text; +} + +function limitedOptionalString(value: unknown, key: string, maxLength: number) { + const text = nullableString(value); + if (!text) return null; + if (text.length > maxLength) { + throw new HttpError(400, `${key} is too long`, 'FIELD_TOO_LONG'); + } + return text; +} + +function redactSensitiveString(value: string) { + return value + .replace(/("(?:access[_-]?token|token|secret|sign|signature|key|password)"\s*:\s*")[^"]+(")/gi, '$1[redacted]$2') + .replace(/(access[_-]?token|token|secret|sign|signature|key|password)=([^&\s]+)/gi, '$1=[redacted]') + .replace(/(Bearer\s+)[A-Za-z0-9._~+/=-]+/gi, '$1[redacted]') + .replace(/(https?:\/\/[^/\s]+\/(?:robot\/send|bot|webhook|open-apis\/bot)\/)[^?\s/]+/gi, '$1[redacted]'); +} + +function isSensitiveFieldKey(key: string) { + const normalized = key.replace(/[^a-z0-9]/gi, '').toLowerCase(); + return new Set([ + 'accesstoken', + 'refreshtoken', + 'token', + 'secret', + 'clientsecret', + 'appsecret', + 'password', + 'credential', + 'credentials', + 'authorization', + 'signature', + 'sign', + 'key', + 'apikey', + 'secretkey', + 'accesskey', + 'privatekey', + ]).has(normalized); +} + +function redactSensitiveValue(value: unknown): unknown { + if (typeof value === 'string') return redactSensitiveString(value); + if (Array.isArray(value)) return value.map(item => redactSensitiveValue(item)); + if (!value || typeof value !== 'object') return value; + const output: Record = {}; + for (const [key, raw] of Object.entries(value as Record)) { + if (isSensitiveFieldKey(key)) { + output[key] = '[redacted]'; + } else { + output[key] = redactSensitiveValue(raw); + } + } + return output; +} + +function safeUrlForResponse(value: unknown) { + const text = nullableString(value); + if (!text) return null; + try { + const url = new URL(text); + url.username = ''; + url.password = ''; + for (const key of Array.from(url.searchParams.keys())) { + if (/token|secret|password|credential|authorization|signature|sign|key/i.test(key)) { + url.searchParams.set(key, '[redacted]'); + } + } + if (/\/(robot\/send|bot|webhook|open-apis\/bot)\//i.test(url.pathname)) { + url.pathname = url.pathname.replace(/(\/(?:robot\/send|bot|webhook|open-apis\/bot)\/)[^/]+/i, '$1[redacted]'); + } + return url.toString(); + } catch { + return redactSensitiveString(text); + } +} + function stringArray(value: unknown) { if (!Array.isArray(value)) return []; const seen = new Set(); @@ -543,6 +638,43 @@ function crmPermission(auth: TenantAdminAuth) { requireTenantPermission(auth, 'crm:read'); } +async function recordCrmAudit( + client: pg.PoolClient, + auth: TenantAdminAuth, + action: string, + targetId: string | null, + details: Record = {}, +) { + await client.query( + ` + insert into public.audit_logs (tenant_id, actor_user_id, action, target_type, target_id, details) + values ($1, $2, $3, 'crm_webhook_queue', $4, $5::jsonb) + `, + [auth.tenantId, auth.userId, action, targetId, JSON.stringify(redactSensitiveValue(details))], + ); +} + +function crmQueueResponseItem>(item: T) { + return { + ...item, + targetUrl: safeUrlForResponse(item.targetUrl), + payload: redactSensitiveValue(item.payload || {}) as Record, + operatorMetadata: redactSensitiveValue(item.operatorMetadata || {}) as Record, + lastError: item.lastError ? redactSensitiveString(String(item.lastError)) : item.lastError, + lastResponseSummary: item.lastResponseSummary ? redactSensitiveString(String(item.lastResponseSummary)) : item.lastResponseSummary, + }; +} + +function crmLogResponseItem>(item: T) { + return { + ...item, + requestBody: item.requestBody ? redactSensitiveString(String(item.requestBody)) : item.requestBody, + requestPayload: redactSensitiveValue(item.requestPayload || {}) as Record, + responseSummary: item.responseSummary ? redactSensitiveString(String(item.responseSummary)) : item.responseSummary, + errorMessage: item.errorMessage ? redactSensitiveString(String(item.errorMessage)) : item.errorMessage, + }; +} + export async function referralInviteCodeRoute(ctx: RequestContext) { const tenantId = await tenantIdFrom(ctx); const userId = await userIdFrom(ctx); @@ -1146,12 +1278,20 @@ export async function crmQueueRoute(ctx: RequestContext) { crmPermission(auth); const limit = intParam(ctx, 'limit', 100, 500); const status = stringParam(ctx, 'status'); + const queueId = optionalUuidString(stringParam(ctx, 'queueId') || stringParam(ctx, 'id'), 'queueId'); const params: unknown[] = [auth.tenantId]; const filters = ['tenant_id = $1']; if (status) { + if (!CRM_QUEUE_STATUSES.includes(status)) { + throw new HttpError(400, `Invalid CRM queue status: ${status}`, 'INVALID_CRM_QUEUE_STATUS'); + } params.push(status); filters.push(`status = $${params.length}`); } + if (queueId) { + params.push(queueId); + filters.push(`id = $${params.length}::uuid`); + } params.push(limit); const items = await query( @@ -1160,6 +1300,11 @@ export async function crmQueueRoute(ctx: RequestContext) { attempts, next_attempt_at as "nextAttemptAt", last_error as "lastError", last_http_code as "lastHttpCode", lead_id as "leadId", sent_at as "sentAt", source, idempotency_key as "idempotencyKey", target_url as "targetUrl", + provider, last_attempt_at as "lastAttemptAt", last_response_summary as "lastResponseSummary", + dead_lettered_at as "deadLetteredAt", ignored_at as "ignoredAt", + ignored_by as "ignoredBy", last_operator_user_id as "lastOperatorUserId", + last_operator_action as "lastOperatorAction", last_operator_note as "lastOperatorNote", + last_operator_at as "lastOperatorAt", operator_metadata as "operatorMetadata", payload, created_at as "createdAt", updated_at as "updatedAt" from public.crm_webhook_queue where ${filters.join(' and ')} @@ -1168,5 +1313,297 @@ export async function crmQueueRoute(ctx: RequestContext) { `, params, ); - return { items }; + return { items: items.map(item => crmQueueResponseItem(item as Record)) }; +} + +export async function crmDeadLettersRoute(ctx: RequestContext) { + const auth = await requireTenantAdmin(ctx); + crmPermission(auth); + const limit = intParam(ctx, 'limit', 100, 500); + const status = stringParam(ctx, 'status'); + const source = stringParam(ctx, 'source'); + const params: unknown[] = [auth.tenantId]; + const filters = [`status = any($2::text[])`]; + params.push(CRM_DEAD_LETTER_STATUSES); + if (status) { + if (!CRM_DEAD_LETTER_STATUSES.includes(status)) { + throw new HttpError(400, `Invalid CRM dead-letter status: ${status}`, 'INVALID_CRM_DEAD_LETTER_STATUS'); + } + params.push(status); + filters.push(`status = $${params.length}`); + } + if (source) { + params.push(source); + filters.push(`source = $${params.length}`); + } + params.push(limit); + + const items = await query( + ` + select id, record_id as "recordId", status, scheduled_at as "scheduledAt", + attempts, next_attempt_at as "nextAttemptAt", last_error as "lastError", + last_http_code as "lastHttpCode", lead_id as "leadId", sent_at as "sentAt", + source, idempotency_key as "idempotencyKey", target_url as "targetUrl", + provider, last_attempt_at as "lastAttemptAt", last_response_summary as "lastResponseSummary", + dead_lettered_at as "deadLetteredAt", ignored_at as "ignoredAt", + ignored_by as "ignoredBy", last_operator_user_id as "lastOperatorUserId", + last_operator_action as "lastOperatorAction", last_operator_note as "lastOperatorNote", + last_operator_at as "lastOperatorAt", operator_metadata as "operatorMetadata", + payload, created_at as "createdAt", updated_at as "updatedAt" + from public.crm_webhook_queue + where tenant_id = $1 and ${filters.join(' and ')} + order by coalesce(dead_lettered_at, updated_at, created_at) desc + limit $${params.length} + `, + params, + ); + const summary = await queryOne<{ + failed: string; + discarded: string; + ignored: string; + total: string; + oldestOpenAt: string | null; + }>( + ` + select count(*) filter (where status = 'failed')::text as failed, + count(*) filter (where status = 'discarded')::text as discarded, + count(*) filter (where ignored_at is not null)::text as ignored, + count(*)::text as total, + min(coalesce(dead_lettered_at, updated_at, created_at)) filter (where ignored_at is null) as "oldestOpenAt" + from public.crm_webhook_queue + where tenant_id = $1 and status in ('failed', 'discarded') + `, + [auth.tenantId], + ); + return { + summary: { + failed: Number(summary?.failed || 0), + discarded: Number(summary?.discarded || 0), + ignored: Number(summary?.ignored || 0), + total: Number(summary?.total || 0), + oldestOpenAt: summary?.oldestOpenAt || null, + }, + items: items.map(item => crmQueueResponseItem(item as Record)), + }; +} + +export async function crmQueueLogsRoute(ctx: RequestContext) { + const auth = await requireTenantAdmin(ctx); + crmPermission(auth); + const queueId = optionalUuidString(stringParam(ctx, 'queueId') || stringParam(ctx, 'id'), 'queueId'); + if (!queueId) throw new HttpError(400, 'queueId is required', 'REQUIRED_FIELD'); + const limit = intParam(ctx, 'limit', 50, 200); + + const task = await queryOne<{ + id: string; + recordId: string | null; + leadId: string | null; + idempotencyKey: string | null; + }>( + ` + select id, record_id as "recordId", lead_id as "leadId", idempotency_key as "idempotencyKey" + from public.crm_webhook_queue + where tenant_id = $1 and id = $2::uuid + limit 1 + `, + [auth.tenantId, queueId], + ); + if (!task) throw new HttpError(404, 'CRM queue task not found', 'CRM_QUEUE_TASK_NOT_FOUND'); + + const items = await query( + ` + select id, queue_id as "queueId", record_id as "recordId", http_code as "httpCode", + outcome, error_message as "errorMessage", lead_id as "leadId", + request_body as "requestBody", request_payload as "requestPayload", + response_summary as "responseSummary", signed_at as "signedAt", + attempt, operator_user_id as "operatorUserId", operation, + created_at as "createdAt", updated_at as "updatedAt" + from public.crm_webhook_log + where tenant_id = $1 + and ( + queue_id = $2::uuid + or ( + queue_id is null + and ( + ($3::text is not null and record_id = $3::text) + or ($4::text is not null and lead_id = $4::text) + ) + ) + ) + order by created_at desc + limit $5 + `, + [auth.tenantId, queueId, task.recordId, task.leadId, limit], + ); + return { item: task, items: items.map(item => crmLogResponseItem(item as Record)) }; +} + +export async function crmQueueActionRoute(ctx: RequestContext) { + const auth = await requireTenantAdmin(ctx); + requireTenantPermission(auth, 'crm:write'); + const body = await readJsonBody(ctx); + const queueId = requiredUuidString(body, 'queueId'); + const action = optionalChoice(body.action, CRM_QUEUE_ACTIONS, 'retry'); + const note = limitedOptionalString(body.note ?? body.reason, 'note', 500); + const metadata = objectValue(body.metadata); + const safeMetadata = redactSensitiveValue(metadata) as Record; + + const item = await transaction(async client => { + const current = await client.query<{ + id: string; + status: string; + attempts: number; + recordId: string | null; + leadId: string | null; + source: string | null; + payload: Record; + }>( + ` + select id, status, attempts, record_id as "recordId", lead_id as "leadId", source, payload + from public.crm_webhook_queue + where tenant_id = $1 and id = $2::uuid + limit 1 + for update + `, + [auth.tenantId, queueId], + ); + const task = current.rows[0]; + if (!task) throw new HttpError(404, 'CRM queue task not found', 'CRM_QUEUE_TASK_NOT_FOUND'); + + if (action === 'retry') { + if (!['failed', 'discarded', 'retrying', 'pending'].includes(task.status)) { + throw new HttpError(409, `CRM task cannot be retried from status ${task.status}`, 'CRM_QUEUE_STATUS_NOT_RETRYABLE'); + } + const result = await client.query( + ` + update public.crm_webhook_queue + set status = 'pending', + next_attempt_at = now(), + scheduled_at = coalesce(scheduled_at, now()), + dead_lettered_at = null, + ignored_at = null, + ignored_by = null, + last_error = null, + last_http_code = null, + last_response_summary = null, + last_operator_user_id = $3::uuid, + last_operator_action = 'retry', + last_operator_note = $4, + last_operator_at = now(), + operator_metadata = coalesce(operator_metadata, '{}'::jsonb) || $5::jsonb, + updated_at = now() + where tenant_id = $1 and id = $2::uuid + returning id, record_id as "recordId", status, scheduled_at as "scheduledAt", + attempts, next_attempt_at as "nextAttemptAt", last_error as "lastError", + last_http_code as "lastHttpCode", lead_id as "leadId", sent_at as "sentAt", + source, idempotency_key as "idempotencyKey", target_url as "targetUrl", + provider, last_attempt_at as "lastAttemptAt", last_response_summary as "lastResponseSummary", + dead_lettered_at as "deadLetteredAt", ignored_at as "ignoredAt", + ignored_by as "ignoredBy", last_operator_user_id as "lastOperatorUserId", + last_operator_action as "lastOperatorAction", last_operator_note as "lastOperatorNote", + last_operator_at as "lastOperatorAt", operator_metadata as "operatorMetadata", + payload, created_at as "createdAt", updated_at as "updatedAt" + `, + [ + auth.tenantId, + queueId, + auth.userId, + note, + JSON.stringify({ lastRetry: { by: auth.userId, at: new Date().toISOString(), note, ...safeMetadata } }), + ], + ); + await client.query( + ` + insert into public.crm_webhook_log ( + tenant_id, queue_id, record_id, http_code, outcome, error_message, lead_id, + request_payload, response_summary, signed_at, attempt, operator_user_id, operation + ) + values ($1, $2::uuid, $3, null, 'operator_retry', $4, $5, $6::jsonb, null, now(), $7, $8::uuid, 'retry') + `, + [ + auth.tenantId, + queueId, + task.recordId, + note || 'Queued for manual retry', + task.leadId, + JSON.stringify({ source: 'tenant_operator', metadata: safeMetadata }), + task.attempts, + auth.userId, + ], + ); + await recordCrmAudit(client, auth, 'crm.queue.retried', queueId, { + previousStatus: task.status, + attempts: task.attempts, + source: task.source, + note, + }); + return result.rows[0]; + } + + if (!CRM_DEAD_LETTER_STATUSES.includes(task.status)) { + throw new HttpError(409, `CRM task cannot be ignored from status ${task.status}`, 'CRM_QUEUE_STATUS_NOT_IGNORABLE'); + } + const result = await client.query( + ` + update public.crm_webhook_queue + set status = 'discarded', + next_attempt_at = null, + dead_lettered_at = coalesce(dead_lettered_at, now()), + ignored_at = now(), + ignored_by = $3::uuid, + last_operator_user_id = $3::uuid, + last_operator_action = 'ignore', + last_operator_note = $4, + last_operator_at = now(), + operator_metadata = coalesce(operator_metadata, '{}'::jsonb) || $5::jsonb, + updated_at = now() + where tenant_id = $1 and id = $2::uuid + returning id, record_id as "recordId", status, scheduled_at as "scheduledAt", + attempts, next_attempt_at as "nextAttemptAt", last_error as "lastError", + last_http_code as "lastHttpCode", lead_id as "leadId", sent_at as "sentAt", + source, idempotency_key as "idempotencyKey", target_url as "targetUrl", + provider, last_attempt_at as "lastAttemptAt", last_response_summary as "lastResponseSummary", + dead_lettered_at as "deadLetteredAt", ignored_at as "ignoredAt", + ignored_by as "ignoredBy", last_operator_user_id as "lastOperatorUserId", + last_operator_action as "lastOperatorAction", last_operator_note as "lastOperatorNote", + last_operator_at as "lastOperatorAt", operator_metadata as "operatorMetadata", + payload, created_at as "createdAt", updated_at as "updatedAt" + `, + [ + auth.tenantId, + queueId, + auth.userId, + note, + JSON.stringify({ lastIgnore: { by: auth.userId, at: new Date().toISOString(), note, ...safeMetadata } }), + ], + ); + await client.query( + ` + insert into public.crm_webhook_log ( + tenant_id, queue_id, record_id, http_code, outcome, error_message, lead_id, + request_payload, response_summary, signed_at, attempt, operator_user_id, operation + ) + values ($1, $2::uuid, $3, null, 'operator_ignore', $4, $5, $6::jsonb, null, now(), $7, $8::uuid, 'ignore') + `, + [ + auth.tenantId, + queueId, + task.recordId, + note || 'Ignored by tenant operator', + task.leadId, + JSON.stringify({ source: 'tenant_operator', metadata: safeMetadata }), + task.attempts, + auth.userId, + ], + ); + await recordCrmAudit(client, auth, 'crm.queue.ignored', queueId, { + previousStatus: task.status, + attempts: task.attempts, + source: task.source, + note, + }); + return result.rows[0]; + }); + + return { item: crmQueueResponseItem(item as Record) }; } diff --git a/apps/taro/src/pages/tenant-admin/admin.css b/apps/taro/src/pages/tenant-admin/admin.css index cf2d53c3..750256f7 100644 --- a/apps/taro/src/pages/tenant-admin/admin.css +++ b/apps/taro/src/pages/tenant-admin/admin.css @@ -99,6 +99,18 @@ grid-template-columns: repeat(3, minmax(0, 1fr)); } +.admin-metric-grid { + display: grid; + grid-template-columns: repeat(3, minmax(0, 1fr)); + gap: 12px; + margin-bottom: 14px; +} + +.admin-metric-grid.compact .admin-metric { + min-height: 88px; + padding: 14px; +} + .admin-metric { min-height: 112px; padding: 18px; @@ -237,6 +249,29 @@ color: #fff; } +.admin-mini-button.danger { + border-color: #dc2626; + background: #fef2f2; + color: #b91c1c; +} + +.admin-sub-list { + display: flex; + flex-direction: column; + gap: 8px; + margin-top: 12px; + padding: 12px; + border-radius: 8px; + background: #f8fafc; +} + +.admin-sub-row { + padding: 10px 12px; + border: 1px solid #e2e8f0; + border-radius: 8px; + background: #fff; +} + .admin-error { display: block; margin-top: 12px; diff --git a/apps/taro/src/pages/tenant-admin/marketing/index.tsx b/apps/taro/src/pages/tenant-admin/marketing/index.tsx index 834efbd4..f2897b2e 100644 --- a/apps/taro/src/pages/tenant-admin/marketing/index.tsx +++ b/apps/taro/src/pages/tenant-admin/marketing/index.tsx @@ -16,7 +16,9 @@ import { loadCouponReport, loadCoupons, loadCrmConfig, + loadCrmDeadLetters, loadCrmQueue, + loadCrmQueueLogs, loadFeedbackReport, loadPointActivityClaims, loadPointActivityTasks, @@ -28,6 +30,7 @@ import { updateCommissionSettings, updateCommissionSettlementProofStatus, updateCommissionSettlementStatus, + updateCrmQueueAction, updateMemberCommissionRate, upsertCoupon, upsertCrmConfig, @@ -43,6 +46,8 @@ import { type CouponRedemptionItem, type CouponReport, type CrmConfigItem, + type CrmDeadLettersPayload, + type CrmQueueLogItem, type CrmQueueItem, type FeedbackReport, type PointActivityTaskItem, @@ -296,6 +301,8 @@ export default function TenantMarketingPage() { const [codes, setCodes] = useState[]>([]); const [crmConfig, setCrmConfig] = useState(null); const [crmQueue, setCrmQueue] = useState([]); + const [crmDeadLetters, setCrmDeadLetters] = useState({}); + const [crmLogs, setCrmLogs] = useState>({}); const [commissionSettings, setCommissionSettings] = useState(null); const [commission, setCommission] = useState(null); const [commissionOrders, setCommissionOrders] = useState([]); @@ -370,6 +377,7 @@ export default function TenantMarketingPage() { codePayload, crmConfigPayload, crmPayload, + crmDeadLettersPayload, commissionSettingsPayload, commissionPayload, orderPayload, @@ -399,6 +407,7 @@ export default function TenantMarketingPage() { loadActivationCodes(20).catch(() => ({ items: [] })), loadCrmConfig().catch(() => ({ item: null })), loadCrmQueue(20, crmStatus || undefined).catch(() => ({ items: [] })), + loadCrmDeadLetters({ limit: 20 }).catch(() => ({ items: [], summary: {} })), loadCommissionSettings().catch(() => ({ item: null })), loadCommissionSummary(period).catch(() => ({ item: null })), loadCommissionOrders({ ...period, limit: 20 }).catch(() => ({ items: [] })), @@ -426,6 +435,7 @@ export default function TenantMarketingPage() { setCodes(codePayload.items || []); setCrmConfig(nextCrm); setCrmQueue(crmPayload.items || []); + setCrmDeadLetters(crmDeadLettersPayload || {}); setCommissionSettings(nextSettings); setCommission(commissionPayload.item || null); setCommissionOrders(orderPayload.items || []); @@ -753,8 +763,12 @@ export default function TenantMarketingPage() { setBusy('crm-queue'); setError(''); try { - const payload = await loadCrmQueue(30, status || undefined); + const [payload, deadLettersPayload] = await Promise.all([ + loadCrmQueue(30, status || undefined), + loadCrmDeadLetters({ limit: 20 }), + ]); setCrmQueue(payload.items || []); + setCrmDeadLetters(deadLettersPayload || {}); } catch (nextError) { setError(nextError instanceof Error ? nextError.message : 'CRM 队列加载失败'); } finally { @@ -762,6 +776,45 @@ export default function TenantMarketingPage() { } } + async function showCrmQueueLogs(item: CrmQueueItem) { + setBusy(`crm-log:${item.id}`); + setError(''); + try { + const payload = await loadCrmQueueLogs(item.id, 20); + setCrmLogs(prev => ({ ...prev, [item.id]: payload.items || [] })); + } catch (nextError) { + setError(nextError instanceof Error ? nextError.message : 'CRM 日志加载失败'); + } finally { + setBusy(''); + } + } + + async function handleCrmQueueAction(item: CrmQueueItem, action: 'retry' | 'ignore') { + const confirmed = await Taro.showModal({ + title: action === 'retry' ? '重试 CRM 任务' : '忽略 CRM 任务', + content: `${item.source || 'CRM'} · ${item.status || 'pending'} · 尝试 ${item.attempts || 0} 次`, + confirmText: action === 'retry' ? '重试' : '忽略', + cancelText: '取消', + }); + if (!confirmed.confirm) return; + setBusy(`crm-action:${item.id}:${action}`); + setError(''); + try { + await updateCrmQueueAction({ + queueId: item.id, + action, + note: action === 'retry' ? '租户后台手动重试' : '租户后台确认忽略', + metadata: { source: 'taro-tenant-admin' }, + }); + Taro.showToast({ title: action === 'retry' ? '已重新入队' : '已忽略', icon: 'success' }); + await refreshCrmQueue(crmStatus); + } catch (nextError) { + setError(nextError instanceof Error ? nextError.message : 'CRM 队列操作失败'); + } finally { + setBusy(''); + } + } + async function refreshUserNotifications(override: Partial = {}) { const nextFilter = { ...userNotificationFilter, ...override }; setUserNotificationFilter(nextFilter); @@ -1366,6 +1419,11 @@ export default function TenantMarketingPage() { CRM 队列 + + 失败池{String(crmDeadLetters.summary?.total || 0)} + 待重试{String(crmDeadLetters.summary?.failed || 0)} + 已忽略{String(crmDeadLetters.summary?.ignored || 0)} + {['', 'pending', 'sent', 'failed', 'discarded'].map(status => ( + {item.status === 'failed' || item.status === 'discarded' ? : null} + {item.status === 'failed' ? : null} + + {crmLogs[item.id]?.length ? ( + + {crmLogs[item.id].slice(0, 5).map(log => ( + + {log.outcome || log.operation || 'log'} · 第 {log.attempt || 0} 次 · HTTP {log.httpCode || '-'} · {shortDate(log.createdAt || log.signedAt)} + {log.errorMessage ? {log.errorMessage} : null} + {log.responseSummary ? {log.responseSummary} : null} + + ))} + + ) : null} ))} {!crmQueue.length ? 暂无 CRM 队列。 : null} + {crmDeadLetters.items?.length ? ( + + 近期失败任务 + {crmDeadLetters.items.slice(0, 5).map(item => ( + + {item.status || 'failed'} · {item.source || 'CRM'} + 进入失败池 {shortDate(item.deadLetteredAt || item.updatedAt || item.createdAt)} · 尝试 {item.attempts || 0} 次 + {item.lastError ? {item.lastError} : null} + + + + {item.status === 'failed' ? : null} + + + ))} + + ) : null} diff --git a/apps/taro/src/services/tenantAdmin.ts b/apps/taro/src/services/tenantAdmin.ts index d8b85c4d..25556c44 100644 --- a/apps/taro/src/services/tenantAdmin.ts +++ b/apps/taro/src/services/tenantAdmin.ts @@ -795,8 +795,49 @@ export interface CrmQueueItem { sentAt?: string | null; source?: string | null; targetUrl?: string | null; + provider?: string | null; + lastAttemptAt?: string | null; + lastResponseSummary?: string | null; + deadLetteredAt?: string | null; + ignoredAt?: string | null; + ignoredBy?: string | null; + lastOperatorUserId?: string | null; + lastOperatorAction?: string | null; + lastOperatorNote?: string | null; + lastOperatorAt?: string | null; + operatorMetadata?: Record; payload?: Record; createdAt?: string; + updatedAt?: string; +} + +export interface CrmDeadLettersPayload { + summary?: { + failed?: number; + discarded?: number; + ignored?: number; + total?: number; + oldestOpenAt?: string | null; + }; + items?: CrmQueueItem[]; +} + +export interface CrmQueueLogItem { + id: string; + queueId?: string | null; + recordId?: string | null; + httpCode?: number | null; + outcome?: string | null; + errorMessage?: string | null; + leadId?: string | null; + requestBody?: string | null; + requestPayload?: Record; + responseSummary?: string | null; + signedAt?: string | null; + attempt?: number | null; + operatorUserId?: string | null; + operation?: string | null; + createdAt?: string; } export interface CommissionSettingsItem { @@ -1599,6 +1640,34 @@ export async function loadCrmQueue(limit = 20, status?: string) { return apiRequest<{ items?: CrmQueueItem[] }>('/api/crm/queue', { query: { limit, status } }); } +export async function loadCrmDeadLetters(input: { + limit?: number; + status?: 'failed' | 'discarded'; + source?: string; +} = {}) { + return apiRequest('/api/crm/dead-letters', { + query: { ...input, limit: input.limit || 30 }, + }); +} + +export async function loadCrmQueueLogs(queueId: string, limit = 30) { + return apiRequest<{ item?: Record; items?: CrmQueueLogItem[] }>('/api/crm/queue/logs', { + query: { queueId, limit }, + }); +} + +export async function updateCrmQueueAction(input: { + queueId: string; + action: 'retry' | 'ignore'; + note?: string | null; + metadata?: Record; +}) { + return apiRequest<{ item?: CrmQueueItem }>('/api/crm/queue/action', { + method: 'POST', + body: input, + }); +} + export async function loadCommissionSettings() { return apiRequest<{ item?: CommissionSettingsItem }>('/api/commission/settings'); } diff --git a/apps/worker/src/jobs/crm.ts b/apps/worker/src/jobs/crm.ts index 9ff77303..906f88b7 100644 --- a/apps/worker/src/jobs/crm.ts +++ b/apps/worker/src/jobs/crm.ts @@ -464,13 +464,14 @@ async function appendLog( await client.query( ` insert into public.crm_webhook_log ( - tenant_id, record_id, http_code, outcome, error_message, lead_id, + tenant_id, queue_id, record_id, http_code, outcome, error_message, lead_id, request_body, request_payload, response_summary, signed_at, attempt ) - values ($1, $2, $3, $4, $5, $6, $7, $8::jsonb, $9, now(), $10) + values ($1, $2::uuid, $3, $4, $5, $6, $7, $8, $9::jsonb, $10, now(), $11) `, [ task.tenantId, + task.id, task.recordId, result.httpCode, result.ok ? 'sent' : 'failed', @@ -524,6 +525,7 @@ async function markTaskResult( attempts = $3, provider = $4, next_attempt_at = case when $2 = 'retrying' then now() + ($5::int * interval '1 second') else null end, + dead_lettered_at = case when $2 = 'failed' then coalesce(dead_lettered_at, now()) else dead_lettered_at end, last_error = $6, last_http_code = $7, last_response_summary = $8, @@ -551,6 +553,7 @@ async function discardTask(client: pg.PoolClient, task: CrmQueueRow, message: st set status = 'discarded', attempts = attempts + 1, last_attempt_at = now(), + dead_lettered_at = coalesce(dead_lettered_at, now()), last_error = $2, updated_at = now() where id = $1 diff --git a/docs/refactor/backend-capability-status.md b/docs/refactor/backend-capability-status.md index 193166d1..392b41f8 100644 --- a/docs/refactor/backend-capability-status.md +++ b/docs/refactor/backend-capability-status.md @@ -162,9 +162,9 @@ | 首绑客资保护 | 可联调 | `/api/referral/bind` | | 手工补绑 | 可联调 | 需要 `referral:write` | | 销售统计/客户列表/团队 | 可联调 | `/api/referral/sales-*`、`team` | -| CRM 配置/队列 | 可联调 | `/api/crm/config`、`/api/crm/queue`;CRM 配置支持 `none/direct/round_robin/referrer` 客资跟进分配策略、租户内候选成员校验、轮询游标、分配审计,客资首绑成功后会写入 `assignedToUserId` 并进入队列 payload | +| CRM 配置/队列 | 可联调 | `/api/crm/config`、`/api/crm/queue`、`/api/crm/dead-letters`、`/api/crm/queue/logs`、`/api/crm/queue/action`;CRM 配置支持 `none/direct/round_robin/referrer` 客资跟进分配策略、租户内候选成员校验、轮询游标、分配审计,客资首绑成功后会写入 `assignedToUserId` 并进入队列 payload;失败/丢弃任务可进入死信运营池,租户管理员可查看脱敏日志、手动重试或忽略,动作写入日志和审计 | | CRM webhook worker | 可联调 | `apps/worker` 已支持 generic webhook、钉钉、飞书、企微群机器人消息体/签名、`lead.created` 客资首绑事件和 `student.crm_push` 学生运营跟进事件、到期任务消费、失败退避重试、最终失败、discarded 和 `crm_webhook_log` | -| CRM 增强 | 待补齐 | 富卡片模板、失败告警、死信运营后台和更细销售转化看板 | +| CRM 增强 | 待补齐 | 富卡片模板、外部失败告警升级和更细销售转化看板 | | 分佣结算基础闭环 | 可联调 | `/api/commission/settings`、`member-rate`、`summary`、`orders`、`settlements`、`settlements/generate`、`settlements/status`、`settlements/export`、`settlements/proofs`;支持订单/激活码归因、批次/成员/默认比例优先级、北京时间账期、结算单生成、审核、打款状态、已打款锁定、CSV/JSON 导出、打款凭证登记/复核、销售/代理本人范围和租户隔离 | | 分佣打款增强 | 待补齐 | 银行/微信/支付宝真实打款 provider、发票管理、批量凭证上传、异常调整单和分佣看板 | diff --git a/docs/refactor/backend-handoff-roadmap.md b/docs/refactor/backend-handoff-roadmap.md index 467a17ca..7ffade8e 100644 --- a/docs/refactor/backend-handoff-roadmap.md +++ b/docs/refactor/backend-handoff-roadmap.md @@ -11,7 +11,7 @@ - 平台侧可以管理租户、SaaS 套餐、订阅、订阅账单候选预览/dry-run/批量生成、自动计费 worker、服务费、逾期催缴、用量台账、月度用量自动采集和用量超额账单。 - 租户侧可以管理品牌、域名、支付账户、登录配置、私密密钥、活动、兑换码、优惠券、成员权限、审计日志、内容入口、分类树、题目集合、练习蓝图、题目、视频、分数线、单词、知识手册、资料资源和题库导出任务。 - 学生侧已经有题库入口、分类树、题目集合、顺序/随机/全真模拟组卷、答题、错题、收藏、背单词进度、个人中心、男女预设头像、站内通知、勋章、排行榜接口(租户默认关闭)、分数线、视频、订单详情/状态轮询、优惠券领取/抵扣、权益、激活码预检查/兑换、资料下载、AI 择校推荐、积分活动任务和积分兑换的基础 API;学生头像不支持上传或第三方头像落库,学生写入口会拒绝头像 URL;签到、积分阈值、反馈解决和积分活动可触发自动勋章发放,反馈处理/奖励、勋章发放和积分兑换会写入用户站内通知;租户后台已具备反馈运营聚合报表和积分风控只读报表。 -- 销售/代理/CRM 已经有邀请码、扫码/分享事件、首绑客资保护、团队关系、统计、CRM 配置、入队、worker 推送和分佣结算基础闭环。 +- 销售/代理/CRM 已经有邀请码、扫码/分享事件、首绑客资保护、团队关系、统计、CRM 配置、入队、worker 推送、失败死信运营、手动重试/忽略和分佣结算基础闭环。 - 旧题库 JSON、单词模板、知识手册嵌套模板、分数线 JSON 和视频绑定 JSON 已经进入后端 preview/import 管线,由后端负责规范化、校验、幂等、审计和租户隔离。 因此,后端现在已经具备进入 Taro 前端第一阶段联调的基础。需要注意的是,它还不是完整生产交付状态,真实云端鉴权、对象存储生产安全、支付/短信/OAuth 生产账号、真实数据 dry-run 迁移仍需要继续补齐或联调;导入后复检、模板下载、字段映射 API 和导入任务详情已可联调,Taro 租户内容页已接入上传/粘贴预览、字段别名覆盖、同步/异步执行、异步轮询和复检详情第一版,租户营销中心已接入 CRM 配置/队列、分佣结算、积分任务/兑换和积分风控只读摘要第一版。 @@ -31,7 +31,7 @@ | 资料下载 | 部分完成 | 资源台账、SVIP 权限校验、`local_dev`/阿里云 OSS/腾讯 COS/Supabase Storage 上传下载签名、上传确认、PDF/图片预览签名、访问审计、动态水印上下文、assets worker 复检、内置安全扫描、外部 HTTP scanner 接入层、题库导出 PDF/Word/每日一练 ZIP 可生成可信 `content_assets` 并走签名下载/预览 | CDN 防盗链、真实 AV/内容安全服务联调、转码级水印、资料前端操作体验 | | 会员与订单 | 可联调 | 下单、订单详情/状态轮询、优惠券领取/抵扣、零元订单自动开通、手工确认权限保护、激活码预检查/兑换、微信支付、支付宝、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、支付/退款补偿 worker、权益发放、资金对账批次/明细/异常查询 API、微信/支付宝官方账单下载 worker、对账差错工单和事件轨迹 | 财务复核报表、异常订单运营台和真实生产账单格式验收 | | 登录认证 | 可联调 | 短信 mock、阿里云/腾讯云短信 adapter、迁移期 session、Supabase Auth JWT、微信小程序登录、微信网页登录、QQ 登录、手机号绑定/换绑、OAuth 配置表 | 真实生产账号和回调域名联调 | -| 销售/代理/CRM | 可联调 | 邀请码、首绑保护、团队关系、销售统计、CRM 配置/队列、钉钉/飞书/企微 worker、分佣规则、成员分佣比例、订单/激活码归因、结算生成、审核和打款状态;Taro 营销中心已接第一版操作台 | 小程序码真实生成、CRM 分配策略、结算导出、真实打款、凭证、财务复核和销售转化看板 | +| 销售/代理/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 执行验收、抽样校验和导入性能压测 | | 数据看板 | 可联调 | 租户 dashboard 聚合接口,收益、注册、学习、内容、激活码、反馈、趋势、24h 活跃、套餐销量和运营动态 | 预聚合 worker、缓存、慢 SQL 监控和销售转化看板 | | AI 择校推荐 | 可联调 | `ai_recommendation_reports`、SVIP 门禁、学生输入 schema、地区/分数线上下文、`local_rules` 稳定 JSON、报告列表/详情、Markdown/HTML 导出和 Taro 学生端基础页 | 真实 AI provider、prompt 版本管理、租户后台配置、报告 PDF worker 渲染和人工复核流程 | @@ -79,7 +79,7 @@ - 对象存储:上传/下载签名已接入阿里云 OSS、腾讯云 COS、Supabase Storage;上传确认、PDF/图片预览签名、动态水印上下文、assets worker 复检、内置安全扫描、外部 HTTP scanner 接入层和题库导出 PDF/Word/每日一练 ZIP worker 已完成,继续补视频播放防盗链、真实 AV/内容安全服务联调和转码/CDN 级水印。 - 真实数据 dry-run:导出 PocketBase 用户、题库、单词、知识手册、分数线、订单、权益,先跑 `npm run pb:import:dry-run -- --profile=production --json --fail-on-warnings`,确认 `migrationReadiness` 的必需集合和关键字段覆盖率通过,再跑迁移和校验报告。 - 生产环境配置:`.env.example` 和 `npm run readiness:production` / `npm run readiness:production:db` 已补;继续补数据库迁移流程、备份恢复、日志、告警和 API 容器部署说明。 -- Taro scaffold:`apps/taro` 地基已建立;学生端、租户后台、平台后台第一批 H5 页面已接真实 API,学生端已接地区选择、错题/收藏复习、阅读理解/案例分析多小题作答、题目反馈、视频解析、练习/模考报告、收银台、订单详情、售后入口、积分任务/兑换/积分明细、个人中心消息摘要和独立消息中心第一版;平台后台关键写操作、租户详情、账务资料编辑、平台员工列表/创建/编辑/禁用恢复、最近平台审计查询/CSV 导出、开放审计告警展示/确认/解决、审计告警外部通知渠道/事件状态摘要、订阅账单候选/dry-run/批量生成、自动计费生成结果查看、逾期预览、催缴记录和催缴外部通知摘要第一版已接入,租户工作台已接权限驱动模块入口,租户学生运营页已接学生创建/更新、禁用/恢复、批量导入、批量分班、备注和跟进任务第一版,租户内容页已接公共题库采纳/同步、冲突查看、单条/批量采纳平台或保留本地、导入问题、字段模板预览/下载、上传/粘贴预览、字段别名覆盖、同步/异步导入、异步轮询和复检详情第一版,租户设置页已接角色模板创建/编辑/停用、成员绑定模板和权限可见性配置第一版,租户营销中心已接 CRM 配置/队列、分佣结算、优惠券规则/核销报表、积分任务/兑换操作台、积分风控摘要和用户通知查看第一版;下一步补公式图片混排、更细数据范围 UI、平台审计告警升级策略、平台催缴通知配置操作台细节和小程序兼容验证。 +- Taro scaffold:`apps/taro` 地基已建立;学生端、租户后台、平台后台第一批 H5 页面已接真实 API,学生端已接地区选择、错题/收藏复习、阅读理解/案例分析多小题作答、题目反馈、视频解析、练习/模考报告、收银台、订单详情、售后入口、积分任务/兑换/积分明细、个人中心消息摘要和独立消息中心第一版;平台后台关键写操作、租户详情、账务资料编辑、平台员工列表/创建/编辑/禁用恢复、最近平台审计查询/CSV 导出、开放审计告警展示/确认/解决、审计告警外部通知渠道/事件状态摘要、订阅账单候选/dry-run/批量生成、自动计费生成结果查看、逾期预览、催缴记录和催缴外部通知摘要第一版已接入,租户工作台已接权限驱动模块入口,租户学生运营页已接学生创建/更新、禁用/恢复、批量导入、批量分班、备注和跟进任务第一版,租户内容页已接公共题库采纳/同步、冲突查看、单条/批量采纳平台或保留本地、导入问题、字段模板预览/下载、上传/粘贴预览、字段别名覆盖、同步/异步导入、异步轮询和复检详情第一版,租户设置页已接角色模板创建/编辑/停用、成员绑定模板和权限可见性配置第一版,租户营销中心已接 CRM 配置/队列/死信运营、分佣结算、优惠券规则/核销报表、积分任务/兑换操作台、积分风控摘要和用户通知查看第一版;下一步补公式图片混排、更细数据范围 UI、平台审计告警升级策略、平台催缴通知配置操作台细节和小程序兼容验证。 ### P1:商用收费和运营能力 @@ -95,7 +95,7 @@ - 租户自定义角色模板基础 API、Taro 权限配置 UI、成员绑定模板和工作台权限驱动入口第一版已完成;继续补班级/教师/学生组合范围 UI、成员批量运营和更完整平台级审计报表。 - 三套默认主题、租户主题预览、Logo/图标/分享图配置。 -- CRM worker:钉钉、飞书、企微机器人发送、签名、失败重试已落地;Taro 营销中心已能配置 CRM 和查看队列;继续补轮询/定向分配、富卡片、失败告警和死信运营台。 +- CRM worker:钉钉、飞书、企微机器人发送、签名、失败重试、死信失败池、脱敏日志、手动重试/忽略和审计已落地;Taro 营销中心已能配置 CRM、查看队列并处理失败任务;继续补富卡片、外部失败告警升级和销售转化看板。 - 销售/代理分佣基础闭环已接 Taro 第一版;继续补销售团队看板、客资跟进效果、结算导出、真实打款、凭证和财务复核。 - AI 择校推荐:`local_rules` JSON 报告、本人报告 Markdown/HTML 导出地基已完成;继续补真实 AI provider、prompt 版本、租户后台配置、PDF worker 报告生成和人工复核。 - 性能压测、慢 SQL 审查、备份恢复演练、灰度发布和回滚预案。 diff --git a/docs/refactor/backend-progress.md b/docs/refactor/backend-progress.md index 8e6ce461..4d6b8850 100644 --- a/docs/refactor/backend-progress.md +++ b/docs/refactor/backend-progress.md @@ -30,7 +30,7 @@ - 已新增 `npm run db:smoke-seed`,用于 `supabase:reset` 后恢复最小烟测数据。 - 已新增 `npm run smoke:core-api`,用于验证个人中心、分数线、题目视频、背单词进度/收藏等学生端核心 API。 - 已新增 `npm run test:api`,自动 seed、构建、启动临时 API,并断言核心学生端接口、内容导航/组卷、租户隔离、资源权限和题目导入。 -- 已新增 `apps/worker` 和 `npm run test:worker:crm`,用于消费 CRM webhook 队列,验证本地 fake webhook、队列状态、日志和密钥不泄露;worker 已支持 `lead.created` 客资首绑事件和 `student.crm_push` 学生运营跟进事件。 +- 已新增 `apps/worker` 和 `npm run test:worker:crm`,用于消费 CRM webhook 队列,验证本地 fake webhook、队列状态、日志和密钥不泄露;worker 已支持 `lead.created` 客资首绑事件和 `student.crm_push` 学生运营跟进事件,失败/丢弃任务会进入死信运营池。 - 已新增 commerce worker 和 `npm run test:worker:commerce`,用于补偿查询微信/支付宝支付、处理中退款和漏通知场景;支付成功会幂等更新订单/支付并开通权益,退款成功会幂等更新退款/订单/支付并在全额退款时撤销订单权益,测试覆盖密钥不泄露和重复执行不重复开通。 - 已新增 platform-billing worker 和 `npm run test:worker:platform-billing`,用于自动处理即将到期且未开票的 SaaS 订阅;worker 使用订阅行锁和账单查重保证幂等,自动生成 `tenant_invoices/tenant_invoice_items` 并写 `platform.invoice.subscription_auto_created` 审计。 - 已新增 platform-dunning worker 和 `npm run test:worker:platform-dunning`,用于扫描已过 `due_date` 且未结清的 SaaS 服务费账单;worker 使用账单行锁和每日唯一催缴约束保证幂等,自动标记 `overdue`、推送租户 `billing_status=past_due`、生成 `tenant_invoice_reminders` 内部催缴记录并写 `platform.invoice.overdue_processed` 审计。 @@ -203,6 +203,9 @@ POST /api/commission/settlements/status GET /api/crm/config PUT /api/crm/config GET /api/crm/queue +GET /api/crm/dead-letters +GET /api/crm/queue/logs +POST /api/crm/queue/action GET /api/tenant-admin/permissions GET /api/tenant-admin/role-templates PUT /api/tenant-admin/role-templates @@ -288,7 +291,7 @@ GET /api/tenant-admin/audit-logs - 班级学生 API 会按 `tenant_memberships.role_template_id -> tenant_role_templates.data_scope`、成员显式权限和 `tenant_class_members` 共同确定可见范围;非全局权限教师只能查看自己负责班级的学生。 - 学生批量导入、批量分班、学生状态、备注和跟进任务都使用独立权限点;教师默认可为范围内学生写备注和跟进任务,但不能批量导入、禁用学生或放大可见班级。 - 销售/代理客资采用首绑保护:普通扫码/分享事件不会覆盖已有归属,只有具备 `referral:write` 的租户成员可手动强制补绑。 -- CRM 当前完成配置、密钥入私密表、客资入队、队列查询和 `apps/worker` 消费;worker 支持 generic webhook、钉钉、飞书、企业微信机器人消息体/签名、失败重试和日志。 +- CRM 当前完成配置、密钥入私密表、客资入队、队列查询、死信失败池、脱敏日志查看、手动重试/忽略和 `apps/worker` 消费;worker 支持 generic webhook、钉钉、飞书、企业微信机器人消息体/签名、失败重试、最终失败/丢弃标记和日志。 - 内容资源当前完成台账、租户后台维护、学生端 SVIP 下载权限,以及 `local_dev`、阿里云 OSS、腾讯 COS、Supabase Storage 的上传/下载签名 provider;上传确认、访问审计、动态水印上下文、租户后台媒体运营报表、assets worker 对象元数据校验/复检、内置 `metadata_rules` 安全扫描和外部 HTTP scanner 接入层已完成。PDF 深度预览体验、防盗链、转码/CDN 级水印和真实 AV/内容安全服务联调仍需继续补。 - 题库内容导航当前以 `content_entries/content_nodes` 为主模型,可表达“入口 -> 多级分类 -> 院校/专业/学科/销售意向标记”;题目集合和练习方式由 `question_collections/practice_blueprints` 管理,练习 session 会保存当次题目 ID 快照。 - 练习访问控制由 `content_entries/content_nodes/question_collections/practice_blueprints` 的 `accessRules` 合并决定;普通用户消耗 `practice_daily_usage`,事件写入 `practice_access_events`,SVIP/staff 不消耗免费额度。 @@ -300,7 +303,7 @@ GET /api/tenant-admin/audit-logs 2. 完成真实短信 provider 联调:阿里云/腾讯云,密钥放 `app_private.tenant_secrets` 或生产 Vault。 3. 完成真实 OAuth provider 联调:微信网页、微信小程序、QQ,确认回调域名、开放平台账号和旧 PocketBase 身份映射策略。 4. 补真实生产账单格式验收、前端财务操作台和营销活动 UI;支付/退款补偿、退款查询确认、退款通知、微信/支付宝官方账单下载、资金对账导入比对、差错工单、异常订单运营台和优惠券核销报表主链路已完成。 -5. 扩展 `apps/worker`:日报统计、CRM 死信告警、学习督导触达联动、公共题库同步失败告警和更完整冲突处理运营台;公共题库同步 worker、学习督导规则 worker 和同步通知已具备基础闭环。 +5. 扩展 `apps/worker`:日报统计、CRM 外部失败告警升级、学习督导触达联动、公共题库同步失败告警和更完整冲突处理运营台;CRM 死信运营、公共题库同步 worker、学习督导规则 worker 和同步通知已具备基础闭环。 6. 继续 Taro H5/小程序兼容验证,重点看公式、题图资源字段化、支付容器和分享链路。 ## 测试命令 diff --git a/docs/refactor/blueprint-coverage.md b/docs/refactor/blueprint-coverage.md index bc01a20b..d3848272 100644 --- a/docs/refactor/blueprint-coverage.md +++ b/docs/refactor/blueprint-coverage.md @@ -27,7 +27,7 @@ | 资料下载/PDF | 可联调 | `content_assets` 资源台账、后台资源管理、OSS/COS/Supabase Storage 上传下载签名、上传确认、PDF/图片预览签名、学生端列表、SVIP 下载权限、动态水印上下文、assets worker 复检、内置安全扫描和外部 HTTP scanner 接入层 | CDN 防盗链、真实 AV/内容安全服务联调、资料前端管理页 | | 营销中心 | 可联调 | SVIP 套餐、激活码批次、激活码生成、优惠券启停/归档、活动分组、最低金额、优惠封顶、单用户限次、首单限制、适用套餐/地区、核销明细、核销报表、Banner/FAQ/公告、勋章管理、手动发放、签到/积分/反馈/活动任务/练习/单词/模考自动发放、积分活动任务、积分兑换商品、兑换订单、优惠券兑换履约和积分风控只读报表 | 营销自动化、活动效果看板和前端活动操作台 | | 销售/代理客资 | 可联调 | 邀请码、扫码/分享事件、首绑保护、销售统计、客资明细、团队关系、手动补绑、分佣比例、归因、结算单、审核和打款状态 | 真实微信小程序码、真实打款、结算导出、销售团队看板 | -| CRM 系统 | 可联调 | CRM 配置、密钥私密存储、客资入队、队列查询、generic/钉钉/飞书/企微 worker、签名、重试和日志 | 定向/轮询分配、富卡片模板、失败告警、死信运营台 | +| CRM 系统 | 可联调 | CRM 配置、密钥私密存储、客资入队、`none/direct/round_robin/referrer` 跟进分配策略、队列查询、死信失败池、脱敏日志、手动重试/忽略、generic/钉钉/飞书/企微 worker、签名、重试和日志 | 富卡片模板、外部失败告警升级、更细销售转化看板 | | 数据看板 | 可联调 | 租户 dashboard 聚合接口,收益、注册、学习、内容、激活码、反馈、趋势、24h 活跃、套餐销量和运营动态 | 预聚合 worker、缓存、慢 SQL 监控和销售转化看板 | | 登录认证 | 可联调 | 短信 mock、阿里云/腾讯云短信 adapter、迁移期 session、Supabase Auth JWT、微信小程序登录、微信网页登录、QQ 登录、手机号绑定/换绑、OAuth 配置表 | 真实生产账号和回调域名联调 | | 支付 | 可联调 | 订单、支付记录、手动确认权限保护、权益发放、租户商户配置、微信支付 JSAPI、支付宝 WAP/H5、webhook 幂等、退款状态机、退款通知、补偿 worker、微信/支付宝官方账单下载 worker、资金对账导入比对、异常查询、差错工单和事件轨迹 | 异常订单运营台、真实生产账单格式验收、服务商/平台代收模式 | @@ -41,4 +41,4 @@ 3. 学习统计增强:专项练习策略和更细题型分析;排行榜默认不进入学生端主线,仅在租户显式购买/开启活动并完成压测后再补防刷和日/周榜预聚合。 4. 视频会员控制:深度防盗链、转码级水印和播放统计。 5. 数据看板预聚合:把实时聚合升级为大租户可承载的日/周/月预聚合。 -6. 真实 provider:短信、微信/QQ 登录、微信支付/支付宝生产账号联调;CRM worker 基础已落地,继续补分配策略和告警。 +6. 真实 provider:短信、微信/QQ 登录、微信支付/支付宝生产账号联调;CRM worker 和死信运营已落地,继续补富卡片模板、外部失败告警升级和销售转化看板。 diff --git a/docs/refactor/crm-worker.md b/docs/refactor/crm-worker.md index bbb31897..5cc40168 100644 --- a/docs/refactor/crm-worker.md +++ b/docs/refactor/crm-worker.md @@ -45,6 +45,9 @@ WORKER_CRM_ALLOW_INSECURE_LOCALHOST=false PUT /api/crm/config GET /api/crm/config GET /api/crm/queue +GET /api/crm/dead-letters +GET /api/crm/queue/logs +POST /api/crm/queue/action POST /api/tenant-admin/students/crm-push ``` @@ -92,6 +95,34 @@ pending -> discarded 每次尝试都会写入 `crm_webhook_log`,日志只记录 provider、目标 host、请求体和响应摘要,不写入 webhook secret。 +## 死信运营 + +CRM 失败运营入口已经落地,租户后台只通过后端 API 处理失败任务: + +```text +GET /api/crm/dead-letters?status=failed&limit=20 +GET /api/crm/queue/logs?queueId= +POST /api/crm/queue/action +``` + +`POST /api/crm/queue/action` 当前支持: + +```json +{ "queueId": "", "action": "retry", "reason": "确认 webhook 已恢复" } +``` + +```json +{ "queueId": "", "action": "ignore", "reason": "租户确认不再推送" } +``` + +安全边界: + +- 读取失败池、查看日志和重试/忽略都要求租户内 `crm:write` 权限。 +- `retry` 会把 `failed/discarded/retrying/pending` 任务重置为 `pending`,清空错误并记录操作者。 +- `ignore` 只允许处理 `failed/discarded` 任务,会标记为 `discarded` 并记录原因。 +- API 会递归脱敏 payload、日志和审计响应中的 secret、token、authorization、password 等敏感字段。 +- 所有动作都会写入 `crm_webhook_log` 和租户审计日志,前端不能直接改 `crm_webhook_queue`。 + ## 安全边界 - 前端不能直接写 `crm_webhook_queue`。 @@ -105,5 +136,5 @@ pending -> discarded ## 后续增强 - 钉钉/飞书/企微富卡片模板。 -- 失败告警和死信运营后台。 +- 外部失败告警升级。 - 更细销售转化看板。 diff --git a/docs/refactor/frontend-handoff-index.md b/docs/refactor/frontend-handoff-index.md index ac09b935..f247a5e8 100644 --- a/docs/refactor/frontend-handoff-index.md +++ b/docs/refactor/frontend-handoff-index.md @@ -99,11 +99,11 @@ | 数据看板 | `apps/taro/src/pages/tenant-admin/dashboard/index.tsx` | `tenant-admin/dashboard` | | 学生运营 | `apps/taro/src/pages/tenant-admin/students/index.tsx` | `tenant-admin/classes`、`tenant-admin/teachers`、`tenant-admin/students`、`students/bulk-upsert`、`students/status`、`classes/members/bulk-assign`、`students/notes`、`students/followups`、`students/followups/report`、`students/supervision/preview`、`students/supervision/generate`、`students/supervision/rules`、`students/crm-push` | | 题库内容 | `apps/taro/src/pages/tenant-admin/content/index.tsx` | `tenant-content/content-entries`、`tenant-content/imports`、`imports/detail`、`imports/issues`、`imports/field-mapping`、`imports/templates`、`imports/post-check`、`tenant-content/exports/questions`、`tenant-content/exports/jobs`、`tenant-content/assets/sign-download`、`tenant-content/assets/sign-preview`、`tenant-content/assets/security-scan-events`、`tenant-content/public-question-banks`、`public-question-banks/adopt`、`public-question-banks/sync`、`public-question-banks/conflicts`、`public-question-banks/conflicts/resolve`、`public-question-banks/conflicts/resolve-batch`、`tenant-content/notifications`、`tenant-content/notifications/status` | -| 营销中心 | `apps/taro/src/pages/tenant-admin/marketing/index.tsx` | `tenant-admin/coupons`、`code-batches`、`activation-codes`、`crm/config`、`crm/queue`、`commission/settings`、`member-rate`、`summary`、`orders`、`settlements`、`settlements/generate`、`settlements/status`、`tenant-admin/feedbacks/report`、`tenant-admin/point-activity-tasks`、`tenant-admin/point-activity-claims`、`tenant-admin/point-exchange-items`、`tenant-admin/point-exchange-orders`、`tenant-admin/points-risk-report`、`tenant-admin/user-notifications`;已接 CRM、分佣、优惠券规则/核销报表、反馈运营摘要、积分任务/兑换配置与记录查看、积分风控摘要、用户通知查看第一版 | +| 营销中心 | `apps/taro/src/pages/tenant-admin/marketing/index.tsx` | `tenant-admin/coupons`、`code-batches`、`activation-codes`、`crm/config`、`crm/queue`、`crm/dead-letters`、`crm/queue/logs`、`crm/queue/action`、`commission/settings`、`member-rate`、`summary`、`orders`、`settlements`、`settlements/generate`、`settlements/status`、`tenant-admin/feedbacks/report`、`tenant-admin/point-activity-tasks`、`tenant-admin/point-activity-claims`、`tenant-admin/point-exchange-items`、`tenant-admin/point-exchange-orders`、`tenant-admin/points-risk-report`、`tenant-admin/user-notifications`;已接 CRM 配置、队列、死信失败池、脱敏日志、重试/忽略、分佣、优惠券规则/核销报表、反馈运营摘要、积分任务/兑换配置与记录查看、积分风控摘要、用户通知查看第一版 | | 财务运营 | `apps/taro/src/pages/tenant-admin/finance/index.tsx` | `commerce/refunds`、`commerce/refunds/status`、`commerce/operations/anomalies`、`commerce/reconciliation/batches`、`commerce/reconciliation/items`、`commerce/reconciliation/issues/create`、`commerce/reconciliation/issues`、`commerce/reconciliation/issues/status`、`commerce/reconciliation/provider-bills/request`、`commerce/reconciliation/provider-bills/jobs`、`commerce/adjustment-vouchers`、`commerce/adjustment-vouchers/status`、`commerce/adjustment-vouchers/report` | | 租户设置 | `apps/taro/src/pages/tenant-admin/settings/index.tsx` | `tenant-admin/overview`、`domains`、`payment-accounts`、`auth-providers`、`theme-templates`、`theme`、`theme/preview`、`theme/publish`、`permissions`、`GET/PUT role-templates`、`POST role-templates/disable`、`GET/PUT members`、`POST members/disable` | -当前租户后台已有第一批运营操作:工作台按权限矩阵隐藏不可见模块;学生运营页支持学生创建/更新、状态禁用/恢复、批量导入、批量分班、学生备注、跟进任务、完成跟进、跟进看板、学习督导候选预览/一键生成跟进任务、保存每日督导规则和 CRM 入队推送;题库内容页支持公共题库采纳/同步、同步通知查看与已读/忽略、同步冲突查看、单条/批量采纳平台版本或保留本地版本、导入任务详情、异步 job 轮询、导入问题查看、字段映射/模板预览/下载、JSON/CSV/Excel 导入预览和执行、字段别名覆盖和导入后复检详情,也可接 JSON/试卷 payload 同步导出与 PDF/Word 异步导出 job 轮询,完成后用 `assetId` 走后台资源签名下载/预览;营销中心支持 CRM 配置保存、队列按状态查看、分佣规则、成员分佣比例、分佣订单明细、结算单生成、审核通过/驳回、标记线下打款、优惠券规则/核销报表、反馈运营摘要、积分任务/兑换配置、领取/兑换记录查看、积分风控只读摘要和用户通知查看第一版;财务运营页支持退款状态流、官方账单任务、对账异常、差错工单和调整凭证,且只通过后端命令层写审计与状态;租户设置页支持主题模板选择、草稿预览、发布、角色模板创建、编辑、停用、权限点、菜单、模块、字段、基础数据范围配置、成员搜索/新建、成员绑定模板、成员状态和额外权限覆盖。下一批需要继续补更精细的学生导入模板体验、更细数据范围 UI、学习督导触达联动/效果归因、主题素材库、真实打款 provider、发票、真实生产账单抽样验收和小程序端兼容。 +当前租户后台已有第一批运营操作:工作台按权限矩阵隐藏不可见模块;学生运营页支持学生创建/更新、状态禁用/恢复、批量导入、批量分班、学生备注、跟进任务、完成跟进、跟进看板、学习督导候选预览/一键生成跟进任务、保存每日督导规则和 CRM 入队推送;题库内容页支持公共题库采纳/同步、同步通知查看与已读/忽略、同步冲突查看、单条/批量采纳平台版本或保留本地版本、导入任务详情、异步 job 轮询、导入问题查看、字段映射/模板预览/下载、JSON/CSV/Excel 导入预览和执行、字段别名覆盖和导入后复检详情,也可接 JSON/试卷 payload 同步导出与 PDF/Word 异步导出 job 轮询,完成后用 `assetId` 走后台资源签名下载/预览;营销中心支持 CRM 配置保存、队列按状态查看、死信失败池、脱敏日志、手动重试/忽略、分佣规则、成员分佣比例、分佣订单明细、结算单生成、审核通过/驳回、标记线下打款、优惠券规则/核销报表、反馈运营摘要、积分任务/兑换配置、领取/兑换记录查看、积分风控只读摘要和用户通知查看第一版;财务运营页支持退款状态流、官方账单任务、对账异常、差错工单和调整凭证,且只通过后端命令层写审计与状态;租户设置页支持主题模板选择、草稿预览、发布、角色模板创建、编辑、停用、权限点、菜单、模块、字段、基础数据范围配置、成员搜索/新建、成员绑定模板、成员状态和额外权限覆盖。下一批需要继续补更精细的学生导入模板体验、更细数据范围 UI、学习督导触达联动/效果归因、主题素材库、真实打款 provider、发票、真实生产账单抽样验收和小程序端兼容。 ## 已落地的 Taro 平台后台页面 diff --git a/docs/refactor/implementation-status.md b/docs/refactor/implementation-status.md index 5e2131d3..51b58ceb 100644 --- a/docs/refactor/implementation-status.md +++ b/docs/refactor/implementation-status.md @@ -35,7 +35,7 @@ | 资料下载/PDF | 已扩展 `content_assets`,新增资源台账和导入任务表 | 旧 `app_assets/images` 兼容导入 | 租户后台资源管理、OSS/COS/Supabase Storage 上传/下载签名、上传确认、PDF/图片预览签名、资源访问审计、动态水印上下文、内置安全扫描、外部 HTTP scanner 接入层、题库导出 PDF/Word/每日一练 ZIP 自动发布可信资源、学生端资料列表/下载权限已实现 | 核心 API 集成测试含 SVIP 资料下载、安全扫描门禁、水印 traceId,assets worker 覆盖内置规则与外部 scanner 通过/失败/不可用 fail-closed,exports worker 测试;Taro 类型检查覆盖学生资料水印预览/下载确认 | 资料资源基础闭环可跑,Taro 学生资料页已接短期签名、可见水印覆盖、追踪码展示和强制水印资源外部打开限制第一版;真实 AV/内容安全服务联调、CDN 防盗链、转码/CDN 级水印待补 | | 个人中心 | 已建 `student_profiles`、会员权益、订单、练习记录、`badges/user_badges`、积分任务/兑换表、`user_notifications` | 已支持部分用户资料和勋章导入 | 个人资料、目标院校/专业、手机号绑定/换绑、会员状态、最近练习、统计聚合、签到积分、积分任务、积分兑换、题目反馈、考试倒计时、勋章和站内通知 API 已实现 | API 集成测试、Taro 类型检查 | 学生端个人中心已接学习报告、14 天趋势、题型表现、最近练习、男女预设头像选择、积分任务/兑换/积分明细、消息中心摘要、激活码和订单入口第一版;独立消息中心已接筛选、批量已读、归档/忽略和站内安全跳转第一版;账号合并和更细学习建议待补;不做头像上传或第三方头像同步 | | 活动/优惠 | 已建优惠券、激活码、激活码批次、banner、FAQ、公告、勋章、积分任务、积分兑换商品、兑换订单表和用户站内通知表 | 部分支持 | banner/FAQ/公告只读与租户后台维护、激活码预检查/兑换、激活码批次、批量生成激活码、优惠券维护、前台领取/下单抵扣、最低金额、优惠封顶、单用户限次、首单限制、适用套餐/地区、活动分组、核销明细、核销报表、勋章维护、手动发放、签到/积分/反馈/活动任务/练习/单词/模考自动发放、积分任务领取、积分兑换、优惠券兑换履约、积分风控只读报表和站内通知已实现 | 核心 API 集成测试 | Taro 租户营销中心已接优惠券、积分任务/兑换、积分风控摘要和用户通知查看第一版;营销自动化、外部订阅消息/短信和更完整活动效果看板继续补 | -| 销售/代理客资追踪 | 已建推荐码、首绑客资、团队关系、小程序码缓存、CRM 队列 | 旧 `referral_tracks` 已有映射基础 | 邀请码、扫码/分享事件、首绑保护、销售统计、客资明细、手动补绑、团队关系、CRM 配置/队列、CRM worker 推送已实现 | 核心 API 集成测试、CRM worker 集成测试 | 增长链路基础可用,真实微信小程序码、CRM 分配策略、富卡片和销售转化看板待补 | +| 销售/代理客资追踪 | 已建推荐码、首绑客资、团队关系、小程序码缓存、CRM 队列 | 旧 `referral_tracks` 已有映射基础 | 邀请码、扫码/分享事件、首绑保护、销售统计、客资明细、手动补绑、团队关系、CRM 配置/队列、跟进分配策略、死信失败池、脱敏日志、重试/忽略和 CRM worker 推送已实现 | 核心 API 集成测试、CRM worker 集成测试 | 增长链路基础可用,真实微信小程序码、富卡片模板、外部失败告警升级和销售转化看板待补 | | 租户后台 | 已建品牌、域名、设置、支付账户、登录 provider、私密密钥表、成员、审计日志、资源台账、导入台账、内容导航台账 | 不适用 | 概览、品牌、设置、域名、支付账户、登录配置、密钥掩码、活动内容、兑换码/优惠券、成员管理、权限矩阵、审计查询、角色模板权限/菜单/模块/字段/数据范围配置、内容入口/分类树/题目集合/练习蓝图维护、资源管理、题目/单词/知识手册/分数线/视频 JSON/CSV/Excel 同步/异步导入已实现 | 核心 API 集成测试含角色/权限/租户隔离/密钥不泄露/导航/组卷/资源与导入断言 | 租户配置与运营闭环可用;Taro 已接角色模板操作台、字段映射操作台和导入复检结果面板第一版;继续补成员绑定模板、权限驱动菜单和更细数据范围 UI | | 平台后台 | 已建 SaaS 套餐、订阅、账单、服务费、用量、审计日志、催缴台账、催缴通知事件、平台权限字段和平台员工状态字段 | 不适用 | 租户管理、租户详情、账务资料维护、平台账号权限目录、平台员工列表/创建/编辑/启停、平台路由细粒度权限强校验、平台审计日志、账单、订阅账单候选预览、dry-run、批量生成、自动计费 worker、重复开票保护、用量超额账单候选预览/dry-run/生成/自动开票 worker、超额开票失败审计告警、收款确认、逾期标记、内部催缴记录、催缴外部通知渠道/事件、用量台账、平台用量自动采集 worker、平台管理员 Supabase JWT 鉴权已实现 | API 集成测试已覆盖平台细粒度权限、平台员工创建/权限目录/JWT 访问/越权拒绝/自降级拒绝/禁用后 JWT 拒绝/审计脱敏、平台租户创建、详情、账务资料更新、状态变更、审计查询、订阅批量开票、重复保护、用量超额账单候选/dry-run/生成/重复保护、逾期 dry-run/处理/提醒查询、催缴通知渠道/事件脱敏、非法输入拒绝和学生越权拒绝;`npm run test:worker:platform-billing` 覆盖自动计费幂等和审计,`npm run test:worker:platform-usage` 覆盖月度用量自动采集、手工调整不覆盖和幂等,`npm run test:worker:platform-usage-overage` 覆盖超额账单生成幂等、失败审计告警和敏感错误脱敏,`npm run test:worker:platform-dunning` 覆盖逾期催缴幂等和审计,`npm run test:worker:platform-dunning-notifications` 覆盖催缴外部通知幂等、联系方式掩码和密钥不泄露;Taro 类型检查覆盖平台账务和平台员工管理页面 | 平台收费、租户运营和员工授权链路骨架可用,平台在线收款和更完整平台审计报表待补 | | 登录认证 | 已建短信验证码、会话、OAuth provider 配置表,并支持 `auth_user_id` 映射 | 旧用户映射已预留 | 短信 mock 登录、迁移期 session、Supabase JWT 验签映射、微信小程序登录主链路、微信网页登录、QQ 登录、手机号绑定/换绑已实现 | API 集成测试 | H5 Supabase Auth 可联调;真实短信/OAuth 生产账号和回调域名联调待补 | @@ -229,6 +229,9 @@ referral/crm: GET /api/crm/config PUT /api/crm/config GET /api/crm/queue + GET /api/crm/dead-letters + GET /api/crm/queue/logs + POST /api/crm/queue/action tenant-admin: GET /api/tenant-admin/permissions diff --git a/docs/refactor/legacy-feature-gap-matrix.md b/docs/refactor/legacy-feature-gap-matrix.md index 61336067..651197cf 100644 --- a/docs/refactor/legacy-feature-gap-matrix.md +++ b/docs/refactor/legacy-feature-gap-matrix.md @@ -66,7 +66,7 @@ | 分数线维护 | 已覆盖 | 字段/院校/专业/记录 CRUD 和 JSON 批量导入已有 | | 视频维护/绑定 | 已覆盖 | video CRUD、question-video 绑定和视频 JSON 批量导入已有 | | 用户站内通知 | 部分覆盖 | 反馈处理、反馈奖励、勋章发放和积分兑换会写入 `user_notifications`,学生可查/标记状态,租户后台可按权限查看租户内通知;Taro 学生个人中心已接消息摘要,独立消息中心已接状态/类型筛选、批量已读、已读/归档/忽略和站内安全跳转第一版,租户营销中心已接用户通知查看第一版 | 外部微信订阅消息/短信、批量统计和运营效果看板待补 | -| CRM 配置和队列 | 部分覆盖 | 配置/队列、`none/direct/round_robin/referrer` 跟进分配策略、学生批量 CRM 推送、跟进效果统计、钉钉/飞书/企微 worker、签名和重试已有;富卡片模板、失败告警、死信运营台和更细销售转化看板待补 | +| CRM 配置和队列 | 部分覆盖 | 配置/队列、`none/direct/round_robin/referrer` 跟进分配策略、学生批量 CRM 推送、跟进效果统计、钉钉/飞书/企微 worker、签名、重试、死信失败池、脱敏日志、手动重试/忽略和审计已有;富卡片模板、外部失败告警升级和更细销售转化看板待补 | | 对象存储配置 | 部分覆盖 | 系统 env provider、上传签名、上传确认、预览下载签名、资源访问审计、水印 traceId、内置安全扫描和外部 HTTP scanner 接入层已有;租户级存储策略、CDN/转码级水印/真实 AV 服务联调待补 | ## 平台 SaaS 后台功能 @@ -102,7 +102,7 @@ 4. 题库导出:服务端 JSON/试卷 payload、PDF/Word 二进制生成、水印、资料发布路径、权限审计、答案脱敏、每日一练基础导出和图片 ZIP 素材包已补;仍缺更精细试卷模板和后台导出操作台体验。 5. 导入扩展:题目/单词/知识手册/分数线/视频已支持 JSON、CSV 和 Excel 预览导入,并可用 `executionMode=async` 进入 imports worker;导入后复检、导入任务详情、模板下载按钮、字段映射 API、Taro 字段别名编辑、异步轮询和 PocketBase JSON dry-run 报告已补,仍缺真实数据执行验收。 6. 公共题库商业化:平台公共/地区题库授权、租户快照采纳、手动同步、自动同步 worker、同步通知、冲突查询、租户自改冲突保护和单条/批量冲突处理已完成基础闭环;还需生产定时调度、失败告警和更完整运营后台消息。 -7. CRM/销售结算:CRM worker、跟进分配策略、学生批量 CRM 推送、跟进效果统计、分佣规则、结算单、审核、打款状态、导出和凭证复核基础闭环已完成;仍缺富卡片模板、失败告警、死信运营台、真实打款 provider、发票和销售结算看板。 +7. CRM/销售结算:CRM worker、跟进分配策略、学生批量 CRM 推送、跟进效果统计、死信失败池、脱敏日志、手动重试/忽略、分佣规则、结算单、审核、打款状态、导出和凭证复核基础闭环已完成;仍缺富卡片模板、外部失败告警升级、真实打款 provider、发票和销售结算看板。 8. AI 择校推荐增强:后端 `local_rules` 地基、SVIP 门禁、报告列表/详情、本人报告 Markdown/HTML 导出和 Taro 基础页已完成;仍缺真实 AI provider、prompt 编排、PDF worker 报告和后台运营配置。 9. 题目反馈增强:站内通知后端第一版已完成;仍缺前端消息中心、外部订阅消息/短信提醒、问题聚合统计和内容修复闭环。 10. 积分活动增强:积分兑换、活动任务、优惠券兑换履约、后台配置和积分风控只读报表第一阶段已完成;仍缺活动效果看板和更细系统任务触发。 diff --git a/docs/refactor/next-development-todo.md b/docs/refactor/next-development-todo.md index 72fdffba..d02ff9a3 100644 --- a/docs/refactor/next-development-todo.md +++ b/docs/refactor/next-development-todo.md @@ -189,12 +189,13 @@ 3. CRM worker - 已完成 `apps/worker` CRM 队列消费、generic webhook、钉钉、飞书、企业微信机器人 adapter、签名、失败重试和日志。 - 已完成 CRM 配置中的跟进分配策略:`none/direct/round_robin/referrer`,后端校验候选人必须是当前租户内 active 的销售/代理/运营/教师/管理员,首绑客资后写入 `assignedToUserId` 并把 assignee 放进 CRM 队列 payload。 - - 继续补富卡片模板、失败告警、死信运营后台和更细销售转化看板。 + - 已完成死信失败池、脱敏日志查看、手动重试/忽略和审计闭环,Taro 租户营销中心已接第一版操作台。 + - 继续补富卡片模板、外部失败告警升级和更细销售转化看板。 4. 运维 - 后台操作审计报表。 - 定时备份、恢复演练。 - - 性能压测、慢 SQL、索引审查;当前已有 API 压测脚本和 4 核 16G runbook,待云服务器部署后输出正式容量报告。 + - 性能压测、慢 SQL、索引审查;当前已有 API 压测脚本、真实刷题读写闭环和 4 核 16G runbook,本地真实迁移数据已完成 30/50/100/150 并发阶梯压测,云服务器部署后还要按同一矩阵复跑并输出正式容量报告。 ## Taro 前端开发 TODO diff --git a/docs/refactor/performance-benchmark-runbook.md b/docs/refactor/performance-benchmark-runbook.md index af19de2d..e6df9716 100644 --- a/docs/refactor/performance-benchmark-runbook.md +++ b/docs/refactor/performance-benchmark-runbook.md @@ -60,7 +60,10 @@ npm run perf:api:local | `PERF_CONCURRENCY` | `6` | 并发 worker 数 | | `PERF_RAMP_SECONDS` | `3` | 并发爬坡时间 | | `PERF_QUESTION_LIMIT` | `20` | 单次合集题目请求数量 | -| `PERF_INCLUDE_WRITES` | `false` | 是否加入创建练习 session 写请求 | +| `PERF_INCLUDE_WRITES` | `false` | 是否加入真实刷题读写闭环 | +| `PERF_PRACTICE_FLOW_RATIO` | `0.15` | 开启写入后,worker 触发刷题闭环的概率,建议 0.05 到 0.15 | +| `PERF_PRACTICE_FLOW_ANSWERS` | `3` | 每个刷题闭环提交的题目数量 | +| `PERF_ENSURE_SVIP_FOR_WRITES` | `true` | 写压测时如果压测学生没有 SVIP,自动创建 1 天测试权益,避免免费额度污染结果 | | `PERF_INCLUDE_LEADERBOARD` | `false` | 是否加入排行榜接口;排行榜租户默认关闭,仅在租户明确开启并需要专项压测时打开 | | `PERF_OUTPUT_DIR` | `docs/refactor/performance-reports` | 报告输出目录 | | `PERF_REQUEST_TIMEOUT_MS` | `15000` | 单请求超时 | @@ -98,7 +101,9 @@ $env:PERF_INCLUDE_WRITES="true" npm run perf:api:local ``` -开启后会加入低权重 `POST /api/learning/practice-sessions`,用于观察组卷、免费额度、权益校验和 session 写入开销。 +开启后会按配置比例执行完整刷题闭环:`POST /api/learning/practice-sessions` 创建 session,`GET /api/learning/practice-sessions/detail` 拉题,`POST /api/learning/answers` 提交若干题答案,`POST /api/learning/practice-sessions/submit` 交卷,再 `GET /api/learning/practice-sessions/report` 读取报告。脚本会优先选择已有 SVIP 权益学生;若测试库中没有权益且 `PERF_ENSURE_SVIP_FOR_WRITES=true`,会给压测学生创建 1 天 `performance_benchmark` 来源的测试权益,避免每日免费额度导致大量 403 干扰容量结论。 + +压测并发 worker 是“无停顿请求流”,不能直接等同于真实在线学生数。真实学生在线刷题会有读题、思考、翻页和网络间隔。容量估算建议先用报告的成功 RPS,再按前端真实埋点得到的单人 RPS 折算。例如每个真实学生平均 0.05 到 0.2 请求/秒,则 700 req/s 理论上约对应 3500 到 14000 名活跃在线学生的请求吞吐;正式承诺仍要以上云 4 核 16G 环境、CDN、对象存储和真实前端埋点复测为准。 ## 4 核 16G 阶梯压测建议 @@ -110,9 +115,10 @@ npm run perf:api:local | baseline-10 | 10 | 2min | 否 | 学生正常浏览/刷题入口 | | baseline-30 | 30 | 5min | 否 | 中小租户晚高峰 | | baseline-50 | 50 | 5min | 否 | 单机读路径压力观察 | -| mixed-30 | 30 | 5min | 是 | 加入少量创建练习 session | +| mixed-30 | 30 | 5min | 是 | 加入完整刷题闭环,观察组卷、答题、交卷和报告 | | mixed-50 | 50 | 5min | 是 | 观察写入、锁和连接池 | | spike-100 | 100 | 2min | 否 | 短峰值和缓存命中观察 | +| spike-100-write | 100 | 2min | 是 | 短峰值刷题闭环,找 P95/P99 拐点 | PowerShell 示例: diff --git a/docs/refactor/performance-benchmark-summary-20260630.md b/docs/refactor/performance-benchmark-summary-20260630.md new file mode 100644 index 00000000..ad252162 --- /dev/null +++ b/docs/refactor/performance-benchmark-summary-20260630.md @@ -0,0 +1,39 @@ +# 真实迁移数据 API 压测摘要 + +更新时间:2026-06-30 + +这份文件只记录脱敏后的聚合指标,便于 README、上线门禁和后续 AI 开发继续引用。原始 JSON/Markdown 报告位于已忽略的 `docs/refactor/performance-reports/`,不要提交到 Git。 + +## 测试口径 + +- 环境:本地 Docker Desktop + 本地 Supabase/PostgreSQL + 本地 API 进程。 +- 数据:PocketBase 真实导出数据导入新 PostgreSQL 后压测。 +- 数据规模:约 74,102 道题、1,597 个题目合集、3,102 个练习蓝图、3,500 个单词、2,676 条知识手册、3,670 个用户。 +- 写入流量:开启真实刷题闭环,包含创建 session、拉取 session detail、提交若干答案、交卷和读取报告。 +- 排行榜:未纳入默认压测,当前产品默认关闭排行榜。 +- 说明:压测 worker 是无停顿请求流,不等同于真实在线学生数。真实在线容量需要前端埋点后按单个学生平均 RPS 折算。 + +## 结果 + +| 并发 worker | 时长 | 刷题写入比例 | 请求数 | 错误率 | 吞吐 | P95 | P99 | +| ---: | ---: | ---: | ---: | ---: | ---: | ---: | ---: | +| 30 | 120s | 10% | 112,896 | 0.00% | 934.74 req/s | 64.26 ms | 81.14 ms | +| 50 | 120s | 10% | 92,737 | 0.00% | 767.24 req/s | 121.51 ms | 159.12 ms | +| 100 | 120s | 8% | 84,608 | 0.00% | 699.27 req/s | 254.69 ms | 331.84 ms | +| 150 | 120s | 6% | 82,693 | 0.00% | 682.40 req/s | 369.90 ms | 493.69 ms | + +## 初步结论 + +- 本地 Docker 环境下,100 个无停顿 worker 内 P95 仍低于 300ms,可以作为当前代码和索引状态的本地舒适区参考。 +- 150 个无停顿 worker 仍然 0 错误,但 P95 上升到约 370ms,已经能看到压力拐点。 +- 如果未来前端真实埋点显示每名在线学生平均 0.05 到 0.2 req/s,则 700 req/s 理论吞吐约对应 3,500 到 14,000 名活跃在线学生请求量;这只是吞吐换算,不是生产 SLA。 +- 正式容量承诺必须在目标 4 核 16G 云服务器、生产 PostgreSQL 参数、对象存储/CDN、真实前端请求节奏和生产网络下复跑。 + +## 后续复测 + +上云后按 `docs/refactor/performance-benchmark-runbook.md` 复跑至少以下场景: + +1. 只读 30/50/100 并发阶梯。 +2. 混合读写 30/50/100 并发阶梯。 +3. 结合 `pg_stat_statements`、慢 SQL、CPU、I/O、连接池和锁等待确认瓶颈。 +4. 把代表性结果写入本地 `docs/refactor/production-launch-evidence.json`,再执行 `npm run launch:gate`。 diff --git a/docs/refactor/taro-frontend-integration.md b/docs/refactor/taro-frontend-integration.md index 932d1189..dfce6cca 100644 --- a/docs/refactor/taro-frontend-integration.md +++ b/docs/refactor/taro-frontend-integration.md @@ -3164,7 +3164,7 @@ GET /api/ai/school-recommendations/export?reportId=&format=html - 勋章:`GET/PUT /api/tenant-admin/badges`、`GET/POST /api/tenant-admin/badge-grants` - 考试日期:`GET/PUT /api/tenant-admin/exam-dates` - 题目反馈:`GET /api/tenant-admin/feedbacks`、`GET /api/tenant-admin/feedbacks/report`、`POST /api/tenant-admin/feedbacks/status`、`GET /api/tenant-admin/feedbacks/events` -- 销售/代理/CRM 队列、CRM 配置、CRM 跟进分配策略、分佣规则、成员分佣比例、分佣订单、结算单审核和线下打款登记;当前 Taro 租户营销中心已接 CRM、分佣、优惠券规则和核销报表第一版,财务运营页已接退款、官方账单、对账异常、差错工单和调整凭证第一版。真实打款 provider、发票、生产账单抽样验收和更完整财务复核体验后续增强。 +- 销售/代理/CRM 队列、CRM 配置、CRM 跟进分配策略、CRM 死信失败池、脱敏日志、手动重试/忽略、分佣规则、成员分佣比例、分佣订单、结算单审核和线下打款登记;当前 Taro 租户营销中心已接 CRM 配置、队列、死信失败池、日志、重试/忽略、分佣、优惠券规则和核销报表第一版,财务运营页已接退款、官方账单、对账异常、差错工单和调整凭证第一版。真实打款 provider、发票、生产账单抽样验收和更完整财务复核体验后续增强。 CRM 分配策略由后端执行,前端只提交配置: @@ -3177,6 +3177,16 @@ CRM 分配策略由后端执行,前端只提交配置: `assignmentMode` 支持 `none`、`direct`、`round_robin`、`referrer`。后端会校验 `assignmentPool` 中的用户必须是当前租户内 active 的销售、代理、运营、教师或管理员;跨租户成员会返回 `CRM_ASSIGNMENT_POOL_INVALID`。首绑客资成功后,`lead.item.assignedToUserId` 是 CRM 跟进负责人,`lead.item.referrerUserId` 仍是受首绑保护的推广/分佣归属,两者不要在前端混用。CRM 队列 `payload.assignee` 可用于展示推送目标,但前端不能自行改写客资归属或分配游标。 +CRM 死信运营统一走后端命令层: + +```text +GET /api/crm/dead-letters +GET /api/crm/queue/logs?queueId= +POST /api/crm/queue/action +``` + +前端只展示后端返回的脱敏 payload/log,不能从日志里拼 webhook URL、secret 或 token。`retry` 用于 webhook 恢复后重新入队,`ignore` 用于租户确认不再处理的失败任务;按钮需要二次确认并要求填写原因,后端会记录操作者和审计。`pending/processing/sent` 任务不能被忽略,前端应按接口错误提示展示,不做本地状态覆盖。 + 租户后台不应在前端自行决定权限;隐藏菜单只是体验优化,接口仍会校验权限。角色模板用于让租户配置“运营、教师、销售、代理”等自定义后台体验,成员绑定模板后,前端按模板的菜单/模块/字段权限渲染,后端按 permission keys 执行真正的访问控制。班级/学生范围权限由后端根据角色、模板 `dataScope.classIds` 和 `tenant_class_members` 计算,教师默认只能看到自己负责班级。 ## 联调顺序 diff --git a/scripts/api-integration-test.js b/scripts/api-integration-test.js index 9247faff..97735348 100644 --- a/scripts/api-integration-test.js +++ b/scripts/api-integration-test.js @@ -76,6 +76,7 @@ const ids = { publicQuestionBankGrant: '00000000-0000-0000-0000-000000000906', partnerSubscription: '00000000-0000-0000-0000-000000000902', platformOverdueInvoice: crypto.randomUUID(), + crmDeadLetterQueue: crypto.randomUUID(), }; const paymentFixture = (() => { @@ -3235,6 +3236,7 @@ async function testProfile() { } async function testLearningLeaderboard() { + await setTenantFeatureFlag(MAIN_TENANT_ID, 'enableLeaderboard', false); const disabled = await request('/api/learning/leaderboard', { query: { metric: 'questions', period: 'all', limit: 10 }, expectStatus: 403, @@ -3242,55 +3244,69 @@ async function testLearningLeaderboard() { assert.equal(disabled.code, 'LEADERBOARD_DISABLED', 'leaderboard should be disabled by default for tenants'); await setTenantFeatureFlag(MAIN_TENANT_ID, 'enableLeaderboard', true); + try { + const questions = await request('/api/learning/leaderboard', { + query: { metric: 'questions', period: 'all', limit: 10 }, + }); + assert.equal(questions.metric, 'questions', 'question leaderboard should echo metric'); + assert.ok(Array.isArray(questions.items), 'question leaderboard should return list items'); + const currentQuestions = questions.currentUser; + assert.equal(currentQuestions?.userId, USER_ID, 'question leaderboard should include current user rank'); + assert.equal(currentQuestions?.isCurrentUser, true, 'current user rank should be flagged'); - const questions = await request('/api/learning/leaderboard', { - query: { metric: 'questions', period: 'all', limit: 10 }, - }); - assert.equal(questions.metric, 'questions', 'question leaderboard should echo metric'); - assert.ok(questions.items?.some(item => item.userId === SECOND_STUDENT_USER_ID), 'question leaderboard should include second student'); - assert.ok(questions.items?.some(item => item.userId === USER_ID), 'question leaderboard should include current student'); - const secondQuestions = questions.items?.find(item => item.userId === SECOND_STUDENT_USER_ID); - const currentQuestions = questions.currentUser; - assert.ok(secondQuestions?.value >= currentQuestions?.value, 'second student should rank at least as high as smoke user by questions'); - assert.equal(currentQuestions?.userId, USER_ID, 'question leaderboard should include current user rank'); - assert.equal(currentQuestions?.isCurrentUser, true, 'current user rank should be flagged'); + const classScoped = await request('/api/learning/leaderboard', { + query: { metric: 'questions', classId: ids.tenantClass, limit: 10 }, + }); + assert.ok(classScoped.items?.some(item => item.userId === USER_ID), 'class scoped leaderboard should include class student'); + assert.ok(!classScoped.items?.some(item => item.userId === SECOND_STUDENT_USER_ID), 'class scoped leaderboard should exclude other class student'); - const score = await request('/api/learning/leaderboard', { - query: { metric: 'score', period: 'all', limit: 10 }, - }); - const scoreLeader = score.items?.find(item => item.userId === SECOND_STUDENT_USER_ID); - assert.ok(scoreLeader?.value >= 30, 'score leaderboard should use platform user score'); - assert.ok(score.currentUser?.value >= 10, 'score leaderboard should include current user score'); + const secondClassScoped = await request('/api/learning/leaderboard', { + query: { metric: 'questions', classId: ids.tenantClassOther, limit: 10 }, + }); + const secondQuestions = secondClassScoped.items?.find(item => item.userId === SECOND_STUDENT_USER_ID); + assert.ok(secondQuestions, 'second student class scoped leaderboard should include second student'); + assert.ok(secondQuestions?.value >= currentQuestions?.value, 'second student should rank at least as high as smoke user by questions'); - const vocabulary = await request('/api/learning/leaderboard', { - query: { metric: 'vocabulary', period: 'all', limit: 10 }, - }); - const vocabularyLeader = vocabulary.items?.find(item => item.userId === SECOND_STUDENT_USER_ID); - assert.equal(vocabularyLeader?.value, 2, 'vocabulary leaderboard should count mastered words'); + const score = await request('/api/learning/leaderboard', { + query: { metric: 'score', period: 'all', classId: ids.tenantClassOther, limit: 10 }, + }); + const scoreLeader = score.items?.find(item => item.userId === SECOND_STUDENT_USER_ID); + assert.ok(scoreLeader?.value >= 30, 'score leaderboard should use platform user score'); - const mockExam = await request('/api/learning/leaderboard', { - query: { metric: 'mock_exam', period: '30d', limit: 10 }, - }); - const mockLeader = mockExam.items?.find(item => item.userId === SECOND_STUDENT_USER_ID); - assert.equal(mockLeader?.value, 95, 'mock exam leaderboard should use best report score'); - assert.ok(mockExam.currentUser?.value >= 70, 'mock exam leaderboard should include current user best score'); + const currentScore = await request('/api/learning/leaderboard', { + query: { metric: 'score', period: 'all', classId: ids.tenantClass, limit: 10 }, + }); + assert.ok(currentScore.currentUser?.value >= 10, 'score leaderboard should include current user score'); - const classScoped = await request('/api/learning/leaderboard', { - query: { metric: 'questions', classId: ids.tenantClass, limit: 10 }, - }); - assert.ok(classScoped.items?.some(item => item.userId === USER_ID), 'class scoped leaderboard should include class student'); - assert.ok(!classScoped.items?.some(item => item.userId === SECOND_STUDENT_USER_ID), 'class scoped leaderboard should exclude other class student'); + const vocabulary = await request('/api/learning/leaderboard', { + query: { metric: 'vocabulary', period: 'all', classId: ids.tenantClassOther, limit: 10 }, + }); + const vocabularyLeader = vocabulary.items?.find(item => item.userId === SECOND_STUDENT_USER_ID); + assert.equal(vocabularyLeader?.value, 2, 'vocabulary leaderboard should count mastered words'); - const trustedLogin = await loginBySms('13800000000'); - const crossTenantDenied = await request('/api/learning/leaderboard', { - tenantId: PARTNER_TENANT_ID, - userId: false, - headers: { authorization: `Bearer ${trustedLogin.session.token}` }, - expectStatus: 403, - }); - assert.equal(crossTenantDenied.code, 'AUTH_TENANT_MISMATCH', 'leaderboard must reject trusted session cross-tenant access'); + const mockExam = await request('/api/learning/leaderboard', { + query: { metric: 'mock_exam', period: '30d', classId: ids.tenantClassOther, limit: 10 }, + }); + const mockLeader = mockExam.items?.find(item => item.userId === SECOND_STUDENT_USER_ID); + assert.equal(mockLeader?.value, 95, 'mock exam leaderboard should use best report score'); - await setTenantFeatureFlag(MAIN_TENANT_ID, 'enableLeaderboard', false); + const currentMockExam = await request('/api/learning/leaderboard', { + query: { metric: 'mock_exam', period: '30d', classId: ids.tenantClass, limit: 10 }, + }); + assert.ok(currentMockExam.currentUser?.value >= 70, 'mock exam leaderboard should include current user best score'); + + const trustedLogin = await loginBySms('13800000000'); + const crossTenantDenied = await request('/api/learning/leaderboard', { + tenantId: PARTNER_TENANT_ID, + userId: false, + headers: { authorization: `Bearer ${trustedLogin.session.token}` }, + expectStatus: 403, + }); + assert.equal(crossTenantDenied.code, 'AUTH_TENANT_MISMATCH', 'leaderboard must reject trusted session cross-tenant access'); + + } finally { + await setTenantFeatureFlag(MAIN_TENANT_ID, 'enableLeaderboard', false); + } const disabledAfterRestore = await request('/api/learning/leaderboard', { query: { metric: 'questions', period: 'all', limit: 10 }, expectStatus: 403, @@ -10363,6 +10379,170 @@ async function testReferralAndCrmGrowth() { ); assert.ok(!JSON.stringify(crmQueue).includes('crm-secret-smoke'), 'CRM queue should not leak secret'); + const crmDeadLetterPool = new pg.Pool({ connectionString: process.env.DATABASE_URL || DEFAULT_DATABASE_URL }); + try { + await crmDeadLetterPool.query( + ` + insert into public.crm_webhook_queue ( + id, tenant_id, record_id, status, scheduled_at, attempts, next_attempt_at, + last_error, last_http_code, lead_id, source, payload, idempotency_key, + target_url, provider, last_attempt_at, last_response_summary, dead_lettered_at + ) + values ( + $1::uuid, $2, $3, 'failed', now(), 3, null, + $4, 500, $5, 'tenant.student.crm_push', + $6::jsonb, $7, $8, 'dingtalk', now(), $9, now() + ) + on conflict (id) do update set status = excluded.status, + last_error = excluded.last_error, + target_url = excluded.target_url, + payload = excluded.payload, + dead_lettered_at = now(), + updated_at = now() + `, + [ + ids.crmDeadLetterQueue, + MAIN_TENANT_ID, + USER_ID, + 'HTTP 500 access_token=dead-letter-secret should be redacted', + manual.lead.item.id, + JSON.stringify({ + eventType: 'student.crm_push', + token: 'payload-secret-token', + student: { userId: USER_ID, avatarPreset: 'male' }, + }), + 'integration-crm-dead-letter', + 'https://oapi.dingtalk.com/robot/send?access_token=dead-letter-secret', + 'response contains access_token=dead-letter-secret', + ], + ); + await crmDeadLetterPool.query( + ` + insert into public.crm_webhook_log ( + tenant_id, queue_id, record_id, http_code, outcome, error_message, lead_id, + request_body, request_payload, response_summary, signed_at, attempt + ) + values ( + $1, $2::uuid, $3, 500, 'failed', $4, $5, + $6, $7::jsonb, $8, now(), 3 + ) + `, + [ + MAIN_TENANT_ID, + ids.crmDeadLetterQueue, + USER_ID, + 'access_token=dead-letter-secret failed', + manual.lead.item.id, + '{"access_token":"dead-letter-secret","student":{"avatarPreset":"male"}}', + JSON.stringify({ url: 'https://oapi.dingtalk.com/robot/send?access_token=dead-letter-secret', token: 'dead-letter-secret' }), + 'remote response access_token=dead-letter-secret', + ], + ); + } finally { + await crmDeadLetterPool.end(); + } + + const crmDeadLetters = await request('/api/crm/dead-letters', { + userId: TENANT_ADMIN_USER_ID, + query: { limit: 10 }, + }); + const deadLetterItem = crmDeadLetters.items?.find(item => item.id === ids.crmDeadLetterQueue); + assert.ok(deadLetterItem, 'CRM dead-letter list should include failed task'); + assert.equal(crmDeadLetters.summary?.failed >= 1, true, 'CRM dead-letter summary should count failed tasks'); + assert.ok(!JSON.stringify(crmDeadLetters).includes('dead-letter-secret'), 'CRM dead-letter list must redact webhook secrets'); + assert.equal(deadLetterItem?.payload?.token, '[redacted]', 'CRM dead-letter payload should redact token fields'); + + const crmDeadLetterLogs = await request('/api/crm/queue/logs', { + userId: TENANT_ADMIN_USER_ID, + query: { queueId: ids.crmDeadLetterQueue, limit: 10 }, + }); + assert.ok(crmDeadLetterLogs.items?.some(item => item.outcome === 'failed'), 'CRM queue logs should include failed attempts'); + assert.ok(!JSON.stringify(crmDeadLetterLogs).includes('dead-letter-secret'), 'CRM queue logs must redact webhook secrets'); + + const crmRetryDenied = await request('/api/crm/queue/action', { + userId: TENANT_OPERATOR_USER_ID, + method: 'POST', + body: { + queueId: ids.crmDeadLetterQueue, + action: 'retry', + }, + expectStatus: 403, + }); + assert.equal(crmRetryDenied.code, 'TENANT_PERMISSION_REQUIRED', 'CRM retry should require crm:write'); + + const crmCrossTenantDenied = await request('/api/crm/queue/logs', { + tenantId: PARTNER_TENANT_ID, + userId: PARTNER_TENANT_ADMIN_USER_ID, + query: { queueId: ids.crmDeadLetterQueue }, + expectStatus: 404, + }); + assert.equal(crmCrossTenantDenied.code, 'CRM_QUEUE_TASK_NOT_FOUND', 'CRM queue logs must be tenant isolated'); + + const crmRetried = await request('/api/crm/queue/action', { + userId: TENANT_ADMIN_USER_ID, + method: 'POST', + body: { + queueId: ids.crmDeadLetterQueue, + action: 'retry', + note: 'integration retry', + metadata: { source: 'integration-test', accessToken: 'dead-letter-secret' }, + }, + }); + assert.equal(crmRetried.item?.status, 'pending', 'CRM retry should reset failed task to pending'); + assert.equal(crmRetried.item?.lastError, null, 'CRM retry should clear last error'); + assert.equal(crmRetried.item?.lastOperatorAction, 'retry', 'CRM retry should record operator action'); + assert.ok(!JSON.stringify(crmRetried).includes('dead-letter-secret'), 'CRM retry response must redact metadata secrets'); + + const crmRetriedLogs = await request('/api/crm/queue/logs', { + userId: TENANT_ADMIN_USER_ID, + query: { queueId: ids.crmDeadLetterQueue, limit: 10 }, + }); + assert.ok(crmRetriedLogs.items?.some(item => item.operation === 'retry'), 'CRM retry should append operator log'); + assert.ok(!JSON.stringify(crmRetriedLogs).includes('dead-letter-secret'), 'CRM retry operator log must redact metadata secrets'); + + const crmIgnorePending = await request('/api/crm/queue/action', { + userId: TENANT_ADMIN_USER_ID, + method: 'POST', + body: { + queueId: ids.crmDeadLetterQueue, + action: 'ignore', + }, + expectStatus: 409, + }); + assert.equal(crmIgnorePending.code, 'CRM_QUEUE_STATUS_NOT_IGNORABLE', 'CRM ignore should only apply to dead-letter tasks'); + + const crmDeadLetterResetPool = new pg.Pool({ connectionString: process.env.DATABASE_URL || DEFAULT_DATABASE_URL }); + try { + await crmDeadLetterResetPool.query( + ` + update public.crm_webhook_queue + set status = 'failed', + attempts = 4, + last_error = 'retry failed with access_token=dead-letter-secret', + last_http_code = 500, + dead_lettered_at = now(), + updated_at = now() + where tenant_id = $1 and id = $2::uuid + `, + [MAIN_TENANT_ID, ids.crmDeadLetterQueue], + ); + } finally { + await crmDeadLetterResetPool.end(); + } + + const crmIgnored = await request('/api/crm/queue/action', { + userId: TENANT_ADMIN_USER_ID, + method: 'POST', + body: { + queueId: ids.crmDeadLetterQueue, + action: 'ignore', + note: 'integration ignored after manual review', + }, + }); + assert.equal(crmIgnored.item?.status, 'discarded', 'CRM ignore should move failed task to discarded'); + assert.equal(crmIgnored.item?.ignoredBy, TENANT_ADMIN_USER_ID, 'CRM ignore should record operator'); + assert.equal(crmIgnored.item?.lastOperatorAction, 'ignore', 'CRM ignore should record operator action'); + const partnerReferralDenied = await request('/api/referral/sales-stats', { tenantId: PARTNER_TENANT_ID, userId: TENANT_SALES_USER_ID, diff --git a/scripts/api-performance-benchmark.js b/scripts/api-performance-benchmark.js index 8a9ec883..10d7e897 100644 --- a/scripts/api-performance-benchmark.js +++ b/scripts/api-performance-benchmark.js @@ -15,6 +15,9 @@ const CONCURRENCY = envNumber('PERF_CONCURRENCY', 6); const RAMP_SECONDS = envNumber('PERF_RAMP_SECONDS', 3, { allowZero: true }); const INCLUDE_WRITES = envBool('PERF_INCLUDE_WRITES', false); const INCLUDE_LEADERBOARD = envBool('PERF_INCLUDE_LEADERBOARD', false); +const PRACTICE_FLOW_RATIO = Math.min(Math.max(envNumber('PERF_PRACTICE_FLOW_RATIO', 0.15, { allowZero: true }), 0), 1); +const PRACTICE_FLOW_ANSWERS = envNumber('PERF_PRACTICE_FLOW_ANSWERS', 3); +const ENSURE_SVIP_FOR_WRITES = envBool('PERF_ENSURE_SVIP_FOR_WRITES', true); const START_PORT = Number(process.env.PERF_API_PORT || 0) || 0; const TENANT_CODE = process.env.PERF_TENANT_CODE || 'master'; const QUESTION_LIMIT = envNumber('PERF_QUESTION_LIMIT', 20); @@ -111,7 +114,7 @@ async function requestJson(endpoint, context, timeoutMs = DEFAULT_TIMEOUT_MS) { message: payload?.error || payload?.message || response.statusText, }; } - return { ok: true, status: response.status, elapsedMs, bytes: JSON.stringify(payload).length }; + return { ok: true, status: response.status, elapsedMs, bytes: JSON.stringify(payload).length, payload }; } catch (error) { return { ok: false, @@ -124,6 +127,164 @@ async function requestJson(endpoint, context, timeoutMs = DEFAULT_TIMEOUT_MS) { } } +function sampleFromResult(name, result, startedAt, finishedAt) { + return { + name, + ok: result.ok, + status: result.status, + elapsedMs: result.elapsedMs, + startedAt, + finishedAt, + }; +} + +function pickAnswerPayload(question) { + const subQuestions = Array.isArray(question.subQuestions) ? question.subQuestions : []; + if (subQuestions.length) { + return { + subAnswers: subQuestions.slice(0, Math.max(PRACTICE_FLOW_ANSWERS, 1)).map((subQuestion, index) => { + const subCorrectIndices = Array.isArray(subQuestion.correctOptionIndices) ? subQuestion.correctOptionIndices : []; + const subCorrectIndex = subQuestion.correctOptionIndex; + const subAnswer = { + subQuestionId: subQuestion.id || subQuestion.subQuestionId || subQuestion.key || `sub_${index + 1}`, + selectedOptions: [], + answerText: '', + }; + if (subCorrectIndices.length) { + subAnswer.selectedOptions = subCorrectIndices.map(item => String(item)); + } else if (subCorrectIndex !== null && subCorrectIndex !== undefined && Number.isFinite(Number(subCorrectIndex))) { + subAnswer.selectedOptions = [String(Number(subCorrectIndex))]; + } else if (Array.isArray(subQuestion.options) && subQuestion.options.length) { + subAnswer.selectedOptions = ['0']; + } else { + subAnswer.answerText = 'performance benchmark answer'; + subAnswer.selfJudgedCorrect = true; + } + return subAnswer; + }), + }; + } + const correctIndices = Array.isArray(question.correctOptionIndices) ? question.correctOptionIndices : []; + const correctIndex = question.correctOptionIndex; + if (correctIndices.length) { + return { selectedOptions: correctIndices.map(item => String(item)) }; + } + if (correctIndex !== null && correctIndex !== undefined && Number.isFinite(Number(correctIndex))) { + return { selectedOptions: [String(Number(correctIndex))] }; + } + if (Array.isArray(question.options) && question.options.length) { + return { selectedOptions: ['0'] }; + } + return { answerText: 'performance benchmark answer', selfJudgedCorrect: true }; +} + +async function runPracticeFlow(context, samples, errors) { + if (!context.blueprint || !context.collection) return false; + const flowStartedAt = performance.now(); + const create = await requestJson({ + name: 'learning.practice_flow.create', + method: 'POST', + path: '/api/learning/practice-sessions', + body: () => ({ + blueprintId: context.blueprint.id, + collectionId: context.collection.id, + mode: context.blueprint.mode || 'sequential', + questionLimit: Math.min(QUESTION_LIMIT, Math.max(PRACTICE_FLOW_ANSWERS, 3)), + metadata: { source: 'api-performance-benchmark' }, + }), + }, context); + const flowFinishedCreate = performance.now(); + samples.push(sampleFromResult('learning.practice_flow.create', create, flowStartedAt, flowFinishedCreate)); + if (!create.ok) { + pushError(errors, 'learning.practice_flow.create', create); + return true; + } + + const practiceSessionId = create.payload?.item?.id; + if (!practiceSessionId) { + pushError(errors, 'learning.practice_flow.create', { + status: 200, + elapsedMs: create.elapsedMs, + error: 'PRACTICE_SESSION_ID_MISSING', + message: 'Practice session create response did not include item.id', + }); + return true; + } + + const detailStartedAt = performance.now(); + const detail = await requestJson({ + name: 'learning.practice_flow.detail', + method: 'GET', + path: '/api/learning/practice-sessions/detail', + query: { practiceSessionId }, + }, context); + const detailFinishedAt = performance.now(); + samples.push(sampleFromResult('learning.practice_flow.detail', detail, detailStartedAt, detailFinishedAt)); + if (!detail.ok) { + pushError(errors, 'learning.practice_flow.detail', detail); + return true; + } + + const questions = Array.isArray(detail.payload?.item?.questions) ? detail.payload.item.questions : []; + for (const question of questions.slice(0, PRACTICE_FLOW_ANSWERS)) { + if (!question?.id) continue; + const answerStartedAt = performance.now(); + const answer = await requestJson({ + name: 'learning.practice_flow.answer', + method: 'POST', + path: '/api/learning/answers', + body: () => ({ + practiceSessionId, + questionId: question.id, + ...pickAnswerPayload(question), + }), + }, context); + const answerFinishedAt = performance.now(); + samples.push(sampleFromResult('learning.practice_flow.answer', answer, answerStartedAt, answerFinishedAt)); + if (!answer.ok) { + pushError(errors, 'learning.practice_flow.answer', answer); + return true; + } + } + + const submitStartedAt = performance.now(); + const submit = await requestJson({ + name: 'learning.practice_flow.submit', + method: 'POST', + path: '/api/learning/practice-sessions/submit', + body: () => ({ practiceSessionId }), + }, context); + const submitFinishedAt = performance.now(); + samples.push(sampleFromResult('learning.practice_flow.submit', submit, submitStartedAt, submitFinishedAt)); + if (!submit.ok) { + pushError(errors, 'learning.practice_flow.submit', submit); + return true; + } + + const reportStartedAt = performance.now(); + const report = await requestJson({ + name: 'learning.practice_flow.report', + method: 'GET', + path: '/api/learning/practice-sessions/report', + query: { practiceSessionId }, + }, context); + const reportFinishedAt = performance.now(); + samples.push(sampleFromResult('learning.practice_flow.report', report, reportStartedAt, reportFinishedAt)); + if (!report.ok) pushError(errors, 'learning.practice_flow.report', report); + return true; +} + +function pushError(errors, endpointName, result) { + if (errors.length >= MAX_ERRORS_TO_KEEP) return; + errors.push({ + endpoint: endpointName, + status: result.status, + error: result.error, + message: result.message, + elapsedMs: round(result.elapsedMs), + }); +} + async function waitForHealth(timeoutMs = 20_000) { const started = Date.now(); let lastError = null; @@ -203,8 +364,14 @@ async function discoverBenchmarkContext() { select pu.id, coalesce(pu.name, pu.username, pu.phone, pu.legacy_id, pu.id::text) as name from public.tenant_memberships tm join public.platform_users pu on pu.id = tm.user_id + left join public.entitlements e on e.tenant_id = tm.tenant_id + and e.user_id = pu.id + and e.entitlement_type = 'svip' + and e.status = 'active' + and (e.expires_at is null or e.expires_at > now()) where tm.tenant_id = $1 and tm.status = 'active' and tm.role = 'student' - order by pu.created_at asc + group by pu.id, pu.name, pu.username, pu.phone, pu.legacy_id, pu.created_at + order by case when count(e.id) > 0 then 0 else 1 end, pu.created_at asc limit 1 `, [tenant.id], @@ -304,6 +471,39 @@ async function discoverBenchmarkContext() { ) : null; + const hasActiveEntitlement = user + ? await one( + pool, + ` + select id + from public.entitlements + where tenant_id = $1 + and user_id = $2 + and entitlement_type = 'svip' + and status = 'active' + and (expires_at is null or expires_at > now()) + limit 1 + `, + [tenant.id, user.id], + ) + : null; + let benchmarkEntitlementCreated = false; + if (INCLUDE_WRITES && ENSURE_SVIP_FOR_WRITES && user && !hasActiveEntitlement) { + await pool.query( + ` + insert into public.entitlements ( + tenant_id, user_id, entitlement_type, scope_type, source_type, + starts_at, expires_at, status, metadata + ) + values ($1, $2, 'svip', 'tenant', 'performance_benchmark', + now(), now() + interval '1 day', 'active', + '{"createdBy":"api-performance-benchmark"}'::jsonb) + `, + [tenant.id, user.id], + ); + benchmarkEntitlementCreated = true; + } + const dbStats = { before: await captureDbStats(pool), tableCounts: await many( @@ -327,7 +527,13 @@ async function discoverBenchmarkContext() { return { tenant, - user, + user: { + ...user, + entitlement: { + hadActiveSvip: Boolean(hasActiveEntitlement), + benchmarkEntitlementCreated, + }, + }, entry, node, collection, @@ -426,24 +632,23 @@ function createEndpointCatalog(context) { { name: 'catalog.handbook_entries', weight: 5, method: 'GET', path: '/api/catalog/handbook-entries', query: { chapterId: context.handbookChapter.id, includeContent: 'true' } }, ); } - if (INCLUDE_WRITES && context.blueprint && context.collection) { - endpoints.push({ - name: 'learning.practice_session.create', - weight: 1, - method: 'POST', - path: '/api/learning/practice-sessions', - body: () => ({ - blueprintId: context.blueprint.id, - collectionId: context.collection.id, - mode: context.blueprint.mode || 'sequential', - questionLimit: Math.min(5, QUESTION_LIMIT), - }), - }); - } - return expandWeightedEndpoints(endpoints); } +function benchmarkEndpointNames(weightedEndpoints, context) { + const names = new Set(weightedEndpoints.map(endpoint => endpoint.name)); + if (INCLUDE_WRITES && context.blueprint && context.collection && PRACTICE_FLOW_RATIO > 0) { + for (const name of [ + 'learning.practice_flow.create', + 'learning.practice_flow.detail', + 'learning.practice_flow.answer', + 'learning.practice_flow.submit', + 'learning.practice_flow.report', + ]) names.add(name); + } + return [...names].sort(); +} + function expandWeightedEndpoints(endpoints) { const weighted = []; for (const endpoint of endpoints) { @@ -542,33 +747,22 @@ async function workerLoop(workerId, weightedEndpoints, context, stopAt, samples, const rampDelay = RAMP_SECONDS > 0 ? (workerId / Math.max(CONCURRENCY, 1)) * RAMP_SECONDS * 1000 : 0; if (rampDelay > 0) await sleep(rampDelay); while (performance.now() < stopAt) { + if (INCLUDE_WRITES && PRACTICE_FLOW_RATIO > 0 && Math.random() < PRACTICE_FLOW_RATIO) { + await runPracticeFlow(context, samples, errors); + continue; + } const endpoint = pickEndpoint(weightedEndpoints); const startedAt = performance.now(); const result = await requestJson(endpoint, context); const finishedAt = performance.now(); - samples.push({ - name: endpoint.name, - ok: result.ok, - status: result.status, - elapsedMs: result.elapsedMs, - startedAt, - finishedAt, - }); - if (!result.ok && errors.length < MAX_ERRORS_TO_KEEP) { - errors.push({ - endpoint: endpoint.name, - status: result.status, - error: result.error, - message: result.message, - elapsedMs: round(result.elapsedMs), - }); - } + samples.push(sampleFromResult(endpoint.name, result, startedAt, finishedAt)); + if (!result.ok) pushError(errors, endpoint.name, result); } } async function runBenchmark(context) { const weightedEndpoints = createEndpointCatalog(context); - const endpointNames = [...new Set(weightedEndpoints.map(endpoint => endpoint.name))].sort(); + const endpointNames = benchmarkEndpointNames(weightedEndpoints, context); const samples = []; const errors = []; const startedAt = new Date(); @@ -596,6 +790,8 @@ async function runBenchmark(context) { concurrency: CONCURRENCY, rampSeconds: RAMP_SECONDS, includeWrites: INCLUDE_WRITES, + practiceFlowRatio: INCLUDE_WRITES ? PRACTICE_FLOW_RATIO : 0, + practiceFlowAnswers: INCLUDE_WRITES ? PRACTICE_FLOW_ANSWERS : 0, questionLimit: QUESTION_LIMIT, timeoutMs: DEFAULT_TIMEOUT_MS, }, @@ -635,6 +831,10 @@ function markdownReport(report) { lines.push(`- 持续时间:${report.config.durationSeconds}s`); lines.push(`- Ramp:${report.config.rampSeconds}s`); lines.push(`- 写入压测:${report.config.includeWrites ? '开启' : '关闭'}`); + if (report.config.includeWrites) { + lines.push(`- 刷题闭环比例:${Math.round((report.config.practiceFlowRatio || 0) * 100)}%`); + lines.push(`- 每个刷题闭环提交答案数:${report.config.practiceFlowAnswers || 0}`); + } lines.push(`- 请求超时:${report.config.timeoutMs}ms`); lines.push(''); lines.push('## 总览'); @@ -682,7 +882,8 @@ function markdownReport(report) { } lines.push('## 说明'); lines.push(''); - lines.push('- 默认压测是只读混合工作负载,适合在真实迁移库上做烟测。'); + lines.push('- 默认压测是只读混合工作负载,适合在真实迁移库上做烟测;开启写入后会执行“创建练习 session -> 读取详情 -> 提交答案 -> 交卷 -> 读取报告”的刷题闭环。'); + lines.push('- 并发 worker 是无停顿请求流,不能直接等同于真实在线人数;真实学生有读题、思考、翻页和网络间隔,应结合前端埋点估算每人 RPS。'); lines.push('- 结果只能代表当前本机 Docker、API 进程和数据库状态;4 核 16G 云服务器需要按 runbook 跑阶梯并发。'); lines.push('- 若 P95 明显升高,下一步应结合 PostgreSQL 慢 SQL、`pg_stat_activity` 和 API 日志定位。'); lines.push(''); diff --git a/supabase/migrations/202606300018_crm_dead_letter_operations.sql b/supabase/migrations/202606300018_crm_dead_letter_operations.sql new file mode 100644 index 00000000..5dddb917 --- /dev/null +++ b/supabase/migrations/202606300018_crm_dead_letter_operations.sql @@ -0,0 +1,42 @@ +alter table public.crm_webhook_queue + add column if not exists dead_lettered_at timestamptz, + add column if not exists ignored_at timestamptz, + add column if not exists ignored_by uuid references public.platform_users(id) on delete set null, + add column if not exists last_operator_user_id uuid references public.platform_users(id) on delete set null, + add column if not exists last_operator_action text, + add column if not exists last_operator_note text, + add column if not exists last_operator_at timestamptz, + add column if not exists operator_metadata jsonb not null default '{}'::jsonb; + +alter table public.crm_webhook_log + add column if not exists queue_id uuid references public.crm_webhook_queue(id) on delete set null, + add column if not exists operator_user_id uuid references public.platform_users(id) on delete set null, + add column if not exists operation text; + +update public.crm_webhook_queue +set dead_lettered_at = coalesce(dead_lettered_at, updated_at, created_at) +where status in ('failed', 'discarded') and dead_lettered_at is null; + +create index if not exists idx_crm_queue_dead_letters + on public.crm_webhook_queue(tenant_id, status, dead_lettered_at desc, updated_at desc) + where status in ('failed', 'discarded'); + +create index if not exists idx_crm_queue_operator + on public.crm_webhook_queue(tenant_id, last_operator_at desc) + where last_operator_at is not null; + +create index if not exists idx_crm_webhook_log_queue + on public.crm_webhook_log(tenant_id, queue_id, created_at desc) + where queue_id is not null; + +comment on column public.crm_webhook_queue.dead_lettered_at is + 'Set when a CRM webhook task reaches a terminal failed/discarded state and needs operator review.'; + +comment on column public.crm_webhook_queue.ignored_at is + 'Set when a tenant operator intentionally dismisses a CRM webhook dead-letter task.'; + +comment on column public.crm_webhook_queue.last_operator_action is + 'Last manual CRM queue operation, for example retry or ignore.'; + +comment on column public.crm_webhook_log.queue_id is + 'Direct reference to the CRM queue task. Older logs may only have record_id or lead_id.';