mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-10 05:00:19 +08:00
Merge branch 'main' of https://github.com/zhefox/Aether
This commit is contained in:
+3
-2
@@ -69,8 +69,9 @@ ADMIN_USERNAME=admin123456
|
||||
# AETHER_DOCKER_UPDATE_COMMAND=./update.sh
|
||||
# AETHER_GATEWAY_DEPLOYMENT_TOPOLOGY=single-node
|
||||
# AETHER_GATEWAY_NODE_ROLE=all
|
||||
# Docker Compose 默认把应用日志输出到 stdout/stderr。
|
||||
# 如需文件日志,可改成 file 或 both,并把可写目录挂载到 /opt/aether/logs。
|
||||
# Docker Compose 默认强制把应用日志输出到 stdout/stderr,并由 Docker 轮转日志。
|
||||
# 如需文件日志,需要在 compose 里把 AETHER_LOG_DESTINATION 改成 file 或 both,
|
||||
# 并把容器用户可写目录挂载到 /opt/aether/logs。
|
||||
# AETHER_LOG_DESTINATION=stdout
|
||||
# AETHER_LOG_FORMAT=pretty
|
||||
# AETHER_LOG_DIR=/opt/aether/logs
|
||||
|
||||
Generated
+1
-1
@@ -478,7 +478,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "aether-tunnel"
|
||||
version = "0.3.12"
|
||||
version = "0.3.13"
|
||||
dependencies = [
|
||||
"aether-contracts",
|
||||
"aether-gateway",
|
||||
|
||||
+2
-2
@@ -25,7 +25,7 @@ RUN ln -s /opt/aether/releases/image /opt/aether/current
|
||||
# --- final stage: distroless runtime ---
|
||||
FROM gcr.io/distroless/static-debian12
|
||||
|
||||
COPY --from=layout --chown=65532:65532 /opt/aether /opt/aether
|
||||
COPY --from=layout /opt/aether /opt/aether
|
||||
|
||||
WORKDIR /opt/aether
|
||||
|
||||
@@ -39,5 +39,5 @@ EXPOSE 8084
|
||||
HEALTHCHECK --interval=30s --timeout=10s --start-period=5s --retries=3 \
|
||||
CMD ["/opt/aether/current/bin/aether-gateway", "--healthcheck"]
|
||||
|
||||
USER 65532:65532
|
||||
USER root
|
||||
ENTRYPOINT ["/opt/aether/current/bin/aether-gateway"]
|
||||
|
||||
@@ -69,7 +69,7 @@ Docker Compose 部署后,可在部署目录直接执行:
|
||||
./update.sh --mode single-node
|
||||
```
|
||||
|
||||
仓库自带的 Docker Compose 默认把应用日志输出到容器 `stdout/stderr`,直接用 `docker compose logs -f app` 查看,避免正式发布镜像切换到非 root 用户后再被宿主机挂载日志目录的权限问题拖垮启动。如果你确实需要文件日志,再显式设置 `AETHER_LOG_DESTINATION=file|both`,并把一个可写目录挂载到 `/opt/aether/logs`(或同步覆盖 `AETHER_LOG_DIR`)。
|
||||
仓库自带的 Docker Compose 默认把应用日志输出到容器 `stdout/stderr`,直接用 `docker compose logs -f app` 查看,并由 Docker 轮转日志,避免正式发布镜像切换到非 root 用户后再被宿主机挂载日志目录的权限问题拖垮启动。如果你确实需要文件日志,需要在 compose 里把 `AETHER_LOG_DESTINATION` 改成 `file|both`,并额外挂载一个容器用户可写的目录到 `/opt/aether/logs`。
|
||||
|
||||
管理后台右上角“版本信息”会检测新版本。Docker Compose 部署只提示版本,实际更新继续执行 `./update.sh`;systemd / launchd / 二进制部署才使用后台自更新,流程是下载对应平台的 GitHub Release 包、强制校验 `SHA256SUMS`、解压到 `/opt/aether/releases/<version>`,再切换 `/opt/aether/current` 并退出进程,交给 systemd / launchd 拉起新版本。
|
||||
|
||||
|
||||
@@ -2,7 +2,8 @@ use super::{
|
||||
ApiKeyLastUsedDelta, DataLayerError, GatewayDataState, GeminiFileMappingListQuery,
|
||||
GeminiFileMappingStats, ProviderCatalogKeyListQuery, PublicHealthStatusCount,
|
||||
PublicHealthTimelineBucket, StoredGeminiFileMapping, StoredGeminiFileMappingListPage,
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
|
||||
StoredProviderCatalogKeyMaintenanceSummary, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogKeyStats, StoredProviderCatalogProvider, StoredRequestCandidate,
|
||||
UpsertGeminiFileMappingRecord, UpsertRequestCandidateRecord,
|
||||
};
|
||||
@@ -285,6 +286,20 @@ impl GatewayDataState {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn list_provider_catalog_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
match &self.provider_catalog_reader {
|
||||
Some(repository) => {
|
||||
repository
|
||||
.list_key_maintenance_summaries_by_provider_ids(provider_ids)
|
||||
.await
|
||||
}
|
||||
None => Ok(Vec::new()),
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn list_provider_catalog_key_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
|
||||
@@ -119,7 +119,8 @@ use aether_data_contracts::repository::pool_scores::{
|
||||
};
|
||||
use aether_data_contracts::repository::provider_catalog::{
|
||||
ProviderCatalogKeyListQuery, ProviderCatalogReadRepository, ProviderCatalogWriteRepository,
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
|
||||
StoredProviderCatalogKeyMaintenanceSummary, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
};
|
||||
use aether_data_contracts::repository::quota::{
|
||||
|
||||
@@ -53,6 +53,7 @@ const POOL_ACTIVE_PROBE_SEALED_SKIP_REASON: &str = "pool_active_probe_sealed";
|
||||
const ROUTING_PROFILE_DISALLOWED_KEY_SKIP_REASON: &str = "routing_profile_disallowed_key";
|
||||
const POOL_SCORE_SCHEDULE_INTEREST_CONCURRENCY: usize = 4;
|
||||
const POOL_SCORE_SCHEDULE_INTEREST_MAX_PER_BATCH: usize = 16;
|
||||
const POOL_SCORE_SCHEDULE_INTEREST_MIN_INTERVAL_SECS: u64 = 60;
|
||||
|
||||
type PoolCatalogKeyContext = PoolMemberSignals;
|
||||
|
||||
@@ -862,30 +863,18 @@ impl<'a> PoolKeyCursor<'a> {
|
||||
return;
|
||||
}
|
||||
|
||||
let Ok(permit) = POOL_SCORE_SCHEDULE_INTEREST_SEMAPHORE
|
||||
.clone()
|
||||
.try_acquire_owned()
|
||||
else {
|
||||
debug!(
|
||||
event_name = "pool_group_score_interest_dropped",
|
||||
log_type = "event",
|
||||
provider_id = %self.group.candidate.provider_id,
|
||||
endpoint_id = %self.group.candidate.endpoint_id,
|
||||
model_id = %self.group.candidate.model_id,
|
||||
score_count = scores.len(),
|
||||
"gateway pool scheduler dropped score schedule interest because the background writer is saturated"
|
||||
);
|
||||
return;
|
||||
};
|
||||
|
||||
let scheduled_at = current_unix_ms() / 1000;
|
||||
let provider_id = self.group.candidate.provider_id.clone();
|
||||
let endpoint_id = self.group.candidate.endpoint_id.clone();
|
||||
let model_id = self.group.candidate.model_id.clone();
|
||||
let score_count = scores.len().min(POOL_SCORE_SCHEDULE_INTEREST_MAX_PER_BATCH);
|
||||
let app = self.state.app().clone();
|
||||
let feedback = scores
|
||||
.iter()
|
||||
.filter(|score| {
|
||||
score.last_scheduled_at.is_none_or(|last_scheduled_at| {
|
||||
scheduled_at.saturating_sub(last_scheduled_at)
|
||||
>= POOL_SCORE_SCHEDULE_INTEREST_MIN_INTERVAL_SECS
|
||||
})
|
||||
})
|
||||
.take(POOL_SCORE_SCHEDULE_INTEREST_MAX_PER_BATCH)
|
||||
.map(|score| PoolMemberScheduleFeedback {
|
||||
identity: PoolMemberIdentity {
|
||||
@@ -912,6 +901,28 @@ impl<'a> PoolKeyCursor<'a> {
|
||||
})),
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
let score_count = feedback.len();
|
||||
if feedback.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
let Ok(permit) = POOL_SCORE_SCHEDULE_INTEREST_SEMAPHORE
|
||||
.clone()
|
||||
.try_acquire_owned()
|
||||
else {
|
||||
debug!(
|
||||
event_name = "pool_group_score_interest_dropped",
|
||||
log_type = "event",
|
||||
provider_id = %self.group.candidate.provider_id,
|
||||
endpoint_id = %self.group.candidate.endpoint_id,
|
||||
model_id = %self.group.candidate.model_id,
|
||||
score_count = scores.len(),
|
||||
"gateway pool scheduler dropped score schedule interest because the background writer is saturated"
|
||||
);
|
||||
return;
|
||||
};
|
||||
|
||||
let app = self.state.app().clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
let _permit = permit;
|
||||
|
||||
@@ -201,6 +201,7 @@ pub(crate) async fn update_existing_provider_oauth_catalog_key(
|
||||
.map(|duration| duration.as_secs())
|
||||
.unwrap_or(0);
|
||||
let mut updated = existing_key.clone();
|
||||
updated.is_active = true;
|
||||
updated.encrypted_api_key = Some(encrypted_api_key);
|
||||
updated.encrypted_auth_config = Some(encrypted_auth_config);
|
||||
updated.api_formats = provider_oauth_catalog_key_api_formats(provider_type, api_formats);
|
||||
|
||||
@@ -1,7 +1,11 @@
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use std::fmt::Debug;
|
||||
use std::future::Future;
|
||||
use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
||||
|
||||
use aether_data_contracts::repository::candidate_selection::StoredMinimalCandidateSelectionRow;
|
||||
use axum::{body::Body, response::Response};
|
||||
use tokio::time::timeout;
|
||||
use tracing::warn;
|
||||
|
||||
use super::models_responses::{
|
||||
build_claude_model_detail_response, build_claude_models_list_response,
|
||||
@@ -15,6 +19,59 @@ use super::models_shared::{
|
||||
};
|
||||
use super::{query_param_value, AppState, GatewayPublicRequestContext};
|
||||
|
||||
#[cfg(not(test))]
|
||||
const MODELS_ROUTE_READ_TIMEOUT: Duration = Duration::from_secs(5);
|
||||
#[cfg(test)]
|
||||
const MODELS_ROUTE_READ_TIMEOUT: Duration = Duration::from_millis(50);
|
||||
|
||||
async fn await_models_route_read<T, E, Fut>(operation: &'static str, future: Fut) -> Option<T>
|
||||
where
|
||||
E: Debug,
|
||||
Fut: Future<Output = Result<T, E>>,
|
||||
{
|
||||
match timeout(MODELS_ROUTE_READ_TIMEOUT, future).await {
|
||||
Ok(Ok(value)) => Some(value),
|
||||
Ok(Err(error)) => {
|
||||
warn!(
|
||||
event_name = "models_route_read_error",
|
||||
log_type = "ops",
|
||||
operation,
|
||||
error = ?error,
|
||||
"gateway local models route read failed"
|
||||
);
|
||||
None
|
||||
}
|
||||
Err(_) => {
|
||||
warn!(
|
||||
event_name = "models_route_read_timeout",
|
||||
log_type = "ops",
|
||||
operation,
|
||||
timeout_ms = MODELS_ROUTE_READ_TIMEOUT.as_millis() as u64,
|
||||
"gateway local models route read timed out"
|
||||
);
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn build_models_read_fallback_response(
|
||||
request_context: &GatewayPublicRequestContext,
|
||||
api_format: &str,
|
||||
) -> Response<Body> {
|
||||
let route_kind = request_context
|
||||
.control_decision
|
||||
.as_ref()
|
||||
.and_then(|decision| decision.route_kind.as_deref());
|
||||
match route_kind {
|
||||
Some("detail") => {
|
||||
let model_id = models_detail_id(&request_context.request_path)
|
||||
.unwrap_or_else(|| "unknown".to_string());
|
||||
build_models_not_found_response(&model_id, api_format)
|
||||
}
|
||||
_ => build_empty_models_list_response(api_format),
|
||||
}
|
||||
}
|
||||
|
||||
fn sort_and_dedup_model_rows(
|
||||
mut rows: Vec<StoredMinimalCandidateSelectionRow>,
|
||||
) -> Vec<StoredMinimalCandidateSelectionRow> {
|
||||
@@ -47,10 +104,11 @@ async fn list_model_rows_for_client_format(
|
||||
) -> Option<Vec<StoredMinimalCandidateSelectionRow>> {
|
||||
let mut collected = Vec::new();
|
||||
for query_format in models_query_api_formats(api_format) {
|
||||
let rows = state
|
||||
.list_minimal_candidate_selection_rows_for_api_format(query_format)
|
||||
.await
|
||||
.ok()?;
|
||||
let rows = await_models_route_read(
|
||||
"candidate_selection_by_api_format",
|
||||
state.list_minimal_candidate_selection_rows_for_api_format(query_format),
|
||||
)
|
||||
.await?;
|
||||
let mut filtered = filter_rows_for_models(rows, auth_snapshot, query_format);
|
||||
collected.append(&mut filtered);
|
||||
}
|
||||
@@ -65,13 +123,14 @@ async fn list_model_rows_for_client_format_and_global_model(
|
||||
) -> Option<Vec<StoredMinimalCandidateSelectionRow>> {
|
||||
let mut collected = Vec::new();
|
||||
for query_format in models_query_api_formats(api_format) {
|
||||
let rows = state
|
||||
.list_minimal_candidate_selection_rows_for_api_format_and_global_model(
|
||||
let rows = await_models_route_read(
|
||||
"candidate_selection_by_global_model",
|
||||
state.list_minimal_candidate_selection_rows_for_api_format_and_global_model(
|
||||
query_format,
|
||||
global_model_name,
|
||||
)
|
||||
.await
|
||||
.ok()?;
|
||||
),
|
||||
)
|
||||
.await?;
|
||||
let mut filtered = filter_rows_for_models(rows, auth_snapshot, query_format);
|
||||
collected.append(&mut filtered);
|
||||
}
|
||||
@@ -96,21 +155,38 @@ pub(super) async fn maybe_build_local_models_route_response(
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
let auth_snapshot = state
|
||||
.data
|
||||
.read_auth_api_key_snapshot(
|
||||
let auth_snapshot = match await_models_route_read(
|
||||
"auth_api_key_snapshot",
|
||||
state.data.read_auth_api_key_snapshot(
|
||||
&auth_context.user_id,
|
||||
&auth_context.api_key_id,
|
||||
now_unix_secs,
|
||||
)
|
||||
.await
|
||||
.ok()
|
||||
.flatten();
|
||||
),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Some(snapshot) => snapshot,
|
||||
None => {
|
||||
return Some(build_models_read_fallback_response(
|
||||
request_context,
|
||||
api_format,
|
||||
))
|
||||
}
|
||||
};
|
||||
let auth_snapshot = auth_snapshot.as_ref();
|
||||
|
||||
match decision.route_kind.as_deref() {
|
||||
Some("list") => {
|
||||
let rows = list_model_rows_for_client_format(state, api_format, auth_snapshot).await?;
|
||||
let rows =
|
||||
match list_model_rows_for_client_format(state, api_format, auth_snapshot).await {
|
||||
Some(rows) => rows,
|
||||
None => {
|
||||
return Some(build_models_read_fallback_response(
|
||||
request_context,
|
||||
api_format,
|
||||
))
|
||||
}
|
||||
};
|
||||
if rows.is_empty() {
|
||||
return Some(build_empty_models_list_response(api_format));
|
||||
}
|
||||
@@ -156,13 +232,22 @@ pub(super) async fn maybe_build_local_models_route_response(
|
||||
}
|
||||
Some("detail") => {
|
||||
let model_id = models_detail_id(&request_context.request_path)?;
|
||||
let rows = list_model_rows_for_client_format_and_global_model(
|
||||
let rows = match list_model_rows_for_client_format_and_global_model(
|
||||
state,
|
||||
api_format,
|
||||
&model_id,
|
||||
auth_snapshot,
|
||||
)
|
||||
.await?;
|
||||
.await
|
||||
{
|
||||
Some(rows) => rows,
|
||||
None => {
|
||||
return Some(build_models_read_fallback_response(
|
||||
request_context,
|
||||
api_format,
|
||||
))
|
||||
}
|
||||
};
|
||||
let Some(row) = rows.first() else {
|
||||
return Some(build_models_not_found_response(&model_id, api_format));
|
||||
};
|
||||
|
||||
@@ -312,17 +312,22 @@ async fn select_keys_for_provider(
|
||||
}
|
||||
|
||||
let result = async {
|
||||
let keys = state
|
||||
.list_provider_catalog_keys_by_provider_ids(std::slice::from_ref(&provider.id))
|
||||
let summaries = state
|
||||
.list_provider_catalog_key_maintenance_summaries_by_provider_ids(std::slice::from_ref(
|
||||
&provider.id,
|
||||
))
|
||||
.await?
|
||||
.into_iter()
|
||||
.filter(|key| key.is_active)
|
||||
.filter(|summary| summary.is_active)
|
||||
.collect::<Vec<_>>();
|
||||
if keys.is_empty() {
|
||||
if summaries.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
let key_ids = keys.iter().map(|key| key.id.clone()).collect::<Vec<_>>();
|
||||
let key_ids = summaries
|
||||
.iter()
|
||||
.map(|summary| summary.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let check_stamps = load_check_timestamps(runtime, &provider.id, &key_ids).await;
|
||||
let selected_ids = select_account_self_check_key_ids(
|
||||
&key_ids,
|
||||
@@ -344,7 +349,9 @@ async fn select_keys_for_provider(
|
||||
)
|
||||
.await;
|
||||
|
||||
let mut keys_by_id = keys
|
||||
let mut keys_by_id = state
|
||||
.list_provider_catalog_keys_by_ids(&selected_ids)
|
||||
.await?
|
||||
.into_iter()
|
||||
.map(|key| (key.id.clone(), key))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
|
||||
@@ -849,13 +849,15 @@ async fn select_keys_for_provider(
|
||||
}
|
||||
|
||||
let result = async {
|
||||
let keys = state
|
||||
.list_provider_catalog_keys_by_provider_ids(std::slice::from_ref(&provider.id))
|
||||
let summaries = state
|
||||
.list_provider_catalog_key_maintenance_summaries_by_provider_ids(std::slice::from_ref(
|
||||
&provider.id,
|
||||
))
|
||||
.await?
|
||||
.into_iter()
|
||||
.filter(|key| key.is_active)
|
||||
.filter(|summary| summary.is_active)
|
||||
.collect::<Vec<_>>();
|
||||
if keys.is_empty() {
|
||||
if summaries.is_empty() {
|
||||
return Ok(PoolQuotaProbeSelectionOutcome::Empty);
|
||||
}
|
||||
|
||||
@@ -864,7 +866,7 @@ async fn select_keys_for_provider(
|
||||
sample_provider_pool_demand(
|
||||
runtime,
|
||||
&provider.id,
|
||||
keys.len(),
|
||||
summaries.len(),
|
||||
config.max_keys_per_provider,
|
||||
)
|
||||
.await
|
||||
@@ -873,14 +875,14 @@ async fn select_keys_for_provider(
|
||||
read_provider_pool_demand_snapshot(
|
||||
runtime,
|
||||
&provider.id,
|
||||
keys.len(),
|
||||
summaries.len(),
|
||||
config.max_keys_per_provider,
|
||||
)
|
||||
.await
|
||||
}
|
||||
};
|
||||
let target_active_count = pool_quota_probe_target_count_for_mode(
|
||||
keys.len(),
|
||||
summaries.len(),
|
||||
pool_config,
|
||||
demand_snapshot.desired_hot,
|
||||
config,
|
||||
@@ -894,7 +896,10 @@ async fn select_keys_for_provider(
|
||||
return Ok(PoolQuotaProbeSelectionOutcome::Empty);
|
||||
}
|
||||
|
||||
let key_ids = keys.iter().map(|key| key.id.clone()).collect::<Vec<_>>();
|
||||
let key_ids = summaries
|
||||
.iter()
|
||||
.map(|summary| summary.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let scores_by_key = load_provider_key_account_scores(state, &provider.id, &key_ids).await;
|
||||
let mut active_member_ids =
|
||||
load_pruned_active_probe_member_ids(runtime, &provider.id, &key_ids, &scores_by_key)
|
||||
@@ -965,7 +970,9 @@ async fn select_keys_for_provider(
|
||||
)
|
||||
.await;
|
||||
|
||||
let mut keys_by_id = keys
|
||||
let mut keys_by_id = state
|
||||
.list_provider_catalog_keys_by_ids(&selected_ids)
|
||||
.await?
|
||||
.into_iter()
|
||||
.map(|key| (key.id.clone(), key))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
|
||||
@@ -230,18 +230,21 @@ pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
.iter()
|
||||
.map(|(provider, _)| provider.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let mut keys_by_provider = BTreeMap::new();
|
||||
for key in state
|
||||
.list_provider_catalog_keys_by_provider_ids(&provider_ids)
|
||||
let mut key_ids_by_provider = BTreeMap::new();
|
||||
for summary in state
|
||||
.list_provider_catalog_key_maintenance_summaries_by_provider_ids(&provider_ids)
|
||||
.await?
|
||||
{
|
||||
keys_by_provider
|
||||
.entry(key.provider_id.clone())
|
||||
if !summary.is_active {
|
||||
continue;
|
||||
}
|
||||
key_ids_by_provider
|
||||
.entry(summary.provider_id.clone())
|
||||
.or_insert_with(Vec::new)
|
||||
.push(key);
|
||||
.push(summary.id);
|
||||
}
|
||||
for keys in keys_by_provider.values_mut() {
|
||||
keys.sort_by(|left, right| left.id.cmp(&right.id));
|
||||
for key_ids in key_ids_by_provider.values_mut() {
|
||||
key_ids.sort();
|
||||
}
|
||||
|
||||
let now = now_unix_secs();
|
||||
@@ -261,25 +264,39 @@ pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
}
|
||||
last_provider_index = Some(provider_index);
|
||||
let (provider, pool_config) = providers[provider_index].clone();
|
||||
let keys = keys_by_provider.remove(&provider.id).unwrap_or_default();
|
||||
let keys = keys
|
||||
.into_iter()
|
||||
.filter(|key| key.is_active)
|
||||
.collect::<Vec<_>>();
|
||||
if keys.is_empty() {
|
||||
let key_ids = key_ids_by_provider.remove(&provider.id).unwrap_or_default();
|
||||
if key_ids.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let total_keys = keys.len();
|
||||
let total_keys = key_ids.len();
|
||||
let provider_cursor_key = provider_offset_cursor_key(&provider.id);
|
||||
let provider_cursor = load_runtime_usize(state, &provider_cursor_key).await % total_keys;
|
||||
let remaining_budget = config
|
||||
.max_upserts_per_tick
|
||||
.saturating_sub(summary.scores_upserted);
|
||||
let provider_budget = remaining_budget.min(total_keys);
|
||||
let mut build_items = Vec::with_capacity(provider_budget);
|
||||
let selected_ids = (0..provider_budget)
|
||||
.map(|offset| {
|
||||
let key_index = (provider_cursor + offset) % total_keys;
|
||||
key_ids[key_index].clone()
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
let mut selected_keys_by_id = state
|
||||
.list_provider_catalog_keys_by_ids(&selected_ids)
|
||||
.await?
|
||||
.into_iter()
|
||||
.filter(|key| key.is_active && key.provider_id == provider.id)
|
||||
.map(|key| (key.id.clone(), key))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
let selected_keys = selected_ids
|
||||
.iter()
|
||||
.filter_map(|key_id| selected_keys_by_id.remove(key_id))
|
||||
.collect::<Vec<_>>();
|
||||
let mut build_items = Vec::with_capacity(selected_keys.len());
|
||||
for offset in 0..provider_budget {
|
||||
let key_index = (provider_cursor + offset) % total_keys;
|
||||
let key = &keys[key_index];
|
||||
let Some(key) = selected_keys.get(offset) else {
|
||||
break;
|
||||
};
|
||||
let draft = build_provider_key_pool_score_upsert(
|
||||
key,
|
||||
provider.provider_type.as_str(),
|
||||
@@ -287,7 +304,7 @@ pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
now,
|
||||
pool_config.score_rules,
|
||||
);
|
||||
build_items.push((key_index, draft.id));
|
||||
build_items.push((offset, draft.id));
|
||||
}
|
||||
if build_items.is_empty() {
|
||||
store_runtime_usize(
|
||||
@@ -319,12 +336,12 @@ pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
.map(|score| (score.id.clone(), score))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
let mut provider_upserts = 0usize;
|
||||
summary.keys_seen = summary.keys_seen.saturating_add(keys.len());
|
||||
summary.keys_seen = summary.keys_seen.saturating_add(total_keys);
|
||||
for (key_index, score_id) in &build_items {
|
||||
if summary.scores_upserted >= config.max_upserts_per_tick {
|
||||
break;
|
||||
}
|
||||
let key = &keys[*key_index];
|
||||
let key = &selected_keys[*key_index];
|
||||
let existing = existing_scores.get(score_id);
|
||||
let upsert = build_provider_key_pool_score_upsert(
|
||||
key,
|
||||
@@ -429,3 +446,207 @@ pub(crate) fn spawn_pool_score_rebuild_worker(
|
||||
}
|
||||
}))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use aether_data::repository::pool_scores::InMemoryPoolMemberScoreRepository;
|
||||
use aether_data::repository::provider_catalog::InMemoryProviderCatalogReadRepository;
|
||||
use aether_data_contracts::repository::pool_scores::PoolMemberIdentity;
|
||||
use aether_data_contracts::repository::provider_catalog::{
|
||||
ProviderCatalogKeyListQuery, ProviderCatalogReadRepository,
|
||||
StoredProviderCatalogKeyMaintenanceSummary, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogKeyStats,
|
||||
};
|
||||
use aether_data_contracts::DataLayerError;
|
||||
use serde_json::json;
|
||||
|
||||
use crate::data::GatewayDataState;
|
||||
|
||||
struct NoWideKeyProviderCatalogReadRepository {
|
||||
inner: InMemoryProviderCatalogReadRepository,
|
||||
}
|
||||
|
||||
impl NoWideKeyProviderCatalogReadRepository {
|
||||
fn seed(
|
||||
providers: Vec<StoredProviderCatalogProvider>,
|
||||
endpoints: Vec<StoredProviderCatalogEndpoint>,
|
||||
keys: Vec<StoredProviderCatalogKey>,
|
||||
) -> Self {
|
||||
Self {
|
||||
inner: InMemoryProviderCatalogReadRepository::seed(providers, endpoints, keys),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl ProviderCatalogReadRepository for NoWideKeyProviderCatalogReadRepository {
|
||||
async fn list_providers(
|
||||
&self,
|
||||
active_only: bool,
|
||||
) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
|
||||
self.inner.list_providers(active_only).await
|
||||
}
|
||||
|
||||
async fn list_providers_by_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
|
||||
self.inner.list_providers_by_ids(provider_ids).await
|
||||
}
|
||||
|
||||
async fn list_endpoints_by_ids(
|
||||
&self,
|
||||
endpoint_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
|
||||
self.inner.list_endpoints_by_ids(endpoint_ids).await
|
||||
}
|
||||
|
||||
async fn list_endpoints_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
|
||||
self.inner
|
||||
.list_endpoints_by_provider_ids(provider_ids)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn list_keys_by_ids(
|
||||
&self,
|
||||
key_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
|
||||
self.inner.list_keys_by_ids(key_ids).await
|
||||
}
|
||||
|
||||
async fn list_keys_by_provider_ids(
|
||||
&self,
|
||||
_provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
|
||||
panic!("pool score rebuild should not read full provider key lists");
|
||||
}
|
||||
|
||||
async fn list_key_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
|
||||
self.inner
|
||||
.list_key_summaries_by_provider_ids(provider_ids)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
self.inner
|
||||
.list_key_maintenance_summaries_by_provider_ids(provider_ids)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
) -> Result<StoredProviderCatalogKeyPage, DataLayerError> {
|
||||
self.inner.list_keys_page(query).await
|
||||
}
|
||||
|
||||
async fn list_key_stats_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyStats>, DataLayerError> {
|
||||
self.inner
|
||||
.list_key_stats_by_provider_ids(provider_ids)
|
||||
.await
|
||||
}
|
||||
}
|
||||
|
||||
fn provider(id: &str) -> StoredProviderCatalogProvider {
|
||||
let mut provider = StoredProviderCatalogProvider::new(
|
||||
id.to_string(),
|
||||
id.to_string(),
|
||||
None,
|
||||
"openai".to_string(),
|
||||
)
|
||||
.expect("provider should build");
|
||||
provider.config = Some(json!({ "pool_advanced": {} }));
|
||||
provider
|
||||
}
|
||||
|
||||
fn key(id: &str, active: bool) -> StoredProviderCatalogKey {
|
||||
StoredProviderCatalogKey::new(
|
||||
id.to_string(),
|
||||
"provider-1".to_string(),
|
||||
id.to_string(),
|
||||
"oauth".to_string(),
|
||||
None,
|
||||
active,
|
||||
)
|
||||
.expect("key should build")
|
||||
}
|
||||
|
||||
fn score_id(provider_id: &str, key_id: &str) -> String {
|
||||
let identity =
|
||||
PoolMemberIdentity::provider_api_key(provider_id.to_string(), key_id.to_string());
|
||||
crate::ai_serving::provider_key_pool_score_id(
|
||||
&identity,
|
||||
&crate::ai_serving::provider_key_pool_score_scope(),
|
||||
)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn pool_score_rebuild_uses_maintenance_summaries_before_full_key_load() {
|
||||
let provider = provider("provider-1");
|
||||
let provider_catalog_repository = Arc::new(NoWideKeyProviderCatalogReadRepository::seed(
|
||||
vec![provider],
|
||||
Vec::new(),
|
||||
vec![
|
||||
key("key-a", true),
|
||||
key("key-b", true),
|
||||
key("key-disabled", false),
|
||||
],
|
||||
));
|
||||
let pool_score_repository = Arc::new(InMemoryPoolMemberScoreRepository::default());
|
||||
let data =
|
||||
GatewayDataState::with_provider_catalog_reader_for_tests(provider_catalog_repository)
|
||||
.with_pool_score_repository_for_tests(Arc::clone(&pool_score_repository));
|
||||
let state = AppState::new()
|
||||
.expect("gateway should build")
|
||||
.with_data_state_for_tests(data);
|
||||
|
||||
let summary = perform_pool_score_rebuild_once_with_config(
|
||||
&state,
|
||||
PoolScoreRebuildWorkerConfig {
|
||||
interval: Duration::from_secs(60),
|
||||
max_upserts_per_tick: 1,
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("rebuild should succeed");
|
||||
|
||||
assert_eq!(
|
||||
summary,
|
||||
PoolScoreRebuildRunSummary {
|
||||
providers_checked: 1,
|
||||
providers_scored: 1,
|
||||
keys_seen: 2,
|
||||
scores_upserted: 1,
|
||||
}
|
||||
);
|
||||
|
||||
let scores = state
|
||||
.data
|
||||
.get_pool_member_scores_by_ids(&GetPoolMemberScoresByIdsQuery {
|
||||
ids: vec![
|
||||
score_id("provider-1", "key-a"),
|
||||
score_id("provider-1", "key-b"),
|
||||
],
|
||||
})
|
||||
.await
|
||||
.expect("scores should load");
|
||||
assert_eq!(scores.len(), 1);
|
||||
assert_eq!(scores[0].member_id, "key-a");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -429,6 +429,17 @@ impl AppState {
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))
|
||||
}
|
||||
|
||||
pub(crate) async fn list_provider_catalog_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<provider_catalog::StoredProviderCatalogKeyMaintenanceSummary>, GatewayError>
|
||||
{
|
||||
self.data
|
||||
.list_provider_catalog_key_maintenance_summaries_by_provider_ids(provider_ids)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))
|
||||
}
|
||||
|
||||
pub(crate) async fn list_provider_catalog_keys_by_ids(
|
||||
&self,
|
||||
key_ids: &[String],
|
||||
|
||||
@@ -7,8 +7,8 @@ use aether_crypto::{
|
||||
use aether_data::repository::provider_catalog::InMemoryProviderCatalogReadRepository;
|
||||
use aether_data_contracts::repository::provider_catalog::{
|
||||
ProviderCatalogKeyListQuery, ProviderCatalogReadRepository, StoredProviderCatalogEndpoint,
|
||||
StoredProviderCatalogKey, StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats,
|
||||
StoredProviderCatalogProvider,
|
||||
StoredProviderCatalogKey, StoredProviderCatalogKeyMaintenanceSummary,
|
||||
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
};
|
||||
use aether_data_contracts::DataLayerError;
|
||||
use axum::body::Body;
|
||||
@@ -107,6 +107,15 @@ impl ProviderCatalogReadRepository for SummaryNullingProviderCatalogReadReposito
|
||||
Ok(keys)
|
||||
}
|
||||
|
||||
async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
self.inner
|
||||
.list_key_maintenance_summaries_by_provider_ids(provider_ids)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
|
||||
@@ -10,7 +10,15 @@ use crate::tests::{
|
||||
to_bytes, AppState, Arc, Body, Json, Mutex, Request, Router, StatusCode, EXECUTION_PATH_HEADER,
|
||||
EXECUTION_PATH_LOCAL_AI_PUBLIC, EXECUTION_PATH_LOCAL_EXECUTION_RUNTIME_MISS,
|
||||
};
|
||||
use aether_data::DataLayerError;
|
||||
use aether_data_contracts::repository::candidate_selection::{
|
||||
MinimalCandidateSelectionReadRepository, StoredMinimalCandidateSelectionRow,
|
||||
StoredPoolKeyCandidateRowsByKeyIdsQuery, StoredPoolKeyCandidateRowsQuery,
|
||||
StoredRequestedModelCandidateRowsQuery,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use axum::response::IntoResponse;
|
||||
use std::future::pending;
|
||||
|
||||
fn gemini_operation_status_label(status: VideoTaskStatus) -> &'static str {
|
||||
match status {
|
||||
@@ -118,6 +126,63 @@ fn sample_gemini_video_task(
|
||||
}
|
||||
}
|
||||
|
||||
struct PendingMinimalCandidateSelectionReadRepository;
|
||||
|
||||
impl PendingMinimalCandidateSelectionReadRepository {
|
||||
async fn pending_rows(
|
||||
&self,
|
||||
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
|
||||
pending::<Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError>>().await
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl MinimalCandidateSelectionReadRepository for PendingMinimalCandidateSelectionReadRepository {
|
||||
async fn list_for_exact_api_format(
|
||||
&self,
|
||||
_api_format: &str,
|
||||
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
|
||||
self.pending_rows().await
|
||||
}
|
||||
|
||||
async fn list_for_exact_api_format_and_global_model(
|
||||
&self,
|
||||
_api_format: &str,
|
||||
_global_model_name: &str,
|
||||
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
|
||||
self.pending_rows().await
|
||||
}
|
||||
|
||||
async fn list_for_exact_api_format_and_requested_model(
|
||||
&self,
|
||||
_api_format: &str,
|
||||
_requested_model_name: &str,
|
||||
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
|
||||
self.pending_rows().await
|
||||
}
|
||||
|
||||
async fn list_for_exact_api_format_and_requested_model_page(
|
||||
&self,
|
||||
_query: &StoredRequestedModelCandidateRowsQuery,
|
||||
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
|
||||
self.pending_rows().await
|
||||
}
|
||||
|
||||
async fn list_pool_key_rows_for_group(
|
||||
&self,
|
||||
_query: &StoredPoolKeyCandidateRowsQuery,
|
||||
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
|
||||
self.pending_rows().await
|
||||
}
|
||||
|
||||
async fn list_pool_key_rows_for_group_key_ids(
|
||||
&self,
|
||||
_query: &StoredPoolKeyCandidateRowsByKeyIdsQuery,
|
||||
) -> Result<Vec<StoredMinimalCandidateSelectionRow>, DataLayerError> {
|
||||
self.pending_rows().await
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_handles_public_openai_models_without_hitting_fallback_probe() {
|
||||
let fallback_probe_hits = Arc::new(Mutex::new(0usize));
|
||||
@@ -175,6 +240,119 @@ async fn gateway_handles_public_openai_models_without_hitting_fallback_probe() {
|
||||
fallback_probe_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_returns_empty_openai_models_when_candidate_rows_stall() {
|
||||
let fallback_probe_hits = Arc::new(Mutex::new(0usize));
|
||||
let fallback_probe_hits_clone = Arc::clone(&fallback_probe_hits);
|
||||
let fallback_probe = Router::new().route(
|
||||
"/{*path}",
|
||||
any(move |_request: Request| {
|
||||
let fallback_probe_hits_inner = Arc::clone(&fallback_probe_hits_clone);
|
||||
async move {
|
||||
*fallback_probe_hits_inner.lock().expect("mutex should lock") += 1;
|
||||
(StatusCode::OK, Body::from("proxied"))
|
||||
}
|
||||
}),
|
||||
);
|
||||
|
||||
let auth_repository = Arc::new(InMemoryAuthApiKeySnapshotRepository::seed(vec![(
|
||||
Some(hash_api_key("sk-openai-models-stalled")),
|
||||
unrestricted_models_snapshot("key-stalled", "user-stalled"),
|
||||
)]));
|
||||
let candidate_repository = Arc::new(PendingMinimalCandidateSelectionReadRepository);
|
||||
|
||||
let (_unused_fallback_probe_url, fallback_probe_handle) = start_server(fallback_probe).await;
|
||||
let gateway = build_router_with_state(
|
||||
AppState::new()
|
||||
.expect("gateway should build")
|
||||
.with_data_state_for_tests(
|
||||
crate::data::GatewayDataState::with_minimal_candidate_selection_and_auth_for_tests(
|
||||
candidate_repository,
|
||||
auth_repository,
|
||||
),
|
||||
),
|
||||
);
|
||||
let (gateway_url, gateway_handle) = start_server(gateway).await;
|
||||
|
||||
let response = reqwest::Client::builder()
|
||||
.timeout(std::time::Duration::from_millis(500))
|
||||
.build()
|
||||
.expect("client should build")
|
||||
.get(format!("{gateway_url}/v1/models"))
|
||||
.header("authorization", "Bearer sk-openai-models-stalled")
|
||||
.send()
|
||||
.await
|
||||
.expect("request should return before client timeout");
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let payload: serde_json::Value = response.json().await.expect("json body should parse");
|
||||
assert_eq!(payload["object"], "list");
|
||||
assert_eq!(
|
||||
payload["data"]
|
||||
.as_array()
|
||||
.expect("data should be an array")
|
||||
.len(),
|
||||
0
|
||||
);
|
||||
assert_eq!(*fallback_probe_hits.lock().expect("mutex should lock"), 0);
|
||||
|
||||
gateway_handle.abort();
|
||||
fallback_probe_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_returns_not_found_for_openai_model_detail_when_candidate_rows_stall() {
|
||||
let fallback_probe_hits = Arc::new(Mutex::new(0usize));
|
||||
let fallback_probe_hits_clone = Arc::clone(&fallback_probe_hits);
|
||||
let fallback_probe = Router::new().route(
|
||||
"/{*path}",
|
||||
any(move |_request: Request| {
|
||||
let fallback_probe_hits_inner = Arc::clone(&fallback_probe_hits_clone);
|
||||
async move {
|
||||
*fallback_probe_hits_inner.lock().expect("mutex should lock") += 1;
|
||||
(StatusCode::OK, Body::from("proxied"))
|
||||
}
|
||||
}),
|
||||
);
|
||||
|
||||
let auth_repository = Arc::new(InMemoryAuthApiKeySnapshotRepository::seed(vec![(
|
||||
Some(hash_api_key("sk-openai-model-detail-stalled")),
|
||||
unrestricted_models_snapshot("key-detail-stalled", "user-detail-stalled"),
|
||||
)]));
|
||||
let candidate_repository = Arc::new(PendingMinimalCandidateSelectionReadRepository);
|
||||
|
||||
let (_unused_fallback_probe_url, fallback_probe_handle) = start_server(fallback_probe).await;
|
||||
let gateway = build_router_with_state(
|
||||
AppState::new()
|
||||
.expect("gateway should build")
|
||||
.with_data_state_for_tests(
|
||||
crate::data::GatewayDataState::with_minimal_candidate_selection_and_auth_for_tests(
|
||||
candidate_repository,
|
||||
auth_repository,
|
||||
),
|
||||
),
|
||||
);
|
||||
let (gateway_url, gateway_handle) = start_server(gateway).await;
|
||||
|
||||
let response = reqwest::Client::builder()
|
||||
.timeout(std::time::Duration::from_millis(500))
|
||||
.build()
|
||||
.expect("client should build")
|
||||
.get(format!("{gateway_url}/v1/models/gpt-stalled"))
|
||||
.header("authorization", "Bearer sk-openai-model-detail-stalled")
|
||||
.send()
|
||||
.await
|
||||
.expect("request should return before client timeout");
|
||||
|
||||
assert_eq!(response.status(), StatusCode::NOT_FOUND);
|
||||
let payload: serde_json::Value = response.json().await.expect("json body should parse");
|
||||
assert_eq!(payload["error"]["code"], "model_not_found");
|
||||
assert_eq!(*fallback_probe_hits.lock().expect("mutex should lock"), 0);
|
||||
|
||||
gateway_handle.abort();
|
||||
fallback_probe_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_handles_public_openai_models_with_cross_format_candidates_without_hitting_fallback_probe(
|
||||
) {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "aether-tunnel"
|
||||
version = "0.3.12"
|
||||
version = "0.3.13"
|
||||
edition = "2021"
|
||||
description = "Tunnel agent for Aether"
|
||||
|
||||
|
||||
@@ -15,13 +15,13 @@ Tunnel 模式下代理节点**无需对外监听端口**,仅需出站连接到
|
||||
<!-- DOWNLOAD_TABLE_START -->
|
||||
| Platform | Download |
|
||||
|----------|----------|
|
||||
| Linux x86_64 (GNU) | [aether-tunnel-linux-amd64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.12/aether-tunnel-linux-amd64.tar.gz) |
|
||||
| Linux ARM64 (GNU) | [aether-tunnel-linux-arm64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.12/aether-tunnel-linux-arm64.tar.gz) |
|
||||
| Linux x86_64 (musl) | [aether-tunnel-linux-musl-amd64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.12/aether-tunnel-linux-musl-amd64.tar.gz) |
|
||||
| Linux ARM64 (musl) | [aether-tunnel-linux-musl-arm64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.12/aether-tunnel-linux-musl-arm64.tar.gz) |
|
||||
| macOS x86_64 | [aether-tunnel-macos-amd64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.12/aether-tunnel-macos-amd64.tar.gz) |
|
||||
| macOS ARM64 | [aether-tunnel-macos-arm64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.12/aether-tunnel-macos-arm64.tar.gz) |
|
||||
| Windows x86_64 | [aether-tunnel-windows-amd64.zip](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.12/aether-tunnel-windows-amd64.zip) |
|
||||
| Linux x86_64 (GNU) | [aether-tunnel-linux-amd64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.13/aether-tunnel-linux-amd64.tar.gz) |
|
||||
| Linux ARM64 (GNU) | [aether-tunnel-linux-arm64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.13/aether-tunnel-linux-arm64.tar.gz) |
|
||||
| Linux x86_64 (musl) | [aether-tunnel-linux-musl-amd64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.13/aether-tunnel-linux-musl-amd64.tar.gz) |
|
||||
| Linux ARM64 (musl) | [aether-tunnel-linux-musl-arm64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.13/aether-tunnel-linux-musl-arm64.tar.gz) |
|
||||
| macOS x86_64 | [aether-tunnel-macos-amd64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.13/aether-tunnel-macos-amd64.tar.gz) |
|
||||
| macOS ARM64 | [aether-tunnel-macos-arm64.tar.gz](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.13/aether-tunnel-macos-arm64.tar.gz) |
|
||||
| Windows x86_64 | [aether-tunnel-windows-amd64.zip](https://github.com/fawney19/Aether/releases/download/tunnel-v0.3.13/aether-tunnel-windows-amd64.zip) |
|
||||
<!-- DOWNLOAD_TABLE_END -->
|
||||
|
||||
上表展示的是最新已发布版本的下载链接。从下一次 `tunnel-v*` 发布开始,表格会自动补上 `Linux x86_64 (musl)` / `Linux ARM64 (musl)` 包,供 Alpine 等 musl 系统直接使用。
|
||||
|
||||
@@ -3,5 +3,6 @@ mod types;
|
||||
pub use types::{
|
||||
ProviderCatalogKeyListOrder, ProviderCatalogKeyListQuery, ProviderCatalogReadRepository,
|
||||
ProviderCatalogWriteRepository, StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
|
||||
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
StoredProviderCatalogKeyMaintenanceSummary, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
};
|
||||
|
||||
@@ -296,6 +296,14 @@ pub struct StoredProviderCatalogKey {
|
||||
pub circuit_breaker_by_format: Option<serde_json::Value>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct StoredProviderCatalogKeyMaintenanceSummary {
|
||||
pub id: String,
|
||||
pub provider_id: String,
|
||||
pub is_active: bool,
|
||||
pub upstream_metadata: Option<serde_json::Value>,
|
||||
}
|
||||
|
||||
impl StoredProviderCatalogKey {
|
||||
pub fn new(
|
||||
id: String,
|
||||
@@ -607,6 +615,11 @@ pub trait ProviderCatalogReadRepository: Send + Sync {
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKey>, crate::DataLayerError>;
|
||||
|
||||
async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, crate::DataLayerError>;
|
||||
|
||||
async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
SET @aether_provider_key_active_priority_index_sql := IF(
|
||||
(
|
||||
SELECT COUNT(*)
|
||||
FROM information_schema.statistics
|
||||
WHERE table_schema = DATABASE()
|
||||
AND table_name = 'provider_api_keys'
|
||||
AND index_name = 'idx_provider_api_keys_provider_active_priority_id'
|
||||
) = 0,
|
||||
'CREATE INDEX idx_provider_api_keys_provider_active_priority_id ON provider_api_keys (provider_id, is_active, internal_priority, id)',
|
||||
'DO 0'
|
||||
);
|
||||
|
||||
PREPARE aether_provider_key_active_priority_index_stmt FROM @aether_provider_key_active_priority_index_sql;
|
||||
EXECUTE aether_provider_key_active_priority_index_stmt;
|
||||
DEALLOCATE PREPARE aether_provider_key_active_priority_index_stmt;
|
||||
|
||||
SET @aether_pool_score_scheduler_rank_index_sql := IF(
|
||||
(
|
||||
SELECT COUNT(*)
|
||||
FROM information_schema.statistics
|
||||
WHERE table_schema = DATABASE()
|
||||
AND table_name = 'pool_member_scores'
|
||||
AND index_name = 'pool_member_scores_scheduler_account_rank_idx'
|
||||
) = 0,
|
||||
'CREATE INDEX pool_member_scores_scheduler_account_rank_idx ON pool_member_scores (pool_kind, pool_id, capability, scope_kind, scope_id, hard_state, score DESC, last_ranked_at DESC, member_id, id)',
|
||||
'DO 0'
|
||||
);
|
||||
|
||||
PREPARE aether_pool_score_scheduler_rank_index_stmt FROM @aether_pool_score_scheduler_rank_index_sql;
|
||||
EXECUTE aether_pool_score_scheduler_rank_index_stmt;
|
||||
DEALLOCATE PREPARE aether_pool_score_scheduler_rank_index_stmt;
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
CREATE INDEX IF NOT EXISTS idx_provider_api_keys_provider_active_priority_id
|
||||
ON public.provider_api_keys USING btree (provider_id, internal_priority, id)
|
||||
WHERE is_active IS TRUE;
|
||||
|
||||
CREATE INDEX IF NOT EXISTS pool_member_scores_scheduler_account_rank_idx
|
||||
ON public.pool_member_scores USING btree (
|
||||
pool_kind,
|
||||
pool_id,
|
||||
capability,
|
||||
scope_kind,
|
||||
score DESC,
|
||||
last_ranked_at DESC NULLS LAST,
|
||||
member_id,
|
||||
id
|
||||
)
|
||||
WHERE scope_id IS NULL AND hard_state IN ('available', 'unknown');
|
||||
@@ -0,0 +1,16 @@
|
||||
CREATE INDEX IF NOT EXISTS idx_provider_api_keys_provider_active_priority_id
|
||||
ON provider_api_keys (provider_id, is_active, internal_priority, id);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS pool_member_scores_scheduler_account_rank_idx
|
||||
ON pool_member_scores (
|
||||
pool_kind,
|
||||
pool_id,
|
||||
capability,
|
||||
scope_kind,
|
||||
scope_id,
|
||||
hard_state,
|
||||
score DESC,
|
||||
last_ranked_at DESC,
|
||||
member_id,
|
||||
id
|
||||
);
|
||||
@@ -197,6 +197,14 @@ CREATE INDEX IF NOT EXISTS idx_provider_api_keys_provider_default_sort ON public
|
||||
|
||||
|
||||
|
||||
--
|
||||
-- Name: idx_provider_api_keys_provider_active_priority_id; Type: INDEX; Schema: public; Owner: -
|
||||
--
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_provider_api_keys_provider_active_priority_id ON public.provider_api_keys USING btree (provider_id, internal_priority, id) WHERE is_active IS TRUE;
|
||||
|
||||
|
||||
|
||||
--
|
||||
-- Name: idx_provider_api_keys_provider_created_at_desc; Type: INDEX; Schema: public; Owner: -
|
||||
--
|
||||
@@ -1037,6 +1045,14 @@ CREATE INDEX IF NOT EXISTS pool_member_scores_rank_idx ON public.pool_member_sco
|
||||
|
||||
|
||||
|
||||
--
|
||||
-- Name: pool_member_scores_scheduler_account_rank_idx; Type: INDEX; Schema: public; Owner: -
|
||||
--
|
||||
|
||||
CREATE INDEX IF NOT EXISTS pool_member_scores_scheduler_account_rank_idx ON public.pool_member_scores USING btree (pool_kind, pool_id, capability, scope_kind, score DESC, last_ranked_at DESC NULLS LAST, member_id, id) WHERE scope_id IS NULL AND hard_state IN ('available', 'unknown');
|
||||
|
||||
|
||||
|
||||
--
|
||||
-- Name: pool_member_scores_member_idx; Type: INDEX; Schema: public; Owner: -
|
||||
--
|
||||
|
||||
@@ -1,8 +1,14 @@
|
||||
pub(super) const SELECT_STATS_DAILY_AGGREGATE_SQL: &str = r#"
|
||||
WITH bucket_usage AS MATERIALIZED (
|
||||
SELECT *
|
||||
FROM usage_billing_facts
|
||||
WHERE created_at >= $1
|
||||
AND created_at < $2
|
||||
)
|
||||
SELECT
|
||||
(
|
||||
SELECT CAST(COUNT(cache_hit_usage.id) AS BIGINT)
|
||||
FROM usage_billing_facts AS cache_hit_usage
|
||||
FROM bucket_usage AS cache_hit_usage
|
||||
WHERE cache_hit_usage.created_at >= $1
|
||||
AND cache_hit_usage.created_at < $2
|
||||
) AS cache_hit_total_requests,
|
||||
@@ -12,13 +18,13 @@ SELECT
|
||||
WHERE GREATEST(COALESCE(cache_hit_usage.cache_read_input_tokens, 0), 0) > 0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS cache_hit_usage
|
||||
FROM bucket_usage AS cache_hit_usage
|
||||
WHERE cache_hit_usage.created_at >= $1
|
||||
AND cache_hit_usage.created_at < $2
|
||||
) AS cache_hit_requests,
|
||||
(
|
||||
SELECT CAST(COUNT(completed_usage.id) AS BIGINT)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -29,7 +35,7 @@ SELECT
|
||||
WHERE GREATEST(COALESCE(completed_usage.cache_read_input_tokens, 0), 0) > 0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -38,7 +44,7 @@ SELECT
|
||||
SELECT CAST(
|
||||
COALESCE(SUM(GREATEST(COALESCE(completed_usage.input_tokens, 0), 0)), 0) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -61,7 +67,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -73,7 +79,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -211,7 +217,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -228,7 +234,7 @@ SELECT
|
||||
0
|
||||
) AS DOUBLE PRECISION
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -242,7 +248,7 @@ SELECT
|
||||
0
|
||||
) AS DOUBLE PRECISION
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -252,7 +258,7 @@ SELECT
|
||||
COALESCE(SUM(COALESCE(CAST(settled_usage.total_cost_usd AS DOUBLE PRECISION), 0)), 0)
|
||||
AS DOUBLE PRECISION
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -260,7 +266,7 @@ SELECT
|
||||
) AS settled_total_cost,
|
||||
(
|
||||
SELECT CAST(COUNT(settled_usage.id) AS BIGINT)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -270,7 +276,7 @@ SELECT
|
||||
SELECT CAST(
|
||||
COALESCE(SUM(GREATEST(COALESCE(settled_usage.input_tokens, 0), 0)), 0) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -280,7 +286,7 @@ SELECT
|
||||
SELECT CAST(
|
||||
COALESCE(SUM(GREATEST(COALESCE(settled_usage.output_tokens, 0), 0)), 0) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -293,7 +299,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -306,7 +312,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -314,7 +320,7 @@ SELECT
|
||||
) AS settled_cache_read_tokens,
|
||||
(
|
||||
SELECT MIN(CAST(EXTRACT(EPOCH FROM settled_usage.finalized_at) AS BIGINT))
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -322,7 +328,7 @@ SELECT
|
||||
) AS settled_first_finalized_at_unix_secs,
|
||||
(
|
||||
SELECT MAX(CAST(EXTRACT(EPOCH FROM settled_usage.finalized_at) AS BIGINT))
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
|
||||
@@ -12,10 +12,16 @@ WHERE created_at >= $1
|
||||
AND provider_name NOT IN ('unknown', 'pending')
|
||||
"#;
|
||||
pub(super) const SELECT_STATS_HOURLY_AGGREGATE_SQL: &str = r#"
|
||||
WITH bucket_usage AS MATERIALIZED (
|
||||
SELECT *
|
||||
FROM usage_billing_facts
|
||||
WHERE created_at >= $1
|
||||
AND created_at < $2
|
||||
)
|
||||
SELECT
|
||||
(
|
||||
SELECT CAST(COUNT(cache_hit_usage.id) AS BIGINT)
|
||||
FROM usage_billing_facts AS cache_hit_usage
|
||||
FROM bucket_usage AS cache_hit_usage
|
||||
WHERE cache_hit_usage.created_at >= $1
|
||||
AND cache_hit_usage.created_at < $2
|
||||
) AS cache_hit_total_requests,
|
||||
@@ -25,13 +31,13 @@ SELECT
|
||||
WHERE GREATEST(COALESCE(cache_hit_usage.cache_read_input_tokens, 0), 0) > 0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS cache_hit_usage
|
||||
FROM bucket_usage AS cache_hit_usage
|
||||
WHERE cache_hit_usage.created_at >= $1
|
||||
AND cache_hit_usage.created_at < $2
|
||||
) AS cache_hit_requests,
|
||||
(
|
||||
SELECT CAST(COUNT(completed_usage.id) AS BIGINT)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -42,7 +48,7 @@ SELECT
|
||||
WHERE GREATEST(COALESCE(completed_usage.cache_read_input_tokens, 0), 0) > 0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -51,7 +57,7 @@ SELECT
|
||||
SELECT CAST(
|
||||
COALESCE(SUM(GREATEST(COALESCE(completed_usage.input_tokens, 0), 0)), 0) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -74,7 +80,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -86,7 +92,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -224,7 +230,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -241,7 +247,7 @@ SELECT
|
||||
0
|
||||
) AS DOUBLE PRECISION
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -255,7 +261,7 @@ SELECT
|
||||
0
|
||||
) AS DOUBLE PRECISION
|
||||
)
|
||||
FROM usage_billing_facts AS completed_usage
|
||||
FROM bucket_usage AS completed_usage
|
||||
WHERE completed_usage.created_at >= $1
|
||||
AND completed_usage.created_at < $2
|
||||
AND completed_usage.status = 'completed'
|
||||
@@ -265,7 +271,7 @@ SELECT
|
||||
COALESCE(SUM(COALESCE(CAST(settled_usage.total_cost_usd AS DOUBLE PRECISION), 0)), 0)
|
||||
AS DOUBLE PRECISION
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -273,7 +279,7 @@ SELECT
|
||||
) AS settled_total_cost,
|
||||
(
|
||||
SELECT CAST(COUNT(settled_usage.id) AS BIGINT)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -283,7 +289,7 @@ SELECT
|
||||
SELECT CAST(
|
||||
COALESCE(SUM(GREATEST(COALESCE(settled_usage.input_tokens, 0), 0)), 0) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -293,7 +299,7 @@ SELECT
|
||||
SELECT CAST(
|
||||
COALESCE(SUM(GREATEST(COALESCE(settled_usage.output_tokens, 0), 0)), 0) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -306,7 +312,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -319,7 +325,7 @@ SELECT
|
||||
0
|
||||
) AS BIGINT
|
||||
)
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -327,7 +333,7 @@ SELECT
|
||||
) AS settled_cache_read_tokens,
|
||||
(
|
||||
SELECT MIN(CAST(EXTRACT(EPOCH FROM settled_usage.finalized_at) AS BIGINT))
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
@@ -335,7 +341,7 @@ SELECT
|
||||
) AS settled_first_finalized_at_unix_secs,
|
||||
(
|
||||
SELECT MAX(CAST(EXTRACT(EPOCH FROM settled_usage.finalized_at) AS BIGINT))
|
||||
FROM usage_billing_facts AS settled_usage
|
||||
FROM bucket_usage AS settled_usage
|
||||
WHERE settled_usage.created_at >= $1
|
||||
AND settled_usage.created_at < $2
|
||||
AND settled_usage.billing_status = 'settled'
|
||||
|
||||
@@ -7,7 +7,7 @@ use tracing::info;
|
||||
// Generated by build.rs from schema/bootstrap/postgres.
|
||||
pub(crate) static EMPTY_DATABASE_SNAPSHOT_SQL: &str =
|
||||
include_str!(concat!(env!("OUT_DIR"), "/empty_database_snapshot.sql"));
|
||||
pub(crate) const EMPTY_DATABASE_SNAPSHOT_CUTOFF_VERSION: i64 = 20260522000000;
|
||||
pub(crate) const EMPTY_DATABASE_SNAPSHOT_CUTOFF_VERSION: i64 = 20260524000000;
|
||||
|
||||
const PUBLIC_BASE_TABLE_COUNT_SQL: &str = r#"
|
||||
SELECT COUNT(*)::BIGINT
|
||||
|
||||
@@ -313,6 +313,7 @@ fn empty_database_snapshot_covers_current_cutoff_versions() {
|
||||
20260520000000,
|
||||
20260520010000,
|
||||
20260522000000,
|
||||
20260524000000,
|
||||
]
|
||||
);
|
||||
}
|
||||
@@ -389,6 +390,10 @@ fn empty_database_snapshot_sql_includes_usage_body_blobs_and_audit_admin_role()
|
||||
assert!(EMPTY_DATABASE_SNAPSHOT_SQL.contains("ix_usage_counter_deltas_unprocessed"));
|
||||
assert!(EMPTY_DATABASE_SNAPSHOT_SQL.contains("idx_entitlement_usage_entitlement_date"));
|
||||
assert!(EMPTY_DATABASE_SNAPSHOT_SQL.contains("idx_provider_api_keys_provider_default_sort"));
|
||||
assert!(
|
||||
EMPTY_DATABASE_SNAPSHOT_SQL.contains("idx_provider_api_keys_provider_active_priority_id")
|
||||
);
|
||||
assert!(EMPTY_DATABASE_SNAPSHOT_SQL.contains("pool_member_scores_scheduler_account_rank_idx"));
|
||||
assert!(EMPTY_DATABASE_SNAPSHOT_SQL.contains("idx_video_tasks_due_poll"));
|
||||
assert!(EMPTY_DATABASE_SNAPSHOT_SQL.contains("request_count bigint DEFAULT 0"));
|
||||
assert!(EMPTY_DATABASE_SNAPSHOT_SQL.contains("usage_count bigint DEFAULT 0 NOT NULL"));
|
||||
@@ -662,6 +667,7 @@ fn mysql_and_sqlite_migrations_include_enabled_incrementals() {
|
||||
20260519130000,
|
||||
20260520000000,
|
||||
20260520010000,
|
||||
20260524000000,
|
||||
]
|
||||
);
|
||||
assert_eq!(
|
||||
@@ -685,6 +691,7 @@ fn mysql_and_sqlite_migrations_include_enabled_incrementals() {
|
||||
20260519130000,
|
||||
20260520000000,
|
||||
20260520010000,
|
||||
20260524000000,
|
||||
]
|
||||
);
|
||||
}
|
||||
@@ -1209,6 +1216,7 @@ fn pending_migrations_from_applied_skips_versions_already_applied() {
|
||||
20260520000000,
|
||||
20260520010000,
|
||||
20260522000000,
|
||||
20260524000000,
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
@@ -54,15 +54,100 @@ SELECT
|
||||
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
|
||||
INNER JOIN LATERAL (
|
||||
SELECT pak.*
|
||||
FROM provider_api_keys pak
|
||||
WHERE pak.provider_id = p.id
|
||||
AND pak.is_active IS TRUE
|
||||
AND (
|
||||
pak.api_formats IS NULL
|
||||
OR EXISTS (
|
||||
SELECT 1
|
||||
FROM json_array_elements_text(pak.api_formats) AS fmt(value)
|
||||
WHERE LOWER(BTRIM(fmt.value)) = ANY($2::text[])
|
||||
)
|
||||
)
|
||||
AND (
|
||||
(
|
||||
LOWER(BTRIM(p.provider_type)) = 'codex'
|
||||
AND LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
AND LOWER($3) IN ('openai:responses', 'openai:responses:compact', 'openai:image')
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'chatgpt_web'
|
||||
AND LOWER(BTRIM(pak.auth_type)) IN ('oauth', 'bearer')
|
||||
AND LOWER($3) = 'openai:image'
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'claude_code'
|
||||
AND LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
AND LOWER($3) = 'claude:messages'
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'kiro'
|
||||
AND LOWER($3) = 'claude:messages'
|
||||
AND (
|
||||
LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
OR (
|
||||
LOWER(BTRIM(pak.auth_type)) = 'bearer'
|
||||
AND pak.auth_config IS NOT NULL
|
||||
AND BTRIM(pak.auth_config) <> ''
|
||||
)
|
||||
)
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'grok'
|
||||
AND LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
AND LOWER($3) IN ('openai:chat', 'openai:responses', 'claude:messages', 'openai:image')
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) IN ('gemini_cli', 'antigravity')
|
||||
AND LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
AND LOWER($3) = 'gemini:generate_content'
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'windsurf'
|
||||
AND LOWER(BTRIM(pak.auth_type)) IN ('oauth', 'api_key', 'bearer')
|
||||
AND LOWER($3) = 'openai:chat'
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'vertex_ai'
|
||||
AND (
|
||||
(
|
||||
LOWER(BTRIM(pak.auth_type)) = 'api_key'
|
||||
AND LOWER($3) IN ('gemini:generate_content', 'gemini:embedding')
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(pak.auth_type)) IN ('service_account', 'vertex_ai')
|
||||
AND LOWER($3) IN ('claude:messages', 'gemini:generate_content', 'gemini:embedding')
|
||||
)
|
||||
)
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) NOT IN (
|
||||
'chatgpt_web',
|
||||
'claude_code',
|
||||
'codex',
|
||||
'gemini_cli',
|
||||
'grok',
|
||||
'vertex_ai',
|
||||
'antigravity',
|
||||
'kiro',
|
||||
'windsurf'
|
||||
)
|
||||
AND LOWER(BTRIM(pak.auth_type)) <> 'oauth'
|
||||
)
|
||||
)
|
||||
ORDER BY pak.internal_priority ASC, pak.id ASC
|
||||
LIMIT CASE WHEN (p.config -> 'pool_advanced') IS NOT NULL THEN 1 ELSE 2147483647 END
|
||||
) pak ON TRUE
|
||||
INNER JOIN models m
|
||||
ON m.provider_id = p.id
|
||||
INNER JOIN global_models gm
|
||||
ON gm.id = m.global_model_id
|
||||
WHERE p.is_active = TRUE
|
||||
AND pe.is_active = TRUE
|
||||
AND pak.is_active = TRUE
|
||||
AND pak.is_active IS TRUE
|
||||
AND m.is_active = TRUE
|
||||
AND m.is_available = TRUE
|
||||
AND gm.is_active = TRUE
|
||||
@@ -248,15 +333,100 @@ SELECT
|
||||
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
|
||||
INNER JOIN LATERAL (
|
||||
SELECT pak.*
|
||||
FROM provider_api_keys pak
|
||||
WHERE pak.provider_id = p.id
|
||||
AND pak.is_active IS TRUE
|
||||
AND (
|
||||
pak.api_formats IS NULL
|
||||
OR EXISTS (
|
||||
SELECT 1
|
||||
FROM json_array_elements_text(pak.api_formats) AS fmt(value)
|
||||
WHERE LOWER(BTRIM(fmt.value)) = ANY($3::text[])
|
||||
)
|
||||
)
|
||||
AND (
|
||||
(
|
||||
LOWER(BTRIM(p.provider_type)) = 'codex'
|
||||
AND LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
AND LOWER($4) IN ('openai:responses', 'openai:responses:compact', 'openai:image')
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'chatgpt_web'
|
||||
AND LOWER(BTRIM(pak.auth_type)) IN ('oauth', 'bearer')
|
||||
AND LOWER($4) = 'openai:image'
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'claude_code'
|
||||
AND LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
AND LOWER($4) = 'claude:messages'
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'kiro'
|
||||
AND LOWER($4) = 'claude:messages'
|
||||
AND (
|
||||
LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
OR (
|
||||
LOWER(BTRIM(pak.auth_type)) = 'bearer'
|
||||
AND pak.auth_config IS NOT NULL
|
||||
AND BTRIM(pak.auth_config) <> ''
|
||||
)
|
||||
)
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'grok'
|
||||
AND LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
AND LOWER($4) IN ('openai:chat', 'openai:responses', 'claude:messages', 'openai:image')
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) IN ('gemini_cli', 'antigravity')
|
||||
AND LOWER(BTRIM(pak.auth_type)) = 'oauth'
|
||||
AND LOWER($4) = 'gemini:generate_content'
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'windsurf'
|
||||
AND LOWER(BTRIM(pak.auth_type)) IN ('oauth', 'api_key', 'bearer')
|
||||
AND LOWER($4) = 'openai:chat'
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) = 'vertex_ai'
|
||||
AND (
|
||||
(
|
||||
LOWER(BTRIM(pak.auth_type)) = 'api_key'
|
||||
AND LOWER($4) IN ('gemini:generate_content', 'gemini:embedding')
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(pak.auth_type)) IN ('service_account', 'vertex_ai')
|
||||
AND LOWER($4) IN ('claude:messages', 'gemini:generate_content', 'gemini:embedding')
|
||||
)
|
||||
)
|
||||
)
|
||||
OR (
|
||||
LOWER(BTRIM(p.provider_type)) NOT IN (
|
||||
'chatgpt_web',
|
||||
'claude_code',
|
||||
'codex',
|
||||
'gemini_cli',
|
||||
'grok',
|
||||
'vertex_ai',
|
||||
'antigravity',
|
||||
'kiro',
|
||||
'windsurf'
|
||||
)
|
||||
AND LOWER(BTRIM(pak.auth_type)) <> 'oauth'
|
||||
)
|
||||
)
|
||||
ORDER BY pak.internal_priority ASC, pak.id ASC
|
||||
LIMIT CASE WHEN (p.config -> 'pool_advanced') IS NOT NULL THEN 1 ELSE 2147483647 END
|
||||
) pak ON TRUE
|
||||
INNER JOIN models m
|
||||
ON m.provider_id = p.id
|
||||
INNER JOIN global_models gm
|
||||
ON gm.id = m.global_model_id
|
||||
WHERE p.is_active = TRUE
|
||||
AND pe.is_active = TRUE
|
||||
AND pak.is_active = TRUE
|
||||
AND pak.is_active IS TRUE
|
||||
AND m.is_active = TRUE
|
||||
AND m.is_available = TRUE
|
||||
AND gm.is_active = TRUE
|
||||
@@ -448,7 +618,7 @@ INNER JOIN global_models gm
|
||||
ON gm.id = m.global_model_id
|
||||
WHERE p.is_active = TRUE
|
||||
AND pe.is_active = TRUE
|
||||
AND pak.is_active = TRUE
|
||||
AND pak.is_active IS TRUE
|
||||
AND m.is_active = TRUE
|
||||
AND m.is_available = TRUE
|
||||
AND gm.is_active = TRUE
|
||||
@@ -1308,6 +1478,25 @@ mod tests {
|
||||
assert!(!sql.contains("AND gm.name = $2\n AND"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn exact_candidate_selection_sql_limits_pool_key_expansion() {
|
||||
for sql in [
|
||||
LIST_FOR_EXACT_API_FORMAT_SQL,
|
||||
LIST_FOR_EXACT_API_FORMAT_AND_GLOBAL_MODEL_SQL,
|
||||
] {
|
||||
assert!(sql.contains("INNER JOIN LATERAL"));
|
||||
assert!(sql.contains("FROM provider_api_keys pak"));
|
||||
assert!(sql.contains("AND pak.is_active IS TRUE"));
|
||||
assert!(!sql.contains("AND pak.is_active = TRUE"));
|
||||
assert!(sql.contains(
|
||||
"LIMIT CASE WHEN (p.config -> 'pool_advanced') IS NOT NULL THEN 1 ELSE 2147483647 END"
|
||||
));
|
||||
}
|
||||
assert!(!LIST_POOL_KEYS_FOR_GROUP_SQL.contains("INNER JOIN LATERAL"));
|
||||
assert!(LIST_POOL_KEYS_FOR_GROUP_SQL.contains("AND pak.is_active IS TRUE"));
|
||||
assert!(!LIST_POOL_KEYS_FOR_GROUP_SQL.contains("AND pak.is_active = TRUE"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn candidate_selection_sql_allows_chatgpt_web_image_auth() {
|
||||
let requested_model_sql = requested_model_selection_sql();
|
||||
|
||||
@@ -8,7 +8,8 @@ use serde_json::{json, Map, Value};
|
||||
use super::{
|
||||
ProviderCatalogKeyListOrder, ProviderCatalogKeyListQuery, ProviderCatalogReadRepository,
|
||||
ProviderCatalogWriteRepository, StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
|
||||
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
StoredProviderCatalogKeyMaintenanceSummary, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
};
|
||||
use crate::repository::usage::{ProviderApiKeyUsageContribution, ProviderApiKeyUsageDelta};
|
||||
use crate::DataLayerError;
|
||||
@@ -431,6 +432,34 @@ impl ProviderCatalogReadRepository for InMemoryProviderCatalogReadRepository {
|
||||
Self::list_keys_by_provider_ids(self, provider_ids).await
|
||||
}
|
||||
|
||||
async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
let index = self.index.read().expect("provider catalog repository lock");
|
||||
let mut keys = index
|
||||
.keys
|
||||
.values()
|
||||
.filter(|key| {
|
||||
provider_ids
|
||||
.iter()
|
||||
.any(|provider_id| provider_id == &key.provider_id)
|
||||
})
|
||||
.map(|key| StoredProviderCatalogKeyMaintenanceSummary {
|
||||
id: key.id.clone(),
|
||||
provider_id: key.provider_id.clone(),
|
||||
is_active: key.is_active,
|
||||
upstream_metadata: key.upstream_metadata.clone(),
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
keys.sort_by(|left, right| {
|
||||
left.provider_id
|
||||
.cmp(&right.provider_id)
|
||||
.then(left.id.cmp(&right.id))
|
||||
});
|
||||
Ok(keys)
|
||||
}
|
||||
|
||||
async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
|
||||
@@ -7,7 +7,8 @@ mod sqlite;
|
||||
pub(crate) use aether_data_contracts::repository::provider_catalog::{
|
||||
ProviderCatalogKeyListOrder, ProviderCatalogKeyListQuery, ProviderCatalogReadRepository,
|
||||
ProviderCatalogWriteRepository, StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
|
||||
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
StoredProviderCatalogKeyMaintenanceSummary, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
};
|
||||
pub use memory::InMemoryProviderCatalogReadRepository;
|
||||
pub use mysql::MysqlProviderCatalogReadRepository;
|
||||
|
||||
@@ -4,8 +4,8 @@ use sqlx::{mysql::MySqlRow, Row};
|
||||
use super::{
|
||||
InMemoryProviderCatalogReadRepository, ProviderCatalogKeyListQuery,
|
||||
ProviderCatalogReadRepository, ProviderCatalogWriteRepository, StoredProviderCatalogEndpoint,
|
||||
StoredProviderCatalogKey, StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats,
|
||||
StoredProviderCatalogProvider,
|
||||
StoredProviderCatalogKey, StoredProviderCatalogKeyMaintenanceSummary,
|
||||
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
};
|
||||
use crate::driver::mysql::MysqlPool;
|
||||
use crate::error::SqlResultExt;
|
||||
@@ -928,6 +928,16 @@ impl ProviderCatalogReadRepository for MysqlProviderCatalogReadRepository {
|
||||
.await
|
||||
}
|
||||
|
||||
async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
self.load_memory()
|
||||
.await?
|
||||
.list_key_maintenance_summaries_by_provider_ids(provider_ids)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
|
||||
@@ -5,7 +5,8 @@ use sqlx::{postgres::PgRow, PgPool, Postgres, QueryBuilder, Row};
|
||||
use super::{
|
||||
ProviderCatalogKeyListOrder, ProviderCatalogKeyListQuery, ProviderCatalogReadRepository,
|
||||
ProviderCatalogWriteRepository, StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
|
||||
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
StoredProviderCatalogKeyMaintenanceSummary, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
};
|
||||
use crate::{
|
||||
error::{postgres_error, SqlxResultExt},
|
||||
@@ -323,6 +324,16 @@ FROM provider_api_keys
|
||||
WHERE provider_id IN (
|
||||
"#;
|
||||
|
||||
const LIST_KEY_MAINTENANCE_SUMMARIES_BY_PROVIDER_IDS_PREFIX: &str = r#"
|
||||
SELECT
|
||||
id,
|
||||
provider_id,
|
||||
is_active,
|
||||
upstream_metadata
|
||||
FROM provider_api_keys
|
||||
WHERE provider_id IN (
|
||||
"#;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct SqlxProviderCatalogReadRepository {
|
||||
pool: PgPool,
|
||||
@@ -513,6 +524,27 @@ impl SqlxProviderCatalogReadRepository {
|
||||
.await
|
||||
}
|
||||
|
||||
pub async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
if provider_ids.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
collect_query_rows(
|
||||
build_list_query(
|
||||
LIST_KEY_MAINTENANCE_SUMMARIES_BY_PROVIDER_IDS_PREFIX,
|
||||
provider_ids,
|
||||
" ORDER BY provider_id ASC, id ASC",
|
||||
)
|
||||
.build()
|
||||
.fetch(&self.pool),
|
||||
map_key_maintenance_summary_row,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
@@ -1917,6 +1949,13 @@ impl ProviderCatalogReadRepository for SqlxProviderCatalogReadRepository {
|
||||
Self::list_key_summaries_by_provider_ids(self, provider_ids).await
|
||||
}
|
||||
|
||||
async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
Self::list_key_maintenance_summaries_by_provider_ids(self, provider_ids).await
|
||||
}
|
||||
|
||||
async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
@@ -2255,6 +2294,17 @@ fn map_key_stats_row(row: &PgRow) -> Result<StoredProviderCatalogKeyStats, DataL
|
||||
)
|
||||
}
|
||||
|
||||
fn map_key_maintenance_summary_row(
|
||||
row: &PgRow,
|
||||
) -> Result<StoredProviderCatalogKeyMaintenanceSummary, DataLayerError> {
|
||||
Ok(StoredProviderCatalogKeyMaintenanceSummary {
|
||||
id: row_get(row, "id")?,
|
||||
provider_id: row_get(row, "provider_id")?,
|
||||
is_active: row_get(row, "is_active")?,
|
||||
upstream_metadata: row_get(row, "upstream_metadata")?,
|
||||
})
|
||||
}
|
||||
|
||||
fn map_key_row(row: &PgRow) -> Result<StoredProviderCatalogKey, DataLayerError> {
|
||||
let rpm_limit = row_get::<Option<i32>>(row, "rpm_limit")?
|
||||
.map(|value| {
|
||||
|
||||
@@ -4,7 +4,8 @@ use sqlx::{sqlite::SqliteRow, QueryBuilder, Row, Sqlite};
|
||||
use super::{
|
||||
ProviderCatalogKeyListOrder, ProviderCatalogKeyListQuery, ProviderCatalogReadRepository,
|
||||
ProviderCatalogWriteRepository, StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
|
||||
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
StoredProviderCatalogKeyMaintenanceSummary, StoredProviderCatalogKeyPage,
|
||||
StoredProviderCatalogKeyStats, StoredProviderCatalogProvider,
|
||||
};
|
||||
use crate::driver::sqlite::{sqlite_optional_real, SqlitePool};
|
||||
use crate::error::SqlResultExt;
|
||||
@@ -278,6 +279,16 @@ FROM provider_api_keys
|
||||
WHERE provider_id IN (
|
||||
"#;
|
||||
|
||||
const LIST_KEY_MAINTENANCE_SUMMARIES_BY_PROVIDER_IDS_PREFIX: &str = r#"
|
||||
SELECT
|
||||
id,
|
||||
provider_id,
|
||||
is_active,
|
||||
upstream_metadata
|
||||
FROM provider_api_keys
|
||||
WHERE provider_id IN (
|
||||
"#;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct SqliteProviderCatalogReadRepository {
|
||||
pool: SqlitePool,
|
||||
@@ -423,6 +434,26 @@ impl SqliteProviderCatalogReadRepository {
|
||||
rows.iter().map(map_key_row).collect()
|
||||
}
|
||||
|
||||
pub async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
if provider_ids.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
let rows = build_list_query(
|
||||
LIST_KEY_MAINTENANCE_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_maintenance_summary_row).collect()
|
||||
}
|
||||
|
||||
pub async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
@@ -1322,6 +1353,13 @@ impl ProviderCatalogReadRepository for SqliteProviderCatalogReadRepository {
|
||||
Self::list_key_summaries_by_provider_ids(self, provider_ids).await
|
||||
}
|
||||
|
||||
async fn list_key_maintenance_summaries_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
|
||||
Self::list_key_maintenance_summaries_by_provider_ids(self, provider_ids).await
|
||||
}
|
||||
|
||||
async fn list_keys_page(
|
||||
&self,
|
||||
query: &ProviderCatalogKeyListQuery,
|
||||
@@ -1805,6 +1843,20 @@ fn map_key_stats_row(row: &SqliteRow) -> Result<StoredProviderCatalogKeyStats, D
|
||||
)
|
||||
}
|
||||
|
||||
fn map_key_maintenance_summary_row(
|
||||
row: &SqliteRow,
|
||||
) -> Result<StoredProviderCatalogKeyMaintenanceSummary, DataLayerError> {
|
||||
Ok(StoredProviderCatalogKeyMaintenanceSummary {
|
||||
id: row.try_get("id").map_sql_err()?,
|
||||
provider_id: row.try_get("provider_id").map_sql_err()?,
|
||||
is_active: row.try_get("is_active").map_sql_err()?,
|
||||
upstream_metadata: optional_json_from_string(
|
||||
row.try_get("upstream_metadata").map_sql_err()?,
|
||||
"provider_api_keys.upstream_metadata",
|
||||
)?,
|
||||
})
|
||||
}
|
||||
|
||||
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() {
|
||||
|
||||
@@ -2,6 +2,7 @@ services:
|
||||
app:
|
||||
image: ${APP_IMAGE:-ghcr.io/fawney19/aether:latest}
|
||||
container_name: aether-app
|
||||
user: "0:0"
|
||||
env_file:
|
||||
- ${AETHER_ENV_FILE:-.env}
|
||||
environment:
|
||||
@@ -14,17 +15,17 @@ services:
|
||||
AETHER_RUNTIME_BACKEND: memory
|
||||
AETHER_GATEWAY_DEPLOYMENT_TOPOLOGY: single-node
|
||||
AETHER_GATEWAY_NODE_ROLE: all
|
||||
AETHER_LOG_DESTINATION: ${AETHER_LOG_DESTINATION:-stdout}
|
||||
AETHER_LOG_DESTINATION: stdout
|
||||
AETHER_LOG_FORMAT: ${AETHER_LOG_FORMAT:-pretty}
|
||||
AETHER_LOG_DIR: ${AETHER_LOG_DIR:-/opt/aether/logs}
|
||||
AETHER_LOG_ROTATION: ${AETHER_LOG_ROTATION:-daily}
|
||||
AETHER_LOG_RETENTION_DAYS: ${AETHER_LOG_RETENTION_DAYS:-7}
|
||||
AETHER_LOG_MAX_FILES: ${AETHER_LOG_MAX_FILES:-30}
|
||||
APP_PORT: ${APP_PORT:-8084}
|
||||
AETHER_GATEWAY_AUTO_PREPARE_DATABASE: ${AETHER_GATEWAY_AUTO_PREPARE_DATABASE:-true}
|
||||
ports:
|
||||
- "${APP_PORT:-8084}:${APP_PORT:-8084}"
|
||||
volumes:
|
||||
- ./data:/opt/aether/data
|
||||
- ./logs:/opt/aether/logs
|
||||
logging:
|
||||
driver: json-file
|
||||
options:
|
||||
max-size: "100m"
|
||||
max-file: "10"
|
||||
restart: unless-stopped
|
||||
|
||||
+7
-7
@@ -76,6 +76,7 @@ services:
|
||||
app:
|
||||
image: ${APP_IMAGE:-ghcr.io/fawney19/aether:latest}
|
||||
container_name: aether-app
|
||||
user: "0:0"
|
||||
env_file:
|
||||
- .env
|
||||
environment:
|
||||
@@ -85,12 +86,8 @@ services:
|
||||
AETHER_BASE_DIR: /opt/aether
|
||||
AETHER_UPDATE_STRATEGY: docker
|
||||
AETHER_DOCKER_UPDATE_COMMAND: ${AETHER_DOCKER_UPDATE_COMMAND:-./update.sh}
|
||||
AETHER_LOG_DESTINATION: ${AETHER_LOG_DESTINATION:-stdout}
|
||||
AETHER_LOG_DESTINATION: stdout
|
||||
AETHER_LOG_FORMAT: ${AETHER_LOG_FORMAT:-pretty}
|
||||
AETHER_LOG_DIR: ${AETHER_LOG_DIR:-/opt/aether/logs}
|
||||
AETHER_LOG_ROTATION: ${AETHER_LOG_ROTATION:-daily}
|
||||
AETHER_LOG_RETENTION_DAYS: ${AETHER_LOG_RETENTION_DAYS:-7}
|
||||
AETHER_LOG_MAX_FILES: ${AETHER_LOG_MAX_FILES:-30}
|
||||
APP_PORT: ${APP_PORT:-8084}
|
||||
AETHER_GATEWAY_AUTO_PREPARE_DATABASE: ${AETHER_GATEWAY_AUTO_PREPARE_DATABASE:-true}
|
||||
depends_on:
|
||||
@@ -100,8 +97,11 @@ services:
|
||||
condition: service_healthy
|
||||
ports:
|
||||
- "${APP_PORT:-8084}:${APP_PORT:-8084}"
|
||||
volumes:
|
||||
- ./logs:/opt/aether/logs
|
||||
logging:
|
||||
driver: json-file
|
||||
options:
|
||||
max-size: "100m"
|
||||
max-file: "10"
|
||||
restart: unless-stopped
|
||||
|
||||
volumes:
|
||||
|
||||
Reference in New Issue
Block a user