From 86657f22cd175ed176d33ffff6a033266d404633 Mon Sep 17 00:00:00 2001 From: Roland Tannous Date: Sun, 12 Apr 2026 14:49:15 +0400 Subject: [PATCH] Split /v1/messages into server-side and client-side tool paths enable_tools=true runs the existing server-side agentic loop with built-in tools (web_search/python/terminal). A bare tools=[...] field now triggers a client-side pass-through: client-provided tools are forwarded to llama-server and any tool_use output is returned to the caller with stop_reason=tool_use for client execution. This fixes Claude Code (and any Anthropic SDK client) which sends tools=[...] expecting client-side execution but was previously routed through execute_tool() and failing with 'Unknown tool'. Adds AnthropicPassthroughEmitter to convert llama-server OpenAI SSE chunks into Anthropic SSE events, plus unit tests covering text blocks, tool_use blocks, mixed, stop reasons, and usage. --- .../core/inference/anthropic_compat.py | 185 ++++++++++++++ studio/backend/routes/inference.py | 235 +++++++++++++++++- .../backend/tests/test_anthropic_messages.py | 191 ++++++++++++++ 3 files changed, 601 insertions(+), 10 deletions(-) diff --git a/studio/backend/core/inference/anthropic_compat.py b/studio/backend/core/inference/anthropic_compat.py index d31b05389a..e7b40a60ce 100644 --- a/studio/backend/core/inference/anthropic_compat.py +++ b/studio/backend/core/inference/anthropic_compat.py @@ -301,3 +301,188 @@ class AnthropicStreamEmitter: "index": self.block_index, }, ) + + +class AnthropicPassthroughEmitter: + """Converts llama-server's OpenAI-format streaming chunks into Anthropic SSE. + + Used for the client-side tool-use pass-through path: the client (e.g. Claude + Code) sends its own tool definitions in the ``tools`` field and expects to + execute them itself. We forward them to llama-server and translate the + streaming response back to Anthropic format without executing anything. + """ + + def __init__(self) -> None: + self.block_index: int = -1 + self._current_block_type: Optional[str] = None # "text" | "tool_use" | None + self._tool_call_states: dict = {} # delta index -> {block_index, id, name} + self._usage: dict = {} + self._stop_reason: str = "end_turn" + + def start(self, message_id: str, model: str) -> list[str]: + return [ + build_anthropic_sse_event( + "message_start", + { + "type": "message_start", + "message": { + "id": message_id, + "type": "message", + "role": "assistant", + "content": [], + "model": model, + "stop_reason": None, + "stop_sequence": None, + "usage": {"input_tokens": 0, "output_tokens": 0}, + }, + }, + ) + ] + + def feed_chunk(self, chunk: dict) -> list[str]: + """Process one OpenAI streaming chat.completion.chunk.""" + events: list[str] = [] + + # usage-only chunks carry token totals + usage = chunk.get("usage") + if usage: + self._usage = usage + + choices = chunk.get("choices") or [] + if not choices: + return events + + choice = choices[0] + delta = choice.get("delta") or {} + finish_reason = choice.get("finish_reason") + + # ── Text content ── + content = delta.get("content") + if content: + if self._current_block_type != "text": + if self._current_block_type is not None: + events.append(self._close_current_block()) + events.extend(self._open_text_block()) + events.append( + build_anthropic_sse_event( + "content_block_delta", + { + "type": "content_block_delta", + "index": self.block_index, + "delta": {"type": "text_delta", "text": content}, + }, + ) + ) + + # ── Tool calls (streaming deltas) ── + tool_calls = delta.get("tool_calls") or [] + for tc in tool_calls: + tc_idx = tc.get("index", 0) + fn = tc.get("function") or {} + if tc_idx not in self._tool_call_states: + # New tool call — close prior block, open tool_use block + if self._current_block_type is not None: + events.append(self._close_current_block()) + tc_id = tc.get("id", "") + tc_name = fn.get("name", "") + self.block_index += 1 + self._current_block_type = "tool_use" + self._tool_call_states[tc_idx] = { + "block_index": self.block_index, + "id": tc_id, + "name": tc_name, + } + events.append( + build_anthropic_sse_event( + "content_block_start", + { + "type": "content_block_start", + "index": self.block_index, + "content_block": { + "type": "tool_use", + "id": tc_id, + "name": tc_name, + "input": {}, + }, + }, + ) + ) + + args_delta = fn.get("arguments", "") + if args_delta: + events.append( + build_anthropic_sse_event( + "content_block_delta", + { + "type": "content_block_delta", + "index": self._tool_call_states[tc_idx]["block_index"], + "delta": { + "type": "input_json_delta", + "partial_json": args_delta, + }, + }, + ) + ) + + # ── Finish reason ── + if finish_reason: + if finish_reason == "tool_calls": + self._stop_reason = "tool_use" + elif finish_reason == "length": + self._stop_reason = "max_tokens" + else: + self._stop_reason = "end_turn" + + return events + + def finish(self) -> list[str]: + events: list[str] = [] + if self._current_block_type is not None: + events.append(self._close_current_block()) + events.append( + build_anthropic_sse_event( + "message_delta", + { + "type": "message_delta", + "delta": { + "stop_reason": self._stop_reason, + "stop_sequence": None, + }, + "usage": { + "output_tokens": self._usage.get("completion_tokens", 0), + }, + }, + ) + ) + events.append( + build_anthropic_sse_event( + "message_stop", + {"type": "message_stop"}, + ) + ) + return events + + def _open_text_block(self) -> list[str]: + self.block_index += 1 + self._current_block_type = "text" + return [ + build_anthropic_sse_event( + "content_block_start", + { + "type": "content_block_start", + "index": self.block_index, + "content_block": {"type": "text", "text": ""}, + }, + ) + ] + + def _close_current_block(self) -> str: + idx = self.block_index + self._current_block_type = None + return build_anthropic_sse_event( + "content_block_stop", + { + "type": "content_block_stop", + "index": idx, + }, + ) diff --git a/studio/backend/routes/inference.py b/studio/backend/routes/inference.py index 281ad4d9e6..8082e4d834 100644 --- a/studio/backend/routes/inference.py +++ b/studio/backend/routes/inference.py @@ -106,6 +106,7 @@ from core.inference.anthropic_compat import ( anthropic_messages_to_openai, anthropic_tools_to_openai, AnthropicStreamEmitter, + AnthropicPassthroughEmitter, ) from auth.authentication import get_current_subject @@ -2254,20 +2255,53 @@ async def anthropic_messages( cancel_event = threading.Event() - # ── Tool-calling path ───────────────────────────────────── - # Two ways to enable tools: - # 1. Anthropic-style: send full tool definitions in payload.tools - # 2. Unsloth shorthand: enable_tools=true + optional enabled_tools list - use_tools = llama_backend.supports_tools and ( - (payload.tools and len(payload.tools) > 0) or payload.enable_tools + # ── Tool routing ────────────────────────────────────────── + # Three paths: + # 1. enable_tools=true → server-side execution of built-in tools (Unsloth shorthand) + # 2. tools=[...] only → client-side pass-through (standard Anthropic behavior) + # 3. neither → plain chat + server_tools = payload.enable_tools and llama_backend.supports_tools + client_tools = ( + not server_tools + and payload.tools + and len(payload.tools) > 0 + and llama_backend.supports_tools ) - if use_tools: + # ── Client-side pass-through path ───────────────────────── + if client_tools: + openai_tools = anthropic_tools_to_openai(payload.tools) + + if payload.stream: + return await _anthropic_passthrough_stream( + request, + cancel_event, + llama_backend, + openai_messages, + openai_tools, + temperature, + top_p, + top_k, + payload.max_tokens, + message_id, + model_name, + ) + return await _anthropic_passthrough_non_streaming( + llama_backend, + openai_messages, + openai_tools, + temperature, + top_p, + top_k, + payload.max_tokens, + message_id, + model_name, + ) + + if server_tools: from core.inference.tools import ALL_TOOLS - if payload.tools and len(payload.tools) > 0: - openai_tools = anthropic_tools_to_openai(payload.tools) - elif payload.enabled_tools is not None: + if payload.enabled_tools is not None: openai_tools = [ t for t in ALL_TOOLS if t["function"]["name"] in payload.enabled_tools ] @@ -2567,3 +2601,184 @@ async def _anthropic_plain_non_streaming(run_gen, message_id, model_name): ), ) return JSONResponse(content = resp.model_dump()) + + +# ===================================================================== +# Client-side tool pass-through (Anthropic-native tools field) +# ===================================================================== + + +def _build_passthrough_payload( + openai_messages, + openai_tools, + temperature, + top_p, + top_k, + max_tokens, + stream, +): + body = { + "messages": openai_messages, + "tools": openai_tools, + "tool_choice": "auto", + "temperature": temperature, + "top_p": top_p, + "top_k": top_k, + "stream": stream, + } + if stream: + body["stream_options"] = {"include_usage": True} + if max_tokens is not None: + body["max_tokens"] = max_tokens + return body + + +async def _anthropic_passthrough_stream( + request, + cancel_event, + llama_backend, + openai_messages, + openai_tools, + temperature, + top_p, + top_k, + max_tokens, + message_id, + model_name, +): + """Streaming client-side pass-through: forward tools to llama-server and + translate its streaming response to Anthropic SSE without executing anything.""" + target_url = f"{llama_backend.base_url}/v1/chat/completions" + body = _build_passthrough_payload( + openai_messages, + openai_tools, + temperature, + top_p, + top_k, + max_tokens, + True, + ) + + async def _stream(): + emitter = AnthropicPassthroughEmitter() + for line in emitter.start(message_id, model_name): + yield line + + try: + async with httpx.AsyncClient() as client: + async with client.stream( + "POST", + target_url, + json = body, + timeout = 600, + ) as resp: + async for raw_line in resp.aiter_lines(): + if await request.is_disconnected(): + cancel_event.set() + return + if not raw_line or not raw_line.startswith("data: "): + continue + data_str = raw_line[6:] + if data_str.strip() == "[DONE]": + break + try: + chunk = json.loads(data_str) + except json.JSONDecodeError: + continue + for line in emitter.feed_chunk(chunk): + yield line + except Exception as e: + logger.error("anthropic_messages passthrough stream error: %s", e) + + for line in emitter.finish(): + yield line + + return StreamingResponse( + _stream(), + media_type = "text/event-stream", + headers = { + "Cache-Control": "no-cache", + "Connection": "keep-alive", + "X-Accel-Buffering": "no", + }, + ) + + +async def _anthropic_passthrough_non_streaming( + llama_backend, + openai_messages, + openai_tools, + temperature, + top_p, + top_k, + max_tokens, + message_id, + model_name, +): + """Non-streaming client-side pass-through.""" + target_url = f"{llama_backend.base_url}/v1/chat/completions" + body = _build_passthrough_payload( + openai_messages, + openai_tools, + temperature, + top_p, + top_k, + max_tokens, + False, + ) + + async with httpx.AsyncClient() as client: + resp = await client.post(target_url, json = body, timeout = 600) + + if resp.status_code != 200: + raise HTTPException( + status_code = resp.status_code, + detail = f"llama-server error: {resp.text[:500]}", + ) + + data = resp.json() + choice = (data.get("choices") or [{}])[0] + message = choice.get("message") or {} + finish_reason = choice.get("finish_reason") + + content_blocks = [] + text = message.get("content") or "" + if text: + text = _TOOL_XML_RE.sub("", text).strip() + if text: + content_blocks.append(AnthropicResponseTextBlock(text = text)) + + tool_calls = message.get("tool_calls") or [] + for tc in tool_calls: + fn = tc.get("function") or {} + try: + args = json.loads(fn.get("arguments", "{}")) + except json.JSONDecodeError: + args = {} + content_blocks.append( + AnthropicResponseToolUseBlock( + id = tc.get("id", ""), + name = fn.get("name", ""), + input = args, + ) + ) + + if tool_calls: + stop_reason = "tool_use" + elif finish_reason == "length": + stop_reason = "max_tokens" + else: + stop_reason = "end_turn" + + usage = data.get("usage") or {} + resp_obj = AnthropicMessagesResponse( + id = message_id, + model = model_name, + content = content_blocks, + stop_reason = stop_reason, + usage = AnthropicUsage( + input_tokens = usage.get("prompt_tokens", 0), + output_tokens = usage.get("completion_tokens", 0), + ), + ) + return JSONResponse(content = resp_obj.model_dump()) diff --git a/studio/backend/tests/test_anthropic_messages.py b/studio/backend/tests/test_anthropic_messages.py index 7042c21e79..501bc74a92 100644 --- a/studio/backend/tests/test_anthropic_messages.py +++ b/studio/backend/tests/test_anthropic_messages.py @@ -30,6 +30,7 @@ from core.inference.anthropic_compat import ( anthropic_tools_to_openai, build_anthropic_sse_event, AnthropicStreamEmitter, + AnthropicPassthroughEmitter, ) @@ -486,3 +487,193 @@ class TestAnthropicStreamEmitter: events = e.feed({"type": "content", "text": "After tool"}) parsed = json.loads(events[0].split("data: ")[1]) assert parsed["delta"]["text"] == "After tool" + + +# ===================================================================== +# Pass-through emitter tests (client-side tool execution path) +# ===================================================================== + + +class TestAnthropicPassthroughEmitter: + + def _parse(self, event_str): + return json.loads(event_str.split("data: ")[1]) + + def test_start_emits_message_start_only(self): + e = AnthropicPassthroughEmitter() + events = e.start("msg_1", "test-model") + assert len(events) == 1 + assert "message_start" in events[0] + parsed = self._parse(events[0]) + assert parsed["message"]["id"] == "msg_1" + assert parsed["message"]["model"] == "test-model" + + def test_text_chunk_opens_text_block_and_emits_delta(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + chunk = {"choices": [{"delta": {"content": "Hello"}}]} + events = e.feed_chunk(chunk) + # content_block_start + content_block_delta + assert len(events) == 2 + assert "content_block_start" in events[0] + assert '"type": "text"' in events[0] + delta = self._parse(events[1]) + assert delta["delta"]["type"] == "text_delta" + assert delta["delta"]["text"] == "Hello" + + def test_sequential_text_chunks_single_block(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + events1 = e.feed_chunk({"choices": [{"delta": {"content": "Hello"}}]}) + events2 = e.feed_chunk({"choices": [{"delta": {"content": " world"}}]}) + # First chunk opens the block, second only emits delta + assert len(events1) == 2 + assert len(events2) == 1 + assert self._parse(events2[0])["delta"]["text"] == " world" + + def test_tool_call_opens_tool_use_block(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + chunk = { + "choices": [{ + "delta": { + "tool_calls": [{ + "index": 0, + "id": "call_1", + "type": "function", + "function": {"name": "Bash", "arguments": ""}, + }] + } + }] + } + events = e.feed_chunk(chunk) + assert len(events) == 1 + parsed = self._parse(events[0]) + assert parsed["type"] == "content_block_start" + assert parsed["content_block"]["type"] == "tool_use" + assert parsed["content_block"]["id"] == "call_1" + assert parsed["content_block"]["name"] == "Bash" + + def test_tool_call_arguments_streamed_as_input_json_delta(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + # Open the tool call + e.feed_chunk({"choices": [{"delta": {"tool_calls": [ + {"index": 0, "id": "c1", "type": "function", + "function": {"name": "Bash", "arguments": ""}} + ]}}]}) + # Stream argument fragments + events1 = e.feed_chunk({"choices": [{"delta": {"tool_calls": [ + {"index": 0, "function": {"arguments": "{\"cmd"}} + ]}}]}) + events2 = e.feed_chunk({"choices": [{"delta": {"tool_calls": [ + {"index": 0, "function": {"arguments": "\": \"ls\"}"}} + ]}}]}) + parsed1 = self._parse(events1[0]) + parsed2 = self._parse(events2[0]) + assert parsed1["delta"]["type"] == "input_json_delta" + assert parsed1["delta"]["partial_json"] == "{\"cmd" + assert parsed2["delta"]["partial_json"] == "\": \"ls\"}" + + def test_text_then_tool_closes_text_block(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + e.feed_chunk({"choices": [{"delta": {"content": "Let me check."}}]}) + events = e.feed_chunk({"choices": [{"delta": {"tool_calls": [ + {"index": 0, "id": "c1", "type": "function", + "function": {"name": "Bash", "arguments": ""}} + ]}}]}) + # Should close text block and open tool_use block + assert "content_block_stop" in events[0] + assert "content_block_start" in events[1] + assert '"type": "tool_use"' in events[1] + + def test_finish_reason_tool_calls_sets_tool_use_stop(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + e.feed_chunk({"choices": [{"delta": {"tool_calls": [ + {"index": 0, "id": "c1", "type": "function", + "function": {"name": "Bash", "arguments": "{}"}} + ]}}]}) + e.feed_chunk({"choices": [{"delta": {}, "finish_reason": "tool_calls"}]}) + events = e.finish() + delta_event = [ev for ev in events if "message_delta" in ev][0] + parsed = self._parse(delta_event) + assert parsed["delta"]["stop_reason"] == "tool_use" + + def test_finish_reason_stop_sets_end_turn(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + e.feed_chunk({"choices": [{"delta": {"content": "Hi"}}]}) + e.feed_chunk({"choices": [{"delta": {}, "finish_reason": "stop"}]}) + events = e.finish() + delta_event = [ev for ev in events if "message_delta" in ev][0] + parsed = self._parse(delta_event) + assert parsed["delta"]["stop_reason"] == "end_turn" + + def test_finish_reason_length_sets_max_tokens(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + e.feed_chunk({"choices": [{"delta": {"content": "Hi"}}]}) + e.feed_chunk({"choices": [{"delta": {}, "finish_reason": "length"}]}) + events = e.finish() + delta_event = [ev for ev in events if "message_delta" in ev][0] + parsed = self._parse(delta_event) + assert parsed["delta"]["stop_reason"] == "max_tokens" + + def test_finish_closes_current_block(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + e.feed_chunk({"choices": [{"delta": {"content": "Hi"}}]}) + events = e.finish() + assert "content_block_stop" in events[0] + assert "message_delta" in events[1] + assert "message_stop" in events[2] + + def test_usage_chunk_captured(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + e.feed_chunk({"choices": [{"delta": {"content": "Hi"}}]}) + e.feed_chunk({ + "choices": [], + "usage": {"prompt_tokens": 10, "completion_tokens": 5}, + }) + events = e.finish() + delta_event = [ev for ev in events if "message_delta" in ev][0] + parsed = self._parse(delta_event) + assert parsed["usage"]["output_tokens"] == 5 + + def test_empty_chunk_returns_no_events(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + events = e.feed_chunk({"choices": []}) + assert events == [] + + def test_no_blocks_at_all_still_produces_valid_finish(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + events = e.finish() + # No content_block_stop because no block was opened + assert not any("content_block_stop" in ev for ev in events) + assert any("message_delta" in ev for ev in events) + assert any("message_stop" in ev for ev in events) + + def test_multiple_tool_calls_distinct_blocks(self): + e = AnthropicPassthroughEmitter() + e.start("msg_1", "m") + # First tool call + e.feed_chunk({"choices": [{"delta": {"tool_calls": [ + {"index": 0, "id": "c1", "type": "function", + "function": {"name": "Bash", "arguments": "{}"}} + ]}}]}) + # Second tool call (different index) + events = e.feed_chunk({"choices": [{"delta": {"tool_calls": [ + {"index": 1, "id": "c2", "type": "function", + "function": {"name": "Read", "arguments": "{}"}} + ]}}]}) + # Should close block 0, open block 1 + assert "content_block_stop" in events[0] + assert "content_block_start" in events[1] + parsed = self._parse(events[1]) + assert parsed["content_block"]["name"] == "Read" + assert parsed["content_block"]["id"] == "c2"