mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-02 02:11:58 +00:00
feat(rust): embed migration folders with a shared migrate! macro (#44104)
* feat(rust): embed migration folders with a shared migrate! macro Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(rust): enable syn proc-macro feature for litellm-migrate-macros Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(rust): reject signed versions and symlinks in migrate! 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: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
1e75147856
commit
e3c15c9d22
17 changed files with 335 additions and 14 deletions
21
litellm-rust/Cargo.lock
generated
21
litellm-rust/Cargo.lock
generated
|
|
@ -4038,6 +4038,26 @@ dependencies = [
|
|||
"strum",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "litellm-migrate"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"litellm-migrate-macros",
|
||||
"rstest",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "litellm-migrate-macros"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"rstest",
|
||||
"syn 2.0.119",
|
||||
"tempfile",
|
||||
"thiserror 2.0.19",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "litellm-model-catalog"
|
||||
version = "0.1.0"
|
||||
|
|
@ -4388,6 +4408,7 @@ dependencies = [
|
|||
"criterion",
|
||||
"flate2",
|
||||
"litellm-http",
|
||||
"litellm-migrate",
|
||||
"litellm-storage-clickhouse",
|
||||
"opentelemetry-proto",
|
||||
"prost",
|
||||
|
|
|
|||
|
|
@ -14,6 +14,8 @@ litellm-router = { path = "crates/router" }
|
|||
litellm-tracing = { path = "crates/tracing" }
|
||||
litellm-traces = { path = "crates/traces" }
|
||||
litellm-storage-clickhouse = { path = "crates/storage-clickhouse" }
|
||||
litellm-migrate = { path = "crates/migrate" }
|
||||
litellm-migrate-macros = { path = "crates/migrate-macros" }
|
||||
litellm-core = { path = "crates/core" }
|
||||
litellm-gateway-mcp = { path = "crates/gateway-mcp" }
|
||||
litellm-gateway = { path = "crates/gateway" }
|
||||
|
|
@ -93,7 +95,10 @@ serde = { version = "1.0", features = ["derive"] }
|
|||
serde_json = { version = "1.0", features = ["float_roundtrip"] }
|
||||
serde_with = { version = "=3.16.1", default-features = false, features = ["std", "macros"] }
|
||||
sha2 = "0.10"
|
||||
syn = { version = "2", default-features = false }
|
||||
sqlx = { version = "0.9.0", default-features = false, features = ["json", "macros", "postgres", "runtime-tokio", "chrono", "tls-rustls-ring-native-roots"] }
|
||||
proc-macro2 = "1"
|
||||
quote = "1"
|
||||
subtle = "2"
|
||||
thiserror = "2.0"
|
||||
tokenizers = { version = "0.23.1", default-features = false, features = ["onig"] }
|
||||
|
|
|
|||
19
litellm-rust/crates/migrate-macros/Cargo.toml
Normal file
19
litellm-rust/crates/migrate-macros/Cargo.toml
Normal file
|
|
@ -0,0 +1,19 @@
|
|||
[package]
|
||||
name = "litellm-migrate-macros"
|
||||
version = "0.1.0"
|
||||
edition.workspace = true
|
||||
license.workspace = true
|
||||
repository.workspace = true
|
||||
|
||||
[lib]
|
||||
proc-macro = true
|
||||
|
||||
[dependencies]
|
||||
proc-macro2.workspace = true
|
||||
quote.workspace = true
|
||||
syn = { workspace = true, features = ["parsing", "printing", "proc-macro"] }
|
||||
thiserror.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
rstest.workspace = true
|
||||
tempfile.workspace = true
|
||||
21
litellm-rust/crates/migrate-macros/src/error.rs
Normal file
21
litellm-rust/crates/migrate-macros/src/error.rs
Normal file
|
|
@ -0,0 +1,21 @@
|
|||
use std::io;
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum Error {
|
||||
#[error("could not read migrations directory `{path}`")]
|
||||
ReadDirectory {
|
||||
path: String,
|
||||
#[source]
|
||||
source: io::Error,
|
||||
},
|
||||
#[error(
|
||||
"migration name `{name}` must be `<digits>_<description>.sql` with a `[a-z0-9_]` description"
|
||||
)]
|
||||
InvalidName { name: String },
|
||||
#[error("migration version `{version}` is declared more than once")]
|
||||
DuplicateVersion { version: u64 },
|
||||
#[error("migrations directory `{path}` contains no migrations")]
|
||||
Empty { path: String },
|
||||
#[error("migration path `{path}` is not valid UTF-8")]
|
||||
NonUtf8Path { path: String },
|
||||
}
|
||||
199
litellm-rust/crates/migrate-macros/src/lib.rs
Normal file
199
litellm-rust/crates/migrate-macros/src/lib.rs
Normal file
|
|
@ -0,0 +1,199 @@
|
|||
mod error;
|
||||
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
use error::Error;
|
||||
use proc_macro::TokenStream;
|
||||
use quote::quote;
|
||||
use syn::LitStr;
|
||||
|
||||
struct Entry {
|
||||
version: u64,
|
||||
description: String,
|
||||
path: PathBuf,
|
||||
}
|
||||
|
||||
fn resolve(dir: &Path) -> Result<Vec<Entry>, Error> {
|
||||
let mut entries = Vec::new();
|
||||
let files = std::fs::read_dir(dir).map_err(|source| Error::ReadDirectory {
|
||||
path: dir.display().to_string(),
|
||||
source,
|
||||
})?;
|
||||
for file in files {
|
||||
let file = file.map_err(|source| Error::ReadDirectory {
|
||||
path: dir.display().to_string(),
|
||||
source,
|
||||
})?;
|
||||
let path = file.path();
|
||||
let name = path
|
||||
.file_name()
|
||||
.and_then(|name| name.to_str())
|
||||
.ok_or_else(|| Error::NonUtf8Path {
|
||||
path: path.display().to_string(),
|
||||
})?
|
||||
.to_owned();
|
||||
let invalid = || Error::InvalidName { name: name.clone() };
|
||||
let stem = name
|
||||
.strip_suffix(".sql")
|
||||
.filter(|_| file.file_type().is_ok_and(|kind| kind.is_file()))
|
||||
.and_then(|stem| stem.split_once('_'))
|
||||
.filter(|(version, description)| {
|
||||
!version.is_empty()
|
||||
&& version.bytes().all(|b| b.is_ascii_digit())
|
||||
&& !description.is_empty()
|
||||
&& description
|
||||
.bytes()
|
||||
.all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'_')
|
||||
})
|
||||
.ok_or_else(invalid)?;
|
||||
let version = stem.0.parse::<u64>().map_err(|_| invalid())?;
|
||||
entries.push(Entry {
|
||||
version,
|
||||
description: stem.1.to_owned(),
|
||||
path,
|
||||
});
|
||||
}
|
||||
if entries.is_empty() {
|
||||
return Err(Error::Empty {
|
||||
path: dir.display().to_string(),
|
||||
});
|
||||
}
|
||||
entries.sort_by_key(|entry| entry.version);
|
||||
for pair in entries.windows(2) {
|
||||
if pair[0].version == pair[1].version {
|
||||
return Err(Error::DuplicateVersion {
|
||||
version: pair[0].version,
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(entries)
|
||||
}
|
||||
|
||||
fn resolve_input(lit: &LitStr) -> Result<Vec<Entry>, Error> {
|
||||
let root = std::env::var("CARGO_MANIFEST_DIR")
|
||||
.map(PathBuf::from)
|
||||
.unwrap_or_default();
|
||||
let dir = root.join(lit.value());
|
||||
let dir = dir.canonicalize().map_err(|source| Error::ReadDirectory {
|
||||
path: dir.display().to_string(),
|
||||
source,
|
||||
})?;
|
||||
if dir.to_str().is_none() {
|
||||
return Err(Error::NonUtf8Path {
|
||||
path: dir.display().to_string(),
|
||||
});
|
||||
}
|
||||
resolve(&dir)
|
||||
}
|
||||
|
||||
#[proc_macro]
|
||||
pub fn migrate(input: TokenStream) -> TokenStream {
|
||||
let lit = syn::parse_macro_input!(input as LitStr);
|
||||
match resolve_input(&lit) {
|
||||
Ok(entries) => {
|
||||
let migrations = entries.iter().map(|entry| {
|
||||
let version = entry.version;
|
||||
let description = &entry.description;
|
||||
let path = entry
|
||||
.path
|
||||
.to_str()
|
||||
.expect("canonical migration path is UTF-8");
|
||||
quote! {
|
||||
::litellm_migrate::Migration {
|
||||
version: #version,
|
||||
description: #description,
|
||||
sql: ::core::include_str!(#path),
|
||||
}
|
||||
}
|
||||
});
|
||||
quote! { &[#(#migrations),*] }.into()
|
||||
}
|
||||
Err(err) => syn::Error::new(lit.span(), err).to_compile_error().into(),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::fs;
|
||||
|
||||
use rstest::rstest;
|
||||
use tempfile::TempDir;
|
||||
|
||||
use super::{Error, resolve};
|
||||
|
||||
fn migrations_dir(files: &[&str]) -> TempDir {
|
||||
let dir = TempDir::new().expect("tempdir");
|
||||
for file in files {
|
||||
fs::write(dir.path().join(file), "SELECT 1").expect("write fixture");
|
||||
}
|
||||
dir
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
fn orders_versions_numerically() {
|
||||
let dir = migrations_dir(&["10_tenth.sql", "2_second.sql", "1_first.sql"]);
|
||||
let entries = resolve(dir.path()).expect("resolves");
|
||||
let versions: Vec<u64> = entries.iter().map(|entry| entry.version).collect();
|
||||
let descriptions: Vec<&str> = entries
|
||||
.iter()
|
||||
.map(|entry| entry.description.as_str())
|
||||
.collect();
|
||||
assert_eq!(versions, [1, 2, 10]);
|
||||
assert_eq!(descriptions, ["first", "second", "tenth"]);
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::dash_in_version(&["0001-dash.sql"])]
|
||||
#[case::not_sql(&["notes.txt"])]
|
||||
#[case::empty_description(&["0001_.sql"])]
|
||||
#[case::non_digit_version(&["x_name.sql"])]
|
||||
#[case::uppercase_description(&["0001_Upper.sql"])]
|
||||
#[case::no_underscore(&["0001.sql"])]
|
||||
#[case::plus_sign_version(&["+10_add.sql"])]
|
||||
fn rejects_invalid_names(#[case] files: &[&str]) {
|
||||
let dir = migrations_dir(files);
|
||||
assert!(matches!(
|
||||
resolve(dir.path()),
|
||||
Err(Error::InvalidName { .. })
|
||||
));
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
fn rejects_subdirectories() {
|
||||
let dir = migrations_dir(&["0001_a.sql"]);
|
||||
fs::create_dir(dir.path().join("0002_b.sql")).expect("subdir");
|
||||
assert!(matches!(
|
||||
resolve(dir.path()),
|
||||
Err(Error::InvalidName { .. })
|
||||
));
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[rstest]
|
||||
fn rejects_symlinks() {
|
||||
let dir = migrations_dir(&["0001_a.sql"]);
|
||||
let target = TempDir::new().expect("tempdir");
|
||||
let target_file = target.path().join("real.sql");
|
||||
fs::write(&target_file, "SELECT 2").expect("write fixture");
|
||||
std::os::unix::fs::symlink(&target_file, dir.path().join("0002_b.sql")).expect("symlink");
|
||||
assert!(matches!(
|
||||
resolve(dir.path()),
|
||||
Err(Error::InvalidName { .. })
|
||||
));
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
fn rejects_duplicate_versions() {
|
||||
let dir = migrations_dir(&["0001_a.sql", "1_b.sql"]);
|
||||
assert!(matches!(
|
||||
resolve(dir.path()),
|
||||
Err(Error::DuplicateVersion { version: 1 })
|
||||
));
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
fn rejects_empty_directory() {
|
||||
let dir = migrations_dir(&[]);
|
||||
assert!(matches!(resolve(dir.path()), Err(Error::Empty { .. })));
|
||||
}
|
||||
}
|
||||
12
litellm-rust/crates/migrate/Cargo.toml
Normal file
12
litellm-rust/crates/migrate/Cargo.toml
Normal file
|
|
@ -0,0 +1,12 @@
|
|||
[package]
|
||||
name = "litellm-migrate"
|
||||
version = "0.1.0"
|
||||
edition.workspace = true
|
||||
license.workspace = true
|
||||
repository.workspace = true
|
||||
|
||||
[dependencies]
|
||||
litellm-migrate-macros.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
rstest.workspace = true
|
||||
5
litellm-rust/crates/migrate/README.md
Normal file
5
litellm-rust/crates/migrate/README.md
Normal file
|
|
@ -0,0 +1,5 @@
|
|||
# Migrations
|
||||
|
||||
`litellm-migrate` exports the `Migration` struct and the `migrate!` macro that embeds a directory of `<digits>_<description>.sql` files at compile time, sorted by numeric version
|
||||
|
||||
The crate does not apply or track migrations; callers decide how and when the embedded SQL runs
|
||||
8
litellm-rust/crates/migrate/src/lib.rs
Normal file
8
litellm-rust/crates/migrate/src/lib.rs
Normal file
|
|
@ -0,0 +1,8 @@
|
|||
pub use litellm_migrate_macros::migrate;
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub struct Migration {
|
||||
pub version: u64,
|
||||
pub description: &'static str,
|
||||
pub sql: &'static str,
|
||||
}
|
||||
1
litellm-rust/crates/migrate/tests/fixtures/migrations/10_tenth.sql
vendored
Normal file
1
litellm-rust/crates/migrate/tests/fixtures/migrations/10_tenth.sql
vendored
Normal file
|
|
@ -0,0 +1 @@
|
|||
SELECT 10;
|
||||
1
litellm-rust/crates/migrate/tests/fixtures/migrations/1_first.sql
vendored
Normal file
1
litellm-rust/crates/migrate/tests/fixtures/migrations/1_first.sql
vendored
Normal file
|
|
@ -0,0 +1 @@
|
|||
SELECT 1;
|
||||
1
litellm-rust/crates/migrate/tests/fixtures/migrations/2_second.sql
vendored
Normal file
1
litellm-rust/crates/migrate/tests/fixtures/migrations/2_second.sql
vendored
Normal file
|
|
@ -0,0 +1 @@
|
|||
SELECT 2;
|
||||
21
litellm-rust/crates/migrate/tests/migrate.rs
Normal file
21
litellm-rust/crates/migrate/tests/migrate.rs
Normal file
|
|
@ -0,0 +1,21 @@
|
|||
use litellm_migrate::Migration;
|
||||
use rstest::rstest;
|
||||
|
||||
const MIGRATIONS: &[Migration] = litellm_migrate::migrate!("tests/fixtures/migrations");
|
||||
|
||||
#[rstest]
|
||||
#[case::first(0, 1, "first", include_str!("fixtures/migrations/1_first.sql"))]
|
||||
#[case::second(1, 2, "second", include_str!("fixtures/migrations/2_second.sql"))]
|
||||
#[case::tenth(2, 10, "tenth", include_str!("fixtures/migrations/10_tenth.sql"))]
|
||||
fn embeds_every_file_sorted_by_numeric_version(
|
||||
#[case] index: usize,
|
||||
#[case] version: u64,
|
||||
#[case] description: &str,
|
||||
#[case] sql: &str,
|
||||
) {
|
||||
assert_eq!(MIGRATIONS.len(), 3);
|
||||
let migration = &MIGRATIONS[index];
|
||||
assert_eq!(migration.version, version);
|
||||
assert_eq!(migration.description, description);
|
||||
assert_eq!(migration.sql, sql);
|
||||
}
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
- Keep OTLP decoding, trace schema, row encoding and named query selection here. Generic ClickHouse connections and HTTP execution belong in `litellm-storage-clickhouse`
|
||||
- Keep this crate independent of Python; PyO3 conversion and public Python exceptions belong in `python-bridge`
|
||||
- Keep the SQL migrations here as the only ClickHouse schema definition
|
||||
- Keep the SQL migrations here as the only ClickHouse schema definition, as `migrations/NNNN_description.sql` files embedded by `litellm_migrate::migrate!`; adding a file is the only step
|
||||
- Use typed query parameters and a dedicated SELECT-only reader with server-side limits
|
||||
- Keep `config/reader.xml` grants on the database the schema is created in (CLICKHOUSE_DATABASE, default `litellm`)
|
||||
- Bound insert time and encoded bytes; make retry deduplication behavior explicit for supported ClickHouse versions
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@ opentelemetry-proto = { workspace = true, features = ["gen-tonic-messages", "tra
|
|||
prost.workspace = true
|
||||
time = { workspace = true, features = ["formatting"] }
|
||||
litellm-http.workspace = true
|
||||
litellm-migrate.workspace = true
|
||||
litellm-storage-clickhouse.workspace = true
|
||||
sha2.workspace = true
|
||||
serde = { workspace = true, features = ["rc"] }
|
||||
|
|
|
|||
3
litellm-rust/crates/traces/build.rs
Normal file
3
litellm-rust/crates/traces/build.rs
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
fn main() {
|
||||
println!("cargo:rerun-if-changed=migrations");
|
||||
}
|
||||
|
|
@ -1,4 +1,5 @@
|
|||
use litellm_http::Client;
|
||||
use litellm_migrate::Migration;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::Connection;
|
||||
|
|
@ -6,17 +7,7 @@ use crate::Error;
|
|||
|
||||
const SCHEMA_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
|
||||
|
||||
const MIGRATIONS: [&str; 9] = [
|
||||
include_str!("../migrations/0001_otel_traces.sql"),
|
||||
include_str!("../migrations/0002_agent_traces.sql"),
|
||||
include_str!("../migrations/0003_agent_traces_mv.sql"),
|
||||
include_str!("../migrations/0004_spend_logs.sql"),
|
||||
include_str!("../migrations/0005_otel_traces_ttl.sql"),
|
||||
include_str!("../migrations/0006_agent_traces_ttl.sql"),
|
||||
include_str!("../migrations/0007_spend_logs_ttl.sql"),
|
||||
include_str!("../migrations/0008_trace_received.sql"),
|
||||
include_str!("../migrations/0009_spend_received.sql"),
|
||||
];
|
||||
const MIGRATIONS: &[Migration] = litellm_migrate::migrate!("migrations");
|
||||
|
||||
pub fn schema_statements(
|
||||
database: &str,
|
||||
|
|
@ -35,8 +26,10 @@ pub fn schema_statements(
|
|||
let database = format!("`{database}`");
|
||||
Ok(
|
||||
std::iter::once(format!("CREATE DATABASE IF NOT EXISTS {database}"))
|
||||
.chain(MIGRATIONS.iter().map(|sql| {
|
||||
sql.replace("{database}", &database)
|
||||
.chain(MIGRATIONS.iter().map(|migration| {
|
||||
migration
|
||||
.sql
|
||||
.replace("{database}", &database)
|
||||
.replace("{trace_retention_days}", &trace_retention_days.to_string())
|
||||
.replace(
|
||||
"{spend_log_retention_days}",
|
||||
|
|
|
|||
|
|
@ -984,6 +984,16 @@ async fn duplicate_span_preview_matches_diagnostic(
|
|||
Ok(())
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
fn schema_includes_every_migration_file() -> TestResult {
|
||||
let files = std::fs::read_dir(concat!(env!("CARGO_MANIFEST_DIR"), "/migrations"))?
|
||||
.filter_map(|entry| entry.ok())
|
||||
.filter(|entry| entry.path().extension().is_some_and(|ext| ext == "sql"))
|
||||
.count();
|
||||
assert_eq!(schema_statements("trace_test", 7, 14)?.len(), 1 + files);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[tokio::test]
|
||||
async fn lens_agent_discovery_and_selection_preserve_scope(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue