Commit 546bad00 by luoqi

feat(plan): 到期自动回池 + 五列不可逆决策快照(按产品裁决)

产品 2026-08-02 的六条裁决,本刀落三条(容量默认 50 / 低意愿门不做 / 不看
plan_executions 是下游口径,随 P3、P5 落)。

## 到期自动回池(产品:也是给主管减负)

新 AssignmentExpiryScheduler,**默认开**(PAC_ASSIGNMENT_EXPIRY=off 可关)——
注意与 PAC_PLAN_AUTO_RECYCLE(默认关)相反,spec 里有断言防搞混。

不做的后果不是"多几条过期单",是**容量口径整体失效**:客服在手只增不减,
几批之后全员触顶、再分就分不下去,而这在主管看来就是分配功能坏了。

四条纪律,全部有断言:
·  **照抄 recycle-scheduler 的 snoozedUntil 守卫** —— 客服约了 6/10 回访、
  plan 已 snooze 到那天,到期也不能收。收了就是系统主动毁掉客服对患者的承诺,
  而且单子还会被别人从池里捞走。已本地实测:一条过期但 snoozed 的单**没被动**。
·  **不写 release_reason** —— 那列只属于客服的处置。退回是"客服看了判断不该我做",
  到期是"客服压根没动",混进一列则退回率的分子分母一起虚高,而"到期未动"这个数
  本身才是主管要的信号(分多了?人不在?)。到期量单独从账本按 reason 数。
·  不清批次归因三列(清了这批的分母就少一条)。
· 新 PlanEventReason.ASSIGNMENT_EXPIRED,不复用 timeout(那是认领超时,另一条路)。

判定时刻可注入(`runExpiry(at?)`),照 sync-incremental 刚立的规矩 ——
测试不靠真实时钟凑时间差。第一版没这么写,当场就复现了同一类 flake。

## 五列不可逆快照(迁移 B)

对抗压测抓出来的:我原以为"拆四值枚举"解决了不可逆问题,**那是错的**。
真正补不回来的是快照 —— 有了它四个值全可推导,反过来不成立。

  dedicated_cs_at_assign        preferences.dedicatedCs 是 upsert 覆盖的当前值
  dedicated_cs_last_visit_at     记**证据不是结论**:"在岗"是滚动窗口判定,
                                今天在岗的人半年后回查变离岗,存 boolean 就无法
                                按分配当时的口径重算(而回访表还在被 reparse 重摄)
  priority_score_at_assign      引擎在 reason 未变时**就地改分、不升版本、不留痕**
  source_confidence_at_assign   它是 score 里的 2× 乘子,不留就分不清"按分选人"
                                实际是不是在"按诊断来源选人"
  selection_mode                rank/explore。探索配额是整套系统里**唯一的因果抓手**
                                (入选本身与结果相关,纯观察数据解不开),不标记等于白留

️ 五列**已同步加进引擎的无条件继承集合**。漏了比不加更危险:引擎重出版本
由数据变化触发 → 与患者活跃度相关 → 与完成率相关,缺失是**系统性偏向**的,
半年后拿到一列 70% 填充率的快照,看着还能用,算出来的结论是错的。

## 写路径改用 VALUES join,不再按桶 updateMany

快照列逐行不同,updateMany 的 data 全桶共用表达不了;逐行 update 是 500 次往返。
一条 `UPDATE ... FROM (VALUES ...)` 两个问题一起解决,且 RETURNING 直接给出
哪些行真落上了,不必再回查一次猜差额。9 变量/行 × 500 = 4,500,远低于上限 32,767。
️ 走原生 SQL 的两个代价手工补上:`@updatedAt` 不触发 → 显式 SET;where 自己写全。

## 一致性硬闸(实测踩出来的)

端到端时发现:一条患者的专属客服是 576,却被标 `dedicated` 分给了 832 ——
标签与事实矛盾,而服务端**刚刚把当时的专属客服取到手**,这是可验证的。
不拦的话 T20 按 assign_strategy 分组时,一批实际铺平的单顶着 dedicated 标签混进
"专属完成率",**没有任何报错**。
️ 只拦这一个方向:spread_*/manual 依赖容量与主管意图,服务端不知道,仍以调用方为准
(错了也能靠 dedicated_cs_at_assign 事后重算 —— 这正是那列的价值)。

## 验证

866 单测(新增 9 条到期断言)+ 本地真实数据端到端:
  快照五列落库     selection_mode=rank/explore、pscore=94.8/91.8、conf=1.0、
                  dcs=832/576、dcs_last=2026-06-24 
  标签不符         → 10001 拒绝,不落库
  到期回收         过期无约定 → 回池(status=active、期限清空、**批次归因三列全在**、
                  release_reason 为空);账本 auto_release + assignment_expired
   snoozed 守卫  过期**但约了回访** → 未被收走,仍 assigned 
  未过期的两条     未动
