From 6085db16cf3636d71d8f267913e838e0232a67bb Mon Sep 17 00:00:00 2001 From: Concedo <39025047+LostRuins@users.noreply.github.com> Date: Sun, 30 Aug 2026 18:04:49 +0800 Subject: [PATCH] fix streaming race condition --- expose.cpp | 8 +++----- expose.h | 3 ++- gpttype_adapter.cpp | 42 +++++++++++++++++++++++++++++++++++------- model_adapter.h | 2 ++ 4 files changed, 42 insertions(+), 13 deletions(-) diff --git a/expose.cpp b/expose.cpp index 822e9c12f..7824b8d2c 100644 --- a/expose.cpp +++ b/expose.cpp @@ -263,17 +263,15 @@ extern "C" } const char * new_token(int idx) { - if (generated_tokens.size() <= idx || idx < 0) return nullptr; - - return generated_tokens[idx].c_str(); + return gpttype_new_token(idx); } int get_stream_count() { - return generated_tokens.size(); + return gpttype_get_stream_count(); } bool has_finished() { - return generation_finished; + return generation_finished.load(); } bool batch_generate_enabled() { return gpttype_batch_generate_enabled(); diff --git a/expose.h b/expose.h index 311bf5111..2f062b22b 100644 --- a/expose.h +++ b/expose.h @@ -1,4 +1,5 @@ #pragma once +#include #include const int tensor_split_max = 16; @@ -396,7 +397,7 @@ extern std::string lora_filename; extern std::string mmproj_filename; extern std::string draftmodel_filename; extern std::vector generated_tokens; -extern bool generation_finished; +extern std::atomic generation_finished; extern bool audio_multimodal_supported; extern bool vision_multimodal_supported; extern float last_eval_time; diff --git a/gpttype_adapter.cpp b/gpttype_adapter.cpp index b39735b07..6ee552316 100644 --- a/gpttype_adapter.cpp +++ b/gpttype_adapter.cpp @@ -77,7 +77,7 @@ std::string lora_filename = ""; std::string mmproj_filename = ""; std::string draftmodel_filename = ""; int speculative_chunk_amt = 4; //do it in chunks of this many tokens -bool generation_finished; +std::atomic generation_finished; bool audio_multimodal_supported = false; bool vision_multimodal_supported = false; float last_process_time = 0; @@ -163,11 +163,12 @@ static std::vector dry_repeat_count; // Indexed as last_n_tokens static std::unordered_map dry_max_token_repeat; static std::vector top_picks_history; static int remaining_tokens = 0; -static bool early_abort = false; +static std::atomic early_abort = false; static std::mutex concat_output_mtx; static std::string concat_output = ""; static std::string concat_output_reader_copy_poll = ""; //for streaming static std::string concat_output_reader_copy_res = ""; //for gen response +static std::string generated_token_reader_copy = ""; //stable copy for streaming token readers static std::vector logit_biases; static bool add_bos_token = true; // if set to false, mmproj handling breaks. dont disable unless you know what you're doing static bool load_guidance = false; //whether to enable cfg for negative prompts @@ -4302,6 +4303,7 @@ struct BatchGenerateRequest bool i_batch_is_prefill = false; llama_sampler * sampler = nullptr; std::vector generated_pieces; + std::string stream_reader_copy; std::string output; int prompt_token_count = 0; int completion_token_count = 0; @@ -5006,7 +5008,8 @@ const char * gpttype_batch_generate_new_token(int request_id, int idx) { return nullptr; } - return req->generated_pieces[idx].c_str(); + req->stream_reader_copy = req->generated_pieces[idx]; + return req->stream_reader_copy.c_str(); } const char * gpttype_batch_generate_pending_output(int request_id) @@ -5256,6 +5259,23 @@ const std::string & gpttype_get_pending_output() return concat_output_reader_copy_poll; } +int gpttype_get_stream_count() +{ + std::lock_guard lock(concat_output_mtx); + return static_cast(generated_tokens.size()); +} + +const char * gpttype_new_token(int idx) +{ + std::lock_guard lock(concat_output_mtx); + if (idx < 0 || idx >= (int) generated_tokens.size()) + { + return nullptr; + } + generated_token_reader_copy = generated_tokens[idx]; + return generated_token_reader_copy.c_str(); +} + const std::vector gpttype_get_top_picks_data() { return top_picks_history; @@ -5599,7 +5619,12 @@ generation_outputs gpttype_generate(const generation_inputs inputs) showed_rnn_warning = false; generation_finished = false; // Set current generation status - generated_tokens.clear(); // New Generation, new tokens + { + std::lock_guard lock(concat_output_mtx); + generated_tokens.clear(); // New Generation, new tokens + generated_tokens.reserve(16); + generated_token_reader_copy = ""; + } delayed_generated_tokens.clear(); concat_output_mtx.lock(); @@ -6442,7 +6467,10 @@ generation_outputs gpttype_generate(const generation_inputs inputs) std::mt19937 rng(kcpp_data->seed); //do some reservation so we don't have to realloc - generated_tokens.reserve(remaining_tokens+16); + { + std::lock_guard lock(concat_output_mtx); + generated_tokens.reserve(remaining_tokens+16); + } //prepare sampler order std::vector sampler_order; @@ -7005,8 +7033,8 @@ generation_outputs gpttype_generate(const generation_inputs inputs) delayed_generated_tokens.push_back(tokenizedstr); while(delayed_generated_tokens.size() > delayed_generated_tokens_limit && delayed_generated_tokens.size() > 0) { - generated_tokens.push_back(delayed_generated_tokens[0]); concat_output_mtx.lock(); + generated_tokens.push_back(delayed_generated_tokens[0]); concat_output += delayed_generated_tokens[0]; concat_output_mtx.unlock(); delayed_generated_tokens.pop_front(); @@ -7380,8 +7408,8 @@ generation_outputs gpttype_generate(const generation_inputs inputs) //flush any remaining delayed tokens while(delayed_generated_tokens.size() > 0) { - generated_tokens.push_back(delayed_generated_tokens[0]); concat_output_mtx.lock(); + generated_tokens.push_back(delayed_generated_tokens[0]); concat_output += delayed_generated_tokens[0]; concat_output_mtx.unlock(); delayed_generated_tokens.pop_front(); diff --git a/model_adapter.h b/model_adapter.h index e2711b0c3..3d0cdad69 100644 --- a/model_adapter.h +++ b/model_adapter.h @@ -101,6 +101,8 @@ generation_outputs gpttype_batch_generate_result(int request_id); bool gpttype_batch_generate_abort(int request_id); void gpttype_batch_generate_release(int request_id); +int gpttype_get_stream_count(); +const char * gpttype_new_token(int idx); const std::string & gpttype_get_pending_output(); std::vector gpttype_get_token_arr(const std::string & input, bool addbos); std::string gpttype_detokenize(const std::vector & input, bool render_special);