From ef7ad161c653abd517655c47d85ed5b1b71745a7 Mon Sep 17 00:00:00 2001 From: rogee Date: Mon, 31 Aug 2026 13:46:57 +0800 Subject: [PATCH] chore: init SIP research repo (Asterisk & LiveKit) --- asterisk/README.md | 168 ++++++++++ asterisk/deploy/conf/ari.conf | 8 + asterisk/deploy/conf/http.conf | 4 + asterisk/deploy/conf/pjsip.conf | 21 ++ asterisk/deploy/conf/rtp.conf | 6 + asterisk/deploy/docker-compose.yml | 47 +++ .../__pycache__/mock-agent.cpython-313.pyc | Bin 0 -> 21667 bytes .../__pycache__/mock-provider.cpython-313.pyc | Bin 0 -> 18950 bytes asterisk/deploy/mock/mock-agent.py | 295 ++++++++++++++++++ asterisk/deploy/mock/mock-provider.py | 244 +++++++++++++++ asterisk/deploy/scripts/call.sh | 18 ++ asterisk/deploy/scripts/health.sh | 44 +++ asterisk/deployment.md | 23 ++ livekit/README.md | 198 ++++++++++++ livekit/deploy/conf/livekit.yaml | 15 + livekit/deploy/conf/sip.yaml | 12 + livekit/deploy/docker-compose.yml | 75 +++++ livekit/deploy/mock/mock-agent.py | 78 +++++ livekit/deploy/mock/mock-provider.py | 237 ++++++++++++++ livekit/deploy/mock/trunk.json | 7 + livekit/deploy/scripts/call.sh | 38 +++ livekit/deploy/scripts/health.sh | 32 ++ livekit/deployment.md | 24 ++ 23 files changed, 1594 insertions(+) create mode 100644 asterisk/README.md create mode 100644 asterisk/deploy/conf/ari.conf create mode 100644 asterisk/deploy/conf/http.conf create mode 100644 asterisk/deploy/conf/pjsip.conf create mode 100644 asterisk/deploy/conf/rtp.conf create mode 100644 asterisk/deploy/docker-compose.yml create mode 100644 asterisk/deploy/mock/__pycache__/mock-agent.cpython-313.pyc create mode 100644 asterisk/deploy/mock/__pycache__/mock-provider.cpython-313.pyc create mode 100755 asterisk/deploy/mock/mock-agent.py create mode 100755 asterisk/deploy/mock/mock-provider.py create mode 100755 asterisk/deploy/scripts/call.sh create mode 100755 asterisk/deploy/scripts/health.sh create mode 100644 asterisk/deployment.md create mode 100644 livekit/README.md create mode 100644 livekit/deploy/conf/livekit.yaml create mode 100644 livekit/deploy/conf/sip.yaml create mode 100644 livekit/deploy/docker-compose.yml create mode 100644 livekit/deploy/mock/mock-agent.py create mode 100644 livekit/deploy/mock/mock-provider.py create mode 100644 livekit/deploy/mock/trunk.json create mode 100644 livekit/deploy/scripts/call.sh create mode 100644 livekit/deploy/scripts/health.sh create mode 100644 livekit/deployment.md diff --git a/asterisk/README.md b/asterisk/README.md new file mode 100644 index 0000000..fb56f21 --- /dev/null +++ b/asterisk/README.md @@ -0,0 +1,168 @@ +# Asterisk + ARI AI 外呼底座 — 部署与运维文档 + +> 目标机器: `39.106.106.246` (阿里云 2vCPU/3.5G, Debian 13) · 交付物: `/opt/asterisk-ari/` +> 项目: https://github.com/asterisk/asterisk · 镜像: andrius/asterisk (Asterisk 22.10.1 LTS) +> 范围: **仅外呼** (ARI originate + externalMedia), Mock SIP 运营商 + Mock LLM/ASR, 已端到端验证双向语音。 +> 前置: livekit 运行时已下线 (lk-* 容器清零、端口释放); Docker 与内核调优沿用其成果。 + +## 1. 水平扩展结论 (先行) + +**结论: 零共享状态的独立节点复制 —— Asterisk 无内置集群, 每节点呼叫状态全在本机内存; 扩容 = 加一个节点, 呼叫发起侧 (agent/业务后端) 在节点名单上选节点发起 ARI originate。不需要 Redis, 不需要任何协调组件。** + +| 层 | 状态存放 | 扩容方式 | 调度机制 | +|---|---|---|---| +| Asterisk 节点 | 本机内存 (通道/桥/endpoint) | 复制 service 块 (独立宿主端口) | 呼叫发起侧选节点 (call.sh 的 NODE 参数 = 生产中的 LB 名单) | +| AI Agent (生产) | 无状态 worker | `--scale`/K8s HPA | 每个外呼任务自带节点选择与 ARI/RTP 会话 | + +对比 livekit 方案: livekit-sip 需要 Redis 存 trunk/呼叫状态并做派发; Asterisk 出站外呼场景连共享层都不需要 —— 每个呼叫自包含 (originate + bridge + externalMedia 全在一个节点内闭环), 节点间零交互。实测 1→2 节点**零配置变更** (共用同一套 conf), 呼叫按发起侧选择分摊 (见 §7)。 + +瓶颈与边界: ① 单节点容量受 CPU (RTP 收发 + 桥混音) 与 RTP 端口段限制 — 本文 10000-10800 共 800 端口, 每呼叫 2 腿 (被叫 + externalMedia) 各占 1 对端口 (RTP+RTCP) ≈ 200 并发/节点; ② 无跨节点媒体互通 — 仅外呼场景无影响 (呼叫不跨节点); ③ 未来要呼入/注册时需前置 Kamailio dispatcher (SIP 负载均衡), 纯出站不需要。 + +## 2. 架构 + +``` +业务后端 / AI Agent (mock-agent, 容器或 host) + │ ① POST /bridges (mixing) ─ ARI REST, 节点可选 (ast1/ast2 轮询 = 调度层) + │ ② POST /channels/externalMedia ─ Asterisk 将桥内混音以 ulaw RTP 发往 agent + │ ③ POST /channels (originate) ─ PJSIP/mock-trunk 外呼被叫 + │ ④ 轮询通道至 Up → addChannel 双通道入桥 + ▼ +┌─ docker net: astari ──────────────────────────────────────────────┐ +│ ast1 (ARI :8088, PJSIP :5060/udp, RTP 10000-10800) ← 独立节点 │ +│ ast2 (ARI :8088→宿主 8089, 其余同上) ← 零共享状态 │ +│ │ INVITE / 100/180 / 200 OK (SDP, PCMU) / ACK │ +│ ▼ │ +│ mock-provider (SIP UAS :5060/udp, RTP :40000, 440Hz 应答音) │ +└────────────────────────────────────────────────────────────────────┘ + agent ⇄ Asterisk: 混音 RTP 双向 (ulaw) + TX = 440Hz 间歇音 (Mock LLM/TTS) → 被叫听到的 "AI 说话" + RX = RMS 统计 (Mock ASR) → 被叫语音进 "识别" +``` + +外呼信令流 (仅出站): +``` +业务后端 → POST /ari/channels (endpoint=PJSIP/+1510...@mock-trunk, app=outbound) + → Asterisk 发 INVITE → mock-provider 100/180 → 200 OK (PCMU) → ACK + → 被叫通道 Up, 与 externalMedia 通道同入 mixing 桥 + → agent 经 RTP 双向收发; 挂断 DELETE 通道 → BYE +``` + +## 3. 端口规划 + +| 端口 | 组件 | 说明 | +|---|---|---| +| 8088/tcp | ast1 | ARI REST + 事件 websocket, 宿主映射 8088 | +| 8089/tcp→8088 | ast2 | 第二节点 ARI (profile `scale`) | +| 5060/udp | ast1/ast2 | PJSIP 信令 (容器网内, 未映射宿主) | +| 10000-10800/udp | ast1/ast2 | RTP 媒体 (各容器独立 netns 不冲突) | +| 5060/udp, 40000/udp | mock-provider | Mock 运营商 SIP + RTP | +| 40001/udp | mock-agent (每次呼叫) | externalMedia 目的端口 | + +生产外网仅需放行: SIP 5060/udp (+tcp/5061 tls) 与 RTP 段; **ARI 8088/8089 只对 agent 内网放行, 严禁公网裸开**。 + +## 4. 从零部署步骤 (全部实测) + +```bash +# ── 4.1 Docker 与内核调优: 沿用 livekit 项目成果 +# (Docker 已装; /etc/sysctl.d/99-livekit-sip.conf 已应用, 见 §8) +# livekit 运行时下线: cd /opt/livekit-sip && docker compose --profile multinode down + +# ── 4.2 部署栈 ───────────────────────────────────────────────────── +mkdir -p /opt/asterisk-ari && cd /opt/asterisk-ari # = 仓库 deploy/ 目录 +docker compose up -d # ast1 + mock-provider +docker compose --profile scale up -d # (可选) ast2 第二节点 + +# ── 4.3 验证 ─────────────────────────────────────────────────────── +./scripts/health.sh # 容器 + ARI REST + trunk qualify + 活动通道 +./scripts/call.sh +15105550123 12 1 # 端到端外呼 (节点1) +./scripts/call.sh +15105550123 10 2 # 端到端外呼 (节点2, 验证扩展) +``` + +配置仅 4 个文件 (~30 行, conf/): `http.conf` (ARI HTTP 8088)、`ari.conf` (用户 outbound)、`pjsip.conf` (出站 trunk endpoint+aor, 无注册, qualify 10s 探活)、`rtp.conf` (RTP 段 10000-10800)。 + +⚠ 踩坑1 (**最关键**): ARI 要求 Stasis app 有**活跃的 websocket 事件订阅**, 否则通道一进 app 就被立即挂断 — 表现为外呼 1.5s 即 BYE、轮询通道 404。mock-agent 内 `AppKeeper` 用原生 socket 保持 `/ari/events` 连接 (只丢事件), 控制面仍纯 REST。 +⚠ 踩坑2: `pjsip show contacts` 显示 `NonQual` 只是首个 qualify 周期前的初始态, 等一个周期 (已设 10s) 即 `Avail`, 不是故障。 +⚠ 踩坑3: agent 回发 RTP 的地址不是 external_host 自己, 而是通道变量 `UNICASTRTP_LOCAL_ADDRESS/PORT` (Asterisk 为该 externalMedia 通道分配的本机 RTP 端口); 拿不到时可用收到的首包源地址兜底。 +⚠ 踩坑4: andrius/asterisk 的 entrypoint 会对 `/etc/asterisk` 做 chown — 按文件 RO 挂载没问题 (chown 带 `|| true` 兜底); 不要挂空目录, 否则内置配置全丢导致 Stasis 初始化失败。 + +## 5. 外呼链路验证 (实测输出) + +`./scripts/call.sh +15105550123 12 1`: +``` +[agent] stasis app 'outbound' running (event ws connected) +[agent] externalMedia channel up: UnicastRTP/mock-agent-12981-... +[agent] asterisk RTP return addr: 172.18.0.2:10550 +[agent] INVITE sent via mock-trunk -> +15105550123 (channel callee-1788153974-699) +[agent] event StasisStart PJSIP/mock-trunk-00000000 +[agent] callee answered: state=Up +[agent] bridged [em-... + callee-...], streaming 12.0s +[agent] RX frames=500 avg_rms=8374 +RESULT number=+15105550123 rx_frames=512 rx_avg_rms=8318 tx_frames=512 bidirectional=YES +``` +mock-provider (被叫侧) 日志: +``` +INCOMING_INVITE from=172.18.0.2:5060 uri="sip:+15105550123@mock-provider:5060" +ANSWERED media=172.18.0.3:40000 codec=PCMU peer=172.18.0.2:10528 ← SDP 协商 PCMU +ACK confirmed, RTP bridging +CALL_DONE reason=bye dur=13.7s sent_rtp=517 recv_rtp=512 recv_avg_amplitude=8167 +``` +**双向语音证明**: agent RX=512 帧 (被叫→ASR 方向, RMS 8318 = 440Hz 检测到) 与 provider recv_rtp=512 (LLM/TTS→被叫方向) **完全对称**; 挂断后 ARI channels=0, 无泄漏。换真实 LLM/ASR 只需替换 mock-agent 的 TX 音源与 RX 消费循环。 + +## 6. 运行状态检测 + +`./scripts/health.sh` (实测输出): +``` +== 容器状态 ast1 / ast2 / ast-mock-provider Up (healthy) +== ARI REST PASS ast1 /asterisk/info (200) PASS ast2 (200) +== PJSIP trunk PASS ast1 mock-provider Avail PASS ast2 Avail (OPTIONS qualify) +== 活动呼叫 PASS ast1 active channels=0 +result: ok=5 fail=0 +``` +持续监控: `watch -n5 ./scripts/health.sh`; 呼叫级状态: provider 日志 `CALL_DONE` 行 (时长/双向计数); 容器级: docker healthcheck (镜像自带)。层级与 livekit 版对齐: 容器 + HTTP/API + 业务功能三层。 + +## 7. 水平扩展实证 (双节点, 零配置变更) + +**① 加节点即扩容**: `docker compose --profile scale up -d` 拉起 ast2, 共用同一套 conf, 无任何既有节点配置改动。 + +**② 3 路并发分摊双节点 (发起侧调度)**: +``` +$ (./scripts/call.sh +15105550123 10 1 &) (./scripts/call.sh +15105550123 10 1 &) (./scripts/call.sh +15105550123 10 2 &) +RESULT ... rx_frames=401 rx_avg_rms=8438 tx_frames=401 bidirectional=YES +RESULT ... rx_frames=398 rx_avg_rms=7991 tx_frames=401 bidirectional=YES +RESULT ... rx_frames=404 rx_avg_rms=8077 tx_frames=405 bidirectional=YES + +$ docker logs ast-mock-provider | grep INCOMING | grep -oE "from=[0-9.]+" | sort | uniq -c + 2 from=172.18.0.2 ← ast1 处理 2 路 (源 IP = ast1) + 1 from=172.18.0.4 ← ast2 处理 1 路 (源 IP = ast2) +``` +→ 加副本即扩容、发起侧调度、双节点各自满血双向, 三项均实证。 +资源基线 (3 并发): ast1 1.4%CPU/53M, ast2 1.1%CPU/52M — 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.conf: 每呼叫 2 腿 × 1 对端口 (RTP+RTCP); 800 端口 ≈ 200 并发/节点, 不够就扩段或加节点 (加节点更快, 见 §7)。 +- pjsip: `direct_media=no` 必须保留 — 媒体必须过 Asterisk 进桥, 否则 externalMedia 拿不到音频; `qualify_frequency=10` 快速摘除故障 trunk。 +- asterisk.conf (如需): `maxcalls`/`maxload` 做背压 (默认无限, 生产按容量设)。 +- 容器: `restart: unless-stopped` 已配; 双节点各 ~52M 内存, 3.5G 机器无需 mem limit。 +- 监控: ARI `GET /channels` 数 + docker healthcheck; 告警项: 活动通道数骤降、trunk `Unavailable`、容器重启、`UdpRcvbufErrors>0`。 + +## 9. 生产化清单 (Mock → 真实运营商) + +1. trunk: pjsip.conf 加 `outbound_auth` (digest 用户名/密码) 与 `from_user/from_domain`, contact 换运营商 SIP 域名/IP; 我方 5060/udp、RTP 段公网放行 (阿里云安全组)。 +2. ARI 安全: 8088/8089 仅对 agent 网段放行; 换强密码; 可开 http.conf TLS。 +3. NAT: 阿里云 ECS 1:1 NAT — PJSIP transport 配 `external_media_address/external_signaling_address` (公网 IP); externalMedia 的 external_host 用 agent 可达地址。 +4. AI 侧: mock-agent 替换为真实 STT/LLM/TTS (RX 帧喂识别、TX 换合成音频, RTP/ARI 接口不变)。 +5. 高可用: 节点 ≥2 (已验证); 调度侧带健康检查摘除故障节点; 需要呼入/注册时前置 Kamailio dispatcher。 diff --git a/asterisk/deploy/conf/ari.conf b/asterisk/deploy/conf/ari.conf new file mode 100644 index 0000000..adc1155 --- /dev/null +++ b/asterisk/deploy/conf/ari.conf @@ -0,0 +1,8 @@ +[general] +enabled = yes +pretty = no + +[outbound] +type = user +read_only = no +password = ari_mock_7c3f9e1b5a2d4806 diff --git a/asterisk/deploy/conf/http.conf b/asterisk/deploy/conf/http.conf new file mode 100644 index 0000000..5b2a817 --- /dev/null +++ b/asterisk/deploy/conf/http.conf @@ -0,0 +1,4 @@ +[general] +enabled = yes +bindaddr = 0.0.0.0 +bindport = 8088 diff --git a/asterisk/deploy/conf/pjsip.conf b/asterisk/deploy/conf/pjsip.conf new file mode 100644 index 0000000..76bdab3 --- /dev/null +++ b/asterisk/deploy/conf/pjsip.conf @@ -0,0 +1,21 @@ +; PJSIP — 出站 trunk (mock-provider = Mock SIP 运营商) +; 仅出站外呼: endpoint + aor(qualify 探活), 无注册/无鉴权; 换真实运营商改 contact 并加 outbound_auth + +[transport-udp] +type = transport +protocol = udp +bind = 0.0.0.0:5060 + +[mock-trunk] +type = endpoint +transport = transport-udp +context = from-external +disallow = all +allow = ulaw +direct_media = no +aors = mock-trunk + +[mock-trunk] +type = aor +contact = sip:mock-provider:5060 +qualify_frequency = 10 ; 首次探测 ≤10s 内完成, 重启后健康检查快速变绿 diff --git a/asterisk/deploy/conf/rtp.conf b/asterisk/deploy/conf/rtp.conf new file mode 100644 index 0000000..fe0f437 --- /dev/null +++ b/asterisk/deploy/conf/rtp.conf @@ -0,0 +1,6 @@ +; RTP 端口段: 每呼叫占 2 个端口对 (被叫腿 + externalMedia 腿, 各含 RTCP) +; 10000-10800 = 800 端口 ≈ 200 并发呼叫/节点 +[general] +rtpstart = 10000 +rtpend = 10800 +strictrtp = no diff --git a/asterisk/deploy/docker-compose.yml b/asterisk/deploy/docker-compose.yml new file mode 100644 index 0000000..dc12fbe --- /dev/null +++ b/asterisk/deploy/docker-compose.yml @@ -0,0 +1,47 @@ +name: asterisk-ari + +services: + asterisk1: + image: andrius/asterisk:latest + container_name: ast1 + volumes: + - ./conf/http.conf:/etc/asterisk/http.conf:ro + - ./conf/ari.conf:/etc/asterisk/ari.conf:ro + - ./conf/pjsip.conf:/etc/asterisk/pjsip.conf:ro + - ./conf/rtp.conf:/etc/asterisk/rtp.conf:ro + ports: ["8088:8088"] + restart: unless-stopped + networks: [astari] + + # 第二个 Asterisk 节点: 无共享状态, 出站外呼各自独立 (水平扩展验证) + asterisk2: + image: andrius/asterisk:latest + container_name: ast2 + volumes: + - ./conf/http.conf:/etc/asterisk/http.conf:ro + - ./conf/ari.conf:/etc/asterisk/ari.conf:ro + - ./conf/pjsip.conf:/etc/asterisk/pjsip.conf:ro + - ./conf/rtp.conf:/etc/asterisk/rtp.conf:ro + ports: ["8089:8088"] + restart: unless-stopped + networks: [astari] + profiles: [scale] + + # Mock SIP 运营商 (UAS): 应答 Asterisk 的出站 INVITE + mock-provider: + image: python:3.12-alpine + container_name: ast-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: [astari] + +networks: + astari: + driver: bridge diff --git a/asterisk/deploy/mock/__pycache__/mock-agent.cpython-313.pyc b/asterisk/deploy/mock/__pycache__/mock-agent.cpython-313.pyc new file mode 100644 index 0000000000000000000000000000000000000000..18e1c418fbebe57fb3318a00bd76d156a143bf30 GIT binary patch literal 21667 zcmcJ13wTt=mFB(ue(BwkkU(%HB%lGk&;ul%0`ZUm0lI1o7Dj2*Ezw45$+ue=+D;W@+1?>kBM>O%y{|`P@UkDu($-dQqG<9=KD_GR$)f z&j<|9@~Yjez)DY*ppu?yK`lKsf<}631+Dbd2|DSi7xdE8AQ+^lQ7}qRlVFmbX2C2y zErLaQrU~iNGb8@Y#8X|!2wA-5NH(uMl3lCeb%o4B8X z{sN|i;mw5%Z>dwsJrWj3acPK4Ps9~SaT$oqJfdo=RgjQ;)~VyU7xGzX%_2SreT=9{aq)xrWUTL$K~;H|1flpYx4SBg4=td zvZ<|w`)Fji!ZdUF`!i#2-}sBiKMJ4WR@`{@<(c7M&Hmt}kKTNhYunDZbE~+v_WkR* znbDujJQcb=F?4<6r!&u9{OH_8&WqArKJF3Mv3ThETTgy;E`0sXXRp8WgBwpj&hb8n z*X`vT0|VU6cLFyrJ-ymA6Zpx@@CA-)zy9{~H=cUx<~uLUJo#e@=A-AoJv;RL?Bi!Y zdf}Hhp8fHSr@wXM`#-lenmDei0tNgKaa^TCa96rcqC4KrD0gH3pzm1!pvOs_qyF7| zk}_sfAh@sw|=#nyFT&E&6j_4`%Tw^YoK9-VDg4P}ghbb|t0j{erv4 z?Q!^A+(&PXa{KS&TlRB4VbF7e`iJ>*^QAwV8T~U9Qj4Mb%U5o``jQ+YH|^$kf5F)Y zQ2CqZKKl7N^n7OY>3E}`{IQJ3x(b?SlrR$*YvkGvZHyOi@fWi%jjgy5zI5ZvMXrf& z<8J=?EN0AVu5B-m#$3Gl$|c)+uKm!)+S;mJes1>sFK_((yFj9nWzU|ymF?}kl&7!Lee~kRnOA!$|TvxSj9RQShQBL!M z+kh!C^U51D7s4}vXErs~R#jD1#|N&w+%wpB%q4KEs_Uw&>gwvMs@K$T9p&ZDLBZj3 z_j|bNx(7|OL^)c!;`Bp$=-}(|yVC~mIi>@2Sav+y5vOGg(*^=F*MZ`x>f+Gp46m+H zSF(2I0J~PlYj|xHThD?%F*aSRU*92hne!XUb~TpmZ7kzO_CZnG-8<+#-fmNi+5y4s z@rfFryU!&Wygs3uLN-=p9bT${6VvWOoT%>Y?-9&MvEt`l25*RY+8((7#rBI2KL232 zF;-L`S+sWi;YdzHB<;S5rVrBYyV@R0yYCJl8IkphI>F@|6g*vt&LW0%b9+Rs;PCXgY#IT5VnmJC-Qy8im&m%kl#}MAKpVekCKA}`@^tk( zUBVK?{TQ`rgL!zzxotsh@Zqz2V}_EUrtyb=QFon9^v6 zt&g`1HBA|eLrv3Z*@5c7v0%yHr4Dj^R zvF?%nP;t2FiuTg#@zQ8{PobkMBCP}nsX4r8mZ0B}Mk zdAa0aa2bIxB;=8o1#gzd=?Y`h3LHf(BCiOZ$o5IN_tI1%eBXs|r%OV38KvRpy${|H zGi5XltA^W0fJmTj*f+9ns5$uXg=L}U3rhkA#!Igp4;N3F%z>sc?a1n2Y1FtlR2?;x z4mJHTD|fhQcuhpJC|C_o#+}ONy$g9@P-NsS142^QkixQHWe_?+ba+)Gx`lN>$d(+h zR>%p~K(~ZYHLoehIww>Di7I6kI2srxW&LP7h=@t9jao3rfl%pEAS7n=Zjf>y5CaC^ zt|4D(M`{@B+ThXJ**2Q(b6f{45XsrKN?zs3F?)|Q46V}+8cP|9RZ?jR3~ol`nyNl8 zSIJe^S8-lP-#{<8kuB>^?DDVXyZn-wK*R6aPw>k zQS#fCcec0(Fr|bNM7bw0G}6kq8Aao<(>|9&5FDq4g-8?C5*_Yy_>PNeFP6A!z%A;) zd$~MLY2}Jq%HXx>1R}bk8Z4Qp>h}nX>5cV?>`5;Tg)|rAlhNYq_qaL-9H)Ex9ZrFG zUdrnA!y96*n{v*q8rn8xvJ87iRs_1D#)6@(w7_egeE@6Qk{&oXRyeY6nEm7WjZbe2 zXwLc0pZeaZ5PP}#<(-#yUNVQfqFnWOQ#teXUadw(gL%|6Vmh}TRAtlO z8A@*Jne+vp7#K@dAm>bDM3XDc-YzsD;b@7GyHg79Ipzq10|?N>vp9IDz=Y#Cz-z$l zvp9HY1znHYrWe&sf_t_a@Ld|p8h3;Mn$ve)fd>M@wgKUD9iNy1cBVIx&Z{~SCACpG zb6g2ySJ75$SEn;I>Ne2$xlUPU%KU1hU2=|7U}#mL_C!6{tTnA7Tj|$r1ux&liEJY} zsS~gvy9Y!K@hPH?{Q6oIHXc%yK zy^r(@&N=&op|jKNar-(uMGFvGp}eg@PVWwQAX?T1=Xa0o4yj|-(qYZh#wqKf7f((u zsr+C`<#_G*vI%x##YADGs(G?%ceHAEtZHv;N$cbidvu9Cw&eb(_5NYaGe&91TkjgF z1hvY%@XZp|Yt*1X>I4M>S+#30Udp_gmqlu$9d(YgC#W3ZU*ZQeDW%Lq zzvdgmLruK8K?QoGosS;z8uTT7PQ`dtI2Mjs%zQl#!l#-p?A+H#U9 zN=TVm9zD{DH=E<^`UL*4xpmM6C%HD-1Lru*b-gl+L3^26l~=_>K)|Q;qLe9SRwsEQ ztfOuf$1Gtu{Pi#!)k~Oewm7w?g>x#wYa4B+bDW-WJc=F1U^jeuFBda_Alf&_{lE5e z33}-O)Dy!_n}OCP2>pUzJLu~!UnlAX z*FdkM%Vo=!>3~i?Wr5*$2p}d#yvB(KARiudd3~b3EgmpAoX*bUE(drfj zP9$ms#ELqv&*2;NikgSK{oqAuckv2)3DCH!y?y_7YKN#H7DzO1Kh@**ot1xn z@z~Z$yN=N}z=^w;Xk?Jqy@(gWk_%bi!P zk^CK>YMJb{lUeon|JP4+OlIEbV6Ox7|{y(-4ur z;QYz4lfSCDta;h=k||tsIb)(ZlHc_E%%(q0=PbIzut}hhKAiJv;g!PiofC&+ww;kB zyCV6!CNp>4QWMa>-7-=QpLxFqe7(DF@wU~>Z}Qe}Td4kRmJ#9KE;Mbk>3>_Mrm#)& zS1;b7RljG-+`*{dTexz&S@k4&8w4i=09td><|54tfu@4w5o6#bBK#fZW0gvKi2c-%p>1Jh zP=caD6Ap7na(5_YF8=}UgM^b>;y@tq)WX-3coM=+J$xf?5KKKrn@Kb_4GipdxdvRq z>@Jiu_yRnVzuo-F<2PRWR;sk@#)~62e)-nT*Plp{>D_$grP<(Q{p@`oS!p7@wJ#+If#wB@UMHhJO zgd@&Jjo+OQ;Nk3$NN7v!v56?AEj4F6+MeL4q!U?YI&{Dg`i?ex>YeG1&%Ojj9pC{fb!J&b-Qnk+RM^kJjVILj>g+&ACjZ6y- zB;cVr2#VA%Ko~;rq8hxnu#Lc*30y@Oj?Sgmf}!p%(5KBX9;)Ff59KFezpEe{i_?s+aiWV#C!#6&oxYy zS4|iqhLw0X1WM1@g6%Ip2qrDTXWexuq=AqI_7(VFHLxYjSFylJ)>&Z>zKjv}K%3B! zIP$g8iGPl3OB~*k31FWxlW$5$sgD>pNC4pfUzJQJ{5cLxCb4YuuuQ~w-c_#_aCvjG z<27EwJd6}lbTwM3&2gE8okg6nM(P0$;0dN^K2`{9*u3>2cAmOwA%5b`4YYg`SOw%f zEK*zX7yDpW&BwqfOr?SaEKrL_&8OMb>p?V06mX9k%Dz|wbTM&43uayeZc1`a9FjVg zE6{v~0^?#2)k+L{vNjNiy{M7Tw5!`}N{-}{&q`{0_-vveQX4ViDdp$b)nH?-8fF6> zD-tqP7epiTVro^gd?f@XINlGA=~%7AwpA#h)@}7Ju_?D!% z`K`0G5d7Ai+o2}`H<+rB@f(*p2Hc$|T&Fh*jRcj&5$jc1T~VE8+6DU6$oW-iCcn)- z&?A76YfSDnX{K%cP@!}YQ`QK5fs_u=9NyJce!uG&FX@}gce_rbj3sHNyQ2@f1h15_ zx+cw(W}2n-E$k+b*0*%Ll;R0T{=!GRTvu-X3Sg^frNHo=1^#7k@{ zitq0ewor7E)N3Ce?1!>w4#D;kkU)8LvSLll>H>S~b9Htl@FNUTy50*RjH9<#AZZaX zs9h(;j3g3eepl-q=ym%!Z%{+ z6_e?8(e%1l`r2XjR7Tm)bgx>jSSD9AMpraW)V{OsjdigV&C!hJsjNk})ofxW-EJidG4;j7laE_%1v^@H@wy`zG|X9lC|}=R+YJpy`^Kao7mfWb!HQL%fMKdhrA(IWYzvi)`8nb4GI`m ze_~>C7fxDNe_&l5cE0Mq;-0X+Q}9MXqQ;>rsv zUs!$3uzV^bdss7-UJ$az(#uFj(K6O@zICiMR2{RH4R8Nr&VuQ@vdFTnvAk`O%xzPa z%#ppr&9{t9Mq$YM`}DF~CdQH-v|cQ{P#D@7J`ls1fQ=TciJ8|-XBS6G8zxF8d{@;I zE8kcgDcKdvZi!f0rj(W}iW(LL+pif)rVUw>h6NuO7KAh}n=hHio1=z?NyFx-Ve{_| zP117cqRFke#9sU!VxVtlLI2~F1ave(i|nkV&}fxtYij;@IrFn$^OaPNwW-i5C^ZrS z6p7|07^kEr2{d=9M03BaU1-r>YMTn;Q+2DU5oQbos>G!oL^9Xl95}LTc-tqK1=(`m zXL>OaGINGa@@OSS1Jw}wC~-~!_(u)TL2Lt&1M}qqMP5bw4hh=-U2rP93cL(zrP1Lv zlAXay>{p3VzlZLR_>k1LqE2$WP8n@Q=n!v@dS4~VS4-t%=lsg@_4g>>aIf-{bISqY z)%>}W!Z(PNjHH)=x{Ue`K^kw6+J?6t25)BR4Wv0mb&z20<-U>+<$^_NQJEe)=}jBg}DprcxiD zwH|2>99=p{fh9RU+pgy`Y0iQx&ACSjR;7fwhqWqZ-WtK(DDw|IR-%Pce_vuqtd`(j zAyz}>BsK&tUV{iRDL71p!3rdpoZP|B-k)Yc10yAEA8bDGns5B!P<6-`ZW?E=Y`^sA zPc|d+U!;1TC@f-Up8~j~--Lqy^neR0WHT=r3&F=h*C@Yi_c+DOZL(cWGX#M`|7n-g zuanGwpo!o_ecb3q)JUgf#+D(RqmDriC=7yP>jpgoj;<2|je$Vz6e=m%=z?8@3nmaQ zc|Ul9QUc^1BJVJHN8pKSFm;kmL0mj2s*Y9RS1qb~J)%CYzXFFh08uFL=yHlG*Gb`r z$a0&c zBDTY^{KL}=76!XRZDHM|2SR(sAq0KB{Mc3?$>^~Se^gwJ) zN2K~lWYL4y799P!lF7;c(|60Ss zziXjuP@J@%(u% zK+3m zVgxrlQQz6=?C(;9v~)I+Xo4V-oC$h;h`gobQJwPI?4=NSyYTat!21q!Tg7U(u^;O+ z+JcWQ+1jn_$1C*OoR6~?YP0^iQLD|8vs0!n{1OgJS#o!l;W1Zg3(4#pxvik@(LoNa zYCAPcwCgU;93v#p4Qdd`KZv~~hQyofAX8DQ!!9I^-xZZ~hVywvqYtM}FG&uu18hYG z;VL}oTpL$d$YMZq+o^>7?>U0N4uc__Dq9(Dju=XC+A#&1pOabYyPAhxNNyg91D<0J zWfCrA7jFz^&bx?E>ZfpO{bEM)MBCX1X@HoGw#8o{=2d%D%kc)J1MM1_6^2MogFX15 z#;dACPp}(9nD5af?M{-6&wkZjpoMl}j}}bj{B>%TIy22m8K2 zc_kKW9Y9hTEDW26O+j%&g%D@1Fp(tYiXrAorkO0H^2fZG<}Xrh?5WgAE?lIz<(Sw` zhtC0_Rv(=nUG9@E=RDk02*6&FpRuy z>!1A$pc1#~H^IiV8>PFv8~sMkak8gV=<{y;OkGjc4O%L(H&+R`3?3|b2K$5pdZp7z zTzw#FsD&`ra5@D6TrVU1K!Sx4$6N)1nDvFuh}QYK_yS7s-i8M`kU4$WccvjAoY@qt zC#eO@aRRk5!-8N_)KG*cIG=zJ@CIzd`l;m=Veh!=(qpmo+Ls;+6+iuGu=?3=TzYI+ z9Z9cM@;1Z_MVE7;r46s~vC@WEal_l&uCnjdT=o8TU2MaiXz`w?VGls2%o(GNBaMNc zsCjYJylkpu+30;E_XQq_T9znD&Zv1|uq$dV4jqV@mrX-2yE52*@qr5ugf-Ek^2nl! zuzN!Dj_D24)rRQ0z3)5!_MyLdD02A0XzS7Fx}%ZKW0AF8kxb|BEiUYCmgWEQF;4Pz zl9=q5x!W2w!aL}HtIaUqekJP?-vfZ~cjV2tIlWB@5>QtvX`~hP1b%lmfb<3kphN%^rK79v3WH%o%NZu|WdBD7K z9z9OBWC~^AAp?NKA0{C!Fef;Uwwz?2Q@3g@(xIhK2U4iA>+JekHK{2dg83XylS<{x zJq(<*$riM$U_S|X-uT5DIHq0qWqKc^>q{8})4CJC1;t+pt(^L#LzJmiT4o<9ZH4UU zeT0So1W#1?P6>aDK$ntAPa20D{O%OO+r*}BOzjEzSE;;fW22-2#^Wj;5~meEjTtud zI{J<|9h<~OeRQLsv$wy0Acb(mQ-4a?Xp8vM5OYoys5>VcH~IBd71AHfHB8dct&>*w zUF|d=*IE4TJcKP|lUkl1DPt1fSOAf5kZ@rBD5q&p>ImKoTAhGqzp`xDiYK5`&EGVM z1|$N>nc-06Kb#w?K+7i_!UA1UPlXswhD3%@Ri98`9D^wB9j19KzNm0 z8)}>DfI&UXaLDvkvLD7uG?LT_@>U$@$$7Nx`uPmG4z<(#SK6%L33J=Pn-q2eCbD<2 z0P=f33-fQpV|T;GFT444nyK=uxvQ8Y?>lDy^tAc}2ND{~Q*m zTefB}(uUnj%Z8kak(QKSSG;xWu9g-_Tp%TW{>{(7sb&NbFaix(&!g5?cK*Hp+-4ZB z>SRR=tc;*=f#Iq^e3fKjD$oic)1IyqkPkp_;_C-nAYbeN$J`4G3FKFMT@TB4tQxOZ z=)%ITl6FE-ecTCaMdjiIbS>C@<+VzSj+5Nbki2H4ripp;Emz;Ep8?$aKj2B?66i)o zJx_Up2SVBySq1&<;HyV}dUT>`VlY;|J*=9}$Q}L0Q{MyAbjkyEIrn3vqmw&H3q={vhO=`<{huuGYMl#!I%>I{8RbNcblzQ$vJPWIN@F z^WUk;tMI(EOgZjr7{M*Fphh$$M^5+zsXrkwIvIUNOxaY9HIx%>n$Siqn})Veo3kg) zOFu9#jhV}anjzTAMqao-3=zE(kC9!iOf1R`EQpn#`S-S zr#xf%UtZ{*nJ4SFE@Izex2iR7XCU;hh1;rU-d)XZ)oR`?*HE}R2l2nrvIzf1pS`s} z{hIjHkvgd;DE+IbpAr_o)2&DZjERje{g zC0ravo9ZRe2AFrdaUKoRe63ikw1SM&J*IzG&FlU@)NHn!61*&=kalz84fEY@O0Gh? zS!r>t%40~ubA#Ok+$^x0@79#2Z9a4)b-t(cb<4fl$~Pfp{{C%kx1yKZ zY7$o#u*QiElH8Y5J(M|5FUpbAlZKg<{_id)O*-QY;=C~b+ND*1R$|S>SGhF%2{ZES z2TinF=3Q~k3Qj%g`_ciweLY%7Un7tx#||ZxW9;etPH829vh1r)Y55u%6RAFE&aPz6 z6xamblF-RL!Y=m`J=lw$O7%U0wt)6XZE3NmOXHwUDLF}N@1xXInsEe%I`&S5n$YWI4yd zN|uH|a)(yTU!p?e6ij$B+H;e5#8_qUcH#p(nfA;C&7d*O;15!~J(EPRpl+F9d#-`Q_DuBR0ecpY`;hz*l=C2!gSeyiOun=6Tzt-C@e8Fe#*qIS%6DK)OyJm# z5x)gPX=kh7l)M`6&yC-WuU@lO5^u=X)7{{zK33;gvQr)8sslJ*kB}RpzADd0ZL-caaZqv-=JL2_ghl$=ljj&<7xsNlXUkz9!}opmkbIZg?SIE=U+wBE_v@twxXS$&*zzf#@L29&yqt@l$+#7=5}xdi-wkuYj_$`HtU#w;tY0H1IP zBwa?-CN0_C@v3$WAqO&43TqQ#bvs?~M0~9pYA0~n5~XI_Rr6U~9by%F4r36|PSRIE zHcIu)aY_DD(htXX0s8WR1Dq>ax9kP2Q6!;s8B`~KJmWJ{W#W@3pBa_UrBnjTU8q9R z8w(Q%+4Mpw9wLjoQ*?98D1TU`)9n`r%%+3wJ2sOdESW|Sf6>x z^y-DICRX(>Yii2SzMIa%7d7{1}DS#^OI;*KJ+1dNd^~Y>8mjsNh>ay?XxbAtm|DPgomRKuuSFaXu@|<>=Dj>2hCxxwn5% z=yGwm^e$b=!iQDRP-PciW4QPsK+%Y;aktXo|o8em;7--t@b z0K#<3g zKaqExvcby6BMBD;68GB{h$;u}(@Sk8Z7D?5gWcWkQ-X%tuELOgjbc@n=)MIr#aYNpBag_bkWO9(JV`hU`R5$BMRIX_ zA0oX%Gk~dJ?`jOs|HiIo%$hWE9~ilqv1DlLlt%a6eUI-8?1`*w8rm1tBqPEr>Afxm zU^_+ZP$G7Q9-#NC_{M2~G)t&F}cz<1FYkNeq zWT<6$FyO|u+w}Z^@BG)tzK)B%$1fcp*SwzgTH1u~ov**~^;p_o6rwSGchBQ{zTG;l z(GRZ%=&7+&frld8rifwlHO-b!HEON(QVNzIA-g^!-0%pEts=N`igLNSQBmvJHrPfmfGJN>aHW!vh$i@R|-yo z;2m(e;$p{zj^N&KR=D=n#w(5C)e-C3II2STk&PoRdW;uL=i~)^lld#6`72^MHry#$ zvx=~$Nnb5t&6LDiXw|qXeCk@!Itgio!P<+B7aD`BLr#>x{>u7r`MCe8?drn!n<9qR zYnpv@HQPFR`l-`r9u2ln7F9&?pHVTi3%5w1h%lv(+McqVSrs(IOx#qQqsdLML$R!e z5#6u~GauhOX_&I+o!>FGBe?eLzLD)yX&HgWSXvQm93xq^*DQ6=HUyf_Y@D)WO_KuR zDVS+oGc5Qdi-9FYypG?SxLetbAuXVJ;_H7fSRxr*C^w`JJse51T{End`coV#3*|=) zwriS|H1^7Dbcc1ID7}2iXrZt#`C4ZJxyX(_M2Sqn(|$CBmCTpGn~>Z3pswu|pKM ziafFf6+R^Iaq@mj9<-#4Fif5vo@l`}VSJsgLi)IZbf+v472hl40Cemm<7sK%lJ-Lq zdP`Wj@v5j7y5nuUH!e1d4jxQ(JmmHvMeQdh*=nUlQ3+M zFVJri?gI$^8SiWGfC`qqtx0D!pDt$DPIi(xn)?6WnCci){SVB2w{&{8i4D};V(1yH z`9yl&%GR(&6p07L;c?4iVDkd@TMQn-!B6P%X`YpBVQ(|!e_GMPvRSto@;^DG%4ECQ w;rvNsanxA+2}5rmm*%rsAJ;lrHt*IRHKVZ%`(x^YKWGe3?!f%h6iQwCKN9C`00000 literal 0 HcmV?d00001 diff --git a/asterisk/deploy/mock/__pycache__/mock-provider.cpython-313.pyc b/asterisk/deploy/mock/__pycache__/mock-provider.cpython-313.pyc new file mode 100644 index 0000000000000000000000000000000000000000..f98b2215129c9d7412d5b67b7ae9977905b26ca9 GIT binary patch literal 18950 zcmcJ1dvqJudFKo;cmo9AqTZlLJz$A1iF#3@Y>K2L+9oA&1W6VZBM<>mut|U#fRY8L z4jrd4XxSlX%N0~RHIz6t6{i)Qq+RB*Yb6_R-Ry1;Xviirev= zy-9&+S;9$$qV>30#9P%1qJ0VXh#(csDsixqVlhu^<76czV*c?`vEaA^Jr~s*#X?qE zCKlmNCl<52#bOD@Un1Tkmg1h5mWqpxm!h_e)h-hkV}vF6*-&Z}m$K63W!x%GT(*Q0 zm#c3qt*GUsmB(GjOOLOp*BS5@$IHbPYjG8Ir&oz9@r(-51z0JT1FjZV0j?1%0N1k7 zDlyvX<9z3V>;lZXY}P9`lAE<$gO+RmK+Ai?br^a5@eQK8;V8%X7IWMR4t7RW3`XyIe$Kkd^ts7h)J$rTL zg{P`^-T29ckN@hz2G@<@XKy4%(CNnTi0eReyQc=A$<=x2K(+P5pZ=*U67}~6PoUq$ z)QA7=Tfh06Kl_)#uYEW^_~H2XsEzC6*RS3f{?3h8pZoa2=to0?AH6(uW8`NyhF@Rz z(X&Gzee1g){UqVSJa0Vv7S|4nB$ZegWHSLhiYAadT zT)U-KHmi+gLwzlhZ0*!K#q@R6o4aRk3-J;Yes?+nUgWy6G!0DRyRvl7ah<~KE=nFa zsZ;4Dj}!Hq^fkPPJIHS}ih^jU<+t$IRnBefh*`VTcXY+ftM@gmZf{sE%KVpP!^z%( z$f-kay=>^0f}yA^M1y?+*&K;VCn@LVW!@K|8M)*VTCAG+D<&>neZ^>_M{Bc3UxnUPn58 zv{1CLc6fuMI?-B%QL8MNi;a3#q14TZcFI?BZb!$=pALTd>CEF_XeAo)T)0kkWju!$A18D50H;;Q83KX|pZ9=R$MRkUafDd&azKuxae{4>w<_ANQqN z-Z$A2<%z=h!70HIFMhIZuxZL{8El$%6ej8tCz7k)b1aYZQ#N}-Nb)0=;d_VfP4UU5 z%bRcMxx5u0L?e^7xj#?1a{{FJ=2=&k8Tu=Y#-2HG*J%w zf_zAr#|MOaK!Y9ubZF^NmXobLEwQozY{l8ym!WKiw&w~h_XcW_Tkd=4mdNg|1!uJ8 z4=`G*MT1q;=zf4fZpvbf>*9xoKtrIKcywsnU`z7!r8TLRODhux$11O!8ZDo)+7eA8 zhM^6~$_dMgRNaKRa=$3WSNH@(=O=e`+D(=I|%%&h! z>k8>S&{nG;fka&cdWQ5})Se|#ugVtC!dZsR#iCGU0~bSm$l%ede9*U@<7RPCLl;Rj zmO5z$^Ae<%aaIeaW6y0;1;oY90eTDs)p^X$gR`4MMvt!3F}E+>0YfZI3&LExi0iTB zF{x%gL`g;Qr*L~X_ z-51;7s;{rz7jtc{MT4qvNX?u)`LQz8&3k*4%m??*q>xkY?tAmTU_X#gszy<89Gk^> zfh>?-zSnP^F>d|`Usf6u0JRnef|A`ke~5%iKtsPP_KyF zB-PSF2cohr9LlnYsLY>zo)Bi{fRc*NTT$?kCwoK%2xo+QjtqvW{wqfg#HRTyS`5o_t9m&%dcKvMq zHQ!IG-!r$}GI4ojpO`s&L89nGFnWv&j;V(oTKdEHf3B;bC)H;v(^Lo~6hzvA8te`n&ZlV`FD%R_P-sX;&VIX?4YMX45 z0{y+d?!fN_ht*;Iz0P5sp;wmq28Nz)qij@|n{4U#Nzq_57=|Je>F*6jrJdB!05vk$ zFB^KpX9JSkA~jQ?a4O*Q%e+?>PK5pEP#FvbBC?4}15!lR^+It8Nlj=Pp(P;s#gd}G z(32toZ(pQG!t}XK)a*O}IBId(#Roi8C4diow>Yy!^qvQ+Y*+ro`^^!MJX!u=wJJkqya6rm$kj6zAiOB_ktR;7hiiY7Cg z4-=Zv?+!^0l4d}XJdPVcoBoQU<;1I)E(y)RqmSrZkW0jWdblfmhg5?;(Cb88WFh41 z3rI}WB#HcfvINCZBHJ2{uC6lg0c^6PGy#} zQIGO$misWT7a?ywI#?Atk3NsvB-HRBlZVH=%d_(y@|dA%t1|}|2FmgnVO|mij|Kf_ zqg_Q7?G|ot$Vni`Ce2GRcZ@l0Oe#5ttR6v=Nth9IcIytz$%i-dSn)0vkAQLXXt(gQ zt<-AvTtB0oZO-%U>hr%)U&oxjd7shOpuGpS_~-z$R7@VDRtov`8FO>0Y{~UmE%I}+ zIwuq7&Rr~6m?^1ai-n{Ec$l8`dF%%4q^KiC=5fu&>ROk~yJY>Dpf6@P9P~A~V#YmE zxUaz_3ny`f_7aW>QD0AE%&>!9ghOF8vNZd8d#l=78eG3ObyGIzj~ScAz-j8(9SNL{ zEpHBoqJdDf>d?9VK!Xc{us7K4BW_d^@%P8_)#?WWp`Pfe2A8xKlb7}YK$#;_VKf6n zHV~49N5T-OM#VsK+h8CEC2%o?`bf(nbJZ6Cqtj!KvLWzTRPsp&X^>q6G}@HBxQ!5_ zCIZbSg-)#}27BjVuO=)NSbz!qwwZOyPH`yV}hwO1ZEXJ%% zs;BIQ!w(ESkkn`F%aXAPyBmgqy>`;JZNj!KV}rbJV62L*@xtQ%1eT8YmB`m{{TpHL zFDUp`0SO-VP10sP!h!-?<5U{0l>pO)g;W4_L{U%TGsL$;efx}2HDy5^mGOI^;OCDV zGIVJi+G7w6^ceJh@}zJaHb!(cTttUAgSi+!`FRvM=lL)IJ;npNcGx&1Z$kQz zu|1*Wj~06Lqz;fYis(YdlbqW$a33E3n>R1qcuy~AG5O%W_xy+d^@X+r@aN6C zSDbSjTv;9R<2QbK<1fzNh+q88%+lT1)&cKay^HC)F>_y|Z@?c6yQF;>pBcMu zgKTAbgrXc+BT?uVk+Z?*Dd|CKvkgEt1tb1oPcRA%m*`K{LDdlY!%|e%_eUcVv@mEO zq?klxzW;7D*Qsa{;1J&M$&&#|ilRd&evuae2DuN4%aWBTUA%j6_f%ftaCj*EUfxP* z3O2`Z<4|LA>7;G-gl+Zc?yLK*?0fzGvC~)Dzb#EUoWoxl`ckrC(y=DvSTj{zdU4;# zz8CI~x4s$~-8@!sW!rcDBqhA z9%^r45Vkk+)=dSF2(}Kk(LwH4f4(DTKhs$2u!bASt>}`DW8jz>K4XR+u_A{x!fG5= z4QFBn4r{cr)`eaw-oz~A60T`Fe83g6Ijp`$DcaxX>u-n^)og?Jl}aLH8U?z~IIL~~ z9NZW3L*SU`N~V;d{ULAL0g0%~t(OLHkzEA}Etn?LB&~ipj z0KkP}%^$9QraD=WY|U6#!9>bimK2k%se;s|)alEcQ)@@dezbaQ(M^uuUN&8}JaHDe zm5ZhqtxPX5#S{0yy(v;{!sYxgL2)*?ke(FOgm-C-&X zWbs#jD)5+eirNPVC_0MI@0Vm#e;^=v!4EAgLrjr*0pKZ!%Es`cUdeY>wseO>&~Ewy ze%aXF8;;;AfKG4*K$Id|33{0iLAJ9zJRRQ1K%YcQ3vI1dw)FZUQKrX$s?fLhWNs$X z(`Yr%`E2%jgP~y5>y;gZkk#s~pPIe|04gdhO`J_xGKKfVg=bCEMI{#-o^MF{UTB;w zs+=gQ%oMFbxp}I3)3`rvUXQ#r(Q@IzROPGgsdXFD<_##ZC!!agm?~LHm8nSb(91`r zN|)fqHtI>XyxehZ^;pY~9+=v*H*MaFQVT8^q79dn@TzmH^7_VROyZ*ZdG`x6zG-Ul z^0@F^9wR?@A$u+9Im}gYt2|3NZTY3A==KEyWRyyA0z(9zAuvo}gaD*0CqWJ2q~{6H zCU(G{#sU!o#{yZs-oCJZpqKJ?ulMwTuQyxc^m>_2N}@3o_If1}S8|Eh3n3W|c6)tM z7~LlZ;QB!w8IOwkL|868p#owBOeB7T`&XUCutI4=)v|$o0JFLvF9@MUW4*2t+yUGlc>Xn%xOME3#q@2s{R@;bk5?1JdRC(t zzr1uBb;>DJmJRm@yC3y3i6^B{CsF-90MMYJ;dVb2ORm>4@2^}HZ%*(-dlGxnwxwxd z^QZ$rVW+cYAaMoSzdMKVt`(f?3rNcxKUB4{HberPLjN=Ti~0_O9iu=S)p#35&n9h^ zXbripSzN7C<7#SKXlY!h#?Nx4Du>1?hQ`CtQsZB_mMGOYSFV&~J!){mH9v{^E{!)q zZVO%H0A~qzhD|srsnJ13n6{f{^Hz*Sl4xF<&2QmKqjQUS>9hKEy5=>+t@_UJ&Q_?k zq-k_ks&~07X&>|6V~tu%+6a0ru&SuBp+?;xwF9oG0e*K9p3O_Jp6x0IU}mH$v4rzU z>@Ovh5-qYMYHgvhKMKo1{|B`)hdR!J_*JWpoV%*i_d|YR&Yy>YpaJqV{YUOid!J%O z^l)ZwFW>QXGuLr=s~&L#c^xxvW0JCN{{yiivdewl(FWH}#0xZ#26&&;Mz#E|*fLKf zAXPQ>Kpr)?+EH~Ndq`hXS6vGw?G)NcK>`N>Vg*ed;*r+Q))p6WFK^J_7&Eyb*@8aG zFhU;3t8$A^Wsmkc)pOOq|H?7g%%xjnV_`>mr}$I5@SZ>(Yb$hAX% zckJiKu05DOa5#OWYw}3<#F6gwVSm~mNS}OU@?>=4WHcQZNSn`G7tT%#j=2Bp``Ov^ zH;=~o>AaHRzGwQ9&6zycq}i1=yQXZ;;a$VK@J3y)bX{3JWp)hPhwO<%G{02kYwpYL z_skn^<6NQib_r*wyj99s%6}IjhVr!)d)$1A-%}}EZmB``t$cpZTH{*$nbFWa6}+%$Q zFD9GjSaM&)?$fKh^$$;+;!CvFx$kpsXy+BkBNk*w-dryO$KS; zX41gX!BoBAr=Q~3^S<3qjyVQ;iJmCw6hml+P$6b4-@SWZTN{gyU@boV>Zf1Ta}pU> zJ&zav_1=GqUAWhsFFl5P2?w}Y{hYQ^;pe-N19y%}1Kluj`hDGx%4YCiZv?_cdX1XJ z36QofoA!&Xk`$H{75%$Z{+9r53tR$$)5=+@5Q#|Lihv*;Ue<@gXJL3L$ELFWlwXpr zpnT3D)}W%kG8qw&?c`#+j{HNnffY<6IjTreMC@GjdPLm?>N{ba1pOuA6og z3_tYDL+88bq^tDe`itu|6DOscD1?4IUAih&|60T4hS5l-bW?oqbYa=W_0O+Qo`RS7 z#JKT|4cBznj{V%3KIq959*#HP)ae~1)A?l=og>ba{<33C_%HT1?Bo77Lq86sYuc~p zcR(KCp>5B!B@2c+o@+}sr|K^4OE!--Uk;{DOcyS`SoM5WO2`zhp0uw{+gIN#;))j0 z#9mnc+fVFV@zOgS?FKAVRdT(_e*Xo^1s~5 zx0VaP+*X48TTUIp0&D9E{aZ^}emPJ16#~j(X3%EJwr)7_qXBQ!7kLzAv`^AFZ5C14 zmla=+G)koem?cY^?QaKRE#(>3e6oTw*e8WW?+A-78DBELE0o`DO+fEb1jAoq7fCF1 zqF{bnmn~qH6rH$ny;_r5v42_G*Q$`oaZ8L$LTFxMT#vMSjZfwpXip{oFSih=9xug<1fv7~0F6+4(rr`N6Ju1wNM?Bn za_S><6(p$kWdJz1!~HlDW_a7f3NRyCkKyO(idN>>@se68y@`gG@QZvOfDGlVItY$H z3f#0cZ`l2e`~3Q(Ku6va)|K#KuCKYe{p$93<)n4#1XMSs$-QG;J6=B6%(OXD8joB$ zl59&I9X&aEbj&lpX*@FShAXp$FQyauwvo2Pj^rcIDi>c_JZc-;c}@5`C-}&{2hs=p z>A%kta2sp}29@Au0T*s`$krID%_QvDwS}U+}V)`>en-=Wn(VVrUJ?5MNN9$3OUAK^VR?w3{ofoZ= z$2`laYR)-TLzXVYBq_;b(RzRxp$DNWOUOgDQ1iL<@R-)3u#3XySP~5sYDVh|@--sQ zVon(OfBv4R?JOP>%xtqd57%d`P?zRlI9e}SMcd{qbDX_K3$N{*o_5i(P|we7W6rs> zyIjf}<1}fNn#)~Rab2s`yQ5!I@hbOLFvN{w9#ft7{5$c>O*_=NMtw@|s^XU)`zkuI zbDDBGYoGFIoyCGpTKFH~6#c`z3fa6K`lIHxV2IwFbu1E#XHf?*3$s_Jnwvb)B`20_ z2Og=b1{}yy6|k^tz1keXDW$nq+1W4pPnmt0xOmnq{-b8UuvIEQeH=BAG3)(3+hLyd|Rb_Pxl1R_z_p@4*_cc^fF9Ao%G-GSa- zWKH`I9BSL&A<6k1fC^uiU7%0x0MCBDW0G2QE7mHNZ=rVZ0rl*s}sMh?+*GE*?k(t5i5X8 zOIZsc9P+{oBu$`>)+OQs+l5nSX=UorOZQF)WrMAWmiNoo#681ZLtW>;l&pKlv3#_7 za1W$-{A|*~P9Q%pITB^pO^c_U#lz=@&LxY^KM@z+w>hVrMZ-@FKam`GaO&RJ~Y-d9(eQ7A3yrjzPFnu zc6Fp1_rLvky6fTeF>l)bc9J&`FC}nJDWrQP(L1h^Rb^imi9a} zQF1JOoZbT~dtphkCAIq9f_3qFn!~n?W$Bzr%o!!Ux$XVh`c%tn`@X$zyz-WwFTIbS zYHIxiy3k(!qwHYo-z}(C2HTDRo;m%;4vn9gDB1J2@UFT2rk=AdSBwx~iiEZw;W zmxWGqO=iEl=0(EXCKnAtKU;mbJWX41Ys}v@vaKHiz@9vR&eK=q4 z4&b{&zq~}n$?Ci6{wEFbHvr&~w!C=s{Pu)&ephk}zB@^j4;@KFhaP6GkA#$nB;0Y+ zhs9+tm%Os<(z2IUyjQ#~-ZBN}!;XtPMs~mnvh(uJvHEfTjjiuiG=R#V`0*3)dZZ5@ znK*d#-FuHtS| z58(4d@%)6S9L49SRj`&S*=wc1JwkSt;Hkuzl;FVx%#M&PBlSX~pd+}$7j*2b_W;#E z0TWon8KRVH;4J|q!ls8mK*MK`rZyt341$xzl?n6(9p1pq9yF!4T!ehV2mrj51AT~bB?8OQo6Mh(^$6Ub_qmcJQxE67iA|y~@;Cqj;>`A8%QKe5 z!J)k8ERdznKX-nA)9BgPcaE=`Sif_kZ0EO}@7oF!Qu6Q%Uzu9JactAr?khVph1=OVpQk^s3hq(3L{j{ve=34>5nE%r{*OH@gx zuZ+$Zdyz79OMt{%gsjwWaJO|FihTaXbw zQ??-JITWF=qB#iegnl0ibW`NDw-29l$Yt7G=GA7YR%mYkk}Plh(pNy zLHFq*FgT5UBBX3VXKPM06RS|jRmYZ+so6MGo-3i|ME)trwM2WLyCl#*(HoTkz_W(a z$1THKE^I;ksQUzDLG^Rcl=l89{O)`O2&?e7F5D`qb2c7i%W~dj`citK*k?N86m>aw5B`FJ ze!Fh|snGz&?x*e0q>uVkY?vkuFpj37ea<)=J^IteoO26NDI_IXg2Td)8-pB7xf3`z zjkJ%hT(g{GqUM}KsXEYtaWc8pSje-i>TUDio%BDPbUlTi4?hb#@iNan*rlZ}!tW08 z`PqIVbJgDG9w26!X zkaYdfRyb*^c*jqhSse`X$wT`T@&;M$S2x&6A( zL3QSa>%xv~yvF&ON$ZLU>xzukH7U3z1lKJm9(YT@c(<&a0Z|I|nDp0bckkqX!0+B7 z{ICMKpO^Bxw+TOA)U+5kzo_MRHweF|Gg5vVPq2X?1sBrmd%bRbM~BS0Yh;tx>koG$ zOpvxl+D_o>1RevxQ7L>VaFtx85@rHyqheDabS5a_6ViX6a$>qHZkqm*MA#Q)@iVdm zn|vx9iJmw|F^jS-drLQL@0496|1kSfl_@OHRicsvi-%*MoGPxHV|4onfhd8m5MT}) zW^oAqAdYeo59teq!W8@%4pGb_!;=&>eVxE40rJl zS;&Y256J}&1kRlZ`y_uGJ};7_f&S>-1gHdQ?vnD+75}-&j{tF;!SlBTJum#5nd7%k za$E2}!)^H|m?qM2Y$!JT#LyEs6j_sT)@5uP2hG2pSFtwZ+?cU#LIu98-FUHKq~YSO zkzJ_`V}?x8mW*TTpmoZeH)&q>j(J(~@Z|EEcb3;=me*y>8wU;l%(dS%n)oI@0feP% za?>a5db7~P*C~Z{N?F~_LV;gSMYvEAE;sFFzBH-7$>EYJVVBX(pU~aSQWIa1Xr;j_ zZWrqL{7%5@BtbpLPpR({1T4gf0Fl4XO2 z!6?Eqi{ouiel^1_h0Gi@(~@nuZO!AWZoBgNRktfm{PNp{YxoZSlLLGq-^$0&PCD04 zIM;r{(cLHO1im^~gqz!TD_@S`1$#V}(J#g+wZTCja@wA;jkLV5=UwwssFlK!-_kdm m%W=YNaQxbA$8H%MzmJf{^_G26>uT<;)yCFp{aaN$;Qt4R7:40001 \ + --number +15105550123 [--duration 15] +""" +import argparse, asyncio, base64, json, math, os, random, socket, struct, sys, threading, time +import urllib.error, urllib.parse, urllib.request + +def log(*a): print(f"[agent {time.strftime('%H:%M:%S')}]", *a, flush=True) + +# ---------- G.711 u-law (与 mock-provider 同源) ---------- +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 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(ts: float) -> bytes: + """20ms / 160 samples @8kHz, 440Hz 500ms-on 500ms-off (模拟 TTS)""" + 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)) + return bytes(out) + +# ---------- ARI REST ---------- +class Ari: + def __init__(self, base, user, password): + self.base = base.rstrip("/") + self.auth = "Basic " + base64.b64encode(f"{user}:{password}".encode()).decode() + + def req(self, method, path, params=None): + url = f"{self.base}/ari{path}" + if params: + url += "?" + urllib.parse.urlencode(params) + r = urllib.request.Request(url, method=method) + r.add_header("Authorization", self.auth) + try: + with urllib.request.urlopen(r, timeout=15) as resp: + raw = resp.read() + return resp.status, (json.loads(raw) if raw else None) + except urllib.error.HTTPError as e: + raw = e.read() + try: + return e.code, (json.loads(raw) if raw else None) + except Exception: + return e.code, {"error": raw.decode("utf-8", "replace")} + +# ---------- ARI Stasis app 保活 (极简事件 websocket, 纯 stdlib) ---------- +class AppKeeper: + """ARI 要求 Stasis app 有活跃 websocket 订阅, 否则进入 app 的通道会被立即挂断. + 本类只保持连接并丢弃事件, 控制面仍走 REST.""" + + def __init__(self, base, app, user, password): + u = urllib.parse.urlparse(base) + self.host, self.port = u.hostname, u.port or 80 + self.app, self.user, self.password = app, user, password + self.sock = None + + def start(self): + key = base64.b64encode(os.urandom(16)).decode() + req = ("GET /ari/events?app=" + urllib.parse.quote(self.app) + + "&api_key=" + urllib.parse.quote(f"{self.user}:{self.password}") + + f" HTTP/1.1\r\nHost: {self.host}:{self.port}\r\n" + + "Upgrade: websocket\r\nConnection: Upgrade\r\n" + + f"Sec-WebSocket-Key: {key}\r\nSec-WebSocket-Version: 13\r\n\r\n") + self.sock = socket.create_connection((self.host, self.port), timeout=10) + self.sock.sendall(req.encode()) + resp = b"" + while b"\r\n\r\n" not in resp: + chunk = self.sock.recv(1) # 逐字节读响应头, 不吞后续 ws 帧 + if not chunk: + raise ConnectionError("ws closed during handshake") + resp += chunk + if b"101" not in resp.split(b"\r\n", 1)[0]: + raise ConnectionError(f"ws handshake failed: {resp[:120]!r}") + self.sock.settimeout(None) + threading.Thread(target=self._drain, daemon=True).start() + + def _recv_exact(self, n): + buf = b"" + while len(buf) < n: + c = self.sock.recv(n - len(buf)) + if not c: + raise ConnectionError("ws closed") + buf += c + return buf + + def _drain(self): + try: + while True: + b0 = self._recv_exact(1)[0] + b1 = self._recv_exact(1)[0] + ln = b1 & 0x7F + if ln == 126: + ln = struct.unpack("!H", self._recv_exact(2))[0] + elif ln == 127: + ln = struct.unpack("!Q", self._recv_exact(8))[0] + payload = self._recv_exact(ln) if ln else b"" + op = b0 & 0x0F + if op == 0x8: # close + break + if op == 0x9: # ping → pong + mask = os.urandom(4) # client 帧必须带 mask + masked = bytes(b ^ mask[i % 4] for i, b in enumerate(payload)) + self.sock.sendall(bytes((0x8A, 0x80 | (len(payload) & 0x7F))) + mask + masked) + continue + if op == 0x1 and payload: + try: + ev = json.loads(payload) + except Exception: + continue + if ev.get("type") in ("StasisStart", "StasisEnd", "ChannelDestroyed"): + log(f"event {ev['type']} {(ev.get('channel') or {}).get('name', '')}") + except Exception: + pass + + def stop(self): + try: + if self.sock: + self.sock.close() + except Exception: + pass + +# ---------- RTP 媒体 ---------- +class Rx(asyncio.DatagramProtocol): + def __init__(self, stats): + self.stats, self.transport = stats, None + + def connection_made(self, transport): + self.transport = transport + + def datagram_received(self, data, addr): + if len(data) < 12: return + self.stats["rx"] += 1 + payload = data[12:] + dec = [ulaw_decode(b) for b in payload[:40]] + self.stats["amp_sum"] += sum(abs(x) for x in dec) / max(1, len(dec)) + if self.stats["tx_to"] is None: # 兜底: 用首包源地址作回程 + 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}") + +async def media_loop(listen_port, tx_to, duration): + stats = {"rx": 0, "amp_sum": 0.0, "tx": 0, "tx_to": tx_to} + loop = asyncio.get_running_loop() + rx_t, _ = await loop.create_datagram_endpoint( + lambda: Rx(stats), local_addr=("0.0.0.0", listen_port)) + + async def tx(): + seq = random.randint(0, 65535); ts = random.randint(0, 10**6) + ssrc = random.randint(1, 2**31) + t0 = time.time() + while time.time() - t0 < duration: + if stats["tx_to"]: + payload = tone_payload(time.time() - t0) + hdr = struct.pack("!BBHII", 0x80, 0, seq & 0xFFFF, ts & 0xFFFFFFFF, ssrc) + rx_t.sendto(hdr + payload, stats["tx_to"]) + stats["tx"] += 1 + seq += 1; ts += 160 + await asyncio.sleep(0.02) + + try: + await asyncio.wait_for(tx(), timeout=duration + 5) + finally: + rx_t.close() + return stats + +# ---------- 主流程 ---------- +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--ari-base", default="http://127.0.0.1:8088") + ap.add_argument("--ari-user", default="outbound") + ap.add_argument("--ari-pass", default="ari_mock_7c3f9e1b5a2d4806") + ap.add_argument("--app", default="outbound") + ap.add_argument("--external-host", required=True, help="本侧 RTP 地址 ip:port (供 Asterisk 回发混音)") + ap.add_argument("--number", default="+15105550123") + ap.add_argument("--caller-id", default="+15109990001") + ap.add_argument("--trunk", default="mock-trunk") + ap.add_argument("--duration", type=float, default=15) + ap.add_argument("--ring-timeout", type=float, default=30) + args = ap.parse_args() + + listen_port = int(args.external_host.rsplit(":", 1)[1]) + ari = Ari(args.ari_base, args.ari_user, args.ari_pass) + suffix = f"{int(time.time())}-{random.randint(100, 999)}" + bridge_id, em_id, callee_id = f"br-{suffix}", f"em-{suffix}", f"callee-{suffix}" + + # 0. 保持 Stasis app 运行 (否则通道进 app 即被挂断) + keeper = AppKeeper(args.ari_base, args.app, args.ari_user, args.ari_pass) + keeper.start() + log(f"stasis app '{args.app}' running (event ws connected)") + + def fail(msg, detail=None): + log(f"FAIL: {msg}", detail or "") + cleanup() + sys.exit(1) + + def cleanup(): + for m, p in (("DELETE", f"/channels/{callee_id}"), ("DELETE", f"/channels/{em_id}"), + ("DELETE", f"/bridges/{bridge_id}")): + try: ari.req(m, p) + except Exception: pass + keeper.stop() + + # 1. 混音桥 + st, r = ari.req("POST", "/bridges", params={"type": "mixing", "bridgeId": bridge_id}) + if st not in (200, 201): fail("create bridge", (st, r)) + + # 2. 外部媒体通道 (Asterisk → 本侧 RTP) + st, em = ari.req("POST", "/channels/externalMedia", params={ + "app": args.app, "external_host": args.external_host, + "format": "ulaw", "channelId": em_id}) + if st not in (200, 201): fail("create externalMedia", (st, r)) + log(f"externalMedia channel up: {em.get('name')}") + + # 3. Asterisk 侧 RTP 回送地址 + tx_to = None + st, v = ari.req("GET", f"/channels/{em_id}/variable", {"variable": "UNICASTRTP_LOCAL_ADDRESS"}) + ip = (v or {}).get("value") if st == 200 else None + st, v = ari.req("GET", f"/channels/{em_id}/variable", {"variable": "UNICASTRTP_LOCAL_PORT"}) + port = (v or {}).get("value") if st == 200 else None + if ip and port: + tx_to = (ip, int(port)) + log(f"asterisk RTP return addr: {ip}:{port}") + else: + log("UNICASTRTP vars not set; will reply-to-source on first RX") + + # 4. 发起外呼 + st, ch = ari.req("POST", "/channels", params={ + "endpoint": f"PJSIP/{args.number}@{args.trunk}", + "app": args.app, "appArgs": "outbound", + "callerId": args.caller_id, "channelId": callee_id, + "timeout": int(args.ring_timeout)}) + if st not in (200, 201): fail("originate", (st, ch)) + log(f"INVITE sent via {args.trunk} -> {args.number} (channel {callee_id})") + + # 5. 等待被叫应答 + deadline = time.time() + args.ring_timeout + state = None + while time.time() < deadline: + st, ch = ari.req("GET", f"/channels/{callee_id}") + if st == 200: + state = ch.get("state") + elif st in (404, 410): + fail("callee channel gone before answer") + if state == "Up": + break + time.sleep(0.25) + if state != "Up": + fail(f"no answer (state={state})") + log(f"callee answered: state={state}") + + # 6. 双通道入桥 → 双向 RTP + st, r = ari.req("POST", f"/bridges/{bridge_id}/addChannel", + params={"channel": f"{em_id},{callee_id}"}) + if st not in (200, 201, 204): fail("addChannel", (st, r)) + log(f"bridged [{em_id} + {callee_id}], streaming {args.duration}s " + f"(RX=ASR mock, TX=440Hz LLM/TTS mock)") + + # 7. 媒体收发 + stats = asyncio.run(media_loop(listen_port, tx_to, args.duration)) + + # 8. 挂断收尾 + cleanup() + 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) + +if __name__ == "__main__": + main() diff --git a/asterisk/deploy/mock/mock-provider.py b/asterisk/deploy/mock/mock-provider.py new file mode 100755 index 0000000..3afc5c7 --- /dev/null +++ b/asterisk/deploy/mock/mock-provider.py @@ -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: ", + "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 diff --git a/asterisk/deploy/scripts/call.sh b/asterisk/deploy/scripts/call.sh new file mode 100755 index 0000000..72a51dd --- /dev/null +++ b/asterisk/deploy/scripts/call.sh @@ -0,0 +1,18 @@ +#!/bin/bash +# 端到端外呼验证: mock-agent(ARI+RTP) → Asterisk → PJSIP trunk → mock-provider(应答) → 双向 RTP 统计 +# 用法: ./call.sh [被叫号码] [时长秒] [asterisk节点 1|2] +set -euo pipefail +cd "$(dirname "$0")/.." + +NUMBER="${1:-+15105550123}" +DURATION="${2:-15}" +NODE="${3:-1}" +NET="asterisk-ari_astari" +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 \ + --ari-base "http://asterisk${NODE}:8088" \ + --external-host "${NAME}:40001" \ + --number "$NUMBER" --duration "$DURATION" diff --git a/asterisk/deploy/scripts/health.sh b/asterisk/deploy/scripts/health.sh new file mode 100755 index 0000000..9109dc1 --- /dev/null +++ b/asterisk/deploy/scripts/health.sh @@ -0,0 +1,44 @@ +#!/bin/bash +# 运行状态检测: 容器 + ARI REST + PJSIP trunk 探活(qualify) + 活动呼叫数 +set -uo pipefail +USER="outbound"; SECRET="ari_mock_7c3f9e1b5a2d4806" +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)" ]; } + +echo "== 容器状态" +docker ps --filter "name=ast" --format "{{.Names}}\t{{.Status}}" + +echo "== ARI REST" +for n in "ast1 8088" "ast2 8089"; do + set -- $n + if running "$1"; then + c=$(curl -s -o /dev/null -w "%{http_code}" -u "$USER:$SECRET" --max-time 3 "http://127.0.0.1:$2/ari/asterisk/info") + [ "$c" = 200 ] && pass "$1 ARI /asterisk/info ($c)" || faill "$1 ARI /asterisk/info ($c)" + else + echo "SKIP $1 (未运行)" + fi +done + +echo "== PJSIP trunk 探活 (OPTIONS qualify)" +for c in ast1 ast2; do + if running "$c"; then + out=$(docker exec "$c" asterisk -rx "pjsip show contacts" 2>/dev/null | grep mock-provider) || true + if echo "$out" | grep -q "Avail"; then + pass "$c trunk mock-provider Avail" + else + faill "$c trunk mock-provider (contact: ${out:-none})" + fi + fi +done + +echo "== 活动呼叫 (ARI channels)" +if running ast1; then + n=$(curl -s -u "$USER:$SECRET" --max-time 3 http://127.0.0.1:8088/ari/channels | python3 -c 'import json,sys; print(len(json.load(sys.stdin)))' 2>/dev/null) + [ -n "$n" ] && pass "ast1 active channels=$n" || faill "ast1 active channels" +fi + +echo +echo "result: ok=$ok fail=$fail" +[ $fail -eq 0 ] diff --git a/asterisk/deployment.md b/asterisk/deployment.md new file mode 100644 index 0000000..71c8f69 --- /dev/null +++ b/asterisk/deployment.md @@ -0,0 +1,23 @@ +目标机器:39.106.106.246 +登录用户: root + +任务目标:在 livekit 运行时下线后,部署 Asterisk 并以 ARI + Mock SIP 实现 SIP 对接与运行状态检测。 + +项目地址: https://github.com/asterisk/asterisk (容器镜像 andrius/asterisk = Asterisk 22 LTS) +项目背景:使用 Asterisk 对接 MOCK SIP 服务商,作为 AI 外呼机器人基础能力底座,使用 MOCK 对接 LLM ASR 实现双向语音通信, 本项目仅实现外呼功能,不实现多人房间与呼入功能调研。给出可实现最快水平扩展的结论方案并赋实践证据。给出详细的部署过程文档与系统调优方案。 + +--- + +## 交付状态 (2026-08-31 全部完成) + +| 目标 | 状态 | 证据 | +|---|---|---| +| livekit 运行时下线 | ✅ | lk-* 容器清零, 7880/7883/8081/8082/9091/9092 端口释放 | +| Asterisk 部署 (双节点) | ✅ | ast1/ast2 healthy, ARI /asterisk/info 200 ×2 (Asterisk 22.10.1) | +| Mock SIP 对接 (外呼) | ✅ | INVITE→180→200 OK(PCMU)→ACK, trunk qualify Avail ×2 | +| 双向语音 (Mock LLM/ASR) | ✅ | `RESULT ... bidirectional=YES` (agent RX 512帧 RMS 8318; provider recv_rtp=512, 双向计数完全对称) | +| 运行状态检测 | ✅ | `scripts/health.sh` 容器+ARI REST+trunk qualify+通道数, ok=5 fail=0 | +| 水平扩展结论+实证 | ✅ | 3 路并发分摊双节点 (INVITE 源 172.18.0.2×2 / 172.18.0.4×1), 零配置变更 | +| 部署文档+调优方案 | ✅ | 本仓库 `README.md` (§4 部署含踩坑, §8 sysctl 沿用已应用) | + +**交付物**: 服务器 `/opt/asterisk-ari/` (compose+conf+mock+scripts+README), 本地 `deploy/` 同源; 原始任务与完成状态见本文件, 详细文档见 `README.md`。 diff --git a/livekit/README.md b/livekit/README.md new file mode 100644 index 0000000..d516fc5 --- /dev/null +++ b/livekit/README.md @@ -0,0 +1,198 @@ +# LiveKit + livekit-sip AI 外呼底座 — 部署与运维文档 + +> 目标机器: `39.106.106.246` (阿里云 2vCPU/3.5G, Debian 13) · 交付物: `/opt/livekit-sip/` +> 项目: https://github.com/livekit/livekit · 插件: https://github.com/livekit/sip +> 范围: **仅外呼** (CreateSIPParticipant 工作流), Mock SIP 运营商 + Mock LLM/ASR, 已端到端验证双向语音。 + +## 1. 水平扩展结论 (先行) + +**结论: 共享 Redis 的无状态水平复制 —— livekit-server 与 livekit-sip 均不在本机存呼叫状态, 扩容 = 加副本, 无需协调器/迁移。** + +| 层 | 状态存放 | 扩容方式 | 调度机制 | +|---|---|---|---| +| livekit-sip 网关 | Redis (trunk/呼叫状态) + 各自持有 RTP 流 | 复制 service 块 (独立端口段) | CreateSIPParticipant 经 Redis **随机派发**到任一网关节点 | +| livekit-server | Redis (rooms/room_node_map/nodes) | 加节点, 客户端连任一节点 | 房间归属节点经 Redis 自动路由, 跨节点媒体互通 | +| AI Agent (生产) | 无状态 worker | `--scale`/K8s HPA | AgentDispatch 按 agent_name 派发 | + +**最快扩容路径**: `docker compose` 里复制一个 `sipN` service 块 (改宿主端口映射即可, 容器网内 5060/10000-10200 天然隔离); K8s 场景用 Deployment + NodePort/hostPort 固定 SIP+RTP 端口段。实测从 1 网关扩到 2 网关**零配置变更** (共用同一 `sip.yaml`), 呼叫自动分摊 (见 §7)。 + +瓶颈与边界: 单网关受 CPU (RTP 打包/转发) 与 RTP 端口段限制; 本文配置 200 端口/网关, 实际每呼叫占 1 个 UDP mux 端口, 单节点并发上限≈端口数。LiveKit 官方口径单节点(4C8G)数百路通话; 水平扩展后容量线性叠加, 上限在 Redis 与出口带宽。 + +## 2. 架构 + +``` + ┌────────────────────────── docker net: lksip ──────────────────────────┐ + lk CLI / 业务后端 │ ┌──────────┐ ws ┌─────────────┐ ┌─────────────┐ │ + (host, :7880 API) ───────┼─►│ livekit │◄───────►│ livekit2 │ │ redis │ │ + CreateSIPParticipant │ │ :7880 │ Redis │ :7883(multinode)│ │ :6379 │ │ + │ trunk/呼叫状态 │ └────┬─────┘ 集群 └─────────────┘ └──────┬──────┘ │ + ▼ 全部经 Redis │ │ RTC(ICE 50100-50200/udp) │ 状态存储 │ + ┌─────────────┐ │ ┌────▼──────────┐ ┌──────────────┐ │ │ + │ redis │◄─────────┼──┤ sip1 (网关) │ │ sip2 (网关) │──► 随机被派发 CreateSIPParticipant│ + └─────────────┘ │ │ 5060/udp+tcp │ │ 5060/udp+tcp │ │ │ + │ │ RTP 10000- │ │ RTP 10000- │ │ │ + mock-agent (host) │ │ 10200/udp │ │ 10200/udp │ │ │ + 订阅=Mock ASR │ └──────┬───────┘ └──────┬───────┘ │ │ + 发布=Mock LLM/TTS ───────┼─────────┘ SIP INVITE / RTP(PCMU) │ │ + │ └──────────►┌────────────────┐ │ │ + │ │ mock-provider │──────┘ (SIP UAS, 440Hz 音) │ + │ │ 5060/udp 40000 │ 纯 stdlib Python │ + │ └────────────────┘ │ + └───────────────────────────────────────────────────────────────────────┘ +``` + +外呼信令流 (仅出站): +``` +业务后端 → lk sip participant create --trunk ST_xxx --call +1510... + → LiveKit API → Redis 派发到某 sip 网关 + → INVITE sip:+1510...@mock-provider → 100/180 → 200 OK(SDP, PCMU) → ACK + → 被叫以 SIP participant 身份进房, 与房间内 agent(=AI) 双向 RTP +``` + +## 3. 端口规划 + +| 端口 | 组件 | 说明 | +|---|---|---| +| 7880/tcp | livekit (node1) | HTTP/WS API, 宿主映射 7880 | +| 7883/tcp→7880 | livekit2 (node2) | 同一 Redis, 双节点集群 | +| 50100-50200/udp | livekit ×2 | RTC ICE (各容器独立 netns 不冲突) | +| 5060/udp+tcp | sip1/sip2 | SIP 信令 (容器网内, 未映射宿主) | +| 10000-10200/udp | sip1/sip2 | RTP 媒体 (同上) | +| 8081/8082 | sip1/sip2 health | `/healthz` + `/` (HTTP 200) | +| 9091/9092 | sip1/sip2 prometheus | `/metrics` | +| 5060/udp, 40000/udp | mock-provider | Mock 运营商 SIP+RTP | + +生产外网仅需放行: SIP 5060/udp+tcp (+5061/tls) 与 RTP 段; LiveKit 客户端走 7880 + RTC 段。 + +## 4. 从零部署步骤 (全部实测) + +```bash +# ── 4.1 Docker (Debian 13, 官方源; 国内需配 mirror) ───────────────────── +apt-get update && apt-get install -y ca-certificates curl gnupg +install -m 0755 -d /etc/apt/keyrings +curl -fsSL https://download.docker.com/linux/debian/gpg -o /etc/apt/keyrings/docker.asc +echo "deb [arch=amd64 signed-by=/etc/apt/keyrings/docker.asc] https://download.docker.com/linux/debian trixie stable" \ + > /etc/apt/sources.list.d/docker.list +apt-get update && apt-get install -y docker-ce docker-ce-cli containerd.io docker-compose-plugin +# ⚠ 踩坑1(Debian 13): 安装后 docker group 缺失 → docker.socket 216/GROUP; 需 groupadd docker +# ⚠ 踩坑2(最小系统): 无 iptables → bridge 驱动失败 "iptables not found"; 需 apt install iptables nftables +cat > /etc/docker/daemon.json <<'EOF' +{ "registry-mirrors": ["https://docker.m.daocloud.io", "https://docker.1ms.run"] } +EOF +systemctl restart docker + +# ── 4.2 内核调优 (见 §8, 已应用 /etc/sysctl.d/99-livekit-sip.conf) ───── + +# ── 4.3 部署栈 ───────────────────────────────────────────────────────── +mkdir -p /opt/livekit-sip && cd /opt/livekit-sip # = 仓库 deploy/ 目录 +docker compose up -d # redis + livekit + sip1 + sip2 + mock-provider +docker compose --profile multinode up -d livekit2 # (可选) LiveKit 第2节点 + +# ── 4.4 Mock agent 依赖 (host python) ────────────────────────────────── +apt-get install -y python3-venv +python3 -m venv /opt/lk-venv +/opt/lk-venv/bin/pip install -i https://mirrors.aliyun.com/pypi/simple/ livekit livekit-api +# ⚠ 踩坑3: SDK 1.x 拆包 — rtc 在 `livekit`, AccessToken 在 `livekit-api`; +# grants 必须传 VideoGrants() 对象; frame.data 是 buffer 非 bytes。 + +# ── 4.5 建 trunk + 发起外呼 ──────────────────────────────────────────── +./scripts/call.sh <房间> <被叫E.164> [时长秒] # 首次自动创建 trunk 并缓存于 Redis +``` + +trunk 定义 (`mock/trunk.json`, 换真实运营商时改 address/auth): +```json +{ "trunk": { "name": "mock-provider", "address": "mock-provider", "numbers": ["+15105550123"] } } +``` + +## 5. 外呼链路验证 (实测输出) + +`./scripts/call.sh test-room-5 +15105550123 12`: +``` +== 2. Mock agent 加入房间 test-room-5 (发布TTS音, 订阅被叫音频) +[agent] joined room=test-room-5 peers=[] +[agent] publishing 440Hz tone (500ms on / 500ms off) +== 3. 发起 SIP 外呼 -> +15105550123 +[agent] subscribed: callee-30395 track=TR_AMu2wATzwBLwqP +SIPCallID: SCL_LEGvhUjpz3N9 ← CreateSIPParticipant 成功, 被叫应答进房 +== 4. 等待 agent 结束并输出双向统计 +[agent] RX frames=1000 avg_rms=4667 +RESULT room=test-room-5 rx_frames=1015 rx_avg_rms=4640 bidirectional=YES +``` +mock-provider (被叫侧) 日志: +``` +INCOMING_INVITE from=172.18.0.6:5060 uri="sip:+15105550123@mock-provider" +ANSWERED media=172.18.0.3:40000 codec=PCMU peer=172.18.0.6:10034 ← SDP 协商 PCMU +ACK confirmed, RTP bridging +CALL_DONE sent_rtp=… recv_rtp=… recv_avg_amplitude=… ← 双向 RTP 计数 +``` +**双向语音证明**: agent RX=被叫→ASR 方向 (1015 帧, RMS 4640 = 440Hz 检测到); provider recv = LLM/TTS→被叫方向 (RTP 计数>0)。换真实 LLM/ASR 只需替换 mock-agent 内的发布/消费循环 (livekit-agents 框架同接口)。 + +## 6. 运行状态检测 + +`./scripts/health.sh` (实测输出): +``` +== 容器状态 lk-livekit2/lk-sip1/lk-sip2/lk-livekit/lk-redis/lk-mock-provider Up +== HTTP 健康 PASS livekit(200) PASS sip1(200) PASS sip2(200) # sip /healthz +== API 功能 PASS lk room list PASS lk sip outbound list +result: ok=5 fail=0 +``` +持续监控: `watch -n5 ./scripts/health.sh`; 指标: `curl :9091/metrics`, `curl :9092/metrics` (prometheus, 含活动呼叫/错误); 呼叫级状态: provider 容器日志 `CALL_DONE` 行; 房间级: `lk room list`。 + +## 7. 水平扩展实证 (同机双网关 + 双 LiveKit 节点) + +**① 双 LiveKit 节点集群 (共 Redis, 自动发现)** +``` +$ docker exec lk-redis redis-cli hlen nodes +2 ← node1 ND_BwShgqdiB8AQ + node2 ND_moj9RpSqq5Ap +``` +并发验证中 agent 全部连 **node2(:7883)**, CreateSIPParticipant 经 **node1(:7880)**, 同房媒体互通 → 跨节点路由成立 (RESULT 见 ③)。 + +**② 双 SIP 网关随机分摊 (无状态验证)** +``` +$ docker logs lk-sip1 | grep -oE '"room": "burst-[0-9]"' | sort | uniq -c + 15 "room": "burst-1" ← sip1 处理 burst-1 +$ docker logs lk-sip2 | ... + 15 "room": "burst-2" ← sip2 处理 burst-2 + 15 "room": "burst-3" ← 及 burst-3 +$ docker logs lk-mock-provider | grep INCOMING | tail -3 +from=172.18.0.5:5060 call_id=lIIip… ← 两个来源 IP = sip1/sip2 +from=172.18.0.6:5060 call_id=mw4Fq… +from=172.18.0.5:5060 call_id=BCVgl… +``` +**③ 3 路并发外呼全部双向 (agent 连 node2)** +``` +RESULT room=burst-1 rx_frames=810 rx_avg_rms=3162 bidirectional=YES +RESULT room=burst-2 rx_frames=866 rx_avg_rms=3066 bidirectional=YES +RESULT room=burst-3 rx_frames=778 rx_avg_rms=2228 bidirectional=YES +``` +→ 加副本即扩容、Redis 自动调度、跨节点房间互通, 三项均实证。资源基线 (3 并发): livekit 2×6%CPU/36M, sip 2×20%CPU/30M, redis 1%/8M。 + +## 8. 系统调优 + +**已应用** (`/etc/sysctl.d/99-livekit-sip.conf`, `sysctl --system` 生效): +```conf +net.core.rmem_max = 16777216 # UDP 收缓冲上限 (pion/WebRTC 推荐 16M; 默认 208K 高并发必丢包) +net.core.wmem_max = 16777216 +net.core.rmem_default = 1048576 +net.core.wmem_default = 1048576 +net.core.netdev_max_backlog = 4096 # 网卡积压队列 (默认 1000) +net.ipv4.ip_local_port_range = 10000 65000 # 出向 SIP/RTP 源端口 (含 RTP 段) +net.ipv4.udp_mem = 8388608 12582912 16777216 # 全局 UDP 页缓存 (3.5G 内存安全值) +net.ipv4.udp_rmem_min/wmem_min = 16384 +``` +调优前基线: `UdpRcvbufErrors=0` (3 并发, 默认缓冲已够); 上述参数为高并发 (>100 路) 预置。验证命令: `nstat -az | grep -i udp`。 + +**组件级**: +- sip 网关: `rtp_port: 10000-10200` (每节点 200 并发余量); `max_cpu_utilization: 0.9` 默认即可, 超 90% 自动拒绝新呼叫 (背压); `log_level: info` 生产勿用 debug。 +- livekit: RTC 段 50100-50200/节点; 房间自动关闭 (empty timeout) 回收端口; `logging.json: true` 便于采集。 +- redis: 纯缓存用法 (`--save ""`), trunk/房间状态可丢 (重建即恢复); 若把 Redis 当持久层需开 AOF。 +- 容器: `restart: unless-stopped` 已配; 单机 3.5G 下建议给 livekit/sip 设 mem limit (如 512M) 防互挤。 +- 监控: 9091/9092 /metrics 接 Prometheus; 告警项: sip 活动呼叫数骤降、livekit 5xx、UdpRcvbufErrors>0、容器重启。 + +## 9. 生产化清单 (Mock → 真实运营商) + +1. trunk: `address` 换运营商 SIP 域名/IP, 加 `auth_username/auth_password`, 号码用 E.164; 我方 5060/udp+tcp、RTP 段需公网放行 (阿里云安全组)。 +2. NAT: 阿里云 ECS 为 1:1 NAT, sip.yaml 改 `use_external_ip: true` (STUN 自动探测公网 IP 写入 SDP), livekit.yaml 同理; 或显式 `nat_1_to_1_ip`。 +3. 安全: SIP over TLS (`tls.port: 5061`) + SRTP (`media_enc`), 或 IP 白名单 ACL; LiveKit API key 轮换。 +4. AI 侧: mock-agent 替换为 livekit-agents 框架 (同房间/同 track 接口), STT/LLM/TTS 三段插拔。 +5. 高可用: Redis 单点 → 云 Redis 主从; 网关 ≥2 跨宿主机; LiveKit 节点 ≥2 (本文已验证集群模式)。 diff --git a/livekit/deploy/conf/livekit.yaml b/livekit/deploy/conf/livekit.yaml new file mode 100644 index 0000000..7ddb107 --- /dev/null +++ b/livekit/deploy/conf/livekit.yaml @@ -0,0 +1,15 @@ +# LiveKit Server 配置 — 单机双节点共享 Redis 组成集群 +port: 7880 +bind_addresses: [""] +redis: + address: redis:6379 +keys: + devkey: "f4f6c1a9e2b74d5893a0c8e17d2b6a45b7c93f1d6e284a05c9b3d7f2e6a81c04" +rtc: + tcp_port: 7881 + port_range_start: 50100 + port_range_end: 50200 + use_external_ip: false +logging: + level: info + json: true diff --git a/livekit/deploy/conf/sip.yaml b/livekit/deploy/conf/sip.yaml new file mode 100644 index 0000000..1d4cebd --- /dev/null +++ b/livekit/deploy/conf/sip.yaml @@ -0,0 +1,12 @@ +# livekit-sip 网关配置 — sip1/sip2 共用(各自独立 netns, 端口不冲突) +api_key: devkey +api_secret: "f4f6c1a9e2b74d5893a0c8e17d2b6a45b7c93f1d6e284a05c9b3d7f2e6a81c04" +ws_url: ws://livekit:7880 +redis: + address: redis:6379 +health_port: 8081 +prometheus_port: 9091 +log_level: info +sip_port: 5060 +rtp_port: 10000-10200 +use_external_ip: false diff --git a/livekit/deploy/docker-compose.yml b/livekit/deploy/docker-compose.yml new file mode 100644 index 0000000..e3ebcc7 --- /dev/null +++ b/livekit/deploy/docker-compose.yml @@ -0,0 +1,75 @@ +name: livekit-sip + +services: + redis: + image: redis:7-alpine + container_name: lk-redis + command: redis-server --maxmemory 256mb --maxmemory-policy allkeys-lru --save "" --appendonly no + restart: unless-stopped + networks: [lksip] + + livekit: + image: livekit/livekit-server:latest + container_name: lk-livekit + command: --config /etc/livekit.yaml + volumes: + - ./conf/livekit.yaml:/etc/livekit.yaml:ro + ports: ["7880:7880"] + restart: unless-stopped + networks: [lksip] + depends_on: [redis] + + # 第二个 LiveKit 节点: 同一 Redis 即自动组成多节点集群(水平扩展验证) + livekit2: + image: livekit/livekit-server:latest + container_name: lk-livekit2 + command: --config /etc/livekit.yaml + volumes: + - ./conf/livekit.yaml:/etc/livekit.yaml:ro + ports: ["7883:7880"] + restart: unless-stopped + networks: [lksip] + depends_on: [redis] + profiles: [multinode] + + # SIP 网关 x2: 无状态, 经 Redis 被随机调度, 天然水平扩展 + sip1: + image: livekit/sip:latest + container_name: lk-sip1 + command: --config /etc/sip.yaml + volumes: + - ./conf/sip.yaml:/etc/sip.yaml:ro + ports: ["8081:8081", "9091:9091"] + restart: unless-stopped + networks: [lksip] + depends_on: [redis, livekit] + + sip2: + image: livekit/sip:latest + container_name: lk-sip2 + command: --config /etc/sip.yaml + volumes: + - ./conf/sip.yaml:/etc/sip.yaml:ro + ports: ["8082:8081", "9092:9091"] + restart: unless-stopped + networks: [lksip] + depends_on: [redis, livekit] + + # Mock SIP 运营商 (UAS): 应答 livekit-sip 的出站 INVITE + mock-provider: + image: python:3.12-alpine + container_name: lk-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: [lksip] + +networks: + lksip: + driver: bridge diff --git a/livekit/deploy/mock/mock-agent.py b/livekit/deploy/mock/mock-agent.py new file mode 100644 index 0000000..e4d8574 --- /dev/null +++ b/livekit/deploy/mock/mock-agent.py @@ -0,0 +1,78 @@ +#!/usr/bin/env python3 +"""Mock AI Agent: 加入房间, 发布 440Hz 间歇音(模拟 TTS 输出), +订阅 SIP 参与者音频并统计 RMS(模拟 ASR 输入), 验证双向语音。 +依赖: pip install livekit +""" +import argparse, array, asyncio, math, struct, time +from livekit import rtc +from livekit.api import AccessToken, VideoGrants + +def rms(data) -> float: + buf = bytes(data)[:4096] + a = array.array("h"); a.frombytes(buf) + if not a: return 0.0 + return (sum(x * x for x in a) / len(a)) ** 0.5 + +async def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--url", default="ws://127.0.0.1:7880") + ap.add_argument("--api-key", default="devkey") + ap.add_argument("--api-secret", default="f4f6c1a9e2b74d5893a0c8e17d2b6a45b7c93f1d6e284a05c9b3d7f2e6a81c04") + ap.add_argument("--room", required=True) + ap.add_argument("--identity", default="mock-agent") + ap.add_argument("--duration", type=float, default=20) + args = ap.parse_args() + + token = (AccessToken(args.api_key, args.api_secret) + .with_identity(args.identity).with_name("MockAgent") + .with_grants(VideoGrants(room_join=True, room=args.room, + can_publish=True, can_subscribe=True))) + + room = rtc.Room() + stats = {"frames": 0, "rms_sum": 0.0, "src": None} + done = asyncio.Event() + + def on_sub(track, pub, participant): + if track.kind != rtc.TrackKind.KIND_AUDIO: + return + print(f"[agent] subscribed: {participant.identity} track={track.sid}", flush=True) + async def consume(): + stream = rtc.AudioStream(track) + async for ev in stream: + f = ev.frame + stats["frames"] += 1 + stats["rms_sum"] += rms(f.data) + if stats["frames"] % 250 == 0: + print(f"[agent] RX frames={stats['frames']} avg_rms={stats['rms_sum']/stats['frames']:.0f}", flush=True) + stats["src"] = asyncio.create_task(consume()) + + room.on("track_subscribed", on_sub) + await room.connect(args.url, token.to_jwt()) + peers = [p.identity for p in room.remote_participants.values()] + print(f"[agent] joined room={args.room} peers={peers}", flush=True) + + src = rtc.AudioSource(48000, 1) + track = rtc.LocalAudioTrack.create_audio_track("agent-tts", src) + await room.local_participant.publish_track(track) + print("[agent] publishing 440Hz tone (500ms on / 500ms off)", flush=True) + + async def tone(): + t, end = 0.0, time.time() + args.duration + while time.time() < end: + samples = [] + on = (t % 1.0) < 0.5 + for _ in range(960): # 20ms @ 48kHz + v = int(6000 * math.sin(2 * math.pi * 440 * t)) if on else 0 + samples.append(v); t += 1 / 48000 + buf = struct.pack("<960h", *samples) + await src.capture_frame(rtc.AudioFrame(buf, 48000, 1, 960)) + await asyncio.sleep(0.02) + + await tone() + n, r = stats["frames"], (stats["rms_sum"] / stats["frames"]) if stats["frames"] else 0.0 + print(f"RESULT room={args.room} rx_frames={n} rx_avg_rms={r:.0f} " + f"bidirectional={'YES' if n > 50 and r > 100 else 'NO'}", flush=True) + await room.disconnect() + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/livekit/deploy/mock/mock-provider.py b/livekit/deploy/mock/mock-provider.py new file mode 100644 index 0000000..bca7018 --- /dev/null +++ b/livekit/deploy/mock/mock-provider.py @@ -0,0 +1,237 @@ +#!/usr/bin/env python3 +"""Mock SIP 运营商 (UAS): 接收 livekit-sip 出站 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: ", + "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:] + for call in list(calls.values()): + if not call.closed and call.peer_ip == addr[0]: + call.recv += 1; call.recv_bytes += len(data) + 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) + return + +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 livekit-sip ...") + await asyncio.Event().wait() + +if __name__ == "__main__": + try: + asyncio.run(main()) + except KeyboardInterrupt: + pass diff --git a/livekit/deploy/mock/trunk.json b/livekit/deploy/mock/trunk.json new file mode 100644 index 0000000..917eafa --- /dev/null +++ b/livekit/deploy/mock/trunk.json @@ -0,0 +1,7 @@ +{ + "trunk": { + "name": "mock-provider", + "address": "mock-provider", + "numbers": ["+15105550123"] + } +} diff --git a/livekit/deploy/scripts/call.sh b/livekit/deploy/scripts/call.sh new file mode 100644 index 0000000..5233c37 --- /dev/null +++ b/livekit/deploy/scripts/call.sh @@ -0,0 +1,38 @@ +#!/bin/bash +# 端到端外呼验证: agent 进房 -> CreateSIPParticipant 呼出 -> mock 运营商应答 -> 双向 RTP 统计 +# 用法: ./call.sh [房间名] [被叫号码] [时长秒] +set -euo pipefail +cd "$(dirname "$0")/.." + +ROOM="${1:-call-$(date +%s)}" +NUMBER="${2:-+15105550123}" +DURATION="${3:-15}" +SECRET="f4f6c1a9e2b74d5893a0c8e17d2b6a45b7c93f1d6e284a05c9b3d7f2e6a81c04" +LK="docker run --rm --network host -e LIVEKIT_URL=http://127.0.0.1:7880 \ + -e LIVEKIT_API_KEY=devkey -e LIVEKIT_API_SECRET=$SECRET livekit/livekit-cli:latest" + +echo "== 1. 确保 outbound trunk 存在 (mock-provider)" +TRUNK_ID=$($LK sip outbound list --json 2>/dev/null | python3 -c ' +import json,sys +for t in json.load(sys.stdin).get("items", []): + if t.get("name") == "mock-provider": + print(t.get("sipTrunkId","")); break') +if [ -z "$TRUNK_ID" ]; then + TRUNK_ID=$($LK sip outbound create mock/trunk.json 2>/dev/null | awk '/^SIPTrunkID:/ {print $2}') + echo " created trunk: $TRUNK_ID" +else + echo " reuse trunk: $TRUNK_ID" +fi + +echo "== 2. Mock agent 加入房间 $ROOM (发布TTS音, 订阅被叫音频)" +/opt/lk-venv/bin/python mock/mock-agent.py --room "$ROOM" --duration "$DURATION" & +AGENT_PID=$! +sleep 2 + +echo "== 3. 发起 SIP 外呼 -> $NUMBER" +$LK sip participant create --room "$ROOM" --identity "callee-$RANDOM" \ + --trunk "$TRUNK_ID" --number "$NUMBER" --call "$NUMBER" \ + --name "MockCallee" --wait --timeout 60s || true + +echo "== 4. 等待 agent 结束并输出双向统计" +wait $AGENT_PID diff --git a/livekit/deploy/scripts/health.sh b/livekit/deploy/scripts/health.sh new file mode 100644 index 0000000..8015d73 --- /dev/null +++ b/livekit/deploy/scripts/health.sh @@ -0,0 +1,32 @@ +#!/bin/bash +# 运行状态检测: 各组件健康 + 集群节点 + trunk 配置 +set -uo pipefail +SECRET="f4f6c1a9e2b74d5893a0c8e17d2b6a45b7c93f1d6e284a05c9b3d7f2e6a81c04" +ok=0; fail=0 + +chk() { # name, command... + local name="$1"; shift + if out=$("$@" 2>&1); then echo "PASS $name"; ok=$((ok+1)) + else echo "FAIL $name: $(echo "$out" | head -1)"; fail=$((fail+1)); fi +} + +echo "== 容器状态" +docker ps --filter "name=lk-" --format "{{.Names}}\t{{.Status}}" + +echo "== HTTP 健康" +code() { curl -s -o /dev/null -w "%{http_code}" --max-time 3 "$1"; } +for p in "livekit http://127.0.0.1:7880/" "sip1 http://127.0.0.1:8081/healthz" "sip2 http://127.0.0.1:8082/healthz"; do + set -- $p + c=$(code "$2") + [ "$c" = "200" ] || [ "$c" = "204" ] && { echo "PASS $1 ($c)"; ok=$((ok+1)); } || { echo "FAIL $1 ($c)"; fail=$((fail+1)); } +done + +echo "== LiveKit API 功能检查" +LK=(docker run --rm --network host -e LIVEKIT_URL=http://127.0.0.1:7880 + -e LIVEKIT_API_KEY=devkey -e LIVEKIT_API_SECRET=$SECRET livekit/livekit-cli:latest) +chk "lk room list" "${LK[@]}" room list +chk "lk sip outbound list" "${LK[@]}" sip outbound list + +echo +echo "result: ok=$ok fail=$fail" +[ $fail -eq 0 ] diff --git a/livekit/deployment.md b/livekit/deployment.md new file mode 100644 index 0000000..cfdd335 --- /dev/null +++ b/livekit/deployment.md @@ -0,0 +1,24 @@ +目标机器:39.106.106.246 +登录用户: root + +任务目标:部署 Livekit 并融合 livekit-sip 插件使用Mock实现SIP对接与运行状态检测。 + +项目地址: https://github.com/livekit/livekit +插件地址: +项目背景:使用livekit对接MOCK SIP服务商,作为AI外呼机器人基础能力底座,使用MOCK对接LLM ASR实现双向语音通信, 本项目仅实现外呼功能,不实现多人房间与呼入功能调研。给出可实现最快水平扩展的结论方案并赋实践证据。给出详细的部署过程文档与系统调优方案。 + +--- + +## 交付状态 (2026-08-31 全部完成) + +| 目标 | 状态 | 证据 | +|---|---|---| +| LiveKit 部署 (双节点集群) | ✅ | Redis `hlen nodes` = 2; health.sh ok=5 fail=0 | +| livekit-sip 网关 ×2 融合 | ✅ | sip1/sip2 `/healthz` 200, 共 Redis 随机派发 | +| Mock SIP 对接 (外呼) | ✅ | INVITE→180→200 OK(PCMU)→ACK, mock-provider 全 stdlib | +| 双向语音 (Mock LLM/ASR) | ✅ | `RESULT ... bidirectional=YES` (agent RX 1015帧 RMS 4640; provider 双向 RTP 计数) | +| 运行状态检测 | ✅ | `scripts/health.sh` 容器+HTTP+API 三层检查 | +| 水平扩展结论+实证 | ✅ | 3 路并发分摊到双网关 (172.18.0.5/0.6 双源 INVITE), agent 连 node2 / 呼叫经 node1 跨节点互通 | +| 部署文档+调优方案 | ✅ | 本仓库 `README.md` (§4 部署含踩坑, §8 sysctl 已应用) | + +**交付物**: 服务器 `/opt/livekit-sip/` (compose+conf+mock+scripts+README), 本地 `deploy/` 同源; 原始任务与完成状态见本文件, 详细文档见 `README.md`。