chore: init SIP research repo (Asterisk & LiveKit)

This commit is contained in:
杨豪
2026-08-31 13:46:57 +08:00
commit 5a54d9fa0c
23 changed files with 1594 additions and 0 deletions
+8
View File
@@ -0,0 +1,8 @@
[general]
enabled = yes
pretty = no
[outbound]
type = user
read_only = no
password = ari_mock_7c3f9e1b5a2d4806
+4
View File
@@ -0,0 +1,4 @@
[general]
enabled = yes
bindaddr = 0.0.0.0
bindport = 8088
+21
View File
@@ -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 内完成, 重启后健康检查快速变绿
+6
View File
@@ -0,0 +1,6 @@
; RTP 端口段: 每呼叫占 2 个端口对 (被叫腿 + externalMedia 腿, 各含 RTCP)
; 10000-10800 = 800 端口 ≈ 200 并发呼叫/节点
[general]
rtpstart = 10000
rtpend = 10800
strictrtp = no
+47
View File
@@ -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
+295
View File
@@ -0,0 +1,295 @@
#!/usr/bin/env python3
"""Mock AI Agent — Asterisk/ARI 版.
外呼流程 (控制面纯 REST + RTP; 另加一个原生 socket websocket 仅用于保持 Stasis app 运行,
否则 ARI 会把进入 app 的通道直接挂断):
0. 连接 /ari/events?app=outbound —— 让 Stasis app 处于运行态 (只丢事件, 不解析业务)
1. POST /bridges 建立 mixing 桥
2. POST /channels/externalMedia 外部媒体通道 (ulaw), Asterisk 把桥内混音发到 external_host,
本侧回发地址 = 通道变量 UNICASTRTP_LOCAL_ADDRESS/PORT (兜底: 首包源地址)
3. POST /channels originate 经 PJSIP trunk 外呼被叫
4. 轮询被叫通道至 Up (应答), 双通道入桥
5. RTP 双向: RX=被叫音频(模拟 ASR 输入, RMS 统计); TX=440Hz 间歇音(模拟 LLM/TTS 输出)
6. DELETE 通道+桥 → BYE
纯 stdlib (urllib + asyncio + 原生 socket)。换真实 LLM/ASR: 替换 TX 音源与 RX 消费即可。
用法: mock-agent.py --ari-base http://asterisk1:8088 --external-host <本容器名>: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()
+244
View File
@@ -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
+18
View File
@@ -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"
+44
View File
@@ -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 ]