fix: align sqlite repositories with postgres behavior

This commit is contained in:
HsungKayphoon
2026-05-16 23:43:40 +08:00
parent f9f1fa928a
commit c5e26a1ed6
4 changed files with 4904 additions and 549 deletions

View File

@@ -45,7 +45,15 @@ SELECT
m.provider_model_mappings AS model_provider_model_mappings,
m.supports_streaming AS model_supports_streaming,
m.is_active AS model_is_active,
m.is_available AS model_is_available
m.is_available AS model_is_available,
CASE
WHEN json_valid(p.config) THEN
CASE
WHEN json_type(p.config, '$.pool_advanced') IS NOT NULL THEN 1
ELSE 0
END
ELSE 0
END AS provider_pool_enabled
FROM providers p
INNER JOIN provider_endpoints pe ON pe.provider_id = p.id
INNER JOIN provider_api_keys pak ON pak.provider_id = p.id
@@ -67,52 +75,139 @@ pub struct SqliteMinimalCandidateSelectionReadRepository {
#[derive(Debug, Clone)]
struct CandidateSelectionRow {
row: StoredMinimalCandidateSelectionRow,
provider_pool_enabled: bool,
key_auth_config: Option<String>,
key_last_used_at_unix_secs: Option<u64>,
}
#[derive(Debug, Clone, Copy)]
enum SelectedRowsOrder {
WithGlobalModel,
WithoutGlobalModel,
}
#[derive(Debug, Clone, Copy)]
enum SelectedRowsFilter<'a> {
None,
GlobalModel(&'a str),
RequestedModel(&'a str),
}
#[derive(Debug, Clone, Copy)]
struct SqlPage {
limit: i64,
offset: i64,
}
impl SqliteMinimalCandidateSelectionReadRepository {
pub fn new(pool: SqlitePool) -> Self {
Self { pool }
}
async fn load_rows_for_api_format(
&self,
api_format: &str,
) -> Result<Vec<CandidateSelectionRow>, DataLayerError> {
let canonical_api_format = normalize_api_format(api_format);
let storage_aliases = api_format_aliases(&canonical_api_format);
let match_aliases = sql_match_aliases(&storage_aliases);
let mut builder = QueryBuilder::<Sqlite>::new(CANDIDATE_SELECTION_COLUMNS);
builder.push(" AND LOWER(pe.api_format) IN (");
{
let mut separated = builder.separated(", ");
for alias in &match_aliases {
separated.push_bind(alias);
}
}
builder.push(")");
let rows = builder.build().fetch_all(&self.pool).await.map_sql_err()?;
let mut items = rows
.iter()
.map(map_candidate_selection_row)
.collect::<Result<Vec<_>, _>>()?;
items.retain(|item| {
api_format_matches(&item.row.endpoint_api_format, &canonical_api_format)
&& item.row.key_supports_api_format(&canonical_api_format)
&& key_auth_channel_matches(item, &canonical_api_format)
});
Ok(items)
}
async fn selected_rows_for_api_format(
&self,
api_format: &str,
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
let rows = self.load_rows_for_api_format(api_format).await?;
Ok(sort_rows(select_pool_rows(rows), true))
self.load_selected_rows_for_api_format(
api_format,
SelectedRowsFilter::None,
SelectedRowsOrder::WithGlobalModel,
None,
)
.await
}
async fn load_selected_rows_for_api_format(
&self,
api_format: &str,
filter: SelectedRowsFilter<'_>,
order: SelectedRowsOrder,
page: Option<SqlPage>,
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
let canonical_api_format = normalize_api_format(api_format);
let storage_aliases = api_format_aliases(&canonical_api_format);
let match_aliases = sql_match_aliases(&storage_aliases);
let mut rows = Vec::new();
for storage_api_format in storage_aliases {
let mut builder = QueryBuilder::<Sqlite>::new("WITH candidate_rows AS (");
builder.push(CANDIDATE_SELECTION_COLUMNS);
push_candidate_sql_filters(&mut builder, &storage_api_format, &match_aliases);
match filter {
SelectedRowsFilter::None => {}
SelectedRowsFilter::GlobalModel(global_model_name) => {
builder.push(" AND gm.name = ");
builder.push_bind(global_model_name);
}
SelectedRowsFilter::RequestedModel(requested_model_name) => {
push_requested_model_sql_filter(
&mut builder,
requested_model_name,
&match_aliases,
);
}
}
builder.push(
r#"
),
pool_rows AS (
SELECT candidate.*
FROM candidate_rows candidate
WHERE candidate.provider_pool_enabled = 1
AND NOT EXISTS (
SELECT 1
FROM candidate_rows other
WHERE other.provider_pool_enabled = 1
AND other.provider_id = candidate.provider_id
AND other.endpoint_id = candidate.endpoint_id
AND other.model_id = candidate.model_id
AND (
other.key_internal_priority < candidate.key_internal_priority
OR (
other.key_internal_priority = candidate.key_internal_priority
AND other.key_id < candidate.key_id
)
)
)
),
selected_rows AS (
SELECT * FROM candidate_rows WHERE provider_pool_enabled = 0
UNION ALL
SELECT * FROM pool_rows
)
SELECT * FROM selected_rows
"#,
);
push_selected_rows_order(&mut builder, order);
if let Some(page) = page {
builder.push(" LIMIT ");
builder.push_bind(page.limit);
builder.push(" OFFSET ");
builder.push_bind(page.offset);
}
let query_rows = builder.build().fetch_all(&self.pool).await.map_sql_err()?;
let mut items = query_rows
.iter()
.map(map_candidate_selection_row)
.collect::<Result<Vec<_>, _>>()?;
items.retain(|item| {
api_format_matches(&item.row.endpoint_api_format, &canonical_api_format)
&& item.row.key_supports_api_format(&canonical_api_format)
&& key_auth_channel_matches(item, &canonical_api_format)
});
rows.extend(items.into_iter().map(|item| item.row));
}
let rows = match filter {
SelectedRowsFilter::RequestedModel(requested_model_name) => rows
.into_iter()
.filter(|row| {
row_matches_requested_model(row, requested_model_name, &canonical_api_format)
})
.collect(),
_ => rows,
};
Ok(dedupe_candidate_selection_rows(rows))
}
}
@@ -130,14 +225,13 @@ impl MinimalCandidateSelectionReadRepository for SqliteMinimalCandidateSelection
api_format: &str,
global_model_name: &str,
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
Ok(sort_rows(
self.selected_rows_for_api_format(api_format)
.await?
.into_iter()
.filter(|row| row.global_model_name == global_model_name)
.collect(),
false,
))
self.load_selected_rows_for_api_format(
api_format,
SelectedRowsFilter::GlobalModel(global_model_name),
SelectedRowsOrder::WithoutGlobalModel,
None,
)
.await
}
async fn list_for_exact_api_format_and_requested_model(
@@ -145,13 +239,11 @@ impl MinimalCandidateSelectionReadRepository for SqliteMinimalCandidateSelection
api_format: &str,
requested_model_name: &str,
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
self.list_for_exact_api_format_and_requested_model_page(
&StoredRequestedModelCandidateRowsQuery {
api_format: api_format.to_string(),
requested_model_name: requested_model_name.to_string(),
offset: 0,
limit: u32::MAX,
},
self.load_selected_rows_for_api_format(
api_format,
SelectedRowsFilter::RequestedModel(requested_model_name),
SelectedRowsOrder::WithGlobalModel,
None,
)
.await
}
@@ -160,42 +252,74 @@ impl MinimalCandidateSelectionReadRepository for SqliteMinimalCandidateSelection
&self,
query: &StoredRequestedModelCandidateRowsQuery,
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
let rows = self
.selected_rows_for_api_format(&query.api_format)
.await?
.into_iter()
.filter(|row| {
row_matches_requested_model(row, &query.requested_model_name, &query.api_format)
})
.collect::<Vec<_>>();
Ok(sort_rows(rows, true)
.into_iter()
.skip(query.offset as usize)
.take(query.limit as usize)
.collect())
self.load_selected_rows_for_api_format(
&query.api_format,
SelectedRowsFilter::RequestedModel(&query.requested_model_name),
SelectedRowsOrder::WithGlobalModel,
Some(SqlPage {
limit: i64::from(query.limit.max(1)),
offset: i64::from(query.offset),
}),
)
.await
}
async fn list_pool_key_rows_for_group(
&self,
query: &StoredPoolKeyCandidateRowsQuery,
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
let rows = self
.load_rows_for_api_format(&query.api_format)
.await?
.into_iter()
.filter(|row| {
row.row.provider_id == query.provider_id
&& row.row.endpoint_id == query.endpoint_id
&& row.row.model_id == query.model_id
})
.collect::<Vec<_>>();
let mut rows = sort_pool_key_rows(rows, &query.order);
Ok(rows
.drain(..)
.skip(query.offset as usize)
.take(query.limit as usize)
.map(|item| item.row)
.collect())
let canonical_api_format = normalize_api_format(&query.api_format);
let storage_aliases = api_format_aliases(&canonical_api_format);
let match_aliases = sql_match_aliases(&storage_aliases);
let mut rows = Vec::<CandidateSelectionRow>::new();
let page_in_sql = !matches!(query.order, StoredPoolKeyCandidateOrder::LoadBalance { .. });
for storage_api_format in storage_aliases {
let mut builder = QueryBuilder::<Sqlite>::new(CANDIDATE_SELECTION_COLUMNS);
push_candidate_sql_filters(&mut builder, &storage_api_format, &match_aliases);
builder.push(" AND p.id = ");
builder.push_bind(&query.provider_id);
builder.push(" AND pe.id = ");
builder.push_bind(&query.endpoint_id);
builder.push(" AND m.id = ");
builder.push_bind(&query.model_id);
if page_in_sql {
push_pool_key_order(&mut builder, &query.order);
builder.push(" LIMIT ");
builder.push_bind(i64::from(query.limit.max(1)));
builder.push(" OFFSET ");
builder.push_bind(i64::from(query.offset));
} else {
builder.push(" ORDER BY pak.id ASC");
}
let query_rows = builder.build().fetch_all(&self.pool).await.map_sql_err()?;
let mut items = query_rows
.iter()
.map(map_candidate_selection_row)
.collect::<Result<Vec<_>, _>>()?;
items.retain(|item| {
api_format_matches(&item.row.endpoint_api_format, &canonical_api_format)
&& item.row.key_supports_api_format(&canonical_api_format)
&& key_auth_channel_matches(item, &canonical_api_format)
});
rows.extend(items);
}
if page_in_sql {
Ok(dedupe_candidate_selection_rows(
rows.into_iter().map(|item| item.row).collect(),
))
} else {
Ok(dedupe_candidate_selection_rows(
sort_pool_key_rows(rows, &query.order)
.into_iter()
.skip(query.offset as usize)
.take(query.limit as usize)
.map(|item| item.row)
.collect(),
))
}
}
async fn list_pool_key_rows_for_group_key_ids(
@@ -211,75 +335,330 @@ impl MinimalCandidateSelectionReadRepository for SqliteMinimalCandidateSelection
.enumerate()
.map(|(index, key_id)| (key_id.as_str(), index))
.collect::<BTreeMap<_, _>>();
let mut rows = self
.load_rows_for_api_format(&query.api_format)
.await?
.into_iter()
.filter(|row| {
row.row.provider_id == query.provider_id
&& row.row.endpoint_id == query.endpoint_id
&& row.row.model_id == query.model_id
&& key_order.contains_key(row.row.key_id.as_str())
})
.map(|item| item.row)
.collect::<Vec<_>>();
let canonical_api_format = normalize_api_format(&query.api_format);
let storage_aliases = api_format_aliases(&canonical_api_format);
let match_aliases = sql_match_aliases(&storage_aliases);
let mut rows = Vec::new();
for storage_api_format in storage_aliases {
let mut builder = QueryBuilder::<Sqlite>::new(CANDIDATE_SELECTION_COLUMNS);
push_candidate_sql_filters(&mut builder, &storage_api_format, &match_aliases);
builder.push(" AND p.id = ");
builder.push_bind(&query.provider_id);
builder.push(" AND pe.id = ");
builder.push_bind(&query.endpoint_id);
builder.push(" AND m.id = ");
builder.push_bind(&query.model_id);
builder.push(" AND pak.id IN (");
{
let mut separated = builder.separated(", ");
for key_id in &query.key_ids {
separated.push_bind(key_id);
}
}
builder.push(")");
builder.push(" ORDER BY CASE pak.id");
for (index, key_id) in query.key_ids.iter().enumerate() {
builder.push(" WHEN ");
builder.push_bind(key_id);
builder.push(" THEN ");
builder.push_bind(i64::try_from(index).map_err(|_| {
DataLayerError::UnexpectedValue("key id order index overflowed".to_string())
})?);
}
builder.push(" ELSE ");
builder.push_bind(i64::try_from(query.key_ids.len()).map_err(|_| {
DataLayerError::UnexpectedValue("key id order length overflowed".to_string())
})?);
builder.push(" END ASC, pak.id ASC");
let query_rows = builder.build().fetch_all(&self.pool).await.map_sql_err()?;
let mut items = query_rows
.iter()
.map(map_candidate_selection_row)
.collect::<Result<Vec<_>, _>>()?;
items.retain(|item| {
api_format_matches(&item.row.endpoint_api_format, &canonical_api_format)
&& item.row.key_supports_api_format(&canonical_api_format)
&& key_auth_channel_matches(item, &canonical_api_format)
});
rows.extend(items.into_iter().map(|item| item.row));
}
let mut rows = dedupe_candidate_selection_rows(rows);
rows.sort_by(|left, right| {
key_order
.get(left.key_id.as_str())
.cmp(&key_order.get(right.key_id.as_str()))
.then(left.key_id.cmp(&right.key_id))
});
Ok(dedupe_candidate_selection_rows(rows))
Ok(rows)
}
}
fn select_pool_rows(rows: Vec<CandidateSelectionRow>) -> Vec<StoredMinimalCandidateSelectionRow> {
let mut selected = Vec::new();
let mut pool_rows =
BTreeMap::<(String, String, String), StoredMinimalCandidateSelectionRow>::new();
for item in rows {
if !item.provider_pool_enabled {
selected.push(item.row);
continue;
}
let key = (
item.row.provider_id.clone(),
item.row.endpoint_id.clone(),
item.row.model_id.clone(),
);
match pool_rows.get(&key) {
Some(existing)
if (existing.key_internal_priority, existing.key_id.as_str())
<= (item.row.key_internal_priority, item.row.key_id.as_str()) => {}
_ => {
pool_rows.insert(key, item.row);
}
}
}
selected.extend(pool_rows.into_values());
dedupe_candidate_selection_rows(selected)
fn push_candidate_sql_filters(
builder: &mut QueryBuilder<'_, Sqlite>,
storage_api_format: &str,
match_aliases: &[String],
) {
builder.push(" AND LOWER(COALESCE(pe.api_format, '')) = ");
builder.push_bind(storage_api_format.trim().to_ascii_lowercase());
push_key_api_format_sql_filter(builder, match_aliases);
push_key_auth_channel_sql_filter(builder, storage_api_format);
}
fn sort_rows(
mut rows: Vec<StoredMinimalCandidateSelectionRow>,
include_global_model: bool,
) -> Vec<StoredMinimalCandidateSelectionRow> {
rows.sort_by(|left, right| {
if include_global_model {
let ordering = left.global_model_name.cmp(&right.global_model_name);
if !ordering.is_eq() {
return ordering;
}
fn push_key_api_format_sql_filter(
builder: &mut QueryBuilder<'_, Sqlite>,
match_aliases: &[String],
) {
builder.push(
r#"
AND (
pak.api_formats IS NULL
OR TRIM(pak.api_formats) = ''
OR CASE
WHEN json_valid(pak.api_formats) THEN
(
(
json_type(pak.api_formats) = 'array'
AND EXISTS (
SELECT 1
FROM json_each(pak.api_formats) AS fmt
WHERE LOWER(TRIM(CAST(fmt.value AS TEXT))) IN (
"#,
);
push_bind_list(builder, match_aliases);
builder.push(
r#"
)
)
)
OR (
json_type(pak.api_formats) = 'text'
AND LOWER(TRIM(CAST(json_extract(pak.api_formats, '$') AS TEXT))) IN (
"#,
);
push_bind_list(builder, match_aliases);
builder.push(
r#"
)
)
OR (
json_type(pak.api_formats) = 'text'
AND EXISTS (
SELECT 1
FROM json_each(
CASE
WHEN json_valid(CAST(json_extract(pak.api_formats, '$') AS TEXT))
THEN CAST(json_extract(pak.api_formats, '$') AS TEXT)
ELSE '[]'
END
) AS fmt
WHERE LOWER(TRIM(CAST(fmt.value AS TEXT))) IN (
"#,
);
push_bind_list(builder, match_aliases);
builder.push(
r#"
)
)
)
)
ELSE 0
END
OR LOWER(TRIM(pak.api_formats)) IN (
"#,
);
push_bind_list(builder, match_aliases);
builder.push(
r#"
)
)
"#,
);
}
fn push_key_auth_channel_sql_filter(
builder: &mut QueryBuilder<'_, Sqlite>,
storage_api_format: &str,
) {
let api_format = normalize_api_format(storage_api_format);
builder.push(
r#"
AND (
(
LOWER(TRIM(p.provider_type)) = 'codex'
AND LOWER(TRIM(pak.auth_type)) = 'oauth'
AND "#,
);
builder.push_bind(api_format.clone());
builder.push(
r#" IN ('openai:responses', 'openai:responses:compact', 'openai:image')
)
OR (
LOWER(TRIM(p.provider_type)) = 'chatgpt_web'
AND LOWER(TRIM(pak.auth_type)) IN ('oauth', 'bearer')
AND "#,
);
builder.push_bind(api_format.clone());
builder.push(
r#" = 'openai:image'
)
OR (
LOWER(TRIM(p.provider_type)) = 'claude_code'
AND LOWER(TRIM(pak.auth_type)) = 'oauth'
AND "#,
);
builder.push_bind(api_format.clone());
builder.push(
r#" = 'claude:messages'
)
OR (
LOWER(TRIM(p.provider_type)) = 'kiro'
AND "#,
);
builder.push_bind(api_format.clone());
builder.push(
r#" = 'claude:messages'
AND (
LOWER(TRIM(pak.auth_type)) = 'oauth'
OR (
LOWER(TRIM(pak.auth_type)) = 'bearer'
AND pak.auth_config IS NOT NULL
AND TRIM(pak.auth_config) <> ''
)
)
)
OR (
LOWER(TRIM(p.provider_type)) IN ('gemini_cli', 'antigravity')
AND LOWER(TRIM(pak.auth_type)) = 'oauth'
AND "#,
);
builder.push_bind(api_format.clone());
builder.push(
r#" = 'gemini:generate_content'
)
OR (
LOWER(TRIM(p.provider_type)) = 'vertex_ai'
AND (
(
LOWER(TRIM(pak.auth_type)) = 'api_key'
AND "#,
);
builder.push_bind(api_format.clone());
builder.push(
r#" = 'gemini:generate_content'
)
OR (
LOWER(TRIM(pak.auth_type)) IN ('service_account', 'vertex_ai')
AND "#,
);
builder.push_bind(api_format.clone());
builder.push(
r#" IN ('claude:messages', 'gemini:generate_content')
)
)
)
OR (
LOWER(TRIM(p.provider_type)) NOT IN (
'chatgpt_web',
'claude_code',
'codex',
'gemini_cli',
'vertex_ai',
'antigravity',
'kiro'
)
AND LOWER(TRIM(pak.auth_type)) <> 'oauth'
)
)
"#,
);
}
fn push_requested_model_sql_filter(
builder: &mut QueryBuilder<'_, Sqlite>,
requested_model_name: &str,
_match_aliases: &[String],
) {
builder.push(
r#"
AND (
gm.name = "#,
);
builder.push_bind(requested_model_name.to_string());
builder.push(
r#"
OR m.provider_model_name = "#,
);
builder.push_bind(requested_model_name.to_string());
builder.push(
r#"
OR (
m.provider_model_mappings IS NOT NULL
AND m.provider_model_mappings LIKE "#,
);
builder.push_bind(format!(
"%{}%",
requested_model_name
.replace('\\', "\\\\")
.replace('%', "\\%")
.replace('_', "\\_")
));
builder.push(
r#"
ESCAPE '\'
)
)
"#,
);
}
fn push_selected_rows_order(builder: &mut QueryBuilder<'_, Sqlite>, order: SelectedRowsOrder) {
builder.push(" ORDER BY ");
if matches!(order, SelectedRowsOrder::WithGlobalModel) {
builder.push("global_model_name ASC, ");
}
builder.push(
"provider_priority ASC, key_internal_priority ASC, provider_id ASC, endpoint_id ASC, key_id ASC, model_id ASC",
);
}
fn push_pool_key_order(
builder: &mut QueryBuilder<'_, Sqlite>,
order: &StoredPoolKeyCandidateOrder,
) {
match order {
StoredPoolKeyCandidateOrder::InternalPriority => {
builder.push(" ORDER BY pak.internal_priority ASC, pak.id ASC");
}
left.provider_priority
.cmp(&right.provider_priority)
.then(left.key_internal_priority.cmp(&right.key_internal_priority))
.then(left.provider_id.cmp(&right.provider_id))
.then(left.endpoint_id.cmp(&right.endpoint_id))
.then(left.key_id.cmp(&right.key_id))
.then(left.model_id.cmp(&right.model_id))
});
rows
StoredPoolKeyCandidateOrder::Lru => {
builder.push(
" ORDER BY pak.last_used_at IS NOT NULL ASC, pak.last_used_at ASC, pak.internal_priority ASC, pak.id ASC",
);
}
StoredPoolKeyCandidateOrder::CacheAffinity => {
builder.push(
" ORDER BY pak.last_used_at IS NULL ASC, pak.last_used_at DESC, pak.internal_priority ASC, pak.id ASC",
);
}
StoredPoolKeyCandidateOrder::SingleAccount => {
builder.push(
" ORDER BY pak.internal_priority ASC, pak.last_used_at IS NULL ASC, pak.last_used_at DESC, pak.id ASC",
);
}
StoredPoolKeyCandidateOrder::LoadBalance { seed } => {
let _ = seed;
builder.push(" ORDER BY pak.id ASC");
}
}
}
fn push_bind_list(builder: &mut QueryBuilder<'_, Sqlite>, values: &[String]) {
let mut separated = builder.separated(", ");
for value in values {
separated.push_bind(value.clone());
}
}
fn sort_pool_key_rows(
@@ -473,9 +852,8 @@ fn dedupe_candidate_selection_rows(
}
fn map_candidate_selection_row(row: &SqliteRow) -> Result<CandidateSelectionRow, DataLayerError> {
let provider_config = parse_json(row.try_get("provider_config").ok().flatten())?;
let _provider_config = parse_json(row.try_get("provider_config").ok().flatten())?;
let global_model_config = parse_json(row.try_get("global_model_config").ok().flatten())?;
let provider_pool_enabled = json_object_field_present(&provider_config, "pool_advanced");
let global_model_mappings = global_model_config
.as_ref()
.and_then(|value| value.get("model_mappings").cloned());
@@ -528,7 +906,6 @@ fn map_candidate_selection_row(row: &SqliteRow) -> Result<CandidateSelectionRow,
model_is_active: row.try_get("model_is_active").map_sql_err()?,
model_is_available: row.try_get("model_is_available").map_sql_err()?,
},
provider_pool_enabled,
key_auth_config: row.try_get("key_auth_config").map_sql_err()?,
key_last_used_at_unix_secs: row
.try_get::<Option<i64>, _>("key_last_used_at_unix_secs")
@@ -550,13 +927,6 @@ fn parse_json(value: Option<String>) -> Result<Option<serde_json::Value>, DataLa
.transpose()
}
fn json_object_field_present(value: &Option<serde_json::Value>, field: &str) -> bool {
value
.as_ref()
.and_then(|value| value.get(field))
.is_some_and(|value| !value.is_null())
}
fn json_bool(value: &serde_json::Value) -> Option<bool> {
value.as_bool().or_else(|| {
value

View File

@@ -1,16 +1,279 @@
use async_trait::async_trait;
use sqlx::{sqlite::SqliteRow, Row};
use sqlx::{sqlite::SqliteRow, QueryBuilder, Row, Sqlite};
use super::{
InMemoryProviderCatalogReadRepository, ProviderCatalogKeyListQuery,
ProviderCatalogReadRepository, ProviderCatalogWriteRepository, StoredProviderCatalogEndpoint,
StoredProviderCatalogKey, StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats,
StoredProviderCatalogProvider,
ProviderCatalogKeyListOrder, ProviderCatalogKeyListQuery, ProviderCatalogReadRepository,
ProviderCatalogWriteRepository, StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
};
use crate::driver::sqlite::{sqlite_optional_real, SqlitePool};
use crate::error::SqlResultExt;
use crate::DataLayerError;
const LIST_PROVIDERS_BY_IDS_PREFIX: &str = r#"
SELECT
id,
name,
description,
website,
provider_type,
billing_type,
CAST(monthly_quota_usd AS REAL) AS monthly_quota_usd,
CAST(monthly_used_usd AS REAL) AS monthly_used_usd,
quota_reset_day,
quota_last_reset_at AS quota_last_reset_at_unix_secs,
quota_expires_at AS quota_expires_at_unix_secs,
provider_priority,
is_active,
keep_priority_on_conversion,
enable_format_conversion,
concurrent_limit,
max_retries,
proxy,
request_timeout,
stream_first_byte_timeout,
config,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs
FROM providers
WHERE id IN (
"#;
const LIST_ENDPOINTS_BY_IDS_PREFIX: &str = r#"
SELECT
id,
provider_id,
api_format,
api_family,
endpoint_kind,
is_active,
health_score,
base_url,
header_rules,
body_rules,
max_retries,
custom_path,
config,
format_acceptance_config,
proxy,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs
FROM provider_endpoints
WHERE id IN (
"#;
const LIST_ENDPOINTS_BY_PROVIDER_IDS_PREFIX: &str = r#"
SELECT
id,
provider_id,
api_format,
api_family,
endpoint_kind,
is_active,
health_score,
base_url,
header_rules,
body_rules,
max_retries,
custom_path,
config,
format_acceptance_config,
proxy,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs
FROM provider_endpoints
WHERE provider_id IN (
"#;
const LIST_KEYS_BY_IDS_PREFIX: &str = r#"
SELECT
id,
provider_id,
name,
auth_type,
capabilities,
is_active,
api_formats,
auth_type_by_format,
allow_auth_channel_mismatch_formats,
COALESCE(api_key, encrypted_key) AS api_key,
auth_config,
note,
internal_priority,
rate_multipliers,
global_priority_by_format,
allowed_models,
expires_at AS expires_at_unix_secs,
cache_ttl_minutes,
max_probe_interval_minutes,
proxy,
fingerprint,
rpm_limit,
concurrent_limit,
learned_rpm_limit,
concurrent_429_count,
rpm_429_count,
last_429_at AS last_429_at_unix_secs,
last_429_type,
adjustment_history,
utilization_samples,
last_probe_increase_at AS last_probe_increase_at_unix_secs,
last_rpm_peak,
request_count,
total_tokens,
CAST(total_cost_usd AS REAL) AS total_cost_usd,
success_count,
error_count,
total_response_time_ms,
last_used_at AS last_used_at_unix_secs,
auto_fetch_models,
last_models_fetch_at AS last_models_fetch_at_unix_secs,
last_models_fetch_error,
locked_models,
model_include_patterns,
model_exclude_patterns,
upstream_metadata,
oauth_invalid_at AS oauth_invalid_at_unix_secs,
oauth_invalid_reason,
status_snapshot,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs,
health_by_format,
circuit_breaker_by_format
FROM provider_api_keys
WHERE id IN (
"#;
const LIST_KEYS_BY_PROVIDER_IDS_PREFIX: &str = r#"
SELECT
id,
provider_id,
name,
auth_type,
capabilities,
is_active,
api_formats,
auth_type_by_format,
allow_auth_channel_mismatch_formats,
COALESCE(api_key, encrypted_key) AS api_key,
auth_config,
note,
internal_priority,
rate_multipliers,
global_priority_by_format,
allowed_models,
expires_at AS expires_at_unix_secs,
cache_ttl_minutes,
max_probe_interval_minutes,
proxy,
fingerprint,
rpm_limit,
concurrent_limit,
learned_rpm_limit,
concurrent_429_count,
rpm_429_count,
last_429_at AS last_429_at_unix_secs,
last_429_type,
adjustment_history,
utilization_samples,
last_probe_increase_at AS last_probe_increase_at_unix_secs,
last_rpm_peak,
request_count,
total_tokens,
CAST(total_cost_usd AS REAL) AS total_cost_usd,
success_count,
error_count,
total_response_time_ms,
last_used_at AS last_used_at_unix_secs,
auto_fetch_models,
last_models_fetch_at AS last_models_fetch_at_unix_secs,
last_models_fetch_error,
locked_models,
model_include_patterns,
model_exclude_patterns,
upstream_metadata,
oauth_invalid_at AS oauth_invalid_at_unix_secs,
oauth_invalid_reason,
status_snapshot,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs,
health_by_format,
circuit_breaker_by_format
FROM provider_api_keys
WHERE provider_id IN (
"#;
const LIST_KEY_SUMMARIES_BY_PROVIDER_IDS_PREFIX: &str = r#"
SELECT
id,
provider_id,
COALESCE(NULLIF(name, ''), id) AS name,
COALESCE(NULLIF(auth_type, ''), 'summary') AS auth_type,
NULL AS capabilities,
is_active,
api_formats,
NULL AS auth_type_by_format,
NULL AS allow_auth_channel_mismatch_formats,
'summary' AS api_key,
CASE
WHEN auth_config IS NULL THEN NULL
ELSE '{}'
END AS auth_config,
NULL AS note,
NULL AS internal_priority,
NULL AS rate_multipliers,
NULL AS global_priority_by_format,
NULL AS allowed_models,
NULL AS expires_at_unix_secs,
NULL AS cache_ttl_minutes,
NULL AS max_probe_interval_minutes,
NULL AS proxy,
NULL AS fingerprint,
NULL AS rpm_limit,
NULL AS concurrent_limit,
NULL AS learned_rpm_limit,
NULL AS concurrent_429_count,
NULL AS rpm_429_count,
NULL AS last_429_at_unix_secs,
NULL AS last_429_type,
NULL AS adjustment_history,
NULL AS utilization_samples,
NULL AS last_probe_increase_at_unix_secs,
NULL AS last_rpm_peak,
NULL AS request_count,
0 AS total_tokens,
0.0 AS total_cost_usd,
NULL AS success_count,
NULL AS error_count,
NULL AS total_response_time_ms,
NULL AS last_used_at_unix_secs,
FALSE AS auto_fetch_models,
NULL AS last_models_fetch_at_unix_secs,
NULL AS last_models_fetch_error,
NULL AS locked_models,
NULL AS model_include_patterns,
NULL AS model_exclude_patterns,
NULL AS upstream_metadata,
NULL AS oauth_invalid_at_unix_secs,
NULL AS oauth_invalid_reason,
NULL AS status_snapshot,
NULL AS created_at_unix_ms,
NULL AS updated_at_unix_secs,
health_by_format,
NULL AS circuit_breaker_by_format
FROM provider_api_keys
WHERE provider_id IN (
"#;
const LIST_KEY_STATS_BY_PROVIDER_IDS_PREFIX: &str = r#"
SELECT
provider_id,
COUNT(*) AS total_keys,
SUM(CASE WHEN is_active THEN 1 ELSE 0 END) AS active_keys
FROM provider_api_keys
WHERE provider_id IN (
"#;
#[derive(Debug, Clone)]
pub struct SqliteProviderCatalogReadRepository {
pool: SqlitePool,
@@ -21,92 +284,340 @@ impl SqliteProviderCatalogReadRepository {
Self { pool }
}
async fn load_memory(&self) -> Result<InMemoryProviderCatalogReadRepository, DataLayerError> {
Ok(InMemoryProviderCatalogReadRepository::seed(
self.load_providers().await?,
self.load_endpoints().await?,
self.load_keys().await?,
))
}
pub async fn list_providers_by_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
if provider_ids.is_empty() {
return Ok(Vec::new());
}
async fn load_providers(&self) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
let rows = sqlx::query(
r#"
SELECT
id, name, description, website, provider_type, billing_type,
CAST(monthly_quota_usd AS REAL) AS monthly_quota_usd,
CAST(monthly_used_usd AS REAL) AS monthly_used_usd,
quota_reset_day,
quota_last_reset_at AS quota_last_reset_at_unix_secs,
quota_expires_at AS quota_expires_at_unix_secs,
provider_priority, is_active, keep_priority_on_conversion,
enable_format_conversion, concurrent_limit, max_retries, proxy,
request_timeout, stream_first_byte_timeout, config,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs
FROM providers
"#,
let rows = build_list_query(
LIST_PROVIDERS_BY_IDS_PREFIX,
provider_ids,
" ORDER BY name ASC",
)
.build()
.fetch_all(&self.pool)
.await
.map_sql_err()?;
rows.iter().map(map_provider_row).collect()
}
async fn load_endpoints(&self) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
pub async fn list_providers(
&self,
active_only: bool,
) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
let rows = sqlx::query(
r#"
SELECT
id, provider_id, api_format, api_family, endpoint_kind, is_active,
health_score, base_url, header_rules, body_rules, max_retries,
custom_path, config, format_acceptance_config, proxy,
id,
name,
description,
website,
provider_type,
billing_type,
CAST(monthly_quota_usd AS REAL) AS monthly_quota_usd,
CAST(monthly_used_usd AS REAL) AS monthly_used_usd,
quota_reset_day,
quota_last_reset_at AS quota_last_reset_at_unix_secs,
quota_expires_at AS quota_expires_at_unix_secs,
provider_priority,
is_active,
keep_priority_on_conversion,
enable_format_conversion,
concurrent_limit,
max_retries,
proxy,
request_timeout,
stream_first_byte_timeout,
config,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs
FROM provider_endpoints
WHERE api_format IS NOT NULL
FROM providers
WHERE (? = FALSE OR is_active = TRUE)
ORDER BY provider_priority ASC, name ASC
"#,
)
.bind(active_only)
.fetch_all(&self.pool)
.await
.map_sql_err()?;
rows.iter().map(map_provider_row).collect()
}
pub async fn list_endpoints_by_ids(
&self,
endpoint_ids: &[String],
) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
if endpoint_ids.is_empty() {
return Ok(Vec::new());
}
let rows = build_list_query(
LIST_ENDPOINTS_BY_IDS_PREFIX,
endpoint_ids,
" ORDER BY api_format ASC, id ASC",
)
.build()
.fetch_all(&self.pool)
.await
.map_sql_err()?;
rows.iter().map(map_endpoint_row).collect()
}
async fn load_keys(&self) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
let rows = sqlx::query(
r#"
SELECT
id, provider_id, name, auth_type, capabilities, is_active, api_formats,
auth_type_by_format, allow_auth_channel_mismatch_formats,
COALESCE(api_key, encrypted_key) AS api_key,
auth_config, note, internal_priority, rate_multipliers,
global_priority_by_format, allowed_models,
expires_at AS expires_at_unix_secs,
cache_ttl_minutes, max_probe_interval_minutes, proxy, fingerprint,
rpm_limit, concurrent_limit, learned_rpm_limit, concurrent_429_count,
rpm_429_count, last_429_at AS last_429_at_unix_secs, last_429_type,
adjustment_history, utilization_samples,
last_probe_increase_at AS last_probe_increase_at_unix_secs,
last_rpm_peak, request_count, total_tokens, total_cost_usd,
success_count, error_count, total_response_time_ms,
last_used_at AS last_used_at_unix_secs, auto_fetch_models,
last_models_fetch_at AS last_models_fetch_at_unix_secs,
last_models_fetch_error, locked_models, model_include_patterns,
model_exclude_patterns, upstream_metadata,
oauth_invalid_at AS oauth_invalid_at_unix_secs,
oauth_invalid_reason, status_snapshot,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs,
health_by_format, circuit_breaker_by_format
FROM provider_api_keys
"#,
pub async fn list_endpoints_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
if provider_ids.is_empty() {
return Ok(Vec::new());
}
let rows = build_list_query(
LIST_ENDPOINTS_BY_PROVIDER_IDS_PREFIX,
provider_ids,
" ORDER BY provider_id ASC, api_format ASC, id ASC",
)
.build()
.fetch_all(&self.pool)
.await
.map_sql_err()?;
rows.iter().map(map_endpoint_row).collect()
}
pub async fn list_keys_by_ids(
&self,
key_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
if key_ids.is_empty() {
return Ok(Vec::new());
}
let rows = build_list_query(
LIST_KEYS_BY_IDS_PREFIX,
key_ids,
" ORDER BY name ASC, id ASC",
)
.build()
.fetch_all(&self.pool)
.await
.map_sql_err()?;
rows.iter().map(map_key_row).collect()
}
pub async fn list_keys_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
if provider_ids.is_empty() {
return Ok(Vec::new());
}
let rows = build_list_query(
LIST_KEYS_BY_PROVIDER_IDS_PREFIX,
provider_ids,
" ORDER BY provider_id ASC, name ASC, id ASC",
)
.build()
.fetch_all(&self.pool)
.await
.map_sql_err()?;
rows.iter().map(map_key_row).collect()
}
pub async fn list_key_summaries_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
if provider_ids.is_empty() {
return Ok(Vec::new());
}
let rows = build_list_query(
LIST_KEY_SUMMARIES_BY_PROVIDER_IDS_PREFIX,
provider_ids,
" ORDER BY provider_id ASC, id ASC",
)
.build()
.fetch_all(&self.pool)
.await
.map_sql_err()?;
rows.iter().map(map_key_row).collect()
}
pub async fn list_keys_page(
&self,
query: &ProviderCatalogKeyListQuery,
) -> Result<StoredProviderCatalogKeyPage, DataLayerError> {
if query.provider_id.trim().is_empty() {
return Err(DataLayerError::InvalidInput(
"provider catalog provider_id is empty".to_string(),
));
}
let offset = i64::try_from(query.offset).map_err(|_| {
DataLayerError::InvalidInput(format!(
"invalid provider catalog key offset: {}",
query.offset
))
})?;
let limit = i64::try_from(query.limit).map_err(|_| {
DataLayerError::InvalidInput(format!(
"invalid provider catalog key limit: {}",
query.limit
))
})?;
let search_pattern = query
.search
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(|value| format!("%{}%", value.to_ascii_lowercase()));
let order_by = match query.order {
ProviderCatalogKeyListOrder::Name => "internal_priority ASC, name ASC, id ASC",
ProviderCatalogKeyListOrder::CreatedAt => {
"internal_priority ASC, COALESCE(created_at, 0) ASC, id ASC"
}
ProviderCatalogKeyListOrder::CreatedAtAsc => {
"created_at IS NULL ASC, created_at ASC, name ASC, id ASC"
}
ProviderCatalogKeyListOrder::CreatedAtDesc => {
"created_at IS NULL ASC, created_at DESC, name ASC, id ASC"
}
ProviderCatalogKeyListOrder::LastUsedAtAsc => {
"last_used_at IS NULL ASC, last_used_at ASC, name ASC, id ASC"
}
ProviderCatalogKeyListOrder::LastUsedAtDesc => {
"last_used_at IS NULL ASC, last_used_at DESC, name ASC, id ASC"
}
};
let count_row = sqlx::query(
r#"
SELECT COUNT(*) AS total
FROM provider_api_keys
WHERE provider_id = ?
AND (? IS NULL OR LOWER(name) LIKE ? OR LOWER(id) LIKE ?)
AND (? IS NULL OR is_active = ?)
"#,
)
.bind(&query.provider_id)
.bind(search_pattern.as_deref())
.bind(search_pattern.as_deref())
.bind(search_pattern.as_deref())
.bind(query.is_active)
.bind(query.is_active)
.fetch_one(&self.pool)
.await
.map_sql_err()?;
let total = count_row.try_get::<i64, _>("total").map_sql_err()?.max(0) as usize;
let sql = format!(
r#"
SELECT
id,
provider_id,
name,
auth_type,
capabilities,
is_active,
api_formats,
auth_type_by_format,
allow_auth_channel_mismatch_formats,
COALESCE(api_key, encrypted_key) AS api_key,
auth_config,
note,
internal_priority,
rate_multipliers,
global_priority_by_format,
allowed_models,
expires_at AS expires_at_unix_secs,
cache_ttl_minutes,
max_probe_interval_minutes,
proxy,
fingerprint,
rpm_limit,
concurrent_limit,
learned_rpm_limit,
concurrent_429_count,
rpm_429_count,
last_429_at AS last_429_at_unix_secs,
last_429_type,
adjustment_history,
utilization_samples,
last_probe_increase_at AS last_probe_increase_at_unix_secs,
last_rpm_peak,
request_count,
total_tokens,
CAST(total_cost_usd AS REAL) AS total_cost_usd,
success_count,
error_count,
total_response_time_ms,
last_used_at AS last_used_at_unix_secs,
auto_fetch_models,
last_models_fetch_at AS last_models_fetch_at_unix_secs,
last_models_fetch_error,
locked_models,
model_include_patterns,
model_exclude_patterns,
upstream_metadata,
oauth_invalid_at AS oauth_invalid_at_unix_secs,
oauth_invalid_reason,
status_snapshot,
created_at AS created_at_unix_ms,
updated_at AS updated_at_unix_secs,
health_by_format,
circuit_breaker_by_format
FROM provider_api_keys
WHERE provider_id = ?
AND (? IS NULL OR LOWER(name) LIKE ? OR LOWER(id) LIKE ?)
AND (? IS NULL OR is_active = ?)
ORDER BY {order_by}
LIMIT ?
OFFSET ?
"#,
);
let rows = sqlx::query(&sql)
.bind(&query.provider_id)
.bind(search_pattern.as_deref())
.bind(search_pattern.as_deref())
.bind(search_pattern.as_deref())
.bind(query.is_active)
.bind(query.is_active)
.bind(limit)
.bind(offset)
.fetch_all(&self.pool)
.await
.map_sql_err()?;
let items = rows
.iter()
.map(map_key_row)
.collect::<Result<Vec<_>, _>>()?;
Ok(StoredProviderCatalogKeyPage { items, total })
}
pub async fn list_key_stats_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKeyStats>, DataLayerError> {
if provider_ids.is_empty() {
return Ok(Vec::new());
}
let rows = build_list_query(
LIST_KEY_STATS_BY_PROVIDER_IDS_PREFIX,
provider_ids,
"\nGROUP BY provider_id\nORDER BY provider_id ASC",
)
.build()
.fetch_all(&self.pool)
.await
.map_sql_err()?;
rows.iter().map(map_key_stats_row).collect()
}
pub async fn create_provider(
&self,
provider: &StoredProviderCatalogProvider,
@@ -865,81 +1376,63 @@ impl ProviderCatalogReadRepository for SqliteProviderCatalogReadRepository {
&self,
active_only: bool,
) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
self.load_memory().await?.list_providers(active_only).await
Self::list_providers(self, active_only).await
}
async fn list_providers_by_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
self.load_memory()
.await?
.list_providers_by_ids(provider_ids)
.await
Self::list_providers_by_ids(self, provider_ids).await
}
async fn list_endpoints_by_ids(
&self,
endpoint_ids: &[String],
) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
self.load_memory()
.await?
.list_endpoints_by_ids(endpoint_ids)
.await
Self::list_endpoints_by_ids(self, endpoint_ids).await
}
async fn list_endpoints_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
self.load_memory()
.await?
.list_endpoints_by_provider_ids(provider_ids)
.await
Self::list_endpoints_by_provider_ids(self, provider_ids).await
}
async fn list_keys_by_ids(
&self,
key_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
self.load_memory().await?.list_keys_by_ids(key_ids).await
Self::list_keys_by_ids(self, key_ids).await
}
async fn list_keys_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
self.load_memory()
.await?
.list_keys_by_provider_ids(provider_ids)
.await
Self::list_keys_by_provider_ids(self, provider_ids).await
}
async fn list_key_summaries_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
self.load_memory()
.await?
.list_key_summaries_by_provider_ids(provider_ids)
.await
Self::list_key_summaries_by_provider_ids(self, provider_ids).await
}
async fn list_keys_page(
&self,
query: &ProviderCatalogKeyListQuery,
) -> Result<StoredProviderCatalogKeyPage, DataLayerError> {
self.load_memory().await?.list_keys_page(query).await
Self::list_keys_page(self, query).await
}
async fn list_key_stats_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKeyStats>, DataLayerError> {
self.load_memory()
.await?
.list_key_stats_by_provider_ids(provider_ids)
.await
Self::list_key_stats_by_provider_ids(self, provider_ids).await
}
}
@@ -1058,6 +1551,21 @@ impl ProviderCatalogWriteRepository for SqliteProviderCatalogReadRepository {
}
}
fn build_list_query<'a>(
prefix: &'static str,
ids: &'a [String],
suffix: &'static str,
) -> QueryBuilder<'a, Sqlite> {
let mut builder = QueryBuilder::<Sqlite>::new(prefix);
let mut separated = builder.separated(", ");
for id in ids {
separated.push_bind(id);
}
separated.push_unseparated(")");
builder.push(suffix);
builder
}
fn current_unix_secs() -> u64 {
chrono::Utc::now().timestamp().max(0) as u64
}
@@ -1346,6 +1854,14 @@ fn map_endpoint_row(row: &SqliteRow) -> Result<StoredProviderCatalogEndpoint, Da
)
}
fn map_key_stats_row(row: &SqliteRow) -> Result<StoredProviderCatalogKeyStats, DataLayerError> {
StoredProviderCatalogKeyStats::new(
row.try_get("provider_id").map_sql_err()?,
row.try_get("total_keys").map_sql_err()?,
row.try_get("active_keys").map_sql_err()?,
)
}
fn map_key_row(row: &SqliteRow) -> Result<StoredProviderCatalogKey, DataLayerError> {
let total_cost_usd = sqlite_optional_real(row, "total_cost_usd")?.unwrap_or(0.0);
if !total_cost_usd.is_finite() {

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff