diff --git a/CLAUDE.md b/CLAUDE.md index 0b56534..009efe0 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -45,6 +45,7 @@ src/fopost/ _http.py HttpClient: headers, retry loop, decode, unwrap() errors.py FopostError + subclasses + error_from_response() models.py pydantic models, PLATFORMS, POST_STATUSES, Page/PageMeta + chat_adapter.py ChatAdapter — the inbox as a send/receive interface, imported on its own resources/ _base.py (Resource, parse_list, UNSET, drop_unset) posts.py accounts.py account_groups.py workspaces.py labels.py ai.py inbox.py contacts.py broadcasts.py ads.py @@ -55,6 +56,21 @@ Request flow: a resource method builds a snake_case body/params dict, calls `HttpClient._decode()` (raises or returns the parsed body) → back in the resource, `unwrap(body)` peels `{"data": ...}` and `Model.model_validate(...)` types it. +**`chat_adapter` is a consumer of the client, not a resource.** `ChatAdapter` holds a +`Fopost` instance and calls `client.inbox.*`; it never touches `HttpClient` and never adds +an endpoint of its own. An endpoint it needs goes into `InboxResource` first. It is not +re-exported from `fopost/__init__.py`: `from fopost.chat_adapter import ChatAdapter` is the +import, matching `@fopost/sdk/chat-adapter` in the TypeScript SDK, and the two surfaces +move together. + +**`inbox.message_received` carries ids only.** The payload is `itemId`, `type`, `platform`, +`accountId`, `receivedAt`, with no text and no author. There is no `GET /v1/inbox/:id` and +no `id` filter on the list, so `receive_one` scans `lookback_pages` of that account's items +and answers `None` when it does not find one. If the API grows a single-item read, switch +to it. Webhook verification prefers `X-FoPost-Signature-256` (HMAC over +`{timestamp}.{body}`, refused past a 300s tolerance) and falls back to the body-only +`X-FoPost-Signature`; both compare with `hmac.compare_digest`. + **The envelope unwrap lives in the resources, not the transport.** `_http.unwrap()` is a free function the resources call; `HttpClient` hands back the whole decoded body so the escape hatch and paginated readers can see `meta`. diff --git a/README.md b/README.md index 440a7cf..3dd1d82 100644 --- a/README.md +++ b/README.md @@ -225,6 +225,73 @@ authenticated call and hands back the decoded body: client.request("GET", "/analytics/summary", params={"workspace_id": workspace.id}) ``` +## Chat adapter + +`fopost.chat_adapter` wraps the inbox conversation and reply endpoints in a send/receive +interface, so a chatbot framework can treat FoPost as one channel across every network +that carries direct messages. + +```python +from fopost import Fopost +from fopost.chat_adapter import ChatAdapter + +chat = ChatAdapter( + Fopost(api_key=os.environ["FOPOST_API_KEY"]), + workspace_id=workspace_id, + webhook_secret=os.environ["FOPOST_WEBHOOK_SECRET"], +) +``` + +**Inbound** is the `inbox.message_received` webhook. Subscribe an endpoint to it in FoPost, +then hand the raw body and the request headers to `parse_webhook`. It verifies the +signature, refuses a replay, and returns the event; the event carries ids only, so +`receive_one` reads the text back: + +```python +@app.post("/webhooks/fopost") +async def inbound(request: Request) -> Response: + event = chat.parse_webhook(await request.body(), request.headers) + message = chat.receive_one(event) + if message is None: + return Response(status_code=204) + + chat.typing(message.conversation_id, message.account_id) + chat.send(reply_to=message.id, text=your_bot(message.text)) + chat.mark_read(message) + return Response(status_code=204) +``` + +No webhook? `receive()` polls the same thing: + +```python +for message in chat.receive(): + chat.send(reply_to=message.id, text=your_bot(message.text)) + chat.mark_read(message) +``` + +**Outbound** takes one of three shapes. Reply to a message, reply into a thread, or open +one by handle: + +```python +chat.send(reply_to=message.id, text="On it.") +chat.send(conversation_id="conv_...", text="Still here.") +chat.send(account_id="acc_...", handle="samrivera", text="Following up.") +``` + +| Method | What it does | +| ------ | ------------ | +| `parse_webhook(body, headers)` | Verifies a delivery and returns the `ChatEvent` | +| `verify_webhook(body, headers)` | Signature check on its own; raises on a forged or stale delivery | +| `receive(...)` | Inbound DMs, unread by default | +| `receive_one(event_or_id)` | The full message behind an event id, or `None` | +| `send(...)` | Reply, reply into a thread, or open one | +| `typing(conversation_id, account_id, on=True)` | Typing indicator | +| `mark_read(message)` | Marks the message read | + +Sending needs the `publish` scope on top of `inbox`. Every failure is a `ChatAdapterError` +with a `code` (`invalid_signature`, `stale_delivery`, `unexpected_event`, +`unsupported_target`, and the rest) or the usual `FopostError` from the API. + ## Example [`examples/create_post.py`](examples/create_post.py) creates a post against a diff --git a/src/fopost/chat_adapter.py b/src/fopost/chat_adapter.py new file mode 100644 index 0000000..b621153 --- /dev/null +++ b/src/fopost/chat_adapter.py @@ -0,0 +1,355 @@ +"""``fopost.chat_adapter``: the FoPost social inbox as a send/receive interface. + +Wraps the inbox conversation and reply endpoints so a chatbot framework can treat +FoPost as one channel across every network that carries direct messages:: + + from fopost import Fopost + from fopost.chat_adapter import ChatAdapter + + chat = ChatAdapter(Fopost(), workspace_id=..., webhook_secret=...) + + event = chat.parse_webhook(raw_body, request.headers) # inbound + message = chat.receive_one(event) + if message: + chat.send(reply_to=message.id, text="On it.") # outbound +""" + +from __future__ import annotations + +import builtins +import hashlib +import hmac +import json +import time +from collections.abc import Mapping, Sequence +from dataclasses import dataclass, field +from typing import Any + +from .client import Fopost +from .models import InboxAttachment, InboxItem + +__all__ = [ + "INBOUND_EVENT", + "SIGNATURE_TOLERANCE_SECONDS", + "ChatAdapter", + "ChatAdapterError", + "ChatAuthor", + "ChatEvent", + "ChatMessage", + "to_chat_message", +] + +#: The webhook event a new direct message raises. +INBOUND_EVENT = "inbox.message_received" + +#: Refuse a delivery signed longer ago than this. Matches the API's own tolerance. +SIGNATURE_TOLERANCE_SECONDS = 300 + +_SIGNATURE_HEADER = "x-fopost-signature" +_TIMESTAMPED_SIGNATURE_HEADER = "x-fopost-signature-256" +_TIMESTAMP_HEADER = "x-fopost-timestamp" +_EVENT_HEADER = "x-fopost-event" + + +class ChatAdapterError(Exception): + """A delivery that could not be trusted, or a target that cannot be reached. + + ``code`` is one of ``invalid_body``, ``unexpected_event``, ``missing_secret``, + ``invalid_signature``, ``stale_delivery`` or ``unsupported_target``. + """ + + def __init__(self, code: str, message: str) -> None: + super().__init__(message) + self.code = code + self.message = message + + +@dataclass(frozen=True) +class ChatAuthor: + name: str | None = None + handle: str | None = None + avatar_url: str | None = None + + +@dataclass(frozen=True) +class ChatMessage: + """One inbound or outbound message, flattened out of an inbox item.""" + + id: str + conversation_id: str | None + account_id: str | None + platform: str + type: str + direction: str + text: str + author: ChatAuthor + received_at: str | None + can_reply: bool + attachments: builtins.list[InboxAttachment] = field(default_factory=list) + #: The untouched inbox item, for anything this shape drops. + raw: InboxItem | None = None + + +@dataclass(frozen=True) +class ChatEvent: + """A verified webhook delivery. Ids only: the text lives behind the API.""" + + event: str + item_id: str + type: str + platform: str + account_id: str + received_at: str | None = None + #: When the API built the envelope, not when it was signed. + timestamp: str | None = None + + +def to_chat_message(item: InboxItem) -> ChatMessage: + received = item.platform_created_at or item.created_at + return ChatMessage( + id=item.id, + conversation_id=item.conversation_id, + account_id=item.account.id if item.account else None, + platform=item.platform, + type=item.type, + direction=item.direction or "inbound", + text=item.text or "", + author=ChatAuthor( + name=item.author_name, + handle=item.author_handle, + avatar_url=item.author_avatar_url, + ), + received_at=received.isoformat() if received is not None else None, + can_reply=bool(item.can_reply), + attachments=list(item.attachments), + raw=item, + ) + + +def _header(headers: Mapping[str, Any], name: str) -> str | None: + getter = getattr(headers, "get", None) + if getter is not None: + value = getter(name) + if value is None: + value = getter(name.title()) + if value is not None: + return value[0] if isinstance(value, (list, tuple)) else str(value) + for key, value in dict(headers).items(): + if key.lower() == name: + return value[0] if isinstance(value, (list, tuple)) else str(value) + return None + + +def _digest(value: str | None) -> str | None: + if value is None: + return None + return value[7:] if value.startswith("sha256=") else value + + +class ChatAdapter: + """Send and receive direct messages through one FoPost workspace.""" + + def __init__( + self, + client: Fopost, + *, + workspace_id: str | None = None, + webhook_secret: str | None = None, + page_size: int = 25, + lookback_pages: int = 4, + ) -> None: + self._client = client + self._workspace_id = workspace_id + self._webhook_secret = webhook_secret + self._page_size = page_size + self._lookback_pages = lookback_pages + + # ── Inbound ── + + def parse_webhook(self, body: str | bytes, headers: Mapping[str, Any]) -> ChatEvent: + """Verify the delivery and return the event. + + Prefers the replay-safe ``X-FoPost-Signature-256`` and falls back to + ``X-FoPost-Signature``. Pass the raw body, never a re-serialized object. + """ + raw = body.decode() if isinstance(body, bytes) else body + self.verify_webhook(raw, headers) + + try: + envelope = json.loads(raw) + except ValueError as err: + raise ChatAdapterError("invalid_body", "Webhook body is not JSON") from err + if not isinstance(envelope, dict): + raise ChatAdapterError("invalid_body", "Webhook body is not an object") + + event = envelope.get("event") or _header(headers, _EVENT_HEADER) + if event != INBOUND_EVENT: + raise ChatAdapterError( + "unexpected_event", f"Expected {INBOUND_EVENT}, got {event or 'nothing'}" + ) + + data = envelope.get("data") or {} + item_id, account_id = data.get("itemId"), data.get("accountId") + if not isinstance(item_id, str) or not isinstance(account_id, str): + raise ChatAdapterError("invalid_body", "Webhook payload carries no item id") + + return ChatEvent( + event=event, + item_id=item_id, + type=data.get("type") or "dm", + platform=data.get("platform") or "", + account_id=account_id, + received_at=data.get("receivedAt"), + timestamp=envelope.get("timestamp"), + ) + + def verify_webhook(self, body: str | bytes, headers: Mapping[str, Any]) -> None: + """Raise ``ChatAdapterError`` when the delivery is unsigned, forged or stale.""" + if not self._webhook_secret: + raise ChatAdapterError("missing_secret", "Pass webhook_secret to verify a delivery") + + raw = body.encode() if isinstance(body, str) else body + secret = self._webhook_secret.encode() + + timestamped = _digest(_header(headers, _TIMESTAMPED_SIGNATURE_HEADER)) + if timestamped is not None: + try: + sent_at = int(_header(headers, _TIMESTAMP_HEADER) or "") + except ValueError as err: + raise ChatAdapterError( + "invalid_signature", "Signed delivery carries no timestamp" + ) from err + age = abs(time.time() - sent_at) + if age > SIGNATURE_TOLERANCE_SECONDS: + raise ChatAdapterError("stale_delivery", f"Delivery is {int(age)}s old") + expected = hmac.new(secret, f"{sent_at}.".encode() + raw, hashlib.sha256).hexdigest() + if not hmac.compare_digest(timestamped, expected): + raise ChatAdapterError("invalid_signature", "Signature does not match the body") + return + + plain = _digest(_header(headers, _SIGNATURE_HEADER)) + if plain is None: + raise ChatAdapterError("invalid_signature", "Delivery carries no signature header") + if not hmac.compare_digest(plain, hmac.new(secret, raw, hashlib.sha256).hexdigest()): + raise ChatAdapterError("invalid_signature", "Signature does not match the body") + + def receive( + self, + *, + account_id: str | None = None, + conversation_id: str | None = None, + platform: str | None = None, + state: str | None = "unread", + limit: int | None = None, + ) -> builtins.list[ChatMessage]: + """Inbound direct messages, newest first. Unread by default. + + Pass ``state=None`` to read every state. + """ + page = self._client.inbox.list( + workspace_id=self._workspace_id, + type="dm", + direction="inbound", + state=state, + account_id=account_id, + conversation_id=conversation_id, + platform=platform, + per_page=limit or self._page_size, + ) + return [to_chat_message(item) for item in page.items] + + def receive_one(self, ref: str | ChatEvent) -> ChatMessage | None: + """The message behind a webhook event, or an item id. + + The event carries ids only, so this reads the text back through the inbox. + It scans ``lookback_pages`` of that account's items and answers ``None`` + when the item has aged past them or was deleted. + """ + item_id = ref if isinstance(ref, str) else ref.item_id + narrow: dict[str, Any] = ( + {} if isinstance(ref, str) else {"account_id": ref.account_id, "type": ref.type} + ) + + for page_number in range(1, self._lookback_pages + 1): + page = self._client.inbox.list( + workspace_id=self._workspace_id, + page=page_number, + per_page=self._page_size, + **narrow, + ) + for item in page.items: + if item.id == item_id: + return to_chat_message(item) + if len(page.items) < self._page_size: + return None + return None + + # ── Outbound ── + + def send( + self, + *, + reply_to: str | None = None, + conversation_id: str | None = None, + account_id: str | None = None, + handle: str | None = None, + text: str | None = None, + media_ids: Sequence[str] | None = None, + quick_replies: Sequence[str] | None = None, + ) -> ChatMessage: + """Send on the platform as the connected account. Needs the ``publish`` scope. + + One of three shapes: ``reply_to`` a message, ``conversation_id`` to reply + into a thread, or ``account_id`` plus ``handle`` to open one. + """ + media = list(media_ids) if media_ids is not None else None + quick = list(quick_replies) if quick_replies is not None else None + + if reply_to is not None: + result = self._client.inbox.reply(reply_to, text, media_ids=media, quick_replies=quick) + return to_chat_message(result.item) + + if conversation_id is not None: + latest = self._latest_in(conversation_id) + if latest is None: + raise ChatAdapterError( + "unsupported_target", + f"Conversation {conversation_id} has no message to reply to", + ) + result = self._client.inbox.reply(latest.id, text, media_ids=media, quick_replies=quick) + return to_chat_message(result.item) + + if account_id is None or handle is None or text is None: + raise ChatAdapterError( + "unsupported_target", + "Pass reply_to, conversation_id, or account_id with handle and text", + ) + + started = self._client.inbox.start_conversation( + text=text, account_id=account_id, handle=handle, media_ids=media + ) + if started.item is None: + raise ChatAdapterError( + "unsupported_target", f"{handle} accepted the message but returned no item" + ) + return to_chat_message(started.item) + + def typing(self, conversation_id: str, account_id: str, on: bool = True) -> None: + """Show (default) or clear the typing indicator. Needs the ``publish`` scope.""" + self._client.inbox.set_typing(conversation_id, account_id=account_id, on=on) + + def mark_read(self, message: ChatMessage | str) -> None: + """Mark one message read, so ``receive()`` stops returning it.""" + item_id = message if isinstance(message, str) else message.id + self._client.inbox.update(item_id, state="read") + + # ── Internals ── + + def _latest_in(self, conversation_id: str) -> InboxItem | None: + page = self._client.inbox.list( + workspace_id=self._workspace_id, + conversation_id=conversation_id, + sort="newest", + per_page=1, + ) + return page.items[0] if page.items else None diff --git a/tests/test_chat_adapter.py b/tests/test_chat_adapter.py new file mode 100644 index 0000000..a373b25 --- /dev/null +++ b/tests/test_chat_adapter.py @@ -0,0 +1,244 @@ +from __future__ import annotations + +import hashlib +import hmac +import json +import time +from typing import Any + +import httpx +import pytest +import respx + +from fopost import Fopost +from fopost.chat_adapter import ( + INBOUND_EVENT, + ChatAdapter, + ChatAdapterError, +) +from tests.conftest import BASE_URL + +SECRET = "whsec_test" + +DM_FIXTURE: dict[str, Any] = { + "id": "itm_1", + "workspaceId": "ws_1", + "platform": "instagram", + "type": "dm", + "state": "unread", + "direction": "inbound", + "conversationId": "conv_1", + "authorName": "Sam Rivera", + "authorHandle": "samrivera", + "text": "Do you ship to Portugal?", + "attachments": [], + "platformCreatedAt": "2026-09-20T09:00:00.000Z", + "createdAt": "2026-09-20T09:00:01.000Z", + "canReply": True, + "account": {"id": "acc_1", "platform": "instagram", "username": "yourbrand"}, +} + + +def item(**over: Any) -> dict[str, Any]: + return {**DM_FIXTURE, **over} + + +def adapter(client: Fopost) -> ChatAdapter: + return ChatAdapter(client, workspace_id="ws_1", webhook_secret=SECRET) + + +def delivery(payload: dict[str, Any]) -> tuple[str, dict[str, str]]: + body = json.dumps( + {"event": INBOUND_EVENT, "data": payload, "timestamp": "2026-09-20T09:00:02.000Z"} + ) + sent_at = int(time.time()) + sign = lambda message: hmac.new( # noqa: E731 + SECRET.encode(), message.encode(), hashlib.sha256 + ).hexdigest() + return body, { + "X-FoPost-Event": INBOUND_EVENT, + "X-FoPost-Timestamp": str(sent_at), + "X-FoPost-Signature": f"sha256={sign(body)}", + "X-FoPost-Signature-256": f"sha256={sign(f'{sent_at}.{body}')}", + } + + +@respx.mock +def test_round_trip_from_delivery_to_reply_to_read(client: Fopost) -> None: + listing = respx.get(f"{BASE_URL}/inbox").mock( + return_value=httpx.Response( + 200, json={"data": [item()], "meta": {"page": 1, "perPage": 25, "total": 1}} + ) + ) + reply = respx.post(f"{BASE_URL}/inbox/itm_1/reply").mock( + return_value=httpx.Response( + 200, + json={ + "data": { + "item": item( + id="itm_2", + direction="outbound", + text="We do, in three to five days.", + state="read", + ), + "reply": {"externalId": "ig_9", "externalUrl": None}, + } + }, + ) + ) + patch = respx.patch(f"{BASE_URL}/inbox/itm_1").mock( + return_value=httpx.Response(200, json={"data": item(state="read")}) + ) + + chat = adapter(client) + body, headers = delivery( + { + "itemId": "itm_1", + "type": "dm", + "platform": "instagram", + "accountId": "acc_1", + "receivedAt": "2026-09-20T09:00:00.000Z", + } + ) + + event = chat.parse_webhook(body, headers) + assert (event.event, event.item_id, event.account_id) == (INBOUND_EVENT, "itm_1", "acc_1") + + message = chat.receive_one(event) + assert message is not None + assert message.text == "Do you ship to Portugal?" + assert message.conversation_id == "conv_1" + assert message.account_id == "acc_1" + assert message.author.handle == "samrivera" + + sent = chat.send(reply_to=message.id, text="We do, in three to five days.") + assert sent.id == "itm_2" + assert sent.direction == "outbound" + + chat.mark_read(message) + + assert dict(listing.calls.last.request.url.params) == { + "workspace_id": "ws_1", + "account_id": "acc_1", + "type": "dm", + "page": "1", + "per_page": "25", + } + assert json.loads(reply.calls.last.request.content) == {"text": "We do, in three to five days."} + assert json.loads(patch.calls.last.request.content) == {"state": "read"} + + +@respx.mock +def test_receive_reads_unread_inbound_dms(client: Fopost) -> None: + route = respx.get(f"{BASE_URL}/inbox").mock( + return_value=httpx.Response( + 200, json={"data": [item()], "meta": {"page": 1, "perPage": 25, "total": 1}} + ) + ) + messages = adapter(client).receive() + assert len(messages) == 1 + params = dict(route.calls.last.request.url.params) + assert params["type"] == "dm" + assert params["direction"] == "inbound" + assert params["state"] == "unread" + + +@respx.mock +def test_send_into_a_conversation_replies_to_its_newest_message(client: Fopost) -> None: + listing = respx.get(f"{BASE_URL}/inbox").mock( + return_value=httpx.Response( + 200, json={"data": [item()], "meta": {"page": 1, "perPage": 1, "total": 1}} + ) + ) + respx.post(f"{BASE_URL}/inbox/itm_1/reply").mock( + return_value=httpx.Response( + 200, + json={ + "data": { + "item": item(id="itm_3", direction="outbound"), + "reply": {"externalId": None, "externalUrl": None}, + } + }, + ) + ) + sent = adapter(client).send(conversation_id="conv_1", text="Still here.") + assert sent.id == "itm_3" + params = dict(listing.calls.last.request.url.params) + assert params["conversation_id"] == "conv_1" + assert params["sort"] == "newest" + + +@respx.mock +def test_send_by_handle_opens_a_conversation(client: Fopost) -> None: + route = respx.post(f"{BASE_URL}/inbox/conversations").mock( + return_value=httpx.Response( + 200, + json={ + "data": { + "conversationId": "conv_9", + "item": item(id="itm_9", direction="outbound"), + } + }, + ) + ) + sent = adapter(client).send(account_id="acc_1", handle="samrivera", text="Following up.") + assert sent.id == "itm_9" + assert json.loads(route.calls.last.request.content) == { + "text": "Following up.", + "account_id": "acc_1", + "handle": "samrivera", + } + + +@respx.mock +def test_receive_one_is_none_past_the_lookback(client: Fopost) -> None: + respx.get(f"{BASE_URL}/inbox").mock( + return_value=httpx.Response( + 200, json={"data": [item(id="other")], "meta": {"page": 1, "perPage": 25, "total": 1}} + ) + ) + assert adapter(client).receive_one("itm_gone") is None + + +@respx.mock +def test_a_forged_body_is_refused(client: Fopost) -> None: + _, headers = delivery({"itemId": "itm_1", "accountId": "acc_1", "type": "dm"}) + with pytest.raises(ChatAdapterError) as err: + adapter(client).parse_webhook('{"event":"inbox.message_received","data":{}}', headers) + assert err.value.code == "invalid_signature" + + +@respx.mock +def test_a_delivery_outside_the_tolerance_is_refused(client: Fopost) -> None: + body, headers = delivery({"itemId": "itm_1", "accountId": "acc_1", "type": "dm"}) + headers["X-FoPost-Timestamp"] = str(int(time.time()) - 4000) + with pytest.raises(ChatAdapterError) as err: + adapter(client).parse_webhook(body, headers) + assert err.value.code == "stale_delivery" + + +@respx.mock +def test_the_compatibility_signature_still_verifies(client: Fopost) -> None: + body, headers = delivery({"itemId": "itm_1", "accountId": "acc_1", "type": "dm"}) + del headers["X-FoPost-Signature-256"] + assert adapter(client).parse_webhook(body, headers).item_id == "itm_1" + + +@respx.mock +def test_another_event_is_refused(client: Fopost) -> None: + body = json.dumps({"event": "post.published", "data": {}, "timestamp": None}) + sent_at = int(time.time()) + signature = hmac.new(SECRET.encode(), f"{sent_at}.{body}".encode(), hashlib.sha256).hexdigest() + headers = { + "X-FoPost-Timestamp": str(sent_at), + "X-FoPost-Signature-256": f"sha256={signature}", + } + with pytest.raises(ChatAdapterError) as err: + adapter(client).parse_webhook(body, headers) + assert err.value.code == "unexpected_event" + + +def test_verification_needs_a_secret(client: Fopost) -> None: + with pytest.raises(ChatAdapterError) as err: + ChatAdapter(client).parse_webhook("{}", {}) + assert err.value.code == "missing_secret"