Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 14 additions & 7 deletions model_gateway/src/routers/grpc/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,8 @@ use super::{
MessagePreparationStage, MessageRequestBuildingStage,
MessageResponseProcessingStage,
},
*,
ChatGeneratePreparationStage, ChatGenerateRequestBuildingStage,
ChatGenerateResponseProcessingStage,
},
streaming,
},
Expand Down Expand Up @@ -132,17 +133,20 @@ impl RequestPipeline {
));

let stages: Vec<Box<dyn PipelineStage>> = vec![
Box::new(PreparationStage::new()),
Box::new(ChatGeneratePreparationStage::new()),
Box::new(WorkerSelectionStage::new(
worker_registry,
policy_registry,
WorkerSelectionMode::Regular,
)),
Box::new(ClientAcquisitionStage),
Box::new(RequestBuildingStage::new(false)), // No PD metadata
Box::new(ChatGenerateRequestBuildingStage::new(false)), // No PD metadata
Box::new(DispatchMetadataStage),
Box::new(RequestExecutionStage::new(ExecutionMode::Single)),
Box::new(ResponseProcessingStage::new(processor, streaming_processor)),
Box::new(ChatGenerateResponseProcessingStage::new(
processor,
streaming_processor,
)),
];

Self {
Expand Down Expand Up @@ -235,17 +239,20 @@ impl RequestPipeline {
));

let stages: Vec<Box<dyn PipelineStage>> = vec![
Box::new(PreparationStage::new()),
Box::new(ChatGeneratePreparationStage::new()),
Box::new(WorkerSelectionStage::new(
worker_registry,
policy_registry,
WorkerSelectionMode::PrefillDecode,
)),
Box::new(ClientAcquisitionStage),
Box::new(RequestBuildingStage::new(true)), // Inject PD metadata
Box::new(ChatGenerateRequestBuildingStage::new(true)), // Inject PD metadata
Box::new(DispatchMetadataStage),
Box::new(RequestExecutionStage::new(ExecutionMode::DualDispatch)),
Box::new(ResponseProcessingStage::new(processor, streaming_processor)),
Box::new(ChatGenerateResponseProcessingStage::new(
processor,
streaming_processor,
)),
];

