Commit 8930c3b6 by luoqi

perf(gap): resolvedTeeth 集合式形态 + 对拍工具(默认仍走 legacy)

形态:cand → scope → resolved 预聚合 → 牙位反连接。核心恒等式
  ∃x∈G: t(x) ⋛ a  ⟺  max{t(x)} ⋛ a,分组 G=(患者,牙位)。
13 个分支对 sig 的相关性只有三类:时间门(10)/ 病历号等值(1)/ 无相关(2),
另有「建议优先」一条多个 sig.type 标量谓词(ndx 标志)。

️ 集合式是**独立重写一份**,刻意不与 legacy 共用片段 —— 共用则重构写错的地方
   两边一起错、对拍互相抵消。等价性靠 verify-gap-equivalence 逐行差分来证。

全口码(K05/K07)不进集合式:它们的 lateral PG 本来就会摘掉(零收益),
实测硬套进 gap_cand 还慢 2~3 倍(1309→4439ms / 1388→3065ms)。

新增:
  · src/cli/verify-gap-equivalence.cli.ts —— 两版同一 REPEATABLE READ 快照,
    双向 EXCEPT ALL 差分到 (患者,信号,牙位);--self 自对拍先证工具可信
  · tests/gap-setbased-parity.spec.ts —— 结构对拍,守「分支集合不许走散」
    (加分支只改一边 = 静默错召,tsc 和现有 spec 都发现不了)
  · scenario.buildScenarioSql() 抽成独立方法,让对拍拿到线上跑的那条 SQL 本身
  · PAC_GAP_VARIANT=setbased 切换;默认 legacy

