mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-03 02:22:24 +00:00
fix(tracing): unify ClickHouse storage configuration (#43941)
* fix(tracing): use ClickHouse URL for reads by default * fix(tracing): unify ClickHouse storage configuration * fix(tracing): update dashboard setup copy for one URL * test(tracing): make tests/unit/tracing a package Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(config): drop legacy string tracing store variant Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(tracing): own ClickHouse defaults in constants and reject unset env references Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(tracing): use raw regex patterns in config tests Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(tracing): read ClickHouse env defaults when tracing config resolves Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(ui): split audit log query guard to fit condition-chain budget Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
a5fef4b4e6
commit
276fc9c63a
37 changed files with 614 additions and 191 deletions
|
|
@ -4,7 +4,20 @@ Lens reviews recorded activity and saves evidence-linked findings in the LiteLLM
|
|||
|
||||
## Start a worker
|
||||
|
||||
Upgrade your existing LiteLLM proxy to a release that includes Lens with PostgreSQL, agent tracing (`general_settings.tracing: {store: clickhouse}`), and ClickHouse configured through `CLICKHOUSE_URL` and a separate SELECT-only `CLICKHOUSE_READER_URL`. Enable the ClickHouse callback and request/response logging to analyze LLM requests. Lens can only inspect content you actually retain
|
||||
Upgrade your existing LiteLLM proxy to a release that includes Lens with PostgreSQL and agent tracing. Configure one ClickHouse URL for trace writes, bounded reads, and Lens queries:
|
||||
|
||||
```yaml
|
||||
general_settings:
|
||||
tracing:
|
||||
store:
|
||||
type: clickhouse
|
||||
url: os.environ/CLICKHOUSE_URL
|
||||
retention_days: 14
|
||||
```
|
||||
|
||||
The URL, database, and retention settings can also come from `CLICKHOUSE_URL`, `CLICKHOUSE_DATABASE`, and `AGENT_TRACING_RETENTION_DAYS` when omitted from YAML. A YAML value wins when both are set. The database defaults to `litellm`. `retention_days` defaults to 14 and applies to both traces and spend logs
|
||||
|
||||
Retention changes require a proxy restart. ClickHouse removes expired rows during background merges, not immediately at startup. Enable request/response logging to analyze LLM requests. Lens can only inspect content you actually retain
|
||||
|
||||
In Lens, click **Set up analysis**, choose an existing virtual key or **Create worker key**, then **Generate setup command**. The LiteLLM address is filled in for you; change it only if the server running Docker needs a different network address. Copy the command and run it on your server. The dialog changes to **Analyzer connected** when the container checks in
|
||||
|
||||
|
|
|
|||
|
|
@ -12,7 +12,6 @@ services:
|
|||
DATABASE_URL: postgresql://litellm:litellm@db:5432/litellm
|
||||
STORE_MODEL_IN_DB: "True"
|
||||
CLICKHOUSE_URL: http://default:local-tracing@clickhouse:8123
|
||||
CLICKHOUSE_READER_URL: http://default:local-tracing@clickhouse:8123
|
||||
CLICKHOUSE_DATABASE: litellm
|
||||
OPENAI_API_KEY: ${OPENAI_API_KEY:-}
|
||||
volumes:
|
||||
|
|
|
|||
|
|
@ -7,4 +7,7 @@ model_list:
|
|||
general_settings:
|
||||
master_key: os.environ/LITELLM_MASTER_KEY
|
||||
tracing:
|
||||
store: clickhouse
|
||||
store:
|
||||
type: clickhouse
|
||||
url: os.environ/CLICKHOUSE_URL
|
||||
retention_days: 14
|
||||
|
|
|
|||
|
|
@ -12,7 +12,10 @@ use serde::Deserialize;
|
|||
pub use error::Error;
|
||||
pub use mcp::{McpAuth, McpServer, McpTransport};
|
||||
pub use model::{LiteLlmParams, Model};
|
||||
pub use settings::{GeneralSettings, LiteLlmSettings, RouterSettings};
|
||||
pub use settings::{
|
||||
ClickHouseStoreSettings, GeneralSettings, LiteLlmSettings, RouterSettings, TracingSettings,
|
||||
TracingStoreSettings,
|
||||
};
|
||||
pub use value::{AdditionalFields, Flag, NumberOrString, Object, OneOrMany, Value};
|
||||
|
||||
#[derive(Clone, Default, Deserialize)]
|
||||
|
|
|
|||
|
|
@ -5,6 +5,47 @@ use serde::Deserialize;
|
|||
|
||||
use crate::{AdditionalFields, Flag, NumberOrString, Object, OneOrMany, Value};
|
||||
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(rename_all = "lowercase")]
|
||||
pub enum TracingStoreKind {
|
||||
Clickhouse,
|
||||
}
|
||||
|
||||
#[derive(Clone, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct ClickHouseStoreSettings {
|
||||
#[serde(rename = "type")]
|
||||
pub kind: TracingStoreKind,
|
||||
pub url: Option<SecretValue>,
|
||||
pub database: Option<String>,
|
||||
pub retention_days: Option<NumberOrString>,
|
||||
}
|
||||
|
||||
impl fmt::Debug for ClickHouseStoreSettings {
|
||||
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
formatter
|
||||
.debug_struct("ClickHouseStoreSettings")
|
||||
.field("kind", &self.kind)
|
||||
.field("database", &self.database)
|
||||
.field("retention_days", &self.retention_days)
|
||||
.finish()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(untagged)]
|
||||
pub enum TracingStoreSettings {
|
||||
ClickHouse(ClickHouseStoreSettings),
|
||||
}
|
||||
|
||||
#[derive(Clone, Default, Debug, Deserialize)]
|
||||
#[serde(default)]
|
||||
pub struct TracingSettings {
|
||||
pub store: Option<TracingStoreSettings>,
|
||||
#[serde(flatten)]
|
||||
pub additional_fields: AdditionalFields,
|
||||
}
|
||||
|
||||
#[derive(Clone, Deserialize)]
|
||||
#[serde(default)]
|
||||
pub struct GeneralSettings {
|
||||
|
|
@ -14,6 +55,7 @@ pub struct GeneralSettings {
|
|||
pub admission_queue_timeout_seconds: f64,
|
||||
pub master_key: Option<SecretValue>,
|
||||
pub database_url: Option<SecretValue>,
|
||||
pub tracing: Option<TracingSettings>,
|
||||
pub database_connection_pool_limit: Option<u64>,
|
||||
pub database_connection_timeout: Option<f64>,
|
||||
pub database_connect_timeout: Option<f64>,
|
||||
|
|
@ -50,6 +92,7 @@ impl Default for GeneralSettings {
|
|||
admission_queue_timeout_seconds: 1.0,
|
||||
master_key: None,
|
||||
database_url: None,
|
||||
tracing: None,
|
||||
database_connection_pool_limit: Some(10),
|
||||
database_connection_timeout: Some(60.0),
|
||||
database_connect_timeout: None,
|
||||
|
|
@ -97,6 +140,7 @@ impl fmt::Debug for GeneralSettings {
|
|||
)
|
||||
.field("master_key", &self.master_key)
|
||||
.field("database_url", &self.database_url)
|
||||
.field("tracing", &self.tracing)
|
||||
.field("store_model_in_db", &self.store_model_in_db)
|
||||
.field("additional_fields", &self.additional_fields.keys())
|
||||
.finish_non_exhaustive()
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
use litellm_config::{Config, Error, Flag, NumberOrString};
|
||||
use litellm_config::{Config, Error, Flag, NumberOrString, TracingStoreSettings};
|
||||
use rstest::{fixture, rstest};
|
||||
use tempfile::TempDir;
|
||||
|
||||
|
|
@ -113,6 +113,54 @@ fn missing_general_settings_has_no_master_key() {
|
|||
assert!(config.general_settings.master_key.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tracing_settings_are_typed_and_redact_the_url() {
|
||||
let config = Config::from_yaml(
|
||||
"general_settings:\n tracing:\n store:\n type: clickhouse\n url: https://writer:password@example.com\n database: analytics\n retention_days: 7\n",
|
||||
)
|
||||
.unwrap();
|
||||
let tracing = config.general_settings.tracing.as_ref().unwrap();
|
||||
let Some(TracingStoreSettings::ClickHouse(store)) = tracing.store.as_ref() else {
|
||||
panic!("expected ClickHouse tracing store")
|
||||
};
|
||||
assert_eq!(
|
||||
store.url.as_ref().unwrap().expose(),
|
||||
"https://writer:password@example.com"
|
||||
);
|
||||
assert_eq!(store.database.as_deref(), Some("analytics"));
|
||||
assert_eq!(store.retention_days, Some(NumberOrString::Number(7.0)));
|
||||
assert!(!format!("{config:?}").contains("password"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tracing_settings_accept_environment_references() {
|
||||
let config = Config::from_yaml(
|
||||
"general_settings:\n tracing:\n store:\n type: clickhouse\n url: os.environ/CLICKHOUSE_URL\n retention_days: os.environ/RETENTION_DAYS\n",
|
||||
)
|
||||
.unwrap();
|
||||
let Some(TracingStoreSettings::ClickHouse(store)) =
|
||||
config.general_settings.tracing.unwrap().store
|
||||
else {
|
||||
panic!("expected ClickHouse tracing store")
|
||||
};
|
||||
assert_eq!(
|
||||
store.retention_days,
|
||||
Some(NumberOrString::String(
|
||||
"os.environ/RETENTION_DAYS".to_owned()
|
||||
))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tracing_settings_reject_string_store() {
|
||||
assert!(Config::from_yaml("general_settings:\n tracing:\n store: clickhouse\n").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tracing_settings_reject_removed_reader_configuration() {
|
||||
assert!(Config::from_yaml("general_settings:\n tracing:\n store:\n type: clickhouse\n reader_url: http://localhost:8123\n").is_err());
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
fn empty_config_matches_python_defaults() {
|
||||
let config = Config::from_yaml("{}").unwrap();
|
||||
|
|
|
|||
|
|
@ -45,7 +45,7 @@ mod _native {
|
|||
use crate::routes::token_counter::TokenCounter;
|
||||
#[pymodule_export]
|
||||
use crate::routes::traces::{
|
||||
NativeTraceStorage, trace_decode_otlp, trace_encode_error,
|
||||
NativeTraceConfig, NativeTraceStorage, trace_decode_otlp, trace_encode_error,
|
||||
trace_normalized_field_definitions,
|
||||
};
|
||||
#[cfg(feature = "huggingface")]
|
||||
|
|
@ -112,6 +112,7 @@ mod tests {
|
|||
"aresponses",
|
||||
"ResponsesWebSocketConnection",
|
||||
"NativeDiagnosticProcessor",
|
||||
"NativeTraceConfig",
|
||||
"NativeTraceStorage",
|
||||
"trace_decode_otlp",
|
||||
"trace_encode_error",
|
||||
|
|
|
|||
|
|
@ -2,9 +2,9 @@ use std::collections::BTreeMap;
|
|||
|
||||
use litellm_host_python::{FromPythonCache, ToPythonCache};
|
||||
use litellm_http::ClientVariant;
|
||||
use litellm_storage_clickhouse::Storage;
|
||||
use litellm_traces::{
|
||||
Error, InsertTable, Parameter, QueryAccessError, QueryReaders, QueryScope, ReadQuery, Shared,
|
||||
Config, Error, InsertTable, Parameter, QueryAccessError, QueryReaders, QueryScope, ReadQuery,
|
||||
Shared,
|
||||
};
|
||||
use prost::Message;
|
||||
use pyo3::{
|
||||
|
|
@ -63,48 +63,49 @@ fn map_query_access_error(error: QueryAccessError) -> PyErr {
|
|||
}
|
||||
}
|
||||
|
||||
#[pyclass(frozen)]
|
||||
pub struct NativeTraceConfig {
|
||||
inner: Config,
|
||||
}
|
||||
|
||||
#[pymethods]
|
||||
impl NativeTraceConfig {
|
||||
#[new]
|
||||
fn new(database: String, url: &str, retention_days: u32) -> PyResult<Self> {
|
||||
Ok(Self {
|
||||
inner: Config::new(database, url, retention_days).map_err(map_error)?,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[pyclass]
|
||||
pub struct NativeTraceStorage {
|
||||
storage: Storage,
|
||||
config: Config,
|
||||
query_readers: QueryReaders,
|
||||
}
|
||||
|
||||
#[pymethods]
|
||||
impl NativeTraceStorage {
|
||||
#[new]
|
||||
#[pyo3(signature = (database, url, reader_url = None))]
|
||||
fn new(database: String, url: &str, reader_url: Option<&str>) -> PyResult<Self> {
|
||||
litellm_traces::schema_statements(&database, 1, 1).map_err(map_error)?;
|
||||
let storage = Storage::new(database, url, reader_url).map_err(map_error)?;
|
||||
fn new(config: PyRef<'_, NativeTraceConfig>) -> PyResult<Self> {
|
||||
Ok(Self {
|
||||
query_readers: QueryReaders::new(
|
||||
storage.writer().clone(),
|
||||
storage.database().to_owned(),
|
||||
config.inner.storage().writer().clone(),
|
||||
config.inner.storage().database().to_owned(),
|
||||
),
|
||||
storage,
|
||||
config: config.inner.clone(),
|
||||
})
|
||||
}
|
||||
|
||||
fn ensure_schema<'py>(
|
||||
&self,
|
||||
py: Python<'py>,
|
||||
trace_retention_days: u32,
|
||||
spend_log_retention_days: u32,
|
||||
) -> PyResult<Bound<'py, PyAny>> {
|
||||
fn ensure_schema<'py>(&self, py: Python<'py>) -> PyResult<Bound<'py, PyAny>> {
|
||||
let client = crate::http::host_client(py, ClientVariant::NoRedirect)?;
|
||||
let connection = self.storage.writer().clone();
|
||||
let database = self.storage.database().to_owned();
|
||||
let connection = self.config.storage().writer().clone();
|
||||
let database = self.config.storage().database().to_owned();
|
||||
let retention_days = self.config.retention_days();
|
||||
crate::execution::run_async(
|
||||
py,
|
||||
async move {
|
||||
litellm_traces::ensure_schema(
|
||||
&client,
|
||||
&connection,
|
||||
&database,
|
||||
trace_retention_days,
|
||||
spend_log_retention_days,
|
||||
)
|
||||
.await
|
||||
litellm_traces::ensure_schema(&client, &connection, &database, retention_days).await
|
||||
},
|
||||
map_error,
|
||||
)
|
||||
|
|
@ -118,8 +119,8 @@ impl NativeTraceStorage {
|
|||
) -> PyResult<Bound<'py, PyAny>> {
|
||||
let table = InsertTable::parse(table).map_err(map_error)?;
|
||||
let client = crate::http::host_client(py, ClientVariant::NoRedirect)?;
|
||||
let connection = self.storage.writer().clone();
|
||||
let database = self.storage.database().to_owned();
|
||||
let connection = self.config.storage().writer().clone();
|
||||
let database = self.config.storage().database().to_owned();
|
||||
crate::execution::run_async(
|
||||
py,
|
||||
async move {
|
||||
|
|
@ -186,9 +187,7 @@ impl NativeTraceStorage {
|
|||
>,
|
||||
) -> PyResult<Bound<'py, PyAny>> {
|
||||
let query = litellm_traces::LensQuery::parse(name).map_err(map_error)?;
|
||||
let connection = self.storage.reader().cloned().ok_or_else(|| {
|
||||
PyRuntimeError::new_err("Trace reads require a separate ClickHouse reader URL")
|
||||
})?;
|
||||
let connection = self.config.storage().reader().clone();
|
||||
let client = crate::http::host_client(py, ClientVariant::NoRedirect)?;
|
||||
crate::execution::run_async(
|
||||
py,
|
||||
|
|
@ -209,9 +208,7 @@ impl NativeTraceStorage {
|
|||
>,
|
||||
) -> PyResult<Bound<'py, PyAny>> {
|
||||
let query = ReadQuery::parse(query).map_err(map_error)?;
|
||||
let connection = self.storage.reader().cloned().ok_or_else(|| {
|
||||
PyRuntimeError::new_err("Trace reads require a separate ClickHouse reader URL")
|
||||
})?;
|
||||
let connection = self.config.storage().reader().clone();
|
||||
let client = crate::http::host_client(py, ClientVariant::NoRedirect)?;
|
||||
crate::execution::run_async(
|
||||
py,
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
# ClickHouse storage
|
||||
|
||||
`litellm-storage-clickhouse` exports `Storage`, a shared writer connection and optional reader connection for one ClickHouse database. It also exports bounded HTTP read and insert execution
|
||||
`litellm-storage-clickhouse` exports `Storage`, a writer and bounded reader derived from one ClickHouse URL and database. It also exports bounded HTTP read and insert execution
|
||||
|
||||
The crate has no trace tables, OTLP types, or named trace queries. `litellm-traces` supplies those rules and uses this storage for both trace rows and spend rows
|
||||
|
|
|
|||
|
|
@ -89,19 +89,17 @@ impl Connection {
|
|||
pub struct Storage {
|
||||
database: String,
|
||||
writer: Connection,
|
||||
reader: Option<Connection>,
|
||||
reader: Connection,
|
||||
}
|
||||
|
||||
impl Storage {
|
||||
pub fn new(database: String, url: &str, reader_url: Option<&str>) -> Result<Self, Error> {
|
||||
pub fn new(database: String, url: &str) -> Result<Self, Error> {
|
||||
if !valid_identifier(&database) {
|
||||
return Err(Error::InvalidSchema);
|
||||
}
|
||||
Ok(Self {
|
||||
writer: Connection::writer(url)?,
|
||||
reader: reader_url
|
||||
.map(|value| Connection::reader(value, &database))
|
||||
.transpose()?,
|
||||
reader: Connection::reader(url, &database)?,
|
||||
database,
|
||||
})
|
||||
}
|
||||
|
|
@ -114,8 +112,8 @@ impl Storage {
|
|||
&self.writer
|
||||
}
|
||||
|
||||
pub fn reader(&self) -> Option<&Connection> {
|
||||
self.reader.as_ref()
|
||||
pub fn reader(&self) -> &Connection {
|
||||
&self.reader
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -10,25 +10,30 @@ fn accepts_only_clickhouse_http_urls(#[case] value: &str, #[case] expected: bool
|
|||
assert_eq!(Connection::parse(value).is_ok(), expected);
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::writer_only(None, false)]
|
||||
#[case::separate_reader(Some("http://localhost:8124"), true)]
|
||||
fn storage_exports_writer_and_optional_reader(
|
||||
#[case] reader_url: Option<&str>,
|
||||
#[case] has_reader: bool,
|
||||
) {
|
||||
let storage = Storage::new("litellm".to_owned(), "http://localhost:8123", reader_url)
|
||||
.expect("valid ClickHouse URLs");
|
||||
#[test]
|
||||
fn storage_uses_one_url_for_writes_and_bounded_reads() {
|
||||
let storage =
|
||||
Storage::new("litellm".to_owned(), "http://localhost:8123").expect("valid ClickHouse URLs");
|
||||
|
||||
assert_eq!(storage.database(), "litellm");
|
||||
assert_eq!(storage.writer().url().host_str(), Some("localhost"));
|
||||
assert_eq!(storage.writer().url().port(), Some(8123));
|
||||
assert_eq!(storage.reader().is_some(), has_reader);
|
||||
assert_eq!(storage.reader().url().port(), Some(8123));
|
||||
assert_eq!(
|
||||
storage
|
||||
.reader()
|
||||
.url()
|
||||
.query_pairs()
|
||||
.find(|(key, _)| key == "database")
|
||||
.unwrap()
|
||||
.1,
|
||||
"litellm"
|
||||
);
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::empty("")]
|
||||
#[case::injection("db; DROP DATABASE default")]
|
||||
fn storage_rejects_invalid_database(#[case] database: &str) {
|
||||
assert!(Storage::new(database.to_owned(), "http://localhost:8123", None).is_err());
|
||||
assert!(Storage::new(database.to_owned(), "http://localhost:8123").is_err());
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1 +1 @@
|
|||
ALTER TABLE {database}.otel_traces MODIFY TTL toDateTime(Timestamp) + INTERVAL {trace_retention_days} DAY
|
||||
ALTER TABLE {database}.otel_traces MODIFY TTL toDateTime(Timestamp) + INTERVAL {retention_days} DAY
|
||||
|
|
|
|||
|
|
@ -1 +1 @@
|
|||
ALTER TABLE {database}.agent_traces_by_key MODIFY TTL toDateTime(StartTs) + INTERVAL {trace_retention_days} DAY
|
||||
ALTER TABLE {database}.agent_traces_by_key MODIFY TTL toDateTime(StartTs) + INTERVAL {retention_days} DAY
|
||||
|
|
|
|||
|
|
@ -1 +1 @@
|
|||
ALTER TABLE {database}.spend_logs MODIFY TTL toDateTime(start_time) + INTERVAL {spend_log_retention_days} DAY
|
||||
ALTER TABLE {database}.spend_logs MODIFY TTL toDateTime(start_time) + INTERVAL {retention_days} DAY
|
||||
|
|
|
|||
25
litellm-rust/crates/traces/src/config.rs
Normal file
25
litellm-rust/crates/traces/src/config.rs
Normal file
|
|
@ -0,0 +1,25 @@
|
|||
use litellm_storage_clickhouse::{Error, Storage};
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Config {
|
||||
storage: Storage,
|
||||
retention_days: u32,
|
||||
}
|
||||
|
||||
impl Config {
|
||||
pub fn new(database: String, url: &str, retention_days: u32) -> Result<Self, Error> {
|
||||
crate::schema_statements(&database, retention_days)?;
|
||||
Ok(Self {
|
||||
storage: Storage::new(database, url)?,
|
||||
retention_days,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn storage(&self) -> &Storage {
|
||||
&self.storage
|
||||
}
|
||||
|
||||
pub fn retention_days(&self) -> u32 {
|
||||
self.retention_days
|
||||
}
|
||||
}
|
||||
|
|
@ -1,3 +1,4 @@
|
|||
mod config;
|
||||
mod error;
|
||||
mod insert;
|
||||
mod normalize;
|
||||
|
|
@ -8,6 +9,7 @@ mod schema;
|
|||
mod shared;
|
||||
mod sql;
|
||||
|
||||
pub use config::Config;
|
||||
pub use error::{DecodeError, QueryAccessError};
|
||||
pub use insert::{InsertRow, InsertTable, encode_rows, insert_rows, insert_shared_rows};
|
||||
pub use litellm_storage_clickhouse::{Connection, Error, Parameter, execute_read};
|
||||
|
|
|
|||
|
|
@ -9,17 +9,12 @@ const SCHEMA_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
|
|||
|
||||
const MIGRATIONS: &[Migration] = litellm_migrate::migrate!("migrations");
|
||||
|
||||
pub fn schema_statements(
|
||||
database: &str,
|
||||
trace_retention_days: u32,
|
||||
spend_log_retention_days: u32,
|
||||
) -> Result<Vec<String>, Error> {
|
||||
pub fn schema_statements(database: &str, retention_days: u32) -> Result<Vec<String>, Error> {
|
||||
if database.is_empty()
|
||||
|| !database
|
||||
.bytes()
|
||||
.all(|c| c.is_ascii_alphanumeric() || c == b'_')
|
||||
|| trace_retention_days == 0
|
||||
|| spend_log_retention_days == 0
|
||||
|| retention_days == 0
|
||||
{
|
||||
return Err(Error::InvalidSchema);
|
||||
}
|
||||
|
|
@ -30,11 +25,7 @@ pub fn schema_statements(
|
|||
migration
|
||||
.sql
|
||||
.replace("{database}", &database)
|
||||
.replace("{trace_retention_days}", &trace_retention_days.to_string())
|
||||
.replace(
|
||||
"{spend_log_retention_days}",
|
||||
&spend_log_retention_days.to_string(),
|
||||
)
|
||||
.replace("{retention_days}", &retention_days.to_string())
|
||||
}))
|
||||
.collect(),
|
||||
)
|
||||
|
|
@ -44,15 +35,13 @@ pub async fn ensure_schema(
|
|||
client: &Client,
|
||||
connection: &Connection,
|
||||
database: &str,
|
||||
trace_retention_days: u32,
|
||||
spend_log_retention_days: u32,
|
||||
retention_days: u32,
|
||||
) -> Result<(), Error> {
|
||||
ensure_schema_with_timeout(
|
||||
client,
|
||||
connection,
|
||||
database,
|
||||
trace_retention_days,
|
||||
spend_log_retention_days,
|
||||
retention_days,
|
||||
SCHEMA_REQUEST_TIMEOUT,
|
||||
)
|
||||
.await
|
||||
|
|
@ -62,11 +51,10 @@ async fn ensure_schema_with_timeout(
|
|||
client: &Client,
|
||||
connection: &Connection,
|
||||
database: &str,
|
||||
trace_retention_days: u32,
|
||||
spend_log_retention_days: u32,
|
||||
retention_days: u32,
|
||||
request_timeout: Duration,
|
||||
) -> Result<(), Error> {
|
||||
for statement in schema_statements(database, trace_retention_days, spend_log_retention_days)? {
|
||||
for statement in schema_statements(database, retention_days)? {
|
||||
let response = client
|
||||
.post(connection.url().clone())
|
||||
.timeout(request_timeout)
|
||||
|
|
|
|||
|
|
@ -109,8 +109,8 @@ async fn schema_supports_span_rollups_and_spend_joins(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
|
||||
let span = serde_json::from_value(serde_json::json!({
|
||||
"Timestamp": timestamp, "TraceId": "trace-1", "SpanId": "span-1", "ParentSpanId": "",
|
||||
|
|
@ -221,7 +221,6 @@ async fn normalized_fields_match_clickhouse_catalog(
|
|||
&Connection::writer(&database.url)?,
|
||||
"trace_test",
|
||||
7,
|
||||
14,
|
||||
)
|
||||
.await?;
|
||||
let catalog = read_json(&database, "SELECT name, type FROM system.columns WHERE database = 'trace_test' AND table = 'otel_traces'").await?;
|
||||
|
|
@ -257,7 +256,7 @@ async fn insert_rejects_unknown_columns_even_if_url_requests_skipping_them(
|
|||
"{}?input_format_skip_unknown_fields=1",
|
||||
database.url
|
||||
))?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let row = BTreeMap::from([
|
||||
(
|
||||
"Timestamp".to_owned(),
|
||||
|
|
@ -291,7 +290,7 @@ async fn retried_trace_insert_does_not_inflate_rollup(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let row: BTreeMap<String, serde_json::Value> = serde_json::from_value(serde_json::json!({
|
||||
"Timestamp": time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64,
|
||||
"TraceId": "retried-trace", "SpanId": "span-1", "ParentSpanId": "",
|
||||
|
|
@ -326,7 +325,7 @@ async fn keyed_rollup_keeps_same_trace_ids_separate_by_api_key(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
|
||||
let rows = vec![
|
||||
serde_json::from_value(serde_json::json!({
|
||||
|
|
@ -370,7 +369,7 @@ async fn rollup_merges_spans_across_days_without_losing_root_fields(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let day_start = time::OffsetDateTime::now_utc()
|
||||
.replace_time(time::Time::MIDNIGHT)
|
||||
.unix_timestamp_nanos() as i64;
|
||||
|
|
@ -417,7 +416,7 @@ async fn spend_deduplication_preserves_subsecond_requests_and_retries(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let now_ms = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64 / 1_000_000;
|
||||
let base_start_time = now_ms / 1000 * 1000;
|
||||
let first_start_time = base_start_time + 100;
|
||||
|
|
@ -469,7 +468,7 @@ async fn retention_changes_materialize_existing_rows_and_remain_idempotent(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 30, 30).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 30).await?;
|
||||
let tables = read_json(
|
||||
&database,
|
||||
"SELECT name FROM system.tables WHERE database = 'trace_test' \
|
||||
|
|
@ -499,7 +498,7 @@ async fn retention_changes_materialize_existing_rows_and_remain_idempotent(
|
|||
insert_rows(&database, "otel_traces", vec![span]).await?;
|
||||
insert_rows(&database, "spend_logs", vec![spend]).await?;
|
||||
assert_eq!(table_rows(&database, "agent_traces_by_key").await?, 1);
|
||||
ensure_schema(&database.client, &writer, "trace_test", 14, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 14).await?;
|
||||
let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
|
||||
loop {
|
||||
let response = read_json(
|
||||
|
|
@ -531,7 +530,7 @@ async fn retention_changes_materialize_existing_rows_and_remain_idempotent(
|
|||
assert_eq!(table_rows(&database, "agent_traces_by_key").await?, 0);
|
||||
assert_eq!(table_rows(&database, "spend_logs").await?, 0);
|
||||
let mutation_count = mutation_rows(&database).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 14, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 14).await?;
|
||||
assert_eq!(mutation_rows(&database).await?, mutation_count);
|
||||
Ok(())
|
||||
}
|
||||
|
|
@ -550,7 +549,7 @@ async fn schema_statement_timeout_maps_to_transport_error() -> TestResult {
|
|||
let writer = Connection::writer(&url)?;
|
||||
let result = tokio::time::timeout(
|
||||
Duration::from_secs(35),
|
||||
ensure_schema(&client, &writer, "trace_test", 7, 14),
|
||||
ensure_schema(&client, &writer, "trace_test", 7),
|
||||
)
|
||||
.await;
|
||||
server.abort();
|
||||
|
|
@ -559,16 +558,11 @@ async fn schema_statement_timeout_maps_to_transport_error() -> TestResult {
|
|||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::empty("", 7, 14)]
|
||||
#[case::sql("db; DROP DATABASE default", 7, 14)]
|
||||
#[case::trace_retention("traces", 0, 14)]
|
||||
#[case::spend_retention("traces", 7, 0)]
|
||||
fn schema_rejects_invalid_configuration(
|
||||
#[case] database: &str,
|
||||
#[case] traces: u32,
|
||||
#[case] spend: u32,
|
||||
) {
|
||||
assert!(schema_statements(database, traces, spend).is_err());
|
||||
#[case::empty("", 7)]
|
||||
#[case::sql("db; DROP DATABASE default", 7)]
|
||||
#[case::retention("traces", 0)]
|
||||
fn schema_rejects_invalid_configuration(#[case] database: &str, #[case] retention_days: u32) {
|
||||
assert!(schema_statements(database, retention_days).is_err());
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
|
|
@ -579,7 +573,7 @@ async fn lens_filters_reads_and_evidence_keep_reused_trace_ids_separate(
|
|||
use litellm_traces::{LensQuery, Parameter};
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
|
||||
for (key, text) in [("one", "timeout"), ("two", "success")] {
|
||||
insert_rows(&database, "otel_traces", vec![serde_json::from_value(serde_json::json!({
|
||||
|
|
@ -686,7 +680,7 @@ async fn lens_request_sample_does_not_trust_caller_tags(
|
|||
use litellm_traces::{LensQuery, Parameter};
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64 / 1_000_000;
|
||||
for (id, internal) in [("external", false), ("internal", true)] {
|
||||
let row = serde_json::from_value(serde_json::json!({
|
||||
|
|
@ -755,7 +749,6 @@ async fn lens_selection_pages_without_losing_or_repeating_runs(
|
|||
&Connection::writer(&database.url)?,
|
||||
"trace_test",
|
||||
7,
|
||||
14,
|
||||
)
|
||||
.await?;
|
||||
execute_write(&database, "INSERT INTO trace_test.spend_logs (request_id,team_id,start_time,end_time) SELECT toString(number),'team',now64(3)-INTERVAL 5 MINUTE,now64(3)-INTERVAL 5 MINUTE FROM numbers(1001)").await?;
|
||||
|
|
@ -836,7 +829,6 @@ async fn lens_content_keeps_output_visible_after_long_input(
|
|||
&Connection::writer(&database.url)?,
|
||||
"trace_test",
|
||||
7,
|
||||
14,
|
||||
)
|
||||
.await?;
|
||||
insert_rows(&database, "spend_logs", vec![serde_json::from_value(serde_json::json!({
|
||||
|
|
@ -902,7 +894,7 @@ async fn trace_error_previews_preserve_paginated_diagnostics(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
|
||||
let rows = (0..span_count)
|
||||
.map(|index| {
|
||||
|
|
@ -992,7 +984,7 @@ async fn duplicate_span_preview_matches_diagnostic(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
|
||||
let message = "a".repeat(200);
|
||||
let rows = [
|
||||
|
|
@ -1041,7 +1033,7 @@ fn schema_includes_every_migration_file() -> TestResult {
|
|||
.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);
|
||||
assert_eq!(schema_statements("trace_test", 7)?.len(), 1 + files);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
|
@ -1053,7 +1045,7 @@ async fn lens_agent_discovery_and_selection_preserve_scope(
|
|||
use litellm_traces::LensQuery;
|
||||
let database = database.await?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
|
||||
for (team, key, trace, agent, span, parent) in [
|
||||
("alpha", "one", "research", "research_agent", "root", ""),
|
||||
|
|
@ -1160,7 +1152,7 @@ async fn query_help_discovers_live_schema_and_runs_its_examples(
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
execute_write(&database, "CREATE USER help_reader").await?;
|
||||
for table in ["otel_traces", "agent_traces_by_key", "spend_logs"] {
|
||||
execute_write(
|
||||
|
|
@ -1371,7 +1363,7 @@ async fn query_help_preserves_schema_and_guide_when_discovery_hits_reader_limits
|
|||
) -> TestResult {
|
||||
let database = database?;
|
||||
let writer = Connection::writer(&database.url)?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
|
||||
ensure_schema(&database.client, &writer, "trace_test", 7).await?;
|
||||
execute_write(
|
||||
&database,
|
||||
"CREATE USER help_reader SETTINGS max_rows_to_read = 1",
|
||||
|
|
|
|||
|
|
@ -33,7 +33,7 @@ async fn database() -> Result<Database, Box<dyn std::error::Error>> {
|
|||
container.get_host_port_ipv4(8123).await?
|
||||
))?;
|
||||
let client = Client::no_redirect_for_test();
|
||||
ensure_schema(&client, &writer, "trace_test", 7, 7).await?;
|
||||
ensure_schema(&client, &writer, "trace_test", 7).await?;
|
||||
for sql in [
|
||||
"INSERT INTO trace_test.otel_traces (TeamId, ApiKeyHash, TraceId, SpanId, Timestamp, SpanAttributes) VALUES ('team-a', 'key-a1', 'shared-trace', 'a1', now(), map('visible', 'a')), ('team-a', 'key-a2', 'shared-trace', 'a2', now(), map('visible', 'a')), ('team-b', 'key-b', 'shared-trace', 'b', now(), map('secret-b', 'b'))",
|
||||
"INSERT INTO trace_test.spend_logs (team_id, api_key, request_id, start_time, end_time, metadata) VALUES ('team-a', 'key-a1', 'a1', now(), now(), '{\"visible\":1}'), ('team-a', 'key-a2', 'a2', now(), now(), '{\"visible\":1}'), ('team-b', 'key-b', 'b', now(), now(), '{\"secret_b\":1}')",
|
||||
|
|
|
|||
|
|
@ -50,8 +50,8 @@ CLICKHOUSE_BATCH_SIZE: Final = get_env_int("CLICKHOUSE_BATCH_SIZE", 10_000)
|
|||
CLICKHOUSE_FLUSH_INTERVAL_SECONDS: Final = float(os.getenv("CLICKHOUSE_FLUSH_INTERVAL_SECONDS", "1.0"))
|
||||
CLICKHOUSE_MAX_BUFFERED_ROWS: Final = get_env_int("CLICKHOUSE_MAX_BUFFERED_ROWS", 200_000)
|
||||
CLICKHOUSE_MAX_RETRIES: Final = get_env_int("CLICKHOUSE_MAX_RETRIES", 3)
|
||||
AGENT_TRACING_RETENTION_DAYS: Final = get_env_int("AGENT_TRACING_RETENTION_DAYS", 30)
|
||||
AGENT_TRACING_SPEND_LOG_RETENTION_DAYS: Final = get_env_int("AGENT_TRACING_SPEND_LOG_RETENTION_DAYS", 90)
|
||||
DEFAULT_CLICKHOUSE_DATABASE: Final = "litellm"
|
||||
DEFAULT_AGENT_TRACING_RETENTION_DAYS: Final = 14
|
||||
OTLP_MAX_BODY_BYTES: Final = get_env_int("OTLP_MAX_BODY_BYTES", 16 * 1024 * 1024)
|
||||
OTLP_MAX_ATTRIBUTE_VALUE_BYTES: Final = get_env_int("OTLP_MAX_ATTRIBUTE_VALUE_BYTES", 64 * 1024)
|
||||
OTLP_RETRY_AFTER_SECONDS: Final = get_env_int("OTLP_RETRY_AFTER_SECONDS", 2)
|
||||
|
|
|
|||
|
|
@ -9,7 +9,6 @@ gzip JSONEachRow insert, either every `CLICKHOUSE_FLUSH_INTERVAL_SECONDS` or as
|
|||
"""
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
from collections.abc import Mapping, Sequence
|
||||
from contextlib import suppress
|
||||
from typing import Any, ClassVar, Final
|
||||
|
|
@ -23,13 +22,11 @@ from litellm.constants import (
|
|||
)
|
||||
from litellm.integrations.custom_batch_logger import CustomBatchLogger
|
||||
from litellm.rust_bridge.traces import ClickHouseStorage
|
||||
from litellm.tracing.config import trace_storage_config
|
||||
|
||||
|
||||
def clickhouse_storage_from_env() -> ClickHouseStorage:
|
||||
return ClickHouseStorage(
|
||||
database=os.getenv("CLICKHOUSE_DATABASE", "litellm"),
|
||||
url=os.getenv("CLICKHOUSE_URL", ""),
|
||||
)
|
||||
return ClickHouseStorage(trace_storage_config({}))
|
||||
|
||||
|
||||
class ClickHouseBatchLogger(CustomBatchLogger):
|
||||
|
|
|
|||
|
|
@ -7,5 +7,5 @@ AGENT_TRACES_BY_KEY_TABLE: Final = "agent_traces_by_key"
|
|||
SPEND_LOGS_TABLE: Final = "spend_logs"
|
||||
|
||||
|
||||
async def ensure_schema(storage: ClickHouseStorage, trace_retention_days: int, spend_log_retention_days: int) -> None:
|
||||
await storage.ensure_schema(trace_retention_days, spend_log_retention_days)
|
||||
async def ensure_schema(storage: ClickHouseStorage) -> None:
|
||||
await storage.ensure_schema()
|
||||
|
|
|
|||
|
|
@ -858,6 +858,7 @@ from litellm.secret_managers.main import (
|
|||
secret_manager_would_be_consulted,
|
||||
str_to_bool,
|
||||
)
|
||||
from litellm.tracing.config import is_clickhouse_tracing_enabled
|
||||
from litellm.types.integrations.slack_alerting import AlertType, SlackAlertingArgs
|
||||
from litellm.types.llms.anthropic import (
|
||||
AnthropicMessagesRequest,
|
||||
|
|
@ -1567,11 +1568,15 @@ async def proxy_startup_event(app: FastAPI) -> AsyncGenerator[ProxyLifespanState
|
|||
|
||||
register_scheduled_sync(scheduler)
|
||||
|
||||
tracing_settings: Final = general_settings.get("tracing")
|
||||
tracing_enabled: Final = TypeAdapter(bool).validate_python(
|
||||
isinstance(tracing_settings, dict) and tracing_settings.get("store") == "clickhouse"
|
||||
tracing_settings: Final = cast( # cast-ok: Pydantic validates the legacy untyped settings value
|
||||
dict[str, object] | None,
|
||||
TypeAdapter(dict[str, object] | None).validate_python(general_settings.get("tracing")),
|
||||
)
|
||||
async with manage_tracing(enabled=tracing_enabled) as receiver:
|
||||
tracing_enabled: Final = is_clickhouse_tracing_enabled(tracing_settings)
|
||||
async with manage_tracing(
|
||||
enabled=tracing_enabled,
|
||||
settings=tracing_settings,
|
||||
) as receiver:
|
||||
state: Final[ProxyLifespanState] = {"tracing_receiver": receiver}
|
||||
yield state
|
||||
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
from collections.abc import AsyncGenerator, Callable
|
||||
from collections.abc import AsyncGenerator, Callable, Mapping
|
||||
from contextlib import asynccontextmanager
|
||||
from typing import Final
|
||||
|
||||
|
|
@ -44,9 +44,12 @@ async def _start_receiver(factory: Callable[[], TraceReceiver]) -> TraceReceiver
|
|||
|
||||
@asynccontextmanager
|
||||
async def manage_tracing(
|
||||
enabled: bool, receiver_factory: Callable[[], TraceReceiver] = TraceReceiver.from_env
|
||||
enabled: bool,
|
||||
receiver_factory: Callable[[], TraceReceiver] | None = None,
|
||||
settings: Mapping[str, object] | None = None,
|
||||
) -> AsyncGenerator[TraceReceiver | None, None]:
|
||||
tracing: Final = await _start_receiver(receiver_factory) if enabled else None
|
||||
factory: Final = receiver_factory or (lambda: TraceReceiver.from_settings(settings or {}))
|
||||
tracing: Final = await _start_receiver(factory) if enabled else None
|
||||
if tracing is None:
|
||||
yield tracing
|
||||
return
|
||||
|
|
|
|||
|
|
@ -25,10 +25,19 @@ def trace_decode_otlp(body: bytes, content_type: str | None) -> list[DecodedSpan
|
|||
def trace_encode_error(message: str) -> bytes: ...
|
||||
def trace_normalized_field_definitions() -> list[dict[str, str]]: ...
|
||||
|
||||
@final
|
||||
class NativeTraceConfig:
|
||||
def __new__(
|
||||
cls,
|
||||
database: str,
|
||||
url: str,
|
||||
retention_days: int,
|
||||
) -> NativeTraceConfig: ...
|
||||
|
||||
@final
|
||||
class NativeTraceStorage:
|
||||
def __new__(cls, database: str, url: str, reader_url: str | None = None) -> NativeTraceStorage: ...
|
||||
def ensure_schema(self, trace_retention_days: int, spend_log_retention_days: int) -> Future[None]: ...
|
||||
def __new__(cls, config: NativeTraceConfig) -> NativeTraceStorage: ...
|
||||
def ensure_schema(self) -> Future[None]: ...
|
||||
def insert_rows(self, table: str, rows: Sequence[Mapping[str, object]]) -> Future[None]: ...
|
||||
def query_sql(self, sql: str, scope: QueryScope, secret: str) -> Future[str]: ...
|
||||
def query_help(self, scope: QueryScope, secret: str) -> Future[str]: ...
|
||||
|
|
@ -329,6 +338,7 @@ __all__ = [
|
|||
"ForkedAfterNativeRuntimeStarted",
|
||||
"HuggingFaceEncoding",
|
||||
"NativeDiagnosticProcessor",
|
||||
"NativeTraceConfig",
|
||||
"NativeTraceStorage",
|
||||
"ProcessReservedForForking",
|
||||
"ResponsesWebSocketConnection",
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
from collections.abc import Awaitable, Mapping, Sequence
|
||||
from dataclasses import dataclass
|
||||
from types import MappingProxyType
|
||||
from typing import Final, Literal, Protocol, TypedDict, cast
|
||||
|
||||
|
|
@ -77,9 +78,9 @@ QueryScope = AdminQueryScope | TeamQueryScope | KeyQueryScope
|
|||
|
||||
|
||||
class NativeStore(Protocol):
|
||||
def __init__(self, database: str, url: str, reader_url: str | None = None) -> None: ...
|
||||
def __init__(self, config: "NativeConfig") -> None: ...
|
||||
|
||||
def ensure_schema(self, trace_retention_days: int, spend_log_retention_days: int) -> Awaitable[None]: ...
|
||||
def ensure_schema(self) -> Awaitable[None]: ...
|
||||
|
||||
def insert_rows(self, table: str, rows: Sequence[Mapping[str, object]]) -> Awaitable[None]: ...
|
||||
|
||||
|
|
@ -93,6 +94,7 @@ class NativeStore(Protocol):
|
|||
|
||||
|
||||
class NativeTraces(Protocol):
|
||||
NativeTraceConfig: type["NativeConfig"]
|
||||
NativeTraceStorage: type[NativeStore]
|
||||
|
||||
def trace_decode_otlp(
|
||||
|
|
@ -115,6 +117,17 @@ QUERY_PARAMETERS: Final = TypeAdapter(dict[str, str | int | list[str]])
|
|||
_FIELD_DEFINITIONS_ADAPTER: Final = TypeAdapter(tuple[NormalizedFieldDefinition, ...])
|
||||
|
||||
|
||||
class NativeConfig(Protocol):
|
||||
def __init__(self, database: str, url: str, retention_days: int) -> None: ...
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True, repr=False)
|
||||
class TraceStorageConfig:
|
||||
url: str
|
||||
database: str = "litellm"
|
||||
retention_days: int = 14
|
||||
|
||||
|
||||
def _native() -> NativeTraces:
|
||||
native: Final = get_native_bridge()
|
||||
if native is None:
|
||||
|
|
@ -143,11 +156,17 @@ def encode_error(message: str) -> bytes:
|
|||
|
||||
|
||||
class ClickHouseStorage:
|
||||
def __init__(self, database: str, url: str, reader_url: str | None = None) -> None:
|
||||
self._native: Final = _native().NativeTraceStorage(database, url, reader_url)
|
||||
def __init__(self, config: TraceStorageConfig) -> None:
|
||||
native: Final = _native()
|
||||
validated: Final = native.NativeTraceConfig(
|
||||
config.database,
|
||||
config.url,
|
||||
config.retention_days,
|
||||
)
|
||||
self._native: Final = native.NativeTraceStorage(validated)
|
||||
|
||||
async def ensure_schema(self, trace_retention_days: int, spend_log_retention_days: int) -> None:
|
||||
await self._native.ensure_schema(trace_retention_days, spend_log_retention_days)
|
||||
async def ensure_schema(self) -> None:
|
||||
await self._native.ensure_schema()
|
||||
|
||||
async def insert_rows(self, table: str, rows: Sequence[Mapping[str, object]]) -> None:
|
||||
await self._native.insert_rows(table, rows)
|
||||
|
|
|
|||
84
litellm/tracing/config.py
Normal file
84
litellm/tracing/config.py
Normal file
|
|
@ -0,0 +1,84 @@
|
|||
import os
|
||||
from collections.abc import Mapping
|
||||
from typing import Final
|
||||
|
||||
from pydantic import TypeAdapter
|
||||
|
||||
from litellm.constants import DEFAULT_AGENT_TRACING_RETENTION_DAYS, DEFAULT_CLICKHOUSE_DATABASE
|
||||
from litellm.rust_bridge.traces import TraceStorageConfig
|
||||
|
||||
STORE_SETTINGS: Final = TypeAdapter(dict[str, object])
|
||||
|
||||
|
||||
def is_clickhouse_tracing_enabled(settings: object) -> bool:
|
||||
if not isinstance(settings, Mapping):
|
||||
return False
|
||||
typed_settings: Final = STORE_SETTINGS.validate_python(settings)
|
||||
store: Final = typed_settings.get("store")
|
||||
if not isinstance(store, Mapping):
|
||||
return False
|
||||
return STORE_SETTINGS.validate_python(store).get("type") == "clickhouse"
|
||||
|
||||
|
||||
def _value(settings: Mapping[str, object], field: str, environ: Mapping[str, str], default: object) -> object:
|
||||
if field not in settings:
|
||||
return default
|
||||
supplied: Final = settings[field]
|
||||
resolved: Final = (
|
||||
environ.get(supplied.removeprefix("os.environ/"))
|
||||
if isinstance(supplied, str) and supplied.startswith("os.environ/")
|
||||
else supplied
|
||||
)
|
||||
if resolved is None:
|
||||
raise ValueError(f"tracing.store.{field} is set but resolved to no value")
|
||||
return resolved
|
||||
|
||||
|
||||
def _retention_days(value: object) -> int:
|
||||
if isinstance(value, bool) or not isinstance(value, (int, str)):
|
||||
raise ValueError("tracing.store.retention_days must be a positive integer")
|
||||
try:
|
||||
days: Final = int(value)
|
||||
except ValueError as error:
|
||||
raise ValueError("tracing.store.retention_days must be a positive integer") from error
|
||||
if not 0 < days <= 2**32 - 1:
|
||||
raise ValueError("tracing.store.retention_days must be a positive integer")
|
||||
return days
|
||||
|
||||
|
||||
def _clickhouse_store(settings: Mapping[str, object]) -> Mapping[str, object]:
|
||||
raw_store: Final = settings.get("store")
|
||||
if raw_store is None:
|
||||
return {}
|
||||
if isinstance(raw_store, Mapping):
|
||||
store: Final = STORE_SETTINGS.validate_python(raw_store)
|
||||
if store.get("type") == "clickhouse":
|
||||
return store
|
||||
raise ValueError("tracing.store.type must be clickhouse")
|
||||
|
||||
|
||||
def trace_storage_config(settings: Mapping[str, object], environ: Mapping[str, str] = os.environ) -> TraceStorageConfig:
|
||||
store: Final = _clickhouse_store(settings)
|
||||
unknown: Final = store.keys() - {"type", "url", "database", "retention_days"}
|
||||
if unknown:
|
||||
raise ValueError(f"unsupported tracing.store settings: {', '.join(sorted(unknown))}")
|
||||
url: Final = _value(store, "url", environ, environ.get("CLICKHOUSE_URL"))
|
||||
database: Final = _value(
|
||||
store, "database", environ, environ.get("CLICKHOUSE_DATABASE", DEFAULT_CLICKHOUSE_DATABASE)
|
||||
)
|
||||
if not isinstance(url, str) or not url:
|
||||
raise ValueError("tracing.store.url or CLICKHOUSE_URL is required")
|
||||
if not isinstance(database, str):
|
||||
raise ValueError("tracing.store.database must be a string")
|
||||
return TraceStorageConfig(
|
||||
url=url,
|
||||
database=database,
|
||||
retention_days=_retention_days(
|
||||
_value(
|
||||
store,
|
||||
"retention_days",
|
||||
environ,
|
||||
environ.get("AGENT_TRACING_RETENTION_DAYS", DEFAULT_AGENT_TRACING_RETENTION_DAYS),
|
||||
)
|
||||
),
|
||||
)
|
||||
|
|
@ -13,21 +13,15 @@ The proxy endpoints are thin wrappers: auth -> build tenant/scope -> call one me
|
|||
"""
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
from collections.abc import AsyncIterable, Callable, Mapping
|
||||
from io import BytesIO
|
||||
from threading import BoundedSemaphore
|
||||
from types import MappingProxyType
|
||||
from typing import Final
|
||||
|
||||
from litellm.constants import (
|
||||
AGENT_TRACING_RETENTION_DAYS,
|
||||
AGENT_TRACING_SPEND_LOG_RETENTION_DAYS,
|
||||
OTLP_MAX_BODY_BYTES,
|
||||
OTLP_MAX_CONCURRENT_INGESTS,
|
||||
)
|
||||
from litellm.integrations.clickhouse.schema import ensure_schema
|
||||
from litellm.constants import OTLP_MAX_BODY_BYTES, OTLP_MAX_CONCURRENT_INGESTS
|
||||
from litellm.rust_bridge.traces import ClickHouseStorage
|
||||
from litellm.tracing.config import trace_storage_config
|
||||
from litellm.tracing.decode import OTLPPayloadTooLargeError, decode_otlp
|
||||
from litellm.tracing.store import TraceStore
|
||||
from litellm.tracing.types import (
|
||||
|
|
@ -103,22 +97,14 @@ class TraceReceiver:
|
|||
|
||||
@classmethod
|
||||
def from_env(cls) -> "TraceReceiver":
|
||||
return cls(
|
||||
store=TraceStore(
|
||||
ClickHouseStorage(
|
||||
database=os.getenv("CLICKHOUSE_DATABASE", "litellm"),
|
||||
url=os.environ["CLICKHOUSE_URL"],
|
||||
reader_url=os.getenv("CLICKHOUSE_READER_URL", os.environ["CLICKHOUSE_URL"]),
|
||||
)
|
||||
)
|
||||
)
|
||||
return cls.from_settings({})
|
||||
|
||||
@classmethod
|
||||
def from_settings(cls, settings: Mapping[str, object]) -> "TraceReceiver":
|
||||
return cls(store=TraceStore(ClickHouseStorage(trace_storage_config(settings))))
|
||||
|
||||
async def start(self) -> None:
|
||||
await ensure_schema(
|
||||
self.store.storage,
|
||||
trace_retention_days=AGENT_TRACING_RETENTION_DAYS,
|
||||
spend_log_retention_days=AGENT_TRACING_SPEND_LOG_RETENTION_DAYS,
|
||||
)
|
||||
await self.store.storage.ensure_schema()
|
||||
|
||||
async def ingest(
|
||||
self,
|
||||
|
|
|
|||
|
|
@ -13,11 +13,18 @@ VIRTUAL_ENV="$repo_root/.venv" uvx --from maturin==1.15.0 maturin develop \
|
|||
config_file="$(mktemp "${TMPDIR:-/tmp}/litellm-tracing-local.XXXXXX.yaml")"
|
||||
trap 'rm -f "$config_file"' EXIT
|
||||
cat > "$config_file" <<'EOF'
|
||||
model_list: []
|
||||
model_list:
|
||||
- model_name: claude-sonnet
|
||||
litellm_params:
|
||||
model: anthropic/claude-sonnet-5-5
|
||||
api_key: os.environ/ANTHROPIC_API_KEY
|
||||
general_settings:
|
||||
master_key: os.environ/LITELLM_MASTER_KEY
|
||||
tracing:
|
||||
store: clickhouse
|
||||
store:
|
||||
type: clickhouse
|
||||
url: os.environ/CLICKHOUSE_URL
|
||||
retention_days: 14
|
||||
EOF
|
||||
|
||||
export LITELLM_MASTER_KEY=sk-local-tracing
|
||||
|
|
@ -25,10 +32,16 @@ export LITELLM_SALT_KEY=sk-local-tracing-salt-key
|
|||
export DATABASE_URL=postgresql://litellm:litellm@127.0.0.1:15432/litellm
|
||||
export STORE_MODEL_IN_DB=True
|
||||
export CLICKHOUSE_URL=http://default:local-tracing@127.0.0.1:18123
|
||||
export CLICKHOUSE_READER_URL="$CLICKHOUSE_URL"
|
||||
export CLICKHOUSE_DATABASE=litellm
|
||||
export LITELLM_LOCAL_MODEL_COST_MAP=True
|
||||
|
||||
printf 'Proxy: http://127.0.0.1:4002/ui\nMaster key: %s\n' "$LITELLM_MASTER_KEY"
|
||||
(
|
||||
cd "$repo_root/ui/litellm-dashboard"
|
||||
"$repo_root/scripts/with_dashboard_node.sh" npm ci
|
||||
NEXT_PUBLIC_BASE_URL= "$repo_root/scripts/with_dashboard_node.sh" npm run build
|
||||
)
|
||||
export LITELLM_UI_PATH="$repo_root/ui/litellm-dashboard/out"
|
||||
|
||||
printf 'Dashboard: http://127.0.0.1:4002/ui/\nProxy: http://127.0.0.1:4002\nMaster key: %s\n' "$LITELLM_MASTER_KEY"
|
||||
"$repo_root/.venv/bin/python" litellm/proxy/proxy_cli.py \
|
||||
--config "$config_file" --host 127.0.0.1 --port 4002
|
||||
|
|
|
|||
|
|
@ -8,21 +8,31 @@ from urllib.parse import parse_qs, urlsplit
|
|||
|
||||
import pytest
|
||||
|
||||
from litellm.rust_bridge._native import NativeTraceStorage, trace_decode_otlp
|
||||
from litellm.rust_bridge.traces import ClickHouseStorage, NormalizedSpan, normalized_field_definitions
|
||||
from litellm.rust_bridge._native import NativeTraceConfig, NativeTraceStorage, trace_decode_otlp
|
||||
from litellm.rust_bridge.traces import (
|
||||
ClickHouseStorage,
|
||||
NormalizedSpan,
|
||||
TraceStorageConfig,
|
||||
normalized_field_definitions,
|
||||
)
|
||||
from litellm.tracing import Tenant, TraceReceiver, TracingPayloadTooLargeError
|
||||
from litellm.tracing.decode import decode_otlp
|
||||
from litellm.tracing.store import TraceStore
|
||||
from litellm.tracing.types import TraceScope
|
||||
from tests.test_litellm_rust.support.recording_server import RecordingServer, ResponseSpec
|
||||
|
||||
pytestmark = pytest.mark.requires_rust_extension
|
||||
|
||||
|
||||
def _native_storage(database: str, url: str, retention_days: int = 14) -> NativeTraceStorage:
|
||||
return NativeTraceStorage(NativeTraceConfig(database, url, retention_days))
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_trace_reader_projects_connection_and_parameters(recording_server: RecordingServer) -> None:
|
||||
recording_server.enqueue(ResponseSpec(body={"data": [{"trace_id": "trace-1"}]}))
|
||||
reader_url: Final = recording_server.base_url.replace("http://", "http://reader:p%40ss%2Fword%25@")
|
||||
storage: Final = NativeTraceStorage("trace_test", recording_server.base_url, reader_url + "?database=wrong")
|
||||
url: Final = recording_server.base_url.replace("http://", "http://reader:p%40ss%2Fword%25@")
|
||||
storage: Final = _native_storage("trace_test", url + "?database=wrong")
|
||||
rows: Final = json.loads(await storage.query("trace_spans", {"trace_id": "trace-1"}))
|
||||
request: Final = recording_server.requests[0]
|
||||
parameters: Final = parse_qs(urlsplit(request.path).query)
|
||||
|
|
@ -39,7 +49,7 @@ async def test_trace_reader_projects_connection_and_parameters(recording_server:
|
|||
@pytest.mark.asyncio
|
||||
async def test_trace_reader_rejects_success_status_with_embedded_error(recording_server: RecordingServer) -> None:
|
||||
recording_server.enqueue(ResponseSpec(body={"data": [], "exception": "query failed"}))
|
||||
storage: Final = NativeTraceStorage("trace_test", recording_server.base_url, recording_server.base_url)
|
||||
storage: Final = _native_storage("trace_test", recording_server.base_url)
|
||||
with pytest.raises(RuntimeError, match="invalid or failed JSON"):
|
||||
await storage.query("trace_spans", {})
|
||||
|
||||
|
|
@ -47,7 +57,7 @@ async def test_trace_reader_rejects_success_status_with_embedded_error(recording
|
|||
@pytest.mark.asyncio
|
||||
async def test_reader_rejects_arbitrary_sql_before_sending(recording_server: RecordingServer) -> None:
|
||||
recording_server.expected_requests = 0
|
||||
storage: Final = NativeTraceStorage("trace_test", recording_server.base_url, recording_server.base_url)
|
||||
storage: Final = _native_storage("trace_test", recording_server.base_url)
|
||||
with pytest.raises(ValueError, match="unknown ClickHouse read query"):
|
||||
await storage.query("SELECT 1", {})
|
||||
|
||||
|
|
@ -55,14 +65,44 @@ async def test_reader_rejects_arbitrary_sql_before_sending(recording_server: Rec
|
|||
@pytest.mark.asyncio
|
||||
async def test_schema_binding_rejects_invalid_database() -> None:
|
||||
with pytest.raises(ValueError, match=r"database.*retention"):
|
||||
NativeTraceStorage("db; DROP DATABASE default", "http://localhost:8123")
|
||||
NativeTraceConfig("db; DROP DATABASE default", "http://localhost:8123", 14)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_schema_binding_rejects_non_positive_retention() -> None:
|
||||
storage: Final = NativeTraceStorage("traces", "http://localhost:8123")
|
||||
with pytest.raises(ValueError, match=r"database.*retention"):
|
||||
await storage.ensure_schema(0, 14)
|
||||
NativeTraceConfig("traces", "http://localhost:8123", 0)
|
||||
|
||||
|
||||
def test_invalid_url_error_does_not_expose_credentials() -> None:
|
||||
with pytest.raises(RuntimeError, match="invalid ClickHouse HTTP URL") as error:
|
||||
NativeTraceConfig("traces", "secret://writer:password@example.com", 7)
|
||||
assert "password" not in str(error.value)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_from_env_reads_with_clickhouse_url(
|
||||
recording_server: RecordingServer, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
recording_server.enqueue(ResponseSpec(body={"data": []}))
|
||||
monkeypatch.setenv("CLICKHOUSE_URL", recording_server.base_url)
|
||||
monkeypatch.delenv("CLICKHOUSE_READER_URL", raising=False)
|
||||
scope: Final[TraceScope] = {"team_ids": (), "api_key_hash": ""}
|
||||
page: Final = await TraceReceiver.from_env().list_traces(scope, 0, 1)
|
||||
assert page == {"data": (), "next_cursor": None}
|
||||
assert len(recording_server.requests) == 1
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_schema_setup_uses_configured_retention(recording_server: RecordingServer) -> None:
|
||||
recording_server.expected_requests = 8
|
||||
storage: Final = _native_storage("trace_test", recording_server.base_url, 7)
|
||||
await storage.ensure_schema()
|
||||
ttl_statements: Final = tuple(
|
||||
request.raw_body for request in recording_server.requests if b"MODIFY TTL" in request.raw_body
|
||||
)
|
||||
assert len(ttl_statements) == 3
|
||||
assert all(b"INTERVAL 7 DAY" in statement for statement in ttl_statements)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
|
|
@ -73,9 +113,9 @@ async def test_schema_setup_uses_writer_credentials_and_rejects_failed_statement
|
|||
recording_server.enqueue(ResponseSpec(body=""))
|
||||
recording_server.enqueue(ResponseSpec(status=403, body="denied"))
|
||||
writer_url: Final = recording_server.base_url.replace("http://", "http://writer:p%40ss%2Fword%25@")
|
||||
storage: Final = NativeTraceStorage("trace_test", writer_url + "?database=wrong&readonly=1")
|
||||
storage: Final = _native_storage("trace_test", writer_url + "?database=wrong&readonly=1", 7)
|
||||
with pytest.raises(RuntimeError, match="schema setup failed with HTTP status 403"):
|
||||
await storage.ensure_schema(7, 14)
|
||||
await storage.ensure_schema()
|
||||
assert len(recording_server.requests) == 2
|
||||
assert recording_server.requests[0].raw_body.startswith(b"CREATE DATABASE IF NOT EXISTS")
|
||||
assert recording_server.requests[1].raw_body.startswith(b"CREATE TABLE IF NOT EXISTS")
|
||||
|
|
@ -89,7 +129,7 @@ async def test_schema_setup_uses_writer_credentials_and_rejects_failed_statement
|
|||
@pytest.mark.asyncio
|
||||
async def test_insert_encodes_and_sends_rows(recording_server: RecordingServer) -> None:
|
||||
recording_server.enqueue(ResponseSpec(body=""))
|
||||
storage: Final = NativeTraceStorage("trace_test", recording_server.base_url)
|
||||
storage: Final = _native_storage("trace_test", recording_server.base_url)
|
||||
before: Final = time.time_ns() // 1_000_000
|
||||
await storage.insert_rows("otel_traces", [{"Timestamp": 1_234_567_890, "Input": "hello", "EngineReceivedMs": -1}])
|
||||
after: Final = time.time_ns() // 1_000_000
|
||||
|
|
@ -169,7 +209,9 @@ def test_normalized_field_contract_matches_decoded_rust_span() -> None:
|
|||
@pytest.mark.asyncio
|
||||
async def test_resource_fanout_reaches_insert_with_identical_values(recording_server: RecordingServer) -> None:
|
||||
body: Final = _resource_export(16 * 1024, 1024)
|
||||
receiver: Final = TraceReceiver(TraceStore(ClickHouseStorage("trace_test", recording_server.base_url)))
|
||||
receiver: Final = TraceReceiver(
|
||||
TraceStore(ClickHouseStorage(TraceStorageConfig(recording_server.base_url, "trace_test")))
|
||||
)
|
||||
tenant: Final = Tenant("team-a", "key-a", "org-a")
|
||||
assert await receiver.ingest(body, "application/json", None, tenant) == 1024
|
||||
encoded: Final = gzip.decompress(recording_server.requests[0].raw_body)
|
||||
|
|
@ -186,7 +228,9 @@ async def test_resource_fanout_reaches_insert_with_identical_values(recording_se
|
|||
async def test_shared_resource_still_hits_insert_limit_before_transport(recording_server: RecordingServer) -> None:
|
||||
recording_server.expected_requests = 0
|
||||
body: Final = _resource_export(64 * 1024, 1024)
|
||||
receiver: Final = TraceReceiver(TraceStore(ClickHouseStorage("trace_test", recording_server.base_url)))
|
||||
receiver: Final = TraceReceiver(
|
||||
TraceStore(ClickHouseStorage(TraceStorageConfig(recording_server.base_url, "trace_test")))
|
||||
)
|
||||
with pytest.raises(TracingPayloadTooLargeError, match="encoded size limit"):
|
||||
await receiver.ingest(body, "application/json", None, Tenant("team-a", "key-a"))
|
||||
assert recording_server.requests == []
|
||||
|
|
@ -194,7 +238,7 @@ async def test_shared_resource_still_hits_insert_limit_before_transport(recordin
|
|||
|
||||
@pytest.mark.asyncio
|
||||
async def test_insert_validates_values_without_pydantic_copy(recording_server: RecordingServer) -> None:
|
||||
storage: Final = ClickHouseStorage("trace_test", recording_server.base_url)
|
||||
storage: Final = ClickHouseStorage(TraceStorageConfig(recording_server.base_url, "trace_test"))
|
||||
invalid: Final = object()
|
||||
with pytest.raises(ValueError, match=type(invalid).__name__):
|
||||
await storage.insert_rows("otel_traces", [{"ResourceAttributes": invalid}])
|
||||
|
|
@ -225,7 +269,7 @@ def test_trace_sql_endpoint_executes_for_admin_and_preserves_clickhouse_envelope
|
|||
for _ in range(11):
|
||||
recording_server.enqueue(ResponseSpec(body=""))
|
||||
recording_server.enqueue(ResponseSpec(body=envelope))
|
||||
storage: Final = ClickHouseStorage("trace_test", recording_server.base_url, recording_server.base_url)
|
||||
storage: Final = ClickHouseStorage(TraceStorageConfig(recording_server.base_url, "trace_test"))
|
||||
app: Final = FastAPI()
|
||||
app.include_router(router)
|
||||
app.dependency_overrides[provide_trace_query_secret] = lambda: "test-master-secret"
|
||||
|
|
@ -260,7 +304,7 @@ def test_trace_help_endpoint_runs_native_schema_and_metadata_discovery(recording
|
|||
{"data": [{"key": "custom.resource"}]},
|
||||
):
|
||||
recording_server.enqueue(ResponseSpec(body=response))
|
||||
storage: Final = ClickHouseStorage("trace_test", recording_server.base_url, recording_server.base_url)
|
||||
storage: Final = ClickHouseStorage(TraceStorageConfig(recording_server.base_url, "trace_test"))
|
||||
app: Final = FastAPI()
|
||||
app.include_router(router)
|
||||
app.dependency_overrides[provide_trace_query_secret] = lambda: "test-master-secret"
|
||||
|
|
@ -299,7 +343,7 @@ def test_trace_sql_endpoint_distinguishes_query_errors_from_reader_failures(
|
|||
recording_server.enqueue(ResponseSpec(status=clickhouse_status, body=b"ClickHouse rejected the query"))
|
||||
envelope: Final = {"meta": [{"name": "answer", "type": "UInt8"}], "data": [{"answer": 42}], "rows": 1}
|
||||
recording_server.enqueue(ResponseSpec(body=envelope))
|
||||
storage: Final = ClickHouseStorage("trace_test", recording_server.base_url, recording_server.base_url)
|
||||
storage: Final = ClickHouseStorage(TraceStorageConfig(recording_server.base_url, "trace_test"))
|
||||
app: Final = FastAPI()
|
||||
app.include_router(router)
|
||||
app.dependency_overrides[provide_trace_query_secret] = lambda: "test-master-secret"
|
||||
|
|
|
|||
|
|
@ -40,10 +40,34 @@ from litellm.proxy.proxy_server import (
|
|||
validate_deployment_complexity_router_placement,
|
||||
validate_deployment_max_agentic_loops,
|
||||
)
|
||||
from litellm.tracing.config import trace_storage_config
|
||||
|
||||
from .conftest import normalize
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_proxy_config_loads_tracing_url_and_retention_from_yaml(tmp_path, monkeypatch) -> None:
|
||||
config_file: Final = tmp_path / "tracing.yaml"
|
||||
config_file.write_text(
|
||||
"model_list: []\ngeneral_settings:\n tracing:\n store:\n"
|
||||
" type: clickhouse\n url: os.environ/TRACING_TEST_URL\n"
|
||||
" database: analytics\n retention_days: 7\n"
|
||||
)
|
||||
monkeypatch.setenv("TRACING_TEST_URL", "http://localhost:8123")
|
||||
monkeypatch.setenv("CLICKHOUSE_URL", "http://unused:8123")
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.prisma_client", None)
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.store_model_in_db", False)
|
||||
monkeypatch.delenv("LITELLM_CONFIG_BUCKET_NAME", raising=False)
|
||||
|
||||
_, _, settings = await ProxyConfig().load_config(router=None, config_file_path=str(config_file))
|
||||
tracing = trace_storage_config(settings["tracing"])
|
||||
assert (tracing.url, tracing.database, tracing.retention_days) == (
|
||||
"http://localhost:8123",
|
||||
"analytics",
|
||||
7,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize("shutdown_error", [False, True])
|
||||
async def test_tracing_config_automatically_logs_spend_without_callback_setting(shutdown_error: bool) -> None:
|
||||
|
|
|
|||
0
tests/unit/tracing/__init__.py
Normal file
0
tests/unit/tracing/__init__.py
Normal file
116
tests/unit/tracing/test_config.py
Normal file
116
tests/unit/tracing/test_config.py
Normal file
|
|
@ -0,0 +1,116 @@
|
|||
import pytest
|
||||
|
||||
from litellm import constants
|
||||
from litellm.tracing.config import is_clickhouse_tracing_enabled, trace_storage_config
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("settings", "enabled"),
|
||||
[
|
||||
({"store": "clickhouse"}, False),
|
||||
({"store": {"type": "clickhouse"}}, True),
|
||||
({"store": {"type": "other"}}, False),
|
||||
(None, False),
|
||||
],
|
||||
)
|
||||
def test_clickhouse_tracing_enablement(settings: object, enabled: bool) -> None:
|
||||
assert is_clickhouse_tracing_enabled(settings) is enabled
|
||||
|
||||
|
||||
def test_yaml_values_override_defaults_and_resolve_nested_references() -> None:
|
||||
config = trace_storage_config(
|
||||
{
|
||||
"store": {
|
||||
"type": "clickhouse",
|
||||
"url": "os.environ/TRACING_URL",
|
||||
"database": "os.environ/TRACING_DATABASE",
|
||||
"retention_days": "os.environ/TRACING_RETENTION_DAYS",
|
||||
},
|
||||
},
|
||||
{
|
||||
"TRACING_URL": "https://writer:password@clickhouse.example:8443",
|
||||
"TRACING_DATABASE": "analytics",
|
||||
"TRACING_RETENTION_DAYS": "7",
|
||||
"CLICKHOUSE_URL": "https://other.example:8443",
|
||||
},
|
||||
)
|
||||
assert config.url == "https://writer:password@clickhouse.example:8443"
|
||||
assert config.database == "analytics"
|
||||
assert config.retention_days == 7
|
||||
assert "password" not in repr(config)
|
||||
|
||||
|
||||
def test_omitted_fields_use_environment() -> None:
|
||||
config = trace_storage_config(
|
||||
{},
|
||||
{
|
||||
"CLICKHOUSE_URL": "http://localhost:8123",
|
||||
"CLICKHOUSE_DATABASE": "env_database",
|
||||
"AGENT_TRACING_RETENTION_DAYS": "11",
|
||||
},
|
||||
)
|
||||
assert (config.url, config.database, config.retention_days) == ("http://localhost:8123", "env_database", 11)
|
||||
|
||||
|
||||
def test_environment_is_read_when_config_is_resolved(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("CLICKHOUSE_URL", "http://localhost:8123")
|
||||
monkeypatch.setenv("CLICKHOUSE_DATABASE", "late_database")
|
||||
monkeypatch.setenv("AGENT_TRACING_RETENTION_DAYS", "9")
|
||||
config = trace_storage_config({})
|
||||
assert (config.database, config.retention_days) == ("late_database", 9)
|
||||
|
||||
|
||||
def test_omitted_fields_without_environment_use_constant_defaults() -> None:
|
||||
config = trace_storage_config({}, {"CLICKHOUSE_URL": "http://localhost:8123"})
|
||||
assert (config.database, config.retention_days) == (
|
||||
constants.DEFAULT_CLICKHOUSE_DATABASE,
|
||||
constants.DEFAULT_AGENT_TRACING_RETENTION_DAYS,
|
||||
)
|
||||
assert (config.database, config.retention_days) == ("litellm", 14)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("field", ["url", "database", "retention_days"])
|
||||
def test_unset_environment_reference_does_not_fall_back(field: str) -> None:
|
||||
store: dict[str, object] = {"type": "clickhouse", "url": "http://localhost:8123", field: "os.environ/MISSING"}
|
||||
with pytest.raises(ValueError, match=rf"tracing.store.{field} is set but resolved to no value") as error:
|
||||
trace_storage_config({"store": store}, {"CLICKHOUSE_URL": "http://fallback:8123"})
|
||||
assert "MISSING" not in str(error.value)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("store", ["clickhouse", {"type": "other"}])
|
||||
def test_non_clickhouse_store_is_rejected(store: object) -> None:
|
||||
with pytest.raises(ValueError, match=r"tracing\.store\.type must be clickhouse"):
|
||||
trace_storage_config({"store": store}, {"CLICKHOUSE_URL": "http://localhost:8123"})
|
||||
|
||||
|
||||
def test_non_string_database_is_rejected() -> None:
|
||||
with pytest.raises(ValueError, match=r"tracing\.store\.database must be a string"):
|
||||
trace_storage_config({"store": {"type": "clickhouse", "url": "http://localhost:8123", "database": 1}}, {})
|
||||
|
||||
|
||||
@pytest.mark.parametrize("value", [0, -1, True, "not-a-number", 2**32])
|
||||
def test_invalid_retention_is_rejected(value: object) -> None:
|
||||
with pytest.raises(ValueError, match=r"tracing.store.retention_days must be a positive integer"):
|
||||
trace_storage_config(
|
||||
{"store": {"type": "clickhouse", "url": "http://localhost:8123", "retention_days": value}}, {}
|
||||
)
|
||||
|
||||
|
||||
def test_missing_url_is_rejected() -> None:
|
||||
with pytest.raises(ValueError, match=r"tracing.store.url or CLICKHOUSE_URL is required"):
|
||||
trace_storage_config({"store": {"type": "clickhouse"}}, {})
|
||||
|
||||
|
||||
def test_legacy_reader_and_split_retention_fields_are_rejected() -> None:
|
||||
with pytest.raises(ValueError, match="reader_url, trace_retention_days"):
|
||||
trace_storage_config(
|
||||
{
|
||||
"store": {
|
||||
"type": "clickhouse",
|
||||
"url": "http://localhost:8123",
|
||||
"reader_url": "http://localhost:8124",
|
||||
"trace_retention_days": 30,
|
||||
}
|
||||
},
|
||||
{},
|
||||
)
|
||||
|
|
@ -53,7 +53,8 @@ export default function AuditLogsPanel({
|
|||
return typeof entry?.value === "string" && entry.value.trim() ? entry.value.trim() : undefined;
|
||||
};
|
||||
|
||||
const canQueryAuditLogs = !!accessToken && !!token && !!userRole && !!userID && isActive && premiumUser;
|
||||
const hasSession = [accessToken, token, userRole, userID].every(Boolean);
|
||||
const canQueryAuditLogs = hasSession && isActive && premiumUser;
|
||||
|
||||
const query = useQuery<AuditLogsResponse>({
|
||||
queryKey: ["audit_logs", pagination.pageIndex, pagination.pageSize, columnFilters, searchTerm],
|
||||
|
|
|
|||
|
|
@ -116,7 +116,8 @@ describe("AgentTracesSection", () => {
|
|||
|
||||
const card = await screen.findByTestId("tracing-setup-card");
|
||||
expect(card).toHaveTextContent("Tracing is not enabled");
|
||||
expect(card).toHaveTextContent("store: clickhouse");
|
||||
expect(card).toHaveTextContent("type: clickhouse");
|
||||
expect(card).toHaveTextContent("url: os.environ/CLICKHOUSE_URL");
|
||||
expect(screen.getByRole("button", { name: "Check setup" })).toBeEnabled();
|
||||
expect(card).not.toHaveTextContent(/langsmith/i);
|
||||
expect(card).toHaveTextContent("ClickHouse and proxy setup");
|
||||
|
|
@ -225,7 +226,7 @@ describe("AgentTracesSection", () => {
|
|||
|
||||
const card = await screen.findByTestId("tracing-setup-card");
|
||||
expect(card).toHaveTextContent("Tracing is not enabled");
|
||||
expect(card).toHaveTextContent("CLICKHOUSE_READER_URL");
|
||||
expect(card).toHaveTextContent("url: os.environ/CLICKHOUSE_URL");
|
||||
});
|
||||
|
||||
it("lists every run with its input, counts and failed column", async () => {
|
||||
|
|
|
|||
|
|
@ -182,7 +182,8 @@ describe("TracingSetupCard", () => {
|
|||
const onCheck = vi.fn();
|
||||
const { card } = renderCard({ detail: "Agent tracing is not enabled", onCheck });
|
||||
expect(screen.getByRole("heading", { name: "Enable tracing" })).toBeVisible();
|
||||
expect(card).toHaveTextContent("store: clickhouse");
|
||||
expect(card).toHaveTextContent("type: clickhouse");
|
||||
expect(card).toHaveTextContent("url: os.environ/CLICKHOUSE_URL");
|
||||
expect(screen.queryByRole("combobox", { name: "Your agent framework" })).not.toBeInTheDocument();
|
||||
expect(screen.queryByRole("button", { name: "Send a test trace" })).not.toBeInTheDocument();
|
||||
await user.click(screen.getByRole("button", { name: "Check setup" }));
|
||||
|
|
|
|||
|
|
@ -231,9 +231,10 @@ export const otlpEndpoints = (proxyUrl: string): readonly (readonly [string, str
|
|||
export const PROXY_CONFIG_SNIPPET = [
|
||||
"general_settings:",
|
||||
" tracing:",
|
||||
" store: clickhouse",
|
||||
"",
|
||||
"# env: CLICKHOUSE_URL (writer) and CLICKHOUSE_READER_URL (read-only user)",
|
||||
" store:",
|
||||
" type: clickhouse",
|
||||
" url: os.environ/CLICKHOUSE_URL",
|
||||
" retention_days: 14",
|
||||
].join("\n");
|
||||
|
||||
function CodeBlock({
|
||||
|
|
@ -575,8 +576,8 @@ function EnableTracing({ checked, checking, onCheck }: { checked: boolean; check
|
|||
<>
|
||||
<Step title="Enable tracing on the proxy">
|
||||
<p className="mb-3 text-sm leading-6 text-muted-foreground">
|
||||
Set your ClickHouse writer and read-only reader URLs, add this to config.yaml, then restart the proxy. Ask
|
||||
your proxy administrator if you don’t manage this deployment.
|
||||
Set your ClickHouse URL, add this to config.yaml, then restart the proxy. Ask your proxy administrator if you
|
||||
don’t manage this deployment.
|
||||
</p>
|
||||
<CodeBlock code={PROXY_CONFIG_SNIPPET} tabs={<FileLabel>config.yaml</FileLabel>} />
|
||||
<a
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue