diff --git a/crates/libsy-llm-client/src/client.rs b/crates/libsy-llm-client/src/client.rs index f2d0301dd..1690a6e58 100644 --- a/crates/libsy-llm-client/src/client.rs +++ b/crates/libsy-llm-client/src/client.rs @@ -52,6 +52,31 @@ const NON_FORWARDABLE_HEADERS: &[&str] = &[ "expect", ]; +/// Phrases that mark a 400/422 as a rejected request *parameter*, not rejected +/// content. A server that refuses an unknown top-level field names it one of these +/// ways. The reasoning-control knobs are listed because operators configure them +/// per target most often. Matching stays parameter-specific on purpose: a bare +/// "unrecognized" would also match unrelated bodies such as "unrecognized model". +const PARAM_REJECT_PHRASES: &[&str] = &[ + "chat_template_kwargs", + "enable_thinking", + "unexpected keyword", + "unknown parameter", + "unknown field", + "unrecognized parameter", + "unrecognized field", + "unrecognized argument", + "unrecognized key", + "extra_forbidden", + "additional properties", + "extra fields not permitted", + "invalid parameter", +]; + +/// Statuses that can mean "unknown parameter". A 429 or 5xx is a transient fault +/// and must never be treated as a parameter reject. +const PARAM_REJECT_STATUSES: &[u16] = &[400, 422]; + const INITIAL_RETRY_DELAY: Duration = Duration::from_millis(250); const MAX_RETRY_BACKOFF: Duration = Duration::from_secs(2); const MAX_RETRY_AFTER: Duration = Duration::from_secs(60); @@ -261,10 +286,15 @@ impl TranslatingLlmClient { strip_unsigned_thinking_blocks(&mut body); } omit_configured_body_fields(&mut body, backend.omit_body_fields()); - merge_extra_body(&mut body, backend.extra_body()); + let mut injected_extra_body = merge_extra_body(&mut body, backend.extra_body()); // After the merge on purpose: the effort override must win over both the caller's // value and any `reasoning` default a target set through `extra_body`. apply_reasoning_effort(&mut body, backend); + // The effort override replaces any injected `reasoning_effort` with the backend + // setting, so a reject of that field must not blame `extra_body`. + if matches!(backend, Backend::OpenAiChat(_)) && backend.reasoning_effort().is_some() { + injected_extra_body.retain(|key| key != "reasoning_effort"); + } if matches!(backend, Backend::Anthropic(_)) { enable_anthropic_prompt_caching(&mut body); } @@ -276,8 +306,28 @@ impl TranslatingLlmClient { let url = endpoint.url(backend); record_gen_ai_request(&url, model, streaming); - self.send_with_retries(&url, backend, &body, metadata, model, streaming) + // A 400/422 naming a parameter this target injected is a deployment + // misconfiguration, not a transient upstream fault. `extra_body` is + // static, so retrying would fail the same way on every later request; + // 400/422 are not retried, so the error below is final. Only a reject + // that names an injected key converts: a reject naming any other field + // is the caller's, and blaming `extra_body` would misdirect the operator. + match self + .send_with_retries(&url, backend, &body, metadata, model, streaming) .await + { + Err(LlmClientError::UpstreamHttp { status, body }) + if !injected_extra_body.is_empty() && is_param_reject(status, &body) => + { + let named = injected_keys_named_by_upstream(&body, &injected_extra_body); + if named.is_empty() { + Err(LlmClientError::UpstreamHttp { status, body }) + } else { + Err(extra_body_rejected(model, &named)) + } + } + result => result, + } } // Sends the encoded body, retrying retryable failures within the backend's retry budget. @@ -1137,12 +1187,112 @@ fn apply_reasoning_effort(body: &mut Value, backend: &Backend) { } // Applies target defaults without overriding fields supplied by the caller. -fn merge_extra_body(body: &mut Value, extra_body: &BTreeMap) { +// Returns the keys it actually injected, so a parameter reject can name exactly +// those and leave the caller's own fields alone. +fn merge_extra_body(body: &mut Value, extra_body: &BTreeMap) -> Vec { let Value::Object(object) = body else { - return; + return Vec::new(); }; + let mut injected = Vec::new(); for (key, value) in extra_body { - object.entry(key.clone()).or_insert_with(|| value.clone()); + if !object.contains_key(key) { + object.insert(key.clone(), value.clone()); + injected.push(key.clone()); + } + } + injected +} + +/// Whether an upstream HTTP status/body pair is the upstream refusing a request +/// *parameter* rather than the request's content. +/// +/// Deliberately narrow: only a 400 or 422 whose body names a parameter qualifies. +/// A rate limit or server error is transient and must surface unchanged. A false +/// positive would turn a real upstream error into a configuration error, so the +/// phrase list stays specific. +fn is_param_reject(status: http::StatusCode, body: &str) -> bool { + if !PARAM_REJECT_STATUSES.contains(&status.as_u16()) { + return false; + } + let lowered = body.to_ascii_lowercase(); + PARAM_REJECT_PHRASES + .iter() + .any(|phrase| lowered.contains(phrase)) +} + +// Attribute only an explicit rejected field, never a key mentioned elsewhere in the body. +fn injected_keys_named_by_upstream(body: &str, injected: &[String]) -> Vec { + let parsed = serde_json::from_str::(body).unwrap_or_default(); + let error = parsed.get("error").unwrap_or(&parsed); + let parameter = error.get("param").and_then(Value::as_str); + let message = error.get("message").and_then(Value::as_str).unwrap_or(body); + let rejected: Vec<&str> = if let Some(parameter) = parameter { + vec![parameter] + } else if let Some(details) = parsed.get("detail").and_then(Value::as_array) { + details + .iter() + .filter_map(|detail| { + if !matches!( + detail.get("type").and_then(Value::as_str), + Some("extra_forbidden" | "value_error.extra") + ) { + return None; + } + let location = detail.get("loc")?.as_array()?; + match location.as_slice() { + [source, field] if source.as_str() == Some("body") => field.as_str(), + _ => None, + } + }) + .collect() + } else { + let lowered = message.to_ascii_lowercase(); + PARAM_REJECT_PHRASES + .iter() + .flat_map(|phrase| { + lowered + .match_indices(*phrase) + .filter_map(move |(start, _)| { + let offset = start + phrase.len(); + let suffix = message[offset..] + .trim_start_matches(|c: char| c.is_ascii_whitespace() || c == ':'); + let suffix = suffix.strip_prefix("argument ").unwrap_or(suffix); + if let Some(quote) = suffix + .chars() + .next() + .filter(|c| matches!(c, '\'' | '"' | '`')) + { + let quoted = &suffix[quote.len_utf8()..]; + return quoted.find(quote).map(|end| "ed[..end]); + } + let end = suffix + .find(|c: char| c.is_ascii_whitespace() || matches!(c, ',' | ';')) + .unwrap_or(suffix.len()); + (end > 0).then_some(&suffix[..end]) + }) + }) + .collect() + }; + injected + .iter() + .filter(|key| rejected.contains(&key.as_str())) + .cloned() + .collect() +} + +/// Builds the error returned when an upstream rejects a parameter this target +/// injected. Names exactly the injected keys the upstream named so the operator +/// can correct `extra_body` in the deployment TOML. The upstream body is left +/// out on purpose: it can quote the request back, and the key names are enough +/// to act on. +fn extra_body_rejected(model: &ModelId, injected: &[String]) -> LlmClientError { + LlmClientError::Configuration { + message: format!( + "upstream rejected a request parameter for model {model}. This target sets \ + extra_body keys [{}], which the model does not accept. Remove or correct them \ + in the target's extra_body.", + injected.join(", ") + ), } } @@ -3213,4 +3363,353 @@ mod tests { .await?; Ok(()) } + + #[test] + fn param_reject_classification_is_narrow() { + let named = r#"{"error":{"message":"unknown parameter: enable_thinking"}}"#; + + for (status, body, expected) in [ + (400, named, true), + (422, named, true), + (400, "extra fields not permitted", true), + // An unrelated 400 is a real request error, not a config problem. + (400, r#"{"error":{"message":"invalid api key"}}"#, false), + // Transient faults must never qualify: reporting a rate limit as a + // configuration error sends the operator to edit a TOML that is correct. + (429, named, false), + (503, named, false), + // Parameter-specific phrases only: "unrecognized model" names no parameter. + (400, r#"{"error":{"message":"unrecognized model"}}"#, false), + ( + 400, + r#"{"error":{"message":"unrecognized parameter: foo"}}"#, + true, + ), + ] { + let status = http::StatusCode::from_u16(status).expect("valid status"); + assert_eq!( + is_param_reject(status, body), + expected, + "HTTP {status} with body {body}" + ); + } + } + + #[test] + fn rejected_parameter_attribution_requires_an_exact_field() { + let injected = vec!["max_tokens".to_string()]; + for body in [ + r#"{"error":{"message":"unknown parameter: max_tokens_per_second"}}"#, + r#"{"error":{"message":"unknown parameter: other; request contained max_tokens"}}"#, + r#"{"error":{"message":"unknown parameter: 'max_tokens per_second'"}}"#, + r#"{"error":{"message":"unknown parameter: 'max_tokens,other'"}}"#, + r#"{"error":{"message":"unknown parameter: 'max_tokens"}}"#, + r#"{"error":{"message":"unknown parameter: max_tokens[0]"}}"#, + r#"{"error":{"message":"unknown parameter: max_tokens/per_second"}}"#, + r#"{"error":{"param":"other","message":"unknown parameter: max_tokens"}}"#, + ] { + assert!(injected_keys_named_by_upstream(body, &injected).is_empty()); + } + assert_eq!( + injected_keys_named_by_upstream( + r#"{"error":{"param":"max_tokens","message":"unknown parameter"}}"#, + &injected, + ), + injected, + ); + assert_eq!( + injected_keys_named_by_upstream( + r#"{"error":{"message":"unknown parameter: 'max_tokens'"}}"#, + &injected, + ), + injected, + ); + } + + #[tokio::test] + async fn structured_and_repeated_rejections_identify_injected_fields() + -> std::result::Result<(), Box> { + for (status, body) in [ + ( + 400, + json!({"error": {"message": "unknown field: caller_flag; unknown field: service_tier"}}), + ), + ( + 422, + json!({"detail": [{"type": "extra_forbidden", "loc": ["body", "service_tier"], "msg": "Extra inputs are not permitted"}]}), + ), + ( + 422, + json!({"detail": [{"type": "value_error.extra", "loc": ["body", "service_tier"], "msg": "extra fields not permitted"}]}), + ), + ] { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/chat/completions")) + .respond_with(ResponseTemplate::new(status).set_body_json(body.clone())) + .expect(1) + .mount(&server) + .await; + let client = TranslatingLlmClient::new(&chat_map_with_extra_body( + &format!("{}/v1", server.uri()), + BTreeMap::from([("service_tier".to_string(), json!("flex"))]), + ))?; + let error = client.call_rewrite_model_raw( + json!({"model": "client-facing", "messages": [{"role": "user", "content": "hi"}]}), + None, Some(&ModelId::new("gpt")), WireFormat::OpenAiChat, + ).await.err().ok_or("expected rejection")?; + let LlmClientError::Configuration { message } = error else { + return Err( + format!("expected configuration error for {body}, got {error:?}").into(), + ); + }; + assert!(message.contains("service_tier"), "{message}"); + assert!(!message.contains("caller_flag"), "{message}"); + } + Ok(()) + } + + #[test] + fn structured_rejections_ignore_nested_and_non_body_fields() { + let injected = vec!["service_tier".to_string()]; + for loc in [ + json!(["body", "options", "service_tier"]), + json!(["query", "service_tier"]), + ] { + let body = json!({"detail": [{"type": "extra_forbidden", "loc": loc, "msg": "Extra inputs are not permitted"}]}).to_string(); + assert!(injected_keys_named_by_upstream(&body, &injected).is_empty()); + } + } + + #[tokio::test] + async fn an_extra_body_reject_is_reported_as_a_configuration_error() + -> std::result::Result<(), Box> { + let reject = || { + ResponseTemplate::new(400).set_body_json(json!({ + "error": {"message": "unknown parameter: enable_thinking"} + })) + }; + let request = + || json!({"model": "client-facing", "messages": [{"role": "user", "content": "hi"}]}); + + // With an injected key, the deployment is at fault. `.expect(1)` pins that the + // request is sent once: extra_body is static, so retrying would fail identically + // on every later request too. + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/chat/completions")) + .respond_with(reject()) + .expect(1) + .mount(&server) + .await; + let extra_body = BTreeMap::from([("enable_thinking".to_string(), json!(false))]); + let client = TranslatingLlmClient::new(&chat_map_with_extra_body( + &format!("{}/v1", server.uri()), + extra_body, + ))?; + + let error = client + .call_rewrite_model_raw( + request(), + None, + Some(&ModelId::new("gpt")), + WireFormat::OpenAiChat, + ) + .await + .err() + .ok_or("expected the call to fail")?; + + let LlmClientError::Configuration { message } = &error else { + return Err(format!("expected a configuration error, got {error:?}").into()); + }; + assert!(message.contains("enable_thinking"), "{message}"); + assert!(message.contains("extra_body"), "{message}"); + // The upstream body is not echoed: it can quote the request back. + assert!(!message.contains("unknown parameter"), "{message}"); + drop(server); + + // With nothing injected, the rejected parameter is the caller's, so the upstream + // error passes through unchanged. + let bare = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/chat/completions")) + .respond_with(reject()) + .mount(&bare) + .await; + let client = TranslatingLlmClient::new(&chat_map_with_extra_body( + &format!("{}/v1", bare.uri()), + BTreeMap::new(), + ))?; + + let error = client + .call_rewrite_model_raw( + request(), + None, + Some(&ModelId::new("gpt")), + WireFormat::OpenAiChat, + ) + .await + .err() + .ok_or("expected the call to fail")?; + + let LlmClientError::UpstreamHttp { status, .. } = &error else { + return Err(format!("expected an upstream error, got {error:?}").into()); + }; + assert_eq!(*status, http::StatusCode::BAD_REQUEST); + Ok(()) + } + + #[tokio::test] + async fn a_reject_naming_an_uninjected_parameter_preserves_the_upstream_error() + -> std::result::Result<(), Box> { + // `extra_body` injects `service_tier`, but the upstream rejects the + // caller-supplied `max_tokens`. Converting would blame `service_tier` + // and lose the original error. + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/chat/completions")) + .respond_with(ResponseTemplate::new(400).set_body_json(json!({ + "error": {"message": "unknown parameter: max_tokens"} + }))) + .mount(&server) + .await; + let client = TranslatingLlmClient::new(&chat_map_with_extra_body( + &format!("{}/v1", server.uri()), + BTreeMap::from([("service_tier".to_string(), json!("flex"))]), + ))?; + + let error = client + .call_rewrite_model_raw( + json!({"model": "client-facing", "messages": [{"role": "user", "content": "hi"}]}), + None, + Some(&ModelId::new("gpt")), + WireFormat::OpenAiChat, + ) + .await + .err() + .ok_or("expected the call to fail")?; + + let LlmClientError::UpstreamHttp { status, body } = &error else { + return Err(format!("expected an upstream error, got {error:?}").into()); + }; + assert_eq!(*status, http::StatusCode::BAD_REQUEST); + assert!(body.contains("max_tokens"), "{body}"); + drop(server); + + // When the upstream does name an injected key, only that key is named: + // the other injected key is not blamed. + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/chat/completions")) + .respond_with(ResponseTemplate::new(400).set_body_json(json!({ + "error": {"message": "unknown parameter: enable_thinking"} + }))) + .mount(&server) + .await; + let client = TranslatingLlmClient::new(&chat_map_with_extra_body( + &format!("{}/v1", server.uri()), + BTreeMap::from([ + ("service_tier".to_string(), json!("flex")), + ("enable_thinking".to_string(), json!(false)), + ]), + ))?; + + let error = client + .call_rewrite_model_raw( + json!({"model": "client-facing", "messages": [{"role": "user", "content": "hi"}]}), + None, + Some(&ModelId::new("gpt")), + WireFormat::OpenAiChat, + ) + .await + .err() + .ok_or("expected the call to fail")?; + + let LlmClientError::Configuration { message } = &error else { + return Err(format!("expected a configuration error, got {error:?}").into()); + }; + assert!(message.contains("enable_thinking"), "{message}"); + assert!(!message.contains("service_tier"), "{message}"); + Ok(()) + } + + #[tokio::test] + async fn an_overridden_reasoning_effort_reject_is_not_blamed_on_extra_body() + -> std::result::Result<(), Box> { + // The backend `reasoning_effort` setting overrides the `extra_body` + // value, so a reject of `reasoning_effort` must not direct the operator + // to edit `extra_body`: the setting is what sends the rejected value. + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/chat/completions")) + .respond_with(ResponseTemplate::new(400).set_body_json(json!({ + "error": {"message": "unknown parameter: reasoning_effort"} + }))) + .mount(&server) + .await; + let mut backend = config(&format!("{}/v1", server.uri())); + backend.reasoning_effort = Some("max".to_string()); + backend.extra_body = BTreeMap::from([("reasoning_effort".to_string(), json!("low"))]); + let client = TranslatingLlmClient::new(&[ModelConfig::new( + "gpt", + Backend::OpenAiChat(backend), + None, + )])?; + + let error = client + .call_rewrite_model_raw( + json!({"model": "client-facing", "messages": [{"role": "user", "content": "hi"}]}), + None, + Some(&ModelId::new("gpt")), + WireFormat::OpenAiChat, + ) + .await + .err() + .ok_or("expected the call to fail")?; + + let LlmClientError::UpstreamHttp { status, body } = &error else { + return Err(format!("expected an upstream error, got {error:?}").into()); + }; + assert_eq!(*status, http::StatusCode::BAD_REQUEST); + assert!(body.contains("reasoning_effort"), "{body}"); + Ok(()) + } + #[tokio::test] + async fn responses_reasoning_effort_injection_remains_attributable() + -> std::result::Result<(), Box> { + // Responses writes reasoning.effort, leaving the injected top-level key intact. + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/responses")) + .respond_with(ResponseTemplate::new(400).set_body_json(json!({ + "error": {"message": "unknown parameter: reasoning_effort"} + }))) + .mount(&server) + .await; + let mut backend = config(&format!("{}/v1", server.uri())); + backend.reasoning_effort = Some("max".to_string()); + backend.extra_body = BTreeMap::from([("reasoning_effort".to_string(), json!("low"))]); + let client = TranslatingLlmClient::new(&[ModelConfig::new( + "gpt", + Backend::OpenAiResponses(backend), + None, + )])?; + + let error = client + .call_rewrite_model_raw( + json!({"model": "client-facing", "messages": [{"role": "user", "content": "hi"}]}), + None, + Some(&ModelId::new("gpt")), + WireFormat::OpenAiChat, + ) + .await + .err() + .ok_or("expected the call to fail")?; + + let LlmClientError::Configuration { message } = &error else { + return Err(format!("expected a configuration error, got {error:?}").into()); + }; + assert!(message.contains("reasoning_effort"), "{message}"); + Ok(()) + } }