diff --git a/.github/workflows/python-package.yml b/.github/workflows/python-package.yml index c80585f..0a971d2 100644 --- a/.github/workflows/python-package.yml +++ b/.github/workflows/python-package.yml @@ -25,14 +25,10 @@ jobs: - name: Install dependencies run: | python -m pip install --upgrade pip - pip install ruff pytest - if [ -f requirements.txt ]; then pip install -r requirements.txt; fi + python -m pip install ruff pytest pytest-asyncio ".[voip]" - name: Lint with ruff run: | - # stop the build if there are Python syntax errors or undefined names - ruff check --output-format=github --select=E9,F63,F7,F82 --target-version=py37 . - # default set of ruff rules with GitHub Annotations - ruff check --output-format=github --target-version=py37 . + ruff check --output-format=github --select=E4,E7,E9,F --target-version=py310 . - name: Test with pytest run: | pytest diff --git a/.gitignore b/.gitignore index 68bc17f..12432fc 100644 --- a/.gitignore +++ b/.gitignore @@ -157,4 +157,4 @@ cython_debug/ # be found at https://github.com/github/gitignore/blob/main/Global/JetBrains.gitignore # and can be added to the global gitignore or merged into this file. For a more nuclear # option (not recommended) you can uncomment the following to ignore the entire idea folder. -#.idea/ +/.idea/ diff --git a/README.md b/README.md index 07d1602..8033116 100644 --- a/README.md +++ b/README.md @@ -40,6 +40,12 @@ authorize the account, go to your [cabinet](https://console.green-api.com/) and python -m pip install whatsapp-api-client-python ``` +VoIP calls require the optional dependencies: + +```shell +python -m pip install "whatsapp-api-client-python[voip]" +``` + ## Import ``` @@ -302,6 +308,8 @@ asyncio.run(main()) | Example of account methods asynchronously | [accountMethodsAsync.py](./examples/async/accountMethodsAsync.py) | | Example of getting last incoming and outgoing calls | [lastCalls.py](./examples/sync/lastCalls.py) | | Example of getting last calls asynchronously | [lastCallsAsync.py](./examples/async/lastCallsAsync.py) | +| Outgoing headless VoIP call | [headless_call.py](./examples/async/voip/headless_call.py) | +| Incoming headless VoIP call | [incoming_call.py](./examples/async/voip/incoming_call.py) | | Example of sending a message with link preview options | [sendMessageWithPreview.py](./examples/sync/sending/sendMessageWithPreview.py) | | Example of sending a message with link preview options asynchronously | [sendMessageWithPreviewAsync.py](./examples/async/sending/sendMessageWithPreviewAsync.py) | | Example of queues methods (counts, clear webhooks queue) | [queuesMethods.py](./examples/sync/queuesMethods.py) | diff --git a/docs/README.md b/docs/README.md index 2ea56a3..12e28a6 100644 --- a/docs/README.md +++ b/docs/README.md @@ -39,6 +39,12 @@ whatsapp-api-client-python - библиотека для интеграции с python -m pip install whatsapp-api-client-python ``` +Для VoIP нужны дополнительные зависимости: + +```shell +python -m pip install "whatsapp-api-client-python[voip]" +``` + ## Импорт ``` @@ -303,6 +309,8 @@ asyncio.run(main()) | Пример асинхронных методов аккаунта | [accountMethodsAsync.py](../examples/async/accountMethodsAsync.py) | | Пример получения журнала звонков | [lastCalls.py](../examples/sync/lastCalls.py) | | Пример асинхронного получения журнала звонков | [lastCallsAsync.py](../examples/async/lastCallsAsync.py) | +| Исходящий VoIP-звонок без аудиоустройств | [headless_call.py](../examples/async/voip/headless_call.py) | +| Входящий VoIP-звонок без аудиоустройств | [incoming_call.py](../examples/async/voip/incoming_call.py) | | Пример отправки сообщения с настройками превью ссылки | [sendMessageWithPreview.py](../examples/sync/sending/sendMessageWithPreview.py) | | Пример асинхронной отправки сообщения с превью ссылки | [sendMessageWithPreviewAsync.py](../examples/async/sending/sendMessageWithPreviewAsync.py) | | Пример методов очереди (счётчики, очистка вебхуков) | [queuesMethods.py](../examples/sync/queuesMethods.py) | diff --git a/examples/async/voip/headless_call.py b/examples/async/voip/headless_call.py new file mode 100644 index 0000000..7f9a78e --- /dev/null +++ b/examples/async/voip/headless_call.py @@ -0,0 +1,58 @@ +"""Outgoing voice call using a WAV file and decoded audio frames, no audio devices.""" + +import asyncio +import os +import sys + +from aiortc.contrib.media import MediaPlayer + +from whatsapp_api_client_python.API import GreenAPI +from whatsapp_api_client_python.tools.voip import CallAudio, FrameAudioSink + + +async def make_audio() -> CallAudio: + player = MediaPlayer(sys.argv[2]) + + if player.audio is None: + raise ValueError("The file has no audio stream") + + async def on_frame(frame): + # Replace this with your application's frame consumer. + print(f"received {frame.samples} samples") + + sink = FrameAudioSink(on_frame) + + async def close(): + await sink.closeAsync() + player.audio.stop() + + return CallAudio(player.audio, sink.attachAsync, close) + + +async def main(): + api = GreenAPI(os.environ["GREEN_API_ID"], os.environ["GREEN_API_TOKEN"]) + calls = api.voip.connect(audio_factory=make_audio) + ended = asyncio.Event() + calls.on("incoming_call", lambda info: print("Incoming call:", info.wid)) + + def on_end(detail): + print("Call ended:", detail) + ended.set() + + calls.on("end_call", on_end) + calls.on("error", lambda detail: print("Call error:", detail)) + + try: + await calls.openAsync(timeout=30) + await api.voip.dialAsync(sys.argv[1]) + await calls.startAudioAsync() + await ended.wait() + finally: + await calls.closeAsync() + + +if __name__ == "__main__": + if len(sys.argv) != 3: + raise SystemExit("Usage: headless_call.py ") + + asyncio.run(main()) diff --git a/examples/async/voip/incoming_call.py b/examples/async/voip/incoming_call.py new file mode 100644 index 0000000..b9d102d --- /dev/null +++ b/examples/async/voip/incoming_call.py @@ -0,0 +1,63 @@ +"""Answer an incoming voice call using a WAV file, without audio devices.""" + +import asyncio +import os +import sys + +from aiortc.contrib.media import MediaPlayer + +from whatsapp_api_client_python.API import GreenAPI +from whatsapp_api_client_python.tools.voip import CallAudio, FrameAudioSink + + +async def make_audio() -> CallAudio: + player = MediaPlayer(sys.argv[1]) + + if player.audio is None: + raise ValueError("The file has no audio stream") + + async def on_frame(frame): + # Replace this with your application's frame consumer. + print(f"received {frame.samples} samples") + + sink = FrameAudioSink(on_frame) + + async def close(): + await sink.closeAsync() + player.audio.stop() + + return CallAudio(player.audio, sink.attachAsync, close) + + +async def main(): + api = GreenAPI(os.environ["GREEN_API_ID"], os.environ["GREEN_API_TOKEN"]) + calls = api.voip.connect(audio_factory=make_audio) + incoming = asyncio.Queue() + ended = asyncio.Event() + + calls.on("incoming_call", incoming.put_nowait) + calls.on("end_call", lambda detail: ended.set()) + calls.on("error", lambda detail: print("Call error:", detail)) + + try: + await calls.openAsync(timeout=30) + + info = await incoming.get() + + print("Incoming call:", info.wid) + + if calls.state is None or calls.state.state != "inc-call": + return + + await api.voip.acceptAsync() + await calls.startAudioAsync() + await ended.wait() + finally: + await calls.closeAsync() + + +if __name__ == "__main__": + if len(sys.argv) != 2: + raise SystemExit("Usage: incoming_call.py ") + + asyncio.run(main()) diff --git a/setup.py b/setup.py index e398c3b..47085f1 100644 --- a/setup.py +++ b/setup.py @@ -5,7 +5,7 @@ setup( name="whatsapp-api-client-python", - version="0.0.54", + version="0.0.55", description=( "This library helps you easily create" " a Python application with WhatsApp API." @@ -40,6 +40,7 @@ "Creative Commons Attribution-NoDerivatives 4.0 International" " (CC BY-ND 4.0)" ), - install_requires=["requests-2.34.2", "aiofiles>=24.1.0", "aiogram>=3.28.2", "aiohttp>=3.13.5"], + install_requires=["requests==2.34.2", "aiofiles>=24.1.0", "aiogram>=3.28.2", "aiohttp>=3.13.5"], + extras_require={"voip": ["aiortc>=1.15,<2", "websockets>=16,<17"]}, python_requires=">=3.10" ) diff --git a/tests/test_voip.py b/tests/test_voip.py new file mode 100644 index 0000000..2266529 --- /dev/null +++ b/tests/test_voip.py @@ -0,0 +1,628 @@ +"""Call scenarios adapted from the standalone WA VoIP client.""" + +import asyncio +import json +from types import SimpleNamespace + +from aiohttp import web +from aiohttp.test_utils import TestServer +import pytest + +from whatsapp_api_client_python.API import GreenAPI, GreenAPIError +from whatsapp_api_client_python.response import Response +from whatsapp_api_client_python.tools.voip import CallAudio, CallsConnection +from whatsapp_api_client_python.tools.voip.signaling import ReconnectingSocket + + +async def until(predicate, turns=60): + for _ in range(turns): + if predicate(): + return + await asyncio.sleep(0) + raise AssertionError("Expected async progress") + + +class FakeSocket: + def __init__(self): + self.listeners = {} + self.sent = [] + self.closed = False + + def on(self, name, callback): + self.listeners.setdefault(name, []).append(callback) + + async def open(self, **kwargs): + pass + + async def send(self, value): + self.sent.append(value) + + async def close(self): + self.closed = True + + async def emit(self, name, detail=None): + for callback in self.listeners.get(name, ()): + value = callback(detail) + if asyncio.iscoroutine(value): + await value + + +class FakeTrack: + kind = "audio" + + +class FakeAudioFactory: + def __init__(self): + self.sessions = [] + + async def __call__(self): + track = FakeTrack() + remote = [] + closed = [] + + async def attach(value): + remote.append(value) + + async def close(): + closed.append(True) + + session = CallAudio(track, attach, close) + session.remote = remote + session.closed = closed + self.sessions.append(session) + return session + + +class FakeBridge: + def __init__(self, ice_servers): + self.ice_servers = ice_servers + self.calls = [] + self.closed = False + self.track_callback = None + self.local_description = {"type": "offer", "sdp": "v=0"} + + def add_track(self, track): + self.calls.append(("track", track)) + + def on_track(self, callback): + self.track_callback = callback + + async def create_offer(self): + return self.local_description + + async def set_local_description(self, offer): + self.calls.append(("local", offer)) + + def local_candidates(self): + return [] + + async def set_remote_description(self, answer): + self.calls.append(("remote", answer)) + + async def add_ice_candidate(self, candidate): + self.calls.append(("candidate", candidate)) + + async def close(self): + self.closed = True + + +class FakeVoip: + def __init__(self): + self.calls = [] + + async def getIceServersAsync(self): + self.calls.append("ice") + return [{"urls": "stun:example.test"}] + + +@pytest.fixture +def harness(): + socket, voip, audio = FakeSocket(), FakeVoip(), FakeAudioFactory() + bridges = [] + + def make_bridge(servers): + bridge = FakeBridge(servers) + bridges.append(bridge) + return bridge + + calls = CallsConnection(voip, socket=socket, bridge_factory=make_bridge, audio_factory=audio) + return calls, socket, voip, audio, bridges + + +async def start(calls, socket): + sent = len(socket.sent) + task = asyncio.create_task(calls.startAudioAsync()) + await until(lambda: len(socket.sent) > sent) + return task + + +@pytest.mark.asyncio +async def test_rest_dial_and_accept_precede_offer(harness, monkeypatch): + calls, socket, _, _, _ = harness + api = GreenAPI("123", "secret") + timeline = [] + + class Transport: + async def request(self, method, url, **kwargs): + endpoint = url.split("/")[-2] + + timeline.append(("REST", endpoint)) + + if endpoint == "callsGetIceServers": + return SimpleNamespace(status_code=200, text="[]") + + return SimpleNamespace(status_code=204, text="") + + async def unexpected_request(*args, **kwargs): + raise AssertionError("VoIP must not use the common requestAsync policy") + + monkeypatch.setattr(api, "requestAsync", unexpected_request) + + api.voip._transport = Transport() + calls._voip = api.voip + original_send = socket.send + + async def send(frame): + timeline.append(("WS", frame["type"])) + await original_send(frame) + + socket.send = send + await api.voip.dialAsync("79991234567") + task = await start(calls, socket) + assert timeline[:3] == [("REST", "callsDial"), ("REST", "callsGetIceServers"), ("WS", "offer")] + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + await calls.stopAudioAsync() + socket.sent.clear() + timeline.clear() + await api.voip.acceptAsync() + task = await start(calls, socket) + assert timeline[:3] == [("REST", "callsAccept"), ("REST", "callsGetIceServers"), ("WS", "offer")] + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + await calls.closeAsync() + + +@pytest.mark.asyncio +async def test_early_candidate_is_applied_after_answer(harness): + calls, socket, _, _, bridges = harness + task = await start(calls, socket) + candidate = {"candidate": "candidate:test", "sdpMid": "0"} + await socket.emit("message", {"type": "ice-candidate", "candidate": candidate}) + assert ("candidate", candidate) not in bridges[0].calls + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + assert bridges[0].calls[-2:] == [("remote", {"type": "answer", "sdp": "v=0"}), ("candidate", candidate)] + await calls.closeAsync() + + +@pytest.mark.asyncio +async def test_stop_and_close_do_not_hang_up(harness): + calls, socket, _, audio, bridges = harness + await socket.emit("message", {"type": "state", "state": {"state": "on-call"}}) + task = await start(calls, socket) + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + await calls.stopAudioAsync() + assert calls.state.state == "on-call" + assert socket.sent[-1] == {"type": "stop"} + assert bridges[0].closed and audio.sessions[0].closed == [True] + await calls.closeAsync() + assert socket.closed + assert socket.sent.count({"type": "stop"}) == 1 + + +@pytest.mark.asyncio +async def test_reconnect_creates_fresh_audio_and_ignores_old_track(harness): + calls, socket, _, audio, bridges = harness + await socket.emit("message", {"type": "state", "state": {"state": "on-call"}}) + task = await start(calls, socket) + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + await socket.emit("disconnect", {"reason": "lost", "code": 1006, "permanent": False}) + assert audio.sessions[0].closed == [True] + await socket.emit("message", {"type": "state", "state": {"state": "on-call"}}) + await until(lambda: len(bridges) == 2 and len(socket.sent) == 2) + assert audio.sessions[1].local_track is not audio.sessions[0].local_track + bridges[0].track_callback(FakeTrack()) + await asyncio.sleep(0) + assert audio.sessions[1].remote == [] + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await calls.closeAsync() + + +@pytest.mark.asyncio +async def test_pending_error_closes_bridge_but_error_after_answer_keeps_it(harness): + calls, socket, _, audio, bridges = harness + errors = [] + + calls.on("error", errors.append) + + task = await start(calls, socket) + await socket.emit("message", {"type": "error", "message": "no active call"}) + with pytest.raises(RuntimeError, match="no active call"): + await task + assert bridges[0].closed and audio.sessions[0].closed == [True] + await socket.emit("message", {"type": "state", "state": {"state": "on-call"}}) + task = await start(calls, socket) + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + await socket.emit("message", {"type": "error", "message": "calls unavailable"}) + assert errors == [{"message": "calls unavailable"}] + assert calls.has_audio_bridge + assert not bridges[1].closed and audio.sessions[1].closed == [] + await socket.emit("message", {"type": "state", "state": {"state": "idle"}}) + assert bridges[1].closed and audio.sessions[1].closed == [True] + await calls.closeAsync() + + +@pytest.mark.asyncio +async def test_idle_reports_end_even_when_stop_send_fails(harness): + calls, socket, _, audio, bridges = harness + ended = [] + calls.on("end_call", ended.append) + await socket.emit("message", {"type": "state", "state": {"state": "on-call"}}) + task = await start(calls, socket) + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + + async def failed_send(frame): + if frame["type"] == "stop": + raise ConnectionError("socket closed") + + socket.send = failed_send + await socket.emit("message", {"type": "state", "state": {"state": "idle", "reason": "hangup"}}) + assert ended == [{"reason": "call-ended", "cause": "hangup"}] + assert bridges[0].closed and audio.sessions[0].closed == [True] + await calls.closeAsync() + + +@pytest.mark.asyncio +async def test_permanent_refusal_does_not_resume(harness): + calls, socket, _, _, bridges = harness + errors, ended = [], [] + calls.on("error", errors.append) + calls.on("end_call", ended.append) + await socket.emit("message", {"type": "state", "state": {"state": "out-call"}}) + task = await start(calls, socket) + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + await socket.emit("message", {"type": "error", "message": "calls disabled"}) + await socket.emit("disconnect", {"reason": "calls disabled", "code": 4001, "permanent": True}) + await socket.emit("message", {"type": "state", "state": {"state": "out-call"}}) + assert errors == [{"message": "calls disabled"}] + assert ended == [{"reason": "connection-lost"}] + assert len(bridges) == 1 + await calls.closeAsync() + + +@pytest.mark.asyncio +async def test_pending_negotiation_fails_on_disconnect(harness): + calls, socket, _, audio, bridges = harness + await socket.emit("message", {"type": "state", "state": {"state": "out-call"}}) + task = await start(calls, socket) + await socket.emit("disconnect", {"reason": "lost", "code": 1006, "permanent": False}) + with pytest.raises(RuntimeError, match="Socket disconnected"): + await task + assert bridges[0].closed and audio.sessions[0].closed == [True] + await calls.closeAsync() + + +@pytest.mark.asyncio +async def test_bad_candidate_after_answer_closes_bridge(harness): + calls, socket, _, audio, bridges = harness + errors = [] + calls.on("error", errors.append) + task = await start(calls, socket) + await socket.emit("message", {"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}}) + await task + + async def bad_candidate(value): + raise ValueError("invalid candidate") + + bridges[0].add_ice_candidate = bad_candidate + await socket.emit("message", {"type": "ice-candidate", "candidate": {"candidate": "broken"}}) + assert errors == [{"message": "invalid candidate"}] + assert bridges[0].closed and audio.sessions[0].closed == [True] + await calls.closeAsync() + + +@pytest.mark.asyncio +async def test_close_cancels_pending_remote_attachment(harness): + calls, socket, _, audio, bridges = harness + entered = asyncio.Event() + cancelled = asyncio.Event() + + async def attach(track): + entered.set() + try: + await asyncio.Event().wait() + except asyncio.CancelledError: + cancelled.set() + raise + + task = await start(calls, socket) + audio.sessions[0].on_remote_track = attach + bridges[0].track_callback(FakeTrack()) + await entered.wait() + await calls.closeAsync() + with pytest.raises(RuntimeError, match="closed"): + await task + assert cancelled.is_set() + assert audio.sessions[0].closed == [True] + + +@pytest.mark.asyncio +async def test_close_while_audio_factory_pending(harness): + calls, _, _, _, _ = harness + waiting = asyncio.Event() + release = asyncio.Event() + audio = FakeAudioFactory() + + async def delayed(): + waiting.set() + await release.wait() + return await audio() + + calls._audio_factory = delayed + task = asyncio.create_task(calls.startAudioAsync()) + await waiting.wait() + await calls.closeAsync() + release.set() + with pytest.raises(RuntimeError, match="closed"): + await task + assert audio.sessions[0].closed == [True] + + +@pytest.mark.asyncio +async def test_websocket_url_and_rest_error(): + api = GreenAPI("123", "secret", host="https://example.test/base/") + assert api.voip.websocketUrl() == "wss://example.test/base/waInstance123/callsRtc/secret" + + api.voip._transport = StaticTransport(403, "denied") + + with pytest.raises(RuntimeError, match="callsDial failed"): + await api.voip.dialAsync("79991234567") + + +class StaticTransport: + def __init__(self, status, body=""): + self.status = status + self.body = body + self.requests = [] + + async def request(self, method, url, **kwargs): + self.requests.append((method, url, kwargs)) + return SimpleNamespace(status_code=self.status, text=self.body) + + +@pytest.mark.parametrize( + ("name_field", "expected"), + [({"name": "Alice"}, "Alice"), ({"name": ""}, ""), ({"name": None}, None), ({}, None)], +) +@pytest.mark.asyncio +async def test_call_name_preserves_empty_and_missing_values(harness, name_field, expected): + payload = {"state": "inc-call", "info": {"id": "call-1", "wid": "200@lid", **name_field}} + api = GreenAPI("123", "secret") + api.voip._transport = StaticTransport(200, json.dumps(payload)) + + assert (await api.voip.getStateAsync()).info.name == expected + + calls, socket, _, _, _ = harness + incoming = [] + calls.on("incoming_call", incoming.append) + await socket.emit("message", {"type": "state", "state": payload}) + + assert calls.state.info.name == expected + assert len(incoming) == 1 + assert incoming[0].wid == "200@lid" + assert incoming[0].name == expected + await calls.closeAsync() + + +@pytest.mark.parametrize("raise_errors", [False, True]) +@pytest.mark.asyncio +async def test_voip_rest_accepts_200_and_204_independently_of_sdk_policy(raise_errors): + api = GreenAPI("123", "secret", raise_errors=raise_errors, host="https://example.test") + transport = StaticTransport(200, '{"state":"idle"}') + api.voip._transport = transport + + assert (await api.voip.getStateAsync()).state == "idle" + + assert transport.requests[0] == ( + "GET", "https://example.test/waInstance123/callsState/secret", + {"headers": {"User-Agent": "GREEN-API_SDK_PY/1.0"}}, + ) + + transport.status, transport.body = 204, "" + + await api.voip.dialAsync("79991234567") + + assert transport.requests[-1] == ( + "POST", "https://example.test/waInstance123/callsDial/secret", + {"headers": {"User-Agent": "GREEN-API_SDK_PY/1.0", "Content-Type": "application/json"}, + "json": {"chatId": "79991234567@c.us"}}, + ) + + for command in (api.voip.acceptAsync, api.voip.rejectAsync, api.voip.hangUpAsync): + await command() + assert transport.requests[-1][0] == "POST" + assert transport.requests[-1][2] == {"headers": {"User-Agent": "GREEN-API_SDK_PY/1.0"}} + + +@pytest.mark.parametrize("raise_errors", [False, True]) +@pytest.mark.asyncio +async def test_voip_rest_rejects_http_error_and_invalid_json(raise_errors): + api = GreenAPI("123", "secret", raise_errors=raise_errors) + transport = StaticTransport(403, "denied") + api.voip._transport = transport + + with pytest.raises(RuntimeError, match="callsAccept failed: 403 denied"): + await api.voip.acceptAsync() + + transport.status, transport.body = 200, "not json" + + with pytest.raises(json.JSONDecodeError): + await api.voip.getIceServersAsync() + + +def test_common_response_keeps_original_200_only_contract(): + assert Response(200, "{}").data == {} + + for code in (201, 204): + response = Response(code, "{}") + assert response.data is None + assert response.error == "{}" + + +def test_common_http_handler_still_rejects_204(caplog): + api = GreenAPI("123", "secret", raise_errors=True) + + with pytest.raises(GreenAPIError, match="status code: 204"): + api._GreenApi__handle_response_async(204, "") + + api.raise_errors = False + + with caplog.at_level("ERROR"): + api._GreenApi__handle_response_async(204, "") + + assert "status code: 204" in caplog.text + + +@pytest.mark.asyncio +async def test_voip_uses_aiohttp_for_204_without_post_body(): + requests = [] + + async def handle(request): + requests.append((request.path, request.headers, await request.read())) + + if request.path.endswith("/callsReject/secret"): + return web.Response(status=302, headers={"Location": "/redirected"}) + + return web.Response(status=204) + + app = web.Application() + app.router.add_post("/{tail:.*}", handle) + + async with TestServer(app) as server: + api = GreenAPI("123", "secret", raise_errors=True, host=str(server.make_url("/"))) + + await api.voip.acceptAsync() + await api.voip.dialAsync("79991234567") + + with pytest.raises(RuntimeError, match="callsReject failed: 302"): + await api.voip.rejectAsync() + + assert requests[0][0] == "/waInstance123/callsAccept/secret" + assert requests[0][2] == b"" + assert requests[0][1]["User-Agent"] == "GREEN-API_SDK_PY/1.0" + assert "Content-Type" not in requests[0][1] + assert requests[1][0] == "/waInstance123/callsDial/secret" + assert requests[1][1]["Content-Type"] == "application/json" + assert json.loads(requests[1][2]) == {"chatId": "79991234567@c.us"} + assert len(requests) == 3 + + +@pytest.mark.asyncio +async def test_no_automatic_answer_timeout(harness): + calls, socket, _, audio, bridges = harness + task = await start(calls, socket) + + with pytest.raises(asyncio.TimeoutError): + await asyncio.wait_for(asyncio.shield(task), timeout=0.02) + + assert not task.done() and calls.has_audio_bridge + await calls.closeAsync() + + with pytest.raises(RuntimeError, match="closed"): + await task + + assert bridges[0].closed and audio.sessions[0].closed == [True] + + +class RawSocket: + def __init__(self): + self.incoming = asyncio.Queue() + self.sent = [] + self.closed = False + + async def recv(self): + value = await self.incoming.get() + if isinstance(value, Exception): + raise value + return value + + async def send(self, text): + self.sent.append(text) + + async def close(self): + self.closed = True + + +class Closed(Exception): + def __init__(self, code): + self.code = code + self.reason = "closed" + + +@pytest.mark.asyncio +async def test_socket_reconnect_policy_and_permanent_refusal(): + sockets, delays = [], [] + + async def factory(url): + raw = RawSocket() + sockets.append(raw) + return raw + + async def sleep(seconds): + delays.append(seconds) + await asyncio.sleep(0) + + socket = ReconnectingSocket("wss://example.test", socket_factory=factory, sleep=sleep) + await socket.open() + sockets[0].incoming.put_nowait(Closed(1006)) + await until(lambda: len(sockets) == 2) + assert delays == [0.5] + sockets[1].incoming.put_nowait(Closed(4001)) + await until(lambda: socket.refused) + assert delays == [0.5] + await socket.close() + + +@pytest.mark.asyncio +async def test_ordered_frames_wait_for_slow_stop(harness): + calls, _, _, _, _ = harness + raw = RawSocket() + socket = ReconnectingSocket("wss://example.test", socket_factory=lambda url: asyncio.sleep(0, result=raw)) + calls._socket = socket + socket.on("message", calls._on_message) + socket.on("disconnect", calls._on_socket_disconnect) + events = [] + calls.on("state", lambda state: events.append(("state", state.state))) + calls.on("end_call", lambda detail: events.append(("end", detail["reason"]))) + calls.on("incoming_call", lambda info: events.append(("incoming", info.wid))) + await socket.open() + raw.incoming.put_nowait(json.dumps({"type": "state", "state": {"state": "on-call"}})) + await until(lambda: calls.state is not None) + task = asyncio.create_task(calls.startAudioAsync()) + await until(lambda: bool(raw.sent)) + raw.incoming.put_nowait(json.dumps({"type": "answer", "answer": {"type": "answer", "sdp": "v=0"}})) + await task + entered, release = asyncio.Event(), asyncio.Event() + original = raw.send + + async def slow_send(text): + if json.loads(text).get("type") == "stop": + entered.set() + await release.wait() + await original(text) + + raw.send = slow_send + raw.incoming.put_nowait(json.dumps({"type": "state", "state": {"state": "idle"}})) + raw.incoming.put_nowait(json.dumps({"type": "state", "state": {"state": "inc-call", "info": {"id": "next", "wid": "200@lid", "name": "Next"}}})) + await entered.wait() + assert events == [("state", "on-call"), ("state", "idle")] + release.set() + await until(lambda: len(events) == 5) + assert events[-3:] == [("end", "call-ended"), ("state", "inc-call"), ("incoming", "200@lid")] + await calls.closeAsync() diff --git a/tests/test_voip_rtc.py b/tests/test_voip_rtc.py new file mode 100644 index 0000000..5d7defc --- /dev/null +++ b/tests/test_voip_rtc.py @@ -0,0 +1,86 @@ +"""Local media and ICE checks without an account or audio hardware.""" + +import asyncio +from fractions import Fraction + +import pytest +from aiortc import AudioStreamTrack +from aiortc.mediastreams import MediaStreamError +from av import AudioFrame + +from whatsapp_api_client_python.tools.voip import FrameAudioSink +from whatsapp_api_client_python.tools.voip.rtc import ( + AiortcBridge, + candidate_from_json, + candidate_to_json, +) + + +class OneFrameTrack(AudioStreamTrack): + def __init__(self): + super().__init__() + self.sent = False + + async def recv(self): + if self.sent: + raise MediaStreamError + self.sent = True + frame = AudioFrame(format="s16", layout="mono", samples=960) + frame.planes[0].update(bytes(1920)) + frame.sample_rate = 48_000 + frame.time_base = Fraction(1, 48_000) + frame.pts = 0 + return frame + + +@pytest.mark.asyncio +async def test_frame_sink_delivers_async_callback_and_closes(): + frames = [] + + async def on_frame(frame): + frames.append(frame) + + sink = FrameAudioSink(on_frame) + await sink.attachAsync(OneFrameTrack()) + for _ in range(30): + if frames: + break + await asyncio.sleep(0) + assert len(frames) == 1 and isinstance(frames[0], AudioFrame) + await sink.closeAsync() + + +def test_candidate_wire_round_trip(): + value = { + "candidate": "candidate:1 1 udp 2122260223 192.0.2.1 54000 typ host", + "sdpMid": "0", "sdpMLineIndex": 0, + } + assert candidate_to_json(candidate_from_json(value)) == value + + +@pytest.mark.asyncio +async def test_two_real_peers_exchange_offer_answer_and_audio(): + caller, callee = AiortcBridge([]), AiortcBridge([]) + caller.add_track(AudioStreamTrack()) + callee.add_track(AudioStreamTrack()) + try: + offer = await caller.create_offer() + try: + await caller.set_local_description(offer) + except PermissionError: + pytest.skip("Network interfaces are unavailable in this sandbox") + assert caller.local_description["type"] == "offer" + assert caller.local_candidates() + await callee.set_remote_description(caller.local_description) + answer = await callee._peer.createAnswer() + await callee._peer.setLocalDescription(answer) + await caller.set_remote_description(callee.local_description) + + async def connected(): + while caller._peer.connectionState != "connected" or callee._peer.connectionState != "connected": + await asyncio.sleep(0.05) + + await asyncio.wait_for(connected(), timeout=10) + finally: + await caller.close() + await callee.close() diff --git a/whatsapp_api_client_python/API.py b/whatsapp_api_client_python/API.py index 78611f2..431376e 100644 --- a/whatsapp_api_client_python/API.py +++ b/whatsapp_api_client_python/API.py @@ -7,6 +7,7 @@ from requests.adapters import HTTPAdapter, Retry from .response import Response as GreenAPIResponse +from .tools.voip import Voip from .tools import ( account, contacts, @@ -68,6 +69,7 @@ def __init__( self.serviceMethods = serviceMethods.ServiceMethods(self) self.webhooks = webhooks.Webhooks(self) self.statuses = statuses.Statuses(self) + self.voip = Voip(self) self.logger = logging.getLogger("whatsapp-api-client-python") self.__prepare_logger() diff --git a/whatsapp_api_client_python/tools/voip/__init__.py b/whatsapp_api_client_python/tools/voip/__init__.py new file mode 100644 index 0000000..ee2d766 --- /dev/null +++ b/whatsapp_api_client_python/tools/voip/__init__.py @@ -0,0 +1,6 @@ +"""Asynchronous GREEN-API call signaling and aiortc support.""" + +from .audio import CallAudio, FrameAudioSink +from .calls import CallInfo, CallState, CallsConnection, Voip + +__all__ = ["CallAudio", "FrameAudioSink", "CallInfo", "CallState", "CallsConnection", "Voip"] diff --git a/whatsapp_api_client_python/tools/voip/audio.py b/whatsapp_api_client_python/tools/voip/audio.py new file mode 100644 index 0000000..b20e91d --- /dev/null +++ b/whatsapp_api_client_python/tools/voip/audio.py @@ -0,0 +1,77 @@ +"""Application-owned audio for one media bridge.""" + +from __future__ import annotations + +import asyncio +import inspect +import logging +from collections.abc import Awaitable, Callable +from dataclasses import dataclass +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from aiortc import MediaStreamTrack + + +@dataclass +class CallAudio: + local_track: MediaStreamTrack + on_remote_track: Callable[[MediaStreamTrack], Awaitable[None]] + close: Callable[[], Awaitable[None]] + + +AudioFactory = Callable[[], Awaitable[CallAudio]] + + +class FrameAudioSink: + """Deliver decoded remote frames to a synchronous or async callback.""" + + def __init__(self, on_frame: Callable): + self._on_frame = on_frame + self._tasks: set[asyncio.Task] = set() + self._closed = False + + async def attachAsync(self, track: MediaStreamTrack) -> None: + if self._closed: + raise RuntimeError("Audio sink is closed") + + if getattr(track, "kind", None) != "audio": + raise ValueError("Expected a remote audio track") + + task = asyncio.create_task(self._consume(track)) + + self._tasks.add(task) + task.add_done_callback(self._tasks.discard) + + async def _consume(self, track: MediaStreamTrack) -> None: + from aiortc.mediastreams import MediaStreamError + + try: + while True: + frame = await track.recv() + result = self._on_frame(frame) + + if inspect.isawaitable(result): + await result + + except MediaStreamError: + pass + except asyncio.CancelledError: + raise + except Exception: + logging.getLogger(__name__).exception("Remote audio frame consumer failed") + + async def closeAsync(self) -> None: + if self._closed: + return + + self._closed = True + tasks = tuple(self._tasks) + + for task in tasks: + task.cancel() + + if tasks: + await asyncio.gather(*tasks, return_exceptions=True) + + self._tasks.clear() diff --git a/whatsapp_api_client_python/tools/voip/calls.py b/whatsapp_api_client_python/tools/voip/calls.py new file mode 100644 index 0000000..2500392 --- /dev/null +++ b/whatsapp_api_client_python/tools/voip/calls.py @@ -0,0 +1,476 @@ +"""REST commands, call state and media bridge lifecycle.""" + +import asyncio +import inspect +import json +import logging +from dataclasses import dataclass +from typing import Literal, Mapping +from urllib.parse import urlsplit, urlunsplit + +import aiohttp + +from .audio import AudioFactory, CallAudio +from .signaling import ReconnectingSocket + +# Types + +CallStateKind = Literal["idle", "inc-call", "out-call", "on-call"] + + +# Constants + +CALL_STATE_KINDS = frozenset(("inc-call", "out-call", "on-call")) + +_MISSING = object() + + +@dataclass(frozen=True) +class CallInfo: + id: str + wid: str + name: str | None + + +@dataclass(frozen=True) +class CallState: + state: CallStateKind + info: CallInfo | None = None + reason: str | None = None + + +def call_state_from_json(value: Mapping) -> CallState: + info = value.get("info") + + if isinstance(info, dict): + info = CallInfo(id=info["id"], wid=info["wid"], name=info.get("name")) + + else: + info = None + + return CallState(state=value["state"], info=info, reason=value.get("reason")) + + +class Voip: + def __init__(self, api, *, transport=None): + self._api = api + self._transport = transport + + def _url(self, method: str) -> str: + return ( + f"{self._api.host.rstrip('/')}/waInstance{self._api.idInstance}" + f"/{method}/{self._api.apiTokenInstance}" + ) + + def websocketUrl(self) -> str: + host = urlsplit(self._api.host.rstrip("/")) + + if host.scheme not in ("http", "https"): + raise ValueError("VoIP host must use http or https") + + scheme = "wss" if host.scheme == "https" else "ws" + path = host.path.rstrip("/") + f"/waInstance{self._api.idInstance}/callsRtc/{self._api.apiTokenInstance}" + + return urlunsplit((scheme, host.netloc, path, "", "")) + + async def _request(self, method: str, endpoint: str, payload=_MISSING): + kwargs = {"headers": {"User-Agent": "GREEN-API_SDK_PY/1.0"}} + + if payload is not _MISSING: + kwargs["headers"]["Content-Type"] = "application/json" + kwargs["json"] = payload + + if self._transport is None: + timeout = aiohttp.ClientTimeout(total=None, connect=5, sock_read=5) + + async with aiohttp.ClientSession( + timeout=timeout, raise_for_status=False, trust_env=True, + skip_auto_headers={"Content-Type"}, + ) as transport: + async with transport.request( + method, self._url(endpoint), allow_redirects=False, **kwargs + ) as response: + status_code = response.status + body = await response.text() + else: + response = await self._transport.request(method, self._url(endpoint), **kwargs) + status_code = response.status_code + body = response.text + + if not 200 <= status_code < 300: + raise RuntimeError(f"{endpoint} failed: {status_code} {body}") + + if status_code == 204 or not body: + return None + + return json.loads(body) + + async def getStateAsync(self) -> CallState: + return call_state_from_json(await self._request("GET", "callsState")) + + async def getIceServersAsync(self): + return await self._request("GET", "callsGetIceServers") + + async def dialAsync(self, target: str) -> None: + chat_id = target if "@" in target else f"{target}@c.us" + await self._request("POST", "callsDial", {"chatId": chat_id}) + + async def acceptAsync(self) -> None: + await self._request("POST", "callsAccept") + + async def rejectAsync(self) -> None: + await self._request("POST", "callsReject") + + async def hangUpAsync(self) -> None: + await self._request("POST", "callsHangUp") + + def connect(self, *, audio_factory: AudioFactory, **callbacks) -> "CallsConnection": + connection = CallsConnection(self, audio_factory=audio_factory) + + for name, callback in callbacks.items(): + if not name.startswith("on_"): + raise TypeError(f"Unknown callback: {name}") + + connection.on(name[3:], callback) + + return connection + + +class CallsConnection: + """One callsRtc connection and its current media bridge.""" + + def __init__( + self, + voip: Voip, + *, + audio_factory: AudioFactory, + socket=None, + bridge_factory=None, + ): + self._voip = voip + self._socket = socket if socket is not None else ReconnectingSocket(voip.websocketUrl()) + self._bridge_factory = bridge_factory + self._audio_factory = audio_factory + self._state: CallState | None = None + self._peer_connection = None + self._audio: CallAudio | None = None + self._pending_remote_candidates: list[dict] = [] + self._remote_description_set = False + self._pending_bridge: asyncio.Future | None = None + self._resume_pending = False + self._refusal_reported = False + self._closed = False + self._resume_task: asyncio.Task | None = None + self._bridge_generation = 0 + self._callbacks: dict[str, list] = {} + self._callback_tasks: set[asyncio.Task] = set() + self._remote_tasks: set[asyncio.Task] = set() + self._socket.on("connect", lambda _: self._emit("connect")) + self._socket.on("disconnect", self._on_socket_disconnect) + self._socket.on("message", self._on_message) + + def on(self, name: str, callback) -> None: + """Register a callback. Async callbacks run separately from frame dispatch.""" + self._callbacks.setdefault(name, []).append(callback) + + def off(self, name: str, callback) -> None: + if callback in self._callbacks.get(name, ()): + self._callbacks[name].remove(callback) + + def _emit(self, name: str, detail=None) -> None: + for callback in tuple(self._callbacks.get(name, ())): + try: + result = callback(detail) + + if inspect.isawaitable(result): + task = asyncio.create_task(result) + self._callback_tasks.add(task) + task.add_done_callback(self._callback_tasks.discard) + task.add_done_callback(self._log_callback_error) + except Exception: + logging.getLogger(__name__).exception("VoIP %s callback failed", name) + + @staticmethod + def _log_callback_error(task: asyncio.Task) -> None: + if not task.cancelled(): + exc = task.exception() + + if exc is not None: + logging.getLogger(__name__).error( + "VoIP callback failed", exc_info=(type(exc), exc, exc.__traceback__) + ) + + @property + def state(self) -> CallState | None: + return self._state + + @property + def has_audio_bridge(self) -> bool: + return self._peer_connection is not None + + async def openAsync(self, *, timeout: float | None = None) -> None: + if self._closed: + raise RuntimeError("CallsConnection is closed") + + await self._socket.open(timeout=timeout) + + async def startAudioAsync(self) -> None: + if self._peer_connection is not None or self._pending_bridge is not None: + raise RuntimeError("Audio bridge already starting or active") + + if self._closed: + raise RuntimeError("CallsConnection is closed") + + self._pending_bridge = asyncio.get_running_loop().create_future() + pending = self._pending_bridge + + try: + await self._start_audio_internal(self._bridge_generation) + except BaseException as exc: + if self._pending_bridge is pending: + self._pending_bridge = None + + await self._teardown_bridge(False) + + if not pending.done(): + pending.set_exception(exc) + + elif pending.exception() is None: + raise + try: + await pending + except BaseException: + if self._pending_bridge is pending: + self._pending_bridge = None + await self._teardown_bridge(False) + + raise + + async def stopAudioAsync(self) -> None: + self._resume_pending = False + self._reject_pending(RuntimeError("Audio bridge stopped")) + await self._teardown_bridge(True) + + async def closeAsync(self) -> None: + if self._closed: + return + + self._closed = True + self._resume_pending = False + + self._reject_pending(RuntimeError("CallsConnection closed")) + await self._teardown_bridge(False) + if self._resume_task is not None: + self._resume_task.cancel() + + await self._socket.close() + + async def _start_audio_internal(self, generation: int) -> None: + audio = await self._audio_factory() + + if generation != self._bridge_generation: + await audio.close() + return + + self._audio = audio + ice_servers = await self._voip.getIceServersAsync() + + if generation != self._bridge_generation: + return + + if self._bridge_factory is None: + from .rtc import AiortcBridge + peer = AiortcBridge(ice_servers) + else: + peer = self._bridge_factory(ice_servers) + + self._peer_connection = peer + self._remote_description_set = False + self._pending_remote_candidates = [] + + peer.add_track(audio.local_track) + peer.on_track(lambda track: self._on_remote_track(track, generation, peer)) + + offer = await peer.create_offer() + + if generation != self._bridge_generation: + return + + await peer.set_local_description(offer) + + if generation != self._bridge_generation: + return + + offer = peer.local_description or offer + + await self._socket.send({"type": "offer", "offer": offer}) + + for candidate in peer.local_candidates(): + if generation != self._bridge_generation: + return + + await self._socket.send({"type": "ice-candidate", "candidate": candidate}) + + def _on_remote_track(self, track, generation: int, peer) -> None: + if generation != self._bridge_generation or peer is not self._peer_connection: + return + + audio = self._audio + + self._emit("remote_track", track) + + if audio is not None: + async def attach(): + if generation == self._bridge_generation and peer is self._peer_connection: + await audio.on_remote_track(track) + + task = asyncio.create_task(attach()) + + self._remote_tasks.add(task) + task.add_done_callback(self._remote_tasks.discard) + task.add_done_callback(self._log_callback_error) + + async def _teardown_bridge(self, notify_server: bool) -> None: + self._bridge_generation += 1 + peer, audio = self._peer_connection, self._audio + self._peer_connection = None + self._audio = None + self._pending_remote_candidates = [] + self._remote_description_set = False + remote_tasks = tuple(self._remote_tasks) + + for task in remote_tasks: + task.cancel() + + if remote_tasks: + await asyncio.gather(*remote_tasks, return_exceptions=True) + + try: + if notify_server and peer is not None: + await self._socket.send({"type": "stop"}) + finally: + try: + if peer is not None: + await peer.close() + finally: + if audio is not None: + await audio.close() + + def _reject_pending(self, exc: Exception) -> None: + pending = self._pending_bridge + self._pending_bridge = None + + if pending is not None and not pending.done(): + pending.set_exception(exc) + + async def _on_message(self, message: dict) -> None: + kind = message.get("type") if isinstance(message, dict) else None + + if kind == "state": + await self._on_state(call_state_from_json(message["state"])) + elif kind == "answer": + peer = self._peer_connection + + if peer is None: + return + + try: + await peer.set_remote_description(message["answer"]) + self._remote_description_set = True + + for candidate in self._pending_remote_candidates: + await peer.add_ice_candidate(candidate) + + self._pending_remote_candidates = [] + except Exception as exc: + self._reject_pending(exc) + await self._teardown_bridge(False) + return + + pending = self._pending_bridge + self._pending_bridge = None + + if pending is not None and not pending.done(): + pending.set_result(None) + elif kind == "ice-candidate": + if self._remote_description_set: + if self._peer_connection is not None: + try: + await self._peer_connection.add_ice_candidate(message["candidate"]) + except Exception as exc: + self._reject_pending(exc) + await self._teardown_bridge(False) + self._emit("error", {"message": str(exc)}) + else: + self._pending_remote_candidates.append(message["candidate"]) + elif kind == "error": + if self._pending_bridge is not None: + self._reject_pending(RuntimeError(message["message"])) + await self._teardown_bridge(False) + else: + self._refusal_reported = True + self._emit("error", {"message": message["message"]}) + + async def _on_state(self, state: CallState) -> None: + previous = self._state.state if self._state else None + self._state = state + + self._emit("state", state) + + if state.state == "inc-call" and previous != "inc-call" and state.info is not None: + self._emit("incoming_call", state.info) + return + + if previous in CALL_STATE_KINDS and state.state not in CALL_STATE_KINDS: + self._resume_pending = False + self._reject_pending(RuntimeError("Call ended during negotiation")) + + try: + await self._teardown_bridge(True) + except Exception: + logging.getLogger(__name__).exception("Could not send callsRtc stop") + + detail = {"reason": "call-ended"} + + if state.reason: + detail["cause"] = state.reason + + self._emit("end_call", detail) + + return + + if self._resume_pending and state.state in CALL_STATE_KINDS and not self._peer_connection and not self._pending_bridge: + self._resume_pending = False + self._resume_task = asyncio.create_task(self._resume_audio()) + + async def _resume_audio(self) -> None: + try: + await self.startAudioAsync() + except Exception as exc: + self._emit("error", {"message": str(exc)}) + + async def _on_socket_disconnect(self, detail: dict) -> None: + previous = self._state.state if self._state else None + permanent = detail.get("permanent", False) + + if previous in CALL_STATE_KINDS: + self._resume_pending = self._peer_connection is not None + + self._reject_pending(RuntimeError("Socket disconnected during negotiation")) + await self._teardown_bridge(False) + + if permanent: + self._resume_pending = False + + if not self._resume_pending: + self._state = None + self._emit("end_call", {"reason": "connection-lost"}) + else: + self._state = None + + if permanent and not self._refusal_reported: + self._emit("error", {"message": detail["reason"]}) + + self._refusal_reported = False + + self._emit("disconnect", detail) diff --git a/whatsapp_api_client_python/tools/voip/rtc.py b/whatsapp_api_client_python/tools/voip/rtc.py new file mode 100644 index 0000000..a53ce94 --- /dev/null +++ b/whatsapp_api_client_python/tools/voip/rtc.py @@ -0,0 +1,121 @@ +"""Thin aiortc boundary; retain the existing SDP and ICE wire behavior.""" + +from aiortc import RTCConfiguration, RTCIceServer, RTCPeerConnection, RTCSessionDescription +from aiortc.sdp import candidate_from_sdp, candidate_to_sdp + + +def candidate_from_json(value: dict): + text = value["candidate"] + + if text.startswith("candidate:"): + text = text[len("candidate:"):] + + candidate = candidate_from_sdp(text) + candidate.sdpMid = value.get("sdpMid") + candidate.sdpMLineIndex = value.get("sdpMLineIndex") + + if candidate.sdpMid is None and candidate.sdpMLineIndex is None: + raise ValueError("ICE candidate needs sdpMid or sdpMLineIndex") + + return candidate + + +def candidate_to_json(candidate) -> dict: + return { + "candidate": "candidate:" + candidate_to_sdp(candidate), + "sdpMid": candidate.sdpMid, + "sdpMLineIndex": candidate.sdpMLineIndex, + } + + +class AiortcBridge: + def __init__(self, ice_servers, *, peer_factory=None): + servers = [RTCIceServer( + urls=entry["urls"], + username=entry.get("username"), + credential=entry.get("credential"), + ) for entry in ice_servers] + + factory = peer_factory or RTCPeerConnection + self._peer = factory(RTCConfiguration(iceServers=servers)) + + self._on_track = None + + if hasattr(self._peer, "on"): + self._peer.on("track", self._handle_track) + + @property + def local_description(self) -> dict | None: + description = getattr(self._peer, "localDescription", None) + + if description is None: + return None + + return {"type": description.type, "sdp": description.sdp} + + def add_track(self, track) -> None: + self._peer.addTrack(track) + + def on_track(self, callback) -> None: + self._on_track = callback + + def _handle_track(self, track) -> None: + if self._on_track is not None: + self._on_track(track) + + def local_candidates(self) -> list[dict]: + """Also send gathered SDP candidates as separate signaling frames.""" + + description = self.local_description + + if description is None: + return [] + + result = [] + mid = None + index = -1 + candidates = [] + + def flush(): + for text in candidates: + result.append({"candidate": text, "sdpMid": mid, "sdpMLineIndex": index}) + + for line in description["sdp"].splitlines(): + if line.startswith("m="): + if index >= 0: + flush() + + index += 1 + mid = None + candidates = [] + elif line.startswith("a=mid:"): + mid = line[len("a=mid:"):] + elif line.startswith("a=candidate:"): + candidates.append(line[len("a="):]) + + if index >= 0: + flush() + + return result + + async def create_offer(self) -> dict: + offer = await self._peer.createOffer() + return {"type": offer.type, "sdp": offer.sdp} + + async def set_local_description(self, offer: dict) -> None: + await self._peer.setLocalDescription(RTCSessionDescription(**offer)) + + async def set_remote_description(self, answer: dict) -> None: + await self._peer.setRemoteDescription(RTCSessionDescription(**answer)) + + async def add_ice_candidate(self, value: dict) -> None: + if value is None or not value.get("candidate"): + await self._peer.addIceCandidate(None) + else: + await self._peer.addIceCandidate(candidate_from_json(value)) + + async def close(self) -> None: + await self._peer.close() + + async def get_stats(self): + return await self._peer.getStats() diff --git a/whatsapp_api_client_python/tools/voip/signaling.py b/whatsapp_api_client_python/tools/voip/signaling.py new file mode 100644 index 0000000..eac0a63 --- /dev/null +++ b/whatsapp_api_client_python/tools/voip/signaling.py @@ -0,0 +1,173 @@ +"""callsRtc WebSocket with the original retry and ordered dispatch policy.""" + +from collections import defaultdict +import asyncio +import inspect +import json +import logging + +# Constants + +INITIAL_BACKOFF_MS = 500 + +MAX_BACKOFF_MS = 10 * 1000 + +PERMANENT_CLOSE_MIN = 4000 + +PERMANENT_CLOSE_MAX = 4999 + + +class ReconnectingSocket: + def __init__(self, url: str, *, socket_factory=None, sleep=None): + self._url = url + self._socket_factory = socket_factory + self._sleep = sleep or asyncio.sleep + self._socket = None + self._closed_by_user = False + self._refused_by_server = False + self._backoff_ms = INITIAL_BACKOFF_MS + self._task: asyncio.Task | None = None + self._connected = asyncio.Event() + self._listeners = defaultdict(list) + + def on(self, name: str, callback) -> None: + self._listeners[name].append(callback) + + @property + def refused(self) -> bool: + return self._refused_by_server + + async def open(self, *, timeout: float | None = None) -> None: + if self._closed_by_user: + raise RuntimeError("ReconnectingSocket is closed") + + if self._task is None: + self._task = asyncio.create_task(self._connect()) + + waiter = asyncio.create_task(self._connected.wait()) + + try: + done, _ = await asyncio.wait((waiter, self._task), timeout=timeout, return_when=asyncio.FIRST_COMPLETED) + + if waiter in done: + return + + if self._task in done: + raise RuntimeError("callsRtc connection was refused or closed") + + raise TimeoutError("Timed out connecting to callsRtc") + finally: + waiter.cancel() + + async def send(self, data) -> None: + if self._socket is None: + raise RuntimeError("ReconnectingSocket is not connected") + + await self._socket.send(json.dumps(data, separators=(",", ":"))) + + async def close(self) -> None: + self._closed_by_user = True + socket = self._socket + self._socket = None + + self._connected.clear() + + if socket is not None: + await socket.close() + + if self._task is not None: + self._task.cancel() + + try: + await self._task + except asyncio.CancelledError: + pass + + async def _dispatch(self, name: str, detail=None) -> None: + # Internal handlers must finish before the next frame is read. + for callback in tuple(self._listeners[name]): + try: + result = callback(detail) + + if inspect.isawaitable(result): + await result + except Exception: + logging.getLogger(__name__).exception("callsRtc %s handler failed", name) + + async def _connect(self) -> None: + if self._socket_factory is None: + from websockets.asyncio.client import connect + self._socket_factory = connect + + while not self._closed_by_user: + try: + socket = await self._socket_factory(self._url) + + if self._closed_by_user: + await socket.close() + return + + self._socket = socket + self._connected.set() + self._backoff_ms = INITIAL_BACKOFF_MS + + await self._dispatch("connect") + + while not self._closed_by_user: + try: + raw = await socket.recv() + except Exception as exc: + code = getattr(exc, "code", None) + + if code is None: + frame = getattr(exc, "rcvd", None) + code = getattr(frame, "code", 1006) + reason = getattr(frame, "reason", "") + else: + reason = getattr(exc, "reason", "") + + self._socket = None + + self._connected.clear() + + if self._closed_by_user: + return + + permanent = PERMANENT_CLOSE_MIN <= code <= PERMANENT_CLOSE_MAX + + self._refused_by_server |= permanent + + await self._dispatch("disconnect", { + "reason": reason or "connection closed", "code": code, + "permanent": permanent, + }) + + if permanent: + return + + break + try: + parsed = json.loads(raw) + except (TypeError, ValueError): + continue + + await self._dispatch("message", parsed) + except asyncio.CancelledError: + raise + except Exception as exc: + if self._closed_by_user: + return + + self._socket = None + + self._connected.clear() + + await self._dispatch("disconnect", { + # Exception text may contain the URL and its API token. + "reason": f"{type(exc).__name__}: connection failed", "code": 1006, + "permanent": False, + }) + if not self._closed_by_user: + delay = self._backoff_ms + self._backoff_ms = min(self._backoff_ms * 2, MAX_BACKOFF_MS) + await self._sleep(delay / 1000)