mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-08 18:37:46 +08:00
fix(models): isolate legacy catalog rows during fetch
This commit is contained in:
@@ -407,7 +407,8 @@ fn is_nonfatal_legacy_provider_key_credential_error(error: &GatewayError) -> boo
|
||||
|| message.contains("provider catalog credential contains reserved framing")
|
||||
|| message.contains("provider catalog credential authentication failed")
|
||||
|| message.contains("provider catalog credential envelope")
|
||||
|| message.contains("provider catalog key provider binding changed during credential migration")
|
||||
|| message
|
||||
.contains("provider catalog key provider binding changed during credential migration")
|
||||
}
|
||||
|
||||
/// Stored provider/endpoint/key proxy secrets are opened independently by the
|
||||
@@ -525,7 +526,9 @@ mod tests {
|
||||
&GatewayError::Internal("stored provider catalog credential is empty".to_string())
|
||||
));
|
||||
assert!(is_nonfatal_legacy_catalog_credential_error(
|
||||
&GatewayError::Internal("Aether secret envelope has the wrong record binding".to_string())
|
||||
&GatewayError::Internal(
|
||||
"Aether secret envelope has the wrong record binding".to_string()
|
||||
)
|
||||
));
|
||||
for field in ["api_formats", "allowed_models"] {
|
||||
assert!(is_nonfatal_legacy_catalog_credential_error(
|
||||
|
||||
@@ -129,7 +129,7 @@ where
|
||||
}
|
||||
|
||||
let providers = state
|
||||
.list_provider_catalog_providers(true)
|
||||
.list_provider_catalog_providers_for_model_fetch(true)
|
||||
.await?
|
||||
.into_iter()
|
||||
.filter(|provider| provider_id_filter.is_none_or(|provider_id| provider.id == provider_id))
|
||||
@@ -144,7 +144,7 @@ where
|
||||
.collect::<Vec<_>>();
|
||||
let mut endpoints_by_provider = HashMap::<String, Vec<StoredProviderCatalogEndpoint>>::new();
|
||||
for endpoint in state
|
||||
.list_provider_catalog_endpoints_by_provider_ids(&provider_ids)
|
||||
.list_provider_catalog_endpoints_for_model_fetch(&provider_ids)
|
||||
.await?
|
||||
{
|
||||
endpoints_by_provider
|
||||
@@ -168,8 +168,12 @@ where
|
||||
for provider in providers {
|
||||
let endpoints = endpoints_by_provider
|
||||
.remove(&provider.id)
|
||||
.unwrap_or_default();
|
||||
.unwrap_or_default()
|
||||
.into_iter()
|
||||
.map(sanitize_model_fetch_endpoint)
|
||||
.collect::<Vec<_>>();
|
||||
let keys = keys_by_provider.remove(&provider.id).unwrap_or_default();
|
||||
let provider = sanitize_model_fetch_provider(provider);
|
||||
for key in keys {
|
||||
if key_id_filter.is_some_and(|key_ids| !key_ids.contains(&key.id)) {
|
||||
continue;
|
||||
@@ -239,6 +243,32 @@ fn sanitize_model_fetch_key(mut key: StoredProviderCatalogKey) -> StoredProvider
|
||||
key
|
||||
}
|
||||
|
||||
/// Keep only non-secret provider metadata in a background fetch target. The
|
||||
/// authoritative transport snapshot is reopened by ID immediately before a
|
||||
/// request, so carrying stored proxy/config JSON here would needlessly retain
|
||||
/// credentials and could expose malformed historical values to later stages.
|
||||
fn sanitize_model_fetch_provider(
|
||||
mut provider: StoredProviderCatalogProvider,
|
||||
) -> StoredProviderCatalogProvider {
|
||||
provider.proxy = None;
|
||||
provider.config = None;
|
||||
provider
|
||||
}
|
||||
|
||||
/// Endpoint selection needs only activity, format, and identity. Clear
|
||||
/// transport rules/proxy data because those are reloaded from the snapshot
|
||||
/// just before execution.
|
||||
fn sanitize_model_fetch_endpoint(
|
||||
mut endpoint: StoredProviderCatalogEndpoint,
|
||||
) -> StoredProviderCatalogEndpoint {
|
||||
endpoint.header_rules = None;
|
||||
endpoint.body_rules = None;
|
||||
endpoint.config = None;
|
||||
endpoint.format_acceptance_config = None;
|
||||
endpoint.proxy = None;
|
||||
endpoint
|
||||
}
|
||||
|
||||
async fn execute_fetch_targets<S>(
|
||||
state: &S,
|
||||
targets: Vec<SelectedFetchTarget>,
|
||||
@@ -569,6 +599,15 @@ fn is_nonfatal_legacy_credential_error(error: &GatewayError) -> bool {
|
||||
}
|
||||
message.contains("provider_api_keys.api_key")
|
||||
|| message.contains("provider_api_keys.auth_config")
|
||||
|| message.contains("stored provider proxy credentials cannot be decrypted")
|
||||
|| message.contains("stored endpoint proxy credentials cannot be decrypted")
|
||||
|| message.contains("stored key proxy credentials cannot be decrypted")
|
||||
|| message.contains("stored provider proxy changed during credential migration")
|
||||
|| message.contains("stored endpoint proxy changed during credential migration")
|
||||
|| message.contains("stored key changed during credential migration")
|
||||
|| message.contains("stored provider proxy credential migration did not stabilize")
|
||||
|| message.contains("stored endpoint proxy credential migration did not stabilize")
|
||||
|| message.contains("stored key proxy credential migration did not stabilize")
|
||||
|| message.contains("legacy provider catalog credential")
|
||||
|| message.contains("provider catalog credential is not an authenticated ciphertext")
|
||||
|| message.contains("provider catalog credential contains reserved framing")
|
||||
@@ -1612,6 +1651,32 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sanitized_model_fetch_provider_and_endpoint_drop_transport_secrets() {
|
||||
let mut provider = sample_provider("provider-sanitize", "openai");
|
||||
provider.proxy = Some(json!({"url": "http://user:[email protected]"}));
|
||||
provider.config = Some(json!({"api_key": "provider-secret"}));
|
||||
let sanitized_provider = super::sanitize_model_fetch_provider(provider);
|
||||
assert_eq!(sanitized_provider.proxy, None);
|
||||
assert_eq!(sanitized_provider.config, None);
|
||||
|
||||
let mut endpoint =
|
||||
sample_endpoint("endpoint-sanitize", "provider-sanitize", "openai:responses");
|
||||
endpoint.header_rules = Some(json!({"authorization": "Bearer endpoint-secret"}));
|
||||
endpoint.body_rules = Some(json!({"token": "endpoint-secret"}));
|
||||
endpoint.config = Some(json!({"password": "endpoint-secret"}));
|
||||
endpoint.format_acceptance_config = Some(json!({"secret": "endpoint-secret"}));
|
||||
endpoint.proxy = Some(json!({"url": "http://user:[email protected]"}));
|
||||
let sanitized_endpoint = super::sanitize_model_fetch_endpoint(endpoint);
|
||||
assert_eq!(sanitized_endpoint.header_rules, None);
|
||||
assert_eq!(sanitized_endpoint.body_rules, None);
|
||||
assert_eq!(sanitized_endpoint.config, None);
|
||||
assert_eq!(sanitized_endpoint.format_acceptance_config, None);
|
||||
assert_eq!(sanitized_endpoint.proxy, None);
|
||||
assert_eq!(sanitized_endpoint.id, "endpoint-sanitize");
|
||||
assert_eq!(sanitized_endpoint.api_format, "openai:responses");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn model_fetch_failure_does_not_persist_upstream_error_body_credentials() {
|
||||
const UPSTREAM_SECRET: &str = "upstream-secret-token-value";
|
||||
|
||||
@@ -26,11 +26,32 @@ pub(crate) trait ModelFetchRuntimeState:
|
||||
active_only: bool,
|
||||
) -> Result<Vec<StoredProviderCatalogProvider>, GatewayError>;
|
||||
|
||||
/// Return provider rows for the background fetcher without opening or
|
||||
/// migrating stored proxy credentials. Production implementations should
|
||||
/// use the raw repository projection so one malformed historical row does
|
||||
/// not abort the entire cycle and a read does not rewrite old data.
|
||||
async fn list_provider_catalog_providers_for_model_fetch(
|
||||
&self,
|
||||
active_only: bool,
|
||||
) -> Result<Vec<StoredProviderCatalogProvider>, GatewayError> {
|
||||
self.list_provider_catalog_providers(active_only).await
|
||||
}
|
||||
|
||||
async fn list_provider_catalog_endpoints_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogEndpoint>, GatewayError>;
|
||||
|
||||
/// Raw endpoint counterpart to
|
||||
/// [`Self::list_provider_catalog_providers_for_model_fetch`].
|
||||
async fn list_provider_catalog_endpoints_for_model_fetch(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogEndpoint>, GatewayError> {
|
||||
self.list_provider_catalog_endpoints_by_provider_ids(provider_ids)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Return raw catalog key rows for the background fetcher. Production
|
||||
/// implementations should avoid the normal bulk credential-opening
|
||||
/// wrapper here: one malformed legacy row must not prevent healthy keys
|
||||
|
||||
@@ -397,6 +397,16 @@ impl ModelFetchRuntimeState for AppState {
|
||||
AppState::list_provider_catalog_providers(self, active_only).await
|
||||
}
|
||||
|
||||
async fn list_provider_catalog_providers_for_model_fetch(
|
||||
&self,
|
||||
active_only: bool,
|
||||
) -> Result<Vec<StoredProviderCatalogProvider>, GatewayError> {
|
||||
self.data
|
||||
.list_provider_catalog_providers(active_only)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))
|
||||
}
|
||||
|
||||
async fn list_provider_catalog_endpoints_by_provider_ids(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
@@ -404,6 +414,16 @@ impl ModelFetchRuntimeState for AppState {
|
||||
AppState::list_provider_catalog_endpoints_by_provider_ids(self, provider_ids).await
|
||||
}
|
||||
|
||||
async fn list_provider_catalog_endpoints_for_model_fetch(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
) -> Result<Vec<StoredProviderCatalogEndpoint>, GatewayError> {
|
||||
self.data
|
||||
.list_provider_catalog_endpoints_by_provider_ids(provider_ids)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))
|
||||
}
|
||||
|
||||
async fn list_provider_catalog_keys_for_model_fetch(
|
||||
&self,
|
||||
provider_ids: &[String],
|
||||
|
||||
Reference in New Issue
Block a user