本地实测(30K 库):11 个子场景全部零差异,行数逐个相同;
牙位级 ×1.05~1.73(库小全热,不作为收益判据,以测试机 585K 为准)。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
parent d2a208a3
...@@ -30,6 +30,8 @@ ...@@ -30,6 +30,8 @@
"recompute-persona": "ts-node --transpile-only src/cli/recompute-persona.cli.ts", "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", "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", "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": "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", "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", "timeline": "ts-node --transpile-only src/cli/timeline.cli.ts",
......
/**
* verify-gap-equivalence — gap 计算形态对拍(legacy ↔ setbased)
*
* 为什么必须有这个工具:
* `buildGapCore` 是【召回】与【画像】共用的单一真理源,改错的后果是**静默少召** ——
* 不报错、不炸测试,几个月后才由一线反馈冒出来。而现有 spec(arch-denture /
* polish-implies / restoration-in-place / treated-evidence / review-implies)
* 全是纯 JS 常量与正则断言,**不碰 SQL 行为**,重写后照样全绿 → 对这次改动的保护 ≈ 0。
* 所以正确性只能靠"两版跑同一份数据、逐 (患者×信号×牙位) 差分"来证。
*
* 怎么保证可信:
* ① 两版在**同一个 REPEATABLE READ 事务**里跑 —— 否则并发增量摄入会造出假差异。
* ② SQL 来自 scenario 自己的 `buildScenarioSql()`,不是这里另抄一份 ——
* 另抄就变成"验证我抄得对不对",而不是验证线上行为。
* ③ `--self` 模式让 legacy 跟自己对拍,先证明工具本身可信(必须零差异)。
*
* Usage:
* pnpm verify-gap-equivalence -- --host=jvs-dw # 全部 11 个子场景
* pnpm verify-gap-equivalence -- --host=jvs-dw --sub=impacted_tooth
* pnpm verify-gap-equivalence -- --host=jvs-dw --self # 自对拍(工具自检)
* pnpm verify-gap-equivalence -- --host=jvs-dw --bench # 只测耗时,不差分
*
* 退出码:0 = 零差异;1 = 有差异或出错(可直接进 CI / 部署脚本)。
*/
import { NestFactory } from '@nestjs/core';
import { Prisma } from '@prisma/client';
import { lookupDxTreatment } from '@pac/types';
import { AppModule } from '../app.module';
import { PrismaService } from '../prisma/prisma.service';
import { TreatmentInitiationRecallScenario } from '../modules/plan/engine/scenarios/treatment-initiation-recall.scenario';
import type { ScenarioScope } from '../modules/plan/engine/scenario.interface';
import type { GapVariant } from '../modules/clinical-gap/potential-treatment-gap.sql';
interface Args {
host: string;
sub?: string;
self: boolean;
bench: boolean;
samples: number;
}
function parseArgs(argv: string[]): Args {
const a: Args = { host: 'demo', self: false, bench: false, samples: 20 };
for (const s of argv) {
if (s.startsWith('--host=')) a.host = s.slice('--host='.length);
else if (s.startsWith('--sub=')) a.sub = s.slice('--sub='.length);
else if (s.startsWith('--samples=')) a.samples = Number(s.slice('--samples='.length)) || 20;
else if (s === '--self') a.self = true;
else if (s === '--bench') a.bench = true;
}
return a;
}
interface DiffRow {
side: string;
patient_id: string;
signal_fact_id: string;
tooth: string | null;
}
async function bootstrap(): Promise<number> {
const args = parseArgs(process.argv.slice(2));
// 报告一律走 console:Nest 的 logger 级别会被 createApplicationContext 全局压掉,
// 而这个 CLI 的输出就是它的全部产物 —— 不能被日志级别吃掉。
const out = (m: string): void => console.log(m);
const bad_ = (m: string): void => console.error(m);
const app = await NestFactory.createApplicationContext(AppModule, {
logger: ['warn', 'error'],
});
let bad = 0;
try {
const prisma = app.get(PrismaService);
const scenario = app.get(TreatmentInitiationRecallScenario);
const host = await prisma.host.findUnique({ where: { name: args.host } });
if (!host) throw new Error(`Host '${args.host}' not found`);
const tenants = await prisma.patient.findMany({
where: { hostId: host.id },
select: { tenantId: true },
distinct: ['tenantId'],
});
if (!tenants.length) throw new Error('No tenants found for host');
// 🔴 now 固定一次:两版必须拿同一个时间锚,否则 cooldown 边界上的信号会来回抖。
const now = new Date();
const entries = Object.entries(TreatmentInitiationRecallScenario.SUB_SCENARIOS).filter(
([k]) => !args.sub || k === args.sub,
);
if (!entries.length) throw new Error(`--sub=${args.sub} 不是已知子场景`);
for (const t of tenants) {
const scope: ScenarioScope = { hostId: host.id, tenantId: t.tenantId, now };
out(`▶ host=${args.host} tenant=${t.tenantId} 子场景=${entries.length} 个`);
for (const [subKey, cfg] of entries) {
const rule = lookupDxTreatment(cfg.primaryCode);
if (!rule) throw new Error(`${subKey}.primaryCode=${cfg.primaryCode} 不在 DiagnosisTreatmentMap`);
const leftVariant: GapVariant = 'legacy';
const rightVariant: GapVariant = args.self ? 'legacy' : 'setbased';
const left = scenario.buildScenarioSql(scope, cfg.primaryCode, rule, leftVariant);
const right = scenario.buildScenarioSql(scope, cfg.primaryCode, rule, rightVariant);
// ── 计时:两版各单独跑一次(不在同一事务里,避免第二次白蹭第一次的缓存判断失真;
// 对拍另开一个事务)。work_mem 与线上一致。
const timeOne = async (sql: Prisma.Sql): Promise<{ ms: number; rows: number }> => {
const t0 = Date.now();
const [, rows] = await prisma.$transaction([
prisma.$executeRaw`SET LOCAL work_mem = '256MB'`,
prisma.$queryRaw<{ n: bigint }[]>(
Prisma.sql`SELECT count(*)::bigint AS n FROM (${sql}) q`,
),
]);
return { ms: Date.now() - t0, rows: Number(rows[0]?.n ?? 0) };
};
const l = await timeOne(left);
const r = await timeOne(right);
const speed = r.ms > 0 ? (l.ms / r.ms).toFixed(2) : 'n/a';
out(
` ${subKey.padEnd(24)} ${leftVariant}=${String(l.ms).padStart(7)}ms/${l.rows}行 ` +
`${rightVariant}=${String(r.ms).padStart(7)}ms/${r.rows}行 ×${speed}`,
);
if (args.bench) continue;
// ── 差分:两版同一快照,双向 EXCEPT ALL(ALL 保留重复,行数不一致也会暴露)──
const key = Prisma.raw('patient_id, signal_fact_id, tooth');
const diffSql = Prisma.sql`
WITH lg AS (${left}), sb AS (${right}),
d1 AS (SELECT ${key} FROM lg EXCEPT ALL SELECT ${key} FROM sb),
d2 AS (SELECT ${key} FROM sb EXCEPT ALL SELECT ${key} FROM lg)
SELECT 'only_legacy'::text AS side, ${key} FROM d1
UNION ALL
SELECT 'only_setbased'::text AS side, ${key} FROM d2`;
const diffs = await prisma.$transaction(
async (tx) => {
await tx.$executeRaw`SET TRANSACTION ISOLATION LEVEL REPEATABLE READ`;
await tx.$executeRaw`SET LOCAL work_mem = '256MB'`;
return tx.$queryRaw<DiffRow[]>(diffSql);
},
{ timeout: 60 * 60 * 1000, maxWait: 60_000 },
);
if (diffs.length === 0) {
out(` ${subKey.padEnd(24)} ✅ 零差异`);
} else {
bad++;
const onlyL = diffs.filter((d) => d.side === 'only_legacy').length;
const onlyR = diffs.filter((d) => d.side === 'only_setbased').length;
bad_(
` ${subKey.padEnd(24)} ❌ 差异 ${diffs.length} 条(只在 legacy=${onlyL} / 只在 setbased=${onlyR})`,
);
for (const d of diffs.slice(0, args.samples)) {
bad_(` ${d.side} patient=${d.patient_id} sig=${d.signal_fact_id} tooth=${d.tooth ?? 'NULL'}`);
}
}
}
}
} catch (e) {
bad_(e instanceof Error ? (e.stack ?? e.message) : String(e));
bad++;
} finally {
await app.close();
}
if (bad === 0) {
console.log('══ 全部零差异 ══');
return 0;
}
console.error(`══ ${bad} 个子场景不一致 / 出错 ══`);
return 1;
}
bootstrap().then((code) => process.exit(code));
...@@ -118,6 +118,58 @@ export interface GapCoreInput { ...@@ -118,6 +118,58 @@ export interface GapCoreInput {
cfgFlags: GapCfgFlags; cfgFlags: GapCfgFlags;
allCodes: readonly string[]; // dxCodes ∪ recCodes allCodes: readonly string[]; // dxCodes ∪ recCodes
resolverCats: readonly string[]; // resolverCategoriesFor(primaryCode) resolverCats: readonly string[]; // resolverCategoriesFor(primaryCode)
/**
* 计算形态。**只换形态,不换口径** —— 两者必须逐 (患者×信号×牙位) 完全一致。
* 'legacy' 逐 (患者×信号) 行跑相关子查询(2026-08 之前唯一形态)
* 'setbased' 先按 (患者,牙位) 预聚合全部 resolved 证据,再做反连接
* 默认 legacy。切换与验证见 docs/design/gap-set-based-rewrite-plan.md,
* 对拍工具 src/cli/verify-gap-equivalence.cli.ts —— **改任何一边都要重跑对拍**。
*/
variant?: GapVariant;
}
export type GapVariant = 'legacy' | 'setbased';
/**
* gap 计算形态开关。默认 legacy(逐行相关子查询);PAC_GAP_VARIANT=setbased 切集合式。
*
* ⚠️ 两种形态**必须产出完全相同的行集** —— 切换前后要跑
* `pnpm verify-gap-equivalence --host=<host>`(逐 患者×信号×牙位 差分,必须零差异)。
* 设成环境开关而不是写死,是为了让同一个进程能在同一份数据上跑两版做对拍,
* 也为了生产上万一发现差异能立刻回退,不用重新发版。
*/
export function gapVariant(): GapVariant {
return process.env.PAC_GAP_VARIANT === 'setbased' ? 'setbased' : 'legacy';
}
/**
* 集合式形态的拼装件。消费方主查询变成:
*
* WITH gap_cand AS MATERIALIZED (
* SELECT <自己的投影>, ${candExtraCols}
* FROM patients p JOIN patient_profiles pp … JOIN patient_facts sig …
* WHERE <自己的闸> ${candWhere}
* )${postCtes}
* SELECT <自己的投影(从 c 取)>, ${toothOutput} AS tooth
* FROM gap_cand c ${remJoin}
* WHERE TRUE ${outerWhere}
*
* 全口码(K05/K07)时 postCtes/remJoin/outerWhere 为空、判定全在 candWhere 里 ——
* 因为全口场景根本不做牙位相减(见方案 §0.1:那条 lateral PG 本来就会自动删掉)。
*/
export interface GapSetBasedPieces {
/// gap_cand 的额外投影列(每项前带逗号)
candExtraCols: Prisma.Sql;
/// gap_cand 的 WHERE add-on(外院已治疗闸;全口码另加两个 NOT EXISTS)
candWhere: Prisma.Sql;
/// gap_cand 之后的 CTE 链(前带逗号),牙位级非空 / 全口码为空
postCtes: Prisma.Sql;
/// 主查询的 gap_rem 连接
remJoin: Prisma.Sql;
/// tooth 输出表达式
toothOutput: Prisma.Sql;
/// 主查询 WHERE add-on(牙位级:剩余非空)
outerWhere: Prisma.Sql;
} }
export interface GapCorePieces { export interface GapCorePieces {
...@@ -135,6 +187,8 @@ export interface GapCorePieces { ...@@ -135,6 +187,8 @@ export interface GapCorePieces {
toothOutput: Prisma.Sql; toothOutput: Prisma.Sql;
/// ⑤a gap 判定(WHERE add-on):全口 → NOT EXISTS 同类治疗;有牙位 → 剩余非空 /// ⑤a gap 判定(WHERE add-on):全口 → NOT EXISTS 同类治疗;有牙位 → 剩余非空
gapWhere: Prisma.Sql; gapWhere: Prisma.Sql;
/// 集合式拼装件(variant='setbased' 时非空;legacy 时 undefined)
setBased?: GapSetBasedPieces;
} }
/** /**
...@@ -571,5 +625,401 @@ export function buildGapCore(input: GapCoreInput): GapCorePieces { ...@@ -571,5 +625,401 @@ export function buildGapCore(input: GapCoreInput): GapCorePieces {
lateralJoin, lateralJoin,
toothOutput, toothOutput,
gapWhere, gapWhere,
setBased: input.variant === 'setbased' ? (buildGapSetBased(input) ?? undefined) : undefined,
};
}
// ═══════════════════════════════════════════════════════════════════════════
// 集合式形态(setbased)—— 与上面的 legacy 形态**独立重写一份**
//
// ⚠️ 刻意不与 legacy 共用分支片段。共用意味着"重构时写错的地方两边一起错、对拍互相抵消",
// 那样对拍就失去意义。这里是**独立再推导一遍**,靠 verify-gap-equivalence 逐
// (患者×信号×牙位) 差分来证明两份等价 —— 同 sql/verify-recall.sql 里
// 「vt_codes 重述一份、不一致即暴露」的思路。
// legacy 分支在集合式全量验收通过、生产稳定跑过一轮之后再删。
//
// 核心恒等式: ∃x∈G: t(x) ⋛ a ⟺ max{t(x) : x∈G} ⋛ a
// 分组 G = (patient_id, tooth);组内其余谓词(category/status/正则/牙位纯度)全与 sig 无关,
// 所以能先按组聚出 max(t),再跟每个信号的锚点比大小 —— 一次算完全部患者。
//
// 13 个分支对 sig 的相关性只有三类:
// ① 时间门(10 条) → 聚合键 (pid, tooth),gate = max(时间)
// ② 病历号等值(1 条) → 聚合键 (pid, tooth, enc),enc = 该病历的 emr_external_id
// ③ 无相关(2 条) → gate = 'infinity'(恒过)
// 另有 1 条((c) 建议优先)多一个 sig.type 标量谓词 → ndx 标志位
//
// 🔴 gate 列**永不为 NULL**:时间门分支一律加 `IS NOT NULL`(等价 —— NULL 本来就过不了
// `>= 锚点`),无门分支写死 'infinity'。若让 NULL 表示"无门",全 NULL 组会被 max()
// 聚成 NULL 而当成恒过 → 误销 → 静默少召。这是本次重写最危险的一个坑。
// ═══════════════════════════════════════════════════════════════════════════
/// 集合式分支的统一行形状(顺序即列序,13 个分支 UNION ALL 必须对齐)
/// patient_id | tooth | gate | enc | strict | ndx
const SB_COLS = Prisma.raw('patient_id, tooth, gate, enc, strict, ndx');
/// 无时间门 → 恒过(不能用 NULL,见上面 🔴)
const SB_NO_GATE = Prisma.sql`'infinity'::timestamptz`;
/// 患者收窄:所有分支都必须挂,否则 CTE 会全表扫 patient_facts(生产 34GB heap)
const sbScope = (alias: string): Prisma.Sql =>
Prisma.sql`${Prisma.raw(alias)}.patient_id IN (SELECT patient_id FROM gap_scope)`;
export function buildGapSetBased(input: GapCoreInput): GapSetBasedPieces | null {
const { rule, cfgFlags, allCodes, resolverCats } = input;
// 不变式:excludeIfEverTreated ⟺ wholeMouth(当前只有 K05/K07 两者同时为真)。
// 牙位级场景因此永远是 sig 锚点形态,集合式路径不必处理 latestDxOfCode 那一支。
// 谁给某个牙位级 rule 加了 excludeIfEverTreated,这里必须炸,而不是静默算错。
if (rule.excludeIfEverTreated && !rule.wholeMouth) {
throw new Error(
'buildGapSetBased: excludeIfEverTreated 目前只在 wholeMouth 规则上出现;' +
'牙位级规则要用它,得先给集合式补 latestDxOfCode 分支并重跑对拍。',
);
}
// ══ 全口码(K05/K07):**不进集合式,原样走 legacy** ══
// §0.1 已实测:全口场景 sigToothExpr=NULL → st={},toothOutput/gapWhere 都是
// `CASE WHEN TRUE`,常量折叠后 lat 无人引用 → PG 的 useless-left-join removal
// 把整个 LATERAL 摘掉。也就是说**它们本来就没在跑 resolvedTeeth**,集合式零收益。
//
// 2026-08-30 本地实测,曾试着把它们也套进 gap_cand 统一形态,结果**变慢 2~3 倍**:
// ortho_no_consult 1309ms → 4439ms(×0.29)
// perio_no_srp 1388ms → 3065ms(×0.45)
// MATERIALIZED 挡住了规划器对这两条(判定全是患者级 NOT EXISTS)的原有安排。
// ⛔ 别再为了"形态统一好看"把它们并进来 —— 没收益、纯风险、还慢。
if (rule.wholeMouth) return null;
const sigToothExpr = Prisma.sql`sig.content->>'tooth_position'`; // 全口码已早退,这里必有牙位
// ── gap_cand 的额外列(前缀 gap_ 避免跟消费方自己的投影撞名)──
const candExtraCols = Prisma.sql`,
p.id AS gap_patient_id,
sig.id AS gap_sig_id,
COALESCE(sig.occurred_at, sig.planned_for) AS gap_anchor,
sig.type AS gap_sig_type,
sig.content->>'source_encounter_external_id' AS gap_sig_enc,
COALESCE(${toothArrSql(sigToothExpr, {
dropDeciduous: cfgFlags.excludeDeciduous === true,
dropThirdMolar: cfgFlags.excludeThirdMolar === true,
})}, ARRAY[]::text[]) AS gap_sig_teeth`;
// ── 外院已治疗(回访 result)患者级排除 —— 与 legacy 同口径,放进 gap_cand ──
const externalTreatmentGate = Prisma.sql`AND NOT EXISTS (
SELECT 1 FROM patient_return_visits rv
WHERE rv.patient_id = p.id
AND rv.result ~ ${EXTERNAL_TREATMENT_VISIT_POS_RE}
AND rv.result !~ ${EXTERNAL_TREATMENT_VISIT_NEG_RE}
AND rv.task_date >= COALESCE(sig.occurred_at, sig.planned_for)::date
)`;
const refusalRe = TREATMENT_REFUSAL_SUBTYPE_PATTERNS.join('|');
const refusalImagingRe = TREATMENT_REFUSAL_IMAGING_EXCLUDE_RE;
const refusalSubtypeMatch = (alias: string): Prisma.Sql =>
Prisma.sql`regexp_replace(${Prisma.raw(alias)}.content->>'subtype', ${refusalImagingRe}, '', 'g') ~ ${refusalRe}`;
// ══ 牙位级:13 个分支预聚合 ══
const branches: Prisma.Sql[] = [];
// ① (a) 治疗家族 resolver —— 同牙做了 resolverCats 家族里任一治疗
branches.push(Prisma.sql`
SELECT rtx.patient_id, rtt AS tooth, rtx.occurred_at AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts rtx
CROSS JOIN unnest(${toothArrSql(Prisma.sql`rtx.content->>'tooth_position'`)}) AS rtt
WHERE ${sbScope('rtx')}
AND rtx.type = 'treatment_record' AND rtx.kind = 'actual'
AND rtx.status IN ('active', 'fulfilled')
AND COALESCE(NULLIF(trim(rtx.content->>'tooth_position'), ''), '') != ''
AND rtx.content->>'category' = ANY(${resolverCats}::text[])
AND rtx.occurred_at IS NOT NULL`);
// ② (a'') 桥类修复的牙位盲点 —— 基牙跨度区间内全部牙位计入
branches.push(Prisma.sql`
SELECT btx.patient_id, cov.arch_tooth AS tooth, btx.occurred_at AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts btx
CROSS JOIN LATERAL (
SELECT min(z.idx) AS mi, max(z.idx) AS ma, z.arch
FROM (
SELECT CASE substr(bt,1,1)
WHEN '1' THEN 9 - substr(bt,2,1)::int
WHEN '2' THEN 8 + substr(bt,2,1)::int
WHEN '4' THEN 9 - substr(bt,2,1)::int
WHEN '3' THEN 8 + substr(bt,2,1)::int END AS idx,
CASE WHEN substr(bt,1,1) IN ('1','2') THEN 'U' ELSE 'L' END AS arch
FROM unnest(${toothArrSql(Prisma.sql`btx.content->>'tooth_position'`)}) AS bt
WHERE bt ~ '^[1-4][1-8]$'
) z GROUP BY z.arch HAVING count(*) >= 2
) span
CROSS JOIN LATERAL (
SELECT CASE WHEN span.arch = 'U'
THEN CASE WHEN gi <= 8 THEN '1' || (9 - gi)::text ELSE '2' || (gi - 8)::text END
ELSE CASE WHEN gi <= 8 THEN '4' || (9 - gi)::text ELSE '3' || (gi - 8)::text END
END AS arch_tooth
FROM generate_series(span.mi, span.ma) AS gi
) cov
WHERE ${sbScope('btx')}
AND btx.type = 'treatment_record' AND btx.kind = 'actual'
AND btx.status IN ('active', 'fulfilled')
AND btx.content->>'category' = 'prosthodontic'
AND btx.occurred_at IS NOT NULL`);
// ③ (a''') 检查所见"缺牙但间隙关闭/无修复间隙" → 该牙无修复指征
const noRestorMsgRe = NO_RESTORATION_GAP_EXAM_PATTERNS.join('|');
branches.push(Prisma.sql`
SELECT src.patient_id, nrt AS tooth, src.occurred_at AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM (
SELECT emrx.patient_id, emrx.content->>'exam_findings' AS ef_text, emrx.occurred_at
FROM patient_facts emrx
WHERE ${sbScope('emrx')}
AND emrx.type = 'emr_record' AND emrx.status IN ('active', 'fulfilled')
AND emrx.content->>'exam_findings' ~ '^\\['
AND emrx.content->>'exam_findings' ~ '缺[牙失]'
AND emrx.content->>'exam_findings' ~ ${noRestorMsgRe}
AND emrx.occurred_at IS NOT NULL
) src
CROSS JOIN LATERAL jsonb_array_elements(src.ef_text::jsonb) AS ef
CROSS JOIN unnest(${toothArrSql(Prisma.sql`ef->>'toothPosition'`)}) AS nrt
WHERE (ef->>'message') ~ '缺[牙失]'
AND (ef->>'message') ~ ${noRestorMsgRe}`);
// ④ (a'''') 患者无意愿:该 category 治疗被患者拒绝
branches.push(Prisma.sql`
SELECT rfx.patient_id, rft AS tooth,
COALESCE(rfx.occurred_at, rfx.planned_for) AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts rfx
CROSS JOIN unnest(${toothArrSql(Prisma.sql`rfx.content->>'tooth_position'`)}) AS rft
WHERE ${sbScope('rfx')}
AND rfx.type = 'treatment_record' AND rfx.status IN ('active', 'fulfilled')
AND ${refusalSubtypeMatch('rfx')}
AND rfx.content->>'category' = ANY(${resolverCats}::text[])
AND COALESCE(rfx.occurred_at, rfx.planned_for) IS NOT NULL`);
// ⑤ (b) 同牙位以【最新诊断】为准 —— 🔴 严格 >(其余分支都是 >=)
branches.push(Prisma.sql`
SELECT ldx.patient_id, ldt AS tooth, ldx.occurred_at AS gate,
NULL::text AS enc, TRUE AS strict, FALSE AS ndx
FROM patient_facts ldx
CROSS JOIN unnest(${toothArrSql(Prisma.sql`ldx.content->>'tooth_position'`)}) AS ldt
WHERE ${sbScope('ldx')}
AND ldx.type = 'diagnosis_record' AND ldx.status = 'active'
AND ldx.content->>'code' = ANY(${[...STRUCTURAL_DX_CODE_LIST]}::text[])
AND ldx.occurred_at IS NOT NULL`);
// ⑥ (c) 诊断 vs 建议冲突以建议为准 —— 🔴 只对 sig.type='diagnosis_record' 生效(ndx)
branches.push(Prisma.sql`
SELECT rdx.patient_id, rdt AS tooth,
COALESCE(rdx.occurred_at, rdx.planned_for) AS gate,
NULL::text AS enc, FALSE AS strict, TRUE AS ndx
FROM patient_facts rdx
CROSS JOIN unnest(${toothArrSql(Prisma.sql`rdx.content->>'tooth_position'`)}) AS rdt
WHERE ${sbScope('rdx')}
AND rdx.type = 'recommendation_record' AND rdx.status = 'active'
AND rdx.content->>'code' = ANY(${[...STRUCTURAL_DX_CODE_LIST]}::text[])
AND COALESCE(rdx.occurred_at, rdx.planned_for) IS NOT NULL`);
// ⑦ (a''''') 「复查/复诊」= 修复体在位的证据(按类目过滤,不收整类 review)
const reviewRules = REVIEW_IMPLIES_TREATMENT.filter((r) =>
(resolverCats as readonly string[]).includes(r.category),
);
if (reviewRules.length) {
branches.push(Prisma.sql`
SELECT rvx.patient_id, rvt AS tooth, rvx.occurred_at AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts rvx
CROSS JOIN unnest(${toothArrSql(Prisma.sql`rvx.content->>'tooth_position'`)}) AS rvt
WHERE ${sbScope('rvx')}
AND rvx.type = 'treatment_record' AND rvx.kind = 'actual'
AND rvx.status IN ('active', 'fulfilled')
AND rvx.content->>'category' = 'review'
AND rvx.content->>'subtype' ~ ${reviewRules.map((r) => r.pattern).join('|')}
AND rvx.occurred_at IS NOT NULL`);
}
// ⑧ (a'''''') 缺牙位上的裸「抛光」= 修复体在位(仅 K08 开闸)
if (cfgFlags.polishImpliesRestoration) {
branches.push(Prisma.sql`
SELECT plx.patient_id, plt AS tooth, plx.occurred_at AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts plx
CROSS JOIN unnest(${toothArrSql(Prisma.sql`plx.content->>'tooth_position'`)}) AS plt
WHERE ${sbScope('plx')}
AND plx.type = 'treatment_record' AND plx.kind = 'actual'
AND plx.status IN ('active', 'fulfilled')
AND plx.content->>'category' = ${MISSING_TOOTH_POLISH_EVIDENCE.category}
AND plx.content->>'subtype' ~ ${MISSING_TOOTH_POLISH_EVIDENCE.subtypePattern}
AND COALESCE(array_length(${toothArrSql(Prisma.sql`plx.content->>'tooth_position'`)}, 1), 0)
BETWEEN 1 AND ${MISSING_TOOTH_POLISH_EVIDENCE.maxTeeth}
AND plx.occurred_at IS NOT NULL`);
}
// ⑨ (a''''''') 病历自由文本自证已治疗(治疗记录 ⋈ 同次病历;与 sig 无关)
const evidenceTerms = TREATED_EVIDENCE_RESTORATION_TERMS.filter((t) =>
(resolverCats as readonly string[]).includes(t.category),
);
if (cfgFlags.treatedEvidenceFromEmrText && evidenceTerms.length) {
const emrTextSql = Prisma.raw(
TREATED_EVIDENCE_EMR_FIELDS.map((f) => `COALESCE(emx.content->>'${f}', '')`).join(" || ' ' || "),
);
branches.push(Prisma.sql`
SELECT tvx.patient_id, tvt AS tooth, tvx.occurred_at AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts tvx
CROSS JOIN unnest(${toothArrSql(Prisma.sql`tvx.content->>'tooth_position'`)}) AS tvt
JOIN patient_facts emx
ON emx.patient_id = tvx.patient_id
AND emx.type = 'emr_record'
AND emx.status IN ('active', 'fulfilled')
AND emx.content->>'emr_external_id' = tvx.content->>'source_encounter_external_id'
WHERE ${sbScope('tvx')}
AND tvx.type = 'treatment_record' AND tvx.kind = 'actual'
AND tvx.status IN ('active', 'fulfilled')
AND COALESCE(NULLIF(trim(tvx.content->>'tooth_position'), ''), '') != ''
AND tvx.content->>'category' = ANY(${[
...treatedEvidenceTriggersFor(resolverCats),
]}::text[])
AND (
tvx.content->>'category' <> ALL(${[...TREATED_EVIDENCE_BATCH_CATEGORIES]}::text[])
OR COALESCE(array_length(${toothArrSql(Prisma.sql`tvx.content->>'tooth_position'`)}, 1), 0)
<= ${TREATED_EVIDENCE_BATCH_MAX_TEETH}
)
AND (${emrTextSql}) ~ ${evidenceTerms.map((t) => t.pattern).join('|')}
AND (${emrTextSql}) ~ ${TREATED_EVIDENCE_COMPLETION_RE}
AND (${emrTextSql}) !~ ${TREATED_EVIDENCE_INTENT_EXCLUDE_RE}
${
TREATED_EVIDENCE_SINGLE_ARCH_ONLY
? Prisma.sql`AND NOT (
EXISTS (SELECT 1 FROM unnest(${toothArrSql(Prisma.sql`tvx.content->>'tooth_position'`)}) au WHERE au ~ '^[12]')
AND
EXISTS (SELECT 1 FROM unnest(${toothArrSql(Prisma.sql`tvx.content->>'tooth_position'`)}) al WHERE al ~ '^[34]')
)`
: Prisma.empty
}
AND tvx.occurred_at IS NOT NULL`);
}
// ⑩ (a'''''''') 检查所见写着修复体在位(同颌闸 + 失效词一票否决)
if (cfgFlags.restorationInPlaceFromExam) {
branches.push(Prisma.sql`
SELECT rsrc.patient_id, ript AS tooth, rsrc.occurred_at AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM (
SELECT rix.patient_id, rix.content->>'exam_findings' AS ef_text, rix.occurred_at
FROM patient_facts rix
WHERE ${sbScope('rix')}
AND rix.type = 'emr_record' AND rix.status IN ('active', 'fulfilled')
AND rix.content->>'exam_findings' ~ '^\\['
AND rix.content->>'exam_findings' ~ ${RESTORATION_IN_PLACE_TERMS_RE}
AND rix.occurred_at IS NOT NULL
) rsrc
CROSS JOIN LATERAL jsonb_array_elements(rsrc.ef_text::jsonb) AS rief
CROSS JOIN unnest(${toothArrSql(Prisma.sql`rief->>'toothPosition'`)}) AS ript
WHERE (rief->>'message') ~ ${RESTORATION_IN_PLACE_TERMS_RE}
AND (rief->>'message') ~ ${RESTORATION_IN_PLACE_STATE_RE}
AND regexp_replace(rief->>'message', ${RESTORATION_IN_PLACE_NEG_STRIP_RE}, '', 'g')
!~ ${RESTORATION_IN_PLACE_FAIL_RE}
AND NOT (
EXISTS (SELECT 1 FROM unnest(${toothArrSql(Prisma.sql`rief->>'toothPosition'`)}) ru WHERE ru ~ '^[12]')
AND
EXISTS (SELECT 1 FROM unnest(${toothArrSql(Prisma.sql`rief->>'toothPosition'`)}) rl WHERE rl ~ '^[34]')
)`);
}
// ⑪ (a''''''''') 整颌活动义齿 —— 🔴 唯一走【病历号等值】相关的分支(enc 列)
// legacy: adx.emr_external_id = sig.source_encounter_external_id(无时间门)
// setbased: enc 进聚合键,gate 恒过
if (cfgFlags.archDentureIsRestored) {
branches.push(Prisma.sql`
SELECT adx.patient_id, adt AS tooth, ${SB_NO_GATE} AS gate,
adx.content->>'emr_external_id' AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts adx
CROSS JOIN LATERAL jsonb_array_elements((adx.content->>'exam_findings')::jsonb) AS ade
CROSS JOIN unnest(${toothArrSql(Prisma.sql`ade->>'toothPosition'`)}) AS adt
WHERE ${sbScope('adx')}
AND adx.type = 'emr_record' AND adx.status IN ('active', 'fulfilled')
AND adx.content->>'exam_findings' ~ '^\\['
AND adx.content->>'exam_findings' ~ ${ARCH_DENTURE_PREFILTER_RE}
AND adx.content->>'emr_external_id' IS NOT NULL
AND (ade->>'message') !~ ${ARCH_DENTURE_INTENT_EXCLUDE_RE}
AND (
((ade->>'message') ~ ${ARCH_DENTURE_UPPER_RE} AND adt ~ ${UPPER_ARCH_FIRST_DIGITS_RE})
OR
((ade->>'message') ~ ${ARCH_DENTURE_LOWER_RE} AND adt ~ ${LOWER_ARCH_FIRST_DIGITS_RE})
)`);
}
// ⑫ §E 正畸减数位 —— 与 sig 无关,gate 恒过
if (cfgFlags.excludeOrthoExtractionSites) {
const exTeeth = toothArrSql(Prisma.sql`ex.content->>'tooth_position'`);
branches.push(Prisma.sql`
SELECT ex.patient_id, eet AS tooth, ${SB_NO_GATE} AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts ex
CROSS JOIN unnest(${exTeeth}) AS eet
WHERE ${sbScope('ex')}
AND ex.type = 'treatment_record' AND ex.kind = 'actual' AND ex.status IN ('active','fulfilled')
AND ex.content->>'category' = 'surgical'
AND eet ~ '^[1-4][45]$'
AND NOT EXISTS (
SELECT 1 FROM unnest(${exTeeth}) AS xt WHERE xt !~ '^[1-4][45]$'
)
AND (${exTeeth} @> ARRAY['14','24']::text[]
OR ${exTeeth} @> ARRAY['34','44']::text[])
AND EXISTS (SELECT 1 FROM patient_facts oc WHERE oc.patient_id = ex.patient_id
AND ((oc.type='diagnosis_record' AND oc.status='active' AND oc.content->>'code'='K07')
OR (oc.type='treatment_record' AND oc.content->>'category'='orthodontic')))`);
}
// ⑬ 「建议拔除」让位同牙编码病种诊断 —— 与 sig 无关,gate 恒过
if (cfgFlags.deferToToothDx) {
branches.push(Prisma.sql`
SELECT dxx.patient_id, ddt AS tooth, ${SB_NO_GATE} AS gate,
NULL::text AS enc, FALSE AS strict, FALSE AS ndx
FROM patient_facts dxx
CROSS JOIN unnest(${toothArrSql(Prisma.sql`dxx.content->>'tooth_position'`)}) AS ddt
WHERE ${sbScope('dxx')}
AND dxx.type = 'diagnosis_record' AND dxx.status = 'active'
AND dxx.content->>'code' = ANY(ARRAY['K00','K01','K02','K03','K04','K06','K08','K09']::text[])`);
}
const branchUnion = Prisma.join(branches, '\n UNION ALL\n');
const postCtes = Prisma.sql`,
gap_scope AS MATERIALIZED (
SELECT DISTINCT gap_patient_id AS patient_id FROM gap_cand
),
gap_resolved AS MATERIALIZED (
-- 按 (患者, 牙位, 病历号, 严格性, 需诊断信号) 聚合出每组的 max(时间门)
SELECT patient_id, tooth, max(gate) AS gate, enc, strict, ndx
FROM ( ${branchUnion} ) sb(${SB_COLS})
GROUP BY patient_id, tooth, enc, strict, ndx
),
gap_rem AS MATERIALIZED (
-- 牙位级反连接。WITH ORDINALITY + ORDER BY ord:保持 legacy 的 unnest 自然顺序,
-- 也保留重复牙位 —— tooth 串会落进 plan_reasons 给客服看,顺序变了就是 diff 噪音。
SELECT c.gap_sig_id AS sig_id,
array_agg(u.x ORDER BY u.ord) AS remaining_teeth
FROM gap_cand c
CROSS JOIN LATERAL unnest(c.gap_sig_teeth) WITH ORDINALITY AS u(x, ord)
WHERE NOT EXISTS (
SELECT 1 FROM gap_resolved r
WHERE r.patient_id = c.gap_patient_id
AND r.tooth = u.x
AND (CASE WHEN r.strict THEN r.gate > c.gap_anchor ELSE r.gate >= c.gap_anchor END)
AND (r.enc IS NULL OR r.enc = c.gap_sig_enc)
AND (NOT r.ndx OR c.gap_sig_type = 'diagnosis_record')
)
GROUP BY c.gap_sig_id
)`;
return {
candExtraCols,
candWhere: externalTreatmentGate,
postCtes,
remJoin: Prisma.sql`LEFT JOIN gap_rem ON gap_rem.sig_id = c.gap_sig_id`,
// 🔴 gap_rem 无行 ≠ 空数组:牙位全被解决的 sig 在 gap_rem 里没有行 → COALESCE 补 {}
toothOutput: Prisma.sql`array_to_string(COALESCE(gap_rem.remaining_teeth, ARRAY[]::text[]), ';')`,
outerWhere: Prisma.sql`AND cardinality(COALESCE(gap_rem.remaining_teeth, ARRAY[]::text[])) > 0`,
}; };
} }
import { Injectable } from '@nestjs/common'; import { Injectable } from '@nestjs/common';
import { Prisma } from '@prisma/client';
import { lookupDxTreatment, resolverCategoriesFor } from '@pac/types'; import { lookupDxTreatment, resolverCategoriesFor } from '@pac/types';
import { PrismaService } from '../../prisma/prisma.service'; import { PrismaService } from '../../prisma/prisma.service';
import { import {
buildGapCore, buildGapCore,
GAP_FLAGS_BY_PRIMARY, GAP_FLAGS_BY_PRIMARY,
GAP_PRIMARY_GROUPS, GAP_PRIMARY_GROUPS,
gapVariant,
} from './potential-treatment-gap.sql'; } from './potential-treatment-gap.sql';
/** /**
...@@ -43,21 +45,23 @@ export class PotentialTreatmentSelector { ...@@ -43,21 +45,23 @@ export class PotentialTreatmentSelector {
if (!rule) continue; if (!rule) continue;
const resolverCats = resolverCategoriesFor(primaryCode) as readonly string[]; const resolverCats = resolverCategoriesFor(primaryCode) as readonly string[];
const cfgFlags = GAP_FLAGS_BY_PRIMARY[primaryCode] ?? {}; const cfgFlags = GAP_FLAGS_BY_PRIMARY[primaryCode] ?? {};
const gap = buildGapCore({ rule, cfgFlags, allCodes, resolverCats }); const gap = buildGapCore({ rule, cfgFlags, allCodes, resolverCats, variant: gapVariant() });
const rows = await this.prisma.$queryRaw<RawGapRow[]>` // 投影列(两形态共用;tooth 单列,取法不同)
SELECT const projection = Prisma.sql`
sig.id AS fact_id, sig.id AS fact_id,
sig.content->>'code' AS code, sig.content->>'code' AS code,
sig.content->>'name_zh' AS name_zh, sig.content->>'name_zh' AS name_zh,
sig.type AS signal_type, sig.type AS signal_type,
${gap.toothOutput} AS tooth,
sig.content->>'confidence' AS confidence, sig.content->>'confidence' AS confidence,
EXTRACT(DAY FROM ${now}::timestamptz - COALESCE(sig.occurred_at, sig.planned_for))::int AS days_since, 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 FROM patients p
JOIN patient_facts sig ON sig.patient_id = p.id JOIN patient_facts sig ON sig.patient_id = p.id
${gap.lateralJoin} ${joinAddon}
WHERE p.host_id = ${hostId}::uuid WHERE p.host_id = ${hostId}::uuid
AND p.tenant_id = ${tenantId} AND p.tenant_id = ${tenantId}
AND p.id = ${patientId}::uuid AND p.id = ${patientId}::uuid
...@@ -68,8 +72,26 @@ export class PotentialTreatmentSelector { ...@@ -68,8 +72,26 @@ export class PotentialTreatmentSelector {
AND COALESCE(sig.occurred_at, sig.planned_for) IS NOT NULL AND COALESCE(sig.occurred_at, sig.planned_for) IS NOT NULL
${gap.restorationIneligibleFrag} ${gap.restorationIneligibleFrag}
${gap.congenitalFrag} ${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) { for (const r of rows) {
out.push({ out.push({
primaryCode, primaryCode,
......
...@@ -10,6 +10,7 @@ import { ...@@ -10,6 +10,7 @@ import {
treatmentCategoryNameZhFor, treatmentCategoryNameZhFor,
recommendedCategoriesForAge, recommendedCategoriesForAge,
refineCategoriesForDiagnosis, refineCategoriesForDiagnosis,
type DxTreatmentRule,
} from '@pac/types'; } from '@pac/types';
import { PrismaService } from '../../../../prisma/prisma.service'; import { PrismaService } from '../../../../prisma/prisma.service';
import type { import type {
...@@ -20,7 +21,14 @@ import type { ...@@ -20,7 +21,14 @@ import type {
import { calcPriority } from '../priority-scorer'; import { calcPriority } from '../priority-scorer';
import { toothSet } from '../../../sync/pipeline/parsers/tooth-position.util'; import { toothSet } from '../../../sync/pipeline/parsers/tooth-position.util';
// ⭐ gap 核心单一真理源(召回 + 潜在治疗画像共用;SQL 逻辑搬此,本文件只组装) // ⭐ gap 核心单一真理源(召回 + 潜在治疗画像共用;SQL 逻辑搬此,本文件只组装)
import { buildGapCore, GAP_FLAGS_BY_PRIMARY, GAP_PRIMARY_GROUPS } from '../../../clinical-gap/potential-treatment-gap.sql'; import {
buildGapCore,
GAP_FLAGS_BY_PRIMARY,
GAP_PRIMARY_GROUPS,
gapVariant,
type GapVariant,
} from '../../../clinical-gap/potential-treatment-gap.sql';
/** /**
* 潜在治疗新链召回(treatment_initiation_recall)— v2.1 重写 * 潜在治疗新链召回(treatment_initiation_recall)— v2.1 重写
...@@ -255,35 +263,31 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin { ...@@ -255,35 +263,31 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin {
// 子场景跑 SQL + 算 6 因子分 // 子场景跑 SQL + 算 6 因子分
// ───────────────────────────────────────────────────────── // ─────────────────────────────────────────────────────────
private async runSubScenario(
/**
* 构造单个子场景的召回 SQL。**抽成独立方法是为了让对拍工具能拿到线上跑的那条 SQL 本身** ——
* verify-gap-equivalence 用同一份代码生成 legacy / setbased 两版再逐行差分,
* 不另写一份查询(另写就变成"验证我抄得对不对",而不是验证线上行为)。
*/
buildScenarioSql(
scope: ScenarioScope, scope: ScenarioScope,
subKey: string, primaryCode: string,
cfg: (typeof TreatmentInitiationRecallScenario.SUB_SCENARIOS)[keyof typeof TreatmentInitiationRecallScenario.SUB_SCENARIOS], rule: DxTreatmentRule,
): Promise<ScenarioHit[]> { variant: GapVariant,
// 临床窗口 / 类别 / 紧迫临界从 canonical-codes.DiagnosisTreatmentMap 单一真理源读 ): Prisma.Sql {
const rule = lookupDxTreatment(cfg.primaryCode);
if (!rule) {
throw new Error(
`SUB_SCENARIOS[${subKey}].primaryCode=${cfg.primaryCode} 在 DiagnosisTreatmentMap 中找不到 — ` +
`检查 canonical-codes.ts(单一真理源)`,
);
}
const start = rule.cooldownDays;
const goldenRange: [number, number] = [rule.cooldownDays, rule.windowDays];
// 码分组 → 单一真理源 GAP_PRIMARY_GROUPS(召回 + 潜在治疗画像共用,不在 SUB_SCENARIOS 内联) // 码分组 → 单一真理源 GAP_PRIMARY_GROUPS(召回 + 潜在治疗画像共用,不在 SUB_SCENARIOS 内联)
const grp = GAP_PRIMARY_GROUPS[cfg.primaryCode] ?? { dxCodes: [], recCodes: [] }; const grp = GAP_PRIMARY_GROUPS[primaryCode] ?? { dxCodes: [], recCodes: [] };
const dxCodes = grp.dxCodes as readonly string[]; const dxCodes = grp.dxCodes as readonly string[];
const recCodes = grp.recCodes as readonly string[]; const recCodes = grp.recCodes as readonly string[];
const allCodes = [...dxCodes, ...recCodes]; const allCodes = [...dxCodes, ...recCodes];
// §E gap 修正 flag → 单一真理源 GAP_FLAGS_BY_PRIMARY(召回 + 潜在治疗画像共用) // §E gap 修正 flag → 单一真理源 GAP_FLAGS_BY_PRIMARY(召回 + 潜在治疗画像共用)
const cfgFlags = GAP_FLAGS_BY_PRIMARY[cfg.primaryCode] ?? {}; const cfgFlags = GAP_FLAGS_BY_PRIMARY[primaryCode] ?? {};
// ⭐ 两个口径分开(单一真理源 canonical-codes): // ⭐ 两个口径分开(单一真理源 canonical-codes):
// expectedCats = rule.categories(窄,主治疗)→ 展示"未启动 X" + 触发预期 + ⑤d 主诉匹配 // expectedCats = rule.categories(窄,主治疗)→ 展示"未启动 X" + 触发预期 + ⑤d 主诉匹配
// resolverCats = resolverCategoriesFor(宽,治疗家族)→ ⑤a "已解决" 判定 // resolverCats = resolverCategoriesFor(宽,治疗家族)→ ⑤a "已解决" 判定
// 结构码(K02/K03/K08…)= 任何局部结构治疗都算(充填/根管/冠桥/种植/外科/美学/儿牙); // 结构码(K02/K03/K08…)= 任何局部结构治疗都算(充填/根管/冠桥/种植/外科/美学/儿牙);
// 牙周/正畸(K05/K06/K07)沿用各自 categories。见 canonical-codes.resolverCategoriesFor。 // 牙周/正畸(K05/K06/K07)沿用各自 categories。见 canonical-codes.resolverCategoriesFor。
const expectedCats = rule.categories as readonly string[]; const resolverCats = resolverCategoriesFor(primaryCode) as readonly string[];
const resolverCats = resolverCategoriesFor(cfg.primaryCode) as readonly string[];
// 收窄(可空,两种粒度): // 收窄(可空,两种粒度):
// - scope.patientId 单患者(详情页"刷新"):O(全租户) → O(1) // - scope.patientId 单患者(详情页"刷新"):O(全租户) → O(1)
...@@ -307,7 +311,7 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin { ...@@ -307,7 +311,7 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin {
// ⭐ gap 核心(sig 牙位 / resolved / remaining + ⑤a 判定 + 废用牙/先天剔除)抽到共享模块 // ⭐ gap 核心(sig 牙位 / resolved / remaining + ⑤a 判定 + 废用牙/先天剔除)抽到共享模块
// potential-treatment-gap.sql —— 召回与潜在治疗画像【单一真理源】,SQL 逻辑零改动只搬家。 // potential-treatment-gap.sql —— 召回与潜在治疗画像【单一真理源】,SQL 逻辑零改动只搬家。
// 召回在此基础上再加时间门(④ cooldown / ⑤b 预约 / ⑤d entered / ⑤f 到诊)+ 6 因子打分。 // 召回在此基础上再加时间门(④ cooldown / ⑤b 预约 / ⑤d entered / ⑤f 到诊)+ 6 因子打分。
const gap = buildGapCore({ rule, cfgFlags, allCodes, resolverCats }); const gap = buildGapCore({ rule, cfgFlags, allCodes, resolverCats, variant });
// ╔═════════════════════════════════════════════════════════════════════╗ // ╔═════════════════════════════════════════════════════════════════════╗
// ║ 召回 SQL 完整解读(initiation = 潜在治疗新链召回) ║ // ║ 召回 SQL 完整解读(initiation = 潜在治疗新链召回) ║
...@@ -374,8 +378,8 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin { ...@@ -374,8 +378,8 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin {
// ║ 输出:每个命中 (patient × sig) 一行,后段 byPatient Map 去重 ║ // ║ 输出:每个命中 (patient × sig) 一行,后段 byPatient Map 去重 ║
// ║ 只留 daysSince 最大那条(最早诊断 = 最有召回价值) ║ // ║ 只留 daysSince 最大那条(最早诊断 = 最有召回价值) ║
// ╚═════════════════════════════════════════════════════════════════════╝ // ╚═════════════════════════════════════════════════════════════════════╝
const scenarioSql = Prisma.sql` // ── 投影列(legacy / setbased 两形态共用;tooth 单列,两边取法不同)──
SELECT const projection = Prisma.sql`
p.id AS patient_id, p.id AS patient_id,
p.external_id AS patient_external_id, p.external_id AS patient_external_id,
-- 建议治疗的年龄适配用(只影响"建议做什么",不参与召回筛选/排除;无生日 → NULL 走默认) -- 建议治疗的年龄适配用(只影响"建议做什么",不参与召回筛选/排除;无生日 → NULL 走默认)
...@@ -386,19 +390,39 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin { ...@@ -386,19 +390,39 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin {
-- 诊断原文词(K00 主类目细分用:滞留→外科 / 早失→正畸 / 先天缺→修复…)。 -- 诊断原文词(K00 主类目细分用:滞留→外科 / 早失→正畸 / 先天缺→修复…)。
-- 纯投影,不进任何 WHERE —— 返回行集与加它之前逐行相同。 -- 纯投影,不进任何 WHERE —— 返回行集与加它之前逐行相同。
sig.content->>'name_zh' AS signal_name_zh, sig.content->>'name_zh' AS signal_name_zh,
-- ⭐ 牙位级相减:有牙位信号 → 剩余未治牙位;全口信号 → 原样 NULL(gap 核心,共享模块)
${gap.toothOutput} AS tooth,
sig.content->>'extracted_by' AS extracted_by, sig.content->>'extracted_by' AS extracted_by,
sig.content->>'confidence' AS confidence, sig.content->>'confidence' AS confidence,
sig.content->>'code_source' AS code_source, -- 置信度因子:std_code/name_map=医生 / image_ai / null sig.content->>'code_source' AS code_source, -- 置信度因子:std_code/name_map=医生 / image_ai / null
sig.clinic_id AS clinic_id, sig.clinic_id AS clinic_id,
COALESCE(sig.occurred_at, sig.planned_for) AS signal_occurred_at, COALESCE(sig.occurred_at, sig.planned_for) AS signal_occurred_at,
EXTRACT(DAY FROM ${scope.now}::timestamptz - COALESCE(sig.occurred_at, sig.planned_for))::int AS days_since EXTRACT(DAY FROM ${scope.now}::timestamptz - COALESCE(sig.occurred_at, sig.planned_for))::int AS days_since`;
// setbased 外层从 gap_cand 取同名列(顺序无关,消费方按列名映射)
const outerProjection = Prisma.raw(
[
'patient_id',
'patient_external_id',
'patient_age',
'signal_fact_id',
'signal_type',
'signal_code',
'signal_name_zh',
'extracted_by',
'confidence',
'code_source',
'clinic_id',
'signal_occurred_at',
'days_since',
]
.map((c) => `c.${c}`)
.join(', '),
);
// ── FROM + 全部闸(①隔离 ②合规 ③信号 ④cooldown ⑤b/⑤f/⑤g)——
// joinAddon = legacy 的 gap lateral(setbased 为空);gapAddon = gap 判定
const queryBody = (joinAddon: Prisma.Sql, gapAddon: Prisma.Sql): Prisma.Sql => Prisma.sql`
FROM patients p FROM patients p
JOIN patient_profiles pp ON pp.patient_id = p.id JOIN patient_profiles pp ON pp.patient_id = p.id
JOIN patient_facts sig ON sig.patient_id = p.id JOIN patient_facts sig ON sig.patient_id = p.id
-- ⭐ 按牙相减(sig 牙位 / 已解决 / 剩余未治)— gap 核心,共享模块单一真理源 ${joinAddon}
${gap.lateralJoin}
WHERE p.host_id = ${scope.hostId}::uuid -- ① 隔离闸 WHERE p.host_id = ${scope.hostId}::uuid -- ① 隔离闸
AND p.tenant_id = ${scope.tenantId} -- ① 隔离闸 AND p.tenant_id = ${scope.tenantId} -- ① 隔离闸
${patientFilter} -- 单刷 / 子集收窄(可空) ${patientFilter} -- 单刷 / 子集收窄(可空)
...@@ -413,12 +437,12 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin { ...@@ -413,12 +437,12 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin {
AND sig.type IN ('diagnosis_record', 'recommendation_record') -- ③ 信号类型 AND sig.type IN ('diagnosis_record', 'recommendation_record') -- ③ 信号类型
AND sig.content->>'code' = ANY(${allCodes}::text[]) -- ③ 信号 code 命中 AND sig.content->>'code' = ANY(${allCodes}::text[]) -- ③ 信号 code 命中
AND COALESCE(sig.occurred_at, sig.planned_for) IS NOT NULL -- ④ 时间不为空 AND COALESCE(sig.occurred_at, sig.planned_for) IS NOT NULL -- ④ 时间不为空
AND COALESCE(sig.occurred_at, sig.planned_for) <= ${this.daysAgo(scope.now, start)}::timestamptz -- ④ 过 cooldown AND COALESCE(sig.occurred_at, sig.planned_for) <= ${this.daysAgo(scope.now, rule.cooldownDays)}::timestamptz -- ④ 过 cooldown
-- ④' 废用牙/无功能牙剔除 + ④' §E 先天缺失剔除 + ⑤a 牙位级 gap 判定(全口 NOT EXISTS / 有牙位剩余非空) -- ④' 废用牙/无功能牙剔除 + ④' §E 先天缺失剔除 + ⑤a 牙位级 gap 判定(全口 NOT EXISTS / 有牙位剩余非空)
-- —— 全部抽到 gap 核心(共享模块),召回与潜在治疗画像口径一致 -- —— 全部抽到 gap 核心(共享模块),召回与潜在治疗画像口径一致
${gap.restorationIneligibleFrag} ${gap.restorationIneligibleFrag}
${gap.congenitalFrag} ${gap.congenitalFrag}
${gap.gapWhere} ${gapAddon}
-- (⑤c 同牙位拔除 已折进 resolved_teeth 的 surgical 分支 — 拔了的牙从 remaining 减掉) -- (⑤c 同牙位拔除 已折进 resolved_teeth 的 surgical 分支 — 拔了的牙从 remaining 减掉)
AND NOT EXISTS ( -- ⑤b 排除:患者已有未来预约 AND NOT EXISTS ( -- ⑤b 排除:患者已有未来预约
-- 召回目的 = 让客服建预约。患者已经有未来预约 → 客服不需要再 push,医生到诊现场处理即可 -- 召回目的 = 让客服建预约。患者已经有未来预约 → 客服不需要再 push,医生到诊现场处理即可
...@@ -463,8 +487,50 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin { ...@@ -463,8 +487,50 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin {
SELECT 1 FROM patient_return_visits rv SELECT 1 FROM patient_return_visits rv
WHERE rv.patient_id = p.id WHERE rv.patient_id = p.id
AND rv.task_date > ${scope.now}::date AND rv.task_date > ${scope.now}::date
) )`;
const sb = gap.setBased;
return sb
? // ══ 集合式:候选 → 患者域 → resolved 预聚合 → 牙位反连接 ══
// 见 docs/design/gap-set-based-rewrite-plan.md §4。gap_cand 必须 MATERIALIZED:
// 它被下游扫三次(gap_scope / gap_rem / 主查询),内联会让 ①②③④ 闸重算三遍。
Prisma.sql`
WITH gap_cand AS MATERIALIZED (
SELECT ${projection}${sb.candExtraCols}
${queryBody(Prisma.empty, sb.candWhere)}
)${sb.postCtes}
SELECT ${outerProjection},
${sb.toothOutput} AS tooth
FROM gap_cand c
${sb.remJoin}
WHERE TRUE ${sb.outerWhere}
`
: // ══ 逐行相关子查询(legacy)══
Prisma.sql`
SELECT ${projection},
-- ⭐ 牙位级相减:有牙位信号 → 剩余未治牙位;全口信号 → 原样 NULL(gap 核心,共享模块)
${gap.toothOutput} AS tooth
${queryBody(gap.lateralJoin, gap.gapWhere)}
`; `;
}
private async runSubScenario(
scope: ScenarioScope,
subKey: string,
cfg: (typeof TreatmentInitiationRecallScenario.SUB_SCENARIOS)[keyof typeof TreatmentInitiationRecallScenario.SUB_SCENARIOS],
): Promise<ScenarioHit[]> {
// 临床窗口 / 类别 / 紧迫临界从 canonical-codes.DiagnosisTreatmentMap 单一真理源读
const rule = lookupDxTreatment(cfg.primaryCode);
if (!rule) {
throw new Error(
`SUB_SCENARIOS[${subKey}].primaryCode=${cfg.primaryCode} 在 DiagnosisTreatmentMap 中找不到 — ` +
`检查 canonical-codes.ts(单一真理源)`,
);
}
const start = rule.cooldownDays;
const goldenRange: [number, number] = [rule.cooldownDays, rule.windowDays];
const expectedCats = rule.categories as readonly string[];
const scenarioSql = this.buildScenarioSql(scope, cfg.primaryCode, rule, gapVariant());
// ══════════════════════════════════════════════════════════════════ // ══════════════════════════════════════════════════════════════════
// 📊 2026-08-29 全量重算耗时归因(测试服 585K 患者 / 87,661 命中,空闲机) // 📊 2026-08-29 全量重算耗时归因(测试服 585K 患者 / 87,661 命中,空闲机)
...@@ -673,7 +739,7 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin { ...@@ -673,7 +739,7 @@ export class TreatmentInitiationRecallScenario implements PlanScenarioPlugin {
// 也不知道是 SQL 还是后处理慢。2026-08-29 生产 plan 段从 12 分钟涨到 159 分钟 // 也不知道是 SQL 还是后处理慢。2026-08-29 生产 plan 段从 12 分钟涨到 159 分钟
// 仍未跑完,是临时加采样器才定位到子场景粒度的 —— 那种事不该再来一次。 // 仍未跑完,是临时加采样器才定位到子场景粒度的 —— 那种事不该再来一次。
this.logger.log( this.logger.log(
`[recall] sub=${subKey} code=${cfg.primaryCode} ` + `[recall] sub=${subKey} code=${cfg.primaryCode} gap=${gapVariant()} ` +
`sql=${sqlMs}ms rows=${rows.length} post=${Date.now() - postStart}ms hits=${hits.length}`, `sql=${sqlMs}ms rows=${rows.length} post=${Date.now() - postStart}ms hits=${hits.length}`,
); );
return hits; return hits;
......
/**
* 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/);
});
});
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