diff --git a/Makefile b/Makefile index 82873c015..2b3050cb6 100644 --- a/Makefile +++ b/Makefile @@ -59,6 +59,9 @@ target/release/paddler_gui: $(PADDLER_SOURCES) esbuild-meta.json target/vulkan/release/paddler: $(PADDLER_SOURCES) esbuild-meta.json cargo build --release -p paddler_cli --features vulkan,web_admin_panel --target-dir target/vulkan +target/vulkan/release/paddler_gui: $(PADDLER_SOURCES) esbuild-meta.json + cargo build --release -p paddler_gui --features vulkan,web_admin_panel --target-dir target/vulkan + # ----------------------------------------------------------------------------- # Phony targets # ----------------------------------------------------------------------------- diff --git a/paddler_agent/src/continuous_batch_arbiter.rs b/paddler_agent/src/continuous_batch_arbiter.rs index a63c897b0..5af2045d0 100644 --- a/paddler_agent/src/continuous_batch_arbiter.rs +++ b/paddler_agent/src/continuous_batch_arbiter.rs @@ -52,6 +52,7 @@ pub struct ContinuousBatchArbiter { pub agent_name: Option, pub chat_template_override: Option, pub desired_slots_total: i32, + pub gpu_devices: Vec, pub inference_parameters: InferenceParameters, pub multimodal_projection_path: Option, pub model_metadata_holder: Arc, @@ -66,6 +67,7 @@ impl ContinuousBatchArbiter { applicable_state: AgentApplicableState, agent_name: Option, desired_slots_total: i32, + gpu_devices: Vec, model_metadata_holder: Arc, slot_aggregated_status_manager: Arc, ) -> ContinuousBatchArbiterBuildOutcome { @@ -79,6 +81,7 @@ impl ContinuousBatchArbiter { agent_name, chat_template_override: applicable_state.chat_template_override, desired_slots_total, + gpu_devices, inference_parameters: applicable_state.inference_parameters, multimodal_projection_path: applicable_state.multimodal_projection_path, model_metadata_holder, @@ -108,6 +111,7 @@ impl ContinuousBatchArbiter { let agent_name_clone = self.agent_name.clone(); let desired_slots_total = self.desired_slots_total; + let gpu_devices = self.gpu_devices.clone(); let inference_parameters = self.inference_parameters.clone(); let model_metadata_holder = self.model_metadata_holder.clone(); let multimodal_projection_path = self.multimodal_projection_path.clone(); @@ -148,14 +152,18 @@ impl ContinuousBatchArbiter { .to_llama_kv_cache_dtype(), ); + let mut model_params = + LlamaModelParams::default().with_n_gpu_layers(inference_parameters.n_gpu_layers); + + if !gpu_devices.is_empty() { + model_params = model_params + .with_devices(&gpu_devices) + .context("Invalid --gpu-devices index")?; + } + let model = Arc::new( - LlamaModel::load_from_file( - &llama_backend, - model_path.clone(), - &LlamaModelParams::default() - .with_n_gpu_layers(inference_parameters.n_gpu_layers), - ) - .context("Unable to load model from file")?, + LlamaModel::load_from_file(&llama_backend, model_path.clone(), &model_params) + .context("Unable to load model from file")?, ); send_startup_signal( @@ -306,6 +314,7 @@ impl ContinuousBatchArbiter { inference_parameters, model_path: model_path.clone(), multimodal_context, + slot_aggregated_status: slot_aggregated_status_manager.slot_aggregated_status.clone(), token_bos_str: model.token_to_piece( &SampledToken::Content(model.token_bos()), &mut special_token_decoder, diff --git a/paddler_agent/src/continuous_batch_scheduler/advance_generating_phase.rs b/paddler_agent/src/continuous_batch_scheduler/advance_generating_phase.rs index 75b25ab62..cb14f0374 100644 --- a/paddler_agent/src/continuous_batch_scheduler/advance_generating_phase.rs +++ b/paddler_agent/src/continuous_batch_scheduler/advance_generating_phase.rs @@ -49,7 +49,13 @@ impl AdvanceGeneratingPhase<'_> { }) .run(request, batch_index) { - SampleOutcome::Sampled(token) => token, + SampleOutcome::Sampled(token) => { + self.scheduler_context + .slot_aggregated_status + .record_generated_token(); + + token + } SampleOutcome::AllCandidatesEliminated => { error!( "{:?}: sequence {} sampling exhausted candidates", diff --git a/paddler_agent/src/continuous_batch_scheduler_context.rs b/paddler_agent/src/continuous_batch_scheduler_context.rs index d9473563c..0034850ff 100644 --- a/paddler_agent/src/continuous_batch_scheduler_context.rs +++ b/paddler_agent/src/continuous_batch_scheduler_context.rs @@ -6,6 +6,7 @@ use llama_cpp_bindings::mtmd::MtmdContext; use paddler_messaging::inference_parameters::InferenceParameters; use crate::chat_template_renderer::ChatTemplateRenderer; +use crate::slot_aggregated_status::SlotAggregatedStatus; pub struct ContinuousBatchSchedulerContext { pub agent_name: Option, @@ -15,6 +16,7 @@ pub struct ContinuousBatchSchedulerContext { pub model: Arc, pub model_path: PathBuf, pub multimodal_context: Option>, + pub slot_aggregated_status: Arc, pub token_bos_str: String, pub token_eos_str: String, pub token_nl_str: String, diff --git a/paddler_agent/src/lib.rs b/paddler_agent/src/lib.rs index f95c51a75..f0d7d941e 100644 --- a/paddler_agent/src/lib.rs +++ b/paddler_agent/src/lib.rs @@ -32,6 +32,7 @@ pub mod embedding_input_tokenized; mod from_request_params; pub mod generate_embedding_batch_request; pub mod grammar_sampler; +pub mod list_gpu_devices; pub mod llamacpp_arbiter_service; pub mod management_socket_client_service; pub mod model_metadata_holder; @@ -59,6 +60,7 @@ pub mod slot_aggregated_status; pub mod slot_aggregated_status_download_progress; pub mod slot_aggregated_status_manager; pub mod slot_guard; +pub mod token_throughput_meter; pub mod tool_call_buffer; pub mod tool_call_event; pub mod tool_call_pipeline; diff --git a/paddler_agent/src/list_gpu_devices.rs b/paddler_agent/src/list_gpu_devices.rs new file mode 100644 index 000000000..ff60b2da3 --- /dev/null +++ b/paddler_agent/src/list_gpu_devices.rs @@ -0,0 +1,18 @@ +use anyhow::Context; +use anyhow::Result; +use llama_cpp_bindings::llama_backend::LlamaBackend; + +pub use llama_cpp_bindings::LlamaBackendDevice; +pub use llama_cpp_bindings::LlamaBackendDeviceType; + +/// Lists every backend device (GPU, integrated GPU, accelerator, or CPU) that llama.cpp can see +/// on this machine, in the same index order accepted by `--gpu-devices`. +/// +/// # Errors +/// Returns an error if the llama.cpp backend fails to initialize. +pub fn list_gpu_devices() -> Result> { + let _llama_backend = + LlamaBackend::init().context("Unable to initialize llama.cpp backend")?; + + Ok(llama_cpp_bindings::list_llama_ggml_backend_devices()) +} diff --git a/paddler_agent/src/llamacpp_arbiter_service.rs b/paddler_agent/src/llamacpp_arbiter_service.rs index 1cc4ad8b3..e9f3a9c3c 100644 --- a/paddler_agent/src/llamacpp_arbiter_service.rs +++ b/paddler_agent/src/llamacpp_arbiter_service.rs @@ -33,6 +33,7 @@ async fn apply_state( agent_applicable_state: Option<&AgentApplicableState>, agent_name: Option<&str>, desired_slots_total: i32, + gpu_devices: &[usize], model_metadata_holder: &Arc, slot_aggregated_status_manager: &Arc, continuous_batch_arbiter_handle: &mut Option, @@ -52,6 +53,7 @@ async fn apply_state( applicable_state, agent_name.map(str::to_owned), desired_slots_total, + gpu_devices.to_vec(), model_metadata_holder.clone(), slot_aggregated_status_manager.clone(), ) { @@ -110,6 +112,7 @@ async fn try_to_apply_state( agent_applicable_state: Option<&AgentApplicableState>, agent_name: Option<&str>, desired_slots_total: i32, + gpu_devices: &[usize], model_metadata_holder: &Arc, slot_aggregated_status_manager: &Arc, continuous_batch_arbiter_handle: &mut Option, @@ -119,6 +122,7 @@ async fn try_to_apply_state( agent_applicable_state, agent_name, desired_slots_total, + gpu_devices, model_metadata_holder, slot_aggregated_status_manager, continuous_batch_arbiter_handle, @@ -150,6 +154,7 @@ pub struct LlamaCppArbiterService { pub continue_from_raw_prompt_request_rx: mpsc::UnboundedReceiver, pub desired_slots_total: i32, pub generate_embedding_batch_request_rx: mpsc::UnboundedReceiver, + pub gpu_devices: Vec, pub continuous_batch_arbiter_handle: Option, pub model_metadata_holder: Arc, pub slot_aggregated_status_manager: Arc, @@ -170,6 +175,7 @@ impl Service for LlamaCppArbiterService { mut continue_from_raw_prompt_request_rx, desired_slots_total, mut generate_embedding_batch_request_rx, + gpu_devices, mut continuous_batch_arbiter_handle, model_metadata_holder, slot_aggregated_status_manager, @@ -203,6 +209,7 @@ impl Service for LlamaCppArbiterService { agent_applicable_state.as_ref(), agent_name.as_deref(), desired_slots_total, + &gpu_devices, &model_metadata_holder, &slot_aggregated_status_manager, &mut continuous_batch_arbiter_handle, @@ -220,6 +227,7 @@ impl Service for LlamaCppArbiterService { agent_applicable_state.as_ref(), agent_name.as_deref(), desired_slots_total, + &gpu_devices, &model_metadata_holder, &slot_aggregated_status_manager, &mut continuous_batch_arbiter_handle, @@ -353,6 +361,7 @@ mod tests { None, None, 1, + &[], &model_metadata_holder, &slot_aggregated_status_manager, &mut continuous_batch_arbiter_handle, @@ -421,6 +430,7 @@ mod tests { continue_from_raw_prompt_request_rx, desired_slots_total: 1, generate_embedding_batch_request_rx, + gpu_devices: Vec::new(), continuous_batch_arbiter_handle: None, model_metadata_holder: Arc::new(ModelMetadataHolder::default()), slot_aggregated_status_manager: Arc::new(SlotAggregatedStatusManager::new(1)), diff --git a/paddler_agent/src/slot_aggregated_status.rs b/paddler_agent/src/slot_aggregated_status.rs index a66ce0a17..b97480bb8 100644 --- a/paddler_agent/src/slot_aggregated_status.rs +++ b/paddler_agent/src/slot_aggregated_status.rs @@ -12,6 +12,7 @@ use tokio::sync::watch; use crate::agent_issue_fix::AgentIssueFix; use crate::dispenses_slots::DispensesSlots; +use crate::token_throughput_meter::TokenThroughputMeter; use paddler_messaging::atomic_value::AtomicValue; use paddler_messaging::produces_snapshot::ProducesSnapshot; use paddler_messaging::subscribes_to_updates::SubscribesToUpdates; @@ -27,6 +28,7 @@ pub struct SlotAggregatedStatus { slots_processing: AtomicValue, slots_total: AtomicValue, state_application_status_code: AtomicValue, + token_throughput_meter: TokenThroughputMeter, update_tx: watch::Sender<()>, uses_chat_template_override: AtomicValue, version: AtomicValue, @@ -50,6 +52,7 @@ impl SlotAggregatedStatus { ), slots_processing: AtomicValue::::new(0), slots_total: AtomicValue::::new(0), + token_throughput_meter: TokenThroughputMeter::new(), update_tx, uses_chat_template_override: AtomicValue::::new(false), version: AtomicValue::::new(0), @@ -174,6 +177,21 @@ impl SlotAggregatedStatus { pub fn slots_processing_count(&self) -> i32 { self.slots_processing.get() } + + /// Records that this agent has just generated one token, contributing to + /// the tokens-per-second rate reported in status snapshots. + /// + /// This does not bump `version` or notify subscribers: it happens once + /// per generated token, far too often to treat as a state change worth + /// pushing immediately. The periodic status update (sent once a second + /// regardless of `version`) is what picks the latest rate up. + pub fn record_generated_token(&self) { + self.token_throughput_meter.record_token(); + } + + pub fn tokens_per_second(&self) -> f64 { + self.token_throughput_meter.tokens_per_second() + } } impl DispensesSlots for SlotAggregatedStatus { @@ -211,6 +229,7 @@ impl ProducesSnapshot for SlotAggregatedStatus { slots_processing: self.slots_processing.get(), slots_total: self.slots_total.get(), state_application_status: self.state_application_status_code.get().try_into()?, + tokens_per_second: self.tokens_per_second(), uses_chat_template_override: self.uses_chat_template_override.get(), version: self.version.get(), }) @@ -341,6 +360,36 @@ mod tests { ); } + #[test] + fn make_snapshot_reports_zero_tokens_per_second_before_any_generation() { + let status = SlotAggregatedStatus::new(1); + + let snapshot = status.make_snapshot().unwrap(); + + assert_eq!(snapshot.tokens_per_second, 0.0); + } + + #[test] + fn record_generated_token_is_reflected_in_snapshot_after_a_window_closes() { + let status = SlotAggregatedStatus::new(1); + + for _ in 0..5 { + status.record_generated_token(); + } + + std::thread::sleep(std::time::Duration::from_millis(1050)); + + status.record_generated_token(); + + let snapshot = status.make_snapshot().unwrap(); + + assert!( + snapshot.tokens_per_second > 0.0, + "expected a positive tokens_per_second, got {}", + snapshot.tokens_per_second + ); + } + #[test] fn get_state_application_status_reflects_set_value() { let status = SlotAggregatedStatus::new(2); diff --git a/paddler_agent/src/token_throughput_meter.rs b/paddler_agent/src/token_throughput_meter.rs new file mode 100644 index 000000000..cc274565d --- /dev/null +++ b/paddler_agent/src/token_throughput_meter.rs @@ -0,0 +1,124 @@ +use std::time::Duration; +use std::time::Instant; + +use parking_lot::Mutex; + +/// How long a window stays open before its token count is turned into a +/// tokens/sec measurement. One second keeps the reported rate intuitive +/// (it is, literally, "tokens counted in the last second"). +const WINDOW_DURATION: Duration = Duration::from_secs(1); + +/// A measurement is considered stale, and reported as `0.0`, once this much +/// time has passed without a new token. Without this, a slot that produced a +/// quick burst and then went idle would keep reporting its last burst's rate +/// forever. +const STALE_AFTER: Duration = Duration::from_secs(2); + +struct ThroughputWindow { + tokens_per_second: f64, + window_started_at: Instant, + window_token_count: u64, +} + +/// Tracks generated tokens for a single agent and derives a continuously +/// updated, approximate tokens-per-second throughput value out of them. +/// +/// The meter is intentionally simple: it counts tokens in a rolling +/// one-second window and, once the window closes, turns that count into the +/// reported rate. There is nothing to configure. +pub struct TokenThroughputMeter { + window: Mutex, +} + +impl TokenThroughputMeter { + #[must_use] + pub fn new() -> Self { + Self { + window: Mutex::new(ThroughputWindow { + tokens_per_second: 0.0, + window_started_at: Instant::now(), + window_token_count: 0, + }), + } + } + + /// Records a single generated token, closing out and measuring the + /// current window if it has been open for at least [`WINDOW_DURATION`]. + pub fn record_token(&self) { + let mut window = self.window.lock(); + let elapsed = window.window_started_at.elapsed(); + + window.window_token_count += 1; + + if elapsed >= WINDOW_DURATION { + window.tokens_per_second = window.window_token_count as f64 / elapsed.as_secs_f64(); + window.window_token_count = 0; + window.window_started_at = Instant::now(); + } + } + + /// Returns the most recently measured tokens-per-second rate, or `0.0` + /// if generation has been idle for longer than [`STALE_AFTER`]. + #[must_use] + pub fn tokens_per_second(&self) -> f64 { + let window = self.window.lock(); + + if window.window_started_at.elapsed() > STALE_AFTER && window.window_token_count == 0 { + return 0.0; + } + + window.tokens_per_second + } +} + +impl Default for TokenThroughputMeter { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests { + use std::thread::sleep; + + use super::*; + + #[test] + fn reports_zero_before_any_tokens_are_recorded() { + let meter = TokenThroughputMeter::new(); + + assert_eq!(meter.tokens_per_second(), 0.0); + } + + #[test] + fn measures_rate_once_a_window_closes() { + let meter = TokenThroughputMeter::new(); + + for _ in 0..5 { + meter.record_token(); + } + + sleep(WINDOW_DURATION + Duration::from_millis(50)); + + meter.record_token(); + + let rate = meter.tokens_per_second(); + + assert!(rate > 0.0, "expected a positive rate, got {rate}"); + } + + #[test] + fn decays_to_zero_after_being_idle() { + let meter = TokenThroughputMeter::new(); + + for _ in 0..5 { + meter.record_token(); + } + + sleep(WINDOW_DURATION + Duration::from_millis(50)); + meter.record_token(); + sleep(STALE_AFTER + Duration::from_millis(50)); + + assert_eq!(meter.tokens_per_second(), 0.0); + } +} diff --git a/paddler_balancer/src/agent_controller.rs b/paddler_balancer/src/agent_controller.rs index be5d17170..6c1be5005 100644 --- a/paddler_balancer/src/agent_controller.rs +++ b/paddler_balancer/src/agent_controller.rs @@ -59,6 +59,7 @@ pub struct AgentController { pub slots_processing: AtomicValue, pub slots_total: AtomicValue, pub state_application_status_code: AtomicValue, + pub tokens_per_second: RwLock, pub uses_chat_template_override: AtomicValue, } @@ -95,6 +96,10 @@ impl AgentController { self.model_path.read().clone() } + pub fn get_tokens_per_second(&self) -> f64 { + *self.tokens_per_second.read() + } + pub fn set_download_filename(&self, filename: Option) { let mut locked_filename = self.download_filename.write(); @@ -113,6 +118,12 @@ impl AgentController { *locked_path = model_path; } + pub fn set_tokens_per_second(&self, tokens_per_second: f64) { + let mut locked_tokens_per_second = self.tokens_per_second.write(); + + *locked_tokens_per_second = tokens_per_second; + } + pub async fn stop_responding_to(&self, request_id: String) -> Result<()> { self.send_rpc_message(AgentJsonRpcMessage::Notification( AgentJsonRpcNotification::StopRespondingTo(request_id), @@ -134,6 +145,7 @@ impl AgentController { model_path, slots_total, state_application_status, + tokens_per_second, uses_chat_template_override, version, .. @@ -184,6 +196,12 @@ impl AgentController { self.set_model_path(model_path); } + if (tokens_per_second - self.get_tokens_per_second()).abs() > 0.01 { + changed = true; + + self.set_tokens_per_second(tokens_per_second); + } + if changed { AgentControllerUpdateResult::Updated } else { @@ -304,6 +322,7 @@ impl ProducesSnapshot for AgentController { slots_processing: self.slots_processing.get(), slots_total: self.slots_total.get(), state_application_status: self.state_application_status_code.get().try_into()?, + tokens_per_second: self.get_tokens_per_second(), uses_chat_template_override: self.uses_chat_template_override.get(), }) } @@ -368,6 +387,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), } } @@ -387,6 +407,7 @@ mod tests { slots_processing: 0, slots_total: 4, state_application_status: AgentStateApplicationStatus::Fresh, + tokens_per_second: 0.0, uses_chat_template_override: true, version: 1, }; @@ -418,6 +439,7 @@ mod tests { slots_processing: 0, slots_total: 0, state_application_status: AgentStateApplicationStatus::Fresh, + tokens_per_second: 0.0, uses_chat_template_override: false, version: 1, }; @@ -448,6 +470,7 @@ mod tests { slots_processing: 0, slots_total: 0, state_application_status: AgentStateApplicationStatus::Fresh, + tokens_per_second: 0.0, uses_chat_template_override: false, version: 1, }; @@ -481,6 +504,7 @@ mod tests { slots_processing: 0, slots_total: 0, state_application_status: AgentStateApplicationStatus::Fresh, + tokens_per_second: 0.0, uses_chat_template_override: false, version: 1, }; diff --git a/paddler_balancer/src/agent_controller_pool.rs b/paddler_balancer/src/agent_controller_pool.rs index 35922e6c1..4eb63b239 100644 --- a/paddler_balancer/src/agent_controller_pool.rs +++ b/paddler_balancer/src/agent_controller_pool.rs @@ -131,6 +131,16 @@ impl AgentControllerPool { slots_total, } } + + /// Sums the most recently reported tokens-per-second rate across every + /// registered agent, giving the balancer's overall generation throughput. + #[must_use] + pub fn total_tokens_per_second(&self) -> f64 { + self.agents + .iter() + .map(|entry| entry.value().get_tokens_per_second()) + .sum() + } } impl Default for AgentControllerPool { @@ -236,6 +246,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }) } @@ -297,6 +308,24 @@ mod tests { assert_eq!(total_slots.slots_total, 12); } + #[test] + fn total_tokens_per_second_sums_rate_across_agents() { + let pool = AgentControllerPool::default(); + + let first = agent_controller_with_slots(1, 4); + first.set_tokens_per_second(12.5); + + let second = agent_controller_with_slots(2, 8); + second.set_tokens_per_second(7.5); + + pool.register_agent_controller("first".to_owned(), first) + .unwrap(); + pool.register_agent_controller("second".to_owned(), second) + .unwrap(); + + assert_eq!(pool.total_tokens_per_second(), 20.0); + } + #[test] fn make_snapshot_includes_each_registered_agent() { let pool = AgentControllerPool::default(); diff --git a/paddler_balancer/src/buffered_request_manager.rs b/paddler_balancer/src/buffered_request_manager.rs index 85f117409..46304d744 100644 --- a/paddler_balancer/src/buffered_request_manager.rs +++ b/paddler_balancer/src/buffered_request_manager.rs @@ -143,6 +143,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }); @@ -219,6 +220,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }); @@ -268,6 +270,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }); diff --git a/paddler_balancer/src/controls_manages_senders_endpoint.rs b/paddler_balancer/src/controls_manages_senders_endpoint.rs index c4b6cefbc..86f47284e 100644 --- a/paddler_balancer/src/controls_manages_senders_endpoint.rs +++ b/paddler_balancer/src/controls_manages_senders_endpoint.rs @@ -110,6 +110,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }), ) diff --git a/paddler_balancer/src/inference_service/http_route/api/post_generate_embedding_batch.rs b/paddler_balancer/src/inference_service/http_route/api/post_generate_embedding_batch.rs index bd8d3472b..721ada85e 100644 --- a/paddler_balancer/src/inference_service/http_route/api/post_generate_embedding_batch.rs +++ b/paddler_balancer/src/inference_service/http_route/api/post_generate_embedding_batch.rs @@ -239,6 +239,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }) } diff --git a/paddler_balancer/src/inference_service/http_route/api/ws_inference_socket/mod.rs b/paddler_balancer/src/inference_service/http_route/api/ws_inference_socket/mod.rs index e3797efdc..aadd04882 100644 --- a/paddler_balancer/src/inference_service/http_route/api/ws_inference_socket/mod.rs +++ b/paddler_balancer/src/inference_service/http_route/api/ws_inference_socket/mod.rs @@ -342,6 +342,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }); diff --git a/paddler_balancer/src/management_service/http_route/api/get_agents.rs b/paddler_balancer/src/management_service/http_route/api/get_agents.rs index 2f2414652..752d1d343 100644 --- a/paddler_balancer/src/management_service/http_route/api/get_agents.rs +++ b/paddler_balancer/src/management_service/http_route/api/get_agents.rs @@ -83,6 +83,7 @@ mod tests { slots_processing: AtomicValue::::new(0), slots_total: AtomicValue::::new(0), state_application_status_code: AtomicValue::::new(status_code), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }) } diff --git a/paddler_balancer/src/management_service/http_route/api/ws_agent_socket/agent_socket_controller_context.rs b/paddler_balancer/src/management_service/http_route/api/ws_agent_socket/agent_socket_controller_context.rs index 9631b28b6..24045e5d3 100644 --- a/paddler_balancer/src/management_service/http_route/api/ws_agent_socket/agent_socket_controller_context.rs +++ b/paddler_balancer/src/management_service/http_route/api/ws_agent_socket/agent_socket_controller_context.rs @@ -92,6 +92,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }), ) diff --git a/paddler_balancer/src/management_service/http_route/api/ws_agent_socket/mod.rs b/paddler_balancer/src/management_service/http_route/api/ws_agent_socket/mod.rs index 2eb425c55..655b8bdf3 100644 --- a/paddler_balancer/src/management_service/http_route/api/ws_agent_socket/mod.rs +++ b/paddler_balancer/src/management_service/http_route/api/ws_agent_socket/mod.rs @@ -124,6 +124,7 @@ impl ControlsWebSocketEndpoint for AgentSocketController { slots_processing, slots_total, state_application_status, + tokens_per_second, uses_chat_template_override, version, }, @@ -159,6 +160,7 @@ impl ControlsWebSocketEndpoint for AgentSocketController { state_application_status_code: AtomicValue::::new( state_application_status as i32, ), + tokens_per_second: RwLock::new(tokens_per_second), uses_chat_template_override: AtomicValue::::new( uses_chat_template_override, ), @@ -409,6 +411,7 @@ mod tests { slots_processing: 0, slots_total: 1, state_application_status: AgentStateApplicationStatus::Fresh, + tokens_per_second: 0.0, uses_chat_template_override: false, version: 0, }, diff --git a/paddler_balancer/src/management_service/http_route/get_metrics.rs b/paddler_balancer/src/management_service/http_route/get_metrics.rs index 276c01990..2d4fe9ff7 100644 --- a/paddler_balancer/src/management_service/http_route/get_metrics.rs +++ b/paddler_balancer/src/management_service/http_route/get_metrics.rs @@ -24,6 +24,7 @@ async fn respond(app_data: Data) -> Result) -> Result::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }) } diff --git a/paddler_balancer/src/request_from_agent.rs b/paddler_balancer/src/request_from_agent.rs index 58f450132..985e848dc 100644 --- a/paddler_balancer/src/request_from_agent.rs +++ b/paddler_balancer/src/request_from_agent.rs @@ -453,6 +453,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }); diff --git a/paddler_balancer/src/statsd_service/mod.rs b/paddler_balancer/src/statsd_service/mod.rs index 7de8c06b5..5280c94c5 100644 --- a/paddler_balancer/src/statsd_service/mod.rs +++ b/paddler_balancer/src/statsd_service/mod.rs @@ -38,16 +38,24 @@ impl StatsdService { slots_total, } = self.agent_controller_pool.total_slots(); let requests_buffered = self.buffered_request_manager.buffered_request_counter.get(); + let tokens_per_second = self.agent_controller_pool.total_tokens_per_second(); let slots_processing = u64::try_from(slots_processing).context("slots_processing count is negative")?; let slots_total = u64::try_from(slots_total).context("slots_total count is negative")?; let requests_buffered = u64::try_from(requests_buffered).context("requests_buffered count is negative")?; + #[expect( + clippy::cast_possible_truncation, + clippy::cast_sign_loss, + reason = "StatsD gauges are unsigned integers; sub-token precision is not meaningful" + )] + let tokens_per_second = tokens_per_second.round() as u64; client.gauge("slots_processing", slots_processing)?; client.gauge("slots_total", slots_total)?; client.gauge("requests_buffered", requests_buffered)?; + client.gauge("tokens_per_second", tokens_per_second)?; client.flush()?; Ok(()) @@ -152,6 +160,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }), ) @@ -199,7 +208,7 @@ mod tests { let mut received_lines: Vec = Vec::new(); let mut datagram = [0_u8; 1024]; - for _ in 0..3 { + for _ in 0..4 { let byte_count = receiver.recv(&mut datagram).await.unwrap(); received_lines.push(String::from_utf8(datagram[..byte_count].to_vec()).unwrap()); @@ -208,6 +217,7 @@ mod tests { assert!(received_lines.contains(&"paddler.slots_processing:0|g".to_owned())); assert!(received_lines.contains(&"paddler.slots_total:0|g".to_owned())); assert!(received_lines.contains(&"paddler.requests_buffered:0|g".to_owned())); + assert!(received_lines.contains(&"paddler.tokens_per_second:0|g".to_owned())); } #[tokio::test] @@ -227,6 +237,7 @@ mod tests { "paddler.slots_processing:0|g".to_owned(), "paddler.slots_total:0|g".to_owned(), "paddler.requests_buffered:0|g".to_owned(), + "paddler.tokens_per_second:0|g".to_owned(), ]; assert!(expected_first_tick_lines.contains(&first_line)); diff --git a/paddler_bootstrap/src/agent_runner.rs b/paddler_bootstrap/src/agent_runner.rs index aa5aa2dea..6ce9c1f75 100644 --- a/paddler_bootstrap/src/agent_runner.rs +++ b/paddler_bootstrap/src/agent_runner.rs @@ -13,6 +13,7 @@ use crate::service_thread::ServiceThread; pub struct AgentRunnerParams { pub agent_name: Option, pub cancellation_token: CancellationToken, + pub gpu_devices: Vec, pub management_address: String, pub slots: i32, } @@ -28,11 +29,12 @@ impl AgentRunner { AgentRunnerParams { agent_name, cancellation_token, + gpu_devices, management_address, slots, }: AgentRunnerParams, ) -> Self { - let bundle = AgentServiceBundle::new(agent_name, &management_address, slots); + let bundle = AgentServiceBundle::new(agent_name, &management_address, slots, gpu_devices); let slot_aggregated_status = bundle.slot_aggregated_status.clone(); let thread = ServiceThread::spawn(cancellation_token, move |task_shutdown| { diff --git a/paddler_bootstrap/src/agent_service_bundle.rs b/paddler_bootstrap/src/agent_service_bundle.rs index 3e7019d9c..a80c97c86 100644 --- a/paddler_bootstrap/src/agent_service_bundle.rs +++ b/paddler_bootstrap/src/agent_service_bundle.rs @@ -27,7 +27,12 @@ pub struct AgentServiceBundle { impl AgentServiceBundle { #[must_use] - pub fn new(agent_name: Option, management_address: &str, slots: i32) -> Self { + pub fn new( + agent_name: Option, + management_address: &str, + slots: i32, + gpu_devices: Vec, + ) -> Self { let (agent_desired_state_tx, agent_desired_state_rx) = mpsc::unbounded_channel::(); let ( @@ -54,6 +59,7 @@ impl AgentServiceBundle { continue_from_raw_prompt_request_rx, desired_slots_total: slots, generate_embedding_batch_request_rx, + gpu_devices, continuous_batch_arbiter_handle: None, model_metadata_holder: model_metadata_holder.clone(), slot_aggregated_status_manager, diff --git a/paddler_bootstrap/src/lib.rs b/paddler_bootstrap/src/lib.rs index c49f74145..cdbed7edf 100644 --- a/paddler_bootstrap/src/lib.rs +++ b/paddler_bootstrap/src/lib.rs @@ -2,5 +2,6 @@ pub mod agent_runner; pub mod agent_service_bundle; pub mod balancer_runner; pub mod balancer_service_bundle; +pub mod list_gpu_devices; pub mod run_service_manager; pub mod service_thread; diff --git a/paddler_bootstrap/src/list_gpu_devices.rs b/paddler_bootstrap/src/list_gpu_devices.rs new file mode 100644 index 000000000..026c919bd --- /dev/null +++ b/paddler_bootstrap/src/list_gpu_devices.rs @@ -0,0 +1,3 @@ +pub use paddler_agent::list_gpu_devices::LlamaBackendDevice; +pub use paddler_agent::list_gpu_devices::LlamaBackendDeviceType; +pub use paddler_agent::list_gpu_devices::list_gpu_devices; diff --git a/paddler_bootstrap/tests/runners.rs b/paddler_bootstrap/tests/runners.rs index fbf6aaeae..22c51746c 100644 --- a/paddler_bootstrap/tests/runners.rs +++ b/paddler_bootstrap/tests/runners.rs @@ -83,6 +83,7 @@ fn make_agent_runner_params( agent_name: Some("test-agent".to_owned()), management_address: management_addr.to_string(), cancellation_token, + gpu_devices: Vec::new(), slots: 1, } } diff --git a/paddler_cli/src/cmd/agent.rs b/paddler_cli/src/cmd/agent.rs index 79a9a9895..b64bf9497 100644 --- a/paddler_cli/src/cmd/agent.rs +++ b/paddler_cli/src/cmd/agent.rs @@ -12,6 +12,12 @@ use super::value_parser::parse_socket_addr::parse_socket_addr; #[derive(Parser)] pub struct Agent { + #[arg(long, value_delimiter = ',')] + /// Restrict inference to specific backend device indices (comma-separated, e.g. "0" or + /// "0,1"). Run `paddler list-gpu-devices` to see available indices and their names. Omit + /// this flag to use every detected device, as before. + gpu_devices: Vec, + #[arg(long, value_parser = parse_socket_addr)] /// Address of the management server that the agent will connect to management_addr: ResolvedSocketAddr, @@ -32,6 +38,7 @@ impl Handler for Agent { self.name.clone(), &self.management_addr.socket_addr.to_string(), self.slots, + self.gpu_devices, ); let mut service_manager = ServiceManager::default(); diff --git a/paddler_cli/src/cmd/list_gpu_devices.rs b/paddler_cli/src/cmd/list_gpu_devices.rs new file mode 100644 index 000000000..d670e235c --- /dev/null +++ b/paddler_cli/src/cmd/list_gpu_devices.rs @@ -0,0 +1,51 @@ +use anyhow::Result; +use async_trait::async_trait; +use clap::Parser; +use command_handler::handler::Handler; +use paddler_bootstrap::list_gpu_devices::LlamaBackendDeviceType; +use paddler_bootstrap::list_gpu_devices::list_gpu_devices; +use tokio_util::sync::CancellationToken; + +#[derive(Parser)] +/// Lists the backend devices llama.cpp detects on this machine, in the index order accepted by +/// `paddler agent --gpu-devices` +pub struct ListGpuDevices; + +#[async_trait(?Send)] +impl Handler for ListGpuDevices { + async fn handle(self, _shutdown: CancellationToken) -> Result<()> { + let devices = list_gpu_devices()?; + + if devices.is_empty() { + println!("No backend devices detected."); + + return Ok(()); + } + + for device in devices { + let device_type = match device.device_type { + LlamaBackendDeviceType::Accelerator => "accelerator", + LlamaBackendDeviceType::Cpu => "cpu", + LlamaBackendDeviceType::Gpu => "gpu", + LlamaBackendDeviceType::IntegratedGpu => "integrated gpu", + LlamaBackendDeviceType::Unknown => "unknown", + }; + + println!( + "{index}: {name} ({device_type}, {backend} backend, {memory_free_mib} MiB free / {memory_total_mib} MiB total) - {description}", + index = device.index, + name = device.name, + backend = device.backend, + memory_free_mib = device.memory_free / 1024 / 1024, + memory_total_mib = device.memory_total / 1024 / 1024, + description = device.description, + ); + } + + println!( + "\nPass one or more of the indices above to `paddler agent --gpu-devices` to restrict inference to those devices, e.g. --gpu-devices 0 or --gpu-devices 0,1." + ); + + Ok(()) + } +} diff --git a/paddler_cli/src/cmd/mod.rs b/paddler_cli/src/cmd/mod.rs index 1ca021419..cc2991de2 100644 --- a/paddler_cli/src/cmd/mod.rs +++ b/paddler_cli/src/cmd/mod.rs @@ -1,3 +1,4 @@ pub mod agent; pub mod balancer; +pub mod list_gpu_devices; pub mod value_parser; diff --git a/paddler_cli/src/lib.rs b/paddler_cli/src/lib.rs index 0e3fb4700..44dc54166 100644 --- a/paddler_cli/src/lib.rs +++ b/paddler_cli/src/lib.rs @@ -5,6 +5,7 @@ use clap::Parser; use clap::Subcommand; use cmd::agent::Agent; use cmd::balancer::Balancer; +use cmd::list_gpu_devices::ListGpuDevices; use command_handler::handler::Handler as _; use command_handler::shutdown_signal::register_shutdown_signals; #[cfg(feature = "web_admin_panel")] @@ -34,6 +35,8 @@ enum Commands { Agent(Agent), /// Distributes incoming requests among agents Balancer(Box), + /// Lists the backend devices llama.cpp detects on this machine + ListGpuDevices(ListGpuDevices), } pub fn run() -> Result<()> { @@ -50,6 +53,7 @@ pub fn run() -> Result<()> { (*handler).handle(shutdown).await } + Some(Commands::ListGpuDevices(handler)) => handler.handle(shutdown).await, None => Ok(()), } }) diff --git a/paddler_client_javascript/src/schemas/Agent.ts b/paddler_client_javascript/src/schemas/Agent.ts index 1ad87eaad..5788f20a1 100644 --- a/paddler_client_javascript/src/schemas/Agent.ts +++ b/paddler_client_javascript/src/schemas/Agent.ts @@ -22,6 +22,7 @@ export const AgentSchema = z "Fresh", "Stuck", ]), + tokens_per_second: z.number(), uses_chat_template_override: z.boolean(), }) .strict(); diff --git a/paddler_client_javascript/tests/schemas/Agent.test.ts b/paddler_client_javascript/tests/schemas/Agent.test.ts index 3ad0ba4ef..647667910 100644 --- a/paddler_client_javascript/tests/schemas/Agent.test.ts +++ b/paddler_client_javascript/tests/schemas/Agent.test.ts @@ -17,6 +17,7 @@ test("parses a fully populated agent payload", function () { slots_processing: 1, slots_total: 4, state_application_status: "Applied", + tokens_per_second: 12.5, uses_chat_template_override: false, }); @@ -39,6 +40,7 @@ test("rejects an unknown state_application_status", function () { slots_processing: 0, slots_total: 1, state_application_status: "Unknown", + tokens_per_second: 0, uses_chat_template_override: false, }); }); diff --git a/paddler_gui/Cargo.toml b/paddler_gui/Cargo.toml index 1dc8ead4d..ae0301639 100644 --- a/paddler_gui/Cargo.toml +++ b/paddler_gui/Cargo.toml @@ -39,6 +39,7 @@ workspace = true default = [] cuda = ["paddler_bootstrap/cuda"] metal = ["paddler_bootstrap/metal"] +vulkan = ["paddler_bootstrap/vulkan"] web_admin_panel = [ "dep:esbuild-metafile", "paddler_balancer/web_admin_panel", diff --git a/paddler_gui/src/agent_running_data.rs b/paddler_gui/src/agent_running_data.rs index dae4782a2..ede9e0f25 100644 --- a/paddler_gui/src/agent_running_data.rs +++ b/paddler_gui/src/agent_running_data.rs @@ -23,6 +23,7 @@ impl AgentRunningData { slots_processing: status.slots_processing, slots_total: status.slots_total, state_application_status: status.state_application_status, + tokens_per_second: status.tokens_per_second, uses_chat_template_override: status.uses_chat_template_override, }; } diff --git a/paddler_gui/src/app.rs b/paddler_gui/src/app.rs index 26b741ee0..fa19b56f4 100644 --- a/paddler_gui/src/app.rs +++ b/paddler_gui/src/app.rs @@ -377,6 +377,7 @@ impl App { Task::stream(iced::stream::channel(1, async move |mut output| { let mut runner = AgentRunner::start(AgentRunnerParams { agent_name, + gpu_devices: Vec::new(), management_address, cancellation_token: cancel, slots, diff --git a/paddler_gui/src/running_balancer_snapshot.rs b/paddler_gui/src/running_balancer_snapshot.rs index 91f6960e4..f591dca73 100644 --- a/paddler_gui/src/running_balancer_snapshot.rs +++ b/paddler_gui/src/running_balancer_snapshot.rs @@ -98,6 +98,7 @@ mod tests { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), }) } diff --git a/paddler_gui/src/screen.rs b/paddler_gui/src/screen.rs index 0d180df6b..ff5c4b3af 100644 --- a/paddler_gui/src/screen.rs +++ b/paddler_gui/src/screen.rs @@ -84,6 +84,7 @@ impl Screen { slots_processing: 0, slots_total: 0, state_application_status: AgentStateApplicationStatus::Fresh, + tokens_per_second: 0.0, uses_chat_template_override: false, }, } diff --git a/paddler_gui/src/ui/view_agent_card.rs b/paddler_gui/src/ui/view_agent_card.rs index e50cdef32..41cbb5663 100644 --- a/paddler_gui/src/ui/view_agent_card.rs +++ b/paddler_gui/src/ui/view_agent_card.rs @@ -96,9 +96,20 @@ pub fn view_agent_card( snapshot.slots_processing, snapshot.slots_total, snapshot.desired_slots_total, ); + let mut status_row_right = column![].spacing(SPACING_HALF); + + status_row_right = + status_row_right.push(text(format!("Slots: {slots_label}")).font(REGULAR)); + + if snapshot.tokens_per_second > 0.0 { + status_row_right = status_row_right.push( + text(format!("{:.1} tok/s", snapshot.tokens_per_second)).font(REGULAR), + ); + } + let status_row_content = row![ container(status_row_left).width(Fill), - text(format!("Slots: {slots_label}")).font(REGULAR), + status_row_right, ]; let card_content = column![name_row, status_row_content,].spacing(SPACING_BASE); diff --git a/paddler_messaging/src/agent_controller_snapshot.rs b/paddler_messaging/src/agent_controller_snapshot.rs index 559ffefa3..bcf7bd8c1 100644 --- a/paddler_messaging/src/agent_controller_snapshot.rs +++ b/paddler_messaging/src/agent_controller_snapshot.rs @@ -21,5 +21,6 @@ pub struct AgentControllerSnapshot { pub slots_processing: i32, pub slots_total: i32, pub state_application_status: AgentStateApplicationStatus, + pub tokens_per_second: f64, pub uses_chat_template_override: bool, } diff --git a/paddler_messaging/src/management_socket/balancer/notification_params/register_agent_params.rs b/paddler_messaging/src/management_socket/balancer/notification_params/register_agent_params.rs index e6153ebc3..20384abde 100644 --- a/paddler_messaging/src/management_socket/balancer/notification_params/register_agent_params.rs +++ b/paddler_messaging/src/management_socket/balancer/notification_params/register_agent_params.rs @@ -2,7 +2,7 @@ use crate::slot_aggregated_status_snapshot::SlotAggregatedStatusSnapshot; use serde::Deserialize; use serde::Serialize; -#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)] #[serde(deny_unknown_fields)] pub struct RegisterAgentParams { pub name: Option, diff --git a/paddler_messaging/src/management_socket/balancer/notification_params/update_agent_status_params.rs b/paddler_messaging/src/management_socket/balancer/notification_params/update_agent_status_params.rs index f6e17c39f..3d3c91c6b 100644 --- a/paddler_messaging/src/management_socket/balancer/notification_params/update_agent_status_params.rs +++ b/paddler_messaging/src/management_socket/balancer/notification_params/update_agent_status_params.rs @@ -2,7 +2,7 @@ use crate::slot_aggregated_status_snapshot::SlotAggregatedStatusSnapshot; use serde::Deserialize; use serde::Serialize; -#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)] #[serde(deny_unknown_fields)] pub struct UpdateAgentStatusParams { pub slot_aggregated_status_snapshot: SlotAggregatedStatusSnapshot, diff --git a/paddler_messaging/src/slot_aggregated_status_snapshot.rs b/paddler_messaging/src/slot_aggregated_status_snapshot.rs index 94dfdbc7f..6b685dbce 100644 --- a/paddler_messaging/src/slot_aggregated_status_snapshot.rs +++ b/paddler_messaging/src/slot_aggregated_status_snapshot.rs @@ -6,7 +6,7 @@ use serde::Serialize; use crate::agent_issue::AgentIssue; use crate::agent_state_application_status::AgentStateApplicationStatus; -#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)] +#[derive(Clone, Debug, Default, Deserialize, PartialEq, Serialize)] #[serde(deny_unknown_fields)] pub struct SlotAggregatedStatusSnapshot { pub desired_slots_total: i32, @@ -19,6 +19,7 @@ pub struct SlotAggregatedStatusSnapshot { pub slots_processing: i32, pub slots_total: i32, pub state_application_status: AgentStateApplicationStatus, + pub tokens_per_second: f64, pub uses_chat_template_override: bool, pub version: i32, } diff --git a/paddler_test_cluster_harness/src/agents_status/assert_slots_total_at_least.rs b/paddler_test_cluster_harness/src/agents_status/assert_slots_total_at_least.rs index 5a37df6ad..2c3690b87 100644 --- a/paddler_test_cluster_harness/src/agents_status/assert_slots_total_at_least.rs +++ b/paddler_test_cluster_harness/src/agents_status/assert_slots_total_at_least.rs @@ -39,6 +39,7 @@ mod tests { slots_processing: 0, slots_total, state_application_status: AgentStateApplicationStatus::Applied, + tokens_per_second: 0.0, uses_chat_template_override: false, }], } diff --git a/paddler_test_cluster_harness/src/agents_stream_watcher.rs b/paddler_test_cluster_harness/src/agents_stream_watcher.rs index 1c733b93b..d0dfcd043 100644 --- a/paddler_test_cluster_harness/src/agents_stream_watcher.rs +++ b/paddler_test_cluster_harness/src/agents_stream_watcher.rs @@ -234,6 +234,7 @@ mod tests { slots_processing: 0, slots_total, state_application_status: AgentStateApplicationStatus::Fresh, + tokens_per_second: 0.0, uses_chat_template_override: false, } } diff --git a/paddler_tests/src/in_process_agent_spawner.rs b/paddler_tests/src/in_process_agent_spawner.rs index 26d7ddb1e..a617c98e6 100644 --- a/paddler_tests/src/in_process_agent_spawner.rs +++ b/paddler_tests/src/in_process_agent_spawner.rs @@ -24,6 +24,7 @@ impl AgentSpawner for InProcessAgentSpawner { let runner = AgentRunner::start(AgentRunnerParams { agent_name: Some(config.name.clone()), cancellation_token: CancellationToken::new(), + gpu_devices: Vec::new(), management_address: self.management_address.clone(), slots: config.slot_count, }); diff --git a/paddler_tests/src/make_agent_controller_without_remote_agent.rs b/paddler_tests/src/make_agent_controller_without_remote_agent.rs index 908d590a2..6cc4f0559 100644 --- a/paddler_tests/src/make_agent_controller_without_remote_agent.rs +++ b/paddler_tests/src/make_agent_controller_without_remote_agent.rs @@ -43,6 +43,7 @@ pub fn make_agent_controller_without_remote_agent(id: &str) -> AgentController { state_application_status_code: AtomicValue::::new( AgentStateApplicationStatus::Fresh as i32, ), + tokens_per_second: RwLock::new(0.0), uses_chat_template_override: AtomicValue::::new(false), } } diff --git a/paddler_tests/tests/arbiter_spawn_is_cancelled_when_shutdown_fires_before_the_model_loads.rs b/paddler_tests/tests/arbiter_spawn_is_cancelled_when_shutdown_fires_before_the_model_loads.rs index 00fc29f99..80160ec56 100644 --- a/paddler_tests/tests/arbiter_spawn_is_cancelled_when_shutdown_fires_before_the_model_loads.rs +++ b/paddler_tests/tests/arbiter_spawn_is_cancelled_when_shutdown_fires_before_the_model_loads.rs @@ -51,6 +51,7 @@ async fn arbiter_spawn_is_cancelled_when_shutdown_fires_before_the_model_loads() applicable_state, None, 1, + Vec::new(), Arc::new(ModelMetadataHolder::new()), Arc::new(SlotAggregatedStatusManager::new(1)), ) { diff --git a/paddler_tests/tests/harness_agents_watcher.rs b/paddler_tests/tests/harness_agents_watcher.rs index c6527c608..3990282e1 100644 --- a/paddler_tests/tests/harness_agents_watcher.rs +++ b/paddler_tests/tests/harness_agents_watcher.rs @@ -26,6 +26,7 @@ fn make_snapshot(agent_id: &str, slots_total: i32) -> AgentControllerPoolSnapsho slots_processing: 0, slots_total, state_application_status: AgentStateApplicationStatus::Applied, + tokens_per_second: 0.0, uses_chat_template_override: false, }], } diff --git a/resources/ts/components/AgentListAgentStatus.module.css b/resources/ts/components/AgentListAgentStatus.module.css index 351886a0f..060823781 100644 --- a/resources/ts/components/AgentListAgentStatus.module.css +++ b/resources/ts/components/AgentListAgentStatus.module.css @@ -2,7 +2,7 @@ align-items: center; display: grid; gap: var(--spacing-half); - grid-template-columns: 1fr auto; + grid-template-columns: 1fr auto auto; justify-items: flex-start; overflow: hidden; text-overflow: ellipsis; diff --git a/resources/ts/components/AgentListAgentStatus.tsx b/resources/ts/components/AgentListAgentStatus.tsx index ceb51f143..6630b262c 100644 --- a/resources/ts/components/AgentListAgentStatus.tsx +++ b/resources/ts/components/AgentListAgentStatus.tsx @@ -10,6 +10,7 @@ export function AgentListAgentStatus({ slots_processing, slots_total, state_application_status, + tokens_per_second, }, }: { agent: Agent; @@ -33,6 +34,11 @@ export function AgentListAgentStatus({ {slots_processing}/{slots_total}/{desired_slots_total} + {tokens_per_second > 0 && ( + + {tokens_per_second.toFixed(1)} tok/s + + )} ); case "AttemptedAndNotAppliable":