1use crate::{
2 anthropic::json_error,
3 logging::{Logger, REDACT_KEYS, create_logger},
4 monitor::{EndpointKind, MonitorHandle},
5 openai_compat::{
6 MAX_OPENAI_REQUEST_BYTES, OpenAiError, OpenAiSurface,
7 request::{extract_model, parse_request},
8 stream::openai_response as render_openai_response,
9 },
10 project,
11 provider::RequestContext,
12 providers::codex::{
13 chat_completions::{ChatCompletionsBackend, request::translate_request},
14 images::{
15 CodexImagesBackend, ImageOperation, ImageRequestError, MAX_EDIT_REQUEST_BYTES,
16 MAX_GENERATION_REQUEST_BYTES, MultipartEditInput, UploadedImage, image_error_response,
17 prepare_json_request, prepare_multipart_edit,
18 },
19 native::{
20 CodexNativeBackend, NativeResponseOutcome, openai_error, validate_native_request_model,
21 },
22 transcription::{
23 CodexTranscriptionBackend, MAX_TRANSCRIPTION_REQUEST_BYTES, TranscriptionRequestError,
24 prepare_transcription, transcription_error_response,
25 },
26 },
27 registry::{Registry, normalize_incoming_model},
28 request_identity::{CLAUDE_AGENT_HEADER, CLAUDE_PARENT_AGENT_HEADER, ConversationIdentity},
29 session::{self, SessionState},
30 traffic::{TrafficCaptureOptions, create_traffic_capture},
31};
32use axum::{
33 Json, Router,
34 body::Body,
35 extract::{DefaultBodyLimit, FromRequest, Multipart, Query, State},
36 http::{Request, StatusCode},
37 response::Response,
38 routing::{get, post},
39};
40use http_body_util::{BodyExt, StreamBody};
41use serde::de::DeserializeOwned;
42use serde_json::{Map, Value, json};
43use std::fs::{self, File};
44use std::future::Future;
45use std::io::Write;
46use std::path::{Path, PathBuf};
47use std::sync::Arc;
48use std::time::Instant;
49use tokio::net::TcpListener;
50use uuid::Uuid;
51
52const CLAUDE_AUTO_REVIEW_SYSTEM_PREFIX: &str =
53 "You are a security monitor for autonomous AI coding agents.";
54const CODEX_AUTO_REVIEW_MODEL: &str = "gpt-5.6-luna";
55
56#[derive(Debug, Clone, PartialEq, Eq)]
57struct AutoReviewRoute {
58 requested_model: String,
59 override_model: String,
60}
61
62fn is_claude_auto_review_request(body: &crate::anthropic::schema::MessagesRequest) -> bool {
63 if body.stream {
64 return false;
65 }
66
67 let has_tools = body
68 .extra
69 .get("tools")
70 .and_then(Value::as_array)
71 .is_some_and(|tools| !tools.is_empty());
72 if has_tools {
73 return false;
74 }
75
76 body.extra
77 .get("system")
78 .and_then(Value::as_array)
79 .is_some_and(|blocks| {
80 blocks.iter().any(|block| {
81 block
82 .get("text")
83 .and_then(Value::as_str)
84 .is_some_and(|text| text.starts_with(CLAUDE_AUTO_REVIEW_SYSTEM_PREFIX))
85 })
86 })
87}
88
89fn apply_auto_review_model(
90 body: &mut crate::anthropic::schema::MessagesRequest,
91 count_tokens: bool,
92 configured_model: Option<&str>,
93 original_provider: &str,
94) -> Option<AutoReviewRoute> {
95 if count_tokens || !is_claude_auto_review_request(body) {
96 return None;
97 }
98
99 let override_model = configured_model
100 .filter(|model| !model.is_empty())
101 .or((original_provider == "codex").then_some(CODEX_AUTO_REVIEW_MODEL))?;
102 let route = AutoReviewRoute {
103 requested_model: body.model.clone()?,
104 override_model: override_model.to_string(),
105 };
106 body.model = Some(route.override_model.clone());
107 Some(route)
108}
109
110pub struct ServerConfig {
111 pub bind_address: String,
112 pub port: u16,
113 pub monitor: Option<MonitorHandle>,
114}
115
116pub async fn serve(config: ServerConfig) -> anyhow::Result<()> {
117 serve_inner(config, std::future::pending::<()>()).await
118}
119
120pub async fn serve_with_shutdown(
121 config: ServerConfig,
122 shutdown: impl Future<Output = ()> + Send + 'static,
123) -> anyhow::Result<()> {
124 serve_inner(config, shutdown).await
125}
126
127async fn serve_inner(
128 config: ServerConfig,
129 shutdown: impl Future<Output = ()> + Send + 'static,
130) -> anyhow::Result<()> {
131 let listener = bind_proxy_listener(&config.bind_address, config.port).await?;
132 serve_listener(listener, config.monitor, shutdown).await
133}
134
135pub async fn bind_proxy_listener(bind_address: &str, port: u16) -> anyhow::Result<TcpListener> {
136 let ip = bind_address
137 .parse::<std::net::IpAddr>()
138 .map_err(|err| anyhow::anyhow!("invalid proxy bind address {bind_address:?}: {err}"))?;
139 let addr = std::net::SocketAddr::new(ip, port);
140 TcpListener::bind(addr)
141 .await
142 .map_err(|err| anyhow::anyhow!("failed to bind proxy listener on {addr}: {err}"))
143}
144
145pub async fn serve_listener(
146 listener: TcpListener,
147 monitor: Option<MonitorHandle>,
148 shutdown: impl Future<Output = ()> + Send + 'static,
149) -> anyhow::Result<()> {
150 let local_addr = listener.local_addr()?;
151 let port = local_addr.port();
152 create_logger("server").info(
153 "server listening",
154 Some(serde_json::Map::from_iter([
155 ("port".to_string(), json!(port)),
156 (
157 "bindAddress".to_string(),
158 json!(local_addr.ip().to_string()),
159 ),
160 (
161 "logDir".to_string(),
162 json!(
163 crate::paths::log_file()
164 .parent()
165 .map(|path| path.display().to_string())
166 ),
167 ),
168 ])),
169 );
170 let app = app_with_monitor(Arc::new(Registry::with_default_alias()), monitor);
171 axum::serve(listener, app)
172 .with_graceful_shutdown(shutdown)
173 .await?;
174 Ok(())
175}
176
177pub fn app(registry: Arc<Registry>) -> Router {
178 app_with_features(
179 registry,
180 None,
181 AppFeatures {
182 responses_api: crate::config::codex_responses_api(),
183 images_api: crate::config::codex_images_api(),
184 transcriptions_api: crate::config::codex_transcriptions_api(),
185 },
186 )
187}
188
189pub fn app_with_monitor(registry: Arc<Registry>, monitor: Option<MonitorHandle>) -> Router {
190 app_with_features(
191 registry,
192 monitor,
193 AppFeatures {
194 responses_api: crate::config::codex_responses_api(),
195 images_api: crate::config::codex_images_api(),
196 transcriptions_api: crate::config::codex_transcriptions_api(),
197 },
198 )
199}
200
201#[derive(Debug, Clone, Copy, Default)]
202pub struct AppFeatures {
203 pub responses_api: bool,
204 pub images_api: bool,
205 pub transcriptions_api: bool,
206}
207
208pub fn app_with_options(
209 registry: Arc<Registry>,
210 monitor: Option<MonitorHandle>,
211 responses_api: bool,
212) -> Router {
213 app_with_features(
214 registry,
215 monitor,
216 AppFeatures {
217 responses_api,
218 images_api: false,
219 transcriptions_api: false,
220 },
221 )
222}
223
224pub fn app_with_features(
225 registry: Arc<Registry>,
226 monitor: Option<MonitorHandle>,
227 features: AppFeatures,
228) -> Router {
229 let native_responses = features
230 .responses_api
231 .then(|| Arc::new(CodexNativeBackend::new()));
232 let chat_completions = features
233 .responses_api
234 .then(|| Arc::new(ChatCompletionsBackend::new()));
235 let images = if features.images_api {
236 match CodexImagesBackend::new() {
237 Ok(backend) => Some(Arc::new(backend)),
238 Err(error) => {
239 create_logger("server").warn(
240 "codex images backend disabled by invalid configuration",
241 Some(Map::from_iter([("error".to_string(), json!(error))])),
242 );
243 None
244 }
245 }
246 } else {
247 None
248 };
249 let transcriptions = features
250 .transcriptions_api
251 .then(|| Arc::new(CodexTranscriptionBackend::new()));
252 let state = Arc::new(AppState {
253 registry,
254 monitor,
255 native_responses,
256 chat_completions,
257 images,
258 transcriptions,
259 });
260 let router = Router::new()
261 .route("/healthz", get(healthz))
262 .route("/v1/messages", post(handler_messages))
263 .route("/v1/messages/count_tokens", post(handler_count_tokens))
264 .route("/v1/models", get(handler_models));
265 let router = if features.responses_api {
266 router
267 .route("/v1/responses", post(handler_responses))
268 .route("/v1/chat/completions", post(handler_chat_completions))
269 } else {
270 router
271 };
272 let router = if features.images_api {
273 router
274 .route("/v1/images/generations", post(handler_image_generation))
275 .route(
276 "/v1/images/edits",
277 post(handler_image_edit).layer(DefaultBodyLimit::max(MAX_EDIT_REQUEST_BYTES)),
278 )
279 } else {
280 router
281 };
282 let router = if features.transcriptions_api {
283 router.route(
284 "/v1/audio/transcriptions",
285 post(handler_transcription)
286 .layer(DefaultBodyLimit::max(MAX_TRANSCRIPTION_REQUEST_BYTES)),
287 )
288 } else {
289 router
290 };
291 router.fallback(fallback_handler).with_state(state)
292}
293
294#[derive(Clone)]
295struct AppState {
296 registry: Arc<Registry>,
297 monitor: Option<MonitorHandle>,
298 native_responses: Option<Arc<CodexNativeBackend>>,
299 chat_completions: Option<Arc<ChatCompletionsBackend>>,
300 images: Option<Arc<CodexImagesBackend>>,
301 transcriptions: Option<Arc<CodexTranscriptionBackend>>,
302}
303
304async fn healthz() -> Json<serde_json::Value> {
305 Json(json!({ "ok": true }))
306}
307
308#[derive(serde::Deserialize)]
309struct ModelsQuery {
310 limit: Option<usize>,
311}
312
313async fn handler_models(
319 State(state): State<Arc<AppState>>,
320 Query(query): Query<ModelsQuery>,
321) -> Json<serde_json::Value> {
322 let mut data: Vec<Value> = state
323 .registry
324 .all_supported_models()
325 .into_iter()
326 .map(|(model, provider)| {
327 json!({
328 "type": "model",
329 "object": "model",
330 "id": model,
331 "display_name": format!("{model} ({provider})"),
332 })
333 })
334 .collect();
335 let has_more = query.limit.is_some_and(|limit| data.len() > limit);
336 if let Some(limit) = query.limit {
337 data.truncate(limit);
338 }
339 Json(json!({
340 "object": "list",
341 "data": data,
342 "has_more": has_more,
343 "first_id": data.first().and_then(|entry| entry.get("id")).cloned(),
344 "last_id": data.last().and_then(|entry| entry.get("id")).cloned(),
345 }))
346}
347
348async fn handler_messages(State(state): State<Arc<AppState>>, req: Request<Body>) -> Response {
349 dispatch_request(state, req, false).await
350}
351
352async fn handler_count_tokens(State(state): State<Arc<AppState>>, req: Request<Body>) -> Response {
353 dispatch_request(state, req, true).await
354}
355
356async fn handler_transcription(State(state): State<Arc<AppState>>, req: Request<Body>) -> Response {
357 let started_at = Instant::now();
358 let log = create_logger("server");
359 let req_id = Uuid::new_v4().to_string();
360 let headers = req.headers().clone();
361 log.info(
362 "request",
363 Some(Map::from_iter([
364 ("reqId".to_string(), json!(&req_id)),
365 ("method".to_string(), json!("POST")),
366 ("path".to_string(), json!("/v1/audio/transcriptions")),
367 ("query".to_string(), json!({})),
368 ])),
369 );
370 let session_id = native_session_id(&headers);
371 if let Some(monitor) = state.monitor.as_ref() {
372 monitor.request_started(
373 &req_id,
374 session_id.clone(),
375 None,
376 EndpointKind::Transcriptions,
377 );
378 monitor.provider_selected(&req_id, "codex", "codex-transcribe", None);
379 }
380 let request_guard = RequestMonitorGuard::new(state.monitor.clone(), req_id.clone());
381 if req.uri().query().is_some() {
382 let response = transcription_error_response(TranscriptionRequestError::invalid(
383 "Transcription endpoint does not accept query parameters",
384 None,
385 ));
386 log_native_request_completed(
387 &log,
388 &req_id,
389 "transcriptions",
390 Some("codex-transcribe"),
391 response.status(),
392 started_at,
393 );
394 return monitor_response_body(response, request_guard);
395 }
396 let content_type = headers
397 .get(http::header::CONTENT_TYPE)
398 .and_then(|value| value.to_str().ok())
399 .unwrap_or_default()
400 .to_ascii_lowercase();
401 if !content_type.starts_with("multipart/form-data") {
402 let response = transcription_error_response(TranscriptionRequestError {
403 status: StatusCode::UNSUPPORTED_MEDIA_TYPE,
404 message: "Transcription request must use multipart/form-data".to_string(),
405 param: None,
406 code: "unsupported_media_type",
407 });
408 log_native_request_completed(
409 &log,
410 &req_id,
411 "transcriptions",
412 Some("codex-transcribe"),
413 response.status(),
414 started_at,
415 );
416 return monitor_response_body(response, request_guard);
417 }
418 let mut multipart = match Multipart::from_request(req, &()).await {
419 Ok(multipart) => multipart,
420 Err(error) => {
421 let status = axum::response::IntoResponse::into_response(error).status();
422 let response = transcription_error_response(TranscriptionRequestError {
423 status: if status == StatusCode::PAYLOAD_TOO_LARGE {
424 StatusCode::PAYLOAD_TOO_LARGE
425 } else {
426 StatusCode::BAD_REQUEST
427 },
428 message: if status == StatusCode::PAYLOAD_TOO_LARGE {
429 "Transcription request exceeded the size limit".to_string()
430 } else {
431 "Invalid multipart transcription request".to_string()
432 },
433 param: None,
434 code: if status == StatusCode::PAYLOAD_TOO_LARGE {
435 "request_too_large"
436 } else {
437 "invalid_multipart"
438 },
439 });
440 log_native_request_completed(
441 &log,
442 &req_id,
443 "transcriptions",
444 Some("codex-transcribe"),
445 response.status(),
446 started_at,
447 );
448 return monitor_response_body(response, request_guard);
449 }
450 };
451 let mut audio = None;
452 let mut filename = None;
453 let mut audio_content_type = None;
454 let mut language = None;
455 while let Some(field) = match multipart.next_field().await {
456 Ok(field) => field,
457 Err(_) => {
458 let response = transcription_error_response(TranscriptionRequestError::invalid(
459 "Invalid multipart transcription request",
460 None,
461 ));
462 return monitor_response_body(response, request_guard);
463 }
464 } {
465 let name = field.name().unwrap_or_default().to_string();
466 match name.as_str() {
467 "file" => {
468 if audio.is_some() {
469 let response =
470 transcription_error_response(TranscriptionRequestError::invalid(
471 "Transcription request must contain one 'file' field",
472 Some("file"),
473 ));
474 return monitor_response_body(response, request_guard);
475 }
476 filename = field.file_name().map(str::to_string);
477 audio_content_type = field.content_type().map(str::to_string);
478 audio = match field.bytes().await {
479 Ok(bytes) => Some(bytes),
480 Err(_) => {
481 let response =
482 transcription_error_response(TranscriptionRequestError::invalid(
483 "Failed to read uploaded audio",
484 Some("file"),
485 ));
486 return monitor_response_body(response, request_guard);
487 }
488 };
489 }
490 "language" => {
491 language = match multipart_text(field, "language").await {
492 Ok(value) => Some(value),
493 Err(error) => {
494 return monitor_response_body(
495 transcription_error_response(TranscriptionRequestError {
496 status: error.status,
497 message: error.message,
498 param: error.param,
499 code: error.code.unwrap_or("invalid_multipart"),
500 }),
501 request_guard,
502 );
503 }
504 };
505 }
506 "model" => {
507 if multipart_text(field, "model").await.is_err() {
508 let response =
509 transcription_error_response(TranscriptionRequestError::invalid(
510 "Invalid multipart 'model' field",
511 Some("model"),
512 ));
513 return monitor_response_body(response, request_guard);
514 }
515 }
516 _ => {
517 let response = transcription_error_response(TranscriptionRequestError::invalid(
518 format!("Unsupported multipart field '{name}'"),
519 None,
520 ));
521 return monitor_response_body(response, request_guard);
522 }
523 }
524 }
525 let prepared = match prepare_transcription(audio, filename, audio_content_type, language) {
526 Ok(prepared) => prepared,
527 Err(error) => {
528 let response = transcription_error_response(error);
529 log_native_request_completed(
530 &log,
531 &req_id,
532 "transcriptions",
533 Some("codex-transcribe"),
534 response.status(),
535 started_at,
536 );
537 return monitor_response_body(response, request_guard);
538 }
539 };
540 let context = RequestContext {
541 req_id: req_id.clone(),
542 session_id,
543 session_seq: None,
544 provider: "codex".to_string(),
545 traffic: None,
546 monitor: state.monitor.clone(),
547 passthrough: None,
548 };
549 let response = match state.transcriptions.as_ref() {
550 Some(backend) => backend.handle(prepared, context).await,
551 None => transcription_error_response(TranscriptionRequestError {
552 status: StatusCode::SERVICE_UNAVAILABLE,
553 message: "Codex transcription API is unavailable".to_string(),
554 param: None,
555 code: "transcriptions_api_unavailable",
556 }),
557 };
558 log_native_request_completed(
559 &log,
560 &req_id,
561 "transcriptions",
562 Some("codex-transcribe"),
563 response.status(),
564 started_at,
565 );
566 monitor_response_body(response, request_guard)
567}
568
569async fn handler_image_generation(
570 State(state): State<Arc<AppState>>,
571 req: Request<Body>,
572) -> Response {
573 dispatch_image_request(state, req, ImageOperation::Generation).await
574}
575
576async fn handler_image_edit(State(state): State<Arc<AppState>>, req: Request<Body>) -> Response {
577 dispatch_image_request(state, req, ImageOperation::Edit).await
578}
579
580async fn dispatch_image_request(
581 state: Arc<AppState>,
582 req: Request<Body>,
583 operation: ImageOperation,
584) -> Response {
585 let started_at = Instant::now();
586 let log = create_logger("server");
587 let req_id = Uuid::new_v4().to_string();
588 let uri = req.uri().clone();
589 let headers = req.headers().clone();
590 let path = uri.path().to_string();
591 log.info(
592 "request",
593 Some(Map::from_iter([
594 ("reqId".to_string(), json!(&req_id)),
595 ("method".to_string(), json!("POST")),
596 ("path".to_string(), json!(&path)),
597 ("query".to_string(), json!(redacted_query(&uri))),
598 ])),
599 );
600 let session_id = native_session_id(&headers);
601 if let Some(monitor) = state.monitor.as_ref() {
602 monitor.request_started(&req_id, session_id.clone(), None, EndpointKind::Images);
603 }
604 let request_guard = RequestMonitorGuard::new(state.monitor.clone(), req_id.clone());
605
606 if uri.query().is_some() {
607 let response = image_error_response(ImageRequestError {
608 status: StatusCode::BAD_REQUEST,
609 message: "Image endpoints do not accept query parameters".to_string(),
610 param: None,
611 code: Some("invalid_request"),
612 });
613 log_native_request_completed(
614 &log,
615 &req_id,
616 operation.label(),
617 None,
618 response.status(),
619 started_at,
620 );
621 return monitor_response_body(response, request_guard);
622 }
623 let content_type = headers
624 .get(http::header::CONTENT_TYPE)
625 .and_then(|value| value.to_str().ok())
626 .unwrap_or_default()
627 .to_ascii_lowercase();
628 let prepared = if content_type.starts_with("application/json") {
629 let limit = match operation {
630 ImageOperation::Generation => MAX_GENERATION_REQUEST_BYTES,
631 ImageOperation::Edit => MAX_EDIT_REQUEST_BYTES,
632 };
633 let body = match axum::body::to_bytes(req.into_body(), limit).await {
634 Ok(body) => body,
635 Err(_) => {
636 let response = image_error_response(ImageRequestError {
637 status: StatusCode::PAYLOAD_TOO_LARGE,
638 message: "Image request exceeded the size limit".to_string(),
639 param: None,
640 code: Some("request_too_large"),
641 });
642 log_native_request_completed(
643 &log,
644 &req_id,
645 operation.label(),
646 None,
647 response.status(),
648 started_at,
649 );
650 return monitor_response_body(response, request_guard);
651 }
652 };
653 prepare_json_request(operation, &body)
654 } else if operation == ImageOperation::Edit && content_type.starts_with("multipart/form-data") {
655 let multipart = match Multipart::from_request(req, &()).await {
656 Ok(multipart) => multipart,
657 Err(_) => {
658 let response = image_error_response(ImageRequestError {
659 status: StatusCode::BAD_REQUEST,
660 message: "Invalid multipart image edit request".to_string(),
661 param: None,
662 code: Some("invalid_multipart"),
663 });
664 log_native_request_completed(
665 &log,
666 &req_id,
667 operation.label(),
668 None,
669 response.status(),
670 started_at,
671 );
672 return monitor_response_body(response, request_guard);
673 }
674 };
675 parse_multipart_image_edit(multipart).await
676 } else {
677 let response = image_error_response(ImageRequestError {
678 status: StatusCode::UNSUPPORTED_MEDIA_TYPE,
679 message: if operation == ImageOperation::Edit {
680 "Image edit request must use application/json or multipart/form-data".to_string()
681 } else {
682 "Image generation request must use application/json".to_string()
683 },
684 param: None,
685 code: Some("unsupported_media_type"),
686 });
687 log_native_request_completed(
688 &log,
689 &req_id,
690 operation.label(),
691 None,
692 response.status(),
693 started_at,
694 );
695 return monitor_response_body(response, request_guard);
696 };
697 let prepared = match prepared {
698 Ok(prepared) => prepared,
699 Err(error) => {
700 let response = image_error_response(error);
701 log_native_request_completed(
702 &log,
703 &req_id,
704 operation.label(),
705 None,
706 response.status(),
707 started_at,
708 );
709 return monitor_response_body(response, request_guard);
710 }
711 };
712 let model = prepared.model.clone();
713 if let Some(monitor) = state.monitor.as_ref() {
714 monitor.provider_selected(&req_id, "codex", &model, None);
715 }
716 let context = RequestContext {
717 req_id: req_id.clone(),
718 session_id,
719 session_seq: None,
720 provider: "codex".to_string(),
721 traffic: None,
722 monitor: state.monitor.clone(),
723 passthrough: None,
724 };
725 let response = match state.images.as_ref() {
726 Some(backend) => backend.handle(operation, prepared, context).await,
727 None => image_error_response(ImageRequestError {
728 status: StatusCode::SERVICE_UNAVAILABLE,
729 message: "Codex Images API is unavailable".to_string(),
730 param: None,
731 code: Some("images_api_unavailable"),
732 }),
733 };
734 log_native_request_completed(
735 &log,
736 &req_id,
737 operation.label(),
738 Some(&model),
739 response.status(),
740 started_at,
741 );
742 monitor_response_body(response, request_guard)
743}
744
745async fn parse_multipart_image_edit(
746 mut multipart: Multipart,
747) -> Result<crate::providers::codex::images::PreparedImageRequest, ImageRequestError> {
748 let mut input = MultipartEditInput::default();
749 while let Some(field) = multipart
750 .next_field()
751 .await
752 .map_err(|_| ImageRequestError {
753 status: StatusCode::BAD_REQUEST,
754 message: "Invalid multipart image edit request".to_string(),
755 param: None,
756 code: Some("invalid_multipart"),
757 })?
758 {
759 let name = field.name().unwrap_or_default().to_string();
760 match name.as_str() {
761 "image" | "image[]" => {
762 let bytes = field.bytes().await.map_err(|_| ImageRequestError {
763 status: StatusCode::BAD_REQUEST,
764 message: "Failed to read uploaded image".to_string(),
765 param: Some("image"),
766 code: Some("invalid_image"),
767 })?;
768 input.images.push(UploadedImage {
769 bytes, });
771 }
772 "prompt" => {
773 input.prompt = Some(multipart_text(field, "prompt").await?);
774 }
775 "model" => {
776 input.model = Some(multipart_text(field, "model").await?);
777 }
778 "background" => {
779 input.background = Some(multipart_text(field, "background").await?);
780 }
781 "quality" => {
782 input.quality = Some(multipart_text(field, "quality").await?);
783 }
784 "size" => {
785 input.size = Some(multipart_text(field, "size").await?);
786 }
787 "n" => {
788 let raw = multipart_text(field, "n").await?;
789 input.n = Some(raw.parse::<u8>().map_err(|_| ImageRequestError {
790 status: StatusCode::BAD_REQUEST,
791 message: "'n' must be an integer between 1 and 10".to_string(),
792 param: Some("n"),
793 code: Some("invalid_request"),
794 })?);
795 }
796 "mask" => {
797 return Err(ImageRequestError {
798 status: StatusCode::BAD_REQUEST,
799 message: "Image masks are not supported by the Codex image backend".to_string(),
800 param: Some("mask"),
801 code: Some("unsupported_parameter"),
802 });
803 }
804 _ => {
805 return Err(ImageRequestError {
806 status: StatusCode::BAD_REQUEST,
807 message: format!("Unsupported multipart field '{name}'"),
808 param: None,
809 code: Some("unsupported_parameter"),
810 });
811 }
812 }
813 }
814 prepare_multipart_edit(input)
815}
816
817async fn multipart_text(
818 field: axum::extract::multipart::Field<'_>,
819 param: &'static str,
820) -> Result<String, ImageRequestError> {
821 field.text().await.map_err(|_| ImageRequestError {
822 status: StatusCode::BAD_REQUEST,
823 message: format!("Invalid multipart '{param}' field"),
824 param: Some(param),
825 code: Some("invalid_multipart"),
826 })
827}
828
829async fn handler_responses(State(state): State<Arc<AppState>>, req: Request<Body>) -> Response {
830 let started_at = Instant::now();
831 let log = create_logger("server");
832 let req_id = Uuid::new_v4().to_string();
833 let method = req.method().clone();
834 let uri = req.uri().clone();
835 let headers = req.headers().clone();
836 let path = uri.path().to_string();
837 let query = redacted_query(&uri);
838 log.info(
839 "request",
840 Some(serde_json::Map::from_iter([
841 ("reqId".to_string(), json!(&req_id)),
842 ("method".to_string(), json!(method.as_str())),
843 ("path".to_string(), json!(&path)),
844 ("query".to_string(), json!(&query)),
845 ])),
846 );
847
848 let session_id = native_session_id(&headers);
849 if let Some(monitor) = state.monitor.as_ref() {
850 monitor.request_started(&req_id, session_id.clone(), None, EndpointKind::Responses);
851 }
852 let request_guard = RequestMonitorGuard::new(state.monitor.clone(), req_id.clone());
853 let body_bytes = match axum::body::to_bytes(req.into_body(), MAX_OPENAI_REQUEST_BYTES).await {
854 Ok(bytes) => bytes,
855 Err(_) => {
856 let response = openai_error(
857 StatusCode::PAYLOAD_TOO_LARGE,
858 "invalid_request_error",
859 "Request body exceeded the size limit".to_string(),
860 None,
861 Some("request_too_large"),
862 );
863 log_native_request_completed(
864 &log,
865 &req_id,
866 "responses",
867 None,
868 response.status(),
869 started_at,
870 );
871 return monitor_response_body(response, request_guard);
872 }
873 };
874 let body: Value = match parse_native_json_body(&body_bytes) {
875 Ok(body) => body,
876 Err(response) => {
877 log_native_request_completed(
878 &log,
879 &req_id,
880 "responses",
881 None,
882 response.status(),
883 started_at,
884 );
885 return monitor_response_body(response, request_guard);
886 }
887 };
888 let (requested_model, normalized_model) = match extract_model(&body) {
889 Ok(model) => model,
890 Err(error) => return monitor_response_body(error.response(), request_guard),
891 };
892 let now = current_millis();
893 let session_state = session::existing_session(session_id.as_deref(), now);
894 let affinity = session_state
895 .as_ref()
896 .and_then(|session| session.affinity_provider.as_ref());
897 let Some(provider) = state
898 .registry
899 .provider_for_model(&normalized_model, affinity)
900 else {
901 let response = OpenAiError::invalid(
902 format!(
903 "Unknown model \"{normalized_model}\". {}",
904 state.registry.unknown_model_message()
905 ),
906 Some("model"),
907 )
908 .response();
909 return monitor_response_body(response, request_guard);
910 };
911 let parsed = if provider.name() == "codex" {
912 if let Err(response) = validate_native_request_model(&body) {
913 return monitor_response_body(response, request_guard);
914 }
915 None
916 } else {
917 match parse_request(
918 OpenAiSurface::Responses,
919 body.clone(),
920 provider.name(),
921 session_id.as_deref(),
922 ) {
923 Ok(parsed) => Some(parsed),
924 Err(error) => return monitor_response_body(error.response(), request_guard),
925 }
926 };
927 if provider.name() != "codex"
928 && let Some(session_id) = session_id.as_deref()
929 {
930 crate::providers::codex::clear_session_compaction(session_id);
931 }
932 let current = session::record_session_request(
933 session_id.as_deref(),
934 session_state.as_ref(),
935 provider.name(),
936 &normalized_model,
937 now,
938 );
939 let effort = parsed
940 .as_ref()
941 .and_then(|parsed| {
942 parsed
943 .messages
944 .extra
945 .get("output_config")
946 .and_then(|value| value.get("effort"))
947 .and_then(Value::as_str)
948 })
949 .or_else(|| body.pointer("/reasoning/effort").and_then(Value::as_str))
950 .map(str::to_string);
951 if let Some(monitor) = state.monitor.as_ref() {
952 if let Some(current) = current.as_ref() {
953 monitor.session_sequence_resolved(&req_id, current.seq);
954 }
955 monitor.provider_selected(&req_id, provider.name(), &normalized_model, effort);
956 }
957 let traffic = create_traffic_capture(TrafficCaptureOptions {
958 req_id: req_id.clone(),
959 session_id: session_id.clone(),
960 session_seq: current.as_ref().map(|session| session.seq),
961 provider: Some(provider.name().to_string()),
962 state_dir_override: None,
963 })
964 .map(Arc::new);
965 if let Some(capture) = traffic.as_ref() {
966 if let Some(monitor) = state.monitor.as_ref() {
967 monitor.traffic_capture_path(&req_id, capture.root().to_path_buf());
968 }
969 capture.write_json(
970 "000-metadata",
971 &json!({
972 "reqId": &req_id,
973 "sessionId": &session_id,
974 "sessionSeq": current.as_ref().map(|session| session.seq),
975 "kind": "responses",
976 "provider": provider.name(),
977 "model": &normalized_model,
978 "requestedModel": &requested_model,
979 "method": method.as_str(),
980 "path": &path,
981 "query": &query,
982 "headers": headers_to_record(&headers),
983 }),
984 );
985 capture.write_json("010-openai-responses-request", &body);
986 if let Some(parsed) = parsed.as_ref() {
987 capture.write_json(
988 "015-anthropic-request",
989 &serde_json::to_value(&parsed.messages).unwrap_or_else(|_| json!({})),
990 );
991 }
992 }
993 let context = RequestContext {
994 req_id: req_id.clone(),
995 session_id,
996 session_seq: current.map(|session| session.seq),
997 provider: provider.name().to_string(),
998 traffic: traffic.clone(),
999 monitor: state.monitor.clone(),
1000 passthrough: None,
1001 };
1002 let response = if let Some(parsed) = parsed {
1003 match provider
1004 .generate_anthropic_stream(parsed.messages, context)
1005 .await
1006 {
1007 Ok(generation) => {
1008 if let Some(capture) = traffic.as_ref() {
1009 capture.write_json(
1010 "016-model-resolution",
1011 &json!({"requestedModel":requested_model,"resolvedModel":generation.resolved_model}),
1012 );
1013 }
1014 match render_openai_response(
1015 OpenAiSurface::Responses,
1016 generation,
1017 parsed.stream,
1018 parsed.include_usage,
1019 parsed.response_metadata,
1020 traffic.clone(),
1021 )
1022 .await
1023 {
1024 Ok(response) => response,
1025 Err(error) => error.response(),
1026 }
1027 }
1028 Err(error) => OpenAiError::from(error).response(),
1029 }
1030 } else {
1031 match state.native_responses.as_ref() {
1032 Some(backend) => backend.handle(body, context).await,
1033 None => openai_error(
1034 StatusCode::NOT_FOUND,
1035 "not_found_error",
1036 "Native Responses API is disabled",
1037 None,
1038 None,
1039 ),
1040 }
1041 };
1042 log_routed_openai_request_completed(
1043 &log,
1044 &req_id,
1045 "responses",
1046 provider.name(),
1047 Some(&normalized_model),
1048 response.status(),
1049 started_at,
1050 );
1051 monitor_response_body(response, request_guard)
1052}
1053
1054async fn handler_chat_completions(
1055 State(state): State<Arc<AppState>>,
1056 req: Request<Body>,
1057) -> Response {
1058 let started_at = Instant::now();
1059 let log = create_logger("server");
1060 let req_id = Uuid::new_v4().to_string();
1061 let method = req.method().clone();
1062 let uri = req.uri().clone();
1063 let headers = req.headers().clone();
1064 let path = uri.path().to_string();
1065 let query = redacted_query(&uri);
1066 log.info(
1067 "request",
1068 Some(serde_json::Map::from_iter([
1069 ("reqId".to_string(), json!(&req_id)),
1070 ("method".to_string(), json!(method.as_str())),
1071 ("path".to_string(), json!(&path)),
1072 ("query".to_string(), json!(&query)),
1073 ])),
1074 );
1075
1076 let session_id = native_session_id(&headers);
1077 if let Some(monitor) = state.monitor.as_ref() {
1078 monitor.request_started(
1079 &req_id,
1080 session_id.clone(),
1081 None,
1082 EndpointKind::ChatCompletions,
1083 );
1084 }
1085 let request_guard = RequestMonitorGuard::new(state.monitor.clone(), req_id.clone());
1086 let body_bytes = match axum::body::to_bytes(req.into_body(), MAX_OPENAI_REQUEST_BYTES).await {
1087 Ok(bytes) => bytes,
1088 Err(_) => {
1089 let response = openai_error(
1090 StatusCode::PAYLOAD_TOO_LARGE,
1091 "invalid_request_error",
1092 "Request body exceeded the size limit".to_string(),
1093 None,
1094 Some("request_too_large"),
1095 );
1096 log_native_request_completed(
1097 &log,
1098 &req_id,
1099 "chat_completions",
1100 None,
1101 response.status(),
1102 started_at,
1103 );
1104 return monitor_response_body(response, request_guard);
1105 }
1106 };
1107 let body = match parse_native_json_body(&body_bytes) {
1108 Ok(body) => body,
1109 Err(response) => {
1110 log_native_request_completed(
1111 &log,
1112 &req_id,
1113 "chat_completions",
1114 None,
1115 response.status(),
1116 started_at,
1117 );
1118 return monitor_response_body(response, request_guard);
1119 }
1120 };
1121 let (requested_model, normalized_model) = match extract_model(&body) {
1122 Ok(model) => model,
1123 Err(error) => return monitor_response_body(error.response(), request_guard),
1124 };
1125 let now = current_millis();
1126 let session_state = session::existing_session(session_id.as_deref(), now);
1127 let affinity = session_state
1128 .as_ref()
1129 .and_then(|session| session.affinity_provider.as_ref());
1130 let Some(provider) = state
1131 .registry
1132 .provider_for_model(&normalized_model, affinity)
1133 else {
1134 let response = OpenAiError::invalid(
1135 format!(
1136 "Unknown model \"{normalized_model}\". {}",
1137 state.registry.unknown_model_message()
1138 ),
1139 Some("model"),
1140 )
1141 .response();
1142 return monitor_response_body(response, request_guard);
1143 };
1144 let (translated, parsed) = if provider.name() == "codex" {
1145 let translated = match translate_request(body.clone()) {
1146 Ok(translated) => translated,
1147 Err(error) => return monitor_response_body(error.response(), request_guard),
1148 };
1149 (Some(translated), None)
1150 } else {
1151 let parsed = match parse_request(
1152 OpenAiSurface::ChatCompletions,
1153 body.clone(),
1154 provider.name(),
1155 session_id.as_deref(),
1156 ) {
1157 Ok(parsed) => parsed,
1158 Err(error) => return monitor_response_body(error.response(), request_guard),
1159 };
1160 (None, Some(parsed))
1161 };
1162 if provider.name() != "codex"
1163 && let Some(session_id) = session_id.as_deref()
1164 {
1165 crate::providers::codex::clear_session_compaction(session_id);
1166 }
1167 let current = session::record_session_request(
1168 session_id.as_deref(),
1169 session_state.as_ref(),
1170 provider.name(),
1171 &normalized_model,
1172 now,
1173 );
1174 let effort = translated
1175 .as_ref()
1176 .and_then(|translated| translated.effort.clone())
1177 .or_else(|| {
1178 parsed.as_ref().and_then(|parsed| {
1179 parsed
1180 .messages
1181 .extra
1182 .get("output_config")
1183 .and_then(|value| value.get("effort"))
1184 .and_then(Value::as_str)
1185 .map(str::to_string)
1186 })
1187 });
1188 if let Some(monitor) = state.monitor.as_ref() {
1189 if let Some(current) = current.as_ref() {
1190 monitor.session_sequence_resolved(&req_id, current.seq);
1191 }
1192 monitor.provider_selected(&req_id, provider.name(), &normalized_model, effort.clone());
1193 }
1194 let traffic = create_traffic_capture(TrafficCaptureOptions {
1195 req_id: req_id.clone(),
1196 session_id: session_id.clone(),
1197 session_seq: current.as_ref().map(|session| session.seq),
1198 provider: Some(provider.name().to_string()),
1199 state_dir_override: None,
1200 })
1201 .map(Arc::new);
1202 if let Some(capture) = traffic.as_ref() {
1203 if let Some(monitor) = state.monitor.as_ref() {
1204 monitor.traffic_capture_path(&req_id, capture.root().to_path_buf());
1205 }
1206 capture.write_json(
1207 "000-metadata",
1208 &json!({
1209 "reqId": &req_id,
1210 "sessionId": &session_id,
1211 "sessionSeq": current.as_ref().map(|session| session.seq),
1212 "kind": "chat_completions",
1213 "provider": provider.name(),
1214 "model": &normalized_model,
1215 "requestedModel": &requested_model,
1216 "effort": &effort,
1217 "method": method.as_str(),
1218 "path": &path,
1219 "query": &query,
1220 "headers": headers_to_record(&headers),
1221 }),
1222 );
1223 capture.write_json("010-openai-chat-completions-request", &body);
1224 if let Some(translated) = translated.as_ref() {
1225 capture.write_json("020-upstream-request", &translated.upstream);
1226 }
1227 if let Some(parsed) = parsed.as_ref() {
1228 capture.write_json(
1229 "015-anthropic-request",
1230 &serde_json::to_value(&parsed.messages).unwrap_or_else(|_| json!({})),
1231 );
1232 }
1233 }
1234 let context = RequestContext {
1235 req_id: req_id.clone(),
1236 session_id,
1237 session_seq: current.map(|session| session.seq),
1238 provider: provider.name().to_string(),
1239 traffic: traffic.clone(),
1240 monitor: state.monitor.clone(),
1241 passthrough: None,
1242 };
1243 let response = if let Some(translated) = translated {
1244 match state.chat_completions.as_ref() {
1245 Some(backend) => backend.handle(translated, context).await,
1246 None => openai_error(
1247 StatusCode::NOT_FOUND,
1248 "not_found_error",
1249 "Chat Completions API is disabled",
1250 None,
1251 None,
1252 ),
1253 }
1254 } else {
1255 let parsed = parsed.expect("non-Codex request was parsed");
1256 match provider
1257 .generate_anthropic_stream(parsed.messages, context)
1258 .await
1259 {
1260 Ok(generation) => {
1261 if let Some(capture) = traffic.as_ref() {
1262 capture.write_json(
1263 "016-model-resolution",
1264 &json!({"requestedModel":requested_model,"resolvedModel":generation.resolved_model}),
1265 );
1266 }
1267 match render_openai_response(
1268 OpenAiSurface::ChatCompletions,
1269 generation,
1270 parsed.stream,
1271 parsed.include_usage,
1272 parsed.response_metadata,
1273 traffic.clone(),
1274 )
1275 .await
1276 {
1277 Ok(response) => response,
1278 Err(error) => error.response(),
1279 }
1280 }
1281 Err(error) => OpenAiError::from(error).response(),
1282 }
1283 };
1284 log_routed_openai_request_completed(
1285 &log,
1286 &req_id,
1287 "chat_completions",
1288 provider.name(),
1289 Some(&normalized_model),
1290 response.status(),
1291 started_at,
1292 );
1293 monitor_response_body(response, request_guard)
1294}
1295
1296fn native_session_id(headers: &http::HeaderMap) -> Option<String> {
1297 [
1298 "x-claude-code-session-id",
1299 "session_id",
1300 "x-client-request-id",
1301 ]
1302 .into_iter()
1303 .find_map(|name| {
1304 headers
1305 .get(name)
1306 .and_then(|value| value.to_str().ok())
1307 .filter(|value| !value.is_empty())
1308 .map(str::to_string)
1309 })
1310}
1311
1312#[allow(clippy::result_large_err)]
1313fn parse_native_json_body(body: &[u8]) -> Result<Value, Response> {
1314 if body.is_empty() {
1315 return Err(openai_error(
1316 StatusCode::BAD_REQUEST,
1317 "invalid_request_error",
1318 "Invalid JSON: empty body",
1319 None,
1320 Some("invalid_json"),
1321 ));
1322 }
1323 serde_json::from_slice(body).map_err(|error| {
1324 openai_error(
1325 StatusCode::BAD_REQUEST,
1326 "invalid_request_error",
1327 format!("Invalid JSON: {error}"),
1328 None,
1329 Some("invalid_json"),
1330 )
1331 })
1332}
1333
1334fn log_routed_openai_request_completed(
1335 log: &Logger,
1336 req_id: &str,
1337 endpoint: &str,
1338 provider: &str,
1339 model: Option<&str>,
1340 status: StatusCode,
1341 started_at: Instant,
1342) {
1343 log.info(
1344 "request_completed",
1345 Some(serde_json::Map::from_iter([
1346 ("reqId".to_string(), json!(req_id)),
1347 ("endpoint".to_string(), json!(endpoint)),
1348 ("provider".to_string(), json!(provider)),
1349 ("model".to_string(), json!(model)),
1350 ("countTokens".to_string(), json!(false)),
1351 ("status".to_string(), json!(status.as_u16())),
1352 ("ms".to_string(), json!(started_at.elapsed().as_millis())),
1353 ])),
1354 );
1355}
1356
1357fn log_native_request_completed(
1358 log: &Logger,
1359 req_id: &str,
1360 endpoint: &str,
1361 model: Option<&str>,
1362 status: StatusCode,
1363 started_at: Instant,
1364) {
1365 log.info(
1366 "request_completed",
1367 Some(serde_json::Map::from_iter([
1368 ("reqId".to_string(), json!(req_id)),
1369 ("endpoint".to_string(), json!(endpoint)),
1370 ("provider".to_string(), json!("codex")),
1371 ("model".to_string(), json!(model)),
1372 ("countTokens".to_string(), json!(false)),
1373 ("status".to_string(), json!(status.as_u16())),
1374 ("ms".to_string(), json!(started_at.elapsed().as_millis())),
1375 ])),
1376 );
1377}
1378
1379async fn dispatch_request(
1380 state: Arc<AppState>,
1381 req: Request<Body>,
1382 count_tokens: bool,
1383) -> Response {
1384 let started_at = Instant::now();
1385 let log = create_logger("server");
1386 let req_id = Uuid::new_v4().to_string();
1387 let method = req.method().clone();
1388 let uri = req.uri().clone();
1389 let headers = req.headers().clone();
1390 let conversation_identity = (!count_tokens)
1391 .then(|| ConversationIdentity::from_headers(&headers))
1392 .flatten();
1393 let path = uri.path().to_string();
1394 let query = redacted_query(&uri);
1395 let endpoint = if count_tokens {
1396 EndpointKind::CountTokens
1397 } else {
1398 EndpointKind::Messages
1399 };
1400 log.info(
1401 "request",
1402 Some(serde_json::Map::from_iter([
1403 ("reqId".to_string(), json!(&req_id)),
1404 ("method".to_string(), json!(method.as_str())),
1405 ("path".to_string(), json!(&path)),
1406 ("query".to_string(), json!(&query)),
1407 ])),
1408 );
1409 let session_id = req
1410 .headers()
1411 .get("x-claude-code-session-id")
1412 .and_then(|value| value.to_str().ok())
1413 .map(std::string::ToString::to_string);
1414 if let Some(monitor) = state.monitor.as_ref() {
1415 monitor.request_started(&req_id, session_id.clone(), None, endpoint);
1416 }
1417 let request_guard = RequestMonitorGuard::new(state.monitor.clone(), req_id.clone());
1418 let now = current_millis();
1419 let body_bytes = match axum::body::to_bytes(req.into_body(), MAX_OPENAI_REQUEST_BYTES).await {
1420 Ok(bytes) => bytes,
1421 Err(err) => {
1422 let response = json_error(
1423 StatusCode::BAD_REQUEST,
1424 "invalid_request_error",
1425 format!("Invalid JSON: {err}"),
1426 );
1427 log_request_completed(
1428 &log,
1429 RequestLogContext {
1430 req_id: &req_id,
1431 provider: None,
1432 model: None,
1433 count_tokens,
1434 status: response.status(),
1435 started_at,
1436 },
1437 );
1438 let (response, details) = record_failed_response(
1439 &log,
1440 FailedResponseLogContext {
1441 req_id: &req_id,
1442 provider: None,
1443 model: None,
1444 count_tokens,
1445 started_at,
1446 },
1447 response,
1448 )
1449 .await;
1450 monitor_failed(
1451 state.monitor.as_ref(),
1452 &req_id,
1453 Some(response.status()),
1454 details
1455 .as_ref()
1456 .map(|details| details.message.as_str())
1457 .unwrap_or("Invalid JSON"),
1458 );
1459 return response;
1460 }
1461 };
1462
1463 let mut body: crate::anthropic::schema::MessagesRequest = match parse_json_body(&body_bytes) {
1464 Ok(body) => body,
1465 Err(response) => {
1466 let status = response.status();
1467 log_request_completed(
1468 &log,
1469 RequestLogContext {
1470 req_id: &req_id,
1471 provider: None,
1472 model: None,
1473 count_tokens,
1474 status: response.status(),
1475 started_at,
1476 },
1477 );
1478 let (response, details) = record_failed_response(
1479 &log,
1480 FailedResponseLogContext {
1481 req_id: &req_id,
1482 provider: None,
1483 model: None,
1484 count_tokens,
1485 started_at,
1486 },
1487 *response,
1488 )
1489 .await;
1490 monitor_failed(
1491 state.monitor.as_ref(),
1492 &req_id,
1493 Some(status),
1494 details
1495 .as_ref()
1496 .map(|details| details.message.as_str())
1497 .unwrap_or("Invalid JSON"),
1498 );
1499 return response;
1500 }
1501 };
1502
1503 if let Some(project) = project::name_from_request(
1504 body.extra.get("system"),
1505 body.messages.iter().rev().map(|message| &message.content),
1506 ) && let Some(monitor) = state.monitor.as_ref()
1507 {
1508 monitor.project_resolved(&req_id, project);
1509 }
1510
1511 let model = match body.model.as_deref() {
1512 Some(model) => model,
1513 None => {
1514 let response = json_error(
1515 StatusCode::BAD_REQUEST,
1516 "invalid_request_error",
1517 format!(
1518 "Missing \"model\" in request body. {}",
1519 state.registry.unknown_model_message()
1520 ),
1521 );
1522 log_request_completed(
1523 &log,
1524 RequestLogContext {
1525 req_id: &req_id,
1526 provider: None,
1527 model: None,
1528 count_tokens,
1529 status: response.status(),
1530 started_at,
1531 },
1532 );
1533 let (response, details) = record_failed_response(
1534 &log,
1535 FailedResponseLogContext {
1536 req_id: &req_id,
1537 provider: None,
1538 model: None,
1539 count_tokens,
1540 started_at,
1541 },
1542 response,
1543 )
1544 .await;
1545 monitor_failed(
1546 state.monitor.as_ref(),
1547 &req_id,
1548 Some(response.status()),
1549 details
1550 .as_ref()
1551 .map(|details| details.message.as_str())
1552 .unwrap_or("Missing model"),
1553 );
1554 return response;
1555 }
1556 };
1557
1558 let mut normalized_model = normalize_incoming_model(model);
1559 body.model = Some(normalized_model.clone());
1560 let session_state = if let Some(session_id) = session_id.as_deref() {
1561 session::existing_session(Some(session_id), now)
1562 } else {
1563 None
1564 };
1565 let session_affinity = session_state
1566 .as_ref()
1567 .and_then(|state| state.affinity_provider.as_ref());
1568 let original_provider = state
1569 .registry
1570 .provider_for_model(&normalized_model, session_affinity);
1571 let configured_auto_review_model = crate::config::auto_review_model();
1572 let auto_review_route = original_provider.as_ref().and_then(|provider| {
1573 apply_auto_review_model(
1574 &mut body,
1575 count_tokens,
1576 configured_auto_review_model.as_deref(),
1577 provider.name(),
1578 )
1579 });
1580 if auto_review_route.is_some() {
1581 normalized_model = normalize_incoming_model(body.model.as_deref().expect("override model"));
1582 body.model = Some(normalized_model.clone());
1583 }
1584
1585 let provider = if auto_review_route.is_some() {
1586 state.registry.provider_for_model(&normalized_model, None)
1587 } else {
1588 original_provider
1589 };
1590
1591 let provider = match provider {
1592 Some(provider) => provider,
1593 None => {
1594 log.warn(
1595 "unknown model",
1596 Some(serde_json::Map::from_iter([
1597 ("reqId".to_string(), json!(&req_id)),
1598 ("model".to_string(), json!(&normalized_model)),
1599 ])),
1600 );
1601 let response = json_error(
1602 StatusCode::BAD_REQUEST,
1603 "invalid_request_error",
1604 format!(
1605 "Unknown model \"{normalized_model}\". {}",
1606 state.registry.unknown_model_message()
1607 ),
1608 );
1609 log_request_completed(
1610 &log,
1611 RequestLogContext {
1612 req_id: &req_id,
1613 provider: None,
1614 model: Some(&normalized_model),
1615 count_tokens,
1616 status: response.status(),
1617 started_at,
1618 },
1619 );
1620 let (response, details) = record_failed_response(
1621 &log,
1622 FailedResponseLogContext {
1623 req_id: &req_id,
1624 provider: None,
1625 model: Some(&normalized_model),
1626 count_tokens,
1627 started_at,
1628 },
1629 response,
1630 )
1631 .await;
1632 monitor_failed(
1633 state.monitor.as_ref(),
1634 &req_id,
1635 Some(response.status()),
1636 details
1637 .as_ref()
1638 .map(|details| details.message.as_str())
1639 .unwrap_or("Unknown model"),
1640 );
1641 return response;
1642 }
1643 };
1644
1645 body.bypass_provider_model_override = auto_review_route.is_some() && provider.name() == "codex";
1646
1647 if let Some(route) = auto_review_route.as_ref() {
1648 log.info(
1649 "auto-review route selected",
1650 Some(Map::from_iter([
1651 ("reqId".to_string(), json!(&req_id)),
1652 ("requestedModel".to_string(), json!(&route.requested_model)),
1653 ("overrideModel".to_string(), json!(&route.override_model)),
1654 ("provider".to_string(), json!(provider.name())),
1655 ])),
1656 );
1657 }
1658
1659 if !count_tokens
1660 && auto_review_route.is_none()
1661 && provider.name() != "codex"
1662 && let Some(session_id) = session_id.as_deref()
1663 {
1664 crate::providers::codex::clear_session_compaction(session_id);
1665 }
1666
1667 let effort = crate::providers::translate_shared::read_effort(&body)
1668 .ok()
1669 .flatten()
1670 .map(str::to_string);
1671 let current = session::record_session_request_with_affinity_update(
1672 session_id.as_deref(),
1673 session_state.as_ref(),
1674 provider.name(),
1675 &normalized_model,
1676 auto_review_route.is_none(),
1677 now,
1678 );
1679 if let Some(monitor) = state.monitor.as_ref() {
1680 if let Some(current) = current.as_ref() {
1681 monitor.session_sequence_resolved(&req_id, current.seq);
1682 }
1683 monitor.provider_selected(&req_id, provider.name(), &normalized_model, effort);
1684 }
1685
1686 let traffic = create_traffic_capture(TrafficCaptureOptions {
1687 req_id: req_id.clone(),
1688 session_id: session_id.clone(),
1689 session_seq: current.as_ref().map(|s| s.seq),
1690 provider: Some(provider.name().to_string()),
1691 state_dir_override: None,
1692 })
1693 .map(Arc::new);
1694
1695 if let Some(capture) = traffic.as_ref() {
1696 if let Some(monitor) = state.monitor.as_ref() {
1697 monitor.traffic_capture_path(&req_id, capture.root().to_path_buf());
1698 }
1699 capture.write_json(
1700 "000-metadata",
1701 &json!({
1702 "reqId": &req_id,
1703 "sessionId": &session_id,
1704 "sessionSeq": current.as_ref().map(|s| s.seq),
1705 "kind": if count_tokens { "count_tokens" } else { "messages" },
1706 "provider": provider.name(),
1707 "model": &normalized_model,
1708 "method": method.as_str(),
1709 "path": &path,
1710 "query": &query,
1711 "headers": headers_to_record(&headers),
1712 }),
1713 );
1714 capture.write_json(
1715 "010-anthropic-request",
1716 &serde_json::to_value(&body).unwrap_or_else(|_| json!({})),
1717 );
1718 }
1719
1720 let context = RequestContext {
1721 req_id: req_id.clone(),
1722 session_id,
1723 session_seq: current.map(|s| s.seq),
1724 provider: provider.name().to_string(),
1725 traffic,
1726 monitor: state.monitor.clone(),
1727 passthrough: Some(crate::provider::Passthrough {
1728 raw_body: body_bytes,
1729 headers,
1730 path_and_query: uri
1731 .path_and_query()
1732 .map(|pq| pq.as_str().to_string())
1733 .unwrap_or_else(|| path.clone()),
1734 }),
1735 };
1736
1737 let response = if count_tokens {
1738 provider.handle_count_tokens(body, context).await
1739 } else {
1740 provider
1741 .handle_messages_with_conversation_identity(
1742 body,
1743 context,
1744 if auto_review_route.is_some() {
1745 None
1746 } else {
1747 conversation_identity
1748 },
1749 )
1750 .await
1751 };
1752 log_request_completed(
1753 &log,
1754 RequestLogContext {
1755 req_id: &req_id,
1756 provider: Some(provider.name()),
1757 model: Some(&normalized_model),
1758 count_tokens,
1759 status: response.status(),
1760 started_at,
1761 },
1762 );
1763 let status = response.status();
1764 if status.is_success() {
1765 return monitor_response_body(response, request_guard);
1766 }
1767
1768 let (response, details) = record_failed_response(
1769 &log,
1770 FailedResponseLogContext {
1771 req_id: &req_id,
1772 provider: Some(provider.name()),
1773 model: Some(&normalized_model),
1774 count_tokens,
1775 started_at,
1776 },
1777 response,
1778 )
1779 .await;
1780 if let Some(details) = details.as_ref() {
1781 monitor_failed(
1782 state.monitor.as_ref(),
1783 &req_id,
1784 Some(status),
1785 details.message.as_str(),
1786 );
1787 } else {
1788 monitor_failed(
1789 state.monitor.as_ref(),
1790 &req_id,
1791 Some(status),
1792 format!("HTTP {}", status.as_u16()),
1793 );
1794 }
1795 response
1796}
1797
1798fn monitor_response_body(response: Response, guard: RequestMonitorGuard) -> Response {
1799 let status = response.status();
1800 let outcome = response
1801 .extensions()
1802 .get::<NativeResponseOutcome>()
1803 .cloned();
1804 let (parts, body) = response.into_parts();
1805 let stream = futures_util::stream::unfold(
1806 (body, guard, outcome),
1807 move |(mut body, mut guard, outcome)| async move {
1808 match body.frame().await {
1809 Some(Ok(frame)) => Some((Ok(frame), (body, guard, outcome))),
1810 Some(Err(err)) => {
1811 guard.failed(status, err.to_string());
1812 Some((Err(err), (body, guard, outcome)))
1813 }
1814 None => {
1815 if let Some(message) = outcome.as_ref().and_then(NativeResponseOutcome::failure)
1816 {
1817 guard.failed(status, message);
1818 } else if status.is_success() {
1819 guard.completed(status);
1820 } else {
1821 guard.failed(status, format!("HTTP {}", status.as_u16()));
1822 }
1823 None
1824 }
1825 }
1826 },
1827 );
1828 Response::from_parts(parts, Body::new(StreamBody::new(stream)))
1829}
1830
1831struct RequestLogContext<'a> {
1832 req_id: &'a str,
1833 provider: Option<&'a str>,
1834 model: Option<&'a str>,
1835 count_tokens: bool,
1836 status: StatusCode,
1837 started_at: Instant,
1838}
1839
1840fn log_request_completed(log: &Logger, ctx: RequestLogContext<'_>) {
1841 log.info(
1842 "request_completed",
1843 Some(serde_json::Map::from_iter([
1844 ("reqId".to_string(), json!(ctx.req_id)),
1845 ("provider".to_string(), json!(ctx.provider)),
1846 ("model".to_string(), json!(ctx.model)),
1847 ("countTokens".to_string(), json!(ctx.count_tokens)),
1848 ("status".to_string(), json!(ctx.status.as_u16())),
1849 (
1850 "ms".to_string(),
1851 json!(ctx.started_at.elapsed().as_millis()),
1852 ),
1853 ])),
1854 );
1855}
1856
1857struct FailedResponseLogContext<'a> {
1858 req_id: &'a str,
1859 provider: Option<&'a str>,
1860 model: Option<&'a str>,
1861 count_tokens: bool,
1862 started_at: Instant,
1863}
1864
1865struct FailedResponseDetails {
1866 message: String,
1867}
1868
1869async fn record_failed_response(
1870 log: &Logger,
1871 ctx: FailedResponseLogContext<'_>,
1872 response: Response,
1873) -> (Response, Option<FailedResponseDetails>) {
1874 if response.status().is_success() {
1875 return (response, None);
1876 }
1877
1878 let status = response.status();
1879 let (parts, body) = response.into_parts();
1880 let bytes = match body.collect().await {
1881 Ok(collected) => collected.to_bytes(),
1882 Err(err) => {
1883 log.info(
1884 "request_failed",
1885 Some(serde_json::Map::from_iter([
1886 ("reqId".to_string(), json!(ctx.req_id)),
1887 ("provider".to_string(), json!(ctx.provider)),
1888 ("model".to_string(), json!(ctx.model)),
1889 ("countTokens".to_string(), json!(ctx.count_tokens)),
1890 ("status".to_string(), json!(status.as_u16())),
1891 (
1892 "ms".to_string(),
1893 json!(ctx.started_at.elapsed().as_millis()),
1894 ),
1895 ("bodyReadError".to_string(), json!(err.to_string())),
1896 ])),
1897 );
1898 return (Response::from_parts(parts, Body::empty()), None);
1899 }
1900 };
1901
1902 let response_body = response_body_value(&bytes);
1903 let message = error_message_from_response(&response_body)
1904 .unwrap_or_else(|| format!("HTTP {}", status.as_u16()));
1905 let document = json!({
1906 "reqId": ctx.req_id,
1907 "provider": ctx.provider,
1908 "model": ctx.model,
1909 "countTokens": ctx.count_tokens,
1910 "status": status.as_u16(),
1911 "elapsedMs": ctx.started_at.elapsed().as_millis(),
1912 "message": message,
1913 "response": response_body,
1914 });
1915 let error_file = write_error_capture(ctx.req_id, &redact_error_value(document));
1916
1917 let mut fields = serde_json::Map::from_iter([
1918 ("reqId".to_string(), json!(ctx.req_id)),
1919 ("provider".to_string(), json!(ctx.provider)),
1920 ("model".to_string(), json!(ctx.model)),
1921 ("countTokens".to_string(), json!(ctx.count_tokens)),
1922 ("status".to_string(), json!(status.as_u16())),
1923 (
1924 "ms".to_string(),
1925 json!(ctx.started_at.elapsed().as_millis()),
1926 ),
1927 ("message".to_string(), json!(message)),
1928 ]);
1929 if let Some(path) = error_file.as_ref() {
1930 fields.insert("errorFile".to_string(), json!(path.display().to_string()));
1931 }
1932 log.info("request_failed", Some(fields));
1933
1934 (
1935 Response::from_parts(parts, Body::from(bytes)),
1936 Some(FailedResponseDetails { message }),
1937 )
1938}
1939
1940fn response_body_value(bytes: &[u8]) -> Value {
1941 match serde_json::from_slice::<Value>(bytes) {
1942 Ok(value) => json!({ "json": value }),
1943 Err(_) => json!({ "text": String::from_utf8_lossy(bytes) }),
1944 }
1945}
1946
1947fn error_message_from_response(response_body: &Value) -> Option<String> {
1948 response_body
1949 .get("json")
1950 .and_then(|body| body.get("error"))
1951 .and_then(|error| error.get("message"))
1952 .and_then(Value::as_str)
1953 .or_else(|| {
1954 response_body
1955 .get("text")
1956 .and_then(Value::as_str)
1957 .map(str::trim)
1958 .filter(|text| !text.is_empty())
1959 })
1960 .map(std::string::ToString::to_string)
1961}
1962
1963fn write_error_capture(req_id: &str, document: &Value) -> Option<PathBuf> {
1964 let dir = crate::paths::state_dir().join("errors");
1965 fs::create_dir_all(&dir).ok()?;
1966 set_mode(&dir, 0o700);
1967 let path = dir.join(format!(
1968 "{}-{}.json",
1969 current_millis(),
1970 sanitize_path_part(req_id)
1971 ));
1972 let mut file = File::create(&path).ok()?;
1973 set_mode(&path, 0o600);
1974 let payload = serde_json::to_vec_pretty(document).ok()?;
1975 file.write_all(&payload).ok()?;
1976 file.write_all(b"\n").ok()?;
1977 Some(path)
1978}
1979
1980fn sanitize_path_part(raw: &str) -> String {
1981 let sanitized: String = raw
1982 .chars()
1983 .map(|ch| {
1984 if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
1985 ch
1986 } else {
1987 '_'
1988 }
1989 })
1990 .collect();
1991 if sanitized.is_empty() {
1992 "unknown".to_string()
1993 } else {
1994 sanitized
1995 }
1996}
1997
1998fn redact_error_value(value: Value) -> Value {
1999 match value {
2000 Value::Array(values) => Value::Array(values.into_iter().map(redact_error_value).collect()),
2001 Value::Object(fields) => {
2002 let mut out = Map::new();
2003 for (key, value) in fields {
2004 if REDACT_KEYS.contains(&key.to_lowercase().as_str()) {
2005 out.insert(key, redact_error_key(value));
2006 } else {
2007 out.insert(key, redact_error_value(value));
2008 }
2009 }
2010 Value::Object(out)
2011 }
2012 value => value,
2013 }
2014}
2015
2016fn redact_error_key(value: Value) -> Value {
2017 match value {
2018 Value::String(value) => Value::String(format!("[redacted len={}]", value.len())),
2019 _ => Value::String("[redacted]".to_string()),
2020 }
2021}
2022
2023struct RequestMonitorGuard {
2024 monitor: Option<MonitorHandle>,
2025 req_id: String,
2026}
2027
2028impl RequestMonitorGuard {
2029 fn new(monitor: Option<MonitorHandle>, req_id: String) -> Self {
2030 Self { monitor, req_id }
2031 }
2032
2033 fn completed(&mut self, status: StatusCode) {
2034 if let Some(monitor) = self.monitor.take() {
2035 monitor.request_completed(&self.req_id, status.as_u16(), None, None);
2036 }
2037 }
2038
2039 fn failed(&mut self, status: StatusCode, error: String) {
2040 if let Some(monitor) = self.monitor.take() {
2041 monitor.request_failed(&self.req_id, Some(status.as_u16()), error);
2042 }
2043 }
2044}
2045
2046impl Drop for RequestMonitorGuard {
2047 fn drop(&mut self) {
2048 if let Some(monitor) = self.monitor.as_ref() {
2049 monitor.request_abandoned(&self.req_id, "Request future ended before completion");
2050 }
2051 }
2052}
2053
2054fn monitor_failed(
2055 monitor: Option<&MonitorHandle>,
2056 req_id: &str,
2057 status: Option<StatusCode>,
2058 error: impl Into<String>,
2059) {
2060 if let Some(monitor) = monitor {
2061 monitor.request_failed(req_id, status.map(|status| status.as_u16()), error);
2062 }
2063}
2064
2065fn headers_to_record(headers: &http::HeaderMap) -> Value {
2066 let mut out = Map::new();
2067 for (key, value) in headers {
2068 if let Ok(raw) = value.to_str() {
2069 let recorded = if REDACT_KEYS.contains(&key.as_str().to_lowercase().as_str())
2070 || matches!(
2071 key.as_str(),
2072 CLAUDE_AGENT_HEADER | CLAUDE_PARENT_AGENT_HEADER
2073 ) {
2074 format!("[redacted len={}]", raw.len())
2075 } else {
2076 raw.to_string()
2077 };
2078 out.insert(key.as_str().to_string(), Value::String(recorded));
2079 }
2080 }
2081 Value::Object(out)
2082}
2083
2084fn redacted_query(uri: &http::Uri) -> Value {
2085 let mut out = Map::new();
2086 let Some(query) = uri.query() else {
2087 return Value::Object(out);
2088 };
2089 for (key, value) in url::form_urlencoded::parse(query.as_bytes()) {
2090 let key = key.into_owned();
2091 let lower = key.to_lowercase();
2092 let value = if REDACT_KEYS.contains(&lower.as_str()) {
2093 Value::String(format!("[redacted len={}]", value.len()))
2094 } else {
2095 Value::String(value.into_owned())
2096 };
2097 out.insert(key, value);
2098 }
2099 Value::Object(out)
2100}
2101
2102fn parse_json_body<T>(body: &[u8]) -> Result<T, Box<Response>>
2103where
2104 T: DeserializeOwned,
2105{
2106 if body.is_empty() {
2107 return Err(Box::new(json_error(
2108 StatusCode::BAD_REQUEST,
2109 "invalid_request_error",
2110 "Invalid JSON: empty body",
2111 )));
2112 }
2113
2114 serde_json::from_slice::<T>(body).map_err(|err| {
2115 Box::new(json_error(
2116 StatusCode::BAD_REQUEST,
2117 "invalid_request_error",
2118 format!("Invalid JSON: {err}"),
2119 ))
2120 })
2121}
2122
2123async fn fallback_handler(method: axum::http::Method, uri: axum::http::Uri) -> Response {
2124 json_error(
2125 StatusCode::NOT_FOUND,
2126 "not_found",
2127 format!("No route for {method} {}", uri.path()),
2128 )
2129}
2130
2131fn current_millis() -> u64 {
2132 use std::time::{SystemTime, UNIX_EPOCH};
2133 SystemTime::now()
2134 .duration_since(UNIX_EPOCH)
2135 .unwrap_or_default()
2136 .as_millis() as u64
2137}
2138
2139fn set_mode(path: &Path, mode: u32) {
2140 #[cfg(unix)]
2141 {
2142 use std::os::unix::fs::PermissionsExt;
2143 if let Ok(meta) = fs::metadata(path) {
2144 let mut perm = meta.permissions();
2145 perm.set_mode(mode);
2146 let _ = fs::set_permissions(path, perm);
2147 }
2148 }
2149 #[cfg(not(unix))]
2150 {
2151 let _ = (path, mode);
2152 }
2153}
2154
2155#[allow(dead_code)]
2156fn _unused(session_state: Option<&SessionState>) {
2157 let _ = session_state;
2158}
2159
2160#[cfg(test)]
2161mod auto_review_tests {
2162 use super::{apply_auto_review_model, headers_to_record, is_claude_auto_review_request};
2163 use crate::anthropic::schema::MessagesRequest;
2164 use crate::request_identity::{CLAUDE_AGENT_HEADER, CLAUDE_PARENT_AGENT_HEADER};
2165 use http::{HeaderMap, HeaderValue};
2166 use serde_json::json;
2167
2168 fn request(system: &str, stream: bool, tools: serde_json::Value) -> MessagesRequest {
2169 serde_json::from_value(json!({
2170 "model": "gpt-5.6-sol",
2171 "max_tokens": 2112,
2172 "stream": stream,
2173 "system": [{"type": "text", "text": system}],
2174 "messages": [{"role": "user", "content": "review this Bash command"}],
2175 "tools": tools
2176 }))
2177 .unwrap()
2178 }
2179
2180 #[test]
2181 fn detects_claude_auto_review_classifier() {
2182 let body = request(
2183 "You are a security monitor for autonomous AI coding agents.\n\n## Context",
2184 false,
2185 json!([]),
2186 );
2187 assert!(is_claude_auto_review_request(&body));
2188 }
2189
2190 #[test]
2191 fn ignores_normal_streaming_and_tool_using_requests() {
2192 assert!(!is_claude_auto_review_request(&request(
2193 "You are an interactive coding agent.",
2194 false,
2195 json!([]),
2196 )));
2197 assert!(!is_claude_auto_review_request(&request(
2198 "You are a security monitor for autonomous AI coding agents.",
2199 true,
2200 json!([]),
2201 )));
2202 assert!(!is_claude_auto_review_request(&request(
2203 "You are a security monitor for autonomous AI coding agents.",
2204 false,
2205 json!([{"name": "Bash"}]),
2206 )));
2207 }
2208
2209 #[test]
2210 fn codex_classifier_defaults_to_luna() {
2211 let mut classifier = request(
2212 "You are a security monitor for autonomous AI coding agents.",
2213 false,
2214 json!([]),
2215 );
2216 let route = apply_auto_review_model(&mut classifier, false, None, "codex")
2217 .expect("classifier should be routed");
2218 assert_eq!(route.requested_model, "gpt-5.6-sol");
2219 assert_eq!(route.override_model, "gpt-5.6-luna");
2220 assert_eq!(classifier.model.as_deref(), Some("gpt-5.6-luna"));
2221 }
2222
2223 #[test]
2224 fn non_codex_classifier_keeps_requested_model_without_override() {
2225 let mut classifier = request(
2226 "You are a security monitor for autonomous AI coding agents.",
2227 false,
2228 json!([]),
2229 );
2230 classifier.model = Some("kimi-for-coding".to_string());
2231 assert!(apply_auto_review_model(&mut classifier, false, None, "kimi").is_none());
2232 assert_eq!(classifier.model.as_deref(), Some("kimi-for-coding"));
2233 }
2234
2235 #[test]
2236 fn configured_model_overrides_provider_default() {
2237 let mut classifier = request(
2238 "You are a security monitor for autonomous AI coding agents.",
2239 false,
2240 json!([]),
2241 );
2242 let route = apply_auto_review_model(&mut classifier, false, Some("grok-4.5"), "codex")
2243 .expect("configured classifier should be routed");
2244 assert_eq!(route.override_model, "grok-4.5");
2245 assert_eq!(classifier.model.as_deref(), Some("grok-4.5"));
2246 }
2247
2248 #[test]
2249 fn traffic_headers_redact_agent_lineage_ids() {
2250 let mut headers = HeaderMap::new();
2251 headers.insert(
2252 "x-claude-code-session-id",
2253 HeaderValue::from_static("session-visible"),
2254 );
2255 headers.insert(
2256 CLAUDE_AGENT_HEADER,
2257 HeaderValue::from_static("agent-secret"),
2258 );
2259 headers.insert(
2260 CLAUDE_PARENT_AGENT_HEADER,
2261 HeaderValue::from_static("parent-secret"),
2262 );
2263
2264 let captured = headers_to_record(&headers);
2265 assert_eq!(
2266 captured["x-claude-code-session-id"],
2267 json!("session-visible")
2268 );
2269 assert_eq!(captured[CLAUDE_AGENT_HEADER], json!("[redacted len=12]"));
2270 assert_eq!(
2271 captured[CLAUDE_PARENT_AGENT_HEADER],
2272 json!("[redacted len=13]")
2273 );
2274 let serialized = captured.to_string();
2275 assert!(!serialized.contains("agent-secret"));
2276 assert!(!serialized.contains("parent-secret"));
2277 }
2278
2279 #[test]
2280 fn count_tokens_keeps_requested_model() {
2281 let mut classifier = request(
2282 "You are a security monitor for autonomous AI coding agents.",
2283 false,
2284 json!([]),
2285 );
2286 assert!(
2287 apply_auto_review_model(&mut classifier, true, Some("grok-4.5"), "codex").is_none()
2288 );
2289 assert_eq!(classifier.model.as_deref(), Some("gpt-5.6-sol"));
2290 }
2291}