Skip to content

fix(chat): Claude Code/OpenCode correctness and /v1/messages latency (prefix cache 48→95%, overhead 2.7→0.3 s) - #2400

Open
jaluma wants to merge 23 commits into
mainfrom
feat/zylon-code
Open

jaluma wants to merge 23 commits into
mainfrom
feat/zylon-code

Conversation

@jaluma

@jaluma jaluma commented Oct 5, 2026 •

Copy link
Copy Markdown
Collaborator

Fixes found while serving Zylon Code (Claude Code / OpenCode through /v1/messages) from one DGX Spark. Measured on test3 with a 39-turn recorded Claude Code session (~36k-token prompts).

Correctness

  • Hidden client tool returned as "not found" (b067d826c). When an interceptor hides the tools for one iteration (loop detection replaces the context stack), a call to a tool the client declared produced a server-side "Tool 'bash' not found." tool_result, which is an invalid message for Anthropic clients. It now goes back to the caller as tool_use. Sync and async engines.
  • Inline system messages kept in place (9df383a75). Claude Code sends a role: system message after every user turn. extract_system_messages hoisted all of them into the system prompt, so the prompt prefix changed every turn and the engine's prefix cache only ever hit the tools block. Later system messages now become a MidConvSystemBlock on the nearest user message; leading ones still join the system list. Prefix-cache hit on the session: 48.6% → 94.6%.

Latency (per request, inside the cluster)

Change Before After
Condensation min-duration wait returns when the producer finishes (a87e45384) +1.5 s fixed 0
/v1/messages ends on the terminal event instead of a 1 s Redis timeout + 2 s status poll (9d86d8a19) +1.1 s ~0.01 s
Skip token counts a byte bound proves unnecessary; exact count near the limit (9d86d8a19) 3 remote tokenize calls skipped when far from the limit
arq poll 0.5 s → 0.05 s (PGPT_ARQ_POLL_DELAY), drop unused timeline snapshots, cache tool-schema models (7855dbdee) 0.25 s pickup, 278 ms CPU ~0.03 s, 143 ms CPU

End to end on the probe: 6.7 s → 3.5 s p50 (the rest is the model plus client RTT). zylon-gpt's own overhead went from ~2.7 s to ~0.3 s.

Throughput (chat worker CPU)

Under load the chat workers pinned one core each at 16 sessions while the engine sat idle. Every streamed token deep-copied the whole growing assistant message, so a request's CPU grew with the square of its output length.

  • Fold stream chunks into the assistant message in place, and extend the token-id list instead of rebuilding it (d0c32ed1d). Sync and async engines.
  • get_tool_calls_from_response (a dict lookup) is called inline instead of through a thread per chunk (d0c32ed1d).
  • BrokerEventChannel sends pending events as one pipelined RPUSH + EXPIRE, and the Redis listener drains queued events in one LPOP (d0c32ed1d).
  • The HTTP disconnect check runs once a second instead of on every event (d0c32ed1d); the per-event check was ~7% of API CPU.
  • The stream processor serializes events inline instead of through a thread per event (d1fddf41a).

On test3 (4 chat workers × 24, 1 API pod), together with the zylon-gpt context-copy fix, chat worker CPU at 8 sessions fell from 700–1000m to ~160m each.

A second round (f3a9de5e7 state fork, 7dc5c0dc1 direct /v1/messages stream) was pushed and then reverted (c54ab831a): on test3 it doubled client TTFT at 8 sessions (p95 3.9 → 12.4 s) with the engine unchanged. The branch is back to d1fddf41a in content; the regression is being bisected before anything from that round returns.

Tests

New regression tests for each change (each fails without its fix). tests/arq tests/engines tests/server/chat tests/components tests/context: 2,057 passed, 2 skipped. ruff, format and ty clean on changed files.

Notes: ChatState.timeline stays (empty) so stored states load. PGPT_ARQ_POLL_DELAY defaults to 0.05.

🤖 Generated with Claude Code

jaluma and others added 6 commits October 2, 2026 14:13
…a not-found tool_result

When an interceptor hides the tools for one iteration (loop detection replaces the
context stack), a model that still calls a client tool such as bash got a server-side
tool_result "Tool 'bash' not found.". That block is not valid in an Anthropic response,
so Anthropic clients (OpenCode) drop the session. If the request's original tools
declare the name as a client tool, stop with tool_use and let the caller execute it.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…to the system prompt

Claude Code sends a role=system message after each user turn. Moving all of them
into the system prompt changed the prompt prefix on every turn, so the engine's
prefix cache only ever hit the tools block (~50% on a 39-turn session vs ~94%
with the messages kept in place). Later system messages now become a
MidConvSystemBlock on the nearest user message; leading ones still join the
system list.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ducer finishes

