Files
Aether/tests/api/handlers/base/test_upstream_stream_bridge.py
fawney19 e0286aebe3 refactor: 共享请求管道、按需懒加载、流式内存护栏与连接池治理
- 抽取 ApiRequestPipeline 单例,44 个路由文件共享同一实例
- Handler/Adapter 模块级 __getattr__ 延迟导入,减少启动时间
- 新增 ensure_stream_buffer_limit() 流式内存护栏(16MB 单行 / 32MB 总量)
- HTTP 空闲连接清理与 curl_cffi LRU 会话池
- ensure_providers_bootstrapped 按需引导指定 provider_types
- Usage 事件序列化迁移至 msgpack,Redis codec 隔离
- 启动预热任务(/readyz 就绪门控)与优雅关闭
- 通知邮件模块独立开关与 SMTP 配置校验
- CryptoService DCL 线程安全修复
- 通知模块开关 DB 查询 30s 内存缓存
- /readyz 对 unknown 状态返回 503
- 预热关闭 5s 超时保护
- 预热适配器逐个 try-except 容错
- FormatConversionRegistry 哨兵模式防并发重复物化
- 流式缓冲检查无条件执行

Closes #230

Co-authored-by: AAEE86 <ppk0227@hotmail.com>
2026-03-14 11:59:07 +08:00

146 lines
4.7 KiB
Python

from __future__ import annotations
import json
from collections.abc import AsyncIterator
import pytest
from src.api.handlers.base.upstream_stream_bridge import (
aggregate_upstream_stream_to_internal_response,
)
from src.config.constants import StreamDefaults
from src.core.api_format.conversion import register_default_normalizers
from src.core.api_format.conversion.internal import TextBlock
from src.core.exceptions import ProviderNotAvailableException
async def _iter_stream_lines(lines: list[str]) -> AsyncIterator[bytes]:
for line in lines:
yield line.encode("utf-8")
@pytest.mark.asyncio
async def test_aggregate_claude_stream_uses_message_start_usage_when_message_delta_absent() -> None:
register_default_normalizers()
lines = [
"data: "
+ json.dumps(
{
"type": "message_start",
"message": {
"id": "msg_bridge_usage",
"type": "message",
"role": "assistant",
"model": "claude-sonnet-4-5",
"content": [],
"usage": {
"input_tokens": 120,
"output_tokens": 0,
"cache_read_input_tokens": 11,
},
},
},
ensure_ascii=False,
)
+ "\n",
"data: "
+ json.dumps(
{
"type": "content_block_start",
"index": 0,
"content_block": {"type": "text", "text": ""},
},
ensure_ascii=False,
)
+ "\n",
"data: "
+ json.dumps(
{
"type": "content_block_delta",
"index": 0,
"delta": {"type": "text_delta", "text": "hello"},
},
ensure_ascii=False,
)
+ "\n",
"data: "
+ json.dumps({"type": "content_block_stop", "index": 0}, ensure_ascii=False)
+ "\n",
]
internal = await aggregate_upstream_stream_to_internal_response(
_iter_stream_lines(lines),
provider_api_format="claude:cli",
provider_name="claude_code",
model="claude-sonnet-4-5",
request_id="req_bridge_usage",
)
assert internal.usage is not None
assert internal.usage.input_tokens == 120
assert internal.usage.output_tokens == 0
assert internal.usage.cache_read_tokens == 11
assert len(internal.content) == 1
assert isinstance(internal.content[0], TextBlock)
assert internal.content[0].text == "hello"
@pytest.mark.asyncio
async def test_aggregate_stream_raises_when_buffer_exceeds_limit() -> None:
register_default_normalizers()
async def _iter_overflow_bytes() -> AsyncIterator[bytes]:
yield b"x" * (StreamDefaults.MAX_STREAM_BUFFER_BYTES + 1)
with pytest.raises(ProviderNotAvailableException):
await aggregate_upstream_stream_to_internal_response(
_iter_overflow_bytes(),
provider_api_format="claude:cli",
provider_name="claude_code",
model="claude-sonnet-4-5",
request_id="req_bridge_overflow",
)
@pytest.mark.asyncio
async def test_aggregate_stream_raises_when_total_buffer_exceeds_hard_limit(
monkeypatch: pytest.MonkeyPatch,
) -> None:
register_default_normalizers()
monkeypatch.setattr(StreamDefaults, "MAX_STREAM_BUFFER_BYTES", 64)
monkeypatch.setattr(StreamDefaults, "MAX_STREAM_BUFFER_TOTAL_BYTES", 80)
async def _iter_total_overflow_bytes() -> AsyncIterator[bytes]:
yield b":" + (b"a" * 30) + b"\n" + b":" + (b"b" * 30) + b"\n" + b":" + (b"c" * 30) + b"\n"
with pytest.raises(ProviderNotAvailableException):
await aggregate_upstream_stream_to_internal_response(
_iter_total_overflow_bytes(),
provider_api_format="claude:cli",
provider_name="claude_code",
model="claude-sonnet-4-5",
request_id="req_bridge_total_overflow",
)
@pytest.mark.asyncio
async def test_aggregate_stream_allows_large_chunk_with_multiple_complete_lines(
monkeypatch: pytest.MonkeyPatch,
) -> None:
register_default_normalizers()
monkeypatch.setattr(StreamDefaults, "MAX_STREAM_BUFFER_BYTES", 64)
async def _iter_multiline_bytes() -> AsyncIterator[bytes]:
yield b":" + (b"a" * 30) + b"\n" + b":" + (b"b" * 30) + b"\n" + b":" + (b"c" * 30) + b"\n"
internal = await aggregate_upstream_stream_to_internal_response(
_iter_multiline_bytes(),
provider_api_format="claude:cli",
provider_name="claude_code",
model="claude-sonnet-4-5",
request_id="req_bridge_multiline",
)
assert internal is not None