# Duty Processor — 日常值班事件处理流水线

基于 Nanobot Agent Loop 框架的 AI 值班事件处理系统。处理**被叫号码为值班室电话**的来电录音，经过 5 步流水线自动生成结构化事件记录。

> 本项目与 emergency_command 是两个独立的并行流程，分别由不同的被叫号码触发，不存在路由关系。

## V8.2 当前基线

- 客户侧部署候选：`Qwen2.5-32B-Instruct`；`GLM-5.1-FP8` 仅作大模型效果参照。
- 第9批正式HTTP直驱：470例（普通292、自然续报178），Qwen综合90%、平均25.929秒、接口失败0；GLM参照综合92%、平均41.008秒、接口失败2。
- Qwen续报判断、续报类型和历史事件匹配均为99.79%，初步验证可用于客户环境部署。
- 业务接口版本仍为v3.0；V8.2是当前算法优化与评测基线。

## 处理流程

```
录音 URL（或 asr_text 直传文本）
    │
    ▼
[Step 1] ASR 转写 + 医疗术语校准 + 文本规范化 ──► normalized_text + 历史事件
    │   ① 通用ASR转写 → ② 医疗术语校准(纠正同音/近音的科室/药品/检查/病名错字)
    │   → ③ 抹平说话人标签 + 数字/日期/时间归一化   (传了 asr_text 则跳过 ASR)
    ▼
[Step 2] 事件类型分类           ──► report_type + report_subtype
    │                               依据：《日常应急值班规则》接报类型
    ▼
[Step 3] 名称摘要生成           ──► event_name + summary
    │                               依据：《日常应急值班规则》名称/摘要模板
    ▼
[Step 4] 接续判断               ──► is_continuation + matched_event + candidates + llm_verdict
    │   规则召回候选：四维内容加权 + 电话一致加分(电话不同不扣分)排序取 top-K
    │   LLM「同一事件」语义裁决主决策(默认；DUTY_LLM_CONTINUATION=0 退为纯规则≥0.7消融)
    ▼
[Step 5] 组装最终结果            ──► daily_duty_event JSON + 全字段溯源
```

> **ASR 说话人区分**：百炼 qwen3-asr-flash 返回不区分说话人的纯文本（无 diarization）；
> 火山引擎支持角色分离。流水线不依赖说话人标签，`strip_speaker_tags` 会抹平
> `说话人1:/坐席:/客户:` 等标签，故"有无说话人标注"下均稳健。

## 事件分类体系

> 严格依据《日常应急值班规则》接报类型定义（与 prompts/type_classify.py 一致）。

| 主类型 | 子类型 |
|--------|--------|
| 日常咨询 | 政策咨询 / 流程咨询 / 信息查询 / 其他咨询 |
| 日常事件处理 | 协调 / 投诉 / 记录 |
| 日常院内医疗急救 | 急会诊 / 卒中 / 紧急用血 / 创伤 / 中毒 / 产科急 / 幼儿急 / 透析急 / 复苏 / 其他 |
| 应急指挥紧急医学救援 | 紧急医学救援（创伤/中毒且同一事件≥10人伤亡） |

## API 接口

| 方法 | 路径 | 说明 |
|------|------|------|
| `POST` | `/api/v3/duty/pipeline` | 编排模式（nanobot agentloop 调度） |
| `POST` | `/api/v3/duty/pipeline_direct` | 直驱模式（确定性按序调工具，推荐） |
| `POST` | `/api/v3/duty/pipeline_direct/stream` | 直驱 SSE 流式（逐 step 推送，demo 体验页用） |
| `POST` | `/api/v3/duty/pipeline/stream` | 编排 SSE 流式（逐工具步推送） |
| `GET` | `/health` | 健康检查 |

### 请求参数

请求体对齐平台**通话记录**字段。工程接入时应完整传入 `caller`/`called`、`callerName`/`calledName`、`callType`、`tmAnswer`/`tmHangup`、`callLength`、`businessId`、`id`；为支持文本冒烟与离线评测，API模型层允许这些字段为空。录音来源 `recordList`/`recording_url` 与 `asr_text` 二选一。旧字段 `recording_time`→`tmAnswer`、`caller_phone`→`caller`、`receiver_phone`→`called`、`receiver_name`→`calledName` 仍兼容。