_consume_and_emit_with_min_duration slept min_duration (tldr_minimum_threshold_seconds,
1.5 s by default) before reading the queue, and the interceptor awaits it at
BEFORE_ITERATION, so every request and every agent iteration paid 1.5 s even when no
condensation was needed. It now buffers for at most min_duration and returns as soon
as the producer's sentinel arrives.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…t, skip avoidable token counts

- Stream readers (direct and multiplexed) stop as soon as a batch ends with a terminal
  event (message_stop / error). Before, they waited for a blocking XREAD to time out
  (1 s) and the periodic status check (2 s), so a non-streaming request returned ~1 s
  after the worker had finished.
- The validator skips tokenizing the user text and system prompt when their UTF-8 byte
  count (plus special-token slack) already fits the token limit; a token covers at
  least one byte, so the exact count cannot exceed it.
- Condensation skips both token counts (system messages, whole conversation: remote
  chat-template renders on Triton) when a byte upper bound fits max_length, or when
  half the bytes fit within 3/4 of it. Near the limit the exact count still decides.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…els, poll the arq queue faster

Measured in-process on a 36k-token Claude Code request (LLM mocked): non-LLM CPU per
request 278 ms -> 143 ms.
- _snapshot deep-copied the whole chat state ~6 times per iteration to append a
  timeline entry that nothing reads; removed from both engines (ChatTimelineEntry and
  ChatState.timeline stay so stored states and fixtures still validate).
- create_model_from_json_schema rebuilt a pydantic model per tool on every request
  (~3 ms each, 18 tools); models are now cached by name + canonical schema (LRU, 512),
  and model_json_schema returns a copy so the shared schema cannot be mutated.
- arq worker poll_delay 0.5 s -> 0.05 s (PGPT_ARQ_POLL_DELAY): a chat job waited
  ~0.25 s on average before pickup.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@jaluma
jaluma marked this pull request as ready for review October 5, 2026 11:01
Copilot AI balanced review requested due to automatic review settings October 5, 2026 11:01
@jaluma jaluma self-assigned this Oct 5, 2026

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Unsafe token-count shortcuts and schema-cache isolation and concurrency defects can cause oversized prompts or request failures.

Review effort: Balanced
Findings: 1 High severity · 3 Medium severity

Open (4)
What changed in this PR

Improves /v1/messages correctness and latency for Claude Code/OpenCode workloads.

Changes:

  • Preserves inline system-message placement and client tool calls.
  • Reduces streaming, condensation, polling, tokenization, and schema-generation overhead.
  • Adds regression coverage for the optimized paths.
File Description
private_gpt/​arq/​runner.py Configures faster worker polling.
private_gpt/​chat/​input_models.py Preserves mid-conversation system messages.
private_gpt/​chat/​schema_models.py Adds bounded schema-model caching.
private_gpt/​components/​chat/​processors/​chat_history/​memory/​tldr_processor.py Adds tokenization shortcuts.
private_gpt/​components/​engines/​chat/​async_chat_engine.py Handles hidden client tools and removes snapshots.
private_gpt/​components/​engines/​chat/​chat_engine.py Mirrors synchronous engine changes.
private_gpt/​components/​engines/​chat/​models/​chat_state.py Documents retained timeline compatibility.
private_gpt/​components/​streaming/​stream/​stream_reader.py Stops readers on terminal events.
private_gpt/​server/​chat/​interceptors/​condensation_interceptor.py Ends buffering when production completes.
private_gpt/​server/​chat/​interceptors/​validator_request_interceptor.py Skips provably unnecessary tokenization.
tests/​arq/​test_runner.py Tests polling configuration.
tests/​components/​chat/​test_tldr_processor_token_count.py Tests condensation token shortcuts.
tests/​components/​streaming/​test_stream_multiplexer.py Tests terminal-event completion.
tests/​engines/​test_chat_agent_engine.py Tests hidden client tools.
tests/​server/​chat/​interceptors/​test_condensation_interceptor.py Tests early producer completion.
tests/​server/​chat/​interceptors/​test_validation_request_interceptor.py Tests validator shortcuts.
tests/​server/​chat/​test_input_models.py Tests inline system placement.
tests/​server/​chat/​test_schema_models.py Tests schema cache behavior.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread private_gpt/chat/schema_models.py Outdated
Comment on lines +403 to +406
cached = _model_cache.get(key)
if cached is not None:
_model_cache.move_to_end(key)
return cached
Comment thread private_gpt/chat/schema_models.py Outdated
Comment on lines +407 to +408
model = _build_model_from_json_schema(copy.deepcopy(schema), model_name)
_model_cache[key] = model
Comment on lines +116 to +119
if (
byte_bound <= max_length
or byte_bound // _MIN_BYTES_PER_TOKEN <= max_length * _HEURISTIC_FILL
):
# skip the tokenizer, which may be a remote call per text.
system_prompt = self._system_prompt_text(context, request)
text_bytes = len(user_text.encode()) + len((system_prompt or "").encode())
if text_bytes + _SPECIAL_TOKENS_SLACK <= token_limit:
@jaluma
jaluma requested a review from pabloogc October 5, 2026 11:06
jaluma and others added 16 commits October 5, 2026 22:59
…le disconnect checks

