mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 17:30:23 +08:00
fix monitoring trace lookup fallback
This commit is contained in:
@@ -81,6 +81,133 @@ async fn admin_monitoring_trace_request_returns_local_payload() {
|
||||
assert_eq!(payload["candidates"][0]["status_code"], json!(502));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn admin_monitoring_trace_request_resolves_usage_id_to_header_trace_id() {
|
||||
let request_candidates = Arc::new(InMemoryRequestCandidateRepository::seed(vec![
|
||||
sample_candidate(
|
||||
"cand-used",
|
||||
"trace-1",
|
||||
0,
|
||||
RequestCandidateStatus::Success,
|
||||
Some(101),
|
||||
Some(33),
|
||||
Some(200),
|
||||
),
|
||||
]));
|
||||
let provider_catalog = Arc::new(InMemoryProviderCatalogReadRepository::seed(
|
||||
vec![sample_provider()],
|
||||
vec![sample_endpoint()],
|
||||
vec![sample_key()],
|
||||
));
|
||||
let mut usage = sample_usage(
|
||||
"usage-request-1",
|
||||
"provider-1",
|
||||
"OpenAI",
|
||||
40,
|
||||
0.02,
|
||||
"completed",
|
||||
Some(200),
|
||||
100,
|
||||
);
|
||||
usage.id = "usage-row-1".to_string();
|
||||
usage.candidate_id = Some("cand-used".to_string());
|
||||
usage.request_headers = Some(json!({
|
||||
"x-trace-id": "trace-1"
|
||||
}));
|
||||
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![usage]));
|
||||
let data_state =
|
||||
crate::data::GatewayDataState::with_request_candidate_and_usage_repository_for_tests(
|
||||
request_candidates,
|
||||
usage_repository,
|
||||
)
|
||||
.with_provider_catalog_reader(provider_catalog);
|
||||
let state = AppState::new()
|
||||
.expect("state should build")
|
||||
.with_data_state_for_tests(data_state);
|
||||
let context = request_context(
|
||||
http::Method::GET,
|
||||
"/api/admin/monitoring/trace/usage-row-1?attempted_only=true",
|
||||
);
|
||||
|
||||
let response = local_monitoring_response(&state, &context)
|
||||
.await
|
||||
.expect("handler should not error")
|
||||
.expect("route should be handled locally");
|
||||
|
||||
assert_eq!(response.status(), http::StatusCode::OK);
|
||||
let body = to_bytes(response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("body should read");
|
||||
let payload: serde_json::Value = serde_json::from_slice(&body).expect("json body should parse");
|
||||
assert_eq!(payload["request_id"], json!("trace-1"));
|
||||
assert_eq!(payload["candidates"][0]["id"], json!("cand-used"));
|
||||
assert_eq!(
|
||||
payload["candidates"][0]["extra_data"]["first_byte_time_ms"],
|
||||
json!(30)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn admin_monitoring_trace_request_resolves_usage_request_id_to_metadata_trace_id() {
|
||||
let request_candidates = Arc::new(InMemoryRequestCandidateRepository::seed(vec![
|
||||
sample_candidate(
|
||||
"cand-used",
|
||||
"trace-2",
|
||||
0,
|
||||
RequestCandidateStatus::Success,
|
||||
Some(101),
|
||||
Some(33),
|
||||
Some(200),
|
||||
),
|
||||
]));
|
||||
let provider_catalog = Arc::new(InMemoryProviderCatalogReadRepository::seed(
|
||||
vec![sample_provider()],
|
||||
vec![sample_endpoint()],
|
||||
vec![sample_key()],
|
||||
));
|
||||
let mut usage = sample_usage(
|
||||
"usage-request-2",
|
||||
"provider-1",
|
||||
"OpenAI",
|
||||
40,
|
||||
0.02,
|
||||
"completed",
|
||||
Some(200),
|
||||
100,
|
||||
);
|
||||
usage.candidate_id = Some("cand-used".to_string());
|
||||
usage.request_metadata = Some(json!({
|
||||
"trace_id": "trace-2"
|
||||
}));
|
||||
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![usage]));
|
||||
let data_state =
|
||||
crate::data::GatewayDataState::with_request_candidate_and_usage_repository_for_tests(
|
||||
request_candidates,
|
||||
usage_repository,
|
||||
)
|
||||
.with_provider_catalog_reader(provider_catalog);
|
||||
let state = AppState::new()
|
||||
.expect("state should build")
|
||||
.with_data_state_for_tests(data_state);
|
||||
let context = request_context(
|
||||
http::Method::GET,
|
||||
"/api/admin/monitoring/trace/usage-request-2",
|
||||
);
|
||||
|
||||
let response = local_monitoring_response(&state, &context)
|
||||
.await
|
||||
.expect("handler should not error")
|
||||
.expect("route should be handled locally");
|
||||
|
||||
assert_eq!(response.status(), http::StatusCode::OK);
|
||||
let body = to_bytes(response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("body should read");
|
||||
let payload: serde_json::Value = serde_json::from_slice(&body).expect("json body should parse");
|
||||
assert_eq!(payload["request_id"], json!("trace-2"));
|
||||
assert_eq!(payload["candidates"][0]["id"], json!("cand-used"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn admin_monitoring_trace_request_returns_oauth_account_label_from_auth_config() {
|
||||
let request_candidates = Arc::new(InMemoryRequestCandidateRepository::seed(vec![
|
||||
|
||||
@@ -12,6 +12,7 @@ use aether_admin::observability::monitoring::{
|
||||
use aether_data_contracts::repository::{
|
||||
candidates::{DecisionTrace, RequestCandidateStatus},
|
||||
provider_catalog::StoredProviderCatalogKey,
|
||||
usage::StoredRequestUsageAudit,
|
||||
};
|
||||
use axum::{
|
||||
body::Body,
|
||||
@@ -21,12 +22,16 @@ use serde_json::{Map, Value};
|
||||
use std::collections::BTreeMap;
|
||||
use tracing::debug;
|
||||
|
||||
struct ResolvedAdminMonitoringTrace {
|
||||
trace: DecisionTrace,
|
||||
usage: Option<StoredRequestUsageAudit>,
|
||||
}
|
||||
|
||||
pub(super) async fn build_admin_monitoring_trace_request_response(
|
||||
state: &AdminAppState<'_>,
|
||||
request_context: &AdminRequestContext<'_>,
|
||||
) -> Result<Response<Body>, GatewayError> {
|
||||
let admin_state = state;
|
||||
let state = state.as_ref();
|
||||
let Some(request_id) =
|
||||
admin_monitoring_trace_request_id_from_path(&request_context.request_path)
|
||||
else {
|
||||
@@ -39,11 +44,8 @@ pub(super) async fn build_admin_monitoring_trace_request_response(
|
||||
Err(detail) => return Ok(admin_monitoring_bad_request_response(detail)),
|
||||
};
|
||||
|
||||
let Some(trace) = state
|
||||
.data
|
||||
.read_decision_trace(&request_id, attempted_only)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))?
|
||||
let Some(resolved) =
|
||||
resolve_admin_monitoring_trace(admin_state, &request_id, attempted_only).await?
|
||||
else {
|
||||
debug!(
|
||||
event_name = "admin_monitoring_request_trace_not_found",
|
||||
@@ -58,22 +60,113 @@ pub(super) async fn build_admin_monitoring_trace_request_response(
|
||||
attempted_only,
|
||||
));
|
||||
};
|
||||
let usage = state
|
||||
.data
|
||||
.read_request_usage_audit(&request_id)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))?;
|
||||
let key_accounts = build_admin_monitoring_key_account_display_map(admin_state, &trace).await?;
|
||||
let key_accounts =
|
||||
build_admin_monitoring_key_account_display_map(admin_state, &resolved.trace).await?;
|
||||
|
||||
Ok(
|
||||
build_admin_monitoring_trace_request_payload_response_with_key_accounts(
|
||||
&trace,
|
||||
usage.as_ref(),
|
||||
&resolved.trace,
|
||||
resolved.usage.as_ref(),
|
||||
&key_accounts,
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
async fn resolve_admin_monitoring_trace(
|
||||
state: &AdminAppState<'_>,
|
||||
request_id: &str,
|
||||
attempted_only: bool,
|
||||
) -> Result<Option<ResolvedAdminMonitoringTrace>, GatewayError> {
|
||||
let app = state.as_ref();
|
||||
if let Some(trace) = app
|
||||
.data
|
||||
.read_decision_trace(request_id, attempted_only)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))?
|
||||
{
|
||||
let usage = app
|
||||
.data
|
||||
.read_request_usage_audit(request_id)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))?;
|
||||
return Ok(Some(ResolvedAdminMonitoringTrace { trace, usage }));
|
||||
}
|
||||
|
||||
let mut usage_candidates = Vec::new();
|
||||
if let Some(usage) = app
|
||||
.data
|
||||
.read_request_usage_audit(request_id)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))?
|
||||
{
|
||||
usage_candidates.push(usage);
|
||||
}
|
||||
if let Some(usage) = state.find_request_usage_by_id(request_id).await? {
|
||||
if !usage_candidates.iter().any(|item| item.id == usage.id) {
|
||||
usage_candidates.push(usage);
|
||||
}
|
||||
}
|
||||
|
||||
for usage in usage_candidates {
|
||||
for trace_request_id in admin_monitoring_usage_trace_request_ids(&usage) {
|
||||
if trace_request_id == request_id {
|
||||
continue;
|
||||
}
|
||||
if let Some(trace) = app
|
||||
.data
|
||||
.read_decision_trace(&trace_request_id, attempted_only)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(err.to_string()))?
|
||||
{
|
||||
return Ok(Some(ResolvedAdminMonitoringTrace {
|
||||
trace,
|
||||
usage: Some(usage),
|
||||
}));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
fn admin_monitoring_usage_trace_request_ids(usage: &StoredRequestUsageAudit) -> Vec<String> {
|
||||
let mut ids = Vec::new();
|
||||
push_non_empty_unique(&mut ids, usage.request_id.as_str());
|
||||
if let Some(trace_id) = usage.trace_id() {
|
||||
push_non_empty_unique(&mut ids, trace_id);
|
||||
}
|
||||
if let Some(trace_id) = usage_trace_id_from_headers(usage.request_headers.as_ref()) {
|
||||
push_non_empty_unique(&mut ids, trace_id.as_str());
|
||||
}
|
||||
if let Some(trace_id) = usage_trace_id_from_headers(usage.provider_request_headers.as_ref()) {
|
||||
push_non_empty_unique(&mut ids, trace_id.as_str());
|
||||
}
|
||||
ids
|
||||
}
|
||||
|
||||
fn usage_trace_id_from_headers(headers: Option<&Value>) -> Option<String> {
|
||||
let object = headers?.as_object()?;
|
||||
object.iter().find_map(|(key, value)| {
|
||||
key.eq_ignore_ascii_case(crate::constants::TRACE_ID_HEADER)
|
||||
.then(|| {
|
||||
value
|
||||
.as_str()
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
})
|
||||
.flatten()
|
||||
.map(ToOwned::to_owned)
|
||||
})
|
||||
}
|
||||
|
||||
fn push_non_empty_unique(values: &mut Vec<String>, value: &str) {
|
||||
let value = value.trim();
|
||||
if value.is_empty() || values.iter().any(|existing| existing == value) {
|
||||
return;
|
||||
}
|
||||
values.push(value.to_string());
|
||||
}
|
||||
|
||||
async fn build_admin_monitoring_key_account_display_map(
|
||||
state: &AdminAppState<'_>,
|
||||
trace: &DecisionTrace,
|
||||
|
||||
Reference in New Issue
Block a user