use crate::params::chat::{ ChatCompletionChunkResponse, ChatCompletionParameters, ChatCompletionResponse, }; use anyhow::{Result, anyhow}; use candle_core::{DType, Device, Tensor}; use candle_nn::VarBuilder; use rocket::async_stream::stream; use rocket::futures::Stream; use crate::models::qwen3::config::{Qwen3Config, Qwen3GenerationConfig}; use crate::models::qwen3::model::Qwen3Model; // use crate::models::GenerateStream; use crate::utils::{ build_completion_chunk_response, build_completion_response, find_type_files, get_device, get_dtype, get_logit_processor, }; use crate::{chat_template::ChatTemplate, models::GenerateModel, tokenizer::TokenizerModel}; pub struct Qwen3GenerateModel<'a> { chat_template: ChatTemplate<'a>, tokenizer: TokenizerModel, qwen3: Qwen3Model, device: Device, eos_token_id1: u32, eos_token_id2: u32, generation_config: Qwen3GenerationConfig, model_name: String, } impl<'a> Qwen3GenerateModel<'a> { pub fn init(path: &str, device: Option<&Device>, dtype: Option) -> Result { let chat_template = ChatTemplate::init(path)?; let tokenizer = TokenizerModel::init(path)?; let config_path = path.to_string() + "/config.json"; let cfg: Qwen3Config = serde_json::from_slice(&std::fs::read(config_path)?)?; let device = &get_device(device); let cfg_dtype = cfg.torch_dtype.as_str(); let dtype = get_dtype(dtype, cfg_dtype); let model_list = find_type_files(path, "safetensors")?; let vb = unsafe { VarBuilder::from_mmaped_safetensors(&model_list, dtype, device)? }; let qwen3 = Qwen3Model::new(&cfg, vb)?; let generation_config_path = path.to_string() + "/generation_config.json"; let generation_config: Qwen3GenerationConfig = serde_json::from_slice(&std::fs::read(generation_config_path)?)?; let model_name = std::path::Path::new(path) .file_name() .and_then(|s| s.to_str()) .unwrap_or("qwen3") .to_string(); Ok(Qwen3GenerateModel { chat_template, tokenizer, qwen3, device: device.clone(), eos_token_id1: generation_config.eos_token_id[0] as u32, eos_token_id2: generation_config.eos_token_id[1] as u32, generation_config, model_name, }) } } impl<'a> GenerateModel for Qwen3GenerateModel<'a> { fn generate(&mut self, mes: ChatCompletionParameters) -> Result { let temperature = mes .temperature .unwrap_or(self.generation_config.temperature); let top_p = mes.top_p.unwrap_or(self.generation_config.top_p); let top_k = self.generation_config.top_k; let seed = mes.seed.unwrap_or(34562) as u64; let mut logit_processor = get_logit_processor(Some(temperature), Some(top_p), Some(top_k), seed); let mes_render = self.chat_template.apply_chat_template(&mes)?; // let enable_thinking = extract_metadata_value::(&mes.metadata, "enable_thinking"); // let mes_render = self // .chat_template // .apply_chat_temp_think(&mes, enable_thinking)?; 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 sample_len = mes.max_tokens.unwrap_or(2048); for _ in 0..sample_len { let logits = self.qwen3.forward(Some(&input_ids), None, 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.eos_token_id1 || next_token == self.eos_token_id2 { 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.qwen3.clear_kv_cache(); let response = build_completion_response(res, &self.model_name, Some(num_token), Some(prompt_tokens)); Ok(response) } fn generate_stream( &mut self, mes: ChatCompletionParameters, ) -> Result< Box< dyn Stream> + Send + Unpin + '_, >, > { let temperature = mes .temperature .unwrap_or(self.generation_config.temperature); let top_p = mes.top_p.unwrap_or(self.generation_config.top_p); let top_k = self.generation_config.top_k; let seed = mes.seed.unwrap_or(34562) as u64; let mut logit_processor = get_logit_processor(Some(temperature), Some(top_p), Some(top_k), seed); let mes_render = self.chat_template.apply_chat_template(&mes)?; // let enable_thinking = extract_metadata_value::(&mes.metadata, "enable_thinking"); // let mes_render = self // .chat_template // .apply_chat_temp_think(&mes, enable_thinking)?; 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 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.qwen3.forward( Some(&input_ids), None, 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.eos_token_id1 || next_token == self.eos_token_id2 { break; } seqlen_offset += seq_len; seq_len = 1; input_ids = Tensor::from_vec(vec![next_token], (1, 1), &self.device)?; } self.qwen3.clear_kv_cache(); }; Ok(Box::new(Box::pin(stream))) } }