mirror of
https://github.com/Routstr/routstr-core.git
synced 2026-10-05 20:28:23 +00:00
fix: handle incomplete upstream streams
This commit is contained in:
@@ -1426,6 +1426,14 @@ class BaseUpstreamProvider:
|
|||||||
if done_seen:
|
if done_seen:
|
||||||
yield b"data: [DONE]\n\n"
|
yield b"data: [DONE]\n\n"
|
||||||
|
|
||||||
|
except httpx.RemoteProtocolError as stream_error:
|
||||||
|
logger.warning(
|
||||||
|
"Upstream stream ended before the response was complete",
|
||||||
|
extra={
|
||||||
|
"error": str(stream_error),
|
||||||
|
"key_hash": key.hashed_key[:8] + "...",
|
||||||
|
},
|
||||||
|
)
|
||||||
except Exception as stream_error:
|
except Exception as stream_error:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Streaming interrupted; finalizing before closing upstream",
|
"Streaming interrupted; finalizing before closing upstream",
|
||||||
@@ -1869,6 +1877,14 @@ class BaseUpstreamProvider:
|
|||||||
if done_seen:
|
if done_seen:
|
||||||
yield b"data: [DONE]\n\n"
|
yield b"data: [DONE]\n\n"
|
||||||
|
|
||||||
|
except httpx.RemoteProtocolError as stream_error:
|
||||||
|
logger.warning(
|
||||||
|
"Upstream Responses API stream ended before the response was complete",
|
||||||
|
extra={
|
||||||
|
"error": str(stream_error),
|
||||||
|
"key_hash": key.hashed_key[:8] + "...",
|
||||||
|
},
|
||||||
|
)
|
||||||
except Exception as stream_error:
|
except Exception as stream_error:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Responses API streaming interrupted; finalizing before closing upstream",
|
"Responses API streaming interrupted; finalizing before closing upstream",
|
||||||
|
|||||||
@@ -404,11 +404,8 @@ async def test_partial_remote_protocol_error_finalizes_and_closes_once(
|
|||||||
client=client,
|
client=client,
|
||||||
)
|
)
|
||||||
emitted = bytearray()
|
emitted = bytearray()
|
||||||
with pytest.raises(httpx.RemoteProtocolError):
|
async for chunk in response.body_iterator:
|
||||||
async for chunk in response.body_iterator:
|
emitted.extend(chunk.encode() if isinstance(chunk, str) else bytes(chunk))
|
||||||
emitted.extend(
|
|
||||||
chunk.encode() if isinstance(chunk, str) else bytes(chunk)
|
|
||||||
)
|
|
||||||
|
|
||||||
adjust.assert_awaited_once()
|
adjust.assert_awaited_once()
|
||||||
if finalization_fails:
|
if finalization_fails:
|
||||||
@@ -423,7 +420,7 @@ async def test_partial_remote_protocol_error_finalizes_and_closes_once(
|
|||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
@pytest.mark.parametrize("api", ["chat", "responses"])
|
@pytest.mark.parametrize("api", ["chat", "responses"])
|
||||||
async def test_partial_stream_preserves_transport_error_when_billing_db_is_down(
|
async def test_partial_stream_closes_when_billing_db_is_down(
|
||||||
api: str,
|
api: str,
|
||||||
) -> None:
|
) -> None:
|
||||||
provider = BaseUpstreamProvider(
|
provider = BaseUpstreamProvider(
|
||||||
@@ -476,9 +473,8 @@ async def test_partial_stream_preserves_transport_error_when_billing_db_is_down(
|
|||||||
reservation_snapshot=snapshot,
|
reservation_snapshot=snapshot,
|
||||||
client=client,
|
client=client,
|
||||||
)
|
)
|
||||||
with pytest.raises(httpx.RemoteProtocolError, match="incomplete chunked read"):
|
async for _ in response.body_iterator:
|
||||||
async for _ in response.body_iterator:
|
pass
|
||||||
pass
|
|
||||||
|
|
||||||
upstream_response.aclose.assert_awaited_once()
|
upstream_response.aclose.assert_awaited_once()
|
||||||
client.aclose.assert_awaited_once()
|
client.aclose.assert_awaited_once()
|
||||||
|
|||||||
Reference in New Issue
Block a user