- Fold stream chunks into the assistant message in place; the per-chunk
  deep copy made worker CPU quadratic in output length.
- Call get_tool_calls_from_response inline (a dict lookup) instead of a
  thread hop per chunk.
- BrokerEventChannel publishes pending events in one pipelined RPUSH+EXPIRE;
  the Redis listener drains queued events in one LPOP.
- Check HTTP disconnection once a second instead of per event.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…token

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…r loop for enqueue

Unreviewed follow-up from the CPU review: ChatState.fork() shares message
contents between engine states. Check that no interceptor mutates content
blocks in place before merging.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The API relayed every token through a second Redis stream (XADD, XREAD,
re-parse) on top of the worker's event list: ~0.5 ms of API CPU per token.
/v1/messages now reads the engine events directly; client disconnects still
cancel the execution. The async /chat endpoints keep the observable stream.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
On test3 they doubled client TTFT at 8 sessions (p50 1.6 -> 3.0 s, p95 3.9 ->
12.5 s) with the engine unchanged (0.2 s TTFT, 88% prefix hit), so the extra
time is in Zylon. Back to d1fddf4 until the cause is found.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The API relayed every token through a second Redis stream (XADD, XREAD,
re-parse) on top of the worker's event list: ~0.5 ms of API CPU per token.
/v1/messages now reads the engine events directly; client disconnects still
cancel the execution. The async /chat endpoints keep the observable stream.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
On test3 (8 vCPU, 1 API pod) it raised TTFT p95 at 68 sessions from 9.5 s to 22.0 s.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
ARQ keeps a job in the queue zset until it finishes, so every poll returned
all in-flight chat jobs and start_jobs ran WATCH/EXISTS/ZSCORE on each
(~660 claims/s per worker at 48 sessions). New jobs waited behind that
loop and it cost ~7% of a core per worker.

Skip jobs this worker already runs and drop jobs locked elsewhere with one
pipelined EXISTS; the WATCH/MULTI claim is unchanged and stays the
authority.

Local stack (fake Triton 0.2 s TTFT, 1 ms RTT, 4 workers, 2 runs each):
  ARQ pickup p50 @48: 87 -> 12 ms
  client TTFT p50 @16/48/68: 0.31/0.88/1.27 -> 0.29/0.63/1.10 s
  worker CPU @16: 1.24-1.47 -> 0.88 cores; @48-68: 3.1-3.5 -> 2.6-2.7

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ight

