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.
This commit is contained in:
Roland Tannous 2026-04-12 14:49:15 +04:00
commit 86657f22cd
3 changed files with 601 additions and 10 deletions

View file

@ -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,
},
)

View file

@ -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())

View file

@ -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"