Add Anthropic live tool continuation

This commit is contained in:
antirez
2026-05-14 12:29:59 +02:00
parent 646798fd1c
commit 43535e1393
2 changed files with 630 additions and 84 deletions
+501 -84
View File
@@ -545,6 +545,9 @@ typedef struct {
char *content;
char *reasoning;
char *tool_call_id;
char **tool_call_ids;
int tool_call_ids_len;
int tool_call_ids_cap;
tool_calls calls;
} chat_msg;
@@ -569,7 +572,9 @@ typedef struct {
static void stop_list_clear(stop_list *stops);
static bool id_list_contains(const stop_list *ids, const char *id);
static void id_list_push_unique(stop_list *ids, const char *id);
static void id_list_free(stop_list *ids);
static bool responses_live_has_call_id(server *s, const char *id);
static bool anthropic_live_has_call_id(server *s, const char *id);
typedef struct {
req_kind kind;
@@ -614,6 +619,9 @@ typedef struct {
bool responses_requires_live_reasoning;
stop_list responses_live_call_ids;
char *responses_live_suffix_text;
bool anthropic_requires_live_tool_state;
stop_list anthropic_live_call_ids;
char *anthropic_live_suffix_text;
tool_replay_stats tool_replay;
} request;
@@ -639,11 +647,27 @@ static void tool_calls_push(tool_calls *calls, tool_call tc) {
calls->v[calls->len++] = tc;
}
static void chat_msg_add_tool_call_id(chat_msg *m, const char *id) {
if (!m || !id || !id[0]) return;
if (!m->tool_call_id) m->tool_call_id = xstrdup(id);
for (int i = 0; i < m->tool_call_ids_len; i++) {
if (m->tool_call_ids[i] && !strcmp(m->tool_call_ids[i], id)) return;
}
if (m->tool_call_ids_len == m->tool_call_ids_cap) {
m->tool_call_ids_cap = m->tool_call_ids_cap ? m->tool_call_ids_cap * 2 : 2;
m->tool_call_ids = xrealloc(m->tool_call_ids,
(size_t)m->tool_call_ids_cap * sizeof(m->tool_call_ids[0]));
}
m->tool_call_ids[m->tool_call_ids_len++] = xstrdup(id);
}
static void chat_msg_free(chat_msg *m) {
free(m->role);
free(m->content);
free(m->reasoning);
free(m->tool_call_id);
for (int i = 0; i < m->tool_call_ids_len; i++) free(m->tool_call_ids[i]);
free(m->tool_call_ids);
tool_calls_free(&m->calls);
memset(m, 0, sizeof(*m));
}
@@ -735,6 +759,9 @@ static void request_free(request *r) {
stop_list_clear(&r->responses_live_call_ids);
free(r->responses_live_call_ids.v);
free(r->responses_live_suffix_text);
stop_list_clear(&r->anthropic_live_call_ids);
free(r->anthropic_live_call_ids.v);
free(r->anthropic_live_suffix_text);
tool_schema_orders_free(&r->tool_orders);
memset(r, 0, sizeof(*r));
}
@@ -1588,11 +1615,13 @@ static bool parse_messages(const char **p, chat_msgs *msgs) {
goto fail;
}
} else if (!strcmp(key, "tool_call_id")) {
free(msg.tool_call_id);
if (!json_string(p, &msg.tool_call_id)) {
char *id = NULL;
if (!json_string(p, &id)) {
free(key);
goto fail;
}
chat_msg_add_tool_call_id(&msg, id);
free(id);
} else if (!strcmp(key, "tool_calls")) {
tool_calls_free(&msg.calls);
if (!parse_tool_calls_value(p, &msg.calls)) {
@@ -1721,6 +1750,7 @@ static bool parse_anthropic_content_block(const char **p, const char *role, chat
tc.arguments = input ? xstrdup(input) : xstrdup("{}");
tool_calls_push(&msg->calls, tc);
} else if (type && !strcmp(type, "tool_result")) {
chat_msg_add_tool_call_id(msg, id);
buf b = {0};
buf_puts(&b, msg->content ? msg->content : "");
buf_puts(&b, "<tool_result>");
@@ -2286,14 +2316,14 @@ static char *render_chat_prompt_text(const chat_msgs *msgs, const char *tool_sch
}
/* Render only the semantic tail that must be appended to the live KV for a
* Responses continuation.
* tool-result continuation.
*
* In the common Codex/pi tool path, the previous assistant tool-call turn is
* already in the model session, including hidden thinking. The next request
* only provides function_call_output/custom_tool_call_output items, sometimes
* inside a full visible replay and sometimes as a short tail. Re-rendering the
* assistant call here would duplicate it and destroy cache alignment, so this
* function starts at the first new item and emits only:
* In the common agent tool path, the previous assistant tool-call turn is
* already in the model session, including hidden thinking and exact sampled
* DSML. The next request provides only the tool results, either as OpenAI
* Responses tool-output items or Anthropic user content blocks. Re-rendering
* the assistant call here would duplicate it and destroy cache alignment, so
* this function starts at the first new item and emits only:
*
* previous EOS, tool results, and the next assistant prefix.
*
@@ -2301,8 +2331,8 @@ static char *render_chat_prompt_text(const chat_msgs *msgs, const char *tool_sch
* suffix tokenization happens later after the cache decision, using the live
* token prefix as the boundary. That avoids BPE merges across the visible
* replay/live-KV boundary. */
static char *render_responses_live_tail(const chat_msgs *msgs, int start,
ds4_think_mode think_mode) {
static char *render_live_tool_tail(const chat_msgs *msgs, int start,
ds4_think_mode think_mode) {
const bool think = ds4_think_mode_enabled(think_mode);
buf out = {0};
buf_puts(&out, "<end▁of▁sentence>");
@@ -2359,6 +2389,14 @@ static bool chat_msg_has_call_id(const chat_msg *m, const char *id) {
return false;
}
static void chat_msg_collect_tool_call_ids(const chat_msg *m, stop_list *ids) {
if (!m || !ids) return;
id_list_push_unique(ids, m->tool_call_id);
for (int i = 0; i < m->tool_call_ids_len; i++) {
id_list_push_unique(ids, m->tool_call_ids[i]);
}
}
static const chat_msg *responses_find_prior_call_msg(const chat_msgs *msgs,
int before,
const char *id) {
@@ -2395,25 +2433,31 @@ static bool responses_validate_tool_outputs(server *s, const chat_msgs *msgs,
for (int i = 0; i < msgs->len; i++) {
const chat_msg *m = &msgs->v[i];
if (strcmp(m->role, "tool") && strcmp(m->role, "function")) continue;
if (!m->tool_call_id || !m->tool_call_id[0]) continue;
const bool live_known = responses_live_has_call_id(s, m->tool_call_id);
const chat_msg *prior = responses_find_prior_call_msg(msgs, i, m->tool_call_id);
if (!live_known && !prior) {
snprintf(err, errlen,
"unknown tool output call_id %s; replay full Responses history",
m->tool_call_id);
return false;
}
if (!prior) {
if (requires_live_tool_state) *requires_live_tool_state = true;
continue;
}
if (needs_reasoning &&
(!prior->reasoning || !prior->reasoning[0]))
{
if (requires_live_reasoning) *requires_live_reasoning = true;
stop_list ids = {0};
chat_msg_collect_tool_call_ids(m, &ids);
for (int j = 0; j < ids.len; j++) {
const char *id = ids.v[j];
const bool live_known = responses_live_has_call_id(s, id);
const chat_msg *prior = responses_find_prior_call_msg(msgs, i, id);
if (!live_known && !prior) {
snprintf(err, errlen,
"unknown tool output call_id %s; replay full Responses history",
id);
id_list_free(&ids);
return false;
}
if (!prior) {
if (requires_live_tool_state) *requires_live_tool_state = true;
continue;
}
if (needs_reasoning &&
(!prior->reasoning || !prior->reasoning[0]))
{
if (requires_live_reasoning) *requires_live_reasoning = true;
}
}
id_list_free(&ids);
}
return true;
}
@@ -2446,14 +2490,92 @@ static void responses_prepare_live_continuation(request *r,
}
} else {
for (int i = tail_start; i < msgs->len; i++) {
id_list_push_unique(&r->responses_live_call_ids, msgs->v[i].tool_call_id);
chat_msg_collect_tool_call_ids(&msgs->v[i], &r->responses_live_call_ids);
}
}
if (r->responses_live_call_ids.len == 0) return;
free(r->responses_live_suffix_text);
r->responses_live_suffix_text =
render_responses_live_tail(msgs, tail_start, r->think_mode);
render_live_tool_tail(msgs, tail_start, r->think_mode);
}
static bool anthropic_msg_is_tool_result_tail(const chat_msg *m) {
return m && !strcmp(m->role, "user") &&
((m->tool_call_id && m->tool_call_id[0]) ||
m->tool_call_ids_len > 0);
}
/* Validate Anthropic tool results before rendering.
*
* A tool_result.tool_use_id is valid if it is either still bound to the live
* Anthropic assistant tool-call frontier or the same request replays the prior
* assistant tool_use block. The first case is the fast path: keep the sampled
* KV and append only the tool-result suffix. The second case is a normal
* stateless replay, where exact DSML tool memory can restore the sampled tool
* bytes before prefix matching. A tool-result-only request with an unknown
* live id has no safe prefix to reconstruct, so report a clear client error. */
static bool anthropic_validate_tool_results(server *s, const chat_msgs *msgs,
bool *requires_live_tool_state,
char *err, size_t errlen) {
if (requires_live_tool_state) *requires_live_tool_state = false;
if (!msgs) return true;
for (int i = 0; i < msgs->len; i++) {
const chat_msg *m = &msgs->v[i];
if (!anthropic_msg_is_tool_result_tail(m)) continue;
stop_list ids = {0};
chat_msg_collect_tool_call_ids(m, &ids);
for (int j = 0; j < ids.len; j++) {
const char *id = ids.v[j];
const bool live_known = anthropic_live_has_call_id(s, id);
const chat_msg *prior = responses_find_prior_call_msg(msgs, i, id);
if (!live_known && !prior) {
snprintf(err, errlen,
"unknown Anthropic tool_result tool_use_id %s; replay full messages",
id);
id_list_free(&ids);
return false;
}
if (!prior && requires_live_tool_state) {
*requires_live_tool_state = true;
}
}
id_list_free(&ids);
}
return true;
}
/* Prepare the Anthropic live-tool fast path.
*
* Anthropic's visible replay normally includes the assistant tool_use JSON and
* the user tool_result. That replay is still only a description of what the
* model sampled. If the incoming tool_result IDs match the live sampled
* frontier, generate_job() can skip replay matching entirely and append just
* EOS + tool_result + next assistant prefix to the real KV. */
static void anthropic_prepare_live_continuation(request *r,
const chat_msgs *msgs) {
if (!r || r->api != API_ANTHROPIC || !msgs || msgs->len == 0) return;
int tail_end = msgs->len;
while (tail_end > 0 && role_is_system(msgs->v[tail_end - 1].role)) tail_end--;
int tail_start = tail_end;
while (tail_start > 0 &&
anthropic_msg_is_tool_result_tail(&msgs->v[tail_start - 1]))
{
tail_start--;
}
if (tail_start == tail_end) return;
stop_list_clear(&r->anthropic_live_call_ids);
for (int i = tail_start; i < msgs->len; i++) {
chat_msg_collect_tool_call_ids(&msgs->v[i], &r->anthropic_live_call_ids);
}
if (r->anthropic_live_call_ids.len == 0) return;
free(r->anthropic_live_suffix_text);
r->anthropic_live_suffix_text =
render_live_tool_tail(msgs, tail_start, r->think_mode);
}
/* The API parsers are intentionally selective JSON parsers: they keep only
@@ -2809,8 +2931,19 @@ static bool parse_anthropic_request(ds4_engine *e, server *s, const char *body,
if (!got_thinking && model_alias_enables_thinking(r->model)) thinking_enabled = true;
r->think_mode = ds4_think_mode_for_context(
think_mode_from_enabled(thinking_enabled, reasoning_effort), ctx_size);
if (!anthropic_validate_tool_results(s, &msgs,
&r->anthropic_requires_live_tool_state,
err, errlen))
{
chat_msgs_free(&msgs);
free(system);
free(tool_schemas);
request_free(r);
return false;
}
kv_cache_restore_tool_memory_for_messages(s, &msgs);
tool_memory_attach_to_messages(s, &msgs, &r->tool_replay);
anthropic_prepare_live_continuation(r, &msgs);
const char *active_tool_schemas = r->has_tools ? tool_schemas : NULL;
r->prompt_preserves_reasoning =
chat_history_uses_tool_context(&msgs, active_tool_schemas);
@@ -3276,9 +3409,7 @@ item_fail:
msg.content = output ? output : xstrdup("");
output = NULL;
if (call_id || item_id) {
msg.tool_call_id = call_id ? call_id : item_id;
if (call_id) call_id = NULL;
else item_id = NULL;
chat_msg_add_tool_call_id(&msg, call_id ? call_id : item_id);
}
chat_msgs_push(msgs, msg);
} else if (!strcmp(t, "reasoning")) {
@@ -3371,9 +3502,7 @@ item_fail:
tools_json ? tools_json : "";
msg.content = xstrdup(body);
if (call_id || item_id) {
msg.tool_call_id = call_id ? call_id : item_id;
if (call_id) call_id = NULL;
else item_id = NULL;
chat_msg_add_tool_call_id(&msg, call_id ? call_id : item_id);
}
chat_msgs_push(msgs, msg);
} else if (!is_bookkeeping) {
@@ -6943,23 +7072,20 @@ typedef struct {
typedef struct {
bool valid;
/* Token frontier of the live model session after the last successful
* Responses assistant turn. Continuing from this point preserves hidden
* thinking and sampled DSML bytes that are not necessarily present in the
* client-visible replay. */
/* Token frontier of a live assistant tool-call turn. Continuing from this
* point preserves hidden thinking and sampled DSML bytes that are not
* necessarily present in the client-visible replay. */
int live_tokens;
/* Rendered conversation text that the client is expected to replay. This
* is deliberately the visible protocol transcript, not a dump of live KV:
* when reasoning summaries are not requested, hidden thinking is omitted.
* If a later request starts with this text, DS4 can continue from
* live_tokens and tokenize only the new suffix. */
/* Optional rendered conversation text that the client is expected to replay.
* Responses uses this because visible replay can omit hidden reasoning.
* Anthropic currently uses only the call-id side of the state. */
char *visible_text;
size_t visible_len;
/* Tool-call ids generated at the same live frontier. A following
* function_call_output/custom_tool_call_output for these ids is a direct
* protocol continuation and should not trigger prompt-prefix matching. */
/* Tool-call ids generated at the same live frontier. A following tool
* result for these ids is a direct protocol continuation and should not
* trigger prompt-prefix matching or checkpoint canonicalization. */
stop_list call_ids;
} responses_live_state;
} live_tool_state;
typedef struct {
bool valid;
@@ -6981,7 +7107,8 @@ struct server {
int default_tokens;
kv_disk_cache kv;
tool_memory tool_mem;
responses_live_state responses_live;
live_tool_state responses_live;
live_tool_state anthropic_live;
visible_live_state thinking_live;
bool disable_exact_dsml_tool_replay;
pthread_mutex_t tool_mu;
@@ -7198,15 +7325,14 @@ static void tool_memory_free(tool_memory *m) {
memset(m, 0, sizeof(*m));
}
/* Single live Responses state.
/* Single live protocol-tool state.
*
* This is not an implementation of persistent previous_response_id or the
* Conversations API. DS4 currently rejects those durable-state parameters.
* The state below is only an in-memory optimization for the current local
* session, where the next request is expected to arrive immediately from the
* same agent loop. If it does not match, DS4 falls back to the same prefix and
* disk-cache machinery used by chat/completions. */
static void responses_live_clear_locked(responses_live_state *st) {
* This is not an implementation of durable remote conversation storage. It is
* only an in-memory binding from protocol tool-call IDs to the current sampled
* KV frontier. If it does not match, DS4 falls back to the same prefix and
* disk-cache machinery used by chat/completions, or returns a clear error for
* tool-result-only requests that have no replayable prefix. */
static void live_tool_state_clear_locked(live_tool_state *st) {
if (!st) return;
stop_list_clear(&st->call_ids);
free(st->visible_text);
@@ -7216,9 +7342,9 @@ static void responses_live_clear_locked(responses_live_state *st) {
st->live_tokens = 0;
}
static void responses_live_free(responses_live_state *st) {
static void live_tool_state_free(live_tool_state *st) {
if (!st) return;
responses_live_clear_locked(st);
live_tool_state_clear_locked(st);
free(st->call_ids.v);
memset(st, 0, sizeof(*st));
}
@@ -7260,7 +7386,7 @@ static void responses_live_remember(server *s, const char *visible_text,
const tool_calls *calls) {
if (!s || !visible_text || !visible_text[0]) return;
pthread_mutex_lock(&s->tool_mu);
responses_live_clear_locked(&s->responses_live);
live_tool_state_clear_locked(&s->responses_live);
s->responses_live.visible_text = xstrdup(visible_text);
s->responses_live.visible_len = strlen(visible_text);
if (calls) {
@@ -7273,10 +7399,29 @@ static void responses_live_remember(server *s, const char *visible_text,
pthread_mutex_unlock(&s->tool_mu);
}
static void anthropic_live_remember(server *s, const tool_calls *calls) {
if (!s || !calls || calls->len == 0) return;
pthread_mutex_lock(&s->tool_mu);
live_tool_state_clear_locked(&s->anthropic_live);
for (int i = 0; i < calls->len; i++) {
id_list_push_unique(&s->anthropic_live.call_ids, calls->v[i].id);
}
s->anthropic_live.live_tokens = ds4_session_pos(s->session);
s->anthropic_live.valid = s->anthropic_live.call_ids.len > 0;
pthread_mutex_unlock(&s->tool_mu);
}
static void responses_live_clear(server *s) {
if (!s) return;
pthread_mutex_lock(&s->tool_mu);
responses_live_clear_locked(&s->responses_live);
live_tool_state_clear_locked(&s->responses_live);
pthread_mutex_unlock(&s->tool_mu);
}
static void anthropic_live_clear(server *s) {
if (!s) return;
pthread_mutex_lock(&s->tool_mu);
live_tool_state_clear_locked(&s->anthropic_live);
pthread_mutex_unlock(&s->tool_mu);
}
@@ -7289,6 +7434,15 @@ static bool responses_live_has_call_id(server *s, const char *id) {
return found;
}
static bool anthropic_live_has_call_id(server *s, const char *id) {
if (!s || !id || !id[0]) return false;
pthread_mutex_lock(&s->tool_mu);
bool found = s->anthropic_live.valid &&
id_list_contains(&s->anthropic_live.call_ids, id);
pthread_mutex_unlock(&s->tool_mu);
return found;
}
static bool responses_live_matches_request(server *s, const stop_list *ids,
int live_tokens) {
if (!s || !ids || ids->len == 0) return false;
@@ -7303,6 +7457,20 @@ static bool responses_live_matches_request(server *s, const stop_list *ids,
return ok;
}
static bool anthropic_live_matches_request(server *s, const stop_list *ids,
int live_tokens) {
if (!s || !ids || ids->len == 0) return false;
pthread_mutex_lock(&s->tool_mu);
bool ok = s->anthropic_live.valid &&
s->anthropic_live.live_tokens == live_tokens &&
s->anthropic_live.call_ids.len == ids->len;
for (int i = 0; ok && i < ids->len; i++) {
ok = id_list_contains(&s->anthropic_live.call_ids, ids->v[i]);
}
pthread_mutex_unlock(&s->tool_mu);
return ok;
}
static bool tool_memory_has_id(server *s, const char *id) {
if (!s || s->disable_exact_dsml_tool_replay || !id || !id[0]) return false;
pthread_mutex_lock(&s->tool_mu);
@@ -7436,7 +7604,7 @@ static void apply_openai_stream_tool_ids(tool_calls *calls,
}
/* =========================================================================
* Disk KV Cache.
* KV Cache.
* =========================================================================
*
* The server has one live Metal session. We persist reusable DS4 session
@@ -7689,6 +7857,10 @@ static void id_list_free(stop_list *ids) {
static void collect_tool_call_ids(const chat_msgs *msgs, stop_list *ids) {
if (!msgs || !ids) return;
for (int i = 0; i < msgs->len; i++) {
id_list_push_unique(ids, msgs->v[i].tool_call_id);
for (int j = 0; j < msgs->v[i].tool_call_ids_len; j++) {
id_list_push_unique(ids, msgs->v[i].tool_call_ids[j]);
}
const tool_calls *calls = &msgs->v[i].calls;
for (int j = 0; j < calls->len; j++) {
id_list_push_unique(ids, calls->v[j].id);
@@ -8150,6 +8322,12 @@ static bool byte_prefix_match(const char *text, size_t text_len,
(prefix_len == 0 || memcmp(text, prefix, prefix_len) == 0);
}
static const char *kv_cache_key_kind(uint8_t ext_flags) {
if (ext_flags & KV_EXT_RESPONSES_VISIBLE) return "responses-visible";
if (ext_flags & KV_EXT_THINKING_VISIBLE) return "thinking-visible";
return "token-text";
}
static void tokens_copy_prefix(ds4_tokens *dst, const ds4_tokens *src, int n) {
dst->len = 0;
if (n > src->len) n = src->len;
@@ -8610,9 +8788,7 @@ static int kv_cache_try_load_text(server *s, const char *prompt_text,
if (loaded_path_out) *loaded_path_out = xstrdup(path);
if (loaded_ext_flags_out) *loaded_ext_flags_out = hdr.ext_flags;
kc->continued_last_store_tokens = loaded;
const char *key_kind = "token-text";
if (hdr.ext_flags & KV_EXT_RESPONSES_VISIBLE) key_kind = "responses-visible";
else if (hdr.ext_flags & KV_EXT_THINKING_VISIBLE) key_kind = "thinking-visible";
const char *key_kind = kv_cache_key_kind(hdr.ext_flags);
if (kc->opt.cold_max_tokens > 0 && loaded > kc->opt.cold_max_tokens) {
unlink(path);
server_log(DS4_LOG_KVCACHE,
@@ -8692,6 +8868,32 @@ static int responses_live_continuation_prompt(server *s, const request *req,
return live_tokens->len;
}
/* Tool-result Anthropic continuation.
*
* /v1/messages has no server-side response object like the OpenAI Responses
* API, but its tool_use_id is still a precise continuation handle inside a live
* local agent loop. When the IDs and live token frontier match, continue from
* the sampled DSML state and append only the user tool_result suffix. */
static int anthropic_live_continuation_prompt(server *s, const request *req,
int live_pos,
ds4_tokens *effective_prompt,
int *matched_ids) {
if (!s || !req || !effective_prompt) return 0;
if (req->api != API_ANTHROPIC || !req->anthropic_live_suffix_text) return 0;
if (req->anthropic_live_call_ids.len == 0) return 0;
if (!anthropic_live_matches_request(s, &req->anthropic_live_call_ids,
live_pos)) return 0;
const ds4_tokens *live_tokens = ds4_session_tokens(s->session);
if (!live_tokens || live_tokens->len != live_pos) return 0;
build_prompt_from_exact_prefix_and_text_suffix(
s->engine, live_tokens, req->anthropic_live_suffix_text,
effective_prompt);
if (matched_ids) *matched_ids = req->anthropic_live_call_ids.len;
return live_tokens->len;
}
/* Visible-replay Responses continuation.
*
* Other clients send the full visible transcript on every turn even though the
@@ -9193,15 +9395,21 @@ static bool should_remember_thinking_checkpoint(const request *r,
static void log_tool_calls_summary(const char *ctx, const tool_calls *calls) {
if (!calls || calls->len == 0) return;
buf names = {0};
buf ids = {0};
for (int i = 0; i < calls->len; i++) {
if (i) buf_putc(&names, ',');
if (i) buf_putc(&ids, ',');
buf_puts(&names, calls->v[i].name ? calls->v[i].name : "?");
buf_puts(&ids, calls->v[i].id ? calls->v[i].id : "?");
}
server_log(DS4_LOG_TOOL,
"ds4-server: tool calls ctx=%s n=%d names=[%s]",
"ds4-server: tool calls ctx=%s n=%d raw_dsml=%d ids=[%s] names=[%s]",
ctx,
calls->len,
calls->raw_dsml && calls->raw_dsml[0] ? 1 : 0,
ids.ptr ? ids.ptr : "",
names.ptr ? names.ptr : "");
buf_free(&ids);
buf_free(&names);
}
@@ -9212,7 +9420,9 @@ static void server_progress_cb(void *ud, const char *event, int current, int tot
double now = now_sec();
double elapsed = now - p->t0;
if (p->seen && current == p->last_current) {
if (p->srv && current > p->cached_tokens) kv_cache_maybe_store_continued(p->srv);
if (p->srv && current > p->cached_tokens) {
kv_cache_maybe_store_continued(p->srv);
}
return;
}
int display_start = p->cached_tokens;
@@ -9250,7 +9460,9 @@ static void server_progress_cb(void *ud, const char *event, int current, int tot
chunk_tps,
avg_tps,
elapsed);
if (p->srv && current > p->cached_tokens) kv_cache_maybe_store_continued(p->srv);
if (p->srv && current > p->cached_tokens) {
kv_cache_maybe_store_continued(p->srv);
}
}
static char *build_tool_checkpoint_suffix(const request *r, const char *content,
@@ -9496,13 +9708,14 @@ static bool should_canonicalize_tool_checkpoint(const server *s, const tool_call
*
* Clients resend full prompts as text. The worker first tries the old exact
* token-prefix hit, then a rendered-text prefix hit for the live checkpoint,
* then the disk text-prefix index, then a cold prefill. On text-prefix hits we
* build a fresh effective prompt from the checkpoint's exact token history plus
* a newly tokenized string suffix; the canonical full-prompt tokens are not
* sliced because BPE may merge across the byte boundary. Cold prompt caching is
* handled before generation: if the stable checkpoint is shorter than the full
* prompt, we prefill to that boundary, store it, and immediately continue to the
* real prompt. The live graph therefore always moves forward. */
* then disk text-prefix restart snapshots, then a cold prefill. On text-prefix
* hits we build a fresh effective prompt from the checkpoint's exact token
* history plus a newly tokenized string suffix; the canonical full-prompt
* tokens are not sliced because BPE may merge across the byte boundary. Cold
* prompt caching is handled before generation: if the stable checkpoint is
* shorter than the full prompt, we prefill to that boundary, store it, and
* immediately continue to the real prompt. The live graph therefore always
* moves forward. */
static void generate_job(server *s, job *j) {
char err[160];
err[0] = '\0';
@@ -9514,9 +9727,11 @@ static void generate_job(server *s, job *j) {
ds4_tokens effective_prompt = {0};
const ds4_tokens *prompt_for_sync = &j->req.prompt;
bool responses_live_continuation = false;
bool anthropic_live_continuation = false;
bool thinking_live_continuation = false;
const char *responses_live_match = NULL;
int responses_live_match_ids = 0;
int anthropic_live_match_ids = 0;
/* Responses gets the first chance to continue from live state. This is
* the whole point of the API shape: a request that is bound to prior live
* output by visible transcript or tool call ids does not need to prove an
@@ -9544,8 +9759,18 @@ static void generate_job(server *s, job *j) {
if (cached > 0) {
responses_live_continuation = true;
prompt_for_sync = &effective_prompt;
} else if (j->req.api == API_RESPONSES &&
j->req.responses_requires_live_tool_state)
} else {
cached = anthropic_live_continuation_prompt(s, &j->req, old_pos,
&effective_prompt,
&anthropic_live_match_ids);
if (cached > 0) {
anthropic_live_continuation = true;
cache_source = "anthropic-tool-output";
prompt_for_sync = &effective_prompt;
}
}
if (cached == 0 && j->req.api == API_RESPONSES &&
j->req.responses_requires_live_tool_state)
{
/* The parser saw a valid live call_id, but by worker execution time the
* live frontier no longer matches. Since the request did not replay
@@ -9555,7 +9780,14 @@ static void generate_job(server *s, job *j) {
http_error(j->fd, 400,
"Responses tool output requires live call state; replay full input instead");
return;
} else {
} else if (cached == 0 && j->req.api == API_ANTHROPIC &&
j->req.anthropic_requires_live_tool_state)
{
ds4_tokens_free(&effective_prompt);
http_error(j->fd, 400,
"Anthropic tool_result requires live tool state; replay full messages instead");
return;
} else if (cached == 0) {
cached = common == old_pos && j->req.prompt.len >= old_pos ? common : 0;
cache_source = cached > 0 ? "memory-token" : "none";
}
@@ -9650,6 +9882,12 @@ static void generate_job(server *s, job *j) {
responses_live_match_ids,
cached,
prompt_tokens);
} else if (anthropic_live_continuation) {
server_log(DS4_LOG_PREFILL,
"ds4-server: anthropic live continuation match=tool-output-ids ids=%d cached=%d prompt=%d",
anthropic_live_match_ids,
cached,
prompt_tokens);
} else if (thinking_live_continuation) {
server_log(DS4_LOG_PREFILL,
"ds4-server: thinking live continuation match=visible-prefix cached=%d prompt=%d",
@@ -9701,11 +9939,13 @@ static void generate_job(server *s, job *j) {
http_error(j->fd, 500, err);
return;
}
/* Once a non-live request wins, the old Responses live binding is stale.
* Keep it only when the current request explicitly continued from it. */
/* Once a non-live request wins, old protocol live bindings are stale. Keep
* a binding only when this request explicitly continued from it. */
if (!responses_live_continuation) responses_live_clear(s);
if (!anthropic_live_continuation) anthropic_live_clear(s);
if (!thinking_live_continuation) thinking_live_clear(s);
ds4_session_set_progress(s->session, NULL, NULL);
kv_cache_maybe_store_continued(s);
server_log(DS4_LOG_PREFILL,
"ds4-server: %s ctx=%s%s%s prompt done %.3fs",
j->req.kind == REQ_CHAT ? "chat" : "completion",
@@ -10079,6 +10319,15 @@ static void generate_job(server *s, job *j) {
responses_live_clear(s);
}
}
if (j->req.api == API_ANTHROPIC) {
if (parsed_calls.len && strcmp(final_finish, "error") &&
strcmp(final_finish, "length"))
{
anthropic_live_remember(s, &parsed_calls);
} else {
anthropic_live_clear(s);
}
}
if (j->req.kind == REQ_CHAT && parsed_calls.len &&
j->req.api != API_RESPONSES &&
@@ -10620,7 +10869,8 @@ static void server_close_resources(server *s) {
}
kv_cache_close(&s->kv);
tool_memory_free(&s->tool_mem);
responses_live_free(&s->responses_live);
live_tool_state_free(&s->responses_live);
live_tool_state_free(&s->anthropic_live);
visible_live_free(&s->thinking_live);
pthread_mutex_destroy(&s->tool_mu);
pthread_mutex_destroy(&s->trace_mu);
@@ -12304,6 +12554,169 @@ static void test_tool_memory_replays_sampled_dsml(void) {
pthread_mutex_destroy(&s.tool_mu);
}
static void test_anthropic_tool_memory_replays_sampled_dsml(void) {
const char *sampled_dsml =
"\n\n" DS4_TOOL_CALLS_START "\n"
"<DSMLinvoke name=\"Bash\">\n"
"<DSMLparameter name=\"command\" string=\"true\">ls -la</DSMLparameter>\n"
"<DSMLparameter name=\"description\" string=\"true\">list files</DSMLparameter>\n"
"</DSMLinvoke>\n"
DS4_TOOL_CALLS_END;
server s;
memset(&s, 0, sizeof(s));
pthread_mutex_init(&s.tool_mu, NULL);
tool_memory_put(&s, "toolu_exact", sampled_dsml);
const char *json =
"["
"{\"role\":\"assistant\",\"content\":["
"{\"type\":\"tool_use\",\"id\":\"toolu_exact\",\"name\":\"Bash\","
"\"input\":{\"description\":\"list files\",\"command\":\"ls -la\"}}"
"]},"
"{\"role\":\"user\",\"content\":["
"{\"type\":\"tool_result\",\"tool_use_id\":\"toolu_exact\",\"content\":\"ok\"}"
"]}"
"]";
const char *p = json;
chat_msgs msgs = {0};
TEST_ASSERT(parse_anthropic_messages(&p, &msgs));
TEST_ASSERT(msgs.len == 2);
TEST_ASSERT(msgs.v[1].tool_call_id && !strcmp(msgs.v[1].tool_call_id, "toolu_exact"));
stop_list ids = {0};
collect_tool_call_ids(&msgs, &ids);
TEST_ASSERT(id_list_contains(&ids, "toolu_exact"));
id_list_free(&ids);
tool_replay_stats stats = {0};
tool_memory_attach_to_messages(&s, &msgs, &stats);
TEST_ASSERT(msgs.v[0].calls.raw_dsml != NULL);
TEST_ASSERT(stats.mem == 1);
TEST_ASSERT(stats.canonical == 0);
char *prompt = render_chat_prompt_text(&msgs, NULL, NULL, DS4_THINK_HIGH);
const char *command = strstr(prompt, "name=\"command\"");
const char *description = strstr(prompt, "name=\"description\"");
TEST_ASSERT(command != NULL);
TEST_ASSERT(description != NULL);
TEST_ASSERT(command < description);
free(prompt);
chat_msgs_free(&msgs);
tool_memory_free(&s.tool_mem);
pthread_mutex_destroy(&s.tool_mu);
}
static void test_anthropic_live_tail_renders_tool_results_only(void) {
request r;
request_init(&r, REQ_CHAT, 128);
r.api = API_ANTHROPIC;
r.think_mode = DS4_THINK_HIGH;
chat_msgs msgs = {0};
chat_msg assistant = {0};
assistant.role = xstrdup("assistant");
tool_call tc = {0};
tc.id = xstrdup("toolu_live");
tc.name = xstrdup("Bash");
tc.arguments = xstrdup("{\"command\":\"pwd\"}");
tool_calls_push(&assistant.calls, tc);
chat_msgs_push(&msgs, assistant);
chat_msg user = {0};
user.role = xstrdup("user");
user.content = xstrdup("<tool_result>/tmp</tool_result>");
chat_msg_add_tool_call_id(&user, "toolu_live");
chat_msgs_push(&msgs, user);
/* Anthropic system text is parsed separately and appended to chat_msgs for
* rendering. The live-tail finder must ignore it when locating the final
* tool_result run. */
chat_msg system = {0};
system.role = xstrdup("system");
system.content = xstrdup("You are terse.");
chat_msgs_push(&msgs, system);
anthropic_prepare_live_continuation(&r, &msgs);
TEST_ASSERT(r.anthropic_live_call_ids.len == 1);
TEST_ASSERT(!strcmp(r.anthropic_live_call_ids.v[0], "toolu_live"));
TEST_ASSERT(r.anthropic_live_suffix_text != NULL);
TEST_ASSERT(!strncmp(r.anthropic_live_suffix_text,
"<end▁of▁sentence><User><tool_result>",
strlen("<end▁of▁sentence><User><tool_result>")));
TEST_ASSERT(strstr(r.anthropic_live_suffix_text, "/tmp</tool_result>") != NULL);
TEST_ASSERT(strstr(r.anthropic_live_suffix_text, "<Assistant><think>") != NULL);
TEST_ASSERT(strstr(r.anthropic_live_suffix_text, "Bash") == NULL);
chat_msgs_free(&msgs);
request_free(&r);
}
static void test_anthropic_tool_result_id_validation(void) {
server s = {0};
pthread_mutex_init(&s.tool_mu, NULL);
chat_msgs msgs = {0};
chat_msg user = {0};
user.role = xstrdup("user");
user.content = xstrdup("<tool_result>out</tool_result>");
chat_msg_add_tool_call_id(&user, "toolu_missing");
chat_msgs_push(&msgs, user);
char err[160] = {0};
TEST_ASSERT(!anthropic_validate_tool_results(&s, &msgs, NULL,
err, sizeof(err)));
TEST_ASSERT(strstr(err, "unknown Anthropic tool_result") != NULL);
pthread_mutex_lock(&s.tool_mu);
s.anthropic_live.valid = true;
s.anthropic_live.live_tokens = 10;
id_list_push_unique(&s.anthropic_live.call_ids, "toolu_missing");
pthread_mutex_unlock(&s.tool_mu);
bool needs_live_tool_state = false;
err[0] = '\0';
TEST_ASSERT(anthropic_validate_tool_results(&s, &msgs,
&needs_live_tool_state,
err, sizeof(err)));
TEST_ASSERT(needs_live_tool_state);
chat_msgs_free(&msgs);
live_tool_state_free(&s.anthropic_live);
pthread_mutex_destroy(&s.tool_mu);
}
static void test_anthropic_full_replay_allows_unknown_live_id(void) {
server s = {0};
pthread_mutex_init(&s.tool_mu, NULL);
chat_msgs msgs = {0};
chat_msg assistant = {0};
assistant.role = xstrdup("assistant");
tool_call tc = {0};
tc.id = xstrdup("toolu_replay");
tc.name = xstrdup("Bash");
tc.arguments = xstrdup("{\"command\":\"pwd\"}");
tool_calls_push(&assistant.calls, tc);
chat_msgs_push(&msgs, assistant);
chat_msg user = {0};
user.role = xstrdup("user");
user.content = xstrdup("<tool_result>/tmp</tool_result>");
chat_msg_add_tool_call_id(&user, "toolu_replay");
chat_msgs_push(&msgs, user);
bool needs_live_tool_state = false;
char err[160] = {0};
TEST_ASSERT(anthropic_validate_tool_results(&s, &msgs,
&needs_live_tool_state,
err, sizeof(err)));
TEST_ASSERT(!needs_live_tool_state);
chat_msgs_free(&msgs);
pthread_mutex_destroy(&s.tool_mu);
}
static void test_tool_checkpoint_canonicalization_gate_exact_replay(void) {
server s;
memset(&s, 0, sizeof(s));
@@ -12399,7 +12812,7 @@ static void test_responses_tool_output_id_validation(void) {
TEST_ASSERT(needs_live_tool_state);
chat_msgs_free(&msgs);
responses_live_free(&s.responses_live);
live_tool_state_free(&s.responses_live);
pthread_mutex_destroy(&s.tool_mu);
}
@@ -12473,7 +12886,7 @@ static void test_responses_stateless_tool_replay_requires_reasoning(void) {
TEST_ASSERT(!needs_live_reasoning);
chat_msgs_free(&msgs);
responses_live_free(&s.responses_live);
live_tool_state_free(&s.responses_live);
pthread_mutex_destroy(&s.tool_mu);
}
@@ -13513,6 +13926,10 @@ static void ds4_server_unit_tests_run(void) {
test_tool_checkpoint_suffix_is_future_prompt_canonical();
test_tool_checkpoint_minifies_json_parameters();
test_tool_memory_replays_sampled_dsml();
test_anthropic_tool_memory_replays_sampled_dsml();
test_anthropic_live_tail_renders_tool_results_only();
test_anthropic_tool_result_id_validation();
test_anthropic_full_replay_allows_unknown_live_id();
test_tool_checkpoint_canonicalization_gate_exact_replay();
test_responses_live_tail_renders_tool_outputs_only();
test_responses_tool_output_id_validation();
+129
View File
@@ -0,0 +1,129 @@
# Anthropic Live Tool Continuation
## Lessons To Reuse
From `misc/SIMPLIFIED_NO_CANONICALIZATION_MODEL.md`:
- The sampled KV state is the most faithful state DS4 has.
- Client-visible protocol objects should select that sampled state; they should
not force DS4 to rewrite it unless there is no live or persisted match.
- Tool calls are a natural instance of this rule:
```text
tool_use_id -> exact sampled DSML/KV frontier
```
From `misc/RESPONSE_API.md`:
- Protocol IDs exist to avoid proving that an already-known prefix tokenizes the
same way again.
- Tool-result requests bound to live call IDs should append only the tool-result
suffix to the live KV.
- Prefix matching and disk checkpoints remain fallbacks for restart, stateless
replay, edits, or client/session switches.
- Unknown IDs are not silently repairable if the request does not replay enough
prior history.
## Anthropic Contract
Anthropic tool loops use `tool_use.id` on assistant output and
`tool_result.tool_use_id` on the following user turn. For a local server with one
live KV session, that ID is enough to prove that the next request is continuing
the just-sampled assistant tool-call frontier, provided the live session is still
at the remembered token position.
The fast path is:
1. The model samples an Anthropic tool call.
2. DS4 assigns/remembers the `toolu_...` IDs and keeps the live KV exactly as
sampled, including hidden thinking and raw DSML.
3. The next `/v1/messages` request contains `tool_result` blocks with those IDs.
4. DS4 verifies that the live token position and the ID set match.
5. DS4 tokenizes only:
```text
<end▁of▁sentence><User><tool_result>...<Assistant><think-or-/think>
```
and appends it to the live sampled prefix.
This avoids checkpoint canonicalization after the tool call and avoids relying
on a replayed JSON tool block to match the sampled DSML byte-for-byte.
## Fallbacks
If the live Anthropic ID binding is gone, DS4 should use normal replay:
- exact token prefix,
- live rendered-text prefix with suffix retokenization,
- disk string-prefix KV checkpoints,
- cold prefill.
This is valid when the request replays the assistant `tool_use` block before the
`tool_result`, because exact DSML tool memory can restore the sampled tool-call
bytes from the tool ID. If the request contains only a `tool_result` for an
unknown live ID and does not replay the prior assistant call, DS4 should return a
clear 400 asking the client to replay the full history.
## What Not To Do
- Do not add a RAM snapshot cache to paper over live-continuation mismatches.
The specific Anthropic issue is not lack of a rewind checkpoint; it is that
the protocol already identifies the live frontier by tool ID.
- Do not canonicalize a live tool-call checkpoint just to match the client JSON.
If DS4 has raw sampled DSML and a matching live ID, the live state should win.
- Do not reuse already-tokenized suffix tokens from the full rendered request at
a string boundary. Tokenize the suffix after the cache/continuation decision.
## QA Checklist
- Unit test: Anthropic `tool_result.tool_use_id` is parsed and collected.
- Unit test: Anthropic live suffix renders only EOS, tool result, and assistant
prefix, not the previous assistant tool call.
- Integration test: a matching live Anthropic ID builds an effective prompt
whose cached prefix length is the live token count.
- Unit test: a tool-result-only request with an unknown ID is rejected.
- Integration test: Claude Code can execute a multi-turn tool session without
post-tool checkpoint canonicalization or needless re-prefill.
- Regression check: Responses live continuation still works.
## Current Implementation Notes
- The speculative RAM KV restart cache was removed. It was solving the wrong
problem for Anthropic tool loops: the protocol already gives DS4 a precise
live continuation key in `tool_use_id`.
- `chat_msg` now keeps all tool-result IDs seen in a message, not just the first
or last one. Anthropic can return multiple `tool_result` blocks in one user
message, and a partial ID set would make live continuation unsafe.
- Anthropic live state stores only the sampled token frontier and generated
tool-call IDs. Responses still also stores visible replay text because that
protocol has a visible-prefix continuation mode.
- Anthropic `tool_result` requests are accepted in two cases:
- the IDs match the current live Anthropic tool-call frontier;
- the request replays the prior assistant `tool_use` blocks, so normal exact
DSML replay and prefix matching can reconstruct the prompt.
- A `tool_result` request with neither live state nor replayed prior tool calls
returns HTTP 400 and asks for full messages.
## QA Log
- `make ds4_test`: pass.
- `./ds4_test --server`: pass.
- `make ds4-server`: pass.
- Claude Code through `~/bin/claude-ds4` against `/v1/messages`: pass.
- Prompt forced two Bash tool calls in one assistant turn.
- Server logged:
```text
anthropic live continuation match=tool-output-ids ids=2 cached=24002 prompt=24041
chat ctx=24002..24041:39 TOOLS prompt done 0.873s
```
- The trace contained no post-tool `tool checkpoint canonicalization` event.
- Negative HTTP check: a `tool_result` with unknown `tool_use_id` and no replayed
prior assistant call returned:
```text
HTTP/1.1 400 Bad Request
unknown Anthropic tool_result tool_use_id toolu_unknown; replay full messages
```