Self {
Expand Down
8 changes: 4 additions & 4 deletions model_gateway/src/routers/grpc/regular/stages/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ pub(crate) mod preparation;
pub(crate) mod request_building;
pub(crate) mod response_processing;

// Re-export main stages used by pipeline
pub(crate) use preparation::PreparationStage;
pub(crate) use request_building::RequestBuildingStage;
pub(crate) use response_processing::ResponseProcessingStage;
// Re-export chat+generate dispatcher stages used by new_regular() / new_pd()
Comment thread
CatherineSue marked this conversation as resolved.
pub(crate) use preparation::ChatGeneratePreparationStage;
pub(crate) use request_building::ChatGenerateRequestBuildingStage;
pub(crate) use response_processing::ChatGenerateResponseProcessingStage;
30 changes: 12 additions & 18 deletions model_gateway/src/routers/grpc/regular/stages/preparation.rs
Original file line number Diff line number Diff line change
@@ -1,16 +1,13 @@
//! Preparation stage that delegates to endpoint-specific implementations
//! Preparation stage for the chat + generate pipeline
//!
//! This stage checks RequestType at runtime and delegates to the appropriate
//! endpoint-specific stage. (ChatPreparationStage, CompletionPreparationStage or GeneratePreparationStage).
//! Dispatches to ChatPreparationStage or GeneratePreparationStage based on
//! request type. Only used by new_regular() and new_pd() pipelines.

use async_trait::async_trait;
use axum::response::Response;
use tracing::error;

use super::{
chat::ChatPreparationStage, completion::CompletionPreparationStage,
generate::GeneratePreparationStage,
};
use super::{chat::ChatPreparationStage, generate::GeneratePreparationStage};
use crate::routers::{
error as grpc_error,
grpc::{
Expand All @@ -19,41 +16,38 @@ use crate::routers::{
},
};

/// Preparation stage (delegates to endpoint-specific implementations)
pub(crate) struct PreparationStage {
/// Preparation stage for chat + generate pipelines
pub(crate) struct ChatGeneratePreparationStage {
chat_stage: ChatPreparationStage,
generate_stage: GeneratePreparationStage,
completion_stage: CompletionPreparationStage,
}

impl PreparationStage {
impl ChatGeneratePreparationStage {
pub fn new() -> Self {
Self {
chat_stage: ChatPreparationStage,
generate_stage: GeneratePreparationStage,
completion_stage: CompletionPreparationStage,
}
}
}

impl Default for PreparationStage {
impl Default for ChatGeneratePreparationStage {
fn default() -> Self {
Self::new()
}
}

#[async_trait]
impl PipelineStage for PreparationStage {
impl PipelineStage for ChatGeneratePreparationStage {
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
match &ctx.input.request_type {
RequestType::Chat(_) => self.chat_stage.execute(ctx).await,
RequestType::Generate(_) => self.generate_stage.execute(ctx).await,
RequestType::Completion(_) => self.completion_stage.execute(ctx).await,
request_type => {
error!(
function = "PreparationStage::execute",
function = "ChatGeneratePreparationStage::execute",
request_type = %request_type,
"{request_type} request type reached regular preparation stage"
"{request_type} should not reach this stage"
);
Err(grpc_error::internal_error(
"wrong_pipeline",
Expand All @@ -64,6 +58,6 @@ impl PipelineStage for PreparationStage {
}

fn name(&self) -> &'static str {
"Preparation"
"ChatGeneratePreparation"
}
}
33 changes: 14 additions & 19 deletions model_gateway/src/routers/grpc/regular/stages/request_building.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,10 @@
//! Request building stage that delegates to endpoint-specific implementations
//! Request building stage for chat and generate endpoints

use async_trait::async_trait;
use axum::response::Response;
use tracing::error;

use super::{
chat::ChatRequestBuildingStage, embedding::request_building::EmbeddingRequestBuildingStage,
generate::GenerateRequestBuildingStage,
};
use super::{chat::ChatRequestBuildingStage, generate::GenerateRequestBuildingStage};
use crate::routers::{
error as grpc_error,
grpc::{
Expand All @@ -16,38 +13,36 @@ use crate::routers::{
},
};

/// Request building stage (delegates to endpoint-specific implementations)
pub(crate) struct RequestBuildingStage {
/// Request building stage for chat and generate pipelines
///
/// These two request types share a single pipeline instance (`new_regular` /
/// `new_pd`) and are dispatched here. All other request types have
/// dedicated pipelines and wire their own request building stages directly.
Comment thread
coderabbitai[bot] marked this conversation as resolved.
pub(crate) struct ChatGenerateRequestBuildingStage {
chat_stage: ChatRequestBuildingStage,
generate_stage: GenerateRequestBuildingStage,
embedding_stage: EmbeddingRequestBuildingStage,
}

impl RequestBuildingStage {
impl ChatGenerateRequestBuildingStage {
pub fn new(inject_pd_metadata: bool) -> Self {
Self {
chat_stage: ChatRequestBuildingStage::new(inject_pd_metadata),
generate_stage: GenerateRequestBuildingStage::new(inject_pd_metadata),
embedding_stage: EmbeddingRequestBuildingStage::new(),
}
}
}

#[async_trait]
impl PipelineStage for RequestBuildingStage {
impl PipelineStage for ChatGenerateRequestBuildingStage {
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
match &ctx.input.request_type {
RequestType::Chat(_) => self.chat_stage.execute(ctx).await,
RequestType::Generate(_) => self.generate_stage.execute(ctx).await,
RequestType::Embedding(_) => self.embedding_stage.execute(ctx).await,
RequestType::Classify(_) => self.embedding_stage.execute(ctx).await,
request_type @ (RequestType::Completion(_)
| RequestType::Responses(_)
| RequestType::Messages(_)) => {
request_type => {
error!(
function = "RequestBuildingStage::execute",
function = "ChatGenerateRequestBuildingStage::execute",
request_type = %request_type,
"{request_type} request type reached regular request building stage"
"{request_type} should not reach this stage"
);
Err(grpc_error::internal_error(
"wrong_pipeline",
Expand All @@ -58,6 +53,6 @@ impl PipelineStage for RequestBuildingStage {
}

fn name(&self) -> &'static str {
"RequestBuilding"
"ChatGenerateRequestBuilding"
}
Comment on lines 55 to 57

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick | 🔵 Trivial

Stage name change affects log labels — update any log-based alerts if needed.

The stage name returned by name() is used throughout pipeline.rs for error and debug logging (as shown in the relevant code snippets). Any log aggregation, dashboards, or alerts filtering on the old "RequestBuilding" label will need to be updated to "ChatGenerateRequestBuilding".

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@model_gateway/src/routers/grpc/regular/stages/request_building.rs` around
lines 55 - 57, The stage name returned by fn name(&self) -> &'static str in the
ChatGenerateRequestBuilding stage was changed to "ChatGenerateRequestBuilding",
so update any log-based filters, dashboards, and alert rules that previously
matched the old "RequestBuilding" label to now match
"ChatGenerateRequestBuilding"; verify any usages in pipeline.rs and related
alerting/aggregation configurations that rely on stage labels and adjust them
accordingly to avoid missing logs or false alerts.

}
Original file line number Diff line number Diff line change
@@ -1,16 +1,15 @@
//! Response processing stage that delegates to endpoint-specific implementations
//! Response processing stage for the chat + generate pipeline
//!
//! Dispatches to ChatResponseProcessingStage or GenerateResponseProcessingStage
//! based on request type. Only used by new_regular() and new_pd() pipelines.

use std::sync::Arc;

use async_trait::async_trait;
use axum::response::Response;
use tracing::error;

use super::{
chat::ChatResponseProcessingStage, classify::ClassifyResponseProcessingStage,
embedding::response_processing::EmbeddingResponseProcessingStage,
generate::GenerateResponseProcessingStage,
};
use super::{chat::ChatResponseProcessingStage, generate::GenerateResponseProcessingStage};
use crate::routers::{
error,
grpc::{
Expand All @@ -20,15 +19,13 @@ use crate::routers::{
},
};

/// Response processing stage (delegates to endpoint-specific implementations)
pub(crate) struct ResponseProcessingStage {
/// Response processing stage for chat + generate pipelines
pub(crate) struct ChatGenerateResponseProcessingStage {
chat_stage: ChatResponseProcessingStage,
generate_stage: GenerateResponseProcessingStage,
embedding_stage: EmbeddingResponseProcessingStage,
classify_stage: ClassifyResponseProcessingStage,
}

impl ResponseProcessingStage {
impl ChatGenerateResponseProcessingStage {
pub fn new(
processor: processor::ResponseProcessor,
streaming_processor: Arc<streaming::StreamingProcessor>,
Expand All @@ -39,27 +36,21 @@ impl ResponseProcessingStage {
streaming_processor.clone(),
),
generate_stage: GenerateResponseProcessingStage::new(processor, streaming_processor),
embedding_stage: EmbeddingResponseProcessingStage::new(),
classify_stage: ClassifyResponseProcessingStage::new(),
}
}
}

#[async_trait]
impl PipelineStage for ResponseProcessingStage {
impl PipelineStage for ChatGenerateResponseProcessingStage {
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
match &ctx.input.request_type {
RequestType::Chat(_) => self.chat_stage.execute(ctx).await,
RequestType::Generate(_) => self.generate_stage.execute(ctx).await,
RequestType::Embedding(_) => self.embedding_stage.execute(ctx).await,
RequestType::Classify(_) => self.classify_stage.execute(ctx).await,
request_type @ (RequestType::Completion(_)
| RequestType::Responses(_)
| RequestType::Messages(_)) => {
request_type => {
error!(
function = "ResponseProcessingStage::execute",
function = "ChatGenerateResponseProcessingStage::execute",
request_type = %request_type,
"{request_type} request type reached regular response processing stage"
"{request_type} should not reach this stage"
);
Err(error::internal_error(
"wrong_pipeline",
Expand All @@ -70,6 +61,6 @@ impl PipelineStage for ResponseProcessingStage {
}

fn name(&self) -> &'static str {
"ResponseProcessing"
"ChatGenerateResponseProcessing"
}
}
Loading