From 5dd102a016c32436145b8317201c79f651c475e5 Mon Sep 17 00:00:00 2001 From: worldmozara Date: Wed, 15 Jul 2026 20:59:47 +0800 Subject: [PATCH] core/agent: route OpenAI provider through Responses API ResponsesChatClient adapts chat-shaped agent history/tools to POST /responses; compat and DeepSeek unchanged. --- CLAUDE.md | 3 +- src/plyngent/agent/__init__.py | 2 + src/plyngent/agent/responses_bridge.py | 295 ++++++++++++++++++++++ src/plyngent/agent/responses_client.py | 106 ++++++++ src/plyngent/agent/usage.py | 10 +- src/plyngent/runtime/client_factory.py | 13 +- tests/test_agent/test_responses_bridge.py | 157 ++++++++++++ tests/test_agent/test_usage.py | 8 + tests/test_runtime/test_client_factory.py | 6 +- 9 files changed, 594 insertions(+), 6 deletions(-) create mode 100644 src/plyngent/agent/responses_bridge.py create mode 100644 src/plyngent/agent/responses_client.py create mode 100644 tests/test_agent/test_responses_bridge.py diff --git a/CLAUDE.md b/CLAUDE.md index d9963a8..88c7512 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -55,7 +55,8 @@ Async SQLAlchemy + aiosqlite. `MemoryStore`: schema init (+ lightweight SQLite ` ### Agent (`agent/`) -- **`ChatClient`** Protocol for `chat_completions`. +- **`ChatClient`** Protocol for `chat_completions` (agent history stays chat-shaped). +- **OpenAI Responses integration**: `ResponsesChatClient` adapts platform `OpenAIClient.responses` to `ChatClient`; selected automatically for `OpenAIProvider` in `create_client`. Compat/DeepSeek stay on chat completions. - **`@tool` / `ToolRegistry`**: decorator infers JSON Schema from type hints; execute tools by name. - **`run_chat_loop`**: multi-round tool loop; default **streaming** text deltas + stream tool-call merge; parallel tools; tool-result char budget; soft context compact on request (**API-calibrated** after first usage when available); cooperative cancel points; optional `on_limit`. - **`ChatAgent`**: optional `MemoryStore` (user message persisted immediately; **completed tool batches checkpointed** mid-turn; unfinished assistant suffix rolled back on failure); `stream`; system prompt; `retry()` continues incomplete turns (user-only **or** after committed tools — does not re-run those tools). diff --git a/src/plyngent/agent/__init__.py b/src/plyngent/agent/__init__.py index b12c0ac..21db963 100644 --- a/src/plyngent/agent/__init__.py +++ b/src/plyngent/agent/__init__.py @@ -2,6 +2,8 @@ from .chat import ChatAgent as ChatAgent from .client import ChatClient as ChatClient from .compact import build_compacted_seed_messages as build_compacted_seed_messages from .compact import summarize_messages as summarize_messages +from .responses_client import ResponsesChatClient as ResponsesChatClient +from .responses_client import wrap_openai_for_agent as wrap_openai_for_agent from .events import AgentEvent as AgentEvent from .events import AssistantMessageEvent as AssistantMessageEvent from .events import CancelledEvent as CancelledEvent diff --git a/src/plyngent/agent/responses_bridge.py b/src/plyngent/agent/responses_bridge.py new file mode 100644 index 0000000..e56503d --- /dev/null +++ b/src/plyngent/agent/responses_bridge.py @@ -0,0 +1,295 @@ +"""Convert agent chat history/tools to OpenAI Responses API shapes and back. + +Agent memory and events stay chat-completions-shaped; only the transport uses +Responses. DeepSeek / openai-compatible paths never enter this module. +""" + +from __future__ import annotations + +from typing import TYPE_CHECKING, Any, cast + +from msgspec import UNSET + +from plyngent.lmproto.openai.model import ( + Response, + ResponseEasyInputMessage, + ResponseFunctionTool, + ResponseFunctionToolCallOutput, + response_function_calls, + response_output_text, +) +from plyngent.lmproto.openai_compatible.model import ( + AnyAssistantToolCall, + AssistantChatMessage, + AssistantFunctionTool, + AssistantFunctionToolCall, + ChatCompletionChoice, + ChatCompletionChunk, + ChatCompletionResponse, + ChatCompletionsParam, + ChunkChoice, + DeltaMessage, + StreamFunctionDelta, + StreamToolCallDelta, + SystemChatMessage, + ToolChatMessage, + ToolFunctionItem, + UserChatMessage, +) + +if TYPE_CHECKING: + from collections.abc import Sequence + + from plyngent.lmproto.openai_compatible.model import AnyChatMessage, AnyToolItem + from plyngent.typedef import Unset + + +def tool_items_to_response_tools( + tools: Sequence[AnyToolItem] | None, +) -> list[ResponseFunctionTool]: + """Map chat ``ToolFunctionItem`` list to flat Responses function tools.""" + if not tools: + return [] + result: list[ResponseFunctionTool] = [] + for item in tools: + if not isinstance(item, ToolFunctionItem): + continue + fn = item.function + result.append( + ResponseFunctionTool( + name=fn.name, + description=fn.description if fn.description is not UNSET else UNSET, + parameters=fn.parameters if fn.parameters is not UNSET else UNSET, + strict=fn.strict if fn.strict is not UNSET else UNSET, + ) + ) + return result + + +def _assistant_to_input_items( + message: AssistantChatMessage, +) -> list[dict[str, Any] | ResponseEasyInputMessage]: + items: list[dict[str, Any] | ResponseEasyInputMessage] = [] + if message.tool_calls is not UNSET and message.tool_calls: + items.extend( + { + "type": "function_call", + "call_id": call.id, + "name": call.function.name, + "arguments": call.function.arguments, + } + for call in message.tool_calls + if isinstance(call, AssistantFunctionToolCall) + ) + if isinstance(message.content, str) and message.content: + items.append(ResponseEasyInputMessage(role="assistant", content=message.content)) + return items + + +def chat_messages_to_responses_input( + messages: Sequence[AnyChatMessage], +) -> tuple[str | None, list[dict[str, Any] | ResponseEasyInputMessage | ResponseFunctionToolCallOutput]]: + """Split system prompts into ``instructions``; rest become Responses ``input`` items.""" + instructions_parts: list[str] = [] + items: list[dict[str, Any] | ResponseEasyInputMessage | ResponseFunctionToolCallOutput] = [] + + for message in messages: + if isinstance(message, SystemChatMessage): + if message.content.strip(): + instructions_parts.append(message.content) + elif isinstance(message, UserChatMessage): + items.append(ResponseEasyInputMessage(role="user", content=message.content)) + elif isinstance(message, AssistantChatMessage): + items.extend(_assistant_to_input_items(message)) + elif isinstance(message, ToolChatMessage): + items.append( + ResponseFunctionToolCallOutput( + call_id=message.tool_call_id, + output=message.content, + ) + ) + else: + content = getattr(message, "content", None) + if isinstance(content, str) and content: + items.append(ResponseEasyInputMessage(role="user", content=content)) + + instructions = "\n\n".join(instructions_parts) if instructions_parts else None + return instructions, items + + +def response_to_assistant_message(response: Response) -> AssistantChatMessage: + """Map a completed Responses object to agent ``AssistantChatMessage``.""" + text = response_output_text(response) + calls = response_function_calls(response) + tool_calls: list[AnyAssistantToolCall] | Unset = UNSET + if calls: + tool_calls = [ + AssistantFunctionToolCall( + id=call.call_id, + function=AssistantFunctionTool(name=call.name, arguments=call.arguments), + ) + for call in calls + ] + reasoning = _reasoning_summary_text(response) + return AssistantChatMessage( + content=text or None, + tool_calls=tool_calls, + reasoning_content=reasoning or UNSET, + ) + + +def _reasoning_summary_text(response: Response) -> str: + parts: list[str] = [] + for raw in response.output: + if raw.get("type") != "reasoning": + continue + summary = raw.get("summary") + if not isinstance(summary, list): + continue + for block in summary: + if not isinstance(block, dict): + continue + block_map = cast("dict[str, object]", block) + if block_map.get("type") in {"summary_text", "output_text"}: + text = block_map.get("text") + if isinstance(text, str) and text: + parts.append(text) + return "".join(parts) + + +def response_to_chat_completion(response: Response) -> ChatCompletionResponse: + """Wrap Responses result as a synthetic chat completion for the agent loop.""" + assistant = response_to_assistant_message(response) + finish: str | None = "tool_calls" if assistant.tool_calls is not UNSET else "stop" + usage = response.usage if response.usage is not UNSET else UNSET + created = int(response.created_at) + return ChatCompletionResponse( + id=response.id, + object="chat.completion", + created=created, + model=response.model, + choices=[ + ChatCompletionChoice( + index=0, + message=assistant, + finish_reason=cast("Any", finish), + ) + ], + usage=cast("Any", usage) if usage is not UNSET else UNSET, + ) + + +def text_delta_chunk(*, model: str, content: str, created: int = 0) -> ChatCompletionChunk: + return ChatCompletionChunk( + id="resp-stream", + object="chat.completion.chunk", + created=created, + model=model, + choices=[ + ChunkChoice( + index=0, + delta=DeltaMessage(content=content), + ) + ], + ) + + +def reasoning_delta_chunk(*, model: str, content: str, created: int = 0) -> ChatCompletionChunk: + return ChatCompletionChunk( + id="resp-stream", + object="chat.completion.chunk", + created=created, + model=model, + choices=[ + ChunkChoice( + index=0, + delta=DeltaMessage(reasoning_content=content), + ) + ], + ) + + +def tool_call_chunks_from_response( + response: Response, + *, + model: str, + created: int = 0, +) -> list[ChatCompletionChunk]: + """Emit complete tool-call stream deltas (one chunk per call) for loop merge.""" + calls = response_function_calls(response) + chunks: list[ChatCompletionChunk] = [] + for index, call in enumerate(calls): + chunks.append( + ChatCompletionChunk( + id=response.id, + object="chat.completion.chunk", + created=created, + model=model, + choices=[ + ChunkChoice( + index=0, + delta=DeltaMessage( + tool_calls=[ + StreamToolCallDelta( + index=index, + id=call.call_id, + type="function", + function=StreamFunctionDelta( + name=call.name, + arguments=call.arguments, + ), + ) + ] + ), + ) + ], + ) + ) + return chunks + + +def usage_chunk_from_response(response: Response, *, model: str) -> ChatCompletionChunk | None: + if response.usage is UNSET or response.usage is None: + return None + created = int(response.created_at) + return ChatCompletionChunk( + id=response.id, + object="chat.completion.chunk", + created=created, + model=model, + choices=[], + usage=cast("dict[str, Any]", response.usage), + ) + + +def chat_param_to_responses_kwargs(param: ChatCompletionsParam) -> dict[str, Any]: + """Build keyword args for :class:`ResponsesCreateParam` from a chat param.""" + instructions, input_items = chat_messages_to_responses_input(param.messages) + tools = tool_items_to_response_tools( + param.tools if param.tools is not UNSET else None, + ) + kwargs: dict[str, Any] = { + "model": param.model, + "input": input_items or "", + "store": False, + } + if instructions: + kwargs["instructions"] = instructions + if tools: + kwargs["tools"] = tools + if param.temperature is not UNSET: + kwargs["temperature"] = param.temperature + if param.top_p is not UNSET: + kwargs["top_p"] = param.top_p + if param.max_completion_tokens is not UNSET: + kwargs["max_output_tokens"] = param.max_completion_tokens + elif param.max_tokens is not UNSET: + kwargs["max_output_tokens"] = param.max_tokens + if param.parallel_tool_calls is not UNSET: + kwargs["parallel_tool_calls"] = param.parallel_tool_calls + if param.tool_choice is not UNSET: + # "auto"/"none"/"required" strings pass through; structured choices left as-is if str. + choice = param.tool_choice + if isinstance(choice, str): + kwargs["tool_choice"] = choice + return kwargs diff --git a/src/plyngent/agent/responses_client.py b/src/plyngent/agent/responses_client.py new file mode 100644 index 0000000..bf49955 --- /dev/null +++ b/src/plyngent/agent/responses_client.py @@ -0,0 +1,106 @@ +"""ChatClient adapter: agent chat loop over OpenAI Responses API.""" + +from __future__ import annotations + +from typing import TYPE_CHECKING, Literal, overload + +from msgspec import UNSET + +from plyngent.agent.responses_bridge import ( + chat_param_to_responses_kwargs, + reasoning_delta_chunk, + response_to_chat_completion, + text_delta_chunk, + tool_call_chunks_from_response, + usage_chunk_from_response, +) +from plyngent.lmproto.openai.model import ResponsesCreateParam + +if TYPE_CHECKING: + from collections.abc import AsyncIterator + + from plyngent.lmproto.openai.client import OpenAIClient + from plyngent.lmproto.openai.model import Response + from plyngent.lmproto.openai_compatible.model import ( + ChatCompletionChunk, + ChatCompletionResponse, + ChatCompletionsParam, + ) + + +class ResponsesChatClient: + """Present OpenAI Responses as :class:`~plyngent.agent.client.ChatClient`. + + History and tool results remain chat-completions-shaped; only the HTTP call + uses ``POST /responses``. + """ + + def __init__(self, client: OpenAIClient) -> None: + self._client = client + + async def models(self) -> list[str]: + return await self._client.models() + + @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]: + kwargs = chat_param_to_responses_kwargs(param) + create = ResponsesCreateParam(**kwargs) + if stream: + return self._stream_as_chat_chunks(create, model=param.model) + response = await self._client.responses(create, stream=False) + return response_to_chat_completion(response) + + async def _stream_as_chat_chunks( + self, + create: ResponsesCreateParam, + *, + model: str, + ) -> AsyncIterator[ChatCompletionChunk]: + stream = await self._client.responses(create, stream=True) + final: Response | None = None + async for event in stream: + etype = event.type + if etype == "response.output_text.delta" and isinstance(event.delta, str) and event.delta: + yield text_delta_chunk(model=model, content=event.delta) + continue + if etype in { + "response.reasoning_summary_text.delta", + "response.reasoning_text.delta", + } and isinstance(event.delta, str) and event.delta: + yield reasoning_delta_chunk(model=model, content=event.delta) + continue + if etype == "response.completed" and event.response is not UNSET and isinstance( + event.response, dict + ): + # Decode full response for tools + usage + import msgspec + + from plyngent.lmproto.openai.model import Response as ResponseModel + + try: + final = msgspec.convert(event.response, ResponseModel) + except (TypeError, ValueError, msgspec.ValidationError): + final = None + + if final is not None: + for chunk in tool_call_chunks_from_response(final, model=model): + yield chunk + usage = usage_chunk_from_response(final, model=model) + if usage is not None: + yield usage + + +def wrap_openai_for_agent(client: OpenAIClient) -> ResponsesChatClient: + """Wrap a platform OpenAI client so the agent uses Responses by default.""" + return ResponsesChatClient(client) diff --git a/src/plyngent/agent/usage.py b/src/plyngent/agent/usage.py index de6be16..116c209 100644 --- a/src/plyngent/agent/usage.py +++ b/src/plyngent/agent/usage.py @@ -86,14 +86,22 @@ def _as_nonneg_int(value: object) -> int: def token_usage_from_api(usage: object) -> TokenUsage | None: - """Parse OpenAI-style usage dict; return None if missing/empty.""" + """Parse OpenAI-style usage dict; return None if missing/empty. + + Accepts chat completions fields (``prompt_tokens`` / ``completion_tokens``) + and Responses fields (``input_tokens`` / ``output_tokens``). + """ if usage is None or usage is UNSET: return None if not isinstance(usage, dict): return None raw = cast("dict[str, object]", usage) prompt = _as_nonneg_int(raw.get("prompt_tokens")) + if prompt == 0: + prompt = _as_nonneg_int(raw.get("input_tokens")) completion = _as_nonneg_int(raw.get("completion_tokens")) + if completion == 0: + completion = _as_nonneg_int(raw.get("output_tokens")) total = _as_nonneg_int(raw.get("total_tokens")) if total == 0 and (prompt or completion): total = prompt + completion diff --git a/src/plyngent/runtime/client_factory.py b/src/plyngent/runtime/client_factory.py index cbd856b..9c61814 100644 --- a/src/plyngent/runtime/client_factory.py +++ b/src/plyngent/runtime/client_factory.py @@ -15,10 +15,14 @@ from plyngent.lmproto.openai_compatible import OpenAICompatibleClient, OpenAICon if TYPE_CHECKING: from collections.abc import Mapping + from plyngent.agent.responses_client import ResponsesChatClient + DEFAULT_OPENAI_BASE_URL = "https://api.openai.com/v1" DEFAULT_DEEPSEEK_BASE_URL = "https://api.deepseek.com" -type ProtocolClient = OpenAIClient | OpenAICompatibleClient | DeepseekOpenAIClient +type ProtocolClient = ( + OpenAIClient | OpenAICompatibleClient | DeepseekOpenAIClient | ResponsesChatClient +) # Backward-compatible name used by older imports/tests. type OpenAICompatibleClientUnion = ProtocolClient @@ -51,12 +55,17 @@ def _deepseek_convention(extras: Mapping[str, str]) -> str: def create_client(provider: Provider) -> ProtocolClient: """Build a protocol client for the given provider config entry. + OpenAI platform providers are wrapped so the agent uses the Responses API + while still exposing a chat-completions-shaped interface. + Raises: ProviderNotSupportedError: When the provider preset (or DeepSeek convention) has no implemented client yet. """ if isinstance(provider, OpenAIProvider): - return OpenAIClient(provider_to_openai_config(provider)) + from plyngent.agent.responses_client import wrap_openai_for_agent + + return wrap_openai_for_agent(OpenAIClient(provider_to_openai_config(provider))) if isinstance(provider, OpenAICompatibleProvider): return OpenAICompatibleClient(provider_to_openai_config(provider)) if isinstance(provider, DeepseekProvider): diff --git a/tests/test_agent/test_responses_bridge.py b/tests/test_agent/test_responses_bridge.py new file mode 100644 index 0000000..e5a65ad --- /dev/null +++ b/tests/test_agent/test_responses_bridge.py @@ -0,0 +1,157 @@ +from __future__ import annotations + +import msgspec +import pytest + +from plyngent.agent.responses_bridge import ( + chat_messages_to_responses_input, + chat_param_to_responses_kwargs, + response_to_assistant_message, + response_to_chat_completion, + tool_items_to_response_tools, +) +from plyngent.agent.responses_client import ResponsesChatClient +from plyngent.lmproto.openai.model import Response, ResponsesCreateParam +from plyngent.lmproto.openai_compatible.model import ( + AssistantChatMessage, + AssistantFunctionTool, + AssistantFunctionToolCall, + ChatCompletionsParam, + SystemChatMessage, + ToolChatMessage, + ToolFunction, + ToolFunctionItem, + UserChatMessage, +) + + +def test_tool_items_to_response_tools() -> None: + items = [ + ToolFunctionItem( + function=ToolFunction( + name="read_file", + description="Read a file", + parameters={"type": "object", "properties": {"path": {"type": "string"}}}, + ) + ) + ] + tools = tool_items_to_response_tools(items) + assert len(tools) == 1 + assert tools[0].name == "read_file" + + +def test_chat_messages_to_input_and_instructions() -> None: + messages = [ + SystemChatMessage(content="You are helpful."), + UserChatMessage(content="hi"), + AssistantChatMessage( + content=None, + tool_calls=[ + AssistantFunctionToolCall( + id="call_1", + function=AssistantFunctionTool(name="read_file", arguments='{"path":"a"}'), + ) + ], + ), + ToolChatMessage(content="file body", tool_call_id="call_1"), + ] + instructions, items = chat_messages_to_responses_input(messages) + assert instructions == "You are helpful." + assert len(items) == 3 # user, function_call, function_call_output + + +def test_response_to_assistant_with_tools() -> None: + raw = { + "id": "resp_1", + "object": "response", + "created_at": 1, + "model": "gpt-test", + "status": "completed", + "output": [ + { + "id": "msg_1", + "type": "message", + "role": "assistant", + "status": "completed", + "content": [{"type": "output_text", "text": "done", "annotations": []}], + }, + { + "type": "function_call", + "call_id": "call_9", + "name": "add", + "arguments": '{"a":1}', + "status": "completed", + }, + ], + "usage": {"input_tokens": 10, "output_tokens": 5, "total_tokens": 15}, + } + response = msgspec.convert(raw, Response) + from msgspec import UNSET + + assistant = response_to_assistant_message(response) + assert assistant.content == "done" + assert assistant.tool_calls is not UNSET + assert isinstance(assistant.tool_calls, list) + call0 = assistant.tool_calls[0] + assert isinstance(call0, AssistantFunctionToolCall) + assert call0.id == "call_9" + assert call0.function.name == "add" + completion = response_to_chat_completion(response) + assert completion.choices[0].message.content == "done" + assert isinstance(completion.usage, dict) + assert completion.usage["input_tokens"] == 10 + + +def test_chat_param_to_responses_kwargs() -> None: + param = ChatCompletionsParam( + model="gpt-test", + messages=[SystemChatMessage(content="sys"), UserChatMessage(content="hi")], + tools=[ToolFunctionItem(function=ToolFunction(name="t", parameters={"type": "object"}))], + temperature=0.2, + ) + kwargs = chat_param_to_responses_kwargs(param) + assert kwargs["model"] == "gpt-test" + assert kwargs["instructions"] == "sys" + assert kwargs["store"] is False + assert kwargs["temperature"] == 0.2 + assert len(kwargs["tools"]) == 1 + + +@pytest.mark.asyncio +async def test_responses_chat_client_non_stream(monkeypatch: pytest.MonkeyPatch) -> None: + from plyngent.lmproto.openai.client import OpenAIClient + from plyngent.lmproto.openai_compatible.config import OpenAIConfig + + platform = OpenAIClient(OpenAIConfig(access_key_or_token="sk", base_url="https://example/v1")) + + body = { + "id": "resp_x", + "object": "response", + "created_at": 1, + "model": "gpt-test", + "status": "completed", + "output": [ + { + "id": "msg_1", + "type": "message", + "role": "assistant", + "status": "completed", + "content": [{"type": "output_text", "text": "hello", "annotations": []}], + } + ], + "usage": {"input_tokens": 3, "output_tokens": 1, "total_tokens": 4}, + } + + async def fake_responses(param: ResponsesCreateParam, *, stream: bool = False): + assert stream is False + assert param.model == "gpt-test" + assert param.store is False + return msgspec.convert(body, Response) + + monkeypatch.setattr(platform, "responses", fake_responses) + client = ResponsesChatClient(platform) + result = await client.chat_completions( + ChatCompletionsParam(model="gpt-test", messages=[UserChatMessage(content="hi")]), + stream=False, + ) + assert result.choices[0].message.content == "hello" diff --git a/tests/test_agent/test_usage.py b/tests/test_agent/test_usage.py index c2e276b..54bdf82 100644 --- a/tests/test_agent/test_usage.py +++ b/tests/test_agent/test_usage.py @@ -50,6 +50,14 @@ def test_token_usage_from_api_infers_total() -> None: assert u.total_tokens == 5 +def test_token_usage_from_api_responses_fields() -> None: + u = token_usage_from_api({"input_tokens": 7, "output_tokens": 2, "total_tokens": 9}) + assert u is not None + assert u.prompt_tokens == 7 + assert u.completion_tokens == 2 + assert u.total_tokens == 9 + + def test_chars_to_tokens() -> None: assert chars_to_tokens(0) == 0 assert chars_to_tokens(1) == 1 diff --git a/tests/test_runtime/test_client_factory.py b/tests/test_runtime/test_client_factory.py index 9ea11d5..c7e84ea 100644 --- a/tests/test_runtime/test_client_factory.py +++ b/tests/test_runtime/test_client_factory.py @@ -19,9 +19,11 @@ def test_openai_provider_defaults_base_url() -> None: config = provider_to_openai_config(provider) assert config.access_key_or_token == "sk-test" assert config.base_url == "https://api.openai.com/v1" + from plyngent.agent.responses_client import ResponsesChatClient + client = create_client(provider) - assert isinstance(client, OpenAIClient) - assert hasattr(client, "responses") + assert isinstance(client, ResponsesChatClient) + assert hasattr(client, "chat_completions") def test_openai_compatible_requires_url() -> None: