diff --git a/litellm-rust/Cargo.lock b/litellm-rust/Cargo.lock index ff0eafee47e..0c31f12cdc3 100644 --- a/litellm-rust/Cargo.lock +++ b/litellm-rust/Cargo.lock @@ -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", diff --git a/litellm-rust/Cargo.toml b/litellm-rust/Cargo.toml index 8d837c2d31b..84b758da538 100644 --- a/litellm-rust/Cargo.toml +++ b/litellm-rust/Cargo.toml @@ -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"] } diff --git a/litellm-rust/crates/migrate-macros/Cargo.toml b/litellm-rust/crates/migrate-macros/Cargo.toml new file mode 100644 index 00000000000..5cd68415ca2 --- /dev/null +++ b/litellm-rust/crates/migrate-macros/Cargo.toml @@ -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 diff --git a/litellm-rust/crates/migrate-macros/src/error.rs b/litellm-rust/crates/migrate-macros/src/error.rs new file mode 100644 index 00000000000..9833009517b --- /dev/null +++ b/litellm-rust/crates/migrate-macros/src/error.rs @@ -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 `_.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 }, +} diff --git a/litellm-rust/crates/migrate-macros/src/lib.rs b/litellm-rust/crates/migrate-macros/src/lib.rs new file mode 100644 index 00000000000..501f59e6fc2 --- /dev/null +++ b/litellm-rust/crates/migrate-macros/src/lib.rs @@ -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, 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::().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, 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 = 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 { .. }))); + } +} diff --git a/litellm-rust/crates/migrate/Cargo.toml b/litellm-rust/crates/migrate/Cargo.toml new file mode 100644 index 00000000000..bb1ecaa3128 --- /dev/null +++ b/litellm-rust/crates/migrate/Cargo.toml @@ -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 diff --git a/litellm-rust/crates/migrate/README.md b/litellm-rust/crates/migrate/README.md new file mode 100644 index 00000000000..4817029451c --- /dev/null +++ b/litellm-rust/crates/migrate/README.md @@ -0,0 +1,5 @@ +# Migrations + +`litellm-migrate` exports the `Migration` struct and the `migrate!` macro that embeds a directory of `_.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 diff --git a/litellm-rust/crates/migrate/src/lib.rs b/litellm-rust/crates/migrate/src/lib.rs new file mode 100644 index 00000000000..f4e065e1b53 --- /dev/null +++ b/litellm-rust/crates/migrate/src/lib.rs @@ -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, +} diff --git a/litellm-rust/crates/migrate/tests/fixtures/migrations/10_tenth.sql b/litellm-rust/crates/migrate/tests/fixtures/migrations/10_tenth.sql new file mode 100644 index 00000000000..31807719e9c --- /dev/null +++ b/litellm-rust/crates/migrate/tests/fixtures/migrations/10_tenth.sql @@ -0,0 +1 @@ +SELECT 10; diff --git a/litellm-rust/crates/migrate/tests/fixtures/migrations/1_first.sql b/litellm-rust/crates/migrate/tests/fixtures/migrations/1_first.sql new file mode 100644 index 00000000000..e0ac49d1ecf --- /dev/null +++ b/litellm-rust/crates/migrate/tests/fixtures/migrations/1_first.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/litellm-rust/crates/migrate/tests/fixtures/migrations/2_second.sql b/litellm-rust/crates/migrate/tests/fixtures/migrations/2_second.sql new file mode 100644 index 00000000000..e7f8100648d --- /dev/null +++ b/litellm-rust/crates/migrate/tests/fixtures/migrations/2_second.sql @@ -0,0 +1 @@ +SELECT 2; diff --git a/litellm-rust/crates/migrate/tests/migrate.rs b/litellm-rust/crates/migrate/tests/migrate.rs new file mode 100644 index 00000000000..61c80351cf4 --- /dev/null +++ b/litellm-rust/crates/migrate/tests/migrate.rs @@ -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); +} diff --git a/litellm-rust/crates/traces/AGENTS.md b/litellm-rust/crates/traces/AGENTS.md index 645e88dfae1..d0181c53308 100644 --- a/litellm-rust/crates/traces/AGENTS.md +++ b/litellm-rust/crates/traces/AGENTS.md @@ -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 diff --git a/litellm-rust/crates/traces/Cargo.toml b/litellm-rust/crates/traces/Cargo.toml index 74de400764c..41b6037d015 100644 --- a/litellm-rust/crates/traces/Cargo.toml +++ b/litellm-rust/crates/traces/Cargo.toml @@ -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"] } diff --git a/litellm-rust/crates/traces/build.rs b/litellm-rust/crates/traces/build.rs new file mode 100644 index 00000000000..3a8149ef075 --- /dev/null +++ b/litellm-rust/crates/traces/build.rs @@ -0,0 +1,3 @@ +fn main() { + println!("cargo:rerun-if-changed=migrations"); +} diff --git a/litellm-rust/crates/traces/src/schema.rs b/litellm-rust/crates/traces/src/schema.rs index 4943f00f7c9..cf11c564931 100644 --- a/litellm-rust/crates/traces/src/schema.rs +++ b/litellm-rust/crates/traces/src/schema.rs @@ -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}", diff --git a/litellm-rust/crates/traces/tests/migrations.rs b/litellm-rust/crates/traces/tests/migrations.rs index 238c5e671c6..1d21010b0d8 100644 --- a/litellm-rust/crates/traces/tests/migrations.rs +++ b/litellm-rust/crates/traces/tests/migrations.rs @@ -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(