refactor generate code

This commit is contained in:
jhqxxx
2026-04-02 22:22:52 +08:00
parent b254d21efc
commit bd8eee6520
59 changed files with 1745 additions and 1800 deletions
+46 -80
View File
@@ -1,20 +1,19 @@
use crate::models::common::generate::get_logit_processor;
use crate::models::common::MultiModalData;
use crate::models::common::generate::{
GenerationContext, generate_generic, generate_stream_generic,
};
use crate::params::chat::{
ChatCompletionChunkResponse, ChatCompletionParameters, ChatCompletionResponse,
};
use anyhow::{Result, anyhow};
use candle_core::{DType, Device, Tensor};
use anyhow::Result;
use candle_core::{DType, Device};
use candle_nn::VarBuilder;
use rocket::async_stream::stream;
use rocket::futures::Stream;
use crate::models::minicpm4::config::MiniCPM4Config;
use crate::models::minicpm4::model::MiniCPMModel;
// use crate::models::GenerateStream;
use crate::utils::{
build_completion_chunk_response, build_completion_response, find_type_files, get_device,
get_dtype,
};
use crate::utils::{find_type_files, get_device, get_dtype};
use crate::{chat_template::ChatTemplate, models::GenerateModel, tokenizer::TokenizerModel};
pub struct MiniCPMGenerateModel<'a> {
@@ -22,8 +21,6 @@ pub struct MiniCPMGenerateModel<'a> {
tokenizer: TokenizerModel,
minicpm: MiniCPMModel,
device: Device,
endoftext_id: u32,
im_end_id: u32,
model_name: String,
}
@@ -40,7 +37,8 @@ impl<'a> MiniCPMGenerateModel<'a> {
let im_end_id = cfg.eos_token_id[1];
let model_list = find_type_files(path, "safetensors")?;
let vb = unsafe { VarBuilder::from_mmaped_safetensors(&model_list, dtype, device)? };
let minicpm = MiniCPMModel::new(vb, cfg)?;
let eos_ids = vec![endoftext_id, im_end_id];
let minicpm = MiniCPMModel::new(vb, cfg, eos_ids)?;
let model_name = std::path::Path::new(path)
.file_name()
.and_then(|s| s.to_str())
@@ -51,8 +49,6 @@ impl<'a> MiniCPMGenerateModel<'a> {
tokenizer,
minicpm,
device: device.clone(),
endoftext_id,
im_end_id,
model_name,
})
}
@@ -60,33 +56,29 @@ impl<'a> MiniCPMGenerateModel<'a> {
impl<'a> GenerateModel for MiniCPMGenerateModel<'a> {
fn generate(&mut self, mes: ChatCompletionParameters) -> Result<ChatCompletionResponse> {
let seed = mes.seed.unwrap_or(34562) as u64;
let mut logit_processor = get_logit_processor(mes.temperature, mes.top_p, None, seed);
let mes_render = self.chat_template.apply_chat_template(&mes)?;
let mut input_ids = self.tokenizer.text_encode(mes_render, &self.device)?;
let mut seq_len = input_ids.dim(1)?;
let prompt_tokens = seq_len as u32;
let mut seqlen_offset = 0;
let mut generate = Vec::new();
let input_ids = self.tokenizer.text_encode(mes_render, &self.device)?;
let seed = mes.seed.unwrap_or(34562) as u64;
let sample_len = mes.max_tokens.unwrap_or(2048);
for _ in 0..sample_len {
let logits = self.minicpm.forward_with_cache(&input_ids, seqlen_offset)?;
let logits = logits.squeeze(0)?.squeeze(0)?.to_dtype(DType::F32)?;
let next_token = logit_processor.sample(&logits)?;
generate.push(next_token);
if next_token == self.endoftext_id || next_token == self.im_end_id {
break;
}
seqlen_offset += seq_len;
seq_len = 1;
input_ids = Tensor::from_vec(vec![next_token], (1, 1), &self.device)?;
}
let num_token = generate.len() as u32;
let res = self.tokenizer.token_decode(generate)?;
self.minicpm.clear_kv_cache();
let response =
build_completion_response(res, &self.model_name, Some(num_token), Some(prompt_tokens));
Ok(response)
let mut ctx = GenerationContext::new(
mes.temperature,
mes.top_p,
None,
seed,
input_ids.dim(1)?,
sample_len,
self.device.clone(),
);
let data = MultiModalData::new(vec![]);
generate_generic(
&mut self.minicpm,
&self.tokenizer,
input_ids,
data,
&mut ctx,
&self.model_name,
)
}
fn generate_stream(
&mut self,
@@ -100,50 +92,24 @@ impl<'a> GenerateModel for MiniCPMGenerateModel<'a> {
>,
> {
let seed = mes.seed.unwrap_or(34562) as u64;
let mut logit_processor = get_logit_processor(mes.temperature, mes.top_p, None, seed);
let mes_render = self.chat_template.apply_chat_template(&mes)?;
let mut input_ids = self.tokenizer.text_encode(mes_render, &self.device)?;
let mut seq_len = input_ids.dim(1)?;
let mut seqlen_offset = 0;
let input_ids = self.tokenizer.text_encode(mes_render, &self.device)?;
let data = MultiModalData::new(vec![]);
let sample_len = mes.max_tokens.unwrap_or(512);
let stream = stream! {
let mut error_tokens = Vec::new();
for _ in 0..sample_len {
let logits = self.minicpm.forward_with_cache(
&input_ids,
seqlen_offset,
)?;
let logits = logits.squeeze(0)?.squeeze(0)?.to_dtype(DType::F32)?;
let next_token = logit_processor.sample(&logits)?;
let mut decode_ids = Vec::new();
if !error_tokens.is_empty(){
decode_ids.extend_from_slice(&error_tokens);
}
decode_ids.push(next_token);
let decoded_token = self.tokenizer.token_decode(decode_ids).map_err(|e| anyhow!(format!("stream decode error{e}")))?;
if decoded_token.contains("�") {
error_tokens.push(next_token);
if error_tokens.len() > 3 {
error_tokens.clear();
}
seqlen_offset += seq_len;
seq_len = 1;
input_ids = Tensor::from_vec(vec![next_token], (1, 1), &self.device)?;
continue;
}
error_tokens.clear();
let chunk = build_completion_chunk_response(decoded_token, &self.model_name, None, None);
yield Ok(chunk);
if next_token == self.endoftext_id || next_token == self.im_end_id {
break;
}
seqlen_offset += seq_len;
seq_len = 1;
input_ids = Tensor::from_vec(vec![next_token], (1, 1), &self.device)?;
}
self.minicpm.clear_kv_cache();
};
let stream = generate_stream_generic(
&mut self.minicpm,
&self.tokenizer,
input_ids,
data,
mes.temperature,
mes.top_p,
None,
seed,
sample_len,
false,
&self.device,
&self.model_name,
)?;
Ok(Box::new(Box::pin(stream)))
}
}
+29 -6
View File
@@ -4,7 +4,10 @@ use candle_nn::{Embedding, Linear, Module, RmsNorm, VarBuilder, embedding, rms_n
use crate::{
models::{
common::modules::{GateUpDownMLP, NaiveAttention},
common::{
InferenceModel,
modules::{GateUpDownMLP, NaiveAttention},
},
minicpm4::config::MiniCPM4Config,
},
position_embed::rope::compute_default_rope_parameters,
@@ -207,10 +210,11 @@ pub struct MiniCPMModel {
norm: RmsNorm,
rope_emb: MiniCPMLongRoPE,
lm_head: Linear,
stop_token_ids: Vec<u32>,
}
impl MiniCPMModel {
pub fn new(vb: VarBuilder, cfg: MiniCPM4Config) -> Result<Self> {
pub fn new(vb: VarBuilder, cfg: MiniCPM4Config, eos_ids: Vec<u32>) -> Result<Self> {
let vb = vb.pp("model");
let embed_tokens = embedding(cfg.vocab_size, cfg.hidden_size, vb.pp("embed_tokens"))?;
let mut layers = Vec::with_capacity(cfg.num_hidden_layers);
@@ -229,10 +233,11 @@ impl MiniCPMModel {
norm,
rope_emb,
lm_head,
stop_token_ids: eos_ids,
})
}
pub fn forward(&mut self, input_ids: &Tensor, position_id: usize) -> Result<Tensor> {
pub fn forward(&mut self, input_ids: &Tensor, seqlen_offset: usize) -> Result<Tensor> {
let (bs, seq_len) = input_ids.dims2()?;
let input_embeds = self
.embed_tokens
@@ -251,7 +256,7 @@ impl MiniCPMModel {
}
};
let (cos, sin) = self.rope_emb.forward(position_id, seq_len)?;
let (cos, sin) = self.rope_emb.forward(seqlen_offset, seq_len)?;
let mut hidden_states = input_embeds;
for decode_layer in &self.layers {
hidden_states =
@@ -267,7 +272,11 @@ impl MiniCPMModel {
Ok(logits)
}
pub fn forward_with_cache(&mut self, input_ids: &Tensor, position_id: usize) -> Result<Tensor> {
pub fn forward_with_cache(
&mut self,
input_ids: &Tensor,
seqlen_offset: usize,
) -> Result<Tensor> {
let (bs, seq_len) = input_ids.dims2()?;
let input_embeds = self
.embed_tokens
@@ -285,7 +294,7 @@ impl MiniCPMModel {
)?)
}
};
let (cos, sin) = self.rope_emb.forward(position_id, seq_len)?;
let (cos, sin) = self.rope_emb.forward(seqlen_offset, seq_len)?;
let mut hidden_states = input_embeds;
for decode_layer in &mut self.layers {
hidden_states = decode_layer.forward_with_cache(
@@ -311,3 +320,17 @@ impl MiniCPMModel {
}
}
}
impl InferenceModel for MiniCPMModel {
fn forward_step(&mut self, input_ids: &Tensor, seqlen_offset: usize) -> Result<Tensor> {
self.forward_with_cache(input_ids, seqlen_offset)
}
fn clear_cache(&mut self) {
self.clear_kv_cache();
}
fn stop_token_ids(&self) -> Vec<u32> {
self.stop_token_ids.clone()
}
}