diff --git a/dare_framework/agent/base_agent.py b/dare_framework/agent/base_agent.py index f4cc12fd..b6e6f04d 100644 --- a/dare_framework/agent/base_agent.py +++ b/dare_framework/agent/base_agent.py @@ -5,6 +5,7 @@ from abc import ABC, abstractmethod import asyncio import contextlib +from contextvars import ContextVar from dataclasses import replace import logging from typing import TYPE_CHECKING, Any @@ -237,10 +238,20 @@ async def _execute_polled_message( envelope_id: str | None, ) -> None: """Execute one polled message and send response envelope through channel.""" - result = await self.execute(task, transport=channel) + # Keep transport-loop execution state task-local so concurrent execute() calls on + # the same agent instance do not leak loop state across tasks. + transport_loop_token = _TRANSPORT_LOOP_EXECUTION_CTX.set(True) + try: + result = await self.execute(task, transport=channel) + finally: + _TRANSPORT_LOOP_EXECUTION_CTX.reset(transport_loop_token) result = self._with_normalized_output_text(result) await self._send_transport_result(result, task=task, transport=channel, reply_to=envelope_id) + def _is_transport_loop_execution(self, *, transport: AgentChannel | None) -> bool: + """Return whether execute() is currently running under the transport loop.""" + return bool(transport is not None and _TRANSPORT_LOOP_EXECUTION_CTX.get()) + def _with_normalized_output_text(self, result: RunResult) -> RunResult: """Ensure RunResult.output_text is filled for downstream consumers.""" if result.output_text is not None: @@ -355,6 +366,7 @@ def get_agent_control_handler(self) -> None: _NO_OP_AGENT_CHANNEL = _NoOpAgentChannel() +_TRANSPORT_LOOP_EXECUTION_CTX: ContextVar[bool] = ContextVar("_transport_loop_execution", default=False) def _coerce_polled_envelopes(polled: Any) -> list[TransportEnvelope]: diff --git a/dare_framework/agent/react_agent.py b/dare_framework/agent/react_agent.py index a320f31e..20884391 100644 --- a/dare_framework/agent/react_agent.py +++ b/dare_framework/agent/react_agent.py @@ -91,7 +91,9 @@ def _print_context_list( from dare_framework.plan.types import RunResult from dare_framework.plan.types import Task from dare_framework.tool import IToolGateway, IToolProvider +from dare_framework.transport.interaction.payloads import build_error_payload, build_success_payload from dare_framework.transport.kernel import AgentChannel +from dare_framework.transport.types import EnvelopeKind, TransportEnvelope, TransportEventType, new_envelope_id class ReactAgent(BaseAgent): @@ -151,7 +153,6 @@ async def _execute_basic( transport: AgentChannel | None = None, ) -> RunResult: """原始基础 ReAct 循环实现。""" - _ = transport task_description = task.description if isinstance(task, Task) else task user_message = Message(role="user", content=task_description) self._context.stm_add(user_message) @@ -204,6 +205,14 @@ async def _execute_basic( latest_usage = usage n_tools = len(response.tool_calls) if response.tool_calls else 0 print(f"[{self.name}] 模型返回, tool_calls={n_tools}", flush=True) + thinking_content = (response.thinking_content or "").strip() + if thinking_content: + await self._emit_transport_success( + transport=transport, + event_type=TransportEventType.THINKING.value, + target="model", + resp={"output": thinking_content}, + ) if usage is not None: tokens = _usage_total_tokens(usage) @@ -217,6 +226,10 @@ async def _execute_basic( final_text = "模型未返回可显示的文本回复。请重试,或明确要求先调用 ask_user 再继续。" assistant_message = Message(role="assistant", content=final_text) self._context.stm_add(assistant_message) + await self._emit_terminal_transport_message( + transport=transport, + output=final_text, + ) output = build_output_envelope( final_text, usage=latest_usage, @@ -237,6 +250,10 @@ async def _execute_basic( if repeated_tool_rounds >= 3: loop_guard = "模型连续重复调用相同工具,已停止自动循环。请换一种描述,或明确要求先调用 ask_user 再继续。" self._context.stm_add(Message(role="assistant", content=loop_guard)) + await self._emit_terminal_transport_message( + transport=transport, + output=loop_guard, + ) output = build_output_envelope(loop_guard, usage=latest_usage) return RunResult(success=True, output=output, output_text=output["content"]) @@ -252,6 +269,16 @@ async def _execute_basic( name = tool_call.get("name", "") tool_call_id = tool_call.get("id", "") params = _normalize_tool_args(tool_call.get("arguments", {})) + await self._emit_transport_success( + transport=transport, + event_type=TransportEventType.TOOL_CALL.value, + target=name or "tool_call", + resp={ + "id": tool_call_id, + "name": name, + "arguments": params, + }, + ) # 打印工具调用信息:名称、参数 params_str = json.dumps(params, ensure_ascii=False) params_preview = params_str[:300] + ("..." if len(params_str) > 300 else "") @@ -260,6 +287,12 @@ async def _execute_basic( try: result = await gateway.invoke(name, envelope=envelope, **params) except Exception as exc: + await self._emit_transport_error( + transport=transport, + target=name or "tool_call", + code="TOOL_CALL_FAILED", + reason=str(exc).strip() or exc.__class__.__name__, + ) result = type("R", (), {"success": False, "output": {}, "error": str(exc)})() success = getattr(result, "success", False) @@ -268,6 +301,18 @@ async def _execute_basic( out_preview = _preview_output(output) print(f"[{self.name}] 工具结果: {name} | success={success} | {out_preview}", flush=True) error = getattr(result, "error", "") or "" + await self._emit_transport_success( + transport=transport, + event_type=TransportEventType.TOOL_RESULT.value, + target=name or "tool_result", + resp={ + "id": tool_call_id, + "name": name, + "success": bool(success), + "output": output, + "error": error, + }, + ) tool_content = json.dumps( {"success": success, "output": output, "error": error} if not success @@ -278,6 +323,10 @@ async def _execute_basic( self._context.stm_add(tool_msg) final_message = "模型在工具循环中未收敛(达到最大轮次)。请缩小范围,或明确要求先调用 ask_user 再继续。" + await self._emit_terminal_transport_message( + transport=transport, + output=final_message, + ) output = build_output_envelope(final_message, usage=latest_usage) return RunResult( success=True, @@ -522,6 +571,73 @@ async def _execute_with_smart_context( output_text=final_message, ) + async def _emit_terminal_transport_message( + self, + *, + transport: AgentChannel | None, + output: str, + ) -> None: + """Emit terminal MESSAGE for direct execute() transport calls.""" + # BaseAgent emits terminal RESULT envelopes in transport-loop execution. + # Avoid emitting duplicate terminal MESSAGE events in that path. + if self._is_transport_loop_execution(transport=transport): + return + await self._emit_transport_success( + transport=transport, + event_type=TransportEventType.MESSAGE.value, + target="prompt", + resp={"output": output}, + ) + + async def _emit_transport_success( + self, + *, + transport: AgentChannel | None, + event_type: str, + target: str, + resp: dict[str, Any], + ) -> None: + """Emit a canonical success payload to transport when available.""" + if transport is None: + return + envelope = TransportEnvelope( + id=new_envelope_id(), + kind=EnvelopeKind.MESSAGE, + event_type=event_type, + payload=build_success_payload(kind="message", target=target, resp=resp), + ) + try: + await transport.send(envelope) + except Exception: + self._logger.exception("react agent transport success emission failed") + + async def _emit_transport_error( + self, + *, + transport: AgentChannel | None, + target: str, + code: str, + reason: str, + ) -> None: + """Emit a canonical error payload to transport when available.""" + if transport is None: + return + envelope = TransportEnvelope( + id=new_envelope_id(), + kind=EnvelopeKind.MESSAGE, + event_type=TransportEventType.ERROR.value, + payload=build_error_payload( + kind="message", + target=target, + code=code, + reason=reason, + ), + ) + try: + await transport.send(envelope) + except Exception: + self._logger.exception("react agent transport error emission failed") + def _preview_output(output: Any, max_len: int = 120) -> str: """生成 output 的简短预览,用于日志打印。""" diff --git a/dare_framework/model/adapters/openai_adapter.py b/dare_framework/model/adapters/openai_adapter.py index 2c57f33b..530add6a 100644 --- a/dare_framework/model/adapters/openai_adapter.py +++ b/dare_framework/model/adapters/openai_adapter.py @@ -103,11 +103,13 @@ async def generate( tool_calls = self._extract_tool_calls(response) usage = self._extract_usage(response) + thinking_content = self._extract_thinking_content(response) return ModelResponse( content=response.content or "", tool_calls=tool_calls, usage=usage, + thinking_content=thinking_content, ) def _ensure_client(self) -> Any: @@ -231,11 +233,55 @@ def _extract_usage(self, response: Any) -> dict[str, Any] | None: """Extract usage information from the response.""" usage = getattr(response, "response_metadata", {}).get("token_usage") if usage: - return { + normalized = { "prompt_tokens": usage.get("prompt_tokens", 0), "completion_tokens": usage.get("completion_tokens", 0), "total_tokens": usage.get("total_tokens", 0), } + reasoning_tokens = self._extract_reasoning_tokens(usage) + if reasoning_tokens is not None: + normalized["reasoning_tokens"] = reasoning_tokens + return normalized + return None + + def _extract_reasoning_tokens(self, usage: dict[str, Any]) -> int | None: + """Extract reasoning token count from provider-specific usage payloads.""" + candidates: list[Any] = [ + usage.get("reasoning_tokens"), + usage.get("output_tokens_details", {}).get("reasoning_tokens") + if isinstance(usage.get("output_tokens_details"), dict) + else None, + usage.get("output_tokens_details", {}).get("reasoning") + if isinstance(usage.get("output_tokens_details"), dict) + else None, + usage.get("completion_tokens_details", {}).get("reasoning_tokens") + if isinstance(usage.get("completion_tokens_details"), dict) + else None, + ] + for candidate in candidates: + try: + if candidate is None: + continue + return int(candidate) + except (TypeError, ValueError): + continue + return None + + def _extract_thinking_content(self, response: Any) -> str | None: + """Extract provider reasoning text into framework-level thinking content.""" + additional_kwargs = getattr(response, "additional_kwargs", {}) + if isinstance(additional_kwargs, dict): + for key in ("reasoning_content", "reasoning", "thinking"): + content = _coerce_text(additional_kwargs.get(key)) + if content: + return content + + response_metadata = getattr(response, "response_metadata", {}) + if isinstance(response_metadata, dict): + for key in ("reasoning_content", "reasoning", "thinking"): + content = _coerce_text(response_metadata.get(key)) + if content: + return content return None def _log_client_config(self, client: Any) -> None: @@ -277,3 +323,25 @@ def _build_http_clients(self) -> tuple[Any | None, Any | None]: __all__ = ["OpenAIModelAdapter"] + + +def _coerce_text(value: Any) -> str | None: + """Coerce heterogenous provider reasoning payloads into a non-empty string.""" + if isinstance(value, str): + text = value.strip() + return text or None + if isinstance(value, dict): + for key in ("text", "content", "reasoning", "thinking"): + text = _coerce_text(value.get(key)) + if text: + return text + return None + if isinstance(value, list): + parts: list[str] = [] + for item in value: + text = _coerce_text(item) + if text: + parts.append(text) + if parts: + return "\n".join(parts) + return None diff --git a/dare_framework/model/adapters/openrouter_adapter.py b/dare_framework/model/adapters/openrouter_adapter.py index 6c401fe2..80ceb20e 100644 --- a/dare_framework/model/adapters/openrouter_adapter.py +++ b/dare_framework/model/adapters/openrouter_adapter.py @@ -90,6 +90,7 @@ async def generate( message = response.choices[0].message content = message.content or "" tool_calls = _extract_tool_calls(message) + thinking_content = _extract_thinking_content(message) usage = None if response.usage: @@ -98,11 +99,17 @@ async def generate( "completion_tokens": response.usage.completion_tokens, "total_tokens": response.usage.total_tokens, } + reasoning_tokens = _extract_reasoning_tokens(response) + if reasoning_tokens is not None: + if usage is None: + usage = {} + usage["reasoning_tokens"] = reasoning_tokens return ModelResponse( content=content, tool_calls=tool_calls, usage=usage, + thinking_content=thinking_content, metadata={ "model": self._model, "finish_reason": response.choices[0].finish_reason, @@ -220,4 +227,75 @@ def _build_async_http_client(options: dict[str, Any]) -> Any | None: return None +def _extract_thinking_content(message: Any) -> str | None: + """Extract reasoning text from OpenRouter message payload variants.""" + for attr in ("reasoning_content", "reasoning", "thinking"): + text = _coerce_text(getattr(message, attr, None)) + if text: + return text + + additional_kwargs = getattr(message, "additional_kwargs", None) + if isinstance(additional_kwargs, dict): + for key in ("reasoning_content", "reasoning", "thinking"): + text = _coerce_text(additional_kwargs.get(key)) + if text: + return text + + model_extra = getattr(message, "model_extra", None) + if isinstance(model_extra, dict): + for key in ("reasoning_content", "reasoning", "thinking"): + text = _coerce_text(model_extra.get(key)) + if text: + return text + return None + + +def _extract_reasoning_tokens(response: Any) -> int | None: + """Extract reasoning token count from OpenRouter/OpenAI-compatible usage.""" + usage = getattr(response, "usage", None) + if usage is None: + return None + candidates: list[Any] = [ + getattr(usage, "reasoning_tokens", None), + _get_nested_value(getattr(usage, "completion_tokens_details", None), "reasoning_tokens"), + _get_nested_value(getattr(usage, "output_tokens_details", None), "reasoning_tokens"), + _get_nested_value(getattr(usage, "output_tokens_details", None), "reasoning"), + ] + for candidate in candidates: + if candidate is None: + continue + try: + return int(candidate) + except (TypeError, ValueError): + continue + return None + + +def _get_nested_value(value: Any, key: str) -> Any: + if isinstance(value, dict): + return value.get(key) + return getattr(value, key, None) + + +def _coerce_text(value: Any) -> str | None: + if isinstance(value, str): + text = value.strip() + return text or None + if isinstance(value, dict): + for key in ("text", "content", "reasoning", "thinking"): + text = _coerce_text(value.get(key)) + if text: + return text + return None + if isinstance(value, list): + parts: list[str] = [] + for item in value: + text = _coerce_text(item) + if text: + parts.append(text) + if parts: + return "\n".join(parts) + return None + + __all__ = ["OpenRouterModelAdapter"] diff --git a/dare_framework/model/types.py b/dare_framework/model/types.py index 0eace042..f03a6893 100644 --- a/dare_framework/model/types.py +++ b/dare_framework/model/types.py @@ -71,6 +71,7 @@ class ModelResponse: tool_calls: list[dict[str, Any]] = field(default_factory=list) usage: dict[str, Any] | None = None metadata: dict[str, Any] = field(default_factory=dict) + thinking_content: str | None = None @dataclass(frozen=True) diff --git a/dare_framework/transport/__init__.py b/dare_framework/transport/__init__.py index 616d0568..69197374 100644 --- a/dare_framework/transport/__init__.py +++ b/dare_framework/transport/__init__.py @@ -2,6 +2,7 @@ from dare_framework.transport.interfaces import AgentChannel, ClientChannel, PollableClientChannel from dare_framework.transport.types import ( + canonicalize_transport_event_type, EnvelopeKind, TransportEventType, TransportEnvelope, @@ -29,6 +30,7 @@ "EnvelopeKind", "TransportEventType", "TransportEnvelope", + "canonicalize_transport_event_type", "normalize_transport_event_type", "new_envelope_id", "Receiver", diff --git a/dare_framework/transport/_internal/adapters.py b/dare_framework/transport/_internal/adapters.py index d0948cc8..68aa2a50 100644 --- a/dare_framework/transport/_internal/adapters.py +++ b/dare_framework/transport/_internal/adapters.py @@ -10,8 +10,8 @@ from dare_framework.transport.interaction.resource_action import ResourceAction from dare_framework.transport.kernel import ClientChannel, PollableClientChannel from dare_framework.transport.types import ( + canonicalize_transport_event_type, EnvelopeKind, - normalize_transport_event_type, Receiver, Sender, TransportEnvelope, @@ -42,7 +42,7 @@ async def recv(msg: TransportEnvelope) -> None: event_type = _resolve_transport_event_type(msg) payload = msg.payload if isinstance(payload, dict): - if event_type == TransportEventType.RESULT.value: + if event_type == TransportEventType.MESSAGE.value: kind = payload.get("kind") resp = payload.get("resp") if kind == "message": @@ -54,24 +54,14 @@ async def recv(msg: TransportEnvelope) -> None: output = resp if resp is not None else payload elif event_type == TransportEventType.ERROR.value: output = payload.get("reason") or payload.get("error") - elif event_type == TransportEventType.APPROVAL_PENDING.value: + elif event_type == TransportEventType.THINKING.value: resp = payload.get("resp") - request_id = None if isinstance(resp, dict): - request = resp.get("request") - if isinstance(request, dict): - request_id = request.get("request_id") - output = f"approval pending: request_id={request_id or '?'}" - elif event_type == TransportEventType.APPROVAL_RESOLVED.value: - resp = payload.get("resp") - request_id = None - decision = None - if isinstance(resp, dict): - request_id = resp.get("request_id") - decision = resp.get("decision") - output = f"approval resolved: request_id={request_id or '?'} decision={decision or '?'}" - elif event_type == TransportEventType.HOOK.value: - output = payload.get("event") + output = resp.get("output") or resp.get("thinking") or resp + else: + output = resp if resp is not None else payload + elif event_type == TransportEventType.STATUS.value: + output = _render_status_output(payload) else: output = payload else: @@ -300,8 +290,30 @@ def _default_deserialize(raw: Any) -> TransportEnvelope: def _resolve_transport_event_type(msg: TransportEnvelope) -> str | None: """Resolve event_type for receiver routing.""" if isinstance(msg.event_type, str): - return normalize_transport_event_type(msg.event_type) + return canonicalize_transport_event_type(msg.event_type) return None +def _render_status_output(payload: dict[str, Any]) -> Any: + """Render canonical status payload into a concise stdio output.""" + resp = payload.get("resp") + if isinstance(resp, dict): + if "phase" in resp: + return resp.get("phase") + request_id = resp.get("request_id") + decision = resp.get("decision") + if isinstance(resp.get("request"), dict): + request_id = resp["request"].get("request_id") + if request_id and decision is not None: + return f"approval resolved: request_id={request_id} decision={decision}" + if request_id: + return f"approval pending: request_id={request_id}" + return "approval update" + if "phase" in payload: + return payload.get("phase") + if "event" in payload: + return payload.get("event") + return payload + + __all__ = ["StdioClientChannel", "WebSocketClientChannel", "DirectClientChannel"] diff --git a/dare_framework/transport/types.py b/dare_framework/transport/types.py index 297cede1..19b793bd 100644 --- a/dare_framework/transport/types.py +++ b/dare_framework/transport/types.py @@ -20,8 +20,15 @@ class EnvelopeKind(StrEnum): class TransportEventType(StrEnum): """Canonical event categories carried by message envelopes.""" - RESULT = "result" + MESSAGE = "message" + TOOL_CALL = "tool_call" + TOOL_RESULT = "tool_result" + THINKING = "thinking" ERROR = "error" + STATUS = "status" + + # Legacy aliases kept for backward compatibility with existing emitters/consumers. + RESULT = "result" HOOK = "hook" APPROVAL_PENDING = "approval.pending" APPROVAL_RESOLVED = "approval.resolved" @@ -31,6 +38,15 @@ class TransportEventType(StrEnum): # Only legacy aliases that differ from canonical event_type values. "approval_pending": TransportEventType.APPROVAL_PENDING.value, "approval_resolved": TransportEventType.APPROVAL_RESOLVED.value, + "tool.result": TransportEventType.TOOL_RESULT.value, + "tool.call": TransportEventType.TOOL_CALL.value, +} + +_LEGACY_TO_CANONICAL_EVENT_TYPE_MAP: dict[str, str] = { + TransportEventType.RESULT.value: TransportEventType.MESSAGE.value, + TransportEventType.HOOK.value: TransportEventType.STATUS.value, + TransportEventType.APPROVAL_PENDING.value: TransportEventType.STATUS.value, + TransportEventType.APPROVAL_RESOLVED.value: TransportEventType.STATUS.value, } @@ -44,6 +60,14 @@ def normalize_transport_event_type(raw: str | None) -> str | None: return _LEGACY_PAYLOAD_EVENT_TYPE_MAP.get(normalized, normalized) +def canonicalize_transport_event_type(raw: str | None) -> str | None: + """Canonicalize legacy event aliases into the stable transport taxonomy.""" + normalized = normalize_transport_event_type(raw) + if normalized is None: + return None + return _LEGACY_TO_CANONICAL_EVENT_TYPE_MAP.get(normalized, normalized) + + @dataclass(frozen=True) class TransportEnvelope: """Transport envelope for agent/client messages.""" @@ -95,6 +119,7 @@ def new_envelope_id() -> str: "EnvelopeKind", "TransportEventType", "normalize_transport_event_type", + "canonicalize_transport_event_type", "TransportEnvelope", "new_envelope_id", "Sender", diff --git a/docs/features/agentscope-d2-d4-thinking-transport.md b/docs/features/agentscope-d2-d4-thinking-transport.md new file mode 100644 index 00000000..76c4f778 --- /dev/null +++ b/docs/features/agentscope-d2-d4-thinking-transport.md @@ -0,0 +1,64 @@ +--- +change_ids: ["agentscope-d2-d4-thinking-transport"] +doc_kind: feature +topics: ["agentscope", "transport", "thinking", "tool-events", "model-response"] +created: 2026-03-02 +updated: 2026-03-02 +status: draft +mode: openspec +--- + +# Feature: agentscope-d2-d4-thinking-transport + +## Scope +补齐 AgentScope 迁移中最先阻断执行闭环的 D2 + D4 能力:统一 transport 中间态事件协议,并在模型与执行循环中保留/输出 thinking 与 reasoning usage 信息。 + +## OpenSpec Artifacts +- Proposal: `openspec/changes/agentscope-d2-d4-thinking-transport/proposal.md` +- Design: `openspec/changes/agentscope-d2-d4-thinking-transport/design.md` +- Specs: + - `openspec/changes/agentscope-d2-d4-thinking-transport/specs/agentscope-thinking-transport/spec.md` + - `openspec/changes/agentscope-d2-d4-thinking-transport/specs/transport-channel/spec.md` + - `openspec/changes/agentscope-d2-d4-thinking-transport/specs/chat-runtime/spec.md` +- Tasks: `openspec/changes/agentscope-d2-d4-thinking-transport/tasks.md` + +## Progress +- 已完成:D2/D4 代码实现(transport canonical 事件、ModelResponse thinking、OpenAI/OpenRouter reasoning 提取、ReAct 中间态事件发射)。 +- 已完成:定向测试与全量回归,OpenSpec tasks 全部打勾。 +- 待完成:提交评审与合并门禁记录补充。 + +## Evidence + +### Commands +- `openspec new change "agentscope-d2-d4-thinking-transport"` +- `openspec status --change agentscope-d2-d4-thinking-transport --json` +- `openspec instructions apply --change "agentscope-d2-d4-thinking-transport" --json` +- `openspec validate --changes "agentscope-d2-d4-thinking-transport"` +- `/Users/lang/workspace/github/Deterministic-Agent-Runtime-Engine/.venv/bin/pytest -q tests/unit/test_transport_types.py tests/unit/test_transport_adapters.py tests/unit/test_openrouter_adapter.py tests/unit/test_openai_model_adapter.py tests/unit/test_react_agent_gateway_injection.py` +- `/Users/lang/workspace/github/Deterministic-Agent-Runtime-Engine/.venv/bin/pytest -q tests/unit/test_transport_channel.py tests/unit/test_base_agent_transport_contract.py tests/unit/test_agent_event_transport_hook.py tests/unit/test_dare_agent_hook_transport_boundary.py tests/unit/test_example_10_agentscope_compat.py` +- `/Users/lang/workspace/github/Deterministic-Agent-Runtime-Engine/.venv/bin/pytest -q` + +### Results +- OpenSpec change 创建成功(schema: spec-driven)。 +- OpenSpec apply status:`13/13 tasks complete`(state=`all_done`)。 +- OpenSpec validate:`9 passed, 0 failed`(包含本 change)。 +- 新增定向红绿回归:`34 passed, 1 warning`。 +- 受影响面回归:`33 passed, 1 warning`。 +- 全量回归:`528 passed, 12 skipped, 1 warning`。 + +### Behavior Verification +- Happy path: + - `ReactAgent` 在有 `thinking_content + tool_calls` 的轮次按序发射 `thinking -> tool_call -> tool_result -> message`; + - `ModelResponse` 保留 `thinking_content`,OpenAI/OpenRouter adapter 将 `reasoning_tokens` 归一化到 `usage`。 +- Error branch: + - transport 对无效 message/control payload 继续返回结构化错误载荷(`code/reason/resp`); + - 工具调用异常时,ReAct transport 发射 `error` 事件并保留 `tool_result` 失败信息。 + +### Risks and Rollback +- 风险:transport 事件枚举扩展可能影响历史消费者。 +- 风险:不同 provider 的 thinking 字段解析口径不一致。 +- 回滚:保留 legacy alias 归一化;必要时关闭 ReAct 中间态事件发射并回退到旧 `result/hook` 消费路径。 + +### Review and Merge Gate Links +- Review request: 待创建 PR 后补充链接。 +- Merge gate: 待 CI + reviewer 通过后补充。 diff --git a/docs/todos/agentscope_domain_execution_todos.md b/docs/todos/agentscope_domain_execution_todos.md index b5e6b14b..d2f37c9a 100644 --- a/docs/todos/agentscope_domain_execution_todos.md +++ b/docs/todos/agentscope_domain_execution_todos.md @@ -21,7 +21,7 @@ | Claim ID | TODO Scope | Owner | Status | Declared At | Expires At | OpenSpec Change | Notes | |---|---|---|---|---|---|---|---| -| CLM-20260302-D2D4 | D2-1~D2-4, D4-1~D4-4 | mindfn | planned | 2026-03-02 | 2026-03-09 | `agentscope-d2-d4-thinking-transport` | 先处理 P0/P1 的 thinking 与 transport 协议统一。 | +| CLM-20260302-D2D4 | D2-1~D2-4, D4-1~D4-4 | mindfn | active | 2026-03-02 | 2026-03-09 | `agentscope-d2-d4-thinking-transport` | 先处理 P0/P1 的 thinking 与 transport 协议统一。 | | CLM-20260302-D5 | D5-1~D5-4 | mindfn | planned | 2026-03-02 | 2026-03-09 | `agentscope-d5-safe-compression` | 压缩链路:tool pair safe + token-aware + auto trigger。 | | CLM-20260302-D7 | D7-1~D7-4 | mindfn | planned | 2026-03-02 | 2026-03-09 | `agentscope-d7-plan-state-tools` | plan 状态机与 finish/revise 原生工具补齐。 | | CLM-20260302-D1D3 | D1-1~D1-4, D3-1~D3-4 | mindfn | planned | 2026-03-02 | 2026-03-09 | `agentscope-d1-d3-message-pipeline` | 多模态输入 schema 与 assemble normalize。 | diff --git a/docs/todos/project_overall_todos.md b/docs/todos/project_overall_todos.md index 71095156..6dedf554 100644 --- a/docs/todos/project_overall_todos.md +++ b/docs/todos/project_overall_todos.md @@ -15,7 +15,7 @@ | Claim ID | TODO Scope | Owner | Status | Declared At | Expires At | OpenSpec Change | Notes | |---|---|---|---|---|---|---|---| -| CLM-20260302-AG1 | T5-2 | mindfn | planned | 2026-03-02 | 2026-03-09 | `agentscope-d2-d4-thinking-transport` | 对齐 D2/D4:thinking + transport 事件链路。 | +| CLM-20260302-AG1 | T5-2 | mindfn | active | 2026-03-02 | 2026-03-09 | `agentscope-d2-d4-thinking-transport` | 对齐 D2/D4:thinking + transport 事件链路。 | | CLM-20260302-AG2 | T2-1 | mindfn | planned | 2026-03-02 | 2026-03-09 | `agentscope-d5-safe-compression` | 对齐 D5:安全压缩与预算收敛。 | | CLM-20260302-AG3 | D7-1~D7-4(关联 T5-5) | mindfn | planned | 2026-03-02 | 2026-03-09 | `agentscope-d7-plan-state-tools` | 先按 AgentScope gap 切片推进 plan 状态机能力。 | | CLM-20260302-AG4 | T5-3 | mindfn | planned | 2026-03-02 | 2026-03-09 | `agentscope-d1-d3-message-pipeline` | 对齐 D1/D3:多模态输入 schema + normalize。 | diff --git a/openspec/changes/agentscope-d2-d4-thinking-transport/.openspec.yaml b/openspec/changes/agentscope-d2-d4-thinking-transport/.openspec.yaml new file mode 100644 index 00000000..fd79bfc5 --- /dev/null +++ b/openspec/changes/agentscope-d2-d4-thinking-transport/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-03-02 diff --git a/openspec/changes/agentscope-d2-d4-thinking-transport/design.md b/openspec/changes/agentscope-d2-d4-thinking-transport/design.md new file mode 100644 index 00000000..31bc1d6e --- /dev/null +++ b/openspec/changes/agentscope-d2-d4-thinking-transport/design.md @@ -0,0 +1,61 @@ +## Context + +当前主干中,模型响应类型仅保留 `content/tool_calls/usage`,无法承载推理内容;transport 事件类型仍以 `result/error/hook` 为主,难以表达 AgentScope 语义中的中间态事件。已有示例(example 10)通过 compat shim 模拟了该能力,但框架本体未提供同等契约。 + +## Goals / Non-Goals + +**Goals:** +- 在框架本体补齐 D2 + D4 需要的最小可用契约。 +- 统一 transport 中间态事件分类与错误 payload。 +- 在 adapter 层保留并规范化 `thinking_content` 与 `reasoning_tokens`。 +- 在 ReAct 执行循环中发射可观测中间态事件,支持端到端消费。 + +**Non-Goals:** +- 不在本切片实现 streaming 生成接口(`generate_stream`)。 +- 不在本切片改造 session 持久化协议(S1/S2)。 +- 不引入新 provider adapter。 + +## Decisions + +1. **事件语义分层** +- Decision: 保留 `TransportEnvelope.event_type` 为字符串字段,但增加一组 canonical 枚举值覆盖 `message/tool_call/tool_result/thinking/error/status`,并维护 legacy alias 到 canonical 的归一化映射。 +- Rationale: 最小化对现有 transport 调度器的侵入,同时让调用侧获得稳定语义。 +- Alternative considered: 直接替换现有 `RESULT/HOOK` 枚举并删除兼容映射;该方案会扩大回归面,暂不采用。 + +2. **模型响应扩展策略** +- Decision: 在 `ModelResponse` 新增可选 `thinking_content` 字段;`usage` 继续使用 dict,但要求 adapter 将 `reasoning_tokens` 规范化写入。 +- Rationale: 维持现有调用接口,避免一次性引入新的 usage 强类型模型。 +- Alternative considered: 引入全新 `Usage` dataclass;当前改动成本与兼容风险较高。 + +3. **中间态事件发射位置** +- Decision: 在 ReAct 主循环中,在模型返回后、工具调用前后发射 thinking/tool_call/tool_result 事件。 +- Rationale: 该位置可精确覆盖执行时序,且无需新增 channel 层。 +- Alternative considered: 在 tool gateway 或 transport channel 内部推导事件;会丢失模型 thinking 阶段信息。 + +4. **测试策略** +- Decision: 采用“单元契约 + 端到端序列”双层验证: + - transport 类型与 payload/error 契约单测; + - adapter thinking/usage 提取单测; + - agent 到 transport 的事件序列回归测试。 +- Rationale: 保证接口稳定且行为可观测。 + +## Risks / Trade-offs + +- [Risk] 旧消费者依赖 `result/hook` 事件语义。 → Mitigation: 保留 alias 归一化并增加回归测试。 +- [Risk] 不同 provider 的 thinking 字段路径不一致。 → Mitigation: 在 adapter 层集中解析并提供降级为 `None` 的一致行为。 +- [Risk] 事件发射增多带来日志噪声。 → Mitigation: 仅输出必要字段,后续在 D8 切片补采样/脱敏。 + +## Migration Plan + +1. 先扩展类型定义(transport/model),保证兼容映射存在。 +2. 再落地 adapter 提取逻辑(thinking + reasoning_tokens)。 +3. 接入 ReAct loop 事件发射,并补充序列测试。 +4. 运行定向测试与全量回归;更新 TODO/feature evidence。 + +Rollback: +- 若出现兼容问题,可回退到旧事件映射并关闭中间态事件发射路径;类型扩展为向后兼容,不影响核心调用。 + +## Open Questions + +- `status` 事件在当前 D2 切片中是否仅作为保留类型,还是立即要求具体 payload 结构? +- `thinking_content` 是否需要在后续切片中加入脱敏策略(例如 prompt 注入片段清理)? diff --git a/openspec/changes/agentscope-d2-d4-thinking-transport/proposal.md b/openspec/changes/agentscope-d2-d4-thinking-transport/proposal.md new file mode 100644 index 00000000..88f0a76e --- /dev/null +++ b/openspec/changes/agentscope-d2-d4-thinking-transport/proposal.md @@ -0,0 +1,36 @@ +## Why + +AgentScope 迁移的首要阻断在于两条链路未闭环:一是模型 `thinking` 内容在响应结构中丢失,二是 transport 无法稳定表达 `thinking/tool_call/tool_result` 中间态事件。若不先补齐该切片,后续压缩、计划、审计等能力会持续依赖 shim,无法进入框架原生实现。 + +## What Changes + +- 新增一组面向 AgentScope 对齐的运行态能力: + - 模型响应保留 `thinking_content`; + - usage 规范化 `reasoning_tokens`; + - ReAct 执行链路输出 `thinking/tool_call/tool_result` 中间态事件。 +- 统一 transport 事件语义,明确 `message/tool_call/tool_result/thinking/error/status` 的规范映射。 +- 统一 payload/error 契约并补齐回归测试,保证 CLI 与 transport 一致消费。 + +## Capabilities + +### New Capabilities +- `agentscope-thinking-transport`: 定义并落地 AgentScope 对齐所需的 thinking + tool 中间态事件契约与模型响应保真能力。 + +### Modified Capabilities +- `transport-channel`: 增加中间态事件分类与错误载荷一致性要求。 +- `chat-runtime`: 增加模型 `thinking_content` 与 `reasoning_tokens` 规范化要求。 + +## Impact + +- Affected code: + - `dare_framework/transport/*` + - `dare_framework/model/types.py` + - `dare_framework/model/adapters/openai_adapter.py` + - `dare_framework/model/adapters/openrouter_adapter.py` + - `dare_framework/agent/react_agent.py` +- Affected tests: + - `tests/unit/test_transport_types.py` + - `tests/unit/test_transport_channel.py` + - `tests/unit/test_openrouter_adapter.py` + - `tests/unit/test_agent_event_transport_hook.py` +- No external dependency change; compatibility path keeps legacy alias mapping where required. diff --git a/openspec/changes/agentscope-d2-d4-thinking-transport/specs/agentscope-thinking-transport/spec.md b/openspec/changes/agentscope-d2-d4-thinking-transport/specs/agentscope-thinking-transport/spec.md new file mode 100644 index 00000000..883a6f49 --- /dev/null +++ b/openspec/changes/agentscope-d2-d4-thinking-transport/specs/agentscope-thinking-transport/spec.md @@ -0,0 +1,31 @@ +## ADDED Requirements + +### Requirement: Canonical thinking and tool intermediate events +The runtime SHALL expose canonical intermediate event types for model/tool execution progress using the following event taxonomy: `message`, `tool_call`, `tool_result`, `thinking`, `error`, `status`. + +#### Scenario: Runtime emits tool lifecycle events +- **WHEN** the model requests one or more tool calls +- **THEN** the runtime emits at least one `tool_call` event before tool execution +- **AND** emits a matching `tool_result` event after tool execution + +#### Scenario: Runtime emits thinking event when available +- **WHEN** adapter extraction returns non-empty `thinking_content` +- **THEN** the runtime emits a `thinking` event before final `message` output + +### Requirement: Model response preserves reasoning content and usage +`ModelResponse` SHALL preserve model reasoning content in `thinking_content` and SHALL normalize `reasoning_tokens` into usage metadata when provider payload includes reasoning token usage. + +#### Scenario: Adapter extracts reasoning content +- **WHEN** provider payload includes a reasoning/thinking field +- **THEN** `ModelResponse.thinking_content` is populated with extracted text + +#### Scenario: Adapter normalizes reasoning token usage +- **WHEN** provider payload includes reasoning token counters +- **THEN** `ModelResponse.usage["reasoning_tokens"]` is present and numeric + +### Requirement: Legacy transport aliases remain consumable +Transport event type normalization SHALL continue to accept legacy aliases and map them to canonical values without breaking existing clients. + +#### Scenario: Legacy approval alias is normalized +- **WHEN** an envelope arrives with a legacy approval alias event type +- **THEN** runtime normalization maps it to the canonical event type value before dispatch diff --git a/openspec/changes/agentscope-d2-d4-thinking-transport/specs/chat-runtime/spec.md b/openspec/changes/agentscope-d2-d4-thinking-transport/specs/chat-runtime/spec.md new file mode 100644 index 00000000..1f1f6e98 --- /dev/null +++ b/openspec/changes/agentscope-d2-d4-thinking-transport/specs/chat-runtime/spec.md @@ -0,0 +1,26 @@ +## MODIFIED Requirements + +### Requirement: LLM-driven execute loop +The runtime SHALL invoke the configured `IModelAdapter` during the execute loop, provide the assembled prompt and available tool definitions, and iterate over tool calls until the model returns a final response. + +The execute loop SHALL preserve adapter-extracted `thinking_content` and normalized `reasoning_tokens` in the model response payload used by runtime hooks/transport emission. + +When runtime transport emission is enabled, the execute loop SHALL emit canonical intermediate events (`thinking`, `tool_call`, `tool_result`) in execution order prior to final `message` output. + +#### Scenario: Model returns a final response +- **WHEN** the model response contains no tool calls +- **THEN** the execute loop returns success and exposes the response content in the run output + +#### Scenario: Model requests a tool call +- **WHEN** the model response includes a tool call +- **THEN** the runtime executes the tool via `ToolRuntime`, appends the result to the message history, and continues + +#### Scenario: Runtime preserves reasoning content and tokens +- **WHEN** adapter response includes thinking content and reasoning token usage +- **THEN** runtime keeps `thinking_content` in the model response object +- **AND** usage contains normalized `reasoning_tokens` + +#### Scenario: Runtime emits ordered intermediate events +- **GIVEN** transport sender is configured +- **WHEN** the runtime performs one round with model thinking and one tool call +- **THEN** emitted events preserve order `thinking -> tool_call -> tool_result -> message` diff --git a/openspec/changes/agentscope-d2-d4-thinking-transport/specs/transport-channel/spec.md b/openspec/changes/agentscope-d2-d4-thinking-transport/specs/transport-channel/spec.md new file mode 100644 index 00000000..9f81ddea --- /dev/null +++ b/openspec/changes/agentscope-d2-d4-thinking-transport/specs/transport-channel/spec.md @@ -0,0 +1,22 @@ +## MODIFIED Requirements + +### Requirement: Envelope kind supports message, action, and control categories +The transport envelope model SHALL provide a primary categorization field for inbound/outbound envelopes that can distinguish: +- `message` (prompt/result/hook style messages) +- `action` (deterministic resource actions) +- `control` (interrupt/pause/retry/reverse) + +For `kind="control"`, implementations SHALL represent control variants via the payload value (e.g. `payload="interrupt"`). +The envelope model MUST NOT require a separate subtype field (e.g. `type`) in order to route inbound envelopes. + +For `kind="message"`, implementations SHALL support canonical `event_type` values `message|tool_call|tool_result|thinking|error|status` so intermediate model/tool progress can be delivered without prompt parsing. + +#### Scenario: Action envelope is distinguishable without slash parsing +- **GIVEN** a client sends `TransportEnvelope(kind="action", payload="tools:list")` +- **WHEN** the channel receives it +- **THEN** the channel can route it deterministically without inspecting prompt text + +#### Scenario: Message envelope carries canonical intermediate event type +- **GIVEN** a runtime emits `TransportEnvelope(kind="message", event_type="tool_call")` +- **WHEN** the envelope is processed by transport consumers +- **THEN** consumers can classify it as a tool invocation intermediate event without content heuristics diff --git a/openspec/changes/agentscope-d2-d4-thinking-transport/tasks.md b/openspec/changes/agentscope-d2-d4-thinking-transport/tasks.md new file mode 100644 index 00000000..cc31c3d1 --- /dev/null +++ b/openspec/changes/agentscope-d2-d4-thinking-transport/tasks.md @@ -0,0 +1,24 @@ +## 1. Transport contract and schema alignment (D2) + +- [x] 1.1 Extend transport canonical event taxonomy to cover `message/tool_call/tool_result/thinking/error/status` with compatibility normalization for legacy aliases. +- [x] 1.2 Align transport payload/error schema so intermediate events and failures can be consumed consistently by CLI and channel integrations. +- [x] 1.3 Add contract-focused unit tests for transport type normalization and payload/error stability. +- [x] 1.4 Add transport sequence tests covering intermediate event ordering behavior. + +## 2. Model response reasoning preservation (D4 core) + +- [x] 2.1 Extend `ModelResponse` with optional `thinking_content` field and preserve backward compatibility for existing callers. +- [x] 2.2 Update OpenAI/OpenRouter adapters to extract reasoning content and normalize `reasoning_tokens` into usage metadata. +- [x] 2.3 Add adapter regression tests verifying thinking preservation and reasoning token normalization. + +## 3. ReAct loop intermediate event emission (D4 execution) + +- [x] 3.1 Emit `thinking` events from the execute loop when `thinking_content` is present. +- [x] 3.2 Emit `tool_call` and `tool_result` events around tool execution rounds in deterministic order. +- [x] 3.3 Add end-to-end agent/transport tests asserting ordered event sequence and non-regression of final message output. + +## 4. Documentation and ledger sync + +- [x] 4.1 Update TODO claim ledger statuses for this slice from `planned` to `active` and keep scope mapping aligned with this change. +- [x] 4.2 Update feature aggregation evidence for this slice with commands/results/behavior verification. +- [x] 4.3 Run targeted tests and full regression, then record evidence links before requesting review. diff --git a/tests/unit/test_base_agent_transport_contract.py b/tests/unit/test_base_agent_transport_contract.py index bf6c2375..0abe4574 100644 --- a/tests/unit/test_base_agent_transport_contract.py +++ b/tests/unit/test_base_agent_transport_contract.py @@ -185,6 +185,22 @@ async def execute(self, task: str | Any, *, transport: Any = None) -> RunResult: raise RuntimeError("simulated model timeout") +class _TransportLoopFlagProbeAgent(BaseAgent): + def __init__(self, name: str) -> None: + super().__init__(name) + self.observed_loop_state: dict[str, bool] = {} + self.loop_execution_started = asyncio.Event() + self.allow_loop_execution_finish = asyncio.Event() + + async def execute(self, task: str | Any, *, transport: Any = None) -> RunResult: + task_text = task.description if hasattr(task, "description") else str(task) + self.observed_loop_state[task_text] = self._is_transport_loop_execution(transport=transport) + if task_text == "loop-task": + self.loop_execution_started.set() + await self.allow_loop_execution_finish.wait() + return RunResult(success=True, output=task_text, output_text=task_text) + + @pytest.mark.asyncio async def test_public_surface_uses_call_not_run() -> None: agent = _CaptureAgent("capture-agent") @@ -273,3 +289,25 @@ async def test_transport_loop_returns_structured_error_when_execute_raises() -> assert isinstance(payload, dict) assert payload.get("code") == "AGENT_EXECUTION_FAILED" assert "simulated model timeout" in str(payload.get("reason")) + + +@pytest.mark.asyncio +async def test_transport_loop_flag_is_task_local_for_concurrent_execute_calls() -> None: + agent = _TransportLoopFlagProbeAgent("probe-agent") + channel = _RecordingSendChannel() + + polled_task = asyncio.create_task( + agent._execute_polled_message( + "loop-task", + channel=channel, + envelope_id="req_1", + ) + ) + await agent.loop_execution_started.wait() + + await agent.execute("direct-task", transport=channel) + agent.allow_loop_execution_finish.set() + await polled_task + + assert agent.observed_loop_state["loop-task"] is True + assert agent.observed_loop_state["direct-task"] is False diff --git a/tests/unit/test_openai_model_adapter.py b/tests/unit/test_openai_model_adapter.py new file mode 100644 index 00000000..94fd2f2e --- /dev/null +++ b/tests/unit/test_openai_model_adapter.py @@ -0,0 +1,69 @@ +from __future__ import annotations + +from types import SimpleNamespace + +from dare_framework.model.adapters.openai_adapter import OpenAIModelAdapter + + +def test_extract_usage_normalizes_reasoning_tokens() -> None: + adapter = OpenAIModelAdapter() + response = SimpleNamespace( + response_metadata={ + "token_usage": { + "prompt_tokens": 11, + "completion_tokens": 22, + "total_tokens": 33, + "completion_tokens_details": { + "reasoning_tokens": 9, + }, + } + }, + additional_kwargs={}, + ) + + usage = adapter._extract_usage(response) + + assert usage == { + "prompt_tokens": 11, + "completion_tokens": 22, + "total_tokens": 33, + "reasoning_tokens": 9, + } + + +def test_extract_usage_reads_reasoning_tokens_from_output_tokens_details() -> None: + adapter = OpenAIModelAdapter() + response = SimpleNamespace( + response_metadata={ + "token_usage": { + "prompt_tokens": 3, + "completion_tokens": 7, + "total_tokens": 10, + "output_tokens_details": { + "reasoning_tokens": 4, + }, + } + }, + additional_kwargs={}, + ) + + usage = adapter._extract_usage(response) + + assert usage == { + "prompt_tokens": 3, + "completion_tokens": 7, + "total_tokens": 10, + "reasoning_tokens": 4, + } + + +def test_extract_thinking_content_from_response_additional_kwargs() -> None: + adapter = OpenAIModelAdapter() + response = SimpleNamespace( + additional_kwargs={ + "reasoning_content": "internal reasoning", + }, + response_metadata={}, + ) + + assert adapter._extract_thinking_content(response) == "internal reasoning" diff --git a/tests/unit/test_openrouter_adapter.py b/tests/unit/test_openrouter_adapter.py index c9650b54..6f56466f 100644 --- a/tests/unit/test_openrouter_adapter.py +++ b/tests/unit/test_openrouter_adapter.py @@ -1,9 +1,14 @@ from __future__ import annotations import json +from types import SimpleNamespace from dare_framework.context.types import Message -from dare_framework.model.adapters.openrouter_adapter import _serialize_messages +from dare_framework.model.adapters.openrouter_adapter import ( + _extract_reasoning_tokens, + _extract_thinking_content, + _serialize_messages, +) def test_serialize_messages_preserves_assistant_tool_calls() -> None: @@ -56,3 +61,21 @@ def test_serialize_messages_preserves_assistant_tool_calls() -> None: assert tool_call["function"]["name"] == "ask_user" assert isinstance(tool_call["function"]["arguments"], str) assert serialized[2]["tool_call_id"] == "call_1" + + +def test_extract_thinking_content_from_openrouter_message_fields() -> None: + message = SimpleNamespace(reasoning="step by step", reasoning_content=None, additional_kwargs={}) + assert _extract_thinking_content(message) == "step by step" + + +def test_extract_reasoning_tokens_from_completion_tokens_details() -> None: + response = SimpleNamespace( + usage=SimpleNamespace( + prompt_tokens=1, + completion_tokens=2, + total_tokens=3, + reasoning_tokens=None, + completion_tokens_details=SimpleNamespace(reasoning_tokens=7), + ) + ) + assert _extract_reasoning_tokens(response) == 7 diff --git a/tests/unit/test_react_agent_gateway_injection.py b/tests/unit/test_react_agent_gateway_injection.py index ef8a45c9..da56c571 100644 --- a/tests/unit/test_react_agent_gateway_injection.py +++ b/tests/unit/test_react_agent_gateway_injection.py @@ -9,6 +9,7 @@ from dare_framework.context import Context from dare_framework.model.types import ModelInput, ModelResponse from dare_framework.tool.types import CapabilityDescriptor, CapabilityType, ToolResult +from dare_framework.transport import TransportEventType class _SequenceModel: @@ -74,6 +75,25 @@ async def generate(self, model_input: ModelInput, *, options: Any | None = None) ) +class _UniqueToolLoopModel: + def __init__(self) -> None: + self._idx = 0 + + async def generate(self, model_input: ModelInput, *, options: Any | None = None) -> ModelResponse: + _ = (model_input, options) + self._idx += 1 + return ModelResponse( + content="still searching", + tool_calls=[ + { + "id": f"tc_{self._idx}", + "name": "tool:echo", + "arguments": {"value": "ping"}, + } + ], + ) + + class _RecordingGateway: def __init__(self, label: str) -> None: self.label = label @@ -98,6 +118,39 @@ async def invoke(self, capability_id: str, *, envelope: Any, **params: Any) -> T return ToolResult(success=True, output={"gateway": self.label, "params": params}) +class _RecordingTransport: + def __init__(self) -> None: + self.sent: list[Any] = [] + + async def send(self, envelope: Any) -> None: + self.sent.append(envelope) + + +class _ThinkingSequenceModel: + def __init__(self) -> None: + self._responses = [ + ModelResponse( + content="calling tool", + thinking_content="need tool data", + tool_calls=[ + { + "id": "tc_1", + "name": "tool:echo", + "arguments": {"value": "ping"}, + } + ], + ), + ModelResponse(content="final answer", tool_calls=[]), + ] + self._idx = 0 + + async def generate(self, model_input: ModelInput, *, options: Any | None = None) -> ModelResponse: + _ = (model_input, options) + response = self._responses[self._idx] + self._idx += 1 + return response + + @pytest.mark.asyncio async def test_react_agent_prefers_injected_gateway_over_context_gateway() -> None: context_gateway = _RecordingGateway("context") @@ -156,3 +209,111 @@ async def test_react_agent_stops_repeated_identical_tool_loop() -> None: assert result.success is True assert "连续重复调用相同工具" in str(result.output_text) assert len(gateway.invoke_calls) == 2 + + +@pytest.mark.asyncio +async def test_react_agent_emits_intermediate_transport_events_in_order() -> None: + context = Context(config=Config()) + gateway = _RecordingGateway("injected") + transport = _RecordingTransport() + + agent = ReactAgent( + name="react-test-transport-events", + model=_ThinkingSequenceModel(), + context=context, + tool_gateway=gateway, + ) + + result = await agent.execute("test", transport=transport) + + assert result.success is True + event_types = [getattr(envelope, "event_type", None) for envelope in transport.sent] + assert event_types == [ + TransportEventType.THINKING.value, + TransportEventType.TOOL_CALL.value, + TransportEventType.TOOL_RESULT.value, + TransportEventType.MESSAGE.value, + ] + payloads = [getattr(envelope, "payload", None) for envelope in transport.sent] + assert payloads[0]["ok"] is True + assert payloads[1]["resp"]["name"] == "tool:echo" + assert payloads[2]["resp"]["success"] is True + assert payloads[3]["resp"]["output"] == "final answer" + + +@pytest.mark.asyncio +async def test_react_agent_transport_loop_emits_single_terminal_result_event() -> None: + context = Context(config=Config()) + gateway = _RecordingGateway("injected") + transport = _RecordingTransport() + + agent = ReactAgent( + name="react-test-transport-loop-terminal", + model=_ThinkingSequenceModel(), + context=context, + tool_gateway=gateway, + agent_channel=transport, + ) + + await agent._execute_polled_message( + "test", + channel=transport, + envelope_id="req_1", + ) + + event_types = [getattr(envelope, "event_type", None) for envelope in transport.sent] + assert event_types == [ + TransportEventType.THINKING.value, + TransportEventType.TOOL_CALL.value, + TransportEventType.TOOL_RESULT.value, + TransportEventType.RESULT.value, + ] + + terminal_payload = getattr(transport.sent[-1], "payload", {}) + assert terminal_payload.get("resp", {}).get("success") is True + assert terminal_payload.get("resp", {}).get("output", {}).get("content") == "final answer" + + +@pytest.mark.asyncio +async def test_react_agent_emits_terminal_message_for_repeated_tool_guard() -> None: + context = Context(config=Config()) + gateway = _RecordingGateway("injected") + transport = _RecordingTransport() + + agent = ReactAgent( + name="react-test-loop-guard-terminal", + model=_RepeatingToolModel(), + context=context, + tool_gateway=gateway, + ) + + result = await agent.execute("test", transport=transport) + + assert result.success is True + assert transport.sent + last_envelope = transport.sent[-1] + assert getattr(last_envelope, "event_type", None) == TransportEventType.MESSAGE.value + assert "连续重复调用相同工具" in str(getattr(last_envelope, "payload", {}).get("resp", {}).get("output", "")) + + +@pytest.mark.asyncio +async def test_react_agent_emits_terminal_message_for_max_round_exit() -> None: + context = Context(config=Config()) + gateway = _RecordingGateway("injected") + transport = _RecordingTransport() + + agent = ReactAgent( + name="react-test-max-round-terminal", + model=_UniqueToolLoopModel(), + context=context, + tool_gateway=gateway, + max_tool_rounds=2, + ) + + result = await agent.execute("test", transport=transport) + + assert result.success is True + assert transport.sent + last_envelope = transport.sent[-1] + assert getattr(last_envelope, "event_type", None) == TransportEventType.MESSAGE.value + assert "达到最大轮次" in str(getattr(last_envelope, "payload", {}).get("resp", {}).get("output", "")) diff --git a/tests/unit/test_transport_adapters.py b/tests/unit/test_transport_adapters.py index 60c77c54..da521088 100644 --- a/tests/unit/test_transport_adapters.py +++ b/tests/unit/test_transport_adapters.py @@ -192,6 +192,23 @@ async def test_stdio_receiver_uses_event_type_without_legacy_payload_type(capsys assert "approval pending: request_id=req-42" in captured.out +@pytest.mark.asyncio +async def test_stdio_receiver_renders_status_phase_from_structured_resp(capsys) -> None: + channel = StdioClientChannel() + receiver = channel.agent_envelope_receiver() + await receiver( + TransportEnvelope( + id="evt-status-phase", + kind=EnvelopeKind.MESSAGE, + event_type=TransportEventType.STATUS.value, + payload={"resp": {"phase": "running"}}, + ) + ) + captured = capsys.readouterr() + assert "running" in captured.out + assert "approval update" not in captured.out + + @pytest.mark.asyncio async def test_stdio_receiver_does_not_route_by_payload_type_without_event_type(capsys) -> None: channel = StdioClientChannel() @@ -212,3 +229,28 @@ async def test_stdio_receiver_does_not_route_by_payload_type_without_event_type( captured = capsys.readouterr() assert "Assistant: {'type': 'result'" in captured.out assert "Assistant: hello" not in captured.out + + +@pytest.mark.asyncio +async def test_stdio_receiver_handles_canonical_thinking_event(capsys) -> None: + channel = StdioClientChannel() + receiver = channel.agent_envelope_receiver() + + await receiver( + TransportEnvelope( + id="evt-thinking", + kind=EnvelopeKind.MESSAGE, + event_type=TransportEventType.THINKING.value, + payload={ + "kind": "message", + "target": "model", + "ok": True, + "resp": { + "output": "need tool data", + }, + }, + ) + ) + + captured = capsys.readouterr() + assert "Assistant: need tool data" in captured.out diff --git a/tests/unit/test_transport_types.py b/tests/unit/test_transport_types.py index 7e0cb349..9b8b40aa 100644 --- a/tests/unit/test_transport_types.py +++ b/tests/unit/test_transport_types.py @@ -2,7 +2,7 @@ import pytest -from dare_framework.transport import normalize_transport_event_type +from dare_framework.transport import canonicalize_transport_event_type, normalize_transport_event_type from dare_framework.transport.types import EnvelopeKind, TransportEnvelope, TransportEventType @@ -48,3 +48,29 @@ def test_transport_envelope_rejects_empty_event_type() -> None: def test_transport_facade_re_exports_event_type_normalizer() -> None: assert normalize_transport_event_type("approval_pending") == TransportEventType.APPROVAL_PENDING.value + + +def test_transport_event_type_includes_canonical_categories() -> None: + assert TransportEventType.MESSAGE.value == "message" + assert TransportEventType.TOOL_CALL.value == "tool_call" + assert TransportEventType.TOOL_RESULT.value == "tool_result" + assert TransportEventType.THINKING.value == "thinking" + assert TransportEventType.ERROR.value == "error" + assert TransportEventType.STATUS.value == "status" + + +@pytest.mark.parametrize( + ("raw", "expected"), + [ + ("result", "message"), + ("tool.result", "tool_result"), + ("tool.call", "tool_call"), + ("hook", "status"), + ("approval.pending", "status"), + ("approval_pending", "status"), + ("approval.resolved", "status"), + ("approval_resolved", "status"), + ], +) +def test_canonicalize_transport_event_type_maps_legacy_aliases(raw: str, expected: str) -> None: + assert canonicalize_transport_event_type(raw) == expected