diff --git a/app/agent_base.py b/app/agent_base.py index 96c2aea4..c9a4b4f9 100644 --- a/app/agent_base.py +++ b/app/agent_base.py @@ -2446,10 +2446,18 @@ async def _handle_external_event(self, payload: Dict) -> None: integration_type = payload.get("integrationType", "").lower() is_self_message = payload.get("is_self_message", False) + message_id = payload.get("messageId", "") + + # The listener capped a long body (PlatformMessage.truncated): + # mark the cut where both the agent and the chat details see it. + truncated = bool(payload.get("truncated")) and bool(message_body) + if truncated: + message_body = f"{message_body}…" + # Normalized attachments (PlatformMessage.attachments) become # descriptor lines with retrieval hints — appended to the body, # or standing in for it on media-only messages so they are no - # longer dropped (docs/plans/attachment-reception-plan.md). + # longer dropped. from app.integrations import format_attachment_descriptors att_lines = format_attachment_descriptors( @@ -2518,11 +2526,28 @@ async def _handle_external_event(self, payload: Dict) -> None: location_parts.append(f"channel {channel_id}") location_str = f" in {' / '.join(location_parts)}" if location_parts else "" + # Tell the agent the body above is partial and how to get the + # rest; whether the rest is needed stays the agent's call. + truncation_note = "" + if truncated: + fetch_hint = ( + f" If the rest matters, read the full message with the " + f"{integration_type} action that gets a message by ID " + f"(message ID: {message_id})." + if message_id + else "" + ) + truncation_note = ( + f"[This message is truncated: only the beginning is shown " + f"above.{fetch_hint}]\n" + ) + if is_self_message: # Self-message = user is directly talking to the agent via their own platform. event_content = ( f"[USER SELF-MESSAGE via {source}]\n" - f"{message_body}\n\n" + f"{message_body}\n" + f"{truncation_note}\n" f"INSTRUCTIONS: Reply to the message to the user on {source}" f"{account_note}" ) @@ -2538,7 +2563,8 @@ async def _handle_external_event(self, payload: Dict) -> None: f"From: {contact_name} ({contact_id}){location_str}\n" f"Platform: {source}\n" f"{received_on}" - f'Message: "{message_body}"\n\n' + f'Message: "{message_body}"\n' + f"{truncation_note}\n" f"INSTRUCTIONS: Notify the user about this message on their " f"preferred platform (check USER.md 'Preferred Messaging " f"Platform'). If USER.md does not name one, notify via " diff --git a/craftos_integrations/README.md b/craftos_integrations/README.md index 5db195cc..f23de060 100644 --- a/craftos_integrations/README.md +++ b/craftos_integrations/README.md @@ -169,9 +169,18 @@ has `listen` enabled, and reconciles whenever accounts change. "messageId": "", "is_self_message": False, "raw": {...}, # full original platform event + "attachments": [...], # normalized non-text payloads ({kind, id, name, mime, size, url, extra}) + "truncated": False, # True when the listener capped messageBody } ``` +`messageBody` is plain, decoded text built from the real message body — +listeners unescape HTML entities and render platform markup before +emitting, and never forward an API preview field (Gmail `snippet`, +Outlook `bodyPreview`) as the body. A listener that caps the length uses +`helpers.clip` and sets `truncated`; the host then marks the cut and tells +the agent the message is truncated and how to read the rest by message id. + --- ## Configuration: OAuth env vars @@ -399,7 +408,7 @@ For a production-level integration, produce in this order: | 2 | Implement `verify_token` (token auth) or `oauth_spec` (OAuth), plus `identity_of` | `provider.py` | | 3 | Optional: `config_class` + `config_fields` for post-connect knobs | `provider.py` | | 4 | Build the client — one method per endpoint, using `helpers.arequest`, returning `Result` | client in `__init__.py` | -| 5 | Optional: `start_listening` / `stop_listening` (webhook / polling / WebSocket) | client | +| 5 | Optional: `start_listening` / `stop_listening` (webhook / polling / WebSocket). Emit **plain, decoded** text from the real body (unescape entities, render platform markup like `<@U…>`; no API preview fields); cap long bodies with `helpers.clip` and pass its flag as `PlatformMessage.truncated` | client | | 6 | Write `INTEGRATION.md` — identifier shape, silent-drop config flags, auth gotchas | integration root | | 7 | Mirror each client method as a `client_op` with sub-set + umbrella tags | `operations.py` | | 8 | Verify — `scripts/verify_integration.py `, then a live smoke test | see "Verification" | diff --git a/craftos_integrations/base.py b/craftos_integrations/base.py index c6009859..dd907af7 100644 --- a/craftos_integrations/base.py +++ b/craftos_integrations/base.py @@ -43,6 +43,10 @@ class PlatformMessage: # The HOST formats these into descriptor text + retrieval hints — # listeners only normalize (docs/plans/attachment-reception-plan.md). attachments: List[Dict[str, Any]] = field(default_factory=list) + # True when the listener capped ``text`` (``helpers.clip``). The HOST + # marks the cut and gives the agent the message id to fetch the rest; + # listeners only declare the fact. + truncated: bool = False MessageCallback = Callable[[PlatformMessage], Awaitable[None]] diff --git a/craftos_integrations/helpers/__init__.py b/craftos_integrations/helpers/__init__.py index d004bbd4..36e51584 100644 --- a/craftos_integrations/helpers/__init__.py +++ b/craftos_integrations/helpers/__init__.py @@ -10,15 +10,19 @@ ``status_code → {ok, result} | {error, details}`` envelope shape. result: ``Result`` / ``Ok`` / ``Err`` TypedDict aliases for the envelope — use as return annotations for static type-checking benefits. + text: ``clip`` — word-boundary length cap for listener message bodies + that reports whether it cut (→ ``PlatformMessage.truncated``). """ from .http import arequest, request from .result import Err, Ok, Result +from .text import clip __all__ = [ "Err", "Ok", "Result", "arequest", + "clip", "request", ] diff --git a/craftos_integrations/helpers/text.py b/craftos_integrations/helpers/text.py new file mode 100644 index 00000000..93deb701 --- /dev/null +++ b/craftos_integrations/helpers/text.py @@ -0,0 +1,29 @@ +"""Text shaping for listener message bodies. + +Listeners that cap a body's length (to keep inbound token cost bounded) +use ``clip`` so every cap cuts the same way and reports that it cut — +the flag goes to ``PlatformMessage.truncated`` and the host marks the cut. +""" + +from __future__ import annotations + +from typing import Tuple + +# How far back from the limit a word boundary may be before we give up +# and hard-cut (fraction of the limit). +_BOUNDARY_WINDOW = 0.2 + + +def clip(text: str, limit: int) -> Tuple[str, bool]: + """Cut ``text`` to at most ``limit`` chars, preferring a word boundary. + + Returns ``(clipped, was_clipped)``. No ellipsis is added — the host + owns presentation of the cut. + """ + if len(text) <= limit: + return text, False + head = text[:limit] + cut = max(head.rfind(" "), head.rfind("\n")) + if cut >= limit * (1 - _BOUNDARY_WINDOW): + head = head[:cut] + return head.rstrip(), True diff --git a/craftos_integrations/providers/_shared.py b/craftos_integrations/providers/_shared.py index 88bb51cf..a7a06c6f 100644 --- a/craftos_integrations/providers/_shared.py +++ b/craftos_integrations/providers/_shared.py @@ -45,6 +45,7 @@ def platform_message_payload(msg: Any) -> Dict[str, Any]: "is_self_message": raw.get("is_self_message", False), "raw": raw, "attachments": list(getattr(msg, "attachments", None) or []), + "truncated": bool(getattr(msg, "truncated", False)), } diff --git a/craftos_integrations/providers/github/client.py b/craftos_integrations/providers/github/client.py index c1a8f3ca..f96c819f 100644 --- a/craftos_integrations/providers/github/client.py +++ b/craftos_integrations/providers/github/client.py @@ -31,7 +31,7 @@ register_client, save_credential, ) -from ...helpers import Result, arequest +from ...helpers import Result, arequest, clip from ...logger import get_logger logger = get_logger(__name__) @@ -40,6 +40,10 @@ POLL_INTERVAL = 15 RETRY_DELAY = 30 +# Inbound notification comments are previewed at this length (token cost); +# a longer comment is flagged as PlatformMessage.truncated. +_COMMENT_PREVIEW_CHARS = 300 + @dataclass class GitHubCredential: @@ -336,8 +340,10 @@ async def _dispatch_notification(self, notif: Dict[str, Any]) -> None: f"[{repo_full}] {subject_type}: {subject_title}", f"Reason: {reason}", ] + truncated = False if comment_body: - text_parts.append(f"Comment by @{comment_author}: {comment_body[:300]}") + preview, truncated = clip(comment_body, _COMMENT_PREVIEW_CHARS) + text_parts.append(f"Comment by @{comment_author}: {preview}") await self._message_callback( PlatformMessage( @@ -350,6 +356,7 @@ async def _dispatch_notification(self, notif: Dict[str, Any]) -> None: message_id=notif.get("id", ""), timestamp=datetime.now(timezone.utc), raw=notif, + truncated=truncated, ) ) diff --git a/craftos_integrations/providers/gmail/client.py b/craftos_integrations/providers/gmail/client.py index 2384c380..ece31715 100644 --- a/craftos_integrations/providers/gmail/client.py +++ b/craftos_integrations/providers/gmail/client.py @@ -17,6 +17,7 @@ import asyncio import base64 +import html import mimetypes import os import re @@ -27,7 +28,7 @@ from email.mime.multipart import MIMEMultipart from email.mime.text import MIMEText from html.parser import HTMLParser -from typing import Any, Dict, List, Optional +from typing import Any, Dict, List, Optional, Tuple from ... import ( BasePlatformClient, @@ -36,7 +37,7 @@ load_config, register_client, ) -from ...helpers import Result, arequest, request as http_request +from ...helpers import Result, arequest, clip, request as http_request from ...logger import get_logger from .._google_common import ( GoogleApiClientMixin, @@ -49,6 +50,11 @@ POLL_INTERVAL = 5 RETRY_DELAY = 10 +# Incoming mail is forwarded to the agent with its body capped at this many +# chars (token cost); a longer body is flagged as PlatformMessage.truncated +# and the agent can read the rest with get_gmail. +_INBOUND_BODY_CHARS = 2000 + # Headers worth showing the agent. Everything else Gmail returns under # format=full is transport/anti-spam machinery (Received chains, ARC-*, DKIM # signatures, Authentication-Results, SPF, List-Unsubscribe JWTs, Feedback-ID, @@ -110,12 +116,16 @@ def _part_charset(part: Dict[str, Any]) -> str: def _decode_part(part: Dict[str, Any]) -> str: - """Decode a part's base64url body using its declared charset. + """Decode a part's base64url body. + + Gmail hands text bodies back already transcoded to UTF-8 but keeps the + original Content-Type charset (a Windows-1252 mail's ``=B7`` arrives as + ``\\xc2\\xb7``), so decoding by the declared charset turns ``·`` into + ``·``. Bytes that are not valid UTF-8 were passed through untranscoded + and are decoded by the declared charset instead. Never raises: an unknown charset or undecodable byte degrades to - replacement characters instead of failing the whole read (a strict - utf-8 decode used to turn any iso-8859-1 / windows-1252 email into an - error that also lost its headers). + replacement characters instead of failing the whole read. """ data = (part.get("body") or {}).get("data", "") if not data: @@ -124,6 +134,10 @@ def _decode_part(part: Dict[str, Any]) -> str: raw = base64.urlsafe_b64decode(data + "=" * (-len(data) % 4)) except Exception: return "" + try: + return raw.decode("utf-8") + except UnicodeDecodeError: + pass try: return raw.decode(_part_charset(part), errors="replace") except LookupError: @@ -190,6 +204,12 @@ def _clean_spaces(text: str) -> str: return _SPACE_RUN.sub(" ", text).strip() +def clean_snippet(raw: str) -> str: + """Gmail's ``snippet`` arrives HTML-escaped (``'`` ``&``); + decode it and collapse filler whitespace.""" + return _clean_spaces(html.unescape(raw or "")) + + def _html_to_text(html: str) -> str: parser = _HtmlText() try: @@ -201,6 +221,51 @@ def _html_to_text(html: str) -> str: return re.sub(r"\n{3,}", "\n\n", "\n".join(lines)).strip() +def _walk_payload( + payload: Dict[str, Any], +) -> Tuple[Dict[str, str], List[Dict[str, Any]]]: + """One pass over a ``format=full`` payload: the first decoded + text/plain + text/html bodies (by mime), and every attachmentId part.""" + texts: Dict[str, str] = {} + attachments: List[Dict[str, Any]] = [] + + def _visit(part: Dict[str, Any]) -> None: + body = part.get("body", {}) or {} + mime = part.get("mimeType", "") + if body.get("attachmentId"): + attachments.append( + { + "filename": part.get("filename", ""), + "attachment_id": body["attachmentId"], + "mimeType": mime, + "size": body.get("size", 0), + } + ) + elif ( + mime in ("text/plain", "text/html") + and "data" in body + and not part.get("filename") # inline-data text attachment + and mime not in texts + ): + texts[mime] = _decode_part(part) + for nested in part.get("parts", []) or []: + _visit(nested) + + _visit(payload) + return texts, attachments + + +def _readable_body(texts: Dict[str, str]) -> Tuple[str, str]: + """``(body, body_format)`` from ``_walk_payload`` texts. Prefers the + plain-text alternative; HTML-only mail (most newsletters/notifications) + is converted rather than coming back empty.""" + if texts.get("text/plain", "").strip(): + return texts["text/plain"], "text" + if texts.get("text/html"): + return _html_to_text(texts["text/html"]), "html_converted" + return "", "none" + + GMAIL = IntegrationSpec( name="gmail", cred_class=GoogleCredential, @@ -368,13 +433,15 @@ async def _fetch_and_dispatch(self, msg_id: str) -> None: if not cfg.process_incoming: return - # format=full + a fields partial-response mask: returns headers, - # snippet, and ONLY the parts skeleton (filename/mimeType/ - # attachmentId/size — no body data), staying ~1-3KB. Quota cost is + # format=full + a fields partial-response mask: headers plus the + # parts skeleton with inline body data (the text/plain + text/html + # bodies we read; part headers carry their charset). Attachment bytes never come inline — those parts + # carry only an attachmentId — so the response stays the size of + # the text. The body, not `snippet`, is what we forward: the snippet + # is an HTML-escaped ~200-char preview (issue #444). Quota cost is # flat regardless of format. Three explicit nesting levels cover - # mixed / mixed-inside-signed / one spare; a bare `payload/parts` - # selector would pull body.data too — keep the sub-selection. - _part_sel = "partId,mimeType,filename,body(attachmentId,size)" + # mixed / mixed-inside-signed / one spare. + _part_sel = "partId,mimeType,filename,headers,body(attachmentId,size,data)" result = await arequest( "GET", f"{GMAIL_API_BASE}/users/me/messages/{msg_id}", @@ -383,7 +450,8 @@ async def _fetch_and_dispatch(self, msg_id: str) -> None: ("format", "full"), ( "fields", - "id,threadId,snippet,labelIds,historyId,payload(mimeType,headers," + "id,threadId,labelIds,historyId,payload(mimeType,headers," + "body(size,data)," f"parts({_part_sel},parts({_part_sel},parts({_part_sel}))))", ), ], @@ -398,7 +466,8 @@ async def _fetch_and_dispatch(self, msg_id: str) -> None: } from_header = headers.get("From", "") subject = headers.get("Subject", "(no subject)") - snippet = msg.get("snippet", "") + texts, _ = _walk_payload(msg.get("payload", {}) or {}) + body, truncated = clip(_readable_body(texts)[0].strip(), _INBOUND_BODY_CHARS) sender_name = from_header sender_email = from_header @@ -421,7 +490,7 @@ async def _fetch_and_dispatch(self, msg_id: str) -> None: except Exception: pass - text = f"Subject: {subject}\n{snippet}" if snippet else f"Subject: {subject}" + text = f"Subject: {subject}\n{body}" if body else f"Subject: {subject}" # Real attachments carry a non-empty filename + attachmentId # (Gmail's own paperclip heuristic); nameless attachmentId parts @@ -459,6 +528,7 @@ def _collect(parts: Any) -> None: timestamp=timestamp, raw=msg, attachments=attachments, + truncated=truncated, ) ) @@ -551,7 +621,7 @@ def _shape(msg): # always present and not set by the sender. "internalDate": msg.get("internalDate"), "sizeEstimate": msg.get("sizeEstimate"), - "snippet": _clean_spaces(msg.get("snippet", "")), + "snippet": clean_snippet(msg.get("snippet", "")), # The `metadataHeaders` request param is honoured ONLY for # format=metadata — with format=full (full_body=True) Gmail returns # the entire raw MIME header block. That is ~55% of the payload and @@ -562,44 +632,8 @@ def _shape(msg): "headers": _filter_headers(msg.get("payload", {}).get("headers", [])), } if full_body: - attachments: List[Dict[str, Any]] = [] - texts: Dict[str, str] = {} - - def _visit(part): - body = part.get("body", {}) or {} - mime = part.get("mimeType", "") - if body.get("attachmentId"): - attachments.append( - { - "filename": part.get("filename", ""), - "attachment_id": body["attachmentId"], - "mimeType": mime, - "size": body.get("size", 0), - } - ) - elif ( - mime in ("text/plain", "text/html") - and "data" in body - and not part.get("filename") # inline-data text attachment - and mime not in texts - ): - texts[mime] = _decode_part(part) - for nested in part.get("parts", []) or []: - _visit(nested) - - _visit(msg.get("payload", {}) or {}) - - # Prefer the plain-text alternative; HTML-only mail (most - # newsletters/notifications) used to come back with no body. - if texts.get("text/plain", "").strip(): - email_info["body"] = texts["text/plain"] - email_info["body_format"] = "text" - elif texts.get("text/html"): - email_info["body"] = _html_to_text(texts["text/html"]) - email_info["body_format"] = "html_converted" - else: - email_info["body"] = "" - email_info["body_format"] = "none" + texts, attachments = _walk_payload(msg.get("payload", {}) or {}) + email_info["body"], email_info["body_format"] = _readable_body(texts) email_info["attachments"] = attachments return email_info diff --git a/craftos_integrations/providers/gmail/operations.py b/craftos_integrations/providers/gmail/operations.py index 64aa960f..5777df02 100644 --- a/craftos_integrations/providers/gmail/operations.py +++ b/craftos_integrations/providers/gmail/operations.py @@ -19,6 +19,7 @@ from ...contracts import Operation from .._shared import client_op +from .client import clean_snippet def _get_gmail_thread_op() -> Operation: @@ -80,7 +81,7 @@ async def fn(client: Any, input_data: Dict[str, Any]) -> Dict[str, Any]: "internalDate": msg.get("internalDate"), "labelIds": labels, "unread": "UNREAD" in labels, - "snippet": msg.get("snippet", ""), + "snippet": clean_snippet(msg.get("snippet", "")), } ) res = { @@ -141,7 +142,7 @@ async def fn(client: Any, input_data: Dict[str, Any]) -> Dict[str, Any]: "message_id": msg.get("id"), "to": headers.get("To", ""), "subject": headers.get("Subject", ""), - "snippet": msg.get("snippet", ""), + "snippet": clean_snippet(msg.get("snippet", "")), }, } return res diff --git a/craftos_integrations/providers/jira/client.py b/craftos_integrations/providers/jira/client.py index a474ca85..30decbc1 100644 --- a/craftos_integrations/providers/jira/client.py +++ b/craftos_integrations/providers/jira/client.py @@ -21,7 +21,7 @@ save_config, register_client, ) -from ...helpers import Result, arequest +from ...helpers import Result, arequest, clip from ...logger import get_logger logger = get_logger(__name__) @@ -30,6 +30,10 @@ POLL_INTERVAL = 10 RETRY_DELAY = 15 +# The latest comment on an inbound issue update is previewed at this length +# (token cost); a longer comment is flagged as PlatformMessage.truncated. +_COMMENT_PREVIEW_CHARS = 200 + # Fields requested from the API when the caller wants a lean issue payload # (include_metadata=False and no explicit fields list). Restricting server-side # avoids the full customfield_* dump Jira returns by default. @@ -429,12 +433,14 @@ async def _dispatch_issue(self, issue: Dict[str, Any]) -> None: ] if labels: text_parts.append(f"Labels: {', '.join(labels)}") + truncated = False if comments: latest = comments[-1] cb = _extract_adf_text(latest.get("body", {})) ca = (latest.get("author") or {}).get("displayName", "") if cb: - text_parts.append(f"Latest comment by {ca}: {cb[:200]}") + preview, truncated = clip(cb, _COMMENT_PREVIEW_CHARS) + text_parts.append(f"Latest comment by {ca}: {preview}") timestamp = None try: @@ -456,6 +462,7 @@ async def _dispatch_issue(self, issue: Dict[str, Any]) -> None: timestamp=timestamp, raw=issue, attachments=attachments, + truncated=truncated, ) ) diff --git a/craftos_integrations/providers/outlook/client.py b/craftos_integrations/providers/outlook/client.py index aafbb4aa..6224e482 100644 --- a/craftos_integrations/providers/outlook/client.py +++ b/craftos_integrations/providers/outlook/client.py @@ -19,7 +19,7 @@ register_client, save_credential, ) -from ...helpers import Result, arequest, request as http_request +from ...helpers import Result, arequest, clip, request as http_request from ...logger import get_logger logger = get_logger(__name__) @@ -31,6 +31,11 @@ POLL_INTERVAL = 5 RETRY_DELAY = 10 +# Incoming mail is forwarded to the agent with its body capped at this many +# chars (token cost); a longer body is flagged as PlatformMessage.truncated +# and the agent can read the rest with get_outlook_email. +_INBOUND_BODY_CHARS = 2000 + @dataclass class OutlookCredential: @@ -203,15 +208,20 @@ async def _poll_loop(self) -> None: async def _check_new_messages(self) -> None: if not self._last_poll_time: return + # The full body as plain text, not `bodyPreview` (a 255-char + # preview that reads as the whole message — issue #444). result = await arequest( "GET", f"{GRAPH_API_BASE}/me/messages", - headers=self._auth_header(), + headers={ + **self._auth_header(), + "Prefer": 'outlook.body-content-type="text"', + }, params={ "$filter": f"receivedDateTime ge {self._last_poll_time}", "$orderby": "receivedDateTime asc", "$top": "50", - "$select": "id,from,subject,bodyPreview,receivedDateTime,conversationId,hasAttachments", + "$select": "id,from,subject,body,receivedDateTime,conversationId,hasAttachments", }, expected=(200,), ) @@ -248,8 +258,11 @@ async def _dispatch_message(self, msg: Dict[str, Any]) -> None: return subject = msg.get("subject", "(no subject)") - snippet = msg.get("bodyPreview", "") - text = f"Subject: {subject}\n{snippet}" if snippet else f"Subject: {subject}" + body, truncated = clip( + ((msg.get("body") or {}).get("content") or "").strip(), + _INBOUND_BODY_CHARS, + ) + text = f"Subject: {subject}\n{body}" if body else f"Subject: {subject}" timestamp = None try: @@ -294,6 +307,7 @@ async def _dispatch_message(self, msg: Dict[str, Any]) -> None: timestamp=timestamp, raw=msg, attachments=attachments, + truncated=truncated, ) ) diff --git a/craftos_integrations/providers/slack/client.py b/craftos_integrations/providers/slack/client.py index 3f8c587c..d132c5dc 100644 --- a/craftos_integrations/providers/slack/client.py +++ b/craftos_integrations/providers/slack/client.py @@ -20,6 +20,7 @@ ) from ...helpers import arequest, request as http_request from ...logger import get_logger +from .formatting import to_plain_text logger = get_logger(__name__) @@ -121,6 +122,8 @@ def __init__(self): self._bot_user_id: Optional[str] = None self._last_timestamps: Dict[str, str] = {} self._catchup_done: bool = False + # user id → display name; one entry per workspace member seen. + self._user_names: Dict[str, str] = {} def has_credentials(self) -> bool: return has_credential(self.spec.cred_file) @@ -307,16 +310,8 @@ async def _process_message(self, msg: Dict[str, Any], channel_id: str) -> None: if (not text and not attachments) or user_id == self._bot_user_id: return - sender_name = user_id - try: - info = self.get_user_info(user_id) - if info.get("ok"): - profile = info.get("user", {}).get("profile", {}) - sender_name = ( - profile.get("display_name") or profile.get("real_name") or user_id - ) - except Exception: - pass + sender_name = self._display_name(user_id) or user_id + text = to_plain_text(text, self._display_name) ts_float = float(msg.get("ts", "0")) timestamp = ( @@ -338,6 +333,27 @@ async def _process_message(self, msg: Dict[str, Any], channel_id: str) -> None: ) ) + def _display_name(self, user_id: str) -> str: + """User id → display name, "" when unknown. Memoized per client: + the sender lookup and every ``<@U…>`` mention share one cache. + Failures are not cached, so a transient API error retries.""" + if not user_id: + return "" + cached = self._user_names.get(user_id) + if cached: + return cached + name = "" + try: + info = self.get_user_info(user_id) + if info.get("ok"): + profile = info.get("user", {}).get("profile", {}) + name = profile.get("display_name") or profile.get("real_name") or "" + except Exception: + pass + if name: + self._user_names[user_id] = name + return name + # ----- API ----- async def send_message(self, recipient: str, text: str, **kwargs) -> Dict[str, Any]: payload: Dict[str, Any] = {"channel": recipient, "text": text} diff --git a/craftos_integrations/providers/slack/formatting.py b/craftos_integrations/providers/slack/formatting.py new file mode 100644 index 00000000..369d2c89 --- /dev/null +++ b/craftos_integrations/providers/slack/formatting.py @@ -0,0 +1,70 @@ +"""Slack message text → readable plain text. + +Slack's ``text`` field is not what the user typed: ``& < >`` arrive +HTML-escaped and every reference is wrapped in angle-bracket markup. +Forwarded as-is, the agent and the chat details show ``<@U0123ABC>`` and +``&``. + +Pure module — the only outside knowledge (user id → display name) is +injected as ``resolve_user``, so it is testable without a client. +""" + +from __future__ import annotations + +import html +import re +from typing import Callable + +# One `<…>` reference. Slack escapes literal `<` / `>` in user text, so a +# raw angle bracket is always markup. +_REFERENCE = re.compile(r"<([^<>]+)>") + +# / / — label-less special mentions. +_SPECIAL_MENTIONS = {"here", "channel", "everyone"} + + +def to_plain_text(text: str, resolve_user: Callable[[str], str]) -> str: + """Convert Slack references + escapes to plain text. Never raises. + + ============================ ========================================= + ``<@U123>`` / ``<@U123|bob>`` ``@`` (via ``resolve_user``; + falls back to the label, then ``@U123``) + ``<#C123|general>`` ``#general`` (``<#C123>`` → ``#C123``) + ```` ``label (https://x)`` + ```` ``https://x`` (``mailto:`` prefix dropped) + ```` etc. ``@here`` / ``@channel`` / ``@everyone`` + ```` ``@team`` + ```` ``fallback`` + ============================ ========================================= + + References are converted BEFORE unescaping: unescaping first would + turn a user's literal ``<@U1>`` into a fake mention. + """ + if not text: + return "" + return html.unescape( + _REFERENCE.sub(lambda m: _render(m.group(1), resolve_user), text) + ) + + +def _render(ref: str, resolve_user: Callable[[str], str]) -> str: + target, _, label = ref.partition("|") + if target.startswith("@"): + user_id = target[1:] + try: + name = resolve_user(user_id) + except Exception: + name = "" + return f"@{name or label or user_id}" + if target.startswith("#"): + return f"#{label or target[1:]}" + if target.startswith("!"): + keyword = target[1:].split("^", 1)[0] + if keyword in _SPECIAL_MENTIONS: + return f"@{keyword}" + return label or target[1:] + if target.startswith("mailto:"): + target = target[len("mailto:") :] + if label and label != target: + return f"{label} ({target})" + return target diff --git a/craftos_integrations/providers/twitter/client.py b/craftos_integrations/providers/twitter/client.py index ab6a27a6..592025cf 100644 --- a/craftos_integrations/providers/twitter/client.py +++ b/craftos_integrations/providers/twitter/client.py @@ -7,6 +7,7 @@ import base64 import hashlib import hmac +import html import secrets as _secrets import time import urllib.parse @@ -335,7 +336,9 @@ async def _dispatch_mention( att["name"] = media["alt_text"] attachments.append(att) - text = tweet.get("text", "") + # v2 tweet text arrives with & < > HTML-escaped. Decode before the + # watch_tag match too — tags are typed as plain text. + text = html.unescape(tweet.get("text", "")) author_id = tweet.get("author_id", "") author_info = users_map.get(author_id, {}) author_username = author_info.get("username", "") diff --git a/tests/integrations/test_listener_text_fidelity.py b/tests/integrations/test_listener_text_fidelity.py new file mode 100644 index 00000000..f2ed0fc7 --- /dev/null +++ b/tests/integrations/test_listener_text_fidelity.py @@ -0,0 +1,391 @@ +"""Inbound message fidelity (issue #444). + +Listeners emit plain, decoded text built from the real message body, cap +it with ``clip`` and declare the cut as ``PlatformMessage.truncated``; the +host marks the cut for the chat and tells the agent the message is +truncated and how to fetch the rest. + +No pytest-asyncio in this repo — async paths are driven with asyncio.run. +""" + +from __future__ import annotations + +import asyncio +import base64 +import time +from types import SimpleNamespace + +import pytest + +import craftos_integrations.providers.gmail.client as gmail_mod +import craftos_integrations.providers.outlook.client as outlook_mod +from craftos_integrations.base import PlatformMessage +from craftos_integrations.helpers import clip +from craftos_integrations.providers._shared import platform_message_payload +from craftos_integrations.providers.gmail.client import clean_snippet +from craftos_integrations.providers.gmail.provider import GmailProvider +from craftos_integrations.providers.outlook.provider import OutlookProvider +from craftos_integrations.providers.slack.formatting import to_plain_text + + +def _b64(text: str) -> str: + return base64.urlsafe_b64encode(text.encode()).decode() + + +def _collect_callback(): + received = [] + + async def callback(msg): + received.append(msg) + + return received, callback + + +# ── clip ──────────────────────────────────────────────────────────────── + + +def test_clip_leaves_short_text_alone(): + assert clip("short", 10) == ("short", False) + assert clip("exactly10!", 10) == ("exactly10!", False) + + +def test_clip_cuts_on_word_boundary(): + text, cut = clip("the quick brown fox jumps", 18) + assert cut is True + assert text == "the quick brown" + + +def test_clip_hard_cuts_without_nearby_boundary(): + text, cut = clip("a " + "x" * 50, 20) + assert cut is True + assert text == ("a " + "x" * 50)[:20] + + +# ── Gmail ─────────────────────────────────────────────────────────────── + + +def test_gmail_clean_snippet_decodes_entities(): + assert clean_snippet("that's Tom & Jerry") == "that's Tom & Jerry" + + +def _gmail_client(monkeypatch, payload): + message = {"id": "m1", "threadId": "t1", "payload": payload} + + async def fake_arequest(method, url, **kwargs): + assert "/users/me/messages/m1" in url + return {"result": message} + + monkeypatch.setattr(gmail_mod, "arequest", fake_arequest) + monkeypatch.setattr( + gmail_mod, "load_config", lambda *a, **k: gmail_mod.GmailConfig() + ) + client = GmailProvider().build_client( + { + "access_token": "tok", + "refresh_token": "ref", + "token_expiry": time.time() + 3600, + "client_id": "cid", + "client_secret": "cs", + "email": "me@x.com", + }, + lambda d: None, + ) + received, client._message_callback = _collect_callback() + asyncio.run(client._fetch_and_dispatch("m1")) + assert len(received) == 1 + return received[0] + + +_GMAIL_HEADERS = [ + {"name": "From", "value": "Joe "}, + {"name": "Subject", "value": "Top 10 users"}, +] + + +def test_gmail_listener_forwards_real_body_not_snippet(monkeypatch): + # The #444 email: an apostrophe that the escaped snippet leaked as '. + body = "Joe here with an email that's automated. You're probably using Stripe." + msg = _gmail_client( + monkeypatch, + { + "headers": _GMAIL_HEADERS, + "parts": [ + {"mimeType": "text/plain", "body": {"data": _b64(body)}}, + {"mimeType": "text/html", "body": {"data": _b64("