StreamProcessor wrote every event with its own XADD + EXPIRE round trip,
so the API paid one Redis write (and the wait_for/Task overhead of
redis-py's socket timeout) per token. The first event is still written
immediately; events that arrive while a write is in flight are queued and
sent in the next pipeline (N XADD + 1 EXPIRE). No timer: at low load every
event goes out alone, under load the batch grows with the backlog.
The writer is flushed before the PROCESSING status update and awaited at
the end; a failed write surfaces as a stream error.

Local stack (1 API, 4 chat workers, fake Triton over 0.5 ms netem, open-loop
Poisson arrivals, 120 s per rate):
  5 req/s: TTFT p50/p95 0.51/1.37 -> 0.33/1.27 s, API 0.76 -> 0.68 core
  7 req/s: TTFT p50/p95 2.01/4.37 -> 0.55/1.64 s, e2e p50 22.0 -> 2.8 s
           (without batching the API saturates and streams back up),
           API 0.40 -> 0.33 ms CPU per delta, 0 errors.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
PingEventInterceptor wrapped every event in asyncio.wait_for(), which creates a
Task and a timeout handle per token on the API loop. It now buffers events in a
deque with one waiter future and arms a single call_later() ping timer only
while it is idle. Order, ping cadence and producer error propagation (raised
after the buffered events) are unchanged; the producer is cancelled on exit.

Local stack (1 API, 4 chat workers, fake Triton over 0.5 ms netem, open-loop
Poisson arrivals, 120 s), with and without this commit on top of ARQ-poll +
batched XADD:
  7 req/s: TTFT p50/p95 1.14/3.97 -> 0.55/1.64 s, e2e p50 7.0 -> 2.8 s.
  5 req/s: TTFT p95 1.75 -> 1.27 s.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A model that called a tool the request did not offer (Qwen3.8 called Grep when
Claude Code only offered Bash, Read, ...) got a server-side tool_result
"Tool 'Grep' not found." inside the assistant message. Tool results never appear
in an Anthropic assistant message; Claude Code and OpenCode reject the turn.

When the request carries any client tool, the caller owns the tool loop: return
the call as tool_use (stop_reason tool_use), so the caller answers with its own
error tool_result and the model self-corrects, as the Anthropic API does. A
request with only server tools keeps the internal not-found result.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ser turn

Claude Code sends a role=system message after every tool result. It became a
MidConvSystemBlock on the tool-result user message, which was then split out as
a standalone user turn. Templates that keep reasoning only after the last user
query (GLM-5.3 clear_thinking, DeepSeek drop_thinking, Qwen) therefore rendered
every earlier tool step of the same agent turn as <think></think>: the model saw
a history of empty reasoning and, on GLM-5.3-Flash, skipped thinking on 62% of
turns (87/230 responses had a thinking block in the h200 live run).

Append the inline system text to the nearest tool result instead (preceding,
else following). It keeps its position, so the prompt prefix stays stable for
the prefix cache, and no user turn is opened mid tool loop.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
enqueue_job opened a new ARQ pool for every chat request (connect, HELLO,
SELECT, PING), and publish_route opened a second one for a SETEX of a route
that has a 24 h TTL. Under load both queue behind the API event loop: at 48
sessions arrive -> enqueued took ~270 ms p50.

Keep one pool per running loop (pools of closed loops are dropped, for the
Celery tasks that run short-lived loops) and refresh the publisher route at
most once a minute; a failed publish is retried on the next enqueue.

Local stack (1 API, 4 chat workers, fake Triton over 0.5 ms netem, perf/ttft-stack
as base, 48 concurrent sessions, two runs each):
  enqueue start -> enqueued p50: 167 / 141 ms -> 32 / 27 ms
  TTFT p50: 1.25 / 1.27 s -> 1.18 / 1.06 s; tokens/s 3801 / 4139 -> 4234 / 4253
At 16 sessions it is within noise.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… inline system text

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The summary workflow computed max_tokens but only passed it to the query
engine's response_synthesizer_kwargs. TreeSummarizeSynthesizer forwards those
to LLM.apredict/astructured_predict as prompt-template args, never to
achat/acomplete, so the cap never reached the LLM. The Triton LLM then defaults
max_tokens to "context window - prompt" (~90k on a 131k Qwen3.6). On test3 a
condenser summary of a 189k-token Claude Code history ran for ~6 minutes at
~220 tok/s until the 300 s condensation timeout. During that time it starved the
Triton backend (O(n^2) stream split, fixed in zylon-chart), and two concurrent
OpenCode turns died with "Triton stream stalled for 60.00s".

- TreeSummarizeSynthesizer takes llm_kwargs and passes them to every LLM call
  (chat/complete and structured predict, sync and async).
- The summary workflow sends max_tokens=SUMMARY_MAX_OUTPUT_TOKENS (4096); this
  also replaces max(4000, num_output * 4) as the reserve for chunk packing.
- The condenser's structured left-side summary uses the same cap; the web-link
  relevance check ({"relevant": bool}) is capped at 64.
  Loop detection already passes max_tokens=128; the conversation-history
  condenser workflow already passes its own max_tokens.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
GLM-5.3 emitted a tool call whose name was the whole call
(`Agent(description="...", prompt="...") run_in_background=false ...`).
ToolUseBlock.name is max 200 chars, so building the tool_use block raised
string_too_long and the whole stream failed with invalid_request_error
(request 32cf7229-89ec-42e8-97a2-27207714da33).

Both chat engines now clamp the name once, where they read the LLM tool
calls (empty -> "unknown", > 200 chars truncated), so the block, the name
map and the tool result agree. Such a name matches no tool: the call is
answered "Tool '...' not found." and the model can retry. zylon-gpt
(fix/parser-tokenize-audit) recovers the common shapes into name + args.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- ping interceptor: cast the dequeued item to Event instead of a mypy-only
  ignore that ty does not honor
- schema model cache: guard the LRU with a lock (reached from to_thread),
  and return a copy of the root-array schema so cached models can't be
  mutated through model_json_schema()
- request validator: reserve special-token slack for both tokenizer calls

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants