238 lines
9.0 KiB
Python
238 lines
9.0 KiB
Python
#!/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: <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:]
|
|
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
|