chore: add FreeSWITCH/ESL research (deploy + mock + docs)
This commit is contained in:
@@ -0,0 +1,177 @@
|
||||
# FreeSWITCH/ESL AI 外呼底座 — 部署与运维文档
|
||||
|
||||
> 目标机器: `39.106.106.246` (阿里云 2vCPU/3.5G, Debian 13) · 交付物: `/opt/freeswitch-esl/`
|
||||
> 项目: https://github.com/signalwire/freeswitch · 镜像: safarov/freeswitch (FreeSWITCH 1.10.12)
|
||||
> 范围: **仅外呼** (ESL originate ×2 + uuid_bridge), Mock SIP 运营商 + Mock LLM/ASR, 已端到端验证双向语音。
|
||||
> 前置: asterisk 运行时已下线 (ast-* 容器清零、8088/8089 释放); Docker 与内核调优沿用 livekit 期成果。
|
||||
|
||||
## 1. 水平扩展结论 (先行)
|
||||
|
||||
**结论: 零共享状态的独立节点复制 —— 与 Asterisk 同构。FreeSWITCH 无内置集群, 每节点呼叫状态全在本机内存; 扩容 = 加一个节点, 呼叫发起侧 (agent/业务后端) 在节点名单上选节点建 ESL 连接。不需要 Redis, 不需要任何协调组件。**
|
||||
|
||||
| 层 | 状态存放 | 扩容方式 | 调度机制 |
|
||||
|---|---|---|---|
|
||||
| FreeSWITCH 节点 | 本机内存 (通道/桥) | 复制 service 块 (独立宿主 ESL 端口) | 呼叫发起侧选节点 (call.sh 的 NODE 参数 = 生产中的 LB 名单) |
|
||||
| AI Agent (生产) | 无状态 worker | `--scale`/K8s HPA | 每个外呼任务自带节点选择与 ESL/SIP/RTP 会话 |
|
||||
|
||||
三方案对比: livekit-sip 需 Redis 存 trunk/呼叫状态并派发; Asterisk (ARI) 与 FreeSWITCH (ESL) 出站外呼均零共享层, 每呼叫自包含 (originate ×2 + bridge 全在一个节点内闭环), 节点间零交互。实测 1→2 节点**零配置变更** (共用同一套 etc-fs), 呼叫按发起侧选择分摊 (见 §7)。
|
||||
|
||||
瓶颈与边界: ① 单节点容量受 CPU 与 RTP 端口段限制 — vanilla 默认 16384-32768 (约 1.6 万端口, 每桥接呼叫 2 腿各 1 对 RTP+RTCP ≈ 8000 并发/节点, 容器 fd/CPU 先到顶); ② ESL 是单连接同步命令式接口, 高并发控制面需连接池或 bgapi; ③ 无跨节点媒体互通 — 仅外呼场景无影响; ④ 未来要呼入/注册时前置 Kamailio dispatcher, 纯出站不需要。
|
||||
|
||||
## 2. 架构
|
||||
|
||||
```
|
||||
业务后端 / AI Agent (mock-agent, 容器, 纯 stdlib Python)
|
||||
│ ① ESL TCP 连接 fs:8021 (auth) — 控制面 (等价 ARI REST)
|
||||
│ ② 本侧 SIP UAS 监听 :5062 + RTP :40001 — 媒体面 (等价 ARI externalMedia)
|
||||
│ ③ originate 腿1 → sofia/external/agent@me:5062 &park()
|
||||
│ ④ originate 腿2 → sofia/gateway/mock-trunk/<被叫> &park()
|
||||
│ ⑤ uuid_bridge 腿1 腿2 → RTP 经 FreeSWITCH 双向
|
||||
▼
|
||||
┌─ docker net: fsnet ─────────────────────────────────────────────┐
|
||||
│ fs1 (ESL :8021, sofia ext :5080/udp, RTP 16384+) ← 独立节点 │
|
||||
│ fs2 (ESL :8021→宿主 8022, 其余同上) ← 零共享状态 │
|
||||
│ │ INVITE / 100/180 / 200 OK (SDP, PCMU) / ACK │
|
||||
│ ▼ │
|
||||
│ mock-provider (SIP UAS :5060/udp, RTP :40000, 440Hz 应答音) │
|
||||
└──────────────────────────────────────────────────────────────────┘
|
||||
agent ⇄ FreeSWITCH: RTP 双向 (PCMU)
|
||||
TX = 440Hz 间歇音 (Mock LLM/TTS) → 被叫听到的 "AI 说话"
|
||||
RX = RMS 统计 (Mock ASR) → 被叫语音进 "识别"
|
||||
```
|
||||
|
||||
**与 Asterisk 方案的关键差异**: FreeSWITCH **没有 ARI externalMedia 式外挂媒体 API**。标准做法是 agent 自身充当一个 SIP 端点 (UAS): FreeSWITCH originate 一条到 agent 的 SIP 腿 (agent 应答 200 OK+SDP), 再 originate 被叫腿, `uuid_bridge` 桥接。控制面从 "ARI REST + 事件 WS" 变为 "ESL TCP 单连接命令式"。agent 的 RTP 收发逻辑与 asterisk 版相同, 仅增加 SIP UAS 部分。
|
||||
|
||||
外呼信令流 (仅出站):
|
||||
```
|
||||
agent → ESL auth → originate sofia/external/agent@host:5062 &park() (腿1, agent 应答)
|
||||
→ originate sofia/gateway/mock-trunk/+1510... &park() (腿2, 被叫应答)
|
||||
→ uuid_bridge <腿1> <腿2> → 双向 RTP 经 FS
|
||||
→ uuid_kill 双腿 → BYE → agent 输出双向统计
|
||||
```
|
||||
|
||||
## 3. 端口规划
|
||||
|
||||
| 端口 | 组件 | 说明 |
|
||||
|---|---|---|
|
||||
| 8021/tcp | fs1 | ESL 控制面, 宿主映射 8021 |
|
||||
| 8022/tcp→8021 | fs2 | 第二节点 ESL (profile `scale`) |
|
||||
| 5080/udp | fs1/fs2 | sofia external profile 信令 (容器网内, 未映射宿主) |
|
||||
| 16384-32768/udp | fs1/fs2 | RTP 媒体 (vanilla 默认段, 各容器独立 netns 不冲突) |
|
||||
| 5060/udp, 40000/udp | mock-provider | Mock 运营商 SIP + RTP |
|
||||
| 5062/udp, 40001/udp | mock-agent (每次呼叫) | 本侧 SIP UAS + RTP |
|
||||
|
||||
生产外网仅需放行: SIP 5060/udp (+5061/tls) 与 RTP 段; **ESL 8021/8022 只对 agent 内网放行 (强密码+ACL), 严禁公网裸开** — ESL 拿到即完全控制 FreeSWITCH (可执行任意系统命令 `system`)。
|
||||
|
||||
## 4. 从零部署步骤 (全部实测)
|
||||
|
||||
```bash
|
||||
# ── 4.1 Docker 与内核调优: 沿用 livekit 项目成果
|
||||
# (/etc/sysctl.d/99-livekit-sip.conf 已应用, 见 §8; asterisk 运行时已下线)
|
||||
|
||||
# ── 4.2 部署栈 ─────────────────────────────────────────────────────
|
||||
mkdir -p /opt/freeswitch-esl && cd /opt/freeswitch-esl # = 仓库 deploy/ 目录
|
||||
./scripts/setup-conf.sh # 生成 etc-fs/ (镜像 vanilla 配置 + 4 处补丁, 见踩坑)
|
||||
docker compose up -d # fs1 + mock-provider
|
||||
docker compose --profile scale up -d # (可选) fs2 第二节点
|
||||
|
||||
# ── 4.3 验证 ───────────────────────────────────────────────────────
|
||||
./scripts/health.sh # 容器 + ESL status + gateway 状态 + 活动呼叫
|
||||
./scripts/call.sh +15105550123 12 1 # 端到端外呼 (节点1)
|
||||
./scripts/call.sh +15105550123 10 2 # 端到端外呼 (节点2, 验证扩展)
|
||||
```
|
||||
|
||||
配置仅 2 个小文件 (conf/): `event_socket.conf.xml` (ESL 0.0.0.0:8021 + 密码 + rfc1918.auto ACL)、`mock-trunk.xml` (sofia 网关: 无注册 + OPTIONS ping 25s 探活); 其余由 `setup-conf.sh` 从镜像 vanilla 配置生成并打补丁。
|
||||
|
||||
⚠ 踩坑1 (**最关键**): safarov 镜像 entrypoint 在 `/etc/freeswitch/freeswitch.xml` 不存在时从 vanilla 种子整目录 cp —— **按文件 RO 挂载会报 "Read-only file system" 且配置落空**; 必须整目录挂载 (`./etc-fs:/etc/freeswitch`), 由 `setup-conf.sh` 生成 (每次全量重建, 幂等)。
|
||||
⚠ 踩坑2: vanilla `event_socket` 的 `listen-ip=::` 在本容器 getaddrinfo 失败 → mod_event_socket 起不来; 改 `0.0.0.0`。
|
||||
⚠ 踩坑3: vanilla `vars.xml` 的 `stun-set` 在配置解析期做外网 STUN 查询, 不通则 `Invalid ext-sip-ip` → **mod_sofia 整体加载失败**; setup-conf.sh 已剥离 stun-set 与 ext-sip-ip/ext-rtp-ip。
|
||||
⚠ 踩坑4 (ESL ACL, 折腾最久): 不设 `apply-inbound-acl` → 默认 loopback.auto, 容器网内 agent 被拒 (`text/rude-rejection`); `localnet.auto` 不含 127.0.0.1 → 容器内 fs_cli 被拒; **自定义 acl.conf.xml 列表在此镜像实测不生效** (127.0.0.1 与 eth0 均被拒)。解: 内置 `rfc1918.auto` (10/8+172.16/12+192.168/16) — agent 与 fs_cli 均放行, fs_cli 需 `-H $(hostname -i)` 走 eth0 源地址。
|
||||
⚠ 踩坑5: 镜像缺 CA 证书, mod_signalwire 每分钟刷 Curl 77 错误日志 → 从 modules.conf.xml 删除加载行。
|
||||
⚠ 踩坑6: fs_cli 默认连 localhost+ClueCon, 需显式 `-H <ip> -P 8021 -p <密码>`; 且 count 类命令输出**首行是空行**, 解析要过滤。
|
||||
⚠ 踩坑7: FS 启动后第一个 originate 可能因网关冷启动未就绪失败; mock-agent 对腿2 已做 3 次重试。
|
||||
|
||||
## 5. 外呼链路验证 (实测输出)
|
||||
|
||||
`./scripts/call.sh +15105550123 12 1`:
|
||||
```
|
||||
[agent] agent up: sip=udp/5062 rtp=udp/40001 ip=172.18.0.4
|
||||
[agent] esl connected: freeswitch1:8021
|
||||
[agent] INVITE from FS 172.18.0.3:5080 call_id=7bdc5fbc-... ← FS originate 腿1 到 agent
|
||||
[agent] leg answered (ACK), fs rtp=('172.18.0.3', 18364) codec=PCMU
|
||||
[agent] leg1 (agent UAS) up: +OK agent-1788159625-782
|
||||
[agent] leg2 (callee) answered: +OK callee-1788159625-782 ← 经 mock-trunk 外呼被叫
|
||||
[agent] bridged [agent-... + callee-...], streaming 12.0s
|
||||
[agent] RX frames=500 avg_rms=8362
|
||||
RESULT number=+15105550123 rx_frames=536 rx_avg_rms=8246 tx_frames=568 bidirectional=YES
|
||||
```
|
||||
mock-provider (被叫侧) 日志:
|
||||
```
|
||||
INCOMING_INVITE from=172.18.0.3:5080 uri="sip:+15105550123@mock-provider:5060"
|
||||
ANSWERED media=172.18.0.2:40000 codec=PCMU peer=172.18.0.3:23888 ← SDP 协商 PCMU
|
||||
ACK confirmed, RTP bridging
|
||||
CALL_DONE reason=bye dur=13.5s sent_rtp=570 recv_rtp=533 recv_avg_amplitude=8236
|
||||
```
|
||||
**双向语音证明**: agent RX=534 帧 (被叫→ASR 方向, RMS 8270 = 440Hz 检测到) 与 provider recv_rtp=533 (LLM/TTS→被叫方向) **计数完全对称**; agent TX=568 ↔ provider sent_rtp=570。挂断后 `show calls count`=0, 无泄漏。换真实 LLM/ASR 只需替换 mock-agent 的 TX 音源与 RX 消费循环。
|
||||
|
||||
## 6. 运行状态检测
|
||||
|
||||
`./scripts/health.sh` (实测输出):
|
||||
```
|
||||
== 容器状态 fs1 / fs2 (healthy) / fs-mock-provider Up
|
||||
== ESL status PASS fs1 (UP, FreeSWITCH 1.10.12) PASS fs2
|
||||
== sofia 网关 PASS fs1 mock-trunk State=NOREG PASS fs2 (NOREG=无注册模式+ping 探活, 正常态)
|
||||
== 活动呼叫 PASS fs1 active calls=0
|
||||
result: ok=5 fail=0
|
||||
```
|
||||
持续监控: `watch -n5 ./scripts/health.sh`; 呼叫级状态: provider 日志 `CALL_DONE` 行 (时长/双向计数); 容器级: docker healthcheck (compose 已配 fs_cli status)。层级与 livekit/asterisk 版对齐: 容器 + API 功能 + 业务功能三层。
|
||||
|
||||
## 7. 水平扩展实证 (双节点, 零配置变更)
|
||||
|
||||
**① 加节点即扩容**: `docker compose --profile scale up -d` 拉起 fs2, 共用同一套 etc-fs, 无任何既有节点配置改动。
|
||||
|
||||
**② 3 路并发分摊双节点 (发起侧调度)**:
|
||||
```
|
||||
$ (./scripts/call.sh +15105550123 10 1 &) (./scripts/call.sh +15105550123 10 1 &) (./scripts/call.sh +15105550123 10 2 &)
|
||||
RESULT ... rx_frames=440 rx_avg_rms=8245 tx_frames=471 bidirectional=YES
|
||||
RESULT ... rx_frames=438 rx_avg_rms=8193 tx_frames=471 bidirectional=YES
|
||||
RESULT ... rx_frames=440 rx_avg_rms=8178 tx_frames=470 bidirectional=YES
|
||||
|
||||
$ docker logs fs-mock-provider | grep INCOMING | grep -oE "from=[0-9.]+" | sort | uniq -c
|
||||
2 from=172.18.0.3 ← fs1 处理 2 路 (源 IP = fs1)
|
||||
1 from=172.18.0.4 ← fs2 处理 1 路 (源 IP = fs2)
|
||||
|
||||
$ (呼中 t=8s) fs1: 2 calls / 4 channels fs2: 1 call / 2 channels ← 每桥接呼叫=2通道
|
||||
$ (provider) CALL_DONE sent_rtp=473 recv_rtp=439 amplitude=8174 ← 双向对称
|
||||
```
|
||||
→ 加副本即扩容、发起侧调度、双节点各自满血双向, 三项均实证。
|
||||
资源基线 (3 并发): fs1 0.3%CPU/48M, fs2 0.3%CPU/44M — 2vCPU 单节点距容量上限还有两个数量级余量。
|
||||
|
||||
## 8. 系统调优
|
||||
|
||||
**内核** (沿用已应用的 `/etc/sysctl.d/99-livekit-sip.conf`, RTP/UDP 通用, 无需改动):
|
||||
```conf
|
||||
net.core.rmem_max = 16777216 # UDP 收缓冲上限 (16M, 默认 208K 高并发必丢包)
|
||||
net.core.wmem_max = 16777216
|
||||
net.core.rmem_default = 1048576
|
||||
net.core.wmem_default = 1048576
|
||||
net.core.netdev_max_backlog = 4096
|
||||
net.ipv4.ip_local_port_range = 10000 65000 # 出向 SIP/RTP 源端口
|
||||
net.ipv4.udp_mem = 8388608 12582912 16777216 # 全局 UDP 页缓存 (3.5G 内存安全值)
|
||||
```
|
||||
验证: `nstat -az | grep -i udp` (`UdpRcvbufErrors` 应为 0)。
|
||||
|
||||
**组件级**:
|
||||
- RTP 段: `switch.conf.xml` 的 `rtp-start-port/rtp-end-port` (vanilla 默认 16384-32768 已注释未改); 生产按并发收窄或分段, 端口不够优先加节点 (见 §7)。
|
||||
- ESL: 高并发控制面改 `bgapi` (异步作业) 或 esl 连接池; 单连接串行 `api` 是同步阻塞的 (本文 mock 场景够用)。`# ponytail: 单 ESL 连接串行 api, >100 并发呼叫控制需 bgapi/连接池`。
|
||||
- sofia 网关: `ping=25` OPTIONS 探活快速摘除故障 trunk; 网关冷启动有秒级窗口 (mock-agent 已重试兜底)。
|
||||
- 容器: `restart: unless-stopped` + healthcheck 已配; 双节点各 ~45M 内存, 3.5G 机器无需 mem limit; `cap_add: SYS_NICE` (FS 定时器优先级)。
|
||||
- 监控: health.sh 三层 + `show calls count` 告警 (骤降); FS 文件日志 `/var/log/freeswitch/freeswitch.log`; 告警项: 容器重启、gateway State 异常、UdpRcvbufErrors>0。
|
||||
|
||||
## 9. 生产化清单 (Mock → 真实运营商)
|
||||
|
||||
1. trunk: mock-trunk.xml 加 `username/password` (digest 注册) 或 IP 白名单模式, proxy 换运营商 SIP 域名/IP; 我方 5060/udp(external profile 5080 改)、RTP 段公网放行 (阿里云安全组)。
|
||||
2. ESL 安全: 8021/8022 仅对 agent 网段放行; 换强密码; 保持 rfc1918.auto ACL (生产可再收紧到容器网段); **ESL 等同 root shell, 公网暴露=事故**。
|
||||
3. NAT: 阿里云 ECS 1:1 NAT — external profile 配 `ext-sip-ip/ext-rtp-ip` 为公网 IP (本文已剥离的参数届时显式回填, 不要用 stun); agent 侧 SIP/RTP 地址用宿主可达地址。
|
||||
4. AI 侧: mock-agent 替换为真实 STT/LLM/TTS (RX 帧喂识别、TX 换合成音频; SIP UAS + ESL 接口不变)。
|
||||
5. 高可用: 节点 ≥2 (已验证); 调度侧带健康检查摘除故障节点; 需要呼入/注册时前置 Kamailio dispatcher。
|
||||
@@ -0,0 +1,15 @@
|
||||
<configuration name="event_socket.conf" description="Socket Client">
|
||||
<settings>
|
||||
<param name="nat-map" value="false"/>
|
||||
<param name="listen-ip" value="0.0.0.0"/>
|
||||
<param name="listen-port" value="8021"/>
|
||||
<param name="password" value="esl_mock_a95a037399153e9092b653fc"/>
|
||||
<!-- ACL 实测 (safarov/freeswitch 1.10.12):
|
||||
- 不设此参数 → 默认 loopback.auto, 容器网内 agent 被拒
|
||||
- localnet.auto 不含 127.0.0.1 → fs_cli(容器内) 被拒
|
||||
- 自定义 ACL (autoload_configs/acl.conf.xml) 启动时不生效
|
||||
→ 用内置 rfc1918.auto (10/8+172.16/12+192.168/16): agent 与 fs_cli 均放行,
|
||||
fs_cli 需 -H $(hostname -i) 走 eth0 源地址 -->
|
||||
<param name="apply-inbound-acl" value="rfc1918.auto"/>
|
||||
</settings>
|
||||
</configuration>
|
||||
@@ -0,0 +1,10 @@
|
||||
<include>
|
||||
<!-- Mock SIP 运营商 trunk: 无注册 (register=false), OPTIONS ping 25s 探活 -->
|
||||
<gateway name="mock-trunk">
|
||||
<param name="proxy" value="mock-provider:5060"/>
|
||||
<param name="register" value="false"/>
|
||||
<param name="ping" value="25"/>
|
||||
<param name="ping-min" value="5"/>
|
||||
<param name="expire-seconds" value="3600"/>
|
||||
</gateway>
|
||||
</include>
|
||||
@@ -0,0 +1,55 @@
|
||||
name: freeswitch-esl
|
||||
|
||||
services:
|
||||
freeswitch1:
|
||||
image: safarov/freeswitch:latest
|
||||
container_name: fs1
|
||||
volumes:
|
||||
- ./etc-fs:/etc/freeswitch
|
||||
ports: ["8021:8021"]
|
||||
cap_add: [SYS_NICE]
|
||||
healthcheck:
|
||||
test: ["CMD-SHELL", "fs_cli -H $(hostname -i) -P 8021 -p esl_mock_a95a037399153e9092b653fc -x status | grep -q '^UP'"]
|
||||
interval: 15s
|
||||
timeout: 5s
|
||||
retries: 3
|
||||
start_period: 30s
|
||||
restart: unless-stopped
|
||||
networks: [fsnet]
|
||||
|
||||
# 第二个 FreeSWITCH 节点: 无共享状态, 出站外呼各自独立 (水平扩展验证)
|
||||
freeswitch2:
|
||||
image: safarov/freeswitch:latest
|
||||
container_name: fs2
|
||||
volumes:
|
||||
- ./etc-fs:/etc/freeswitch
|
||||
ports: ["8022:8021"]
|
||||
cap_add: [SYS_NICE]
|
||||
healthcheck:
|
||||
test: ["CMD-SHELL", "fs_cli -H $(hostname -i) -P 8021 -p esl_mock_a95a037399153e9092b653fc -x status | grep -q '^UP'"]
|
||||
interval: 15s
|
||||
timeout: 5s
|
||||
retries: 3
|
||||
start_period: 30s
|
||||
restart: unless-stopped
|
||||
networks: [fsnet]
|
||||
profiles: [scale]
|
||||
|
||||
# Mock SIP 运营商 (UAS): 应答 FreeSWITCH 的出站 INVITE
|
||||
mock-provider:
|
||||
image: python:3.12-alpine
|
||||
container_name: fs-mock-provider
|
||||
working_dir: /app
|
||||
command: python mock/mock-provider.py
|
||||
volumes:
|
||||
- ./mock:/app/mock:ro
|
||||
environment:
|
||||
SIP_PORT: "5060"
|
||||
RTP_PORT: "40000"
|
||||
RING_DELAY: "1.5"
|
||||
restart: unless-stopped
|
||||
networks: [fsnet]
|
||||
|
||||
networks:
|
||||
fsnet:
|
||||
driver: bridge
|
||||
@@ -0,0 +1,316 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Mock AI Agent — FreeSWITCH/ESL 版.
|
||||
|
||||
方案: FreeSWITCH 无 ARI externalMedia 式外挂媒体 API, 标准做法是本 agent 自身充当
|
||||
SIP UAS 端点 (等价于 asterisk 版的 externalMedia 腿), 控制面走 ESL:
|
||||
0. ESL TCP 连接 fs:8021 (auth)
|
||||
1. 本侧 SIP UAS 监听 (5062) + RTP 监听 (40001)
|
||||
2. ESL originate 腿1: sofia/external/agent@<本容器>:5062 &park() ← 本侧 UAS 应答
|
||||
3. ESL originate 腿2: sofia/gateway/mock-trunk/<被叫> &park() ← mock-provider 应答
|
||||
4. ESL uuid_bridge 腿1 腿2 → 桥接, RTP 经 FreeSWITCH 双向
|
||||
5. RTP: RX=被叫音频(模拟 ASR 输入, RMS 统计); TX=440Hz 间歇音(模拟 LLM/TTS 输出)
|
||||
6. ESL uuid_kill 双腿 → BYE, 输出双向统计
|
||||
纯 stdlib (asyncio + socket)。换真实 LLM/ASR: 替换 TX 音源与 RX 消费即可。
|
||||
用法: mock-agent.py --esl-host fs1 --sip-host <本容器名>:5062 \
|
||||
--number +15105550123 [--duration 15]
|
||||
"""
|
||||
import argparse, asyncio, math, os, random, socket, struct, time
|
||||
|
||||
def log(*a): print(f"[agent {time.strftime('%H:%M:%S')}]", *a, flush=True)
|
||||
|
||||
# ---------- G.711 编码 ----------
|
||||
def ulaw_encode(s: int) -> int:
|
||||
s = max(-32635, min(32635, s)); sign = 0x80 if s < 0 else 0
|
||||
if sign: s = -s
|
||||
s += 0x84
|
||||
e = 7
|
||||
for i in range(7, -1, -1):
|
||||
if s & (0x40 << i): e = i; break
|
||||
return (~(sign | (e << 4) | ((s >> (e + 3)) & 0x0F))) & 0xFF
|
||||
|
||||
def alaw_encode(s: int) -> int:
|
||||
s = max(-32635, min(32635, s)); sign = 0x80 if s < 0 else 0
|
||||
if sign: s = -s
|
||||
e = 7
|
||||
for i in range(7, -1, -1):
|
||||
if s & (0x40 << i): e = i; break
|
||||
m = (s >> (e + 3)) & 0x0F
|
||||
return ((sign | (e << 4) | m) ^ 0x55) & 0xFF
|
||||
|
||||
def ulaw_decode(u: int) -> int:
|
||||
u = ~u & 0xFF
|
||||
sign = u & 0x80; e = (u >> 4) & 0x07; m = u & 0x0F
|
||||
s = ((m << (e + 3)) + 0x84) << 2
|
||||
return -s if sign else s
|
||||
|
||||
def alaw_decode(a: int) -> int:
|
||||
a ^= 0x55
|
||||
t = (a & 0x0F) << 4
|
||||
seg = (a & 0x70) >> 4
|
||||
if seg == 0: t += 8
|
||||
else:
|
||||
t += 0x108
|
||||
if seg > 1: t <<= (seg - 1)
|
||||
return -t if a & 0x80 else t
|
||||
|
||||
def tone_payload(codec: str, ts: float) -> bytes:
|
||||
"""20ms / 160 samples @8kHz, 440Hz 500ms-on 500ms-off (模拟 TTS)"""
|
||||
enc = ulaw_encode if codec == "PCMU" else alaw_encode
|
||||
out = bytearray()
|
||||
on = (ts % 1.0) < 0.5
|
||||
for i in range(160):
|
||||
t = ts + i / 8000.0
|
||||
v = int(6000 * math.sin(2 * math.pi * 440 * t)) if on else 0
|
||||
out.append(enc(v))
|
||||
return bytes(out)
|
||||
|
||||
# ---------- SIP 报文 (与 mock-provider 同源) ----------
|
||||
def parse_msg(data: bytes):
|
||||
head, _, body = data.decode("utf-8", "replace").partition("\r\n\r\n")
|
||||
lines = head.split("\r\n")
|
||||
headers = []
|
||||
for ln in lines[1:]:
|
||||
n, _, v = ln.partition(":")
|
||||
headers.append((n.strip().lower(), v.strip()))
|
||||
return lines[0], headers, body
|
||||
|
||||
def hget(headers, name):
|
||||
for n, v in headers:
|
||||
if n == name: return v
|
||||
return None
|
||||
|
||||
def build_resp(first, headers, code, reason, tag=None, extra=None, body=""):
|
||||
out = [f"SIP/2.0 {code} {reason}"]
|
||||
for n, v in headers:
|
||||
if n == "via": out.append(f"Via: {v}")
|
||||
out.append(f"From: {hget(headers,'from')}")
|
||||
to = hget(headers, "to")
|
||||
if tag and "tag=" not in to: to = f"{to};tag={tag}"
|
||||
out.append(f"To: {to}")
|
||||
out.append(f"Call-ID: {hget(headers,'call-id')}")
|
||||
out.append(f"CSeq: {hget(headers,'cseq')}")
|
||||
if extra:
|
||||
out += extra
|
||||
if body:
|
||||
out.append("Content-Type: application/sdp")
|
||||
out.append(f"Content-Length: {len(body)}")
|
||||
return ("\r\n".join(out) + "\r\n\r\n" + body).encode()
|
||||
|
||||
def parse_sdp_offer(body: str):
|
||||
ip, port, pts = None, None, []
|
||||
for ln in body.splitlines():
|
||||
if ln.startswith("c=IN IP4 "): ip = ln.split()[2]
|
||||
elif ln.startswith("m=audio "):
|
||||
parts = ln.split(); port = int(parts[1]); pts = [p for p in parts[3:] if p.isdigit()]
|
||||
return ip, port, pts
|
||||
|
||||
def pick_codec(pts):
|
||||
if "0" in pts: return "0", "PCMU"
|
||||
if "8" in pts: return "8", "PCMA"
|
||||
return (pts[0], "PCMU") if pts else ("0", "PCMU")
|
||||
|
||||
# ---------- SIP UAS (被 FreeSWITCH originate 的"externalMedia 腿") ----------
|
||||
class SipUas(asyncio.DatagramProtocol):
|
||||
"""应答 FS 的 INVITE (200+SDP), 处理 ACK/BYE/OPTIONS/re-INVITE/重传."""
|
||||
def __init__(self, my_ip, sip_port, rtp_port):
|
||||
self.my_ip, self.sip_port, self.rtp_port = my_ip, sip_port, rtp_port
|
||||
self.transport = None
|
||||
self.tag = os.urandom(5).hex()
|
||||
self.cid = None
|
||||
self.fs_rtp = None # FS 侧 RTP 地址 (INVITE SDP)
|
||||
self.pt, self.codec = "0", "PCMU"
|
||||
self.confirmed = False
|
||||
self.ended = asyncio.Event()
|
||||
|
||||
def connection_made(self, transport):
|
||||
self.transport = transport
|
||||
|
||||
def datagram_received(self, data, addr):
|
||||
first, headers, body = parse_msg(data)
|
||||
method = first.split()[0] if first else "?"
|
||||
if method == "INVITE":
|
||||
cid = hget(headers, "call-id")
|
||||
new_call = self.cid is None
|
||||
if new_call:
|
||||
self.cid = cid
|
||||
log(f"INVITE from FS {addr[0]}:{addr[1]} call_id={cid}")
|
||||
ip, port, pts = parse_sdp_offer(body)
|
||||
if ip and port:
|
||||
self.fs_rtp = (ip, port) # re-INVITE 时同步更新
|
||||
if new_call:
|
||||
self.pt, self.codec = pick_codec(pts)
|
||||
# 重传与 re-INVITE 统一: 按当前请求头重建 200 OK (CSeq 正确)
|
||||
sdp = (f"v=0\r\no=agent {random.randint(1,10**9)} {random.randint(1,10**9)} IN IP4 {self.my_ip}\r\n"
|
||||
f"s=agent\r\nc=IN IP4 {self.my_ip}\r\nt=0 0\r\n"
|
||||
f"m=audio {self.rtp_port} RTP/AVP {self.pt}\r\n"
|
||||
f"a=rtpmap:{self.pt} {self.codec}/8000\r\na=sendrecv\r\n")
|
||||
self.transport.sendto(build_resp(
|
||||
first, headers, 200, "OK", tag=self.tag,
|
||||
extra=[f"Contact: <sip:agent@{self.my_ip}:{self.sip_port}>",
|
||||
"User-Agent: MockAgent/1.0"], body=sdp), addr)
|
||||
elif method == "ACK":
|
||||
if not self.confirmed:
|
||||
self.confirmed = True
|
||||
log(f"leg answered (ACK), fs rtp={self.fs_rtp} codec={self.codec}")
|
||||
elif method == "BYE":
|
||||
self.transport.sendto(build_resp(first, headers, 200, "OK", tag=self.tag), addr)
|
||||
self.ended.set()
|
||||
elif method == "OPTIONS":
|
||||
self.transport.sendto(build_resp(first, headers, 200, "OK", tag=self.tag), addr)
|
||||
|
||||
# ---------- RTP 媒体 ----------
|
||||
class Rx(asyncio.DatagramProtocol):
|
||||
def __init__(self, stats):
|
||||
self.stats = stats
|
||||
|
||||
def datagram_received(self, data, addr):
|
||||
if len(data) < 12: return
|
||||
pt = data[1] & 0x7F
|
||||
self.stats["rx"] += 1
|
||||
payload = data[12:]
|
||||
dec = [(ulaw_decode(b) if pt == 0 else alaw_decode(b)) for b in payload[:40]]
|
||||
self.stats["amp_sum"] += sum(abs(x) for x in dec) / max(1, len(dec))
|
||||
self.stats["tx_to"] = addr # 始终跟随最新来源 (自适应媒体路径)
|
||||
if self.stats["rx"] % 250 == 0:
|
||||
log(f"RX frames={self.stats['rx']} avg_rms={self.stats['amp_sum']/self.stats['rx']:.0f}")
|
||||
|
||||
# ---------- ESL 客户端 (极简: auth + api, 同步命令式) ----------
|
||||
class Esl:
|
||||
def __init__(self, host, port, password):
|
||||
self.host, self.port, self.password = host, port, password
|
||||
|
||||
async def connect(self):
|
||||
self.r, self.w = await asyncio.open_connection(self.host, self.port)
|
||||
m = await self._read()
|
||||
if m.get("content-type") != "auth/request":
|
||||
raise ConnectionError(f"unexpected esl greeting: {m}")
|
||||
self.w.write(f"auth {self.password}\n\n".encode())
|
||||
await self.w.drain()
|
||||
m = await self._read()
|
||||
if "+OK" not in m.get("reply-text", ""):
|
||||
raise ConnectionError(f"esl auth failed: {m}")
|
||||
|
||||
async def _read(self):
|
||||
headers = {}
|
||||
while True:
|
||||
line = await self.r.readline()
|
||||
if not line: raise ConnectionError("esl closed")
|
||||
line = line.decode("utf-8", "replace").rstrip("\r\n")
|
||||
if not line: break
|
||||
k, _, v = line.partition(":")
|
||||
headers[k.strip().lower()] = v.strip()
|
||||
n = int(headers.get("content-length", 0) or 0)
|
||||
headers["_body"] = (await self.r.readexactly(n)).decode("utf-8", "replace") if n else ""
|
||||
return headers
|
||||
|
||||
async def api(self, cmd):
|
||||
self.w.write(f"api {cmd}\n\n".encode())
|
||||
await self.w.drain()
|
||||
return (await self._read())["_body"].strip()
|
||||
|
||||
async def close(self):
|
||||
try: self.w.close()
|
||||
except Exception: pass
|
||||
|
||||
# ---------- 主流程 ----------
|
||||
async def run(args):
|
||||
my_ip = socket.gethostbyname(socket.gethostname())
|
||||
sip_port = int(args.sip_host.rsplit(":", 1)[1])
|
||||
suffix = f"{int(time.time())}-{random.randint(100, 999)}"
|
||||
agent_uuid, callee_uuid = f"agent-{suffix}", f"callee-{suffix}"
|
||||
tmo = int(args.ring_timeout)
|
||||
|
||||
stats = {"rx": 0, "amp_sum": 0.0, "tx": 0, "tx_to": None, "codec": "PCMU"}
|
||||
loop = asyncio.get_running_loop()
|
||||
uas_t, uas = await loop.create_datagram_endpoint(
|
||||
lambda: SipUas(my_ip, sip_port, args.rtp_port), local_addr=("0.0.0.0", sip_port))
|
||||
rtp_t, _ = await loop.create_datagram_endpoint(
|
||||
lambda: Rx(stats), local_addr=("0.0.0.0", args.rtp_port))
|
||||
log(f"agent up: sip=udp/{sip_port} rtp=udp/{args.rtp_port} ip={my_ip}")
|
||||
|
||||
esl = Esl(args.esl_host, args.esl_port, args.esl_password)
|
||||
await esl.connect()
|
||||
log(f"esl connected: {args.esl_host}:{args.esl_port}")
|
||||
|
||||
def fail(msg, detail=None):
|
||||
log(f"FAIL: {msg}", detail or "")
|
||||
raise SystemExit(1)
|
||||
|
||||
try:
|
||||
# 腿1: 本 agent (SIP UAS, 等价 externalMedia): 应答后 park
|
||||
r = await esl.api(
|
||||
f"originate {{origination_uuid={agent_uuid},"
|
||||
f"origination_caller_id_number={args.caller_id},originate_timeout={tmo}}}"
|
||||
f"sofia/external/agent@{args.sip_host} &park()")
|
||||
if not r.startswith("+OK"): fail("originate agent leg", r)
|
||||
log(f"leg1 (agent UAS) up: {r}")
|
||||
|
||||
# 腿2: 被叫经 mock-trunk 网关 (网关冷启动自动重试)
|
||||
for attempt in range(3):
|
||||
r = await esl.api(
|
||||
f"originate {{origination_uuid={callee_uuid},"
|
||||
f"origination_caller_id_number={args.caller_id},originate_timeout={tmo}}}"
|
||||
f"sofia/gateway/mock-trunk/{args.number} &park()")
|
||||
if r.startswith("+OK"): break
|
||||
log(f"leg2 originate retry {attempt + 1}: {r}")
|
||||
await asyncio.sleep(2)
|
||||
if not r.startswith("+OK"): fail("originate callee leg", r)
|
||||
log(f"leg2 (callee) answered: {r}")
|
||||
|
||||
stats["codec"] = uas.codec
|
||||
if uas.fs_rtp and not stats["tx_to"]:
|
||||
stats["tx_to"] = uas.fs_rtp
|
||||
log(f"fs rtp peer (INVITE SDP): {uas.fs_rtp[0]}:{uas.fs_rtp[1]}")
|
||||
|
||||
# 桥接双腿 → RTP 经 FS 双向
|
||||
r = await esl.api(f"uuid_bridge {agent_uuid} {callee_uuid}")
|
||||
if not r.startswith("+OK"): fail("uuid_bridge", r)
|
||||
log(f"bridged [{agent_uuid} + {callee_uuid}], streaming {args.duration}s "
|
||||
f"(RX=ASR mock, TX=440Hz LLM/TTS mock)")
|
||||
|
||||
# 媒体收发
|
||||
async def tx():
|
||||
seq = random.randint(0, 65535); ts = random.randint(0, 10**6)
|
||||
ssrc = random.randint(1, 2**31)
|
||||
pt = 0 if stats["codec"] == "PCMU" else 8
|
||||
t0 = time.time()
|
||||
while time.time() - t0 < args.duration:
|
||||
if stats["tx_to"]:
|
||||
hdr = struct.pack("!BBHII", 0x80, pt, seq & 0xFFFF,
|
||||
ts & 0xFFFFFFFF, ssrc)
|
||||
rtp_t.sendto(hdr + tone_payload(stats["codec"], time.time() - t0),
|
||||
stats["tx_to"])
|
||||
stats["tx"] += 1
|
||||
seq += 1; ts += 160
|
||||
await asyncio.sleep(0.02)
|
||||
await asyncio.wait_for(tx(), timeout=args.duration + 5)
|
||||
finally:
|
||||
for u in (callee_uuid, agent_uuid):
|
||||
try: await esl.api(f"uuid_kill {u}")
|
||||
except Exception: pass
|
||||
await esl.close()
|
||||
uas_t.close(); rtp_t.close()
|
||||
await asyncio.sleep(0.3) # 等 BYE 应答落地
|
||||
|
||||
n = stats["rx"]
|
||||
amp = stats["amp_sum"] / n if n else 0.0
|
||||
print(f"RESULT number={args.number} rx_frames={n} rx_avg_rms={amp:.0f} "
|
||||
f"tx_frames={stats['tx']} bidirectional={'YES' if n > 50 and amp > 100 else 'NO'}",
|
||||
flush=True)
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser()
|
||||
ap.add_argument("--esl-host", default="fs1")
|
||||
ap.add_argument("--esl-port", type=int, default=8021)
|
||||
ap.add_argument("--esl-password", default="esl_mock_a95a037399153e9092b653fc")
|
||||
ap.add_argument("--sip-host", required=True, help="本侧 SIP 地址 host:port (供 FS INVITE)")
|
||||
ap.add_argument("--rtp-port", type=int, default=40001)
|
||||
ap.add_argument("--number", default="+15105550123")
|
||||
ap.add_argument("--caller-id", default="+15109990001")
|
||||
ap.add_argument("--duration", type=float, default=15)
|
||||
ap.add_argument("--ring-timeout", type=float, default=30)
|
||||
args = ap.parse_args()
|
||||
asyncio.run(run(args))
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Executable
+244
@@ -0,0 +1,244 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Mock SIP 运营商 (UAS): 接收 Asterisk(PJSIP) 出站 INVITE, 振铃->应答, 双向收发 PCMU/PCMA RTP.
|
||||
仅 stdlib, 无依赖。一个 RTP 端口复用所有呼叫(按来源 IP 匹配呼叫)。
|
||||
用法: python3 mock-provider.py (环境变量 SIP_PORT/RTP_PORT/RING_DELAY)
|
||||
"""
|
||||
import asyncio, math, os, random, socket, struct, time
|
||||
|
||||
SIP_PORT = int(os.environ.get("SIP_PORT", "5060"))
|
||||
RTP_PORT = int(os.environ.get("RTP_PORT", "40000"))
|
||||
RING_DELAY = float(os.environ.get("RING_DELAY", "1.5"))
|
||||
MY_IP = socket.gethostbyname(socket.gethostname())
|
||||
|
||||
def log(*a): print(f"[provider {time.strftime('%H:%M:%S')}]", *a, flush=True)
|
||||
|
||||
# ---------- G.711 编码 ----------
|
||||
def ulaw_encode(s: int) -> int:
|
||||
s = max(-32635, min(32635, s)); sign = 0x80 if s < 0 else 0
|
||||
if sign: s = -s
|
||||
s += 0x84
|
||||
e = 7
|
||||
for i in range(7, -1, -1):
|
||||
if s & (0x40 << i): e = i; break
|
||||
return (~(sign | (e << 4) | ((s >> (e + 3)) & 0x0F))) & 0xFF
|
||||
|
||||
def alaw_encode(s: int) -> int:
|
||||
s = max(-32635, min(32635, s)); sign = 0x80 if s < 0 else 0
|
||||
if sign: s = -s
|
||||
e = 7
|
||||
for i in range(7, -1, -1):
|
||||
if s & (0x40 << i): e = i; break
|
||||
m = (s >> (e + 3)) & 0x0F
|
||||
return ((sign | (e << 4) | m) ^ 0x55) & 0xFF
|
||||
|
||||
def ulaw_decode(u: int) -> int:
|
||||
u = ~u & 0xFF
|
||||
sign = u & 0x80; e = (u >> 4) & 0x07; m = u & 0x0F
|
||||
s = ((m << (e + 3)) + 0x84) << 2
|
||||
return -s if sign else s
|
||||
|
||||
def tone_payload(codec: str, ts: float) -> bytes:
|
||||
"""20ms / 160 samples @8kHz, 440Hz 500ms-on 500ms-off"""
|
||||
out = bytearray()
|
||||
on = (ts % 1.0) < 0.5
|
||||
for i in range(160):
|
||||
t = ts + i / 8000.0
|
||||
v = int(6000 * math.sin(2 * math.pi * 440 * t)) if on else 0
|
||||
out.append(ulaw_encode(v) if codec == "PCMU" else alaw_encode(v))
|
||||
return bytes(out)
|
||||
|
||||
# ---------- SIP 报文 ----------
|
||||
def parse_msg(data: bytes):
|
||||
head, _, body = data.decode("utf-8", "replace").partition("\r\n\r\n")
|
||||
lines = head.split("\r\n")
|
||||
headers = []
|
||||
for ln in lines[1:]:
|
||||
n, _, v = ln.partition(":")
|
||||
headers.append((n.strip().lower(), v.strip()))
|
||||
return lines[0], headers, body
|
||||
|
||||
def hget(headers, name):
|
||||
for n, v in headers:
|
||||
if n == name: return v
|
||||
return None
|
||||
|
||||
def build_resp(first, headers, code, reason, tag=None, extra=None, body=""):
|
||||
out = [f"SIP/2.0 {code} {reason}"]
|
||||
for n, v in headers:
|
||||
if n == "via": out.append(f"Via: {v}")
|
||||
out.append(f"From: {hget(headers,'from')}")
|
||||
to = hget(headers, "to")
|
||||
if tag and "tag=" not in to: to = f"{to};tag={tag}"
|
||||
out.append(f"To: {to}")
|
||||
out.append(f"Call-ID: {hget(headers,'call-id')}")
|
||||
out.append(f"CSeq: {hget(headers,'cseq')}")
|
||||
if extra:
|
||||
out += extra
|
||||
if body:
|
||||
out.append(f"Content-Type: application/sdp")
|
||||
out.append(f"Content-Length: {len(body)}")
|
||||
return ("\r\n".join(out) + "\r\n\r\n" + body).encode()
|
||||
|
||||
def parse_sdp_offer(body: str):
|
||||
"""返回 (对端媒体IP, 对端媒体端口, payload类型列表)"""
|
||||
ip, port, pts = None, None, []
|
||||
for ln in body.splitlines():
|
||||
if ln.startswith("c=IN IP4 "): ip = ln.split()[2]
|
||||
elif ln.startswith("m=audio "):
|
||||
parts = ln.split(); port = int(parts[1]); pts = [p for p in parts[3:] if p.isdigit()]
|
||||
return ip, port, pts
|
||||
|
||||
def sdp_answer(codec_pt: str, codec_name: str) -> str:
|
||||
sid = random.randint(1, 10**9)
|
||||
return (f"v=0\r\no=mock {sid} {sid+1} IN IP4 {MY_IP}\r\ns=mock\r\n"
|
||||
f"c=IN IP4 {MY_IP}\r\nt=0 0\r\n"
|
||||
f"m=audio {RTP_PORT} RTP/AVP {codec_pt}\r\n"
|
||||
f"a=rtpmap:{codec_pt} {codec_name}/8000\r\na=sendrecv\r\n")
|
||||
|
||||
# ---------- 呼叫对象 ----------
|
||||
class Call:
|
||||
def __init__(self, first, headers, addr):
|
||||
self.id = hget(headers, "call-id")
|
||||
self.tag = os.urandom(5).hex()
|
||||
self.first, self.headers, self.addr = first, headers, addr
|
||||
self.peer_ip = self.peer_port = None
|
||||
self.codec = "PCMU"; self.pt = "0"
|
||||
self.ok_raw = None; self.confirmed = False; self.closed = False
|
||||
self.sent = self.recv = self.recv_bytes = self.recv_amp_sum = 0
|
||||
self.t0 = time.time()
|
||||
self.last_resp = None
|
||||
|
||||
calls: dict[str, Call] = {}
|
||||
rtp_sock = None
|
||||
|
||||
def pick_codec(pts):
|
||||
if "0" in pts: return "0", "PCMU"
|
||||
if "8" in pts: return "8", "PCMA"
|
||||
return (pts[0], "PCMU") if pts else ("0", "PCMU")
|
||||
|
||||
async def ringing_then_answer(call, sip_t):
|
||||
await asyncio.sleep(RING_DELAY)
|
||||
if call.closed: return
|
||||
body = sdp_answer(call.pt, call.codec)
|
||||
call.ok_raw = build_resp(call.first, call.headers, 200, "OK", tag=call.tag,
|
||||
extra=[f"Contact: <sip:mock@{MY_IP}:{SIP_PORT}>",
|
||||
"User-Agent: MockProvider/1.0"], body=body)
|
||||
call.last_resp = call.ok_raw
|
||||
sip_t.sendto(call.ok_raw, call.addr)
|
||||
log(f"ANSWERED call_id={call.id} media={MY_IP}:{RTP_PORT} codec={call.codec} peer={call.peer_ip}:{call.peer_port}")
|
||||
start_rtp(call)
|
||||
for _ in range(5): # 重传 200 OK 直到 ACK (UDP 可靠性)
|
||||
await asyncio.sleep(1.0)
|
||||
if call.confirmed or call.closed: break
|
||||
sip_t.sendto(call.ok_raw, call.addr)
|
||||
|
||||
def start_rtp(call):
|
||||
async def sender():
|
||||
seq = random.randint(0, 65535); ts = random.randint(0, 10**6)
|
||||
ssrc = random.randint(1, 2**31)
|
||||
pt = 0 if call.codec == "PCMU" else 8
|
||||
t0 = time.time()
|
||||
while not call.closed:
|
||||
now = time.time() - t0
|
||||
payload = tone_payload(call.codec, now)
|
||||
hdr = struct.pack("!BBHII", 0x80, pt, seq & 0xFFFF, ts & 0xFFFFFFFF, ssrc)
|
||||
try:
|
||||
rtp_sock.sendto(hdr + payload, (call.peer_ip, call.peer_port))
|
||||
except OSError:
|
||||
pass
|
||||
call.sent += 1; seq += 1; ts += 160
|
||||
await asyncio.sleep(0.02)
|
||||
asyncio.create_task(sender())
|
||||
|
||||
def close_call(call, reason):
|
||||
if call.closed: return
|
||||
call.closed = True
|
||||
dur = time.time() - call.t0
|
||||
avg = call.recv_amp_sum / call.recv if call.recv else 0
|
||||
log(f"CALL_DONE call_id={call.id} reason={reason} dur={dur:.1f}s "
|
||||
f"sent_rtp={call.sent} recv_rtp={call.recv} recv_bytes={call.recv_bytes} recv_avg_amplitude={avg:.0f}")
|
||||
calls.pop(call.id, None)
|
||||
|
||||
def on_sip(data, addr, sip_t):
|
||||
first, headers, body = parse_msg(data)
|
||||
method = first.split()[0] if first else "?"
|
||||
cid = hget(headers, "call-id")
|
||||
|
||||
if method == "INVITE":
|
||||
call = calls.get(cid)
|
||||
if call: # 重传: 回最后一条响应
|
||||
if call.last_resp: sip_t.sendto(call.last_resp, addr)
|
||||
return
|
||||
call = Call(first, headers, addr)
|
||||
call.peer_ip, call.peer_port, pts = parse_sdp_offer(body)
|
||||
call.pt, call.codec = pick_codec(pts)
|
||||
calls[cid] = call
|
||||
log(f"INCOMING_INVITE from={addr[0]}:{addr[1]} uri=\"{first.split()[1]}\" call_id={cid}")
|
||||
sip_t.sendto(build_resp(first, headers, 100, "Trying"), addr)
|
||||
sip_t.sendto(build_resp(first, headers, 180, "Ringing"), addr)
|
||||
call.last_resp = build_resp(first, headers, 180, "Ringing")
|
||||
asyncio.create_task(ringing_then_answer(call, sip_t))
|
||||
|
||||
elif method == "ACK":
|
||||
call = calls.get(cid)
|
||||
if call and not call.confirmed:
|
||||
call.confirmed = True
|
||||
log(f"ACK call_id={cid} confirmed, RTP bridging")
|
||||
|
||||
elif method == "BYE":
|
||||
call = calls.get(cid)
|
||||
sip_t.sendto(build_resp(first, headers, 200, "OK", tag=call.tag if call else None), addr)
|
||||
if call: close_call(call, "bye")
|
||||
|
||||
elif method == "CANCEL":
|
||||
call = calls.get(cid)
|
||||
sip_t.sendto(build_resp(first, headers, 487, "Request Terminated",
|
||||
tag=call.tag if call else None), addr)
|
||||
if call: close_call(call, "cancelled")
|
||||
|
||||
elif method == "OPTIONS":
|
||||
sip_t.sendto(build_resp(first, headers, 200, "OK", tag="mock"), addr)
|
||||
|
||||
class RtpProto(asyncio.DatagramProtocol):
|
||||
def datagram_received(self, data, addr):
|
||||
if len(data) < 12: return
|
||||
pt = data[1] & 0x7F
|
||||
payload = data[12:]
|
||||
# 优先 ip:port 精确匹配 (同 IP 多并发呼叫时按 RTP 端口区分), 兜底仅 IP
|
||||
for call in list(calls.values()):
|
||||
if not call.closed and call.peer_ip == addr[0] and call.peer_port == addr[1]:
|
||||
self._tally(call, pt, payload); return
|
||||
for call in list(calls.values()):
|
||||
if not call.closed and call.peer_ip == addr[0]:
|
||||
self._tally(call, pt, payload); return
|
||||
|
||||
@staticmethod
|
||||
def _tally(call, pt, payload):
|
||||
call.recv += 1; call.recv_bytes += len(payload) + 12
|
||||
if payload and pt in (0, 8):
|
||||
dec = [ulaw_decode(b) for b in payload[:40]]
|
||||
call.recv_amp_sum += sum(abs(x) for x in dec) / len(dec)
|
||||
|
||||
class SipProto(asyncio.DatagramProtocol):
|
||||
def __init__(self): self.t = None
|
||||
def connection_made(self, transport): self.t = transport
|
||||
def datagram_received(self, data, addr):
|
||||
if self.t: on_sip(data, addr, self.t)
|
||||
|
||||
async def main():
|
||||
global rtp_sock
|
||||
loop = asyncio.get_running_loop()
|
||||
sip_t, _ = await loop.create_datagram_endpoint(
|
||||
lambda: SipProto(), local_addr=("0.0.0.0", SIP_PORT))
|
||||
rtp_t, _ = await loop.create_datagram_endpoint(
|
||||
RtpProto, local_addr=("0.0.0.0", RTP_PORT))
|
||||
rtp_sock = rtp_t
|
||||
log(f"mock provider up: sip=udp/{SIP_PORT} rtp=udp/{RTP_PORT} ip={MY_IP}")
|
||||
log("waiting for outbound INVITE from asterisk ...")
|
||||
await asyncio.Event().wait()
|
||||
|
||||
if __name__ == "__main__":
|
||||
try:
|
||||
asyncio.run(main())
|
||||
except KeyboardInterrupt:
|
||||
pass
|
||||
Executable
+18
@@ -0,0 +1,18 @@
|
||||
#!/bin/bash
|
||||
# 端到端外呼验证: mock-agent(ESL+SIP UAS+RTP) → FreeSWITCH → mock-trunk → mock-provider(应答) → 双向 RTP 统计
|
||||
# 用法: ./call.sh [被叫号码] [时长秒] [freeswitch节点 1|2]
|
||||
set -euo pipefail
|
||||
cd "$(dirname "$0")/.."
|
||||
|
||||
NUMBER="${1:-+15105550123}"
|
||||
DURATION="${2:-15}"
|
||||
NODE="${3:-1}"
|
||||
NET="freeswitch-esl_fsnet"
|
||||
NAME="mock-agent-$$"
|
||||
|
||||
docker run --rm --name "$NAME" --network "$NET" \
|
||||
-v "$PWD/mock:/app/mock:ro" -w /app \
|
||||
python:3.12-alpine python mock/mock-agent.py \
|
||||
--esl-host "freeswitch${NODE}" --esl-port 8021 \
|
||||
--sip-host "${NAME}:5062" \
|
||||
--number "$NUMBER" --duration "$DURATION"
|
||||
Executable
+53
@@ -0,0 +1,53 @@
|
||||
#!/bin/bash
|
||||
# 运行状态检测: 容器 + ESL 命令 + sofia 网关状态 + 活动呼叫数
|
||||
# 注: ESL ACL=rfc1918.auto 不含 loopback → fs_cli 必须 -H <eth0 IP> 走非 loopback 源地址;
|
||||
# 且需显式 -P 8021 -p <密码> (fs_cli 默认 localhost+ClueCon)
|
||||
set -uo pipefail
|
||||
PASS="esl_mock_a95a037399153e9092b653fc"
|
||||
ok=0; fail=0
|
||||
pass() { echo "PASS $1"; ok=$((ok+1)); }
|
||||
faill() { echo "FAIL $1"; fail=$((fail+1)); }
|
||||
running() { [ -n "$(docker ps -q -f name=$1)" ]; }
|
||||
fsc() { # fsc <容器> <fs_cli 参数...>
|
||||
local ip
|
||||
ip=$(docker exec "$1" hostname -i 2>/dev/null | awk '{print $1}')
|
||||
docker exec "$1" fs_cli -H "$ip" -P 8021 -p "$PASS" "${@:2}" 2>/dev/null
|
||||
}
|
||||
|
||||
echo "== 容器状态"
|
||||
docker ps --filter "name=fs" --format "{{.Names}}\t{{.Status}}"
|
||||
|
||||
echo "== ESL (fs_cli status)"
|
||||
for c in fs1 fs2; do
|
||||
if running "$c"; then
|
||||
out=$(fsc "$c" -x "status" | head -1)
|
||||
if echo "$out" | grep -q "^UP"; then
|
||||
pass "$c ESL status: ${out:0:60}"
|
||||
else
|
||||
faill "$c ESL status: ${out:-none}"
|
||||
fi
|
||||
else
|
||||
echo "SKIP $c (未运行)"
|
||||
fi
|
||||
done
|
||||
|
||||
echo "== sofia 网关 (mock-trunk)"
|
||||
for c in fs1 fs2; do
|
||||
if running "$c"; then
|
||||
state=$(fsc "$c" -x "sofia status gateway mock-trunk" | awk -F'[:\t ]+' '/^State/{print $2}')
|
||||
case "$state" in
|
||||
REGED|NOREG) pass "$c gateway mock-trunk State=$state" ;;
|
||||
*) faill "$c gateway mock-trunk State=${state:-none}" ;;
|
||||
esac
|
||||
fi
|
||||
done
|
||||
|
||||
echo "== 活动呼叫"
|
||||
if running fs1; then
|
||||
n=$(fsc fs1 -x "show calls count" | grep total | awk '{print $1}')
|
||||
[ -n "$n" ] && pass "fs1 active calls=$n" || faill "fs1 active calls"
|
||||
fi
|
||||
|
||||
echo
|
||||
echo "result: ok=$ok fail=$fail"
|
||||
[ $fail -eq 0 ]
|
||||
Executable
+29
@@ -0,0 +1,29 @@
|
||||
#!/bin/bash
|
||||
# 生成 etc-fs 运行配置: 镜像 vanilla 配置 + 本仓库覆盖/补丁
|
||||
# 每次全量重建 (rm -rf → 拷 vanilla → 打补丁), 天然幂等, 不会叠加损坏 XML
|
||||
# 产出: ../etc-fs/ (整目录挂载到容器 /etc/freeswitch, 见 docker-compose.yml)
|
||||
set -euo pipefail
|
||||
cd "$(dirname "$0")/.."
|
||||
IMG="safarov/freeswitch:latest"
|
||||
|
||||
# 1. 从镜像导出干净 vanilla 配置
|
||||
rm -rf etc-fs && mkdir -p etc-fs
|
||||
docker run --rm --entrypoint sh -v "$PWD/etc-fs:/out" "$IMG" \
|
||||
-c "cp -a /usr/share/freeswitch/conf/vanilla/. /out/"
|
||||
|
||||
# 2. ESL: 0.0.0.0:8021 + 自定义密码 + 内置 rfc1918.auto ACL
|
||||
# (vanilla listen-ip=:: 在本容器 getaddrinfo 失败; ACL 详见文件内注释)
|
||||
cp conf/event_socket.conf.xml etc-fs/autoload_configs/event_socket.conf.xml
|
||||
|
||||
# 3. Mock 运营商 trunk 网关 (无注册 + OPTIONS ping 探活)
|
||||
cp conf/mock-trunk.xml etc-fs/sip_profiles/external/mock-trunk.xml
|
||||
|
||||
# 4. 去 stun: 依赖 (容器网内无 NAT; stun-set 在配置解析期做外网 STUN 查询,
|
||||
# 不通则 Invalid ext-sip-ip → mod_sofia 加载失败)
|
||||
sed -i '/cmd="stun-set"/d' etc-fs/vars.xml
|
||||
sed -i '/<param name="ext-sip-ip"/d;/<param name="ext-rtp-ip"/d' etc-fs/sip_profiles/*.xml
|
||||
|
||||
# 5. 删除 mod_signalwire (镜像缺 CA 证书, Curl 77 每分钟刷错误日志)
|
||||
sed -i '/<load module="mod_signalwire"\/>/d' etc-fs/autoload_configs/modules.conf.xml
|
||||
|
||||
echo "etc-fs rebuilt (freeswitch.xml: $(ls etc-fs/freeswitch.xml))"
|
||||
@@ -0,0 +1,25 @@
|
||||
目标机器:39.106.106.246
|
||||
登录用户: root
|
||||
|
||||
任务目标:在 asterisk 运行时下线后,调研部署 FreeSWITCH 并以 ESL + Mock SIP 实现 SIP 对接与运行状态检测。
|
||||
|
||||
项目地址: https://github.com/signalwire/freeswitch (容器镜像 safarov/freeswitch = FreeSWITCH 1.10.12)
|
||||
项目背景:使用 FreeSWITCH 对接 MOCK SIP 服务商,作为AI外呼机器人基础能力底座,使用MOCK对接LLM ASR实现双向语音通信, 本项目仅实现外呼功能,不实现多人房间与呼入功能调研。给出可实现最快水平扩展的结论方案并赋实践证据。给出详细的部署过程文档与系统调优方案。
|
||||
|
||||
---
|
||||
|
||||
## 交付状态 (2026-08-31 全部完成)
|
||||
|
||||
| 目标 | 状态 | 证据 |
|
||||
|---|---|---|
|
||||
| asterisk 运行时下线 | ✅ | ast-* 容器清零, 8088/8089 端口释放; `/opt/asterisk-ari/` 交付物保留 |
|
||||
| FreeSWITCH 部署 (双节点) | ✅ | fs1/fs2 healthy, ESL status UP ×2 (FreeSWITCH 1.10.12), health.sh ok=5 fail=0 |
|
||||
| Mock SIP 对接 (外呼) | ✅ | ESL originate ×2 → INVITE→180→200 OK(PCMU)→ACK, gateway NOREG+ping 探活 ×2 |
|
||||
| 双向语音 (Mock LLM/ASR) | ✅ | `RESULT ... bidirectional=YES` (agent RX 536帧 RMS 8246; provider recv_rtp=533, 双向计数对称) |
|
||||
| 运行状态检测 | ✅ | `scripts/health.sh` 容器+ESL status+gateway+呼叫数, ok=5 fail=0 |
|
||||
| 水平扩展结论+实证 | ✅ | 3 路并发分摊双节点 (INVITE 源 172.18.0.3×2 / 172.18.0.4×1, 呼中 fs1=2 calls/4ch, fs2=1/2), 零配置变更 |
|
||||
| 部署文档+调优方案 | ✅ | 本仓库 `README.md` (§4 部署含 7 处踩坑, §8 sysctl 沿用已应用) |
|
||||
|
||||
**调研要点 (vs Asterisk/ARI)**: FreeSWITCH 无 externalMedia 式外挂媒体 API —— mock-agent 自身充当 SIP UAS 端点 (FS originate 至 agent 并应答), 控制面走 ESL TCP (单连接命令式) 而非 ARI REST+WS; 水平扩展模型与 Asterisk 相同 (零共享状态节点复制, 不需要 livekit 那样的 Redis)。
|
||||
|
||||
**交付物**: 服务器 `/opt/freeswitch-esl/` (compose+conf+mock+scripts+etc-fs 生成器+README), 本地 `deploy/` 同源; 原始任务与完成状态见本文件, 详细文档见 `README.md`。
|
||||
Reference in New Issue
Block a user