Merge task-configured AI protocols with SIP diagnostics

This commit is contained in:
2026-10-08 12:44:45 +08:00
21 changed files with 748 additions and 77 deletions
+2 -2
View File
@@ -94,8 +94,8 @@
- 独立 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 或真实接通。
- 2026-10-08 用户批准本地配置规则调整及新模型调用支持:SaaS SIP `revision` 可省略,完整读取覆盖旧配置;`transport/auth_mode` 判定忽略大小写,但不补其它缺失线路参数。Dispatcher 持久分配内部 SIP 加载代次并核对 Agent/Asterisk,旧通话排空与未知占用保留规则不变。任务列表仅忽略 `schema_version/control_seq`。任务和连接记录统一使用 `provider_id`;任务决定厂商/模型/ASR–LLM–TTS 用途,provider 只提供连接信息,不要求 `role/adapter/enabled`,未知供应商/模型/连接参数显式拒绝。新增本地调用代码及模拟测试不授权部署、拨号或真实 AI 请求。SQLite 当前布局因新增持久 SIP 快照变为 3;旧布局只读拒绝,不自动迁移或改写旧文件。
- AI 使用任务内不可变授权快照:仅经获批准百炼/火山 ASR(新增 `fun-asr-flash-8k-realtime-2026-01-28`)、OpenAI 兼容 LLM、百炼 TTS(`qwen3-tts-flash`/`Cherry`/`Chinese` 及 `cosyvoice-v3-flash` 的任务音色和速度)能表达的参数进入每通话实例;ASR-only 不启动 LLM/TTS,完整 AI 不借旧语音测试的授权或参数。只有最终用户 ASR 文本的明确字面关键词可触发拒联/挂断;`agent.conversation.hangup_keywords` 仅接受必填非空 `name`/`triggers`/`closingRemark` 的对象数组,多组命中按配置顺序取第一组,经获批 TTS 完整播放该组结束语后主动挂断,不调用 LLM、不开始新一轮对话;合成/播放失败须明确报错并尝试结束,不重播,旧字符串数组直接拒绝;不由 SDK 默认值、环境、CLI、metadata 或宽松 Schema 改写业务参数,不因 SDK 重试产生第二次发起/收费或重播。日志只存脱敏版本/摘要/计数,不存密钥、prompt、完整对话或音频。
- 2026-10-08 用户批准本地配置规则调整及新模型调用支持:SaaS SIP `revision` 可省略,完整读取覆盖旧配置;`transport/auth_mode` 判定忽略大小写,但不补其它缺失线路参数。Dispatcher 持久分配内部 SIP 加载代次并核对 Agent/Asterisk,旧通话排空与未知占用保留规则不变。任务列表仅忽略 `schema_version/control_seq`。任务和连接记录统一使用 `provider_id`;任务决定厂商/模型/ASR–LLM–TTS 用途,provider 只提供连接信息,不要求 `role/adapter/enabled`,未知供应商/协议及不可表达的连接参数显式拒绝。2026-10-08 追加用户确认:移除精确型号/音色业务硬编码;同协议模型及可表达参数从任务取得。TTS 必须明确 protocol(HTTP 或 task-based WebSocket),不按 model 名/地址存在与否猜协议、不失败回退。模型实际可用性由供应商调用返回事实决定,不以本地模拟代签。新增本地调用代码及模拟测试不授权部署、拨号或真实 AI 请求。SQLite 当前布局因新增持久 SIP 快照变为 3;旧布局只读拒绝,不自动迁移或改写旧文件。
- AI 使用任务内不可变授权快照:仅当前已实现且获批准的协议(火山 ASR、OpenAI 兼容 LLM、百炼 task-based ASR/TTS、百炼 HTTP TTS)可表达的参数进入每通话实例;不限制精确型号或音色,不从 SDK 默认值补任务参数。百炼 task ASR 支持任务明确选定 8/16 kHz;TTS 电话输出仍为单声道 16 kHz PCM16,HTTP speed=1、task speed=0.5–2;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 对象。
+2 -1
View File
@@ -176,8 +176,9 @@
},
"tts": {
"type": "object", "additionalProperties": false,
"required": ["provider_id", "model", "voice", "language_type", "format"],
"required": ["provider_id", "protocol", "model", "voice", "language_type", "format"],
"properties": {
"protocol": {"enum": ["dashscope_tts_http", "dashscope_task_websocket"], "$comment": "Task selects the protocol; never infer it from model names or endpoint availability."},
"provider_id": {"type": "string", "minLength": 1}, "model": {"type": "string", "minLength": 1}, "voice": {"type": "string", "minLength": 1}, "language_type": {"type": "string", "minLength": 1}, "speed": {"type": "number", "minimum": 0.25, "maximum": 3}, "timeout_ms": {"type": "integer", "minimum": 1}, "format": {"type": "object", "additionalProperties": false, "required": ["encoding", "sample_rate_hz", "channels"], "properties": {"encoding": {"enum": ["pcm_s16le", "pcma"]}, "sample_rate_hz": {"type": "integer", "minimum": 8000}, "channels": {"const": 1}}}
}
},
@@ -8,7 +8,7 @@
"immutable":true,"mode":"full_ai",
"asr":{"provider_id":"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_id":"llm-example","model":"example-chat","temperature":0,"max_tokens":256,"timeout_ms":5000},
"tts":{"provider_id":"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},
"tts":{"provider_id":"tts-example","protocol":"dashscope_tts_http","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":[{"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}
}
@@ -0,0 +1,15 @@
{
"resource":"task_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,
"task_id":"task-protocol-models","task_revision":2,"status":"running","name":"Isolated protocol fixture",
"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":[]},
"agent":{
"immutable":true,"mode":"full_ai",
"asr":{"provider_id":"bailian-example","model":"fixture-asr-vnext","language":"en-US","input":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1,"sample_width_bytes":2},"timeout_ms":5000},
"llm":{"provider_id":"bailian-example","model":"fixture-chat-vnext","temperature":0,"max_tokens":256,"timeout_ms":5000},
"tts":{"provider_id":"bailian-example","protocol":"dashscope_task_websocket","model":"fixture-tts-vnext","voice":"fixture-speaker","language_type":"English","speed":1.25,"format":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1},"timeout_ms":5000},
"prompt":{"text":"Isolated protocol fixture.","allowed_variables":[],"max_bytes":32768},
"conversation":{"opening":"Hello.","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}
}
}
@@ -0,0 +1,15 @@
{
"resource":"task_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,
"task_id":"task-missing-protocol","task_revision":2,"status":"running","name":"Isolated invalid fixture",
"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":[]},
"agent":{
"immutable":true,"mode":"full_ai",
"asr":{"provider_id":"bailian-example","model":"fixture-asr-vnext","language":"en-US","input":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1,"sample_width_bytes":2}},
"llm":{"provider_id":"bailian-example","model":"fixture-chat-vnext"},
"tts":{"provider_id":"bailian-example","model":"cosyvoice-v3-flash","voice":"fixture-speaker","language_type":"English","speed":1,"format":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1}},
"prompt":{"text":"Isolated invalid fixture.","allowed_variables":[]},
"conversation":{"opening":"Hello.","hangup_keywords":[]}
}
}
+2 -2
View File
@@ -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": "63c76130981b9daf119ee3caa47b502feb0fc3b6cff09b06efd238c0bb9cdde7"
"docs/thirds/saas-dispatcher.md": "fbe73272923f505b00e1b0f6177cfa13c4a29ffab27c1cac15bed5b460e2cec2"
},
"bundle_sha256": "412ee384ce66eaaa37a89ae32ff3a2e2fb667f1ac1edf24f29a75231f4ac4cef",
"bundle_sha256": "f23f612db4863ceec7c3be0c8e71e636dee18e9f0d89cf0dc0db948e439be2ea",
"bundle_algorithm": "sha256 of sorted relative-path + space + sha256(file) + newline; only root-level JSON and examples/**/*.json, excluding manifest.json"
}
@@ -0,0 +1,43 @@
# AI 按协议复用、模型由任务配置(2026-10-08)
## 范围与结论
用户批准取消精确型号/音色限制,并明确选择由任务增加 TTS `protocol` 来区分已有 HTTP 与 task-based WebSocket 调用。provider 仍仅提供连接信息,不恢复 role/adapter。
本轮基线为 `98b3fd4`;全部修改在独立工作区 `/home/rogee/Workspace/sip-go-agent-protocol-models`、分支 `change/ai-protocol-model-config-20261008`,不写入原工作区的未提交修改。没有部署、拨号、真实 AI/OSS 请求或 SaaS 业务数据修改;没有再次读取实际 SaaS 配置,也未复制任何私有 env。
**同协议新增模型/音色不再要求修改本方代码。** 模型 ID 原样进入请求;这不代表任意模型都支持该协议,实际服务端拒绝时如实失败,不换模型或协议重试。
当前规范见 [SaaS↔Dispatcher](../thirds/saas-dispatcher.md),字段与正反例以 `contracts/local/` 为准。本记录不另建字段合同;此前真实 SaaS 剩余缺口见 [前轮证据](saas-config-model-adjustments-20261008.md),本轮没有宣称已修复或签收。
## 提供方需调整的字段
在任务详情 `agent.tts` 增加必填 `protocol`:
| 值 | 本方选择 |
| --- | --- |
| `dashscope_tts_http` | 已有 HTTP 生成接口、音频引用下载及格式转换;连接必须提供正确的 HTTP 生成 endpoint |
| `dashscope_task_websocket` | 已有 task-based WebSocket TTS;连接必须提供 ws_endpoint |
字段缺失、未知值或选定 endpoint 缺失直接拒绝,不依据 model 名或备用地址猜测。其它由提供方修正的 SIP、任务详情和额度字段要求不变。
模型、音色、语言、速度等来自任务,可表达的参数不被本地型号列表替换。协议本身仍有约束:HTTP 调用目前不表达倍速,须明确 speed=1;task TTS 支持 speed=0.5–2。电话播放目标仍为单声道 16 kHz PCM16。语言名称或 BCP 47 标签转为任务协议需要的显式语言提示;指定 Auto 则按协议省略提示,未指定语言不自动推测。
百炼 task ASR 接受任务指定的模型及 8/16 kHz 输入目标:16 kHz 发送原始 Agent PCM,8 kHz 才转换。当前适配不能表达 interim 开关,显式提供即报错,不忽略。火山 ASR 保留其 SDK 能表达的参数约束;没有新增供应商或协议。
## 本地验证
- TDD:新增同协议任意模型/音色及明确协议用例,初始校验失败;实现后通过。新增 interim 不可表达用例先失败,再改为显式拒绝。
- HTTP 和 WebSocket 实际本地模拟请求验证 model/voice/language/rate 原值;16 kHz ASR 字节完全保持、8 kHz 仅显式转换。使用未出现在任何型号列表中的模型与音色,未访问真实供应商。
- 熟悉型号名称搭配相反协议时仍按任务明确选择请求;服务端 400 只失败一次,不访问另一协议地址。
- 缺协议、缺选定 endpoint、未知参数、不支持的倍速/采样率/语言、未明确 ASR 模型和 interim 选择均明确拒绝;不插入 SDK/环境默认值。
- 新增完整开场→ASR→LLM→TTS 回合,三个未列入任何型号白名单的模型 ID 均原样进入各自本地协议请求;一次 ASR、一次 LLM、两次 TTS(开场/回复),没有额外请求。原关键词不调用 LLM/结束语播放后挂断、完整音频/最终文本、错身份、超时、错误、无隐式重试等用例仍通过。
- `make check` 全部通过:格式、当前合同及来源/hash、Proto、`go test -race ./...`、vet、构建、实际隔离 RabbitMQ、既有 HTTPS 与双向 TLS 端到端测试。
- 42 个合同 JSON 正反例/hash 通过。新增缺 protocol 反例已单独确认只有这一项结构错误。
- `make coverage` 业务语句覆盖率 **70.7%**(门槛 65%)。本地日志 `.local/protocol-make-check-final.log`、`.local/protocol-coverage-final.log`。
- 隔离 MQ 验收曾一次出现 ready backlog 数为 2 的检查失败;该次缺少 pending 身份诊断,不能据此推断消息丢失。检查发现测试把 socket 写完当作消息已入队,并把准入探测消息当成必然留在队列的 backlog;生产发布者已有 mandatory/confirm,业务代码不需要改变。本轮修正测试发布者确认边界,并在消费者确已停止后单独投递三条已确认 backlog,保留原 purge、零消费者、回执与不执行检查;异常时记录队列状态及 pending event IDs。修正后该用例实际 RabbitMQ、race 连续 **50 次通过**,随后完整验收再次通过。重复验证中的一次容器启动未就绪单独失败,未跳过 gate 或延长 readiness timeout;只读资源诊断未证明资源不足,重新独立启动后完成 50 次验证,不声称修复了 broker 本身。
- A01–A12/K01–K16 既有归属、冻结快照、额度、控制、幂等、录音恢复、未知占用等测试未放宽;范围与外部缺口继续遵循 [P08 证据](saas-dispatcher-p08-acceptance.md)。本轮不代签真实供应商、模型、Asterisk 媒体、通话或生产门禁。
## 依赖与数据边界
语言解析使用项目已有 `golang.org/x/text v0.41.0`,仅从间接声明改为直接声明,没有升级版本或新增库;HTTP、WebSocket、音频转换继续复用现有依赖。没有新数据库表、自动迁移或历史路径兼容;未动任何现存数据库、恢复文件、spool/outbox 或外部资源。
+2 -2
View File
@@ -9,7 +9,7 @@
| 资源 | 方法和路径 | 应用规则 | 正例 |
| --- | --- | --- | --- |
| SIP 全量 | `GET /internal/v1/dispatcher/sip` | SaaS `revision` 可省略,最新读取的完整快照覆盖旧配置,不以其大小拒绝更新;D 持久分配内部加载代次,重启/重试不丢失待加载内容,并用该代次核对 Agent/Asterisk 实际加载。`transport/auth_mode` 判定忽略大小写;未知 `transport/auth_mode/registration_required/max_concurrent_calls` 为 `null`,不能作为可执行线路默认值;变更时关准入、排空并核验 Agent/Asterisk 实际加载 | [`sip`](../../contracts/local/examples/config-read-sip.json) |
| AI provider 全量 | `GET /internal/v1/dispatcher/ai-providers` | 启动时一次读取全局列表,本进程全部任务共用到下一次重启;运行中 start/resume 不重复拉取,服务商变更须重启后才生效。任务以 `provider_id` 关联连接记录并决定 ASR/LLM/TTS 用途及模型;provider 只提供供应商标识和 `api_endpoint/ws_endpoint/api_key` 等连接信息,不要求 `role/adapter/enabled`。按明确供应商及任务模型选择已有调用实现,不按名称猜测协议;缺失连接、未知供应商/模型或不支持的连接参数显式拒绝 | [`providers`](../../contracts/local/examples/config-read-providers.json) |
| AI provider 全量 | `GET /internal/v1/dispatcher/ai-providers` | 启动时一次读取全局列表,本进程全部任务共用到下一次重启;运行中 start/resume 不重复拉取,服务商变更须重启后才生效。任务以 `provider_id` 关联连接记录并决定 ASR/LLM/TTS 用途及模型;provider 只提供供应商标识和 `api_endpoint/ws_endpoint/api_key` 等连接信息,不要求 `role/adapter/enabled`。按明确供应商及任务调用协议选择已有实现,模型 ID 是原样请求参数,不作本地型号白名单或协议选择依据;TTS 必须在任务中明确 `protocol`,provider 不增加 role/adapter。缺失连接、未知供应商/协议或不支持的参数显式拒绝;服务端拒绝模型时如实失败,不切换协议/模型重试 | [`providers`](../../contracts/local/examples/config-read-providers.json) |
| 任务 | `GET /internal/v1/dispatcher/task/{task_id}` | 归属、状态和 `task_revision` 校验后持久绑定同一 `agent`/provider 快照;不可变摘要由本项目计算。同一 revision 内容不同拒绝;ASR-only 只需 ASR,不强制 LLM/TTS;0 和 false 原样保留 | [`ASR-only`](../../contracts/local/examples/config-read-task-asr.json) · [`full AI`](../../contracts/local/examples/config-read-task-full.json) |
| 任务列表 | `GET /internal/v1/dispatcher/tasks`,后续 `?after=<cursor>` | `schema_version/control_seq` 忽略且不作版本或去重依据;其它未知字段仍拒绝。首次/重启完整读到 **带 cursor 且 tasks=[]** 的终止页;非空短页不能提前结束。非空页持久成功后才使用下一 cursor;控制队列积压处理前不开新任务准入。不使用旧 snapshot_id/watermark/mode=snapshot | [`page`](../../contracts/local/examples/task-discovery-page.json) · [`end`](../../contracts/local/examples/task-discovery-end.json) |
| 租户额度 | `GET /internal/v1/dispatcher/tenant/{tenant_id}/quota` | `quota_revision` 为业务修订号;未知占用不得算成已释放 | [`quota`](../../contracts/local/examples/config-read-quota.json) |
@@ -33,7 +33,7 @@ RabbitMQ 是 Topic,**SaaS 独占创建、绑定、退役 exchange/queue,D
## 调度、AI 与真实结果(K01–K09、K11–K14)
- AI 的用途及模型由任务指定,provider 只提供连接信息。保留火山 ASR、OpenAI 兼容 LLM 与百炼 `qwen3-tts-flash`(`Cherry/Chinese/1`)调用;新增百炼 `fun-asr-flash-8k-realtime-2026-01-28` 和 `cosyvoice-v3-flash` 的正式 WebSocket 调用代码。Fun-ASR 显式把 Agent 16 kHz PCM16 转为模型的 8 kHz,只在 task-finished 后返回最终用户文本;CosyVoice 使用任务原模型、音色和速度(0.5–2),显式请求单声道 16 kHz PCM16,并收到 task-finished 后才交付完整音频。每段只启动一次 task,不隐式重试;超时、身份不匹配、失败或不完整结果显式报错,不返回半成品,不回退其它模型。仅本地模拟验证,不构成真实模型、语音或通话签收。
- AI 的用途、模型及可表达参数由任务指定,provider 只提供连接信息;复用已实现的火山 ASR、OpenAI 兼容 LLM、百炼 task-based ASR/TTS 和百炼 HTTP TTS 协议,不为每个型号或音色维护代码白名单。TTS 的 `protocol` 必填:`dashscope_tts_http` 选择 HTTP 生成与音频引用下载,`dashscope_task_websocket` 选择 WebSocket task;不从 model 名、地址是否存在或失败结果猜测/切换协议。HTTP 协议当前不能表达倍速,只接受明确 speed=1;task TTS 原样提交 voice/rate(0.5–2),语言标签/BCP 47 转为明确 ISO 语言提示,显式 Auto 不发送提示,不能表达的语言拒绝。百炼 task ASR 的 model/language 必须明确,当前协议适配不能表达 interim 开关,显式提供即拒绝而不忽略;支持任务选定 8/16 kHz PCM16:16 kHz 发送原始 Agent PCM,8 kHz 才显式转换。电话 TTS 输出仍要求单声道 16 kHz PCM16;转换或完整性失败明确报错,不截断/重标格式。task-based 请求只在 task-finished 后返回最终用户文本/完整音频。每段只启动一次请求/task,不隐式重试;超时、身份不匹配、失败或不完整结果不返回半成品,不回退其它协议或模型。新增同协议模型仅修改任务配置,但模型实际可用性由供应商返回事实决定;本地 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={}`。
- `agent.conversation.hangup_keywords` 使用对象数组,字段及正例以现行 Schema 为准;旧字符串数组直接拒绝。只有**用户侧 ASR 最终识别文本**字面包含某组 `triggers` 才触发,同时命中多组按配置顺序取第一组;使用该组 `closingRemark` 合成语音,完整播放后主动挂断,不调用 LLM 生成结束语。匹配后不再开启新一轮对话,重复识别不得重播或重复挂断;中间识别、助手回复、开场白、TTS 均不能触发。合成或播放失败时明确报错并尝试结束通话,不重试、不把未播放音频报告为已播放;挂断结果不明仍保留未知事实。组名、触发词列表及结束语均必填且非空;ASR-only 不配置或合成结束语。同一任务 revision 不同内容拒绝准入;provider 连接缺失或选定模型无法表达获批参数时不可调用。Mock 参数验证不等于真实供应商验收。
+1 -1
View File
@@ -15,6 +15,7 @@ require (
github.com/shirou/gopsutil/v4 v4.26.8
github.com/spf13/cobra v1.10.1
github.com/zaf/g711 v1.4.0
golang.org/x/text v0.41.0
golang.org/x/time v0.4.0
google.golang.org/grpc v1.83.2
google.golang.org/protobuf v1.36.12
@@ -57,7 +58,6 @@ require (
golang.org/x/net v0.58.0 // indirect
golang.org/x/sync v0.22.0 // indirect
golang.org/x/sys v0.47.0 // indirect
golang.org/x/text v0.41.0 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect
modernc.org/libc v1.75.7 // indirect
modernc.org/mathutil v1.7.1 // indirect
+2
View File
@@ -12,6 +12,7 @@ func bailianFixture(t *testing.T) Binding {
task, _ := currentFixture(t, "full_ai")
task = changeCurrentAgent(t, task, func(a map[string]any) {
asr := a["asr"].(map[string]any)
delete(asr, "interim")
asr["provider_id"] = "bailian"
delete(asr, "provider_ref")
asr["model"] = "fun-asr-flash-8k-realtime-2026-01-28"
@@ -23,6 +24,7 @@ func bailianFixture(t *testing.T) Binding {
tts := a["tts"].(map[string]any)
tts["provider_id"] = "bailian"
delete(tts, "provider_ref")
tts["protocol"] = TTSProtocolDashScopeTask
tts["model"] = "cosyvoice-v3-flash"
tts["voice"] = "longanyang"
tts["speed"] = 1.25
+49 -34
View File
@@ -150,60 +150,67 @@ func runBailianTask(ctx context.Context, endpoint, key string, payload any, send
}
func recognizeBailian(ctx context.Context, approved ASRConfig, pcm16 []byte) (string, error) {
if approved.Model != "fun-asr-flash-8k-realtime-2026-01-28" || approved.Language != "zh-CN" || approved.Request.SampleRate != 8000 {
return "", errors.New("unapproved Fun-ASR settings")
if strings.TrimSpace(approved.Model) == "" || (approved.Request.SampleRate != 8000 && approved.Request.SampleRate != 16000) {
return "", errors.New("Bailian task ASR needs an explicit model and 8000/16000 Hz PCM16")
}
if len(pcm16) > maxBailianPCMBytes {
return "", errors.New("Fun-ASR input exceeds approved PCM bound")
}
// Agent conversation media is PCM16/16000; the selected 8k model needs an
// explicit conversion, never a sample-rate label change or byte truncation.
cmd := exec.CommandContext(ctx, "ffmpeg", "-hide_banner", "-loglevel", "error", "-f", "s16le", "-ar", "16000", "-ac", "1", "-i", "pipe:0", "-f", "s16le", "-ar", "8000", "-ac", "1", "pipe:1")
cmd.Stdin = bytes.NewReader(pcm16)
input, err := cmd.Output()
hint, err := bailianLanguageHint(approved.Language)
if err != nil {
exitCode := -1
var exitErr *exec.ExitError
if errors.As(err, &exitErr) {
exitCode = exitErr.ExitCode()
return "", err
}
if len(pcm16) == 0 || len(pcm16)%2 != 0 || len(pcm16) > maxBailianPCMBytes {
return "", errors.New("Bailian task ASR input is empty, incomplete or oversized PCM16")
}
// Agent media is PCM16/16000. Only convert when the task explicitly selects
// 8000 Hz; a 16000 Hz request sends the original bytes, without relabeling.
input := pcm16
if approved.Request.SampleRate == 8000 {
cmd := exec.CommandContext(ctx, "ffmpeg", "-hide_banner", "-loglevel", "error", "-f", "s16le", "-ar", "16000", "-ac", "1", "-i", "pipe:0", "-f", "s16le", "-ar", "8000", "-ac", "1", "pipe:1")
cmd.Stdin = bytes.NewReader(pcm16)
input, err = cmd.Output()
if err != nil {
exitCode := -1
var exitErr *exec.ExitError
if errors.As(err, &exitErr) {
exitCode = exitErr.ExitCode()
}
return "", fmt.Errorf("Bailian task ASR PCM16 conversion failed (cause=%T, exit_code=%d, input_bytes=%d)", err, exitCode, len(pcm16))
}
return "", fmt.Errorf("Fun-ASR PCM16 16000-to-8000 conversion failed (cause=%T, exit_code=%d, input_bytes=%d)", err, exitCode, len(pcm16))
if len(input) == 0 || len(input)%2 != 0 {
return "", errors.New("Bailian task ASR resampler returned incomplete PCM16")
}
slog.Debug("Bailian task ASR PCM resampling completed", "input_sample_rate", 16000, "output_sample_rate", 8000, "input_bytes", len(pcm16), "output_bytes", len(input))
}
if len(input) == 0 || len(input)%2 != 0 {
return "", errors.New("Fun-ASR resampler returned incomplete PCM16")
}
slog.Debug("Fun-ASR PCM resampling completed", "input_sample_rate", 16000, "output_sample_rate", 8000, "input_bytes", len(pcm16), "output_bytes", len(input))
payload := map[string]any{"task_group": "audio", "task": "asr", "function": "recognition", "model": approved.Model, "parameters": map[string]any{"format": "pcm", "sample_rate": 8000, "language_hints": []string{"zh"}}, "input": map[string]any{}}
payload := map[string]any{"task_group": "audio", "task": "asr", "function": "recognition", "model": approved.Model, "parameters": map[string]any{"format": "pcm", "sample_rate": int(approved.Request.SampleRate), "language_hints": []string{hint}}, "input": map[string]any{}}
var finals []string
finalBytes := 0
seen := map[string]string{}
err = runBailianTask(ctx, approved.Provider.Endpoint, approved.Provider.Credential, payload, func(t bailianTask) error {
for off := 0; off < len(input); off += 3200 {
if err := t.conn.Write(ctx, websocket.MessageBinary, input[off:min(off+3200, len(input))]); err != nil {
return errors.New("Fun-ASR audio delivery failed")
return errors.New("Bailian task ASR audio delivery failed")
}
}
return t.command(ctx, "finish-task", map[string]any{"input": map[string]any{}})
}, func(kind websocket.MessageType, _ []byte, event bailianEvent) error {
if kind != websocket.MessageText {
return errors.New("Fun-ASR returned binary output")
return errors.New("Bailian task ASR returned binary output")
}
s := event.Payload.Output.Sentence
if !s.Final {
return nil
}
if s.Begin == nil || s.End == nil || *s.End < *s.Begin || strings.TrimSpace(s.Text) == "" {
return errors.New("Fun-ASR returned an incomplete final sentence")
return errors.New("Bailian task ASR returned an incomplete final sentence")
}
id := strconv.FormatInt(*s.Begin, 10) + ":" + strconv.FormatInt(*s.End, 10)
if old, ok := seen[id]; ok {
if old != s.Text {
return errors.New("Fun-ASR changed an already final sentence")
return errors.New("Bailian task ASR changed an already final sentence")
}
return nil
}
if len(finals) >= 1024 || len(s.Text) > maxBailianPCMBytes-finalBytes {
return errors.New("Fun-ASR final text exceeds per-turn protocol bound")
return errors.New("Bailian task ASR final text exceeds per-turn protocol bound")
}
finalBytes += len(s.Text)
seen[id] = s.Text
@@ -214,18 +221,26 @@ func recognizeBailian(ctx context.Context, approved ASRConfig, pcm16 []byte) (st
return "", err
}
if len(finals) == 0 {
return "", errors.New("Fun-ASR finished without final user text")
return "", errors.New("Bailian task ASR finished without final user text")
}
return strings.Join(finals, ""), nil
}
func synthesizeCosyVoice(ctx context.Context, approved TTSConfig, text string) ([]byte, error) {
if approved.Model != "cosyvoice-v3-flash" || approved.Voice == "" || approved.LanguageType != "Chinese" || approved.SampleRate != 16000 || approved.Speed < 0.5 || approved.Speed > 2 || strings.TrimSpace(text) == "" {
return nil, errors.New("unapproved CosyVoice settings or empty text")
func synthesizeBailianTaskTTS(ctx context.Context, approved TTSConfig, text string) ([]byte, error) {
if approved.Protocol != TTSProtocolDashScopeTask || strings.TrimSpace(approved.Model) == "" || strings.TrimSpace(approved.Voice) == "" || approved.SampleRate != 16000 || approved.Speed < 0.5 || approved.Speed > 2 || strings.TrimSpace(text) == "" {
return nil, errors.New("unapproved Bailian task TTS settings or empty text")
}
payload := map[string]any{"task_group": "audio", "task": "tts", "function": "SpeechSynthesizer", "model": approved.Model, "parameters": map[string]any{"text_type": "PlainText", "voice": approved.Voice, "format": "pcm", "sample_rate": approved.SampleRate, "rate": approved.Speed, "language_hints": []string{"zh"}}, "input": map[string]any{}}
hints, err := bailianTTSLanguageHints(approved.LanguageType)
if err != nil {
return nil, err
}
parameters := map[string]any{"text_type": "PlainText", "voice": approved.Voice, "format": "pcm", "sample_rate": approved.SampleRate, "rate": approved.Speed}
if len(hints) != 0 {
parameters["language_hints"] = hints
}
payload := map[string]any{"task_group": "audio", "task": "tts", "function": "SpeechSynthesizer", "model": approved.Model, "parameters": parameters, "input": map[string]any{}}
var audio []byte
err := runBailianTask(ctx, approved.Provider.WSEndpoint, approved.Provider.Credential, payload, func(t bailianTask) error {
err = runBailianTask(ctx, approved.Provider.WSEndpoint, approved.Provider.Credential, payload, func(t bailianTask) error {
if err := t.command(ctx, "continue-task", map[string]any{"input": map[string]any{"text": text}}); err != nil {
return err
}
@@ -233,7 +248,7 @@ func synthesizeCosyVoice(ctx context.Context, approved TTSConfig, text string) (
}, func(kind websocket.MessageType, body []byte, _ bailianEvent) error {
if kind == websocket.MessageBinary {
if len(body) > maxBailianPCMBytes-len(audio) {
return errors.New("CosyVoice audio exceeds approved PCM bound")
return errors.New("Bailian task TTS audio exceeds approved PCM bound")
}
audio = append(audio, body...)
}
@@ -243,8 +258,8 @@ func synthesizeCosyVoice(ctx context.Context, approved TTSConfig, text string) (
return nil, err
}
if len(audio) == 0 || len(audio)%2 != 0 {
return nil, errors.New("CosyVoice finished without complete PCM16 audio")
return nil, errors.New("Bailian task TTS finished without complete PCM16 audio")
}
slog.Debug("CosyVoice complete PCM16 received", "sample_rate", approved.SampleRate, "bytes", len(audio))
slog.Debug("Bailian task TTS complete PCM16 received", "sample_rate", approved.SampleRate, "bytes", len(audio))
return audio, nil
}
+1 -1
View File
@@ -28,7 +28,7 @@ func synthesizeBailianTTS(ctx context.Context, approved TTSConfig, text string)
if err := bailianURL(approved.Provider.Endpoint, false); err != nil {
return nil, fmt.Errorf("approved Bailian endpoint: %w", err)
}
if approved.Model != "qwen3-tts-flash" || approved.Voice != "Cherry" || approved.LanguageType != "Chinese" || approved.Speed != 1 || approved.SampleRate != 16000 || approved.Provider.Credential == "" || text == "" {
if approved.Protocol != TTSProtocolDashScopeHTTP || strings.TrimSpace(approved.Model) == "" || strings.TrimSpace(approved.Voice) == "" || strings.TrimSpace(approved.LanguageType) == "" || approved.Speed != 1 || approved.SampleRate != 16000 || approved.Provider.Credential == "" || text == "" {
return nil, errors.New("approved Bailian TTS settings or text are incomplete")
}
// Never follow a provider or audio redirect into an unapproved second request.
+49 -26
View File
@@ -48,6 +48,7 @@ type LLMConfig struct {
type TTSConfig struct {
Provider configread.Provider
Protocol string
Model string
Voice string
LanguageType string
@@ -91,6 +92,7 @@ type currentAgentSettings struct {
} `json:"llm"`
TTS *struct {
ProviderRef string `json:"provider_id"`
Protocol string `json:"protocol"`
Model string `json:"model"`
Voice string `json:"voice"`
LanguageType string `json:"language_type"`
@@ -140,7 +142,7 @@ func Bind(task configread.Task, providers map[string]configread.Provider) (Bindi
return Binding{}, fmt.Errorf("decode immutable AI settings: %w", err)
}
bound := Binding{Mode: settings.Mode}
asrProvider, err := currentProvider(providers, settings.ASR.ProviderRef, "asr", settings.ASR.Model)
asrProvider, err := currentProvider(providers, settings.ASR.ProviderRef, "asr", "")
if err != nil {
return Binding{}, err
}
@@ -151,18 +153,26 @@ func Bind(task configread.Task, providers map[string]configread.Provider) (Bindi
if err != nil {
return Binding{}, fmt.Errorf("ASR input: %w", err)
}
if asrProvider.Code == "ali_bailian" {
if settings.ASR.Model != "fun-asr-flash-8k-realtime-2026-01-28" || sampleRate != 8000 || settings.ASR.Language != "zh-CN" {
return Binding{}, errors.New("Bailian ASR requires the approved Fun-ASR 8k model, 8000 Hz PCM16 and zh-CN")
}
} else if sampleRate != 16000 {
return Binding{}, errors.New("ASR media requires 16000 Hz PCM16")
}
lang := doubaospeech.Language(settings.ASR.Language)
switch lang {
case doubaospeech.LanguageZhCN, doubaospeech.LanguageEnUS, doubaospeech.LanguageJaJP, doubaospeech.LanguageKoKR:
default:
return Binding{}, errors.New("ASR language is unsupported by selected SDK")
if asrProvider.Code == "ali_bailian" {
if settings.ASR.Interim != nil {
return Binding{}, errors.New("Bailian task ASR protocol cannot express interim selection")
}
if strings.TrimSpace(settings.ASR.Model) == "" || (sampleRate != 8000 && sampleRate != 16000) {
return Binding{}, errors.New("Bailian task ASR requires an explicit model and 8000/16000 Hz PCM16")
}
if _, err := bailianLanguageHint(settings.ASR.Language); err != nil {
return Binding{}, fmt.Errorf("ASR language: %w", err)
}
} else {
if sampleRate != 16000 {
return Binding{}, errors.New("ASR media requires 16000 Hz PCM16")
}
switch lang {
case doubaospeech.LanguageZhCN, doubaospeech.LanguageEnUS, doubaospeech.LanguageJaJP, doubaospeech.LanguageKoKR:
default:
return Binding{}, errors.New("ASR language is unsupported by selected SDK")
}
}
asrTimeout, err := currentTimeout(settings.ASR.TimeoutMS)
if err != nil {
@@ -195,7 +205,7 @@ func Bind(task configread.Task, providers map[string]configread.Provider) (Bindi
default:
return Binding{}, errors.New("AI mode is unsupported")
}
llmProvider, err := currentProvider(providers, settings.LLM.ProviderRef, "llm", settings.LLM.Model)
llmProvider, err := currentProvider(providers, settings.LLM.ProviderRef, "llm", "")
if err != nil {
return Binding{}, err
}
@@ -204,7 +214,7 @@ func Bind(task configread.Task, providers map[string]configread.Provider) (Bindi
return Binding{}, fmt.Errorf("LLM timeout: %w", err)
}
bound.LLM = &LLMConfig{Provider: llmProvider, Model: settings.LLM.Model, Temperature: settings.LLM.Temperature, MaxTokens: settings.LLM.MaxTokens, Timeout: llmTimeout}
ttsProvider, err := currentProvider(providers, settings.TTS.ProviderRef, "tts", settings.TTS.Model)
ttsProvider, err := currentProvider(providers, settings.TTS.ProviderRef, "tts", settings.TTS.Protocol)
if err != nil {
return Binding{}, err
}
@@ -218,24 +228,30 @@ func Bind(task configread.Task, providers map[string]configread.Provider) (Bindi
if ttsSampleRate != doubaospeech.SampleRate(16000) {
return Binding{}, errors.New("TTS media requires 16000 Hz PCM16")
}
switch settings.TTS.Model {
case "qwen3-tts-flash":
if settings.TTS.Voice != "Cherry" || settings.TTS.LanguageType != "Chinese" || settings.TTS.Speed == nil || *settings.TTS.Speed != 1 {
return Binding{}, errors.New("approved Bailian TTS voice, language or speed is unsupported")
if strings.TrimSpace(settings.TTS.Model) == "" || strings.TrimSpace(settings.TTS.Voice) == "" || strings.TrimSpace(settings.TTS.LanguageType) == "" || settings.TTS.Speed == nil {
return Binding{}, errors.New("TTS model, voice, language and speed must be explicit")
}
switch settings.TTS.Protocol {
case TTSProtocolDashScopeHTTP:
if *settings.TTS.Speed != 1 {
return Binding{}, errors.New("TTS HTTP protocol cannot express the requested speed")
}
case "cosyvoice-v3-flash":
if settings.TTS.LanguageType != "Chinese" || settings.TTS.Speed == nil || *settings.TTS.Speed < 0.5 || *settings.TTS.Speed > 2 || ttsProvider.WSEndpoint == "" {
return Binding{}, errors.New("CosyVoice TTS language, speed or WebSocket connection is unsupported")
case TTSProtocolDashScopeTask:
if *settings.TTS.Speed < 0.5 || *settings.TTS.Speed > 2 {
return Binding{}, errors.New("TTS task protocol speed is outside 0.5–2")
}
if _, err := bailianTTSLanguageHints(settings.TTS.LanguageType); err != nil {
return Binding{}, fmt.Errorf("TTS language: %w", err)
}
default:
return Binding{}, errors.New("approved Bailian TTS model is unsupported")
return Binding{}, errors.New("TTS protocol is unsupported")
}
ttsTimeout, err := currentTimeout(settings.TTS.TimeoutMS)
if err != nil {
return Binding{}, fmt.Errorf("TTS timeout: %w", err)
}
bound.TTS = &TTSConfig{
Provider: ttsProvider, Model: settings.TTS.Model, Voice: settings.TTS.Voice, LanguageType: settings.TTS.LanguageType,
Provider: ttsProvider, Protocol: settings.TTS.Protocol, Model: settings.TTS.Model, Voice: settings.TTS.Voice, LanguageType: settings.TTS.LanguageType,
Speed: *settings.TTS.Speed, SampleRate: int(ttsSampleRate), Timeout: ttsTimeout,
}
bound.Prompt = settings.Prompt.Text
@@ -272,7 +288,7 @@ func Bind(task configread.Task, providers map[string]configread.Provider) (Bindi
return bound, nil
}
func currentProvider(providers map[string]configread.Provider, ref, role, model string) (configread.Provider, error) {
func currentProvider(providers map[string]configread.Provider, ref, role, protocol string) (configread.Provider, error) {
p, found := providers[ref]
if !found || ref == "" || p.ProviderRef != ref || p.Credential == "" {
return configread.Provider{}, fmt.Errorf("%s provider connection is missing or incomplete", role)
@@ -299,8 +315,12 @@ func currentProvider(providers map[string]configread.Provider, ref, role, model
if p.Code != "ali_bailian" {
return configread.Provider{}, errors.New("TTS provider is unsupported")
}
if model == "cosyvoice-v3-flash" {
switch protocol {
case TTSProtocolDashScopeTask:
p.Endpoint = p.WSEndpoint
case TTSProtocolDashScopeHTTP:
default:
return configread.Provider{}, errors.New("TTS protocol is unsupported")
}
default:
return configread.Provider{}, errors.New("AI purpose is unsupported")
@@ -309,9 +329,12 @@ func currentProvider(providers map[string]configread.Provider, ref, role, model
if err != nil || u.Host == "" || u.User != nil || u.RawQuery != "" || u.Fragment != "" || (u.Scheme != "https" && u.Scheme != "http" && u.Scheme != "wss" && u.Scheme != "ws") || strings.ContainsAny(p.Credential, "\r\n") {
return configread.Provider{}, fmt.Errorf("%s provider endpoint or credential is invalid", role)
}
if (role == "asr" && p.Code == "ali_bailian" || role == "tts" && model == "cosyvoice-v3-flash") && u.Scheme != "wss" && u.Scheme != "ws" {
if (role == "asr" && p.Code == "ali_bailian" || role == "tts" && protocol == TTSProtocolDashScopeTask) && u.Scheme != "wss" && u.Scheme != "ws" {
return configread.Provider{}, errors.New("selected speech model requires a WebSocket connection")
}
if (role == "llm" || role == "tts" && protocol == TTSProtocolDashScopeHTTP) && u.Scheme != "http" && u.Scheme != "https" {
return configread.Provider{}, errors.New("selected protocol requires an HTTP connection")
}
return p, nil
}
-1
View File
@@ -106,7 +106,6 @@ func TestBindCurrentRejectsSDKUnsupportedTTSWithoutChangingSchema(t *testing.T)
{"speed-below", func(tts map[string]any) { tts["speed"] = 0.25 }},
{"speed-above", func(tts map[string]any) { tts["speed"] = 3.0 }},
{"speed-unrepresentable", func(tts map[string]any) { tts["speed"] = 1.005 }},
{"language-type", func(tts map[string]any) { tts["language_type"] = "English" }},
{"sample-rate", func(tts map[string]any) { tts["format"].(map[string]any)["sample_rate_hz"] = 12345 }},
} {
t.Run(tc.name, func(t *testing.T) {
+6 -3
View File
@@ -262,10 +262,13 @@ func (b Binding) synthesize(ctx context.Context, text string, alreadyPending int
// missing or failed response never counts as queued audio.
var audio []byte
var err error
if b.TTS.Model == "cosyvoice-v3-flash" {
audio, err = synthesizeCosyVoice(ctx, *b.TTS, text)
} else {
switch b.TTS.Protocol {
case TTSProtocolDashScopeTask:
audio, err = synthesizeBailianTaskTTS(ctx, *b.TTS, text)
case TTSProtocolDashScopeHTTP:
audio, err = synthesizeBailianTTS(ctx, *b.TTS, text)
default:
return nil, 0, errors.New("TTS protocol is missing or unsupported")
}
if err != nil {
kind, status := bailianFailureSummary(err)
+73
View File
@@ -0,0 +1,73 @@
package ai
import (
"git.ipao.vip/rogee/go-sip/internal/configread"
"testing"
)
func TestProtocolCapabilitiesRejectUnexpressibleSettingsWithoutFallback(t *testing.T) {
for _, tc := range []struct {
name, protocol string
change func(map[string]any)
connection func(map[string]configread.Provider)
}{
{name: "HTTP-speed", protocol: TTSProtocolDashScopeHTTP, change: func(a map[string]any) { a["tts"].(map[string]any)["speed"] = 1.5 }},
{name: "HTTP-missing-endpoint", protocol: TTSProtocolDashScopeHTTP, connection: func(p map[string]configread.Provider) { x := p["tts-example"]; x.Endpoint = ""; p[x.ProviderRef] = x }},
{name: "WS-missing-endpoint", protocol: TTSProtocolDashScopeTask, connection: func(p map[string]configread.Provider) { x := p["tts-example"]; x.WSEndpoint = ""; p[x.ProviderRef] = x }},
{name: "WS-speed", protocol: TTSProtocolDashScopeTask, change: func(a map[string]any) { a["tts"].(map[string]any)["speed"] = 2.5 }},
{name: "WS-language", protocol: TTSProtocolDashScopeTask, change: func(a map[string]any) { a["tts"].(map[string]any)["language_type"] = "not a language" }},
{name: "WS-unspecified-language", protocol: TTSProtocolDashScopeTask, change: func(a map[string]any) { a["tts"].(map[string]any)["language_type"] = "und-US" }},
{name: "ASR-interim", protocol: TTSProtocolDashScopeTask, change: func(a map[string]any) { a["asr"].(map[string]any)["interim"] = false }},
{name: "missing-ASR-model", protocol: TTSProtocolDashScopeTask, change: func(a map[string]any) { delete(a["asr"].(map[string]any), "model") }},
{name: "ASR-rate", protocol: TTSProtocolDashScopeTask, change: func(a map[string]any) { a["asr"].(map[string]any)["input"].(map[string]any)["sample_rate_hz"] = 48000 }},
{name: "ASR-language", protocol: TTSProtocolDashScopeTask, change: func(a map[string]any) { a["asr"].(map[string]any)["language"] = "bad language" }},
{name: "unknown-parameter", protocol: TTSProtocolDashScopeTask, change: func(a map[string]any) { a["tts"].(map[string]any)["extra_parameter"] = true }},
} {
t.Run(tc.name, func(t *testing.T) {
task, providers := protocolTask(t, tc.protocol)
if tc.change != nil {
task = changeCurrentAgent(t, task, tc.change)
}
if tc.connection != nil {
tc.connection(providers)
}
if _, err := Bind(task, providers); err == nil {
t.Fatal("unsupported request silently altered or accepted")
}
})
}
}
func TestKnownModelNamesDoNotSelectTheProtocol(t *testing.T) {
for _, tc := range []struct{ protocol, model string }{{TTSProtocolDashScopeHTTP, "cosyvoice-v3-flash"}, {TTSProtocolDashScopeTask, "qwen3-tts-flash"}} {
t.Run(tc.protocol, func(t *testing.T) {
task, providers := protocolTask(t, tc.protocol)
task = changeCurrentAgent(t, task, func(a map[string]any) { a["tts"].(map[string]any)["model"] = tc.model })
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
if bound.TTS.Protocol != tc.protocol || bound.TTS.Model != tc.model {
t.Fatal("explicit selection replaced by a model-name mapping")
}
})
}
}
func TestLanguageHintsKeepExplicitLanguageAndNeverGuessFromRegion(t *testing.T) {
for _, tc := range []struct{ value, want string }{{"Chinese", "zh"}, {"English", "en"}, {"Japanese", "ja"}, {"Korean", "ko"}, {"German", "de"}, {"French", "fr"}, {"Russian", "ru"}, {"Italian", "it"}, {"Spanish", "es"}, {"Portuguese", "pt"}, {"en-US", "en"}, {"zh-CN", "zh"}, {"fr-FR", "fr"}} {
got, err := bailianLanguageHint(tc.value)
if err != nil || got != tc.want {
t.Fatalf("hint %s = %s %v", tc.value, got, err)
}
}
for _, value := range []string{"", "und", "und-US", "not a language"} {
if _, err := bailianLanguageHint(value); err == nil {
t.Fatalf("language inferred from %q", value)
}
}
hints, err := bailianTTSLanguageHints("Auto")
if err != nil || hints != nil {
t.Fatal("explicit Auto must omit the hint")
}
}
+79
View File
@@ -0,0 +1,79 @@
package ai
import (
"git.ipao.vip/rogee/go-sip/internal/configread"
"testing"
)
func protocolTask(t *testing.T, protocol string) (configread.Task, map[string]configread.Provider) {
t.Helper()
task, providers := currentFixture(t, "full_ai")
task = changeCurrentAgent(t, task, func(a map[string]any) {
asr := a["asr"].(map[string]any)
delete(asr, "interim")
asr["model"] = "unlisted-asr-model-2029"
asr["language"] = "en-US"
asr["input"].(map[string]any)["sample_rate_hz"] = 16000
llm := a["llm"].(map[string]any)
llm["model"] = "unlisted-chat-model-2029"
tts := a["tts"].(map[string]any)
tts["protocol"] = protocol
tts["model"] = "unlisted-tts-model-2029"
tts["voice"] = "unlisted-voice"
tts["language_type"] = "English"
tts["speed"] = 1.0
if protocol == "dashscope_task_websocket" {
tts["speed"] = 1.75
}
})
asr := providers["asr-example"]
asr.Code = "ali_bailian"
asr.WSEndpoint = "wss://speech.example.invalid/inference"
providers[asr.ProviderRef] = asr
tts := providers["tts-example"]
tts.WSEndpoint = asr.WSEndpoint
providers[tts.ProviderRef] = tts
return task, providers
}
func TestModelsAndVoicesAreDataWithinExplicitTTSProtocol(t *testing.T) {
for _, protocol := range []string{"dashscope_tts_http", "dashscope_task_websocket"} {
t.Run(protocol, func(t *testing.T) {
task, providers := protocolTask(t, protocol)
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
if bound.ASR.Model != "unlisted-asr-model-2029" || bound.ASR.Request.SampleRate != 16000 || bound.ASR.Language != "en-US" || bound.LLM.Model != "unlisted-chat-model-2029" || bound.TTS.Model != "unlisted-tts-model-2029" || bound.TTS.Voice != "unlisted-voice" || bound.TTS.LanguageType != "English" {
t.Fatal("task-selected model/voice/language changed")
}
expected := providers["tts-example"].Endpoint
if protocol == "dashscope_task_websocket" {
expected = providers["tts-example"].WSEndpoint
}
if bound.TTS.Provider.Endpoint != expected {
t.Fatal("connection selected from model name instead of explicit protocol")
}
})
}
}
func TestMissingOrUnsupportedProtocolNeverFallsBack(t *testing.T) {
for _, protocol := range []string{"", "unknown-protocol"} {
t.Run(protocol, func(t *testing.T) {
task, providers := protocolTask(t, "dashscope_task_websocket")
task = changeCurrentAgent(t, task, func(a map[string]any) {
tts := a["tts"].(map[string]any)
if protocol == "" {
delete(tts, "protocol")
} else {
tts["protocol"] = protocol
}
tts["model"] = "cosyvoice-v3-flash"
})
if _, err := Bind(task, providers); err == nil {
t.Fatal("protocol guessed from a familiar model")
}
})
}
}
+212
View File
@@ -0,0 +1,212 @@
package ai
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"time"
"git.ipao.vip/rogee/go-sip/internal/media"
doubaospeech "github.com/GizClaw/doubao-speech-go"
"github.com/coder/websocket"
)
func TestTaskTTSProtocolForwardsOpaqueModelVoiceSpeedAndLanguage(t *testing.T) {
for _, language := range []string{"English", "fr-FR", "Auto"} {
t.Run(language, func(t *testing.T) {
task, providers := protocolTask(t, TTSProtocolDashScopeTask)
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
bound.TTS.LanguageType = language
var calls atomic.Int32
expected := []byte{1, 0, 2, 0}
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
calls.Add(1)
c, err := websocket.Accept(w, r, nil)
if err != nil {
t.Error(err)
return
}
defer c.CloseNow()
ctx, cancel := context.WithTimeout(r.Context(), 2*time.Second)
defer cancel()
cmd := mockRead(t, c, ctx)
p := cmd.Payload["parameters"].(map[string]any)
if cmd.Payload["model"] != "unlisted-tts-model-2029" || p["voice"] != "unlisted-voice" || p["rate"] != 1.75 || p["sample_rate"] != float64(16000) || p["format"] != "pcm" {
t.Errorf("approved parameters changed: %+v", cmd.Payload)
}
if language == "Auto" {
if _, ok := p["language_hints"]; ok {
t.Error("explicit Auto must not insert a language hint")
}
} else {
expectedHint := "en"
if language == "fr-FR" {
expectedHint = "fr"
}
hints, ok := p["language_hints"].([]any)
if !ok || len(hints) != 1 || hints[0] != expectedHint {
t.Error("language was replaced with a model-specific default")
}
}
mockEvent(c, ctx, cmd.Header.TaskID, "task-started", nil)
text := mockRead(t, c, ctx)
finish := mockRead(t, c, ctx)
if text.Payload["input"].(map[string]any)["text"] != "Protocol fixture." || finish.Header.Action != "finish-task" {
t.Error("TTS text/lifecycle changed")
}
c.Write(ctx, websocket.MessageBinary, expected)
mockEvent(c, ctx, cmd.Header.TaskID, "task-finished", nil)
}))
defer server.Close()
bound.TTS.Provider.WSEndpoint = "ws" + strings.TrimPrefix(server.URL, "http")
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
audio, err := bound.Synthesize(ctx, "Protocol fixture.")
if err != nil || !bytes.Equal(audio, expected) || calls.Load() != 1 {
t.Fatalf("task TTS: audio=%d err=%v requests=%d", len(audio), err, calls.Load())
}
})
}
}
func TestTaskASRProtocolForwardsOpaqueModelAndSelectedSampleRate(t *testing.T) {
for _, rate := range []int{8000, 16000} {
t.Run(fmt.Sprint(rate), func(t *testing.T) {
task, providers := protocolTask(t, TTSProtocolDashScopeTask)
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
bound.ASR.Request.SampleRate = doubaospeech.SampleRate(rate)
pcm := make([]byte, 32000)
for i := range pcm {
pcm[i] = byte(i % 251)
}
var calls atomic.Int32
captured := make(chan []byte, 1)
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var received []byte
defer func() { captured <- received }()
calls.Add(1)
c, err := websocket.Accept(w, r, nil)
if err != nil {
t.Error(err)
return
}
defer c.CloseNow()
ctx, cancel := context.WithTimeout(r.Context(), 2*time.Second)
defer cancel()
cmd := mockRead(t, c, ctx)
p := cmd.Payload["parameters"].(map[string]any)
if cmd.Payload["model"] != "unlisted-asr-model-2029" || p["sample_rate"] != float64(rate) || p["language_hints"].([]any)[0] != "en" {
t.Error("ASR task model/rate/language changed")
}
mockEvent(c, ctx, cmd.Header.TaskID, "task-started", nil)
for {
kind, raw, err := c.Read(ctx)
if err != nil {
t.Error(err)
return
}
if kind == websocket.MessageBinary {
received = append(received, raw...)
continue
}
var finish mockCommand
json.Unmarshal(raw, &finish)
if finish.Header.Action != "finish-task" {
t.Error("ASR not finished explicitly")
}
break
}
mockEvent(c, ctx, cmd.Header.TaskID, "result-generated", map[string]any{"output": map[string]any{"sentence": map[string]any{"text": "Final user text.", "sentence_end": true, "begin_time": 0, "end_time": 1000}}})
mockEvent(c, ctx, cmd.Header.TaskID, "task-finished", nil)
}))
defer server.Close()
bound.ASR.Provider.Endpoint = "ws" + strings.TrimPrefix(server.URL, "http")
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
text, err := bound.Recognize(ctx, pcm)
received := <-captured
if err != nil || text != "Final user text." || len(received) != rate*2 || calls.Load() != 1 {
t.Fatalf("ASR: text=%q err=%v sent=%d calls=%d", text, err, len(received), calls.Load())
}
if rate == 16000 && !bytes.Equal(received, pcm) {
t.Fatal("16k PCM was transformed despite the selected 16k protocol rate")
}
})
}
}
func TestHTTPProtocolForwardsOpaqueModelVoiceAndLanguage(t *testing.T) {
wav, _, err := media.EncodeMonoWAV([]byte{0, 0, 2, 0, 3, 0, 4, 0}, 1024)
if err != nil {
t.Fatal(err)
}
var generation, downloads atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/audio" {
downloads.Add(1)
w.Header().Set("Content-Type", "audio/wav")
w.Write(wav)
return
}
generation.Add(1)
var req struct {
Model string `json:"model"`
Input struct {
Text, Voice string
LanguageType string `json:"language_type"`
} `json:"input"`
}
if json.NewDecoder(r.Body).Decode(&req) != nil || req.Model != "unlisted-tts-model-2029" || req.Input.Voice != "unlisted-voice" || req.Input.LanguageType != "English" || req.Input.Text != "Protocol fixture." {
t.Error("HTTP task parameters changed")
}
w.Header().Set("Content-Type", "application/json")
fmt.Fprintf(w, `{"output":{"audio":{"url":%q}}}`, "http://"+r.Host+"/audio")
}))
defer server.Close()
task, providers := protocolTask(t, TTSProtocolDashScopeHTTP)
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
bound.TTS.Provider.Endpoint = server.URL + "/api/v1/services/aigc/multimodal-generation/generation"
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
audio, err := bound.Synthesize(ctx, "Protocol fixture.")
if err != nil || len(audio) == 0 || generation.Load() != 1 || downloads.Load() != 1 {
t.Fatalf("HTTP TTS: audio=%d err=%v requests=%d downloads=%d", len(audio), err, generation.Load(), downloads.Load())
}
}
func TestExplicitHTTPProtocolNeverGuessesWebSocketFromFamiliarModel(t *testing.T) {
var httpCalls, wsCalls atomic.Int32
ws := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { wsCalls.Add(1); w.WriteHeader(500) }))
defer ws.Close()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { httpCalls.Add(1); w.WriteHeader(400) }))
defer server.Close()
task, providers := protocolTask(t, TTSProtocolDashScopeHTTP)
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
bound.TTS.Model = "cosyvoice-v3-flash"
bound.TTS.Provider.Endpoint = server.URL + "/api/v1/services/aigc/multimodal-generation/generation"
bound.TTS.Provider.WSEndpoint = "ws" + strings.TrimPrefix(ws.URL, "http")
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
audio, err := bound.Synthesize(ctx, "No fallback.")
if err == nil || len(audio) != 0 || httpCalls.Load() != 1 || wsCalls.Load() != 0 {
t.Fatalf("model/protocol fallback: audio=%d err=%v HTTP=%d WS=%d", len(audio), err, httpCalls.Load(), wsCalls.Load())
}
}
+106
View File
@@ -0,0 +1,106 @@
package ai
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/coder/websocket"
)
func TestProtocolModelsCompleteTurnWithoutPerModelCode(t *testing.T) {
task, providers := protocolTask(t, TTSProtocolDashScopeTask)
task = changeCurrentAgent(t, task, func(a map[string]any) {
a["conversation"].(map[string]any)["opening"] = "Hello."
a["conversation"].(map[string]any)["hangup_keywords"] = []any{}
})
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
var asrCalls, ttsCalls, llmCalls atomic.Int32
expected := []byte{1, 0, 2, 0}
speech := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
c, err := websocket.Accept(w, r, nil)
if err != nil {
t.Error(err)
return
}
defer c.CloseNow()
ctx, cancel := context.WithTimeout(r.Context(), 2*time.Second)
defer cancel()
cmd := mockRead(t, c, ctx)
mockEvent(c, ctx, cmd.Header.TaskID, "task-started", nil)
switch cmd.Payload["task"] {
case "asr":
asrCalls.Add(1)
if cmd.Payload["model"] != "unlisted-asr-model-2029" {
t.Error("ASR model replaced")
}
for {
kind, _, err := c.Read(ctx)
if err != nil {
t.Error(err)
return
}
if kind == websocket.MessageText {
break
}
}
mockEvent(c, ctx, cmd.Header.TaskID, "result-generated", map[string]any{"output": map[string]any{"sentence": map[string]any{"text": "Tell me more.", "sentence_end": true, "begin_time": 0, "end_time": 100}}})
case "tts":
ttsCalls.Add(1)
if cmd.Payload["model"] != "unlisted-tts-model-2029" || cmd.Payload["parameters"].(map[string]any)["voice"] != "unlisted-voice" {
t.Error("TTS model/voice replaced")
}
mockRead(t, c, ctx)
finish := mockRead(t, c, ctx)
if finish.Header.Action != "finish-task" {
t.Error("TTS lifecycle changed")
}
c.Write(ctx, websocket.MessageBinary, expected)
default:
t.Error("unknown role")
return
}
mockEvent(c, ctx, cmd.Header.TaskID, "task-finished", nil)
}))
defer speech.Close()
llm := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
llmCalls.Add(1)
var request map[string]any
if r.Method != "POST" || r.URL.Path != "/v1/chat/completions" || json.NewDecoder(r.Body).Decode(&request) != nil || request["model"] != "unlisted-chat-model-2029" {
t.Error("LLM model/protocol changed")
}
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"id":"mock","object":"chat.completion","created":0,"model":"unlisted-chat-model-2029","choices":[{"index":0,"message":{"role":"assistant","content":"Generic response."},"finish_reason":"stop"}]}`)
}))
defer llm.Close()
bound.ASR.Provider.Endpoint = "ws" + strings.TrimPrefix(speech.URL, "http")
bound.TTS.Provider.WSEndpoint = bound.ASR.Provider.Endpoint
bound.LLM.Provider.Endpoint = llm.URL + "/v1"
call, err := NewCall(bound, func(context.Context) error { t.Error("ordinary turn attempted hangup"); return nil })
if err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
opening, err := call.Open(ctx)
if err != nil || !bytes.Equal(opening, expected) {
t.Fatalf("opening: %v", err)
}
result, err := call.RunTurn(ctx, make([]byte, 3200))
if err != nil || result.Transcript != "Tell me more." || result.Reply != "Generic response." || !bytes.Equal(result.AudioPCM16, expected) || result.EndedByKeyword {
t.Fatalf("generic model turn: %+v %v", result, err)
}
if asrCalls.Load() != 1 || llmCalls.Load() != 1 || ttsCalls.Load() != 2 {
t.Fatalf("unexpected/retried requests: ASR=%d LLM=%d TTS=%d", asrCalls.Load(), llmCalls.Load(), ttsCalls.Load())
}
}
+40
View File
@@ -0,0 +1,40 @@
package ai
import (
"errors"
"golang.org/x/text/language"
)
const (
TTSProtocolDashScopeHTTP = "dashscope_tts_http"
TTSProtocolDashScopeTask = "dashscope_task_websocket"
)
// Task language labels/locales express the language; DashScope task parameters
// use ISO language hints. Do not infer an unspecified language from a region.
func bailianLanguageHint(value string) (string, error) {
names := map[string]string{"Chinese": "zh", "English": "en", "Japanese": "ja", "Korean": "ko", "German": "de", "French": "fr", "Russian": "ru", "Italian": "it", "Spanish": "es", "Portuguese": "pt"}
if code, ok := names[value]; ok {
return code, nil
}
tag, err := language.Parse(value)
if err != nil {
return "", errors.New("language cannot be expressed by the selected task protocol")
}
base, _, _ := tag.Raw()
if base.String() == "und" {
return "", errors.New("task language must be explicit")
}
return base.String(), nil
}
func bailianTTSLanguageHints(value string) ([]string, error) {
if value == "Auto" {
return nil, nil
} // Explicit user selection, not a missing-value default.
hint, err := bailianLanguageHint(value)
if err != nil {
return nil, err
}
return []string{hint}, nil
}
@@ -86,11 +86,32 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) {
}
defer func(q, key, exchange string) { _ = admin.QueueUnbind(q, key, exchange, nil) }(queue, route.BindingKey, route.Exchange)
}
// Socket write completion is not broker acceptance. Establish the same
// mandatory/confirm boundary required of the simulated SaaS publisher.
if err := admin.Confirm(false); err != nil {
t.Fatal(err)
}
returned := admin.NotifyReturn(make(chan amqp.Return, 1))
publish := func(route tenant.Route, body []byte) {
t.Helper()
if err := admin.PublishWithContext(context.Background(), route.Exchange, route.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
confirmation, err := admin.PublishWithDeferredConfirmWithContext(ctx, route.Exchange, route.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body})
if err != nil {
t.Fatal(err)
}
if confirmation == nil {
t.Fatal("isolated SaaS publisher did not enter confirm mode")
}
accepted, err := confirmation.WaitContext(ctx)
if err != nil || !accepted {
t.Fatalf("isolated SaaS message not confirmed: accepted=%v error=%v", accepted, err)
}
select {
case msg := <-returned:
t.Fatalf("isolated SaaS message returned: code=%d", msg.ReplyCode)
default:
}
}
publish(controlRoute, controlBody(t, "control-example", "pause", "drain"))
publish(taskRoute, executeBody(t, "call-example", "15803300952"))
@@ -366,11 +387,35 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) {
t.Fatalf("pending SIP change originated call %s", spec.EventID)
case <-time.After(120 * time.Millisecond):
}
for _, eventID := range []string{"stop-backlog-1", "stop-backlog-2"} {
// The preceding instruction probes the durable admission fence and may
// legitimately have reached the inbox before asynchronous consumer stop.
// Do not count it as ready-queue backlog. Establish consumer quiescence
// before publishing an independent, confirmed three-message purge fixture.
deadline = time.After(5 * time.Second)
for {
state, err := admin.QueueInspect(taskRoute.Queue)
if err != nil {
t.Fatal(err)
}
if state.Consumers == 0 {
break
}
select {
case <-deadline:
t.Fatalf("consumer did not stop behind SIP admission fence: %+v", state)
case <-time.After(20 * time.Millisecond):
}
}
for _, eventID := range []string{"stop-backlog-0", "stop-backlog-1", "stop-backlog-2"} {
publish(taskRoute, executeBody(t, eventID, "15003164745"))
}
if state, err := admin.QueueInspect(taskRoute.Queue); err != nil || state.Messages < 3 {
t.Fatalf("expected isolated stop backlog: %+v %v", state, err)
pendingCommands, pendingErr := db.ListPendingExecute(id)
var pendingIDs []string
for _, command := range pendingCommands {
pendingIDs = append(pendingIDs, command.EventID)
}
t.Fatalf("expected isolated stop backlog: queue=%+v inspect_error=%v pending_event_ids=%v pending_error=%v", state, err, pendingIDs, pendingErr)
}
publish(controlRoute, controlBody(t, "stop-after-sip", "stop", ""))
select {