diff --git a/tests/plugins/platforms/photon/test_inbound.py b/tests/plugins/platforms/photon/test_inbound.py index c604158574..c624d9d1c7 100644 --- a/tests/plugins/platforms/photon/test_inbound.py +++ b/tests/plugins/platforms/photon/test_inbound.py @@ -1,8 +1,7 @@ """Inbound dispatch + dedup tests for PhotonAdapter. -These bypass the loopback HTTP stream — they call ``_dispatch_inbound`` / -``_on_inbound_line`` / ``_dedup`` directly, exercising the -sidecar-event parsing without spawning the Node sidecar or binding ports. +These exercise the sidecar-event stream and parsing without spawning the +Node sidecar or binding ports. """ from __future__ import annotations @@ -16,7 +15,7 @@ import pytest from gateway.config import Platform, PlatformConfig from gateway.platforms.event import MessageEvent, MessageType -from plugins.platforms.photon.adapter import PhotonAdapter, _aiter_ndjson_lines +from plugins.platforms.photon.adapter import PhotonAdapter def _make_adapter(monkeypatch: pytest.MonkeyPatch) -> PhotonAdapter: @@ -155,23 +154,47 @@ 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")), ] + stream = payloads[0] + "\n\n" + payloads[1] + received: List[str] = [] class ChunkedResponse: + status_code = 200 + + async def __aenter__(self): + return self + + async def __aexit__(self, *_args): + return None + 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) + async def aiter_lines(self): + for line in stream.splitlines(): + yield line - assert [event.text for event in captured] == [separated, "ordinary"] + class Client: + def stream(self, *_args, **_kwargs): + return ChunkedResponse() + + async def capture_line(line: str) -> None: + received.append(line) + if len(received) == 2: + adapter._inbound_running = False + + adapter._http_client = Client() + adapter._inbound_running = True + monkeypatch.setattr(adapter, "_on_inbound_line", capture_line) + + await adapter._inbound_loop() + + assert received == payloads def test_is_duplicate_window(monkeypatch: pytest.MonkeyPatch) -> None: