Commit 2cfc043c by luoqi

fix(sync): fact 变更检测纳入四时间锚 — 宿主纠正事件时间可传导为新版本

FactWriter 此前只比 content hash + status:occurredAt/plannedFor/validFrom/validUntil
不参与 → 宿主纠正事件时间(如 FRIDAY 伪 UTC 整体 -8h)重摄永远 evidence_appended,
时间锚永不更新。时间锚是时间轴/召回窗口 COALESCE(occurred_at,planned_for)/过期 cron
的直接输入,属事实实质 → 纳入等值判定(毫秒精度,双 null 视为等),单写与 bulk 两路同步。
不纳入 title/summary(展示层)。含 supersede 传导测试。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
parent e1697a0d
...@@ -14,6 +14,11 @@ function factVersionKey(sourceUnit: string, subjectId: string): string { ...@@ -14,6 +14,11 @@ function factVersionKey(sourceUnit: string, subjectId: string): string {
return `${sourceUnit}#${subjectId}`; return `${sourceUnit}#${subjectId}`;
} }
/** 时刻相等 — 双方皆缺(null/undefined)视为等;毫秒精度(列是 timestamptz(3))。 */
function sameInstant(a: Date | null | undefined, b: Date | null | undefined): boolean {
return (a?.getTime() ?? null) === (b?.getTime() ?? null);
}
/** /**
* FactWriter — patient_facts 版本流写入器 * FactWriter — patient_facts 版本流写入器
* *
...@@ -21,9 +26,9 @@ function factVersionKey(sourceUnit: string, subjectId: string): string { ...@@ -21,9 +26,9 @@ function factVersionKey(sourceUnit: string, subjectId: string): string {
* *
* 行为(对单 subject_id): * 行为(对单 subject_id):
* 1. 找当前 active 版本(partial UNIQUE 保证唯一) * 1. 找当前 active 版本(partial UNIQUE 保证唯一)
* 2. 计算 draft.content 的稳定 hash; * 2. 计算 draft.content 的稳定 hash;比较 hash + status + 时间锚(occurredAt/plannedFor/validFrom/validUntil)
* ┌─ 跟 active 一致 → 幂等(把当前 transactionId 追加到 active.transactionIds,如未存在) * ┌─ 一致 → 幂等(把当前 transactionId 追加到 active.transactionIds,如未存在)
* └─ 不一致 / 无 active → supersede 旧版本 → 插入新版本 version = old.version + 1(无则 1) * └─ 任一不同 / 无 active → supersede 旧版本 → 插入新版本 version = old.version + 1(无则 1)
* *
* 并发说明: * 并发说明:
* - W2 是单进程顺序写,不需要锁 * - W2 是单进程顺序写,不需要锁
...@@ -117,8 +122,9 @@ export class FactWriter { ...@@ -117,8 +122,9 @@ export class FactWriter {
const latestHash = this.hashContent(latest.content as Prisma.InputJsonValue); const latestHash = this.hashContent(latest.content as Prisma.InputJsonValue);
const sameContent = latestHash === newHash; const sameContent = latestHash === newHash;
const sameStatus = latest.status === draftStatus; const sameStatus = latest.status === draftStatus;
const sameTemporal = this.sameTemporalAnchors(latest, validatedDraft);
if (sameContent && sameStatus) { if (sameContent && sameStatus && sameTemporal) {
// 完全一致 — 幂等。把当前 transaction 加进证据数组(去重) // 完全一致 — 幂等。把当前 transaction 加进证据数组(去重)
if (!latest.transactionIds.includes(transactionId)) { if (!latest.transactionIds.includes(transactionId)) {
await tx.patientFact.update({ await tx.patientFact.update({
...@@ -140,7 +146,7 @@ export class FactWriter { ...@@ -140,7 +146,7 @@ export class FactWriter {
}; };
} }
// 有变化(content 或 status)— 如果 latest 是 active,supersede 它; // 有变化(content / status / 时间锚)— 如果 latest 是 active,supersede 它;
// 终态版本(cancelled/fulfilled/expired/invalidated/superseded)不动,只追新版本。 // 终态版本(cancelled/fulfilled/expired/invalidated/superseded)不动,只追新版本。
if (latest.status === FactStatus.ACTIVE) { if (latest.status === FactStatus.ACTIVE) {
await tx.patientFact.update({ await tx.patientFact.update({
...@@ -258,6 +264,10 @@ export class FactWriter { ...@@ -258,6 +264,10 @@ export class FactWriter {
hash: string; hash: string;
transactionIds: string[]; transactionIds: string[];
sourceUpdatedAt: Date | null; sourceUpdatedAt: Date | null;
occurredAt: Date | null;
plannedFor: Date | null;
validFrom: Date | null;
validUntil: Date | null;
}>(); }>();
for (const [k, latest] of latestBySubject.entries()) { for (const [k, latest] of latestBySubject.entries()) {
liveLatest.set(k, { liveLatest.set(k, {
...@@ -267,6 +277,10 @@ export class FactWriter { ...@@ -267,6 +277,10 @@ export class FactWriter {
hash: this.hashContent(latest.content as Prisma.InputJsonValue), hash: this.hashContent(latest.content as Prisma.InputJsonValue),
transactionIds: latest.transactionIds, transactionIds: latest.transactionIds,
sourceUpdatedAt: latest.sourceUpdatedAt ?? null, sourceUpdatedAt: latest.sourceUpdatedAt ?? null,
occurredAt: latest.occurredAt ?? null,
plannedFor: latest.plannedFor ?? null,
validFrom: latest.validFrom ?? null,
validUntil: latest.validUntil ?? null,
}); });
} }
...@@ -294,7 +308,12 @@ export class FactWriter { ...@@ -294,7 +308,12 @@ export class FactWriter {
continue; continue;
} }
if (live && live.hash === entry.hash && live.status === draftStatus) { if (
live &&
live.hash === entry.hash &&
live.status === draftStatus &&
this.sameTemporalAnchors(live, entry.draft)
) {
// 内容一致 + 状态一致 // 内容一致 + 状态一致
if (live.id && !live.transactionIds.includes(entry.transactionId)) { if (live.id && !live.transactionIds.includes(entry.transactionId)) {
toEvidenceAppend.push({ factId: live.id, transactionId: entry.transactionId }); toEvidenceAppend.push({ factId: live.id, transactionId: entry.transactionId });
...@@ -354,6 +373,10 @@ export class FactWriter { ...@@ -354,6 +373,10 @@ export class FactWriter {
hash: entry.hash, hash: entry.hash,
transactionIds: [entry.transactionId], transactionIds: [entry.transactionId],
sourceUpdatedAt: entry.sourceUpdatedAt ?? null, sourceUpdatedAt: entry.sourceUpdatedAt ?? null,
occurredAt: entry.draft.occurredAt ?? null,
plannedFor: entry.draft.plannedFor ?? null,
validFrom: entry.draft.validFrom ?? null,
validUntil: entry.draft.validUntil ?? null,
}); });
results.push({ results.push({
action: live ? 'superseded' : 'created', action: live ? 'superseded' : 'created',
...@@ -394,6 +417,33 @@ export class FactWriter { ...@@ -394,6 +417,33 @@ export class FactWriter {
} }
/** /**
* 时间锚相等判定 — 变更检测的组成部分(与 content hash / status 并列)。
*
* occurredAt / plannedFor / validFrom / validUntil 是事实实质:时间轴排序、
* 召回窗口(COALESCE(occurred_at, planned_for))、过期 cron(valid_until)都直接消费。
* 宿主纠正事件时间(如时区修复整体平移)时,content 往往不变 — 只比 hash+status
* 会让纠正永远停在 evidence_appended,occurred_at 无法更新(时间锚不参与检测的旧 bug)。
*
* 不纳入:title / summary(展示层,规则引擎不读)、clinicId(如需纳入另行评估)。
*/
private sameTemporalAnchors(
latest: {
occurredAt: Date | null;
plannedFor: Date | null;
validFrom: Date | null;
validUntil: Date | null;
},
draft: Pick<FactDraft, 'occurredAt' | 'plannedFor' | 'validFrom' | 'validUntil'>,
): boolean {
return (
sameInstant(latest.occurredAt, draft.occurredAt) &&
sameInstant(latest.plannedFor, draft.plannedFor) &&
sameInstant(latest.validFrom, draft.validFrom) &&
sameInstant(latest.validUntil, draft.validUntil)
);
}
/**
* 稳定 JSON hash — 递归按 key 排序后 sha256。 * 稳定 JSON hash — 递归按 key 排序后 sha256。
* JSON.stringify 默认按 insertion order,key 顺序不同会算出不同 hash。 * JSON.stringify 默认按 insertion order,key 顺序不同会算出不同 hash。
*/ */
......
/**
* FactWriter 时间锚变更检测 — occurredAt/plannedFor/validFrom/validUntil 纳入 supersede 判定
*
* 背景 bug:变更检测原来只比 content hash + status,宿主纠正事件时间
* (如 FRIDAY Mongo 伪 UTC 修复,occurredAt 整体 -8h)时 content 不变 →
* 重摄永远 evidence_appended,occurred_at 无法更新。
*
* 覆盖:
* - 回归:content+status+时间锚全同 → evidence_appended / unchanged(幂等语义不变)
* - occurredAt 平移(-8h)→ superseded,新版本携带纠正后时间
* - plannedFor / null↔值 变化 → superseded(planned 事实的时间锚在 plannedFor)
* - stale gate 优先:sourceUpdatedAt 更旧 → stale_skipped,时间锚差异不放行
* - 终态版本不动:latest 是 fulfilled + 时间差 → 追新版本,终态行保持原状
* - bulkWrite 与 writeDraft 同语义(混合批 + batch 内链式)
*
* 跑:pnpm test -- fact-writer-temporal
*/
import { FactKind, FactStatus, FactType } from '@pac/types';
import { FactWriter } from '../src/modules/sync/pipeline/fact-writer.service';
import type { FactDraft } from '../src/modules/sync/pipeline/parsers/parser.interface';
// ─────────────────────────────────────────────
// 忠实内存 mock:findFirst/findMany 真按 where 过滤、orderBy version desc、
// $transaction 直接跑 callback,update/create/createMany/updateMany 真改内存行
// ─────────────────────────────────────────────
interface Row {
id: string;
hostId: string;
tenantId: string;
sourceUnit: string;
patientId: string;
subjectId: string;
kind: string;
type: string;
status: string;
version: number;
clinicId: string | null;
occurredAt: Date | null;
sourceUpdatedAt: Date | null;
plannedFor: Date | null;
validFrom: Date | null;
validUntil: Date | null;
title: string | null;
summary: string | null;
content: unknown;
transactionIds: string[];
supersededAt: Date | null;
}
function makeStore(seed: Partial<Row>[] = []) {
let nextId = 1;
const rows: Row[] = [];
const insert = (data: Record<string, unknown>): Row => {
const row: Row = {
id: `fact-${nextId++}`,
hostId: data.hostId as string,
tenantId: data.tenantId as string,
sourceUnit: (data.sourceUnit as string) ?? '',
patientId: data.patientId as string,
subjectId: data.subjectId as string,
kind: data.kind as string,
type: data.type as string,
status: data.status as string,
version: data.version as number,
clinicId: (data.clinicId as string | null) ?? null,
occurredAt: (data.occurredAt as Date | null) ?? null,
sourceUpdatedAt: (data.sourceUpdatedAt as Date | null) ?? null,
plannedFor: (data.plannedFor as Date | null) ?? null,
validFrom: (data.validFrom as Date | null) ?? null,
validUntil: (data.validUntil as Date | null) ?? null,
title: (data.title as string | null) ?? null,
summary: (data.summary as string | null) ?? null,
content: data.content,
transactionIds: (data.transactionIds as string[]) ?? [],
supersededAt: null,
};
rows.push(row);
return row;
};
for (const s of seed) insert(s as Record<string, unknown>);
const matches = (row: Row, where: Record<string, unknown>): boolean => {
for (const [k, v] of Object.entries(where)) {
const rv = (row as unknown as Record<string, unknown>)[k];
if (v !== null && typeof v === 'object' && 'in' in (v as object)) {
if (!(v as { in: unknown[] }).in.includes(rv)) return false;
} else if (rv !== v) {
return false;
}
}
return true;
};
const model = {
findFirst: async (q: { where: Record<string, unknown> }) => {
const hit = rows
.filter((r) => matches(r, q.where))
.sort((a, b) => b.version - a.version)[0];
return hit ? { ...hit } : null;
},
findMany: async (q: { where: Record<string, unknown> }) => {
return rows
.filter((r) => matches(r, q.where))
.sort((a, b) => a.subjectId.localeCompare(b.subjectId) || b.version - a.version)
.map((r) => ({ ...r }));
},
update: async (q: { where: { id: string }; data: Record<string, unknown> }) => {
const row = rows.find((r) => r.id === q.where.id)!;
Object.assign(row, q.data);
return { ...row };
},
updateMany: async (q: { where: { id: { in: string[] } }; data: Record<string, unknown> }) => {
for (const r of rows) if (q.where.id.in.includes(r.id)) Object.assign(r, q.data);
},
create: async (q: { data: Record<string, unknown> }) => ({ ...insert(q.data) }),
createMany: async (q: { data: Record<string, unknown>[] }) => {
for (const d of q.data) insert(d);
},
};
const tx = {
patientFact: model,
// bulkWrite 的 evidence append 走 raw SQL(array_append 去重);内存 mock 等价实现
$executeRaw: async (strings: TemplateStringsArray, ...values: unknown[]) => {
const [transactionId, factId] = values as [string, string];
const row = rows.find((r) => r.id === factId);
if (row && !row.transactionIds.includes(transactionId)) {
row.transactionIds = [...row.transactionIds, transactionId];
}
},
};
const prisma = {
patientFact: model,
$transaction: async (cb: (t: typeof tx) => Promise<unknown>) => cb(tx),
};
return { prisma, rows };
}
// ─────────────────────────────────────────────
// 夹具
// ─────────────────────────────────────────────
const HOST = 'host-1';
const TENANT = 'tenant-1';
const UNIT = 'unit-a';
const PATIENT = 'patient-1';
const T0 = new Date('2026-07-01T10:00:00.000Z');
const T0_MINUS_8H = new Date('2026-07-01T02:00:00.000Z');
/// 库里存的是 zod 校验后的 content(nullableString().default(null) 会补齐 null 字段)
/// — 种子必须用同一形态,否则 hash 恒不等,测不出幂等路径
const VALIDATED_CONTENT = {
encounter_external_id: 'enc-1',
doctor_id: null,
doctor_name: null,
chief_complaint: null,
encounter_type: null,
notes: null,
};
function encounterDraft(overrides: Partial<FactDraft> = {}): FactDraft {
return {
subjectId: 'enc-1',
kind: FactKind.ACTUAL,
type: FactType.ENCOUNTER_RECORD,
occurredAt: T0,
content: { encounter_external_id: 'enc-1' },
...overrides,
};
}
function writer(prisma: unknown): FactWriter {
return new FactWriter(prisma as never);
}
function seedActive(overrides: Partial<Row> = {}): Partial<Row> {
return {
hostId: HOST,
tenantId: TENANT,
sourceUnit: UNIT,
patientId: PATIENT,
subjectId: 'enc-1',
kind: FactKind.ACTUAL,
type: FactType.ENCOUNTER_RECORD,
status: FactStatus.ACTIVE,
version: 1,
occurredAt: T0,
content: { ...VALIDATED_CONTENT },
transactionIds: ['tx-1'],
...overrides,
};
}
const baseInput = {
hostId: HOST,
tenantId: TENANT,
sourceUnit: UNIT,
patientId: PATIENT,
};
// ─────────────────────────────────────────────
// writeDraft
// ─────────────────────────────────────────────
describe('FactWriter.writeDraft | 时间锚变更检测', () => {
test('回归:content+status+时间锚全同 → evidence_appended(新 tx)/ unchanged(同 tx)', async () => {
const { prisma, rows } = makeStore([seedActive()]);
const w = writer(prisma);
const r1 = await w.writeDraft({ ...baseInput, draft: encounterDraft(), transactionId: 'tx-2' });
expect(r1.action).toBe('evidence_appended');
expect(rows).toHaveLength(1);
expect(rows[0]!.transactionIds).toEqual(['tx-1', 'tx-2']);
const r2 = await w.writeDraft({ ...baseInput, draft: encounterDraft(), transactionId: 'tx-2' });
expect(r2.action).toBe('unchanged');
expect(rows).toHaveLength(1);
});
test('occurredAt 平移 -8h(content 不变)→ superseded,新版本带纠正后时间', async () => {
const { prisma, rows } = makeStore([seedActive()]);
const w = writer(prisma);
const r = await w.writeDraft({
...baseInput,
draft: encounterDraft({ occurredAt: T0_MINUS_8H }),
transactionId: 'tx-2',
});
expect(r.action).toBe('superseded');
expect(r.version).toBe(2);
expect(rows).toHaveLength(2);
const [old, next] = rows;
expect(old!.status).toBe(FactStatus.SUPERSEDED);
expect(old!.supersededAt).toBeInstanceOf(Date);
expect(next!.status).toBe(FactStatus.ACTIVE);
expect(next!.occurredAt?.getTime()).toBe(T0_MINUS_8H.getTime());
});
test('plannedFor 变化(occurredAt 双 null)→ superseded — planned 事实的时间锚', async () => {
const { prisma, rows } = makeStore([
seedActive({ occurredAt: null, plannedFor: T0 }),
]);
const w = writer(prisma);
const r = await w.writeDraft({
...baseInput,
draft: encounterDraft({ occurredAt: null, plannedFor: T0_MINUS_8H }),
transactionId: 'tx-2',
});
expect(r.action).toBe('superseded');
expect(rows[1]!.plannedFor?.getTime()).toBe(T0_MINUS_8H.getTime());
});
test('occurredAt null → 值 / 值 → null 都算变化 → superseded', async () => {
const s1 = makeStore([seedActive({ occurredAt: null })]);
const r1 = await writer(s1.prisma).writeDraft({
...baseInput,
draft: encounterDraft({ occurredAt: T0 }),
transactionId: 'tx-2',
});
expect(r1.action).toBe('superseded');
const s2 = makeStore([seedActive({ occurredAt: T0 })]);
const r2 = await writer(s2.prisma).writeDraft({
...baseInput,
draft: encounterDraft({ occurredAt: null }),
transactionId: 'tx-2',
});
expect(r2.action).toBe('superseded');
});
test('stale gate 优先:sourceUpdatedAt 更旧 → stale_skipped,时间锚差异不放行', async () => {
const { prisma, rows } = makeStore([
seedActive({ sourceUpdatedAt: new Date('2026-07-02T00:00:00Z') }),
]);
const r = await writer(prisma).writeDraft({
...baseInput,
draft: encounterDraft({ occurredAt: T0_MINUS_8H }),
transactionId: 'tx-2',
sourceUpdatedAt: new Date('2026-07-01T00:00:00Z'),
});
expect(r.action).toBe('stale_skipped');
expect(rows).toHaveLength(1);
expect(rows[0]!.occurredAt?.getTime()).toBe(T0.getTime());
});
test('终态 latest(fulfilled)+ 时间差 → 追新版本,终态行不动', async () => {
const { prisma, rows } = makeStore([
seedActive({ status: FactStatus.FULFILLED }),
]);
const r = await writer(prisma).writeDraft({
...baseInput,
draft: encounterDraft({ status: FactStatus.FULFILLED, occurredAt: T0_MINUS_8H }),
transactionId: 'tx-2',
});
expect(r.action).toBe('superseded');
expect(rows[0]!.status).toBe(FactStatus.FULFILLED); // 终态版本保持原状
expect(rows[1]!.version).toBe(2);
expect(rows[1]!.occurredAt?.getTime()).toBe(T0_MINUS_8H.getTime());
});
});
// ─────────────────────────────────────────────
// bulkWrite — 与 writeDraft 同语义
// ─────────────────────────────────────────────
describe('FactWriter.bulkWrite | 时间锚变更检测', () => {
test('混合批:全同 → evidence_appended;occurredAt 平移 → superseded', async () => {
const { prisma, rows } = makeStore([
seedActive({ subjectId: 'enc-1' }),
seedActive({ subjectId: 'enc-2', version: 3 }),
]);
const w = writer(prisma);
const results = await w.bulkWrite([
{ ...baseInput, draft: encounterDraft({ subjectId: 'enc-1', content: { encounter_external_id: 'enc-1' } }), transactionId: 'tx-9' },
{ ...baseInput, draft: encounterDraft({ subjectId: 'enc-2', content: { encounter_external_id: 'enc-1' }, occurredAt: T0_MINUS_8H }), transactionId: 'tx-9' },
]);
expect(results.map((r) => r.action)).toEqual(['evidence_appended', 'superseded']);
const enc1 = rows.filter((r) => r.subjectId === 'enc-1');
expect(enc1).toHaveLength(1);
expect(enc1[0]!.transactionIds).toContain('tx-9');
const enc2 = rows.filter((r) => r.subjectId === 'enc-2').sort((a, b) => a.version - b.version);
expect(enc2).toHaveLength(2);
expect(enc2[0]!.status).toBe(FactStatus.SUPERSEDED);
expect(enc2[1]!.version).toBe(4);
expect(enc2[1]!.occurredAt?.getTime()).toBe(T0_MINUS_8H.getTime());
});
test('batch 内链式:同 subject 两 draft 仅时间不同 → 连升两版', async () => {
const { prisma, rows } = makeStore([seedActive()]);
const w = writer(prisma);
const results = await w.bulkWrite([
{ ...baseInput, draft: encounterDraft({ occurredAt: T0_MINUS_8H }), transactionId: 'tx-2' },
{ ...baseInput, draft: encounterDraft({ occurredAt: new Date('2026-07-01T03:00:00Z') }), transactionId: 'tx-3' },
]);
expect(results.map((r) => r.action)).toEqual(['superseded', 'superseded']);
expect(results.map((r) => r.version)).toEqual([2, 3]);
expect(rows).toHaveLength(3);
});
test('bulk stale gate 优先于时间锚差异', async () => {
const { prisma, rows } = makeStore([
seedActive({ sourceUpdatedAt: new Date('2026-07-02T00:00:00Z') }),
]);
const results = await writer(prisma).bulkWrite([
{
...baseInput,
draft: encounterDraft({ occurredAt: T0_MINUS_8H }),
transactionId: 'tx-2',
sourceUpdatedAt: new Date('2026-07-01T00:00:00Z'),
},
]);
expect(results[0]!.action).toBe('stale_skipped');
expect(rows).toHaveLength(1);
});
});
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