agentic_server/handler/http/
responses.rs1use 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}