From b47fbfa9cf1c3ce255c0599ecc702a21a148d96e Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 7 Oct 2026 12:40:17 +0800 Subject: [PATCH] Bind each SIP trunk to its approved caller identity --- AGENTS.md | 6 +-- contracts/local/config-read.schema.json | 7 ++-- contracts/local/examples/config-read-sip.json | 2 +- .../local/examples/config-read-task-asr.json | 2 +- .../local/examples/config-read-task-full.json | 2 +- .../config-read-task-asr-with-llm.json | 2 +- .../mq-result-failure-missing-message.json | 2 +- .../invalid/mq-result-invented-identity.json | 2 +- .../mq-result-invented-recording-state.json | 2 +- .../examples/mq-result-no-recording.json | 2 +- .../local/examples/mq-result-uploaded.json | 2 +- contracts/local/manifest.json | 4 +- contracts/local/mq.schema.json | 4 +- docs/thirds/saas-dispatcher.md | 4 +- internal/agent/recording_delivery_test.go | 4 +- internal/asterisk/loader.go | 31 +++++++++++---- internal/asterisk/loader_test.go | 17 ++++++++- internal/asterisk/sip.go | 3 +- internal/asterisk/sip_test.go | 38 +++++++++++++++++-- internal/callflow/result_payload.go | 6 +-- internal/callflow/result_payload_test.go | 20 +++++----- internal/configread/snapshots.go | 1 - internal/contract/current_test.go | 22 +++++++++++ internal/dispatcher/execute.go | 3 +- internal/dispatcher/policy.go | 20 ++-------- internal/dispatcher/policy_test.go | 29 ++++++++++++++ internal/rpc/approved_execution.go | 3 +- internal/rpc/approved_recorded_mock.go | 4 +- internal/rpc/approved_recorded_mock_test.go | 2 +- internal/rpc/approved_recorded_real.go | 8 ++-- internal/rpc/approved_snapshot_test.go | 2 +- .../recording_delivery_integration_test.go | 4 +- internal/rpc/recording_server_flow_test.go | 4 +- internal/store/result.go | 15 ++++---- internal/store/result_test.go | 6 +-- 35 files changed, 188 insertions(+), 97 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 92eb8d8..a9bea84 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -88,9 +88,9 @@ - SaaS→Dispatcher 的**五类只读配置**为 `GET /internal/v1/dispatcher/sip`、`/task/:task_id`、`/tasks`、`/tenant/:tenant_id/quota`、`/ai-providers`;路径前缀固定,均须校验 Dispatcher UUID/资源归属、数字 `tenant_id`、完整快照、来源、有效授权和版本。配置读失败、过期、矛盾或不确定时关新准入;没有旧 MQ 配置回退、通用业务 HTTP、ETag 兜底或偷偷启用旧执行字段。已接纳任务持久绑定原快照。 - 呼叫、控制、必要回执和**每通话唯一最终结果**经固定 `v1` RabbitMQ Topic/队列,任务与控制队列由 SaaS 预建,Dispatcher 不可自行建/删/绑定;stop 时先停该任务消费者并关闭持久准入,再 purge 仅该任务队列的待投递消息,失败不回成功;SaaS 停止继续投递已 stop 任务,后续误投递不执行。结果进入指定共享 durable 队列。独立 D UUID 和接收队列不能广播后正文过滤;`tenant_key` 原值保留,任务只能由归属 D 执行。入站先校验和持久 inbox 后 ACK,状态/outbox 同事务;persistent、mandatory、无 return、publisher confirm 成功才记交付,confirm **不是** SaaS 应用收讫。失败/确认丢失与重启只重发同身份消息,不重复拨号或捏造结果。 -- `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 须重新核验。历史命令不因消息年龄过期,但实际呼出前仍检查任务归属/事件号码格式/时段/授权/额度和有效通话上限;任务结束不删未确认结果、恢复、幂等或未知占用。 +- `call.execute` 只带获批 `task_id/callee`,调用线路、AI 和时限由该任务快照固定;主叫由绑定线路的 SIP 快照唯一 `caller_id` 决定,不再由任务挑选主叫档案;`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/租约到期不得自动清除未知执行;明确 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 配置的唯一编辑/审批面仍是 management。 +- Dispatcher↔Agent 复用 Unary gRPC 和受控 Endpoint;Agent 预绑定 D UUID 与服务端证书指纹,激活/会话代际、peer mTLS/SAN/SNI 和已签发期限须核对,新 boot 不清未知占用。Agent 不自行向 SaaS 取任务/AI/OSS 授权;Dispatcher 只用已经核验的 Agent `GetLoadedSIP` revision 开执行准入。本机 Mock 的加载回报不证明 Asterisk 已实际加载;仅隔离 SIP-only Agent 在配置原子写入、PJSIP reload 和运行态 endpoint/AOR/UDP transport 一致后持久标记 revision,每次加载查询重新核验。只支持明确的 UDP、IP/none 鉴权、无需 REGISTER 的 IPv4/PCMA 线路;未知字段和其他传输/鉴权/注册方式拒绝,不热更静态 transport。每条 SIP trunk 只允许一个原值 `caller_id`,Agent 把该值同时配置为 Asterisk endpoint 的 `from_user` 和 `callerid`,并在 `GetLoadedSIP` 前核对实际运行值;旧任务 `caller_profile_id` 与线路 `caller_profiles` 必须直接拒绝,不保留兼容路径。SIP 配置的唯一编辑/审批面仍是 management。 - 非生产真实 Agent 的 Asterisk `res_hep`/`res_hep_pjsip` 仅镜像至本机 UDP:在 ARI Dial 前读取本次 SIP Call-ID 并绑定执行,按 Call-ID 和原始 INVITE transaction 只采集真实最终响应的状态码、状态行与完整 `raw`;丢失/无法关联以 `sip_capture_error` 显式报告,不从 ARI、号码或挂断原因推测。镜像缺失不阻止已证实终结的额度释放;未经确认的占用仍保留。配置及回滚见 [`deploys/cell/README.md`](deploys/cell/README.md),不得将完整 SIP 报文写入普通日志、提交或聊天;本地测试通过不等于测试机已部署 HEP 或真实接通。 - AI 使用任务内不可变授权快照:仅经获批准百炼/火山 ASR、OpenAI 兼容 LLM、百炼 TTS(`qwen3-tts-flash`/`Cherry`/`Chinese`)能表达的参数进入每通话实例;ASR-only 不启动 LLM/TTS,完整 AI 不借旧语音测试的授权或参数。只有最终用户 ASR 文本的明确字面关键词可触发拒联/挂断;不由 SDK 默认值、环境、CLI、metadata 或宽松 Schema 改写业务参数,不因 SDK 重试产生第二次发起/收费或重播。日志只存脱敏版本/摘要/计数,不存密钥、prompt、完整对话或音频。 - **私有配置位置(本机路径相对本仓库根目录,只读,绝不提交)**:`.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`。 @@ -108,7 +108,7 @@ - 全线路的原始被叫号码仅由校验归属及任务后的 SaaS `call.execute.payload.callee` 确定;不在 Dispatcher、SaaS Mock 或抓证脚本设置固定号码/日期特判,不向任务快照增设号码列表。只接受 1–32 位 ASCII 数字;格式校验不等于真实拨号授权。不再另设本地固定的 `09:00`–`20:00` 窗口或每线路每号码每日 3 次上限;任务、线路时段和额度以 SaaS 配置快照校验为准,不等候/自动延迟/自动重试/静默换线。每次真实试拨仍须使用者明确安排,并由专用主机脚本在拨号前启用抓包和 PJSIP logger、签发与该通 `event_id`/trunk/原始号码绑定的短时有效活跃抓包凭证;Agent 拒绝缺失/失效/不匹配的凭证。SaaS Mock 投递和本地时段测试不构成真实拨号授权。 - 任务按周一至周日多个时段与指定排除日期配置,线路只有每周允许时段(**没有线路排除日期**),Asia/Shanghai 左闭右开、跨日拆分;缺失或不确定 fail-closed,不自动重拨。由 Dispatcher 在持久接纳与实际发出指令前判定,并取任务/获批 AI 较小通话时限;Agent 仅校验会话和签发期限,不重算外呼策略。本地 SaaS 快照策略测试与主机抓证/逐次授权须分别报告。 -- `BD` 等主叫原值不得清洗或当作 Digest 用户名;业务原始被叫号码不变,仅被选定数企 trunk 按规则构造 `7089<原号>`(其它线路使用自己的前缀),不重复加前缀。三条 trunk 独立,不能把同地址伪造为备用线路或换线重拨;服务商反馈 PCMA,对应 Asterisk `allow=alaw`,传输/注册/鉴权/并发仍待真实签收。不以 sipgo/diago 另造 Asterisk 替代架构。 +- 当前修复仅在本地核验,未更新测试机上的旧多主叫 Mock 快照,也未部署或授权新呼叫。`BD` 等主叫原值不得清洗或当作 Digest 用户名;业务原始被叫号码不变,仅被选定数企 trunk 按规则构造 `7089<原号>`(其它线路使用自己的前缀),不重复加前缀。三条 trunk 独立,不能把同地址伪造为备用线路或换线重拨;服务商反馈 PCMA,对应 Asterisk `allow=alaw`,传输/注册/鉴权/并发仍待真实签收。不以 sipgo/diago 另造 Asterisk 替代架构。 ## 运行环境、诊断与开发门禁 diff --git a/contracts/local/config-read.schema.json b/contracts/local/config-read.schema.json index fdfa651..4b45699 100644 --- a/contracts/local/config-read.schema.json +++ b/contracts/local/config-read.schema.json @@ -25,7 +25,7 @@ }, "trunk": { "type": "object", "additionalProperties": false, - "required": ["trunk_id", "provider_id", "codec", "dial_prefix", "enabled", "server_host", "server_port", "transport", "auth_mode", "registration_required", "max_concurrent_calls", "caller_profiles", "schedule"], + "required": ["trunk_id", "provider_id", "codec", "dial_prefix", "enabled", "server_host", "server_port", "transport", "auth_mode", "registration_required", "max_concurrent_calls", "caller_id", "schedule"], "properties": { "trunk_id": {"type": "string", "minLength": 1}, "provider_id": {"type": "string", "minLength": 1}, @@ -38,7 +38,7 @@ "auth_mode": {"enum": ["ip", "digest", "none", null]}, "registration_required": {"type": ["boolean", "null"]}, "max_concurrent_calls": {"type": ["integer", "null"], "minimum": 1}, - "caller_profiles": {"type": "array", "items": {"type": "object", "additionalProperties": false, "required": ["caller_profile_id", "caller_id"], "properties": {"caller_profile_id": {"type": "string", "minLength": 1}, "caller_id": {"type": "string", "minLength": 1}}}}, + "caller_id": {"type": "string", "pattern": "^[A-Za-z0-9][A-Za-z0-9._+-]{0,63}$"}, "schedule": {"$ref": "#/$defs/weekly_schedule"} } }, @@ -65,7 +65,7 @@ }, "task": { "type": "object", "additionalProperties": false, - "required": ["resource", "dispatcher_id", "tenant_id", "task_id", "task_revision", "status", "max_concurrent_calls", "ring_timeout_ms", "max_call_duration_ms", "route_policy_id", "caller_profile_id", "allowed_trunk_ids", "schedule", "agent"], + "required": ["resource", "dispatcher_id", "tenant_id", "task_id", "task_revision", "status", "max_concurrent_calls", "ring_timeout_ms", "max_call_duration_ms", "route_policy_id", "allowed_trunk_ids", "schedule", "agent"], "properties": { "resource": {"const": "task_config"}, "dispatcher_id": {"$ref": "#/$defs/dispatcher_id"}, @@ -78,7 +78,6 @@ "ring_timeout_ms": {"type": "integer", "minimum": 1}, "max_call_duration_ms": {"type": "integer", "minimum": 1}, "route_policy_id": {"type": "string", "minLength": 1}, - "caller_profile_id": {"type": "string", "minLength": 1}, "allowed_trunk_ids": {"type": "array", "minItems": 1, "uniqueItems": true, "items": {"type": "string", "minLength": 1}}, "schedule": {"$ref": "#/$defs/task_schedule"}, "agent": {"$ref": "#/$defs/agent"} diff --git a/contracts/local/examples/config-read-sip.json b/contracts/local/examples/config-read-sip.json index 588c12f..96eb4b0 100644 --- a/contracts/local/examples/config-read-sip.json +++ b/contracts/local/examples/config-read-sip.json @@ -6,7 +6,7 @@ "trunk_id":"trunk-mock","provider_id":"provider-mock","codec":"PCMA","dial_prefix":"","enabled":true, "server_host":"sip.example.invalid","server_port":5060, "transport":null,"auth_mode":null,"registration_required":null,"max_concurrent_calls":null, - "caller_profiles":[{"caller_profile_id":"caller-profile-mock","caller_id":"BD00000000"}], + "caller_id":"BD00000000", "schedule":{"time_zone":"Asia/Shanghai","weekly_windows":{"monday":[{"start":"09:00","end":"20:00"}],"tuesday":[],"wednesday":[],"thursday":[],"friday":[],"saturday":[],"sunday":[]}} }] } diff --git a/contracts/local/examples/config-read-task-asr.json b/contracts/local/examples/config-read-task-asr.json index 796aff2..653b010 100644 --- a/contracts/local/examples/config-read-task-asr.json +++ b/contracts/local/examples/config-read-task-asr.json @@ -2,7 +2,7 @@ "resource":"task_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001, "task_id":"task-asr","task_revision":1,"status":"running","name":"ASR example", "max_concurrent_calls":2,"ring_timeout_ms":30000,"max_call_duration_ms":120000, - "route_policy_id":"route-mock","caller_profile_id":"caller-profile-mock","allowed_trunk_ids":["trunk-mock"], + "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":"asr_only","asr":{"provider_ref":"asr-example","language":"zh-CN","input":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1,"sample_width_bytes":2},"interim":false,"timeout_ms":5000}} } diff --git a/contracts/local/examples/config-read-task-full.json b/contracts/local/examples/config-read-task-full.json index 8efc346..eebd23c 100644 --- a/contracts/local/examples/config-read-task-full.json +++ b/contracts/local/examples/config-read-task-full.json @@ -2,7 +2,7 @@ "resource":"task_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001, "task_id":"task-full","task_revision":2,"status":"running","name":"Full AI example", "max_concurrent_calls":2,"ring_timeout_ms":30000,"max_call_duration_ms":120000, - "route_policy_id":"route-mock","caller_profile_id":"caller-profile-mock","allowed_trunk_ids":["trunk-mock"], + "route_policy_id":"route-mock","allowed_trunk_ids":["trunk-mock"], "schedule":{"time_zone":"Asia/Shanghai","starts_at":"2026-09-21T00:00:00+08:00","ends_at":null,"weekly_windows":{"monday":[{"start":"09:00","end":"20:00"}],"tuesday":[],"wednesday":[],"thursday":[],"friday":[],"saturday":[],"sunday":[]},"excluded_dates":["2026-10-01"]}, "agent":{ "immutable":true,"mode":"full_ai", diff --git a/contracts/local/examples/invalid/config-read-task-asr-with-llm.json b/contracts/local/examples/invalid/config-read-task-asr-with-llm.json index fd43173..29ff37d 100644 --- a/contracts/local/examples/invalid/config-read-task-asr-with-llm.json +++ b/contracts/local/examples/invalid/config-read-task-asr-with-llm.json @@ -1 +1 @@ -{"resource":"task_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"task_id":"task-asr","task_revision":1,"status":"running","max_concurrent_calls":2,"ring_timeout_ms":30000,"max_call_duration_ms":120000,"route_policy_id":"route-mock","caller_profile_id":"caller-profile-mock","allowed_trunk_ids":["trunk-mock"],"schedule":{"time_zone":"Asia/Shanghai","starts_at":null,"ends_at":null,"weekly_windows":{"monday":[],"tuesday":[],"wednesday":[],"thursday":[],"friday":[],"saturday":[],"sunday":[]},"excluded_dates":[]},"agent":{"immutable":true,"mode":"asr_only","asr":{"provider_ref":"asr-example","language":"zh-CN","input":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1,"sample_width_bytes":2}},"llm":{"provider_ref":"llm-example","model":"must-not-be-used"}}} +{"resource":"task_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"task_id":"task-asr","task_revision":1,"status":"running","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":null,"ends_at":null,"weekly_windows":{"monday":[],"tuesday":[],"wednesday":[],"thursday":[],"friday":[],"saturday":[],"sunday":[]},"excluded_dates":[]},"agent":{"immutable":true,"mode":"asr_only","asr":{"provider_ref":"asr-example","language":"zh-CN","input":{"encoding":"pcm_s16le","sample_rate_hz":16000,"channels":1,"sample_width_bytes":2}},"llm":{"provider_ref":"llm-example","model":"must-not-be-used"}}} diff --git a/contracts/local/examples/invalid/mq-result-failure-missing-message.json b/contracts/local/examples/invalid/mq-result-failure-missing-message.json index 025cbf6..cc528b6 100644 --- a/contracts/local/examples/invalid/mq-result-failure-missing-message.json +++ b/contracts/local/examples/invalid/mq-result-failure-missing-message.json @@ -1 +1 @@ -{"event_id":"call-result-example","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","caller_profile_id":"caller-profile-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"failed","reason_code":null,"transcript":[],"opt_out":false,"recording":{}}} +{"event_id":"call-result-example","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"failed","reason_code":null,"transcript":[],"opt_out":false,"recording":{}}} diff --git a/contracts/local/examples/invalid/mq-result-invented-identity.json b/contracts/local/examples/invalid/mq-result-invented-identity.json index 1fa31d6..cd9d48b 100644 --- a/contracts/local/examples/invalid/mq-result-invented-identity.json +++ b/contracts/local/examples/invalid/mq-result-invented-identity.json @@ -1 +1 @@ -{"event_id":"call-result-example","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","call_id":"legacy-call-id","caller_profile_id":"caller-profile-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"failed","reason_code":null,"reason_message":"storage unavailable","transcript":[],"opt_out":false,"recording":{}}} +{"event_id":"call-result-example","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","call_id":"legacy-call-id","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"failed","reason_code":null,"reason_message":"storage unavailable","transcript":[],"opt_out":false,"recording":{}}} diff --git a/contracts/local/examples/invalid/mq-result-invented-recording-state.json b/contracts/local/examples/invalid/mq-result-invented-recording-state.json index 5b53df1..2657689 100644 --- a/contracts/local/examples/invalid/mq-result-invented-recording-state.json +++ b/contracts/local/examples/invalid/mq-result-invented-recording-state.json @@ -1 +1 @@ -{"event_id":"call-result-example","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","caller_profile_id":"caller-profile-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"failed","reason_code":null,"reason_message":"storage unavailable","transcript":[],"opt_out":false,"recording":{"status":"unavailable"}}} +{"event_id":"call-result-example","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"failed","reason_code":null,"reason_message":"storage unavailable","transcript":[],"opt_out":false,"recording":{"status":"unavailable"}}} diff --git a/contracts/local/examples/mq-result-no-recording.json b/contracts/local/examples/mq-result-no-recording.json index d994d49..9650b65 100644 --- a/contracts/local/examples/mq-result-no-recording.json +++ b/contracts/local/examples/mq-result-no-recording.json @@ -1 +1 @@ -{"event_id":"call-result-example","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","caller_profile_id":"caller-profile-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"failed","reason_code":null,"status_line":null,"raw":null,"sip_capture_error":null,"reason_message":"recording storage unavailable","transcript":[],"opt_out":false,"recording":{}}} +{"event_id":"call-result-example","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"failed","reason_code":null,"status_line":null,"raw":null,"sip_capture_error":null,"reason_message":"recording storage unavailable","transcript":[],"opt_out":false,"recording":{}}} diff --git a/contracts/local/examples/mq-result-uploaded.json b/contracts/local/examples/mq-result-uploaded.json index c179448..4a5b65e 100644 --- a/contracts/local/examples/mq-result-uploaded.json +++ b/contracts/local/examples/mq-result-uploaded.json @@ -1 +1 @@ -{"event_id":"call-result-example-2","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","caller_profile_id":"caller-profile-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"answered","reason_code":null,"status_line":null,"raw":null,"sip_capture_error":null,"reason_message":"answered, SIP status unavailable in isolated Mock","transcript":[{"turn_id":"turn-1","segment_id":"segment-1","role":"user","text":"Example utterance","start_ms":0,"end_ms":1000}],"opt_out":true,"recording":{"status":"uploaded","bucket":"example-bucket","object_key":"example/recording.wav","format":"wav","channels":1,"sample_rate_hz":16000,"duration_ms":60000,"size_bytes":64000,"checksum_sha256":"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}}} +{"event_id":"call-result-example-2","event_type":"call.execute.result","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:02:00Z","payload":{"task_id":"task-asr","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:00:00Z","ended_at":"2026-09-21T01:01:00Z","duration_ms":60000,"outcome":"answered","reason_code":null,"status_line":null,"raw":null,"sip_capture_error":null,"reason_message":"answered, SIP status unavailable in isolated Mock","transcript":[{"turn_id":"turn-1","segment_id":"segment-1","role":"user","text":"Example utterance","start_ms":0,"end_ms":1000}],"opt_out":true,"recording":{"status":"uploaded","bucket":"example-bucket","object_key":"example/recording.wav","format":"wav","channels":1,"sample_rate_hz":16000,"duration_ms":60000,"size_bytes":64000,"checksum_sha256":"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}}} diff --git a/contracts/local/manifest.json b/contracts/local/manifest.json index fe61045..ea0e336 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": "e9e780a75c203ecfc36e43db93bf22a12b52c4c337bd6d2dd410e105ca7326f9" + "docs/thirds/saas-dispatcher.md": "f0dc3c419a3f2de2ad08e00a65e6259f7810035fe8c59feffcfcc024ed4a6a53" }, - "bundle_sha256": "e5ac2b46cb6544778774cce4616ae0d9b6da941206805f81e1a4a6aaa3e3a3d5", + "bundle_sha256": "29e590fbca2a389e8c003ef475398f16968c69fa4bdc24d9bc35ea07e6934bd1", "bundle_algorithm": "sha256 of sorted relative-path + space + sha256(file) + newline; only root-level JSON and examples/**/*.json, excluding manifest.json" } diff --git a/contracts/local/mq.schema.json b/contracts/local/mq.schema.json index b48d0a4..9ef72ec 100644 --- a/contracts/local/mq.schema.json +++ b/contracts/local/mq.schema.json @@ -72,9 +72,9 @@ "properties": { "event_id": {"$ref": "#/$defs/event_id"}, "event_type": {"const": "call.execute.result"}, "dispatcher_id": {"$ref": "#/$defs/dispatcher_id"}, "tenant_id": {"$ref": "#/$defs/tenant_id"}, "issued_at": {"$ref": "#/$defs/issued_at"}, "payload": {"type": "object", "additionalProperties": false, - "required": ["task_id", "caller_profile_id", "callee", "trunk_id", "started_at", "ended_at", "duration_ms", "outcome", "reason_code", "status_line", "raw", "sip_capture_error", "transcript", "opt_out", "recording"], + "required": ["task_id", "callee", "trunk_id", "started_at", "ended_at", "duration_ms", "outcome", "reason_code", "status_line", "raw", "sip_capture_error", "transcript", "opt_out", "recording"], "properties": { - "task_id": {"$ref": "#/$defs/task_id"}, "caller_profile_id": {"type": "string", "minLength": 1}, "callee": {"type": "string", "minLength": 1}, "trunk_id": {"type": "string", "minLength": 1}, + "task_id": {"$ref": "#/$defs/task_id"}, "callee": {"type": "string", "minLength": 1}, "trunk_id": {"type": "string", "minLength": 1}, "started_at": {"$ref": "#/$defs/issued_at"}, "ended_at": {"$ref": "#/$defs/issued_at"}, "duration_ms": {"type": "integer", "minimum": 0}, "outcome": {"enum": ["answered", "no_answer", "failed"]}, "reason_code": {"type": ["integer", "null"], "minimum": 100, "maximum": 699}, diff --git a/docs/thirds/saas-dispatcher.md b/docs/thirds/saas-dispatcher.md index 1ca5423..e437186 100644 --- a/docs/thirds/saas-dispatcher.md +++ b/docs/thirds/saas-dispatcher.md @@ -26,7 +26,7 @@ RabbitMQ 是 Topic,**SaaS 独占创建、绑定、退役 exchange/queue,D | --- | --- | --- | | `sip.config` | SaaS→D;`revision` 触发全量重新读取和实际加载核验,不用通知正文代替全量 | [`notification`](../../contracts/local/examples/mq-sip-change.json) | | `task.control` | SaaS→D;start/pause/resume/stop,无控制去重/CAS 字段。start 读取新任务配置及租户配额并持久绑定,无 Agent 控制动作;resume 重读最新任务配置和额度后解除人工暂停;均沿用启动时 AI 服务商列表及已生效 SIP 快照,不重复拉取全局服务商或 SIP,也不额外核验 SIP;pause/stop 发给 Agent,省略 `active_call_policy` 默认 **hangup**,显式仅 drain/hangup。配置读取、归属或 AI 核验失败时不接纳,SIP 待处理时新呼叫仍关闭;stop 先停止该任务消费者(使在途未确认消息回队列),持久关闭任务准入并抑制本地待执行指令,再对该任务队列执行 purge 清除当时待投递消息;purge 失败不发送成功回执、保持准入关闭并重试原控制消息。SaaS 须停止继续投递已 stop 的任务;purge 不阻止之后新投递,这些消息不执行。成功操作后才回应用回执,回执不是已完成活跃通话排空/挂断、SaaS 已停止投递或未来消息不存在的证据;stop 同 ID 不可恢复 | [`start`](../../contracts/local/examples/mq-control-start.json) · [`control`](../../contracts/local/examples/mq-control.json) · [`ack`](../../contracts/local/examples/mq-control-ack.json) | -| `call.execute` | SaaS→D 仅 `{task_id,callee}`;一次指令保留独立消息/执行身份,路由/主叫/AI/时限从绑定任务读取。`callee` 原值由通过归属及任务校验的 SaaS 事件确定,不从任务快照或本地号码列表另选;`dispatched` 表示已发出呼叫指令,**不表示接通**;号码格式不合规则回 `rejected,reason_code:null,reason_message`、不拨号不发最终结果、不暂停整任务 | [`execute`](../../contracts/local/examples/mq-execute.json) · [`dispatched`](../../contracts/local/examples/mq-execute-ack.json) · [`rejected`](../../contracts/local/examples/mq-execute-rejected.json) | +| `call.execute` | SaaS→D 仅 `{task_id,callee}`;一次指令保留独立消息/执行身份,路由/AI/时限从绑定任务读取,主叫由绑定 SIP 线路快照的唯一 `caller_id` 决定。`callee` 原值由通过归属及任务校验的 SaaS 事件确定,不从任务快照或本地号码列表另选;`dispatched` 表示已发出呼叫指令,**不表示接通**;号码格式不合规则回 `rejected,reason_code:null,reason_message`、不拨号不发最终结果、不暂停整任务 | [`execute`](../../contracts/local/examples/mq-execute.json) · [`dispatched`](../../contracts/local/examples/mq-execute-ack.json) · [`rejected`](../../contracts/local/examples/mq-execute-rejected.json) | | `call.execute.result` | D→SaaS;按 `task_id` + 原号码关联,每次呼叫仅一份最终结果;不新增外部 call_id/source_command_id;真正终结且录音成功上传、无录音或预期录音生成失败后才发送 | [`uploaded`](../../contracts/local/examples/mq-result-uploaded.json) · [`empty`](../../contracts/local/examples/mq-result-no-recording.json) | 旧 `command.result`、`call.result`、分散通话/转写/拒联事件、`recording.uploaded` 不再作为对外并行通知或兼容别名。D 在 inbox 持久后 ACK;状态/outbox 同事务;结果发布使用原消息身份可靠交付;confirm 不是 SaaS 应用收讫。消息年龄不让旧命令绕过准入;未来 issued_at 不提前接纳。重复投递/未知执行不触发再次拨号。 @@ -34,7 +34,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 不构成真实百炼/通话验收。 -- 被叫号码只来自归属当前 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 收件均不构成真实拨号授权。 +- 每条 SaaS SIP 线路只配置一个原值 `caller_id`;任务只选获准线路,不再带 `caller_profile_id`,线路不再带 `caller_profiles`。Dispatcher 将线路主叫绑定到每次执行;Agent 在 Asterisk PJSIP 线路同时设置 `from_user` 与 `callerid`,核对实际加载的出站身份后才准入,不以 ARI 创建通道后的 `CALLERID(num)` 代替 SIP `From`。主叫不是 Digest 用户名,不从环境、CLI 或被叫前缀推断。被叫号码只来自归属当前 D、租户及已接纳任务的 `call.execute.payload.callee` 原值;不在 Dispatcher、SaaS Mock 或抓证脚本另设固定号码列表,也不在任务快照增加号码列表。Dispatcher 与非生产发布/抓证入口只接受 1–32 位 ASCII 数字的原始号码,不能把线路前缀当成该号码的本地替代值。已选 SIP trunk、任务与线路每周时段、任务排除日期、任务/租户/线路额度、任务与 AI 较小通话时限均在接纳及实际发呼叫指令前检查。线路字段未知则 fail-closed;选线后固定、不自动重拨/换线。隔离 Mock 中规则暂不满足时保留待执行指令、暂停该任务的调度,规则允许后重验;与人工 pause/stop 分离,不能自动解除人为停止。非生产真实路径不另设固定的 `09:00`–`20:00` 时间门禁或每线路每原始号码每日 3 次上限;任务/线路时段与额度以 SaaS 已校验快照为准。逐通抓证、Agent 活跃凭证及使用者逐次授权仍须满足;格式有效或本地 Mock 收件均不构成真实拨号授权。 - 接通事实为真时 `outcome=answered`(后续异常不抹掉接通);已发起但忙线、拒接、无人接听且确定结束为 `no_answer`;确认未接通并由 Agent/Asterisk 执行故障终结为 `failed`;未知状态保持未知占用,不能伪造结束、结果或自动重拨。真实 SIP 状态码原样数字写入 `reason_code`,无真实 SIP 码则 `null` 并以 `reason_message` 说明;禁止本地虚构数字错误码。`call.execute.result.payload` 的 `status_line`、`raw` 与 `sip_capture_error` 始终存在:仅将经同一 ARI 通道拨号前取得的 SIP Call-ID 与 HEP INVITE 事务严格关联的最终响应写入原样状态行、完整原样报文与状态码;`raw` 不拼装、不截断,不能从目标号码、时间、挂断原因或 ARI HTTP 状态猜测。无 SIP 响应时前两项为 `null`;若已发起 SIP 但镜像/关联/解码失败,第三项须写明确错误,已确认结束仍报告真实结果并释放额度,不以原文缺失伪装为通话未知。原始报文只进入受控结果通道,不写日志、仓库或长期测试证据。无应答且没有录音时 `transcript=[]`、`opt_out=false`、`recording={}`。 - 只有**用户侧 ASR 最终识别文本**包含任一 `hangup_keywords` 字面字符串才挂断;中间识别、助手回复、开场白、TTS 均不能触发;重复结果不可反复终结。同一任务 revision 不同内容拒绝准入;provider 禁用/角色不符不可调用。Mock 参数验证不等于真实供应商验收。 diff --git a/internal/agent/recording_delivery_test.go b/internal/agent/recording_delivery_test.go index 08f78c3..5d7eb8e 100644 --- a/internal/agent/recording_delivery_test.go +++ b/internal/agent/recording_delivery_test.go @@ -72,7 +72,7 @@ func testRecordingDelivery(stub *recordingDeliveryRPC) *RecordingDelivery { func TestRecordingDeliveryNoRecordingReportsOnlyAfterConfirmedEnd(t *testing.T) { stub := &recordingDeliveryRPC{} delivery := testRecordingDelivery(stub) - payload := []byte(`{"task_id":"task-asr","caller_profile_id":"caller-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:30:00Z","ended_at":"2026-09-21T01:30:05Z","duration_ms":5000,"outcome":"no_answer","reason_code":480,"reason_message":"no answer","status_line":null,"raw":null,"sip_capture_error":null,"transcript":[],"opt_out":false,"recording":{}}`) + payload := []byte(`{"task_id":"task-asr","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:30:00Z","ended_at":"2026-09-21T01:30:05Z","duration_ms":5000,"outcome":"no_answer","reason_code":480,"reason_message":"no answer","status_line":null,"raw":null,"sip_capture_error":null,"transcript":[],"opt_out":false,"recording":{}}`) if err := delivery.Complete(context.Background(), CompletedRecording{ResultPayload: payload}); err != nil { t.Fatal(err) } @@ -209,7 +209,7 @@ func testDirectRecordingDelivery(t *testing.T, putStatus int) (*RecordingDeliver } delivery.Recovery = &RecordingRecovery{Root: root, Upload: UploadClient{AllowInsecureHTTP: true}} completed := CompletedRecording{ - ResultPayload: []byte(`{"task_id":"task-asr","caller_profile_id":"caller-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:30:00Z","ended_at":"2026-09-21T01:30:20Z","duration_ms":20000,"outcome":"answered","reason_code":null,"reason_message":"answered","status_line":null,"raw":null,"sip_capture_error":null,"transcript":[],"opt_out":false,"recording":{}}`), + ResultPayload: []byte(`{"task_id":"task-asr","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:30:00Z","ended_at":"2026-09-21T01:30:20Z","duration_ms":20000,"outcome":"answered","reason_code":null,"reason_message":"answered","status_line":null,"raw":null,"sip_capture_error":null,"transcript":[],"opt_out":false,"recording":{}}`), Expected: true, RecordingID: "recording-mock", UploadID: "upload-mock", WAV: wav, DurationMS: durationMS, } return delivery, stub, completed, &puts, root diff --git a/internal/asterisk/loader.go b/internal/asterisk/loader.go index 8323c9e..32316e3 100644 --- a/internal/asterisk/loader.go +++ b/internal/asterisk/loader.go @@ -28,9 +28,10 @@ type Loader struct { } type loadedTrunk struct { - ID string `json:"id"` - Host string `json:"host"` - Port int `json:"port"` + ID string `json:"id"` + Host string `json:"host"` + Port int `json:"port"` + CallerID string `json:"caller_id"` } type appliedState struct { @@ -42,6 +43,16 @@ type appliedState struct { var endpointLine = regexp.MustCompile(`(?m)^\s*Endpoint:\s+(\S+)`) +func endpointOption(output []byte, name string) string { + for line := range strings.SplitSeq(string(output), "\n") { + key, value, ok := strings.Cut(line, ":") + if ok && strings.EqualFold(strings.TrimSpace(key), name) { + return strings.TrimSpace(value) + } + } + return "" +} + func (l Loader) paths() (string, string, error) { if l.ConfigDir == "" || l.Asterisk == "" || l.LibraryDir == "" { return "", "", errors.New("native Asterisk executable, library and config paths are required") @@ -115,6 +126,9 @@ func (l Loader) observe(ctx context.Context, state appliedState) error { if err != nil || !bytes.Contains(endpoint, []byte(t.ID+"-aor")) || !bytes.Contains(endpoint, []byte("alaw")) || !bytes.Contains(endpoint, []byte("go-sip-udp")) || !bytes.Contains(endpoint, []byte("go-sip-no-inbound")) { return fmt.Errorf("SIP endpoint %q did not load its approved AOR/codec", t.ID) } + if endpointOption(endpoint, "from_user") != t.CallerID || endpointOption(endpoint, "callerid") != fmt.Sprintf("\"%s\" <%s>", t.CallerID, t.CallerID) { + return fmt.Errorf("SIP endpoint %q did not load its approved caller identity", t.ID) + } aor, err := l.cli(ctx, "pjsip show aor "+t.ID+"-aor") if err != nil || !bytes.Contains(aor, []byte(fmt.Sprintf("sip:%s:%d", t.Host, t.Port))) { return fmt.Errorf("SIP AOR %q did not load its approved contact", t.ID) @@ -173,10 +187,11 @@ func (l Loader) Apply(ctx context.Context, approvedJSON []byte) (map[string]int6 return nil, err } var trunks []struct { - ID string `json:"trunk_id"` - Host string `json:"server_host"` - Port int `json:"server_port"` - Enabled bool `json:"enabled"` + ID string `json:"trunk_id"` + Host string `json:"server_host"` + Port int `json:"server_port"` + Enabled bool `json:"enabled"` + CallerID string `json:"caller_id"` } if err := json.Unmarshal(sip.Trunks, &trunks); err != nil { return nil, err @@ -188,7 +203,7 @@ func (l Loader) Apply(ctx context.Context, approvedJSON []byte) (map[string]int6 state := appliedState{Revision: sip.Revision, Hash: fmt.Sprintf("%x", sha256.Sum256([]byte(text))), SnapshotHash: fmt.Sprintf("%x", sha256.Sum256(canonical))} for _, t := range trunks { if t.Enabled { - state.Trunks = append(state.Trunks, loadedTrunk{ID: t.ID, Host: t.Host, Port: t.Port}) + state.Trunks = append(state.Trunks, loadedTrunk{ID: t.ID, Host: t.Host, Port: t.Port, CallerID: t.CallerID}) } } sort.Slice(state.Trunks, func(i, j int) bool { return state.Trunks[i].ID < state.Trunks[j].ID }) diff --git a/internal/asterisk/loader_test.go b/internal/asterisk/loader_test.go index cd98d00..0701d48 100644 --- a/internal/asterisk/loader_test.go +++ b/internal/asterisk/loader_test.go @@ -19,7 +19,7 @@ func TestLoaderPersistsOnlyVerifiedNativeReload(t *testing.T) { case "$*" in *"pjsip show transports"*) echo 'Transport: go-sip-udp udp 0 0 0.0.0.0:5060' ;; *"pjsip show endpoints"*) echo 'Endpoint: '; echo 'Endpoint: trunk-shuqi Not in use' ;; - *"pjsip show endpoint trunk-shuqi"*) echo 'Aor: trunk-shuqi-aor allow: alaw transport: go-sip-udp context: go-sip-no-inbound' ;; + *"pjsip show endpoint trunk-shuqi"*) echo 'Aor: trunk-shuqi-aor allow: alaw transport: go-sip-udp context: go-sip-no-inbound'; echo 'from_user : BD93205882'; echo 'callerid : "BD93205882" ' ;; *"pjsip show aor trunk-shuqi-aor"*) echo 'Contact: sip:61.132.228.221:5060' ;; *) exit 9 ;; esac @@ -66,6 +66,21 @@ esac if err := os.WriteFile(fake, []byte(script), 0700); err != nil { t.Fatal(err) } + if err := os.WriteFile(fake, []byte(strings.Replace(script, "from_user : BD93205882", "from_user : anonymous", 1)), 0700); err != nil { + t.Fatal(err) + } + if _, err := loader.LoadedSIP(context.Background()); err == nil { + t.Fatal("accepted Asterisk endpoint with anonymous outbound From user") + } + if err := os.WriteFile(fake, []byte(strings.Replace(script, "callerid : \"BD93205882\" ", "callerid : ", 1)), 0700); err != nil { + t.Fatal(err) + } + if _, err := loader.LoadedSIP(context.Background()); err == nil { + t.Fatal("accepted Asterisk endpoint without approved caller ID") + } + if err := os.WriteFile(fake, []byte(script), 0700); err != nil { + t.Fatal(err) + } changed := testSIP(t) changed.Trunks = []byte(strings.Replace(string(changed.Trunks), "BD93205882", "BD93205883", 1)) changedBody, err := json.Marshal(changed) diff --git a/internal/asterisk/sip.go b/internal/asterisk/sip.go index c33b1f4..6bea573 100644 --- a/internal/asterisk/sip.go +++ b/internal/asterisk/sip.go @@ -23,6 +23,7 @@ type trunk struct { AuthMode *string `json:"auth_mode"` RegistrationRequired *bool `json:"registration_required"` Codec string `json:"codec"` + CallerID string `json:"caller_id"` Enabled bool `json:"enabled"` } @@ -64,7 +65,7 @@ func Render(sip configread.SIP) (string, error) { if !t.Enabled { continue } - _, _ = fmt.Fprintf(&text, "\n[%s]\ntype=endpoint\ntransport=go-sip-udp\ncontext=go-sip-no-inbound\ndisallow=all\nallow=alaw\naors=%s-aor\ndirect_media=no\n\n[%s-aor]\ntype=aor\ncontact=sip:%s:%d\n", t.ID, t.ID, t.ID, t.Host, t.Port) + _, _ = fmt.Fprintf(&text, "\n[%s]\ntype=endpoint\ntransport=go-sip-udp\ncontext=go-sip-no-inbound\ndisallow=all\nallow=alaw\naors=%s-aor\nfrom_user=%s\ncallerid=\"%s\" <%s>\ndirect_media=no\n\n[%s-aor]\ntype=aor\ncontact=sip:%s:%d\n", t.ID, t.ID, t.CallerID, t.CallerID, t.CallerID, t.ID, t.Host, t.Port) } if text.Len() == 0 { return "", errors.New("empty generated SIP configuration") diff --git a/internal/asterisk/sip_test.go b/internal/asterisk/sip_test.go index 90bc8a4..f4e57c7 100644 --- a/internal/asterisk/sip_test.go +++ b/internal/asterisk/sip_test.go @@ -11,7 +11,7 @@ import ( func testSIP(t *testing.T) configread.SIP { t.Helper() var sip configread.SIP - if err := json.Unmarshal([]byte(`{"resource":"sip_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","revision":9,"trunks":[{"trunk_id":"trunk-shuqi","provider_id":"shuqi","codec":"PCMA","dial_prefix":"7089","enabled":true,"server_host":"61.132.228.221","server_port":5060,"transport":"udp","auth_mode":"ip","registration_required":false,"max_concurrent_calls":1,"caller_profiles":[{"caller_profile_id":"caller-shuqi","caller_id":"BD93205882"}],"schedule":{"time_zone":"Asia/Shanghai","weekly_windows":{"monday":[{"start":"09:00","end":"20:00"}],"tuesday":[],"wednesday":[],"thursday":[],"friday":[],"saturday":[],"sunday":[]}}}]}`), &sip); err != nil { + if err := json.Unmarshal([]byte(`{"resource":"sip_config","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","revision":9,"trunks":[{"trunk_id":"trunk-shuqi","provider_id":"shuqi","codec":"PCMA","dial_prefix":"7089","enabled":true,"server_host":"61.132.228.221","server_port":5060,"transport":"udp","auth_mode":"ip","registration_required":false,"max_concurrent_calls":1,"caller_id":"BD93205882","schedule":{"time_zone":"Asia/Shanghai","weekly_windows":{"monday":[{"start":"09:00","end":"20:00"}],"tuesday":[],"wednesday":[],"thursday":[],"friday":[],"saturday":[],"sunday":[]}}}]}`), &sip); err != nil { t.Fatal(err) } return sip @@ -22,13 +22,41 @@ func TestRenderApprovedIPTrunk(t *testing.T) { if err != nil { t.Fatal(err) } - for _, want := range []string{"[trunk-shuqi]", "type=endpoint", "transport=go-sip-udp", "allow=alaw", "aors=trunk-shuqi-aor", "contact=sip:61.132.228.221:5060"} { + for _, want := range []string{"[trunk-shuqi]", "type=endpoint", "transport=go-sip-udp", "allow=alaw", "aors=trunk-shuqi-aor", "contact=sip:61.132.228.221:5060", "from_user=BD93205882", "callerid=\"BD93205882\" "} { if !strings.Contains(text, want) { t.Errorf("rendered SIP is missing %q", want) } } - if strings.Contains(text, "type=transport") || strings.Contains(text, "7089") || strings.Contains(text, "BD93205882") || strings.Contains(text, "register=") { - t.Fatal("per-call prefix/caller or registration leaked into static endpoint") + if strings.Contains(text, "type=transport") || strings.Contains(text, "7089") || strings.Contains(text, "register=") { + t.Fatal("per-call prefix or registration leaked into static endpoint") + } +} + +func TestRenderKeepsDistinctTrunkCallerIdentities(t *testing.T) { + sip := testSIP(t) + var trunks []map[string]any + if err := json.Unmarshal(sip.Trunks, &trunks); err != nil { + t.Fatal(err) + } + second := make(map[string]any, len(trunks[0])) + for key, value := range trunks[0] { + second[key] = value + } + second["trunk_id"], second["caller_id"] = "trunk-zhongding", "mbkq" + trunks = append(trunks, second) + var err error + sip.Trunks, err = json.Marshal(trunks) + if err != nil { + t.Fatal(err) + } + text, err := Render(sip) + if err != nil { + t.Fatal(err) + } + for _, want := range []string{"[trunk-shuqi]\ntype=endpoint", "from_user=BD93205882\ncallerid=\"BD93205882\" ", "[trunk-zhongding]\ntype=endpoint", "from_user=mbkq\ncallerid=\"mbkq\" "} { + if !strings.Contains(text, want) { + t.Fatalf("missing separately approved line identity %q in %s", want, text) + } } } @@ -41,6 +69,8 @@ func TestRenderRejectsUnsupportedOrAmbiguousSIP(t *testing.T) { {"missing_auth", `"auth_mode":"ip"`, `"auth_mode":null`}, {"bad_id", `"trunk_id":"trunk-shuqi"`, `"trunk_id":"bad]\n[attacker"`}, {"bad_host", `"server_host":"61.132.228.221"`, `"server_host":"oops\npassword=bad"`}, + {"missing_caller", `"caller_id":"BD93205882"`, `"caller_id":""`}, + {"unsafe_caller", `"caller_id":"BD93205882"`, `"caller_id":"BD93205882\nfrom_user=anonymous"`}, } { t.Run(mutation.name, func(t *testing.T) { sip := testSIP(t) diff --git a/internal/callflow/result_payload.go b/internal/callflow/result_payload.go index a501ad3..87bb1c5 100644 --- a/internal/callflow/result_payload.go +++ b/internal/callflow/result_payload.go @@ -14,7 +14,6 @@ import ( // from the Agent SDK. The reason code is nil if no actual SIP code is known. type FinalCallFacts struct { TaskID string - CallerProfileID string Callee string TrunkID string StartedAt time.Time @@ -31,7 +30,7 @@ type FinalCallFacts struct { // user-side ASR text with observed capture timing is included; unplayed LLM // replies, interim ASR and hypothetical recordings are never transcribed. func FinalResultPayload(facts FinalCallFacts, call Result) ([]byte, error) { - if strings.TrimSpace(facts.TaskID) == "" || strings.TrimSpace(facts.CallerProfileID) == "" || strings.TrimSpace(facts.Callee) == "" || strings.TrimSpace(facts.TrunkID) == "" || strings.TrimSpace(facts.ReasonMessage) == "" { + if strings.TrimSpace(facts.TaskID) == "" || strings.TrimSpace(facts.Callee) == "" || strings.TrimSpace(facts.TrunkID) == "" || strings.TrimSpace(facts.ReasonMessage) == "" { return nil, errors.New("final call result is missing approved identity or observed reason") } if facts.StartedAt.IsZero() || facts.EndedAt.IsZero() || facts.EndedAt.Before(facts.StartedAt) { @@ -102,7 +101,6 @@ func FinalResultPayload(facts FinalCallFacts, call Result) ([]byte, error) { } return json.Marshal(struct { TaskID string `json:"task_id"` - CallerProfileID string `json:"caller_profile_id"` Callee string `json:"callee"` TrunkID string `json:"trunk_id"` StartedAt string `json:"started_at"` @@ -118,7 +116,7 @@ func FinalResultPayload(facts FinalCallFacts, call Result) ([]byte, error) { OptOut bool `json:"opt_out"` Recording map[string]any `json:"recording"` }{ - TaskID: facts.TaskID, CallerProfileID: facts.CallerProfileID, Callee: facts.Callee, TrunkID: facts.TrunkID, + TaskID: facts.TaskID, Callee: facts.Callee, TrunkID: facts.TrunkID, StartedAt: facts.StartedAt.UTC().Format(time.RFC3339Nano), EndedAt: facts.EndedAt.UTC().Format(time.RFC3339Nano), DurationMS: callDuration, Outcome: facts.Outcome, ReasonCode: facts.ReasonCode, ReasonMessage: facts.ReasonMessage, StatusLine: statusLine, Raw: raw, SIPCaptureError: captureError, diff --git a/internal/callflow/result_payload_test.go b/internal/callflow/result_payload_test.go index 2b9af4e..e5d9011 100644 --- a/internal/callflow/result_payload_test.go +++ b/internal/callflow/result_payload_test.go @@ -12,7 +12,7 @@ import ( func approvedResultFixture() (FinalCallFacts, Result) { start := time.Date(2026, 9, 21, 1, 0, 0, 0, time.UTC) facts := FinalCallFacts{ - TaskID: "task-asr", CallerProfileID: "caller-profile-mock", Callee: "15003164745", TrunkID: "trunk-mock", + TaskID: "task-asr", Callee: "15003164745", TrunkID: "trunk-mock", StartedAt: start, EndedAt: start.Add(3 * time.Second), Outcome: "answered", ReasonMessage: "completed in isolated mock", } result := Result{ @@ -46,14 +46,13 @@ func TestFinalResultPayloadUsesOnlyObservedUserASRAndMediaTimes(t *testing.T) { t.Fatalf("locally produced result violates the sole current MQ schema: %v", err) } var parsed struct { - TaskID string `json:"task_id"` - CallerProfileID string `json:"caller_profile_id"` - DurationMS int64 `json:"duration_ms"` - Outcome string `json:"outcome"` - ReasonCode *int `json:"reason_code"` - OptOut bool `json:"opt_out"` - Recording map[string]any `json:"recording"` - Transcript []struct { + TaskID string `json:"task_id"` + DurationMS int64 `json:"duration_ms"` + Outcome string `json:"outcome"` + ReasonCode *int `json:"reason_code"` + OptOut bool `json:"opt_out"` + Recording map[string]any `json:"recording"` + Transcript []struct { SegmentID string `json:"segment_id"` TurnID string `json:"turn_id"` Role string `json:"role"` @@ -65,7 +64,7 @@ func TestFinalResultPayloadUsesOnlyObservedUserASRAndMediaTimes(t *testing.T) { if err := json.Unmarshal(payload, &parsed); err != nil { t.Fatal(err) } - if parsed.TaskID != facts.TaskID || parsed.CallerProfileID != facts.CallerProfileID || parsed.DurationMS != 3000 || parsed.Outcome != "answered" || parsed.ReasonCode != nil || !parsed.OptOut || len(parsed.Recording) != 0 || len(parsed.Transcript) != 2 { + if parsed.TaskID != facts.TaskID || parsed.DurationMS != 3000 || parsed.Outcome != "answered" || parsed.ReasonCode != nil || !parsed.OptOut || len(parsed.Recording) != 0 || len(parsed.Transcript) != 2 { t.Fatalf("result changed approved facts or invented a recording: %+v", parsed) } if first, second := parsed.Transcript[0], parsed.Transcript[1]; first.SegmentID != "segment-1" || first.TurnID != "turn-1" || first.Role != "user" || first.StartMS != 100 || first.EndMS != 1200 || first.Text != result.Turns[0].Transcript || second.SegmentID != "segment-2" || second.TurnID != "turn-2" || second.StartMS != 1600 || second.EndMS != 2250 || second.Text != result.Turns[1].Transcript { @@ -121,7 +120,6 @@ func TestFinalResultPayloadRejectsMissingOrInventedCallFacts(t *testing.T) { name string change func(*FinalCallFacts, *Result) }{ - {"missing_approved_caller", func(f *FinalCallFacts, _ *Result) { f.CallerProfileID = "" }}, {"unapproved_outcome", func(f *FinalCallFacts, _ *Result) { f.Outcome = "new_state" }}, {"no_sip_reason", func(f *FinalCallFacts, _ *Result) { f.ReasonMessage = "" }}, {"ended_before_started", func(f *FinalCallFacts, _ *Result) { f.EndedAt = f.StartedAt.Add(-time.Second) }}, diff --git a/internal/configread/snapshots.go b/internal/configread/snapshots.go index b52a3de..c435d54 100644 --- a/internal/configread/snapshots.go +++ b/internal/configread/snapshots.go @@ -79,7 +79,6 @@ type Task struct { RingTimeoutMS int64 `json:"ring_timeout_ms"` MaxCallDurationMS int64 `json:"max_call_duration_ms"` RoutePolicyID string `json:"route_policy_id"` - CallerProfileID string `json:"caller_profile_id"` AllowedTrunkIDs []string `json:"allowed_trunk_ids"` Schedule json.RawMessage `json:"schedule"` Agent Agent `json:"agent"` diff --git a/internal/contract/current_test.go b/internal/contract/current_test.go index b277068..f15d8d6 100644 --- a/internal/contract/current_test.go +++ b/internal/contract/current_test.go @@ -7,6 +7,28 @@ import ( "testing" ) +func TestRejectsRetiredCallerProfiles(t *testing.T) { + for _, tc := range []struct { + name, file, from, to string + }{ + {"config-read", "config-read-sip.json", `"caller_id":"BD00000000"`, `"caller_profiles":[{"caller_profile_id":"legacy","caller_id":"BD00000000"}]`}, + {"config-read", "config-read-task-asr.json", `"route_policy_id":"route-mock"`, `"route_policy_id":"route-mock","caller_profile_id":"legacy"`}, + {"mq", "mq-result-no-recording.json", `"task_id":"task-asr"`, `"task_id":"task-asr","caller_profile_id":"legacy"`}, + } { + raw, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", tc.file)) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(raw), tc.from) { + t.Fatalf("missing test fixture field in %s", tc.file) + } + changed := strings.Replace(string(raw), tc.from, tc.to, 1) + if err := ValidateCurrent(tc.name, []byte(changed)); err == nil { + t.Fatalf("retired caller profile accepted by %s", tc.file) + } + } +} + func TestValidateCurrent(t *testing.T) { for _, tc := range []struct { name string diff --git a/internal/dispatcher/execute.go b/internal/dispatcher/execute.go index d514bf6..4eb6cfb 100644 --- a/internal/dispatcher/execute.go +++ b/internal/dispatcher/execute.go @@ -23,7 +23,6 @@ type CallSpec struct { EventID string TenantID int64 TaskID string - CallerProfileID string CallerID string Callee string DialedCallee string @@ -215,7 +214,7 @@ func (c *ExecuteController) dispatchPending(ctx context.Context, cmd store.Execu } spec := CallSpec{ DispatcherID: cmd.DispatcherID, EventID: cmd.EventID, TenantID: cmd.TenantID, - TaskID: cmd.TaskID, CallerProfileID: snapshot.Task.CallerProfileID, + TaskID: cmd.TaskID, CallerID: choice.CallerID, Callee: cmd.Callee, DialedCallee: choice.DialedCallee, TrunkID: choice.TrunkID, RingTimeoutMS: snapshot.Task.RingTimeoutMS, MaxCallDurationMS: choice.MaxCallDurationMS, diff --git a/internal/dispatcher/policy.go b/internal/dispatcher/policy.go index 8055c35..e772ff6 100644 --- a/internal/dispatcher/policy.go +++ b/internal/dispatcher/policy.go @@ -35,15 +35,10 @@ type trunkConfig struct { AuthMode *string `json:"auth_mode"` RegistrationRequired *bool `json:"registration_required"` MaxConcurrentCalls *int64 `json:"max_concurrent_calls"` - CallerProfiles []callerProfile `json:"caller_profiles"` + CallerID string `json:"caller_id"` Schedule callwindow.WeeklySchedule `json:"schedule"` } -type callerProfile struct { - ID string `json:"caller_profile_id"` - CallerID string `json:"caller_id"` -} - // validRawCallee checks only what the native dial route can safely represent; // the SaaS call.execute event, not a local list, chooses the business number. func validRawCallee(callee string) bool { @@ -122,15 +117,8 @@ func SelectTrunk(snapshot configread.Snapshot, callee string, at time.Time, trun invalid = fmt.Errorf("%w: SIP trunk %q has unknown required execution fields", ErrRuleInvalid, id) continue } - var callerID string - for _, profile := range trunk.CallerProfiles { - if profile.ID == snapshot.Task.CallerProfileID { - callerID = profile.CallerID - break - } - } - if callerID == "" { - invalid = fmt.Errorf("%w: caller profile is missing on trunk %q", ErrRuleInvalid, id) + if trunk.CallerID == "" { + invalid = fmt.Errorf("%w: caller ID is missing on trunk %q", ErrRuleInvalid, id) continue } end, err := callwindow.Evaluate(taskSchedule, trunk.Schedule, at) @@ -152,7 +140,7 @@ func SelectTrunk(snapshot configread.Snapshot, callee string, at time.Time, trun continue } return SelectedTrunk{ - TrunkID: id, CallerID: callerID, Callee: callee, + TrunkID: id, CallerID: trunk.CallerID, Callee: callee, DialedCallee: trunk.DialPrefix + callee, Deadline: deadline, MaxCallDurationMS: maxDurationMS, }, nil diff --git a/internal/dispatcher/policy_test.go b/internal/dispatcher/policy_test.go index 852aebf..4faaf34 100644 --- a/internal/dispatcher/policy_test.go +++ b/internal/dispatcher/policy_test.go @@ -71,6 +71,35 @@ func TestSelectTrunkSaaSEventCalleeScheduleAndCaller(t *testing.T) { } } +func TestSelectTrunkUsesEachLinesOwnCaller(t *testing.T) { + snapshot := policySnapshot(t) + var trunks []trunkConfig + if err := json.Unmarshal(snapshot.SIP.Trunks, &trunks); err != nil { + t.Fatal(err) + } + second := trunks[0] + second.TrunkID, second.CallerID = "trunk-second", "mbkq" + trunks = append(trunks, second) + var err error + snapshot.SIP.Trunks, err = json.Marshal(trunks) + if err != nil { + t.Fatal(err) + } + snapshot.Task.AllowedTrunkIDs = []string{second.TrunkID} + chosen, err := SelectTrunk(snapshot, "15003164745", monday(9, 30), nil, map[string]int64{second.TrunkID: 8}) + if err != nil || chosen.TrunkID != second.TrunkID || chosen.CallerID != second.CallerID { + t.Fatalf("second line did not keep its own approved caller: %+v, %v", chosen, err) + } + trunks[1].CallerID = "" + snapshot.SIP.Trunks, err = json.Marshal(trunks) + if err != nil { + t.Fatal(err) + } + if _, err := SelectTrunk(snapshot, "15003164745", monday(9, 30), nil, map[string]int64{second.TrunkID: 8}); !errors.Is(err, ErrRuleInvalid) { + t.Fatalf("second line with missing caller did not fail closed: %v", err) + } +} + func TestSelectTrunkFailsClosedForUnknownLineAndLimits(t *testing.T) { snapshot := policySnapshot(t) at := monday(9, 30) diff --git a/internal/rpc/approved_execution.go b/internal/rpc/approved_execution.go index 13d3fb2..8c2e7ab 100644 --- a/internal/rpc/approved_execution.go +++ b/internal/rpc/approved_execution.go @@ -28,7 +28,6 @@ type ApprovedExecution struct { DispatcherID string TenantID int64 TaskID string - CallerProfileID string SourceEventID string CallID string SelectedTrunkID string @@ -152,7 +151,7 @@ func (s *Server) ExecuteApproved(ctx context.Context, req *agentpb.ExecuteApprov } approved := ApprovedExecution{ DispatcherID: dispatcherID, TenantID: req.TenantId, TaskID: req.TaskId, - CallerProfileID: task.CallerProfileID, SourceEventID: req.SourceEventId, CallID: req.CallId, + SourceEventID: req.SourceEventId, CallID: req.CallId, SelectedTrunkID: req.SelectedTrunkId, CallerID: req.CallerId, Callee: req.Callee, DialedCallee: req.DialedCallee, SIPRevision: req.SipRevision, RingTimeout: time.Duration(req.RingTimeoutMs) * time.Millisecond, diff --git a/internal/rpc/approved_recorded_mock.go b/internal/rpc/approved_recorded_mock.go index bd60c98..9464428 100644 --- a/internal/rpc/approved_recorded_mock.go +++ b/internal/rpc/approved_recorded_mock.go @@ -38,7 +38,7 @@ func (r *ApprovedRecordedMockCall) Prepare(approved ApprovedExecution) (func(con return nil, ErrApprovedMockDeliveryUnavailable } if approved.DispatcherID == "" || approved.TenantID <= 0 || approved.SourceEventID == "" || approved.CallID != approved.SourceEventID || - approved.TaskID == "" || approved.CallerProfileID == "" || approved.Callee == "" || approved.SelectedTrunkID == "" || r.MaxWAVBytes <= 44 || strings.TrimSpace(r.ReasonMessage) == "" { + approved.TaskID == "" || approved.Callee == "" || approved.SelectedTrunkID == "" || r.MaxWAVBytes <= 44 || strings.TrimSpace(r.ReasonMessage) == "" { return nil, errors.New("approved Mock call identity, media limit or observed reason is incomplete") } if r.Outcome != "answered" && r.Outcome != "no_answer" && r.Outcome != "failed" { @@ -81,7 +81,7 @@ func (r *ApprovedRecordedMockCall) Prepare(approved ApprovedExecution) (func(con resultOutcome, resultReason = "failed", "isolated Mock media did not complete" } payload, err := callflow.FinalResultPayload(callflow.FinalCallFacts{ - TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID, + TaskID: approved.TaskID, Callee: approved.Callee, TrunkID: approved.SelectedTrunkID, StartedAt: startedAt, EndedAt: endedAt, Outcome: resultOutcome, ReasonMessage: resultReason, }, observed) diff --git a/internal/rpc/approved_recorded_mock_test.go b/internal/rpc/approved_recorded_mock_test.go index 5c23139..3cc58dd 100644 --- a/internal/rpc/approved_recorded_mock_test.go +++ b/internal/rpc/approved_recorded_mock_test.go @@ -79,7 +79,7 @@ func recordedMockFixture(t *testing.T, pcm []byte, maxDuration time.Duration) (* stub := &recordedMockRPC{targetURL: localOSS.URL + "/object"} approved := ApprovedExecution{ DispatcherID: "c046b893-8628-4589-ae50-619d049248a6", TenantID: 42, - TaskID: "task-asr", CallerProfileID: "caller-mock", SourceEventID: "event-fixture", CallID: "event-fixture", + TaskID: "task-asr", SourceEventID: "event-fixture", CallID: "event-fixture", Callee: "15003164745", SelectedTrunkID: "trunk-mock", MaxCallDuration: maxDuration, DialBefore: time.Now().Add(time.Minute), AI: ai.Binding{ Mode: "asr_only", Conversation: ai.ConversationConfig{SilenceTimeout: 100 * time.Millisecond, MaxDuration: maxDuration}, diff --git a/internal/rpc/approved_recorded_real.go b/internal/rpc/approved_recorded_real.go index 5e3cb6d..0f55809 100644 --- a/internal/rpc/approved_recorded_real.go +++ b/internal/rpc/approved_recorded_real.go @@ -44,7 +44,7 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con return nil, errors.New("real call requires bounded media and per-call delivery") } if approved.DispatcherID == "" || approved.TenantID <= 0 || approved.SourceEventID == "" || approved.CallID != approved.SourceEventID || approved.TaskID == "" || - approved.CallerProfileID == "" || approved.CallerID == "" || approved.Callee == "" || approved.DialedCallee == "" || approved.SelectedTrunkID == "" || + approved.CallerID == "" || approved.Callee == "" || approved.DialedCallee == "" || approved.SelectedTrunkID == "" || approved.RingTimeout <= 0 || approved.MaxCallDuration <= 0 || approved.DialBefore.IsZero() || (approved.AI.Mode != string(ai.ModeASROnly) && approved.AI.Mode != string(ai.ModeFullAI)) { return nil, errors.New("approved real call identity, timing or AI mode is incomplete") } @@ -115,7 +115,7 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con return err // uncertain origination or end is never retried or reported as completed } facts := callflow.FinalCallFacts{ - TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID, + TaskID: approved.TaskID, Callee: approved.Callee, TrunkID: approved.SelectedTrunkID, StartedAt: attemptedAt, EndedAt: time.Now().UTC(), Outcome: "failed", ReasonMessage: "call ended before answer", } @@ -140,7 +140,7 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con return errors.Join(cause, err) } facts := callflow.FinalCallFacts{ - TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID, + TaskID: approved.TaskID, Callee: approved.Callee, TrunkID: approved.SelectedTrunkID, StartedAt: startedAt, EndedAt: time.Now().UTC(), Outcome: "failed", ReasonMessage: reason, } @@ -178,7 +178,7 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con outcome, reason = "failed", "real AI or media flow failed" } facts := callflow.FinalCallFacts{ - TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID, + TaskID: approved.TaskID, Callee: approved.Callee, TrunkID: approved.SelectedTrunkID, StartedAt: startedAt, EndedAt: endedAt, Outcome: outcome, ReasonMessage: reason, } diff --git a/internal/rpc/approved_snapshot_test.go b/internal/rpc/approved_snapshot_test.go index 5c719c1..5659321 100644 --- a/internal/rpc/approved_snapshot_test.go +++ b/internal/rpc/approved_snapshot_test.go @@ -21,7 +21,7 @@ func TestApprovedAgentReceivesFrozenASROnlySnapshot(t *testing.T) { s := activatedApprovedServer(t, now, filepath.Join(t.TempDir(), "session.json"), req.DispatcherId, 1, func(_ context.Context, call ApprovedExecution) error { attempts++ - if call.DispatcherID != req.DispatcherId || call.TenantID != req.TenantId || call.TaskID != req.TaskId || call.SIPRevision != 8 || call.SelectedTrunkID != req.SelectedTrunkId || call.CallerProfileID != "caller-profile-mock" { + if call.DispatcherID != req.DispatcherId || call.TenantID != req.TenantId || call.TaskID != req.TaskId || call.SIPRevision != 8 || call.SelectedTrunkID != req.SelectedTrunkId { t.Fatal("Agent did not receive Dispatcher-authorized call binding") } if call.AI.Mode != "asr_only" || call.AI.LLM != nil || call.AI.TTS != nil || call.AI.ASR.Provider.Credential != "example-only-not-a-real-secret" { diff --git a/internal/rpc/recording_delivery_integration_test.go b/internal/rpc/recording_delivery_integration_test.go index d589586..a069253 100644 --- a/internal/rpc/recording_delivery_integration_test.go +++ b/internal/rpc/recording_delivery_integration_test.go @@ -39,7 +39,7 @@ func TestRecordingDeliveryRealMutualTLSAndLocalOSSCommitsOneSQLiteResult(t *test } startedAt := time.Now().UTC() callResult, err := RunApprovedCall(context.Background(), ApprovedExecution{ - AI: bound, MaxCallDuration: time.Second, TaskID: snapshot.Task.TaskID, CallerProfileID: snapshot.Task.CallerProfileID, + AI: bound, MaxCallDuration: time.Second, TaskID: snapshot.Task.TaskID, }, mediaSession, hangup, mock) endedAt := time.Now().UTC() if err != nil || len(callResult.Turns) != 1 || len(callResult.OutboundTurns) != 0 || callResult.Turns[0].Transcript != "隔离 Mock 最终识别文本" { @@ -104,7 +104,7 @@ func TestRecordingDeliveryRealMutualTLSAndLocalOSSCommitsOneSQLiteResult(t *test } payload, err := callflow.FinalResultPayload(callflow.FinalCallFacts{ - TaskID: snapshot.Task.TaskID, CallerProfileID: snapshot.Task.CallerProfileID, + TaskID: snapshot.Task.TaskID, Callee: "15003164745", TrunkID: "trunk-mock", StartedAt: startedAt, EndedAt: endedAt, Outcome: "answered", ReasonMessage: "isolated Mock media completed", }, callResult) diff --git a/internal/rpc/recording_server_flow_test.go b/internal/rpc/recording_server_flow_test.go index 7efe8f1..eacff2b 100644 --- a/internal/rpc/recording_server_flow_test.go +++ b/internal/rpc/recording_server_flow_test.go @@ -108,8 +108,8 @@ func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.Store, context. func recordingResultPayload(t *testing.T, task configread.Snapshot, grant *agentpb.UploadGrant) []byte { t.Helper() result := map[string]any{ - "task_id": "task-asr", "caller_profile_id": task.Task.CallerProfileID, - "callee": "15003164745", "trunk_id": "trunk-mock", + "task_id": "task-asr", + "callee": "15003164745", "trunk_id": "trunk-mock", "started_at": "2026-09-21T01:30:01Z", "ended_at": "2026-09-21T01:30:06Z", "duration_ms": 5000, "outcome": "no_answer", "reason_code": 486, "reason_message": "busy", "status_line": nil, "raw": nil, "sip_capture_error": nil, diff --git a/internal/store/result.go b/internal/store/result.go index 7c0ed38..d815808 100644 --- a/internal/store/result.go +++ b/internal/store/result.go @@ -42,13 +42,12 @@ func (s *Store) recordCallResult(dispatcherID, sourceEventID string, payload []b return OutboxEvent{}, false, fmt.Errorf("%w: durable call identity and payload are required", ErrResultInvalid) } var result struct { - TaskID string `json:"task_id"` - CallerProfileID string `json:"caller_profile_id"` - Callee string `json:"callee"` - TrunkID string `json:"trunk_id"` - StartedAt string `json:"started_at"` - EndedAt string `json:"ended_at"` - Recording json.RawMessage `json:"recording"` + TaskID string `json:"task_id"` + Callee string `json:"callee"` + TrunkID string `json:"trunk_id"` + StartedAt string `json:"started_at"` + EndedAt string `json:"ended_at"` + Recording json.RawMessage `json:"recording"` } if err := json.Unmarshal(payload, &result); err != nil { return OutboxEvent{}, false, fmt.Errorf("%w: decode JSON: %v", ErrResultInvalid, err) @@ -85,7 +84,7 @@ func (s *Store) recordCallResult(dispatcherID, sourceEventID string, payload []b if err := json.Unmarshal(snapshotJSON, &snapshot); err != nil { return OutboxEvent{}, false, fmt.Errorf("decode frozen call task: %w", err) } - if snapshot.Task.DispatcherID != dispatcherID || snapshot.Task.TenantID != tenantID || snapshot.Task.TaskID != taskID || result.TaskID != taskID || result.Callee != callee || result.TrunkID != trunkID || result.CallerProfileID != snapshot.Task.CallerProfileID { + if snapshot.Task.DispatcherID != dispatcherID || snapshot.Task.TenantID != tenantID || snapshot.Task.TaskID != taskID || result.TaskID != taskID || result.Callee != callee || result.TrunkID != trunkID { return OutboxEvent{}, false, errors.New("final result identity differs from the approved call") } var recording struct { diff --git a/internal/store/result_test.go b/internal/store/result_test.go index 0bc76c5..2dc6fc6 100644 --- a/internal/store/result_test.go +++ b/internal/store/result_test.go @@ -15,8 +15,8 @@ import ( func currentResultPayload(t *testing.T) []byte { t.Helper() payload := map[string]any{ - "task_id": "task-asr", "caller_profile_id": currentStoreSnapshot(t).Task.CallerProfileID, - "callee": "15003164745", "trunk_id": "trunk-mock", + "task_id": "task-asr", + "callee": "15003164745", "trunk_id": "trunk-mock", "started_at": "2026-09-21T01:30:01Z", "ended_at": "2026-09-21T01:30:06Z", "duration_ms": 5000, "outcome": "no_answer", "reason_code": 486, "reason_message": "busy", "status_line": nil, "raw": nil, "sip_capture_error": nil, @@ -133,7 +133,7 @@ func TestFinalResultRejectsDifferentBoundTaskCallOrTrunk(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "result-bound-1", true) for _, change := range []struct{ field, value string }{ - {"call_id", "other-call"}, {"task_id", "other-task"}, {"callee", "15830461047"}, {"trunk_id", "other-trunk"}, {"caller_profile_id", "other-caller"}, + {"call_id", "other-call"}, {"task_id", "other-task"}, {"callee", "15830461047"}, {"trunk_id", "other-trunk"}, } { var claim map[string]any if err := json.Unmarshal(currentResultPayload(t), &claim); err != nil {