Skip to main content

agentic_server/handler/http/
responses.rs

1use axum::extract::{Request, State};
2use axum::http::request::Parts;
3use axum::response::{IntoResponse, Response};
4use bytes::Bytes;
5use either::Either;
6use tracing::debug;
7
8use std::sync::Arc;
9
10use agentic_core::executor::{ExecuteRequest, compact_response as execute_compaction};
11use agentic_core::proxy::{ProxyRequest, proxy_request};
12use agentic_core::types::request_response::{CompactRequest, RequestPayload};
13use agentic_core::types::tools::ResponsesTool;
14
15use super::super::common::{
16    convert_response, executor_error_response, extract_bearer, read_and_parse, read_json, sse_response,
17};
18use crate::app::AppState;
19
20async fn proxy_responses(state: &AppState, parts: Parts, body: Bytes) -> Response {
21    let proxy_req = ProxyRequest {
22        headers: parts.headers,
23        body,
24        query: parts.uri.query().map(str::to_string),
25    };
26    convert_response(proxy_request(proxy_req, &state.proxy_state).await)
27}
28
29async fn execute_responses(state: &AppState, parts: Parts, payload: RequestPayload) -> Response {
30    let auth = extract_bearer(&parts.headers, state.openai_api_key.as_deref());
31    match ExecuteRequest::new(payload, Arc::clone(&state.exec_ctx))
32        .with_auth(auth)
33        .run()
34        .await
35    {
36        Ok(Either::Left(response_payload)) => axum::Json(response_payload).into_response(),
37        Ok(Either::Right(stream)) => sse_response(stream),
38        Err(e) => executor_error_response(e),
39    }
40}
41
42fn has_gateway_tools(payload: &RequestPayload) -> bool {
43    payload
44        .tools
45        .as_ref()
46        .is_some_and(|tools| tools.iter().any(|tool| !matches!(tool, ResponsesTool::Function(_))))
47}
48
49pub async fn responses(State(state): State<AppState>, req: Request) -> Response {
50    let (parts, body) = req.into_parts();
51    let (bytes, payload) = match read_and_parse(body).await {
52        Ok(v) => v,
53        Err(e) => return e,
54    };
55
56    let should_execute = payload.store
57        || payload.previous_response_id.is_some()
58        || payload.conversation_id.is_some()
59        || payload.input.contains_compaction()
60        || payload.input.has_compaction_trigger()
61        || payload
62            .context_management
63            .as_ref()
64            .is_some_and(|entries| !entries.is_empty())
65        || has_gateway_tools(&payload);
66    debug!(
67        route = if should_execute { "executor" } else { "proxy" },
68        store = payload.store,
69        stream = payload.stream,
70        has_previous_response_id = payload.previous_response_id.is_some(),
71        has_conversation_id = payload.conversation_id.is_some(),
72        has_compaction = payload.input.contains_compaction(),
73        has_compaction_trigger = payload.input.has_compaction_trigger(),
74        context_management = payload.context_management.as_ref().map_or(0, Vec::len),
75        tools = payload.tools.as_ref().map_or(0, Vec::len),
76        "routing HTTP responses request"
77    );
78
79    if should_execute {
80        execute_responses(&state, parts, payload).await
81    } else {
82        proxy_responses(&state, parts, bytes).await
83    }
84}
85
86pub async fn compact_response(State(state): State<AppState>, req: Request) -> Response {
87    let (parts, body) = req.into_parts();
88    let request: CompactRequest = match read_json(body).await {
89        Ok(request) => request,
90        Err(response) => return response,
91    };
92    let auth = extract_bearer(&parts.headers, state.openai_api_key.as_deref());
93    match execute_compaction(request, state.exec_ctx.as_ref(), auth.as_deref()).await {
94        Ok(response) => axum::Json(response).into_response(),
95        Err(error) => executor_error_response(error),
96    }
97}