Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
34 changes: 23 additions & 11 deletions pywa/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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,
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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:
Expand Down
31 changes: 28 additions & 3 deletions pywa/types/message_status.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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:
Expand Down
28 changes: 18 additions & 10 deletions pywa_async/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -13,10 +13,12 @@
from .errors import PywaDeprecationWarning
from .handlers import (
Handler,
MessageStatusHandler,
RawUpdateHandler,
)
from .types import (
ContinueHandling,
MessageStatus,
RawUpdate,
StopHandling,
)
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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:
Expand Down
35 changes: 35 additions & 0 deletions tests/test_async.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import copy
import inspect
import json
import logging
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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"

Expand All @@ -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"),
]
30 changes: 30 additions & 0 deletions tests/test_listeners.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand Down
33 changes: 33 additions & 0 deletions tests/test_server.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import copy
import io
import json
import logging
Expand Down Expand Up @@ -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"
Expand All @@ -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")
Expand Down