| """anuma.ai Responses SSE 事件 → IREvent(原生协议与 OpenAI Responses 一致)。 |
| |
| 上游:``POST portal.anuma.ai/api/v1/responses``,返回标准 OpenAI Responses SSE: |
| ``response.created`` / ``response.output_text.delta`` / ``response.completed`` 等。 |
| 正文增量由 ``response.output_text.delta`` 产出;思维链摘要增量由 |
| ``response.reasoning_summary_text.delta`` 产出(必须产出 kind="thinking")。 |
| |
| **原生 function calling**:上游接受客户端 tools 并原生返回 function_call 事件 |
| (``response.output_item.added`` type=function_call + ``response.function_call_arguments.delta/done``)。 |
| 本 parser 把 function_call 统一转成 ``<tool_call>{...}</tool_call>`` 围栏文本 |
| (携带上游 call_id),走宿主既有解析链路(parse_tool_calls / ToolCallStreamParser)。 |
| """ |
| from __future__ import annotations |
|
|
| import json |
| from typing import Any |
|
|
| from app.events import IREvent, Usage |
| from app.upstream.base import EventParser |
|
|
|
|
| class DefaultParser(EventParser): |
| def __init__(self) -> None: |
| |
| self._pending_calls: dict[str, dict[str, str]] = {} |
| |
| self._saw_text_delta = False |
|
|
| def parse(self, raw: Any) -> list[IREvent]: |
| """raw 为单个 SSE 事件 dict(``{"type": ..., ...}``),返回 0..n 个 IREvent。""" |
| if not isinstance(raw, dict): |
| return [] |
| etype = raw.get("type") or "" |
| delta = raw.get("delta") or {} |
| if not isinstance(delta, dict): |
| delta = {} |
|
|
| events: list[IREvent] = [] |
|
|
| |
| if etype == "response.created": |
| self._pending_calls.clear() |
| self._saw_text_delta = False |
| return events |
|
|
| |
| if etype == "response.output_item.added": |
| item = raw.get("item") or {} |
| if isinstance(item, dict) and item.get("type") == "function_call": |
| call_id = str(item.get("call_id") or item.get("id") or "") |
| if call_id: |
| self._pending_calls[call_id] = { |
| "name": str(item.get("name") or ""), |
| "arguments": str(item.get("arguments") or ""), |
| } |
| return events |
|
|
| |
| if etype == "response.function_call_arguments.delta": |
| |
| for call_id, pending in self._pending_calls.items(): |
| pending["arguments"] += str(delta if isinstance(delta, str) else delta.get("delta") or "") |
| return events |
|
|
| |
| if etype == "response.function_call_arguments.done": |
| for call_id, pending in self._pending_calls.items(): |
| fence = _function_call_fence(pending["name"], pending["arguments"], call_id) |
| if fence: |
| events.append(IREvent(kind="text", text=fence)) |
| self._pending_calls.clear() |
| return events |
|
|
| |
| |
| |
| |
| |
| |
| |
| if etype == "response.reasoning_summary_text.delta": |
| t = delta.get("OfString") if isinstance(delta, dict) else None |
| if not t: |
| s = delta.get("OfResponseReasoningSummaryDeltaEventDelta") if isinstance(delta, dict) else None |
| if isinstance(s, dict): |
| t = s.get("text") or "" |
| else: |
| t = s or "" |
| if t: |
| events.append(IREvent(kind="thinking", thinking=str(t))) |
| return events |
|
|
| |
| |
| |
| |
| |
| if etype == "response.output_text.delta": |
| text = delta.get("OfString") if isinstance(delta, dict) else "" |
| if isinstance(text, dict): |
| text = text.get("text") or "" |
| if text: |
| self._saw_text_delta = True |
| events.append(IREvent(kind="text", text=str(text))) |
| return events |
|
|
| |
| |
| |
| |
| |
| |
| |
| |
| |
| if etype == "response.output_text.done": |
| text = raw.get("text") |
| if isinstance(text, str) and text: |
| if self._saw_text_delta and "portal.anuma.ai/api/v1/media" not in text: |
| return events |
| events.append(IREvent(kind="text", text=text)) |
| return events |
|
|
| |
| |
| |
| |
| |
| if etype == "response.completed": |
| resp = raw.get("response") or {} |
| if isinstance(resp, dict): |
| for item in resp.get("output") or []: |
| if isinstance(item, dict) and item.get("type") == "function_call": |
| fence = _function_call_fence( |
| str(item.get("name") or ""), |
| item.get("arguments") or "{}", |
| str(item.get("call_id") or item.get("id") or ""), |
| ) |
| if fence: |
| events.append(IREvent(kind="text", text=fence)) |
| |
| |
| for tce in resp.get("tool_call_events") or []: |
| if not isinstance(tce, dict): |
| continue |
| output = tce.get("output") |
| if not isinstance(output, str) or not output: |
| continue |
| try: |
| out = json.loads(output) |
| except json.JSONDecodeError: |
| continue |
| urls = [] |
| for img in out.get("output_images") or []: |
| u = img.get("url") if isinstance(img, dict) else None |
| if u: |
| urls.append(u) |
| if urls: |
| text = "\n".join(f"" for u in urls) |
| events.append(IREvent(kind="text", text=text)) |
| usage = self._parse_usage(raw) |
| if usage: |
| events.append(IREvent(kind="finish", usage_delta=usage, finish_reason="stop")) |
| else: |
| events.append(IREvent(kind="finish", finish_reason="stop")) |
| return events |
|
|
| |
| if etype == "" or etype == "response.usage": |
| usage = self._parse_usage(raw) |
| if usage: |
| events.append(IREvent(kind="finish", usage_delta=usage, finish_reason="stop")) |
| return events |
|
|
| |
| if etype == "response.failed": |
| err = raw.get("error") or raw.get("message") or "upstream response failed" |
| msg = err.get("message") if isinstance(err, dict) else str(err) |
| events.append(IREvent(kind="error", error=msg or "upstream response failed")) |
| return events |
|
|
| return events |
|
|
| @staticmethod |
| def _parse_usage(raw: dict[str, Any]) -> Usage | None: |
| usage = raw.get("usage") |
| if not isinstance(usage, dict): |
| return None |
| |
| |
| input_tokens = usage.get("input_tokens") or usage.get("prompt_tokens") or 0 |
| output_tokens = usage.get("output_tokens") or usage.get("completion_tokens") or 0 |
| if not input_tokens and not output_tokens: |
| return None |
| details = usage.get("output_tokens_details") or {} |
| thinking = details.get("reasoning_tokens") if isinstance(details, dict) else 0 |
| cached = 0 |
| in_details = usage.get("input_tokens_details") or {} |
| if isinstance(in_details, dict): |
| cached = in_details.get("cached_tokens") or 0 |
| return Usage( |
| input_tokens=int(input_tokens), |
| output_tokens=int(output_tokens), |
| thinking_tokens=int(thinking or 0), |
| cached_tokens=int(cached or 0), |
| model=str(usage.get("model") or "") or None, |
| provider="anuma", |
| ) |
|
|
|
|
| def _function_call_fence(name: str, arguments: str, call_id: str) -> str | None: |
| """function_call → ``<tool_call>{"name","arguments","id"}</tool_call>`` 围栏文本。""" |
| if not name: |
| return None |
| try: |
| args = json.loads(arguments) if arguments.strip() else {} |
| except json.JSONDecodeError: |
| args = {"value": arguments} |
| if not isinstance(args, dict): |
| args = {"value": args} |
| return ( |
| f"<tool_call>{json.dumps({'name': name, 'arguments': args, 'id': call_id}, |
| ensure_ascii=False)}</tool_call>" |
| ) |
|
|