fix(photon): preserve Unicode NDJSON separators
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user