diff --git a/plugins/platforms/photon/adapter.py b/plugins/platforms/photon/adapter.py index efbf6fb288..3f1497f520 100644 --- a/plugins/platforms/photon/adapter.py +++ b/plugins/platforms/photon/adapter.py @@ -21,7 +21,7 @@ import sys import time from datetime import datetime, timezone from pathlib import Path -from typing import TYPE_CHECKING, Any, Callable, Dict, List, Optional, Tuple +from typing import TYPE_CHECKING, Any, AsyncIterator, Callable, Dict, List, Optional, Tuple from urllib.parse import urlparse if TYPE_CHECKING: # type checkers see httpx as always-imported; runtime keeps it optional @@ -86,6 +86,18 @@ _TARGET_NOT_ALLOWED_MESSAGE = ( "targets — upgrade to a dedicated line or use another delivery channel") +async def _aiter_ndjson_lines(response: Any) -> AsyncIterator[str]: + """Split the sidecar stream on its protocol delimiter, LF, and nothing else.""" + pending = "" + async for chunk in response.aiter_text(): + lines = (pending + chunk).split("\n") + pending = lines.pop() + for line in lines: + yield line + if pending: + yield pending + + # -- Sidecar runtime record ---------------------------------------------------- def _runtime_record_path() -> Path: @@ -621,7 +633,7 @@ class PhotonAdapter(BasePlatformAdapter): if resp.status_code != 200: raise RuntimeError(f"/inbound returned {resp.status_code}") backoff = 1.0 - async for line in resp.aiter_lines(): + async for line in _aiter_ndjson_lines(resp): if not self._inbound_running: break line = line.strip() diff --git a/tests/plugins/platforms/photon/test_inbound.py b/tests/plugins/platforms/photon/test_inbound.py index 0bee1d169c..c604158574 100644 --- a/tests/plugins/platforms/photon/test_inbound.py +++ b/tests/plugins/platforms/photon/test_inbound.py @@ -16,7 +16,7 @@ import pytest from gateway.config import Platform, PlatformConfig from gateway.platforms.event import MessageEvent, MessageType -from plugins.platforms.photon.adapter import PhotonAdapter +from plugins.platforms.photon.adapter import PhotonAdapter, _aiter_ndjson_lines def _make_adapter(monkeypatch: pytest.MonkeyPatch) -> PhotonAdapter: @@ -150,6 +150,30 @@ async def test_on_inbound_line_dispatches_and_dedups( assert captured[0].text == "ping" +@pytest.mark.asyncio +async def test_ndjson_stream_preserves_unicode_line_separators( + monkeypatch: pytest.MonkeyPatch, +) -> None: + adapter = _make_adapter(monkeypatch) + captured = _capture(adapter, monkeypatch) + separated = "first\u2028second\u2029third\u0085fourth" + payloads = [ + json.dumps(_dm_event(separated, msg_id="unicode-lines"), ensure_ascii=False), + json.dumps(_dm_event("ordinary", msg_id="ordinary-line")), + ] + + class ChunkedResponse: + async def aiter_text(self): + stream = "\n".join(payloads) + "\n" + for chunk in (stream[:17], stream[17:43], stream[43:]): + yield chunk + + async for line in _aiter_ndjson_lines(ChunkedResponse()): + await adapter._on_inbound_line(line) + + assert [event.text for event in captured] == [separated, "ordinary"] + + def test_is_duplicate_window(monkeypatch: pytest.MonkeyPatch) -> None: adapter = _make_adapter(monkeypatch) assert adapter._dedup.is_duplicate("id-1") is False