use crate::error::{RealizarError, Result};
use crate::gguf::forward_qwen35::{
qwen35_route, qwen35_route_notice, Qwen35Model, Qwen35Route, Qwen35State,
QWEN35_GPU_FALLBACK_PREFIX,
};
use crate::gguf::{MappedGGUFModel, OwnedQuantizedModel};
const MIN_CAPACITY: usize = 4096;
pub const QWEN35_SESSION_PREFILL_ENV: &str = "APR_QWEN35_SESSION_PREFILL";
pub type Qwen35Turn = crate::session::Turn;
enum Step {
Gpu(String),
Fatal(RealizarError),
}
impl From<RealizarError> for Step {
fn from(e: RealizarError) -> Self {
Self::Fatal(e)
}
}
#[cfg(feature = "cuda")]
struct GpuBackend {
model: crate::gguf::cuda::Qwen35CudaModel<'static>,
state: Option<crate::gguf::cuda::Qwen35CudaState>,
device_name: String,
hash: crate::gguf::forward_qwen35::Qwen35ModelHash,
validated: bool,
}
#[cfg(feature = "cuda")]
impl GpuBackend {
fn build(
qwen: &'static Qwen35Model<'static>,
mapped: &MappedGGUFModel,
notices: &mut Vec<String>,
plan_positions: Option<usize>,
) -> std::result::Result<Self, GpuBuild> {
use crate::gguf::forward_qwen35::{Qwen35ModelHash, QWEN35_F2_PROBE_MAX};
let executor = crate::cuda::CudaExecutor::new(0)
.map_err(|e| format!("CUDA initialization failed: {e}"))?;
let planned = match plan_positions {
Some(positions) => Some(plan_capacity(qwen, &executor, positions)?),
None => None,
};
let device_name = executor
.device_name()
.unwrap_or_else(|_| "Unknown GPU".to_string());
let vram_mb = executor.memory_info().unwrap_or((0, 0)).1 / (1024 * 1024);
let model = crate::gguf::cuda::Qwen35CudaModel::with_max_seq_len(
qwen,
executor,
QWEN35_F2_PROBE_MAX + 2,
)
.map_err(|e| GpuBuild::Fallback(format!("the CUDA model would not build: {e}")))?;
let mut model = model;
if let Some((attention, rows)) = planned {
model.set_prefill_chunk_rows(rows);
model.set_prefill_attention(attention);
}
say(
notices,
format!(
"Backend: GPU (CUDA, {device_name}, {vram_mb} MB VRAM) [qwen35 hybrid forward, #3090]"
),
);
Ok(Self {
model,
state: None,
device_name,
hash: Qwen35ModelHash::of(mapped.data()),
validated: false,
})
}
}
#[cfg(feature = "cuda")]
enum GpuBuild {
Fallback(String),
Refused(Box<crate::capacity::CapacityRefusal>),
}
#[cfg(feature = "cuda")]
impl From<String> for GpuBuild {
fn from(reason: String) -> Self {
Self::Fallback(reason)
}
}
#[cfg(feature = "cuda")]
fn plan_capacity(
qwen: &Qwen35Model<'_>,
executor: &crate::cuda::CudaExecutor,
positions: usize,
) -> std::result::Result<(crate::gguf::cuda::PrefillAttention, usize), GpuBuild> {
let device_memory = crate::capacity::measure_device_memory(executor)?;
let (gpu_free, gpu_total) = device_memory.plan_free_total();
let attention_paths =
crate::gguf::cuda::Qwen35CudaModel::prefill_attention_candidates_for(qwen, executor);
let chunk_rows_to_try: &[usize] = match device_memory {
crate::capacity::DeviceMemory::Unified { .. } => &[
crate::gguf::cuda::UNIFIED_PREFILL_CHUNK_ROWS,
crate::gguf::cuda::PREFILL_MAX_CHUNK_ROWS,
],
crate::capacity::DeviceMemory::Discrete { .. } => {
&[crate::gguf::cuda::PREFILL_MAX_CHUNK_ROWS]
},
};
let mut passed_over = Vec::new();
let planned =
crate::capacity::plan_first_fit(&attention_paths, chunk_rows_to_try, |attention, rows| {
let verdict = crate::capacity::plan(&crate::capacity::CapacityInputs {
memory: Some(device_memory),
..crate::gguf::cuda::Qwen35CudaModel::capacity_inputs(
qwen, positions, gpu_free, gpu_total, attention, rows,
)
});
if let crate::capacity::CapacityVerdict::Refused(r) = &verdict {
passed_over.push(crate::capacity::passed_over_line(
attention.as_str(),
rows,
r,
));
}
verdict
});
match planned {
Ok(fit) => {
if !passed_over.is_empty() {
eprintln!(
"[qwen35] prefill plan: {} did not fit; using {} at {} rows ({:.0} MiB)",
passed_over.join("; "),
fit.0.as_str(),
fit.1,
fit.2.total_mb
);
}
if fit.2.kv_dtype != crate::capacity::KvDtype::F32 {
return Err(GpuBuild::Fallback(format!(
"the capacity plan chose a {:?} KV cache, which this build cannot allocate",
fit.2.kv_dtype
)));
}
Ok((fit.0, fit.1))
},
Err(Some(refusal)) => Err(GpuBuild::Refused(refusal)),
Err(None) => Err(GpuBuild::Fallback(
"no prefill attention path to plan".to_string(),
)),
}
}
enum Checkpoint {
#[cfg(feature = "cuda")]
Gpu(crate::gguf::cuda::Qwen35CudaCheckpoint),
Cpu(crate::gguf::forward_qwen35::Qwen35Checkpoint),
}
enum Backend {
#[cfg(feature = "cuda")]
Gpu(Box<GpuBackend>),
Cpu(Option<Qwen35State>),
}
pub type Qwen35Session = crate::session::Session<Qwen35Forward>;
impl crate::session::Session<Qwen35Forward> {
pub fn load(mapped: &MappedGGUFModel, no_gpu: bool) -> Result<Self> {
Ok(Self::new(Qwen35Forward::load(mapped, no_gpu)?))
}
pub fn load_for_run(
qwen: &'static Qwen35Model<'static>,
mapped: &MappedGGUFModel,
no_gpu: bool,
positions: usize,
) -> Result<Self> {
Ok(Self::new(Qwen35Forward::from_host(
qwen,
mapped,
no_gpu,
Some(positions),
)?))
}
#[must_use]
pub fn num_layers(&self) -> usize {
self.engine().qwen.layers.len()
}
}
pub struct Qwen35Forward {
qwen: &'static Qwen35Model<'static>,
backend: Backend,
capacity: usize,
turn_positions: usize,
allocations: u64,
min_capacity: usize,
context_length: usize,
notices: Vec<String>,
per_token_prefill: bool,
batched_prefills: usize,
checkpoint: Option<Checkpoint>,
im_start: Option<u32>,
}
impl Qwen35Forward {
pub fn load(mapped: &MappedGGUFModel, no_gpu: bool) -> Result<Self> {
Self::from_host(Self::leak_host(mapped)?, mapped, no_gpu, None)
}
pub fn leak_host(mapped: &MappedGGUFModel) -> Result<&'static Qwen35Model<'static>> {
let base = Qwen35Model::create_base_model(&mapped.model, mapped.data())?;
let base: &'static OwnedQuantizedModel = Box::leak(Box::new(base));
Ok(Box::leak(Box::new(Qwen35Model::from_model_and_layers(
base,
&mapped.model,
mapped.data(),
)?)))
}
pub fn cached_host(
path: &std::path::Path,
mapped: &MappedGGUFModel,
) -> Result<&'static Qwen35Model<'static>> {
type Key = (std::path::PathBuf, u64, Option<std::time::SystemTime>);
static HOSTS: std::sync::Mutex<Vec<(Key, &'static Qwen35Model<'static>)>> =
std::sync::Mutex::new(Vec::new());
let meta = std::fs::metadata(path).ok();
let key: Key = (
std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf()),
meta.as_ref().map_or(0, std::fs::Metadata::len),
meta.and_then(|m| m.modified().ok()),
);
let mut hosts = HOSTS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some((_, host)) = hosts.iter().find(|(k, _)| *k == key) {
return Ok(host);
}
let host = Self::leak_host(mapped)?;
hosts.push((key, host));
Ok(host)
}
fn from_host(
qwen: &'static Qwen35Model<'static>,
mapped: &MappedGGUFModel,
no_gpu: bool,
plan_positions: Option<usize>,
) -> Result<Self> {
let context_length = qwen.base.config.context_length.max(1);
let mut notices = Vec::new();
let per_token_prefill =
std::env::var(QWEN35_SESSION_PREFILL_ENV).as_deref() == Ok("per-token");
if per_token_prefill {
say(
&mut notices,
format!("[qwen35] {QWEN35_SESSION_PREFILL_ENV}=per-token: prompts prefill one token at a time, not batched"),
);
}
let im_start = mapped
.model
.vocabulary()
.and_then(|v| v.iter().position(|t| t == "<|im_start|>"))
.and_then(|i| u32::try_from(i).ok());
let route = qwen35_route(no_gpu, cfg!(feature = "cuda"));
if let Some(notice) = qwen35_route_notice(route) {
say(&mut notices, notice.to_string());
}
let backend = match route {
#[cfg(feature = "cuda")]
Qwen35Route::Gpu => match GpuBackend::build(qwen, mapped, &mut notices, plan_positions)
{
Ok(gpu) => Backend::Gpu(Box::new(gpu)),
Err(GpuBuild::Refused(refusal)) => {
return Err(RealizarError::CapacityRefused(refusal));
},
Err(GpuBuild::Fallback(reason)) => {
say(&mut notices, fallback_line(&reason));
Backend::Cpu(None)
},
},
_ => Backend::Cpu(None),
};
let min_capacity = if plan_positions.is_some() {
0
} else {
MIN_CAPACITY
};
Ok(Self {
qwen,
backend,
capacity: 0,
turn_positions: 0,
allocations: 0,
min_capacity,
context_length,
notices,
per_token_prefill,
batched_prefills: 0,
checkpoint: None,
im_start,
})
}
#[must_use]
pub fn base(&self) -> &'static OwnedQuantizedModel {
self.qwen.base
}
fn try_forward(&mut self, tokens: &[u32], start: usize) -> std::result::Result<Vec<f32>, Step> {
#[cfg(feature = "cuda")]
self.validate_gpu_once(tokens)?;
if start == 0 {
self.reset_state()?;
}
#[cfg(feature = "cuda")]
if tokens.len().saturating_sub(start) > 1 {
if let Some(logits) = self.try_batched_prefill(&tokens[start..], start)? {
return Ok(logits);
}
}
let mut logits = Vec::new();
for (pos, &token) in tokens.iter().enumerate().skip(start) {
logits = self.forward_one(token, pos)?;
}
Ok(logits)
}
#[cfg(feature = "cuda")]
fn try_batched_prefill(
&mut self,
new: &[u32],
pos0: usize,
) -> std::result::Result<Option<Vec<f32>>, Step> {
let qwen = self.qwen;
let Backend::Gpu(gpu) = &mut self.backend else {
return Ok(None);
};
if self.per_token_prefill {
return Ok(None);
}
let end = pos0 + new.len();
if let Err(why) = fit_prefill_plan(qwen, &mut gpu.model, end) {
eprintln!(
"[qwen35] batched prefill: {} tokens at position {pos0} would not fit ({why}); \
prefilling one token at a time on the GPU",
new.len()
);
return Ok(None);
}
let state = gpu
.state
.as_mut()
.ok_or_else(|| Step::Gpu("the device state was never allocated".to_string()))?;
let t0 = std::time::Instant::now();
let logits = batched_prefill_outcome(gpu.model.prefill(new, state, pos0), new.len(), pos0)?;
let ms = t0.elapsed().as_secs_f64() * 1000.0;
eprintln!(
"[qwen35] batched prefill: {} tokens in {ms:.0} ms ({:.0} tok/s, chunk {} rows, attention {}, from position {pos0})",
new.len(),
new.len() as f64 * 1000.0 / ms.max(1e-9),
gpu.model.prefill_chunk_rows(end),
gpu.model.prefill_attention_mode().as_str(),
);
self.batched_prefills += 1;
Ok(Some(logits))
}
#[cfg(feature = "cuda")]
fn validate_gpu_once(&mut self, probe: &[u32]) -> std::result::Result<(), Step> {
let qwen = self.qwen;
let Backend::Gpu(gpu) = &mut self.backend else {
return Ok(());
};
if gpu.validated {
return Ok(());
}
let outcome = crate::gguf::forward_qwen35::f2_validate_qwen35_receipted_hashed(
&mut gpu.model,
qwen,
probe,
&gpu.hash,
&gpu.device_name,
);
if !outcome.accepted {
return Err(Step::Gpu(
"the F2 CPU-parity guard rejected the GPU path".to_string(),
));
}
gpu.validated = true;
Ok(())
}
fn forward_one(&mut self, token: u32, pos: usize) -> std::result::Result<Vec<f32>, Step> {
let qwen = self.qwen;
match &mut self.backend {
#[cfg(feature = "cuda")]
Backend::Gpu(gpu) => {
let state = gpu
.state
.as_mut()
.ok_or_else(|| Step::Gpu("the device state was never allocated".to_string()))?;
gpu.model.forward_single(token, state, pos).map_err(|e| {
Step::Gpu(format!("the GPU forward failed at position {pos}: {e}"))
})
},
Backend::Cpu(state) => {
let state = state.as_mut().ok_or_else(|| RealizarError::InvalidShape {
reason: "qwen35 session: the CPU state was never allocated".to_string(),
})?;
Ok(qwen.forward_single_qwen35(token, state, pos)?)
},
}
}
fn reset_state(&mut self) -> std::result::Result<(), Step> {
match &mut self.backend {
#[cfg(feature = "cuda")]
Backend::Gpu(gpu) => {
if let Some(state) = gpu.state.as_mut() {
gpu.model
.reset_state(state)
.map_err(|e| Step::Gpu(format!("the device state would not reset: {e}")))?;
}
},
Backend::Cpu(state) => {
if let Some(state) = state.as_mut() {
state.reset();
}
},
}
Ok(())
}
fn bind_cuda_context_or_fall_back(&mut self) -> Result<()> {
#[cfg(feature = "cuda")]
if let Backend::Gpu(gpu) = &self.backend {
if let Err(e) = gpu.model.make_current() {
self.fall_back_to_cpu(&format!(
"the CUDA context would not bind to this thread: {e}"
))?;
}
}
Ok(())
}
fn ensure_capacity_or_fall_back(&mut self) -> Result<()> {
match self.ensure_capacity() {
Ok(()) => Ok(()),
Err(Step::Gpu(reason)) => self.fall_back_to_cpu(&reason),
Err(Step::Fatal(e)) => Err(e),
}
}
fn ensure_capacity(&mut self) -> std::result::Result<(), Step> {
let positions = self.turn_positions;
let have_state = match &self.backend {
#[cfg(feature = "cuda")]
Backend::Gpu(gpu) => gpu.state.is_some(),
Backend::Cpu(state) => state.is_some(),
};
if have_state && positions <= self.capacity {
return Ok(());
}
let capacity = grown_capacity(
positions,
self.capacity,
self.min_capacity,
self.context_length,
);
self.allocations += 1;
self.checkpoint = None;
let qwen = self.qwen;
match &mut self.backend {
#[cfg(feature = "cuda")]
Backend::Gpu(gpu) => {
gpu.state = None;
gpu.state = Some(gpu.model.new_state_with_capacity(capacity).map_err(|e| {
Step::Gpu(format!(
"a decode state for {capacity} positions would not allocate: {e}"
))
})?);
},
Backend::Cpu(state) => *state = Some(qwen.new_state(capacity)),
}
self.capacity = capacity;
Ok(())
}
fn fall_back_to_cpu(&mut self, reason: &str) -> Result<()> {
say(&mut self.notices, fallback_line(reason));
self.backend = Backend::Cpu(None);
self.checkpoint = None;
self.capacity = 0;
match self.ensure_capacity() {
Ok(()) => Ok(()),
Err(Step::Fatal(e)) => Err(e),
Err(Step::Gpu(reason)) => Err(RealizarError::UnsupportedOperation {
operation: "qwen35_session".to_string(),
reason: format!("the CPU backend reported a GPU failure: {reason}"),
}),
}
}
}
impl Qwen35Forward {
fn after_failed_restore(&mut self, on_gpu: bool, e: RealizarError) -> Result<bool> {
if on_gpu {
self.fall_back_to_cpu(&format!("the device checkpoint would not restore: {e}"))?;
Ok(false)
} else {
Err(e)
}
}
}
impl crate::session::ArchForward for Qwen35Forward {
fn arch(&self) -> &'static str {
"qwen35"
}
fn on_gpu(&self) -> bool {
match self.backend {
#[cfg(feature = "cuda")]
Backend::Gpu(_) => true,
Backend::Cpu(_) => false,
}
}
fn context_length(&self) -> usize {
self.context_length
}
fn batched_prefills(&self) -> usize {
self.batched_prefills
}
fn notices(&self) -> &[String] {
&self.notices
}
fn reserve(&mut self, positions: usize) -> Result<bool> {
let before = self.allocations;
self.turn_positions = positions;
self.bind_cuda_context_or_fall_back()?;
self.ensure_capacity_or_fall_back()?;
Ok(self.allocations != before)
}
fn checkpoint_at(&self, prompt: &[u32]) -> Option<usize> {
self.im_start
.and_then(|id| prompt.iter().rposition(|&t| t == id))
.filter(|&k| k > 0)
.or_else(|| prompt.len().checked_sub(1))
}
fn save_checkpoint(&mut self) -> Result<()> {
let saved = match &mut self.backend {
#[cfg(feature = "cuda")]
Backend::Gpu(gpu) => {
let Some(state) = gpu.state.as_ref() else {
self.checkpoint = None;
return Ok(());
};
let mut slot = match self.checkpoint.take() {
Some(Checkpoint::Gpu(c)) => Some(c),
_ => None,
};
gpu.model.save_checkpoint(state, &mut slot)?;
slot.map(Checkpoint::Gpu)
},
Backend::Cpu(state) => state
.as_ref()
.map(|state| Checkpoint::Cpu(state.checkpoint())),
};
self.checkpoint = saved;
Ok(())
}
fn restore_checkpoint(&mut self) -> Result<bool> {
let restored = match (&mut self.backend, &self.checkpoint) {
#[cfg(feature = "cuda")]
(Backend::Gpu(gpu), Some(Checkpoint::Gpu(c))) => match gpu.state.as_mut() {
Some(state) => gpu.model.restore_checkpoint(c, state),
None => Ok(false),
},
(Backend::Cpu(Some(state)), Some(Checkpoint::Cpu(c))) => Ok(state.restore(c)),
_ => Ok(false),
};
match restored {
Ok(restored) => Ok(restored),
Err(e) => self.after_failed_restore(self.on_gpu(), e),
}
}
fn validate(&mut self, probe: &[u32]) -> Result<()> {
#[cfg(feature = "cuda")]
match self.validate_gpu_once(probe) {
Ok(()) => {},
Err(Step::Gpu(reason)) => self.fall_back_to_cpu(&reason)?,
Err(Step::Fatal(e)) => return Err(e),
}
#[cfg(not(feature = "cuda"))]
let _ = probe;
Ok(())
}
fn forward(&mut self, tokens: &[u32], start: usize) -> Result<Vec<f32>> {
let mut start = start;
loop {
match self.try_forward(tokens, start) {
Ok(logits) => return Ok(logits),
Err(Step::Gpu(reason)) => {
self.fall_back_to_cpu(&reason)?;
start = 0;
},
Err(Step::Fatal(e)) => return Err(e),
}
}
}
}
#[cfg(feature = "cuda")]
fn fit_prefill_plan(
qwen: &Qwen35Model<'_>,
model: &mut crate::gguf::cuda::Qwen35CudaModel<'static>,
end: usize,
) -> std::result::Result<(), String> {
let memory = crate::capacity::measure_device_memory(model.executor_mut())?;
let (free, _) = memory.plan_free_total();
let rows_to_try: &[usize] = match memory {
crate::capacity::DeviceMemory::Unified { .. } => &[
crate::gguf::cuda::UNIFIED_PREFILL_CHUNK_ROWS,
crate::gguf::cuda::PREFILL_MAX_CHUNK_ROWS,
],
crate::capacity::DeviceMemory::Discrete { .. } => {
&[crate::gguf::cuda::PREFILL_MAX_CHUNK_ROWS]
},
};
let attentions = crate::gguf::cuda::Qwen35CudaModel::prefill_attention_candidates_for(
qwen,
model.executor_mut(),
);
let (attention, rows) = choose_prefill_plan(
free,
&attentions,
rows_to_try,
|attention, rows| {
model.set_prefill_attention(attention);
model.set_prefill_chunk_rows(rows);
model.prefill_workspace_bytes(end) as u64 + crate::capacity::OVERHEAD_BYTES
},
|attention| attention.as_str(),
)?;
model.set_prefill_attention(attention);
model.set_prefill_chunk_rows(rows);
Ok(())
}
#[cfg_attr(not(feature = "cuda"), allow(dead_code))]
fn choose_prefill_plan<A: Copy>(
free: u64,
attentions: &[A],
rows_to_try: &[usize],
mut need: impl FnMut(A, usize) -> u64,
name: impl Fn(A) -> &'static str,
) -> std::result::Result<(A, usize), String> {
let mut refused = Vec::new();
for &attention in attentions {
for &rows in rows_to_try {
let need = need(attention, rows);
if need <= free {
return Ok((attention, rows));
}
refused.push(format!(
"{} at {rows} rows needs {} MiB",
name(attention),
need >> 20
));
}
}
Err(format!("{}; {} MiB free", refused.join(", "), free >> 20))
}
#[cfg_attr(not(feature = "cuda"), allow(dead_code))]
fn batched_prefill_outcome<E: std::fmt::Display>(
result: std::result::Result<Vec<f32>, E>,
tokens: usize,
pos0: usize,
) -> std::result::Result<Vec<f32>, Step> {
result.map_err(|e| {
Step::Gpu(format!(
"the GPU batched prefill of {tokens} tokens at position {pos0} failed: {e}"
))
})
}
fn say(notices: &mut Vec<String>, line: String) {
eprintln!("{line}");
notices.push(line);
}
fn fallback_line(reason: &str) -> String {
format!("{QWEN35_GPU_FALLBACK_PREFIX}, falling back to CPU: {reason}")
}
fn grown_capacity(
positions: usize,
current: usize,
min_capacity: usize,
context_length: usize,
) -> usize {
if min_capacity == 0 {
return positions
.max(current.saturating_mul(2))
.min(context_length)
.max(positions);
}
let mut capacity = min_capacity.max(current.saturating_mul(2));
while capacity < positions {
capacity = capacity.saturating_mul(2);
}
capacity.min(context_length).max(positions)
}
#[cfg(test)]
#[path = "qwen35_session_tests.rs"]
mod tests;