feat(ai): bind approved Bailian TTS and convert returned audio

This commit is contained in:
2026-10-04 10:17:20 +08:00
parent 5ab4cbb3be
commit 28a2cdd495
17 changed files with 406 additions and 172 deletions
+1 -1
View File
@@ -90,7 +90,7 @@
- `call.execute` 只带获批 `task_id/callee`,调用线路、主叫、AI 和时限由该任务快照固定;`task.control` 的 start/pause/resume/stop 无旧 CAS/版本字段,控制回 task/action/status 处理结果,stop 不可恢复。任务启动/重启完整读取列表,运行中不定时轮询任务列表;新建任务由 start 读取任务配置和租户额度,暂停后修改的任务由 resume 重读;全局 AI 服务商列表每次启动只读取一次,本进程全部任务复用,变更须重启后才生效;没有 edit 事件。启动时全量读取并核验 SIP,运行中只由 `sip.config` 通知触发全量 SIP 读取,无常规定时 SIP 轮询;任务 start/resume 使用已生效 SIP 快照,不自行拉取或向 Agent 核验 SIP。SIP 修订变化待旧呼叫结束且 Agent 已加载后,仅在本地重新绑定任务快照,不重读任务列表;待处理期间新准入关闭。发现分页同快照先全部校验后提交,MQ 控制持久状态高于偶发 HTTP 状态;冲突关准入、resume 须重新核验。历史命令不因消息年龄过期,但实际呼出前仍检查任务/白名单/时段/授权/额度和有效通话上限;任务结束不删未确认结果、恢复、幂等或未知占用。
- 独立 Dispatcher 的 SQLite 是任务、额度、inbox/outbox 的权威数据;Agent 无业务数据库,录音、执行与上传恢复只写受控私有文件。额度包含未知占用,新 boot/租约到期不得自动清除未知执行;不实现双活数据库、自动跨机热备、多 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 配置的唯一编辑/审批面仍是 management。
- AI 使用任务内不可变授权快照:仅经获批准百炼/火山 ASR、OpenAI 兼容 LLM、火山 TTS 能表达的参数进入每通话实例;ASR-only 不启动 LLM/TTS,完整 AI 不借旧语音测试的 LLM/TTS。只有最终用户 ASR 文本的明确字面关键词可触发拒联/挂断;不由 SDK 默认值、环境、CLI、metadata 或宽松 Schema 改写业务参数,不因 SDK 重试产生第二次发起/收费或重播。日志只存脱敏版本/摘要/计数,不存密钥、prompt、完整对话或音频。
- AI 使用任务内不可变授权快照:仅经获批准百炼/火山 ASR、OpenAI 兼容 LLM、百炼 TTS(`qwen3-tts-flash`/`Cherry`/`Chinese`)能表达的参数进入每通话实例;ASR-only 不启动 LLM/TTS,完整 AI 不借旧语音测试的授权或参数。只有最终用户 ASR 文本的明确字面关键词可触发拒联/挂断;不由 SDK 默认值、环境、CLI、metadata 或宽松 Schema 改写业务参数,不因 SDK 重试产生第二次发起/收费或重播。日志只存脱敏版本/摘要/计数,不存密钥、prompt、完整对话或音频。
- **本项目既有私有配置的取用入口(仅在相应真实服务获得明确授权时使用,不是运行时回退)**:工作区 `.local/provider-ai.env` 为权限 `0600` 的历史百炼/火山凭据文件,按明确字段读取 `BAILIAN_API_KEY`、`BAILIAN_BASE_URL` 等所需项;仓库根 `aliyun-oss.env` 为权限 `0600` 的冒号分隔文件,按 `Bucket`、`Endpoint`、`Region`、`AccessKeyId`、`AccessKeySecret` 字段解析。一次性外部诊断仅以进程内存向现有 `AGENT_CALL_REAL_AI_INTEGRATION=1`、`AGENT_CALL_OSS_INTEGRATION=1` 测试注入获准字段;常规构建/测试不读这些文件,绝不 `source` 不可信内容、输出密钥/签名 URL、提交文件或覆盖/清理旧 OSS 对象。现行业务 AI 参数仍须来自 SaaS 批准的 task/providers 完整快照,OSS 须由 Dispatcher 私有 `DISPATCHER_OSS_CONFIG_FILE` 与获批上传授权提供;上述历史文件**不能**直接当 SaaS 快照、任务授权或真实通话准入。旧 `.local/asterisk-*/ari.conf` 仅为历史本机文件,不证明重装后主机已有 ARI 配置。
- Agent 录音经受控双向 TLS 向 D 领取短期 OSS 上传授权,每次尝试只作**一次 HTTPS PUT**;正常上传成功不写录音/结果业务文件。首次明确失败须先完整保存录音与结果两份恢复文件,才从该时刻启动 48 小时重试;按 1、2、4、8、16、32、60 分钟及其后每 60 分钟的固定节奏显式重新申请授权,同一 OSS 目标、同一消息身份。PUT 结果未知不得盲目重传;48 小时届满仍失败时保留文件待人工,**不伪造最终结果或自动清理**。D 不转发文件,已确认结束的通话及时释放执行占用;未知执行仍占用。只有真实终结后才通过唯一 `call.execute.result` 回报录音路径、最终转写和拒联事实;无录音或录音生成失败以空 `recording={}` 和真实结果收口,生成失败须说明原因。不能恢复的录音不声称零丢失,也不伪造 OSS/SaaS 应用回执。凭据/TOKEN/签名 URL 不写入样例、日志、源码或证据。
+2 -2
View File
@@ -172,9 +172,9 @@
},
"tts": {
"type": "object", "additionalProperties": false,
"required": ["provider_ref", "model", "voice", "format"],
"required": ["provider_ref", "model", "voice", "language_type", "format"],
"properties": {
"provider_ref": {"type": "string", "minLength": 1}, "model": {"type": "string", "minLength": 1}, "voice": {"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}}}
"provider_ref": {"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}}}
}
},
"prompt": {
@@ -4,6 +4,6 @@
"providers":[
{"provider_ref":"asr-example","role":"asr","enabled":true,"adapter":"volcengine_asr","endpoint":"https://asr.example.invalid","credential":"example-only-not-a-real-secret"},
{"provider_ref":"llm-example","role":"llm","enabled":true,"adapter":"openai_compatible","endpoint":"https://llm.example.invalid/v1","credential":"example-only-not-a-real-secret"},
{"provider_ref":"tts-example","role":"tts","enabled":true,"adapter":"volcengine_tts","endpoint":"https://tts.example.invalid","credential":"example-only-not-a-real-secret"}
{"provider_ref":"tts-example","role":"tts","enabled":true,"adapter":"bailian_tts","endpoint":"https://tts.example.invalid/api/v1/services/aigc/multimodal-generation/generation","credential":"example-only-not-a-real-secret"}
]
}
@@ -8,7 +8,7 @@
"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":"example-tts","voice":"example-neutral","speed":1,"format":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1},"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}
}
+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": "973ef69b1f71747d3e90f9798a7bf75d2fe2aac7a56d91dad97bbb8d4b9c566e"
"docs/thirds/saas-dispatcher.md": "dfa9d264d48d2c2641a06c54d4d6f077e7ecea90a7274b967b179428fbc35ddc"
},
"bundle_sha256": "d89302b79228c884ce5e0c413bbe9c688c725724de73cb84892e1d38d638e2a5",
"bundle_sha256": "130e8ba1b1d87f38265b4125451c01cb0ef0b4b90e108040891aa41d9f0ed473",
"bundle_algorithm": "sha256 of sorted relative-path + space + sha256(file) + newline; only root-level JSON and examples/**/*.json, excluding manifest.json"
}
@@ -63,7 +63,9 @@
## P05:Agent 快照与 SDK 隔离链路(项目内组件通过,主入口待 P07)
- Dispatcher 在启动、增量发现、恢复任务、SIP 更新及发指令前对完整 AI/provider 快照执行能力校验;不支持的任务关闭准入而不占额度或呼叫。已按用户确认保留现行 TTS Schema:火山 TTS V2 SDK 无法表达的 PCMA、0.5–2 倍以外或整数 speech_rate 无法精确表达的速度均明确拒绝,不静默修改配置;不支持的插话配置同样拒绝。
> 以下为 P05 当时的火山 TTS 历史验收记录,不再代表当前 TTS 实现。现行合同已改为任务明确授权的百炼 `qwen3-tts-flash`/`Cherry`/`Chinese`:本地 Mock 已核验生成请求、签名音频引用下载、16 kHz PCM16 转换、一次请求与失败不重试;未发起真实百炼 TTS 或真实通话,旧凭据不构成新任务授权。当前依据以 [`SaaS↔Dispatcher 规范`](../thirds/saas-dispatcher.md) 及合同为准。
- Dispatcher 在启动、增量发现、恢复任务、SIP 更新及发指令前对完整 AI/provider 快照执行能力校验;不支持的任务关闭准入而不占额度或呼叫。当时按已确认的 TTS Schema:火山 TTS V2 SDK 无法表达的 PCMA、0.5–2 倍以外或整数 speech_rate 无法精确表达的速度均明确拒绝,不静默修改配置;不支持的插话配置同样拒绝。
- `ExecuteApproved` 只在隔离 Mock 中签发:携带任务原始 JSON、引用的 provider 明文凭据、SIP revision、选定路由/主叫/原始被叫、独立的拨号期限及最大通话时限;任务+provider 原始字节+SIP revision 计算绑定摘要。Agent 对照已激活的 D 会话身份、绑定摘要、SDK 能力及**实际观察到的** SIP 加载结果;发出 Mock 指令前仅持久保存摘要和未知占用,失败/重复/重启不自动重拨,文件中不留明文凭据。当前 `LoadedSIP` 为可注入的隔离 Mock 观察器,**不等于真实 Asterisk 已加载核验**。
- ASR-only 不启动 LLM/TTS;full-AI 使用冻结的 model、显式 temperature=0/max_tokens、TTS voice/format/speech_rate 与 provider 原值凭据。真实 ASR SDK WebSocket 发出的音频输入/识别配置、LLM 与 TTS SDK 向隔离 HTTP Mock 发出的实际参数均已捕获核对;已配置的开场白由 TTS 单次合成,失败不自动重播且未完成时不进入对话轮次;空开场白不调用 TTS。关键词仅匹配用户侧最终 ASR,失败/未知挂断不自动重试。`ApprovedOriginator` 经生成的 Unary gRPC Stub 交付完整快照,SIP 全量 revision 不吻合即拒绝。
- 对话控制隔离验证:共享 `callflow` 不再根据转写或回复内的硬编码词推断拒联,仅处理明确的关键词/拒联事实;ASR-only 不播报开场或 TTS 回复。配置的首语音/静默期限、整通话期限与最大轮数交给媒体控制器;回复按 Unicode 字符数分片,火山 TTS 的缓存音频块数超限明确失败而不重试;不支持的插话配置在准入前拒绝。该控制器尚未接入新主 CLI,不能宣称真实通话媒体已验证。
+2 -1
View File
@@ -9,7 +9,7 @@
| 资源 | 方法和路径 | 应用规则 | 正例 |
| --- | --- | --- | --- |
| SIP 全量 | `GET /internal/v1/dispatcher/sip` | `revision` 是实际加载版本核对依据;未知 `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_ref` 只标识供应商;将明文 `credential` 原值交给选定角色的 SDK,不引入 ref 查找或交换;缺失、禁用、角色不匹配、无效凭据拒绝执行 | [`providers`](../../contracts/local/examples/config-read-providers.json) |
| AI provider 全量 | `GET /internal/v1/dispatcher/ai-providers` | 启动时一次读取全局列表,本进程全部任务共用到下一次重启;运行中 start/resume 不重复拉取,服务商变更须重启后才生效。`provider_ref` 只标识供应商;将明文 `credential` 原值交给选定角色的已核验适配器,不引入 ref 查找或交换;缺失、禁用、角色不匹配、无效凭据拒绝执行 | [`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>` | 首次/重启完整读到 **带 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,6 +33,7 @@ RabbitMQ 是 Topic,**SaaS 独占创建、绑定、退役 exchange/queue,D
## 调度、AI 与真实结果(K01–K09、K11–K14)
- 当前 TTS 唯一获批适配器是 `bailian_tts`:任务快照须明确提供 `qwen3-tts-flash`、`Cherry`、`Chinese`、速度 `1` 和单声道 16 kHz PCM16 目标格式;使用已核验的 provider 凭据及生成端点。每段仅发起一次生成请求,下载返回的短期音频引用后转换为电话可用的 PCM16;不可用、超时、缺少转换工具或参数不支持时显式失败,不回退旧火山 TTS、不隐式重试或记录签名音频 URL。历史测试凭据不是任务授权,本地转换 Mock 不构成真实百炼/通话验收。
- 白名单仅含 `15003164745`、`15830461047` 原值;已选 SIP trunk、任务与线路每周时段、任务排除日期、任务/租户/线路额度、任务与 AI 较小通话时限均在接纳及实际发呼叫指令前检查。线路字段未知则 fail-closed;选线后固定、不自动重拨/换线。隔离 Mock 中规则暂不满足时保留待执行指令、暂停该任务的调度,规则允许后重验;与人工 pause/stop 分离,不能自动解除人为停止。本规则**不**放宽真实路径 Asia/Shanghai `09:00`–`20:00` 固定门禁或授权真实拨号。
- 接通事实为真时 `outcome=answered`(后续异常不抹掉接通);已发起但忙线、拒接、无人接听且确定结束为 `no_answer`;确认未接通并由 Agent/Asterisk 执行故障终结为 `failed`;未知状态保持未知占用,不能伪造结束、结果或自动重拨。真实 SIP 状态码原样数字写入 `reason_code`,无真实 SIP 码则 `null` 并以 `reason_message` 说明;禁止本地虚构数字错误码。无应答且没有录音时 `transcript=[]`、`opt_out=false`、`recording={}`。
- 只有**用户侧 ASR 最终识别文本**包含任一 `hangup_keywords` 字面字符串才挂断;中间识别、助手回复、开场白、TTS 均不能触发;重复结果不可反复终结。同一任务 revision 不同内容拒绝准入;provider 禁用/角色不符不可调用。Mock 参数验证不等于真实供应商验收。
+136
View File
@@ -0,0 +1,136 @@
package ai
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net"
"net/http"
"net/url"
"os/exec"
"strings"
)
const (
maxBailianReplyBytes = 64 << 10
maxBailianAudioBytes = 8 << 20
maxBailianPCMBytes = 2 << 20 // ponytail: enough for about one minute; raise only with an approved longer sentence.
)
func synthesizeBailianTTS(ctx context.Context, approved TTSConfig, text string) ([]byte, error) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
return nil, errors.New("Bailian TTS requires the configured ffmpeg audio converter")
}
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 == "" {
return nil, errors.New("approved Bailian TTS settings or text are incomplete")
}
// Never follow a provider or audio redirect into an unapproved second request.
client := http.Client{CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }}
request := struct {
Model string `json:"model"`
Input struct {
Text string `json:"text"`
Voice string `json:"voice"`
LanguageType string `json:"language_type"`
} `json:"input"`
}{Model: approved.Model}
request.Input.Text, request.Input.Voice, request.Input.LanguageType = text, approved.Voice, approved.LanguageType
body, err := json.Marshal(request)
if err != nil {
return nil, errors.New("encode approved Bailian request")
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, approved.Provider.Endpoint, bytes.NewReader(body))
if err != nil {
return nil, errors.New("create approved Bailian request")
}
req.Header.Set("Authorization", "Bearer "+approved.Provider.Credential)
req.Header.Set("Content-Type", "application/json")
response, err := client.Do(req)
if err != nil {
return nil, fmt.Errorf("Bailian TTS request failed: type=%T", err)
}
defer response.Body.Close()
if response.StatusCode/100 != 2 {
return nil, fmt.Errorf("Bailian TTS rejected request: HTTP %d", response.StatusCode)
}
raw, err := io.ReadAll(io.LimitReader(response.Body, maxBailianReplyBytes+1))
if err != nil || len(raw) > maxBailianReplyBytes {
return nil, errors.New("Bailian TTS response is unreadable or oversized")
}
var result struct {
Output struct {
Audio struct {
URL string `json:"url"`
} `json:"audio"`
} `json:"output"`
}
if err := json.Unmarshal(raw, &result); err != nil || result.Output.Audio.URL == "" {
return nil, errors.New("Bailian TTS response has no complete audio reference")
}
if err := bailianURL(result.Output.Audio.URL, true); err != nil {
return nil, errors.New("Bailian TTS audio reference is invalid")
}
audioRequest, err := http.NewRequestWithContext(ctx, http.MethodGet, result.Output.Audio.URL, nil)
if err != nil {
return nil, errors.New("create Bailian audio request")
}
audioResponse, err := client.Do(audioRequest)
if err != nil {
return nil, fmt.Errorf("Bailian audio download failed: type=%T", err)
}
defer audioResponse.Body.Close()
if audioResponse.StatusCode/100 != 2 {
return nil, fmt.Errorf("Bailian audio download rejected: HTTP %d", audioResponse.StatusCode)
}
encoded, err := io.ReadAll(io.LimitReader(audioResponse.Body, maxBailianAudioBytes+1))
if err != nil || len(encoded) == 0 || len(encoded) > maxBailianAudioBytes {
return nil, errors.New("Bailian audio is empty, unreadable or oversized")
}
cmd := exec.CommandContext(ctx, "ffmpeg", "-hide_banner", "-loglevel", "error", "-nostdin", "-i", "pipe:0", "-f", "s16le", "-ac", "1", "-ar", "16000", "pipe:1")
cmd.Stdin = bytes.NewReader(encoded)
cmd.Stderr = io.Discard
var pcm boundedPCM
cmd.Stdout = &pcm
if err := cmd.Run(); err != nil {
return nil, fmt.Errorf("Bailian audio conversion failed: type=%T", err)
}
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("Bailian TTS deadline or cancellation: %w", err)
}
if pcm.Len() == 0 || pcm.Len()%2 != 0 {
return nil, errors.New("Bailian TTS returned incomplete PCM16 audio")
}
return pcm.Bytes(), nil
}
type boundedPCM struct{ bytes.Buffer }
func (w *boundedPCM) Write(data []byte) (int, error) {
if len(data) > maxBailianPCMBytes-w.Len() {
return 0, errors.New("Bailian TTS PCM exceeds approved media ceiling")
}
return w.Buffer.Write(data)
}
func bailianURL(raw string, signedAudio bool) error {
parsed, err := url.Parse(raw)
if err != nil || parsed.Host == "" || parsed.User != nil || parsed.Fragment != "" {
return errors.New("HTTPS URL without embedded identity required")
}
if parsed.Scheme != "https" {
host := parsed.Hostname()
if parsed.Scheme != "http" || (host != "localhost" && !net.ParseIP(host).IsLoopback()) {
return errors.New("HTTPS URL required outside isolated local Mock")
}
}
if !signedAudio && (parsed.RawQuery != "" || !strings.HasSuffix(parsed.Path, "/api/v1/services/aigc/multimodal-generation/generation")) {
return errors.New("approved Bailian generation endpoint required")
}
return nil
}
+109
View File
@@ -0,0 +1,109 @@
package ai
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"os/exec"
"strings"
"sync/atomic"
"testing"
"git.ipao.vip/rogee/go-sip/internal/media"
)
func mockBailianTTS(t *testing.T, pcm []byte, observe func(string)) (string, *atomic.Int32) {
t.Helper()
wav, _, err := media.EncodeMonoWAV(pcm, 1024)
if err != nil {
t.Fatal(err)
}
calls := &atomic.Int32{}
var server *httptest.Server
server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/audio" && r.Method == http.MethodGet {
_, _ = w.Write(wav)
return
}
if r.URL.Path != "/api/v1/services/aigc/multimodal-generation/generation" || r.Method != http.MethodPost {
http.Error(w, "unexpected endpoint", http.StatusNotFound)
return
}
calls.Add(1)
var request struct {
Input struct {
Text string `json:"text"`
} `json:"input"`
}
if err := json.NewDecoder(r.Body).Decode(&request); err != nil {
http.Error(w, "invalid request", http.StatusBadRequest)
return
}
if observe != nil {
observe(request.Input.Text)
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]any{"output": map[string]any{"audio": map[string]any{"url": server.URL + "/audio?sig=fixture"}}})
}))
t.Cleanup(server.Close)
return server.URL + "/api/v1/services/aigc/multimodal-generation/generation", calls
}
func TestBailianTTSDownloadFailureDoesNotLeakSignedURLOrRetry(t *testing.T) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
t.Skip("Bailian TTS conversion requires ffmpeg")
}
var calls, downloads atomic.Int32
var server *httptest.Server
server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/audio" {
downloads.Add(1)
http.Error(w, "private-data", http.StatusForbidden)
return
}
calls.Add(1)
_ = json.NewEncoder(w).Encode(map[string]any{"output": map[string]any{"audio": map[string]any{"url": server.URL + "/audio?sig=private-signature"}}})
}))
defer server.Close()
task, providers := currentFixture(t, "full_ai")
p := providers["tts-example"]
p.Endpoint = server.URL + "/api/v1/services/aigc/multimodal-generation/generation"
providers[p.ProviderRef] = p
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
_, err = bound.Synthesize(context.Background(), "批准的回复")
if err == nil || !strings.Contains(err.Error(), "HTTP 403") || strings.Contains(err.Error(), "private-signature") || strings.Contains(err.Error(), "private-data") || strings.Contains(err.Error(), p.Credential) || calls.Load() != 1 || downloads.Load() != 1 {
t.Fatalf("download failure must be explicit without leaking/retrying: calls=%d downloads=%d err=%v", calls.Load(), downloads.Load(), err)
}
}
func TestBailianTTSDoesNotFollowChargeableRedirect(t *testing.T) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
t.Skip("Bailian TTS conversion requires ffmpeg")
}
var calls, redirected atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/charged-again" {
redirected.Add(1)
return
}
calls.Add(1)
http.Redirect(w, r, "/charged-again", http.StatusTemporaryRedirect)
}))
defer server.Close()
task, providers := currentFixture(t, "full_ai")
p := providers["tts-example"]
p.Endpoint = server.URL + "/api/v1/services/aigc/multimodal-generation/generation"
providers[p.ProviderRef] = p
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
_, err = bound.Synthesize(context.Background(), "批准的回复")
if err == nil || !strings.Contains(err.Error(), "HTTP 307") || calls.Load() != 1 || redirected.Load() != 0 {
t.Fatalf("TTS redirect cannot create a second provider request: calls=%d redirected=%d err=%v", calls.Load(), redirected.Load(), err)
}
}
+23 -24
View File
@@ -44,9 +44,13 @@ type LLMConfig struct {
}
type TTSConfig struct {
Provider configread.Provider
Request doubaospeech.TTSV2Request
Timeout time.Duration
Provider configread.Provider
Model string
Voice string
LanguageType string
Speed float64
SampleRate int
Timeout time.Duration
}
// ConversationConfig preserves the task's explicit dialogue limits. An
@@ -83,12 +87,13 @@ type currentAgentSettings struct {
TimeoutMS *int64 `json:"timeout_ms"`
} `json:"llm"`
TTS *struct {
ProviderRef string `json:"provider_ref"`
Model string `json:"model"`
Voice string `json:"voice"`
Speed *float64 `json:"speed"`
TimeoutMS *int64 `json:"timeout_ms"`
Format struct {
ProviderRef string `json:"provider_ref"`
Model string `json:"model"`
Voice string `json:"voice"`
LanguageType string `json:"language_type"`
Speed *float64 `json:"speed"`
TimeoutMS *int64 `json:"timeout_ms"`
Format struct {
Encoding string `json:"encoding"`
SampleRateHz int `json:"sample_rate_hz"`
Channels int `json:"channels"`
@@ -111,9 +116,8 @@ type currentAgentSettings struct {
} `json:"conversation"`
}
// Bind rejects schema-valid settings which the selected SDK cannot
// express. In particular, the published TTS schema is intentionally not
// silently narrowed to the SDK's speed/format capabilities.
// Bind rejects schema-valid settings which the approved provider cannot
// express; neither defaults nor lossy format/speed conversions are permitted.
func Bind(task configread.Task, providers map[string]configread.Provider) (Binding, error) {
if len(task.Raw) == 0 {
return Binding{}, errors.New("approved task snapshot is missing")
@@ -193,7 +197,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", "volcengine_tts")
ttsProvider, err := currentProvider(providers, settings.TTS.ProviderRef, "tts", "bailian_tts")
if err != nil {
return Binding{}, err
}
@@ -207,22 +211,17 @@ 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")
}
ttsRate := 0
if settings.TTS.Speed != nil {
floatRate := (*settings.TTS.Speed - 1) * 100
if floatRate < -50 || floatRate > 100 || math.Abs(floatRate-math.Round(floatRate)) > 1e-9 {
return Binding{}, errors.New("TTS speed is unsupported by selected SDK")
}
ttsRate = int(math.Round(floatRate))
if settings.TTS.Model != "qwen3-tts-flash" || settings.TTS.Voice != "Cherry" || settings.TTS.LanguageType != "Chinese" || settings.TTS.Speed == nil || *settings.TTS.Speed != 1 {
return Binding{}, errors.New("approved Bailian TTS model, voice, language or speed is unsupported")
}
ttsTimeout, err := currentTimeout(settings.TTS.TimeoutMS)
if err != nil {
return Binding{}, fmt.Errorf("TTS timeout: %w", err)
}
bound.TTS = &TTSConfig{Provider: ttsProvider, Request: doubaospeech.TTSV2Request{
Speaker: settings.TTS.Voice, ResourceID: settings.TTS.Model,
Format: doubaospeech.FormatPCMS16LE, SampleRate: ttsSampleRate, SpeechRate: ttsRate,
}, Timeout: ttsTimeout}
bound.TTS = &TTSConfig{
Provider: ttsProvider, 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
bound.AllowedVariables = append([]string(nil), settings.Prompt.AllowedVariables...)
if settings.Prompt.MaxBytes != nil {
+13 -2
View File
@@ -76,8 +76,8 @@ func TestBindCurrentFullAIUsesApprovedSDKFields(t *testing.T) {
if bound.LLM.Provider.Credential != providers["llm-example"].Credential || bound.LLM.Model != "example-chat" || bound.LLM.Temperature == nil || *bound.LLM.Temperature != 0 || bound.LLM.MaxTokens == nil || *bound.LLM.MaxTokens != 256 || bound.LLM.Timeout != 5*time.Second {
t.Fatal("LLM model/explicit zero/limit/credential/timeout not bound")
}
if bound.TTS.Provider.Credential != providers["tts-example"].Credential || bound.TTS.Request.ResourceID != "example-tts" || bound.TTS.Request.Speaker != "example-neutral" || bound.TTS.Request.Format != doubaospeech.FormatPCMS16LE || bound.TTS.Request.SampleRate != 16000 || bound.TTS.Request.SpeechRate != 0 || bound.TTS.Timeout != 5*time.Second {
t.Fatal("TTS model/voice/speed/format/credential/timeout not bound to SDK request")
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" {
t.Fatal("immutable prompt and keyword behavior not bound")
@@ -106,6 +106,7 @@ 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) {
@@ -119,6 +120,16 @@ func TestBindCurrentRejectsSDKUnsupportedTTSWithoutChangingSchema(t *testing.T)
}
}
func TestBindCurrentRejectsRemovedVolcTTSAdapter(t *testing.T) {
task, providers := currentFixture(t, "full_ai")
provider := providers["tts-example"]
provider.Adapter = "volcengine_tts"
providers[provider.ProviderRef] = provider
if _, err := Bind(task, providers); err == nil || !strings.Contains(err.Error(), "provider") {
t.Fatalf("removed TTS adapter must fail closed, got %v", err)
}
}
func TestBindCurrentRejectsUnauthorizedProviderBeforeCall(t *testing.T) {
for _, tc := range []struct {
name string
+18 -37
View File
@@ -2,46 +2,41 @@ package ai
import (
"context"
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"os/exec"
"strings"
"sync/atomic"
"testing"
)
func TestTTSRejectsExcessPendingAudioWithoutRetry(t *testing.T) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
t.Skip("Bailian audio conversion requires ffmpeg")
}
task, providers := currentFixture(t, "full_ai")
task = changeCurrentAgent(t, task, func(agent map[string]any) {
agent["conversation"].(map[string]any)["max_pending_audio_chunks"] = 1
conversation := agent["conversation"].(map[string]any)
conversation["sentence_max_chars"] = 2
conversation["max_pending_audio_chunks"] = 1
})
var requests atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
requests.Add(1)
for _, pcm := range [][]byte{{1, 0}, {2, 0}} {
_, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString(pcm))
}
_, _ = fmt.Fprintln(w, `{"code":20000000,"message":"ok","data":null}`)
}))
defer server.Close()
endpoint, requests := mockBailianTTS(t, []byte{1, 0}, nil)
p := providers["tts-example"]
p.Endpoint = server.URL
p.Endpoint = endpoint
providers[p.ProviderRef] = p
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
if _, err := bound.Synthesize(context.Background(), "批准回复"); err == nil || !strings.Contains(err.Error(), "pending audio") || strings.Contains(err.Error(), p.Credential) {
if _, err := bound.SynthesizeReply(context.Background(), "你好世界"); err == nil || !strings.Contains(err.Error(), "pending audio") || strings.Contains(err.Error(), p.Credential) {
t.Fatalf("explicit pending-audio bound was ignored or leaked credentials: %v", err)
}
if requests.Load() != 1 {
t.Fatalf("SDK must not automatically retry/bill twice: requests=%d", requests.Load())
t.Fatalf("provider must not retry/bill a second time: requests=%d", requests.Load())
}
}
func TestReplyChunksUnicodeByApprovedSentenceLimit(t *testing.T) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
t.Skip("Bailian audio conversion requires ffmpeg")
}
task, providers := currentFixture(t, "full_ai")
task = changeCurrentAgent(t, task, func(agent map[string]any) {
conversation := agent["conversation"].(map[string]any)
@@ -49,31 +44,17 @@ func TestReplyChunksUnicodeByApprovedSentenceLimit(t *testing.T) {
conversation["max_pending_audio_chunks"] = 3
})
texts := make(chan string, 3)
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var payload struct {
Params struct {
Text string `json:"text"`
} `json:"req_params"`
}
if err := json.NewDecoder(r.Body).Decode(&payload); err != nil {
http.Error(w, "invalid SDK request", http.StatusBadRequest)
return
}
texts <- payload.Params.Text
_, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString([]byte{1, 0}))
_, _ = fmt.Fprintln(w, `{"code":20000000,"message":"ok","data":null}`)
}))
defer server.Close()
endpoint, requests := mockBailianTTS(t, []byte{1, 0}, func(text string) { texts <- text })
p := providers["tts-example"]
p.Endpoint = server.URL
p.Endpoint = endpoint
providers[p.ProviderRef] = p
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
audio, err := bound.SynthesizeReply(context.Background(), "你好世界,再见")
if err != nil || len(audio) != 6 {
t.Fatalf("reply must synthesize all bounded chunks: len=%d err=%v", len(audio), err)
if err != nil || len(audio) != 6 || requests.Load() != 3 {
t.Fatalf("reply must synthesize all bounded chunks: len=%d requests=%d err=%v", len(audio), requests.Load(), err)
}
for _, want := range []string{"你好世", "界,再", "见"} {
if got := <-texts; got != want {
+7 -22
View File
@@ -2,38 +2,23 @@ package ai
import (
"context"
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"os/exec"
"strings"
"sync/atomic"
"testing"
)
func TestCallOpeningUsesApprovedTTSOnce(t *testing.T) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
t.Skip("Bailian audio conversion requires ffmpeg")
}
task, providers := currentFixture(t, "full_ai")
var calls atomic.Int32
text := make(chan string, 1)
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
calls.Add(1)
var payload struct {
Params struct {
Text string `json:"text"`
} `json:"req_params"`
}
if err := json.NewDecoder(r.Body).Decode(&payload); err != nil {
http.Error(w, "invalid SDK request", http.StatusBadRequest)
return
}
text <- payload.Params.Text
_, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString([]byte{1, 0, 2, 0}))
_, _ = fmt.Fprintln(w, `{"code":20000000,"message":"ok","data":null}`)
}))
defer server.Close()
endpoint, calls := mockBailianTTS(t, []byte{1, 0, 2, 0}, func(value string) { text <- value })
p := providers["tts-example"]
p.Endpoint = server.URL
p.Endpoint = endpoint
providers[p.ProviderRef] = p
bound, err := Bind(task, providers)
if err != nil {
@@ -61,7 +46,7 @@ func TestCallOpeningFailureIsVisibleAndNeverRetried(t *testing.T) {
}))
defer server.Close()
p := providers["tts-example"]
p.Endpoint = server.URL
p.Endpoint = server.URL + "/api/v1/services/aigc/multimodal-generation/generation"
providers[p.ProviderRef] = p
bound, err := Bind(task, providers)
if err != nil {
+6 -31
View File
@@ -228,38 +228,13 @@ func (b Binding) synthesize(ctx context.Context, text string, alreadyPending int
}
ctx, cancel := currentDeadline(ctx, b.TTS.Timeout)
defer cancel()
request := b.TTS.Request // per-call copy: concurrent calls never share mutable SDK parameters
request.Text = text
client := doubaospeech.NewClient("",
doubaospeech.WithAPIKey(b.TTS.Provider.Credential),
doubaospeech.WithBaseURL(b.TTS.Provider.Endpoint),
)
var audio []byte
chunks := 0
completed := false
for chunk, err := range client.TTSV2.Stream(ctx, &request) {
if err != nil {
return nil, chunks, err
}
if chunk == nil {
return nil, chunks, errors.New("TTS returned an empty stream chunk")
}
if len(chunk.Audio) > 0 {
chunks++
if limit > 0 && alreadyPending+chunks > limit {
return nil, chunks, errors.New("TTS exceeded approved pending audio chunk limit")
}
audio = append(audio, chunk.Audio...)
}
if chunk.IsLast {
completed = true
break
}
// Bailian produces one complete response for this approved sentence; a
// missing or failed response never counts as queued audio.
audio, err := synthesizeBailianTTS(ctx, *b.TTS, text)
if err != nil {
return nil, 0, err
}
if !completed || len(audio) == 0 || len(audio)%2 != 0 {
return nil, chunks, errors.New("TTS stream ended without complete PCM16 audio")
}
return audio, chunks, nil
return audio, 1, nil
}
func currentDeadline(ctx context.Context, timeout time.Duration) (context.Context, context.CancelFunc) {
+36 -24
View File
@@ -2,52 +2,64 @@ package ai
import (
"context"
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"os/exec"
"strings"
"testing"
"git.ipao.vip/rogee/go-sip/internal/media"
)
func TestTTSPassesApprovedParametersToSDK(t *testing.T) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
t.Skip("Bailian audio conversion requires ffmpeg")
}
want := []byte{1, 0, 2, 0}
wav, _, err := media.EncodeMonoWAV(want, 1024)
if err != nil {
t.Fatal(err)
}
task, providers := currentFixture(t, "full_ai")
task = changeCurrentAgent(t, task, func(agent map[string]any) { agent["tts"].(map[string]any)["speed"] = 1.3 })
var captured map[string]any
var key, resource string
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
key, resource = r.Header.Get("X-Api-Key"), r.Header.Get("X-Api-Resource-Id")
if r.Method != http.MethodPost || r.URL.Path != "/api/v3/tts/unidirectional" {
http.Error(w, "unexpected SDK endpoint", http.StatusBadRequest)
return
var credential string
postCount, getCount := 0, 0
var server *httptest.Server
server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/api/v1/services/aigc/multimodal-generation/generation":
postCount++
credential = r.Header.Get("Authorization")
if r.Method != http.MethodPost || json.NewDecoder(r.Body).Decode(&captured) != nil {
http.Error(w, "invalid generation request", http.StatusBadRequest)
return
}
w.Header().Set("Content-Type", "application/json")
_, _ = fmt.Fprintf(w, `{"output":{"audio":{"url":%q}}}`, server.URL+"/audio?test-token=redacted")
case "/audio":
getCount++
_, _ = w.Write(wav)
default:
http.Error(w, "unexpected endpoint", http.StatusNotFound)
}
if err := json.NewDecoder(r.Body).Decode(&captured); err != nil {
http.Error(w, "bad SDK request", http.StatusBadRequest)
return
}
_, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString([]byte{1, 0, 2, 0}))
_, _ = fmt.Fprintln(w, `{"code":20000000,"message":"ok","data":null}`)
}))
defer server.Close()
p := providers["tts-example"]
p.Endpoint = server.URL
p.Endpoint = server.URL + "/api/v1/services/aigc/multimodal-generation/generation"
providers[p.ProviderRef] = p
bound, err := Bind(task, providers)
if err != nil {
t.Fatal(err)
}
audio, err := bound.Synthesize(context.Background(), "批准的回复")
if err != nil || string(audio) != string([]byte{1, 0, 2, 0}) {
t.Fatalf("SDK TTS response: length=%d err=%v", len(audio), err)
if err != nil || string(audio) != string(want) {
t.Fatalf("Bailian TTS response: length=%d err=%v", len(audio), err)
}
if key != p.Credential || resource != "example-tts" {
t.Fatal("provider credential/resource not passed to official request")
}
params := captured["req_params"].(map[string]any)
format := params["audio_params"].(map[string]any)
if params["text"] != "批准的回复" || params["speaker"] != "example-neutral" || format["format"] != "pcm_s16le" || format["sample_rate"] != float64(16000) || format["speech_rate"] != float64(30) {
t.Fatalf("SDK TTS approved speed/voice/format not preserved: %v", format)
input, ok := captured["input"].(map[string]any)
if !ok || captured["model"] != "qwen3-tts-flash" || input["voice"] != "Cherry" || input["language_type"] != "Chinese" || input["text"] != "批准的回复" || credential != "Bearer "+p.Credential || postCount != 1 || getCount != 1 {
t.Fatal("approved Bailian model/voice/text/credential or single-request bound was lost")
}
}
@@ -2,12 +2,12 @@ package rpc
import (
"context"
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"os"
"os/exec"
"path/filepath"
"testing"
"time"
@@ -16,15 +16,18 @@ import (
"git.ipao.vip/rogee/go-sip/internal/ai"
"git.ipao.vip/rogee/go-sip/internal/configread"
"git.ipao.vip/rogee/go-sip/internal/dispatcher"
"git.ipao.vip/rogee/go-sip/internal/media"
)
type approvedSDKObservation struct {
Body map[string]any
Authorization string
ResourceID string
}
func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
t.Skip("Bailian TTS conversion requires ffmpeg")
}
now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC)
req := approvedTestRequest(t, now)
fullJSON, err := os.ReadFile("../../contracts/local/examples/config-read-task-full.json")
@@ -40,7 +43,16 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T
req.SourceEventId, req.CallId = "event-full", "event-full"
llmObserved := make(chan approvedSDKObservation, 1)
ttsObserved := make(chan approvedSDKObservation, 2)
mockSDKs := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
wav, _, err := media.EncodeMonoWAV([]byte{1, 0, 2, 0}, 1024)
if err != nil {
t.Fatal(err)
}
var mockSDKs *httptest.Server
mockSDKs = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/audio" && r.Method == http.MethodGet {
_, _ = w.Write(wav)
return
}
var body map[string]any
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, "bad SDK request", http.StatusBadRequest)
@@ -51,10 +63,10 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T
llmObserved <- approvedSDKObservation{Body: body, Authorization: r.Header.Get("Authorization")}
w.Header().Set("Content-Type", "application/json")
_, _ = fmt.Fprintln(w, `{"id":"mock","object":"chat.completion","created":0,"model":"example-chat","choices":[{"index":0,"message":{"role":"assistant","content":"回答"},"finish_reason":"stop"}]}`)
case "/api/v3/tts/unidirectional":
ttsObserved <- approvedSDKObservation{Body: body, Authorization: r.Header.Get("X-Api-Key"), ResourceID: r.Header.Get("X-Api-Resource-Id")}
_, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString([]byte{1, 0, 2, 0}))
_, _ = fmt.Fprintln(w, `{"code":20000000,"message":"ok","data":null}`)
case "/api/v1/services/aigc/multimodal-generation/generation":
ttsObserved <- approvedSDKObservation{Body: body, Authorization: r.Header.Get("Authorization")}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]any{"output": map[string]any{"audio": map[string]any{"url": mockSDKs.URL + "/audio?sig=fixture"}}})
default:
http.Error(w, "unexpected SDK request", http.StatusNotFound)
}
@@ -68,7 +80,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T
llm.Endpoint = mockSDKs.URL + "/v1"
providers[llm.ProviderRef] = llm
tts := providers["tts-example"]
tts.Endpoint = mockSDKs.URL
tts.Endpoint = mockSDKs.URL + "/api/v1/services/aigc/multimodal-generation/generation"
providers[tts.ProviderRef] = tts
req.ProvidersJson, err = json.Marshal(providers)
if err != nil {
@@ -148,11 +160,11 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T
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")
}
if ttsCall.Authorization != tts.Credential || ttsCall.ResourceID != "example-tts" {
t.Fatal("Agent did not use approved TTS resource and credential")
if ttsCall.Authorization != "Bearer "+tts.Credential || ttsCall.Body["model"] != "qwen3-tts-flash" {
t.Fatal("Agent did not use approved Bailian TTS model and credential")
}
params := ttsCall.Body["req_params"].(map[string]any)
if openingCall.Body["req_params"].(map[string]any)["text"] != "Example greeting" || params["speaker"] != "example-neutral" || params["text"] != "回答" || params["audio_params"].(map[string]any)["format"] != "pcm_s16le" {
t.Fatal("Agent did not send the approved opening then reply via TTS")
params := ttsCall.Body["input"].(map[string]any)
if openingCall.Body["input"].(map[string]any)["text"] != "Example greeting" || params["voice"] != "Cherry" || params["text"] != "回答" {
t.Fatal("Agent did not send the approved opening then reply via Bailian TTS")
}
}
+21 -10
View File
@@ -3,12 +3,11 @@ package rpc
import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"os"
"os/exec"
"strings"
"sync/atomic"
"testing"
@@ -66,6 +65,9 @@ func TestApprovedCallRunnerASROnlyHonorsSignedTimeout(t *testing.T) {
}
func TestApprovedCallRunnerSynthesizesOpeningBeforeMediaCapture(t *testing.T) {
if _, err := exec.LookPath("ffmpeg"); err != nil {
t.Skip("Bailian TTS conversion requires ffmpeg")
}
raw, err := os.ReadFile("../../contracts/local/examples/config-read-task-full.json")
if err != nil {
t.Fatal(err)
@@ -80,23 +82,32 @@ func TestApprovedCallRunnerSynthesizesOpeningBeforeMediaCapture(t *testing.T) {
t.Fatal(err)
}
var sdkCalls atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
wav, _, err := media.EncodeMonoWAV([]byte{1, 0, 2, 0}, 1024)
if err != nil {
t.Fatal(err)
}
var server *httptest.Server
server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/audio" {
_, _ = w.Write(wav)
return
}
sdkCalls.Add(1)
var payload struct {
Params struct {
Input struct {
Text string `json:"text"`
} `json:"req_params"`
} `json:"input"`
}
if err := json.NewDecoder(r.Body).Decode(&payload); err != nil || payload.Params.Text != "Example greeting" {
if err := json.NewDecoder(r.Body).Decode(&payload); err != nil || payload.Input.Text != "Example greeting" {
http.Error(w, "opening text changed", http.StatusBadRequest)
return
}
_, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString([]byte{1, 0, 2, 0}))
_, _ = fmt.Fprintln(w, `{"code":20000000,"message":"ok","data":null}`)
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]any{"output": map[string]any{"audio": map[string]any{"url": server.URL + "/audio?sig=fixture"}}})
}))
defer server.Close()
tts := providers["tts-example"]
tts.Endpoint = server.URL
tts.Endpoint = server.URL + "/api/v1/services/aigc/multimodal-generation/generation"
providers[tts.ProviderRef] = tts
bound, err := ai.Bind(task, providers)
if err != nil {
@@ -108,7 +119,7 @@ func TestApprovedCallRunnerSynthesizesOpeningBeforeMediaCapture(t *testing.T) {
if err != nil {
t.Fatal(err)
}
_, err = RunApprovedCall(context.Background(), ApprovedExecution{AI: bound, MaxCallDuration: 100 * time.Millisecond}, mediaSession, hangup, call)
_, err = RunApprovedCall(context.Background(), ApprovedExecution{AI: bound, MaxCallDuration: time.Second}, mediaSession, hangup, call)
if err == nil || !strings.Contains(err.Error(), "deadline") || sdkCalls.Load() != 1 || mediaSession.sent != 1 || mediaSession.rate != 16000 {
t.Fatalf("approved opening must be played once before timed capture: requests=%d sent=%d rate=%d err=%v", sdkCalls.Load(), mediaSession.sent, mediaSession.rate, err)
}