mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-29 01:42:19 +00:00
* refactor(rust): align the cache crates with Python and activate every backend The cache port had drifted: lifecycle and Redis-only operations sat on `BaseCache`, counters were pinned to `f64`, each semantic backend defined its own embedder and prompt handling, and only the in-memory backend could be selected natively. - Split `disconnect` and `test_connection` out of `BaseCache` into optional capabilities, implemented only where the Python class defines them, and give every Redis-only operation its own capability trait. - Decouple counters from the stored value type, so one backend can serve both responses and counters as Python's `RedisCache` does. - Share one `Embedder` and prompt contract in `litellm_cache::semantic`, and make the Redis and Valkey semantic backends generic over their codec. - Port the Python operations that were missing: `async_refresh_ttl`, `async_rpush_and_trim`, `async_set_cache_pipeline_with_ttls`, the DualCache pipeline, sadd, bulk delete and TTL reads, and the semantic-similarity write-back. - Take the HTTP client from the host pool in the GCS, S3 and Azure backends. - Activate all nine backends through the Rust catalog, whose rules all stay `PYTHON_ONLY`, and route the `Cache` facade's storage calls to the native runtime when one is selected. - Give every crate the same layout, move all tests to `tests/` on rstest, and add the shared `litellm-cache-testing` contract suite. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: freeze native cache request kwargs and batch entries for type discipline Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix: declare semantic lookup methods in the native stub Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(rust): align the cache crates with Python and activate every backend The cache port had drifted: lifecycle and Redis-only operations sat on `BaseCache`, counters were pinned to `f64`, each semantic backend defined its own embedder and prompt handling, and only the in-memory backend could be selected natively. - Split `disconnect` and `test_connection` out of `BaseCache` into optional capabilities, implemented only where the Python class defines them, and give every Redis-only operation its own capability trait. - Decouple counters from the stored value type, so one backend can serve both responses and counters as Python's `RedisCache` does. - Share one `Embedder` and prompt contract in `litellm_cache::semantic`, and make the Redis and Valkey semantic backends generic over their codec. - Port the Python operations that were missing: `async_refresh_ttl`, `async_rpush_and_trim`, `async_set_cache_pipeline_with_ttls`, the DualCache pipeline, sadd, bulk delete and TTL reads, and the semantic-similarity write-back. - Take the HTTP client from the host pool in the GCS, S3 and Azure backends. - Activate all nine backends through the Rust catalog, whose rules all stay `PYTHON_ONLY`, and route the `Cache` facade's storage calls to the native runtime when one is selected. - Give every crate the same layout, move all tests to `tests/` on rstest, and add the shared `litellm-cache-testing` contract suite. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: freeze native cache request kwargs and batch entries for type discipline Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix: declare semantic lookup methods in the native stub Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(rust): opt the native Messages and tokenizer suites into Rust explicitly #42517 made the Messages, token counter and tokenizer routes Python-only, so tests/test_litellm_rust silently exercised the Python path or failed outright. Each suite now prepends a RUST_OPT_IN rule for its route, keeping native coverage without changing the shipped default. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(rust): pop one at a time in the Redis 6 lpop pipeline and drop explanatory comments Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Yujong Lee <yujong@berri.ai> Co-authored-by: Claude Opus 5 <noreply@anthropic.com> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
192 lines
6 KiB
Rust
192 lines
6 KiB
Rust
use std::{
|
|
sync::{Arc, OnceLock},
|
|
time::Duration,
|
|
};
|
|
|
|
use litellm_cache::{BatchEntry, CacheCodec, Error};
|
|
|
|
use crate::{
|
|
connection::{ConnectionRef, Connections},
|
|
topology::RedisTopology,
|
|
};
|
|
|
|
const DEFAULT_TTL: Duration = Duration::from_secs(600);
|
|
|
|
pub struct RedisCache<S, C = redis::Connection> {
|
|
pub(crate) connections: Arc<Connections<C>>,
|
|
pub(crate) default_ttl: Duration,
|
|
pub(crate) codec: S,
|
|
pub(crate) namespace: Option<String>,
|
|
pub(crate) topology: RedisTopology,
|
|
/// The server's major version, read from `INFO` once, like Python's `redis_version`.
|
|
pub(crate) major_version: Arc<OnceLock<u32>>,
|
|
}
|
|
|
|
impl<S: CacheCodec> RedisCache<S> {
|
|
pub fn new(url: &str, default_ttl: Option<Duration>, codec: S) -> Result<Self, Error> {
|
|
Self::connect(url, &RedisTopology::Standalone, default_ttl, codec)
|
|
}
|
|
|
|
pub fn connect(
|
|
url: &str,
|
|
topology: &RedisTopology,
|
|
default_ttl: Option<Duration>,
|
|
codec: S,
|
|
) -> Result<Self, Error> {
|
|
let connections = Connections::open(url, topology)?;
|
|
Ok(Self {
|
|
connections: Arc::new(connections),
|
|
default_ttl: default_ttl.unwrap_or(DEFAULT_TTL),
|
|
codec,
|
|
namespace: None,
|
|
topology: topology.clone(),
|
|
major_version: Arc::default(),
|
|
})
|
|
}
|
|
}
|
|
|
|
impl<S, C> RedisCache<S, C>
|
|
where
|
|
S: CacheCodec,
|
|
C: redis::ConnectionLike + Send + 'static,
|
|
{
|
|
pub fn with_connection(connection: C, default_ttl: Option<Duration>, codec: S) -> Self {
|
|
Self {
|
|
connections: Arc::new(Connections::fixed(connection)),
|
|
default_ttl: default_ttl.unwrap_or(DEFAULT_TTL),
|
|
codec,
|
|
namespace: None,
|
|
topology: RedisTopology::Standalone,
|
|
major_version: Arc::default(),
|
|
}
|
|
}
|
|
|
|
pub fn with_namespace(self, namespace: Option<String>) -> Self {
|
|
Self {
|
|
namespace: namespace.filter(|value| !value.is_empty()),
|
|
..self
|
|
}
|
|
}
|
|
|
|
pub fn namespace(&self) -> Option<&str> {
|
|
self.namespace.as_deref()
|
|
}
|
|
|
|
pub fn topology(&self) -> &RedisTopology {
|
|
&self.topology
|
|
}
|
|
|
|
pub(crate) fn namespaced_key(&self, key: &str) -> String {
|
|
namespaced_key(self.namespace.as_deref(), key)
|
|
}
|
|
|
|
pub(crate) fn namespaced_keys(&self, keys: &[String]) -> Vec<String> {
|
|
keys.iter().map(|key| self.namespaced_key(key)).collect()
|
|
}
|
|
|
|
/// Whole seconds for `ttl`, falling back to the default TTL like Python's `get_ttl`.
|
|
pub(crate) fn ttl_or_default(&self, ttl: Option<Duration>) -> u64 {
|
|
ttl_seconds(ttl.unwrap_or(self.default_ttl))
|
|
}
|
|
|
|
pub(crate) fn execute<T>(
|
|
&self,
|
|
operation: impl FnOnce(&mut ConnectionRef<'_>) -> Result<T, Error>,
|
|
) -> Result<T, Error> {
|
|
self.connections.execute(operation)
|
|
}
|
|
|
|
/// `_parse_redis_major_version`: the major version from `INFO`, or
|
|
/// `DEFAULT_REDIS_MAJOR_VERSION` when `INFO` fails or its version does not parse. The first
|
|
/// answer is kept, as Python reads `redis_version` once at construction.
|
|
pub(crate) async fn major_version(&self) -> u32 {
|
|
if let Some(version) = self.major_version.get() {
|
|
return *version;
|
|
}
|
|
let info = self
|
|
.run(|connection| connection.node_text(&redis::cmd("INFO")))
|
|
.await;
|
|
let version = info
|
|
.ok()
|
|
.and_then(|info| parse_major_version(&info))
|
|
.unwrap_or_else(default_major_version);
|
|
*self.major_version.get_or_init(|| version)
|
|
}
|
|
|
|
pub(crate) async fn run<T, F>(&self, operation: F) -> Result<T, Error>
|
|
where
|
|
T: Send + 'static,
|
|
F: FnOnce(&mut ConnectionRef<'_>) -> Result<T, Error> + Send + 'static,
|
|
{
|
|
Connections::run_blocking(Arc::clone(&self.connections), operation).await
|
|
}
|
|
|
|
pub(crate) fn namespaced_pattern(&self) -> Result<String, Error> {
|
|
let namespace = self.namespace.as_ref().ok_or(Error::UnscopedFlush)?;
|
|
let escaped: String = namespace
|
|
.chars()
|
|
.flat_map(|ch| {
|
|
if matches!(ch, '*' | '?' | '[' | ']' | '\\') {
|
|
vec!['\\', ch]
|
|
} else {
|
|
vec![ch]
|
|
}
|
|
})
|
|
.collect();
|
|
Ok(format!("{escaped}:*"))
|
|
}
|
|
|
|
pub(crate) fn decode_response(&self, value: redis::Value) -> Result<Option<S::Value>, Error> {
|
|
match value {
|
|
redis::Value::Nil => Ok(None),
|
|
redis::Value::BulkString(bytes) => self.codec.decode(&bytes).map(Some),
|
|
redis::Value::SimpleString(text) => self.codec.decode(text.as_bytes()).map(Some),
|
|
_ => Err(Error::InvalidEntry),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn decode_batch_response(
|
|
&self,
|
|
value: redis::Value,
|
|
) -> Result<BatchEntry<S::Value>, Error> {
|
|
match self.decode_response(value) {
|
|
Ok(Some(value)) => Ok(BatchEntry::Hit(value)),
|
|
Ok(None) => Ok(BatchEntry::Miss),
|
|
Err(Error::InvalidEntry) => Ok(BatchEntry::Invalid),
|
|
Err(error) => Err(error),
|
|
}
|
|
}
|
|
}
|
|
|
|
pub(crate) fn namespaced_key(namespace: Option<&str>, key: &str) -> String {
|
|
match namespace {
|
|
Some(namespace) if !key.starts_with(&format!("{namespace}:")) => {
|
|
format!("{namespace}:{key}")
|
|
}
|
|
_ => key.into(),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn ttl_seconds(ttl: Duration) -> u64 {
|
|
ttl.as_secs()
|
|
.saturating_add(u64::from(ttl.subsec_nanos() > 0))
|
|
.max(1)
|
|
}
|
|
|
|
fn parse_major_version(info: &str) -> Option<u32> {
|
|
let version = info
|
|
.lines()
|
|
.find_map(|line| line.trim().strip_prefix("redis_version:"))?
|
|
.trim();
|
|
match version.split_once('.') {
|
|
Some((major, _)) => major.parse().ok(),
|
|
None => version.parse::<f64>().ok().map(|major| major as u32),
|
|
}
|
|
}
|
|
|
|
fn default_major_version() -> u32 {
|
|
std::env::var("DEFAULT_REDIS_MAJOR_VERSION")
|
|
.ok()
|
|
.and_then(|value| value.parse().ok())
|
|
.unwrap_or(7)
|
|
}
|