零 drift(diff 只剩两条与本次无关的既有项)。测试数据已清理。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
parent 300b3f50
-- 分配决策快照五列 —— 记录「当初为什么选了这个人」
--
-- 【为什么必须现在加】这五列记的东西**都会被就地覆盖,事后无法重建**:
-- · patients.preferences.dedicatedCs 摄入时 upsert 覆盖,只留当前值
-- · 「在岗」判定 "近 N 月有回访"是滚动窗口,今天在岗半年后变离岗
-- · priority_score 引擎在 reason 未变时**就地改分、不升版本**
-- · confidenceFactor 同上,且它是 score 里的 2× 乘子
-- · 探索配额标记 不记则随机探索配额白留(无法与正常入选区分)
-- 分配一旦开始跑,每一批没记的都是永久空洞。几十字节 vs 一次不可逆的信息丢失。
--
-- 【为什么不重写 25 万行】五列全部可空、全部无默认 → PG 11+ 纯元数据操作。
-- 与 20260802090000 同一条纪律:任何一列加 DEFAULT 或 NOT NULL,前提当场失效。
-- 存量 NULL 即语义正确(= 快照上线前分配的单 / 自助认领单),**不回填**。
--
-- 【为什么可以用普通 DDL】本次不建任何索引:这五列是**取数时读**的快照,
-- 不是筛选维度(要按它筛才该建索引 —— T17 同一条口径)。
-- 单条 ALTER 带五个 ADD COLUMN,一次拿锁。锁窗口与迁移 A 同量级(1-3 秒)。
--
-- 【失败后】整份 SQL 在一个隐式事务里,不会出现半吊子 schema。
-- npx prisma migrate resolve --rolled-back 20260802140000_add_assignment_decision_snapshot
SET LOCAL lock_timeout = '10s';
ALTER TABLE "followup_plans"
ADD COLUMN "dedicated_cs_at_assign" TEXT,
ADD COLUMN "dedicated_cs_last_visit_at" TIMESTAMPTZ(3),
ADD COLUMN "priority_score_at_assign" DOUBLE PRECISION,
ADD COLUMN "source_confidence_at_assign" DOUBLE PRECISION,
ADD COLUMN "selection_mode" TEXT;
......@@ -1078,12 +1078,45 @@ model FollowupPlan {
releaseReason String? @map("release_reason")
releaseNote String? @map("release_note") @db.Text
/// 这条是怎么落到该客服头上的:dedicated(命中专属客服)/ spread(专属容量不足转铺平)/ manual(主管指定)
/// **必须在分配当时就记,事后永远补不回来** —— 唯一的反推来源
/// patients.preferences.dedicatedCs upsert 覆盖的"当前值",历史不留痕。
/// 用途:反推"专属 vs 铺平的完成率差是否显著",进而决定默认策略要不要改。
/// 这条是怎么落到该客服头上的, AssignStrategy 枚举(四值)
/// **必须在分配当时就记,事后永远补不回来**。用途:反推"专属 vs 铺平的完成率差"
assignStrategy String? @map("assign_strategy")
// ── 分配当时的决策快照(2026-08)——————————————————————————————
// ⚠️⚠️ 这四列的存在理由只有一个:**它们记录的东西会被就地覆盖,事后无法重建**
// 不是为了现在用,是为了半年后回答"当初为什么选了这个人"。几十字节换一次
// 不可逆的信息保全 —— 漏了不会报错,只会在需要它的那天发现整批数据不可用。
// ⚠️ 新增任何一列都**必须同步加进 plan-engine 的「归因列无条件继承」集合**
// 漏了比不加更危险:引擎重出版本由数据变化触发 与患者活跃度相关 与完成率相关,
// 缺失是**系统性偏向**,你会拿到一列 70% 填充率的快照并以为它能用。
/// 分配当时该患者的专属客服(宿主 user id)
/// 有了它,assign_strategy 的四个值全部可推导(空→无专属 / ==assigneededicated /
/// !=assignee 且在岗→溢出 / 不在岗→离岗);反过来不成立。
/// 来源 patients.preferences.dedicatedCs 是摄入时 upsert 覆盖的**当前值**,历史不留痕。
dedicatedCsAtAssign String? @map("dedicated_cs_at_assign")
/// 该专属客服**最近一次回访的时刻**(分配当时算出来的)
/// **证据不是结论**:「在岗」判定是"近 N 月有回访"这个滚动窗口,
/// 今天在岗的人半年后回查会变成离岗 —— boolean **无法按分配当时的口径重算**,
/// 而回访表还在被 reparse 反复重摄。存时间戳成本一样,信息量高一个量级。
dedicatedCsLastVisitAt DateTime? @map("dedicated_cs_last_visit_at") @db.Timestamptz(3)
/// 分配当时驱动排序的那个分。
/// ⚠️ priority_score 会被引擎**就地改分、不升版本、不留痕**(reason 未变时的刷新路径),
/// 且分值口径已被迁移改过一次 —— 半年后库里那个分不是当初做选择时的分。
priorityScoreAtAssign Float? @map("priority_score_at_assign")
/// 分配当时该 plan 头部 reason 的证据可信度档(0.5 / 0.8 / 1.0)
/// 它是 priority_score 里的 2× 乘子,在大格子里能压倒意愿度 ——
/// 不单独留一份就分不清"按分选人"实际是不是在"按诊断来源选人"
sourceConfidenceAtAssign Float? @map("source_confidence_at_assign")
/// 这条是怎么被选进批次的:rank(按排序键正常入选)/ explore(随机探索配额)
/// 探索配额是整套系统里**唯一的因果抓手** —— 没有它,排序键好不好永远只能靠推断
/// (入选本身就与结果相关,观察数据无法解耦)。不记则配额白留。
selectionMode String? @map("selection_mode")
createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz(3)
updatedAt DateTime @updatedAt @map("updated_at") @db.Timestamptz(3)
/// 被新版本取代的时间;status='active'/'assigned' 等其他状态时为 null
......
import { Injectable, Logger, OnModuleInit } from '@nestjs/common';
import { Cron, CronExpression } from '@nestjs/schedule';
import { PlanEventType, PlanEventReason } from '@pac/types';
import { PrismaService } from '../../prisma/prisma.service';
import { recordPlanEventsBulk, computeHeldSeconds } from './plan-event.recorder';
/**
* AssignmentExpiryScheduler —— 分配单到期,系统收回池子。
*
* ── 为什么这件事必须做,而且默认开 ────────────────────────────
* 时效是分配单的一部分(T11:不存在无限期批次),但**光写一个到期时刻不会让任何事发生**。
* 不收回的后果不是"多了几条过期单",是**容量口径整体失效**:
* 客服的「在手」只增不减,几批之后全员触顶,再分就分不下去 ——
* 而这在主管看来就是分配功能坏了。
*
* 产品定调(2026-08):到期自动回池**也是给主管减负** —— 让他不必去追"这单还要不要"。
* 主管在确认单上确认时效的那一下,就是对到期行为的预授权,不违 T8
* (T8 要防的是"助手替主管做决定",不是"主管定好的规则到点执行")。
*
* ── 与既有 RecycleSchedulerService 的关系:两条互不干扰的路 ────
* · RecycleScheduler 看 `recycle_at`,是**认领**的 24h 兜底,生产**未启用**
* · 本服务 看 `assignment_expires_at`,是**分配**的时效,默认启用
* 刻意不合并:两者的语义、开关、口径都不同,合了之后想单独关一边就得加分支,
* 而那个分支迟早写错。
*/
/** 关掉的方法:PAC_ASSIGNMENT_EXPIRY=off。默认开(产品已定)。 */
function isEnabled(): boolean {
return (process.env.PAC_ASSIGNMENT_EXPIRY ?? '').trim().toLowerCase() !== 'off';
}
/** 单轮上限 —— 防积压时一次性打爆事务;剩下的下一轮继续 */
const BATCH_LIMIT = 500;
@Injectable()
export class AssignmentExpiryScheduler implements OnModuleInit {
private readonly logger = new Logger(AssignmentExpiryScheduler.name);
constructor(private readonly prisma: PrismaService) {}
onModuleInit(): void {
this.logger.log(
isEnabled()
? '分配单到期自动回池:已启用(每 10 分钟扫一次);关闭设 PAC_ASSIGNMENT_EXPIRY=off'
: '分配单到期自动回池:已关闭(PAC_ASSIGNMENT_EXPIRY=off)—— 在手量将只增不减',
);
}
/**
* @param at 判定时刻(默认此刻)。显式可注入是为了让判据可测 ——
* 测试不必靠真实时钟凑时间差(高负载下事件循环被拖慢会越过阈值边界,产生间歇性假失败)。
* 同款做法见 `sync-incremental.scheduler.reapStaleRunningLocks`,那条是踩过之后立的规矩。
*/
@Cron(CronExpression.EVERY_10_MINUTES, { name: 'plan-assignment-expiry' })
async runExpiry(at?: Date): Promise<void> {
if (!isEnabled()) return;
const now = at ?? new Date();
const due = await this.prisma.followupPlan.findMany({
where: {
status: 'assigned',
assignmentExpiresAt: { not: null, lt: now },
supersededAt: null,
// ⭐⭐ 照抄 RecycleScheduler 的守卫,理由完全相同:
// 客服约了 6/10 回访、plan 已 snooze 到 6/10 —— 在那之前**绝不能**因到期被收走,
// 否则客服丢了已经对患者承诺过的回访关系,6/10 一到单子还会被别人从池里捞走。
// 回访日过后仍未处理,才允许收。
// ⚠️ 这条不是可选优化。漏了它,分配功能会主动破坏客服已经做出的承诺。
OR: [{ snoozedUntil: null }, { snoozedUntil: { lte: now } }],
},
select: {
id: true, hostId: true, tenantId: true, patientId: true,
assigneeUserId: true, assignedAt: true,
},
take: BATCH_LIMIT,
});
if (due.length === 0) return;
// 分片进事务:每片状态变更与账本同生共死,一片失败不影响其余
const CHUNK = 200;
let recycled = 0;
for (let i = 0; i < due.length; i += CHUNK) {
const chunk = due.slice(i, i + CHUNK);
try {
await this.prisma.$transaction(async (tx) => {
const res = await tx.followupPlan.updateMany({
// 带状态条件 → 并发下若客服刚好提交了执行/自己退回了,本次不生效(幂等)
where: { id: { in: chunk.map((p) => p.id) }, status: 'assigned' },
data: {
status: 'active',
assigneeUserId: null,
assignedAt: null,
recycleAt: null,
assignmentExpiresAt: null,
// ⛔ **不写 release_reason** —— 那一列只属于客服的处置。
// 到期是"客服压根没动",不是"客服判断不该我做";混进去退回率的分子分母一起脏。
// 到期的量单独从 plan_event_logs 按 reason 数。
// ⛔ **不清 assignment_id / assigned_by / assign_strategy** —— 批次归因是历史事实,
// 清了这批的分母就少一条,"分了 60 条其中 8 条到期没人动"就算不出来了。
// ⛔ **绝不动 snoozedUntil** —— 与退回同一条纪律。
},
});
recycled += res.count;
await recordPlanEventsBulk(
tx,
chunk.map((p) => ({
hostId: p.hostId,
tenantId: p.tenantId,
planId: p.id,
patientId: p.patientId,
event: PlanEventType.AUTO_RELEASE,
assigneeUserId: null, // 释放后无人归属
actorUserId: null, // 系统行为
// ⭐ 必须在清空 assignedAt **之前**算(上面 findMany 取的就是清空前的值)
heldSeconds: computeHeldSeconds(p.assignedAt, now),
reason: PlanEventReason.ASSIGNMENT_EXPIRED,
})),
);
});
} catch (err) {
this.logger.error(
`分配到期回收失败(${chunk.length} 条): ${err instanceof Error ? err.message : err}`,
);
}
}
if (recycled > 0) {
this.logger.log(
`分配到期:${recycled} 条超期未处理的分配单已退回召回池(已记账本 reason=assignment_expired)` +
(due.length === BATCH_LIMIT ? `;本轮已达单轮上限 ${BATCH_LIMIT},剩余下一轮继续` : ''),
);
}
}
}
......@@ -651,6 +651,16 @@ export class PlanEngineService {
assignStrategy: latest?.assignStrategy ?? null,
releaseReason: latest?.releaseReason ?? null,
releaseNote: latest?.releaseNote ?? null,
// ⚠️⚠️ 决策快照五列同样**无条件继承**,而且漏了比归因列更隐蔽:
// 引擎重出版本是**数据变化触发**的 → 与患者活跃度相关 → 与完成率相关。
// 所以漏继承丢掉的不是随机的一批,是**系统性偏向活跃患者**的一批 ——
// 半年后拿到一列 70% 填充率的快照,看着还能用,算出来的结论是错的。
// ⭐ 今后**每加一列快照,这里必须同步加一行**(spec 有断言锁死)。
dedicatedCsAtAssign: latest?.dedicatedCsAtAssign ?? null,
dedicatedCsLastVisitAt: latest?.dedicatedCsLastVisitAt ?? null,
priorityScoreAtAssign: latest?.priorityScoreAtAssign ?? null,
sourceConfidenceAtAssign: latest?.sourceConfidenceAtAssign ?? null,
selectionMode: latest?.selectionMode ?? null,
// ⚠️ 唯独 assignmentExpiresAt **跟随归属**,不无条件继承:
// 它不是归因,是"当前这次分配的截止时刻"。归属都没了还留着期限,
// 会让列自身失去自洽(非空 ⟺ 有一次在办的分配),超期口径全得靠调用方
......
......@@ -163,6 +163,32 @@ export class PlanAssignmentService {
);
}
// ── 决策快照取数(事务外,只读)────────────────────────────
// ⭐ 这三份数据都会被覆盖,分配当时不取就永久没了。宁可多三次只读查询,
// 也不要留一列半年后才发现是空的快照。
const patientIds = [...new Set(applicable.map(({ row }) => row.patientId))];
const dedicatedByPatient = await this.dedicatedCsOf(patientIds);
const lastVisitByCs = await this.lastVisitOf([...new Set([...dedicatedByPatient.values()])]);
const confidenceByPlan = await this.headConfidenceOf(applicable.map(({ row }) => row.id));
// ⭐ 一致性硬闸:`dedicated` 的字面意思就是"分给了他自己的专属客服",
// 而服务端刚刚把当时的专属客服取到手 —— 这是**可验证的**,就不该放过去。
// 不拦的后果:T20 按 assign_strategy 分组时,一批实际是铺平的单顶着 dedicated 的标签,
// 算出来的"专属完成率"里混进了铺平的样本,而**没有任何报错**。
// ⚠️ 只拦这一个方向:spread_* 与 manual 的判定依赖容量与主管意图,服务端不知道,
// 仍以调用方为准(错了也能靠 dedicated_cs_at_assign 快照事后重算,这正是那列的价值)。
const mislabeled = applicable.filter(
({ row, item }) =>
item.assignStrategy === 'dedicated' &&
dedicatedByPatient.get(row.patientId) !== item.assigneeUserId,
);
if (mislabeled.length > 0) {
throw new BadRequestException(
`有 ${mislabeled.length} 条标记为「专属承接」,但承接人并不是该患者当前的专属客服。` +
`请让助手重新出确认单(标记与事实不符会污染后续的策略效果分析)。`,
);
}
const now = new Date();
try {
const result = await this.prisma.$transaction(
......@@ -182,74 +208,67 @@ export class PlanAssignmentService {
select: { id: true, expiresAt: true },
});
// ⭐ 按 (客服, 策略, 时效) 分桶 —— 同桶的行 data 完全一样,才能一条 updateMany 打完。
// 不按客服分组、逐条 update 的话,500 条 = 500 次往返。
// ⚠️ 类型必须是 Unchecked* —— Checked 版把 assignmentId 当关系字段(要求写成
// `assignment: { connect: ... }`),而 updateMany 不支持关系写入,只认标量外键。
const buckets = new Map<
string,
{ data: Prisma.FollowupPlanUncheckedUpdateManyInput; ids: string[] }
>();
for (const { row, item } of applicable) {
// ⭐ **一条 SQL 打完全部 500 行**,而不是按客服分桶发多条 updateMany。
// 换法的原因是快照列:priority_score_at_assign / dedicated_cs_at_assign 等
// **逐行不同**,updateMany 的 data 是全桶共用的,表达不了。
// 逐行 update 则是 500 次往返。VALUES join 两个问题一起解决,
// 而且 RETURNING 直接给出**哪些行真的落上了**,不必再回查一次去猜差额。
// bind 变量:每行 9 个 × 500 = 4,500,远低于 PG 上限 32,767。
//
// ⚠️ 走原生 SQL 的两个代价,都必须手工补上:
// ① Prisma 的 `@updatedAt` **不会触发** → 显式 SET updated_at
// ② where 条件要自己写全 —— 这里的 `status='active' AND assignee IS NULL
// AND superseded_at IS NULL` 就是并发安全的全部依据:
// 从上面 findMany 到这一刻若有人抢先认领,该行不满足条件、不被更新,
// RETURNING 里也就没有它。
const values = applicable.map(({ row, item }) => {
const perItemExpires =
item.expiresInDays != null
? endOfDayInHostTimezone(now, item.expiresInDays, tz)
: head.expiresAt;
const key = `${item.assigneeUserId}|${item.assignStrategy}|${perItemExpires.toISOString()}`;
let b = buckets.get(key);
if (!b) {
b = {
data: {
status: 'assigned',
assigneeUserId: item.assigneeUserId,
assignedAt: now,
assignedBy: actor.userId,
assignmentId: head.id,
assignmentExpiresAt: perItemExpires,
assignStrategy: item.assignStrategy,
// 重新有人接手 → 上一次的退回原因不再是当前状态(与 PlanService.assign 同口径)
releaseReason: null,
releaseNote: null,
},
ids: [],
};
buckets.set(key, b);
}
b.ids.push(row.id);
}
const dcs = dedicatedByPatient.get(row.patientId) ?? null;
return Prisma.sql`(
${row.id}::uuid, ${item.assigneeUserId}::text, ${perItemExpires}::timestamptz,
${item.assignStrategy}::text, ${dcs}::text,
${dcs ? (lastVisitByCs.get(dcs) ?? null) : null}::timestamptz,
${row.priorityScore}::double precision,
${confidenceByPlan.get(row.id) ?? null}::double precision,
${item.selectionMode ?? 'rank'}::text
)`;
});
const applied: string[] = [];
const lost: string[] = [];
for (const b of buckets.values()) {
for (let i = 0; i < b.ids.length; i += UPDATE_CHUNK) {
const chunk = b.ids.slice(i, i + UPDATE_CHUNK);
// ⭐ **where 里带状态条件 = 并发安全**:从上面 findMany 到这里之间,
// 若有人抢先认领了其中某条,它的 assignee 已非 null → 本次更新不到它,
// count 与 chunk 长度的差就是被抢走的数量。
// (同款技巧见 recycle-scheduler:带 status 条件的 updateMany。)
const res = await tx.followupPlan.updateMany({
where: {
id: { in: chunk },
status: 'active',
assigneeUserId: null,
supersededAt: null,
},
data: b.data,
});
if (res.count === chunk.length) {
applied.push(...chunk);
} else {
// 差额是谁不知道(updateMany 不返回行),回查一次把落上的挑出来
const ok = await tx.followupPlan.findMany({
where: { id: { in: chunk }, assignmentId: head.id },
select: { id: true },
});
const okSet = new Set(ok.map((r) => r.id));
applied.push(...chunk.filter((id) => okSet.has(id)));
lost.push(...chunk.filter((id) => !okSet.has(id)));
}
}
}
const appliedRows = await tx.$queryRaw<Array<{ id: string }>>(Prisma.sql`
UPDATE followup_plans AS fp SET
status = 'assigned',
assignee_user_id = v.assignee,
assigned_at = ${now}::timestamptz,
assigned_by = ${actor.userId}::text,
assignment_id = ${head.id}::uuid,
assignment_expires_at = v.expires_at,
assign_strategy = v.strategy,
release_reason = NULL,
release_note = NULL,
dedicated_cs_at_assign = v.dedicated_cs,
dedicated_cs_last_visit_at = v.dedicated_cs_last_visit,
priority_score_at_assign = v.priority_score,
source_confidence_at_assign = v.source_confidence,
selection_mode = v.selection_mode,
updated_at = ${now}::timestamptz
FROM (VALUES ${Prisma.join(values)}) AS v(
plan_id, assignee, expires_at, strategy,
dedicated_cs, dedicated_cs_last_visit,
priority_score, source_confidence, selection_mode
)
WHERE fp.id = v.plan_id
AND fp.status = 'active'
AND fp.assignee_user_id IS NULL
AND fp.superseded_at IS NULL
RETURNING fp.id
`);
const applied = appliedRows.map((r) => r.id);
const appliedSet = new Set(applied);
const lost = applicable.map(({ row }) => row.id).filter((id) => !appliedSet.has(id));
if (applied.length === 0) {
// 全被抢走 → 回滚,库里不留空批次(与上面 applicable=0 同一条纪律)
......@@ -504,6 +523,73 @@ export class PlanAssignmentService {
};
}
/**
* 分配当时每个患者的专属客服。
* ⚠️ 只此一刻取得到:`patients.preferences.dedicatedCs` 是摄入 upsert **覆盖**的当前值。
*/
private async dedicatedCsOf(patientIds: string[]): Promise<Map<string, string>> {
if (patientIds.length === 0) return new Map();
const rows = await this.prisma.patient.findMany({
where: { id: { in: patientIds } },
select: { id: true, preferences: true },
});
const out = new Map<string, string>();
for (const r of rows) {
const id = (r.preferences as { dedicatedCs?: { id?: string } } | null)?.dedicatedCs?.id;
if (typeof id === 'string' && id) out.set(r.id, id);
}
return out;
}
/**
* 每个客服**最近一次回访的时刻**。
*
* ⚠️ 必须用 `source_created_at`,**绝不能用 `task_date`** —— 后者含未来排程
* (生产实测最远 2033 年、DW 侧甚至 2121 年),拿它算"最近活跃"会把早已离职的人
* 判成在岗。这条坑 schema.prisma 的 task_date 注释里已写过一次。
*
* ⚠️ 存**时刻**而不是"是否在岗":在岗判定是滚动窗口,今天在岗的人半年后回查会变离岗,
* 存 boolean 就无法按分配当时的口径重算了。
*/
private async lastVisitOf(csIds: string[]): Promise<Map<string, Date>> {
if (csIds.length === 0) return new Map();
const rows = await this.prisma.patientReturnVisit.groupBy({
by: ['taskDirectorId'],
where: { taskDirectorId: { in: csIds }, sourceCreatedAt: { not: null } },
_max: { sourceCreatedAt: true },
});
const out = new Map<string, Date>();
for (const r of rows) {
if (r.taskDirectorId && r._max.sourceCreatedAt) out.set(r.taskDirectorId, r._max.sourceCreatedAt);
}
return out;
}
/**
* 每条 plan 的**头部 reason** 的证据可信度(0.5 影像AI / 0.8 建议 / 1.0 诊断)。
*
* 头部 = priority_score 最高的那条 reason —— plan.priorityScore 就是对它们取 MAX,
* 所以驱动排序的是这一条的分,可信度也该取这一条的。
* ⚠️ 它是分里的 **2× 乘子**,不单独留一份,事后就分不清「按分选人」实际是不是
* 在「按诊断来源选人」。
*/
private async headConfidenceOf(planIds: string[]): Promise<Map<string, number>> {
if (planIds.length === 0) return new Map();
const rows = await this.prisma.planReason.findMany({
where: { planId: { in: planIds } },
select: { planId: true, priorityScore: true, breakdown: true },
orderBy: { priorityScore: 'desc' },
});
const out = new Map<string, number>();
for (const r of rows) {
if (out.has(r.planId)) continue; // 已排序,第一条即头部
const cf = (r.breakdown as { priority?: { confidenceFactor?: number } } | null)?.priority
?.confidenceFactor;
if (typeof cf === 'number') out.set(r.planId, cf);
}
return out;
}
/** host 时区(缺省 Asia/Shanghai —— 现有宿主全在东八区,manifest 里也是必填项) */
private async hostTimezone(hostId: string): Promise<string> {
const host = await this.prisma.host.findUnique({
......
......@@ -6,6 +6,7 @@ import { PlanAssignmentService } from './plan-assignment.service';
import { ExecutionService } from './execution.service';
import { ExecutionCallbackService } from './execution-callback.service';
import { RecycleSchedulerService } from './recycle-scheduler.service';
import { AssignmentExpiryScheduler } from './assignment-expiry.scheduler';
import { PlanEngineService } from './engine/plan-engine.service';
import { ChainComposerService } from './engine/chain-composer.service';
import { TreatmentInitiationRecallScenario } from './engine/scenarios/treatment-initiation-recall.scenario';
......@@ -30,6 +31,7 @@ import { RecallDebugService } from './recall-debug/recall-debug.service';
ExecutionService,
ExecutionCallbackService,
RecycleSchedulerService,
AssignmentExpiryScheduler,
PlanEngineService,
ChainComposerService,
TreatmentInitiationRecallScenario,
......
import { AssignmentExpiryScheduler } from '../src/modules/plan/assignment-expiry.scheduler';
import type { PrismaService } from '../src/prisma/prisma.service';
/**
* 分配单到期自动回池回归。
*
* 产品定调:到期回池是给主管减负(他不必再去追"这单还要不要")。
* 但它是**系统主动把单从客服手里收走**,三条红线错一条都会造成静默的数据/信任损失:
* ① 约好回访的单绝不能被收(snoozedUntil 守卫)—— 收了就是系统主动毁客服对患者的承诺
* ② 不写 release_reason —— 那列只属于客服的处置,到期混进去退回率分子分母一起脏
* ③ 不清批次归因三列 —— 清了这批的分母就少一条
*/
const NOW = new Date('2026-08-10T03:00:00Z');
function makeService(rows: Array<Record<string, unknown>>) {
const captured: {
where?: Record<string, unknown>;
data?: Record<string, unknown>;
updateWhere?: Record<string, unknown>;
} = {};
const events: Array<Record<string, unknown>> = [];
const findMany = jest.fn(async ({ where }: { where: Record<string, unknown> }) => {
captured.where = where;
return rows;
});
const tx = {
followupPlan: {
updateMany: jest.fn(
async (args: { where: Record<string, unknown>; data: Record<string, unknown> }) => {
captured.data = args.data;
captured.updateWhere = args.where;
return { count: rows.length };
},
),
},
planEventLog: {
createMany: jest.fn(async ({ data }: { data: Array<Record<string, unknown>> }) => {
events.push(...data);
return { count: data.length };
}),
},
};
const prisma = {
followupPlan: { findMany },
$transaction: jest.fn(async (fn: (t: typeof tx) => Promise<unknown>) => fn(tx)),
} as unknown as PrismaService;
return { svc: new AssignmentExpiryScheduler(prisma), captured, events, findMany, tx };
}
const PLAN = {
id: 'p1',
hostId: 'h1',
tenantId: 't1',
patientId: 'pat1',
assigneeUserId: 'u-staff',
assignedAt: new Date(NOW.getTime() - 3 * 86400_000), // 3 天前分的
};
describe('AssignmentExpiryScheduler', () => {
const OLD = process.env.PAC_ASSIGNMENT_EXPIRY;
afterEach(() => {
if (OLD === undefined) delete process.env.PAC_ASSIGNMENT_EXPIRY;
else process.env.PAC_ASSIGNMENT_EXPIRY = OLD;
});
test('⭐⭐ 红线①:查询条件必须带 snoozedUntil 守卫(约好回访的单不能被收)', async () => {
delete process.env.PAC_ASSIGNMENT_EXPIRY;
const { svc, captured } = makeService([PLAN]);
await svc.runExpiry(NOW);
// 客服约了 6/10 回访、plan snooze 到 6/10 —— 到期也不能收,
// 收了客服就丢了已对患者承诺的回访关系,而且单子还会被别人从池里捞走
expect(captured.where?.OR).toEqual([
{ snoozedUntil: null },
{ snoozedUntil: { lte: expect.any(Date) } },
]);
// 只收已过期的 assigned
expect(captured.where?.status).toBe('assigned');
expect(captured.where?.assignmentExpiresAt).toMatchObject({ not: null });
expect(captured.where?.supersededAt).toBeNull();
});
test('⭐⭐ 红线②:**不写 release_reason** —— 到期不是客服的处置', async () => {
const { svc, captured } = makeService([PLAN]);
await svc.runExpiry(NOW);
// 退回(客服看了判断"不该我做")与到期(客服压根没动)是两件事。
// 混进同一列,退回率的分子分母一起虚高,而"到期未动"这个数本身才是主管要的信号。
expect(captured.data).not.toHaveProperty('releaseReason');
expect(captured.data).not.toHaveProperty('releaseNote');
});
test('⭐⭐ 红线③:**不清批次归因三列**,但清在办期限', async () => {
const { svc, captured } = makeService([PLAN]);
await svc.runExpiry(NOW);
expect(captured.data).not.toHaveProperty('assignmentId');
expect(captured.data).not.toHaveProperty('assignedBy');
expect(captured.data).not.toHaveProperty('assignStrategy');
expect(captured.data).toMatchObject({
status: 'active',
assigneeUserId: null,
assignmentExpiresAt: null,
});
});
test('⭐ 红线④:绝不动 snoozedUntil(与退回同一条纪律)', async () => {
const { svc, captured } = makeService([PLAN]);
await svc.runExpiry(NOW);
expect(captured.data).not.toHaveProperty('snoozedUntil');
});
test('落账本:auto_release + reason=assignment_expired + 持有时长算得出来', async () => {
const { svc, events } = makeService([PLAN]);
await svc.runExpiry(NOW);
expect(events).toHaveLength(1);
expect(events[0]).toMatchObject({
event: 'auto_release',
reason: 'assignment_expired',
assigneeUserId: null,
actorUserId: null, // 系统行为
});
// 时间界注入 → 精确 3 天,不给宽容区间也不会 flake
expect(events[0]!.heldSeconds).toBe(3 * 86400);
});
test('并发安全:updateMany 的 where 带 status 条件(客服刚提交执行则本次不生效)', async () => {
const { svc, captured } = makeService([PLAN]);
await svc.runExpiry(NOW);
expect(captured.updateWhere?.status).toBe('assigned');
});
test('无到期单 → 不发任何写(纯净轮次零副作用)', async () => {
const { svc, tx } = makeService([]);
await svc.runExpiry(NOW);
expect(tx.followupPlan.updateMany).not.toHaveBeenCalled();
expect(tx.planEventLog.createMany).not.toHaveBeenCalled();
});
test('⭐ 开关 off → 一行都不碰(连查询都不发)', async () => {
process.env.PAC_ASSIGNMENT_EXPIRY = 'off';
const { svc, findMany } = makeService([PLAN]);
await svc.runExpiry(NOW);
expect(findMany).not.toHaveBeenCalled();
});
test('⭐ 默认是**开**的 —— 与 PAC_PLAN_AUTO_RECYCLE(默认关)相反,别搞混', async () => {
delete process.env.PAC_ASSIGNMENT_EXPIRY;
const { svc, findMany } = makeService([PLAN]);
await svc.runExpiry(NOW);
expect(findMany).toHaveBeenCalled();
// 不收回的后果不是"多几条过期单",是容量口径整体失效:在手只增不减,几批后全员触顶
});
});
......@@ -1038,6 +1038,17 @@ export const PlanEventReason = {
CLINIC_MOVED: 'clinic_moved',
/// 该患者本轮召回信号全部消失 → plan 直接关闭,归属一并释放
SIGNALS_CLEARED: 'signals_cleared',
/**
* 分配单到期,系统收回池子。
*
* ⚠️ 与客服**主动退回**是两件事,统计时必须分开:
* 退回(release + ReleaseReason)= 客服看了、判断"不该我做",是**处置**
* 到期(auto_release + 本值) = 客服压根没动,是**没处置**
* 把它算进退回率会让分母虚高、分子虚高,而"到期未动"这个数本身才是主管要的信号
* (分多了?人不在?)。所以到期**不写 followup_plans.release_reason** ——
* 那一列只属于客服的处置。
*/
ASSIGNMENT_EXPIRED: 'assignment_expired',
} as const;
export type PlanEventReason = (typeof PlanEventReason)[keyof typeof PlanEventReason];
......@@ -1047,7 +1058,6 @@ export type PlanEventReasonValue = PlanEventReason | ReleaseReason;
/// ⚠️ 还没登记、等产品决策落地后再加(**加的时候必须补进上面的枚举,不要现场写字符串**):
/// · 分配批次被主管撤销 —— 教条 T21 / 开发规划 P6.4
/// · 分配单到期被收回 —— 开发规划 Q-7(到期后行为未定,MVS 只标红不收单)
// =============================================================
// Plan Scripts / Summaries(异步生成,状态机)
......
......@@ -22,6 +22,16 @@ export const AssignmentItemSchema = z.object({
assignStrategy: AssignStrategySchema,
/// 单条覆盖批次时效(不传 = 跟批次)。相对天数,与批次同口径,服务端统一转绝对时刻
expiresInDays: z.number().int().positive().max(90).optional(),
/**
* 这条是**怎么被选进批次**的。
* rank 按排序键正常入选(默认)
* explore 随机探索配额抽中的
*
* ⭐ 探索配额是整套系统里**唯一的因果抓手**:入选本身与结果相关(高分的人本来就更容易成),
* 纯观察数据永远解不开"是排序键选得准,还是这批人本来就好"。留了配额却不标记,
* 等于白留 —— 事后无法把两组分开。
*/
selectionMode: z.enum(['rank', 'explore']).optional(),
});
export type AssignmentItem = z.infer<typeof AssignmentItemSchema>;
......
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