**方式一：传录音文件（完整流水线）**
```json
{
  "recording_url": "https://example.com/recording.mp4",
  "id": 10086,
  "businessId": "BIZ-10086",
  "caller": "13800000000",
  "callerName": "8A病区医生",
  "called": "01012345678",
  "calledName": "总值班",
  "callType": "callin",
  "tmAnswer": "2025-01-01 08:00:00",
  "tmHangup": "2025-01-01 08:02:30",
  "callLength": 150,
  "asr_provider": "qwen3_asr"
}
```

**方式二：传 ASR 文本（跳过录音/ASR，通话记录仍需带全；推荐测试用）**
```json
{
  "asr_text": "8A病区医生来电，高钙血症患者病情危重，需请肾内科急会诊",
  "id": 10086,
  "businessId": "BIZ-10086",
  "caller": "13800000000",
  "callerName": "8A病区医生",
  "called": "01012345678",
  "calledName": "总值班",
  "callType": "callin",
  "tmAnswer": "2025-01-01 08:00:00",
  "tmHangup": "2025-01-01 08:02:30",
  "callLength": 150
}
```

### 响应示例（直驱模式）

```json
{
  "code": 200,
  "message": "success",
  "data": {
    "daily_duty_event": {
      "report_type": "日常院内医疗急救",
      "report_subtype": "急会诊",
      "event_name": "急会诊@急诊抢救室3床|高钙血症",
      "summary": "2025-03-08 14:30，8A病区医生来电，报告1位高钙血症患者病情危重...",
      "continuation_info": {"is_continuation": false, "status": "NEW_EVENT"},
      "metadata": {
        "execution_mode": "direct",
        "pipeline_duration_ms": 5600,
        "field_sources": {
          "report_type": "classify_report_type·依据《日常应急值班规则》接报类型",
          "event_name": "generate_name_summary·依据《日常应急值班规则》名称模板",
          "summary": "generate_name_summary·依据《日常应急值班规则》摘要模板",
          "continuation_info": "judge_continuation·规则召回候选(四维加权+电话加分排序top-K)→LLM同一事件语义主决策(DUTY_LLM_CONTINUATION=0退纯规则≥0.7消融)"
        }
      }
    }
  }
}
```

### 错误码

| 错误码 | 说明 |
|--------|------|
| 40001 | 录音内容为空（asr_text 与录音来源均为空/纯空白时也返回此码） |
| 40002 | 录音无法识别 |
| 40003 | ASR 服务调用超时 |
| 40004 | ASR 服务内部错误 / 录音下载失败 / 未配置对应 ASR 密钥 |
| 50001 | 流水线总超时（当前服务默认120秒，`/pipeline`、`/pipeline_direct`及流式接口均有护栏） |
| 50002 | 内部处理错误 |

## 源表锚定

| 工具 | 依据的源表 | 说明 |
|------|-----------|------|
| calibrate_asr_text | ASR后处理（通用能力） | 医疗术语同音纠错，只纠错不增删事实，非源表判定 |
| classify_report_type | 《日常应急值班规则》 | 4大类+子类接报类型体系 |
| generate_name_summary | 《日常应急值班规则》 | 名称/摘要模板规则 |
| judge_continuation | 工程规则（无源表） | 规则召回候选：四维内容加权(类型/对象/地点/描述,和1.0)+电话一致加分(0.10)排序取 top-K → LLM「同一事件」语义裁决**主决策**(默认；DUTY_LLM_CONTINUATION=0 退为纯规则 top-1≥0.7 消融) |

## 模型配置

工具模型在 `config.yaml` 配置（与 nanobot 编排模型 `config.json` 独立）：

```yaml
llm:
  base_url: "http://192.168.68.92:8986/v1"
  api_key: "Empty"
  default_model: "Qwen2.5-32B-Instruct"
  models:
    calibrate_asr_text: "Qwen2.5-32B-Instruct"
    classify_report_type: "Qwen2.5-32B-Instruct"
    generate_name_summary: "Qwen2.5-32B-Instruct"
```

> 仓库默认 profile 与第9批 V8.2 Qwen评测一致；`GLM-5.1-FP8` 仅作为效果参照。

正式切换自建模型端点时，工具模型使用完整原子配置：

```bash
export EXP_TOOL_BASE_URL=http://192.168.68.92:8986/v1
export EXP_TOOL_API_KEY=Empty
export EXP_TOOL_MODEL=Qwen2.5-32B-Instruct
```

