-
Notifications
You must be signed in to change notification settings - Fork 46.7k
fix(mattermost): classify API errors, escalate fatals, lock single-instance, surface audio/slash attachments #35645
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
lambertian
wants to merge
1
commit into
NousResearch:main
Choose a base branch
from
lambertian:fix/mattermost-reliability
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -50,6 +50,24 @@ | |
| _RECONNECT_JITTER = 0.2 | ||
|
|
||
|
|
||
| class _MMApiError(Exception): | ||
| """Internal Mattermost API failure carrying transient/permanent class. | ||
|
|
||
| ``retryable`` is True for network failures, request timeouts, HTTP 429 | ||
| rate-limits, and 5xx server errors — conditions that may succeed on a | ||
| retry. It is False for genuine 4xx client errors (400/401/403/404), | ||
| which will not. ``send()`` / ``edit_message()`` translate this into | ||
| ``SendResult.retryable`` so ``BasePlatformAdapter._send_with_retry`` | ||
| (network retry with backoff) and the streaming progress-edit consumer | ||
| (which keeps ``can_edit`` alive across transient edit failures) make the | ||
| right call instead of treating every failure as permanent. | ||
| """ | ||
|
|
||
| def __init__(self, message: str, *, retryable: bool): | ||
| super().__init__(message) | ||
| self.retryable = retryable | ||
|
|
||
|
|
||
| def check_mattermost_requirements() -> bool: | ||
| """Return True if the Mattermost adapter can be used.""" | ||
| token = os.getenv("MATTERMOST_TOKEN", "") | ||
|
|
@@ -109,58 +127,74 @@ def _headers(self) -> Dict[str, str]: | |
| "Content-Type": "application/json", | ||
| } | ||
|
|
||
| async def _api_get(self, path: str) -> Dict[str, Any]: | ||
| """GET /api/v4/{path}.""" | ||
| async def _request_json( | ||
| self, method: str, path: str, *, payload: Optional[Dict[str, Any]] = None | ||
| ) -> Dict[str, Any]: | ||
| """Issue an /api/v4 request, raising :class:`_MMApiError` on failure. | ||
|
|
||
| Distinguishes transient failures (network errors, timeouts, HTTP 429 | ||
| and 5xx) from permanent 4xx client errors so callers that surface a | ||
| ``SendResult`` can set ``retryable`` correctly. Callers that want the | ||
| legacy empty-dict sentinel use :meth:`_api_get` / :meth:`_api_post` / | ||
| :meth:`_api_put`, which wrap this. | ||
| """ | ||
| import aiohttp | ||
| url = f"{self._base_url}/api/v4/{path.lstrip('/')}" | ||
| kwargs: Dict[str, Any] = { | ||
| "headers": self._headers(), | ||
| "timeout": aiohttp.ClientTimeout(total=30), | ||
| } | ||
| if payload is not None: | ||
| kwargs["json"] = payload | ||
| verb = method.upper() | ||
| sender = { | ||
| "GET": self._session.get, | ||
| "POST": self._session.post, | ||
| "PUT": self._session.put, | ||
| }[verb] | ||
| try: | ||
| async with self._session.get(url, headers=self._headers(), timeout=aiohttp.ClientTimeout(total=30)) as resp: | ||
| async with sender(url, **kwargs) as resp: | ||
| if resp.status >= 400: | ||
| body = await resp.text() | ||
| logger.error("MM API GET %s → %s: %s", path, resp.status, body[:200]) | ||
| return {} | ||
| logger.error("MM API %s %s → %s: %s", verb, path, resp.status, body[:200]) | ||
| # 429 (rate-limit) and 5xx (server) are transient; retrying | ||
| # may succeed. 4xx client errors are permanent. | ||
| retryable = resp.status == 429 or resp.status >= 500 | ||
| raise _MMApiError( | ||
| f"Mattermost API {verb} {path} returned HTTP {resp.status}", | ||
| retryable=retryable, | ||
| ) | ||
| return await resp.json() | ||
| except aiohttp.ClientError as exc: | ||
| logger.error("MM API GET %s network error: %s", path, exc) | ||
| except (aiohttp.ClientError, asyncio.TimeoutError) as exc: | ||
| logger.error("MM API %s %s network error: %s", verb, path, exc) | ||
| raise _MMApiError( | ||
| f"Mattermost API {verb} {path} network error: {exc}", | ||
| retryable=True, | ||
| ) from exc | ||
|
|
||
| async def _api_get(self, path: str) -> Dict[str, Any]: | ||
| """GET /api/v4/{path}; returns {} on any failure (legacy sentinel).""" | ||
| try: | ||
| return await self._request_json("GET", path) | ||
| except _MMApiError: | ||
| return {} | ||
|
|
||
| async def _api_post( | ||
| self, path: str, payload: Dict[str, Any] | ||
| ) -> Dict[str, Any]: | ||
| """POST /api/v4/{path} with JSON body.""" | ||
| import aiohttp | ||
| url = f"{self._base_url}/api/v4/{path.lstrip('/')}" | ||
| """POST /api/v4/{path}; returns {} on any failure (legacy sentinel).""" | ||
| try: | ||
| async with self._session.post( | ||
| url, headers=self._headers(), json=payload, | ||
| timeout=aiohttp.ClientTimeout(total=30) | ||
| ) as resp: | ||
| if resp.status >= 400: | ||
| body = await resp.text() | ||
| logger.error("MM API POST %s → %s: %s", path, resp.status, body[:200]) | ||
| return {} | ||
| return await resp.json() | ||
| except aiohttp.ClientError as exc: | ||
| logger.error("MM API POST %s network error: %s", path, exc) | ||
| return await self._request_json("POST", path, payload=payload) | ||
| except _MMApiError: | ||
| return {} | ||
|
|
||
| async def _api_put( | ||
| self, path: str, payload: Dict[str, Any] | ||
| ) -> Dict[str, Any]: | ||
| """PUT /api/v4/{path} with JSON body.""" | ||
| import aiohttp | ||
| url = f"{self._base_url}/api/v4/{path.lstrip('/')}" | ||
| """PUT /api/v4/{path}; returns {} on any failure (legacy sentinel).""" | ||
| try: | ||
| async with self._session.put( | ||
| url, headers=self._headers(), json=payload | ||
| ) as resp: | ||
| if resp.status >= 400: | ||
| body = await resp.text() | ||
| logger.error("MM API PUT %s → %s: %s", path, resp.status, body[:200]) | ||
| return {} | ||
| return await resp.json() | ||
| except aiohttp.ClientError as exc: | ||
| logger.error("MM API PUT %s network error: %s", path, exc) | ||
| return await self._request_json("PUT", path, payload=payload) | ||
| except _MMApiError: | ||
| return {} | ||
|
|
||
| async def _upload_file( | ||
|
|
@@ -198,18 +232,68 @@ async def connect(self) -> bool: | |
|
|
||
| if not self._base_url or not self._token: | ||
| logger.error("Mattermost: URL or token not configured") | ||
| # Missing config never succeeds on retry — escalate as a | ||
| # non-retryable fatal error so the gateway drops the platform | ||
| # from its reconnect queue instead of looping forever. | ||
| self._set_fatal_error( | ||
| "config_missing", | ||
| "MATTERMOST_URL and MATTERMOST_TOKEN must be set", | ||
| retryable=False, | ||
| ) | ||
| return False | ||
|
|
||
| # Single-instance guard: without a scoped lock, two gateway processes | ||
| # configured with the same token each open their own /api/v4/websocket | ||
| # listener and both receive the same 'posted' events. The per-process | ||
| # MessageDeduplicator (self._dedup) is in-memory and cannot dedup | ||
| # across processes, so every inbound post would be handled twice | ||
| # (duplicate agent runs/replies). Key the lock on URL+token like the | ||
| # Discord/Signal/Slack adapters; _acquire_platform_lock records a | ||
| # non-retryable fatal error on conflict. | ||
| if not self._acquire_platform_lock( | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
|
||
| "mattermost-token", | ||
| f"{self._base_url}|{self._token}", | ||
| "Mattermost bot token", | ||
| ): | ||
| return False | ||
|
|
||
| self._session = aiohttp.ClientSession( | ||
| timeout=aiohttp.ClientTimeout(total=30) | ||
| ) | ||
| self._closing = False | ||
|
|
||
| # Verify credentials and fetch bot identity. | ||
| me = await self._api_get("users/me") | ||
| # Verify credentials and fetch bot identity. Distinguish a permanent | ||
| # auth/permission failure (bad token, wrong URL) from a transient | ||
| # network blip reaching users/me so the gateway can stop retrying a | ||
| # revoked token but keep retrying a flaky connection. | ||
| try: | ||
| me = await self._request_json("GET", "users/me") | ||
| except _MMApiError as exc: | ||
| await self._session.close() | ||
| # Release the lock acquired above so a permanent auth failure | ||
| # (or a transient blip) does not leak it; mirrors Discord's | ||
| # _release_platform_lock() on its post-acquire failure exits. | ||
| self._release_platform_lock() | ||
| if exc.retryable: | ||
| logger.error("Mattermost: transient error reaching server — will retry: %s", exc) | ||
| self._set_fatal_error("connect_failed", str(exc), retryable=True) | ||
| else: | ||
| logger.error("Mattermost: authentication failed — check MATTERMOST_TOKEN and MATTERMOST_URL") | ||
| self._set_fatal_error( | ||
| "auth_failed", | ||
| "Mattermost authentication failed — check MATTERMOST_TOKEN/MATTERMOST_URL", | ||
| retryable=False, | ||
| ) | ||
| return False | ||
| if not me or "id" not in me: | ||
| logger.error("Mattermost: failed to authenticate — check MATTERMOST_TOKEN and MATTERMOST_URL") | ||
| await self._session.close() | ||
| self._release_platform_lock() | ||
| self._set_fatal_error( | ||
| "auth_failed", | ||
| "Mattermost authentication failed — check MATTERMOST_TOKEN/MATTERMOST_URL", | ||
| retryable=False, | ||
| ) | ||
| return False | ||
|
|
||
| self._bot_user_id = me["id"] | ||
|
|
@@ -247,6 +331,14 @@ async def disconnect(self) -> None: | |
| if self._session and not self._session.closed: | ||
| await self._session.close() | ||
|
|
||
| # Release the single-instance lock acquired in connect() so another | ||
| # gateway process can take over the token (no-op if never acquired). | ||
| self._release_platform_lock() | ||
|
|
||
| # Flip _running=False and write the 'disconnected' runtime status so | ||
| # is_connected and status tooling reflect the shutdown. No-op-safe | ||
| # when a fatal error is already recorded (base guards on that). | ||
| self._mark_disconnected() | ||
| logger.info("Mattermost: disconnected") | ||
|
|
||
|
|
||
|
|
@@ -293,7 +385,13 @@ async def send( | |
| resolved_root = await self._resolve_root_id(reply_to) | ||
| payload["root_id"] = resolved_root | ||
|
|
||
| data = await self._api_post("posts", payload) | ||
| try: | ||
| data = await self._request_json("POST", "posts", payload=payload) | ||
| except _MMApiError as exc: | ||
| # Propagate transient/permanent class so _send_with_retry | ||
| # retries network/5xx/429 failures with backoff instead of | ||
| # silently dropping the response after one plain-text attempt. | ||
| return SendResult(success=False, error=str(exc), retryable=exc.retryable) | ||
| if not data or "id" not in data: | ||
| return SendResult(success=False, error="Failed to create post") | ||
| last_id = data["id"] | ||
|
|
@@ -328,10 +426,17 @@ async def edit_message( | |
| ) -> SendResult: | ||
| """Edit an existing post.""" | ||
| formatted = self.format_message(content) | ||
| data = await self._api_put( | ||
| f"posts/{message_id}/patch", | ||
| {"message": formatted}, | ||
| ) | ||
| try: | ||
| data = await self._request_json( | ||
| "PUT", | ||
| f"posts/{message_id}/patch", | ||
| payload={"message": formatted}, | ||
| ) | ||
| except _MMApiError as exc: | ||
| # A transient edit failure must not permanently disable in-place | ||
| # progress editing for the rest of a streamed response: surface | ||
| # retryable so the consumer keeps can_edit=True. | ||
| return SendResult(success=False, error=str(exc), retryable=exc.retryable) | ||
| if not data or "id" not in data: | ||
| return SendResult(success=False, error="Failed to edit post") | ||
| return SendResult(success=True, message_id=data["id"]) | ||
|
|
@@ -615,6 +720,19 @@ async def send_multiple_images( | |
| # WebSocket | ||
| # ------------------------------------------------------------------ | ||
|
|
||
| async def _escalate_ws_fatal(self, message: str) -> None: | ||
| """Record a non-retryable fatal error and notify the gateway. | ||
|
|
||
| When the WebSocket listener gives up on a permanent auth/permission | ||
| failure it must bridge that state back to the gateway: otherwise | ||
| ``_ws_task`` ends while ``_running`` stays True, ``is_connected`` | ||
| keeps returning True, and the bot is a zombie that the gateway never | ||
| reconnects. Mirrors IRC's receive-loop ``finally`` which calls | ||
| ``_set_fatal_error`` + ``_notify_fatal_error``. | ||
| """ | ||
| self._set_fatal_error("ws_auth_failed", message, retryable=False) | ||
| await self._notify_fatal_error() | ||
|
|
||
| async def _ws_loop(self) -> None: | ||
| """Connect to the WebSocket and listen for events, reconnecting on failure.""" | ||
| delay = _RECONNECT_BASE_DELAY | ||
|
|
@@ -634,9 +752,11 @@ async def _ws_loop(self) -> None: | |
| err_str = str(exc).lower() | ||
| if isinstance(exc, aiohttp.WSServerHandshakeError) and exc.status in {401, 403}: | ||
| logger.error("Mattermost WS auth failed (HTTP %d) — stopping reconnect", exc.status) | ||
| await self._escalate_ws_fatal(f"Mattermost WebSocket auth failed (HTTP {exc.status})") | ||
| return | ||
| if "401" in err_str or "403" in err_str or "unauthorized" in err_str: | ||
| logger.error("Mattermost WS permanent error: %s — stopping reconnect", exc) | ||
| await self._escalate_ws_fatal(f"Mattermost WebSocket permanent error: {exc}") | ||
| return | ||
| logger.warning("Mattermost WS error: %s — reconnecting in %.0fs", exc, delay) | ||
|
|
||
|
|
@@ -834,12 +954,23 @@ async def _handle_ws_event(self, event: Dict[str, Any]) -> None: | |
| except Exception as exc: | ||
| logger.warning("Mattermost: error downloading file %s: %s", fid, exc) | ||
|
|
||
| # Set message type based on downloaded media types. | ||
| if media_types and msg_type == MessageType.TEXT: | ||
| # Set message type based on downloaded media types. An attachment | ||
| # whose caption happens to start with '/' was provisionally tagged | ||
| # COMMAND above; the media type wins so downstream document/image/ | ||
| # audio surfacing (which keys on the message type) still fires. | ||
| if media_types and msg_type in {MessageType.TEXT, MessageType.COMMAND}: | ||
| if any(m.startswith("image/") for m in media_types): | ||
| msg_type = MessageType.PHOTO | ||
| elif any(m.startswith("audio/") for m in media_types): | ||
| msg_type = MessageType.VOICE | ||
| # Mattermost has no distinct 'voice note' concept — an uploaded | ||
| # audio file is just a file. Classify it AUDIO, not VOICE, so | ||
| # run.py surfaces the cached path to the agent as a referenceable | ||
| # file (the audio_file_paths branch) instead of force-running STT | ||
| # on what may be music, a podcast, or a non-speech clip and never | ||
| # handing the agent the actual file. Mirrors Discord, which only | ||
| # tags true voice-message attachments VOICE and ordinary audio | ||
| # files AUDIO. | ||
| msg_type = MessageType.AUDIO | ||
| elif media_types: | ||
| msg_type = MessageType.DOCUMENT | ||
|
|
||
|
|
||
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Please retain the
..path rejection before constructing this URL. Current main added it ind836b2bacbecause WebSocket-event IDs can be attacker-controlled and otherwise steer authenticated bearer-token requests; the legacy wrappers no longer protect this shared request path.