diff --git a/AGENTS.md b/AGENTS.md index 818412f..90f4764 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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 不写入样例、日志、源码或证据。 diff --git a/contracts/local/config-read.schema.json b/contracts/local/config-read.schema.json index 171ace8..fdfa651 100644 --- a/contracts/local/config-read.schema.json +++ b/contracts/local/config-read.schema.json @@ -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": { diff --git a/contracts/local/examples/config-read-providers.json b/contracts/local/examples/config-read-providers.json index cf89a5c..f68af7d 100644 --- a/contracts/local/examples/config-read-providers.json +++ b/contracts/local/examples/config-read-providers.json @@ -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"} ] } diff --git a/contracts/local/examples/config-read-task-full.json b/contracts/local/examples/config-read-task-full.json index 25be20c..8efc346 100644 --- a/contracts/local/examples/config-read-task-full.json +++ b/contracts/local/examples/config-read-task-full.json @@ -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} } diff --git a/contracts/local/manifest.json b/contracts/local/manifest.json index 132bd5c..1176a3a 100644 --- a/contracts/local/manifest.json +++ b/contracts/local/manifest.json @@ -4,8 +4,8 @@ "sources": { "docs/archive/sources/v0.5-proposal.md": "612fdaee50aff6aa7fbef16c2d469d99857646c6d2235617d0e67f6098cd7ada", "docs/archive/sources/plan-saas-dispatcher-v05-v0.1.md": "666f39e56ea9f4b55661efcac82edd6f9729848e2d60e5f24cdf5aa3ac97ee87", - "docs/thirds/saas-dispatcher.md": "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" } diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index d7e97e2..57a0223 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -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,不能宣称真实通话媒体已验证。 diff --git a/docs/thirds/saas-dispatcher.md b/docs/thirds/saas-dispatcher.md index 87820c1..c85e9a1 100644 --- a/docs/thirds/saas-dispatcher.md +++ b/docs/thirds/saas-dispatcher.md @@ -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 且 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 参数验证不等于真实供应商验收。 diff --git a/internal/ai/bailian_tts.go b/internal/ai/bailian_tts.go new file mode 100644 index 0000000..ddfa003 --- /dev/null +++ b/internal/ai/bailian_tts.go @@ -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 +} diff --git a/internal/ai/bailian_tts_test.go b/internal/ai/bailian_tts_test.go new file mode 100644 index 0000000..1c66eee --- /dev/null +++ b/internal/ai/bailian_tts_test.go @@ -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) + } +} diff --git a/internal/ai/binding.go b/internal/ai/binding.go index 1db6bfb..6ac8dc6 100644 --- a/internal/ai/binding.go +++ b/internal/ai/binding.go @@ -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 { diff --git a/internal/ai/binding_test.go b/internal/ai/binding_test.go index f2f323c..fda0cc5 100644 --- a/internal/ai/binding_test.go +++ b/internal/ai/binding_test.go @@ -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 diff --git a/internal/ai/limits_test.go b/internal/ai/limits_test.go index 90e6426..e4754fa 100644 --- a/internal/ai/limits_test.go +++ b/internal/ai/limits_test.go @@ -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 { diff --git a/internal/ai/opening_test.go b/internal/ai/opening_test.go index 6248273..5bb6b59 100644 --- a/internal/ai/opening_test.go +++ b/internal/ai/opening_test.go @@ -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 { diff --git a/internal/ai/pipeline.go b/internal/ai/pipeline.go index a1487ab..c049462 100644 --- a/internal/ai/pipeline.go +++ b/internal/ai/pipeline.go @@ -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) { diff --git a/internal/ai/pipeline_test.go b/internal/ai/pipeline_test.go index 72e0151..346127a 100644 --- a/internal/ai/pipeline_test.go +++ b/internal/ai/pipeline_test.go @@ -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") } } diff --git a/internal/rpc/approved_full_ai_integration_test.go b/internal/rpc/approved_full_ai_integration_test.go index f077b35..344eac9 100644 --- a/internal/rpc/approved_full_ai_integration_test.go +++ b/internal/rpc/approved_full_ai_integration_test.go @@ -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") } } diff --git a/internal/rpc/approved_runner_test.go b/internal/rpc/approved_runner_test.go index f689a85..52b0512 100644 --- a/internal/rpc/approved_runner_test.go +++ b/internal/rpc/approved_runner_test.go @@ -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) }