Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
251 changes: 221 additions & 30 deletions app/adapters/responses_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,18 @@ def _text_format_to_response_format(fmt) -> dict | None:

def responses_request_to_chat(body: dict) -> dict:
"""Convert Responses input, instructions and tools to a Chat request."""
# Collect declared tools from both top-level tools and input additional_tools items.
raw_tools = []
if isinstance(body.get("tools"), list):
raw_tools.extend(body["tools"])
inp = body.get("input", [])
if isinstance(inp, list):
for item in inp:
if isinstance(item, dict) and item.get("type") == "additional_tools" and isinstance(item.get("tools"), list):
raw_tools.extend(item["tools"])

registry, chat_tools = _build_tool_registry_and_chat_tools(raw_tools)

messages: list[dict] = []

# instructions → system message
Expand All @@ -55,11 +67,10 @@ def responses_request_to_chat(body: dict) -> dict:
messages.append({"role": "system", "content": instructions})

# input → messages
inp = body.get("input", [])
if isinstance(inp, str):
messages.append({"role": "user", "content": inp})
elif isinstance(inp, list):
messages.extend(_convert_input_items(inp))
messages.extend(_convert_input_items(inp, tool_registry=registry))

# Build the Chat request body.
chat: dict[str, Any] = {"messages": messages, "stream": True}
Expand All @@ -68,12 +79,12 @@ def responses_request_to_chat(body: dict) -> dict:
if "model" in body:
chat["model"] = body["model"]

# Normalize function tool definitions.
tools = body.get("tools")
if tools:
chat["tools"] = _convert_tools_for_chat(tools)
if chat_tools:
chat["tools"] = chat_tools
if registry.upstream_to_identity:
chat["_tool_registry"] = registry.to_dict()
if "tool_choice" in body:
chat["tool_choice"] = body["tool_choice"]
chat["tool_choice"] = _convert_tool_choice_for_chat(body["tool_choice"], tool_registry=registry)

