From b39cec30c991c06fd8ee32f53c8d992642f80347 Mon Sep 17 00:00:00 2001 From: kerthcet Date: Sun, 20 Sep 2026 23:31:50 +0100 Subject: [PATCH] fix the module name Signed-off-by: kerthcet --- docs/configuration.md | 2 +- docs/fsm_architecture.md | 2 +- src/api/chat.rs | 2 +- src/api/routes.rs | 2 +- src/api/tests.rs | 8 ++--- src/backend/engine.rs | 30 ---------------- src/backend/mock.rs | 38 ++++++++++---------- src/backend/mod.rs | 33 +++++++++++++++-- src/cli/chat.rs | 2 +- src/cli/commands.rs | 10 +++--- src/cli/serve.rs | 13 ++++--- src/{backend/llm_engine.rs => engine/mod.rs} | 16 ++++----- src/lib.rs | 1 + src/main.rs | 1 + 14 files changed, 79 insertions(+), 81 deletions(-) delete mode 100644 src/backend/engine.rs rename src/{backend/llm_engine.rs => engine/mod.rs} (98%) diff --git a/docs/configuration.md b/docs/configuration.md index 0f0fba8..0569bf6 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -33,4 +33,4 @@ puma run qwen/qwen2.5-0.5b --tokens-per-block 32 - **`--default-max-tokens`** only applies when a request omits `max_tokens`; an explicit per-request `max_tokens` always takes precedence. - The flag defaults are sourced from `EngineConfig` in - `src/backend/llm_engine.rs`, so the CLI and the library stay in sync. + `src/engine/mod.rs`, so the CLI and the library stay in sync. diff --git a/docs/fsm_architecture.md b/docs/fsm_architecture.md index 2ddccf5..d115597 100644 --- a/docs/fsm_architecture.md +++ b/docs/fsm_architecture.md @@ -207,7 +207,7 @@ PUMA uses a two-level event system (inspired by TokenSpeed): - `src/sequence_manager/mod.rs` - Owns memory + FSM transition logic (`advance()`) - `src/scheduler/core.rs` - Scheduling policy that drives the SequenceManager - `src/scheduler/events.rs` - External scheduler events -- `src/backend/llm_engine.rs` - Event coordinator +- `src/engine/mod.rs` - Event coordinator ## Usage Example diff --git a/src/api/chat.rs b/src/api/chat.rs index e155733..a94c014 100644 --- a/src/api/chat.rs +++ b/src/api/chat.rs @@ -15,7 +15,7 @@ use crate::api::types::{ ChatChoice, ChatChoiceDelta, ChatCompletionChunk, ChatCompletionRequest, ChatCompletionResponse, ChatMessage, ChatMessageDelta, ErrorResponse, Usage, }; -use crate::backend::EngineHandle; +use crate::engine::EngineHandle; /// Main handler for chat completions pub async fn chat_completions( diff --git a/src/api/routes.rs b/src/api/routes.rs index 15875be..834739d 100644 --- a/src/api/routes.rs +++ b/src/api/routes.rs @@ -10,7 +10,7 @@ use tower_http::{ LatencyUnit, }; -use crate::backend::EngineHandle; +use crate::engine::EngineHandle; use crate::registry::model_registry::ModelRegistry; use super::{chat, completions, models}; diff --git a/src/api/tests.rs b/src/api/tests.rs index 41a18a1..b3c9b74 100644 --- a/src/api/tests.rs +++ b/src/api/tests.rs @@ -13,16 +13,16 @@ use tempfile::TempDir; use tower::util::ServiceExt; // for `oneshot` and `ready` use super::routes::create_router; -use crate::backend::mock::MockEngine; -use crate::backend::{engine, EngineConfig}; +use crate::backend::mock::MockBackend; +use crate::engine::{self, EngineConfig}; use crate::registry::model_registry::{CacheInfo, ModelInfo, ModelMetadata, ModelRegistry}; /// Helper to create test app with a pre-registered test model /// Returns the router and the temp directory (which must be kept alive) fn create_test_app() -> (axum::Router, TempDir) { // Build the engine and spawn its runner; the handle drives the router - let (handle, runner) = engine( - MockEngine::new(), + let (handle, runner) = engine::spawn( + MockBackend::new(), create_test_tokenizer(), "test-model".to_string(), EngineConfig::default(), diff --git a/src/backend/engine.rs b/src/backend/engine.rs deleted file mode 100644 index a3bb6f2..0000000 --- a/src/backend/engine.rs +++ /dev/null @@ -1,30 +0,0 @@ -use crate::block_manager::types::TokenId; -use std::io; -use std::pin::Pin; -use tokio_stream::Stream; - -/// Backend trait - low-level inference that works with token IDs -/// -/// LLMEngine handles tokenization (text → tokens) -/// Backend handles inference (tokens → tokens) -pub trait Backend: Send + Sync { - /// Generate tokens from input token_ids - /// Returns generated token IDs - fn generate( - &self, - token_ids: Vec, - max_tokens: usize, - temperature: f32, - ) -> impl std::future::Future, io::Error>> + Send; - - /// Generate tokens with streaming - /// Returns stream of token IDs as they're generated - fn generate_stream( - &self, - token_ids: Vec, - max_tokens: usize, - temperature: f32, - ) -> impl std::future::Future< - Output = Result + Send>>, io::Error>, - > + Send; -} diff --git a/src/backend/mock.rs b/src/backend/mock.rs index 3639aa7..035db68 100644 --- a/src/backend/mock.rs +++ b/src/backend/mock.rs @@ -3,7 +3,7 @@ use std::io; use std::pin::Pin; use tokio_stream::Stream; -use super::engine::Backend; +use super::Backend; use crate::block_manager::types::TokenId; /// Default vocab size the mock samples completion tokens from. @@ -12,7 +12,7 @@ use crate::block_manager::types::TokenId; /// (so they decode to real text), large enough for varied output. const DEFAULT_VOCAB_SIZE: u32 = 1000; -/// Mock inference engine that behaves like a real autoregressive model. +/// Mock inference backend that behaves like a real autoregressive model. /// /// Unlike a naive echo, this consumes the prompt as *context* and emits only /// **new** completion tokens — never the prompt back. Tokens are produced by a @@ -23,11 +23,11 @@ const DEFAULT_VOCAB_SIZE: u32 = 1000; /// /// Generation stops at `max_tokens` (the mock has no EOS concept yet). #[derive(Clone)] -pub struct MockEngine { +pub struct MockBackend { vocab_size: u32, } -impl MockEngine { +impl MockBackend { pub fn new() -> Self { Self { vocab_size: DEFAULT_VOCAB_SIZE, @@ -73,7 +73,7 @@ fn hash_tokens(tokens: &[TokenId]) -> u64 { hash } -impl Backend for MockEngine { +impl Backend for MockBackend { async fn generate( &self, token_ids: Vec, @@ -101,7 +101,7 @@ impl Backend for MockEngine { } } -impl Default for MockEngine { +impl Default for MockBackend { fn default() -> Self { Self::new() } @@ -113,9 +113,9 @@ mod tests { #[tokio::test] async fn returns_only_completion_tokens() { - let engine = MockEngine::new(); + let backend = MockBackend::new(); let prompt = vec![5, 9, 2]; - let out = engine.generate(prompt.clone(), 4, 0.0).await.unwrap(); + let out = backend.generate(prompt.clone(), 4, 0.0).await.unwrap(); // Completion-only: exactly max_tokens, and it does not start with the // prompt (a real model returns a continuation, not an echo). assert_eq!(out.len(), 4); @@ -124,24 +124,24 @@ mod tests { #[tokio::test] async fn is_deterministic() { - let engine = MockEngine::new(); - let a = engine.generate(vec![1, 2, 3], 8, 0.0).await.unwrap(); - let b = engine.generate(vec![1, 2, 3], 8, 0.0).await.unwrap(); + let backend = MockBackend::new(); + let a = backend.generate(vec![1, 2, 3], 8, 0.0).await.unwrap(); + let b = backend.generate(vec![1, 2, 3], 8, 0.0).await.unwrap(); assert_eq!(a, b, "same prompt must yield same completion"); } #[tokio::test] async fn is_input_dependent() { - let engine = MockEngine::new(); - let a = engine.generate(vec![1, 2, 3], 8, 0.0).await.unwrap(); - let b = engine.generate(vec![3, 2, 1], 8, 0.0).await.unwrap(); + let backend = MockBackend::new(); + let a = backend.generate(vec![1, 2, 3], 8, 0.0).await.unwrap(); + let b = backend.generate(vec![3, 2, 1], 8, 0.0).await.unwrap(); assert_ne!(a, b, "different prompts should yield different completions"); } #[tokio::test] async fn tokens_stay_within_vocab() { - let engine = MockEngine::with_vocab_size(50); - let out = engine.generate(vec![7, 7, 7], 32, 0.0).await.unwrap(); + let backend = MockBackend::with_vocab_size(50); + let out = backend.generate(vec![7, 7, 7], 32, 0.0).await.unwrap(); assert!( out.iter().all(|&t| t < 50), "ids must be within vocab range" @@ -150,11 +150,11 @@ mod tests { #[tokio::test] async fn stream_matches_generate() { - let engine = MockEngine::new(); + let backend = MockBackend::new(); let prompt = vec![10, 20, 30]; - let batched = engine.generate(prompt.clone(), 6, 0.0).await.unwrap(); + let batched = backend.generate(prompt.clone(), 6, 0.0).await.unwrap(); let mut streamed = Vec::new(); - let mut s = engine.generate_stream(prompt, 6, 0.0).await.unwrap(); + let mut s = backend.generate_stream(prompt, 6, 0.0).await.unwrap(); while let Some(tok) = s.next().await { streamed.push(tok); } diff --git a/src/backend/mod.rs b/src/backend/mod.rs index c2c7de0..32dd4cd 100644 --- a/src/backend/mod.rs +++ b/src/backend/mod.rs @@ -1,5 +1,32 @@ -pub mod engine; -pub mod llm_engine; pub mod mock; -pub use llm_engine::{engine, EngineConfig, EngineHandle}; +use crate::block_manager::types::TokenId; +use std::io; +use std::pin::Pin; +use tokio_stream::Stream; + +/// Backend trait - low-level inference that works with token IDs +/// +/// The engine handles tokenization (text → tokens) +/// Backend handles inference (tokens → tokens) +pub trait Backend: Send + Sync { + /// Generate tokens from input token_ids + /// Returns generated token IDs + fn generate( + &self, + token_ids: Vec, + max_tokens: usize, + temperature: f32, + ) -> impl std::future::Future, io::Error>> + Send; + + /// Generate tokens with streaming + /// Returns stream of token IDs as they're generated + fn generate_stream( + &self, + token_ids: Vec, + max_tokens: usize, + temperature: f32, + ) -> impl std::future::Future< + Output = Result + Send>>, io::Error>, + > + Send; +} diff --git a/src/cli/chat.rs b/src/cli/chat.rs index 087251e..f0e338b 100644 --- a/src/cli/chat.rs +++ b/src/cli/chat.rs @@ -5,7 +5,7 @@ use rustyline_derive::{Completer, Helper, Highlighter, Validator}; use std::io::{self, Write}; use tokio_stream::StreamExt; -use crate::backend::EngineHandle; +use crate::engine::EngineHandle; #[derive(Clone)] struct PlaceholderHint { diff --git a/src/cli/commands.rs b/src/cli/commands.rs index e55457c..95c30bb 100644 --- a/src/cli/commands.rs +++ b/src/cli/commands.rs @@ -4,10 +4,10 @@ use prettytable::{format, row, Table}; use tokenizers::Tokenizer; -use crate::backend::mock::MockEngine; -use crate::backend::{engine, EngineConfig}; +use crate::backend::mock::MockBackend; use crate::cli::{chat, inspect, ls, rm}; use crate::downloader::{self, Provider}; +use crate::engine::{self, EngineConfig}; use crate::registry::model_registry::ModelRegistry; use crate::system::system_info::SystemInfo; use crate::utils::format::{format_size_decimal, format_time_ago}; @@ -315,12 +315,12 @@ pub async fn run(cli: Cli) { }; // Load inference backend - // TODO: Replace MockEngine with real backend that loads model files + // TODO: Replace MockBackend with real backend that loads model files // Real backend will use: registry.get_model(&args.model)?.metadata.cache.path - let backend = MockEngine::new(); + let backend = MockBackend::new(); // Create engine: cheap send-side handle + runner that owns the scheduler - let (handle, runner) = engine( + let (handle, runner) = engine::spawn( backend, tokenizer, args.model.clone(), diff --git a/src/cli/serve.rs b/src/cli/serve.rs index a2538a5..729aaca 100644 --- a/src/cli/serve.rs +++ b/src/cli/serve.rs @@ -5,9 +5,8 @@ use tokenizers::Tokenizer; use tracing::{debug, info}; use crate::api::routes::create_router; -use crate::backend::engine; -use crate::backend::mock::MockEngine; -use crate::backend::EngineConfig; +use crate::backend::mock::MockBackend; +use crate::engine::{self, EngineConfig}; use crate::registry::model_registry::ModelRegistry; /// Execute the serve command @@ -34,15 +33,15 @@ pub async fn execute( ); info!("Starting PUMA to serve model: {}", model_name); - // Initialize backend (MockEngine for now, replace with MLX later) - let backend = MockEngine::new(); - debug!("Using MockEngine backend"); + // Initialize backend (MockBackend for now, replace with MLX later) + let backend = MockBackend::new(); + debug!("Using MockBackend backend"); // TODO: Load the model's real tokenizer; placeholder BPE for now let tokenizer = Tokenizer::new(BPE::default()); // Create engine: cheap send-side handle + runner that owns the scheduler - let (handle, runner) = engine(backend, tokenizer, model_name.to_string(), config); + let (handle, runner) = engine::spawn(backend, tokenizer, model_name.to_string(), config); // Spawn the runner's event loop; the handle submits work via events tokio::spawn(runner.serve()); diff --git a/src/backend/llm_engine.rs b/src/engine/mod.rs similarity index 98% rename from src/backend/llm_engine.rs rename to src/engine/mod.rs index b5eae91..9f48f11 100644 --- a/src/backend/llm_engine.rs +++ b/src/engine/mod.rs @@ -5,7 +5,7 @@ use tokenizers::Tokenizer; use tokio::sync::mpsc; use tokio_stream::StreamExt; -use super::engine::Backend; +use crate::backend::Backend; use crate::block_manager::allocator::CpuAllocator; use crate::block_manager::manager::BlockManager; use crate::block_manager::types::{SequenceId, TokenId}; @@ -345,11 +345,11 @@ fn stream_flush(decoded: &str, sent_len: usize) -> Option<&str> { /// Tunable engine parameters. /// -/// These were previously hardcoded inside [`engine`]. Construct via +/// These were previously hardcoded inside [`spawn`]. Construct via /// [`EngineConfig::default`] and override fields as needed: /// /// ``` -/// # use puma::backend::llm_engine::EngineConfig; +/// # use puma::engine::EngineConfig; /// let cfg = EngineConfig { max_batch_size: 64, ..Default::default() }; /// ``` #[derive(Debug, Clone)] @@ -397,7 +397,7 @@ impl Default for EngineConfig { /// to tune memory/batching. Spawn `runner.serve()` on a task and share the /// returned handle with the API / CLI. All request submission goes through /// events, so the handle never touches the scheduler directly. -pub fn engine( +pub fn spawn( backend: B, tokenizer: Tokenizer, model: String, @@ -440,7 +440,7 @@ pub fn engine( #[cfg(test)] mod tests { use super::*; - use crate::backend::mock::MockEngine; + use crate::backend::mock::MockBackend; fn create_test_tokenizer() -> Tokenizer { use tokenizers::models::bpe::BPE; @@ -451,10 +451,10 @@ mod tests { } #[tokio::test] - async fn test_llm_engine() { - let backend = MockEngine::new(); + async fn test_engine_spawn() { + let backend = MockBackend::new(); let tokenizer = create_test_tokenizer(); - let (handle, runner) = engine( + let (handle, runner) = spawn( backend, tokenizer, "test-model".to_string(), diff --git a/src/lib.rs b/src/lib.rs index 909efec..b71d309 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -5,6 +5,7 @@ pub mod backend; pub mod block_manager; pub mod cli; pub mod downloader; +pub mod engine; pub mod fsm; pub mod registry; pub mod scheduler; diff --git a/src/main.rs b/src/main.rs index a15b72d..7c38914 100644 --- a/src/main.rs +++ b/src/main.rs @@ -5,6 +5,7 @@ mod backend; mod block_manager; mod cli; mod downloader; +mod engine; mod fsm; mod registry; mod scheduler;