三个变量必须同时提供，显式环境变量优先于默认 profile。后续切换模型时同时替换模型名、URL 和 Key，并重启服务。编排模型另用 `AGENTLOOP_BASE_URL`、`AGENTLOOP_API_KEY`、`AGENTLOOP_MODEL`；直驱不会读取编排模型配置。Qwen与GLM调用均关闭思考模式。

Demo 使用正式 SSE 流式端点，评测使用正式非流式端点；除此之外两者使用相同
profile、预处理、业务处理链和自然续报历史机制。

## 对比评测摘要

完整结果见 `experiments/REPORT.md`，数据在 `experiments/runs/`。

第9批V8.2直驱全量结果（`cases_duty_v2.json`，470例）：

| 模型 | 综合 | 平均耗时 | 接口失败 | 续报判断/类型/匹配 |
|---|---:|---:|---:|---:|
| Qwen2.5-32B-Instruct | 90% | 25.929s | 0/470 | 99.79% |
| GLM-5.1-FP8（效果参照） | 92% | 41.008s | 2/470 | 98.09% |

Qwen与参照模型综合差距为2个百分点，同时在耗时、接口稳定性和续报链路上满足继续客户侧部署验证的条件。

## 快速开始

```bash
pip install -e .
duty-processor serve --host 0.0.0.0 --port 8001

# 最简调用：只传 asr_text
curl -X POST http://localhost:8001/api/v3/duty/pipeline_direct \
  -H "Content-Type: application/json" \
  -d '{"asr_text":"8A病区医生来电，高钙血症患者病情危重，需请肾内科急会诊"}'
```

## 环境变量

| 变量 | 说明 | 必需 |
|------|------|------|
| `DASHSCOPE_API_KEY` | 默认 `qwen3_asr` 等 DashScope ASR 能力的 API 密钥；不作为工具模型 Key | 按 ASR provider |
| `EXP_TOOL_BASE_URL` / `EXP_TOOL_API_KEY` / `EXP_TOOL_MODEL` | 直驱工具模型完整端点配置，必须同时设置 | 可选（默认使用 Qwen2.5 profile） |
| `AGENTLOOP_BASE_URL` / `AGENTLOOP_API_KEY` / `AGENTLOOP_MODEL` | 编排模型完整端点配置 | 可选 |
| `EXP_MOCK_DUTY_TEXT` | Mock ASR 返回文本（评测用） | 可选 |
| `DUTY_DISABLE_CALIBRATION` | 设 1 关闭 ASR 医疗术语校准（消融对照） | 可选 |
| `DUTY_LLM_CONTINUATION` | 接续的 LLM「同一事件」语义裁决，**默认开启**；设 0 关闭（仅留确定性词面匹配，消融/离线用） | 可选 |
| `DUTY_USE_MOCK_EVENTS` | 设 0 关闭 Mock 历史事件库，改走工程历史事件 repository；未配置 repository 时续接判定退化为「新事件」 | 可选 |
| `VOLCENGINE_ASR_API_KEY` / `VOLCENGINE_ASR_APP_ID` | 仅当走 `volcengine` 音频ASR时需要；默认 `qwen3_asr` 使用 `DASHSCOPE_API_KEY` | 可选 |

## 工程历史事件数据源对接

续报判定只需要读取「今日/昨日历史事件」。工程最终提供数据库直连、HTTP/RPC 或 SDK 均可，本项目内部统一通过 `EventRepository` 适配，不把数据源细节写进业务链路。

生产接入时：

1. 设置 `DUTY_USE_MOCK_EVENTS=0` 关闭 Mock 历史事件。
2. 在服务启动阶段调用 `configure_event_repository(repository)` 注入工程数据源适配器。
3. 若未注入或查询异常，历史事件返回空列表，续报判定退化为新事件，避免误合并。

适配器最小接口：

```python
from datetime import date
from typing import Any

from duty_processor.services.event_db import configure_event_repository


class EngineeringEventRepository:
    async def list_events_on_date(self, target_date: date) -> list[dict[str, Any]]:
        """返回 event_time 落在 target_date（北京日期）的历史事件。"""
        ...


configure_event_repository(EngineeringEventRepository())
```

返回字段契约：