# Forward supported parameters.
for key in ("temperature", "top_p", "stop", "seed",
Expand Down Expand Up @@ -102,7 +113,7 @@ def responses_request_to_chat(body: dict) -> dict:
return chat


def _convert_input_items(items: list) -> list[dict]:
def _convert_input_items(items: list, tool_registry: ToolRegistry | None = None) -> list[dict]:
"""Convert input items and merge adjacent assistant messages with tool calls."""
messages: list[dict] = []
# Buffer adjacent assistant text and function calls.
Expand All @@ -127,6 +138,17 @@ def _flush_assistant():
item_type = item.get("type")
role = item.get("role", "")

# Ignore additional_tools in message sequence.
if item_type == "additional_tools":
continue

# Agent message from multi-agent collaboration.
if item_type == "agent_message":
_flush_assistant()
content = _extract_content(item.get("content", ""))
messages.append({"role": "user", "content": content})
continue

# Untyped role messages
if item_type is None and role in ("user", "system", "developer"):
_flush_assistant()
Expand Down Expand Up @@ -165,11 +187,14 @@ def _flush_assistant():
raise ValueError("function_call.arguments must be a JSON string")
if pending_assistant_content is None:
pending_assistant_content = ""
raw_name = item.get("name", "")
ns = item.get("namespace")
mapped_name = tool_registry.get_upstream_name(ns, raw_name) if tool_registry else raw_name
pending_tool_calls.append({
"id": item.get("call_id", item.get("id", _rand_id("call_"))),
"type": "function",
"function": {
"name": item.get("name", ""),
"name": mapped_name,
"arguments": arguments,
},
})
Expand Down Expand Up @@ -208,6 +233,10 @@ def _extract_content(content) -> str | list[dict]:
kind = p.get("type")
if kind in ("input_text", "text", "output_text"):
parts.append({"type": "text", "text": p.get("text", "")})
elif kind == "encrypted_content":
text = p.get("encrypted_content") or p.get("text", "")
if isinstance(text, str) and text:
parts.append({"type": "text", "text": text})
elif kind == "input_image":
if p.get("file_id"):
raise ValueError("Responses input_image file_id is not supported; provide image_url instead")
Expand Down Expand Up @@ -242,28 +271,169 @@ def _extract_output_text(content_parts: list) -> str | list[dict]:
return "".join(texts)


def _convert_tools_for_chat(tools: list) -> list:
"""Convert Responses tool definitions to Chat function objects."""
result = []
def _sanitize_schema(schema: Any) -> Any:
"""Recursively strip 'encrypted' client-only marker from parameter schemas."""
if isinstance(schema, dict):
return {k: _sanitize_schema(v) for k, v in schema.items() if k != "encrypted"}
if isinstance(schema, list):
return [_sanitize_schema(item) for item in schema]
return schema


class ToolRegistry:
"""Bidirectional mapping between Responses (namespace, name) and upstream Chat function names."""

def __init__(self, mappings: dict[str, dict[str, Any]] | None = None):
self.upstream_to_identity: dict[str, tuple[str | None, str]] = {}
self.identity_to_upstream: dict[tuple[str | None, str], str] = {}
self._name_to_identities: dict[str, list[tuple[str | None, str]]] = {}

if mappings:
for up_name, ident in mappings.items():
ns = ident.get("namespace") if isinstance(ident, dict) else None
nm = ident.get("name", up_name) if isinstance(ident, dict) else up_name
self._record(up_name, ns, nm)

def _record(self, upstream_name: str, namespace: str | None, name: str) -> None:
ident = (namespace, name)
self.upstream_to_identity[upstream_name] = ident
self.identity_to_upstream[ident] = upstream_name
self._name_to_identities.setdefault(name, []).append(ident)

def register(self, namespace: str | None, name: str) -> str:
"""Register a tool identity (namespace, name) and return its unique upstream name."""
ident = (namespace, name)
if ident in self.identity_to_upstream:
return self.identity_to_upstream[ident]

if namespace is None:
upstream_name = name
else:
base = f"{namespace}__{name}"
candidate = base
counter = 1
while candidate in self.upstream_to_identity:
candidate = f"{base}_{counter}"
counter += 1
upstream_name = candidate

self._record(upstream_name, namespace, name)
return upstream_name

def get_identity(self, upstream_name: str) -> tuple[str | None, str]:
"""Map upstream function name back to (namespace, name). Never guess by prefix."""
if upstream_name in self.upstream_to_identity:
return self.upstream_to_identity[upstream_name]
return None, upstream_name

def get_upstream_name(self, namespace: str | None, name: str) -> str:
"""Map (namespace, name) to upstream name, with fallback for historical calls."""
ident = (namespace, name)
if ident in self.identity_to_upstream:
return self.identity_to_upstream[ident]
if namespace is None:
if (None, name) in self.identity_to_upstream:
return self.identity_to_upstream[(None, name)]
matches = self._name_to_identities.get(name, [])
if len(matches) == 1:
return self.identity_to_upstream[matches[0]]
return name
return f"{namespace}__{name}"

def to_dict(self) -> dict:
return {
up_name: {"namespace": ident[0], "name": ident[1]}
for up_name, ident in self.upstream_to_identity.items()
}

@classmethod
def from_dict(cls, data: dict | None) -> ToolRegistry:
if isinstance(data, cls):
return data
if not isinstance(data, dict):
return cls()
return cls(data)


def _collect_raw_tools(tools: list, current_ns: str = "") -> list[tuple[str | None, str, dict]]:
"""Recursively collect (namespace, name, tool_def) from Responses tools."""
collected = []
for t in tools:
if not isinstance(t, dict):
continue
if t.get("type") != "function":
continue
# Already in Chat format.
if "function" in t:
result.append(t)
tool_type = t.get("type")
if tool_type == "namespace" and isinstance(t.get("tools"), list):
ns_name = t.get("name", "")
full_ns = f"{current_ns}.{ns_name}" if current_ns else ns_name
collected.extend(_collect_raw_tools(t["tools"], full_ns))
elif tool_type == "function" or "function" in t:
fn = t.get("function") if isinstance(t.get("function"), dict) else t
name = fn.get("name") or t.get("name")
if isinstance(name, str) and name:
collected.append((current_ns or None, name, t))
return collected


def _format_chat_tool(upstream_name: str, tool_def: dict) -> dict:
"""Format tool definition for Chat Completions, ensuring sanitized schema."""
if "function" in tool_def:
tool_obj = json.loads(json.dumps(tool_def))
fn_dict = tool_obj.get("function")
if isinstance(fn_dict, dict):
fn_dict["name"] = upstream_name
if "parameters" in fn_dict:
fn_dict["parameters"] = _sanitize_schema(fn_dict["parameters"])
return tool_obj

fn: dict[str, Any] = {"name": upstream_name}
if "description" in tool_def:
fn["description"] = tool_def["description"]
if "parameters" in tool_def:
fn["parameters"] = _sanitize_schema(tool_def["parameters"])
if "strict" in tool_def:
fn["strict"] = tool_def["strict"]
return {"type": "function", "function": fn}


def _build_tool_registry_and_chat_tools(raw_tools: list) -> tuple[ToolRegistry, list[dict]]:
"""Build ToolRegistry and convert declared tools into unique Chat function tools."""
collected = _collect_raw_tools(raw_tools)
registry = ToolRegistry()
result = []

# Register plain tools first so their upstream name strictly preserves original name.
plain = [item for item in collected if item[0] is None]
namespaced = [item for item in collected if item[0] is not None]

seen_identities: set[tuple[str | None, str]] = set()
for ns, name, tool_def in plain + namespaced:
ident = (ns, name)
if ident in seen_identities:
continue
# Nest the flat Responses function fields.
fn: dict[str, Any] = {"name": t.get("name", "")}
if "description" in t:
fn["description"] = t["description"]
if "parameters" in t:
fn["parameters"] = t["parameters"]
if "strict" in t:
fn["strict"] = t["strict"]
result.append({"type": "function", "function": fn})
return result
seen_identities.add(ident)
upstream_name = registry.register(ns, name)
result.append(_format_chat_tool(upstream_name, tool_def))

return registry, result


def _convert_tools_for_chat(tools: list) -> list:
"""Backwards-compatible helper for converting Responses tools."""
_, chat_tools = _build_tool_registry_and_chat_tools(tools)
return chat_tools


def _convert_tool_choice_for_chat(choice: Any, tool_registry: ToolRegistry | None = None) -> Any:
"""Convert Responses tool_choice to Chat Completions format."""
if not isinstance(choice, dict):
return choice
fn_obj = choice.get("function") if isinstance(choice.get("function"), dict) else choice
name = fn_obj.get("name") if isinstance(fn_obj, dict) else None
if not isinstance(name, str) or not name.strip():
return choice
namespace = fn_obj.get("namespace", choice.get("namespace")) if isinstance(fn_obj, dict) else choice.get("namespace")
mapped_name = tool_registry.get_upstream_name(namespace, name) if tool_registry else name
return {"type": "function", "function": {"name": mapped_name}}


# ---------------------------------------------------------------------------
Expand All @@ -275,7 +445,9 @@ class ResponsesStreamConverter:

def __init__(self, model: str = "unknown", parallel_tool_calls: bool = True, *,
realtime: bool = False, budget: StreamOutputBudget | None = None,
tool_states: dict | None = None, declared_names=None):
tool_states: dict | None = None, declared_names=None,
tool_registry: ToolRegistry | dict | None = None,
tool_namespaces: dict[str, str] | None = None):
self.resp_id = _rand_id("resp_")
self.msg_id = _rand_id("msg_")
self.model = model
Expand All @@ -284,6 +456,13 @@ def __init__(self, model: str = "unknown", parallel_tool_calls: bool = True, *,
self._budget = budget if budget is not None else StreamOutputBudget(0)
self._tool_states = tool_states
self._local_tool_states: dict[int, dict] = {}
if isinstance(tool_registry, ToolRegistry):
self._tool_registry = tool_registry
elif isinstance(tool_registry, dict):
self._tool_registry = ToolRegistry.from_dict(tool_registry)
else:
self._tool_registry = None
self._tool_namespaces = dict(tool_namespaces or {})
self._declared_names = frozenset(
value for value in (declared_names or ()) if isinstance(value, str) and value)
self.created_at = int(time.time())
Expand Down Expand Up @@ -698,14 +877,26 @@ def _msg_item(self, status: str = "in_progress", empty: bool = False) -> dict:
}