x

")}}, + ], + }, + ) + assert msg.text == f"Subject: Top 10 users\n{body}" + assert msg.truncated is False + + +def test_gmail_listener_converts_html_only_mail(monkeypatch): + msg = _gmail_client( + monkeypatch, + { + "headers": _GMAIL_HEADERS, + "mimeType": "text/html", + "body": {"data": _b64("

Tom & Jerry

")}, + }, + ) + assert msg.text == "Subject: Top 10 users\nTom & Jerry" + + +_CP1252_PLAIN = [ + {"name": "Content-Type", "value": 'text/plain; charset="Windows-1252"'} +] + + +def test_gmail_listener_decodes_transcoded_body_as_utf8(monkeypatch): + # Gmail returns the body as UTF-8 but keeps the sender's Windows-1252 + # header; decoding by the header turned "·" into "·". + data = base64.urlsafe_b64encode("01 · Map".encode("utf-8")).decode() + msg = _gmail_client( + monkeypatch, + { + "headers": _GMAIL_HEADERS, + "parts": [ + { + "mimeType": "text/plain", + "headers": _CP1252_PLAIN, + "body": {"data": data}, + } + ], + }, + ) + assert msg.text == "Subject: Top 10 users\n01 · Map" + + +def test_gmail_listener_decodes_untranscoded_body_by_declared_charset(monkeypatch): + data = base64.urlsafe_b64encode("01 · Map".encode("cp1252")).decode() + msg = _gmail_client( + monkeypatch, + { + "headers": _GMAIL_HEADERS, + "parts": [ + { + "mimeType": "text/plain", + "headers": _CP1252_PLAIN, + "body": {"data": data}, + } + ], + }, + ) + assert msg.text == "Subject: Top 10 users\n01 · Map" + + +def test_gmail_listener_caps_long_body_and_flags_it(monkeypatch): + long_body = "word " * 1000 # 5000 chars + msg = _gmail_client( + monkeypatch, + { + "headers": _GMAIL_HEADERS, + "parts": [{"mimeType": "text/plain", "body": {"data": _b64(long_body)}}], + }, + ) + assert msg.truncated is True + body = msg.text.split("\n", 1)[1] + assert len(body) <= gmail_mod._INBOUND_BODY_CHARS + + +# ── Outlook ───────────────────────────────────────────────────────────── + + +def _outlook_dispatch(content): + client = OutlookProvider().build_client( + { + "access_token": "tok", + "refresh_token": "ref", + "token_expiry": time.time() + 3600, + "client_id": "cid", + "email": "me@o.com", + }, + lambda d: None, + ) + received, client._message_callback = _collect_callback() + msg = { + "id": "om1", + "from": {"emailAddress": {"address": "bob@x.com", "name": "Bob"}}, + "subject": "Yo", + "body": {"contentType": "text", "content": content}, + "receivedDateTime": "2026-08-12T10:00:00Z", + } + asyncio.run(client._dispatch_message(msg)) + assert len(received) == 1 + return received[0] + + +def test_outlook_forwards_full_text_body(): + msg = _outlook_dispatch("x" * 300) # past the old 255-char bodyPreview + assert msg.text == "Subject: Yo\n" + "x" * 300 + assert msg.truncated is False + + +def test_outlook_caps_long_body_and_flags_it(): + msg = _outlook_dispatch("word " * 1000) + assert msg.truncated is True + assert len(msg.text.split("\n", 1)[1]) <= outlook_mod._INBOUND_BODY_CHARS + + +# ── Slack ─────────────────────────────────────────────────────────────── + +_NAMES = {"U1": "Ada"} + + +@pytest.mark.parametrize( + "raw, expected", + [ + ("hi <@U1>", "hi @Ada"), + ("hi <@U9>", "hi @U9"), # unresolvable → id + ("hi <@U9|bob>", "hi @bob"), # unresolvable → label + ("see <#C1|general>", "see #general"), + ("see <#C1>", "see #C1"), + ("", "docs (https://x.test)"), + ("", "https://x.test"), + ("", "a@b.co"), + (" ", "@here @channel @everyone"), + ("", "@here"), + ("", "@devs"), + ("", "Nov 14"), + ("Tom & Jerry <3", "Tom & Jerry <3"), + # A literal "<@U1>" the user typed arrives escaped and must stay text. + ("<@U1>", "<@U1>"), + ("", ""), + ], +) +def test_slack_to_plain_text(raw, expected): + assert to_plain_text(raw, lambda uid: _NAMES.get(uid, "")) == expected + + +def test_slack_to_plain_text_survives_resolver_error(): + def boom(uid): + raise RuntimeError("api down") + + assert to_plain_text("hi <@U1>", boom) == "hi @U1" + + +def test_slack_display_name_is_memoized(): + from craftos_integrations.providers.slack.client import SlackClient + + client = SlackClient() + calls = [] + + def fake_user_info(uid): + calls.append(uid) + return {"ok": True, "user": {"profile": {"display_name": "Ada"}}} + + client.get_user_info = fake_user_info + assert client._display_name("U1") == "Ada" + assert client._display_name("U1") == "Ada" + assert calls == ["U1"] + + +def test_slack_display_name_does_not_cache_failures(): + from craftos_integrations.providers.slack.client import SlackClient + + client = SlackClient() + calls = [] + + def flaky(uid): + calls.append(uid) + return {"ok": False} + + client.get_user_info = flaky + assert client._display_name("U1") == "" + assert client._display_name("U1") == "" + assert len(calls) == 2 + + +# ── Twitter ───────────────────────────────────────────────────────────── + + +def test_twitter_decodes_text_before_tag_match(): + from craftos_integrations.providers.twitter.client import ( + TwitterClient, + TwitterConfig, + ) + + client = TwitterClient() + client._config = lambda: TwitterConfig(watch_tag="R&D") + received, client._message_callback = _collect_callback() + tweet = {"id": "1", "author_id": "a1", "text": "R&D check this <now>"} + users = {"a1": {"username": "ada", "name": "Ada"}} + asyncio.run(client._dispatch_mention(tweet, users)) + + assert len(received) == 1 + assert "&" not in received[0].text + assert "" in received[0].text + + +# ── payload contract ──────────────────────────────────────────────────── + + +def test_payload_forwards_truncated(): + msg = PlatformMessage(platform="gmail", sender_id="a", text="x", truncated=True) + assert platform_message_payload(msg)["truncated"] is True + assert platform_message_payload(PlatformMessage("gmail", "a"))["truncated"] is False + + +def test_payload_tolerates_legacy_message_without_truncated(): + class OldMessage: + platform = "slack" + sender_id = "u" + sender_name = "" + text = "hi" + channel_id = "" + channel_name = "" + message_id = "" + raw = {} + + assert platform_message_payload(OldMessage())["truncated"] is False + + +# ── host: the chat sees the cut, the agent gets the facts ──────────────── + + +def _ingest(payload): + """Drive AgentBase._handle_external_event; return the chat payload.""" + from app.agent_base import AgentBase + + captured = [] + + async def fake_chat(p): + captured.append(p) + + agent = SimpleNamespace(_handle_chat_message=fake_chat) + asyncio.run(AgentBase._handle_external_event(agent, payload)) + assert len(captured) == 1 + return captured[0] + + +_EMAIL_EVENT = { + "source": "Gmail", + "integrationType": "gmail", + "contactId": "joe@posthog.com", + "contactName": "Joe", + "messageBody": "Subject: Hi\nYou probably also", + "messageId": "m1", +} + + +def test_truncated_body_marks_cut_for_chat_and_gives_agent_facts(): + chat = _ingest({**_EMAIL_EVENT, "truncated": True}) + # Chat details: the text and a visible cut — no agent instructions. + assert chat["message_body"] == "Subject: Hi\nYou probably also…" + # Agent prompt: the same text plus the facts to fetch the rest. + assert "You probably also…" in chat["text"] + assert "[This message is truncated" in chat["text"] + assert "gmail action that gets a message by ID (message ID: m1)" in chat["text"] + + +def test_truncated_body_without_id_still_tells_agent(): + chat = _ingest({**_EMAIL_EVENT, "messageId": "", "truncated": True}) + assert "[This message is truncated: only the beginning is shown above.]" in ( + chat["text"] + ) + + +def test_untruncated_body_is_unchanged(): + chat = _ingest(_EMAIL_EVENT) + assert chat["message_body"] == "Subject: Hi\nYou probably also" + assert "truncated" not in chat["text"].lower() diff --git a/tests/integrations/test_provider_listeners.py b/tests/integrations/test_provider_listeners.py index 8d7fe514..58849d02 100644 --- a/tests/integrations/test_provider_listeners.py +++ b/tests/integrations/test_provider_listeners.py @@ -68,21 +68,21 @@ async def emit(event): GMAIL_MESSAGE = { "id": "m1", "threadId": "t1", - "snippet": "hello there", "payload": { "headers": [ {"name": "From", "value": "Alice "}, {"name": "Subject", "value": "Hi"}, {"name": "Date", "value": "Tue, 11 Aug 2026 10:00:00 +0000"}, ], - # fields-mask shape: parts skeleton only, no body.data. The + # fields-mask shape: parts skeleton with inline text bodies only + # (attachment parts carry an attachmentId, never data). The # nameless attachmentId part is an inline image — not reported. "parts": [ { "partId": "0", "mimeType": "text/plain", "filename": "", - "body": {"size": 20}, + "body": {"size": 11, "data": "aGVsbG8gdGhlcmU="}, # "hello there" }, { "partId": "1", @@ -178,6 +178,7 @@ async def scenario(): "extra": {"message_id": "m1"}, } ], + "truncated": False, } ] assert cursor == {"history_id": "101", "seen_ids": ["m1"]} @@ -235,7 +236,7 @@ async def scenario(): "id": "om1", "from": {"emailAddress": {"address": "bob@x.com", "name": "Bob"}}, "subject": "Yo", - "bodyPreview": "preview text", + "body": {"contentType": "text", "content": "preview text"}, "receivedDateTime": "2026-08-12T10:00:00Z", "conversationId": "conv1", } @@ -296,6 +297,7 @@ async def scenario(): "is_self_message": False, "raw": OUTLOOK_MESSAGE, "attachments": [], + "truncated": False, } ] # Watermark advanced to the newest receivedDateTime; dedup ids kept. @@ -451,6 +453,7 @@ async def scenario(): "is_self_message": False, "raw": message, "attachments": [], + "truncated": False, } ] assert cursor == {"last_timestamps": {"C1": msg_ts}}