```json
{
  "event_id": "事件唯一ID",
  "event_time": "接报/录音时间，ISO 8601",
  "report_type": "接报大类",
  "report_subtype": "接报子类",
  "caller_phone": "主叫号码",
  "receiver_phone": "被叫号码",
  "event_name": "事件名称",
  "summary": "事件摘要",
  "normalized_text": "归一化文本"
}
```

本项目只消费历史事件用于续报判定；当前处理结果的持久化仍由工程系统负责。若判定为续报，响应中的 `daily_duty_event.continuation_info.matched_event` 即被关联的历史 `event_id`。

## 续报判定后续维护

续报这块按“数据来源 / 规则召回 / 语义裁决 / 结果关联”分层，后面比较好改。原则是：业务流程只调用 `judge_continuation`，不要把数据库、提示词、打分权重散落到 API 层。

### 当前链路

1. 预处理阶段按 `recording_time` 查询今日和昨日历史事件：`query_today_events()` / `query_yesterday_events()`。
2. 名称摘要生成后，`judge_continuation` 组装当前事件的 `event_name`、`summary`、`report_type`、`report_subtype`、`caller_phone`。
3. 规则层先按四维内容分 + 电话加分召回候选，默认取 top 5，丢弃明显无关候选。
4. 默认开启 LLM “同一事件”裁决，在候选里判断是否为同一起真实事件/事态的续报；`DUTY_LLM_CONTINUATION=0` 时退回纯规则 top-1 且分数不低于 0.7。
5. 最终结果写入 `daily_duty_event.continuation_info`，其中 `matched_event` 是被关联的历史 `event_id`。

### 常见改动点

| 要改什么 | 改哪里 | 说明 |
|----------|--------|------|
| 接工程历史库 | `src/duty_processor/services/event_db.py` | 实现并注入 `EventRepository.list_events_on_date(target_date)`；不要在 API 或 tool 里直接查库。 |
| 调整召回数量或候选门槛 | `src/duty_processor/tools/judge_continuation.py` | 当前是 `rank_candidates(..., k=5, floor=0.1)`；候选太少会漏召，太多会增加 LLM 误并和耗时。 |
| 调整规则权重/阈值 | `src/duty_processor/prompts/continuation.py` | 改 `CONTINUATION_WEIGHTS`、`PHONE_BONUS`、`CONTINUATION_THRESHOLD` 或各 `match_*` 函数。改后补 `tests/test_phase2.py` 用例，防止过召/漏召。 |
| 调整语义裁决口径 | `src/duty_processor/prompts/continuation_llm.py` | 适合改“同一事件”的判定边界，例如多部门问询、同患者多次处置、同类但不同患者不应合并。 |
| 调整 LLM 调用/降级 | `src/duty_processor/services/continuation_service.py` | 控制是否启用 LLM、调用哪个模型、异常时如何降级。异常默认“不命中”，避免误合并。 |
| 调整 demo 多通来电串联 | `src/duty_processor/api/app.py` | `DUTY_RECORD_EVENTS=1` 时 `_maybe_record_event()` 只把新事件写入运行期内存；`/_debug/reset_events` 可清空演示残留。 |
| 调整最终关联字段 | `src/duty_processor/tools/assemble_duty_result.py` | 当前通过 `continuation_info.matched_event` 表示父/主事件 ID；工程落库时用它建立续报关系。 |

### 修改护栏

- 历史事件缺失、工程库未接、工程库异常时，续报应退化为新事件，不能阻断主流程。
- 电话一致只加分，电话不同不扣分；很多真实续报会跨科室、跨号码上报。
- 同类型、同科室、同诉求不等于同一事件；涉及患者/事故/事项不一致时要保守判新事件。
- 如果新增历史字段参与匹配，要同时更新工程数据源字段契约、demo 运行期写回字段和续报测试。
- 改权重、阈值、prompt 后至少跑 `tests/test_phase2.py`，全量回归用 `.venv/bin/python -m pytest -q`。

## 技术栈

- **Agent Loop**: Nanobot (Skill-Tool-Hook 架构)
- **LLM**: DashScope/百炼 Qwen3.7 (plus/max)，思考模式已关闭
- **ASR**: 百炼 qwen3-asr-flash / 火山引擎 / Mock
- **Web**: FastAPI + Uvicorn
- **Python**: 3.11
