Commit 0b3b7cd0 by luoqi

merge: main(98 个提交) → feat/friday-ai

零冲突。两个重叠文件(host-admin.cli.ts / canonical-codes.ts)自动合并,
已核对无残留标记、本分支新增的码表全在。

合并后重建 @pac/types 与 @pac/pac-client,全量门禁通过:
  friday-ai   374 测试 · tsc ×3 · lint 0 warning · 边界闸 5 条
  pac-service 1943 测试 · tsc

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
parents 9d962648 e19135b4
Pipeline #3651 failed in 0 seconds
......@@ -205,3 +205,18 @@ SENTRY_ENVIRONMENT=
SENTRY_TRACES_SAMPLE_RATE=0.1
# release(可选):部署时注入 git SHA
SENTRY_RELEASE=
# ── plan 批量重算性能 ────────────────────────────────────────────────
# 召回子场景并发度。**生产设 4**(2026-08-30 起)。
# 生产全量实测(113 万患者,同一台 RDS,相邻两轮):
# =1 场景段 2h00m / 整轮 2h24m
# =4 场景段 51m56s / 整轮 1h20m → 场景段 ×2.31,省 64 分钟
# 瓶颈是**延迟**(逐次索引探查在等 page 返回)不是磁盘吞吐,所以并发 2 就超线性(×1.50)。
# ⚠️ 判据只能看 `[plan] 阶段耗时` 的**墙钟** —— 并发下单条查询的 sql= 会被争抢拉长,
# 看单条会误判成"变慢了"。详见 treatment-initiation-recall.scenario.ts 里 conc 处注释。
PAC_RECALL_SUBSCENARIO_CONCURRENCY=4
# gap 计算形态:legacy(默认,逐行相关子查询)/ setbased(集合式)。
# ⚠️ 生产**未启用**。集合式已在测试机验证零差异且再省 22~31%,但未经生产验证。
# 见 docs/design/gap-set-based-rewrite-plan.md。
PAC_GAP_VARIANT=
......@@ -24,6 +24,22 @@ field_mapping:
# 不是 PAC appointment_type 通用语义"预约目的"(咨询/治疗/检查/拔牙等);
# 且"初诊/复诊"是衍生信号,PAC 自己从 encounter_record 聚合算更可信(单一源 = PAC)。
doctorId: appo_doc_id
# ⭐ 预约医生姓名 —— DW 一直在同一行里给,只是历来没映射(预约事实因此只有 id 没有名字)。
# 列名叫 resource_name(排班资源名)容易让人以为是诊室/椅位,**不是**。行为核实(2026-08-31):
# ① 椅位另有列 appo_chair(仅 3 个取值),诊室/资源另有 resource_id;
# ② resource_name 与 appo_doc_id 的绑定比与 resource_id 更紧
# (本地:(doc_id,name) 140 组 < (resource_id,name) 151 组)——资源名不会有这个方向;
# ③ 拿「已到诊」预约按 (患者,日期) 对当天病历的 doctor_name:
# 生产近 90 天 99,170/101,192 = 98.0% 命中;144,148 条明细里
# 「id 对上而名字不对」和「名字对上而 id 不对」**各 0 条**(本地全量各 1 条,均为资源改派期历史行)。
# → resource_name 就是 appo_doc_id 这个人的姓名,不是恰好像人名的资源名。
# ⚠️ 语义是**约号时约的那位医生**,不是实际接诊医生(剩下 2% 是改派,业务事实非数据问题);
# 要实际接诊医生走病历 doctor_name,两个口径别混。
# ⚠️ 少数行排的**不是人**而是房间/服务/台席("预约"/"学前街手术室"/"正畸咨询"…):
# 生产近 180 天 9,838/1,047,901 = 0.94%。canonical 层原样透传,
# 由 AppointmentParser.isPersonResource 决定写不写 content.doctor_name
# (原值另存 content.resource_name,不丢)。词表与误伤核验记在那个函数上。
doctorName: resource_name
status: appo_status
complaintCategory: appo_complaint_category # 预约科目 / 就诊意向(种植/正畸/…)
complaintText: appo_complaint # 预约主诉自由文本(跟 category 配对)
......
......@@ -572,6 +572,27 @@ transforms:
- 拒绝拍片
- 治疗中
- 已交付纸质病历
# ── 姑息处置 → 同样落 _treatment_review_raw(category=review)──
# 2026-08-28 重评「⛔ 刻意不收」那几个词(见 treatment-category-actual-rules.yaml 文末)。
# 那个块的根因写着「category 被治疗史与 resolver 两处消费,只有一个旋钮,只能二选一
# (现选都不要=丢弃)」—— **前提不成立**:第二个旋钮一直都在,就是 review 类目,
# 而且它就是为这件事设计的(见 PACTreatmentCategories.review 注释:「医生做了某个流程
# 节点/临床判断本次不动手」「不该塞 preventive,也不该丢弃」),且 review 出现在
# **0 个** resolver 家族里(结构家族注释明写「刻意排除 …preventive / review」)。
# 所以 review 精确地给①治疗史完整性、不给②解缺口 —— 冠周炎冲洗上药之后仍照常召拔牙。
# 触发个案:季炎萍 TS0B010543 牙41~37,下颌活动义齿戴数年,本次「调磨」过长边缘,
# treat_plan 只有「调磨」被丢 → 她库里**零条治疗事实**,系统看成"从没治过"。
# 影响面:DW 全量 treat_plan 全是这五词的 EMR 4,336 份 → 进治疗史,不解任何缺口。
#
# ⛔ **单独一条 route,不并进上面的 &review_terms** —— 那个 anchor 被 C.9 dispose 闸门
# 复用(blank_or_all_in),并进去会同时打开闸门:那 4,336 份处置 100% 非空,实测
# TOP400 处置里约 216 条会落进结构 resolver 家族,且含「冠周冲洗,派力奥上药」
# →prosthodontic(含"冠"字)、「去暂封…玻璃离子暂封」→restorative 这类误判,
# 会误销缺口。dispose 闸门维持现状(不开=不回退),要开另行评估。
# 两个消费方语义本就不同:route 问"哪些名字不是真治疗",闸门问"何时可安全抽处置"。
- output: _treatment_review_raw
when:
equals: [调磨, 试戴, 冲洗, 换药, 上药]
# 真治疗动作 → treatment_actual_rows ⭐ kind=actual(临床真实)
- output: _treatment_actual_raw_emr
when:
......
......@@ -30,6 +30,8 @@
"recompute-persona": "ts-node --transpile-only src/cli/recompute-persona.cli.ts",
"backfill-plan-labels": "ts-node --transpile-only src/cli/backfill-plan-labels.cli.ts",
"recompute-persona:prod": "node --max-old-space-size=8192 dist/cli/recompute-persona.cli.js",
"verify-gap-equivalence": "ts-node --transpile-only src/cli/verify-gap-equivalence.cli.ts",
"verify-gap-equivalence:prod": "node --max-old-space-size=8192 dist/cli/verify-gap-equivalence.cli.js",
"recompute-plans": "ts-node --transpile-only src/cli/recompute-plans.cli.ts",
"recompute-plans:prod": "node --max-old-space-size=8192 dist/cli/recompute-plans.cli.js",
"timeline": "ts-node --transpile-only src/cli/timeline.cli.ts",
......
......@@ -15,6 +15,7 @@ import { Logger } from '@nestjs/common';
import { AppModule } from '../app.module';
import { PrismaService } from '../prisma/prisma.service';
import { PlanScriptOrchestrator } from '../modules/ai/orchestrators/plan-script.orchestrator';
import { disableSchedulersForCli } from './bootstrap-flags';
interface CliArgs {
planId?: string;
......@@ -81,6 +82,8 @@ async function bootstrap() {
}
const logger = new Logger('ai:gen-script');
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['warn', 'error', 'log'],
});
......
......@@ -25,6 +25,7 @@ import { NestFactory } from '@nestjs/core';
import { Logger } from '@nestjs/common';
import { AppModule } from '../app.module';
import { PlanLabelService } from '../modules/plan/plan-label.service';
import { disableSchedulersForCli } from './bootstrap-flags';
async function main(): Promise<void> {
const log = new Logger('backfill-plan-labels');
......@@ -32,6 +33,10 @@ async function main(): Promise<void> {
const all = argv.includes('--all');
const checkOnly = argv.includes('--check');
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['error', 'warn', 'log'],
});
......
/**
* CLI 启动前的写侧总闸 —— **必须在 `NestFactory.createApplicationContext` 之前调用**。
*
* 🔴 为什么存在(2026-08-30 生产事故):
* 每个 CLI 都会 `createApplicationContext(AppModule)`,于是把**整个应用的 onModuleInit
* 全跑一遍** —— 包括 SyncIncrementalSchedulerService 的僵尸锁回收和 cron 注册。
* 在长驻的 pac-service 容器里 `docker exec` 跑一个 CLI,那个 CLI 就会把 service 里
* **正在跑**的那轮同步的锁当成僵尸清掉,那轮增量当场夭折。
* 实测:08:17:11 跑 recompute-plans → 08:17:13 正常跑着的 08:15 那轮被标 failed。
*
* 所以:凡是**不以调度器为目的**的 CLI,一律先调本函数。
* ⚠️ 例外只有 `sync-incremental.cli`(它本身就是要触发同步)—— 但它也不该回收别人的锁,
* 那一层由 scheduler 的年龄阈值(REAP_MIN_AGE_MS)兜底。
*/
export function disableSchedulersForCli(): void {
process.env.PAC_SCHEDULER_DISABLED = '1';
}
......@@ -15,6 +15,7 @@ import {
ColdImportService,
SyncAlreadyRunningError,
} from '../modules/sync/cold-import/cold-import.service';
import { disableSchedulersForCli } from './bootstrap-flags';
interface CliArgs {
dir?: string;
......@@ -112,6 +113,10 @@ async function bootstrap() {
`${args.since ? `, since=${args.since}${args.months ? `(--months=${args.months})` : ''}` : ''})`,
);
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['log', 'warn', 'error'],
});
......
......@@ -22,6 +22,7 @@ import { NestFactory } from '@nestjs/core';
import { Logger } from '@nestjs/common';
import { AppModule } from '../app.module';
import { HostsService } from '../modules/admin/hosts.service';
import { disableSchedulersForCli } from './bootstrap-flags';
interface ParsedArgs {
command: string;
......@@ -92,6 +93,8 @@ async function main() {
}
const logger = new Logger('pac:host');
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['warn', 'error'],
});
......
......@@ -20,6 +20,7 @@ import { ColdImportService } from '../modules/sync/cold-import/cold-import.servi
import { PersonaService } from '../modules/persona/persona.service';
import { PlanEngineService } from '../modules/plan/engine/plan-engine.service';
import * as path from 'node:path';
import { disableSchedulersForCli } from './bootstrap-flags';
interface Args {
host: string;
......@@ -73,6 +74,8 @@ async function bootstrap() {
}
const logger = new Logger('import-patient');
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['log', 'warn', 'error'],
});
......
......@@ -16,6 +16,7 @@ import { NestFactory } from '@nestjs/core';
import { Logger } from '@nestjs/common';
import { AppModule } from '../app.module';
import { PrismaService } from '../prisma/prisma.service';
import { disableSchedulersForCli } from './bootstrap-flags';
interface Row {
fileNum: string;
......@@ -48,6 +49,10 @@ async function main(): Promise<void> {
const rows = parseCsv(csvArg);
logger.log(`对照表 ${rows.length} 行(去表头/垃圾后)`);
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, { logger: ['error', 'warn'] });
const prisma = app.get(PrismaService);
try {
......
......@@ -17,6 +17,7 @@ import { PrismaService } from '../prisma/prisma.service';
import { RedisService } from '../redis/redis.service';
import { doctorOptionsCacheKey } from '../modules/plan/plan.service';
import { runPool } from '../common/run-pool';
import { disableSchedulersForCli } from './bootstrap-flags';
interface Args {
host: string;
......@@ -60,6 +61,8 @@ async function bootstrap() {
if (args.concurrency > 1 && !process.env.PAC_DB_CONCURRENCY) {
process.env.PAC_DB_CONCURRENCY = String(args.concurrency);
}
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['log', 'warn', 'error'],
});
......
......@@ -16,6 +16,7 @@ import { Logger } from '@nestjs/common';
import { AppModule } from '../app.module';
import { PlanEngineService } from '../modules/plan/engine/plan-engine.service';
import { PrismaService } from '../prisma/prisma.service';
import { disableSchedulersForCli } from './bootstrap-flags';
interface Args {
host: string;
......@@ -47,6 +48,8 @@ async function bootstrap() {
const n = Math.max(1, Number(process.env.PAC_PLAN_BATCH_CONCURRENCY) || 8);
if (n > 1) process.env.PAC_DB_CONCURRENCY = String(n);
}
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['log', 'warn', 'error'],
});
......
......@@ -25,6 +25,7 @@ import { createClient } from '@clickhouse/client';
import { AppModule } from '../app.module';
import { PrismaService } from '../prisma/prisma.service';
import { ColdImportManifestSchema } from '../modules/sync/cold-import/manifest.schema';
import { disableSchedulersForCli } from './bootstrap-flags';
// CLI 是短命进程,不需要 org-tree 启动预热(且会在 app.close() 时跟后台 warmAll 抢连接报噪音)。
process.env.PAC_ORGTREE_WARM_ON_BOOT = 'false';
......@@ -118,6 +119,10 @@ async function main(): Promise<void> {
process.exit(0);
}
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, { logger: ['warn', 'error'] });
try {
const prisma = app.get(PrismaService);
......
......@@ -38,6 +38,7 @@ import { PrismaService } from '../prisma/prisma.service';
import { ColdImportService } from '../modules/sync/cold-import/cold-import.service';
import { PersonaService } from '../modules/persona/persona.service';
import { PlanEngineService } from '../modules/plan/engine/plan-engine.service';
import { disableSchedulersForCli } from './bootstrap-flags';
interface CliArgs {
dir: string;
......@@ -81,6 +82,10 @@ async function main(): Promise<void> {
return;
}
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, { logger: ['log', 'warn', 'error'] });
try {
const prisma = app.get(PrismaService);
......
......@@ -34,6 +34,7 @@ import {
} from '@pac/types';
import { AppModule } from '../app.module';
import { PrismaService } from '../prisma/prisma.service';
import { disableSchedulersForCli } from './bootstrap-flags';
interface Args {
host: string;
......@@ -77,6 +78,8 @@ function mulberry32(seed: number): () => number {
async function main(): Promise<void> {
const logger = new Logger('SeedAssignment');
const args = parseArgs(process.argv.slice(2));
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, { logger: ['warn', 'error'] });
const prisma = app.get(PrismaService);
......
......@@ -13,9 +13,12 @@ import { NestFactory } from '@nestjs/core';
import { Logger } from '@nestjs/common';
import { AppModule } from '../app.module';
import { StaleScanService } from '../queues/stale-scan.service';
import { disableSchedulersForCli } from './bootstrap-flags';
async function main(): Promise<void> {
const logger = new Logger('stale-scan-cli');
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, { logger });
try {
......
......@@ -22,6 +22,7 @@ import { PullOrchestrator } from '../modules/sync/pull/pull.orchestrator';
import { ReconcileOrchestrator } from '../modules/sync/reconcile/reconcile.orchestrator';
import { HmacVerifier } from '../modules/sync/push/hmac-verifier.service';
import { randomUUID, createHash } from 'node:crypto';
import { disableSchedulersForCli } from './bootstrap-flags';
interface CliArgs {
cmd: 'pull-setup' | 'pull' | 'reconcile' | 'push' | 'help';
......@@ -70,6 +71,8 @@ async function bootstrap() {
process.exit(0);
}
const logger = new Logger('sync:test');
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, { logger: ['warn', 'error'] });
try {
......
......@@ -15,6 +15,7 @@ import { AppModule } from '../app.module';
import { PrismaService } from '../prisma/prisma.service';
import { PatientService } from '../modules/patient/patient.service';
import type { TenantScopeContext } from '../common/decorators/tenant-scope.decorator';
import { disableSchedulersForCli } from './bootstrap-flags';
interface CliArgs {
pid?: string; // host externalId
......@@ -91,6 +92,8 @@ async function bootstrap() {
}
const logger = new Logger('timeline:cli');
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['warn', 'error'],
});
......
......@@ -6,6 +6,7 @@ import { NestFactory } from '@nestjs/core';
import { AppModule } from '../app.module';
import { PrismaService } from '../prisma/prisma.service';
import { ChainComposerService } from '../modules/plan/engine/chain-composer.service';
import { disableSchedulersForCli } from './bootstrap-flags';
async function main() {
const id = process.argv.find((a) => a.startsWith('--id='))?.slice('--id='.length);
......@@ -13,6 +14,8 @@ async function main() {
console.error('Usage: --id=<patientId>');
process.exit(1);
}
// ⚠️ 必须在建应用上下文之前:否则会清掉 service 里正在跑的同步锁(见 bootstrap-flags)
disableSchedulersForCli();
const app = await NestFactory.createApplicationContext(AppModule, { logger: ['error'] });
const prisma = app.get(PrismaService);
const composer = app.get(ChainComposerService);
......
import { Injectable } from '@nestjs/common';
import { Prisma } from '@prisma/client';
import { lookupDxTreatment, resolverCategoriesFor } from '@pac/types';
import { PrismaService } from '../../prisma/prisma.service';
import {
buildGapCore,
GAP_FLAGS_BY_PRIMARY,
GAP_PRIMARY_GROUPS,
gapVariant,
type GapVariant,
} from './potential-treatment-gap.sql';
/**
......@@ -31,8 +34,11 @@ export class PotentialTreatmentSelector {
patientId: string;
now: Date;
activeCodes: Set<string>;
/// 仅对拍工具用:强制 gap 计算形态。生产路径不传,走 gapVariant() 的环境开关。
variant?: GapVariant;
}): Promise<PotentialGap[]> {
const { hostId, tenantId, patientId, now, activeCodes } = opts;
const variant = opts.variant ?? gapVariant();
const out: PotentialGap[] = [];
for (const [primaryCode, group] of Object.entries(GAP_PRIMARY_GROUPS)) {
......@@ -43,21 +49,23 @@ export class PotentialTreatmentSelector {
if (!rule) continue;
const resolverCats = resolverCategoriesFor(primaryCode) as readonly string[];
const cfgFlags = GAP_FLAGS_BY_PRIMARY[primaryCode] ?? {};
const gap = buildGapCore({ rule, cfgFlags, allCodes, resolverCats });
const gap = buildGapCore({ rule, cfgFlags, allCodes, resolverCats, variant });
const rows = await this.prisma.$queryRaw<RawGapRow[]>`
SELECT
// 投影列(两形态共用;tooth 单列,取法不同)
const projection = Prisma.sql`
sig.id AS fact_id,
sig.content->>'code' AS code,
sig.content->>'name_zh' AS name_zh,
sig.type AS signal_type,
${gap.toothOutput} AS tooth,
sig.content->>'confidence' AS confidence,
EXTRACT(DAY FROM ${now}::timestamptz - COALESCE(sig.occurred_at, sig.planned_for))::int AS days_since,
COALESCE(sig.occurred_at, sig.planned_for) AS anchor_at
COALESCE(sig.occurred_at, sig.planned_for) AS anchor_at`;
// ⚠️ 画像是**逐患者**调用(全量 54.7 万次),这里的 scope 恒为 1 个患者 ——
// gap_scope 只有一行,各分支走 (patient_id, type, status) 索引,形态不会退化成全表扫。
const queryBody = (joinAddon: Prisma.Sql, gapAddon: Prisma.Sql): Prisma.Sql => Prisma.sql`
FROM patients p
JOIN patient_facts sig ON sig.patient_id = p.id
${gap.lateralJoin}
${joinAddon}
WHERE p.host_id = ${hostId}::uuid
AND p.tenant_id = ${tenantId}
AND p.id = ${patientId}::uuid
......@@ -68,8 +76,26 @@ export class PotentialTreatmentSelector {
AND COALESCE(sig.occurred_at, sig.planned_for) IS NOT NULL
${gap.restorationIneligibleFrag}
${gap.congenitalFrag}
${gap.gapWhere}
`;
${gapAddon}`;
const sb = gap.setBased;
const sql = sb
? Prisma.sql`
WITH gap_cand AS MATERIALIZED (
SELECT ${projection}${sb.candExtraCols}
${queryBody(Prisma.empty, sb.candWhere)}
)${sb.postCtes}
SELECT c.fact_id, c.code, c.name_zh, c.signal_type, c.confidence, c.days_since, c.anchor_at,
${sb.toothOutput} AS tooth
FROM gap_cand c
${sb.remJoin}
WHERE TRUE ${sb.outerWhere}`
: Prisma.sql`
SELECT ${projection},
${gap.toothOutput} AS tooth
${queryBody(gap.lateralJoin, gap.gapWhere)}`;
const rows = await this.prisma.$queryRaw<RawGapRow[]>(sql);
for (const r of rows) {
out.push({
primaryCode,
......
......@@ -142,8 +142,12 @@ export function temperatureBucketCaseSql(
* 关联条件变成恒真且 SQL 不报错)—— 那种场景用下方的 `labelExistsSql` / `labelTemperatureExistsSql`。
*
* ⚠️ `f.status = 'active'`:治完的诊断是 `fulfilled` 不是删除。拿它当证据 =
* 对着一个已经做完的诊断说"您还没做"。当前数据 5,034 条证据全是 active,
* 这条过滤是**不变量守卫** —— 治疗落库后引擎还没重算的窗口期里,它就是唯一防线。
* 对着一个已经做完的诊断说"您还没做"。这条过滤是**不变量守卫** —— 治疗落库后
* 引擎还没重算的窗口期里,它就是唯一防线。
* 2026-08-28 生产实测(全量,不再是早期 5,034 条的小样本):在跑计划的 1,092,580 条
* active 理由 + 1,265 条 assigned + 41 条 completed,**无一条**取不到 active 证据;
* 只有 abandoned 里 6/282 取不到。即引擎重算追得上事实换代,守卫没有在挡真数据。
* ⚠️ 这个数掉下来 = 重算落后于摄入,格位会静默变空(不报错),值得当信号看。
* 🔴 `anc.at IS NOT NULL` **必须留着**,即便锚点已经换成末诊、不再用它定档:
* 它是「这条证据有日期」的守卫,决定**哪些 plan 进得来**。删掉行集就变了 ——
* 而多出来的那些人不会报错,只会悄悄出现在格子里。
......
......@@ -117,6 +117,7 @@ export class MockPullStrategy implements PullStrategy {
occurredAt: baseTime.toISOString(),
status: 'scheduled',
doctorId: `mock-d${(i % 3) + 1}`,
doctorName: `模拟医生${(i % 3) + 1}`,
treatmentCategory: '复查',
});
}
......
......@@ -185,10 +185,22 @@ export class ColdImportService {
// 2a. dryRun:报 scope(患者数 + 各资源 txn 数),不写、不分批装配(便宜)。
if (opts.dryRun) {
// ⛔ 同样分块 —— 早先这里也是整份清单塞 in,>3.2 万患者时 dry-run 直接崩,
// 而 --patients-file 的文档恰恰说「按受影响患者收窄是最有效的提速手段(可达 250 倍)」:
// 最需要先 dry-run 探一探的大清单场景,正好是它唯一不工作的场景。
const DRY_CHUNK = 3000;
for (const cfg of reparseableCfgs) {
const n = await this.prisma.patientTransaction.count({
where: { hostId: host.id, subjectType: cfg.emits!.subjectType, ...(opts.patientIds?.length ? { patientId: { in: opts.patientIds } } : {}) },
});
let n = 0;
const chunks: Array<string[] | null> = opts.patientIds?.length
? Array.from({ length: Math.ceil(opts.patientIds.length / DRY_CHUNK) }, (_, i) =>
opts.patientIds!.slice(i * DRY_CHUNK, (i + 1) * DRY_CHUNK),
)
: [null];
for (const chunk of chunks) {
n += await this.prisma.patientTransaction.count({
where: { hostId: host.id, subjectType: cfg.emits!.subjectType, ...(chunk ? { patientId: { in: chunk } } : {}) },
});
}
this.logger.log(`reparse[dry]: ${cfg.canonical}(${cfg.emits!.subjectType}) txns=${n} 实跑按版本流 supersede 变更的、跳过不变的`);
}
this.logger.log(`reparse[dry]: 范围 ${scopePatientIds.length} 患者;去掉 --dry-run 实跑(非破坏)`);
......@@ -276,16 +288,34 @@ export class ColdImportService {
}
// 3. 受影响 patientId = 本次真正被 supersede(内容变了)的 fact 的 distinct patient → 只重算这些。
const changed = await this.prisma.patientFact.findMany({
where: {
hostId: host.id,
supersededAt: { gte: runStart },
...(opts.patientIds?.length ? { patientId: { in: opts.patientIds } } : {}),
},
select: { patientId: true },
distinct: ['patientId'],
});
const affectedPatientIds = changed.map((a) => a.patientId).filter((x): x is string => !!x);
//
// ⛔ 必须**分块**查 —— `patientId: { in: [...] }` 直接塞完整清单会撞 PG 的
// 32767 bind 变量上限。2026-08-29 生产实测:18 万患者的 reparse 跑满 61/61 批、
// 写完全部事实之后,**倒在这最后一步**:
// `Assertion violation: too many bind variables ... received 32769`
// 6.6 小时的活全干完了,只因收尾统计炸掉而 exit 1 —— 最难受的一种失败。
// (同族的另一处在 dryRun 分支的 count,见下方注释。)
// 分块大小取 BATCH 同款 3000:每块 3001 个变量,离上限很远。
const CHANGED_CHUNK = 3000;
const affectedSet = new Set<string>();
const scanChunks: Array<string[] | null> = opts.patientIds?.length
? Array.from({ length: Math.ceil(opts.patientIds.length / CHANGED_CHUNK) }, (_, i) =>
opts.patientIds!.slice(i * CHANGED_CHUNK, (i + 1) * CHANGED_CHUNK),
)
: [null]; // 不限定患者 → 一次全查(where 里没有 in,无变量上限问题)
for (const chunk of scanChunks) {
const changed = await this.prisma.patientFact.findMany({
where: {
hostId: host.id,
supersededAt: { gte: runStart },
...(chunk ? { patientId: { in: chunk } } : {}),
},
select: { patientId: true },
distinct: ['patientId'],
});
for (const c of changed) if (c.patientId) affectedSet.add(c.patientId);
}
const affectedPatientIds = [...affectedSet];
return { perResource, affectedPatientIds, dryRunDiffs };
}
......@@ -2476,16 +2506,24 @@ export function traceRawSourceTable(primaryTable: string, transforms: ReadonlyAr
input?: string;
output?: string;
inputs?: string[];
outputs?: Array<{ output?: string }>;
/// ⚠️ route_by_pattern 的字段名是 `routes`(见 transforms.schema.ts RouteByPatternOpSchema),
/// 不是 `outputs` —— 早先这里写成 outputs,恒为 undefined,导致**凡链路经过 route 的资源
/// 回溯都停在中间表**(_treatment_actual_raw_emr 等)。后果:reparse 把 rawPayload 灌进中间表,
/// 随即被 transform 链用空结果覆盖 → 治疗类 reparse 恒 0 变更**且报成功**(exit 0)。
/// diagnosis 没踩到只因它的链 split→derive→derive 不过 route。
routes?: Array<{ output?: string }>;
};
if (t.kind === 'union' && t.output && Array.isArray(t.inputs)) {
unionInputs.set(t.output, t.inputs);
continue;
}
if (t.output && t.input) byOutput.set(t.output, t.input);
// ⚠️ 跳过**原地 derive**(output === input,如 `_treat_plan_raw → _treat_plan_raw` 补 treat_name):
// 它不改变这张表的来源。登记进去会让 byOutput 指向自己 → resolve 撞环保护当场返回,
// 回溯同样停在中间表(与 routes 那条是**两个独立 bug**,只修一个仍然失效)。
if (t.output && t.input && t.output !== t.input) byOutput.set(t.output, t.input);
// route_by_pattern 多 output:每个 output 都回到同一 input
if (Array.isArray(t.outputs) && t.input) {
for (const o of t.outputs) if (o?.output) byOutput.set(o.output, t.input);
if (Array.isArray(t.routes) && t.input) {
for (const o of t.routes) if (o?.output) byOutput.set(o.output, t.input);
}
}
const resolve = (tbl: string, seen: Set<string>): string => {
......
......@@ -262,6 +262,14 @@ const AppointmentRecordContent = z
arrived_at: isoDateString.nullable().optional().default(null),
appointment_type: nullableString(),
doctor_id: nullableString(),
// 约号时约的医生姓名(源自 jvs-dw resource_name;核实记录见 appointment.yaml)。
// ⚠️ 口径 = **约的**医生,不是实际接诊医生(改派时不同,实际接诊看 emr_record.doctor_name)。
// ⚠️ 已过「资源是不是人」这道闸(AppointmentParser.isPersonResource):排房间/服务/台席的
// 那 0.94% 在这里是 null,原值仍在 resource_name。所以本字段可以直接拼「X医生」。
doctor_name: nullableString(),
// 排班资源原名(未过滤)。多数等于 doctor_name;不是人时("预约"/"学前街手术室"/"正畸咨询")
// doctor_name 为 null 而本字段保留原值 —— 既不丢 host 快照,也便于审计上面那道闸。
resource_name: nullableString(),
// 预约科目/就诊意向(常规/正畸/种植/修复/拔牙/牙周…),host appo_complaint_category
complaint_category: nullableString(),
// 预约主诉自由文本(跟 complaint_category 配对:分类 + 原文)— Layer C 源
......
......@@ -44,6 +44,12 @@ export class AppointmentParser implements Parser {
const arrivedAt = c.arrivedAt ? new Date(c.arrivedAt as string) : null;
const appointmentType = (c.appointmentType as string | undefined) ?? null;
const doctorId = (c.doctorId as string | undefined) ?? null;
// 排班资源名(host 行内快照,DW resource_name)。约 99% 是医生本人,详见下方 isPersonResource。
const resourceName = (c.doctorName as string | undefined)?.trim() || null;
// 约号时约的医生姓名。见 appointment.yaml 里的核实记录。
// ⚠️ 不是实际接诊医生 —— 改派时两者不同,实际接诊以病历 doctor_name 为准。
const doctorName =
resourceName && AppointmentParser.isPersonResource(resourceName) ? resourceName : null;
return [
{
......@@ -77,6 +83,8 @@ export class AppointmentParser implements Parser {
arrived_at: arrivedAt ? arrivedAt.toISOString() : null,
appointment_type: appointmentType,
doctor_id: doctorId,
doctor_name: doctorName,
resource_name: resourceName,
complaint_category: (c.complaintCategory as string | undefined) ?? null,
complaint_text: (c.complaintText as string | undefined) ?? null,
duration_minutes:
......@@ -93,6 +101,42 @@ export class AppointmentParser implements Parser {
];
}
/**
* 排班资源是不是**一个人**。
*
* DW 的 `resource_name`(→ canonical doctorName)是「排班资源」的名字,绝大多数资源就是医生本人,
* 但少数排的是房间 / 服务 / 台席。生产近 180 天 1,047,901 条预约里实测 9,838 条(**0.94%**)不是人:
* 预约(5,696)· 学前街手术室(1,218)· 种植手术(天使)(396)· 会诊室(351)· 显微镜管理(328)
* 舒适治疗(罗院/徐院/梁博)(738)· 种植室(欢乐)(260)· 方寸诊所客服(243)· 正畸咨询(205)
* 种植手术(金融街)(158)· 华侨城/三星手术室(141)· 公共诊室(49+19)· 种植手术(30)
* 华贸诊所B诊区客服(4)· 主诉初诊(2)
* 不拦的话,时间轴上每 106 条预约就有 1 条写着「预约医生」「学前街手术室医生」,
* 更糟的是详情页「主治医生」的最高频兜底可能解析成「学前街手术室」。
*
* 词表是**从上面这份实测清单反推**出来的(九个词覆盖全部 19 个取值),不是凭感觉列的。
* 误伤核验(生产 2026-08-31):临床事实里 1,314 个 distinct doctor_name 过一遍本规则,
* 命中 **1** 个 —— 「公共诊室」,它本身就是漏进病历的房间名,不是医生。**真人零误伤**。
*
* ⚠️ 被拦下的原值不丢:仍完整写在 content.resource_name(以及 raw_payload)里,
* `WHERE resource_name IS NOT NULL AND doctor_name IS NULL` 就能审计本规则拦了什么。
* ⚠️ 新诊所接入 / host 改排班命名后要复查这份词表 —— 用上面那条审计 SQL 看有没有新形态漏网。
*/
private static readonly NON_PERSON_TOKENS = [
'预约',
'手术',
'室',
'治疗',
'咨询',
'客服',
'管理',
'诊区',
'初诊',
];
static isPersonResource(name: string): boolean {
return !AppointmentParser.NON_PERSON_TOKENS.some((t) => name.includes(t));
}
// 已发生事实(actual)的 host 状态:到诊 / 接诊 / 结算 + walk-in 变更(已到店)
private static readonly ACTUAL_STATUSES = new Set([
'arrived',
......
......@@ -37,10 +37,34 @@ import { schedulerDisabled } from './scheduler-switch';
*
* 跑失败:cursor 不前进 → 下次自动 catchup;log ERROR 不抛。
*/
/**
* 僵尸锁的最小年龄。比「一轮同步的正常耗时」留足余量 —— 生产实测单轮摄入 28~52 分钟,
* 取 3 小时:真崩溃留下的锁必然远超此值,而正常在跑的绝不会。
*/
const REAP_MIN_AGE_MS = 3 * 60 * 60 * 1000;
@Injectable()
export class SyncIncrementalSchedulerService implements OnModuleInit {
private readonly logger = new Logger(SyncIncrementalSchedulerService.name);
/**
* 本进程内「该 host 正在跑」的闸 —— **防套圈**。
*
* 🔴 2026-08-29 生产事故:plan 段耗时涨到 2 小时以上后,cron(每 2 小时)照常触发下一轮,
* 两轮的 plan 段并发抢同一批 I/O → 两轮都更慢 → 更容易被再下一轮套圈 → 雪崩。
* 实测 08-28 20:15 起连续多轮 plan 段一次都没跑完,直到 08-29 上午仍有两轮在并行。
*
* ⚠️ 为什么现有的锁挡不住:
* ① NestJS 的 CronJob **默认不防重入** —— 上一次回调还在 await,下一次照样进;
* ② `sync_logs` 的 partial UNIQUE(host_id) WHERE status='running' 只覆盖**摄入段**,
* 摄入一结束锁就放了,而 persona / plan 段还在跑,恰恰是最慢的部分。
* 所以必须在**回调入口**挡,不能靠库里的锁。
*
* 跳过而不是排队:摄入是游标增量,跳过这轮的数据下轮自然 catchup;
* plan 是时间驱动的全量,跳一轮只是晚 2 小时评估,远好过雪崩。
*/
private readonly runningHosts = new Set<string>();
constructor(
private readonly prisma: PrismaService,
private readonly coldImport: ColdImportService,
......@@ -126,11 +150,21 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
/**
* 回收僵尸同步锁 —— 把 startedAt 早于本进程启动的 running sync_log 标 failed。
*
* 为什么这样判据安全:sync 只在两处跑 —— 本 service 进程的 cron 回调,或一次性 CLI
* (`sync-incremental.cli` / `cold-import.cli`,跑完即退)。两者的进程都不可能比本进程
* 启动得更早还活着。所以 `startedAt < PROCESS_STARTED_AT` 的 running 行 = 上一个已死进程的残留。
* 反过来,本进程启动后新建的 running 行(startedAt >= PROCESS_STARTED_AT)绝不碰 —— 那可能是
* 正在跑的真锁(例如运维手动触发的 CLI 与本进程并存),误清会把在跑的同步" orphan"掉。
* 🔴 2026-08-30 生产事故:本判据**曾经是错的**,理由写反了方向。
* 原注释说「sync 只在 service 进程的 cron 或一次性 CLI 里跑,两者的进程都不可能比本进程
* 启动得更早还活着」——【但长驻的 pac-service 恰恰就是「启动得更早还活着」的那个】。
* 任何 CLI(recompute-plans / recompute-persona / reparse …)都会
* `createApplicationContext(AppModule)`,于是也跑一遍本 onModuleInit;
* 此时 CLI 进程的 PROCESS_STARTED_AT = 现在,而 service 里**正在跑**的那轮 sync
* startedAt 更早 → 被当成僵尸锁清掉 → 那一轮增量当场夭折。
* 实测:08:17:11 在生产容器里跑 recompute-plans,08:17:13(Nest 启动 2 秒后)
* 08:15 那轮正常同步就被标 failed。前 19 轮全 success,只死了撞上的这一轮。
* 数据没丢(cursor_after=null,下轮按同一水位 catchup),但白丢一轮、晚 2 小时落库。
*
* 现在的判据:`startedAt < 本进程启动` **且** `startedAt < now - REAP_MIN_AGE_MS`。
* 后半条是真正的防线 —— 真僵尸锁是上个进程崩溃留下的,必然已经躺了很久;
* 而被误伤的那种,是"刚起没多久还在正常跑"的。用年龄区分,不靠进程身份猜。
* 另外 CLI 侧统一设 PAC_SCHEDULER_DISABLED=1(见 src/cli/bootstrap-flags.ts),双保险。
*
* 幂等:被标 failed 只是让并发锁释放;数据侧不受影响(游标没推进,下次增量靠 48h 回看窗补齐)。
*
......@@ -139,14 +173,17 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
*/
private async reapStaleRunningLocks(processStartedAt: Date = PROCESS_STARTED_AT): Promise<void> {
try {
// 双条件取更早的那个界:既要早于本进程启动,又要已经躺够 REAP_MIN_AGE_MS。
const ageCutoff = new Date(Date.now() - REAP_MIN_AGE_MS);
const cutoff = ageCutoff < processStartedAt ? ageCutoff : processStartedAt;
const stale = await this.prisma.syncLog.findMany({
where: { status: SyncStatus.RUNNING, startedAt: { lt: processStartedAt } },
where: { status: SyncStatus.RUNNING, startedAt: { lt: cutoff } },
select: { id: true, hostId: true, startedAt: true, triggeredBy: true },
});
if (stale.length === 0) return;
const { count } = await this.prisma.syncLog.updateMany({
where: { status: SyncStatus.RUNNING, startedAt: { lt: processStartedAt } },
where: { status: SyncStatus.RUNNING, startedAt: { lt: cutoff } },
data: {
status: SyncStatus.FAILED,
endedAt: new Date(),
......@@ -168,6 +205,15 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
/// 单 host 跑一轮(cron 回调用,吞异常不影响该 host 下次 / 别的 host)
private async runHostSafe(host: string): Promise<void> {
// ⛔ 上一轮还没跑完就跳过本轮 —— 见 runningHosts 的注释(防套圈雪崩)
if (this.runningHosts.has(host)) {
this.logger.warn(
`sync-incremental: host=${host} **跳过本轮** —— 上一轮仍在运行(防套圈)。` +
`连续出现说明单轮已撑不下 cron 间隔,需要查 plan 段耗时。`,
);
return;
}
this.runningHosts.add(host);
try {
await this.runOne(path.join(this.dataDir(), host));
} catch (err) {
......@@ -176,6 +222,9 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
} else {
this.logger.error(`sync-incremental: host=${host} failed: ${(err as Error).message}`);
}
} finally {
// ⚠️ 必须在 finally —— 抛异常时不释放会把该 host 永久锁死到进程重启
this.runningHosts.delete(host);
}
}
......
/**
* 预约医生姓名(jvs-dw resource_name → content.doctor_name)。
*
* 背景:预约事实历来只有 doctor_id 没有姓名,前端/话术拿到的是裸 ID。
* DW 其实一直在同一行里给了姓名,列名叫 `resource_name`(排班资源名)——名字有误导性,
* 核实过程与判据记在 data/jvs-dw/assemblers/appointment.yaml 的注释里。
*
* 本文件锁三件事:
* ① yaml 确实把 resource_name 映到 canonical doctorName(改名/删映射会红);
* ② parser 把它写进 content.doctor_name,空白归一成 null;
* ③ 口径不漂移 —— doctor_name 与 doctor_id 是**同一个人**的两面,别一个来自预约、
* 另一个来自别处。
*/
import { readFileSync } from 'node:fs';
import { join } from 'node:path';
import * as yaml from 'js-yaml';
import { Action } from '@pac/types';
import { AppointmentParser } from '../src/modules/sync/pipeline/parsers/appointment.parser';
import type { ParserContext } from '../src/modules/sync/pipeline/parsers/parser.interface';
const YAML_PATH = join(
__dirname,
'../data/jvs-dw/assemblers/appointment.yaml',
);
describe('appointment.yaml | resource_name → doctorName', () => {
const cfg = yaml.load(readFileSync(YAML_PATH, 'utf-8')) as {
field_mapping: Record<string, string>;
};
test('doctorName 映射到 host 列 resource_name', () => {
expect(cfg.field_mapping.doctorName).toBe('resource_name');
});
// ⚠️ 这两列是**不同的人**,历史上极易混:
// appo_doc_id = 约的医生(id) ←→ resource_name 是它的姓名
// director_id = 另一个角色(仅 34% 有值,与 appo_doc_id 几乎从不相同)
// create_name / appo_handler = 呼叫中心约号人(值形如 "400-叶玉娇")
test('doctorId 仍取 appo_doc_id,没有被 director/handler 之类顶替', () => {
expect(cfg.field_mapping.doctorId).toBe('appo_doc_id');
});
});
describe('AppointmentParser | content.doctor_name', () => {
const parser = new AppointmentParser();
const ctx = (row: Record<string, unknown>): ParserContext => ({
transaction: {
id: 'tx-1',
hostId: 'host-1',
tenantId: 'tenant-1',
patientId: 'p-1',
action: Action.APPOINTMENT_CREATED,
subjectType: 'appointment',
subjectId: 'appt-1',
occurredAt: new Date('2026-08-01T02:00:00Z'),
clinicId: 'c-1',
},
canonicalRow: {
externalId: 'appt-1',
patientExternalId: 'p-1',
clinicId: 'c-1',
scheduledAt: '2026-08-01T02:00:00Z',
status: 'scheduled',
doctorId: '5624',
...row,
},
});
const contentOf = (row: Record<string, unknown>) =>
parser.parse(ctx(row))[0]!.content as Record<string, unknown>;
test('有姓名 → 写进 content.doctor_name,与 doctor_id 成对', () => {
const c = contentOf({ doctorName: '赵茜' });
expect(c.doctor_name).toBe('赵茜');
expect(c.doctor_id).toBe('5624');
});
test('两侧空白被裁掉(host 常见尾随空格)', () => {
expect(contentOf({ doctorName: ' 李闻 ' }).doctor_name).toBe('李闻');
});
test.each([
['缺字段', {}],
['空串', { doctorName: '' }],
['纯空白', { doctorName: ' ' }],
])('%s → null(不写空串,免得下游把 "" 当姓名渲染成「医生」)', (_label, row) => {
expect(contentOf(row).doctor_name).toBeNull();
});
test('resource_name 始终保留原值(host 快照不丢)', () => {
expect(contentOf({ doctorName: '赵茜' }).resource_name).toBe('赵茜');
expect(contentOf({ doctorName: '学前街手术室' }).resource_name).toBe('学前街手术室');
});
});
/**
* 「排班资源是不是人」这道闸。
*
* 生产近 180 天 0.94% 的预约排的是房间/服务/台席而非医生 —— 不拦的话时间轴上每 106 条
* 就有一条写着「预约医生」「学前街手术室医生」。词表是从那份实测清单反推的,不是凭感觉列的,
* 所以下面逐个锁住:清单变了(新诊所换了命名)这里就该红。
*/
describe('AppointmentParser.isPersonResource | 非人资源闸', () => {
const parser = new AppointmentParser();
const doctorNameOf = (nm: string) =>
(parser.parse({
transaction: {
id: 'tx-1', hostId: 'h', tenantId: 't', patientId: 'p',
action: Action.APPOINTMENT_CREATED, subjectType: 'appointment', subjectId: 'a',
occurredAt: new Date('2026-08-01T02:00:00Z'), clinicId: 'c',
},
canonicalRow: {
externalId: 'a', patientExternalId: 'p', clinicId: 'c',
scheduledAt: '2026-08-01T02:00:00Z', status: 'scheduled', doctorName: nm,
},
})[0]!.content as Record<string, unknown>).doctor_name;
// 生产实测出现过的 19 个非人取值,逐个锁
test.each([
'预约',
'学前街手术室',
'种植手术(天使)',
'会诊室',
'显微镜管理',
'舒适治疗(罗院)',
'种植室(欢乐)',
'方寸诊所客服',
'正畸咨询',
'种植手术(金融街)',
'华侨城手术室',
'公共诊室',
'三星手术室',
'顺义瑞捷公共诊室',
'华贸诊所B诊区客服',
'主诉初诊',
])('%s → 不是人,doctor_name 留空', (nm) => {
expect(AppointmentParser.isPersonResource(nm)).toBe(false);
expect(doctorNameOf(nm)).toBeNull();
});
// 真人零误伤 —— 含生产里那些容易被规则误伤的形态:
// 带消歧后缀的(吕晓玉(Y))、四字名(孙汪心悦)、外籍全名(SABELLI MARIA LUZ)
test.each([
'赵茜',
'李闻',
'刘柳',
'吕晓玉(Y)',
'姚志远(Z)',
'孙汪心悦',
'SABELLI MARIA LUZ',
'TAN CHYI YANN',
'薛玫(L)',
])('%s → 是人,doctor_name 照写', (nm) => {
expect(AppointmentParser.isPersonResource(nm)).toBe(true);
expect(doctorNameOf(nm)).toBe(nm);
});
});
/**
* 「活动义齿 = 假缺失」按颌判定。
*
* 戴活动义齿的患者诊断栏**永远**写「缺失牙」—— 天然牙确实没了,诊断没错;
* 但那个位置已被义齿盖住,「缺失牙未**启动**修复」不成立(治疗做的就是义齿)。
* 活动义齿是整颌跨度的修复体 → 判出"这一颌有义齿",那一颌的缺牙位全部解除。
*
* 🔴 本条**天然绕开混合句** —— 王志荣/童然夫的检查所见牙位横跨上下颌,
* 按条目判的那条(RESTORATION_IN_PLACE_*)整条跳过,按颌判直接可用。
*
* 跑:
* pnpm test -- arch-denture-false-missing
*/
import {
ARCH_DENTURE_UPPER_RE,
ARCH_DENTURE_LOWER_RE,
ARCH_DENTURE_INTENT_EXCLUDE_RE,
ARCH_DENTURE_PREFILTER_RE,
UPPER_ARCH_FIRST_DIGITS_RE,
LOWER_ARCH_FIRST_DIGITS_RE,
} from '@pac/types';
import { GAP_FLAGS_BY_PRIMARY } from '../src/modules/clinical-gap/potential-treatment-gap.sql';
const UP = new RegExp(ARCH_DENTURE_UPPER_RE);
const LOW = new RegExp(ARCH_DENTURE_LOWER_RE);
const INTENT = new RegExp(ARCH_DENTURE_INTENT_EXCLUDE_RE);
const PRE = new RegExp(ARCH_DENTURE_PREFILTER_RE);
/** 模拟:该句判出哪几颌有义齿(未发生词一票否决) */
const arches = (msg: string): string[] => {
if (INTENT.test(msg)) return [];
const out: string[] = [];
if (UP.test(msg)) out.push('上');
if (LOW.test(msg)) out.push('下');
return out;
};
describe('🔴 两条待定个案 —— 检查所见跨上下颌,按条目判会整条跳过,按颌判可用', () => {
it('王志荣 TS0K051842:上下颌都有义齿(召回下颌 → 解)', () => {
expect(
arches('上下颌吸附性义齿修复,下颌固位可,上颌固位稍差,说话、喝水时义齿易脱落,牙槽嵴黏膜萎缩,无疼痛不适'),
).toEqual(['上', '下']);
});
it('童然夫 TS0M013276:上下颌都有义齿(召回上颌 → 解)', () => {
expect(
arches('牙缺失,口内活动义齿修复,上颌义齿卡环紧,不易取戴,下颌义齿无法完全就位(三个多月未佩戴)'),
).toEqual(['上', '下']);
});
});
describe('按颌判定 —— 生产真实句子', () => {
it.each([
['缺失,上颌活动义齿修复', ['上']],
['缺失,下颌活动义齿修复', ['下']],
['上颌义齿卡环折断', ['上']], // ⭐ 状态不好仍算已治疗:义齿存在 = 修复启动过
['下颌义齿压痛', ['下']],
['上颌义齿固位欠佳', ['上']],
['上下颌活动义齿修复', ['上', '下']],
['上下颌全口义齿,咬合接触均匀。牙龈色粉,无溃疡。', ['上', '下']],
['U全口义齿固位不良,咬合关系欠佳', ['上', '下']],
['上下颌种植临时义齿存,牙龈未见异常', ['上', '下']],
['右侧下颌义齿舌侧粘膜有压痕', ['下']],
])('%s → %s', (msg, expected) => {
expect(arches(msg)).toEqual(expected);
});
});
describe('⛔ 不作数', () => {
it.each([
'建议上颌活动义齿修复', // 还没做
'患者要求下颌义齿修复',
'拟行上下颌全口义齿修复',
'缺牙区粘膜无异常,牙槽嵴有吸收', // 压根没提义齿
])('%s', (msg) => {
expect(arches(msg)).toEqual([]);
});
it('🔴 颌词与义齿词距离过远不得相连', () => {
// 「上颌」讲的是残根,「义齿」是下次的打算 —— 中间隔了十几个字
expect(arches('上颌见残根,牙龈红肿,牙槽嵴吸收明显,下次考虑做义齿')).toEqual([]);
});
});
describe('🔴 牙位由条目自己给,不整颌铺开', () => {
const UPD = new RegExp(UPPER_ARCH_FIRST_DIGITS_RE);
const LOWD = new RegExp(LOWER_ARCH_FIRST_DIGITS_RE);
/** 模拟 SQL:条目牙位 ∩ 句中点名有义齿的那一颌 */
const resolved = (msg: string, teeth: string[]): string[] => {
if (INTENT.test(msg)) return [];
return teeth.filter((t) => (UP.test(msg) && UPD.test(t)) || (LOW.test(msg) && LOWD.test(t)));
};
it('童然夫:条目 17 颗跨颌,上颌句 → 只解上颌那 10 颗', () => {
const teeth = '13;16;17;21;22;23;24;25;26;27;41;42;45;46;47;31;32'.split(';');
expect(resolved('牙缺失,口内活动义齿修复,上颌义齿卡环紧,不易取戴', teeth)).toEqual(
['13', '16', '17', '21', '22', '23', '24', '25', '26', '27'],
);
});
it('🔴 只提上颌 → 下颌牙位不得被解(义齿没覆盖到的那一颌)', () => {
expect(resolved('上颌活动义齿修复', ['14', '15', '16', '36', '37'])).toEqual(['14', '15', '16']);
});
it('🔴 ⛔ 不得铺到条目之外 —— 同颌里义齿没补的缺牙位必须仍能召回', () => {
// 上颌局部义齿只补了 14;15;16;24 也缺着但不在这条记录里 → 24 不该被解
expect(resolved('上颌活动义齿修复', ['14', '15', '16'])).not.toContain('24');
});
it('上下颌句 → 两边都解,但仍限条目内', () => {
expect(resolved('上下颌活动义齿修复', ['15', '16', '36'])).toEqual(['15', '16', '36']);
});
});
describe('表自身自洽', () => {
it('颌 → 牙位首位:上颌 1/2,下颌 3/4', () => {
expect(new RegExp(UPPER_ARCH_FIRST_DIGITS_RE).test('16')).toBe(true);
expect(new RegExp(UPPER_ARCH_FIRST_DIGITS_RE).test('36')).toBe(false);
expect(new RegExp(LOWER_ARCH_FIRST_DIGITS_RE).test('46')).toBe(true);
expect(new RegExp(LOWER_ARCH_FIRST_DIGITS_RE).test('26')).toBe(false);
});
it('🔴 只对缺失牙(K08)开闸', () => {
const on = Object.entries(GAP_FLAGS_BY_PRIMARY)
.filter(([, f]) => f.archDentureIsRestored === true)
.map(([code]) => code);
expect(on).toEqual(['K08']);
});
});
/**
* ⚡ SQL 侧的预过滤(potential-treatment-gap.sql:archDentureBranch)在展开 JSON 数组之前
* 先用 ARCH_DENTURE_PREFILTER_RE 把整份 exam_findings 筛一道。那是纯剪枝,**前提是
* 该词必须始终是上下颌两条 RE 的必要条件** —— 一旦不是,预过滤就会静默丢掉真命中
* (少召,不报错)。这组用例就是锁这个蕴含关系的。
*
* 实测依据(2026-08-29 测试库):exam_findings 是数组的 emr 事实 1,519,829 份,
* 含「义齿|假牙」的仅 7,396 份(0.49%),预过滤剪掉 99.5% 的无用展开。
*/
describe('⚡ 预过滤词必须是两条 RE 的必要条件(改 RE 时这组会先炸)', () => {
const 会命中的句子 = [
'上下颌吸附性义齿修复,下颌固位可,上颌固位稍差',
'牙缺失,口内活动义齿修复,上颌义齿卡环紧,不易取戴,下颌义齿无法完全就位',
'上颌活动义齿在位',
'下半口假牙尚可',
'全口义齿修复',
'上下全口假牙使用中',
'上颌见活动义齿,基托边缘密合',
];
it.each(会命中的句子)('命中 RE 的句子必然含预过滤词:%s', (msg) => {
expect(UP.test(msg) || LOW.test(msg)).toBe(true); // 前提:这些确实命中
expect(PRE.test(msg)).toBe(true); // 结论:那就一定含预过滤词
});
it('⛔ 反向:不含预过滤词的文本,两条 RE 都不可能命中(否则预过滤会丢真命中)', () => {
for (const msg of [
'上颌见残根,下颌牙列完整',
'全口牙石(+),牙龈红肿',
'上颌种植体冠在位,边缘密合', // 修复体在位但非义齿 → 归按条目那条分支管
'下颌固定桥完好',
]) {
expect(PRE.test(msg)).toBe(false);
expect(UP.test(msg)).toBe(false);
expect(LOW.test(msg)).toBe(false);
}
});
});
/**
* CLI 启动不得清掉正在跑的同步锁(2026-08-30 生产事故回归测试)
*
* 事故:在生产容器里 `docker exec` 跑 recompute-plans,CLI 会
* `createApplicationContext(AppModule)` → 跑一遍 SyncIncrementalScheduler.onModuleInit
* → reapStaleRunningLocks 把 service 里**正在跑**的那轮同步当僵尸锁清掉,那轮增量夭折。
* 实测 08:17:11 起 CLI,08:17:13 正常跑着的 08:15 那轮被标 failed(前 19 轮全 success)。
*
* 两道防线,这里各锁一条。
*/
import * as fs from 'node:fs';
import * as path from 'node:path';
const CLI_DIR = path.resolve(__dirname, '../src/cli');
/// 唯一豁免:它本身就是要触发同步(靠 scheduler 的年龄阈值兜底)
const EXEMPT = new Set(['sync-incremental.cli.ts']);
describe('CLI 调度器总闸', () => {
const cliFiles = fs
.readdirSync(CLI_DIR)
.filter((f) => f.endsWith('.cli.ts'))
.filter((f) => fs.readFileSync(path.join(CLI_DIR, f), 'utf8').includes('createApplicationContext'));
it('存在会启动完整应用上下文的 CLI(否则本 spec 形同虚设)', () => {
expect(cliFiles.length).toBeGreaterThan(5);
});
it.each(cliFiles.filter((f) => !EXEMPT.has(f)))(
'%s 在建上下文之前调用 disableSchedulersForCli()',
(file) => {
const src = fs.readFileSync(path.join(CLI_DIR, file), 'utf8');
const guard = src.indexOf('disableSchedulersForCli()');
const ctx = src.indexOf('NestFactory.createApplicationContext');
expect(guard).toBeGreaterThanOrEqual(0);
// 顺序也要对:晚于建上下文就没意义了(onModuleInit 已经跑完)
expect(guard).toBeLessThan(ctx);
},
);
it('僵尸锁回收带年龄阈值 —— 不能只靠"比本进程早"这一个判据', () => {
const sched = fs.readFileSync(
path.resolve(__dirname, '../src/queues/sync-incremental.scheduler.ts'),
'utf8',
);
expect(sched).toContain('REAP_MIN_AGE_MS');
// 阈值必须显著大于一轮同步的正常耗时(生产实测 28~52 分钟)
const m = sched.match(/const REAP_MIN_AGE_MS = ([^;]+);/);
expect(m).not.toBeNull();
// eslint-disable-next-line no-eval
const ms = eval(m![1]!) as number;
expect(ms).toBeGreaterThanOrEqual(2 * 60 * 60 * 1000);
});
});
/**
* gap 集合式形态的**结构对拍**(纯 SQL 文本层,不连库)
*
* 定位:数据层的等价性由 `pnpm verify-gap-equivalence`(逐 患者×信号×牙位 差分)证明,
* 本 spec 只守一件单元测试能守住的事 —— **两种形态的分支集合不许走散**。
* 典型事故:后来人给 legacy 加了第 14 条 resolved 分支,忘了同步 setbased →
* 线上悄悄少销一类证据 → 静默多召 / 少召。那种漏法 tsc 和现有 spec 全都发现不了,
* 但分支计数会当场炸。
*
* ⚠️ 本 spec 断言的是"两边都改了",不是"改对了"。改完仍必须跑 verify-gap-equivalence。
*/
import { lookupDxTreatment, resolverCategoriesFor } from '@pac/types';
import {
buildGapCore,
GAP_FLAGS_BY_PRIMARY,
GAP_PRIMARY_GROUPS,
} from '../src/modules/clinical-gap/potential-treatment-gap.sql';
const PRIMARY_CODES = Object.keys(GAP_PRIMARY_GROUPS);
const WHOLE_MOUTH = ['K05', 'K07'];
function core(primaryCode: string, variant: 'legacy' | 'setbased') {
const rule = lookupDxTreatment(primaryCode);
if (!rule) throw new Error(`no rule for ${primaryCode}`);
const grp = GAP_PRIMARY_GROUPS[primaryCode];
return buildGapCore({
rule,
cfgFlags: GAP_FLAGS_BY_PRIMARY[primaryCode] ?? {},
allCodes: [...grp.dxCodes, ...grp.recCodes],
resolverCats: resolverCategoriesFor(primaryCode) as readonly string[],
variant,
});
}
const count = (hay: string, needle: RegExp): number => (hay.match(needle) ?? []).length;
describe('gap 集合式 ↔ 逐行形态:结构对拍', () => {
it('variant 默认 legacy;只有显式 setbased 才产出集合式拼装件', () => {
const g = core('K08', 'legacy');
expect(g.setBased).toBeUndefined();
expect(core('K08', 'setbased').setBased).toBeDefined();
});
describe.each(PRIMARY_CODES)('%s', (code) => {
const isWhole = WHOLE_MOUTH.includes(code);
it('全口码不进集合式(原样走 legacy),牙位码必须有 resolved 预聚合链', () => {
const g = core(code, 'setbased');
if (isWhole) {
// 全口场景 legacy 的 lateral 本来就会被 PG 的 useless-left-join removal 摘掉 →
// 集合式零收益;实测硬套进来还慢 2~3 倍(见 buildGapSetBased 里的早退注释)。
expect(g.setBased).toBeUndefined();
} else {
const sb = g.setBased!;
expect(sb.postCtes.sql).toContain('gap_scope');
expect(sb.postCtes.sql).toContain('gap_resolved');
expect(sb.postCtes.sql).toContain('gap_rem');
expect(sb.remJoin.sql).toContain('LEFT JOIN gap_rem');
}
});
if (!WHOLE_MOUTH.includes(code)) {
it('两种形态的 resolved 分支条数必须一致(加分支只改一边 = 静默错召)', () => {
const legacy = core(code, 'legacy').lateralJoin.sql;
const setbased = core(code, 'setbased').setBased!.postCtes.sql;
// legacy 分支用裸 UNION 分隔;setbased 用 UNION ALL(先聚合后去重,不需要 UNION 的排序去重)
const legacyBranches = count(legacy, /\bUNION\b(?!\s+ALL)/g) + 1;
const setBranches = count(setbased, /\bUNION ALL\b/g) + 1;
expect(setBranches).toBe(legacyBranches);
});
it('每个分支都挂了患者收窄(漏一个就全表扫 34GB patient_facts)', () => {
const sql = core(code, 'setbased').setBased!.postCtes.sql;
const branches = sql
.slice(sql.indexOf('FROM ('), sql.indexOf(') sb('))
.split(/\bUNION ALL\b/);
expect(branches.length).toBeGreaterThan(1);
for (const b of branches) {
expect(b).toContain('IN (SELECT patient_id FROM gap_scope)');
}
});
it('gate 列永不为 NULL —— 时间门分支带 IS NOT NULL,无门分支写死 infinity', () => {
const sql = core(code, 'setbased').setBased!.postCtes.sql;
const branches = sql
.slice(sql.indexOf('FROM ('), sql.indexOf(') sb('))
.split(/\bUNION ALL\b/);
for (const b of branches) {
const ungated = b.includes(`'infinity'::timestamptz AS gate`);
const gated = /IS NOT NULL/.test(b);
// 二者必居其一:否则全 NULL 组会被 max() 聚成 NULL、当成"无门恒过"→ 误销 → 静默少召
expect(ungated || gated).toBe(true);
}
});
it('严格 > 只出现在「更晚结构诊断」一条分支上', () => {
const sql = core(code, 'setbased').setBased!.postCtes.sql;
expect(count(sql, /TRUE AS strict/g)).toBe(1);
expect(sql).toContain('CASE WHEN r.strict THEN r.gate > c.gap_anchor ELSE r.gate >= c.gap_anchor END');
});
it('牙位顺序与重复原样保留(tooth 串会落进 plan_reasons 给客服看)', () => {
const sql = core(code, 'setbased').setBased!.postCtes.sql;
expect(sql).toContain('WITH ORDINALITY');
expect(sql).toContain('array_agg(u.x ORDER BY u.ord)');
});
it('gap_rem 无行要补空数组,不能留 NULL', () => {
const sb = core(code, 'setbased').setBased!;
expect(sb.toothOutput.sql).toContain("COALESCE(gap_rem.remaining_teeth, ARRAY[]::text[])");
expect(sb.outerWhere.sql).toContain("COALESCE(gap_rem.remaining_teeth, ARRAY[]::text[])");
});
}
});
it('病历号等值相关(整颌活动义齿)只在 K08 出现,且走 enc 列而非时间门', () => {
const k08 = core('K08', 'setbased').setBased!.postCtes.sql;
expect(k08).toContain("adx.content->>'emr_external_id' AS enc");
expect(k08).toContain("r.enc IS NULL OR r.enc = c.gap_sig_enc");
const k02 = core('K02', 'setbased').setBased!.postCtes.sql;
expect(k02).not.toContain("AS enc,\n FALSE AS strict");
expect(count(k02, /adx\./g)).toBe(0);
});
it('「建议优先于诊断」分支只对 diagnosis_record 信号生效(ndx 标志)', () => {
const sql = core('K08', 'setbased').setBased!.postCtes.sql;
expect(count(sql, /TRUE AS ndx/g)).toBe(1);
expect(sql).toContain("NOT r.ndx OR c.gap_sig_type = 'diagnosis_record'");
});
it('牙位级规则若被加上 excludeIfEverTreated,集合式必须直接炸而不是静默算错', () => {
const rule = { ...lookupDxTreatment('K08')!, excludeIfEverTreated: true };
expect(() =>
buildGapCore({
rule,
cfgFlags: GAP_FLAGS_BY_PRIMARY.K08,
allCodes: ['K08'],
resolverCats: resolverCategoriesFor('K08') as readonly string[],
variant: 'setbased',
}),
).toThrow(/excludeIfEverTreated/);
});
});
/**
* 姑息处置词(调磨/试戴/冲洗/换药/上药)→ route 到 review,不再整条丢弃。
*
* ── 为什么改 ──
* 原「⛔ 刻意不收」块的根因写着「category 被治疗史与 resolver 两处消费,只有一个旋钮,
* 只能二选一(现选都不要=丢弃)」—— 前提不成立:第二个旋钮就是 review 类目,
* 且它就是为这件事设计的(PACTreatmentCategories.review:「医生做了某个流程节点/
* 临床判断本次不动手」「不该塞 preventive,也不该丢弃」),review 在 0 个 resolver 家族里。
* 触发个案:季炎萍 TS0B010543 —— treat_plan 只有「调磨」被丢 → 库里零条治疗事实。
*
* ── 🔴 本测试锁两件事,任一被"顺手优化"掉都会静默出事 ──
* ① 这五个词必须落 _treatment_review_raw(既不被丢,也不进真治疗表)
* ② **它们不得进 &review_terms 那个 anchor** —— 那张表被 C.9 dispose 闸门复用,
* 并进去会同时打开闸门。实测那 4,336 份 EMR 处置 100% 非空,TOP400 里约 216 条
* 会落进结构 resolver 家族,含「冠周冲洗,派力奥上药」→prosthodontic(含"冠"字)、
* 「去暂封…玻璃离子暂封」→restorative 这类误判 → 误销缺口。
*
* 跑:
* pnpm test -- palliative-route-to-review
*/
import * as fs from 'node:fs';
import * as path from 'node:path';
import * as yaml from 'js-yaml';
import { runRouteByPattern } from '../src/modules/sync/transforms/operators/route-by-pattern.op';
const PALLIATIVE = ['调磨', '试戴', '冲洗', '换药', '上药'] as const;
const manifest = yaml.load(
fs.readFileSync(path.join(__dirname, '../data/jvs-dw/manifest.yaml'), 'utf-8'),
) as { transforms: Array<Record<string, unknown>> };
const treatPlanRoute = manifest.transforms.find(
(t) => t.kind === 'route_by_pattern' && t.input === '_treat_plan_raw',
) as { routes: Array<{ output: string; when: Record<string, unknown> }> };
const disposeGate = manifest.transforms.find(
(t) => t.kind === 'filter' && t.output === '_emr_dispose_gate',
) as { where: { treat_plan: { blank_or_all_in: string[] } } };
const route = (names: string[]): Record<string, string[]> => {
const { outputs } = runRouteByPattern(
treatPlanRoute as never,
names.map((n) => ({ treat_name: n })) as never,
);
return Object.fromEntries(
Object.entries(outputs).map(([k, rows]) => [
k,
(rows as Array<{ treat_name: string }>).map((r) => r.treat_name),
]),
);
};
describe('① 姑息处置 → _treatment_review_raw', () => {
it('五个词全部落 review 表', () => {
const out = route([...PALLIATIVE]);
expect(out['_treatment_review_raw']).toEqual([...PALLIATIVE]);
expect(out['_treatment_actual_raw_emr'] ?? []).toEqual([]);
});
it('⛔ 真治疗不受影响,仍走 actual', () => {
const out = route(['树脂充填', '调合', '抛光', '种植戴牙']);
expect(out['_treatment_actual_raw_emr']).toEqual(['树脂充填', '调合', '抛光', '种植戴牙']);
expect(out['_treatment_review_raw'] ?? []).toEqual([]);
});
it('原有复查词行为不变(同落 review 表)', () => {
expect(route(['常规复查'])['_treatment_review_raw']).toContain('常规复查');
});
it('「建议…」仍走 recommendation,不被本条 route 抢走', () => {
expect(route(['建议种植修复'])['_recommendation_raw']).toEqual(['建议种植修复']);
});
});
describe('② 🔴 dispose 闸门必须维持原状(anchor 不得合并)', () => {
it.each([...PALLIATIVE])('%s 不得出现在 dispose 闸门词表里', (w) => {
expect(disposeGate.where.treat_plan.blank_or_all_in).not.toContain(w);
});
it('闸门词表仍是原来那 34 个复查/流程词', () => {
expect(disposeGate.where.treat_plan.blank_or_all_in).toHaveLength(34);
expect(disposeGate.where.treat_plan.blank_or_all_in).toContain('常规复查');
});
it('⭐ 两条 route 指向同一个 review 表 —— 这是拆 anchor 的手段,不是笔误', () => {
const toReview = treatPlanRoute.routes.filter((r) => r.output === '_treatment_review_raw');
expect(toReview).toHaveLength(2);
});
});
......@@ -460,6 +460,75 @@ describe('runAllForHost 批量路径 — 5 种结局等价', () => {
expect(plans.filter((p) => p.patientId === 'pat-d' && p.status === 'active')).toHaveLength(0);
});
// ═══ 2026-08-28 抑制改为「全有或全无」(陆雪 TS0K090243)═══
// 处置是计划级的(客服点一次「已完成治疗」压掉该 plan 全部 reason),抑制却曾是逐 hit 的,
// 于是出现"部分抑制":一部分理由复活、另一部分还压着 → 卡片跟病历对不上。
test('⭐ 部分逃逸 → 被压的理由**一起复活**(不是只回来逃逸那条)', async () => {
const { prisma, plans } = makeStore({
plans: [
{
id: 'p-term-partial',
patientId: 'pat-partial',
version: 1,
status: 'abandoned',
snoozedUntil: new Date('2026-12-01T00:00:00Z'),
// 客服一次「已完成治疗」压掉这三条
reasons: [
{ scenario: SCEN, subKey: 'missing_tooth@17' },
{ scenario: SCEN, subKey: 'missing_tooth@36;37;46;47' },
{ scenario: SCEN, subKey: 'perio_no_srp@whole' },
],
},
],
});
const res = await engine(
prisma,
makeScenario([
// 同一条诊断被改判 → subKey 变了,不在抑制集里 → 逃逸
hit('pat-partial', 'hard_tissue_damage@17', 60),
// 这两条仍在抑制集里 —— 旧口径会被过滤掉,新口径必须跟着一起回来
hit('pat-partial', 'missing_tooth@36;37;46;47', 40),
hit('pat-partial', 'perio_no_srp@whole', 30),
]),
).runAllForHost({ hostId: HOST, tenantId: TENANT, now: NOW });
expect(res.plansSuppressed).toBe(0);
// 已有 v1(终态)→ 计数走"升版本"而非"新建",这是既有计数口径,不是本次改动引入的
expect(res.plansSuperseded).toBe(1);
const created = plans.find((p) => p.patientId === 'pat-partial' && p.status === 'active');
expect(created!.reasons.map((r) => r.subKey).sort()).toEqual(
['hard_tissue_damage@17', 'missing_tooth@36;37;46;47', 'perio_no_srp@whole'].sort(),
);
// 顶层分数取全部 hit 的最大值(不是只看逃逸那条)
expect(created!.priorityScore).toBe(60);
});
test('⭐ 一条都逃不出去 → 整个 plan 不生成(一起抑制,口径不变)', async () => {
const { prisma, plans } = makeStore({
plans: [
{
id: 'p-term-all',
patientId: 'pat-all',
version: 1,
status: 'abandoned',
snoozedUntil: new Date('2026-12-01T00:00:00Z'),
reasons: [
{ scenario: SCEN, subKey: 'missing_tooth@17' },
{ scenario: SCEN, subKey: 'perio_no_srp@whole' },
],
},
],
});
const res = await engine(
prisma,
makeScenario([hit('pat-all', 'missing_tooth@17'), hit('pat-all', 'perio_no_srp@whole')]),
).runAllForHost({ hostId: HOST, tenantId: TENANT, now: NOW });
expect(res.plansSuppressed).toBe(1);
expect(res.plansCreated).toBe(0);
expect(plans.filter((p) => p.patientId === 'pat-all' && p.status === 'active')).toHaveLength(0);
});
test('stale-close:有 active 但本轮 0 命中 → plansClosed=1', async () => {
const { prisma, plans } = makeStore({
plans: [
......
/**
* MISSING_TOOTH_POLISH_EVIDENCE —— 「缺失牙位上的抛光」作为修复体在位的**证据**。
*
* 蕴含成立的前提是那颗牙**已经不在了**:牙都拔了,抛的只能是义齿/种植冠。
* 所以守三条闸,任何一条松掉都会静默误销召回(少召不报错,一线只会觉得"系统没提醒过"):
* ① 词形闸:只收光秃秃一个「抛光」——「全口龈上洁治,抛光」是洗牙(针对天然牙),
* 「树脂充填…修整抛光」是充填的一个步骤(本来就带 restorative,本来就解得开)
* ② 牙位闸:牙位数 ≤ maxTeeth —— 挂满一口牙的裸抛光是洁治语境漏进来的
* ③ 场景闸:只对缺失牙(K08)开 —— 牙还在时抛光就是抛天然牙除渍,什么也不蕴含
*
* 正则用 PG 语义书写(\s 在 PG ARE 与 JS 等价),这里用 JS RegExp 校验词形。
*
* 跑:
* pnpm test -- polish-implies-restoration
*/
import { MISSING_TOOTH_POLISH_EVIDENCE } from '@pac/types';
import { GAP_FLAGS_BY_PRIMARY } from '../src/modules/clinical-gap/potential-treatment-gap.sql';
const RE = new RegExp(MISSING_TOOTH_POLISH_EVIDENCE.subtypePattern);
const hits = (subtype: string) => RE.test(subtype);
describe('MISSING_TOOTH_POLISH_EVIDENCE · ① 词形闸(收)', () => {
it.each([
['抛光'],
['抛光。'],
['抛光.'],
['抛光,'],
['抛光,'],
['抛光;'],
['抛光、'],
[' 抛光 '],
['\n抛光\n'],
])('收:%s', (subtype) => {
expect(hits(subtype)).toBe(true);
});
});
describe('MISSING_TOOTH_POLISH_EVIDENCE · ① 词形闸(不收)', () => {
it.each([
// 洁治流程 —— 针对天然牙,放行会把洗过牙的患者所有缺牙位一次性解光
['全口龈上洁治,抛光'],
['全口龈上洁治,抛光。'],
['全口洁治,抛光'],
['龈上洁治+抛光'],
['全口超声洁治,抛光。'],
['洁牙,抛光,口腔卫生宣教'],
// 涂氟 —— 全牙列预防处置,针对天然牙,不指向修复体
['抛光,涂氟'],
['抛光涂氟'],
['抛光+涂氟'],
['抛光 涂氟'],
['全口抛光涂氟'],
// 别的治疗里的一个步骤 —— 本来就带对的类目,本来就解得开
['去净腐质,酒消毒牙面,酸蚀,冲洗干燥,粘接剂涂布,树脂充填,修整抛光,嘱充填后注意事项。'],
['去腐净,GIC垫,Z250充填,调合抛光。'],
['试戴全瓷冠,精确就位,调整邻面接触点及咬合至合适,抛光,富士I玻璃离子粘固。'],
['去除愈合基台,上氧化锆基台一体冠,中心螺丝加力至35N/cm,暂封,调合抛光'],
// 义齿类 —— 已经是 prosthodontic,不该也不需要从本表走
['义齿抛光'],
['活动义齿抛光'],
['上颌种植义齿调合抛光。'],
['义齿组织面调磨抛光'],
// 全口语境的裸词变体
['全口抛光'],
['全口牙列抛光,清洁牙面,轻干燥,全牙列涂布氟保护漆。'],
// 空 / 无关
[''],
[' '],
['调磨'],
])('不收:%s', (subtype) => {
expect(hits(subtype)).toBe(false);
});
});
describe('MISSING_TOOTH_POLISH_EVIDENCE · ② 牙位闸', () => {
it('maxTeeth 必须是个小数字 —— 针对某颗牙的操作不会写满一口牙', () => {
expect(MISSING_TOOTH_POLISH_EVIDENCE.maxTeeth).toBeGreaterThanOrEqual(1);
expect(MISSING_TOOTH_POLISH_EVIDENCE.maxTeeth).toBeLessThanOrEqual(8);
});
});
describe('MISSING_TOOTH_POLISH_EVIDENCE · ③ 场景闸', () => {
it('只对缺失牙 K08 开闸', () => {
const on = Object.entries(GAP_FLAGS_BY_PRIMARY)
.filter(([, f]) => f.polishImpliesRestoration === true)
.map(([code]) => code);
expect(on).toEqual(['K08']);
});
});
describe('MISSING_TOOTH_POLISH_EVIDENCE · 表自身自洽', () => {
it('只认 preventive —— 别的类目里的抛光都是某治疗的步骤,本来就解得开', () => {
expect(MISSING_TOOTH_POLISH_EVIDENCE.category).toBe('preventive');
});
it('正则可编译且锚定首尾(避免退化成"含抛光即可")', () => {
expect(MISSING_TOOTH_POLISH_EVIDENCE.subtypePattern.startsWith('^')).toBe(true);
expect(MISSING_TOOTH_POLISH_EVIDENCE.subtypePattern.endsWith('$')).toBe(true);
expect(() => new RegExp(MISSING_TOOTH_POLISH_EVIDENCE.subtypePattern)).not.toThrow();
});
it('why 说清蕴含的依据', () => {
expect(MISSING_TOOTH_POLISH_EVIDENCE.why.length).toBeGreaterThan(6);
});
});
......@@ -47,12 +47,27 @@ function flatten(frag: Prisma.Sql | SqlLike | unknown): string {
/// 返回 [] 让调用方走空结果短路,不需要再 mock 后续查询。
function makeSqlCapturingPrisma(): { prisma: PrismaService; sql: () => string } {
const captured: string[] = [];
const queryRaw = (strings: TemplateStringsArray, ...values: unknown[]) => {
const parts = Array.isArray(strings?.raw) ? strings.raw : [String(strings)];
const queryRaw = (strings: TemplateStringsArray | { sql?: string }, ...values: unknown[]) => {
// 两种调用形态都要认:
// ① 标签模板 `$queryRaw`...`` → strings 是 TemplateStringsArray
// ② 传 Prisma.sql 对象 `$queryRaw(sql)` → strings 是 Sql,已摊平好的 .sql 直接用
// (2026-08-29 起主查询走 ② —— 先建 Prisma.sql 对象才能在 PAC_RECALL_DUMP_SQL=1 时
// 打出完整 SQL 做 EXPLAIN;pg_stat_activity 会把它截断在 1024 字节。)
if (typeof (strings as { sql?: string })?.sql === 'string') {
captured.push((strings as { sql: string }).sql);
return Promise.resolve([]);
}
const t = strings as TemplateStringsArray;
const parts = Array.isArray(t?.raw) ? t.raw : [String(strings)];
captured.push(parts.map((s, i) => s + (i < values.length ? flatten(values[i]) : '')).join(''));
return Promise.resolve([]);
};
const prisma = { $queryRaw: queryRaw } as unknown as PrismaService;
// 主查询走 `$transaction([SET LOCAL work_mem, $queryRaw(sql)])`(2026-08-29 起,
// 为抬高 work_mem;见 scenario 里那段实测注释)。桩要同时认这三个:
// $queryRaw 负责捕获 SQL,$executeRaw 吞掉 SET LOCAL,$transaction 把数组按序解析。
const noop = () => Promise.resolve(0);
const tx = (arr: unknown[]) => Promise.all(arr as Promise<unknown>[]);
const prisma = { $queryRaw: queryRaw, $executeRaw: noop, $transaction: tx } as unknown as PrismaService;
return { prisma, sql: () => captured.join('\n/* --- next query --- */\n') };
}
......
/**
* reconcileDiagnosisCode — 宿主 std_code 错标 / 截码失真的纠偏规则。
*
* 两条规则各有一个**反向闸**,测试重点全在闸上 —— 闸失效不会报错,
* 只会让一批患者悄悄换场景(K08 缺牙 base 60 ↔ K03 牙体损伤 base 35),
* 表现为矩阵某一列莫名多/少几百格。
*
* 跑:
* pnpm test -- reconcile-diagnosis-code
*/
import { reconcileDiagnosisCode } from '@pac/types';
describe('reconcileDiagnosisCode', () => {
describe('K00 → K08(后天缺失被错标成发育障碍)', () => {
it.each(['牙齿缺少', '牙缺失', '缺牙', '后天性牙齿缺失'])('%s → K08', (name) => {
expect(reconcileDiagnosisCode('K00', name)).toBe('K08');
});
it('带「先天」→ 保持 K00(真发育障碍)', () => {
expect(reconcileDiagnosisCode('K00', '先天性牙齿缺失')).toBe('K00');
expect(reconcileDiagnosisCode('K00', '先天缺牙')).toBe('K00');
});
it('其它 K00 不动', () => {
expect(reconcileDiagnosisCode('K00', '乳牙滞留')).toBe('K00');
expect(reconcileDiagnosisCode('K00', '多生牙')).toBe('K00');
});
});
describe('K08 → K03(K08.3 牙根残留被截成 K08 缺牙)', () => {
// 生产实测:名字恰为这几个的 932 条应翻 K03
it.each(['残根', '残冠', '残根残冠', '残根/残冠'])('%s → K03', (name) => {
expect(reconcileDiagnosisCode('K08', name)).toBe('K03');
});
it('无法保留 / 不能保留 → K03', () => {
expect(reconcileDiagnosisCode('K08', '无法保留')).toBe('K03');
expect(reconcileDiagnosisCode('K08', '患牙不能保留')).toBe('K03');
});
// ⛔ 反向闸:复合诊断(生产 325 条)两个病挤在一个 name,翻 K03 会丢掉缺牙那一半
it.each([
'牙列缺损,残根',
'缺失,残根',
'残根,牙列缺损',
'牙缺失;#44残根',
'牙缺失, 残根残冠',
'16,12,26缺失,17,11,21,22,23,25,27残根',
])('复合诊断「%s」→ 保持 K08', (name) => {
expect(reconcileDiagnosisCode('K08', name)).toBe('K08');
});
it('纯缺牙类 K08 不动', () => {
expect(reconcileDiagnosisCode('K08', '后天性牙齿缺失')).toBe('K08');
expect(reconcileDiagnosisCode('K08', '牙列缺损')).toBe('K08');
expect(reconcileDiagnosisCode('K08', '牙齿缺少')).toBe('K08');
});
it('残根本来就是 K03 的不受影响(医生没填 std_code 的多数派)', () => {
expect(reconcileDiagnosisCode('K03', '残根')).toBe('K03');
expect(reconcileDiagnosisCode('K03', '残冠')).toBe('K03');
});
});
describe('边界', () => {
it('code 为 null → null', () => {
expect(reconcileDiagnosisCode(null, '残根')).toBeNull();
});
it('name 为 null → 原样返回,不误翻', () => {
expect(reconcileDiagnosisCode('K08', null)).toBe('K08');
expect(reconcileDiagnosisCode('K00', null)).toBe('K00');
});
it('无关码不动', () => {
expect(reconcileDiagnosisCode('K05', '慢性牙周炎')).toBe('K05');
expect(reconcileDiagnosisCode('K02', '龋病')).toBe('K02');
});
});
});
/**
* refineCategoriesForDiagnosis —— K03 同码内主类目重排(残根/残冠 → 外科)。
*
* 守的是**「分配的潜在治疗项目」与「卡片上的目标」必须对应**这条口径:
* 分配矩阵把残根患者放进「拔牙治疗」列(POTENTIAL_LABEL_RULES: K03 + 含词 → extraction),
* 卡片的「目标 · X」就不能写「充填 / 嵌体」。两处共用 EXTRACTION_NAME_KEYWORDS 同一份词表。
*
* 跑:
* pnpm test -- refine-categories-k03
*/
import {
refineCategoriesForDiagnosis,
recommendedCategoriesForAge,
DiagnosisTreatmentMap,
EXTRACTION_NAME_KEYWORDS,
classifyCodeToLabel,
} from '@pac/types';
/** K03 在字典里的候选类目(单一真理源,不在测试里硬编码顺序) */
const K03_CATS = DiagnosisTreatmentMap.K03!.categories as readonly string[];
/** 复现 scenario 里 focusCategory 的算法:语义重排 → 年龄重排 → 取首项 */
const focusOf = (code: string, nameZh: string, age: number | null): string | null =>
recommendedCategoriesForAge(refineCategoriesForDiagnosis(code, nameZh, K03_CATS), age)[0] ?? null;
describe('refineCategoriesForDiagnosis · K03', () => {
it('K03 候选类目本身含 surgical(否则重排无处可挪)', () => {
expect(K03_CATS).toContain('surgical');
expect(K03_CATS).toContain('restorative');
});
describe('残根 / 残冠 → 外科打头', () => {
it.each(EXTRACTION_NAME_KEYWORDS)('「%s」→ surgical', (name) => {
expect(refineCategoriesForDiagnosis('K03', name, K03_CATS)[0]).toBe('surgical');
});
it('66 岁残根患者(孙海燕 TS0M012582 牙16;18;28;47)目标 = 外科,不是充填', () => {
expect(focusOf('K03', '残根', 66)).toBe('surgical');
});
it('年龄闸不会把它挪回去(K03 无 implant,recommendedCategoriesForAge 空转)', () => {
for (const age of [8, 18, 40, 66, 88, null]) {
expect(focusOf('K03', '残冠', age)).toBe('surgical');
}
});
});
describe('⛔ 只挪位,不增删', () => {
it('返回集合恒等于入参集合', () => {
const out = refineCategoriesForDiagnosis('K03', '残根', K03_CATS);
expect([...out].sort()).toEqual([...K03_CATS].sort());
});
it('其余相对顺序保持', () => {
const out = refineCategoriesForDiagnosis('K03', '残根', K03_CATS);
const rest = out.filter((c) => c !== 'surgical');
expect(rest).toEqual(K03_CATS.filter((c) => c !== 'surgical'));
});
});
describe('不该动的不动', () => {
it('楔状缺损 / 牙体缺损 → 保持原序(补得回来)', () => {
expect(refineCategoriesForDiagnosis('K03', '牙齿楔状缺损', K03_CATS)).toEqual(K03_CATS);
expect(refineCategoriesForDiagnosis('K03', '牙体缺损', K03_CATS)).toEqual(K03_CATS);
expect(focusOf('K03', '牙齿楔状缺损', 66)).toBe('restorative');
});
it('诊断名为空 → 原样返回(历史数据行为不变)', () => {
expect(refineCategoriesForDiagnosis('K03', null, K03_CATS)).toEqual(K03_CATS);
expect(refineCategoriesForDiagnosis('K03', ' ', K03_CATS)).toEqual(K03_CATS);
});
it('别的码不受影响', () => {
const k08 = DiagnosisTreatmentMap.K08!.categories as readonly string[];
expect(refineCategoriesForDiagnosis('K08', '残根', k08)).toEqual(k08);
expect(refineCategoriesForDiagnosis(null, '残根', K03_CATS)).toEqual(K03_CATS);
});
});
describe('K00 原有行为回归(表驱动改写后不能退化)', () => {
const k00 = DiagnosisTreatmentMap.K00!.categories as readonly string[];
it.each([
['乳牙滞留', 'surgical'],
['乳牙早失', 'orthodontic'],
['萌出障碍', 'orthodontic'],
['先天缺牙', 'prosthodontic'],
['釉质发育不全', 'prosthodontic'],
])('K00「%s」→ %s 打头', (name, lead) => {
expect(refineCategoriesForDiagnosis('K00', name, k00)[0]).toBe(lead);
});
});
describe('⭐ 标签与目标同源(本次修复的核心口径)', () => {
it.each(EXTRACTION_NAME_KEYWORDS)('「%s」:分配标签 extraction ⟺ 目标 surgical', (name) => {
expect(classifyCodeToLabel('K03', name, 66)).toBe('extraction');
expect(focusOf('K03', name, 66)).toBe('surgical');
});
it('非拔除类 K03:标签 restoration ⟺ 目标 restorative', () => {
expect(classifyCodeToLabel('K03', '牙齿楔状缺损', 66)).toBe('restoration');
expect(focusOf('K03', '牙齿楔状缺损', 66)).toBe('restorative');
});
});
});
/**
* RESTORATION_IN_PLACE_* —— 检查所见写着修复体在位(缺牙缺口的第三条举证路线)。
*
* 前两条(治疗名 / 主诉现病史)都要求患者**为那副修复体而来**;这条接的是
* "为别的牙来、医生顺带记了全口状况"的那一类 —— 周燕芬 TS0M013273 即是。
*
* 判据 = 修复体名词 ∧ 在位状态词 ∧ ¬失效词,且条目牙位不跨颌。
* 🔴 失效词是命门:修复体**做了但不能用**仍然该召回。
*
* 跑:
* pnpm test -- restoration-in-place-from-exam
*/
import {
RESTORATION_IN_PLACE_TERMS_RE,
RESTORATION_IN_PLACE_STATE_RE,
RESTORATION_IN_PLACE_FAIL_RE,
RESTORATION_IN_PLACE_NEG_STRIP_RE,
} from '@pac/types';
import { GAP_FLAGS_BY_PRIMARY } from '../src/modules/clinical-gap/potential-treatment-gap.sql';
const TERMS = new RegExp(RESTORATION_IN_PLACE_TERMS_RE);
const STATE = new RegExp(RESTORATION_IN_PLACE_STATE_RE);
const FAIL = new RegExp(RESTORATION_IN_PLACE_FAIL_RE);
const NEG = new RegExp(RESTORATION_IN_PLACE_NEG_STRIP_RE, 'g');
/** 模拟三条件(不含同颌闸 —— 那是牙位层面的,由 SQL 负责) */
const inPlace = (msg: string): boolean =>
TERMS.test(msg) && STATE.test(msg) && !FAIL.test(msg.replace(NEG, ''));
describe('正例 —— 生产真实检查所见,判为修复体在位', () => {
it.each([
// 🔴 周燕芬 TS0M013273 上颌 11~27(单颌条目)—— 本条路线的触发个案
'见活动义齿,基托边缘密合,伸展范围良好,固位力良好,无压痛',
'下颌覆盖牙齿在位,无松动,就位可,边缘贴合;杆卡无松动,牙龈无红肿',
'可见全瓷冠修复,边缘密合邻接良好,咬合合适,牙龈及黏膜未见明显异常',
'见全冠修复体,叩(-),无松动,牙龈未见明显异常',
'义齿密合良好,牙龈无红肿',
// ⛔ 「固位良好/固位稳定」不得被「固位…不良」那条贬义模式误否
'铸造活动义齿,固位稳定,咬合关系良好。', // 邢五一 BJ0A078837
'缺失,粘膜无异常,牙槽嵴有明显吸收。活动义齿固位良好。', // 孙玉娟 GZ0A031993
// 🔴 「无松动」类是好话,抹掉后不得被裸「松动」否掉
'金属烤瓷冠修复桥体,边缘密合,未见明显松动,冷叩诊无不适', // 陈铭燕 GZ0A016429
'烤瓷冠,边缘密合,探痛(—),冷热(—),叩(—)松动(—)', // 马琦 SC10794
'活动义齿在位,牙龈粘膜未见明显异常,37牙体缺损至平龈,叩(-),未见明显松动', // 李太花
'46口内未探及,45-47烤瓷固定桥修复,边缘密合度一般,无明显松动。', // 丁亚茹
// ⛔ 「尚可/尚密合/较密合」是可接受状态,不得被贬义模式误否
'已行固定桥修复,边缘密合度尚可,冷不敏,叩(-),不松动', // 陈晓晖 BJ0D033690
'冠桥修复体,边缘较密合,叩诊无不适', // 李科伟 BD13007
'全瓷冠桥修复,近中单端悬臂,25#缺失,冠边缘位于龈上尚密合,叩-,龈-,松-', // 周斯男 SH0L005714
// ⛔ 贬义词跨标点不得误否:「边缘密合」后面另起一句说别的差
'种植冠修复,边缘密合,牙龈无红肿,口腔卫生差'
])('%s', (msg) => {
expect(inPlace(msg)).toBe(true);
});
});
describe('🔴 反例 —— 做了但不能用 / 还没做,必须仍然召回', () => {
it.each([
// 王茜 BJ0F028868:外院假牙做好了却装不上
'外院制作右下活动假牙无法就位一周',
'活动义齿基托折断,无法佩戴',
'烤瓷冠脱落,牙体缺损,边缘尚密合', // 有"密合"但也有"脱落" → 失效词一票否决
'义齿破损,固位力良好但影响进食',
'牙列缺损,建议活动义齿修复,患者考虑中',
'缺牙区黏膜良好,要求义齿修复',
'上颌活动义齿不贴合,需重衬', // 童然夫/周根娣同类:不贴合 = 需处理
// 🔴 生产实测第一版误销的那批 —— 「密合」子串会匹配它的否定形式
'固定桥修复,边缘欠密合,龈缘红肿', // 赵河 BA27761
'锤造冠固定桥修复,密合度差,继发龋坏', // 陈德家 BA40818
'固定桥修复,边缘密合度欠佳,牙龈(-)叩(-)松动(-)', // 张华 BJ0A054066
'固定桥修复,边缘密合性不佳,46近中可探入', // 王竞男 BJ0A090886
'金属烤瓷冠桥修复长桥,修复体松动。边缘欠密合,部分崩瓷,咬合欠佳', // 翟健民 BJ0C021369
'种植修复体,松动I,边缘不密合', // 李露 BJ0D049016
'缺失,活动义齿修复,不可摘除,不密合', // 王秀梅 BJ0D030624
'胶连活动义齿修复,义齿与软组织不密合,客人自觉舌侧基托厚,漏风,影响发音', // 朱少宇 BJ0A081328
'固定桥,崩瓷,边缘不密合,叩(-),不松,龈水肿', // 余健 BJ0D024381
// 🔴 否定前缀第二次踩到:「无义齿修复」含「义齿」却是说没有
'缺失,无义齿修复,46,47均向近中移位,46烤瓷冠修复,边缘密合,无异常', // 徐磊 BJ0U001139
// 基牙松动度的另两种写法
'烤瓷冠修复体,边缘较密合,11Ⅰ°松,21Ⅱ°松,22Ⅰ-Ⅱ°松', // 王秋枫 BJ0D053568
'可见全冠,边缘密合,松动二度,叩(+)', // 侯永生 BJ0R008817
'瓷修复体存,边缘密合,牙松动2度,PD:8mm', // 康振英 GZ0A019990
'固定桥完好,边缘密合较差,叩痛(-),有明显松动', // 吴实 SH0T009924
// 🔴 第四轮:「密合」的否定形式列不完,改用模式匹配
'单端固定桥;44-47烤瓷牙+树脂基托;牙冠形态不良,密合度不佳,局部牙龈红肿', // 张顺宝 SH0K018368
'金属全冠修复体,边缘密合差,卡探针', // 姜华 ZJ0B010604
'覆盖义齿修复,固位不良,port基台稳定无松动,牙龈未见异常', // 周锡英 TS0B001674
])('%s → 不作数', (msg) => {
expect(inPlace(msg)).toBe(false);
});
it('没有修复体名词,只是牙龈好 → 不作数', () => {
expect(inPlace('牙龈未见明显异常,无压痛')).toBe(false);
});
it('有修复体但没说状态 → 不作数', () => {
expect(inPlace('见活动义齿')).toBe(false);
});
});
describe('表自身自洽', () => {
it('⛔ 不收裸「冠」—— 天然牙冠也叫冠', () => {
expect(RESTORATION_IN_PLACE_TERMS_RE).not.toMatch(/(^|\|)(\||$)/);
expect(inPlace('牙冠完好,无松动')).toBe(false);
// 代价:省略定语的「冠边缘密合」接不住。那类患者通常另有治疗名/主诉证据
// (陶美玉 TS0K070812 就被「种植复查」和处置「调合」两条路各接住一次),
// 不值得为它放宽到裸「冠」。
expect(inPlace('冠边缘密合,牙龈未见异常')).toBe(false);
});
it('🔴 失效词必须含"无法就位/脱落/折断" —— 做了但不能用仍该召回', () => {
for (const w of ['无法就位', '脱落', '折断', '不贴合', '不密合', '欠密合', '崩瓷']) {
expect(RESTORATION_IN_PLACE_FAIL_RE).toContain(w);
}
});
it('🔴 只对缺失牙(K08)开闸 —— 判据是拿 K08 的检查所见核出来的', () => {
const on = Object.entries(GAP_FLAGS_BY_PRIMARY)
.filter(([, f]) => f.restorationInPlaceFromExam === true)
.map(([code]) => code);
expect(on).toEqual(['K08']);
});
it('三个正则都可编译', () => {
for (const re of [
RESTORATION_IN_PLACE_TERMS_RE,
RESTORATION_IN_PLACE_STATE_RE,
RESTORATION_IN_PLACE_FAIL_RE,
]) {
expect(() => new RegExp(re)).not.toThrow();
}
});
});
/**
* REVIEW_IMPLIES_TREATMENT —— 「复查/复诊」作为修复体在位的**证据**。
*
* 守两条闸,任何一条松掉都会静默误销召回(少召不报错,一线只会觉得"系统没提醒过"):
* ① 词形闸:必须「治疗词 + 复查/复诊」——「正畸复诊」(在做)vs「正畸会诊」(还在谈)只差一字
* ② 类目闸:必须与 resolverCategoriesFor 交集 ——「牙周复查」不能解 K02 龋齿
*
* 正则用 PG 语义书写,这里用 JS RegExp 近似校验词形(两者对本表用到的语法等价)。
*
* 跑:
* pnpm test -- review-implies-treatment
*/
import { REVIEW_IMPLIES_TREATMENT, resolverCategoriesFor } from '@pac/types';
/** 某个 subtype 命中哪些规则 */
const rulesFor = (subtype: string) =>
REVIEW_IMPLIES_TREATMENT.filter((r) => new RegExp(r.pattern).test(subtype));
/** 模拟消费方:该诊断码下,这个 subtype 会不会解除缺口 */
const resolves = (code: string, subtype: string): boolean => {
const cats = resolverCategoriesFor(code) as readonly string[];
return rulesFor(subtype).some((r) => cats.includes(r.category));
};
describe('REVIEW_IMPLIES_TREATMENT · 词形闸', () => {
it.each([
['种植复查', 'implant'],
['牙周复查', 'periodontic'],
['正畸复诊', 'orthodontic'],
['正畸复查', 'orthodontic'],
['保持器复诊', 'orthodontic'],
])('生产实有词「%s」→ %s', (subtype, cat) => {
expect(rulesFor(subtype).map((r) => r.category)).toContain(cat);
});
// ⛔ 这些是生产带牙位 review 里的大头,收进来就是误销
it.each([
'观察', '暂观', '观察,必要时拔除', '暂观,必要时拔除', '随访观察', '观察随访', '随诊观察',
'无治疗', '初诊检查', '检查', '方案沟通', '沟通治疗方案', '听方案', '取资料', '缴费',
'转诊', '请全科医生会诊', '治疗中',
])('「%s」不命中任何规则', (subtype) => {
expect(rulesFor(subtype)).toHaveLength(0);
});
// ⛔ 泛指复查也不收 —— 它不指明治疗对象(宽口径 391 条里的争议区)
it.each(['复查', '常规复查', '定期复查'])('泛指「%s」不命中', (subtype) => {
expect(rulesFor(subtype)).toHaveLength(0);
});
// ⛔ 一字之差,语义相反
it.each(['正畸会诊', '正畸咨询', '正畸检查'])('「%s」不命中(还在谈,没在做)', (subtype) => {
expect(rulesFor(subtype)).toHaveLength(0);
});
});
describe('REVIEW_IMPLIES_TREATMENT · 类目闸', () => {
it('K08 缺失牙 + 种植复查 → 解除(宋志宏 TS0M001982 牙46;36;37)', () => {
expect(resolves('K08', '种植复查')).toBe(true);
});
it('K04 根尖周炎 + 种植复查 → 解除(种了牙就没牙髓了)', () => {
expect(resolves('K04', '种植复查')).toBe(true);
});
// ⛔ 生产实测的 9 条跨类目错配,逐条锁死
it.each([
['K02', '牙周复查', '洗牙不补龋'],
['K03', '牙周复查', '洗牙不修牙体'],
['K08', '牙周复查', '洗牙不修缺牙'],
['K06', '种植复查', 'K06 只认牙周/外科'],
['K01', '正畸复查', 'K01 走结构家族,不含正畸'],
])('%s + %s → 不解除(%s)', (code, subtype) => {
expect(resolves(code, subtype)).toBe(false);
});
});
describe('REVIEW_IMPLIES_TREATMENT · 表自身自洽', () => {
it('每条 pattern 都能编译成正则', () => {
for (const r of REVIEW_IMPLIES_TREATMENT) {
expect(() => new RegExp(r.pattern)).not.toThrow();
}
});
it('每条都写了 why(改表的人得说明临床依据)', () => {
for (const r of REVIEW_IMPLIES_TREATMENT) {
expect(r.why.length).toBeGreaterThan(4);
}
});
it('每条 pattern 都强制要求出现 复查/复诊', () => {
for (const r of REVIEW_IMPLIES_TREATMENT) {
expect(r.pattern).toMatch(/复查\|复诊/);
}
});
});
import { SyncIncrementalSchedulerService } from '../src/queues/sync-incremental.scheduler';
/**
* cron 回调防重入(runHostSafe 的 runningHosts 闸)回归。
*
* 🔴 2026-08-29 生产事故:plan 段耗时涨过 2 小时后,cron(每 2 小时)照常触发下一轮,
* 两轮的 plan 段并发抢同一批 I/O → 都更慢 → 更易被再下一轮套圈 → 雪崩。
* 实测 08-28 20:15 起连续多轮 plan 段一次都没跑完,到 08-29 上午仍有两轮在并行。
*
* 为什么现有的锁挡不住(这两条是本用例存在的理由):
* ① NestJS CronJob **默认不防重入** —— 上一次回调还在 await,下一次照样进;
* ② `sync_logs` 的 partial UNIQUE(host_id) WHERE status='running' 只覆盖**摄入段**,
* 摄入一结束锁就放了,而最慢的 persona / plan 段还在跑。
*
* 跑:
* pnpm test -- scheduler-reentrancy-guard
*/
/** 造一个只关心 runOne 编排的实例;runOne 用可控的 promise 替换 */
function makeService() {
const svc = new SyncIncrementalSchedulerService(
{} as never, {} as never, {} as never, {} as never, {} as never, {} as never,
);
const calls: string[] = [];
let release!: () => void;
const gate = new Promise<void>((r) => { release = r; });
(svc as unknown as { runOne: (dir: string) => Promise<void> }).runOne = async (dir: string) => {
calls.push(dir);
await gate; // 卡住不返回 = 模拟"上一轮还在跑"
};
const run = (host: string) =>
(svc as unknown as { runHostSafe: (h: string) => Promise<void> }).runHostSafe(host);
return { svc, calls, run, release };
}
describe('runHostSafe 防重入', () => {
test('🔴 上一轮未结束时,同 host 的下一轮被跳过(不并发进 runOne)', async () => {
const { calls, run, release } = makeService();
const first = run('jvs-dw'); // 卡在 gate 上
await run('jvs-dw'); // 第二次应立即返回
await run('jvs-dw'); // 第三次同样
expect(calls).toHaveLength(1); // ⛔ 只有第一轮真正进了 runOne
release();
await first;
});
test('上一轮结束后,下一轮正常放行', async () => {
const { calls, run, release } = makeService();
const first = run('jvs-dw');
release();
await first;
await run('jvs-dw');
expect(calls).toHaveLength(2);
});
test('不同 host 互不阻塞(闸是按 host 的,不是全局)', async () => {
const { calls, run, release } = makeService();
const a = run('jvs-dw');
const b = run('friday');
// runOne 收到的是 path.join(dataDir, host) 的完整路径,按后缀断言
expect(calls.map((d) => d.split('/').pop())).toEqual(['jvs-dw', 'friday']);
release();
await Promise.all([a, b]);
});
test('⛔ runOne 抛异常也必须释放闸 —— 否则该 host 被永久锁死到进程重启', async () => {
const svc = new SyncIncrementalSchedulerService(
{} as never, {} as never, {} as never, {} as never, {} as never, {} as never,
);
let n = 0;
(svc as unknown as { runOne: (dir: string) => Promise<void> }).runOne = async () => {
n++;
throw new Error('boom');
};
const run = (h: string) =>
(svc as unknown as { runHostSafe: (h: string) => Promise<void> }).runHostSafe(h);
await run('jvs-dw'); // runHostSafe 吞异常
await run('jvs-dw'); // 闸已释放 → 应能再进
expect(n).toBe(2);
});
});
/**
* traceRawSourceTable —— reparse 的源表回溯。
*
* 回溯错了不会报错:rawPayload 被灌进**中间表**,随即被 transform 链用空结果覆盖 →
* reparse 恒 0 变更、exit 0、日志写"重衍完成"。2026-08-27 实测治疗类三个资源全中,
* 意味着任何治疗字典修补都无法回填存量,而且没人会发现。
*
* 跑:
* pnpm test -- trace-raw-source-table
*/
import * as fs from 'node:fs';
import * as path from 'node:path';
import * as yaml from 'js-yaml';
import { traceRawSourceTable } from '../src/modules/sync/cold-import/cold-import.service';
const MANIFEST = path.join(__dirname, '../data/jvs-dw/manifest.yaml');
const manifest = yaml.load(fs.readFileSync(MANIFEST, 'utf8')) as { transforms?: unknown[] };
const TF = manifest.transforms ?? [];
const trace = (t: string) => traceRawSourceTable(t, TF);
describe('traceRawSourceTable · 真实 manifest', () => {
// ⛔ 这四个必须都回溯到真正的 DW 原始表,任何一个停在中间表 = reparse 静默失效
it.each([
'treatment_actual_rows',
'treatment_planned_rows',
'treatment_review_rows',
'diagnosis_rows',
])('%s → fact_emr_treatment_out', (tbl) => {
expect(trace(tbl)).toBe('fact_emr_treatment_out');
});
it('⛔ 不能停在 route 的输出表上(本次 bug 的指纹)', () => {
for (const tbl of ['treatment_actual_rows', 'treatment_planned_rows', 'treatment_review_rows']) {
expect(trace(tbl)).not.toMatch(/^_/); // 下划线开头 = manifest 里的中间表
}
});
it('原始表回溯到自身(reparse 据此判定"非 transform 产出"并跳过)', () => {
expect(trace('fact_emr_treatment_out')).toBe('fact_emr_treatment_out');
});
});
describe('traceRawSourceTable · 合成用例', () => {
it('route_by_pattern 用的是 routes 字段,不是 outputs', () => {
const tf = [
{ kind: 'split_json_array', input: 'raw_src', output: '_mid' },
{ kind: 'route_by_pattern', input: '_mid', field: 'x', routes: [{ output: '_a' }, { output: '_b' }] },
{ kind: 'derive', input: '_a', output: 'final_rows' },
];
expect(traceRawSourceTable('final_rows', tf)).toBe('raw_src');
});
it('union 各分支回溯不一致 → 停在 union 输出(既有语义不变)', () => {
const tf = [
{ kind: 'derive', input: 'raw_a', output: 't1' },
{ kind: 'union', inputs: ['t1', 'raw_b'], output: 'merged' },
];
expect(traceRawSourceTable('merged', tf)).toBe('merged');
});
});
......@@ -17,7 +17,10 @@ const transforms = [
{ kind: 'filter', input: 'patient_settlement_spec', output: '_refund_item_raw', where: {} },
{ kind: 'lookup', input: '_refund_item_raw', output: 'refund_item_rows', from: 'patient_settlement', select: {} },
// route_by_pattern 多输出回同一 input
{ kind: 'route_by_pattern', input: '_treat_raw', outputs: [{ output: '_actual_raw' }, { output: '_rec_raw' }] },
// ⚠️ 2026-08-27 修正:字段名是 `routes`(见 transforms.schema.ts RouteByPatternOpSchema),
// 本 fixture 原先写 `outputs` —— 跟当时的实现犯了同一个错,于是测试一直绿、生产一直坏
// (治疗类三个资源 reparse 恒 0 变更且报成功)。fixture 必须照**契约**写,不能照实现写。
{ kind: 'route_by_pattern', input: '_treat_raw', routes: [{ output: '_actual_raw' }, { output: '_rec_raw' }] },
];
describe('traceRawSourceTable', () => {
......
/**
* TREATED_EVIDENCE_* —— 病历自由文本自证已治疗(补充举证路线)。
*
* 判据 = 修复体词 ∧ 完成标记 ∧ ¬未发生标记。三个条件缺一不可,
* 少一个就会把"还没做"判成"已做" → 少召,而少召不报错。
* 下面正反例**全部取自生产真实串**(主诉 / 现病史 / 处置)。
*
* 跑:
* pnpm test -- treated-evidence-from-emr-text
*/
import { GAP_FLAGS_BY_PRIMARY } from '../src/modules/clinical-gap/potential-treatment-gap.sql';
import {
TREATED_EVIDENCE_RESTORATION_TERMS,
TREATED_EVIDENCE_COMPLETION_RE,
TREATED_EVIDENCE_INTENT_EXCLUDE_RE,
TREATED_EVIDENCE_EMR_FIELDS,
treatedEvidenceTriggersFor,
TREATED_EVIDENCE_BATCH_CATEGORIES,
resolverCategoriesFor,
} from '@pac/types';
const COMPLETION = new RegExp(TREATED_EVIDENCE_COMPLETION_RE);
const INTENT = new RegExp(TREATED_EVIDENCE_INTENT_EXCLUDE_RE);
/** 模拟 SQL 三条件:返回命中的类目集合(空 = 不作数) */
const evidenceCats = (text: string): string[] => {
if (!COMPLETION.test(text)) return [];
if (INTENT.test(text)) return [];
return TREATED_EVIDENCE_RESTORATION_TERMS.filter((t) => new RegExp(t.pattern).test(text)).map(
(t) => t.category,
);
};
describe('正例 —— 生产真实主诉/现病史,必须判为已治疗', () => {
it.each([
['种植戴牙3个月余按约复查', 'prosthodontic'], // 王作荣 TS0M002910 主诉
['3个月前在本院行左下后牙种植戴牙,无不适,今按约复查', 'prosthodontic'], // 同上 现病史
['种植戴牙1年按约复查', 'prosthodontic'], // 陈惠英 TS0K028757
['种植戴牙3月余', 'prosthodontic'], // ⭐ 不含"复查"二字 —— 现有 REVIEW 词表漏的就是这类
['右侧后牙种植戴牙后一年', 'prosthodontic'],
['一年前患者于本院行右上后牙种植戴牙,现按计划复诊', 'prosthodontic'],
['复诊,主诉上颌义齿咬物疼痛,冷热不适近月余', 'prosthodontic'],
['种植戴牙半年余', 'prosthodontic'],
['种植戴牙4月余按约复查', 'prosthodontic'],
['上颌活动义齿不贴合1周余,今来就诊', 'prosthodontic'],
// ⭐ 量词含「数/多/几」—— 季炎萍 TS0B010543,不认就漏
['下颌活动义齿戴牙数年余', 'prosthodontic'],
['患者下颌活动义齿戴牙数年余,有压痛,今来就诊', 'prosthodontic'],
['全口义齿戴用多年', 'prosthodontic'],
['戴牙后2周复查,自述无特殊不适。', 'prosthodontic'],
])('%s → %s', (text, cat) => {
expect(evidenceCats(text)).toContain(cat);
});
});
describe('反例 —— 必须判为不作数(判错就是少召,不报错)', () => {
it.each([
// ① 诉求,不是已发生
['要求窝沟封闭'],
['患者要求窝沟封闭。(保险收费给孩子正畸用)'],
['今行检查,未行处置,口腔卫生宣教,建议择期种植修复'],
// 🔴 ② 王志荣本人 2024-06-14 现病史 —— 与他 2026 那句字面几乎一样,语义相反
['下颌活动牙齿修复后牙龈反复疼痛,要求拔除口内剩余牙齿后全口义齿修复'],
// ③ 有完成标记但**没有修复体词** —— 光凭"复查"不能算已治疗
['定期复查'],
['定期3个月复查涂氟'],
['按计划复诊'],
['左上后牙拔除术后3月余'], // 拔牙3个月正是该种植,必须照常召回
['右下后牙松动数月。'],
['右下后牙折断1月'], // 童然夫主诉:折断的是天然牙,不是修复体
// ④ 有修复体词但没有完成标记
['要求种植修复'],
['咨询义齿修复方案'],
// 🔴 ⑤ 生产实测踩出来的:种植是多阶段过程,没到「戴牙」都不算修复完成
['右上后牙种植体拔除后两个月'], // BJ0E015462 —— 种植失败取出,最该召回
['右上后牙12年前行种植修复,1个月前种植体及牙冠一起脱落'], // SC15894
['复查种植三周后复查。三周前于我处行左上后牙种植手术,今来复查。'], // SH0Q019876 只做了一期
['拔牙后三个月,种植检查'], // BJ0V001864 正在准备种植
['拔牙3个月后常规复诊。转诊种植科。'], // TS0M004268
['右上后牙自行脱落,前来检查。6个月后种植前检查'], // SH0L009117
// 🔴 ⑥ 修复体做了但不能用 —— 仍该召回
['外院制作右下活动假牙无法就位一周。'], // BJ0F028868
['前牙牙冠脱落数日。患者自述十余年前曾在我院牙冠修复,几日前牙冠脱落,今来诊。'], // BA14863
])('%s → 不作数', (text) => {
expect(evidenceCats(text)).toEqual([]);
});
});
describe('表自身自洽', () => {
it('⭐「计划」不在未发生词里 —— 否则误伤「按计划复诊」', () => {
expect(INTENT.test('按计划复诊')).toBe(false);
expect(TREATED_EVIDENCE_INTENT_EXCLUDE_RE).not.toMatch(/(^|\|)计划(\||$)/);
});
it('⭐ 触发条件 = 该动作解不开本场景的缺口(不是硬清单)', () => {
const k08 = resolverCategoriesFor('K08') as readonly string[];
const trig = treatedEvidenceTriggersFor(k08);
// 治疗类走现有 resolver,不进本表
for (const c of ['implant', 'prosthodontic', 'restorative', 'endodontic', 'surgical']) {
expect(trig).not.toContain(c);
}
// 🔴 牙周必须在内 —— 王祖妹 TS0B006805 那次动作是「转洁牙中心全口牙洁治」(periodontic),
// 硬写 ['preventive','review'] 会漏掉她
for (const c of ['preventive', 'review', 'periodontic', 'orthodontic']) {
expect(trig).toContain(c);
}
});
it('⛔ 牙周/正畸是全口批量记录,额外加牙位上限', () => {
expect([...TREATED_EVIDENCE_BATCH_CATEGORIES].sort()).toEqual(['orthodontic', 'periodontic']);
});
it('⛔ 暂不扫检查所见(混合描述) / 处置(夹带 base64)', () => {
expect(TREATED_EVIDENCE_EMR_FIELDS).not.toContain('exam_findings');
expect(TREATED_EVIDENCE_EMR_FIELDS).not.toContain('disposal');
expect([...TREATED_EVIDENCE_EMR_FIELDS]).toEqual(['illness_desc', 'pre_illness']);
});
it('🔴 ⛔ 不认裸「种植」—— 只认「戴牙」(种植是多阶段,一期/二期/取模都不算完成)', () => {
const pats = TREATED_EVIDENCE_RESTORATION_TERMS.map((t) => t.pattern).join('|');
expect(pats).not.toMatch(/(^|\|)种植(\||$)/);
expect(pats).toContain('戴牙');
});
it('⭐ 量词收「数/多/几」—— 主诉常写"数年余/多年"而非确切数字', () => {
expect(TREATED_EVIDENCE_COMPLETION_RE).toContain('数');
expect(new RegExp(TREATED_EVIDENCE_COMPLETION_RE).test('戴牙数年余')).toBe(true);
expect(new RegExp(TREATED_EVIDENCE_COMPLETION_RE).test('戴用多年')).toBe(true);
});
it('🔴 ⛔ 症状词不能当完成标记(脱落=需要重做,该召回)', () => {
expect(TREATED_EVIDENCE_COMPLETION_RE).not.toContain('脱落');
expect(TREATED_EVIDENCE_COMPLETION_RE).not.toContain('松动');
});
it('类目必须落在缺失牙(K08)的 resolver 家族里,否则本表白配', () => {
const cats = resolverCategoriesFor('K08') as readonly string[];
for (const t of TREATED_EVIDENCE_RESTORATION_TERMS) expect(cats).toContain(t.category);
});
it('🔴 ⛔ 只对缺失牙(K08)开闸 —— 词表是拿 K08 的命中人工核出来的,别的场景没验过', () => {
const on = Object.entries(GAP_FLAGS_BY_PRIMARY)
.filter(([, f]) => f.treatedEvidenceFromEmrText === true)
.map(([code]) => code);
expect(on).toEqual(['K08']);
});
it('三个正则都可编译', () => {
expect(() => new RegExp(TREATED_EVIDENCE_COMPLETION_RE)).not.toThrow();
expect(() => new RegExp(TREATED_EVIDENCE_INTENT_EXCLUDE_RE)).not.toThrow();
for (const t of TREATED_EVIDENCE_RESTORATION_TERMS) {
expect(() => new RegExp(t.pattern)).not.toThrow();
}
});
});
......@@ -102,6 +102,31 @@ services:
environment:
NODE_ENV: production
PORT: 3101
# 🔴 **主进程堆上限** —— ⛔ 别删,删了每 ~4 小时崩一次(2026-08-26 生产实证)。
#
# 症状:`FATAL ERROR: Reached heap limit — JavaScript heap out of memory`,
# 容器 exitCode=139 自动重启。08-23~08-26 崩了 **11 次**,期间**一次完整的
# 全量 plans 都没跑成过**(日志里 `Result(ruier-grp)` 出现 0 次)。
# 根因:每轮增量同步结束会在**主进程内**跑一次 `runAllForHost` 全量召回。
# 补摄把患者灌到 111 万、召回池 54 万后,`selectHits` 一次出 96 万条命中,
# 连同 prefetch 的 latest/snoozed/persona 全在堆里 —— 而主进程启动命令
# (`node --enable-source-maps dist/main.js`)**没带 max-old-space-size**,
# 取 Node 默认 **4144 MB**,装不下。
# ⚠️ 容器**没有内存上限**(HostConfig.Memory=0),宿主 15.36G 也一直有 7-11G 空闲 ——
# ⛔ 别往"容器内存不够"上找,撑爆的是 V8 自己的堆,与宿主余量无关。
# ⚠️ 同一容器里 CLI 那条路(cold-import / recompute-plans)一直是 `-max-old-space-size=8192`,
# **两条路配置不一致**,崩的一直是没配的这条。这里补齐。
#
# ⚠️ 宿主内存账:15.36G 总量,本进程 8G 上限 + CLI 8G 上限 = 16G > 总量。
# 实测两者不会同时顶到上限(补摄期间增量被 sync 并发锁挡下,实测 12 小时零崩溃),
# 且 max-old-space 是**天花板不是预留**(实测峰值:服务 4G / CLI 4.9G)。
# ⛔ 但**别手工并发跑** `recompute-plans --host=<全量>`(它也吃 8G)——
# 那条 CLI **不走 sync 锁**,与服务自己那轮全量撞上时才真有可能把宿主吃满。
#
# ⚠️ 这是**止血不是根治**:池子还在涨,8G 迟早也会到顶。
# 根治是「增量之后别跑全量」—— `recompute-plans` 已经有 `--clinics` 收窄能力,
# 同样的思路搬到 SyncIncrementalSchedulerService 那一处即可。
NODE_OPTIONS: --max-old-space-size=8192
# 覆盖 .env 的 localhost URL → 走 docker 内部网络
DATABASE_URL: postgresql://${POSTGRES_USER:-pac}:${POSTGRES_PASSWORD:-pac}@postgres:5432/${POSTGRES_DB:-pac}?schema=public
REDIS_URL: redis://redis:6379
......
# 增量摄入的回看窗:为什么是 48 小时,能不能缩
> 结论先行:**不能缩,而且根因不在我们这边**。真正的解法是让 DW 提供「入仓时间」列。
> 测量日期 2026-08-31,生产 jvs-dw(113 万患者)。全程只读。
## 1. 现象
每轮增量(2 小时一轮)的账:
```
fetched = 248,112 行 其中 dup = 208,556(84%)是重复
实际只写 transactions = 16,448 / facts = 12,854
耗时 31~35 分钟,占整轮(约 1h57m)的 27%
```
直觉是「84% 白拉,把回看窗从 48h 缩到 24h 就能砍一半」。**测下来这个直觉是错的。**
## 2. 根因:游标跟的量与数据到达顺序无关
- 游标跟的是 **`updated_date`** —— **源 HIS 系统**记录被改动的时刻,≈ 事件发生时间
- 数据何时能被 PAC 查到,取决于 **DW 自己的 ETL**,与 `updated_date` 无关
于是必然存在「`updated_date` 比游标旧、但此刻才到达」的行 —— 游标无论怎么设都追不上。
**回看窗不是冗余,是对上游不可控延迟的唯一防线。**
最干净的证据(退费记录):
```
created_date 08-23 13:05
updated_date 08-23 13:13 ← 源系统 8 分钟内当场写完
received_at 08-25 12:39 ← PAC 47 小时后才看到
```
中间那 47 小时 `updated_date` 一动没动。
## 3. 落库滞后实测(7 天,281,760 条**增量**记录)
| 分位 | 滞后 |
|---|---|
| p50 | 2.39h |
| p95 | 5.80h |
| p99.9 | 15.39h |
| max | **47.43h** |
| 阈值 | 超过的行数 |
|---|---|
| >12h | 4,389 |
| >24h | **31** |
| >36h | **19** |
| >48h | **0** |
**按表拆开后,两类行为截然不同:**
| 表 | 条数 | max 滞后 | >24h |
|---|---|---|---|
| **refund** | 150 | **47.4h** | **27** |
| treatment | 48,880 | 42.2h | 4 |
| image | 48,346 | 19.8h | 0 |
| emr / diagnosis / recommendation | 60,748 | 15.5~15.6h | 0 |
| appointment / encounter / payment | 123,402 | 14.9~15.0h | 0 |
**召回真正依赖的临床表(diagnosis / emr / recommendation / encounter)7 天内最大滞后 15.5h,超 24h 一条没有。**
顶到 48h 边界的是 **refund**,而召回逻辑不看退费。
## 4. ⛔ 为什么仍然不能全局缩窗
1. 缩到 24h 会丢 **31 条**、缩到 36h 会丢 **19 条**。游标一旦推过去就是**永久漏拉**,
下一轮补不回来(见 [[incremental-ingest-history-gap]])。
2. **max = 47.43h 恰好顶在 48h 边界、>48h 为 0** —— 这是**截断的指纹**,不是"刚好够"。
真实分布的尾巴可能更长,我们只是观察不到。
> 🔴 **这个测量的方向是单边的**:它能证明「48h 绰绰有余」,**不能**证明「48h 不够」。
> 而现在数据显示边界被顶到了,连"绰绰有余"都谈不上。别拿单边结论去支撑双边决策。
## 4b. 证据链:确认是 DW 侧没有,不是我们没去问
**决定性判据**(只用 PAC 自己的数据,可复现):一行出现在第 N 轮,而它的 `updated_date`
在第 N−1、N−2… 轮的窗口里**就已经被覆盖** → 证明 DW 当时确实没有这行。
以那条 47.4 小时的退费为例(`updated_date=08-23 13:13`,`received_at=08-25 12:39`):
```
08-23 14:15 sync fetched=231,515 窗口下界 08-21 14:15 ✅ 覆盖 → 没拉到
08-23 16:15 sync fetched=235,101 ✅ → 没拉到
…(共 17 轮增量,全部 success)…
08-24 20:09 FULL fetched=4,914,444 ← 全表重读,不看游标 → 仍没拉到
08-25 02:22 FULL fetched=4,255,188 ← 再次全表重读 → 仍没拉到
08-25 12:15 sync fetched=213,108 窗口下界 08-23 12:15 → **拉到了**
```
**19 次查询(含 2 次把 DW 整个重读一遍)都看不见,第 20 次看见了。**
最后那轮的窗口下界是 08-23 12:15,记录是 13:13 —— **只剩 58 分钟余量**,再晚一轮就永久丢。
### 游标的真实语义(容易读错)
```ts
// cold-import.service.ts
// ⭐ cursor_after = run_start ISO(关键!不是 max(updated_date))
const ignoreCursor = options.incremental === false; // full: 忽略游标,全表重读
```
游标推进用的是**墙钟(本轮启动时刻)**,不是数据里的 `max(updated_date)`
每轮查询下界 = `上轮 run_start − 48h`
⚠️ 原注释的理由「run_start 之后 DW 任何写入,下次增量 `WHERE > run_start` 都能捞回」
**有漏洞**:过滤的是 `updated_date`,不是"DW 何时写的"。若 DW 写入时该行
`updated_date` 已早于 `(游标 − 48h)`,下次也捞不回来 —— 这正是漏拉的机制。
## 4c. 「DW 承诺 2 小时更新」与实测的差距
若承诺成立,滞后应封顶 ~4h(2h 刷新 + 2h 等下一轮)。实测(7 天 / 28.2 万条):
| 表 | 条数 | >6h | >12h | >24h |
|---|---|---|---|---|
| **refund** | 150 | **18.7%** | **18.7%** | **27** |
| **recommendation** | 2,631 | **18.9%** | 4.1% | 0 |
| image | 48,348 | 7.8% | 1.3% | 0 |
| diagnosis | 32,763 | 7.5% | 2.1% | 0 |
| treatment | 48,884 | 6.9% | 2.1% | **4** |
| emr | 25,362 | 6.2% | 1.7% | 0 |
| encounter | 30,490 | 2.5% | 1.8% | 0 |
| payment | 8,548 | 2.5% | 0.2% | 0 |
| appointment | 84,380 | 2.3% | 1.0% | 0 |
**中位数守住了(p50 2.39h),尾巴没守住。****分表差异极大**
(appointment 2.3% vs recommendation 18.9% vs refund 18.7%)——
指向 DW 内部不同表走不同的 ETL 链路/调度,不是整体延迟。
`recommendation` 是召回的第二信号源,18.9% 超 6 小时,直接影响召回时效。
## 4d. 🔴 数据时效与丢失边界是**同一个数**
```
PAC 能拿到的数据:最坏 T+2(47.43h 实测)
丢失边界 :48h
```
不是"DW 最慢两天",而是**比 T+2 更晚的根本进不来**。所以:
- **能看见的迟到**:0~48 小时
- **看不见的迟到**:>48 小时 —— 从每轮窗口掉出去,**永远不会知道它存在过**
### 漏掉的数据不会自动修复
`full:` 全量补摄忽略游标、全表重读,机制上能捞回所有迟到数据。
**但生产的全量补摄是人工临时跑的,没有定时任务**:
```
08-21 ×4 08-22 ×2 08-24~08-26 ×5(批量补摄期间)
最近一次:08-26 07:08
```
**08-26 之后漏掉的至今还漏着。**
### 两个未采纳的兜底(记录备查)
| 方案 | 作用 | 代价 |
|---|---|---|
| **定期全量补摄**(如每周一次凌晨) | 把"可能永久漏"变成"最多漏一周" | 单次拉取 300~500 万行,需挑窗口 |
| **跟 DW 对账**:要某历史时段的行数与 PAC 库比对 | **唯一能回答"到底漏了多少"的办法**,其余都是推断 | 只读,需 DW 配合提供计数 |
## 5. 根治办法(不在我们这边):请 DW 提供入仓时间列
**DW 当前一个入仓时间字段都没有。** 十张表的时间列只有 `created_date` / `updated_date`,
全是源 HIS 的;没有 `etl_date` / `load_time` / `dw_insert_time`,也没有分区列
(`rq` 是业务日期,只到天,已排除)。
若 DW 增加一列由其 ETL 写入、只增不改的时间戳:
| | 现在 | 有入仓时间之后 |
|---|---|---|
| 游标语义 | 事件时间,与到达顺序无关 | **在到达顺序上单调** |
| 漏拉风险 | 靠窗口够大来赌 | **结构上不可能** |
| 每轮 fetched | 248,112 行(84% 重复) | 约 16,000 行 |
| 摄入耗时 | 31~35 分钟 | 预计 5~10 分钟 |
## 6. 我们这边的缓解手段(未采纳)
**按表分设回看窗**:临床表 24h、退费/支付 48h+。拉取量约降一半,摄入 31m → ~20m。
**没做,理由**:
- 引入「每张表一个窗口」的配置复杂度,且必须持续盯着 DW 各表的延迟特性有没有变
- `treatment` 有 4 条 >24h(最大 42.2h),它是判「已治疗」的关键表,漏了会**误召**,
需要单独定窗或先查清成因
- 在根治方案可能拿得到的情况下,性价比不高
## 7. ⛔ 取样陷阱(第一版测错了)
第一版按 `received_at > now() - 7 days` 一刀切,得出 p50 滞后 **43,889 小时(5 年)**
max 14.7 年的荒谬结果。原因:窗口内混进了 08-24~08-26 的 **5 次 `full:` 全量补摄**,
那些行的 `updated_date` 是几年前的历史数据。
**在既有增量又有补摄的系统里,按时间窗取样是不成立的**,必须按事件来源限定
(`sync_logs.triggered_by LIKE 'sync:%'`)。修正后样本从 689 万降到 28 万,结论才成立。
## 8. 复现方式(只读,约 9 秒)
```sql
SELECT count(*), percentile_cont(0.999) WITHIN GROUP (ORDER BY lag_h), max(lag_h),
count(*) FILTER (WHERE lag_h > 24)
FROM (
SELECT EXTRACT(EPOCH FROM (t.received_at - (t.raw_payload->>'updated_date')::timestamptz))/3600 AS lag_h
FROM patient_transactions t JOIN sync_logs s ON s.id = t.sync_log_id
WHERE s.triggered_by LIKE 'sync:%' AND s.started_at > now() - interval '7 days'
AND t.raw_payload ? 'updated_date'
) x WHERE lag_h >= 0;
```
......@@ -151,6 +151,10 @@ export const AppointmentCanonicalSchema = z
]),
appointmentType: z.string().optional().nullable(),
doctorId: z.string().optional().nullable(),
/// 排班资源名 = 约号时约的那位医生的姓名(host 行内快照,jvs-dw resource_name)。
/// ⚠️ 不等于实际接诊医生 —— 改派时以病历为准。
/// ⚠️ 约 1% 排的不是人而是房间/服务,canonical 层原样透传,由 parser 判定后再落 doctor_name。
doctorName: z.string().optional().nullable(),
complaintCategory: z.string().optional().nullable(),
durationMinutes: z.coerce.number().optional().nullable(),
arrivedAt: optionalIsoDateTime,
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment