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
12 changes: 0 additions & 12 deletions crates/protocols/src/responses.rs
Original file line number Diff line number Diff line change
Expand Up @@ -573,18 +573,6 @@ impl ResponseUsage {
}
}

#[derive(Debug, Clone, Default, Deserialize, Serialize, schemars::JsonSchema)]
pub struct ResponsesGetParams {
#[serde(default)]
pub include: Vec<String>,
#[serde(default)]
pub include_obfuscation: Option<bool>,
#[serde(default)]
pub starting_after: Option<i64>,
#[serde(default)]
pub stream: Option<bool>,
}

impl ResponsesUsage {
pub fn to_response_usage(&self) -> ResponseUsage {
match self {
Expand Down
23 changes: 1 addition & 22 deletions model_gateway/src/routers/grpc/common/responses/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,33 +2,12 @@
//!
//! These handlers are used by both pipelines for retrieving and cancelling responses.

use axum::response::{IntoResponse, Response};
use axum::response::Response;
use smg_data_connector::ResponseId;

use super::ResponsesContext;
use crate::routers::error;

/// Implementation for GET /v1/responses/{response_id}
///
/// Retrieves a stored response from the database.
/// Used by both regular and harmony implementations.
pub(crate) async fn get_response_impl(ctx: &ResponsesContext, response_id: &str) -> Response {
let resp_id = ResponseId::from(response_id);

// Retrieve response from storage
match ctx.response_storage.get_response(&resp_id).await {
Ok(Some(stored_response)) => axum::Json(stored_response.raw_response).into_response(),
Ok(None) => error::not_found(
"response_not_found",
format!("Response with id '{response_id}' not found"),
),
Err(e) => error::internal_error(
"retrieve_response_failed",
format!("Failed to retrieve response: {e}"),
),
}
}

/// Implementation for POST /v1/responses/{response_id}/cancel
///
/// Background mode is no longer supported, so this endpoint always returns
Expand Down
21 changes: 3 additions & 18 deletions model_gateway/src/routers/grpc/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,20 +6,14 @@ use axum::{
response::{IntoResponse, Response},
};
use openai_protocol::{
chat::ChatCompletionRequest,
classify::ClassifyRequest,
embedding::EmbeddingRequest,
generate::GenerateRequest,
messages::CreateMessageRequest,
responses::{ResponsesGetParams, ResponsesRequest},
chat::ChatCompletionRequest, classify::ClassifyRequest, embedding::EmbeddingRequest,
generate::GenerateRequest, messages::CreateMessageRequest, responses::ResponsesRequest,
};
use tracing::debug;

use super::{
common::responses::{
handlers::{cancel_response_impl, get_response_impl},
utils::validate_worker_availability,
ResponsesContext,
handlers::cancel_response_impl, utils::validate_worker_availability, ResponsesContext,
},
context::SharedComponents,
harmony::{serve_harmony_responses, serve_harmony_responses_stream, HarmonyDetector},
Expand Down Expand Up @@ -482,15 +476,6 @@ impl RouterTrait for GrpcRouter {
self.route_responses_impl(headers, body, model_id).await
}

async fn get_response(
&self,
_headers: Option<&HeaderMap>,
response_id: &str,
_params: &ResponsesGetParams,
) -> Response {
get_response_impl(&self.responses_context, response_id).await
}

async fn cancel_response(&self, _headers: Option<&HeaderMap>, response_id: &str) -> Response {
cancel_response_impl(&self.responses_context, response_id).await
}
Expand Down
18 changes: 1 addition & 17 deletions model_gateway/src/routers/http/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ use openai_protocol::{
embedding::EmbeddingRequest,
generate::GenerateRequest,
rerank::{RerankRequest, RerankResponse, RerankResult},
responses::{ResponsesGetParams, ResponsesRequest},
responses::ResponsesRequest,
};
use reqwest::Client;
use tokio::sync::mpsc;
Expand Down Expand Up @@ -450,12 +450,6 @@ impl Router {
.unwrap_or_else(|| error::bad_gateway("no_worker_response", "No worker response"))
}

// Route a GET request with provided headers to a specific endpoint
async fn route_get_request(&self, headers: Option<&HeaderMap>, endpoint: &str) -> Response {
self.route_simple_request(headers, endpoint, Method::GET)
.await
}

// Route a POST request with empty body to a specific endpoint
async fn route_post_empty_request(
&self,
Expand Down Expand Up @@ -732,16 +726,6 @@ impl RouterTrait for Router {
.await
}

async fn get_response(
&self,
headers: Option<&HeaderMap>,
response_id: &str,
_params: &ResponsesGetParams,
) -> Response {
let endpoint = format!("v1/responses/{response_id}");
self.route_get_request(headers, &endpoint).await
}

async fn cancel_response(&self, headers: Option<&HeaderMap>, response_id: &str) -> Response {
let endpoint = format!("v1/responses/{response_id}/cancel");
self.route_post_empty_request(headers, &endpoint).await
Expand Down
35 changes: 2 additions & 33 deletions model_gateway/src/routers/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ use openai_protocol::{
RealtimeTranscriptionSessionCreateRequest,
},
rerank::RerankRequest,
responses::{ResponsesGetParams, ResponsesRequest},
responses::ResponsesRequest,
};

pub mod anthropic;
Expand All @@ -38,6 +38,7 @@ pub mod mesh;
pub mod openai;
pub mod parse;
pub mod persistence_utils;
pub mod responses;
pub mod router_manager;
pub mod tokenize;
pub mod worker_selection;
Expand Down Expand Up @@ -139,16 +140,6 @@ pub trait RouterTrait: Send + Sync + Debug {
.into_response()
}

/// Retrieve a stored/background response by id
async fn get_response(
&self,
_headers: Option<&HeaderMap>,
_response_id: &str,
_params: &ResponsesGetParams,
) -> Response {
(StatusCode::NOT_IMPLEMENTED, "Get response not implemented").into_response()
}

/// Cancel a background response by id
async fn cancel_response(&self, _headers: Option<&HeaderMap>, _response_id: &str) -> Response {
(
Expand All @@ -158,28 +149,6 @@ pub trait RouterTrait: Send + Sync + Debug {
.into_response()
}

/// Delete a response by id
async fn delete_response(&self, _headers: Option<&HeaderMap>, _response_id: &str) -> Response {
(
StatusCode::NOT_IMPLEMENTED,
"Responses delete endpoint not implemented",
)
.into_response()
}

/// List input items of a response by id
async fn list_response_input_items(
&self,
_headers: Option<&HeaderMap>,
_response_id: &str,
) -> Response {
(
StatusCode::NOT_IMPLEMENTED,
"Responses list input items endpoint not implemented",
)
.into_response()
}

/// Route embedding requests (OpenAI-compatible /v1/embeddings)
async fn route_embeddings(
&self,
Expand Down
81 changes: 1 addition & 80 deletions model_gateway/src/routers/openai/responses/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
//! - Response accumulation for persistence
//! - Tool call detection and output index remapping
//! - Input history loading from conversations and response chains
//! - Storage query handlers for response retrieval
//! - Shared helpers for response retrieval-related logic

mod accumulator;
mod common;
Expand All @@ -19,85 +19,6 @@ mod utils;

// Re-exported for openai::mcp::tool_handler (cross-module dependency)
pub(crate) use accumulator::StreamingResponseAccumulator;
// --- Storage query handlers (extracted from router.rs) ---
use axum::{
http::StatusCode,
response::{IntoResponse, Response},
Json,
};
pub(crate) use common::{extract_output_index, get_event_type};
pub use non_streaming::handle_non_streaming_response;
use openai_protocol::responses::generate_id;
use serde_json::{json, Value};
use smg_data_connector::ResponseId;
pub use streaming::handle_streaming_response;
use tracing::warn;

use super::context::ResponsesComponents;
use crate::routers::error;

/// Fetch a single stored response by ID.
pub(crate) async fn get_response(components: &ResponsesComponents, response_id: &str) -> Response {
let id = ResponseId::from(response_id);
match components.response_storage.get_response(&id).await {
Ok(Some(stored)) => {
let mut response_json = stored.raw_response;
if let Some(obj) = response_json.as_object_mut() {
obj.insert("id".to_string(), json!(id.0));
}
(StatusCode::OK, Json(response_json)).into_response()
}
Ok(None) => error::not_found(
"not_found",
format!("No response found with id '{response_id}'"),
),
Err(e) => error::internal_error("storage_error", format!("Failed to get response: {e}")),
}
}

/// List input items for a stored response.
pub(crate) async fn list_response_input_items(
components: &ResponsesComponents,
response_id: &str,
) -> Response {
let resp_id = ResponseId::from(response_id);

match components.response_storage.get_response(&resp_id).await {
Ok(Some(stored)) => {
let items = stored.input.as_array().cloned().unwrap_or_default();

let items_with_ids: Vec<Value> = items
.into_iter()
.map(|mut item| {
if item.get("id").is_none() {
if let Some(obj) = item.as_object_mut() {
obj.insert("id".to_string(), json!(generate_id("msg")));
}
}
item
})
.collect();

let response_body = json!({
"object": "list",
"data": items_with_ids,
"first_id": items_with_ids.first().and_then(|v| v.get("id").and_then(|i| i.as_str())),
"last_id": items_with_ids.last().and_then(|v| v.get("id").and_then(|i| i.as_str())),
"has_more": false
});

(StatusCode::OK, Json(response_body)).into_response()
}
Ok(None) => error::not_found(
"not_found",
format!("No response found with id '{response_id}'"),
),
Err(e) => {
warn!("Failed to retrieve input items for {}: {}", response_id, e);
error::internal_error(
"storage_error",
format!("Failed to retrieve input items: {e}"),
)
}
}
}
19 changes: 1 addition & 18 deletions model_gateway/src/routers/openai/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use openai_protocol::{
RealtimeClientSecretCreateRequest, RealtimeSessionCreateRequest,
RealtimeTranscriptionSessionCreateRequest,
},
responses::{ResponsesGetParams, ResponsesRequest},
responses::ResponsesRequest,
};

use super::{
Expand Down Expand Up @@ -173,23 +173,6 @@ impl crate::routers::RouterTrait for OpenAIRouter {
responses_route::route_responses(&deps, headers, body, model_id).await
}

async fn get_response(
&self,
_headers: Option<&HeaderMap>,
response_id: &str,
_params: &ResponsesGetParams,
) -> Response {
super::responses::get_response(&self.responses_components, response_id).await
}

async fn list_response_input_items(
&self,
_headers: Option<&HeaderMap>,
response_id: &str,
) -> Response {
super::responses::list_response_input_items(&self.responses_components, response_id).await
}

async fn route_realtime_session(
&self,
headers: Option<&HeaderMap>,
Expand Down
Loading
Loading