def _fc_item(self, tc: dict, status: str, *, include_arguments: bool = True) -> dict:
return {
raw_name = tc["name"] or (tc.get("state", {}).get("name") or "")
name = raw_name
ns = tc.get("namespace")
if not ns:
if self._tool_registry:
ns, name = self._tool_registry.get_identity(raw_name)
elif self._tool_namespaces and raw_name in self._tool_namespaces:
ns = self._tool_namespaces[raw_name]

item: dict[str, Any] = {
"type": "function_call",
"id": tc["fc_id"],
"call_id": tc["id"] or (tc.get("state", {}).get("id") or ""),
"name": tc["name"] or (tc.get("state", {}).get("name") or ""),
"name": name,
"arguments": self._tool_arguments(tc) if include_arguments else "",
"status": status,
}
if ns:
item["namespace"] = ns
return item

def _response_obj(self, status: str, incomplete_reason: str | None = None) -> dict:
output = []
Expand Down
19 changes: 14 additions & 5 deletions converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -3644,10 +3644,15 @@ def attempt(routed, cred, headers, url):
async def _nonstream_adapted(url, headers, body, model_name, t0, rid, cred, *, anthropic=False,
payload=None, canonical=None, request=None, policy=None):
policy = policy or _snapshot_stream_policy("messages" if anthropic else "responses", body)
converter = (AnthropicStreamConverter(model=model_name) if anthropic else ResponsesStreamConverter(model=model_name, parallel_tool_calls=body.get("parallel_tool_calls", True)))
tool_registry = body.get("_tool_registry")
converter = (AnthropicStreamConverter(model=model_name) if anthropic else
ResponsesStreamConverter(model=model_name,
parallel_tool_calls=body.get("parallel_tool_calls", True),
tool_registry=tool_registry))

async def fetch(routed, cred, headers, url):
return await _fetch_checked_chat(url, headers, routed, model_name, rid, cred,
routed_body = {k: v for k, v in routed.items() if not k.startswith("_")}
return await _fetch_checked_chat(url, headers, routed_body, model_name, rid, cred,
filter_retry=True, max_collect_bytes=policy.max_collect_bytes)
try:
collected = await await_or_hangup(
Expand All @@ -3671,6 +3676,7 @@ async def _stream_adapted(url, headers, body, model_name, t0, rid, cred=None, *,
"""Map protocol events while sharing connection, aggregation and failure handling."""
protocol = "messages" if anthropic else "responses"
policy = body.pop(_REQUEST_POLICY_KEY, None) or _snapshot_stream_policy(protocol, body)
tool_registry = body.get("_tool_registry")
state = {}
tracker = None
declared_names = _declared_tool_names(body)
Expand All @@ -3684,14 +3690,17 @@ async def _stream_adapted(url, headers, body, model_name, t0, rid, cred=None, *,
ResponsesStreamConverter(model=model_name,
parallel_tool_calls=body.get("parallel_tool_calls", True),
realtime=True, budget=budget, tool_states=tracker.tools,
declared_names=declared_names))
declared_names=declared_names,
tool_registry=tool_registry))
else:
converter = (AnthropicStreamConverter(model=model_name) if anthropic else
ResponsesStreamConverter(model=model_name,
parallel_tool_calls=body.get("parallel_tool_calls", True)))
parallel_tool_calls=body.get("parallel_tool_calls", True),
tool_registry=tool_registry))
sent = False
upstream_body = {k: v for k, v in body.items() if not k.startswith("_")}
upstream = _chat_sse_lines(
url, headers, body, model_name, t0, rid, cred,
url, headers, upstream_body, model_name, t0, rid, cred,
policy=policy, tracker=tracker, state=state)
try:
try:
Expand Down
Loading
Loading