mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-08 03:08:45 +00:00
refactor(rust): extract inference-messages crate (#44818)
Co-authored-by: Yujong Lee <yujong@berri.ai> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
785cb12f37
commit
d7f5b40ab4
38 changed files with 645 additions and 126 deletions
31
litellm-rust/Cargo.lock
generated
31
litellm-rust/Cargo.lock
generated
|
|
@ -3836,6 +3836,7 @@ dependencies = [
|
|||
"litellm-host-http",
|
||||
"litellm-http",
|
||||
"litellm-inference",
|
||||
"litellm-inference-messages",
|
||||
"litellm-inference-responses",
|
||||
"litellm-inference-transcription",
|
||||
"litellm-llms",
|
||||
|
|
@ -4041,6 +4042,34 @@ dependencies = [
|
|||
"wiremock",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "litellm-inference-messages"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"bytes",
|
||||
"futures-util",
|
||||
"litellm-auth",
|
||||
"litellm-auth-aws",
|
||||
"litellm-cache-memory",
|
||||
"litellm-cache-response",
|
||||
"litellm-core-utils",
|
||||
"litellm-host",
|
||||
"litellm-host-native",
|
||||
"litellm-http",
|
||||
"litellm-inference",
|
||||
"litellm-llms",
|
||||
"litellm-llms-types",
|
||||
"litellm-secrets",
|
||||
"litellm-tracing",
|
||||
"reqwest 0.12.28",
|
||||
"rstest",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"tokio",
|
||||
"tracing",
|
||||
"wiremock",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "litellm-inference-responses"
|
||||
version = "0.1.0"
|
||||
|
|
@ -4170,6 +4199,7 @@ dependencies = [
|
|||
"litellm-host-python",
|
||||
"litellm-http",
|
||||
"litellm-inference",
|
||||
"litellm-inference-messages",
|
||||
"litellm-inference-responses",
|
||||
"litellm-inference-transcription",
|
||||
"litellm-llms",
|
||||
|
|
@ -4223,6 +4253,7 @@ version = "0.1.0"
|
|||
dependencies = [
|
||||
"litellm-config",
|
||||
"litellm-inference",
|
||||
"litellm-inference-messages",
|
||||
"rstest",
|
||||
]
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@ litellm-storage-clickhouse = { path = "crates/storage-clickhouse" }
|
|||
litellm-inference = { path = "crates/inference" }
|
||||
litellm-inference-transcription = { path = "crates/inference-transcription" }
|
||||
litellm-inference-responses = { path = "crates/inference-responses" }
|
||||
litellm-inference-messages = { path = "crates/inference-messages" }
|
||||
litellm-gateway-mcp = { path = "crates/gateway-mcp" }
|
||||
litellm-gateway = { path = "crates/gateway" }
|
||||
litellm-gateway-inference = { path = "crates/gateway-inference" }
|
||||
|
|
|
|||
|
|
@ -15,6 +15,7 @@ litellm-gateway-auth.workspace = true
|
|||
litellm-inference.workspace = true
|
||||
litellm-inference-transcription.workspace = true
|
||||
litellm-inference-responses.workspace = true
|
||||
litellm-inference-messages.workspace = true
|
||||
litellm-host-http.workspace = true
|
||||
litellm-host.workspace = true
|
||||
litellm-http.workspace = true
|
||||
|
|
|
|||
|
|
@ -17,9 +17,9 @@ use std::sync::Arc;
|
|||
use axum::{Router, routing::post};
|
||||
use litellm_http::{ClientVariant, HttpClientConfig, media::UrlPolicy};
|
||||
use litellm_inference::{
|
||||
chat_completions::ChatCompletionsRoute, messages::MessagesRoute, ocr::OcrRoute,
|
||||
resources::CoreResources,
|
||||
chat_completions::ChatCompletionsRoute, ocr::OcrRoute, resources::CoreResources,
|
||||
};
|
||||
use litellm_inference_messages::MessagesRoute;
|
||||
use litellm_inference_responses::ResponsesRoute;
|
||||
use litellm_inference_transcription::AudioTranscriptionRoute;
|
||||
use litellm_llms::base_llm::ocr::{handler::OcrClient, settings::OcrSettings};
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ use axum::{
|
|||
response::{IntoResponse, Response},
|
||||
};
|
||||
use litellm_host_http::Sse;
|
||||
use litellm_inference::messages::{MessagesCall, messages_body, route::Messages};
|
||||
use litellm_inference_messages::{MessagesCall, messages_body, route::Messages};
|
||||
use litellm_llms_types::headers::{ProviderSpecificHeader, ProviderSpecificHeaders};
|
||||
use serde_json::{Map, Value};
|
||||
|
||||
|
|
|
|||
33
litellm-rust/crates/inference-messages/Cargo.toml
Normal file
33
litellm-rust/crates/inference-messages/Cargo.toml
Normal file
|
|
@ -0,0 +1,33 @@
|
|||
[package]
|
||||
name = "litellm-inference-messages"
|
||||
version = "0.1.0"
|
||||
edition.workspace = true
|
||||
license.workspace = true
|
||||
repository.workspace = true
|
||||
|
||||
[dependencies]
|
||||
bytes.workspace = true
|
||||
futures-util.workspace = true
|
||||
litellm-auth = { workspace = true, features = ["aws", "azure", "gcp"] }
|
||||
litellm-auth-aws.workspace = true
|
||||
litellm-cache-response.workspace = true
|
||||
litellm-core-utils.workspace = true
|
||||
litellm-host.workspace = true
|
||||
litellm-http.workspace = true
|
||||
litellm-inference.workspace = true
|
||||
litellm-llms.workspace = true
|
||||
litellm-llms-types.workspace = true
|
||||
litellm-secrets.workspace = true
|
||||
litellm-tracing.workspace = true
|
||||
reqwest.workspace = true
|
||||
serde.workspace = true
|
||||
serde_json = { workspace = true, features = ["preserve_order"] }
|
||||
tracing.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
litellm-cache-memory.workspace = true
|
||||
litellm-host-native.workspace = true
|
||||
litellm-inference = { workspace = true, features = ["test-support"] }
|
||||
rstest.workspace = true
|
||||
tokio.workspace = true
|
||||
wiremock.workspace = true
|
||||
|
|
@ -67,7 +67,7 @@ mod tests {
|
|||
use rstest::rstest;
|
||||
|
||||
use super::{MessagesProvider, messages_provider, string_headers, truncate_error_body};
|
||||
use crate::messages::Error;
|
||||
use crate::Error;
|
||||
use litellm_core_utils::get_llm_provider_logic::LlmProviders;
|
||||
|
||||
#[rstest]
|
||||
4
litellm-rust/crates/inference-messages/src/constants.rs
Normal file
4
litellm-rust/crates/inference-messages/src/constants.rs
Normal file
|
|
@ -0,0 +1,4 @@
|
|||
/// Full-request timeout ceiling for Anthropic Messages provider calls, in
|
||||
/// seconds. Mirrors the Python Anthropic Messages default. The per-request
|
||||
/// timeout from the caller still overrides this on the request builder.
|
||||
pub(crate) const MESSAGES_TIMEOUT_SECS: u64 = 600;
|
||||
|
|
@ -19,7 +19,8 @@ use super::{
|
|||
Error, MessagesCallResponse, MessagesRoute, common_utils::truncate_error_body,
|
||||
prepare::ProviderMessagesRequest,
|
||||
};
|
||||
use crate::{constants::MESSAGES_TIMEOUT_SECS, context::CallContext, outbound::outbound_request};
|
||||
use crate::constants::MESSAGES_TIMEOUT_SECS;
|
||||
use litellm_inference::{context::CallContext, outbound::outbound_request};
|
||||
|
||||
pub(super) struct ProviderCall {
|
||||
pub identity: ProviderIdentity,
|
||||
|
|
@ -166,7 +167,9 @@ async fn send(
|
|||
body,
|
||||
Some(timeout.unwrap_or(Duration::from_secs(MESSAGES_TIMEOUT_SECS))),
|
||||
)?;
|
||||
crate::outbound::send(request, http).await.map_err(network)
|
||||
litellm_inference::outbound::send(request, http)
|
||||
.await
|
||||
.map_err(network)
|
||||
}
|
||||
|
||||
async fn provider_error(response: reqwest::Response) -> Error {
|
||||
|
|
@ -1,4 +1,5 @@
|
|||
mod common_utils;
|
||||
mod constants;
|
||||
mod handler;
|
||||
mod prepare;
|
||||
pub mod route;
|
||||
|
|
@ -8,11 +9,11 @@ use futures_util::FutureExt;
|
|||
use litellm_auth::AuthServices;
|
||||
use litellm_host::interceptors::{ExecutionFacts, Interceptors, ResultSource};
|
||||
|
||||
use crate::{caching::CallCache, context::CallContext};
|
||||
use litellm_inference::{caching::CallCache, context::CallContext};
|
||||
use litellm_secrets::source::SecretSource;
|
||||
use std::sync::Arc;
|
||||
|
||||
pub use crate::error::RouteError as Error;
|
||||
pub use litellm_inference::RouteError as Error;
|
||||
pub use types::{MessagesCall, MessagesCallResponse, MessagesShaping, messages_body};
|
||||
|
||||
#[derive(Clone)]
|
||||
|
|
@ -49,7 +50,7 @@ impl MessagesRoute {
|
|||
&self,
|
||||
call: MessagesCall,
|
||||
interceptors: &impl litellm_host::interceptors::Interceptors<Error>,
|
||||
options: impl Into<crate::CallOptions>,
|
||||
options: impl Into<litellm_inference::CallOptions>,
|
||||
) -> Result<MessagesCallResponse, Error> {
|
||||
let context = CallContext::new(interceptors, options.into());
|
||||
litellm_host::lifecycle::observe_call(context.observers.clone(), self.run(call, context))
|
||||
|
|
@ -69,9 +70,12 @@ impl MessagesRoute {
|
|||
call: MessagesCall,
|
||||
context: CallContext<'_, impl Interceptors<Error>>,
|
||||
) -> Result<MessagesCallResponse, Error> {
|
||||
crate::diagnostic::call(async {
|
||||
litellm_inference::diagnostic::call(async {
|
||||
let prepared = prepare::prepare(call, self.secrets.as_ref()).await?;
|
||||
crate::diagnostic::provider(&prepared.body.model, prepared.provider.as_str());
|
||||
litellm_inference::diagnostic::provider(
|
||||
&prepared.body.model,
|
||||
prepared.provider.as_str(),
|
||||
);
|
||||
let request = self.prepare_outbound(prepared, &context).boxed().await?;
|
||||
let cache = CallCache::<route::Messages>::from_wire(
|
||||
self.cache.as_ref().filter(|_| request.cacheable()),
|
||||
|
|
@ -17,7 +17,7 @@ use super::{
|
|||
common_utils::{MessagesProvider, messages_provider, string_headers},
|
||||
types::invalid_request,
|
||||
};
|
||||
use crate::provider::resolve_llm_provider;
|
||||
use litellm_inference::provider::resolve_llm_provider;
|
||||
|
||||
struct ResolvedProvider {
|
||||
model: String,
|
||||
|
|
@ -148,7 +148,7 @@ mod tests {
|
|||
use serde_json::{Map, Value, json};
|
||||
|
||||
use super::*;
|
||||
use crate::messages::MessagesShaping;
|
||||
use crate::MessagesShaping;
|
||||
|
||||
#[fixture]
|
||||
fn shaping() -> MessagesShaping {
|
||||
|
|
@ -33,9 +33,9 @@ impl super::MessagesRoute {
|
|||
pub fn machine(
|
||||
self,
|
||||
request: super::MessagesCall,
|
||||
options: impl Into<crate::CallOptions>,
|
||||
options: impl Into<litellm_inference::CallOptions>,
|
||||
) -> MessagesMachine {
|
||||
let crate::CallOptions {
|
||||
let litellm_inference::CallOptions {
|
||||
cache: cache_options,
|
||||
observers,
|
||||
} = options.into();
|
||||
|
|
@ -43,9 +43,9 @@ impl super::MessagesRoute {
|
|||
request,
|
||||
observers,
|
||||
move |call, _, interceptors, observers| async move {
|
||||
let context = crate::context::CallContext::new(
|
||||
let context = litellm_inference::context::CallContext::new(
|
||||
&interceptors,
|
||||
crate::CallOptions {
|
||||
litellm_inference::CallOptions {
|
||||
cache: cache_options,
|
||||
observers,
|
||||
},
|
||||
|
|
@ -56,11 +56,11 @@ impl super::MessagesRoute {
|
|||
}
|
||||
}
|
||||
|
||||
impl crate::caching::Cachable for Messages {
|
||||
impl litellm_inference::caching::Cachable for Messages {
|
||||
const SURFACE: &'static str = "messages";
|
||||
}
|
||||
|
||||
impl crate::caching::StreamCachable for Messages {
|
||||
impl litellm_inference::caching::StreamCachable for Messages {
|
||||
const TERMINAL_EVENT: &'static str = "message_stop";
|
||||
|
||||
fn replay(data: bytes::Bytes) -> Option<litellm_host::call::OutputOf<Self>> {
|
||||
243
litellm-rust/crates/inference-messages/tests/caching.rs
Normal file
243
litellm-rust/crates/inference-messages/tests/caching.rs
Normal file
|
|
@ -0,0 +1,243 @@
|
|||
mod support;
|
||||
|
||||
use std::{
|
||||
sync::{
|
||||
Arc,
|
||||
atomic::{AtomicUsize, Ordering},
|
||||
},
|
||||
time::Duration,
|
||||
};
|
||||
|
||||
use litellm_cache_memory::InMemoryCache;
|
||||
use litellm_cache_response::{
|
||||
CacheOptions, CacheScope, ResponseCache, ResponseCacheConfig, ResponseCacheService, ScopedCache,
|
||||
};
|
||||
use litellm_host::interceptors::{
|
||||
ExecutionFacts, Interceptors, ProviderIdentity, RawResponse, RequestContext, ResultSource,
|
||||
WireRequest,
|
||||
};
|
||||
use litellm_inference::{
|
||||
RouteError,
|
||||
caching::{CacheRequest, execute_unary},
|
||||
};
|
||||
use rstest::{fixture, rstest};
|
||||
use serde_json::{Value, json};
|
||||
|
||||
#[fixture]
|
||||
fn cache() -> Arc<dyn ResponseCacheService> {
|
||||
Arc::new(
|
||||
ResponseCache::new(Arc::new(InMemoryCache::new(
|
||||
Some(100),
|
||||
Some(Duration::from_secs(60)),
|
||||
)))
|
||||
.with_config(ResponseCacheConfig {
|
||||
namespace: "test".into(),
|
||||
max_entry_bytes: 4096,
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::system("system", json!("answer ALPHA"), json!("answer BETA"))]
|
||||
#[case::stop_sequences("stop_sequences", json!(["STOP"]), json!(["END"]))]
|
||||
#[case::top_k("top_k", json!(5), json!(10))]
|
||||
#[case::tools("tools", json!([{"name":"a","input_schema":{"type":"object"}}]), json!([{"name":"b","input_schema":{"type":"object"}}]))]
|
||||
#[case::tool_choice("tool_choice", json!({"type":"auto"}), json!({"type":"none"}))]
|
||||
#[tokio::test]
|
||||
async fn messages_cache_identity_includes_provider_native_parameters(
|
||||
cache: Arc<dyn ResponseCacheService>,
|
||||
#[case] field: &str,
|
||||
#[case] original: Value,
|
||||
#[case] changed: Value,
|
||||
) {
|
||||
use litellm_inference_messages::route::Messages;
|
||||
use litellm_llms_types::formats::messages::MessagesResponse;
|
||||
|
||||
let calls = AtomicUsize::new(0);
|
||||
for (value, expected_call) in [(original.clone(), 0), (changed, 1), (original, 0)] {
|
||||
let response =
|
||||
execute_unary::<Messages, _, _>(
|
||||
CacheRequest::from_wire(
|
||||
ProviderIdentity {
|
||||
model: "test".into(),
|
||||
provider: "anthropic".into(),
|
||||
},
|
||||
Some(&WireRequest {
|
||||
url: "https://example.test/v1/messages".into(),
|
||||
headers: vec![],
|
||||
body: json!({
|
||||
"model":"test", "messages":[{"role":"user","content":"hello"}],
|
||||
"max_tokens":32, (field):value
|
||||
}),
|
||||
}),
|
||||
),
|
||||
Some(cache.clone()),
|
||||
Some(CacheOptions::new(CacheScope::Shared)),
|
||||
&(),
|
||||
None,
|
||||
|| async {
|
||||
let call = calls.fetch_add(1, Ordering::SeqCst);
|
||||
Ok(Box::new(serde_json::from_value::<MessagesResponse>(json!({
|
||||
"id":call.to_string(), "type":"message", "role":"assistant", "model":"test",
|
||||
"content":[{"type":"text","text":format!("answer {call}")}],
|
||||
"stop_reason":"end_turn", "stop_sequence":null
|
||||
})).unwrap()))
|
||||
},
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(response.id, expected_call.to_string());
|
||||
assert_eq!(
|
||||
response.content[0]["text"],
|
||||
format!("answer {expected_call}")
|
||||
);
|
||||
}
|
||||
assert_eq!(calls.load(Ordering::SeqCst), 2);
|
||||
}
|
||||
|
||||
struct ChangingSecrets {
|
||||
revision: AtomicUsize,
|
||||
endpoints: [String; 2],
|
||||
change_credentials: bool,
|
||||
}
|
||||
|
||||
impl litellm_secrets::source::SecretSource for ChangingSecrets {
|
||||
fn get_secret_str<'a>(
|
||||
&'a self,
|
||||
name: &'a str,
|
||||
) -> futures_util::future::BoxFuture<
|
||||
'a,
|
||||
Result<Option<litellm_secrets::SecretValue>, litellm_secrets::Error>,
|
||||
> {
|
||||
Box::pin(async move {
|
||||
let revision = self.revision.load(Ordering::SeqCst);
|
||||
let value = if name.ends_with("_API_KEY") {
|
||||
Some(format!(
|
||||
"key-{}",
|
||||
if self.change_credentials { revision } else { 0 }
|
||||
))
|
||||
} else if name.ends_with("_API_BASE") {
|
||||
Some(self.endpoints[revision].clone())
|
||||
} else {
|
||||
None
|
||||
};
|
||||
Ok(value.map(litellm_secrets::SecretValue::new))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct ChangingHooks {
|
||||
calls: AtomicUsize,
|
||||
rewrite: bool,
|
||||
facts: std::sync::Mutex<Vec<ExecutionFacts>>,
|
||||
}
|
||||
|
||||
impl Interceptors<RouteError> for ChangingHooks {
|
||||
async fn before_provider_request(
|
||||
&self,
|
||||
mut wire: WireRequest,
|
||||
_: RequestContext,
|
||||
) -> Result<WireRequest, RouteError> {
|
||||
let call = self.calls.fetch_add(1, Ordering::SeqCst);
|
||||
if self.rewrite {
|
||||
wire.body["temperature"] = json!(if call < 2 { 0.1 } else { 0.8 });
|
||||
}
|
||||
Ok(wire)
|
||||
}
|
||||
|
||||
async fn after_provider_response(&self, _: RawResponse) -> Result<(), RouteError> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn result_ready(&self, facts: ExecutionFacts) -> Result<(), RouteError> {
|
||||
self.facts.lock().unwrap().push(facts);
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::credentials("credentials")]
|
||||
#[case::endpoint("endpoint")]
|
||||
#[case::callback("callback")]
|
||||
#[tokio::test]
|
||||
async fn messages_cache_identity_follows_resolved_configuration_and_request_callbacks(
|
||||
cache: Arc<dyn ResponseCacheService>,
|
||||
#[case] change: &str,
|
||||
) {
|
||||
use litellm_inference_messages::MessagesCall;
|
||||
use wiremock::{Mock, MockServer, ResponseTemplate, matchers::method};
|
||||
|
||||
let first = MockServer::start().await;
|
||||
let second = MockServer::start().await;
|
||||
let response = json!({"id":"message-test", "type":"message", "role":"assistant", "model":"test",
|
||||
"content":[{"type":"text", "text":"answer"}], "stop_reason":"end_turn", "stop_sequence":null,
|
||||
"usage":{"input_tokens":3,"output_tokens":2}});
|
||||
Mock::given(method("POST"))
|
||||
.respond_with(ResponseTemplate::new(200).set_body_json(response.clone()))
|
||||
.expect(if change == "endpoint" { 1 } else { 2 })
|
||||
.mount(&first)
|
||||
.await;
|
||||
Mock::given(method("POST"))
|
||||
.respond_with(ResponseTemplate::new(200).set_body_json(response))
|
||||
.expect(if change == "endpoint" { 1 } else { 0 })
|
||||
.mount(&second)
|
||||
.await;
|
||||
let secrets = Arc::new(ChangingSecrets {
|
||||
revision: AtomicUsize::new(0),
|
||||
endpoints: [
|
||||
first.uri(),
|
||||
if change == "endpoint" {
|
||||
second.uri()
|
||||
} else {
|
||||
first.uri()
|
||||
},
|
||||
],
|
||||
change_credentials: change == "credentials",
|
||||
});
|
||||
let hooks = ChangingHooks {
|
||||
rewrite: change == "callback",
|
||||
..Default::default()
|
||||
};
|
||||
for call in 0..4 {
|
||||
secrets
|
||||
.revision
|
||||
.store(usize::from(call >= 2), Ordering::SeqCst);
|
||||
let cache = ScopedCache::new(cache.clone(), CacheScope::Shared);
|
||||
support::messages_route(secrets.clone()).with_cache(cache).execute(MessagesCall {
|
||||
body: serde_json::from_value(json!({"model":"anthropic/cache-test-model","messages":[{"role":"user","content":"hello"}],"max_tokens":32})).unwrap(),
|
||||
api_key:None,api_base:None,custom_llm_provider:None,extra_headers:None,provider_specific_header:None,timeout:None,shaping:Default::default(),
|
||||
}, &hooks, None).await.unwrap();
|
||||
}
|
||||
assert_eq!(hooks.calls.load(Ordering::SeqCst), 4);
|
||||
{
|
||||
let facts = hooks.facts.lock().unwrap();
|
||||
assert_eq!(facts[0].source, ResultSource::Provider);
|
||||
assert_eq!(facts[2].source, ResultSource::Provider);
|
||||
let (ResultSource::Cache { key: first_key }, ResultSource::Cache { key: second_key }) =
|
||||
(&facts[1].source, &facts[3].source)
|
||||
else {
|
||||
panic!("unchanged effective requests must hit the cache");
|
||||
};
|
||||
assert_ne!(first_key, second_key);
|
||||
}
|
||||
let requests = first.received_requests().await.unwrap();
|
||||
if change == "credentials" {
|
||||
assert_ne!(
|
||||
requests[0].headers["x-api-key"],
|
||||
requests[1].headers["x-api-key"]
|
||||
);
|
||||
}
|
||||
if change == "callback" {
|
||||
assert_eq!(
|
||||
serde_json::from_slice::<Value>(&requests[0].body).unwrap()["temperature"],
|
||||
0.1
|
||||
);
|
||||
assert_eq!(
|
||||
serde_json::from_slice::<Value>(&requests[1].body).unwrap()["temperature"],
|
||||
0.8
|
||||
);
|
||||
}
|
||||
first.verify().await;
|
||||
second.verify().await;
|
||||
}
|
||||
|
|
@ -5,8 +5,8 @@ use litellm_host::{
|
|||
interceptors::{ExecutionFacts, RequestContext, ResultSource, WireRequest},
|
||||
lifecycle::CallEvent,
|
||||
};
|
||||
use litellm_inference::messages::{MessagesCallResponse, route::Messages};
|
||||
use litellm_inference::test_support::{RecordingSecrets, no_secrets};
|
||||
use litellm_inference_messages::{MessagesCallResponse, route::Messages};
|
||||
use litellm_llms::base_llm::messages::context::MessagesModelCapabilities as AnthropicModelCapabilities;
|
||||
use rstest::rstest;
|
||||
|
||||
|
|
@ -4,11 +4,11 @@ use std::{
|
|||
};
|
||||
|
||||
use litellm_http::{HttpSettings, Resolution};
|
||||
use litellm_inference::messages::{
|
||||
use litellm_inference::test_support::RecordingSecrets;
|
||||
use litellm_inference_messages::{
|
||||
Error, MessagesCall, MessagesShaping,
|
||||
route::{Messages, MessagesMachine, MessagesOutput},
|
||||
};
|
||||
use litellm_inference::test_support::RecordingSecrets;
|
||||
use litellm_llms_types::formats::messages::{MessagesRequest, MessagesResponse};
|
||||
use litellm_secrets::source::SecretSource;
|
||||
use rstest::fixture;
|
||||
|
|
@ -3,10 +3,10 @@ use litellm_host::{
|
|||
lifecycle::ExecutionEvent,
|
||||
};
|
||||
use litellm_http::transport::Error as TransportError;
|
||||
use litellm_inference::messages::{MessagesCallResponse, messages_body};
|
||||
use litellm_inference::test_support::{
|
||||
RecordingSecrets, http_config, no_secrets, provider_http, resources,
|
||||
};
|
||||
use litellm_inference_messages::{MessagesCallResponse, messages_body};
|
||||
use rstest::rstest;
|
||||
|
||||
use super::*;
|
||||
|
|
@ -272,7 +272,7 @@ async fn the_facade_sends_through_the_injected_http_pool_configuration(call: Mes
|
|||
};
|
||||
|
||||
let resources = resources();
|
||||
let response = litellm_inference::messages::MessagesRoute::new(
|
||||
let response = litellm_inference_messages::MessagesRoute::new(
|
||||
provider_http(&resources, &Resolution::from(&settings).config),
|
||||
resources.auth,
|
||||
no_secrets(),
|
||||
|
|
@ -351,7 +351,7 @@ async fn route_uses_injected_dependencies_and_optional_cache(
|
|||
) {
|
||||
use litellm_cache_memory::InMemoryCache;
|
||||
use litellm_cache_response::{CacheScope, ResponseCache, ScopedCache};
|
||||
use litellm_inference::messages::MessagesRoute;
|
||||
use litellm_inference_messages::MessagesRoute;
|
||||
|
||||
let upstream = upstream([message_response(), message_response()]).await;
|
||||
let resources = resources();
|
||||
|
|
@ -5,11 +5,11 @@ use std::{
|
|||
|
||||
use bytes::Bytes;
|
||||
use futures_util::{StreamExt, TryStreamExt};
|
||||
use litellm_inference::messages::{
|
||||
use litellm_inference::test_support::{RecordingSecrets, no_secrets};
|
||||
use litellm_inference_messages::{
|
||||
MessagesCallResponse,
|
||||
route::{Messages, MessagesStreamHead},
|
||||
};
|
||||
use litellm_inference::test_support::{RecordingSecrets, no_secrets};
|
||||
use litellm_tracing::{Logger, Metadata, Record, Sink};
|
||||
use rstest::rstest;
|
||||
use tokio::{
|
||||
|
|
@ -36,7 +36,7 @@ struct TraceSink(mpsc::Sender<(String, Value)>);
|
|||
|
||||
impl Sink for TraceSink {
|
||||
fn enabled(&self, metadata: &Metadata<'_>) -> bool {
|
||||
metadata.target().starts_with("litellm_inference::messages")
|
||||
metadata.target().starts_with("litellm_inference_messages")
|
||||
}
|
||||
|
||||
fn emit(&self, record: &Record) {
|
||||
282
litellm-rust/crates/inference-messages/tests/support/mod.rs
Normal file
282
litellm-rust/crates/inference-messages/tests/support/mod.rs
Normal file
|
|
@ -0,0 +1,282 @@
|
|||
//! Shared fixtures for route integration tests: a scripted upstream and a recording
|
||||
//! secret source.
|
||||
|
||||
#![allow(dead_code)] // each test binary compiles this module on its own and uses a different subset
|
||||
|
||||
use std::{
|
||||
ops::ControlFlow,
|
||||
sync::{Arc, Mutex},
|
||||
};
|
||||
|
||||
use litellm_inference::test_support::{http_config, provider_http, resources};
|
||||
use litellm_secrets::source::SecretSource;
|
||||
use serde_json::Value;
|
||||
use wiremock::{Mock, MockServer, Request, ResponseTemplate, matchers::any};
|
||||
|
||||
/// A port nothing listens on, for calls that must fail before any request is sent.
|
||||
pub const UNREACHABLE_BASE: &str = "http://127.0.0.1:1";
|
||||
|
||||
pub fn messages_route(secrets: Arc<dyn SecretSource>) -> litellm_inference_messages::MessagesRoute {
|
||||
let resources = resources();
|
||||
litellm_inference_messages::MessagesRoute::new(
|
||||
provider_http(&resources, &http_config()),
|
||||
resources.auth,
|
||||
secrets,
|
||||
)
|
||||
}
|
||||
|
||||
/// Starts an upstream that answers its n-th request with the n-th response and 404s after.
|
||||
pub async fn upstream(responses: impl IntoIterator<Item = ResponseTemplate>) -> MockServer {
|
||||
let server = MockServer::start().await;
|
||||
respond_in_order(&server, responses).await;
|
||||
server
|
||||
}
|
||||
|
||||
/// Scripts responses on a started server, for responses that need its address.
|
||||
pub async fn respond_in_order(
|
||||
server: &MockServer,
|
||||
responses: impl IntoIterator<Item = ResponseTemplate>,
|
||||
) {
|
||||
for response in responses {
|
||||
Mock::given(any())
|
||||
.respond_with(response)
|
||||
.up_to_n_times(1)
|
||||
.mount(server)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn received(server: &MockServer) -> Vec<Request> {
|
||||
server
|
||||
.received_requests()
|
||||
.await
|
||||
.expect("request recording is on")
|
||||
}
|
||||
|
||||
pub async fn only_request(server: &MockServer) -> Request {
|
||||
let [request] = <[Request; 1]>::try_from(received(server).await)
|
||||
.unwrap_or_else(|requests| panic!("expected one request, got {}", requests.len()));
|
||||
request
|
||||
}
|
||||
|
||||
pub fn json_response(body: Value) -> ResponseTemplate {
|
||||
ResponseTemplate::new(200).set_body_json(body)
|
||||
}
|
||||
|
||||
pub fn status_response(status: u16, body: Value) -> ResponseTemplate {
|
||||
ResponseTemplate::new(status).set_body_json(body)
|
||||
}
|
||||
|
||||
pub trait ReceivedRequest {
|
||||
fn header(&self, name: &str) -> Option<&str>;
|
||||
fn header_values(&self, name: &str) -> Vec<&str>;
|
||||
fn json(&self) -> Value;
|
||||
fn body_text(&self) -> String;
|
||||
/// The path and query, as the request line carried them.
|
||||
fn target(&self) -> String;
|
||||
fn query(&self, name: &str) -> Option<String>;
|
||||
}
|
||||
|
||||
impl ReceivedRequest for Request {
|
||||
fn header(&self, name: &str) -> Option<&str> {
|
||||
self.headers.get(name).and_then(|value| value.to_str().ok())
|
||||
}
|
||||
|
||||
fn header_values(&self, name: &str) -> Vec<&str> {
|
||||
self.headers
|
||||
.get_all(name)
|
||||
.iter()
|
||||
.filter_map(|value| value.to_str().ok())
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn json(&self) -> Value {
|
||||
serde_json::from_slice(&self.body).expect("request body is json")
|
||||
}
|
||||
|
||||
fn body_text(&self) -> String {
|
||||
String::from_utf8_lossy(&self.body).into_owned()
|
||||
}
|
||||
|
||||
fn target(&self) -> String {
|
||||
match self.url.query() {
|
||||
Some(query) => format!("{}?{query}", self.url.path()),
|
||||
None => self.url.path().to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
fn query(&self, name: &str) -> Option<String> {
|
||||
self.url
|
||||
.query_pairs()
|
||||
.find_map(|(key, value)| (key == name).then(|| value.into_owned()))
|
||||
}
|
||||
}
|
||||
|
||||
pub struct RecordingCall<P: litellm_host::protocol::Protocol> {
|
||||
pub request: Mutex<Option<P::Request>>,
|
||||
pub events: Arc<CallEvents>,
|
||||
pub chunks: Mutex<Vec<P::Chunk>>,
|
||||
pub head: Mutex<Option<P::StreamHead>>,
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
pub struct CallEvents(pub Observations);
|
||||
pub struct Observations {
|
||||
pub sender: litellm_host::observation::ObservationSender,
|
||||
receiver: Mutex<tokio::sync::mpsc::Receiver<litellm_host::lifecycle::CallEvent>>,
|
||||
recorded: Mutex<Vec<litellm_host::lifecycle::CallEvent>>,
|
||||
}
|
||||
|
||||
impl Default for Observations {
|
||||
fn default() -> Self {
|
||||
let (sender, receiver) = litellm_host::observation::observation_channel(
|
||||
std::num::NonZeroUsize::new(128).unwrap(),
|
||||
);
|
||||
Self {
|
||||
sender,
|
||||
receiver: Mutex::new(receiver),
|
||||
recorded: Mutex::new(Vec::new()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Observations {
|
||||
pub fn lock(
|
||||
&self,
|
||||
) -> std::sync::LockResult<std::sync::MutexGuard<'_, Vec<litellm_host::lifecycle::CallEvent>>>
|
||||
{
|
||||
let mut events = self.recorded.lock()?;
|
||||
let mut receiver = self.receiver.lock().unwrap();
|
||||
while let Ok(event) = receiver.try_recv() {
|
||||
events.push(event);
|
||||
}
|
||||
Ok(events)
|
||||
}
|
||||
}
|
||||
|
||||
impl litellm_host::lifecycle::CallObserver for CallEvents {
|
||||
fn observe(&self, event: litellm_host::lifecycle::CallEvent) {
|
||||
self.0.sender.emit(event);
|
||||
}
|
||||
}
|
||||
|
||||
impl<P: litellm_host::protocol::Protocol> RecordingCall<P> {
|
||||
pub fn new(request: P::Request) -> Self {
|
||||
Self {
|
||||
request: Mutex::new(Some(request)),
|
||||
events: Arc::new(CallEvents::default()),
|
||||
chunks: Mutex::new(Vec::new()),
|
||||
head: Mutex::new(None),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<P: litellm_host::protocol::Protocol> litellm_host::interceptors::Interceptors<P::Error>
|
||||
for RecordingCall<P>
|
||||
{
|
||||
async fn before_provider_request(
|
||||
&self,
|
||||
wire: litellm_host::interceptors::WireRequest,
|
||||
_: litellm_host::interceptors::RequestContext,
|
||||
) -> Result<litellm_host::interceptors::WireRequest, P::Error> {
|
||||
Ok(litellm_host::interceptors::WireRequest {
|
||||
headers: wire
|
||||
.headers
|
||||
.into_iter()
|
||||
.chain([("x-hook".into(), "called".into())])
|
||||
.collect(),
|
||||
..wire
|
||||
})
|
||||
}
|
||||
|
||||
async fn after_provider_response(
|
||||
&self,
|
||||
_: litellm_host::interceptors::RawResponse,
|
||||
) -> Result<(), P::Error> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl<P> RecordingCall<P>
|
||||
where
|
||||
P: litellm_host::protocol::Protocol<HostCall = std::convert::Infallible>,
|
||||
P::Error: From<litellm_host::machine::MachineFault>,
|
||||
{
|
||||
pub fn request(&self) -> Result<P::Request, P::Error> {
|
||||
self.request
|
||||
.lock()
|
||||
.unwrap()
|
||||
.take()
|
||||
.ok_or_else(|| litellm_host::machine::MachineFault::Abandoned.into())
|
||||
}
|
||||
pub fn runtime(&self) -> litellm_host_native::in_process::Host<'_, (), Self, Self> {
|
||||
litellm_host_native::in_process::Host {
|
||||
services: &(),
|
||||
interceptors: self,
|
||||
stream: self,
|
||||
observers: Some(&self.events.0.sender),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<P> litellm_host_native::in_process::StreamConsumer<P> for RecordingCall<P>
|
||||
where
|
||||
P: litellm_host::protocol::Protocol<HostCall = std::convert::Infallible>,
|
||||
P::Error: From<litellm_host::machine::MachineFault>,
|
||||
{
|
||||
async fn open_stream(&self, head: P::StreamHead) -> Result<ControlFlow<()>, P::Error> {
|
||||
*self.head.lock().unwrap() = Some(head);
|
||||
Ok(ControlFlow::Continue(()))
|
||||
}
|
||||
async fn send_chunk(&self, chunk: P::Chunk) -> Result<ControlFlow<()>, P::Error> {
|
||||
self.chunks.lock().unwrap().push(chunk);
|
||||
Ok(ControlFlow::Continue(()))
|
||||
}
|
||||
}
|
||||
impl<P> litellm_host::lifecycle::CallObserver for RecordingCall<P>
|
||||
where
|
||||
P: litellm_host::protocol::Protocol<HostCall = std::convert::Infallible>,
|
||||
P::Error: From<litellm_host::machine::MachineFault>,
|
||||
{
|
||||
fn observe(&self, event: litellm_host::lifecycle::CallEvent) {
|
||||
self.events.0.sender.emit(event);
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
pub struct TraceCapture(Arc<Mutex<Vec<Value>>>);
|
||||
|
||||
impl TraceCapture {
|
||||
pub fn logger(&self) -> litellm_tracing::Logger {
|
||||
litellm_tracing::Logger::new(self.clone())
|
||||
}
|
||||
|
||||
pub fn records(&self) -> Vec<Value> {
|
||||
self.0.lock().unwrap().clone()
|
||||
}
|
||||
|
||||
pub fn summaries(&self, name: &str) -> Vec<Value> {
|
||||
self.records()
|
||||
.into_iter()
|
||||
.filter(|record| record["span_name"] == name)
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
impl litellm_tracing::Sink for TraceCapture {
|
||||
fn enabled(&self, metadata: &litellm_tracing::Metadata<'_>) -> bool {
|
||||
metadata.target().starts_with("litellm_inference")
|
||||
}
|
||||
|
||||
fn emit(&self, record: &litellm_tracing::Record) {
|
||||
self.0
|
||||
.lock()
|
||||
.unwrap()
|
||||
.push(Value::Object(record.fields.clone()));
|
||||
}
|
||||
}
|
||||
|
||||
#[rstest::fixture]
|
||||
pub fn traces() -> TraceCapture {
|
||||
TraceCapture::default()
|
||||
}
|
||||
|
|
@ -1,10 +1,5 @@
|
|||
pub const OPENAI_DEFAULT_API_BASE: &str = "https://api.openai.com";
|
||||
|
||||
/// Full-request timeout ceiling for Anthropic Messages provider calls, in
|
||||
/// seconds. Mirrors the Python Anthropic Messages default. The per-request
|
||||
/// timeout from the caller still overrides this on the request builder.
|
||||
pub(crate) const MESSAGES_TIMEOUT_SECS: u64 = 600;
|
||||
|
||||
/// Full-request timeout ceiling for chat completions provider calls, in
|
||||
/// seconds. Mirrors the Python chat completions default.
|
||||
pub(crate) const CHAT_COMPLETIONS_TIMEOUT_SECS: u64 = 600;
|
||||
|
|
|
|||
|
|
@ -5,7 +5,6 @@ pub mod caching;
|
|||
pub mod chat_completions;
|
||||
pub mod constants;
|
||||
pub mod error;
|
||||
pub mod messages;
|
||||
pub mod ocr;
|
||||
pub mod outbound;
|
||||
pub mod provider;
|
||||
|
|
|
|||
|
|
@ -402,64 +402,6 @@ async fn an_invalid_cached_envelope_is_replaced_by_a_provider_result(#[case] poi
|
|||
assert_eq!(calls.load(Ordering::SeqCst), 1);
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::system("system", json!("answer ALPHA"), json!("answer BETA"))]
|
||||
#[case::stop_sequences("stop_sequences", json!(["STOP"]), json!(["END"]))]
|
||||
#[case::top_k("top_k", json!(5), json!(10))]
|
||||
#[case::tools("tools", json!([{"name":"a","input_schema":{"type":"object"}}]), json!([{"name":"b","input_schema":{"type":"object"}}]))]
|
||||
#[case::tool_choice("tool_choice", json!({"type":"auto"}), json!({"type":"none"}))]
|
||||
#[tokio::test]
|
||||
async fn messages_cache_identity_includes_provider_native_parameters(
|
||||
cache: Arc<dyn ResponseCacheService>,
|
||||
#[case] field: &str,
|
||||
#[case] original: Value,
|
||||
#[case] changed: Value,
|
||||
) {
|
||||
use litellm_inference::messages::route::Messages;
|
||||
use litellm_llms_types::formats::messages::MessagesResponse;
|
||||
|
||||
let calls = AtomicUsize::new(0);
|
||||
for (value, expected_call) in [(original.clone(), 0), (changed, 1), (original, 0)] {
|
||||
let response =
|
||||
execute_unary::<Messages, _, _>(
|
||||
CacheRequest::from_wire(
|
||||
ProviderIdentity {
|
||||
model: "test".into(),
|
||||
provider: "anthropic".into(),
|
||||
},
|
||||
Some(&WireRequest {
|
||||
url: "https://example.test/v1/messages".into(),
|
||||
headers: vec![],
|
||||
body: json!({
|
||||
"model":"test", "messages":[{"role":"user","content":"hello"}],
|
||||
"max_tokens":32, (field):value
|
||||
}),
|
||||
}),
|
||||
),
|
||||
Some(cache.clone()),
|
||||
Some(CacheOptions::new(CacheScope::Shared)),
|
||||
&(),
|
||||
None,
|
||||
|| async {
|
||||
let call = calls.fetch_add(1, Ordering::SeqCst);
|
||||
Ok(Box::new(serde_json::from_value::<MessagesResponse>(json!({
|
||||
"id":call.to_string(), "type":"message", "role":"assistant", "model":"test",
|
||||
"content":[{"type":"text","text":format!("answer {call}")}],
|
||||
"stop_reason":"end_turn", "stop_sequence":null
|
||||
})).unwrap()))
|
||||
},
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(response.id, expected_call.to_string());
|
||||
assert_eq!(
|
||||
response.content[0]["text"],
|
||||
format!("answer {expected_call}")
|
||||
);
|
||||
}
|
||||
assert_eq!(calls.load(Ordering::SeqCst), 2);
|
||||
}
|
||||
|
||||
struct UnavailableCache;
|
||||
|
||||
impl litellm_cache::BaseCache for UnavailableCache {
|
||||
|
|
@ -943,9 +885,6 @@ impl Interceptors<RouteError> for ChangingHooks {
|
|||
#[case::chat_credentials("chat", "credentials")]
|
||||
#[case::chat_endpoint("chat", "endpoint")]
|
||||
#[case::chat_callback("chat", "callback")]
|
||||
#[case::messages_credentials("messages", "credentials")]
|
||||
#[case::messages_endpoint("messages", "endpoint")]
|
||||
#[case::messages_callback("messages", "callback")]
|
||||
#[tokio::test]
|
||||
async fn cache_identity_follows_resolved_configuration_and_request_callbacks(
|
||||
cache: Arc<dyn ResponseCacheService>,
|
||||
|
|
@ -953,9 +892,8 @@ async fn cache_identity_follows_resolved_configuration_and_request_callbacks(
|
|||
#[case] change: &str,
|
||||
) {
|
||||
use litellm_cache_response::ScopedCache;
|
||||
use litellm_inference::{
|
||||
chat_completions::{ChatCompletionsRoute, types::ChatCompletionsRequest},
|
||||
messages::MessagesCall,
|
||||
use litellm_inference::chat_completions::{
|
||||
ChatCompletionsRoute, types::ChatCompletionsRequest,
|
||||
};
|
||||
use wiremock::{Mock, MockServer, ResponseTemplate, matchers::method};
|
||||
|
||||
|
|
@ -1020,12 +958,6 @@ async fn cache_identity_follows_resolved_configuration_and_request_callbacks(
|
|||
.await
|
||||
.unwrap();
|
||||
}
|
||||
"messages" => {
|
||||
support::messages_route(secrets.clone()).with_cache(cache).execute(MessagesCall {
|
||||
body: serde_json::from_value(json!({"model":"anthropic/cache-test-model","messages":[{"role":"user","content":"hello"}],"max_tokens":32})).unwrap(),
|
||||
api_key:None,api_base:None,custom_llm_provider:None,extra_headers:None,provider_specific_header:None,timeout:None,shaping:Default::default(),
|
||||
}, &hooks, None).await.unwrap();
|
||||
}
|
||||
_ => unreachable!(),
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -17,17 +17,6 @@ use wiremock::{Mock, MockServer, Request, ResponseTemplate, matchers::any};
|
|||
/// A port nothing listens on, for calls that must fail before any request is sent.
|
||||
pub const UNREACHABLE_BASE: &str = "http://127.0.0.1:1";
|
||||
|
||||
pub fn messages_route(
|
||||
secrets: Arc<dyn SecretSource>,
|
||||
) -> litellm_inference::messages::MessagesRoute {
|
||||
let resources = resources();
|
||||
litellm_inference::messages::MessagesRoute::new(
|
||||
provider_http(&resources, &http_config()),
|
||||
resources.auth,
|
||||
secrets,
|
||||
)
|
||||
}
|
||||
|
||||
pub fn chat_completions_route() -> litellm_inference::chat_completions::ChatCompletionsRoute {
|
||||
let resources = resources();
|
||||
litellm_inference::chat_completions::ChatCompletionsRoute::new(
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
This directory owns shared Messages API data contracts and serialization: request and response bodies, messages, content blocks, usage, and stream-event payloads. Messages is a format independent of the provider that originated it. Keep one canonical definition of each shared contract here
|
||||
|
||||
Adapter contracts and execution inputs such as `MessagesTransformContext` belong in `llms/src/base_llm/messages`. Provider rewriting and interpretation belong in `llms/src/<provider>/messages`. Call envelopes, live streams, and call orchestration belong in `core/src/messages`
|
||||
Adapter contracts and execution inputs such as `MessagesTransformContext` belong in `llms/src/base_llm/messages`. Provider rewriting and interpretation belong in `llms/src/<provider>/messages`. Call envelopes, live streams, and call orchestration belong in `inference-messages`
|
||||
|
||||
Represent web-search results, encrypted-content fields, and thinking configuration as data here. Decisions to flatten results, remove encrypted content, select thinking budgets, or require beta headers belong to provider transformations. Data-shape validation belongs here, while model capability checks and request adaptation do not
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
This directory owns Anthropic's implementation of the Messages adapter contract in `base_llm/messages`. Shared Messages API data contracts belong in `litellm-llms-types::formats::messages`, and call orchestration belongs in `core/src/messages`. Sharing the `llms` crate with `base_llm/messages` does not erase this boundary
|
||||
This directory owns Anthropic's implementation of the Messages adapter contract in `base_llm/messages`. Shared Messages API data contracts belong in `litellm-llms-types::formats::messages`, and call orchestration belongs in `inference-messages`. Sharing the `llms` crate with `base_llm/messages` does not erase this boundary
|
||||
|
||||
Payload shaping, metadata filtering, tool-ID rewriting, web-search replay handling, thinking translation, and beta selection are provider policy. Keep them here or in Anthropic helpers shared by its operations. Pure payload shaping belongs with transformations, even if an existing file is named `handler.rs`
|
||||
|
||||
|
|
|
|||
|
|
@ -1,3 +1,3 @@
|
|||
This directory owns Azure's Messages adapter: its endpoints, authentication policy, headers, and transformations. Implement the shared adapter contract from `base_llm/messages`, consume API data contracts from `litellm-llms-types::formats::messages`, and leave call orchestration to `core/src/messages`
|
||||
This directory owns Azure's Messages adapter: its endpoints, authentication policy, headers, and transformations. Implement the shared adapter contract from `base_llm/messages`, consume API data contracts from `litellm-llms-types::formats::messages`, and leave call orchestration to `inference-messages`
|
||||
|
||||
The Claude adapter may explicitly reuse payload policy from `anthropic/messages` when it applies to Azure's Claude backend. Keep Azure-specific differences here. Sharing that helper does not make Anthropic policy a format-wide default or justify a dependency from `base_llm/messages` on provider implementations
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
This directory owns the shared Messages provider adapter contract, its execution inputs such as `MessagesTransformContext`, and provider-independent transformation machinery. Public request, response, content-block, and event schemas belong in `litellm-llms-types::formats::messages`. Call orchestration belongs in `core/src/messages`, and provider implementations belong in `llms/src/<provider>/messages`
|
||||
This directory owns the shared Messages provider adapter contract, its execution inputs such as `MessagesTransformContext`, and provider-independent transformation machinery. Public request, response, content-block, and event schemas belong in `litellm-llms-types::formats::messages`. Call orchestration belongs in `inference-messages`, and provider implementations belong in `llms/src/<provider>/messages`
|
||||
|
||||
Do not import provider implementations or embed their policy in shared trait defaults, normalization, or context defaults. A context carries inputs the shared adapter contract needs, not every provider's settings. Thinking-budget choices and model-specific restrictions do not become format rules merely because several providers host Claude
|
||||
|
||||
|
|
|
|||
|
|
@ -1,3 +1,3 @@
|
|||
This directory owns Bedrock's Messages adapter: its endpoints, authentication policy, wire adaptation, and response decoding. Implement the shared adapter contract from `base_llm/messages`, consume API data contracts from `litellm-llms-types::formats::messages`, and leave call orchestration to `core/src/messages`
|
||||
This directory owns Bedrock's Messages adapter: its endpoints, authentication policy, wire adaptation, and response decoding. Implement the shared adapter contract from `base_llm/messages`, consume API data contracts from `litellm-llms-types::formats::messages`, and leave call orchestration to `inference-messages`
|
||||
|
||||
The Claude adapter may explicitly reuse payload policy from `anthropic/messages` when it applies to Bedrock's Claude backend. Keep Bedrock-specific differences here. Sharing that helper does not make Anthropic policy a format-wide default or justify a dependency from `base_llm/messages` on provider implementations
|
||||
|
|
|
|||
|
|
@ -85,8 +85,8 @@ GIL handling to `litellm-host-python`.
|
|||
`messages(...)`, calling the matching `litellm-inference` entrypoint.
|
||||
- Do not add one exported PyO3 function per provider helper unless there is a
|
||||
measured reason.
|
||||
- Provider dispatch belongs in the `litellm-inference` route module (e.g.
|
||||
`litellm_inference::messages`), not in this PyO3 crate.
|
||||
- Provider dispatch belongs in the `litellm-inference-*` route crate (e.g.
|
||||
`litellm_inference_messages`), not in this PyO3 crate.
|
||||
- Python owns rollout state and fallback. Rust should return errors; Python
|
||||
decides whether to raise or fall back. For a rust-only provider/route (no
|
||||
Python reference), the Python side is a thin dispatch that calls Rust and
|
||||
|
|
|
|||
|
|
@ -39,6 +39,7 @@ litellm-callbacks-legacy-python.workspace = true
|
|||
litellm-inference.workspace = true
|
||||
litellm-inference-transcription.workspace = true
|
||||
litellm-inference-responses.workspace = true
|
||||
litellm-inference-messages.workspace = true
|
||||
litellm-core-utils.workspace = true
|
||||
litellm-http.workspace = true
|
||||
litellm-llms.workspace = true
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ use litellm_host_python::{PythonHostCalls, PythonOwned};
|
|||
use bytes::Bytes;
|
||||
use litellm_host_python::{InvokeError, PythonBinding, from_py, lookup, to_py};
|
||||
use litellm_http::transport::Error as TransportError;
|
||||
use litellm_inference::messages::{
|
||||
use litellm_inference_messages::{
|
||||
Error, MessagesCall, MessagesShaping, messages_body,
|
||||
route::{Messages, MessagesStreamHead},
|
||||
};
|
||||
|
|
|
|||
|
|
@ -26,7 +26,7 @@ fn run_messages(
|
|||
py,
|
||||
arguments,
|
||||
move |py, arguments, request| {
|
||||
let route = litellm_inference::messages::MessagesRoute::new(
|
||||
let route = litellm_inference_messages::MessagesRoute::new(
|
||||
crate::http::provider_client(py, arguments, asynchronous)?
|
||||
.map_err(crate::http::client_error)?,
|
||||
crate::http::resources().auth.clone(),
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ repository.workspace = true
|
|||
[dependencies]
|
||||
litellm-config.workspace = true
|
||||
litellm-inference.workspace = true
|
||||
litellm-inference-messages.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
rstest.workspace = true
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
use std::time::Duration;
|
||||
|
||||
use litellm_inference::messages::MessagesShaping;
|
||||
use litellm_inference_messages::MessagesShaping;
|
||||
|
||||
#[derive(Clone, Debug, Default)]
|
||||
pub struct Deployment {
|
||||
|
|
|
|||
|
|
@ -1,7 +1,7 @@
|
|||
use std::time::Duration;
|
||||
|
||||
use litellm_config::Config;
|
||||
use litellm_inference::messages::MessagesShaping;
|
||||
use litellm_inference_messages::MessagesShaping;
|
||||
use litellm_router::{Deployment, Router};
|
||||
use rstest::rstest;
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue