diff --git a/CHANGELOG.md b/CHANGELOG.md index 1ae7d3aa..80ce69a8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,7 @@ #### 4.4.0 (2026-11-08) **Latest** +- [server] dispatch every message status included in a batched webhook update - [server] enhance logging infrastructure with update hash and context binding - [cli] add commands to list and download official example bots from GitHub - [cli] add `--access-log` option to command line arguments for production and development diff --git a/pywa/server.py b/pywa/server.py index e3389eef..9e9c9760 100644 --- a/pywa/server.py +++ b/pywa/server.py @@ -8,7 +8,7 @@ import warnings from collections import OrderedDict from collections.abc import Callable -from typing import TYPE_CHECKING +from typing import TYPE_CHECKING, cast from . import _helpers as helpers from . import errors, handlers, utils @@ -20,7 +20,13 @@ setup_console_logging, ) from .errors import PywaDeprecationWarning, PywaWarning -from .types import AccountUpdate, MessageType, RawUpdate, UserPreferenceCategory +from .types import ( + AccountUpdate, + MessageStatus, + MessageType, + RawUpdate, + UserPreferenceCategory, +) from .types.base_update import ( BaseUpdate, ContinueHandling, @@ -400,14 +406,20 @@ def _call_handlers(self: "WhatsApp", raw_update: RawUpdate) -> None: return log.debug("Dispatched to %s", handler_type.__name__) try: - constructed_update: BaseUpdate = self._handlers_to_updates[ - handler_type - ].from_update(client=self, update=raw_update) - if log.isEnabledFor(logging.DEBUG): - log.debug("Constructed update: %s", constructed_update) - if self._process_listener(constructed_update): - return - self._invoke_callbacks(handler_type, constructed_update) + update_type = self._handlers_to_updates[handler_type] + constructed_updates = ( + cast(type[MessageStatus], update_type).from_updates( + client=self, update=raw_update + ) + if handler_type is handlers.MessageStatusHandler + else (update_type.from_update(client=self, update=raw_update),) + ) + for constructed_update in constructed_updates: + if log.isEnabledFor(logging.DEBUG): + log.debug("Constructed update: %s", constructed_update) + if self._process_listener(constructed_update): + continue + self._invoke_callbacks(handler_type, constructed_update) except Exception: log.exception("Failed to construct update (field=%s)", raw_update.field) finally: @@ -473,7 +485,7 @@ def _process_listener(self: "WhatsApp", update: BaseUpdate) -> bool: ) for identifier in listener_identifiers: listener = self._listeners.get(identifier) - if listener is not None: + if listener is not None and not listener.is_set(): log.debug("Found matching listener") break else: diff --git a/pywa/types/message_status.py b/pywa/types/message_status.py index bb657141..b10c3199 100644 --- a/pywa/types/message_status.py +++ b/pywa/types/message_status.py @@ -188,9 +188,8 @@ def from_update( contact_idx: int = 0, status_idx: int = 0, ) -> MessageStatus: - status = (value := (entry := update["entry"][0])["changes"][0]["value"])[ - "statuses" - ][status_idx] + value = (entry := update["entry"][0])["changes"][0]["value"] + status = value["statuses"][status_idx] error = value.get("errors", status.get("errors", (None,)))[0] return cls( _client=client, @@ -213,6 +212,32 @@ def from_update( error=WhatsAppError.from_dict(error=error) if error else None, ) + @classmethod + def from_updates( + cls, client: WhatsApp, update: RawUpdate + ) -> tuple[MessageStatus, ...]: + value = update["entry"][0]["changes"][0]["value"] + statuses = value["statuses"] + contacts = value["contacts"] + contact_indexes = { + identifier: index + for index, contact in enumerate(contacts) + for identifier in (contact.get("user_id"), contact.get("wa_id")) + if identifier is not None + } + return tuple( + cls.from_update( + client=client, + update=update, + contact_idx=contact_indexes.get( + status.get("recipient_participant_user_id", status["recipient_id"]), + min(status_idx, len(contacts) - 1), + ), + status_idx=status_idx, + ) + for status_idx, status in enumerate(statuses) + ) + @dataclasses.dataclass(frozen=True, slots=True) class Conversation: diff --git a/pywa_async/server.py b/pywa_async/server.py index e7c6113f..8bfde5f9 100644 --- a/pywa_async/server.py +++ b/pywa_async/server.py @@ -3,7 +3,7 @@ import logging import time import warnings -from typing import TYPE_CHECKING +from typing import TYPE_CHECKING, cast from pywa._logging import bind_update_logger, get_update_hash from pywa.server import _logger, _update_hash_of @@ -13,10 +13,12 @@ from .errors import PywaDeprecationWarning from .handlers import ( Handler, + MessageStatusHandler, RawUpdateHandler, ) from .types import ( ContinueHandling, + MessageStatus, RawUpdate, StopHandling, ) @@ -144,14 +146,20 @@ async def _call_handlers(self: "WhatsApp", raw_update: RawUpdate) -> None: return log.debug("Dispatched to %s", handler_type.__name__) try: - constructed_update: BaseUpdate = self._handlers_to_updates[ - handler_type - ].from_update(client=self, update=raw_update) - if log.isEnabledFor(logging.DEBUG): - log.debug("Constructed update: %s", constructed_update) - if await self._process_listener(constructed_update): - return - await self._invoke_callbacks(handler_type, constructed_update) + update_type = self._handlers_to_updates[handler_type] + constructed_updates = ( + cast(type[MessageStatus], update_type).from_updates( + client=self, update=raw_update + ) + if handler_type is MessageStatusHandler + else (update_type.from_update(client=self, update=raw_update),) + ) + for constructed_update in constructed_updates: + if log.isEnabledFor(logging.DEBUG): + log.debug("Constructed update: %s", constructed_update) + if await self._process_listener(constructed_update): + continue + await self._invoke_callbacks(handler_type, constructed_update) except Exception: log.exception("Failed to construct update (field=%s)", raw_update.field) finally: @@ -219,7 +227,7 @@ async def _process_listener(self: "WhatsApp", update: BaseUpdate) -> bool: ) for identifier in listener_identifiers: listener = self._listeners.get(identifier) - if listener is not None: + if listener is not None and not listener.is_set(): log.info("Found matching listener") break else: diff --git a/tests/test_async.py b/tests/test_async.py index a39f87f5..49fd0396 100644 --- a/tests/test_async.py +++ b/tests/test_async.py @@ -1,3 +1,4 @@ +import copy import inspect import json import logging @@ -440,6 +441,7 @@ def test_all_methods_are_overwritten_in_async(overrides): "_msg_cls", "_group_participant_cls", "_msg_status_cls", + "from_updates", "_httpx_client", "_flow_req_cls", "_api_fields", @@ -638,6 +640,9 @@ def test_import_pywa_async_succeeds(): _ASYNC_MESSAGE_UPDATE = json.loads( pathlib.Path("tests/data/updates/message.json").read_text() )["text"] +_ASYNC_MESSAGE_STATUS_UPDATE = json.loads( + pathlib.Path("tests/data/updates/message_status.json").read_text() +)["sent"] _ASYNC_PHONE_NUMBER = "972987654321" _ASYNC_MESSAGE_TEXT = "Body Text" @@ -660,3 +665,33 @@ async def test_pii_present_at_debug_async(caplog): combined = "\n".join(r.getMessage() for r in caplog.records) assert _ASYNC_PHONE_NUMBER in combined assert _ASYNC_MESSAGE_TEXT in combined + + +@pytest.mark.asyncio +async def test_batched_message_statuses_are_all_dispatched_async(): + wa = WhatsAppAsync(phone_id="1122334455667", token="xyz", filter_updates=False) + seen = [] + wa.on_message_status( + lambda _, status: seen.append((status.id, status.from_user.wa_id)) + ) + + update = copy.deepcopy(_ASYNC_MESSAGE_STATUS_UPDATE) + value = update["entry"][0]["changes"][0]["value"] + statuses = value["statuses"] + for index, status_id in enumerate(("wamid.second", "wamid.third"), start=2): + status = copy.deepcopy(statuses[0]) + status["id"] = status_id + status["recipient_id"] = f"97298765432{index}" + statuses.append(status) + contact = copy.deepcopy(value["contacts"][0]) + contact["wa_id"] = status["recipient_id"] + value["contacts"].append(contact) + value["contacts"].reverse() + + await wa.webhook_update_handler(json.dumps(update).encode()) + + assert seen == [ + ("wamid.xyzxyz", "972987654321"), + ("wamid.second", "972987654322"), + ("wamid.third", "972987654323"), + ] diff --git a/tests/test_listeners.py b/tests/test_listeners.py index ea23638c..d9d99703 100644 --- a/tests/test_listeners.py +++ b/tests/test_listeners.py @@ -7,12 +7,14 @@ from pywa import WhatsApp as WhatsAppSync from pywa import filters from pywa.listeners import ( + Listener, ListenerCanceled, ListenerStopped, ListenerTimeout, UserUpdateListenerIdentifier, ) from pywa_async import WhatsApp as WhatsAppAsync +from pywa_async.listeners import Listener as ListenerAsync class DummyUpdate: @@ -47,6 +49,17 @@ def emit_update(): assert isinstance(result, DummyUpdate) +def test_completed_listener_ignores_later_update_sync(wa_sync: WhatsAppSync): + identifier = next(DummyUpdate().listener_identifiers) + listener = Listener(filters=filters.true, cancelers=filters.false) + wa_sync._listeners[identifier] = listener + first = DummyUpdate() + + assert wa_sync._process_listener(first) is True + assert wa_sync._process_listener(DummyUpdate()) is False + assert listener.result is first + + @pytest.mark.asyncio async def test_listener_success_async(wa_async: WhatsAppAsync): identifiers = DummyUpdate().listener_identifiers @@ -64,6 +77,23 @@ async def emit_update(): assert isinstance(result, DummyUpdate) +@pytest.mark.asyncio +async def test_completed_listener_ignores_later_update_async(wa_async: WhatsAppAsync): + identifier = next(DummyUpdate().listener_identifiers) + listener = ListenerAsync( + wa=wa_async, + identifier=identifier, + filters=filters.true, + cancelers=filters.false, + ) + wa_async._listeners[identifier] = listener + first = DummyUpdate() + + assert await wa_async._process_listener(first) is True + assert await wa_async._process_listener(DummyUpdate()) is False + assert listener.future.result() is first + + def test_listener_timeout_sync(wa_sync: WhatsAppSync): identifiers = DummyUpdate().listener_identifiers first_id = next(identifiers) diff --git a/tests/test_server.py b/tests/test_server.py index bc400f9f..1623a976 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -1,3 +1,4 @@ +import copy import io import json import logging @@ -279,6 +280,9 @@ def test_setup_console_logging_pins_cli_banner_to_info(clean_logging): _MESSAGE_UPDATE = json.loads( pathlib.Path("tests/data/updates/message.json").read_text() )["text"] +_MESSAGE_STATUS_UPDATE = json.loads( + pathlib.Path("tests/data/updates/message_status.json").read_text() +)["sent"] _PHONE_NUMBER = "972987654321" _MESSAGE_TEXT = "Body Text" _CONTACT_NAME = "Test Name" @@ -288,6 +292,35 @@ def _make_client() -> WhatsApp: return WhatsApp(phone_id="1122334455667", token="xyz", filter_updates=False) +def test_batched_message_statuses_are_all_dispatched(): + wa = _make_client() + seen = [] + wa.on_message_status( + lambda _, status: seen.append((status.id, status.from_user.wa_id)) + ) + + update = copy.deepcopy(_MESSAGE_STATUS_UPDATE) + value = update["entry"][0]["changes"][0]["value"] + statuses = value["statuses"] + for index, status_id in enumerate(("wamid.second", "wamid.third"), start=2): + status = copy.deepcopy(statuses[0]) + status["id"] = status_id + status["recipient_id"] = f"97298765432{index}" + statuses.append(status) + contact = copy.deepcopy(value["contacts"][0]) + contact["wa_id"] = status["recipient_id"] + value["contacts"].append(contact) + value["contacts"].reverse() + + wa.webhook_update_handler(json.dumps(update).encode()) + + assert seen == [ + ("wamid.xyzxyz", "972987654321"), + ("wamid.second", "972987654322"), + ("wamid.third", "972987654323"), + ] + + def test_pii_absent_at_info(caplog): wa = _make_client() caplog.set_level(logging.INFO, logger="pywa")