From 330342680665390c6c34dc0f9e0846b75195e1fd Mon Sep 17 00:00:00 2001 From: Codex Date: Tue, 30 Jun 2026 10:04:48 +0800 Subject: [PATCH] feat: normalize legacy migration learning data --- README.md | 2 +- docs/refactor/backend-capability-status.md | 1 + docs/refactor/next-development-todo.md | 5 +- .../pocketbase-real-data-migration-runbook.md | 23 +- .../import-pocketbase/src/dry-run-report.ts | 4 + scripts/import-pocketbase/src/import-json.ts | 457 +++++++++++++++++- .../import-pocketbase/src/validate-import.ts | 114 +++++ 7 files changed, 584 insertions(+), 22 deletions(-) diff --git a/README.md b/README.md index a10564d5..8041c2cb 100644 --- a/README.md +++ b/README.md @@ -439,7 +439,7 @@ dry-run 会检查导出目录、JSON 形态、核心集合缺失、重复/缺失 默认 dry-run 使用 `development` profile;正式迁移、预生产验收和 CI 应使用 `production` profile。生产 profile 会额外输出 `migrationReadiness`,检查用户、题目、科目、分类、订单、套餐、激活码、单词和知识手册等必需集合,以及用户手机号、题目归属、订单套餐、激活码、单词和手册归属等关键字段覆盖率。`--profile` 只接受 `development` 或 `production`,拼写错误会按 blocker 失败。 -当前真实 SQLite 基线已经跑通只读导出:58 个业务 collection、248555 条记录、9 个 storage 原始资源文件。production dry-run 目前仍有 2 个真实数据 blocker:30 个订单缺用户、22 个知识手册章节缺所属手册;还有 `user_answer_records`、`mock_exam_configs`、`referral_qrcodes`、`commission_settings` 等 mapper gap,详见: +当前真实 SQLite 基线已经跑通只读导出:58 个业务 collection、248555 条记录、9 个 storage 原始资源文件。`user_answer_records`、`mock_exam_configs`、`referral_qrcodes`、`commission_settings` 已进入标准化导入。production dry-run 目前仍有 2 个真实数据 blocker:30 个订单缺用户、22 个知识手册章节缺所属手册;导入器会把它们隔离到财务复核/迁移待复核手册并写入 `pb_import_issues`,但正式切换前仍必须人工确认,详见: ```text docs/refactor/pocketbase-real-data-migration-runbook.md diff --git a/docs/refactor/backend-capability-status.md b/docs/refactor/backend-capability-status.md index fead0ae9..3f110083 100644 --- a/docs/refactor/backend-capability-status.md +++ b/docs/refactor/backend-capability-status.md @@ -172,6 +172,7 @@ | --- | --- | --- | | PocketBase schema/导出分析 | 可联调 | `scripts/import-pocketbase` 支持 schema summary/risk、`npm run pb:import:dry-run` 导出目录静态迁移报告 | | PocketBase JSON dry-run | 可联调 | 不写数据库,检查导出目录、JSON 形态、核心集合、旧 ID、敏感字段、schema relation、未映射集合和关键业务计数;`--profile=production` 会额外检查生产迁移必需集合和关键字段覆盖率,正式切换建议配合 `--fail-on-warnings` | +| PocketBase SQLite 标准化导入 | 迁移期 | 已支持真实 SQLite 只读导出后的核心集合标准化,新增覆盖 `user_answer_records -> answer_records/wrong_questions`、`mock_exam_configs -> practice_blueprints`、`referral_qrcodes -> referral_qrcodes/referral_codes`、`commission_settings -> tenant_commission_settings`;缺用户订单会进入财务复核且不自动开权益,缺归属手册章节会进入“迁移待复核手册”并写 `pb_import_issues`;全量真实导入仍需继续做批量化性能优化 | | 题目 JSON preview/import | 可联调 | 后端负责规范化、issue、幂等、审计 | | 公共题库采纳、手动同步和自动同步 | 可联调 | 平台授权后,租户可采纳公共题库并复制已发布题目快照;同步 API 和 `public-banks` worker 支持新增/更新题目、重新校验授权、跨租户拒绝、审计记录、租户内容通知和租户自改冲突保护;冲突处理 API 已支持单条/批量采纳平台版本和保留租户本地版本;worker 失败会写 `public_question_bank_sync_failed` 租户通知并保留稳定错误码,恢复成功会自动关闭失败通知;平台可通过 `/api/platform-admin/question-bank-sync-status` 查看跨租户同步运营状态;已覆盖跨租户、重复采纳、采纳后组卷、同步新增题、通知隔离/已读/自动 resolved、冲突不覆盖、单条/批量冲突处理、worker 自动同步/失败通知/恢复关闭、starter 单地区不可见/不可采纳第二地区题库、pro 全国套餐可见第二地区题库等测试 | | 单词 JSON preview/import | 可联调 | 兼容旧模板 | diff --git a/docs/refactor/next-development-todo.md b/docs/refactor/next-development-todo.md index 78c5e8cc..e9302cb3 100644 --- a/docs/refactor/next-development-todo.md +++ b/docs/refactor/next-development-todo.md @@ -76,8 +76,9 @@ - 从 PocketBase 导出现有用户、题库、单词、知识手册、分数线、订单、权益数据。 - 已补 `npm run pb:export:sqlite` 只读 SQLite 导出工具,默认从 `F:\project\参考\旧题库数据库文件\data.db` 导出到 `.gitignore` 覆盖的 `pb_export/`,并生成 `sqlite-export-manifest.json`、`storage-manifest.json` 和脱敏统计。 - 已补 `npm run pb:import:dry-run` 静态迁移报告工具、`--profile=production` 生产迁移门禁、关键集合/关键字段覆盖率检查、strict warning 门禁测试和真实数据迁移验收 runbook;真实 SQLite 已导出 58 个业务 collection、248555 条记录,并跑过 production dry-run。 - - 当前真实 dry-run 剩余 blocker:30 个订单缺 `userId`,其中 7 个 paid;22 个知识手册章节缺 `subjectId`。正式导入前必须补自动修复/隔离 mapper 和人工复核报告。 - - 当前真实 dry-run 剩余 mapper gap:`user_answer_records` 85442 条、`mock_exam_configs` 48 条、`referral_qrcodes` 79 条、`commission_settings` 1 条。它们分别对应学习历史/错题归因、全真模拟蓝图、推广码/小程序码、分佣设置,后续要进入标准化导入。 + - 当前真实 dry-run 剩余 blocker:30 个订单缺 `userId`,其中 7 个 paid;22 个知识手册章节缺 `subjectId`。导入器已补隔离策略:缺用户订单进入财务复核且不自动开权益,缺归属手册章节进入“迁移待复核手册”。正式切换前仍必须人工找回/确认这些记录。 + - 已补真实 mapper gap:`user_answer_records` 85442 条标准化到 `answer_records/wrong_questions`,`mock_exam_configs` 48 条标准化到 `practice_blueprints`,`referral_qrcodes` 79 条标准化到 `referral_qrcodes/referral_codes`,`commission_settings` 1 条标准化到 `tenant_commission_settings`。 + - 真实 `pb:import:json` 试跑已验证 raw 导入和新 mapper 开始落库,但 248555 条真实记录的逐条标准化仍超过 10 分钟,需要下一阶段继续做批量化/事务分批优化,重点是 `questions`、`question_versions`、`user_answer_records`、`pb_raw_records`。 - 处理完 blocker 后再跑 `npm run pb:import:dry-run -- --profile=production --json --fail-on-warnings`,再跑 `pb:import:json`、`pb:import:validate` 和业务抽样。 - 对题目 JSON、单词、知识手册、分数线、视频走后端 preview/import API 做二次验证。 diff --git a/docs/refactor/pocketbase-real-data-migration-runbook.md b/docs/refactor/pocketbase-real-data-migration-runbook.md index d2bab6c2..6d54eead 100644 --- a/docs/refactor/pocketbase-real-data-migration-runbook.md +++ b/docs/refactor/pocketbase-real-data-migration-runbook.md @@ -136,18 +136,29 @@ storage-manifest: 9 个原始资源文件 auxiliary.db: _logs 121715 条,仅摘要 ``` -最近一次 production dry-run 仍有 2 个 blocker,需要在正式导入前处理: +最近一次 production dry-run 仍有 2 个真实数据 blocker,需要在正式切换前处理: - `orders.userId` 覆盖率 95.3%,30 个订单缺用户,其中 7 个为已支付订单。处理策略应优先从支付通知、手机号、订单号或人工售后记录找回用户;找不回的已支付订单必须进入人工异常台账,不允许静默开权益。 - `handbook_chapters.subjectId` 覆盖率 94.4%,22 个章节缺所属手册。处理策略应按章节标题和条目归属归并到正确 `handbook_subjects`,无法归并的放入迁移隔离手册并标记待人工复核。 -production dry-run 还有这些 warning,属于后续 mapper 或运营处理项: +当前导入器已经补齐下面 4 个旧集合的标准化 mapper: + +- `user_answer_records` 85442 条会导入到 `answer_records`,并按错误记录重建 `wrong_questions`。 +- `mock_exam_configs` 48 条会导入到 `practice_blueprints(mode='mock_exam', assembly_type='filters')`。 +- `referral_qrcodes` 79 条会导入到 `referral_qrcodes`,并同步补 `referral_codes`。 +- `commission_settings` 1 条会导入到 `tenant_commission_settings`。 + +导入器对两个真实 blocker 采用“隔离 + 审计”策略,不会静默丢数据: + +- 缺用户订单仍导入 `orders/payments/order_items`,`raw_payload.migration.reviewRequired=true`、`entitlementBlocked=true`,并写入 `pb_import_issues`。缺用户订单不会自动开通权益。 +- 缺 `subjectId` 的手册章节会挂到 `handbook_subjects.legacy_id='__migration_orphan_handbook_subject__'` 的“迁移待复核手册”,章节 `metadata.reviewRequired=true`,并写入 `pb_import_issues`。 + +production dry-run 还有这些 warning,属于运营处理项: -- `user_answer_records` 85442 条尚未规范化到新练习历史/答题记录模型。 -- `mock_exam_configs` 48 条尚未规范化到 `practice_blueprints`。 -- `referral_qrcodes` 79 条尚未规范化到新推广码/小程序码模型。 -- `commission_settings` 1 条需要迁移到新分佣设置。 - `dashboard_cache`、空 `customer_messages/smscodes/tenant_config` 可不作为正式数据源,必要时只保留审计摘要。 +- `recent_practices` 和少量 `user_answer_records` 有旧用户引用断裂,需要在导入后复核 `pb_import_issues`。 + +性能注意:真实 `pb:import:json` 试跑已经证明 248555 条真实记录会超过 10 分钟。正式迁移前需要继续优化 importer 的批量写入和分批事务,重点是 `pb_raw_records`、`questions/question_versions`、`user_answer_records`。迁移演练必须记录总耗时、每阶段耗时和失败恢复策略。 ## 阶段 1:静态 Dry-Run diff --git a/scripts/import-pocketbase/src/dry-run-report.ts b/scripts/import-pocketbase/src/dry-run-report.ts index 349b7cdc..c6423977 100644 --- a/scripts/import-pocketbase/src/dry-run-report.ts +++ b/scripts/import-pocketbase/src/dry-run-report.ts @@ -96,6 +96,7 @@ const supportedCollections = new Set([ 'categories', 'code_batches', 'codes', + 'commission_settings', 'coupons', 'coupon_redemptions', 'crm_config', @@ -109,6 +110,7 @@ const supportedCollections = new Set([ 'handbook_subjects', 'images', 'majors', + 'mock_exam_configs', 'module_nodes', 'orders', 'products', @@ -117,6 +119,7 @@ const supportedCollections = new Set([ 'questions', 'recent_practices', 'referral_tracks', + 'referral_qrcodes', 'region_modules', 'regions', 'reports', @@ -132,6 +135,7 @@ const supportedCollections = new Set([ 'svip_plans', 'timelines', 'user_badges', + 'user_answer_records', 'user_word_favorites', 'user_word_progress', 'users', diff --git a/scripts/import-pocketbase/src/import-json.ts b/scripts/import-pocketbase/src/import-json.ts index b8190f3e..18d79121 100644 --- a/scripts/import-pocketbase/src/import-json.ts +++ b/scripts/import-pocketbase/src/import-json.ts @@ -39,6 +39,7 @@ const legacyLookupTables = new Set([ 'majors', 'module_nodes', 'orders', + 'practice_blueprints', 'products', 'questions', 'region_modules', @@ -54,6 +55,12 @@ const legacyLookupTables = new Set([ 'vocabulary_words', ]); +const nonCollectionJsonFiles = new Set([ + 'pb_schema.sqlite.json', + 'sqlite-export-manifest.json', + 'storage-manifest.json', +]); + function asArray(input: unknown): JsonRecord[] { if (Array.isArray(input)) return input as JsonRecord[]; if (input && typeof input === 'object' && Array.isArray((input as { items?: unknown[] }).items)) { @@ -179,6 +186,12 @@ function arrayValue(value: unknown): unknown[] { return []; } +function objectValue(value: unknown): JsonRecord { + const parsed = parseJsonish(value, {}); + if (parsed && typeof parsed === 'object' && !Array.isArray(parsed)) return parsed as JsonRecord; + return {}; +} + function entriesValue(value: unknown): Array<[string, unknown]> { const parsed = parseJsonish(value, {}); if (Array.isArray(parsed)) return parsed.map((item, index) => [String(index), item]); @@ -190,6 +203,20 @@ function dateText(value: unknown): string { return text(value) || ''; } +function nullableDateText(...values: unknown[]): string { + for (const value of values) { + const result = dateText(value); + if (result) return result; + } + return ''; +} + +function clampRate(value: unknown, fallback = 0.2): number { + const parsed = numberValue(value, fallback); + if (parsed === null || !Number.isFinite(parsed)) return fallback; + return Math.min(1, Math.max(0, parsed)); +} + function normalizeTenantRole(roleValue: unknown): string { const role = text(roleValue)?.toLowerCase(); if (role === 'superadmin') return 'platform_admin'; @@ -395,22 +422,45 @@ async function markNormalized(runId: string, collection: string) { ); } +const legacyLookupCache = new Map(); +const userLookupCache = new Map(); +const questionVersionLookupCache = new Map(); + async function legacyId(tableName: string, legacyValue: unknown): Promise { const value = text(legacyValue); if (!value) return null; if (!legacyLookupTables.has(tableName)) throw new Error(`Unsafe legacy lookup table: ${tableName}`); + const cacheKey = `${tableName}:${value}`; + if (legacyLookupCache.has(cacheKey)) return legacyLookupCache.get(cacheKey) || null; const row = await queryOne<{ id: string }>( `select id from public.${tableName} where tenant_id = $1 and legacy_id = $2 limit 1`, [tenantId, value], ); - return row?.id || null; + const result = row?.id || null; + legacyLookupCache.set(cacheKey, result); + return result; } async function userIdByLegacy(value: unknown): Promise { const legacyValue = text(value); if (!legacyValue) return null; + if (userLookupCache.has(legacyValue)) return userLookupCache.get(legacyValue) || null; const row = await queryOne<{ id: string }>('select id from public.platform_users where legacy_id = $1 limit 1', [legacyValue]); - return row?.id || null; + const result = row?.id || null; + userLookupCache.set(legacyValue, result); + return result; +} + +async function currentQuestionVersionId(questionId: string | null): Promise { + if (!questionId) return null; + if (questionVersionLookupCache.has(questionId)) return questionVersionLookupCache.get(questionId) || null; + const row = await queryOne<{ current_version_id: string | null }>( + `select current_version_id from public.questions where tenant_id = $1 and id = $2 limit 1`, + [tenantId, questionId], + ); + const result = row?.current_version_id || null; + questionVersionLookupCache.set(questionId, result); + return result; } async function upsertSecret(secretScope: string, secretKey: string, secretValue: unknown, provider: string | null = null) { @@ -825,6 +875,109 @@ async function normalizeQuestions(records: JsonRecord[]) { } } +async function normalizeUserAnswerRecords(runId: string, records: JsonRecord[]) { + const affectedWrongPairs = new Set(); + + for (const r of records) { + const userId = await userIdByLegacy(r.userId); + if (!userId) { + await issue( + runId, + 'user_answer_records', + r.id, + 'answer_user_not_found', + `Answer record cannot be imported because userId was not resolved: ${text(r.userId) || '(empty)'}`, + 'critical', + ); + continue; + } + + const questionId = await legacyId('questions', r.questionId); + if (!questionId) { + await issue( + runId, + 'user_answer_records', + r.id, + 'answer_question_not_found', + `Answer record kept with legacy_question_id only because questionId was not resolved: ${text(r.questionId) || '(empty)'}`, + 'warning', + ); + } + + const isCorrect = r.isCorrect === null || r.isCorrect === undefined || r.isCorrect === '' + ? null + : boolValue(r.isCorrect); + const answeredAt = nullableDateText(r.answeredAt, r.updated, r.created); + const answerPayload = { + source: 'pocketbase.user_answer_records', + selectedOptions: arrayValue(r.selectedOptions), + legacyCategoryId: text(r.categoryId), + legacyQuestionId: text(r.questionId), + importedAt: new Date().toISOString(), + }; + + await pool.query( + ` + insert into public.answer_records ( + tenant_id, user_id, question_id, question_version_id, legacy_id, + legacy_question_id, legacy_category_id, selected_options, + answer_payload, is_correct, answered_at, created_at + ) + values ($1,$2,$3,$4,$5,$6,$7,$8::jsonb,$9::jsonb,$10, + coalesce(nullif($11::text,'')::timestamptz, now()), + coalesce(nullif($12::text,'')::timestamptz, now()) + ) + on conflict (tenant_id, legacy_id) do update set + user_id = excluded.user_id, + question_id = excluded.question_id, + question_version_id = excluded.question_version_id, + legacy_question_id = excluded.legacy_question_id, + legacy_category_id = excluded.legacy_category_id, + selected_options = excluded.selected_options, + answer_payload = excluded.answer_payload, + is_correct = excluded.is_correct, + answered_at = excluded.answered_at + `, + [ + tenantId, + userId, + questionId, + await currentQuestionVersionId(questionId), + r.id, + text(r.questionId), + text(r.categoryId), + JSON.stringify(arrayValue(r.selectedOptions)), + JSON.stringify(answerPayload), + isCorrect, + answeredAt, + nullableDateText(r.created, answeredAt), + ], + ); + + if (isCorrect === false && questionId) affectedWrongPairs.add(`${userId}:${questionId}`); + } + + for (const pair of affectedWrongPairs) { + const [userId, questionId] = pair.split(':'); + await pool.query( + ` + insert into public.wrong_questions (tenant_id, user_id, question_id, wrong_count, last_wrong_at, resolved_at) + select $1, $2, $3, count(*)::integer, max(answered_at), null + from public.answer_records + where tenant_id = $1 + and user_id = $2 + and question_id = $3 + and is_correct = false + on conflict (tenant_id, user_id, question_id) + do update set wrong_count = excluded.wrong_count, + last_wrong_at = excluded.last_wrong_at, + resolved_at = null + `, + [tenantId, userId, questionId], + ); + } +} + async function normalizeUsers(records: JsonRecord[]) { for (const r of records) { const safeProfile = publicProfile(r); @@ -1166,6 +1319,84 @@ async function normalizeSvipPlans(records: JsonRecord[]) { } } +function normalizeMockExamSections(questionTypes: unknown) { + return arrayValue(questionTypes) + .map((item, index) => { + const object = objectValue(item); + const questionType = text(object.type) || text(object.questionType) || `section_${index + 1}`; + const count = intValue(object.count, 0); + return { + key: text(object.key) || questionType, + title: text(object.title) || text(object.name) || questionType, + questionType, + count, + sortOrder: intValue(object.order, index + 1) - 1, + scoreEach: numberValue(object.scoreEach ?? object.score, null), + }; + }) + .filter(section => section.count > 0) + .sort((a, b) => a.sortOrder - b.sortOrder); +} + +async function normalizeMockExamConfigs(runId: string, records: JsonRecord[]) { + for (const r of records) { + const subjectId = await legacyId('subjects', r.subjectId); + const sections = normalizeMockExamSections(r.questionTypes); + const questionLimit = intValue(r.totalQuestions, sections.reduce((sum, section) => sum + section.count, 0)); + if (!subjectId) { + await issue( + runId, + 'mock_exam_configs', + r.id, + 'mock_exam_subject_not_found', + `Mock exam blueprint imported without resolved subject because subjectId was not found: ${text(r.subjectId) || '(empty)'}`, + 'warning', + ); + } + + await pool.query( + ` + insert into public.practice_blueprints ( + tenant_id, legacy_id, name, mode, assembly_type, question_limit, + duration_minutes, sections, rules, status, sort_order, created_at, updated_at + ) + values ($1,$2,$3,'mock_exam','filters',$4,$5,$6::jsonb,$7::jsonb,$8,$9, + coalesce(nullif($10::text,'')::timestamptz, now()), + coalesce(nullif($11::text,'')::timestamptz, now()) + ) + on conflict (tenant_id, legacy_id) do update set + name = excluded.name, + question_limit = excluded.question_limit, + duration_minutes = excluded.duration_minutes, + sections = excluded.sections, + rules = excluded.rules, + status = excluded.status, + sort_order = excluded.sort_order, + updated_at = excluded.updated_at + `, + [ + tenantId, + r.id, + text(r.name) || `全真模拟-${text(r.subjectId) || text(r.id) || '未命名'}`, + questionLimit > 0 ? questionLimit : null, + intValue(r.duration, 45), + JSON.stringify(sections), + JSON.stringify({ + source: 'pocketbase.mock_exam_configs', + subjectId, + legacySubjectId: text(r.subjectId), + legacyQuestionTypes: arrayValue(r.questionTypes), + randomize: true, + }), + boolValue(r.isActive, true) ? 'active' : 'archived', + intValue(r.order), + dateText(r.created), + dateText(r.updated), + ], + ); + } +} + async function normalizeCodeBatches(records: JsonRecord[]) { for (const r of records) { await pool.query( @@ -1373,8 +1604,19 @@ async function normalizeCouponRedemptions(records: JsonRecord[]) { } } -async function normalizeOrders(records: JsonRecord[]) { +async function normalizeOrders(runId: string, records: JsonRecord[]) { for (const r of records) { + const resolvedUserId = await userIdByLegacy(r.userId); + const sanitized = sanitizeRecord(r); + const rawPayload = { + ...(sanitized as Record), + migration: { + source: 'pocketbase.orders', + reviewRequired: !resolvedUserId, + reviewReason: !resolvedUserId ? 'legacy_user_not_resolved' : null, + entitlementBlocked: !resolvedUserId, + }, + }; const order = await queryOne<{ id: string }>( ` insert into public.orders ( @@ -1410,7 +1652,7 @@ async function normalizeOrders(records: JsonRecord[]) { `, [ tenantId, - await userIdByLegacy(r.userId), + resolvedUserId, r.id, text(r.userId), text(r.orderNo) || `legacy-${text(r.id)}`, @@ -1426,13 +1668,23 @@ async function normalizeOrders(records: JsonRecord[]) { await legacyId('regions', r.regionId), text(r.regionId), dateText(r.paidAt), - JSON.stringify(sanitizeRecord(r)), + JSON.stringify(rawPayload), dateText(r.created), dateText(r.updated), ], ); if (!order) continue; + if (!resolvedUserId) { + await issue( + runId, + 'orders', + r.id, + 'order_user_not_resolved', + `Order imported for finance review only because userId was not resolved: ${text(r.userId) || '(empty)'}`, + normalizeOrderStatus(r.status) === 'paid' ? 'critical' : 'warning', + ); + } await pool.query( ` @@ -1625,17 +1877,60 @@ async function normalizeHandbookSubjects(records: JsonRecord[]) { } } -async function normalizeHandbookChapters(records: JsonRecord[]) { +async function migrationOrphanHandbookSubjectId(): Promise { + const legacyIdValue = '__migration_orphan_handbook_subject__'; + const row = await queryOne<{ id: string }>( + ` + insert into public.handbook_subjects ( + tenant_id, legacy_id, name, type, description, sort_order, is_active, metadata + ) + values ($1,$2,'迁移待复核手册','migration_review','旧 PocketBase 章节缺少 subjectId 时自动挂载到这里,需人工归并。',999999,true,$3::jsonb) + on conflict (tenant_id, legacy_id) do update set + name = excluded.name, + description = excluded.description, + metadata = public.handbook_subjects.metadata || excluded.metadata, + updated_at = now() + returning id + `, + [ + tenantId, + legacyIdValue, + JSON.stringify({ + source: 'pocketbase.handbook_chapters', + reviewRequired: true, + reviewReason: 'legacy_chapter_subject_missing', + }), + ], + ); + if (!row) throw new Error('Failed to create migration orphan handbook subject'); + return row.id; +} + +async function normalizeHandbookChapters(runId: string, records: JsonRecord[]) { for (const r of records) { + let subjectId = await legacyId('handbook_subjects', r.subjectId); + const missingSubject = !subjectId; + if (missingSubject) { + subjectId = await migrationOrphanHandbookSubjectId(); + await issue( + runId, + 'handbook_chapters', + r.id, + 'handbook_chapter_subject_missing', + `Handbook chapter was attached to migration review subject because subjectId was missing or unresolved: ${text(r.subjectId) || '(empty)'}`, + 'critical', + ); + } + await pool.query( ` insert into public.handbook_chapters ( tenant_id, subject_id, legacy_id, name, description, sort_order, - is_active, created_at, updated_at + is_active, metadata, created_at, updated_at ) - values ($1,$2,$3,$4,$5,$6,$7, - coalesce(nullif($8::text,'')::timestamptz, now()), - coalesce(nullif($9::text,'')::timestamptz, now()) + values ($1,$2,$3,$4,$5,$6,$7,$8::jsonb, + coalesce(nullif($9::text,'')::timestamptz, now()), + coalesce(nullif($10::text,'')::timestamptz, now()) ) on conflict (tenant_id, legacy_id) do update set subject_id = excluded.subject_id, @@ -1643,16 +1938,23 @@ async function normalizeHandbookChapters(records: JsonRecord[]) { description = excluded.description, sort_order = excluded.sort_order, is_active = excluded.is_active, + metadata = excluded.metadata, updated_at = excluded.updated_at `, [ tenantId, - await legacyId('handbook_subjects', r.subjectId), + subjectId, r.id, text(r.name) || '未命名手册章节', text(r.description), intValue(r.order), boolValue(r.isActive, true), + JSON.stringify({ + source: 'pocketbase.handbook_chapters', + legacySubjectId: text(r.subjectId), + reviewRequired: missingSubject, + reviewReason: missingSubject ? 'legacy_subject_missing_or_unresolved' : null, + }), dateText(r.created), dateText(r.updated), ], @@ -2273,6 +2575,130 @@ async function normalizeReferralTracks(records: JsonRecord[]) { } } +function refCodeFromScene(sceneValue: unknown, fallback: unknown): string { + const scene = text(sceneValue); + if (scene) { + const match = /^ref[_=-]?(.+)$/i.exec(scene); + const code = (match?.[1] || scene).replace(/[^a-z0-9]/gi, '').toUpperCase(); + if (code) return code; + } + const fallbackText = text(fallback)?.replace(/[^a-z0-9]/gi, '').toUpperCase(); + return fallbackText || 'LEGACY'; +} + +async function normalizeReferralQrcodes(runId: string, records: JsonRecord[]) { + for (const r of records) { + const userId = await userIdByLegacy(r.userId); + const scene = text(r.scene) || `legacy_${text(r.id) || refCodeFromScene(null, r.userId)}`; + const page = text(r.page) || 'pages/index/index'; + const refCode = refCodeFromScene(scene, r.id); + const qrcodeUrl = text(r.qrcodeUrl) || text(r.image); + const metadata = { + source: 'pocketbase.referral_qrcodes', + legacyId: text(r.id), + legacyUserId: text(r.userId), + image: text(r.image), + originalQrcodeUrl: text(r.qrcodeUrl), + }; + + if (!userId) { + await issue( + runId, + 'referral_qrcodes', + r.id, + 'referral_qrcode_user_not_found', + `Referral qrcode cannot be attached to a user because userId was not resolved: ${text(r.userId) || '(empty)'}`, + 'critical', + ); + continue; + } + + await pool.query( + ` + insert into public.referral_codes (tenant_id, user_id, code, status, channel, landing_path, metadata, created_at, updated_at) + values ($1,$2,$3::citext,'active','wechat-miniapp',$4,$5::jsonb, + coalesce(nullif($6::text,'')::timestamptz, now()), + coalesce(nullif($7::text,'')::timestamptz, now()) + ) + on conflict (tenant_id, user_id) do update set + code = excluded.code, + channel = excluded.channel, + landing_path = excluded.landing_path, + metadata = public.referral_codes.metadata || excluded.metadata, + updated_at = excluded.updated_at + `, + [tenantId, userId, refCode, page, JSON.stringify(metadata), dateText(r.created), dateText(r.updated)], + ); + + await pool.query( + ` + insert into public.referral_qrcodes ( + tenant_id, user_id, ref_code, scene, page, provider, qrcode_url, status, + metadata, created_at, updated_at + ) + values ($1,$2,$3::citext,$4,$5,'wechat-miniapp',$6,$7,$8::jsonb, + coalesce(nullif($9::text,'')::timestamptz, now()), + coalesce(nullif($10::text,'')::timestamptz, now()) + ) + on conflict (tenant_id, provider, scene, page) do update set + user_id = excluded.user_id, + ref_code = excluded.ref_code, + qrcode_url = excluded.qrcode_url, + status = excluded.status, + metadata = public.referral_qrcodes.metadata || excluded.metadata, + updated_at = excluded.updated_at + `, + [ + tenantId, + userId, + refCode, + scene, + page, + qrcodeUrl, + qrcodeUrl ? 'ready' : 'pending', + JSON.stringify(metadata), + dateText(r.created), + dateText(r.updated), + ], + ); + } +} + +async function normalizeCommissionSettings(records: JsonRecord[]) { + for (const r of records) { + const config = { + source: 'pocketbase.commission_settings', + legacyId: text(r.id), + remark: text(r.remark), + raw: sanitizeRecord(r), + }; + + await pool.query( + ` + insert into public.tenant_commission_settings ( + tenant_id, default_rate, min_settlement_cents, settlement_cycle, config, + created_at, updated_at + ) + values ($1,$2,0,'monthly',$3::jsonb, + coalesce(nullif($4::text,'')::timestamptz, now()), + coalesce(nullif($5::text,'')::timestamptz, now()) + ) + on conflict (tenant_id) do update set + default_rate = excluded.default_rate, + config = public.tenant_commission_settings.config || excluded.config, + updated_at = excluded.updated_at + `, + [ + tenantId, + clampRate(r.defaultRate, 0.2), + JSON.stringify(config), + dateText(r.created), + dateText(r.updated), + ], + ); + } +} + async function normalizeBadges(records: JsonRecord[]) { for (const r of records) { await pool.query( @@ -2820,6 +3246,7 @@ async function loadCollections(): Promise { const files = fs.readdirSync(exportDir).filter(file => file.toLowerCase().endsWith('.json')); const collections: CollectionMap = {}; for (const file of files) { + if (nonCollectionJsonFiles.has(file)) continue; const collection = collectionNameFromFile(file); collections[collection] = asArray(JSON.parse(fs.readFileSync(path.join(exportDir, file), 'utf8'))); } @@ -2841,14 +3268,16 @@ async function normalizeAll(runId: string, collections: CollectionMap) { await runNormalizer(runId, collections, 'users', normalizeUsers); await normalizeUserEntitlementsAndStats(runId, collections.users || []); if ((collections.users || []).length > 0) await markNormalized(runId, 'users'); + await runNormalizer(runId, collections, 'user_answer_records', records => normalizeUserAnswerRecords(runId, records)); await runNormalizer(runId, collections, 'settings', normalizeSettings); await runNormalizer(runId, collections, 'svip_plans', normalizeSvipPlans); + await runNormalizer(runId, collections, 'mock_exam_configs', records => normalizeMockExamConfigs(runId, records)); await runNormalizer(runId, collections, 'code_batches', normalizeCodeBatches); await runNormalizer(runId, collections, 'coupons', normalizeCoupons); await runNormalizer(runId, collections, 'coupon_redemptions', normalizeCouponRedemptions); await runNormalizer(runId, collections, 'codes', normalizeActivationCodes); - await runNormalizer(runId, collections, 'orders', normalizeOrders); + await runNormalizer(runId, collections, 'orders', records => normalizeOrders(runId, records)); await runNormalizer(runId, collections, 'vocabulary_units', normalizeVocabularyUnits); await runNormalizer(runId, collections, 'vocabulary', normalizeVocabularyWords); @@ -2863,7 +3292,7 @@ async function normalizeAll(runId: string, collections: CollectionMap) { } await runNormalizer(runId, collections, 'handbook_subjects', normalizeHandbookSubjects); - await runNormalizer(runId, collections, 'handbook_chapters', normalizeHandbookChapters); + await runNormalizer(runId, collections, 'handbook_chapters', records => normalizeHandbookChapters(runId, records)); await runNormalizer(runId, collections, 'handbook_entries', normalizeHandbookEntries); await runNormalizer(runId, collections, 'banners', records => normalizeSimpleContent(records, 'banners')); await runNormalizer(runId, collections, 'faqs', records => normalizeSimpleContent(records, 'faqs')); @@ -2880,6 +3309,8 @@ async function normalizeAll(runId: string, collections: CollectionMap) { await runNormalizer(runId, collections, 'scoreline_fields', normalizeScorelineFields); await runNormalizer(runId, collections, 'scoreline_records', normalizeScorelineRecords); await runNormalizer(runId, collections, 'referral_tracks', normalizeReferralTracks); + await runNormalizer(runId, collections, 'referral_qrcodes', records => normalizeReferralQrcodes(runId, records)); + await runNormalizer(runId, collections, 'commission_settings', normalizeCommissionSettings); await runNormalizer(runId, collections, 'badges', normalizeBadges); await runNormalizer(runId, collections, 'user_badges', normalizeUserBadges); await runNormalizer(runId, collections, 'recent_practices', normalizeRecentPractices); diff --git a/scripts/import-pocketbase/src/validate-import.ts b/scripts/import-pocketbase/src/validate-import.ts index d841076c..a209efba 100644 --- a/scripts/import-pocketbase/src/validate-import.ts +++ b/scripts/import-pocketbase/src/validate-import.ts @@ -217,6 +217,120 @@ async function validationChecks(): Promise { ), ); + checks.push( + result( + 'legacy_answers_missing_user', + await scalar( + ` + select count(*) + from public.answer_records ar + left join public.platform_users u on u.id = ar.user_id + where ar.tenant_id = $1 + and ar.legacy_id is not null + and u.id is null + `, + [tenantId], + ), + 'Legacy answer_records exist without a valid user.', + 'Legacy answer_records reference valid users.', + ), + ); + + checks.push( + result( + 'legacy_answers_unresolved_questions', + await scalar( + ` + select count(*) + from public.answer_records + where tenant_id = $1 + and legacy_question_id is not null + and question_id is null + `, + [tenantId], + ), + 'Some legacy answer_records kept only legacy_question_id because questions were not resolved.', + 'Legacy answer_records question references are resolved.', + true, + ), + ); + + checks.push( + result( + 'wrong_questions_invalid_references', + await scalar( + ` + select count(*) + from public.wrong_questions wq + left join public.platform_users u on u.id = wq.user_id + left join public.questions q on q.id = wq.question_id and q.tenant_id = wq.tenant_id + where wq.tenant_id = $1 + and (u.id is null or q.id is null) + `, + [tenantId], + ), + 'Wrong-question rows reference missing users or questions.', + 'Wrong-question rows reference valid users and questions.', + ), + ); + + checks.push( + result( + 'legacy_mock_blueprints_incomplete', + await scalar( + ` + select count(*) + from public.practice_blueprints + where tenant_id = $1 + and mode = 'mock_exam' + and legacy_id is not null + and ( + question_limit is null + or jsonb_array_length(sections) = 0 + or rules->>'legacySubjectId' is null + ) + `, + [tenantId], + ), + 'Legacy mock exam blueprints are missing limits, sections, or legacy subject traceability.', + 'Legacy mock exam blueprints have limits, sections, and legacy subject traceability.', + ), + ); + + checks.push( + result( + 'referral_qrcodes_incomplete', + await scalar( + ` + select count(*) + from public.referral_qrcodes + where tenant_id = $1 + and (ref_code is null or scene = '' or page = '') + `, + [tenantId], + ), + 'Referral qrcodes are missing ref_code, scene, or page.', + 'Referral qrcodes have ref_code, scene, and page.', + ), + ); + + checks.push( + result( + 'commission_settings_invalid_rate', + await scalar( + ` + select count(*) + from public.tenant_commission_settings + where tenant_id = $1 + and (default_rate < 0 or default_rate > 1) + `, + [tenantId], + ), + 'Tenant commission settings contain invalid default_rate values.', + 'Tenant commission settings default_rate values are valid.', + ), + ); + checks.push( result( 'entitlements_unresolved_user',