From 38edc098540294f0dd92190d74c03a27022c0cdd Mon Sep 17 00:00:00 2001 From: worldmozara Date: Wed, 15 Jul 2026 20:52:17 +0800 Subject: [PATCH] core/lmproto: move Responses API into openai package openai_compatible stays chat-completions-only; OpenAI provider uses lmproto.openai.OpenAIClient. --- CLAUDE.md | 8 +- src/plyngent/lmproto/openai/__init__.py | 12 ++ src/plyngent/lmproto/openai/client.py | 133 ++++++++++++++++++ .../responses_model.py => openai/model.py} | 5 +- .../lmproto/openai_compatible/__init__.py | 12 +- .../lmproto/openai_compatible/client.py | 80 ++--------- src/plyngent/runtime/client_factory.py | 11 +- tests/test_lmproto/test_responses.py | 6 +- tests/test_runtime/test_client_factory.py | 7 +- 9 files changed, 180 insertions(+), 94 deletions(-) create mode 100644 src/plyngent/lmproto/openai/__init__.py create mode 100644 src/plyngent/lmproto/openai/client.py rename src/plyngent/lmproto/{openai_compatible/responses_model.py => openai/model.py} (98%) diff --git a/CLAUDE.md b/CLAUDE.md index e06e5ac..d9963a8 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -30,13 +30,15 @@ All protocol models use `msgspec.Struct` — not dataclasses, not Pydantic. Opti ``` lmproto/openai_compatible/ ← Base: model, config, client +lmproto/openai/ ← OpenAI platform: Responses + chat (extends base) lmproto/deepseek/openai_compat/ ← Extends base via inheritance + extra fields ``` - **`openai_compatible/model.py`** — tagged chat messages (`SystemChatMessage`, `UserChatMessage`, …), tools, request/response, streaming chunks. -- **`openai_compatible/responses_model.py`** — OpenAI Responses API (`ResponsesCreateParam`, `Response`, function_call items, stream events). -- **`openai_compatible/client.py`** — `BaseOpenAIClient` / `OpenAIClient` via `niquests` async + SSE; `chat_completions`, `models`, `responses` / `get_response` / `delete_response`. +- **`openai_compatible/client.py`** — `BaseOpenAIClient` / `OpenAICompatibleClient` via `niquests` async + SSE; `chat_completions`, `models` only (`OpenAIClient` is a compat alias). - **`openai_compatible/config.py`** — `OpenAIConfig` (token + base URL). +- **`openai/model.py`** — OpenAI Responses API (`ResponsesCreateParam`, `Response`, function_call items, stream events). +- **`openai/client.py`** — platform `OpenAIClient`: chat completions + `responses` / `get_response` / `delete_response`. - **DeepSeek** — `DeepseekOpenAIClient`; models add `reasoning_content`, `prefix`, `ThinkingOptions`. Config default model ids: `deepseek-v4-flash`, `deepseek-v4-pro`. ### Config (`config/`) @@ -45,7 +47,7 @@ TOML load/store (`ConfigStore`): `[providers]` tagged union presets, `[database] ### Runtime (`runtime/`) -`create_client(provider)` maps config `Provider` → protocol client. OpenAI / openai-compatible / deepseek(openai convention) supported; anthropic and deepseek anthropic convention raise `ProviderNotSupportedError`. +`create_client(provider)` maps config `Provider` → protocol client. OpenAI → `lmproto.openai.OpenAIClient` (Responses-capable); openai-compatible / deepseek(openai convention) → chat-completions clients; anthropic and deepseek anthropic convention raise `ProviderNotSupportedError`. ### Memory (`memory/`) diff --git a/src/plyngent/lmproto/openai/__init__.py b/src/plyngent/lmproto/openai/__init__.py new file mode 100644 index 0000000..ab89d60 --- /dev/null +++ b/src/plyngent/lmproto/openai/__init__.py @@ -0,0 +1,12 @@ +from .client import OpenAIClient as OpenAIClient +from .model import Response as Response +from .model import ResponseDeleted as ResponseDeleted +from .model import ResponseFunctionTool as ResponseFunctionTool +from .model import ResponseFunctionToolCall as ResponseFunctionToolCall +from .model import ResponseFunctionToolCallOutput as ResponseFunctionToolCallOutput +from .model import ResponseOutputMessage as ResponseOutputMessage +from .model import ResponseOutputText as ResponseOutputText +from .model import ResponsesCreateParam as ResponsesCreateParam +from .model import ResponseStreamEvent as ResponseStreamEvent +from .model import response_function_calls as response_function_calls +from .model import response_output_text as response_output_text diff --git a/src/plyngent/lmproto/openai/client.py b/src/plyngent/lmproto/openai/client.py new file mode 100644 index 0000000..03b6430 --- /dev/null +++ b/src/plyngent/lmproto/openai/client.py @@ -0,0 +1,133 @@ +"""OpenAI platform client: chat completions (compat) + Responses API.""" + +from __future__ import annotations + +from typing import TYPE_CHECKING, Any, Literal, cast, overload + +import msgspec + +from plyngent.lmproto.openai_compatible.client import BaseOpenAIClient +from plyngent.lmproto.openai_compatible.config import OpenAIConfig # noqa: TC001 +from plyngent.lmproto.openai_compatible.model import ( # noqa: TC001 + ChatCompletionChunk, + ChatCompletionResponse, + ChatCompletionsParam, +) + +from .model import Response, ResponseDeleted, ResponsesCreateParam, ResponseStreamEvent + +if TYPE_CHECKING: + from collections.abc import AsyncIterator + + +class OpenAIClient(BaseOpenAIClient): + """Full OpenAI HTTP client (``/chat/completions`` + ``/responses`` + ``/models``).""" + + response_decoder: msgspec.json.Decoder[Response] + response_deleted_decoder: msgspec.json.Decoder[ResponseDeleted] + response_event_decoder: msgspec.json.Decoder[ResponseStreamEvent] + + def __init__(self, config: OpenAIConfig) -> None: + super().__init__(config) + self.response_decoder = msgspec.json.Decoder(Response) + self.response_deleted_decoder = msgspec.json.Decoder(ResponseDeleted) + self.response_event_decoder = msgspec.json.Decoder(ResponseStreamEvent) + + @overload + async def chat_completions( + self, param: ChatCompletionsParam, *, stream: Literal[False] = False + ) -> ChatCompletionResponse: ... + + @overload + async def chat_completions( + self, param: ChatCompletionsParam, *, stream: Literal[True] + ) -> AsyncIterator[ChatCompletionChunk]: ... + + async def chat_completions( + self, param: ChatCompletionsParam, *, stream: bool = False + ) -> ChatCompletionResponse | AsyncIterator[ChatCompletionChunk]: + param = msgspec.structs.replace(param, stream=stream) + data = self.encoder.encode(param) + if stream: + resp = await self.session.post( + "/chat/completions", + data=data, + headers={"Content-Type": "application/json"}, + stream=True, + ) + return self._parse_sse(resp) + resp = await self.session.post( + "/chat/completions", + data=data, + headers={"Content-Type": "application/json"}, + stream=False, + ) + body = await self._read_json_body(resp, what="chat completions") + return self.decoder.decode(body) + + @overload + async def responses( + self, param: ResponsesCreateParam, *, stream: Literal[False] = False + ) -> Response: ... + + @overload + async def responses( + self, param: ResponsesCreateParam, *, stream: Literal[True] + ) -> AsyncIterator[ResponseStreamEvent]: ... + + async def responses( + self, param: ResponsesCreateParam, *, stream: bool = False + ) -> Response | AsyncIterator[ResponseStreamEvent]: + """Create a model response via OpenAI ``POST /responses``.""" + param = msgspec.structs.replace(param, stream=stream) + data = self.encoder.encode(param) + if stream: + resp = await self.session.post( + "/responses", + data=data, + headers={"Content-Type": "application/json"}, + stream=True, + ) + return self._parse_response_sse(resp) + resp = await self.session.post( + "/responses", + data=data, + headers={"Content-Type": "application/json"}, + stream=False, + ) + body = await self._read_json_body(resp, what="responses") + return self.response_decoder.decode(body) + + async def get_response(self, response_id: str) -> Response: + """Retrieve a stored response via ``GET /responses/{id}``.""" + resp = await self.session.get(f"/responses/{response_id}", stream=False) + body = await self._read_json_body(resp, what="responses") + return self.response_decoder.decode(body) + + async def delete_response(self, response_id: str) -> ResponseDeleted: + """Delete a stored response via ``DELETE /responses/{id}``.""" + resp = await self.session.delete(f"/responses/{response_id}", stream=False) + body = await self._read_json_body(resp, what="responses") + return self.response_deleted_decoder.decode(body) + + async def _parse_response_sse(self, resp: object) -> AsyncIterator[ResponseStreamEvent]: + """Yield Responses API SSE events; stop at ``data: [DONE]``.""" + from plyngent.lmproto.openai_compatible.client import ( + close_async_response, + sse_data_payload, + ) + + typed = cast("Any", resp) + try: + await self._ensure_ok(typed, what="responses") + async for raw in typed.iter_lines(): + if not raw: + continue + parsed = sse_data_payload(bytes(raw)) + if parsed is None: + continue + if parsed is False: + break + yield self.response_event_decoder.decode(parsed) + finally: + await close_async_response(typed) diff --git a/src/plyngent/lmproto/openai_compatible/responses_model.py b/src/plyngent/lmproto/openai/model.py similarity index 98% rename from src/plyngent/lmproto/openai_compatible/responses_model.py rename to src/plyngent/lmproto/openai/model.py index a7f51a7..b24637a 100644 --- a/src/plyngent/lmproto/openai_compatible/responses_model.py +++ b/src/plyngent/lmproto/openai/model.py @@ -13,10 +13,9 @@ from typing import Any, Literal, cast import msgspec from msgspec import UNSET, Struct, field +from plyngent.lmproto.openai_compatible.model import ServiceTier # noqa: TC001 from plyngent.typedef import JSONSchema, Unset # noqa: TC001 -from .model import ServiceTier # noqa: TC001 - type ResponseStatus = Literal[ "completed", "failed", @@ -294,7 +293,7 @@ def response_output_text(response: Response) -> str: for raw in response.output: if raw.get("type") != "message": continue - content = raw.get("content") + content: list[object] | None = raw.get("content") if not isinstance(content, list): continue for block_obj in content: diff --git a/src/plyngent/lmproto/openai_compatible/__init__.py b/src/plyngent/lmproto/openai_compatible/__init__.py index 7597079..d01c651 100644 --- a/src/plyngent/lmproto/openai_compatible/__init__.py +++ b/src/plyngent/lmproto/openai_compatible/__init__.py @@ -1,5 +1,6 @@ from .client import BaseOpenAIClient as BaseOpenAIClient from .client import OpenAIClient as OpenAIClient +from .client import OpenAICompatibleClient as OpenAICompatibleClient from .config import OpenAIConfig as OpenAIConfig from .model import AnyAssistantToolCall as AnyAssistantToolCall from .model import AnyChatMessage as AnyChatMessage @@ -40,14 +41,3 @@ from .model import ToolFunctionItem as ToolFunctionItem from .model import UserChatMessage as UserChatMessage from .model import Verbosity as Verbosity from .model import VoiceName as VoiceName -from .responses_model import Response as Response -from .responses_model import ResponseDeleted as ResponseDeleted -from .responses_model import ResponseFunctionTool as ResponseFunctionTool -from .responses_model import ResponseFunctionToolCall as ResponseFunctionToolCall -from .responses_model import ResponseFunctionToolCallOutput as ResponseFunctionToolCallOutput -from .responses_model import ResponseOutputMessage as ResponseOutputMessage -from .responses_model import ResponseOutputText as ResponseOutputText -from .responses_model import ResponsesCreateParam as ResponsesCreateParam -from .responses_model import ResponseStreamEvent as ResponseStreamEvent -from .responses_model import response_function_calls as response_function_calls -from .responses_model import response_output_text as response_output_text diff --git a/src/plyngent/lmproto/openai_compatible/client.py b/src/plyngent/lmproto/openai_compatible/client.py index 7d32c54..20ae148 100644 --- a/src/plyngent/lmproto/openai_compatible/client.py +++ b/src/plyngent/lmproto/openai_compatible/client.py @@ -18,7 +18,6 @@ from .model import ( ModelsResponse, StreamToolCallDelta, ) -from .responses_model import Response, ResponseDeleted, ResponsesCreateParam, ResponseStreamEvent if TYPE_CHECKING: from collections.abc import AsyncIterator @@ -31,7 +30,7 @@ _HTTP_ERROR = 400 _ERROR_BODY_PREVIEW = 500 -async def _close_async_response(resp: AsyncResponse) -> None: +async def close_async_response(resp: AsyncResponse) -> None: """Close a streaming response; support sync or async close().""" for name in ("aclose", "close"): method = getattr(resp, name, None) @@ -43,6 +42,10 @@ async def _close_async_response(resp: AsyncResponse) -> None: return +# Internal alias kept for older call sites in this module. +_close_async_response = close_async_response + + async def read_response_body(resp: object) -> bytes | str | None: """Read response body; niquests ``AsyncResponse.content`` is an async property.""" content = getattr(resp, "content", None) @@ -98,14 +101,13 @@ def http_error_message( class BaseOpenAIClient: + """Shared OpenAI-compatible HTTP surface (session, models, chat SSE).""" + session: AsyncSession encoder: msgspec.json.Encoder decoder: msgspec.json.Decoder[ChatCompletionResponse] chunk_decoder: msgspec.json.Decoder[ChatCompletionChunk] models_decoder: msgspec.json.Decoder[ModelsResponse] - response_decoder: msgspec.json.Decoder[Response] - response_deleted_decoder: msgspec.json.Decoder[ResponseDeleted] - response_event_decoder: msgspec.json.Decoder[ResponseStreamEvent] def __init__(self, config: OpenAIConfig) -> None: self.session = niquests.AsyncSession( @@ -116,9 +118,6 @@ class BaseOpenAIClient: self.decoder = msgspec.json.Decoder(ChatCompletionResponse) self.chunk_decoder = msgspec.json.Decoder(ChatCompletionChunk) self.models_decoder = msgspec.json.Decoder(ModelsResponse) - self.response_decoder = msgspec.json.Decoder(Response) - self.response_deleted_decoder = msgspec.json.Decoder(ResponseDeleted) - self.response_event_decoder = msgspec.json.Decoder(ResponseStreamEvent) async def models(self) -> list[str]: """List model ids via OpenAI-compatible ``GET /models``. @@ -164,22 +163,6 @@ class BaseOpenAIClient: finally: await _close_async_response(resp) - async def _parse_response_sse(self, resp: AsyncResponse) -> AsyncIterator[ResponseStreamEvent]: - """Yield Responses API SSE events; stop at ``data: [DONE]``.""" - try: - await self._ensure_ok(resp, what="responses") - async for raw in resp.iter_lines(): - if not raw: - continue - parsed = sse_data_payload(bytes(raw)) - if parsed is None: - continue - if parsed is False: - break - yield self.response_event_decoder.decode(parsed) - finally: - await _close_async_response(resp) - async def _read_json_body(self, resp: object, *, what: str) -> bytes: await self._ensure_ok(resp, what=what) body = await read_response_body(resp) @@ -192,7 +175,9 @@ class BaseOpenAIClient: return bytes(body) -class OpenAIClient(BaseOpenAIClient): +class OpenAICompatibleClient(BaseOpenAIClient): + """Chat Completions only (``POST /chat/completions`` + ``GET /models``).""" + def __init__(self, config: OpenAIConfig) -> None: super().__init__(config) @@ -229,50 +214,9 @@ class OpenAIClient(BaseOpenAIClient): body = await self._read_json_body(resp, what="chat completions") return self.decoder.decode(body) - @overload - async def responses( - self, param: ResponsesCreateParam, *, stream: Literal[False] = False - ) -> Response: ... - @overload - async def responses( - self, param: ResponsesCreateParam, *, stream: Literal[True] - ) -> AsyncIterator[ResponseStreamEvent]: ... - - async def responses( - self, param: ResponsesCreateParam, *, stream: bool = False - ) -> Response | AsyncIterator[ResponseStreamEvent]: - """Create a model response via OpenAI ``POST /responses``.""" - param = msgspec.structs.replace(param, stream=stream) - data = self.encoder.encode(param) - if stream: - resp = await self.session.post( - "/responses", - data=data, - headers={"Content-Type": "application/json"}, - stream=True, - ) - return self._parse_response_sse(resp) - resp = await self.session.post( - "/responses", - data=data, - headers={"Content-Type": "application/json"}, - stream=False, - ) - body = await self._read_json_body(resp, what="responses") - return self.response_decoder.decode(body) - - async def get_response(self, response_id: str) -> Response: - """Retrieve a stored response via ``GET /responses/{id}``.""" - resp = await self.session.get(f"/responses/{response_id}", stream=False) - body = await self._read_json_body(resp, what="responses") - return self.response_decoder.decode(body) - - async def delete_response(self, response_id: str) -> ResponseDeleted: - """Delete a stored response via ``DELETE /responses/{id}``.""" - resp = await self.session.delete(f"/responses/{response_id}", stream=False) - body = await self._read_json_body(resp, what="responses") - return self.response_deleted_decoder.decode(body) +# Backward-compatible alias: many call sites still import OpenAIClient from this package. +OpenAIClient = OpenAICompatibleClient def merge_stream_tool_calls(deltas: list[StreamToolCallDelta]) -> list[AssistantFunctionToolCall]: diff --git a/src/plyngent/runtime/client_factory.py b/src/plyngent/runtime/client_factory.py index 48f1fd0..cbd856b 100644 --- a/src/plyngent/runtime/client_factory.py +++ b/src/plyngent/runtime/client_factory.py @@ -9,7 +9,8 @@ from plyngent.config.models import ( Provider, ) from plyngent.lmproto.deepseek import DeepseekOpenAIClient -from plyngent.lmproto.openai_compatible import OpenAIClient, OpenAIConfig +from plyngent.lmproto.openai import OpenAIClient +from plyngent.lmproto.openai_compatible import OpenAICompatibleClient, OpenAIConfig if TYPE_CHECKING: from collections.abc import Mapping @@ -17,7 +18,9 @@ if TYPE_CHECKING: DEFAULT_OPENAI_BASE_URL = "https://api.openai.com/v1" DEFAULT_DEEPSEEK_BASE_URL = "https://api.deepseek.com" -type OpenAICompatibleClient = OpenAIClient | DeepseekOpenAIClient +type ProtocolClient = OpenAIClient | OpenAICompatibleClient | DeepseekOpenAIClient +# Backward-compatible name used by older imports/tests. +type OpenAICompatibleClientUnion = ProtocolClient class ProviderNotSupportedError(NotImplementedError): @@ -45,7 +48,7 @@ def _deepseek_convention(extras: Mapping[str, str]) -> str: return extras.get("convention", "openai").lower() -def create_client(provider: Provider) -> OpenAICompatibleClient: +def create_client(provider: Provider) -> ProtocolClient: """Build a protocol client for the given provider config entry. Raises: @@ -55,7 +58,7 @@ def create_client(provider: Provider) -> OpenAICompatibleClient: if isinstance(provider, OpenAIProvider): return OpenAIClient(provider_to_openai_config(provider)) if isinstance(provider, OpenAICompatibleProvider): - return OpenAIClient(provider_to_openai_config(provider)) + return OpenAICompatibleClient(provider_to_openai_config(provider)) if isinstance(provider, DeepseekProvider): convention = _deepseek_convention(provider.extras) if convention in {"openai", "openai_compat", "openai-compatible"}: diff --git a/tests/test_lmproto/test_responses.py b/tests/test_lmproto/test_responses.py index d3c506b..a3c4d38 100644 --- a/tests/test_lmproto/test_responses.py +++ b/tests/test_lmproto/test_responses.py @@ -3,9 +3,8 @@ from __future__ import annotations import msgspec import pytest -from plyngent.lmproto.openai_compatible.client import OpenAIClient -from plyngent.lmproto.openai_compatible.config import OpenAIConfig -from plyngent.lmproto.openai_compatible.responses_model import ( +from plyngent.lmproto.openai import ( + OpenAIClient, Response, ResponseDeleted, ResponseFunctionTool, @@ -15,6 +14,7 @@ from plyngent.lmproto.openai_compatible.responses_model import ( response_function_calls, response_output_text, ) +from plyngent.lmproto.openai_compatible.config import OpenAIConfig def _sample_response_body() -> bytes: diff --git a/tests/test_runtime/test_client_factory.py b/tests/test_runtime/test_client_factory.py index d8e20ba..9ea11d5 100644 --- a/tests/test_runtime/test_client_factory.py +++ b/tests/test_runtime/test_client_factory.py @@ -9,7 +9,8 @@ from plyngent.config.models import ( OpenAIProvider, ) from plyngent.lmproto.deepseek import DeepseekOpenAIClient -from plyngent.lmproto.openai_compatible import OpenAIClient +from plyngent.lmproto.openai import OpenAIClient +from plyngent.lmproto.openai_compatible import OpenAICompatibleClient from plyngent.runtime import ProviderNotSupportedError, create_client, provider_to_openai_config @@ -20,6 +21,7 @@ def test_openai_provider_defaults_base_url() -> None: assert config.base_url == "https://api.openai.com/v1" client = create_client(provider) assert isinstance(client, OpenAIClient) + assert hasattr(client, "responses") def test_openai_compatible_requires_url() -> None: @@ -34,7 +36,8 @@ def test_openai_compatible_client() -> None: url="https://example.com/v1", ) client = create_client(provider) - assert isinstance(client, OpenAIClient) + assert isinstance(client, OpenAICompatibleClient) + assert not isinstance(client, OpenAIClient) assert provider_to_openai_config(provider).base_url == "https://example.com/v1"