From 6f0083bd13248ebd3e035efc5d5c18dcc88d51c1 Mon Sep 17 00:00:00 2001 From: Codex Date: Mon, 29 Jun 2026 09:38:11 +0800 Subject: [PATCH] feat: add public question bank sync worker --- .env.example | 6 + README.md | 27 ++- apps/api/src/features/tenant-content/index.ts | 2 + .../features/tenant-content/public-banks.ts | 117 ++++++++- apps/worker/package.json | 3 +- apps/worker/src/config.ts | 8 + apps/worker/src/index.ts | 11 + apps/worker/src/jobs/public-banks.ts | 197 ++++++++++++++++ docs/refactor/architecture.md | 2 +- docs/refactor/backend-capability-status.md | 14 +- docs/refactor/backend-handoff-roadmap.md | 6 +- docs/refactor/backend-progress.md | 3 +- docs/refactor/blueprint-coverage.md | 4 +- docs/refactor/implementation-status.md | 8 +- docs/refactor/legacy-feature-gap-matrix.md | 6 +- docs/refactor/next-development-todo.md | 7 +- docs/refactor/taro-frontend-integration.md | 7 +- package.json | 1 + scripts/api-integration-test.js | 19 ++ .../public-bank-worker-integration-test.js | 223 ++++++++++++++++++ 20 files changed, 632 insertions(+), 39 deletions(-) create mode 100644 apps/worker/src/jobs/public-banks.ts create mode 100644 scripts/public-bank-worker-integration-test.js diff --git a/.env.example b/.env.example index 5c1a672f..a94e9f41 100644 --- a/.env.example +++ b/.env.example @@ -67,6 +67,12 @@ WORKER_ASSET_MIN_AGE_SECONDS=300 WORKER_ASSET_RECHECK_INTERVAL_SECONDS=86400 WORKER_ASSET_REQUEST_TIMEOUT_MS=10000 +# Worker 配置:公共题库自动同步。冲突会保留租户自改题目并等待后台处理。 +WORKER_PUBLIC_BANK_SYNC_BATCH_SIZE=5 +WORKER_PUBLIC_BANK_SYNC_COPY_LIMIT=1000 +WORKER_PUBLIC_BANK_SYNC_ID=public-banks-1 +WORKER_PUBLIC_BANK_SYNC_CLAIM_STALE_SECONDS=900 + # 对象存储配置:API 和 assets worker 共用 STORAGE_DEFAULT_PROVIDER=local_dev STORAGE_DEFAULT_BUCKET=tenant-assets diff --git a/README.md b/README.md index 1f5c6d45..0836e57d 100644 --- a/README.md +++ b/README.md @@ -16,10 +16,10 @@ - 租户内容能力:可配置题库入口、任意深度分类树、考试意向标记、题目集合、顺序/随机/全真模拟蓝图、题目录入/更新、视频绑定、分数线、单词、知识手册、资料资源台账、题目/单词/知识手册/分数线/视频 JSON/CSV/Excel 批量导入。 - 学生端能力:题库入口、分类树、题目集合、顺序/随机/模考 session 组卷快照、答题、错题本、收藏夹、背单词进度、个人中心、勋章、考试倒计时、签到积分、题目反馈、排行榜、分数线、题目视频、订单详情/状态轮询、优惠券领取/抵扣、权益、激活码预检查/兑换、资料下载。 - 平台后台能力:租户管理、SaaS 套餐、订阅、账单、服务费收款、用量记录、公共题库授权。 -- 公共题库商业化能力:租户可采纳平台授权题库为本租户副本,并可手动同步平台新增/更新题目;同步会保护租户自改题目,返回冲突而不覆盖。 +- 公共题库商业化能力:租户可采纳平台授权题库为本租户副本,并可手动或由 worker 自动同步平台新增/更新题目;同步会保护租户自改题目,返回冲突而不覆盖,后台可查询冲突明细。 - 题库导出基础能力:租户内容编辑可按题目集合、内容入口或分类节点导出 JSON、`paper_json` 和打印 payload,后端强制租户隔离、答案/解析开关、复合题子题脱敏、导出 job 和审计。 - 销售/代理/CRM 增长链路:邀请码、扫码/分享事件、首绑客资保护、销售统计、团队关系、CRM 配置和队列。 -- `apps/worker` 后台任务进程:CRM webhook 队列消费、generic/钉钉/飞书/企微机器人发送、签名、失败重试和日志;commerce worker 可补偿查询微信/支付宝支付和退款状态;assets worker 可复检托管资源元数据并自动下架异常资源。 +- `apps/worker` 后台任务进程:CRM webhook 队列消费、generic/钉钉/飞书/企微机器人发送、签名、失败重试和日志;commerce worker 可补偿查询微信/支付宝支付和退款状态;assets worker 可复检托管资源元数据并自动下架异常资源;public-banks worker 可自动同步公共题库采纳副本。 - 销售/代理分佣结算基础闭环:租户默认比例、成员比例、激活码批次比例、订单/激活码归因、结算单生成、审核、线下打款状态和权限隔离。 - 订单售后基础闭环:退款请求、审核、处理状态流、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、退款金额累计、部分/全额退款订单状态、全额退款权益撤销、退款事件和审计日志。 - PocketBase schema/数据导入器雏形和导入后校验脚本。 @@ -32,7 +32,7 @@ - OSS/COS/Supabase Storage 上传下载签名 provider 已接入;上传后校验、PDF/图片预览和资源复检 worker 已完成,CDN 防盗链、杀毒扫描和视频动态水印还没完成。 - Excel/CSV 导入解析已完成并复用 `content_import_jobs/items/issues` 管线;大批量异步导入 worker 基础已接入,支持 queued job 消费、重试和审计;导入后复检、模板下载和字段映射 API 已完成,前端 UI 待接。 - 题库导出目前完成服务端结构化 payload;PDF/Word 二进制生成、导出水印、发布到资料下载和导出 worker 还没完成。 -- 勋章管理/手动发放已可联调;自动发放规则、积分活动联动、分佣真实打款、结算导出、发票/凭证、CRM 轮询/定向分配、富卡片模板、失败告警、销售转化看板、公共题库自动同步 worker 和冲突操作台还没完成。 +- 勋章管理/手动发放已可联调;自动发放规则、积分活动联动、分佣真实打款、结算导出、发票/凭证、CRM 轮询/定向分配、富卡片模板、失败告警、销售转化看板、公共题库版本通知和冲突处理操作台还没完成。 - Taro 跨端前端还没开始 scaffold。 - 根目录已清理为新 Supabase SaaS monorepo 编排层;旧 PocketBase/React 项目和旧构建产物仅保留在 `参考/` 目录作为迁移参考,不进入 Git 提交。 @@ -59,7 +59,7 @@ ```text apps/api/ Node.js 业务 API -apps/worker/ 后台异步任务:CRM webhook、支付/退款补偿、资源复检、后续导入复检等 +apps/worker/ 后台异步任务:CRM webhook、支付/退款补偿、资源复检、导入执行、公共题库同步等 packages/config/ 共享配置 packages/db/ PostgreSQL 连接池和查询封装 packages/domain/ 领域常量和共享类型 @@ -116,6 +116,12 @@ npm --workspace @tiku-saas/worker run assets:once npm --workspace @tiku-saas/worker run imports:once ``` +单次运行公共题库自动同步 worker: + +```bash +npm --workspace @tiku-saas/worker run public-banks:once +``` + 默认本地数据库: ```text @@ -157,6 +163,7 @@ npm run test:worker:crm npm run test:worker:commerce npm run test:worker:assets npm run test:worker:imports +npm run test:worker:public-banks ``` ## API 模块 @@ -205,10 +212,18 @@ API 身份上下文: 最近本地验证命令: ```text +npx supabase db reset npm run check:refactor +npm run test:worker:imports +npm run test:worker:crm +npm run test:worker:commerce +npm run test:worker:assets +npm run test:worker:public-banks +npm audit --audit-level=high +git diff --check ``` -结果:通过。 +结果:通过。`npm audit --audit-level=high` 无高危漏洞;当前依赖树仍有 `exceljs -> uuid` 的 moderate 级提示,修复需要破坏性降级 `exceljs`,后续应在导入 Excel 回归充分后单独处理。 ## 下一步建议 @@ -218,4 +233,4 @@ npm run check:refactor 2. Taro 前端 scaffold,让 H5 和小程序共用同一套 API。 3. 对象存储 CDN 防盗链、杀毒扫描、视频动态水印和生命周期策略。 4. 题库导出 PDF/Word worker、真实数据 dry-run、导入字段映射 UI 和复检结果操作台。 -5. 微信网页/QQ 登录、完整资金流水对账、公共题库自动同步 worker/冲突操作台、积分活动深化,以及排行榜防刷/预聚合。 +5. 微信网页/QQ 登录、完整资金流水对账、公共题库版本通知/冲突处理操作台、积分活动深化,以及排行榜防刷/预聚合。 diff --git a/apps/api/src/features/tenant-content/index.ts b/apps/api/src/features/tenant-content/index.ts index daa96a5a..54301874 100644 --- a/apps/api/src/features/tenant-content/index.ts +++ b/apps/api/src/features/tenant-content/index.ts @@ -46,6 +46,7 @@ import { } from './navigation.js'; import { adoptPublicQuestionBankRoute, + publicQuestionBankConflictsRoute, publicQuestionBanksRoute, syncPublicQuestionBankRoute, } from './public-banks.js'; @@ -80,6 +81,7 @@ export const tenantContentRoutes: RouteDefinition[] = [ ['GET', '/api/tenant-content/public-question-banks', publicQuestionBanksRoute], ['POST', '/api/tenant-content/public-question-banks/adopt', adoptPublicQuestionBankRoute], ['POST', '/api/tenant-content/public-question-banks/sync', syncPublicQuestionBankRoute], + ['GET', '/api/tenant-content/public-question-banks/conflicts', publicQuestionBankConflictsRoute], ['PUT', '/api/tenant-content/content-entries', upsertContentEntryRoute], ['GET', '/api/tenant-content/content-nodes', contentNodesAdminRoute], ['PUT', '/api/tenant-content/content-nodes', upsertContentNodeRoute], diff --git a/apps/api/src/features/tenant-content/public-banks.ts b/apps/api/src/features/tenant-content/public-banks.ts index be013ed4..382c8688 100644 --- a/apps/api/src/features/tenant-content/public-banks.ts +++ b/apps/api/src/features/tenant-content/public-banks.ts @@ -72,6 +72,23 @@ interface SyncQuestionResult { targetHash: string | null; } +interface PublicQuestionBankSyncAuth { + tenantId: string; + userId: string | null; + role: string; + permissions: Record; + templatePermissions: Record; +} + +export interface PublicQuestionBankSyncInput { + tenantId: string; + adoptionId: string; + actorUserId?: string | null; + copyLimit?: number; + triggeredBy?: 'manual' | 'worker'; + workerId?: string | null; +} + function slugFromName(name: string) { const ascii = name .normalize('NFKD') @@ -311,7 +328,7 @@ async function sourceQuestionSnapshots(client: pg.PoolClient, input: { } async function syncQuestionSnapshot(client: pg.PoolClient, input: { - auth: TenantContentAuth; + auth: PublicQuestionBankSyncAuth; sourceTenantId: string; sourceQuestionBankId: string; targetQuestionBankId: string; @@ -508,6 +525,10 @@ function syncCounts(results: SyncQuestionResult[]) { }; } +function objectValue(value: unknown): Record { + return value && typeof value === 'object' && !Array.isArray(value) ? value as Record : {}; +} + function buildSourceSnapshot(input: { sourceTenantId: string; sourceQuestionBankId: string; @@ -539,7 +560,7 @@ function buildSourceSnapshot(input: { } async function syncQuestionsSnapshot(client: pg.PoolClient, input: { - auth: TenantContentAuth; + auth: PublicQuestionBankSyncAuth; sourceTenantId: string; sourceQuestionBankId: string; targetQuestionBankId: string; @@ -873,13 +894,17 @@ export async function adoptPublicQuestionBankRoute(ctx: RequestContext) { return { item }; } -export async function syncPublicQuestionBankRoute(ctx: RequestContext) { - const auth = await requireTenantContentEditor(ctx); - const body = await readJsonBody(ctx); - const adoptionId = requiredString(body, 'adoptionId'); - const copyLimit = Math.max(1, Math.min(intValue(body.copyLimit, 1000), 1000)); +export async function executePublicQuestionBankSync(input: PublicQuestionBankSyncInput) { + const copyLimit = Math.max(1, Math.min(intValue(input.copyLimit, 1000), 1000)); + const auth: PublicQuestionBankSyncAuth = { + tenantId: input.tenantId, + userId: input.actorUserId || null, + role: input.triggeredBy === 'worker' ? 'system_worker' : 'tenant_content_editor', + permissions: { 'content:*': true }, + templatePermissions: {}, + }; - const result = await transaction(async client => { + return transaction(async client => { const adoptionResult = await client.query( ` select id, tenant_id as "tenantId", @@ -900,7 +925,7 @@ export async function syncPublicQuestionBankRoute(ctx: RequestContext) { limit 1 for update `, - [auth.tenantId, adoptionId], + [auth.tenantId, input.adoptionId], ); const adoption = adoptionResult.rows[0]; if (!adoption) { @@ -947,6 +972,21 @@ export async function syncPublicQuestionBankRoute(ctx: RequestContext) { counts: syncResult.counts, conflictCount: conflicts.length, conflicts: conflicts.slice(0, 50), + triggeredBy: input.triggeredBy || 'manual', + workerId: input.workerId || null, + finishedAt: new Date().toISOString(), + }, + publicBankSyncWorker: { + ...( + adoption.metadata?.publicBankSyncWorker + && typeof adoption.metadata.publicBankSyncWorker === 'object' + && !Array.isArray(adoption.metadata.publicBankSyncWorker) + ? adoption.metadata.publicBankSyncWorker as Record + : {} + ), + lastStatus: syncStatus, + lastWorkerId: input.workerId || null, + lastFinishedAt: new Date().toISOString(), }, }; @@ -1003,6 +1043,8 @@ export async function syncPublicQuestionBankRoute(ctx: RequestContext) { syncStatus, syncSummary: syncResult.counts, conflicts, + triggeredBy: input.triggeredBy || 'manual', + workerId: input.workerId || null, }), ], ); @@ -1016,6 +1058,59 @@ export async function syncPublicQuestionBankRoute(ctx: RequestContext) { }, }; }); - - return result; +} + +export async function syncPublicQuestionBankRoute(ctx: RequestContext) { + const auth = await requireTenantContentEditor(ctx); + const body = await readJsonBody(ctx); + return executePublicQuestionBankSync({ + tenantId: auth.tenantId, + actorUserId: auth.userId, + adoptionId: requiredString(body, 'adoptionId'), + copyLimit: intValue(body.copyLimit, 1000), + triggeredBy: 'manual', + }); +} + +export async function publicQuestionBankConflictsRoute(ctx: RequestContext) { + const auth = await requireTenantContentEditor(ctx); + const adoptionId = stringParam(ctx, 'adoptionId'); + if (!adoptionId) { + throw new HttpError(400, 'adoptionId is required', 'REQUIRED_FIELD'); + } + + const rows = await query<{ + id: string; + syncStatus: string; + metadata: Record; + lastSyncedAt: string | null; + updatedAt: string; + }>( + ` + select id, sync_status as "syncStatus", metadata, + last_synced_at as "lastSyncedAt", updated_at as "updatedAt" + from public.tenant_question_bank_adoptions + where tenant_id = $1 and id = $2 and status <> 'archived' + limit 1 + `, + [auth.tenantId, adoptionId], + ); + const adoption = rows[0]; + if (!adoption) { + throw new HttpError(404, 'Question bank adoption not found', 'QUESTION_BANK_ADOPTION_NOT_FOUND'); + } + + const lastSync = objectValue(adoption.metadata?.lastSync); + const conflicts = Array.isArray(lastSync.conflicts) ? lastSync.conflicts : []; + return { + item: { + adoptionId: adoption.id, + syncStatus: adoption.syncStatus, + lastSyncedAt: adoption.lastSyncedAt, + updatedAt: adoption.updatedAt, + conflictCount: Number(lastSync.conflictCount || conflicts.length || 0), + counts: objectValue(lastSync.counts), + conflicts, + }, + }; } diff --git a/apps/worker/package.json b/apps/worker/package.json index 38f5fc83..0e6263af 100644 --- a/apps/worker/package.json +++ b/apps/worker/package.json @@ -11,7 +11,8 @@ "crm:once": "tsx src/index.ts --once --job crm", "commerce:once": "tsx src/index.ts --once --job commerce", "assets:once": "tsx src/index.ts --once --job assets", - "imports:once": "tsx src/index.ts --once --job imports" + "imports:once": "tsx src/index.ts --once --job imports", + "public-banks:once": "tsx src/index.ts --once --job public-banks" }, "dependencies": { "@supabase/storage-js": "^2.108.2", diff --git a/apps/worker/src/config.ts b/apps/worker/src/config.ts index ae5d883a..6f801a85 100644 --- a/apps/worker/src/config.ts +++ b/apps/worker/src/config.ts @@ -20,6 +20,10 @@ export interface WorkerConfig { importBatchSize: number; importWorkerId: string; importBackoffSeconds: number[]; + publicBankSyncBatchSize: number; + publicBankSyncCopyLimit: number; + publicBankSyncWorkerId: string; + publicBankSyncClaimStaleSeconds: number; storageMaxUploadBytes: number; storageAllowedMimePrefixes: string[]; storageAllowedMimeTypes: string[]; @@ -61,6 +65,10 @@ export const config: WorkerConfig = { importBackoffSeconds: envList('WORKER_IMPORT_BACKOFF_SECONDS', '30,120,600,1800') .map((value: string) => Number(value)) .filter((value: number) => Number.isFinite(value) && value > 0), + publicBankSyncBatchSize: envNumber('WORKER_PUBLIC_BANK_SYNC_BATCH_SIZE', 5), + publicBankSyncCopyLimit: envNumber('WORKER_PUBLIC_BANK_SYNC_COPY_LIMIT', 1000), + publicBankSyncWorkerId: envString('WORKER_PUBLIC_BANK_SYNC_ID', `public-banks-${process.pid}`), + publicBankSyncClaimStaleSeconds: envNumber('WORKER_PUBLIC_BANK_SYNC_CLAIM_STALE_SECONDS', 15 * 60), storageMaxUploadBytes: envNumber('STORAGE_MAX_UPLOAD_BYTES', 1024 * 1024 * 500), storageAllowedMimePrefixes: envList('STORAGE_ALLOWED_MIME_PREFIXES', 'image/,video/,audio/'), storageAllowedMimeTypes: envList( diff --git a/apps/worker/src/index.ts b/apps/worker/src/index.ts index bf6ced53..358ac566 100644 --- a/apps/worker/src/index.ts +++ b/apps/worker/src/index.ts @@ -50,6 +50,17 @@ async function runOnce() { ); return; } + if (job === 'public-banks') { + const { closePublicBankSyncExecutorPool, processPublicBankSyncBatch } = await import('./jobs/public-banks.js'); + extraClosers.add(closePublicBankSyncExecutorPool); + const result = await processPublicBankSyncBatch(); + console.log( + `[worker] public-banks batch processed=${result.processed}` + + ` synced=${result.synced} conflicts=${result.conflicts}` + + ` failed=${result.failed} skipped=${result.skipped}`, + ); + return; + } throw new Error(`Unsupported worker job: ${job}`); } diff --git a/apps/worker/src/jobs/public-banks.ts b/apps/worker/src/jobs/public-banks.ts new file mode 100644 index 00000000..61698b93 --- /dev/null +++ b/apps/worker/src/jobs/public-banks.ts @@ -0,0 +1,197 @@ +import crypto from 'node:crypto'; +import { pool } from '../db.js'; +import { config } from '../config.js'; +import { + executePublicQuestionBankSync, +} from '../../../api/src/features/tenant-content/public-banks.js'; +import { closePool as closeApiPublicBankPool } from '../../../api/src/core/db.js'; + +interface PublicBankSyncCandidate { + id: string; + tenantId: string; + createdBy: string | null; + updatedBy: string | null; +} + +interface PublicBankSyncWorkerResult { + processed: number; + synced: number; + conflicts: number; + failed: number; + skipped: number; +} + +function nowIso() { + return new Date().toISOString(); +} + +function errorMessage(error: unknown) { + return error instanceof Error ? error.message : String(error); +} + +function errorCode(error: unknown) { + return typeof error === 'object' && error !== null && 'code' in error + ? String((error as { code?: unknown }).code || 'PUBLIC_BANK_SYNC_WORKER_ERROR') + : 'PUBLIC_BANK_SYNC_WORKER_ERROR'; +} + +function truncate(value: unknown, max = 1900) { + return String(value ?? '').slice(0, max); +} + +async function claimPublicBankSyncCandidates(limit: number, claimId: string) { + const client = await pool.connect(); + try { + await client.query('begin'); + const result = await client.query( + ` + with candidates as ( + select a.id + from public.tenant_question_bank_adoptions a + join public.question_banks qb on qb.id = a.source_question_bank_id + left join public.question_bank_grants g on g.id = a.grant_id + where a.status in ('active', 'sync_pending') + and (a.sync_status <> 'failed' or a.status = 'sync_pending') + and a.grant_id is not null + and a.target_question_bank_id is not null + and a.target_entry_id is not null + and a.target_collection_id is not null + and qb.source_scope = 'platform' + and qb.status = 'active' + and ( + a.sync_status = 'pending' + or a.status = 'sync_pending' + or a.last_synced_at is null + or qb.updated_at > coalesce(a.last_synced_at, '1970-01-01'::timestamptz) + or coalesce(g.updated_at, '1970-01-01'::timestamptz) > coalesce(a.last_synced_at, '1970-01-01'::timestamptz) + or exists ( + select 1 + from public.questions q + where q.tenant_id = qb.tenant_id + and q.question_bank_id = qb.id + and q.status = 'published' + and q.updated_at > coalesce(a.last_synced_at, '1970-01-01'::timestamptz) + ) + ) + and ( + a.status <> 'sync_pending' + or coalesce(nullif(a.metadata #>> '{publicBankSyncWorker,claimedAt}', '')::timestamptz, '1970-01-01'::timestamptz) + <= now() - ($2::integer * interval '1 second') + ) + order by coalesce(a.last_synced_at, '1970-01-01'::timestamptz) asc, a.updated_at asc + limit $1 + for update of a skip locked + ) + update public.tenant_question_bank_adoptions a + set status = 'sync_pending', + sync_status = 'pending', + metadata = jsonb_set( + coalesce(a.metadata, '{}'::jsonb), + '{publicBankSyncWorker}', + coalesce(a.metadata->'publicBankSyncWorker', '{}'::jsonb) || $3::jsonb, + true + ), + updated_at = now() + from candidates + where a.id = candidates.id + returning a.id, + a.tenant_id as "tenantId", + a.created_by as "createdBy", + a.updated_by as "updatedBy" + `, + [ + limit, + Math.max(60, config.publicBankSyncClaimStaleSeconds), + JSON.stringify({ + claimId, + workerId: config.publicBankSyncWorkerId, + claimedAt: nowIso(), + }), + ], + ); + await client.query('commit'); + return result.rows; + } catch (error) { + await client.query('rollback'); + throw error; + } finally { + client.release(); + } +} + +async function markPublicBankSyncFailed(candidate: PublicBankSyncCandidate, error: unknown) { + const details = { + code: errorCode(error), + message: truncate(errorMessage(error)), + workerId: config.publicBankSyncWorkerId, + failedAt: nowIso(), + }; + await pool.query( + ` + update public.tenant_question_bank_adoptions + set status = 'active', + sync_status = 'failed', + metadata = jsonb_set( + coalesce(metadata, '{}'::jsonb), + '{publicBankSyncWorker}', + coalesce(metadata->'publicBankSyncWorker', '{}'::jsonb) || $3::jsonb, + true + ), + updated_at = now() + where tenant_id = $1 and id = $2 + `, + [ + candidate.tenantId, + candidate.id, + JSON.stringify({ + lastStatus: 'failed', + lastError: details, + lastWorkerId: config.publicBankSyncWorkerId, + lastFinishedAt: nowIso(), + }), + ], + ); + await pool.query( + ` + insert into public.audit_logs (tenant_id, actor_user_id, action, target_type, target_id, details) + values ($1, null, 'content.public_question_bank.sync_worker_failed', 'tenant_question_bank_adoption', $2, $3::jsonb) + `, + [candidate.tenantId, candidate.id, JSON.stringify(details)], + ); +} + +export async function processPublicBankSyncBatch(limit = config.publicBankSyncBatchSize): Promise { + const candidates = await claimPublicBankSyncCandidates(limit, crypto.randomUUID()); + const result: PublicBankSyncWorkerResult = { + processed: candidates.length, + synced: 0, + conflicts: 0, + failed: 0, + skipped: 0, + }; + + for (const candidate of candidates) { + try { + const sync = await executePublicQuestionBankSync({ + tenantId: candidate.tenantId, + adoptionId: candidate.id, + actorUserId: candidate.updatedBy || candidate.createdBy || null, + copyLimit: config.publicBankSyncCopyLimit, + triggeredBy: 'worker', + workerId: config.publicBankSyncWorkerId, + }); + if (sync.sync.status === 'conflict') result.conflicts += 1; + else if (sync.sync.counts.inserted || sync.sync.counts.updated || sync.sync.counts.skipped) result.synced += 1; + else result.skipped += 1; + } catch (error) { + await markPublicBankSyncFailed(candidate, error); + result.failed += 1; + } + } + + return result; +} + +export async function closePublicBankSyncExecutorPool() { + await closeApiPublicBankPool(); +} diff --git a/docs/refactor/architecture.md b/docs/refactor/architecture.md index 3a8ca9a6..cb29fd73 100644 --- a/docs/refactor/architecture.md +++ b/docs/refactor/architecture.md @@ -94,5 +94,5 @@ npm run pb:import:validate - `apps/api/src/features` 继续按业务域扩展:退款对账、真实 OAuth provider、平台审计和更多后台任务。 - `src/services/supabaseApi.ts` 逐页替换旧 PB 只读接口,优先学生端和小程序共用页面。 -- 扩展 `apps/worker`:CRM webhook 已落地,后续继续承接支付补偿、日报统计、导入后异步检查和公共题库同步。 +- 扩展 `apps/worker`:CRM webhook、支付/退款补偿、资源复检、异步导入和公共题库同步已落地;后续继续补日报统计、失败告警、版本通知和冲突处理运营台。 - 新增 `apps/taro` 后,Auth/JWT 优先复用 Supabase client;复杂业务命令复用 `apps/api`/RPC/Edge Functions,不单独维护另一套后端逻辑。 diff --git a/docs/refactor/backend-capability-status.md b/docs/refactor/backend-capability-status.md index 21f8de7c..1924ffea 100644 --- a/docs/refactor/backend-capability-status.md +++ b/docs/refactor/backend-capability-status.md @@ -135,7 +135,7 @@ | 平台租户/套餐/订阅/账单/用量 | 可联调 | `/api/platform-admin/*` | | 数据看板聚合接口 | 可联调 | `GET /api/tenant-admin/dashboard`;支持 `7d/30d/90d`、地区筛选、学生/学习/内容/订单/激活码/反馈卡片、趋势、24h 活跃、题型分布、科目排行、地区统计、套餐销量和运营动态 | | 平台公共题库授权 | 可联调 | `/api/platform-admin/question-banks`、`question-bank-grants`;支持按 SaaS 套餐、指定租户或全部活跃租户披露平台公共题库 | -| 租户采纳/同步公共题库 | 可联调 | `/api/tenant-content/public-question-banks`、`public-question-banks/adopt`、`public-question-banks/sync`;租户只能看到自己订阅/授权范围内题库,采纳后生成租户自己的题库、入口、集合和题目快照,可直接进入练习;平台更新后可手动同步,租户自改题目会标记冲突并跳过 | +| 租户采纳/同步公共题库 | 可联调 | `/api/tenant-content/public-question-banks`、`public-question-banks/adopt`、`public-question-banks/sync`、`public-question-banks/conflicts`;租户只能看到自己订阅/授权范围内题库,采纳后生成租户自己的题库、入口、集合和题目快照,可直接进入练习;平台更新后可手动或由 worker 自动同步,租户自改题目会标记冲突并跳过,后台可查询最近一次冲突明细 | | 题库导出基础 | 可联调 | `/api/tenant-content/exports/questions`、`/api/tenant-content/exports/jobs`;支持按题目集合、内容入口或分类节点导出 JSON/试卷 payload,后端校验租户内容编辑权限、跨租户隔离、答案/解析开关、复合题子题脱敏、导出 job 和审计;PDF/Word 二进制与水印 worker 后续补 | ## 销售、代理、CRM @@ -159,7 +159,7 @@ | --- | --- | --- | | PocketBase schema 分析 | 可联调 | `scripts/import-pocketbase` | | 题目 JSON preview/import | 可联调 | 后端负责规范化、issue、幂等、审计 | -| 公共题库采纳与手动同步 | 可联调 | 平台授权后,租户可采纳公共题库并复制已发布题目快照;同步 API 支持新增/更新题目、重新校验授权、跨租户拒绝、审计记录和租户自改冲突保护;已覆盖跨租户、重复采纳、采纳后组卷、同步新增题和冲突不覆盖测试 | +| 公共题库采纳、手动同步和自动同步 | 可联调 | 平台授权后,租户可采纳公共题库并复制已发布题目快照;同步 API 和 `public-banks` worker 支持新增/更新题目、重新校验授权、跨租户拒绝、审计记录和租户自改冲突保护;已覆盖跨租户、重复采纳、采纳后组卷、同步新增题、冲突不覆盖和 worker 自动同步测试 | | 单词 JSON preview/import | 可联调 | 兼容旧模板 | | 知识手册 JSON preview/import | 可联调 | 支持书籍/章节/小节/知识点归一化 | | 分数线 JSON preview/import | 可联调 | 支持 `fields/schools/majors/records` 分桶或 `items` 列表,后端校验租户地区和院校/专业引用 | @@ -167,16 +167,22 @@ | Excel/CSV 导入 | 可联调 | 题目、单词、知识手册、分数线、视频已支持 CSV 和 `.xlsx` 解析,解析后复用 `content_import_jobs/items/issues` 管线并保留 `parser_metadata`;模板下载、字段映射 API 和导入后复检已接入 | | 大批量异步导入 | 可联调 | `executionMode=async` 会将 preview job 置为 `pending`;`apps/worker --job imports` 抢占 queued job,复用 API 导入 executor,支持重试、清锁和审计 | | 题库导出任务 | 可联调 | `content_export_jobs` 记录导出范围、格式、题量、输出 hash、选项和执行人;当前返回 inline base64 JSON 文件,前端可先下载 `.json` 或交给后续 PDF/Word worker 渲染 | -| 公共题库自动同步增强 | 待补齐 | 手动同步 API 已完成;后续需 worker 做定时同步、失败重试、版本升级通知、冲突操作台和批量确认/跳过 | +| 公共题库自动同步增强 | 部分覆盖 | `apps/worker --job public-banks` 已可抢占待同步采纳记录、自动同步平台新增/更新题目、记录失败和审计;后续需接入生产定时调度、版本升级通知、冲突操作台和批量确认/跳过 | ## 当前验证 最近需通过: ```bash -npm audit +npx supabase db reset +npm audit --audit-level=high npm run check:refactor npm run test:worker:imports +npm run test:worker:crm +npm run test:worker:commerce +npm run test:worker:assets +npm run test:worker:public-banks +git diff --check ``` `check:refactor` 包含: diff --git a/docs/refactor/backend-handoff-roadmap.md b/docs/refactor/backend-handoff-roadmap.md index 2d8350b8..e8299a11 100644 --- a/docs/refactor/backend-handoff-roadmap.md +++ b/docs/refactor/backend-handoff-roadmap.md @@ -21,9 +21,9 @@ | 模块 | 当前状态 | 已经具备 | 上线前还要补 | | --- | --- | --- | --- | | 多租户底座 | 可联调 | 租户、域名、品牌、设置、RLS 基础、审计、Supabase JWT/API 身份映射 | 真实云端 Auth/JWKS 回归、生产 RLS 深测 | -| 平台后台 | 基础完成 | 租户、套餐、订阅、账单、服务费、用量、公共题库授权 | 自动计费、平台审计、公共题库自动同步运营台 | +| 平台后台 | 基础完成 | 租户、套餐、订阅、账单、服务费、用量、公共题库授权、公共题库自动同步 worker | 自动计费、平台审计、公共题库版本通知和冲突处理运营台 | | 租户后台 | 可联调 | 品牌、域名、支付账户、登录配置、密钥掩码、活动、兑换码、优惠券、勋章管理/发放、成员权限、角色模板、菜单/模块/字段权限配置 API、班级/教师/学生范围权限 | 前端权限 UI、更细的数据范围组合 | -| 题库与练习 | 可联调 | 内容入口、任意深度分类、题目集合、顺序/随机/全真模拟蓝图、组卷快照、答题、错题、收藏、模考报告、排行榜、公共题库采纳快照和手动同步、JSON/试卷 payload 导出 | 专项策略、PDF/Word 导出 worker、公共题库自动同步 worker/冲突操作台、排行榜防刷/预聚合 | +| 题库与练习 | 可联调 | 内容入口、任意深度分类、题目集合、顺序/随机/全真模拟蓝图、组卷快照、答题、错题、收藏、模考报告、排行榜、公共题库采纳快照、手动同步、自动同步 worker、冲突查询 API、JSON/试卷 payload 导出 | 专项策略、PDF/Word 导出 worker、公共题库版本通知和冲突操作台、排行榜防刷/预聚合 | | 背单词 | 可联调 | 单元、单词、进度、收藏、统计、每日计划、JSON/CSV/Excel 导入、排行榜 | 更细复习参数 | | 知识手册 | 可联调 | 科目、章节、条目、Markdown 内容、嵌套 JSON/CSV/Excel 导入 | 富文本资源、版本管理、附件/PDF 关联 | | 分数线 | 可联调 | 院校、专业、动态字段、记录、年份、趋势、后台维护、JSON/CSV/Excel 导入 | 复杂筛选、AI 择校上下文 | @@ -86,7 +86,7 @@ - 完整资金流水对账、账单下载比对和异常订单运营台。 - XPay 或其它实际支付网关 adapter。 - 阿里云/腾讯云短信、微信小程序登录、微信网页登录、QQ 登录。 -- 公共题库/地区题库自动同步 worker、版本通知、冲突操作台,以及租户按 SaaS 套餐购买地区、科目和题库范围的更细计费策略。 +- 公共题库/地区题库自动同步 worker 已具备单批执行能力;继续补版本通知、冲突操作台,以及租户按 SaaS 套餐购买地区、科目和题库范围的更细计费策略。 - 导入模板、字段映射和复检 API 已可用;前端继续补模板下载按钮、字段映射 UI、job 状态轮询和复检结果面板。 - 视频深度防盗链、动态水印和播放统计。 - 数据看板 API:收益、注册趋势、答题次数、收入趋势、题型分布、题目总量、套餐销量、24h 活跃。 diff --git a/docs/refactor/backend-progress.md b/docs/refactor/backend-progress.md index 7ee9200d..acc5b814 100644 --- a/docs/refactor/backend-progress.md +++ b/docs/refactor/backend-progress.md @@ -261,7 +261,7 @@ GET /api/tenant-admin/audit-logs 2. 接入真实短信 provider:阿里云/腾讯云,密钥放 `app_private.tenant_secrets` 或生产 Vault。 3. 接入真实 OAuth provider:微信网页、微信小程序、QQ,并处理旧 PocketBase 身份映射。 4. 补完整资金流水对账、异常订单运营台和优惠券核销报表;支付/退款补偿、退款查询确认和退款通知主链路已完成。 -5. 扩展 `apps/worker`:日报统计、导入后检查、CRM 死信告警和公共题库同步。 +5. 扩展 `apps/worker`:日报统计、CRM 死信告警、公共题库版本通知和冲突处理运营台;公共题库同步 worker 已具备单批执行能力。 6. 开始 Taro scaffold,把 `supabaseApi` 抽到跨端包或适配层。 ## 测试命令 @@ -272,5 +272,6 @@ npm run test:worker:crm npm run test:worker:commerce npm run test:worker:assets npm run test:worker:imports +npm run test:worker:public-banks npm run check:refactor ``` diff --git a/docs/refactor/blueprint-coverage.md b/docs/refactor/blueprint-coverage.md index b3c175b9..fc8ecb84 100644 --- a/docs/refactor/blueprint-coverage.md +++ b/docs/refactor/blueprint-coverage.md @@ -18,7 +18,7 @@ | 平台超级管理员 | 部分完成 | 租户管理、SaaS 套餐、订阅、账单、服务费收款、用量记录 | 公共题库披露策略、地区/全国套餐权限、平台侧主题模板库、平台审计 | | 租户品牌和域名 | 基础完成 | 品牌、Logo、主题 JSON、公开资源、域名、租户公开配置 | 三套默认主题、主题可视化编辑、图标/图片上传 | | 租户成员权限 | 可联调 | owner/admin/operator/teacher/sales/agent/student,权限矩阵,成员启停,角色模板、菜单/模块/字段权限、班级/学生范围权限和审计查询 | 前端权限 UI、更细的数据范围组合 | -| 题库内容维护 | 可联调 | 内容入口、任意深度分类树、院校/专业/学科/销售意向标记、题目集合、顺序/随机/全真模拟练习蓝图、题目录入/更新、题目/单词/知识手册/分数线/视频 JSON/CSV/Excel 预览导入、`executionMode=async` 导入 worker、导入后复检、模板/字段映射 API、视频绑定、分数线、单词、知识手册后台 API、公共题库授权、采纳快照和手动同步 | 字段映射 UI、公共题库自动同步 worker/冲突操作台、可视化拖拽排序前端 | +| 题库内容维护 | 可联调 | 内容入口、任意深度分类树、院校/专业/学科/销售意向标记、题目集合、顺序/随机/全真模拟练习蓝图、题目录入/更新、题目/单词/知识手册/分数线/视频 JSON/CSV/Excel 预览导入、`executionMode=async` 导入 worker、导入后复检、模板/字段映射 API、视频绑定、分数线、单词、知识手册后台 API、公共题库授权、采纳快照、手动同步、自动同步 worker 和冲突查询 API | 字段映射 UI、公共题库版本通知/冲突操作台、可视化拖拽排序前端 | | 学生刷题 | 基础完成 | 内容入口、分类树、题目集合、顺序刷题、随机刷题、全真模拟 session 题目快照、答题、错题本、收藏夹、模考交卷评分报告、错题复习计划、排行榜 | 专项练习策略、题型统计深度分析、排行榜防刷/预聚合 | | 背单词 | 基础完成 | 单词单元、单词、进度、收藏、统计、每日复习计划、旧模板/新模板 JSON/CSV/Excel 预览导入、内容导航绑定、排行榜 | 更细复习参数 | | 知识手册 | 基础完成 | 科目、章节、条目只读与后台维护、书籍/章节/小节/知识点嵌套 JSON 预览导入、内容导航绑定 | 富文本资源、版本管理、附件/PDF 关联、Excel/Markdown 批量解析 | @@ -37,7 +37,7 @@ ## 接下来优先级 1. 完善内容导入和对象存储:字段映射 UI、真实数据 dry-run、CDN 防盗链、杀毒扫描和视频水印。 -2. 公共题库/地区题库授权:已完成披露、采纳快照和手动同步冲突保护;继续补自动同步 worker、版本通知、冲突操作台和按 SaaS 套餐限制地区。 +2. 公共题库/地区题库授权:已完成披露、采纳快照、手动同步、自动同步 worker、冲突查询和租户自改保护;继续补版本通知、冲突操作台和按 SaaS 套餐限制地区。 3. 学习统计增强:排行榜防刷/预聚合、断点续练、专项练习策略和更细题型分析。 4. 视频会员控制:深度防盗链、水印和播放统计。 5. 数据看板预聚合:把实时聚合升级为大租户可承载的日/周/月预聚合。 diff --git a/docs/refactor/implementation-status.md b/docs/refactor/implementation-status.md index 77180034..6bb0860a 100644 --- a/docs/refactor/implementation-status.md +++ b/docs/refactor/implementation-status.md @@ -24,7 +24,7 @@ | 模块 | 数据模型 | PocketBase 导入 | API | 自动化测试 | 当前状态 | | --- | --- | --- | --- | --- | --- | | 多租户隔离 | 已建 `tenants`、`tenant_domains`、`tenant_branding`、`tenant_settings`、RLS 基础 | 部分支持 | 租户解析、品牌、域名、支付账户、登录 provider、平台建租户已实现 | 核心 API 集成测试含租户隔离断言 | 基础可用,正式 JWT/RLS 权限闭环未完成 | -| 刷题题库 | 已建题库、题目、题目版本、内容入口、任意深度分类树、考试意向标记、题目集合、练习蓝图、导入任务台账、导出任务台账、公共题库授权/采纳表 | 已支持核心映射,JSON/CSV/Excel 导入可落到新入口/节点/集合 | 题目列表、内容入口、分类树、集合题目、顺序/随机/全真模拟 session、答题提交、租户后台题目录入/更新、JSON/CSV/Excel 预览/导入、JSON/试卷 payload 导出、异步导入 worker、平台公共题库授权、租户采纳快照和手动同步已实现 | 核心 API 集成测试含导航、组卷、导入、导出权限/脱敏、公共题库授权、采纳后组卷、同步新增题和租户自改冲突保护断言 | 新题库导航和组卷基础闭环可跑,公共题库采纳/手动同步、导入后复检、模板下载、字段映射 API 和导出基础可联调;PDF/Word 导出 worker、公共题库自动同步 worker、版本通知和冲突操作台仍需补齐 | +| 刷题题库 | 已建题库、题目、题目版本、内容入口、任意深度分类树、考试意向标记、题目集合、练习蓝图、导入任务台账、导出任务台账、公共题库授权/采纳表 | 已支持核心映射,JSON/CSV/Excel 导入可落到新入口/节点/集合 | 题目列表、内容入口、分类树、集合题目、顺序/随机/全真模拟 session、答题提交、租户后台题目录入/更新、JSON/CSV/Excel 预览/导入、JSON/试卷 payload 导出、异步导入 worker、平台公共题库授权、租户采纳快照、手动同步、自动同步 worker 和冲突查询已实现 | 核心 API 集成测试含导航、组卷、导入、导出权限/脱敏、公共题库授权、采纳后组卷、同步新增题、租户自改冲突保护和 worker 自动同步断言 | 新题库导航和组卷基础闭环可跑,公共题库采纳/手动/自动同步、冲突查询、导入后复检、模板下载、字段映射 API 和导出基础可联调;PDF/Word 导出 worker、公共题库版本通知和冲突操作台仍需补齐 | | 错题本 | 已建 `wrong_questions` | 已支持旧错题归一化 | 错题列表、答题自动入错题、移出错题已实现 | 仅烟测 | 基础功能已实现,复习计划和统计未完成 | | 收藏夹 | 已建 `favorite_questions` | 已支持旧收藏归一化 | 收藏/取消收藏、收藏列表已实现 | 仅烟测 | 基础功能已实现 | | 用户订阅/题库会员/SVIP | 已建 `orders`、`payments`、`entitlements`、`svip_plans`、激活码 | 已映射旧 SVIP/会员权益 | 下单、订单详情/状态轮询、手工支付确认权限保护、微信/支付宝支付、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、激活码预检查/兑换、优惠券抵扣、零元订单自动开通、权益查询已实现 | API 集成测试 | 商城主链路可联调,对账、支付补偿和异常订单自动处理待补 | @@ -126,6 +126,10 @@ video: tenant-content: GET /api/tenant-content/content-entries PUT /api/tenant-content/content-entries + GET /api/tenant-content/public-question-banks + POST /api/tenant-content/public-question-banks/adopt + POST /api/tenant-content/public-question-banks/sync + GET /api/tenant-content/public-question-banks/conflicts GET /api/tenant-content/content-nodes PUT /api/tenant-content/content-nodes GET /api/tenant-content/question-collections @@ -280,7 +284,7 @@ platform-admin: 为了先把旧项目核心业务补齐,再进入支付/短信等商用关键模块,建议按下面顺序继续: 1. 完善内容导入和文件上传:字段映射 UI、真实数据 dry-run、CDN 防盗链、杀毒扫描。 -2. 补公共题库自动同步 worker、版本通知/冲突操作台、租户套餐地区/科目/题库范围限制、主题模板系统。 +2. 补公共题库版本通知/冲突操作台、租户套餐地区/科目/题库范围限制、主题模板系统。 3. 补学习统计增强:排行榜防刷/预聚合、断点续练、专项练习策略和更细题型分析。 4. 补视频商用控制:深度防盗链、动态水印和播放统计。 5. 补 AI 择校推荐报告、排行榜防刷/预聚合、勋章自动发放。 diff --git a/docs/refactor/legacy-feature-gap-matrix.md b/docs/refactor/legacy-feature-gap-matrix.md index 3e39edf2..b7b8af7a 100644 --- a/docs/refactor/legacy-feature-gap-matrix.md +++ b/docs/refactor/legacy-feature-gap-matrix.md @@ -76,7 +76,7 @@ | SaaS 套餐 | 部分覆盖 | 已和公共题库授权打通;后续继续补地区数量、科目范围、存储/学生数等组合套餐限制 | | 年费/服务费账单 | 已覆盖 | 真实支付/开票/催缴流程待补 | | 租户用量记录 | 已覆盖 | 自动采集 worker 待补 | -| 公共题库/地区题库 | 部分覆盖 | 已有平台公共题库列表、授权、租户可采纳列表、采纳快照复制、采纳后练习组卷,以及手动同步 API;同步会重新校验授权、复制平台新增/更新题目,并对租户自改题目返回冲突不覆盖 | 缺自动同步 worker、版本通知、冲突操作台和更完整运营 UI | +| 公共题库/地区题库 | 部分覆盖 | 已有平台公共题库列表、授权、租户可采纳列表、采纳快照复制、采纳后练习组卷、手动同步 API、自动同步 worker 和冲突查询 API;同步会重新校验授权、复制平台新增/更新题目,并对租户自改题目返回冲突不覆盖 | 缺版本通知、冲突处理操作台、批量接受/保留策略和更完整运营 UI | | 跨租户运营看板 | 部分覆盖 | overview 有基础;缺完整 BI 聚合 | | 租户安全审计 | 部分覆盖 | audit logs 有;缺平台级审计报表 | @@ -99,7 +99,7 @@ 2. 账号设置完整流:头像上传、绑定/更换手机号、微信/QQ 账号合并、密码/邮箱能力。 3. 题库导出:服务端 JSON/试卷 payload 导出、权限审计和答案脱敏已补;仍缺 PDF/Word 二进制生成、水印、资料发布和后台导出操作台。 4. 导入扩展:题目/单词/知识手册/分数线/视频已支持 JSON、CSV 和 Excel 预览导入,并可用 `executionMode=async` 进入 imports worker;导入后复检、模板下载和字段映射 API 已补,仍缺前端字段映射 UI 和真实数据 dry-run。 -5. 公共题库商业化:平台公共/地区题库授权、租户快照采纳、手动同步和租户自改冲突保护已完成基础闭环;还需自动同步 worker、版本通知、冲突操作台和运营后台 UI。 +5. 公共题库商业化:平台公共/地区题库授权、租户快照采纳、手动同步、自动同步 worker、冲突查询和租户自改冲突保护已完成基础闭环;还需版本通知、冲突处理操作台和运营后台 UI。 6. CRM/销售结算:CRM worker、分佣规则、结算单、审核和打款状态基础闭环已完成;仍缺轮询/定向分配、打款导出、凭证和销售结算看板。 7. 题目反馈增强:处理通知、消息提醒、问题聚合统计和内容修复闭环。 8. 积分活动增强:积分兑换、活动任务、连续签到奖励规则和风控。 @@ -119,7 +119,7 @@ 2. 对象存储 PDF 预览、视频深度防盗链、动态水印。 3. 字段映射 UI、真实数据 dry-run 和导入复检结果操作台。 4. 数据看板预聚合 worker、销售/代理转化看板和分佣结算。 -5. 公共题库自动同步 worker、版本通知、冲突操作台和租户确认/跳过策略。 +5. 公共题库版本通知、冲突处理操作台和租户确认/跳过策略。 ### P2:增强体验 diff --git a/docs/refactor/next-development-todo.md b/docs/refactor/next-development-todo.md index bc668557..204b6d60 100644 --- a/docs/refactor/next-development-todo.md +++ b/docs/refactor/next-development-todo.md @@ -23,7 +23,7 @@ - 旧题库运营缺口已补一批:考试日期/倒计时、题目反馈/纠错处理、每日签到积分和积分流水、学习排行榜已完成接口和集成测试。 - 勋章管理已完成租户后台维护、手动发放、重复发放幂等、学生个人中心展示、权限点和集成测试;后续补自动发放规则和活动联动。 - 旧商城体验已补齐主链路:订单详情、订单状态轮询、激活码预检查、自用激活码拒绝、优惠券前台领取、下单抵扣、零元订单自动支付开通权益,且手工支付确认已限制为租户后台 `tenant:payment:write` 权限。 -- 公共题库商业化基础闭环已完成:平台公共题库可由平台管理员按 SaaS 套餐/指定租户/全部活跃租户授权;租户内容管理员只能看到自己被授权的公共题库,并可采纳为本租户题库、内容入口、题目集合和题目快照,采纳后可直接进入练习 session;平台题库后续新增/更新题目可通过手动同步 API 进入租户副本,租户自改题目会返回冲突并保留原内容。 +- 公共题库商业化基础闭环已完成:平台公共题库可由平台管理员按 SaaS 套餐/指定租户/全部活跃租户授权;租户内容管理员只能看到自己被授权的公共题库,并可采纳为本租户题库、内容入口、题目集合和题目快照,采纳后可直接进入练习 session;平台题库后续新增/更新题目可通过手动同步 API 或 `public-banks` worker 进入租户副本,租户自改题目会返回冲突并保留原内容,后台可查询最近一次冲突明细。 - 租户后台数据看板已完成首版聚合 API:`GET /api/tenant-admin/dashboard`,支持租户/地区维度的收益、注册、学习、内容、激活码、反馈、趋势、24h 活跃、套餐销量和运营动态,前端可直接联调。 - 支付/退款补偿 worker 已完成:`apps/worker --job commerce` 可查询微信/支付宝支付和处理中退款,补偿漏通知订单,支付成功幂等开通权益,退款成功幂等更新退款/订单/支付并在全额退款时撤销订单权益。 - 内容资源复检 worker 已完成:`apps/worker --job assets` 可复检 `content_assets` 中的托管对象元数据,正常资源写回复检证据,异常资源自动置为 `failed + draft` 并写入审计和安全标记。 @@ -87,8 +87,9 @@ 5. 公共题库和租户授权 - 已完成平台公共题库/地区题库的基础授权、租户采纳、题目快照复制和手动同步。 + - 已完成 `public-banks` worker 自动同步、失败记录、审计和冲突查询 API。 - 继续补按 SaaS 套餐限制地区数量、科目范围、题库范围的更细计费策略。 - - 继续补公共题库自动同步 worker、同步失败重试、版本通知、冲突操作台和运营后台 UI。 + - 继续补生产定时调度、版本通知、冲突处理操作台、批量接受/保留策略和运营后台 UI。 6. 视频会员控制 - 已完成视频 SVIP 权限、播放次数扣减、签名播放和播放日志。 @@ -212,5 +213,5 @@ 2. 云服务器部署 Supabase/PostgreSQL 和 API,配置对象存储生产环境变量,跑 `check:refactor` 的远程等价测试。 3. 导出现有 PocketBase 数据,做完整 dry-run 迁移。 4. 开始 `apps/taro`,先接租户解析、首页、题库、背单词、知识手册。 -5. 并行补对象存储、真实登录、完整资金流水对账、题库导出 PDF/Word worker 和公共题库自动同步 worker/冲突操作台。 +5. 并行补对象存储、真实登录、完整资金流水对账、题库导出 PDF/Word worker 和公共题库版本通知/冲突操作台。 6. 前后端联调通过后,再做支付、权限、数据导入、资料下载、视频播放的商用验收。 diff --git a/docs/refactor/taro-frontend-integration.md b/docs/refactor/taro-frontend-integration.md index dfc6389a..d4dd67ec 100644 --- a/docs/refactor/taro-frontend-integration.md +++ b/docs/refactor/taro-frontend-integration.md @@ -184,7 +184,7 @@ tenant::theme | 租户考试日期 | `GET/PUT /api/tenant-admin/exam-dates` | | 租户反馈处理 | `GET /api/tenant-admin/feedbacks`、`POST /api/tenant-admin/feedbacks/status`、`GET /api/tenant-admin/feedbacks/events` | | 租户勋章 | `GET/PUT /api/tenant-admin/badges`、`GET/POST /api/tenant-admin/badge-grants` | -| 公共题库采纳/同步 | `GET /api/tenant-content/public-question-banks`、`POST /api/tenant-content/public-question-banks/adopt`、`POST /api/tenant-content/public-question-banks/sync` | +| 公共题库采纳/同步 | `GET /api/tenant-content/public-question-banks`、`POST /api/tenant-content/public-question-banks/adopt`、`POST /api/tenant-content/public-question-banks/sync`、`GET /api/tenant-content/public-question-banks/conflicts?adoptionId=...` | | 题库导出 | `POST /api/tenant-content/exports/questions`、`GET /api/tenant-content/exports/jobs` | ## 练习访问控制契约 @@ -951,6 +951,7 @@ PUT /api/platform-admin/question-bank-grants GET /api/tenant-content/public-question-banks POST /api/tenant-content/public-question-banks/adopt POST /api/tenant-content/public-question-banks/sync +GET /api/tenant-content/public-question-banks/conflicts?adoptionId=... ``` 采纳请求: @@ -970,6 +971,7 @@ POST /api/tenant-content/public-question-banks/sync - 采纳成功后后端会生成本租户自己的 `questionBankId`、`entryId`、`collectionId` 和题目快照,学生端直接按普通 `/api/catalog/content-entries`、`question-collections`、`practice-sessions` 接入。 - 重复采纳返回 `QUESTION_BANK_ALREADY_ADOPTED`,前端展示“已采纳”即可。 - 已采纳公共题库可以手动同步平台后续新增/更新题目;同步会重新校验当前租户仍有授权,且只写入租户自己的题目副本。 +- 后端也可以由 `apps/worker --job public-banks` 自动同步待更新的采纳题库;前端不需要轮询平台源库,只需要在租户后台展示同步状态、最近同步时间和冲突数量。 同步请求: @@ -1019,9 +1021,10 @@ POST /api/tenant-content/public-question-banks/sync - `sync.status=synced`:刷新公共题库列表、题目集合和题目列表。 - `sync.status=conflict` 或 `item.syncStatus=failed`:展示冲突数量和冲突题目,不要把它当系统异常。冲突表示租户已经改过这道采纳题,后端已跳过并保留租户内容。 - `action=conflict` 的记录可以进入后续“冲突处理”页面:展示平台源题 ID、租户目标题 ID、上次平台 hash、当前平台 hash、租户当前 hash。当前后端只负责保护不覆盖,批量接受平台版本/保留租户版本的操作台后续补。 +- 页面初始化或 worker 后台同步完成后,可以调用 `GET /api/tenant-content/public-question-banks/conflicts?adoptionId=...` 查询最近一次同步状态、`counts` 和 `conflicts`。这个接口只返回当前租户自己的采纳记录,跨租户会返回 `QUESTION_BANK_ADOPTION_NOT_FOUND`。 - `QUESTION_BANK_GRANT_NOT_AVAILABLE`:说明 SaaS 套餐/授权已失效,提示联系平台或升级套餐。 - `QUESTION_BANK_ADOPTION_NOT_FOUND`:说明不是当前租户的采纳记录或记录已归档,前端不要跨租户重试。 -- 后续会补自动同步 worker、版本通知、冲突操作台和批量确认策略;当前租户后台可以先提供手动“同步平台更新”按钮。 +- 后续会补版本通知、冲突操作台和批量确认策略;当前租户后台可以先提供手动“同步平台更新”按钮,并展示 worker 自动同步后的冲突查询结果。 ## 登录对接 diff --git a/package.json b/package.json index 459cac76..bd564baa 100644 --- a/package.json +++ b/package.json @@ -37,6 +37,7 @@ "test:worker:commerce": "npm run db:smoke-seed && npm run build:worker && node scripts/commerce-worker-integration-test.js", "test:worker:assets": "npm run db:smoke-seed && npm run build:worker && node scripts/asset-worker-integration-test.js", "test:worker:imports": "npm run db:smoke-seed && npm run build:worker && node scripts/import-worker-integration-test.js", + "test:worker:public-banks": "npm run db:smoke-seed && npm run build:worker && node scripts/public-bank-worker-integration-test.js", "test:api:remote": "node scripts/api-integration-test.js", "pb:schema:summary": "npm --workspace @tiku-saas/import-pocketbase run schema:summary", "pb:schema:risk": "npm --workspace @tiku-saas/import-pocketbase run schema:risk", diff --git a/scripts/api-integration-test.js b/scripts/api-integration-test.js index 15d7d570..45b7bf0a 100644 --- a/scripts/api-integration-test.js +++ b/scripts/api-integration-test.js @@ -3824,6 +3824,25 @@ async function testPublicQuestionBankAdoption() { 'public bank sync should identify the source question that conflicts with tenant edits', ); + const conflictList = await request('/api/tenant-content/public-question-banks/conflicts', { + tenantId: PARTNER_TENANT_ID, + userId: PARTNER_TENANT_ADMIN_USER_ID, + query: { adoptionId: adopted.item.id }, + }); + assert.equal(conflictList.item?.syncStatus, 'failed', 'conflict list should expose failed adoption sync status'); + assert.ok(conflictList.item?.conflictCount >= 1, 'conflict list should expose conflict count'); + assert.ok( + conflictList.item?.conflicts?.some(item => item.sourceQuestionId === ids.question && item.action === 'conflict'), + 'conflict list should expose source question conflict details', + ); + + const crossTenantConflictListDenied = await request('/api/tenant-content/public-question-banks/conflicts', { + userId: TENANT_ADMIN_USER_ID, + query: { adoptionId: adopted.item.id }, + expectStatus: 404, + }); + assert.equal(crossTenantConflictListDenied.code, 'QUESTION_BANK_ADOPTION_NOT_FOUND', 'public bank conflicts must be tenant isolated'); + const collectionAfterConflict = await request('/api/catalog/question-collections/questions', { tenantId: PARTNER_TENANT_ID, userId: false, diff --git a/scripts/public-bank-worker-integration-test.js b/scripts/public-bank-worker-integration-test.js new file mode 100644 index 00000000..885b9fad --- /dev/null +++ b/scripts/public-bank-worker-integration-test.js @@ -0,0 +1,223 @@ +import assert from 'node:assert/strict'; +import pg from 'pg'; +import { spawn } from 'node:child_process'; + +const databaseUrl = process.env.DATABASE_URL || 'postgresql://postgres:postgres@127.0.0.1:54322/postgres'; +const MAIN_TENANT_ID = '00000000-0000-0000-0000-000000000001'; +const PARTNER_TENANT_ID = '00000000-0000-0000-0000-000000000901'; +const PARTNER_ADMIN_USER_ID = '00000000-0000-0000-0000-000000000907'; + +const ids = { + sourceQuestionBank: '00000000-0000-0000-0000-000000000400', + grant: '00000000-0000-0000-0000-000000000906', + targetQuestionBank: '00000000-0000-0000-0000-00000000c901', + targetEntry: '00000000-0000-0000-0000-00000000c902', + targetCollection: '00000000-0000-0000-0000-00000000c903', + adoption: '00000000-0000-0000-0000-00000000c904', +}; + +function runWorkerOnce() { + const child = spawn(process.execPath, ['apps/worker/dist/apps/worker/src/index.js', '--once', '--job', 'public-banks'], { + cwd: process.cwd(), + env: { + ...process.env, + DATABASE_URL: databaseUrl, + WORKER_PUBLIC_BANK_SYNC_BATCH_SIZE: '5', + WORKER_PUBLIC_BANK_SYNC_COPY_LIMIT: '2', + WORKER_PUBLIC_BANK_SYNC_ID: 'public-bank-sync-integration-test', + }, + stdio: ['ignore', 'pipe', 'pipe'], + windowsHide: true, + }); + let output = ''; + child.stdout.on('data', chunk => { + output += chunk.toString(); + }); + child.stderr.on('data', chunk => { + output += chunk.toString(); + }); + return new Promise((resolve, reject) => { + child.on('error', reject); + child.on('exit', code => { + try { + assert.equal(code, 0, `worker should exit 0\n${output}`); + assert.match(output, /public-banks batch processed=\d+/, 'worker output should include public bank sync summary'); + resolve(output); + } catch (error) { + reject(error); + } + }); + }); +} + +async function cleanup(pool) { + await pool.query('delete from public.audit_logs where tenant_id = $1 and target_id = $2', [PARTNER_TENANT_ID, ids.adoption]); + await pool.query('delete from public.tenant_question_bank_adoptions where tenant_id = $1 and source_question_bank_id = $2', [PARTNER_TENANT_ID, ids.sourceQuestionBank]); + await pool.query( + ` + delete from public.question_collection_items + where tenant_id = $1 + and (collection_id = $2 or metadata->>'source' = 'public_question_bank_adoption') + `, + [PARTNER_TENANT_ID, ids.targetCollection], + ); + await pool.query( + ` + delete from public.question_versions + where tenant_id = $1 + and question_id in ( + select id + from public.questions + where tenant_id = $1 and legacy_id like 'public:%' + ) + `, + [PARTNER_TENANT_ID], + ); + await pool.query('delete from public.questions where tenant_id = $1 and legacy_id like $2', [PARTNER_TENANT_ID, 'public:%']); + await pool.query('delete from public.question_collections where tenant_id = $1 and id = $2', [PARTNER_TENANT_ID, ids.targetCollection]); + await pool.query('delete from public.content_entries where tenant_id = $1 and id = $2', [PARTNER_TENANT_ID, ids.targetEntry]); + await pool.query('delete from public.question_banks where tenant_id = $1 and id = $2', [PARTNER_TENANT_ID, ids.targetQuestionBank]); +} + +async function createPendingAdoption(pool) { + await pool.query( + ` + insert into public.question_banks (id, tenant_id, name, source_scope, status, metadata) + values ($1, $2, 'worker 自动同步目标题库', 'tenant', 'active', '{"source":"public_bank_worker_test"}'::jsonb) + `, + [ids.targetQuestionBank, PARTNER_TENANT_ID], + ); + await pool.query( + ` + insert into public.content_entries ( + id, tenant_id, entry_key, name, entry_type, icon, route, + description, visibility, access_rules, layout_config, sort_order, + is_active, created_by, updated_by + ) + values ( + $1, $2, 'worker-public-bank-sync', 'worker 自动同步题库入口', + 'question_practice', 'book-open', '/practice', + 'worker integration public bank sync', 'public', '{}'::jsonb, + '{}'::jsonb, 100, true, $3, $3 + ) + `, + [ids.targetEntry, PARTNER_TENANT_ID, PARTNER_ADMIN_USER_ID], + ); + await pool.query( + ` + insert into public.question_collections ( + id, tenant_id, entry_id, question_bank_id, name, + collection_type, source_type, filters, question_count, + status, sort_order, metadata, created_by, updated_by + ) + values ( + $1, $2, $3, $4, 'worker 自动同步题目集合', + 'manual', 'manual_questions', '{}'::jsonb, 0, + 'active', 1, '{"source":"public_question_bank_adoption"}'::jsonb, $5, $5 + ) + `, + [ids.targetCollection, PARTNER_TENANT_ID, ids.targetEntry, ids.targetQuestionBank, PARTNER_ADMIN_USER_ID], + ); + await pool.query( + ` + insert into public.tenant_question_bank_adoptions ( + id, tenant_id, source_question_bank_id, grant_id, + target_question_bank_id, target_entry_id, target_collection_id, + adoption_mode, status, sync_status, source_snapshot, + copied_question_count, metadata, created_by, updated_by + ) + values ( + $1, $2, $3, $4, + $5, $6, $7, + 'copied_snapshot', 'active', 'pending', '{}'::jsonb, + 0, '{"source":"public_bank_worker_test"}'::jsonb, $8, $8 + ) + `, + [ + ids.adoption, + PARTNER_TENANT_ID, + ids.sourceQuestionBank, + ids.grant, + ids.targetQuestionBank, + ids.targetEntry, + ids.targetCollection, + PARTNER_ADMIN_USER_ID, + ], + ); +} + +async function main() { + const pool = new pg.Pool({ connectionString: databaseUrl }); + try { + await cleanup(pool); + await createPendingAdoption(pool); + + const output = await runWorkerOnce(); + assert.match(output, /processed=1/, 'worker should claim the pending public bank adoption'); + assert.match(output, /synced=1/, 'worker should sync the pending public bank adoption'); + + const adoption = await pool.query( + ` + select status, sync_status, copied_question_count, metadata, last_synced_at + from public.tenant_question_bank_adoptions + where tenant_id = $1 and id = $2 + `, + [PARTNER_TENANT_ID, ids.adoption], + ); + assert.equal(adoption.rows[0]?.status, 'active', 'worker should return adoption to active status'); + assert.equal(adoption.rows[0]?.sync_status, 'synced', 'worker should mark adoption synced'); + assert.equal(Number(adoption.rows[0]?.copied_question_count), 2, 'worker should respect configured copy limit'); + assert.equal(adoption.rows[0]?.metadata?.lastSync?.triggeredBy, 'worker', 'metadata should record worker-triggered sync'); + assert.ok(adoption.rows[0]?.last_synced_at, 'worker should record last_synced_at'); + + const copied = await pool.query( + ` + select q.id, q.legacy_id, v.content + from public.questions q + join public.question_versions v on v.id = q.current_version_id + where q.tenant_id = $1 + and q.legacy_id like $2 + order by q.created_at asc + `, + [PARTNER_TENANT_ID, `public:${MAIN_TENANT_ID}:%`], + ); + assert.equal(copied.rowCount, 2, 'worker should create copied tenant questions'); + assert.ok(copied.rows.every(row => row.content), 'copied questions should have current version content'); + + const collectionItems = await pool.query( + ` + select count(*)::integer as count + from public.question_collection_items + where tenant_id = $1 and collection_id = $2 + `, + [PARTNER_TENANT_ID, ids.targetCollection], + ); + assert.equal(Number(collectionItems.rows[0]?.count), 2, 'worker should bind copied questions to target collection'); + + const audit = await pool.query( + ` + select action, details + from public.audit_logs + where tenant_id = $1 and target_id = $2 and target_type = 'tenant_question_bank_adoption' + order by created_at desc + limit 1 + `, + [PARTNER_TENANT_ID, ids.adoption], + ); + assert.equal(audit.rows[0]?.action, 'content.public_question_bank.synced', 'worker should write sync audit'); + assert.equal(audit.rows[0]?.details?.triggeredBy, 'worker', 'worker audit should record trigger source'); + + const secondOutput = await runWorkerOnce(); + assert.match(secondOutput, /processed=0/, 'worker should skip already synced public bank when source has not changed'); + + console.log('Public question bank sync worker integration test complete.'); + } finally { + await cleanup(pool).catch(() => {}); + await pool.end(); + } +} + +main().catch(error => { + console.error(error); + process.exit(1); +});