[router][grpc] Move all error logs to their call sites (#12859)
This commit is contained in:
@@ -7,7 +7,7 @@ use axum::{
|
||||
response::{IntoResponse, Response},
|
||||
};
|
||||
use serde_json::{json, to_value};
|
||||
use tracing::{debug, warn};
|
||||
use tracing::{debug, error, warn};
|
||||
|
||||
use crate::{
|
||||
core::WorkerRegistry,
|
||||
@@ -44,6 +44,10 @@ pub async fn ensure_mcp_connection(
|
||||
.await
|
||||
.is_none()
|
||||
{
|
||||
error!(
|
||||
function = "ensure_mcp_connection",
|
||||
"Failed to connect to MCP server"
|
||||
);
|
||||
return Err(error::failed_dependency(
|
||||
"Failed to connect to MCP server. Check server_url and authorization.",
|
||||
));
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use super::PipelineStage;
|
||||
use crate::routers::grpc::{
|
||||
@@ -15,11 +16,13 @@ pub struct ClientAcquisitionStage;
|
||||
#[async_trait]
|
||||
impl PipelineStage for ClientAcquisitionStage {
|
||||
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
|
||||
let workers = ctx
|
||||
.state
|
||||
.workers
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Worker selection not completed"))?;
|
||||
let workers = ctx.state.workers.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "ClientAcquisitionStage::execute",
|
||||
"Worker selection stage not completed"
|
||||
);
|
||||
error::internal_error("Worker selection not completed")
|
||||
})?;
|
||||
|
||||
let clients = match workers {
|
||||
WorkerSelection::Single { worker } => {
|
||||
|
||||
@@ -4,6 +4,7 @@ use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use super::PipelineStage;
|
||||
use crate::routers::grpc::{
|
||||
@@ -17,11 +18,13 @@ pub struct DispatchMetadataStage;
|
||||
#[async_trait]
|
||||
impl PipelineStage for DispatchMetadataStage {
|
||||
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
|
||||
let proto_request = ctx
|
||||
.state
|
||||
.proto_request
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Proto request not built"))?;
|
||||
let proto_request = ctx.state.proto_request.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "DispatchMetadataStage::execute",
|
||||
"Proto request not built"
|
||||
);
|
||||
error::internal_error("Proto request not built")
|
||||
})?;
|
||||
|
||||
let request_id = proto_request.request_id.clone();
|
||||
let model = match &ctx.input.request_type {
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use super::PipelineStage;
|
||||
use crate::{
|
||||
@@ -35,17 +36,21 @@ impl RequestExecutionStage {
|
||||
#[async_trait]
|
||||
impl PipelineStage for RequestExecutionStage {
|
||||
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
|
||||
let proto_request = ctx
|
||||
.state
|
||||
.proto_request
|
||||
.take()
|
||||
.ok_or_else(|| error::internal_error("Proto request not built"))?;
|
||||
let proto_request = ctx.state.proto_request.take().ok_or_else(|| {
|
||||
error!(
|
||||
function = "RequestExecutionStage::execute",
|
||||
"Proto request not built"
|
||||
);
|
||||
error::internal_error("Proto request not built")
|
||||
})?;
|
||||
|
||||
let clients = ctx
|
||||
.state
|
||||
.clients
|
||||
.as_mut()
|
||||
.ok_or_else(|| error::internal_error("Client acquisition not completed"))?;
|
||||
let clients = ctx.state.clients.as_mut().ok_or_else(|| {
|
||||
error!(
|
||||
function = "RequestExecutionStage::execute",
|
||||
"Client acquisition not completed"
|
||||
);
|
||||
error::internal_error("Client acquisition not completed")
|
||||
})?;
|
||||
|
||||
let result = match self.mode {
|
||||
ExecutionMode::Single => self.execute_single(proto_request, clients).await?,
|
||||
@@ -70,14 +75,22 @@ impl RequestExecutionStage {
|
||||
proto_request: proto::GenerateRequest,
|
||||
clients: &mut ClientSelection,
|
||||
) -> Result<ExecutionResult, Response> {
|
||||
let client = clients
|
||||
.single_mut()
|
||||
.ok_or_else(|| error::internal_error("Expected single client but got dual"))?;
|
||||
let client = clients.single_mut().ok_or_else(|| {
|
||||
error!(
|
||||
function = "execute_single",
|
||||
"Expected single client but got dual"
|
||||
);
|
||||
error::internal_error("Expected single client but got dual")
|
||||
})?;
|
||||
|
||||
let stream = client
|
||||
.generate(proto_request)
|
||||
.await
|
||||
.map_err(|e| error::internal_error(format!("Failed to start generation: {}", e)))?;
|
||||
let stream = client.generate(proto_request).await.map_err(|e| {
|
||||
error!(
|
||||
function = "execute_single",
|
||||
error = %e,
|
||||
"Failed to start generation"
|
||||
);
|
||||
error::internal_error(format!("Failed to start generation: {}", e))
|
||||
})?;
|
||||
|
||||
Ok(ExecutionResult::Single { stream })
|
||||
}
|
||||
@@ -87,9 +100,13 @@ impl RequestExecutionStage {
|
||||
proto_request: proto::GenerateRequest,
|
||||
clients: &mut ClientSelection,
|
||||
) -> Result<ExecutionResult, Response> {
|
||||
let (prefill_client, decode_client) = clients
|
||||
.dual_mut()
|
||||
.ok_or_else(|| error::internal_error("Expected dual clients but got single"))?;
|
||||
let (prefill_client, decode_client) = clients.dual_mut().ok_or_else(|| {
|
||||
error!(
|
||||
function = "execute_dual_dispatch",
|
||||
"Expected dual clients but got single"
|
||||
);
|
||||
error::internal_error("Expected dual clients but got single")
|
||||
})?;
|
||||
|
||||
let prefill_request = proto_request.clone();
|
||||
let decode_request = proto_request;
|
||||
@@ -103,6 +120,11 @@ impl RequestExecutionStage {
|
||||
let prefill_stream = match prefill_result {
|
||||
Ok(s) => s,
|
||||
Err(e) => {
|
||||
error!(
|
||||
function = "execute_dual_dispatch",
|
||||
error = %e,
|
||||
"Prefill worker failed to start"
|
||||
);
|
||||
return Err(error::internal_error(format!(
|
||||
"Prefill worker failed to start: {}",
|
||||
e
|
||||
@@ -114,6 +136,11 @@ impl RequestExecutionStage {
|
||||
let decode_stream = match decode_result {
|
||||
Ok(s) => s,
|
||||
Err(e) => {
|
||||
error!(
|
||||
function = "execute_dual_dispatch",
|
||||
error = %e,
|
||||
"Decode worker failed to start"
|
||||
);
|
||||
return Err(error::internal_error(format!(
|
||||
"Decode worker failed to start: {}",
|
||||
e
|
||||
|
||||
@@ -4,7 +4,7 @@ use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::warn;
|
||||
use tracing::{error, warn};
|
||||
|
||||
use super::PipelineStage;
|
||||
use crate::{
|
||||
@@ -47,11 +47,13 @@ impl WorkerSelectionStage {
|
||||
#[async_trait]
|
||||
impl PipelineStage for WorkerSelectionStage {
|
||||
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
|
||||
let prep = ctx
|
||||
.state
|
||||
.preparation
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Preparation stage not completed"))?;
|
||||
let prep = ctx.state.preparation.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "WorkerSelectionStage::execute",
|
||||
"Preparation stage not completed"
|
||||
);
|
||||
error::internal_error("Preparation stage not completed")
|
||||
})?;
|
||||
|
||||
// For Harmony, use selection_text produced during Harmony encoding
|
||||
// Otherwise, use original_text from regular preparation
|
||||
@@ -66,6 +68,12 @@ impl PipelineStage for WorkerSelectionStage {
|
||||
match self.select_single_worker(ctx.input.model_id.as_deref(), text) {
|
||||
Some(w) => WorkerSelection::Single { worker: w },
|
||||
None => {
|
||||
error!(
|
||||
function = "WorkerSelectionStage::execute",
|
||||
mode = "Regular",
|
||||
model_id = ?ctx.input.model_id,
|
||||
"No available workers for model"
|
||||
);
|
||||
return Err(error::service_unavailable(format!(
|
||||
"No available workers for model: {:?}",
|
||||
ctx.input.model_id
|
||||
@@ -77,6 +85,12 @@ impl PipelineStage for WorkerSelectionStage {
|
||||
match self.select_pd_pair(ctx.input.model_id.as_deref(), text) {
|
||||
Some((prefill, decode)) => WorkerSelection::Dual { prefill, decode },
|
||||
None => {
|
||||
error!(
|
||||
function = "WorkerSelectionStage::execute",
|
||||
mode = "PrefillDecode",
|
||||
model_id = ?ctx.input.model_id,
|
||||
"No available PD worker pairs for model"
|
||||
);
|
||||
return Err(error::service_unavailable(format!(
|
||||
"No available PD worker pairs for model: {:?}",
|
||||
ctx.input.model_id
|
||||
|
||||
@@ -9,7 +9,6 @@ use axum::{
|
||||
Json,
|
||||
};
|
||||
use serde_json::json;
|
||||
use tracing::{error, warn};
|
||||
|
||||
/// Create a 500 Internal Server Error response
|
||||
///
|
||||
@@ -21,7 +20,6 @@ use tracing::{error, warn};
|
||||
/// ```
|
||||
pub fn internal_error(message: impl Into<String>) -> Response {
|
||||
let msg = message.into();
|
||||
error!("{}", msg);
|
||||
(
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(json!({
|
||||
@@ -45,7 +43,6 @@ pub fn internal_error(message: impl Into<String>) -> Response {
|
||||
/// ```
|
||||
pub fn bad_request(message: impl Into<String>) -> Response {
|
||||
let msg = message.into();
|
||||
error!("{}", msg);
|
||||
(
|
||||
StatusCode::BAD_REQUEST,
|
||||
Json(json!({
|
||||
@@ -69,7 +66,6 @@ pub fn bad_request(message: impl Into<String>) -> Response {
|
||||
/// ```
|
||||
pub fn not_found(message: impl Into<String>) -> Response {
|
||||
let msg = message.into();
|
||||
warn!("{}", msg);
|
||||
(
|
||||
StatusCode::NOT_FOUND,
|
||||
Json(json!({
|
||||
@@ -93,7 +89,6 @@ pub fn not_found(message: impl Into<String>) -> Response {
|
||||
/// ```
|
||||
pub fn service_unavailable(message: impl Into<String>) -> Response {
|
||||
let msg = message.into();
|
||||
warn!("{}", msg);
|
||||
(
|
||||
StatusCode::SERVICE_UNAVAILABLE,
|
||||
Json(json!({
|
||||
@@ -117,7 +112,6 @@ pub fn service_unavailable(message: impl Into<String>) -> Response {
|
||||
/// ```
|
||||
pub fn failed_dependency(message: impl Into<String>) -> Response {
|
||||
let msg = message.into();
|
||||
warn!("{}", msg);
|
||||
(
|
||||
StatusCode::FAILED_DEPENDENCY,
|
||||
Json(json!({
|
||||
|
||||
@@ -4,6 +4,7 @@ use std::sync::Arc;
|
||||
|
||||
use axum::response::Response;
|
||||
use proto::generate_complete::MatchedStop::{MatchedStopStr, MatchedTokenId};
|
||||
use tracing::error;
|
||||
|
||||
use super::HarmonyParserAdapter;
|
||||
use crate::{
|
||||
@@ -63,6 +64,11 @@ impl HarmonyResponseProcessor {
|
||||
|
||||
// Parse Harmony channels with HarmonyParserAdapter
|
||||
let mut parser = HarmonyParserAdapter::new().map_err(|e| {
|
||||
error!(
|
||||
function = "process_non_streaming_chat_response",
|
||||
error = %e,
|
||||
"Failed to create Harmony parser"
|
||||
);
|
||||
error::internal_error(format!("Failed to create Harmony parser: {}", e))
|
||||
})?;
|
||||
|
||||
@@ -73,7 +79,14 @@ impl HarmonyResponseProcessor {
|
||||
complete.finish_reason.clone(),
|
||||
matched_stop.clone(),
|
||||
)
|
||||
.map_err(|e| error::internal_error(format!("Harmony parsing failed: {}", e)))?;
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "process_non_streaming_chat_response",
|
||||
error = %e,
|
||||
"Harmony parsing failed on complete response"
|
||||
);
|
||||
error::internal_error(format!("Harmony parsing failed: {}", e))
|
||||
})?;
|
||||
|
||||
// Build response message (assistant)
|
||||
let message = ChatCompletionMessage {
|
||||
@@ -171,6 +184,11 @@ impl HarmonyResponseProcessor {
|
||||
|
||||
// Parse Harmony channels
|
||||
let mut parser = HarmonyParserAdapter::new().map_err(|e| {
|
||||
error!(
|
||||
function = "process_responses_iteration",
|
||||
error = %e,
|
||||
"Failed to create Harmony parser"
|
||||
);
|
||||
error::internal_error(format!("Failed to create Harmony parser: {}", e))
|
||||
})?;
|
||||
|
||||
@@ -190,7 +208,14 @@ impl HarmonyResponseProcessor {
|
||||
complete.finish_reason.clone(),
|
||||
matched_stop,
|
||||
)
|
||||
.map_err(|e| error::internal_error(format!("Harmony parsing failed: {}", e)))?;
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "process_responses_iteration",
|
||||
error = %e,
|
||||
"Harmony parsing failed on complete response"
|
||||
);
|
||||
error::internal_error(format!("Harmony parsing failed: {}", e))
|
||||
})?;
|
||||
|
||||
// VALIDATION: Check if model incorrectly generated Tool role messages
|
||||
// This happens when the model copies the format of tool result messages
|
||||
|
||||
@@ -39,7 +39,7 @@ use axum::response::Response;
|
||||
use bytes::Bytes;
|
||||
use serde_json::{from_str, from_value, json, to_string, to_value, Value};
|
||||
use tokio::sync::mpsc;
|
||||
use tracing::{debug, warn};
|
||||
use tracing::{debug, error, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::{
|
||||
@@ -324,6 +324,12 @@ async fn execute_with_mcp_loop(
|
||||
|
||||
// Safety check: prevent infinite loops
|
||||
if iteration_count > MAX_TOOL_ITERATIONS {
|
||||
error!(
|
||||
function = "execute_with_mcp_loop",
|
||||
iteration_count = iteration_count,
|
||||
max_iterations = MAX_TOOL_ITERATIONS,
|
||||
"Maximum tool iterations exceeded"
|
||||
);
|
||||
return Err(error::internal_error(format!(
|
||||
"Maximum tool iterations ({}) exceeded",
|
||||
MAX_TOOL_ITERATIONS
|
||||
@@ -1157,6 +1163,13 @@ async fn execute_mcp_tools(
|
||||
// Parse tool arguments from JSON string
|
||||
let args_str = tool_call.function.arguments.as_deref().unwrap_or("{}");
|
||||
let args: Value = from_str(args_str).map_err(|e| {
|
||||
error!(
|
||||
function = "execute_mcp_tools",
|
||||
tool_name = %tool_call.function.name,
|
||||
call_id = %tool_call.id,
|
||||
error = %e,
|
||||
"Failed to parse tool arguments JSON"
|
||||
);
|
||||
error::internal_error(format!(
|
||||
"Invalid tool arguments JSON for tool '{}': {}",
|
||||
tool_call.function.name, e
|
||||
@@ -1519,6 +1532,12 @@ async fn load_previous_messages(
|
||||
.get_response_chain(&prev_id, None)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "load_previous_messages",
|
||||
prev_id = %prev_id_str,
|
||||
error = %e,
|
||||
"Failed to load previous response chain from storage"
|
||||
);
|
||||
error::internal_error(format!(
|
||||
"Failed to load previous response chain for {}: {}",
|
||||
prev_id_str, e
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use serde_json::json;
|
||||
use tracing::error;
|
||||
|
||||
use super::super::HarmonyBuilder;
|
||||
use crate::{
|
||||
@@ -56,6 +57,10 @@ impl PipelineStage for HarmonyPreparationStage {
|
||||
let request_arc = ctx.responses_request_arc();
|
||||
self.prepare_responses(ctx, &request_arc).await?;
|
||||
} else {
|
||||
error!(
|
||||
function = "HarmonyPreparationStage::execute",
|
||||
"Unsupported request type for Harmony pipeline"
|
||||
);
|
||||
return Err(error::bad_request(
|
||||
"Only Chat and Responses requests supported in Harmony pipeline".to_string(),
|
||||
));
|
||||
@@ -78,6 +83,10 @@ impl HarmonyPreparationStage {
|
||||
) -> Result<Option<Response>, Response> {
|
||||
// Validate - reject logprobs
|
||||
if request.logprobs {
|
||||
error!(
|
||||
function = "prepare_chat",
|
||||
"logprobs requested but not supported for Harmony models"
|
||||
);
|
||||
return Err(error::bad_request(
|
||||
"logprobs are not supported for Harmony models".to_string(),
|
||||
));
|
||||
@@ -94,10 +103,14 @@ impl HarmonyPreparationStage {
|
||||
};
|
||||
|
||||
// Step 3: Build via Harmony
|
||||
let build_output = self
|
||||
.builder
|
||||
.build_from_chat(&body_ref)
|
||||
.map_err(|e| error::bad_request(format!("Harmony build failed: {}", e)))?;
|
||||
let build_output = self.builder.build_from_chat(&body_ref).map_err(|e| {
|
||||
error!(
|
||||
function = "prepare_chat",
|
||||
error = %e,
|
||||
"Harmony build failed for chat request"
|
||||
);
|
||||
error::bad_request(format!("Harmony build failed: {}", e))
|
||||
})?;
|
||||
|
||||
// Step 4: Store results
|
||||
ctx.state.preparation = Some(PreparationOutput {
|
||||
@@ -154,6 +167,10 @@ impl HarmonyPreparationStage {
|
||||
};
|
||||
|
||||
if tool_constraint.is_some() && text_constraint.is_some() {
|
||||
error!(
|
||||
function = "prepare_responses",
|
||||
"Conflicting constraints: both tool_choice and text format specified"
|
||||
);
|
||||
return Err(error::bad_request(
|
||||
"Cannot use both tool_choice (required/function) and text format (json_object/json_schema) simultaneously".to_string(),
|
||||
));
|
||||
@@ -162,10 +179,14 @@ impl HarmonyPreparationStage {
|
||||
let constraint = tool_constraint.or(text_constraint);
|
||||
|
||||
// Step 3: Build via Harmony from responses API request
|
||||
let build_output = self
|
||||
.builder
|
||||
.build_from_responses(request)
|
||||
.map_err(|e| error::bad_request(format!("Harmony build failed: {}", e)))?;
|
||||
let build_output = self.builder.build_from_responses(request).map_err(|e| {
|
||||
error!(
|
||||
function = "prepare_responses",
|
||||
error = %e,
|
||||
"Harmony build failed for responses request"
|
||||
);
|
||||
error::bad_request(format!("Harmony build failed: {}", e))
|
||||
})?;
|
||||
|
||||
// Step 4: Store results with constraint
|
||||
ctx.state.preparation = Some(PreparationOutput {
|
||||
@@ -200,12 +221,25 @@ impl HarmonyPreparationStage {
|
||||
TextFormat::Text => Ok(None),
|
||||
TextFormat::JsonObject => {
|
||||
let tag = build_text_format_structural_tag(&serde_json::json!({"type": "object"}))
|
||||
.map_err(|e| Box::new(error::internal_error(e)))?;
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "generate_text_format_constraint",
|
||||
error = %e,
|
||||
"Failed to build text format structural tag for JsonObject"
|
||||
);
|
||||
Box::new(error::internal_error(e))
|
||||
})?;
|
||||
Ok(Some(("structural_tag".to_string(), tag)))
|
||||
}
|
||||
TextFormat::JsonSchema { schema, .. } => {
|
||||
let tag = build_text_format_structural_tag(schema)
|
||||
.map_err(|e| Box::new(error::internal_error(e)))?;
|
||||
let tag = build_text_format_structural_tag(schema).map_err(|e| {
|
||||
error!(
|
||||
function = "generate_text_format_constraint",
|
||||
error = %e,
|
||||
"Failed to build text format structural tag for JsonSchema"
|
||||
);
|
||||
Box::new(error::internal_error(e))
|
||||
})?;
|
||||
Ok(Some(("structural_tag".to_string(), tag)))
|
||||
}
|
||||
}
|
||||
@@ -266,11 +300,19 @@ impl HarmonyPreparationStage {
|
||||
};
|
||||
|
||||
// Validate specific function exists
|
||||
if specific_function.is_some() && tools_to_use.is_empty() {
|
||||
return Err(Box::new(error::bad_request(format!(
|
||||
"Tool '{}' not found in tools list",
|
||||
specific_function.unwrap()
|
||||
))));
|
||||
match specific_function {
|
||||
Some(tool_name) if tools_to_use.is_empty() => {
|
||||
error!(
|
||||
function = "generate_tool_call_constraint",
|
||||
tool_name = %tool_name,
|
||||
"Specified tool not found in tools list"
|
||||
);
|
||||
return Err(Box::new(error::bad_request(format!(
|
||||
"Tool '{}' not found in tools list",
|
||||
tool_name
|
||||
))));
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
|
||||
// Build tags for each tool - need two patterns per tool for reasoning on/off
|
||||
@@ -312,6 +354,11 @@ impl HarmonyPreparationStage {
|
||||
});
|
||||
|
||||
serde_json::to_string(&structural_tag).map_err(|e| {
|
||||
error!(
|
||||
function = "generate_tool_call_constraint",
|
||||
error = %e,
|
||||
"Failed to serialize structural tag"
|
||||
);
|
||||
Box::new(error::internal_error(format!(
|
||||
"Failed to serialize structural tag: {}",
|
||||
e
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::debug;
|
||||
use tracing::{debug, error};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::routers::grpc::{
|
||||
@@ -30,18 +30,22 @@ impl HarmonyRequestBuildingStage {
|
||||
impl PipelineStage for HarmonyRequestBuildingStage {
|
||||
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
|
||||
// Get preparation output
|
||||
let prep = ctx
|
||||
.state
|
||||
.preparation
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Preparation not completed"))?;
|
||||
let prep = ctx.state.preparation.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "HarmonyRequestBuildingStage::execute",
|
||||
"Preparation stage not completed"
|
||||
);
|
||||
error::internal_error("Preparation not completed")
|
||||
})?;
|
||||
|
||||
// Get clients
|
||||
let clients = ctx
|
||||
.state
|
||||
.clients
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Client acquisition not completed"))?;
|
||||
let clients = ctx.state.clients.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "HarmonyRequestBuildingStage::execute",
|
||||
"Client acquisition stage not completed"
|
||||
);
|
||||
error::internal_error("Client acquisition not completed")
|
||||
})?;
|
||||
let builder_client = match clients {
|
||||
ClientSelection::Single { client } => client,
|
||||
ClientSelection::Dual { prefill, .. } => prefill,
|
||||
@@ -52,6 +56,10 @@ impl PipelineStage for HarmonyRequestBuildingStage {
|
||||
RequestType::Chat(_) => format!("chatcmpl-{}", Uuid::new_v4()),
|
||||
RequestType::Responses(_) => format!("responses-{}", Uuid::new_v4()),
|
||||
RequestType::Generate(_) => {
|
||||
error!(
|
||||
function = "HarmonyRequestBuildingStage::execute",
|
||||
"Generate request type not supported for Harmony models"
|
||||
);
|
||||
return Err(error::bad_request(
|
||||
"Generate requests are not supported with Harmony models".to_string(),
|
||||
));
|
||||
@@ -75,7 +83,14 @@ impl PipelineStage for HarmonyRequestBuildingStage {
|
||||
None,
|
||||
prep.tool_constraints.clone(),
|
||||
)
|
||||
.map_err(|e| error::bad_request(format!("Invalid request parameters: {}", e)))?
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "HarmonyRequestBuildingStage::execute",
|
||||
error = %e,
|
||||
"Failed to build generate request from chat"
|
||||
);
|
||||
error::bad_request(format!("Invalid request parameters: {}", e))
|
||||
})?
|
||||
}
|
||||
RequestType::Responses(request) => builder_client
|
||||
.build_generate_request_from_responses(
|
||||
@@ -86,7 +101,14 @@ impl PipelineStage for HarmonyRequestBuildingStage {
|
||||
prep.harmony_stop_ids.clone(),
|
||||
prep.tool_constraints.clone(),
|
||||
)
|
||||
.map_err(|e| error::bad_request(format!("Invalid request parameters: {}", e)))?,
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "HarmonyRequestBuildingStage::execute",
|
||||
error = %e,
|
||||
"Failed to build generate request from responses"
|
||||
);
|
||||
error::bad_request(format!("Invalid request parameters: {}", e))
|
||||
})?,
|
||||
_ => unreachable!(),
|
||||
};
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use super::super::{HarmonyResponseProcessor, HarmonyStreamingProcessor};
|
||||
use crate::routers::grpc::{
|
||||
@@ -46,19 +47,24 @@ impl PipelineStage for HarmonyResponseProcessingStage {
|
||||
match &ctx.input.request_type {
|
||||
RequestType::Chat(_) => {
|
||||
// Get execution result (output tokens from model)
|
||||
let execution_result = ctx
|
||||
.state
|
||||
.response
|
||||
.execution_result
|
||||
.take()
|
||||
.ok_or_else(|| error::internal_error("No execution result"))?;
|
||||
let execution_result =
|
||||
ctx.state.response.execution_result.take().ok_or_else(|| {
|
||||
error!(
|
||||
function = "HarmonyResponseProcessingStage::execute",
|
||||
request_type = "Chat",
|
||||
"No execution result available"
|
||||
);
|
||||
error::internal_error("No execution result")
|
||||
})?;
|
||||
|
||||
let dispatch = ctx
|
||||
.state
|
||||
.dispatch
|
||||
.as_ref()
|
||||
.cloned()
|
||||
.ok_or_else(|| error::internal_error("Dispatch metadata not set"))?;
|
||||
let dispatch = ctx.state.dispatch.as_ref().cloned().ok_or_else(|| {
|
||||
error!(
|
||||
function = "HarmonyResponseProcessingStage::execute",
|
||||
request_type = "Chat",
|
||||
"Dispatch metadata not set"
|
||||
);
|
||||
error::internal_error("Dispatch metadata not set")
|
||||
})?;
|
||||
|
||||
// For streaming, delegate to streaming processor and return SSE response
|
||||
if is_streaming {
|
||||
@@ -92,19 +98,24 @@ impl PipelineStage for HarmonyResponseProcessingStage {
|
||||
}
|
||||
|
||||
// For non-streaming, process normally
|
||||
let execution_result = ctx
|
||||
.state
|
||||
.response
|
||||
.execution_result
|
||||
.take()
|
||||
.ok_or_else(|| error::internal_error("No execution result"))?;
|
||||
let execution_result =
|
||||
ctx.state.response.execution_result.take().ok_or_else(|| {
|
||||
error!(
|
||||
function = "HarmonyResponseProcessingStage::execute",
|
||||
request_type = "Responses",
|
||||
"No execution result available"
|
||||
);
|
||||
error::internal_error("No execution result")
|
||||
})?;
|
||||
|
||||
let dispatch = ctx
|
||||
.state
|
||||
.dispatch
|
||||
.as_ref()
|
||||
.cloned()
|
||||
.ok_or_else(|| error::internal_error("Dispatch metadata not set"))?;
|
||||
let dispatch = ctx.state.dispatch.as_ref().cloned().ok_or_else(|| {
|
||||
error!(
|
||||
function = "HarmonyResponseProcessingStage::execute",
|
||||
request_type = "Responses",
|
||||
"Dispatch metadata not set"
|
||||
);
|
||||
error::internal_error("Dispatch metadata not set")
|
||||
})?;
|
||||
|
||||
let responses_request = ctx.responses_request_arc();
|
||||
let iteration_result = self
|
||||
@@ -115,9 +126,15 @@ impl PipelineStage for HarmonyResponseProcessingStage {
|
||||
ctx.state.response.responses_iteration_result = Some(iteration_result);
|
||||
Ok(None)
|
||||
}
|
||||
RequestType::Generate(_) => Err(error::internal_error(
|
||||
"Generate requests not supported in Harmony pipeline",
|
||||
)),
|
||||
RequestType::Generate(_) => {
|
||||
error!(
|
||||
function = "HarmonyResponseProcessingStage::execute",
|
||||
"Generate request type not supported in Harmony pipeline"
|
||||
);
|
||||
Err(error::internal_error(
|
||||
"Generate requests not supported in Harmony pipeline",
|
||||
))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -219,9 +219,19 @@ impl RequestPipeline {
|
||||
match ctx.state.response.final_response {
|
||||
Some(FinalResponse::Chat(response)) => axum::Json(response).into_response(),
|
||||
Some(FinalResponse::Generate(_)) => {
|
||||
error!(
|
||||
function = "execute_chat",
|
||||
"Wrong response type: expected Chat, got Generate"
|
||||
);
|
||||
error::internal_error("Internal error: wrong response type")
|
||||
}
|
||||
None => error::internal_error("No response produced"),
|
||||
None => {
|
||||
error!(
|
||||
function = "execute_chat",
|
||||
"No response produced by pipeline"
|
||||
);
|
||||
error::internal_error("No response produced")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -260,9 +270,19 @@ impl RequestPipeline {
|
||||
match ctx.state.response.final_response {
|
||||
Some(FinalResponse::Generate(response)) => axum::Json(response).into_response(),
|
||||
Some(FinalResponse::Chat(_)) => {
|
||||
error!(
|
||||
function = "execute_generate",
|
||||
"Wrong response type: expected Generate, got Chat"
|
||||
);
|
||||
error::internal_error("Internal error: wrong response type")
|
||||
}
|
||||
None => error::internal_error("No response produced"),
|
||||
None => {
|
||||
error!(
|
||||
function = "execute_generate",
|
||||
"No response produced by pipeline"
|
||||
);
|
||||
error::internal_error("No response produced")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -285,6 +305,10 @@ impl RequestPipeline {
|
||||
match stage.execute(&mut ctx).await {
|
||||
Ok(Some(_response)) => {
|
||||
// Streaming not supported for responses sync mode
|
||||
error!(
|
||||
function = "execute_chat_for_responses",
|
||||
"Streaming attempted in responses context"
|
||||
);
|
||||
return Err(error::bad_request(
|
||||
"Streaming is not supported in this context".to_string(),
|
||||
));
|
||||
@@ -308,9 +332,19 @@ impl RequestPipeline {
|
||||
match ctx.state.response.final_response {
|
||||
Some(FinalResponse::Chat(response)) => Ok(response),
|
||||
Some(FinalResponse::Generate(_)) => {
|
||||
error!(
|
||||
function = "execute_chat_for_responses",
|
||||
"Wrong response type: expected Chat, got Generate"
|
||||
);
|
||||
Err(error::internal_error("Internal error: wrong response type"))
|
||||
}
|
||||
None => Err(error::internal_error("No response produced")),
|
||||
None => {
|
||||
error!(
|
||||
function = "execute_chat_for_responses",
|
||||
"No response produced by pipeline"
|
||||
);
|
||||
Err(error::internal_error("No response produced"))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -376,6 +410,10 @@ impl RequestPipeline {
|
||||
.responses_iteration_result
|
||||
.take()
|
||||
.ok_or_else(|| {
|
||||
error!(
|
||||
function = "execute_harmony_responses",
|
||||
"No ResponsesIterationResult produced by pipeline"
|
||||
);
|
||||
error::internal_error("No ResponsesIterationResult produced by pipeline")
|
||||
})
|
||||
}
|
||||
@@ -421,10 +459,12 @@ impl RequestPipeline {
|
||||
}
|
||||
|
||||
// Extract execution_result (the raw stream from workers)
|
||||
ctx.state
|
||||
.response
|
||||
.execution_result
|
||||
.take()
|
||||
.ok_or_else(|| error::internal_error("No ExecutionResult produced by pipeline"))
|
||||
ctx.state.response.execution_result.take().ok_or_else(|| {
|
||||
error!(
|
||||
function = "execute_harmony_responses_streaming",
|
||||
"No ExecutionResult produced by pipeline"
|
||||
);
|
||||
error::internal_error("No ExecutionResult produced by pipeline")
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -43,7 +43,7 @@ use bytes::Bytes;
|
||||
use futures_util::StreamExt;
|
||||
use serde_json::json;
|
||||
use tokio::sync::mpsc;
|
||||
use tracing::{debug, warn};
|
||||
use tracing::{debug, error, warn};
|
||||
use uuid::Uuid;
|
||||
use validator::Validate;
|
||||
|
||||
@@ -644,8 +644,14 @@ async fn execute_without_mcp(
|
||||
response_id: Option<String>,
|
||||
) -> Result<ResponsesResponse, Response> {
|
||||
// Convert ResponsesRequest → ChatCompletionRequest
|
||||
let chat_request = conversions::responses_to_chat(modified_request)
|
||||
.map_err(|e| error::bad_request(format!("Failed to convert request: {}", e)))?;
|
||||
let chat_request = conversions::responses_to_chat(modified_request).map_err(|e| {
|
||||
error!(
|
||||
function = "execute_without_mcp",
|
||||
error = %e,
|
||||
"Failed to convert ResponsesRequest to ChatCompletionRequest"
|
||||
);
|
||||
error::bad_request(format!("Failed to convert request: {}", e))
|
||||
})?;
|
||||
|
||||
// Execute chat pipeline (errors already have proper HTTP status codes)
|
||||
let chat_response = ctx
|
||||
@@ -659,8 +665,14 @@ async fn execute_without_mcp(
|
||||
.await?; // Preserve the Response error as-is
|
||||
|
||||
// Convert ChatCompletionResponse → ResponsesResponse
|
||||
conversions::chat_to_responses(&chat_response, original_request, response_id)
|
||||
.map_err(|e| error::internal_error(format!("Failed to convert to responses format: {}", e)))
|
||||
conversions::chat_to_responses(&chat_response, original_request, response_id).map_err(|e| {
|
||||
error!(
|
||||
function = "execute_without_mcp",
|
||||
error = %e,
|
||||
"Failed to convert ChatCompletionResponse to ResponsesResponse"
|
||||
);
|
||||
error::internal_error(format!("Failed to convert to responses format: {}", e))
|
||||
})
|
||||
}
|
||||
|
||||
/// Load conversation history and response chains, returning modified request
|
||||
@@ -737,7 +749,15 @@ async fn load_conversation_history(
|
||||
.conversation_storage
|
||||
.get_conversation(&conv_id)
|
||||
.await
|
||||
.map_err(|e| error::internal_error(format!("Failed to check conversation: {}", e)))?;
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "load_conversation_history",
|
||||
conversation_id = %conv_id_str,
|
||||
error = %e,
|
||||
"Failed to check conversation existence in storage"
|
||||
);
|
||||
error::internal_error(format!("Failed to check conversation: {}", e))
|
||||
})?;
|
||||
|
||||
if conversation.is_none() {
|
||||
return Err(error::not_found(format!(
|
||||
|
||||
@@ -16,7 +16,7 @@ use futures_util::StreamExt;
|
||||
use serde_json::{json, Value};
|
||||
use tokio::sync::mpsc;
|
||||
use tokio_stream::wrappers::UnboundedReceiverStream;
|
||||
use tracing::{debug, warn};
|
||||
use tracing::{debug, error, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::conversions;
|
||||
@@ -250,8 +250,15 @@ pub(super) async fn execute_tool_loop(
|
||||
|
||||
loop {
|
||||
// Convert to chat request
|
||||
let mut chat_request = conversions::responses_to_chat(¤t_request)
|
||||
.map_err(|e| error::bad_request(format!("Failed to convert request: {}", e)))?;
|
||||
let mut chat_request = conversions::responses_to_chat(¤t_request).map_err(|e| {
|
||||
error!(
|
||||
function = "tool_loop",
|
||||
iteration = state.iteration,
|
||||
error = %e,
|
||||
"Failed to convert ResponsesRequest to ChatCompletionRequest in tool loop"
|
||||
);
|
||||
error::bad_request(format!("Failed to convert request: {}", e))
|
||||
})?;
|
||||
|
||||
// Prepare tools and tool_choice for this iteration
|
||||
prepare_chat_tools_and_choice(&mut chat_request, &mcp_chat_tools, state.iteration);
|
||||
@@ -301,6 +308,13 @@ pub(super) async fn execute_tool_loop(
|
||||
response_id.clone(),
|
||||
)
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "tool_loop",
|
||||
iteration = state.iteration,
|
||||
error = %e,
|
||||
context = "function_tool_calls",
|
||||
"Failed to convert ChatCompletionResponse to ResponsesResponse"
|
||||
);
|
||||
error::internal_error(format!("Failed to convert to responses format: {}", e))
|
||||
})?;
|
||||
|
||||
@@ -331,6 +345,13 @@ pub(super) async fn execute_tool_loop(
|
||||
response_id.clone(),
|
||||
)
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "tool_loop",
|
||||
iteration = state.iteration,
|
||||
error = %e,
|
||||
context = "max_tool_calls_limit",
|
||||
"Failed to convert ChatCompletionResponse to ResponsesResponse"
|
||||
);
|
||||
error::internal_error(format!("Failed to convert to responses format: {}", e))
|
||||
})?;
|
||||
|
||||
@@ -453,6 +474,13 @@ pub(super) async fn execute_tool_loop(
|
||||
response_id.clone(),
|
||||
)
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
function = "tool_loop",
|
||||
iteration = state.iteration,
|
||||
error = %e,
|
||||
context = "final_response",
|
||||
"Failed to convert ChatCompletionResponse to ResponsesResponse"
|
||||
);
|
||||
error::internal_error(format!("Failed to convert to responses format: {}", e))
|
||||
})?;
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ use std::borrow::Cow;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use crate::{
|
||||
protocols::chat::ChatCompletionRequest,
|
||||
@@ -43,18 +44,22 @@ impl ChatPreparationStage {
|
||||
let body_ref = utils::filter_chat_request_by_tool_choice(request);
|
||||
|
||||
// Step 2: Process messages and apply chat template
|
||||
let processed_messages =
|
||||
match utils::process_chat_messages(&body_ref, &*ctx.components.tokenizer) {
|
||||
Ok(msgs) => msgs,
|
||||
Err(e) => {
|
||||
return Err(error::bad_request(e));
|
||||
}
|
||||
};
|
||||
let processed_messages = match utils::process_chat_messages(
|
||||
&body_ref,
|
||||
&*ctx.components.tokenizer,
|
||||
) {
|
||||
Ok(msgs) => msgs,
|
||||
Err(e) => {
|
||||
error!(function = "ChatPreparationStage::execute", error = %e, "Failed to process chat messages");
|
||||
return Err(error::bad_request(e));
|
||||
}
|
||||
};
|
||||
|
||||
// Step 3: Tokenize the processed text
|
||||
let encoding = match ctx.components.tokenizer.encode(&processed_messages.text) {
|
||||
Ok(encoding) => encoding,
|
||||
Err(e) => {
|
||||
error!(function = "ChatPreparationStage::execute", error = %e, "Tokenization failed");
|
||||
return Err(error::internal_error(format!("Tokenization failed: {}", e)));
|
||||
}
|
||||
};
|
||||
@@ -64,7 +69,10 @@ impl ChatPreparationStage {
|
||||
// Step 4: Build tool constraints if needed
|
||||
let tool_call_constraint = if let Some(tools) = body_ref.tools.as_ref() {
|
||||
utils::generate_tool_constraints(tools, &request.tool_choice, &request.model)
|
||||
.map_err(|e| error::bad_request(format!("Invalid tool configuration: {}", e)))?
|
||||
.map_err(|e| {
|
||||
error!(function = "ChatPreparationStage::execute", error = %e, "Invalid tool configuration");
|
||||
error::bad_request(format!("Invalid tool configuration: {}", e))
|
||||
})?
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::routers::grpc::{
|
||||
@@ -26,17 +27,21 @@ impl ChatRequestBuildingStage {
|
||||
#[async_trait]
|
||||
impl PipelineStage for ChatRequestBuildingStage {
|
||||
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
|
||||
let prep = ctx
|
||||
.state
|
||||
.preparation
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Preparation not completed"))?;
|
||||
let prep = ctx.state.preparation.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "ChatRequestBuildingStage::execute",
|
||||
"Preparation not completed"
|
||||
);
|
||||
error::internal_error("Preparation not completed")
|
||||
})?;
|
||||
|
||||
let clients = ctx
|
||||
.state
|
||||
.clients
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Client acquisition not completed"))?;
|
||||
let clients = ctx.state.clients.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "ChatRequestBuildingStage::execute",
|
||||
"Client acquisition not completed"
|
||||
);
|
||||
error::internal_error("Client acquisition not completed")
|
||||
})?;
|
||||
|
||||
let chat_request = ctx.chat_request_arc();
|
||||
|
||||
@@ -63,7 +68,10 @@ impl PipelineStage for ChatRequestBuildingStage {
|
||||
.clone(),
|
||||
prep.tool_constraints.clone(),
|
||||
)
|
||||
.map_err(|e| error::bad_request(format!("Invalid request parameters: {}", e)))?;
|
||||
.map_err(|e| {
|
||||
error!(function = "ChatRequestBuildingStage::execute", error = %e, "Failed to build generate request");
|
||||
error::bad_request(format!("Invalid request parameters: {}", e))
|
||||
})?;
|
||||
|
||||
// Inject PD metadata if needed
|
||||
if self.inject_pd_metadata {
|
||||
|
||||
@@ -7,6 +7,7 @@ use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use crate::routers::grpc::{
|
||||
common::stages::PipelineStage,
|
||||
@@ -54,19 +55,26 @@ impl ChatResponseProcessingStage {
|
||||
let is_streaming = ctx.is_streaming();
|
||||
|
||||
// Extract execution result
|
||||
let execution_result = ctx
|
||||
.state
|
||||
.response
|
||||
.execution_result
|
||||
.take()
|
||||
.ok_or_else(|| error::internal_error("No execution result"))?;
|
||||
let execution_result = ctx.state.response.execution_result.take().ok_or_else(|| {
|
||||
error!(
|
||||
function = "ChatResponseProcessingStage::execute",
|
||||
"No execution result"
|
||||
);
|
||||
error::internal_error("No execution result")
|
||||
})?;
|
||||
|
||||
// Get dispatch metadata (needed by both streaming and non-streaming)
|
||||
let dispatch = ctx
|
||||
.state
|
||||
.dispatch
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Dispatch metadata not set"))?
|
||||
.ok_or_else(|| {
|
||||
error!(
|
||||
function = "ChatResponseProcessingStage::execute",
|
||||
"Dispatch metadata not set"
|
||||
);
|
||||
error::internal_error("Dispatch metadata not set")
|
||||
})?
|
||||
.clone();
|
||||
|
||||
if is_streaming {
|
||||
@@ -85,12 +93,13 @@ impl ChatResponseProcessingStage {
|
||||
|
||||
let chat_request = ctx.chat_request_arc();
|
||||
|
||||
let stop_decoder = ctx
|
||||
.state
|
||||
.response
|
||||
.stop_decoder
|
||||
.as_mut()
|
||||
.ok_or_else(|| error::internal_error("Stop decoder not initialized"))?;
|
||||
let stop_decoder = ctx.state.response.stop_decoder.as_mut().ok_or_else(|| {
|
||||
error!(
|
||||
function = "ChatResponseProcessingStage::execute",
|
||||
"Stop decoder not initialized"
|
||||
);
|
||||
error::internal_error("Stop decoder not initialized")
|
||||
})?;
|
||||
|
||||
let response = self
|
||||
.processor
|
||||
|
||||
@@ -4,6 +4,7 @@ use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use crate::{
|
||||
protocols::{common::InputIds, generate::GenerateRequest},
|
||||
@@ -44,6 +45,7 @@ impl GeneratePreparationStage {
|
||||
let (original_text, token_ids) = match self.resolve_generate_input(ctx, request) {
|
||||
Ok(res) => res,
|
||||
Err(msg) => {
|
||||
error!(function = "GeneratePreparationStage::execute", error = %msg, "Failed to resolve generate input");
|
||||
return Err(error::bad_request(msg));
|
||||
}
|
||||
};
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::routers::grpc::{
|
||||
@@ -26,17 +27,21 @@ impl GenerateRequestBuildingStage {
|
||||
#[async_trait]
|
||||
impl PipelineStage for GenerateRequestBuildingStage {
|
||||
async fn execute(&self, ctx: &mut RequestContext) -> Result<Option<Response>, Response> {
|
||||
let prep = ctx
|
||||
.state
|
||||
.preparation
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Preparation not completed"))?;
|
||||
let prep = ctx.state.preparation.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "GenerateRequestBuildingStage::execute",
|
||||
"Preparation not completed"
|
||||
);
|
||||
error::internal_error("Preparation not completed")
|
||||
})?;
|
||||
|
||||
let clients = ctx
|
||||
.state
|
||||
.clients
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Client acquisition not completed"))?;
|
||||
let clients = ctx.state.clients.as_ref().ok_or_else(|| {
|
||||
error!(
|
||||
function = "GenerateRequestBuildingStage::execute",
|
||||
"Client acquisition not completed"
|
||||
);
|
||||
error::internal_error("Client acquisition not completed")
|
||||
})?;
|
||||
|
||||
let generate_request = ctx.generate_request_arc();
|
||||
|
||||
@@ -59,7 +64,10 @@ impl PipelineStage for GenerateRequestBuildingStage {
|
||||
prep.original_text.clone(),
|
||||
prep.token_ids.clone(),
|
||||
)
|
||||
.map_err(error::bad_request)?;
|
||||
.map_err(|e| {
|
||||
error!(function = "GenerateRequestBuildingStage::execute", error = %e, "Failed to build generate request");
|
||||
error::bad_request(e)
|
||||
})?;
|
||||
|
||||
// Inject PD metadata if needed
|
||||
if self.inject_pd_metadata {
|
||||
|
||||
@@ -4,6 +4,7 @@ use std::{sync::Arc, time::Instant};
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use crate::routers::grpc::{
|
||||
common::stages::PipelineStage,
|
||||
@@ -52,19 +53,26 @@ impl GenerateResponseProcessingStage {
|
||||
let is_streaming = ctx.is_streaming();
|
||||
|
||||
// Extract execution result
|
||||
let execution_result = ctx
|
||||
.state
|
||||
.response
|
||||
.execution_result
|
||||
.take()
|
||||
.ok_or_else(|| error::internal_error("No execution result"))?;
|
||||
let execution_result = ctx.state.response.execution_result.take().ok_or_else(|| {
|
||||
error!(
|
||||
function = "GenerateResponseProcessingStage::execute",
|
||||
"No execution result"
|
||||
);
|
||||
error::internal_error("No execution result")
|
||||
})?;
|
||||
|
||||
// Get dispatch metadata (needed by both streaming and non-streaming)
|
||||
let dispatch = ctx
|
||||
.state
|
||||
.dispatch
|
||||
.as_ref()
|
||||
.ok_or_else(|| error::internal_error("Dispatch metadata not set"))?
|
||||
.ok_or_else(|| {
|
||||
error!(
|
||||
function = "GenerateResponseProcessingStage::execute",
|
||||
"Dispatch metadata not set"
|
||||
);
|
||||
error::internal_error("Dispatch metadata not set")
|
||||
})?
|
||||
.clone();
|
||||
|
||||
if is_streaming {
|
||||
@@ -82,12 +90,13 @@ impl GenerateResponseProcessingStage {
|
||||
let request_logprobs = ctx.generate_request().return_logprob.unwrap_or(false);
|
||||
let generate_request = ctx.generate_request_arc();
|
||||
|
||||
let stop_decoder = ctx
|
||||
.state
|
||||
.response
|
||||
.stop_decoder
|
||||
.as_mut()
|
||||
.ok_or_else(|| error::internal_error("Stop decoder not initialized"))?;
|
||||
let stop_decoder = ctx.state.response.stop_decoder.as_mut().ok_or_else(|| {
|
||||
error!(
|
||||
function = "GenerateResponseProcessingStage::execute",
|
||||
"Stop decoder not initialized"
|
||||
);
|
||||
error::internal_error("Stop decoder not initialized")
|
||||
})?;
|
||||
|
||||
let result_array = self
|
||||
.processor
|
||||
|
||||
@@ -4,6 +4,7 @@ use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use axum::response::Response;
|
||||
use tracing::error;
|
||||
|
||||
use super::{chat::ChatResponseProcessingStage, generate::GenerateResponseProcessingStage};
|
||||
use crate::routers::grpc::{
|
||||
@@ -40,9 +41,15 @@ impl PipelineStage for ResponseProcessingStage {
|
||||
match &ctx.input.request_type {
|
||||
RequestType::Chat(_) => self.chat_stage.execute(ctx).await,
|
||||
RequestType::Generate(_) => self.generate_stage.execute(ctx).await,
|
||||
RequestType::Responses(_) => Err(error::bad_request(
|
||||
"Responses API processing must be handled by responses handler".to_string(),
|
||||
)),
|
||||
RequestType::Responses(_) => {
|
||||
error!(
|
||||
function = "ResponseProcessingStage::execute",
|
||||
"Responses API not supported in regular pipeline"
|
||||
);
|
||||
Err(error::bad_request(
|
||||
"Responses API processing must be handled by responses handler".to_string(),
|
||||
))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -43,8 +43,14 @@ pub async fn get_grpc_client_from_worker(
|
||||
let client_arc = worker
|
||||
.get_grpc_client()
|
||||
.await
|
||||
.map_err(|e| error::internal_error(format!("Failed to get gRPC client: {}", e)))?
|
||||
.ok_or_else(|| error::internal_error("Selected worker is not configured for gRPC"))?;
|
||||
.map_err(|e| {
|
||||
error!(function = "get_grpc_client_from_worker", error = %e, "Failed to get gRPC client");
|
||||
error::internal_error(format!("Failed to get gRPC client: {}", e))
|
||||
})?
|
||||
.ok_or_else(|| {
|
||||
error!(function = "get_grpc_client_from_worker", "Selected worker not configured for gRPC");
|
||||
error::internal_error("Selected worker is not configured for gRPC")
|
||||
})?;
|
||||
|
||||
Ok((*client_arc).clone())
|
||||
}
|
||||
@@ -596,7 +602,7 @@ pub async fn collect_stream_responses(
|
||||
all_responses.push(complete);
|
||||
}
|
||||
Some(Error(err)) => {
|
||||
error!("{} error: {}", worker_name, err.message);
|
||||
error!(function = "collect_stream_responses", worker = %worker_name, error = %err.message, "Worker generation error");
|
||||
// Don't mark as completed - let Drop send abort for error cases
|
||||
return Err(error::internal_error(format!(
|
||||
"{} generation failed: {}",
|
||||
@@ -612,7 +618,7 @@ pub async fn collect_stream_responses(
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!("{} stream error: {:?}", worker_name, e);
|
||||
error!(function = "collect_stream_responses", worker = %worker_name, error = ?e, "Worker stream error");
|
||||
// Don't mark as completed - let Drop send abort for error cases
|
||||
return Err(error::internal_error(format!(
|
||||
"{} stream failed: {}",
|
||||
|
||||
Reference in New Issue
Block a user