From 2435fef5c02a937bbd82a0d320870512df1efa55 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 7 Oct 2026 20:58:49 +0800 Subject: [PATCH] feat(ai): play configured closing remarks before keyword hangup --- AGENTS.md | 2 +- contracts/local/config-read.schema.json | 11 +- .../local/examples/config-read-task-full.json | 2 +- .../config-read-legacy-hangup-keywords.json | 15 ++ contracts/local/manifest.json | 4 +- docs/evidence/hangup-keyword-closing.md | 33 ++++ docs/thirds/saas-dispatcher.md | 2 +- internal/ai/approved_mock.go | 36 +++-- internal/ai/approved_mock_test.go | 16 +- internal/ai/binding.go | 23 +-- internal/ai/binding_test.go | 2 +- internal/ai/hangup_config_test.go | 39 +++++ internal/ai/keyword_hangup.go | 101 +++++++++---- internal/ai/keyword_hangup_test.go | 143 ++++++++++-------- internal/ai/pipeline.go | 46 ++++-- internal/ai/pipeline_test.go | 24 ++- internal/callflow/approved.go | 3 +- internal/callflow/approved_test.go | 29 +++- internal/callflow/flow.go | 33 +++- internal/callflow/keyword_closing_test.go | 110 ++++++++++++++ .../rpc/approved_full_ai_integration_test.go | 43 ++++-- 21 files changed, 554 insertions(+), 163 deletions(-) create mode 100644 contracts/local/examples/invalid/config-read-legacy-hangup-keywords.json create mode 100644 docs/evidence/hangup-keyword-closing.md create mode 100644 internal/ai/hangup_config_test.go create mode 100644 internal/callflow/keyword_closing_test.go diff --git a/AGENTS.md b/AGENTS.md index 3419d82..3978460 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -92,7 +92,7 @@ - 独立 Dispatcher 的 SQLite 是任务、额度、inbox/outbox 的权威数据;Agent 无业务数据库,录音、执行与上传恢复只写受控私有文件。额度包含未知占用,新 boot/租约到期不得自动清除未知执行;明确 Agent 在拨号前拒绝时须持久拒绝并释放,已接纳但证实 ARI originate 从未提交的失败由原 Agent 经正式终结回执释放,RPC 交付未知时仍占用。2026-10-04 修复版 `d2c4a99` 已安装并通过非呼出时段主机诊断,新隔离测试租户额度 2/revision 2 已启用并由 Dispatcher 只读拉取;原快照保留;该条旧执行仅按前述经使用者授权的人工核实记录解除 Dispatcher 占用,绝非新 boot 或超时自动清除未知执行。仅 Agent 用机器可读标记证实本次未发起的拒绝可自动释放;旧执行的同错误码不算证明。不实现双活数据库、自动跨机热备、多 D 共享额度或第二租户公平。本轮不借本机 D1/D2 隔离夹具宣称多 D 运行。不得建立旧表/旧消息/旧 HTTP 执行兼容通道。 - Dispatcher↔Agent 复用 Unary gRPC 和受控 Endpoint;Agent 预绑定 D UUID 与服务端证书指纹,激活/会话代际、peer mTLS/SAN/SNI 和已签发期限须核对,新 boot 不清未知占用。Agent 不自行向 SaaS 取任务/AI/OSS 授权;Dispatcher 只用已经核验的 Agent `GetLoadedSIP` revision 开执行准入。本机 Mock 的加载回报不证明 Asterisk 已实际加载;仅隔离 SIP-only Agent 在配置原子写入、PJSIP reload 和运行态 endpoint/AOR/UDP transport 一致后持久标记 revision,每次加载查询重新核验。只支持明确的 UDP、IP/none 鉴权、无需 REGISTER 的 IPv4/PCMA 线路;未知字段和其他传输/鉴权/注册方式拒绝,不热更静态 transport。每条 SIP trunk 只允许一个原值 `caller_id`,Agent 把该值同时配置为 Asterisk endpoint 的 `from_user` 和 `callerid`,并在 `GetLoadedSIP` 前核对实际运行值;旧任务 `caller_profile_id` 与线路 `caller_profiles` 必须直接拒绝,不保留兼容路径。SIP 配置的唯一编辑/审批面仍是 management。 - 非生产真实 Agent 的 Asterisk `res_hep`/`res_hep_pjsip` 仅镜像至本机 UDP:在 ARI Dial 前读取本次 SIP Call-ID 并绑定执行,按 Call-ID 和原始 INVITE transaction 只采集真实最终响应的状态码、状态行与完整 `raw`;丢失/无法关联以 `sip_capture_error` 显式报告,不从 ARI、号码或挂断原因推测。镜像缺失不阻止已证实终结的额度释放;未经确认的占用仍保留。配置及回滚见 [`deploys/cell/README.md`](deploys/cell/README.md),不得将完整 SIP 报文写入普通日志、提交或聊天;本地测试通过不等于测试机已部署 HEP 或真实接通。 -- AI 使用任务内不可变授权快照:仅经获批准百炼/火山 ASR、OpenAI 兼容 LLM、百炼 TTS(`qwen3-tts-flash`/`Cherry`/`Chinese`)能表达的参数进入每通话实例;ASR-only 不启动 LLM/TTS,完整 AI 不借旧语音测试的授权或参数。只有最终用户 ASR 文本的明确字面关键词可触发拒联/挂断;不由 SDK 默认值、环境、CLI、metadata 或宽松 Schema 改写业务参数,不因 SDK 重试产生第二次发起/收费或重播。日志只存脱敏版本/摘要/计数,不存密钥、prompt、完整对话或音频。 +- AI 使用任务内不可变授权快照:仅经获批准百炼/火山 ASR、OpenAI 兼容 LLM、百炼 TTS(`qwen3-tts-flash`/`Cherry`/`Chinese`)能表达的参数进入每通话实例;ASR-only 不启动 LLM/TTS,完整 AI 不借旧语音测试的授权或参数。只有最终用户 ASR 文本的明确字面关键词可触发拒联/挂断;`agent.conversation.hangup_keywords` 仅接受必填非空 `name`/`triggers`/`closingRemark` 的对象数组,多组命中按配置顺序取第一组,经获批 TTS 完整播放该组结束语后主动挂断,不调用 LLM、不开始新一轮对话;合成/播放失败须明确报错并尝试结束,不重播,旧字符串数组直接拒绝;不由 SDK 默认值、环境、CLI、metadata 或宽松 Schema 改写业务参数,不因 SDK 重试产生第二次发起/收费或重播。日志只存脱敏版本/摘要/计数,不存密钥、prompt、完整对话或音频。 - **私有配置位置(本机路径相对本仓库根目录,只读,绝不提交)**:`.local/provider-ai.env` 是 `0600` 的 `KEY=VALUE` 文件;字段名为 `BAILIAN_API_KEY`、`BAILIAN_BASE_URL`、`BAILIAN_WSS_BASE_URL`、`BAILIAN_TTS_VOICE`、`VOLC_ASR_APP_NAME`、`VOLC_ASR_APP_KEY`、`VOLCENGINE_ACCESS_KEY`、`VOLCENGINE_SECRET_KEY`、`VOLCENGINE_REGION`、`VOLCENGINE_DISABLE_SSL`。根目录 `aliyun-oss.env` 也是 `0600`,**不是 shell env 文件**;它以冒号分隔,字段名准确为 `bucket`、`Endpoint`、`Region`,以及 `RAM` 下的 `username`、`accessKeyId`、`accessKeySecret`(大小写须保持原样)。测试机 `rogee` 用户的现行 ARI 文件位于 `~/.config/go-sip-asterisk/{ari.conf,http.conf,ari-secret}`,不是旧 `.local/asterisk-*/ari.conf`;访问测试机前先核对已登记的 SSH 主机指纹,不展示 `ari-secret`。 - **下次安全读取步骤**:先确认工作目录是本仓库,用 `stat` 仅检查本机两份文件是否存在、所有者与权限 `0600`;不满足即停止。按各自格式在受限本机进程中解析所需字段到内存,不执行 `source`、不打印全文/字段值、不写临时明文副本,不把密钥、签名 URL、音频或完整对话带入聊天、日志、提交及长期证据。AI 的历史文件只可作为**获准凭据来源**,模型/voice/速度等仍由当前获批的 task/providers 快照固定,不能用环境变量覆盖。OSS 历史文件也不能直接传给 `DISPATCHER_OSS_CONFIG_FILE`:该运行配置要求私有 JSON、`dispatcher_id` 和 `oss` 字段,并以环境变量**名称引用**密钥;需按现行合同构造并核验授权后才能使用。普通构建和测试不读取这些私有文件;真实服务测试必须显式启用对应 opt-in 并受现行门禁约束。 - **授权边界**:本轮验收目标是三条已登记线路的真实接通及 LLM 正常应答;旧两个号码已分别在三条线路试拨,六通均为 SIP 480,零接通。新增号码的历史逐次授权不等于本次代码变更获准部署或拨号;本次全局审查本地业务硬编码并以 SaaS 配置快照决定任务、线路和额度,**不部署、不拨号**。上述 AI 与 OSS 私有配置仍仅用于另经明确授权的非生产测试,下次任务须重新确认范围和真实服务调用授权,不能沿用本轮或历史一次性授权。不得把历史配置直接当 SaaS 快照、任务授权或真实呼叫准入,不覆盖/清理旧 OSS 对象。 diff --git a/contracts/local/config-read.schema.json b/contracts/local/config-read.schema.json index 4b45699..23a4f81 100644 --- a/contracts/local/config-read.schema.json +++ b/contracts/local/config-read.schema.json @@ -181,9 +181,18 @@ "required": ["text", "allowed_variables"], "properties": {"text": {"type": "string", "minLength": 1}, "allowed_variables": {"type": "array", "uniqueItems": true, "items": {"type": "string", "minLength": 1}}, "max_bytes": {"type": "integer", "minimum": 1}} }, + "hangup_keyword": { + "type": "object", "additionalProperties": false, + "required": ["name", "triggers", "closingRemark"], + "properties": { + "name": {"type": "string", "minLength": 1, "pattern": "\\S"}, + "triggers": {"type": "array", "minItems": 1, "items": {"type": "string", "minLength": 1, "pattern": "\\S"}}, + "closingRemark": {"type": "string", "minLength": 1, "pattern": "\\S"} + } + }, "conversation": { "type": "object", "additionalProperties": false, - "properties": {"opening": {"type": "string"}, "hangup_keywords": {"type": "array", "items": {"type": "string", "minLength": 1}}, "allow_interrupt": {"type": "boolean"}, "silence_timeout_ms": {"type": "integer", "minimum": 1}, "max_duration_ms": {"type": "integer", "minimum": 1}, "max_turns": {"type": "integer", "minimum": 1}, "sentence_max_chars": {"type": "integer", "minimum": 1}, "max_pending_audio_chunks": {"type": "integer", "minimum": 1}} + "properties": {"opening": {"type": "string"}, "hangup_keywords": {"type": "array", "items": {"$ref": "#/$defs/hangup_keyword"}}, "allow_interrupt": {"type": "boolean"}, "silence_timeout_ms": {"type": "integer", "minimum": 1}, "max_duration_ms": {"type": "integer", "minimum": 1}, "max_turns": {"type": "integer", "minimum": 1}, "sentence_max_chars": {"type": "integer", "minimum": 1}, "max_pending_audio_chunks": {"type": "integer", "minimum": 1}} } } } diff --git a/contracts/local/examples/config-read-task-full.json b/contracts/local/examples/config-read-task-full.json index eebd23c..414609c 100644 --- a/contracts/local/examples/config-read-task-full.json +++ b/contracts/local/examples/config-read-task-full.json @@ -10,6 +10,6 @@ "llm":{"provider_ref":"llm-example","model":"example-chat","temperature":0,"max_tokens":256,"timeout_ms":5000}, "tts":{"provider_ref":"tts-example","model":"qwen3-tts-flash","voice":"Cherry","language_type":"Chinese","speed":1,"format":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1},"timeout_ms":5000}, "prompt":{"text":"Example only","allowed_variables":[],"max_bytes":32768}, - "conversation":{"opening":"Example greeting","hangup_keywords":["不用了"],"allow_interrupt":false,"silence_timeout_ms":3000,"max_duration_ms":120000,"max_turns":20,"sentence_max_chars":80,"max_pending_audio_chunks":32} + "conversation":{"opening":"Example greeting","hangup_keywords":[{"name":"结束通话","triggers":["不用了"],"closingRemark":"好的,祝您生活愉快。"}],"allow_interrupt":false,"silence_timeout_ms":3000,"max_duration_ms":120000,"max_turns":20,"sentence_max_chars":80,"max_pending_audio_chunks":32} } } diff --git a/contracts/local/examples/invalid/config-read-legacy-hangup-keywords.json b/contracts/local/examples/invalid/config-read-legacy-hangup-keywords.json new file mode 100644 index 0000000..eebd23c --- /dev/null +++ b/contracts/local/examples/invalid/config-read-legacy-hangup-keywords.json @@ -0,0 +1,15 @@ +{ + "resource":"task_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001, + "task_id":"task-full","task_revision":2,"status":"running","name":"Full AI example", + "max_concurrent_calls":2,"ring_timeout_ms":30000,"max_call_duration_ms":120000, + "route_policy_id":"route-mock","allowed_trunk_ids":["trunk-mock"], + "schedule":{"time_zone":"Asia/Shanghai","starts_at":"2026-09-21T00:00:00+08:00","ends_at":null,"weekly_windows":{"monday":[{"start":"09:00","end":"20:00"}],"tuesday":[],"wednesday":[],"thursday":[],"friday":[],"saturday":[],"sunday":[]},"excluded_dates":["2026-10-01"]}, + "agent":{ + "immutable":true,"mode":"full_ai", + "asr":{"provider_ref":"asr-example","language":"zh-CN","input":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1,"sample_width_bytes":2},"interim":true,"timeout_ms":5000}, + "llm":{"provider_ref":"llm-example","model":"example-chat","temperature":0,"max_tokens":256,"timeout_ms":5000}, + "tts":{"provider_ref":"tts-example","model":"qwen3-tts-flash","voice":"Cherry","language_type":"Chinese","speed":1,"format":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1},"timeout_ms":5000}, + "prompt":{"text":"Example only","allowed_variables":[],"max_bytes":32768}, + "conversation":{"opening":"Example greeting","hangup_keywords":["不用了"],"allow_interrupt":false,"silence_timeout_ms":3000,"max_duration_ms":120000,"max_turns":20,"sentence_max_chars":80,"max_pending_audio_chunks":32} + } +} diff --git a/contracts/local/manifest.json b/contracts/local/manifest.json index 6c74332..ef73fe0 100644 --- a/contracts/local/manifest.json +++ b/contracts/local/manifest.json @@ -4,8 +4,8 @@ "sources": { "docs/archive/sources/v0.5-proposal.md": "612fdaee50aff6aa7fbef16c2d469d99857646c6d2235617d0e67f6098cd7ada", "docs/archive/sources/plan-saas-dispatcher-v05-v0.1.md": "666f39e56ea9f4b55661efcac82edd6f9729848e2d60e5f24cdf5aa3ac97ee87", - "docs/thirds/saas-dispatcher.md": "478cf69df793e6bacb6d19dcb50c1a9dc56d4cd9cae59ea5744ce94ee9e9d08c" + "docs/thirds/saas-dispatcher.md": "7aa4f38510799c9687bb233c9a15c2e44ad3d905d99ce0058f5074e7417fde91" }, - "bundle_sha256": "d25884427faaa4e096cfb2eb769d47216fef25f8e0359fb352a0c4934c8ed07c", + "bundle_sha256": "5eccbf38f0b6332dba25d2810da674340b660fd56b93b75fb5d8e240dbc6ef36", "bundle_algorithm": "sha256 of sorted relative-path + space + sha256(file) + newline; only root-level JSON and examples/**/*.json, excluding manifest.json" } diff --git a/docs/evidence/hangup-keyword-closing.md b/docs/evidence/hangup-keyword-closing.md new file mode 100644 index 0000000..4bce173 --- /dev/null +++ b/docs/evidence/hangup-keyword-closing.md @@ -0,0 +1,33 @@ +# 挂断关键词结束语:本地验证 + +## 范围与规则 + +使用者确认将 `agent.conversation.hangup_keywords` 改为对象数组。唯一字段规范和正例仍在 `contracts/local/`;业务说明在 `docs/thirds/saas-dispatcher.md`,不另维护平行字段定义。 + +- 仅用户最终 ASR 文本的字面包含匹配有效;同时命中多组时按配置顺序取第一组。 +- 使用获批 TTS 合成对应结束语,完整发送音频后主动挂断,不调用 LLM 生成结束语。 +- 命中后禁止新一轮对话;重复识别、合成失败及未知挂断均不触发重播或重试。 +- 合成或播放失败保留关键词命中事实,明确返回错误并尝试结束通话;多个错误同时返回。未成功发送的音频不报告为完整播放。 +- 旧字符串数组、缺字段、空字段及未知字段直接拒绝;ASR-only 不合成结束语。 + +## 可执行证据 + +| 检查 | 当前测试 | +| --- | --- | +| 新结构接纳、旧结构/空字段/错误大小写拒绝 | `TestBindHangupGroups`、`TestBindRejectsInvalidHangupGroups`、`TestCurrentContractExamples` 及新增旧结构反例 | +| 用户最终文本、字面匹配、组顺序、快照复制、并发单次预留 | `TestKeywordHangupOnlyFinalUserLiteralAndOrderedGroups`、`TestKeywordHangupConcurrentReservations` | +| 播完结束语后才挂断,且不再开始下一轮 | `TestApprovedKeywordPlaysClosingBeforeHangup`、`TestKeywordHangupWaitsForPlaybackReturn` | +| 合成/播放/挂断失败、多个错误并存、禁止重播 | `TestKeywordClosingFailuresRemainVisibleAndHangupOnce`、`TestKeywordHangupFailureVisibleWithoutRetry`、`TestKeywordHangupUsesFinalUserTextOnly` | +| 官方 TTS SDK、批准配置、结束语正文与发送/挂断顺序 | `TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword`:真实 SDK 对接本机模拟 HTTP 服务;最终 ASR 文本为显式测试输入,不代表真实 ASR 识别成功 | +| 隔离 Mock 不捏造结束语音频 | `TestApprovedMockFullAIFreezesScriptAndReturnsKeywordClosing` | + +P08 表保留当时的历史结果;其中旧 `TestKeywordHangupOnlyFinalUserLiteralAndOnlyOnce` 和 `TestKeywordHangupConcurrentDuplicateFinalTriggersOnce` 的现行替代分别为上表的顺序组匹配与并发预留测试。本变更涉及 A01/A04/A05/A08/A11 和 K09;其他业务边界不变,由完整回归及隔离 MQ/HTTPS/双向 TLS 测试核验。 + +## 本轮结果 + +- TDD:新配置和播放顺序测试先失败,完成实现后通过。 +- `PATH=/tmp/sip-go-agent-tools/bin:$PATH make check`:通过;包括格式、现行合同和 Proto 来源/hash、`go vet ./...`、`go test -race ./...`、构建,以及实际启动临时 RabbitMQ 的隔离端到端测试。 +- `bash scripts/coverage.sh`:手写业务语句覆盖率 **70.1%**,仅排除生成目录,超过 65%。 +- `git diff --check`:通过。 + +仅修改并验证本地代码,没有部署、真实拨号、真实 AI/OSS 服务调用,也没有读取私有凭据或处置旧 SQLite、录音、恢复文件。音频发送顺序的本地验证不证明真实线路的远端播放或生产签收。 diff --git a/docs/thirds/saas-dispatcher.md b/docs/thirds/saas-dispatcher.md index 0048560..9f855dc 100644 --- a/docs/thirds/saas-dispatcher.md +++ b/docs/thirds/saas-dispatcher.md @@ -36,7 +36,7 @@ RabbitMQ 是 Topic,**SaaS 独占创建、绑定、退役 exchange/queue,D - 当前 TTS 唯一获批适配器是 `bailian_tts`:任务快照须明确提供 `qwen3-tts-flash`、`Cherry`、`Chinese`、速度 `1` 和单声道 16 kHz PCM16 目标格式;使用已核验的 provider 凭据及生成端点。每段仅发起一次生成请求,下载返回的短期音频引用后转换为电话可用的 PCM16;不可用、超时、缺少转换工具或参数不支持时显式失败,不回退旧火山 TTS、不隐式重试或记录签名音频 URL。历史测试凭据不是任务授权,本地转换 Mock 不构成真实百炼/通话验收。 - 每条 SaaS SIP 线路只配置一个原值 `caller_id`;任务只选获准线路,不再带 `caller_profile_id`,线路不再带 `caller_profiles`。Dispatcher 将线路主叫绑定到每次执行;Agent 在 Asterisk PJSIP 线路同时设置 `from_user` 与 `callerid`,核对实际加载的出站身份后才准入,不以 ARI 创建通道后的 `CALLERID(num)` 代替 SIP `From`。主叫不是 Digest 用户名,不从环境、CLI 或被叫前缀推断。被叫号码只来自归属当前 D、租户及已接纳任务的 `call.execute.payload.callee` 原值;不在 Dispatcher、SaaS Mock 或抓证脚本另设固定号码列表,也不在任务快照增加号码列表。Dispatcher 与非生产发布/抓证入口只接受 1–32 位 ASCII 数字的原始号码,不能把线路前缀当成该号码的本地替代值。已选 SIP trunk、任务与线路每周时段、任务排除日期、任务/租户/线路额度、任务与 AI 较小通话时限均在接纳及实际发呼叫指令前检查。线路字段未知则 fail-closed;选线后固定、不自动重拨/换线。隔离 Mock 中规则暂不满足时保留待执行指令、暂停该任务的调度,规则允许后重验;与人工 pause/stop 分离,不能自动解除人为停止。非生产真实路径不另设固定的 `09:00`–`20:00` 时间门禁或每线路每原始号码每日 3 次上限;任务/线路时段与额度以 SaaS 已校验快照为准。逐通抓证、Agent 活跃凭证及使用者逐次授权仍须满足;格式有效或本地 Mock 收件均不构成真实拨号授权。 - 接通事实为真时 `outcome=answered`(后续异常不抹掉接通);已发起但忙线、拒接、无人接听且确定结束为 `no_answer`;确认未接通并由 Agent/Asterisk 执行故障终结为 `failed`;未知状态保持未知占用,不能伪造结束、结果或自动重拨。真实 SIP 状态码原样数字写入 `reason_code`,无真实 SIP 码则 `null` 并以 `reason_message` 说明;禁止本地虚构数字错误码。`call.execute.result.payload` 的 `status_line`、`raw` 与 `sip_capture_error` 始终存在:仅将经同一 ARI 通道拨号前取得的 SIP Call-ID 与 HEP INVITE 事务严格关联的最终响应写入原样状态行、完整原样报文与状态码;`raw` 不拼装、不截断,不能从目标号码、时间、挂断原因或 ARI HTTP 状态猜测。无 SIP 响应时前两项为 `null`;若已发起 SIP 但镜像/关联/解码失败,第三项须写明确错误,已确认结束仍报告真实结果并释放额度,不以原文缺失伪装为通话未知。原始报文只进入受控结果通道,不写日志、仓库或长期测试证据。无应答且没有录音时 `transcript=[]`、`opt_out=false`、`recording={}`。 -- 只有**用户侧 ASR 最终识别文本**包含任一 `hangup_keywords` 字面字符串才挂断;中间识别、助手回复、开场白、TTS 均不能触发;重复结果不可反复终结。同一任务 revision 不同内容拒绝准入;provider 禁用/角色不符不可调用。Mock 参数验证不等于真实供应商验收。 +- `agent.conversation.hangup_keywords` 使用对象数组,字段及正例以现行 Schema 为准;旧字符串数组直接拒绝。只有**用户侧 ASR 最终识别文本**字面包含某组 `triggers` 才触发,同时命中多组按配置顺序取第一组;使用该组 `closingRemark` 合成语音,完整播放后主动挂断,不调用 LLM 生成结束语。匹配后不再开启新一轮对话,重复识别不得重播或重复挂断;中间识别、助手回复、开场白、TTS 均不能触发。合成或播放失败时明确报错并尝试结束通话,不重试、不把未播放音频报告为已播放;挂断结果不明仍保留未知事实。组名、触发词列表及结束语均必填且非空;ASR-only 不配置或合成结束语。同一任务 revision 不同内容拒绝准入;provider 禁用/角色不符不可调用。Mock 参数验证不等于真实供应商验收。 ## 录音与最终结果(K10、K15、K16) diff --git a/internal/ai/approved_mock.go b/internal/ai/approved_mock.go index 74ef184..07ab9b3 100644 --- a/internal/ai/approved_mock.go +++ b/internal/ai/approved_mock.go @@ -50,7 +50,7 @@ func NewApprovedMockPipeline(bound Binding, script ApprovedMockScript, hangup fu } switch bound.Mode { case string(ModeASROnly): - if bound.Opening != "" || bound.LLM != nil || bound.TTS != nil || len(script.OpeningPCM16) != 0 { + if bound.Opening != "" || bound.LLM != nil || bound.TTS != nil || len(bound.HangupKeywords) != 0 || len(script.OpeningPCM16) != 0 { return nil, errors.New("ASR-only Mock cannot synthesize opening or assistant audio") } case string(ModeFullAI): @@ -73,15 +73,23 @@ func NewApprovedMockPipeline(bound Binding, script ApprovedMockScript, hangup fu return nil, errors.New("ASR-only Mock contains assistant audio") } } else { - keywordTurn := false - for _, literal := range bound.HangupKeywords { - if strings.Contains(turn.Transcript, literal) { - keywordTurn = true + var matched *HangupKeyword + for _, group := range bound.HangupKeywords { + for _, trigger := range group.Triggers { + if strings.Contains(turn.Transcript, trigger) { + matched = &group + break + } + } + if matched != nil { break } } - if !keywordTurn && (turn.Reply == "" || !validMockPCM(turn.ReplyPCM16)) { - return nil, errors.New("full-AI Mock reply is missing or invalid") + if turn.Reply == "" || !validMockPCM(turn.ReplyPCM16) { + return nil, errors.New("full-AI Mock reply or closing audio is missing or invalid") + } + if matched != nil && turn.Reply != matched.ClosingRemark { + return nil, errors.New("Mock closing remark differs from approved hangup group") } if len(turn.ReplyPCM16)%2 != 0 { return nil, errors.New("Mock reply PCM16 has an odd byte count") @@ -92,6 +100,10 @@ func NewApprovedMockPipeline(bound Binding, script ApprovedMockScript, hangup fu return &ApprovedMockPipeline{mode: bound.Mode, openingPCM: bytes.Clone(script.OpeningPCM16), turns: frozen, keyword: keyword}, nil } +func (p *ApprovedMockPipeline) FinishKeyword(ctx context.Context) error { + return p.keyword.Finish(ctx) +} + func validMockPCM(pcm []byte) bool { return len(pcm) > 0 && len(pcm)%2 == 0 } func (p *ApprovedMockPipeline) Open(ctx context.Context) ([]byte, error) { @@ -125,15 +137,13 @@ func (p *ApprovedMockPipeline) RunTurn(ctx context.Context, pcm []byte) (TurnRes } turn := p.turns[p.next] p.next++ // A transport or hangup failure cannot replay this turn. - stopped, err := p.keyword.Handle(ctx, ASRSegment{Source: "user", Text: turn.Transcript, Final: true}) - if stopped { - p.ended = true - } + group, err := p.keyword.Match(ASRSegment{Source: "user", Text: turn.Transcript, Final: true}) if err != nil { return TurnResult{}, err } - if stopped { - return TurnResult{Transcript: turn.Transcript, EndedByKeyword: true}, nil + if group != nil { + p.ended = true + return TurnResult{Transcript: turn.Transcript, Reply: group.ClosingRemark, AudioPCM16: bytes.Clone(turn.ReplyPCM16), EndedByKeyword: true}, nil } if p.mode == string(ModeASROnly) { return TurnResult{Transcript: turn.Transcript}, nil diff --git a/internal/ai/approved_mock_test.go b/internal/ai/approved_mock_test.go index 73934d0..cdd71f3 100644 --- a/internal/ai/approved_mock_test.go +++ b/internal/ai/approved_mock_test.go @@ -36,15 +36,15 @@ func TestApprovedMockAbsentOptionalTurnLimitUsesSharedOneTurnBoundary(t *testing } } -func TestApprovedMockFullAIFreezesScriptAndStopsBeforeKeywordReply(t *testing.T) { +func TestApprovedMockFullAIFreezesScriptAndReturnsKeywordClosing(t *testing.T) { opening := bytes.Repeat([]byte{1, 0}, 320) reply := bytes.Repeat([]byte{2, 0}, 320) script := ApprovedMockScript{OpeningPCM16: opening, Turns: []ApprovedMockTurn{ {Transcript: "正常模拟用户语音", Reply: "模拟助手答复", ReplyPCM16: reply}, - {Transcript: "请不要联系", Reply: "不得播报", ReplyPCM16: reply}, + {Transcript: "请不要联系", Reply: "好的,再见。", ReplyPCM16: reply}, }} var hangups int - bound := Binding{Mode: "full_ai", Opening: "批准的开场", HangupKeywords: []string{"不要联系"}, LLM: &LLMConfig{}, TTS: &TTSConfig{}, Conversation: ConversationConfig{MaxTurns: 2}} + bound := Binding{Mode: "full_ai", Opening: "批准的开场", HangupKeywords: []HangupKeyword{{Name: "结束", Triggers: []string{"不要联系"}, ClosingRemark: "好的,再见。"}}, LLM: &LLMConfig{}, TTS: &TTSConfig{}, Conversation: ConversationConfig{MaxTurns: 2}} pipeline, err := NewApprovedMockPipeline(bound, script, func(context.Context) error { hangups++; return nil }) if err != nil { t.Fatal(err) @@ -62,8 +62,14 @@ func TestApprovedMockFullAIFreezesScriptAndStopsBeforeKeywordReply(t *testing.T) t.Fatalf("full-AI mock changed approved fixture bytes: turn=%+v err=%v", first, err) } last, err := pipeline.RunTurn(context.Background(), bytes.Repeat([]byte{1, 0}, 1600)) - if err != nil || !last.EndedByKeyword || last.Transcript != "请不要联系" || last.Reply != "" || len(last.AudioPCM16) != 0 || hangups != 1 { - t.Fatalf("keyword did not stop before the assistant reply: turn=%+v hangups=%d err=%v", last, hangups, err) + if err != nil || !last.EndedByKeyword || last.Transcript != "请不要联系" || last.Reply != "好的,再见。" || len(last.AudioPCM16) != 640 || last.AudioPCM16[0] != 2 || hangups != 0 { + t.Fatalf("keyword closing was not frozen or hung up early: turn=%+v hangups=%d err=%v", last, hangups, err) + } + if err := pipeline.FinishKeyword(context.Background()); err != nil || hangups != 1 { + t.Fatalf("closing hangup failed: %v count=%d", err, hangups) + } + if _, err := pipeline.RunTurn(context.Background(), bytes.Repeat([]byte{1, 0}, 1600)); err == nil { + t.Fatal("mock continued after closing") } } diff --git a/internal/ai/binding.go b/internal/ai/binding.go index 6ac8dc6..e610481 100644 --- a/internal/ai/binding.go +++ b/internal/ai/binding.go @@ -25,7 +25,7 @@ type Binding struct { PromptMaxBytes int AllowedVariables []string Opening string - HangupKeywords []string + HangupKeywords []HangupKeyword Conversation ConversationConfig } @@ -105,14 +105,14 @@ type currentAgentSettings struct { MaxBytes *int `json:"max_bytes"` } `json:"prompt"` Conversation *struct { - Opening string `json:"opening"` - HangupKeywords []string `json:"hangup_keywords"` - AllowInterrupt *bool `json:"allow_interrupt"` - SilenceTimeoutMS *int64 `json:"silence_timeout_ms"` - MaxDurationMS *int64 `json:"max_duration_ms"` - MaxTurns int `json:"max_turns"` - SentenceMaxChars int `json:"sentence_max_chars"` - MaxPendingAudioChunks int `json:"max_pending_audio_chunks"` + Opening string `json:"opening"` + HangupKeywords []HangupKeyword `json:"hangup_keywords"` + AllowInterrupt *bool `json:"allow_interrupt"` + SilenceTimeoutMS *int64 `json:"silence_timeout_ms"` + MaxDurationMS *int64 `json:"max_duration_ms"` + MaxTurns int `json:"max_turns"` + SentenceMaxChars int `json:"sentence_max_chars"` + MaxPendingAudioChunks int `json:"max_pending_audio_chunks"` } `json:"conversation"` } @@ -249,7 +249,10 @@ func Bind(task configread.Task, providers map[string]configread.Provider) (Bindi MaxPendingAudioChunks: settings.Conversation.MaxPendingAudioChunks, } bound.Opening = settings.Conversation.Opening - bound.HangupKeywords = append([]string(nil), settings.Conversation.HangupKeywords...) + bound.HangupKeywords = cloneHangupKeywords(settings.Conversation.HangupKeywords) + if err := validateHangupKeywords(bound.HangupKeywords); err != nil { + return Binding{}, err + } return bound, nil } diff --git a/internal/ai/binding_test.go b/internal/ai/binding_test.go index fda0cc5..f21c1d1 100644 --- a/internal/ai/binding_test.go +++ b/internal/ai/binding_test.go @@ -79,7 +79,7 @@ func TestBindCurrentFullAIUsesApprovedSDKFields(t *testing.T) { if bound.TTS.Provider.Credential != providers["tts-example"].Credential || bound.TTS.Model != "qwen3-tts-flash" || bound.TTS.Voice != "Cherry" || bound.TTS.LanguageType != "Chinese" || bound.TTS.Speed != 1 || bound.TTS.SampleRate != 16000 || bound.TTS.Timeout != 5*time.Second { t.Fatal("Bailian TTS model/voice/neutral speed/PCM16 output/credential/timeout not bound") } - if len(bound.HangupKeywords) != 1 || bound.HangupKeywords[0] != "不用了" || bound.Prompt != "Example only" || bound.Opening != "Example greeting" { + if len(bound.HangupKeywords) != 1 || bound.HangupKeywords[0].Name != "结束通话" || len(bound.HangupKeywords[0].Triggers) != 1 || bound.HangupKeywords[0].Triggers[0] != "不用了" || bound.HangupKeywords[0].ClosingRemark != "好的,祝您生活愉快。" || bound.Prompt != "Example only" || bound.Opening != "Example greeting" { t.Fatal("immutable prompt and keyword behavior not bound") } } diff --git a/internal/ai/hangup_config_test.go b/internal/ai/hangup_config_test.go new file mode 100644 index 0000000..d855821 --- /dev/null +++ b/internal/ai/hangup_config_test.go @@ -0,0 +1,39 @@ +package ai + +import "testing" + +func TestBindHangupGroups(t *testing.T) { + task, providers := currentFixture(t, "full_ai") + task = changeCurrentAgent(t, task, func(agent map[string]any) { + agent["conversation"].(map[string]any)["hangup_keywords"] = []any{ + map[string]any{"name": "挂机1", "triggers": []string{"不需要", "再见"}, "closingRemark": "好的,祝您生活愉快。"}, + } + }) + if _, err := Bind(task, providers); err != nil { + t.Fatalf("new hangup groups must be accepted: %v", err) + } +} + +func TestBindRejectsInvalidHangupGroups(t *testing.T) { + for name, groups := range map[string]any{ + "legacy": []string{"再见"}, + "missing-name": []any{map[string]any{"triggers": []string{"再见"}, "closingRemark": "再见。"}}, + "blank-name": []any{map[string]any{"name": " ", "triggers": []string{"再见"}, "closingRemark": "再见。"}}, + "empty-triggers": []any{map[string]any{"name": "挂机1", "triggers": []string{}, "closingRemark": "再见。"}}, + "blank-trigger": []any{map[string]any{"name": "挂机1", "triggers": []string{" "}, "closingRemark": "再见。"}}, + "missing-remark": []any{map[string]any{"name": "挂机1", "triggers": []string{"再见"}}}, + "blank-remark": []any{map[string]any{"name": "挂机1", "triggers": []string{"再见"}, "closingRemark": " \n"}}, + "wrong-casing": []any{map[string]any{"name": "挂机1", "triggers": []string{"再见"}, "closing_remark": "再见。"}}, + "unknown-field": []any{map[string]any{"name": "挂机1", "triggers": []string{"再见"}, "closingRemark": "再见。", "fallback": true}}, + } { + t.Run(name, func(t *testing.T) { + task, providers := currentFixture(t, "full_ai") + task = changeCurrentAgent(t, task, func(agent map[string]any) { + agent["conversation"].(map[string]any)["hangup_keywords"] = groups + }) + if _, err := Bind(task, providers); err == nil { + t.Fatal("invalid or removed hangup configuration accepted") + } + }) + } +} diff --git a/internal/ai/keyword_hangup.go b/internal/ai/keyword_hangup.go index f9d14ee..5c5162a 100644 --- a/internal/ai/keyword_hangup.go +++ b/internal/ai/keyword_hangup.go @@ -17,61 +17,106 @@ type ASRSegment struct { Final bool } -// KeywordHangup is private to one call. Once a hangup is attempted its -// outcome may be unknown, so repeated ASR notifications cannot hang up again. +// HangupKeyword is one ordered group from the approved task snapshot. +type HangupKeyword struct { + Name string `json:"name"` + Triggers []string `json:"triggers"` + ClosingRemark string `json:"closingRemark"` +} + +func cloneHangupKeywords(groups []HangupKeyword) []HangupKeyword { + frozen := append([]HangupKeyword(nil), groups...) + for i := range frozen { + frozen[i].Triggers = append([]string(nil), frozen[i].Triggers...) + } + return frozen +} + +func validateHangupKeywords(groups []HangupKeyword) error { + valid := func(text string) bool { return utf8.ValidString(text) && strings.TrimSpace(text) != "" } + for i, group := range groups { + if !valid(group.Name) || !valid(group.ClosingRemark) || len(group.Triggers) == 0 { + return fmt.Errorf("hangup group %d requires name, triggers and closingRemark", i) + } + for _, trigger := range group.Triggers { + if !valid(trigger) { + return fmt.Errorf("hangup group %d requires nonempty UTF-8 triggers", i) + } + } + } + return nil +} + +// KeywordHangup reserves the first matching group before synthesis. Finish is +// called only after closing playback, or on explicit synthesis/playback failure. +// Neither the closing nor an uncertain hangup may be retried. type KeywordHangup struct { - keywords []string + groups []HangupKeyword hangup func(context.Context) error mu sync.Mutex requested bool + attempted bool } -func NewKeywordHangup(keywords []string, hangup func(context.Context) error) (*KeywordHangup, error) { +func NewKeywordHangup(groups []HangupKeyword, hangup func(context.Context) error) (*KeywordHangup, error) { if hangup == nil { return nil, errors.New("keyword termination requires a real hangup action") } - frozen := append([]string(nil), keywords...) - for _, keyword := range frozen { - if keyword == "" || !utf8.ValidString(keyword) { - return nil, errors.New("keyword termination requires nonempty UTF-8 literals") - } + if err := validateHangupKeywords(groups); err != nil { + return nil, err } - return &KeywordHangup{keywords: frozen, hangup: hangup}, nil + return &KeywordHangup{groups: cloneHangupKeywords(groups), hangup: hangup}, nil } -func (k *KeywordHangup) Handle(ctx context.Context, segment ASRSegment) (bool, error) { +func (k *KeywordHangup) Requested() bool { + k.mu.Lock() + defer k.mu.Unlock() + return k.requested +} + +func (k *KeywordHangup) Match(segment ASRSegment) (*HangupKeyword, error) { if k == nil { - return false, errors.New("keyword terminator is unavailable") + return nil, errors.New("keyword terminator is unavailable") } if segment.Source != "user" && segment.Source != "assistant" { - return false, errors.New("ASR text source is unknown") + return nil, errors.New("ASR text source is unknown") } if !segment.Final || segment.Source != "user" { - return false, nil + return nil, nil } if !utf8.ValidString(segment.Text) { - return false, errors.New("final user ASR text is not UTF-8") + return nil, errors.New("final user ASR text is not UTF-8") } k.mu.Lock() + defer k.mu.Unlock() if k.requested { - k.mu.Unlock() - return false, nil + return nil, errors.New("keyword closing already requested; no further turns permitted") } - matched := false - for _, keyword := range k.keywords { - if strings.Contains(segment.Text, keyword) { - matched = true - break + for _, group := range k.groups { + for _, trigger := range group.Triggers { + if strings.Contains(segment.Text, trigger) { + k.requested = true + group.Triggers = append([]string(nil), group.Triggers...) + return &group, nil + } } } - if !matched { - k.mu.Unlock() - return false, nil + return nil, nil +} + +func (k *KeywordHangup) Finish(ctx context.Context) error { + if k == nil { + return errors.New("keyword terminator is unavailable") } - k.requested = true + k.mu.Lock() + if !k.requested || k.attempted { + k.mu.Unlock() + return errors.New("keyword hangup not pending or already attempted") + } + k.attempted = true k.mu.Unlock() if err := k.hangup(ctx); err != nil { - return true, fmt.Errorf("keyword-triggered hangup outcome unknown: %w", err) + return fmt.Errorf("keyword-triggered hangup outcome unknown: %w", err) } - return true, nil + return nil } diff --git a/internal/ai/keyword_hangup_test.go b/internal/ai/keyword_hangup_test.go index 47899d4..1f76c63 100644 --- a/internal/ai/keyword_hangup_test.go +++ b/internal/ai/keyword_hangup_test.go @@ -3,97 +3,122 @@ package ai import ( "context" "errors" - "strings" "sync" - "sync/atomic" "testing" ) -func TestKeywordHangupOnlyFinalUserLiteralAndOnlyOnce(t *testing.T) { - var calls atomic.Int64 - stopper, err := NewKeywordHangup([]string{"停止通话", "STOP"}, func(context.Context) error { calls.Add(1); return nil }) +func keywordGroups() []HangupKeyword { + return []HangupKeyword{ + {Name: "挂机1", Triggers: []string{"不需要", "再见"}, ClosingRemark: "好的,祝您生活愉快。"}, + {Name: "挂机2", Triggers: []string{"再见", "不要联系"}, ClosingRemark: "好的,再见。"}, + } +} + +func TestKeywordHangupOnlyFinalUserLiteralAndOrderedGroups(t *testing.T) { + calls := 0 + groups := keywordGroups() + stopper, err := NewKeywordHangup(groups, func(context.Context) error { calls++; return nil }) if err != nil { t.Fatal(err) } - cases := []struct { - event ASRSegment - expected bool - }{ - {ASRSegment{Source: "user", Text: "请停止通话", Final: false}, false}, - {ASRSegment{Source: "assistant", Text: "请停止通话", Final: true}, false}, - {ASRSegment{Source: "user", Text: "请停通话", Final: true}, false}, - {ASRSegment{Source: "user", Text: "stop", Final: true}, false}, - {ASRSegment{Source: "user", Text: "现在停止通话", Final: true}, true}, - {ASRSegment{Source: "user", Text: "现在停止通话", Final: true}, false}, - {ASRSegment{Source: "user", Text: "STOP", Final: true}, false}, - } - for i, tc := range cases { - triggered, err := stopper.Handle(context.Background(), tc.event) - if err != nil || triggered != tc.expected { - t.Fatalf("event %d: triggered=%v want=%v err=%v", i, triggered, tc.expected, err) + groups[0].Triggers[1], groups[0].ClosingRemark = "改过的词", "改过的结束语" + for _, segment := range []ASRSegment{ + {Source: "user", Text: "再见", Final: false}, + {Source: "assistant", Text: "再见", Final: true}, + {Source: "user", Text: "继续", Final: true}, + {Source: "user", Text: "不 需要", Final: true}, + } { + if group, err := stopper.Match(segment); err != nil || group != nil || calls != 0 { + t.Fatalf("nonfinal/nonliteral/assistant text triggered closing: %v", err) } } - if calls.Load() != 1 { - t.Fatalf("duplicate or interim hangup: calls=%d", calls.Load()) + group, err := stopper.Match(ASRSegment{Source: "user", Text: "再见,不要联系", Final: true}) + if err != nil || group == nil || group.Name != "挂机1" || group.ClosingRemark != "好的,祝您生活愉快。" || calls != 0 { + t.Fatalf("first configured match must reserve frozen closing without hanging up: group=%+v calls=%d err=%v", group, calls, err) + } + if _, err := stopper.Match(ASRSegment{Source: "user", Text: "再见", Final: true}); err == nil { + t.Fatal("duplicate result allowed closing replay") + } + if err := stopper.Finish(context.Background()); err != nil || calls != 1 { + t.Fatalf("finish: %v calls=%d", err, calls) + } + if err := stopper.Finish(context.Background()); err == nil || calls != 1 { + t.Fatal("hangup retried") } } -func TestKeywordHangupFailureIsVisibleButUnknownActionNeverRepeated(t *testing.T) { - var calls atomic.Int64 - stopper, err := NewKeywordHangup([]string{"终止"}, func(context.Context) error { calls.Add(1); return errors.New("mock ARI hangup outcome unknown") }) +func TestKeywordHangupFailureVisibleWithoutRetry(t *testing.T) { + want := errors.New("transport outcome unknown") + calls := 0 + stopper, err := NewKeywordHangup(keywordGroups(), func(context.Context) error { calls++; return want }) if err != nil { t.Fatal(err) } - result := ASRSegment{Source: "user", Text: "请终止", Final: true} - if triggered, err := stopper.Handle(context.Background(), result); !triggered || err == nil || !strings.Contains(err.Error(), "outcome unknown") { - t.Fatalf("hangup failure hidden: %v %v", triggered, err) + if _, err := stopper.Match(ASRSegment{Source: "user", Text: "不需要", Final: true}); err != nil { + t.Fatal(err) } - if triggered, err := stopper.Handle(context.Background(), result); triggered || err != nil || calls.Load() != 1 { - t.Fatalf("unknown hangup retried automatically: %v %v calls=%d", triggered, err, calls.Load()) + if err := stopper.Finish(context.Background()); !errors.Is(err, want) || calls != 1 { + t.Fatalf("failure hidden: %v", err) + } + if err := stopper.Finish(context.Background()); err == nil || calls != 1 { + t.Fatal("uncertain hangup retried") } } -func TestKeywordHangupConcurrentDuplicateFinalTriggersOnce(t *testing.T) { - var calls atomic.Int64 - stopper, err := NewKeywordHangup([]string{"停机"}, func(context.Context) error { calls.Add(1); return nil }) +func TestKeywordHangupConcurrentReservations(t *testing.T) { + stopper, err := NewKeywordHangup(keywordGroups(), func(context.Context) error { return nil }) if err != nil { t.Fatal(err) } var wg sync.WaitGroup - for i := 0; i < 32; i++ { - wg.Add(1) - go func() { - defer wg.Done() - _, _ = stopper.Handle(context.Background(), ASRSegment{Source: "user", Text: "停机", Final: true}) - }() + matched := make(chan bool, 20) + for range 20 { + wg.Go(func() { + group, _ := stopper.Match(ASRSegment{Source: "user", Text: "不需要", Final: true}) + matched <- group != nil + }) } wg.Wait() - if calls.Load() != 1 { - t.Fatalf("concurrent duplicate final ASR hung up %d times", calls.Load()) + close(matched) + count := 0 + for ok := range matched { + if ok { + count++ + } + } + if count != 1 { + t.Fatalf("closing reserved %d times", count) } } -func TestKeywordHangupRejectsInvalidConfigurationAndUnknownSource(t *testing.T) { - callback := func(context.Context) error { return nil } - if _, err := NewKeywordHangup([]string{""}, callback); err == nil { - t.Fatal("accepted empty keyword matching every transcript") +func TestKeywordHangupRejectsInvalidInput(t *testing.T) { + for _, groups := range [][]HangupKeyword{ + {{Name: "", Triggers: []string{"词"}, ClosingRemark: "结束"}}, + {{Name: "组", Triggers: []string{""}, ClosingRemark: "结束"}}, + {{Name: "组", Triggers: []string{string([]byte{0xff})}, ClosingRemark: "结束"}}, + {{Name: "组", Triggers: []string{"词"}, ClosingRemark: " "}}, + } { + if _, err := NewKeywordHangup(groups, func(context.Context) error { return nil }); err == nil { + t.Fatal("invalid group accepted") + } } - if _, err := NewKeywordHangup([]string{string([]byte{0xff})}, callback); err == nil { - t.Fatal("accepted invalid UTF-8 keyword") + if _, err := NewKeywordHangup(nil, nil); err == nil { + t.Fatal("missing hangup accepted") } - if _, err := NewKeywordHangup([]string{"终止"}, nil); err == nil { - t.Fatal("accepted missing actual hangup action") + stopper, _ := NewKeywordHangup(keywordGroups(), func(context.Context) error { return nil }) + for _, segment := range []ASRSegment{{Source: "unknown", Final: true}, {Source: "user", Text: string([]byte{0xff}), Final: true}} { + if _, err := stopper.Match(segment); err == nil { + t.Fatal("invalid ASR input accepted") + } } - words := []string{"终止"} - stopper, err := NewKeywordHangup(words, callback) - if err != nil { - t.Fatal(err) + if err := stopper.Finish(context.Background()); err == nil { + t.Fatal("hangup without keyword reservation") } - words[0] = "改动" - if triggered, err := stopper.Handle(context.Background(), ASRSegment{Source: "user", Text: "终止", Final: true}); err != nil || !triggered { - t.Fatalf("keyword list was mutable after approval: %v %v", triggered, err) + var unavailable *KeywordHangup + if _, err := unavailable.Match(ASRSegment{}); err == nil { + t.Fatal("nil matcher accepted") } - if _, err := stopper.Handle(context.Background(), ASRSegment{Source: "other", Text: "终止", Final: true}); err == nil { - t.Fatal("unknown text source was silently accepted") + if err := unavailable.Finish(context.Background()); err == nil { + t.Fatal("nil finish accepted") } } diff --git a/internal/ai/pipeline.go b/internal/ai/pipeline.go index 06e2019..35ef6c5 100644 --- a/internal/ai/pipeline.go +++ b/internal/ai/pipeline.go @@ -23,11 +23,15 @@ type Call struct { bound Binding keyword *KeywordHangup mu sync.Mutex + turnMu sync.Mutex openingRequested bool openingReady bool } func NewCall(bound Binding, hangup func(context.Context) error) (*Call, error) { + if bound.Mode == "asr_only" && len(bound.HangupKeywords) != 0 { + return nil, errors.New("ASR-only call cannot play keyword closing remarks") + } keyword, err := NewKeywordHangup(bound.HangupKeywords, hangup) if err != nil { return nil, err @@ -65,19 +69,42 @@ func (c *Call) Open(ctx context.Context) ([]byte, error) { return audio, nil } -func (c *Call) HandleFinalASR(ctx context.Context, segment ASRSegment) (bool, error) { +func (c *Call) HandleFinalASR(ctx context.Context, segment ASRSegment) (TurnResult, error) { if c == nil { - return false, errors.New("AI call is unavailable") + return TurnResult{}, errors.New("AI call is unavailable") } - return c.keyword.Handle(ctx, segment) + group, err := c.keyword.Match(segment) + if err != nil || group == nil { + return TurnResult{}, err + } + turn := TurnResult{Transcript: segment.Text, Reply: group.ClosingRemark, EndedByKeyword: true} + slog.Info("approved keyword closing requested", "stage", "closing") + turn.AudioPCM16, err = c.bound.SynthesizeReply(ctx, group.ClosingRemark) + if err != nil { + slog.Error("approved keyword closing synthesis failed", "stage", "closing") + return turn, fmt.Errorf("keyword closing TTS failed: %w", err) + } + return turn, nil } -// RunTurn performs one utterance. No LLM or TTS request is sent after a -// keyword hangup, and only a final user recognition may enter the LLM. +func (c *Call) FinishKeyword(ctx context.Context) error { + if c == nil { + return errors.New("AI call is unavailable") + } + return c.keyword.Finish(ctx) +} + +// RunTurn performs one utterance. A keyword match bypasses the LLM and +// synthesizes only the approved closing; subsequent turns are rejected. func (c *Call) RunTurn(ctx context.Context, pcm16 []byte) (TurnResult, error) { if c == nil { return TurnResult{}, errors.New("AI call is unavailable") } + c.turnMu.Lock() + defer c.turnMu.Unlock() + if c.keyword.Requested() { + return TurnResult{}, errors.New("keyword closing already requested; no further turns permitted") + } if c.bound.Mode == "full_ai" && c.bound.Opening != "" { c.mu.Lock() ready := c.openingReady @@ -90,12 +117,9 @@ func (c *Call) RunTurn(ctx context.Context, pcm16 []byte) (TurnResult, error) { if err != nil { return TurnResult{}, fmt.Errorf("ASR failed: %w", err) } - stopped, err := c.HandleFinalASR(ctx, ASRSegment{Source: "user", Text: text, Final: true}) - if err != nil { - return TurnResult{}, err - } - if stopped { - return TurnResult{Transcript: text, EndedByKeyword: true}, nil + closing, err := c.HandleFinalASR(ctx, ASRSegment{Source: "user", Text: text, Final: true}) + if err != nil || closing.EndedByKeyword { + return closing, err } if c.bound.Mode == "asr_only" { return TurnResult{Transcript: text}, nil diff --git a/internal/ai/pipeline_test.go b/internal/ai/pipeline_test.go index 346127a..a599490 100644 --- a/internal/ai/pipeline_test.go +++ b/internal/ai/pipeline_test.go @@ -153,15 +153,25 @@ func TestKeywordHangupUsesFinalUserTextOnly(t *testing.T) { {Source: "user", Text: "继续说", Final: true}, } { stop, err := call.HandleFinalASR(context.Background(), segment) - if err != nil || stop || calls != 0 { - t.Fatalf("interim/assistant/unmatched text cannot hang up: stopped=%t err=%v calls=%d", stop, err, calls) + if err != nil || stop.EndedByKeyword || calls != 0 { + t.Fatalf("interim/assistant/unmatched text cannot hang up: stopped=%t err=%v calls=%d", stop.EndedByKeyword, err, calls) } } - for i := 0; i < 2; i++ { - stop, err := call.HandleFinalASR(context.Background(), ASRSegment{Source: "user", Text: "我不用了", Final: true}) - if err != nil || stop != (i == 0) || calls != 1 { - t.Fatalf("matched final user text hangs up at most once: stopped=%t err=%v calls=%d", stop, err, calls) - } + // A missing TTS is an explicit closing failure, not permission to call LLM + // or silently drop the keyword fact. No provider request is made here. + call.bound.TTS = nil + stop, err := call.HandleFinalASR(context.Background(), ASRSegment{Source: "user", Text: "我不用了", Final: true}) + if err == nil || !stop.EndedByKeyword || stop.Transcript != "我不用了" || stop.Reply != "好的,祝您生活愉快。" || calls != 0 { + t.Fatalf("closing failure lost keyword fact or hung up before playback: turn=%+v calls=%d err=%v", stop, calls, err) + } + if err := call.FinishKeyword(context.Background()); err != nil || calls != 1 { + t.Fatalf("closing failure cleanup: calls=%d err=%v", calls, err) + } + if _, err := call.HandleFinalASR(context.Background(), ASRSegment{Source: "user", Text: "我不用了", Final: true}); err == nil { + t.Fatal("closing failure was replayed") + } + if _, err := call.RunTurn(context.Background(), nil); err == nil { + t.Fatal("new ASR turn started after keyword closing") } } diff --git a/internal/callflow/approved.go b/internal/callflow/approved.go index dee5127..52f1c14 100644 --- a/internal/callflow/approved.go +++ b/internal/callflow/approved.go @@ -12,6 +12,7 @@ import ( type ApprovedPipeline interface { Open(context.Context) ([]byte, error) RunTurn(context.Context, []byte) (ai.TurnResult, error) + FinishKeyword(context.Context) error } func ExecuteApproved(ctx context.Context, session MediaSession, mode ai.Mode, pipeline ApprovedPipeline, capture CaptureConfig) (Result, error) { @@ -23,7 +24,7 @@ func ExecuteApproved(ctx context.Context, session MediaSession, mode ai.Mode, pi ctx, cancel = context.WithTimeout(ctx, capture.CallDuration) defer cancel() } - return executeFlow(ctx, session, mode, pipeline.Open, pipeline.RunTurn, capture) + return executeFlow(ctx, session, mode, pipeline.Open, pipeline.RunTurn, pipeline.FinishKeyword, capture) } // ApprovedCapture hands explicit task conversation limits to the shared media diff --git a/internal/callflow/approved_test.go b/internal/callflow/approved_test.go index 8a6486f..dff0f00 100644 --- a/internal/callflow/approved_test.go +++ b/internal/callflow/approved_test.go @@ -15,6 +15,15 @@ type approvedFlowPipeline struct { openingCalls int turnCalls int turn ai.TurnResult + turnErr error + finish func(context.Context) error +} + +func (p *approvedFlowPipeline) FinishKeyword(ctx context.Context) error { + if p.finish != nil { + return p.finish(ctx) + } + return nil } func (p *approvedFlowPipeline) Open(context.Context) ([]byte, error) { @@ -24,7 +33,7 @@ func (p *approvedFlowPipeline) Open(context.Context) ([]byte, error) { func (p *approvedFlowPipeline) RunTurn(context.Context, []byte) (ai.TurnResult, error) { p.turnCalls++ - return p.turn, nil + return p.turn, p.turnErr } type scriptedTurnSession struct { @@ -112,17 +121,24 @@ func TestApprovedInvalidCallStopsBeforeReply(t *testing.T) { } } -func TestApprovedKeywordStopsBeforeSendingReply(t *testing.T) { - pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "我不用了", EndedByKeyword: true, AudioPCM16: make([]byte, 6400)}} +func TestApprovedKeywordPlaysClosingBeforeHangup(t *testing.T) { session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}} + hangups := 0 + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "我不用了", Reply: "好的,祝您生活愉快。", EndedByKeyword: true, AudioPCM16: make([]byte, 6400)}, finish: func(context.Context) error { + hangups++ + if session.stats.SentPackets != 2 { + t.Fatal("hangup occurred before closing playback completed") + } + return nil + }} result, err := ExecuteApproved(context.Background(), session, ai.ModeFullAI, pipeline, CaptureConfig{ FirstSpeechTimeout: time.Second, MaxDuration: 2 * time.Millisecond, MaxTurns: 3, }) if err != nil { t.Fatal(err) } - if !result.Turn.EndedByKeyword || len(result.Turns) != 1 || len(result.OutboundTurns) != 1 || session.stats.SentPackets != 1 || pipeline.openingCalls != 1 || pipeline.turnCalls != 1 { - t.Fatalf("keyword hangup must stop before LLM/TTS reply: turns=%d outbound=%d sent=%d", len(result.Turns), len(result.OutboundTurns), session.stats.SentPackets) + if !result.Turn.EndedByKeyword || len(result.Turns) != 1 || len(result.OutboundTurns) != 2 || session.stats.SentPackets != 2 || pipeline.openingCalls != 1 || pipeline.turnCalls != 1 || hangups != 1 { + t.Fatalf("closing playback/hangup incomplete: turns=%d outbound=%d sent=%d hangups=%d", len(result.Turns), len(result.OutboundTurns), session.stats.SentPackets, hangups) } } @@ -170,10 +186,11 @@ func TestApprovedCapturePreservesConversationControls(t *testing.T) { type approvedDeadlineProbe struct{ deadlinePresent bool } func (*approvedDeadlineProbe) Open(context.Context) ([]byte, error) { return nil, nil } +func (*approvedDeadlineProbe) FinishKeyword(context.Context) error { return nil } func (p *approvedDeadlineProbe) RunTurn(ctx context.Context, _ []byte) (ai.TurnResult, error) { deadline, ok := ctx.Deadline() p.deadlinePresent = ok && time.Until(deadline) > 0 && time.Until(deadline) <= time.Second - return ai.TurnResult{Transcript: "最终识别", EndedByKeyword: true}, nil + return ai.TurnResult{Transcript: "最终识别"}, nil } func TestApprovedConversationDurationBoundsWholeCall(t *testing.T) { diff --git a/internal/callflow/flow.go b/internal/callflow/flow.go index a937664..00d5ad1 100644 --- a/internal/callflow/flow.go +++ b/internal/callflow/flow.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "log/slog" "time" "git.ipao.vip/rogee/go-sip/internal/ai" @@ -34,8 +35,8 @@ type Result struct { RTP media.RTPStats } -func executeFlow(ctx context.Context, session MediaSession, mode ai.Mode, open func(context.Context) ([]byte, error), runTurn func(context.Context, []byte) (ai.TurnResult, error), capture CaptureConfig) (Result, error) { - if session == nil || open == nil || runTurn == nil { +func executeFlow(ctx context.Context, session MediaSession, mode ai.Mode, open func(context.Context) ([]byte, error), runTurn func(context.Context, []byte) (ai.TurnResult, error), finishKeyword func(context.Context) error, capture CaptureConfig) (Result, error) { + if session == nil || open == nil || runTurn == nil || finishKeyword == nil { return Result{}, errors.New("media session and AI pipeline are required") } if mode != ai.ModeFullAI && mode != ai.ModeASROnly { @@ -71,12 +72,38 @@ func executeFlow(ctx context.Context, session MediaSession, mode ai.Mode, open f return result, fmt.Errorf("captured audio turn %d is too short", turnIndex+1) } turn, err := runTurn(ctx, inbound) + if turn.EndedByKeyword { + result.Turn = turn + result.Turns = append(result.Turns, turn) + closingErr := err + if closingErr == nil { + if mode != ai.ModeFullAI || len(turn.AudioPCM16) == 0 { + closingErr = errors.New("keyword closing has no playable full-AI audio") + } else if sendErr := session.SendPCM16(ctx, turn.AudioPCM16, 16000); sendErr != nil { + closingErr = fmt.Errorf("send keyword closing: %w", sendErr) + } else { + result.OutboundTurns = append(result.OutboundTurns, clonePCM(turn.AudioPCM16)) + slog.Info("approved keyword closing playback completed", "audio_bytes", len(turn.AudioPCM16)) + } + } + if closingErr != nil { + slog.Error("approved keyword closing failed", "stage", "closing", "turn", turnIndex+1) + } + hangupErr := finishKeyword(ctx) + if hangupErr != nil { + slog.Error("approved keyword hangup failed", "stage", "hangup", "turn", turnIndex+1) + } else { + slog.Info("approved keyword hangup completed", "turn", turnIndex+1) + } + result.RTP = session.Stats() + return result, errors.Join(closingErr, hangupErr) + } if err != nil { return result, fmt.Errorf("run AI turn %d: %w", turnIndex+1, err) } result.Turn = turn result.Turns = append(result.Turns, turn) - if turn.InvalidCall || turn.EndedByKeyword { + if turn.InvalidCall { result.RTP = session.Stats() return result, nil } diff --git a/internal/callflow/keyword_closing_test.go b/internal/callflow/keyword_closing_test.go new file mode 100644 index 0000000..ff200d4 --- /dev/null +++ b/internal/callflow/keyword_closing_test.go @@ -0,0 +1,110 @@ +package callflow + +import ( + "context" + "errors" + "testing" + "time" + + "git.ipao.vip/rogee/go-sip/internal/ai" +) + +type closingSendSession struct { + *scriptedTurnSession + sends int + closing func(context.Context) error +} + +func (s *closingSendSession) SendPCM16(ctx context.Context, pcm []byte, rate int) error { + s.sends++ + if s.sends == 2 && s.closing != nil { + if err := s.closing(ctx); err != nil { + return err + } + } + return s.scriptedTurnSession.SendPCM16(ctx, pcm, rate) +} + +func TestKeywordClosingFailuresRemainVisibleAndHangupOnce(t *testing.T) { + synthesisErr := errors.New("injected closing synthesis failure") + playbackErr := errors.New("injected closing playback failure") + hangupErr := errors.New("injected unknown hangup outcome") + for _, tc := range []struct { + name string + synthErr, sendErr, finishErr error + empty bool + }{ + {name: "synthesis", synthErr: synthesisErr}, + {name: "playback", sendErr: playbackErr}, + {name: "hangup", finishErr: hangupErr}, + {name: "both-playback-and-hangup", sendErr: playbackErr, finishErr: hangupErr}, + {name: "empty-closing", empty: true}, + } { + t.Run(tc.name, func(t *testing.T) { + session := &closingSendSession{scriptedTurnSession: &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}}, closing: func(context.Context) error { return tc.sendErr }} + hangups := 0 + turn := ai.TurnResult{Transcript: "不需要", Reply: "好的,再见。", EndedByKeyword: true, AudioPCM16: make([]byte, 6400)} + if tc.empty { + turn.AudioPCM16 = nil + } + pipeline := &approvedFlowPipeline{turn: turn, turnErr: tc.synthErr, finish: func(context.Context) error { hangups++; return tc.finishErr }} + result, err := ExecuteApproved(context.Background(), session, ai.ModeFullAI, pipeline, CaptureConfig{FirstSpeechTimeout: time.Second, MaxDuration: 2 * time.Millisecond, MaxTurns: 3}) + if err == nil || hangups != 1 || pipeline.turnCalls != 1 || len(result.Turns) != 1 || !result.Turn.EndedByKeyword || result.Turn.Transcript != "不需要" { + t.Fatalf("closing error/facts lost or new turn started: err=%v hangups=%d turns=%d result=%+v", err, hangups, pipeline.turnCalls, result.Turn) + } + for _, want := range []error{tc.synthErr, tc.sendErr, tc.finishErr} { + if want != nil && !errors.Is(err, want) { + t.Fatalf("cause hidden: want=%v got=%v", want, err) + } + } + if (tc.synthErr != nil || tc.sendErr != nil || tc.empty) && len(result.OutboundTurns) != 1 { + t.Fatal("unsent closing reported as played") + } + if tc.synthErr != nil && session.sends != 1 { + t.Fatal("partial failed synthesis was played") + } + if session.sends > 2 { + t.Fatal("closing replayed") + } + }) + } +} + +func TestKeywordHangupWaitsForPlaybackReturn(t *testing.T) { + started, release, hungup := make(chan struct{}), make(chan struct{}), make(chan struct{}) + session := &closingSendSession{scriptedTurnSession: &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}}, closing: func(context.Context) error { + close(started) + <-release + return nil + }} + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "再见", Reply: "再见。", EndedByKeyword: true, AudioPCM16: make([]byte, 6400)}, finish: func(context.Context) error { close(hungup); return nil }} + done := make(chan error, 1) + go func() { + _, err := ExecuteApproved(context.Background(), session, ai.ModeFullAI, pipeline, CaptureConfig{FirstSpeechTimeout: time.Second, MaxDuration: 2 * time.Millisecond, MaxTurns: 3}) + done <- err + }() + select { + case <-started: + case <-time.After(5 * time.Second): + close(release) + t.Fatal("closing playback never started") + } + select { + case <-hungup: + close(release) + t.Fatal("hung up while closing playback was still in progress") + default: + } + close(release) + if err := <-done; err != nil { + t.Fatal(err) + } + select { + case <-hungup: + default: + t.Fatal("playback completed without hangup") + } + if session.sends != 2 || pipeline.turnCalls != 1 { + t.Fatal("new turn/replay after closing") + } +} diff --git a/internal/rpc/approved_full_ai_integration_test.go b/internal/rpc/approved_full_ai_integration_test.go index 344eac9..54484ad 100644 --- a/internal/rpc/approved_full_ai_integration_test.go +++ b/internal/rpc/approved_full_ai_integration_test.go @@ -1,6 +1,7 @@ package rpc import ( + "bytes" "context" "encoding/json" "fmt" @@ -14,6 +15,7 @@ import ( agentpb "git.ipao.vip/rogee/go-sip/gen/agent" "git.ipao.vip/rogee/go-sip/internal/ai" + "git.ipao.vip/rogee/go-sip/internal/callflow" "git.ipao.vip/rogee/go-sip/internal/configread" "git.ipao.vip/rogee/go-sip/internal/dispatcher" "git.ipao.vip/rogee/go-sip/internal/media" @@ -24,6 +26,14 @@ type approvedSDKObservation struct { Authorization string } +// finalUserSDKPipeline substitutes only final ASR text; TTS and call sequencing +// use the actual approved implementations and local SDK/media endpoints. +type finalUserSDKPipeline struct{ *ai.Call } + +func (p finalUserSDKPipeline) RunTurn(ctx context.Context, _ []byte) (ai.TurnResult, error) { + return p.HandleFinalASR(ctx, ai.ASRSegment{Source: "user", Text: "我不用了", Final: true}) +} + func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T) { if _, err := exec.LookPath("ffmpeg"); err != nil { t.Skip("Bailian TTS conversion requires ffmpeg") @@ -42,7 +52,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T req.TaskId = task.TaskID req.SourceEventId, req.CallId = "event-full", "event-full" llmObserved := make(chan approvedSDKObservation, 1) - ttsObserved := make(chan approvedSDKObservation, 2) + ttsObserved := make(chan approvedSDKObservation, 3) wav, _, err := media.EncodeMonoWAV([]byte{1, 0, 2, 0}, 1024) if err != nil { t.Fatal(err) @@ -97,23 +107,23 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T if approved.AI.Mode != "full_ai" || approved.AI.LLM == nil || approved.AI.TTS == nil || approved.AI.Conversation.MaxTurns != 20 || approved.AI.LLM.Temperature == nil || *approved.AI.LLM.Temperature != 0 { t.Fatal("approved full-AI settings were not delivered to the Agent") } - call, err := ai.NewCall(approved.AI, func(context.Context) error { hangups++; return nil }) + session := callflow.NewMemorySession(bytes.Repeat([]byte{1, 0}, 1600)) + call, err := ai.NewCall(approved.AI, func(context.Context) error { + if session.Stats().SentBytes != 8 { + t.Fatal("keyword hung up before opening and closing media were sent") + } + hangups++ + return nil + }) if err != nil { return err } - openingAudio, err := call.Open(ctx) - if err != nil { - return fmt.Errorf("synthesize approved opening: %w", err) - } - if len(openingAudio) != 4 { - t.Fatal("Agent did not synthesize the approved opening audio") - } for _, segment := range []ai.ASRSegment{ {Source: "user", Text: "不用了", Final: false}, {Source: "assistant", Text: "不用了", Final: true}, {Source: "user", Text: "继续", Final: true}, } { - if stopped, err := call.HandleFinalASR(ctx, segment); err != nil || stopped { + if stopped, err := call.HandleFinalASR(ctx, segment); err != nil || stopped.EndedByKeyword { t.Fatal("interim, assistant or unmatched ASR text triggered hangup") } } @@ -125,8 +135,12 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T if err != nil || len(audio) != 4 { return fmt.Errorf("bound TTS SDK request: %w", err) } - if stopped, err := call.HandleFinalASR(ctx, ai.ASRSegment{Source: "user", Text: "我不用了", Final: true}); err != nil || !stopped || hangups != 1 { - t.Fatal("final user keyword did not hang up exactly once") + result, err := callflow.ExecuteApproved(ctx, session, ai.ModeFullAI, finalUserSDKPipeline{call}, callflow.CaptureConfig{FirstSpeechTimeout: time.Second, MaxDuration: 2 * time.Millisecond, MaxTurns: 20}) + if err != nil || len(result.Turns) != 1 || !result.Turn.EndedByKeyword || result.Turn.Reply != "好的,祝您生活愉快。" || len(result.OutboundTurns) != 2 || hangups != 1 { + t.Fatalf("SDK keyword closing must play then hang up without another turn: %v", err) + } + if _, err := call.RunTurn(ctx, nil); err == nil { + t.Fatal("another ASR turn started after closing") } return nil }) @@ -156,7 +170,10 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T t.Fatalf("full-AI isolated Agent call failed: err=%v calls=%d hangups=%d", err, calls, hangups) } llmCall := <-llmObserved - openingCall, ttsCall := <-ttsObserved, <-ttsObserved + ttsCall, openingCall, closingCall := <-ttsObserved, <-ttsObserved, <-ttsObserved + if closingCall.Body["input"].(map[string]any)["text"] != "好的,祝您生活愉快。" || len(llmObserved) != 0 || len(ttsObserved) != 0 { + t.Fatal("keyword closing text changed or extra LLM/TTS requests were sent") + } if llmCall.Authorization != "Bearer "+llm.Credential || llmCall.Body["model"] != "example-chat" || llmCall.Body["temperature"] != float64(0) || llmCall.Body["max_tokens"] != float64(256) { t.Fatal("Agent did not send approved LLM values through SDK") }