From 6e094487178f25636c50ba44347e18383de3699c Mon Sep 17 00:00:00 2001 From: tiammomo <26957354+tiammomo@users.noreply.github.com> Date: Mon, 7 Sep 2026 12:39:51 +0800 Subject: [PATCH 1/4] refactor(ledger): separate reporting and operations helpers --- docs/ARCHITECTURE.md | 7 +- src/enterprise_ledger.rs | 1219 +-------------------------- src/enterprise_ledger/operations.rs | 274 ++++++ src/enterprise_ledger/reporting.rs | 949 +++++++++++++++++++++ 4 files changed, 1229 insertions(+), 1220 deletions(-) create mode 100644 src/enterprise_ledger/reporting.rs diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 751b82b..55f0fea 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -179,8 +179,11 @@ Tool Use verification evidence is maintained separately in the - `src/database.rs`: SQLx PostgreSQL URL/TLS policy, pool bounds, acquisition timeout, and credential-safe location rendering. - `src/enterprise_ledger.rs`: mandatory-tenant request/attempt lifecycle, - operational log, usage-policy aggregation, budget, and append-only audit - repository. PostgreSQL is required at runtime; memory is test-only. + usage admission, budget settlement, retention, and append-only audit writes. + `enterprise_ledger/reporting.rs` owns read-only request/log/dashboard/usage + views; `enterprise_ledger/operations.rs` owns optional incident queries, + observations, validation, and row mapping. PostgreSQL is required at runtime; + memory is test-only. - `src/exchange.rs`: typed client-protocol parsing, capability/fidelity checks, Provider rendering, and cross-protocol response mapping. - `src/stream_lifecycle.rs`: shared upstream terminal state and normalized diff --git a/src/enterprise_ledger.rs b/src/enterprise_ledger.rs index f8c92a4..b085177 100644 --- a/src/enterprise_ledger.rs +++ b/src/enterprise_ledger.rs @@ -1,4 +1,5 @@ mod operations; +mod reporting; use std::{ collections::{BTreeMap, HashMap, HashSet}, @@ -1944,904 +1945,6 @@ impl EnterpriseLedger { } } - pub(crate) async fn overview(&self) -> Result { - let mut overview = EnterpriseLedgerOverview { - backend: self.backend_name(), - location: self.location().to_owned(), - lease_ttl_secs: self.lease_ttl.as_secs(), - reconcile_interval_secs: self.reconcile_interval.as_secs(), - total_requests: 0, - started_requests: 0, - completed_requests: 0, - failed_requests: 0, - cancelled_requests: 0, - unreconciled_requests: 0, - idempotent_requests: 0, - active_leases: 0, - expired_leases: 0, - chargeable_requests: 0, - estimate_only_requests: 0, - total_cost_microunits: 0, - total_billable_cost_microunits: 0, - organization_count: 0, - project_count: 0, - environment_count: 0, - }; - - match self.backend.as_ref() { - LedgerBackend::Memory(ledger) => { - let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); - let now = Instant::now(); - let mut organizations = HashSet::new(); - let mut projects = HashSet::new(); - let mut environments = HashSet::new(); - for request in ledger.requests.values() { - overview.total_requests += 1; - match request.record.state.as_str() { - "started" => overview.started_requests += 1, - "completed" => overview.completed_requests += 1, - "failed" => overview.failed_requests += 1, - "cancelled" => overview.cancelled_requests += 1, - _ => {} - } - if request.record.terminal_reason.as_deref() - == Some("lease_expired_unreconciled") - { - overview.unreconciled_requests += 1; - } - if request.idempotency_key_hash.is_some() { - overview.idempotent_requests += 1; - } - if !request.record.terminal { - if request.record.lease_expires_at > now { - overview.active_leases += 1; - } else { - overview.expired_leases += 1; - } - } - if request.record.chargeable { - overview.chargeable_requests += 1; - if request.record.terminal - && request.record.billable_cost_microunits.is_none() - { - overview.estimate_only_requests += 1; - } - } - overview.total_cost_microunits = overview - .total_cost_microunits - .saturating_add(request.record.cost_amount_microunits); - overview.total_billable_cost_microunits = overview - .total_billable_cost_microunits - .saturating_add(request.record.billable_cost_microunits.unwrap_or(0)); - organizations.insert(request.record.tenant.organization_id.clone()); - projects.insert(( - request.record.tenant.organization_id.clone(), - request.record.tenant.project_id.clone(), - )); - environments.insert(( - request.record.tenant.organization_id.clone(), - request.record.tenant.project_id.clone(), - request.record.tenant.environment_id.clone(), - )); - } - overview.organization_count = usize_to_i64(organizations.len()); - overview.project_count = usize_to_i64(projects.len()); - overview.environment_count = usize_to_i64(environments.len()); - } - LedgerBackend::Postgres(pool) => { - let row = sqlx::query( - "SELECT - count(*)::bigint AS total_requests, - count(*) FILTER (WHERE state = 'started')::bigint AS started_requests, - count(*) FILTER (WHERE state = 'completed')::bigint AS completed_requests, - count(*) FILTER (WHERE state = 'failed')::bigint AS failed_requests, - count(*) FILTER (WHERE state = 'cancelled')::bigint AS cancelled_requests, - count(*) FILTER (WHERE terminal_reason = 'lease_expired_unreconciled')::bigint AS unreconciled_requests, - count(*) FILTER (WHERE idempotency_key_hash IS NOT NULL)::bigint AS idempotent_requests, - count(*) FILTER (WHERE state = 'started' AND lease_expires_at > now())::bigint AS active_leases, - count(*) FILTER (WHERE state = 'started' AND lease_expires_at <= now())::bigint AS expired_leases, - count(*) FILTER (WHERE chargeable)::bigint AS chargeable_requests, - count(*) FILTER ( - WHERE state <> 'started' AND chargeable - AND billable_cost_microunits IS NULL - )::bigint AS estimate_only_requests, - COALESCE(sum(cost_amount_microunits), 0)::bigint AS total_cost_microunits, - COALESCE(sum(billable_cost_microunits), 0)::bigint AS total_billable_cost_microunits, - count(DISTINCT organization_id)::bigint AS organization_count, - count(DISTINCT (organization_id, project_id))::bigint AS project_count, - count(DISTINCT (organization_id, project_id, environment_id))::bigint AS environment_count - FROM modelport_gateway_requests", - ) - .fetch_one(pool) - .await?; - overview.total_requests = row.try_get("total_requests")?; - overview.started_requests = row.try_get("started_requests")?; - overview.completed_requests = row.try_get("completed_requests")?; - overview.failed_requests = row.try_get("failed_requests")?; - overview.cancelled_requests = row.try_get("cancelled_requests")?; - overview.unreconciled_requests = row.try_get("unreconciled_requests")?; - overview.idempotent_requests = row.try_get("idempotent_requests")?; - overview.active_leases = row.try_get("active_leases")?; - overview.expired_leases = row.try_get("expired_leases")?; - overview.chargeable_requests = row.try_get("chargeable_requests")?; - overview.estimate_only_requests = row.try_get("estimate_only_requests")?; - overview.total_cost_microunits = row.try_get("total_cost_microunits")?; - overview.total_billable_cost_microunits = - row.try_get("total_billable_cost_microunits")?; - overview.organization_count = row.try_get("organization_count")?; - overview.project_count = row.try_get("project_count")?; - overview.environment_count = row.try_get("environment_count")?; - } - } - Ok(overview) - } - - pub(crate) async fn list_requests( - &self, - query: &EnterpriseLedgerQuery, - ) -> Result { - let query = query.normalized()?; - match self.backend.as_ref() { - LedgerBackend::Memory(ledger) => { - let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); - let mut requests = ledger - .requests - .iter() - .filter(|(_, request)| query.matches_memory(request)) - .map(|(ledger_id, request)| { - memory_request_row( - ledger_id, - request, - usize_to_i64( - ledger - .attempts - .values() - .filter(|attempt| attempt.request_ledger_id == *ledger_id) - .count(), - ), - ) - }) - .collect::>(); - requests.sort_by(|left, right| { - right - .created_at_ms - .cmp(&left.created_at_ms) - .then_with(|| right.ledger_id.cmp(&left.ledger_id)) - }); - let total = usize_to_i64(requests.len()); - let start = query.offset().min(requests.len()); - let end = start.saturating_add(query.page_size).min(requests.len()); - Ok(EnterpriseRequestPage { - requests: requests[start..end].to_vec(), - total, - page: query.page, - page_size: query.page_size, - }) - } - LedgerBackend::Postgres(pool) => { - let count = sqlx::query_scalar::<_, i64>(REQUEST_COUNT_SQL) - .bind(query.state.as_deref()) - .bind(query.protocol.as_deref()) - .bind(query.organization_id.as_deref()) - .bind(query.project_id.as_deref()) - .bind(query.environment_id.as_deref()) - .bind(query.search.as_deref()) - .bind(query.traffic_class.as_deref()) - .fetch_one(pool) - .await?; - let rows = sqlx::query(REQUEST_LIST_SQL) - .bind(query.state.as_deref()) - .bind(query.protocol.as_deref()) - .bind(query.organization_id.as_deref()) - .bind(query.project_id.as_deref()) - .bind(query.environment_id.as_deref()) - .bind(query.search.as_deref()) - .bind(query.traffic_class.as_deref()) - .bind(usize_to_i64(query.page_size)) - .bind(usize_to_i64(query.offset())) - .bind(None::) - .fetch_all(pool) - .await?; - Ok(EnterpriseRequestPage { - requests: rows - .iter() - .map(request_row_from_pg) - .collect::>()?, - total: count, - page: query.page, - page_size: query.page_size, - }) - } - } - } - - #[cfg(test)] - pub(crate) async fn usage_rows(&self) -> Result, AppError> { - self.usage_rows_since(None).await - } - - pub(crate) async fn usage_rows_since( - &self, - since_ms: Option, - ) -> Result, AppError> { - let since_ms_i64 = since_ms.map(|value| i64::try_from(value).unwrap_or(i64::MAX)); - let mut requests = match self.backend.as_ref() { - LedgerBackend::Memory(ledger) => { - let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); - ledger - .requests - .iter() - .filter(|(_, request)| { - since_ms_i64.is_none_or(|since| request.record.created_at_ms >= since) - }) - .map(|(ledger_id, request)| { - memory_request_row( - ledger_id, - request, - usize_to_i64( - ledger - .attempts - .values() - .filter(|attempt| attempt.request_ledger_id == *ledger_id) - .count(), - ), - ) - }) - .collect::>() - } - LedgerBackend::Postgres(pool) => sqlx::query(REQUEST_LIST_SQL) - .bind(None::<&str>) - .bind(None::<&str>) - .bind(None::<&str>) - .bind(None::<&str>) - .bind(None::<&str>) - .bind(None::<&str>) - .bind(None::<&str>) - .bind(i64::MAX) - .bind(0_i64) - .bind(since_ms_i64) - .fetch_all(pool) - .await? - .iter() - .map(request_row_from_pg) - .collect::, _>>()?, - }; - requests.retain(|request| request.state != "started"); - requests.sort_by(|left, right| { - right - .created_at_ms - .cmp(&left.created_at_ms) - .then_with(|| right.ledger_id.cmp(&left.ledger_id)) - }); - Ok(requests.iter().map(operational_log_row).collect()) - } - - pub(crate) async fn operational_logs( - &self, - query: &OperationalLogQuery, - ) -> Result, AppError> { - let LedgerBackend::Postgres(pool) = self.backend.as_ref() else { - return Ok(None); - }; - - let mut summary_query = QueryBuilder::::new( - "SELECT - count(*)::bigint AS total_requests, - count(*) FILTER (WHERE r.state = 'completed')::bigint AS success_requests, - count(*) FILTER (WHERE r.tool_use_requested)::bigint AS tool_use_requests, - count(*) FILTER ( - WHERE r.tool_use_requested AND r.state = 'completed' - )::bigint AS tool_use_success_requests, - COALESCE(sum(r.input_tokens), 0)::bigint AS total_input_tokens, - COALESCE(sum(r.output_tokens), 0)::bigint AS total_output_tokens, - COALESCE(sum(r.cache_write_tokens), 0)::bigint AS total_cache_write_tokens, - COALESCE(sum(r.cache_read_tokens), 0)::bigint AS total_cache_read_tokens, - COALESCE(sum(r.cost_amount_microunits), 0)::bigint AS total_cost_microunits, - COALESCE(sum(r.actual_cost_microunits), 0)::bigint AS total_actual_cost_microunits, - COALESCE(sum(r.billable_cost_microunits), 0)::bigint AS total_billable_cost_microunits, - count(*) FILTER ( - WHERE r.billable_cost_microunits IS NOT NULL - )::bigint AS billable_requests, - count(*) FILTER ( - WHERE r.billable_cost_microunits IS NULL - )::bigint AS estimate_only_requests, - percentile_disc(0.95) WITHIN GROUP (ORDER BY r.latency_ms) - FILTER (WHERE r.latency_ms IS NOT NULL) AS latency_p95_ms, - count(r.latency_ms)::bigint AS latency_sample_count, - percentile_disc(0.95) WITHIN GROUP (ORDER BY r.first_byte_latency_ms) - FILTER (WHERE r.first_byte_latency_ms IS NOT NULL) - AS first_byte_latency_p95_ms, - count(r.first_byte_latency_ms)::bigint AS first_byte_latency_sample_count, - (EXTRACT(EPOCH FROM min(r.created_at)) * 1000)::bigint AS first_timestamp_ms, - (EXTRACT(EPOCH FROM max(r.created_at)) * 1000)::bigint AS last_timestamp_ms - FROM modelport_gateway_requests r", - ); - push_operational_log_filters(&mut summary_query, query); - let summary_row = summary_query.build().fetch_one(pool).await?; - let total: i64 = summary_row.try_get("total_requests")?; - let total_input_tokens: i64 = summary_row.try_get("total_input_tokens")?; - let total_output_tokens: i64 = summary_row.try_get("total_output_tokens")?; - let total_cache_write_tokens: i64 = summary_row.try_get("total_cache_write_tokens")?; - let total_cache_read_tokens: i64 = summary_row.try_get("total_cache_read_tokens")?; - let total_tokens = total_input_tokens - .saturating_add(total_output_tokens) - .saturating_add(total_cache_write_tokens) - .saturating_add(total_cache_read_tokens); - let first_timestamp: Option = summary_row.try_get("first_timestamp_ms")?; - let last_timestamp: Option = summary_row.try_get("last_timestamp_ms")?; - let minutes = match (first_timestamp, last_timestamp) { - (Some(first), Some(last)) if last > first => { - ((last - first) as f64 / 60_000.0).max(1.0) - } - _ => 1.0, - }; - let summary = json!({ - "totalRequests": nonnegative_u64(total), - "successRequests": nonnegative_u64(summary_row.try_get("success_requests")?), - "toolUseRequests": nonnegative_u64(summary_row.try_get("tool_use_requests")?), - "toolUseSuccessRequests": nonnegative_u64( - summary_row.try_get("tool_use_success_requests")? - ), - "totalInputTokens": nonnegative_u64(total_input_tokens), - "totalOutputTokens": nonnegative_u64(total_output_tokens), - "totalCacheWriteTokens": nonnegative_u64(total_cache_write_tokens), - "totalCacheReadTokens": nonnegative_u64(total_cache_read_tokens), - "totalTokens": nonnegative_u64(total_tokens), - "totalCostEstimate": microunits_usd( - summary_row.try_get("total_cost_microunits")? - ), - "totalActualCost": microunits_usd( - summary_row.try_get("total_actual_cost_microunits")? - ), - "totalBillableCost": microunits_usd( - summary_row.try_get("total_billable_cost_microunits")? - ), - "billableRequests": nonnegative_u64( - summary_row.try_get("billable_requests")? - ), - "estimateOnlyRequests": nonnegative_u64( - summary_row.try_get("estimate_only_requests")? - ), - "latencyP95Ms": summary_row - .try_get::, _>("latency_p95_ms")? - .map(nonnegative_u64) - .unwrap_or(0), - "latencySampleCount": nonnegative_u64( - summary_row.try_get("latency_sample_count")? - ), - "firstByteLatencyP95Ms": summary_row - .try_get::, _>("first_byte_latency_p95_ms")? - .map(nonnegative_u64) - .unwrap_or(0), - "firstByteLatencySampleCount": nonnegative_u64( - summary_row.try_get("first_byte_latency_sample_count")? - ), - "rpm": total.max(0) as f64 / minutes, - "tpm": total_tokens.max(0) as f64 / minutes, - }); - - let mut rows_query = QueryBuilder::::new(OPERATIONAL_LOG_SELECT_SQL); - push_operational_log_filters(&mut rows_query, query); - rows_query - .push(" ORDER BY r.created_at DESC, r.ledger_id DESC LIMIT ") - .push_bind(usize_to_i64(query.page_size)) - .push(" OFFSET ") - .push_bind(usize_to_i64( - query.page.saturating_sub(1).saturating_mul(query.page_size), - )); - let rows = rows_query.build().fetch_all(pool).await?; - let logs = rows - .iter() - .map(request_row_from_pg) - .collect::, _>>()? - .iter() - .map(operational_log_row) - .collect(); - - Ok(Some(OperationalLogPage { - logs, - total, - summary, - })) - } - - pub(crate) async fn dashboard_snapshot( - &self, - start_ms: u64, - end_ms: u64, - bucket_ms: u64, - today_start_ms: u64, - api_keys: (u64, u64), - ) -> Result, AppError> { - let LedgerBackend::Postgres(pool) = self.backend.as_ref() else { - return Ok(None); - }; - let start_ms = i64::try_from(start_ms).unwrap_or(i64::MAX); - let end_ms = i64::try_from(end_ms).unwrap_or(i64::MAX); - let bucket_ms = i64::try_from(bucket_ms.max(1)).unwrap_or(i64::MAX); - let today_start_ms = i64::try_from(today_start_ms).unwrap_or(i64::MAX); - - let provider_rows = sqlx::query( - "SELECT - COALESCE(provider_id, 'unrouted') AS provider_id, - count(*)::bigint AS requests, - count(*) FILTER (WHERE state = 'completed')::bigint AS successes, - COALESCE(sum(latency_ms), 0)::bigint AS duration_ms, - COALESCE(sum(input_tokens), 0)::bigint AS input_tokens, - COALESCE(sum(output_tokens), 0)::bigint AS output_tokens, - COALESCE(sum(cache_write_tokens), 0)::bigint AS cache_write_tokens, - COALESCE(sum(cache_read_tokens), 0)::bigint AS cache_read_tokens, - COALESCE(sum(cost_amount_microunits), 0)::bigint AS cost_microunits - FROM modelport_gateway_requests - WHERE state <> 'started' - AND traffic_class = 'business' - AND created_at >= to_timestamp($1::double precision / 1000.0) - GROUP BY COALESCE(provider_id, 'unrouted')", - ) - .bind(today_start_ms) - .fetch_all(pool) - .await?; - let mut usage_summary = UsageSummary { - api_keys_total: api_keys.0, - api_keys_active: api_keys.1, - ..UsageSummary::default() - }; - let mut provider_usage = BTreeMap::new(); - let mut total_duration_ms = 0u64; - for row in provider_rows { - let requests = nonnegative_u64(row.try_get("requests")?); - let successes = nonnegative_u64(row.try_get("successes")?); - let duration_ms = nonnegative_u64(row.try_get("duration_ms")?); - let input_tokens = nonnegative_u64(row.try_get("input_tokens")?); - let output_tokens = nonnegative_u64(row.try_get("output_tokens")?); - let cache_write_tokens = nonnegative_u64(row.try_get("cache_write_tokens")?); - let cache_read_tokens = nonnegative_u64(row.try_get("cache_read_tokens")?); - let cost_microunits: i64 = row.try_get("cost_microunits")?; - usage_summary.total_requests = usage_summary.total_requests.saturating_add(requests); - usage_summary.total_successes = usage_summary.total_successes.saturating_add(successes); - usage_summary.total_input_tokens = usage_summary - .total_input_tokens - .saturating_add(input_tokens); - usage_summary.total_output_tokens = usage_summary - .total_output_tokens - .saturating_add(output_tokens); - usage_summary.total_cache_write_tokens = usage_summary - .total_cache_write_tokens - .saturating_add(cache_write_tokens); - usage_summary.total_cache_read_tokens = usage_summary - .total_cache_read_tokens - .saturating_add(cache_read_tokens); - usage_summary.total_cost_estimate += microunits_usd(cost_microunits); - total_duration_ms = total_duration_ms.saturating_add(duration_ms); - provider_usage.insert( - row.try_get("provider_id")?, - ProviderUsageStats { - requests_total: requests, - successes_total: successes, - duration_ms_total: duration_ms, - input_tokens_total: input_tokens, - output_tokens_total: output_tokens, - cache_write_tokens_total: cache_write_tokens, - cache_read_tokens_total: cache_read_tokens, - cost_estimate_usd_total: microunits_usd(cost_microunits), - }, - ); - } - usage_summary.average_latency_ms = total_duration_ms - .checked_div(usage_summary.total_requests) - .unwrap_or(0); - - let bucket_count = - usize::try_from((end_ms.saturating_sub(start_ms) / bucket_ms).saturating_add(1)) - .unwrap_or(1) - .max(1); - let mut requests = vec![0u64; bucket_count]; - let mut errors = vec![0u64; bucket_count]; - let mut input_tokens = vec![0u64; bucket_count]; - let mut output_tokens = vec![0u64; bucket_count]; - let mut cache_write_tokens = vec![0u64; bucket_count]; - let mut cache_read_tokens = vec![0u64; bucket_count]; - let bucket_rows = sqlx::query( - "SELECT - floor( - ((EXTRACT(EPOCH FROM created_at) * 1000) - $1::double precision) - / $3::double precision - )::bigint AS bucket_index, - count(*)::bigint AS requests, - count(*) FILTER (WHERE state <> 'completed')::bigint AS errors, - COALESCE(sum(input_tokens), 0)::bigint AS input_tokens, - COALESCE(sum(output_tokens), 0)::bigint AS output_tokens, - COALESCE(sum(cache_write_tokens), 0)::bigint AS cache_write_tokens, - COALESCE(sum(cache_read_tokens), 0)::bigint AS cache_read_tokens, - COALESCE(sum(cost_amount_microunits), 0)::bigint AS cost_microunits - FROM modelport_gateway_requests - WHERE state <> 'started' - AND traffic_class = 'business' - AND created_at >= to_timestamp($1::double precision / 1000.0) - AND created_at <= to_timestamp($2::double precision / 1000.0) - GROUP BY bucket_index - ORDER BY bucket_index", - ) - .bind(start_ms) - .bind(end_ms) - .bind(bucket_ms) - .fetch_all(pool) - .await?; - let mut matched_requests = 0u64; - let mut success_requests = 0u64; - let mut total_input_tokens = 0u64; - let mut total_output_tokens = 0u64; - let mut total_cache_write_tokens = 0u64; - let mut total_cache_read_tokens = 0u64; - let mut total_cost_microunits = 0i64; - for row in bucket_rows { - let index = usize::try_from(row.try_get::("bucket_index")?) - .unwrap_or(bucket_count.saturating_sub(1)) - .min(bucket_count.saturating_sub(1)); - let row_requests = nonnegative_u64(row.try_get("requests")?); - let row_errors = nonnegative_u64(row.try_get("errors")?); - let row_input_tokens = nonnegative_u64(row.try_get("input_tokens")?); - let row_output_tokens = nonnegative_u64(row.try_get("output_tokens")?); - let row_cache_write_tokens = nonnegative_u64(row.try_get("cache_write_tokens")?); - let row_cache_read_tokens = nonnegative_u64(row.try_get("cache_read_tokens")?); - let row_cost_microunits: i64 = row.try_get("cost_microunits")?; - requests[index] = row_requests; - errors[index] = row_errors; - input_tokens[index] = row_input_tokens; - output_tokens[index] = row_output_tokens; - cache_write_tokens[index] = row_cache_write_tokens; - cache_read_tokens[index] = row_cache_read_tokens; - matched_requests = matched_requests.saturating_add(row_requests); - success_requests = - success_requests.saturating_add(row_requests.saturating_sub(row_errors)); - total_input_tokens = total_input_tokens.saturating_add(row_input_tokens); - total_output_tokens = total_output_tokens.saturating_add(row_output_tokens); - total_cache_write_tokens = - total_cache_write_tokens.saturating_add(row_cache_write_tokens); - total_cache_read_tokens = total_cache_read_tokens.saturating_add(row_cache_read_tokens); - total_cost_microunits = total_cost_microunits.saturating_add(row_cost_microunits); - } - - let model_rows = sqlx::query( - "SELECT - COALESCE(resolved_model, requested_model, 'unknown') AS model, - COALESCE(provider_id, 'unknown') AS provider, - count(*)::bigint AS requests, - COALESCE(sum( - input_tokens + output_tokens + cache_write_tokens + cache_read_tokens - ), 0)::bigint AS tokens, - COALESCE(sum(cost_amount_microunits), 0)::bigint AS cost_microunits - FROM modelport_gateway_requests - WHERE state <> 'started' - AND traffic_class = 'business' - AND created_at >= to_timestamp($1::double precision / 1000.0) - AND created_at <= to_timestamp($2::double precision / 1000.0) - GROUP BY - COALESCE(resolved_model, requested_model, 'unknown'), - COALESCE(provider_id, 'unknown') - ORDER BY tokens DESC, requests DESC, model ASC - LIMIT 200", - ) - .bind(start_ms) - .bind(end_ms) - .fetch_all(pool) - .await?; - let model_usage = model_rows - .iter() - .map(|row| { - Ok(json!({ - "model": row.try_get::("model")?, - "provider": row.try_get::("provider")?, - "requests": nonnegative_u64(row.try_get("requests")?), - "tokens": nonnegative_u64(row.try_get("tokens")?), - "cost": microunits_usd(row.try_get("cost_microunits")?), - })) - }) - .collect::, sqlx::Error>>()?; - let request_time_series = dashboard_value_series(&requests, start_ms, bucket_ms); - let error_time_series = dashboard_value_series(&errors, start_ms, bucket_ms); - let token_time_series = (0..bucket_count) - .map(|index| { - let billed_input = input_tokens[index] - .saturating_add(cache_write_tokens[index]) - .saturating_add(cache_read_tokens[index]); - json!({ - "timestamp": dashboard_bucket_timestamp(start_ms, bucket_ms, index), - "inputTokens": input_tokens[index], - "outputTokens": output_tokens[index], - "cacheWriteTokens": cache_write_tokens[index], - "cacheReadTokens": cache_read_tokens[index], - "cacheHitRate": if billed_input == 0 { - 0.0 - } else { - cache_read_tokens[index] as f64 / billed_input as f64 * 100.0 - }, - }) - }) - .collect(); - let total_tokens = total_input_tokens - .saturating_add(total_output_tokens) - .saturating_add(total_cache_write_tokens) - .saturating_add(total_cache_read_tokens); - let minutes = (end_ms.saturating_sub(start_ms) as f64 / 60_000.0).max(1.0); - - Ok(Some(DashboardLedgerSnapshot { - usage_summary, - provider_usage, - matched_requests, - request_time_series, - error_time_series, - token_time_series, - model_usage, - summary: json!({ - "totalRequests": matched_requests, - "successRequests": success_requests, - "totalInputTokens": total_input_tokens, - "totalOutputTokens": total_output_tokens, - "totalCacheWriteTokens": total_cache_write_tokens, - "totalCacheReadTokens": total_cache_read_tokens, - "totalTokens": total_tokens, - "totalCostEstimate": microunits_usd(total_cost_microunits), - "rpm": matched_requests as f64 / minutes, - "tpm": total_tokens as f64 / minutes, - }), - })) - } - - pub(crate) async fn latency_stats_since( - &self, - since_ms: u64, - ) -> Result, AppError> { - let LedgerBackend::Postgres(pool) = self.backend.as_ref() else { - return Ok(None); - }; - let since_ms = i64::try_from(since_ms).unwrap_or(i64::MAX); - let overall = sqlx::query( - "SELECT - percentile_disc(0.50) WITHIN GROUP (ORDER BY latency_ms) AS p50, - percentile_disc(0.90) WITHIN GROUP (ORDER BY latency_ms) AS p90, - percentile_disc(0.95) WITHIN GROUP (ORDER BY latency_ms) AS p95, - percentile_disc(0.99) WITHIN GROUP (ORDER BY latency_ms) AS p99, - floor(COALESCE(avg(latency_ms), 0))::bigint AS avg, - COALESCE(max(latency_ms), 0)::bigint AS max, - count(*)::bigint AS count - FROM modelport_gateway_requests - WHERE state <> 'started' - AND created_at >= to_timestamp($1::double precision / 1000.0)", - ) - .bind(since_ms) - .fetch_one(pool) - .await?; - let by_model_rows = sqlx::query( - "SELECT - COALESCE(resolved_model, requested_model, 'unknown') AS name, - percentile_disc(0.50) WITHIN GROUP (ORDER BY latency_ms) AS p50, - percentile_disc(0.90) WITHIN GROUP (ORDER BY latency_ms) AS p90, - percentile_disc(0.95) WITHIN GROUP (ORDER BY latency_ms) AS p95, - percentile_disc(0.99) WITHIN GROUP (ORDER BY latency_ms) AS p99, - floor(COALESCE(avg(latency_ms), 0))::bigint AS avg, - COALESCE(max(latency_ms), 0)::bigint AS max, - count(*)::bigint AS count - FROM modelport_gateway_requests - WHERE state <> 'started' - AND created_at >= to_timestamp($1::double precision / 1000.0) - GROUP BY COALESCE(resolved_model, requested_model, 'unknown') - ORDER BY count DESC - LIMIT 200", - ) - .bind(since_ms) - .fetch_all(pool) - .await?; - let by_provider_rows = sqlx::query( - "SELECT - COALESCE(provider_id, 'unrouted') AS name, - percentile_disc(0.50) WITHIN GROUP (ORDER BY latency_ms) AS p50, - percentile_disc(0.90) WITHIN GROUP (ORDER BY latency_ms) AS p90, - percentile_disc(0.95) WITHIN GROUP (ORDER BY latency_ms) AS p95, - percentile_disc(0.99) WITHIN GROUP (ORDER BY latency_ms) AS p99, - floor(COALESCE(avg(latency_ms), 0))::bigint AS avg, - COALESCE(max(latency_ms), 0)::bigint AS max, - count(*)::bigint AS count - FROM modelport_gateway_requests - WHERE state <> 'started' - AND created_at >= to_timestamp($1::double precision / 1000.0) - GROUP BY COALESCE(provider_id, 'unrouted') - ORDER BY count DESC - LIMIT 200", - ) - .bind(since_ms) - .fetch_all(pool) - .await?; - let grouped = |rows: Vec| -> Result { - let mut values = serde_json::Map::new(); - for row in rows { - values.insert(row.try_get("name")?, latency_stats_from_pg(&row)?); - } - Ok(Value::Object(values)) - }; - - Ok(Some(json!({ - "p50": optional_nonnegative_u64(&overall, "p50")?, - "p90": optional_nonnegative_u64(&overall, "p90")?, - "p95": optional_nonnegative_u64(&overall, "p95")?, - "p99": optional_nonnegative_u64(&overall, "p99")?, - "avg": nonnegative_u64(overall.try_get("avg")?), - "max": nonnegative_u64(overall.try_get("max")?), - "byModel": grouped(by_model_rows)?, - "byProvider": grouped(by_provider_rows)?, - "sampleCount": nonnegative_u64(overall.try_get("count")?), - "percentilesEstimated": false, - }))) - } - - pub(crate) async fn usage_row(&self, ledger_id: &str) -> Result, AppError> { - let request = match self.backend.as_ref() { - LedgerBackend::Memory(ledger) => { - let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); - ledger.requests.get(ledger_id).map(|request| { - memory_request_row( - ledger_id, - request, - usize_to_i64( - ledger - .attempts - .values() - .filter(|attempt| attempt.request_ledger_id == ledger_id) - .count(), - ), - ) - }) - } - LedgerBackend::Postgres(pool) => sqlx::query(REQUEST_DETAIL_SQL) - .bind(ledger_id) - .fetch_optional(pool) - .await? - .as_ref() - .map(request_row_from_pg) - .transpose()?, - }; - Ok(request - .filter(|request| request.state != "started") - .as_ref() - .map(operational_log_row)) - } - - pub(crate) async fn management_usage(&self) -> Result { - match self.backend.as_ref() { - LedgerBackend::Memory(ledger) => { - let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); - let now = u64::try_from(now_millis()).unwrap_or(u64::MAX); - let day_start = current_period("daily", now).0; - let month_start = current_period("monthly", now).0; - let rolling_day_start = now.saturating_sub(24 * 60 * 60 * 1_000); - let mut stats = ManagementUsageStats::default(); - for request in ledger - .requests - .values() - .filter(|request| request.record.terminal) - { - let created_at = u64::try_from(request.record.created_at_ms).unwrap_or(0); - if created_at >= rolling_day_start { - let requests = stats - .users_24h - .entry(request.principal_id.clone()) - .or_default(); - *requests = requests.saturating_add(1); - } - if let Some(api_key_id) = request.api_key_id.as_deref() - && created_at >= day_start - { - let row = stats.api_keys.entry(api_key_id.to_owned()).or_default(); - row.requests_today = row.requests_today.saturating_add(1); - row.tokens_today = row - .tokens_today - .saturating_add(request_total_tokens(&request.record)); - } - if let Some(team_id) = request.team_id.as_deref() - && created_at >= month_start - { - let row = stats.teams.entry(team_id.to_owned()).or_default(); - let cost = request - .record - .billable_cost_microunits - .map_or(0.0, microunits_usd); - row.monthly_spend_usd += cost; - if created_at >= day_start { - row.requests_today = row.requests_today.saturating_add(1); - row.daily_spend_usd += cost; - } - } - } - Ok(stats) - } - LedgerBackend::Postgres(pool) => { - let api_key_rows = sqlx::query( - "SELECT - api_key_id, - count(*)::bigint AS requests_today, - COALESCE(sum( - input_tokens + output_tokens - + cache_write_tokens + cache_read_tokens - ), 0)::bigint AS tokens_today - FROM modelport_gateway_requests - WHERE state <> 'started' - AND api_key_id IS NOT NULL - AND created_at >= ( - date_trunc('day', now() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' - ) - GROUP BY api_key_id", - ) - .fetch_all(pool) - .await?; - let team_rows = sqlx::query( - "SELECT - team_id, - count(*) FILTER ( - WHERE created_at >= ( - date_trunc('day', now() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' - ) - )::bigint AS requests_today, - COALESCE(sum(billable_cost_microunits) FILTER ( - WHERE chargeable - AND created_at >= ( - date_trunc('day', now() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' - ) - ), 0)::bigint AS daily_spend_microunits, - COALESCE(sum(billable_cost_microunits) FILTER ( - WHERE chargeable - ), 0)::bigint AS monthly_spend_microunits - FROM modelport_gateway_requests - WHERE state <> 'started' - AND team_id IS NOT NULL - AND created_at >= ( - date_trunc('month', now() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' - ) - GROUP BY team_id", - ) - .fetch_all(pool) - .await?; - let user_rows = sqlx::query( - "SELECT principal_id, count(*)::bigint AS requests_24h - FROM modelport_gateway_requests - WHERE state <> 'started' - AND created_at >= now() - interval '24 hours' - GROUP BY principal_id", - ) - .fetch_all(pool) - .await?; - let mut stats = ManagementUsageStats::default(); - for row in api_key_rows { - stats.api_keys.insert( - row.try_get("api_key_id")?, - ApiKeyUsageStats { - requests_today: nonnegative_u64(row.try_get("requests_today")?), - tokens_today: nonnegative_u64(row.try_get("tokens_today")?), - }, - ); - } - for row in team_rows { - stats.teams.insert( - row.try_get("team_id")?, - TeamUsageStats { - requests_today: nonnegative_u64(row.try_get("requests_today")?), - daily_spend_usd: microunits_usd(row.try_get("daily_spend_microunits")?), - monthly_spend_usd: microunits_usd( - row.try_get("monthly_spend_microunits")?, - ), - }, - ); - } - for row in user_rows { - stats.users_24h.insert( - row.try_get("principal_id")?, - nonnegative_u64(row.try_get("requests_24h")?), - ); - } - Ok(stats) - } - } - } - pub(crate) async fn check_usage_policy( &self, policy: &UsagePolicySnapshot, @@ -3181,52 +2284,6 @@ impl EnterpriseLedger { } } - pub(crate) async fn request_detail( - &self, - ledger_id: &str, - ) -> Result, AppError> { - match self.backend.as_ref() { - LedgerBackend::Memory(ledger) => { - let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); - let Some(request) = ledger.requests.get(ledger_id) else { - return Ok(None); - }; - let mut attempts = ledger - .attempts - .iter() - .filter(|(_, attempt)| attempt.request_ledger_id == ledger_id) - .map(|(attempt_id, attempt)| memory_attempt_row(attempt_id, attempt)) - .collect::>(); - attempts.sort_by_key(|attempt| attempt.created_at_ms); - Ok(Some(EnterpriseRequestDetail { - request: memory_request_row(ledger_id, request, usize_to_i64(attempts.len())), - attempts, - })) - } - LedgerBackend::Postgres(pool) => { - let Some(row) = sqlx::query(REQUEST_DETAIL_SQL) - .bind(ledger_id) - .fetch_optional(pool) - .await? - else { - return Ok(None); - }; - let request = request_row_from_pg(&row)?; - let attempt_rows = sqlx::query(ATTEMPT_LIST_SQL) - .bind(ledger_id) - .fetch_all(pool) - .await?; - Ok(Some(EnterpriseRequestDetail { - request, - attempts: attempt_rows - .iter() - .map(attempt_row_from_pg) - .collect::>()?, - })) - } - } - } - pub(crate) async fn budget_view( &self, scope: &EnterpriseBudgetScopeQuery, @@ -6546,280 +5603,6 @@ fn duration_millis_i64(duration: Duration) -> i64 { i64::try_from(duration.as_millis()).unwrap_or(i64::MAX) } -fn validate_ops_observation(observation: &OpsObservation) -> Result<(), AppError> { - let bounded = [ - ("eventKey", observation.event_key.as_str(), 240_usize), - ("detectorType", observation.detector_type.as_str(), 80), - ("title", observation.title.as_str(), 240), - ("summary", observation.summary.as_str(), 2_000), - ( - "recoveryCriteria", - observation.recovery_criteria.as_str(), - 1_000, - ), - ]; - for (name, value, maximum) in bounded { - let length = value.chars().count(); - if length == 0 || length > maximum { - return Err(AppError::InvalidRequest(format!( - "{name} must contain 1 to {maximum} characters" - ))); - } - } - let allowed = match observation.event_key.as_str() { - "readiness:gateway" => { - observation.detector_type == "readiness_storage" - && matches!(observation.severity, OpsSeverity::Sev1 | OpsSeverity::Sev2) - } - "provider:availability" => { - observation.detector_type == "provider_health" - && matches!(observation.severity, OpsSeverity::Sev2 | OpsSeverity::Sev3) - } - "requests:failure-ratio" => { - observation.detector_type == "request_anomaly" - && matches!(observation.severity, OpsSeverity::Sev2 | OpsSeverity::Sev3) - } - "budget:capacity" => { - observation.detector_type == "budget_quota" - && matches!(observation.severity, OpsSeverity::Sev2 | OpsSeverity::Sev3) - } - "ledger:finalization-backlog" => { - observation.detector_type == "ledger_backlog" - && matches!(observation.severity, OpsSeverity::Sev2 | OpsSeverity::Sev3) - } - "change:verification" => { - observation.detector_type == "post_change_verification" - && observation.severity == OpsSeverity::Sev2 - } - _ => false, - }; - if !allowed { - return Err(AppError::Forbidden( - "observation is outside the versioned operations rule allowlist".to_owned(), - )); - } - if !observation.affected_scope.is_object() || !observation.evidence.is_object() { - return Err(AppError::InvalidRequest( - "affectedScope and evidence must be JSON objects".to_owned(), - )); - } - if serde_json::to_vec(&observation.evidence)?.len() > 32 * 1_024 { - return Err(AppError::InvalidRequest( - "incident evidence must not exceed 32 KiB".to_owned(), - )); - } - if serde_json::to_vec(&observation.affected_scope)?.len() > 8 * 1_024 { - return Err(AppError::InvalidRequest( - "incident affectedScope must not exceed 8 KiB".to_owned(), - )); - } - let now = u64::try_from(now_millis()).unwrap_or_default(); - if observation.observed_at_ms == 0 - || observation.observed_at_ms > now.saturating_add(5 * 60 * 1_000) - { - return Err(AppError::InvalidRequest( - "observedAtMs must be a current, non-zero timestamp".to_owned(), - )); - } - Ok(()) -} - -fn validate_ops_heartbeat(heartbeat: &OpsHeartbeat) -> Result<(), AppError> { - for (name, value, maximum) in [ - ("instanceId", heartbeat.instance_id.as_str(), 160_usize), - ("agentVersion", heartbeat.agent_version.as_str(), 80), - ("ruleSetVersion", heartbeat.rule_set_version.as_str(), 80), - ] { - let length = value.chars().count(); - if length == 0 || length > maximum { - return Err(AppError::InvalidRequest(format!( - "{name} must contain 1 to {maximum} characters" - ))); - } - } - if heartbeat.selected_model.as_deref().is_some_and(|value| { - value.is_empty() || value.chars().count() > 320 || value.chars().any(char::is_control) - }) { - return Err(AppError::InvalidRequest( - "heartbeat selectedModel must contain 1 to 320 non-control characters".to_owned(), - )); - } - if !matches!( - heartbeat.model_status.as_str(), - "disabled" | "configured" | "missing_credential" | "error" - ) { - return Err(AppError::InvalidRequest( - "heartbeat modelStatus is unsupported".to_owned(), - )); - } - if !matches!( - heartbeat.mode.as_str(), - "disabled" | "replay" | "shadow" | "read_only" - ) { - return Err(AppError::InvalidRequest( - "agent mode must be disabled, replay, shadow, or read_only".to_owned(), - )); - } - if !(10..=3_600).contains(&heartbeat.interval_seconds) { - return Err(AppError::InvalidRequest( - "agent intervalSeconds must be between 10 and 3600".to_owned(), - )); - } - let now = u64::try_from(now_millis()).unwrap_or_default(); - if heartbeat.observed_at_ms == 0 - || heartbeat.observed_at_ms > now.saturating_add(5 * 60 * 1_000) - { - return Err(AppError::InvalidRequest( - "heartbeat observedAtMs must be a current, non-zero timestamp".to_owned(), - )); - } - Ok(()) -} - -fn validate_ops_actor(actor_id: &str, actor_name: &str) -> Result<(), AppError> { - if actor_id.is_empty() - || actor_name.is_empty() - || actor_id.chars().count() > 160 - || actor_name.chars().count() > 160 - { - return Err(AppError::InvalidRequest( - "operations actor identity is missing or too long".to_owned(), - )); - } - Ok(()) -} - -fn ops_evidence_hash(value: &Value) -> Result { - let encoded = serde_json::to_vec(value)?; - Ok(format!("{:x}", Sha256::digest(encoded))) -} - -async fn fetch_ops_incident_row( - transaction: &mut sqlx::Transaction<'_, Postgres>, - incident_id: &str, -) -> Result { - Ok(sqlx::query( - "SELECT *, - (EXTRACT(EPOCH FROM first_seen_at) * 1000)::bigint AS first_seen_at_ms, - (EXTRACT(EPOCH FROM last_seen_at) * 1000)::bigint AS last_seen_at_ms, - (EXTRACT(EPOCH FROM resolved_at) * 1000)::bigint AS resolved_at_ms - FROM modelport_ops_incidents WHERE incident_id = $1", - ) - .bind(incident_id) - .fetch_one(&mut **transaction) - .await?) -} - -fn ops_incident_summary_from_row(row: &PgRow) -> Result { - let severity: String = row.try_get("severity")?; - let status: String = row.try_get("status")?; - Ok(OpsIncidentSummary { - id: row.try_get("incident_id")?, - event_key: row.try_get("event_key")?, - detector_type: row.try_get("detector_type")?, - severity: parse_ops_severity(&severity)?, - status: parse_ops_status(&status)?, - title: row.try_get("title")?, - summary: row.try_get("summary")?, - affected_scope: row.try_get("affected_scope")?, - recovery_criteria: row.try_get("recovery_criteria")?, - first_seen_at_ms: nonnegative_u64(row.try_get("first_seen_at_ms")?), - last_seen_at_ms: nonnegative_u64(row.try_get("last_seen_at_ms")?), - resolved_at_ms: row - .try_get::, _>("resolved_at_ms")? - .map(nonnegative_u64), - occurrence_count: nonnegative_u64(row.try_get("occurrence_count")?), - }) -} - -fn parse_ops_severity(value: &str) -> Result { - match value { - "SEV-1" => Ok(OpsSeverity::Sev1), - "SEV-2" => Ok(OpsSeverity::Sev2), - "SEV-3" => Ok(OpsSeverity::Sev3), - "SEV-4" => Ok(OpsSeverity::Sev4), - _ => Err(AppError::Database( - "operations incident contains an invalid severity".to_owned(), - )), - } -} - -fn parse_ops_status(value: &str) -> Result { - match value { - "open" => Ok(OpsIncidentStatus::Open), - "acknowledged" => Ok(OpsIncidentStatus::Acknowledged), - "mitigating" => Ok(OpsIncidentStatus::Mitigating), - "monitoring" => Ok(OpsIncidentStatus::Monitoring), - "resolved" => Ok(OpsIncidentStatus::Resolved), - "suppressed" => Ok(OpsIncidentStatus::Suppressed), - _ => Err(AppError::Database( - "operations incident contains an invalid status".to_owned(), - )), - } -} - -fn highest_ops_severity(values: Vec) -> Option { - values.into_iter().min_by_key(|severity| match severity { - OpsSeverity::Sev1 => 1, - OpsSeverity::Sev2 => 2, - OpsSeverity::Sev3 => 3, - OpsSeverity::Sev4 => 4, - }) -} - -fn ops_severity_from_rank(value: i32) -> Option { - match value { - 1 => Some(OpsSeverity::Sev1), - 2 => Some(OpsSeverity::Sev2), - 3 => Some(OpsSeverity::Sev3), - 4 => Some(OpsSeverity::Sev4), - _ => None, - } -} - -fn ops_agent_summary(heartbeat: &OpsHeartbeat) -> OpsAgentSummary { - OpsAgentSummary { - instance_id: heartbeat.instance_id.clone(), - agent_version: heartbeat.agent_version.clone(), - mode: heartbeat.mode.clone(), - rule_set_version: heartbeat.rule_set_version.clone(), - observed_at_ms: heartbeat.observed_at_ms, - queue_depth: heartbeat.queue_depth, - interval_seconds: heartbeat.interval_seconds, - online: u64::try_from(now_millis()) - .unwrap_or_default() - .saturating_sub(heartbeat.observed_at_ms) - <= heartbeat.interval_seconds.saturating_mul(3_000), - analysis_enabled: heartbeat.analysis_enabled, - selected_model: heartbeat.selected_model.clone(), - model_status: heartbeat.model_status.clone(), - model_last_success_at_ms: heartbeat.model_last_success_at_ms, - } -} - -fn ops_agent_summary_from_row(row: &PgRow) -> Result { - let observed_at_ms = nonnegative_u64(row.try_get("observed_at_ms")?); - Ok(OpsAgentSummary { - instance_id: row.try_get("instance_id")?, - agent_version: row.try_get("agent_version")?, - mode: row.try_get("mode")?, - rule_set_version: row.try_get("rule_set_version")?, - observed_at_ms, - queue_depth: nonnegative_u64(row.try_get("queue_depth")?), - interval_seconds: nonnegative_u64(row.try_get("interval_seconds")?), - online: u64::try_from(now_millis()) - .unwrap_or_default() - .saturating_sub(observed_at_ms) - <= nonnegative_u64(row.try_get("interval_seconds")?).saturating_mul(3_000), - analysis_enabled: row.try_get("analysis_enabled")?, - selected_model: row.try_get("selected_model")?, - model_status: row.try_get("model_status")?, - model_last_success_at_ms: row - .try_get::, _>("model_last_success_at_ms")? - .map(nonnegative_u64), - }) -} - fn now_millis() -> i64 { SystemTime::now() .duration_since(UNIX_EPOCH) diff --git a/src/enterprise_ledger/operations.rs b/src/enterprise_ledger/operations.rs index cb7a68a..3d69f0d 100644 --- a/src/enterprise_ledger/operations.rs +++ b/src/enterprise_ledger/operations.rs @@ -897,3 +897,277 @@ impl EnterpriseLedger { Ok(()) } } + +fn validate_ops_observation(observation: &OpsObservation) -> Result<(), AppError> { + let bounded = [ + ("eventKey", observation.event_key.as_str(), 240_usize), + ("detectorType", observation.detector_type.as_str(), 80), + ("title", observation.title.as_str(), 240), + ("summary", observation.summary.as_str(), 2_000), + ( + "recoveryCriteria", + observation.recovery_criteria.as_str(), + 1_000, + ), + ]; + for (name, value, maximum) in bounded { + let length = value.chars().count(); + if length == 0 || length > maximum { + return Err(AppError::InvalidRequest(format!( + "{name} must contain 1 to {maximum} characters" + ))); + } + } + let allowed = match observation.event_key.as_str() { + "readiness:gateway" => { + observation.detector_type == "readiness_storage" + && matches!(observation.severity, OpsSeverity::Sev1 | OpsSeverity::Sev2) + } + "provider:availability" => { + observation.detector_type == "provider_health" + && matches!(observation.severity, OpsSeverity::Sev2 | OpsSeverity::Sev3) + } + "requests:failure-ratio" => { + observation.detector_type == "request_anomaly" + && matches!(observation.severity, OpsSeverity::Sev2 | OpsSeverity::Sev3) + } + "budget:capacity" => { + observation.detector_type == "budget_quota" + && matches!(observation.severity, OpsSeverity::Sev2 | OpsSeverity::Sev3) + } + "ledger:finalization-backlog" => { + observation.detector_type == "ledger_backlog" + && matches!(observation.severity, OpsSeverity::Sev2 | OpsSeverity::Sev3) + } + "change:verification" => { + observation.detector_type == "post_change_verification" + && observation.severity == OpsSeverity::Sev2 + } + _ => false, + }; + if !allowed { + return Err(AppError::Forbidden( + "observation is outside the versioned operations rule allowlist".to_owned(), + )); + } + if !observation.affected_scope.is_object() || !observation.evidence.is_object() { + return Err(AppError::InvalidRequest( + "affectedScope and evidence must be JSON objects".to_owned(), + )); + } + if serde_json::to_vec(&observation.evidence)?.len() > 32 * 1_024 { + return Err(AppError::InvalidRequest( + "incident evidence must not exceed 32 KiB".to_owned(), + )); + } + if serde_json::to_vec(&observation.affected_scope)?.len() > 8 * 1_024 { + return Err(AppError::InvalidRequest( + "incident affectedScope must not exceed 8 KiB".to_owned(), + )); + } + let now = u64::try_from(now_millis()).unwrap_or_default(); + if observation.observed_at_ms == 0 + || observation.observed_at_ms > now.saturating_add(5 * 60 * 1_000) + { + return Err(AppError::InvalidRequest( + "observedAtMs must be a current, non-zero timestamp".to_owned(), + )); + } + Ok(()) +} + +fn validate_ops_heartbeat(heartbeat: &OpsHeartbeat) -> Result<(), AppError> { + for (name, value, maximum) in [ + ("instanceId", heartbeat.instance_id.as_str(), 160_usize), + ("agentVersion", heartbeat.agent_version.as_str(), 80), + ("ruleSetVersion", heartbeat.rule_set_version.as_str(), 80), + ] { + let length = value.chars().count(); + if length == 0 || length > maximum { + return Err(AppError::InvalidRequest(format!( + "{name} must contain 1 to {maximum} characters" + ))); + } + } + if heartbeat.selected_model.as_deref().is_some_and(|value| { + value.is_empty() || value.chars().count() > 320 || value.chars().any(char::is_control) + }) { + return Err(AppError::InvalidRequest( + "heartbeat selectedModel must contain 1 to 320 non-control characters".to_owned(), + )); + } + if !matches!( + heartbeat.model_status.as_str(), + "disabled" | "configured" | "missing_credential" | "error" + ) { + return Err(AppError::InvalidRequest( + "heartbeat modelStatus is unsupported".to_owned(), + )); + } + if !matches!( + heartbeat.mode.as_str(), + "disabled" | "replay" | "shadow" | "read_only" + ) { + return Err(AppError::InvalidRequest( + "agent mode must be disabled, replay, shadow, or read_only".to_owned(), + )); + } + if !(10..=3_600).contains(&heartbeat.interval_seconds) { + return Err(AppError::InvalidRequest( + "agent intervalSeconds must be between 10 and 3600".to_owned(), + )); + } + let now = u64::try_from(now_millis()).unwrap_or_default(); + if heartbeat.observed_at_ms == 0 + || heartbeat.observed_at_ms > now.saturating_add(5 * 60 * 1_000) + { + return Err(AppError::InvalidRequest( + "heartbeat observedAtMs must be a current, non-zero timestamp".to_owned(), + )); + } + Ok(()) +} + +fn validate_ops_actor(actor_id: &str, actor_name: &str) -> Result<(), AppError> { + if actor_id.is_empty() + || actor_name.is_empty() + || actor_id.chars().count() > 160 + || actor_name.chars().count() > 160 + { + return Err(AppError::InvalidRequest( + "operations actor identity is missing or too long".to_owned(), + )); + } + Ok(()) +} + +fn ops_evidence_hash(value: &Value) -> Result { + let encoded = serde_json::to_vec(value)?; + Ok(format!("{:x}", Sha256::digest(encoded))) +} + +async fn fetch_ops_incident_row( + transaction: &mut sqlx::Transaction<'_, Postgres>, + incident_id: &str, +) -> Result { + Ok(sqlx::query( + "SELECT *, + (EXTRACT(EPOCH FROM first_seen_at) * 1000)::bigint AS first_seen_at_ms, + (EXTRACT(EPOCH FROM last_seen_at) * 1000)::bigint AS last_seen_at_ms, + (EXTRACT(EPOCH FROM resolved_at) * 1000)::bigint AS resolved_at_ms + FROM modelport_ops_incidents WHERE incident_id = $1", + ) + .bind(incident_id) + .fetch_one(&mut **transaction) + .await?) +} + +fn ops_incident_summary_from_row(row: &PgRow) -> Result { + let severity: String = row.try_get("severity")?; + let status: String = row.try_get("status")?; + Ok(OpsIncidentSummary { + id: row.try_get("incident_id")?, + event_key: row.try_get("event_key")?, + detector_type: row.try_get("detector_type")?, + severity: parse_ops_severity(&severity)?, + status: parse_ops_status(&status)?, + title: row.try_get("title")?, + summary: row.try_get("summary")?, + affected_scope: row.try_get("affected_scope")?, + recovery_criteria: row.try_get("recovery_criteria")?, + first_seen_at_ms: nonnegative_u64(row.try_get("first_seen_at_ms")?), + last_seen_at_ms: nonnegative_u64(row.try_get("last_seen_at_ms")?), + resolved_at_ms: row + .try_get::, _>("resolved_at_ms")? + .map(nonnegative_u64), + occurrence_count: nonnegative_u64(row.try_get("occurrence_count")?), + }) +} + +fn parse_ops_severity(value: &str) -> Result { + match value { + "SEV-1" => Ok(OpsSeverity::Sev1), + "SEV-2" => Ok(OpsSeverity::Sev2), + "SEV-3" => Ok(OpsSeverity::Sev3), + "SEV-4" => Ok(OpsSeverity::Sev4), + _ => Err(AppError::Database( + "operations incident contains an invalid severity".to_owned(), + )), + } +} + +fn parse_ops_status(value: &str) -> Result { + match value { + "open" => Ok(OpsIncidentStatus::Open), + "acknowledged" => Ok(OpsIncidentStatus::Acknowledged), + "mitigating" => Ok(OpsIncidentStatus::Mitigating), + "monitoring" => Ok(OpsIncidentStatus::Monitoring), + "resolved" => Ok(OpsIncidentStatus::Resolved), + "suppressed" => Ok(OpsIncidentStatus::Suppressed), + _ => Err(AppError::Database( + "operations incident contains an invalid status".to_owned(), + )), + } +} + +fn highest_ops_severity(values: Vec) -> Option { + values.into_iter().min_by_key(|severity| match severity { + OpsSeverity::Sev1 => 1, + OpsSeverity::Sev2 => 2, + OpsSeverity::Sev3 => 3, + OpsSeverity::Sev4 => 4, + }) +} + +fn ops_severity_from_rank(value: i32) -> Option { + match value { + 1 => Some(OpsSeverity::Sev1), + 2 => Some(OpsSeverity::Sev2), + 3 => Some(OpsSeverity::Sev3), + 4 => Some(OpsSeverity::Sev4), + _ => None, + } +} + +fn ops_agent_summary(heartbeat: &OpsHeartbeat) -> OpsAgentSummary { + OpsAgentSummary { + instance_id: heartbeat.instance_id.clone(), + agent_version: heartbeat.agent_version.clone(), + mode: heartbeat.mode.clone(), + rule_set_version: heartbeat.rule_set_version.clone(), + observed_at_ms: heartbeat.observed_at_ms, + queue_depth: heartbeat.queue_depth, + interval_seconds: heartbeat.interval_seconds, + online: u64::try_from(now_millis()) + .unwrap_or_default() + .saturating_sub(heartbeat.observed_at_ms) + <= heartbeat.interval_seconds.saturating_mul(3_000), + analysis_enabled: heartbeat.analysis_enabled, + selected_model: heartbeat.selected_model.clone(), + model_status: heartbeat.model_status.clone(), + model_last_success_at_ms: heartbeat.model_last_success_at_ms, + } +} + +fn ops_agent_summary_from_row(row: &PgRow) -> Result { + let observed_at_ms = nonnegative_u64(row.try_get("observed_at_ms")?); + Ok(OpsAgentSummary { + instance_id: row.try_get("instance_id")?, + agent_version: row.try_get("agent_version")?, + mode: row.try_get("mode")?, + rule_set_version: row.try_get("rule_set_version")?, + observed_at_ms, + queue_depth: nonnegative_u64(row.try_get("queue_depth")?), + interval_seconds: nonnegative_u64(row.try_get("interval_seconds")?), + online: u64::try_from(now_millis()) + .unwrap_or_default() + .saturating_sub(observed_at_ms) + <= nonnegative_u64(row.try_get("interval_seconds")?).saturating_mul(3_000), + analysis_enabled: row.try_get("analysis_enabled")?, + selected_model: row.try_get("selected_model")?, + model_status: row.try_get("model_status")?, + model_last_success_at_ms: row + .try_get::, _>("model_last_success_at_ms")? + .map(nonnegative_u64), + }) +} diff --git a/src/enterprise_ledger/reporting.rs b/src/enterprise_ledger/reporting.rs new file mode 100644 index 0000000..fd7ca3d --- /dev/null +++ b/src/enterprise_ledger/reporting.rs @@ -0,0 +1,949 @@ +//! Read-only ledger views for the console, request evidence, and usage reporting. + +use super::*; + +impl EnterpriseLedger { + pub(crate) async fn overview(&self) -> Result { + let mut overview = EnterpriseLedgerOverview { + backend: self.backend_name(), + location: self.location().to_owned(), + lease_ttl_secs: self.lease_ttl.as_secs(), + reconcile_interval_secs: self.reconcile_interval.as_secs(), + total_requests: 0, + started_requests: 0, + completed_requests: 0, + failed_requests: 0, + cancelled_requests: 0, + unreconciled_requests: 0, + idempotent_requests: 0, + active_leases: 0, + expired_leases: 0, + chargeable_requests: 0, + estimate_only_requests: 0, + total_cost_microunits: 0, + total_billable_cost_microunits: 0, + organization_count: 0, + project_count: 0, + environment_count: 0, + }; + + match self.backend.as_ref() { + LedgerBackend::Memory(ledger) => { + let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); + let now = Instant::now(); + let mut organizations = HashSet::new(); + let mut projects = HashSet::new(); + let mut environments = HashSet::new(); + for request in ledger.requests.values() { + overview.total_requests += 1; + match request.record.state.as_str() { + "started" => overview.started_requests += 1, + "completed" => overview.completed_requests += 1, + "failed" => overview.failed_requests += 1, + "cancelled" => overview.cancelled_requests += 1, + _ => {} + } + if request.record.terminal_reason.as_deref() + == Some("lease_expired_unreconciled") + { + overview.unreconciled_requests += 1; + } + if request.idempotency_key_hash.is_some() { + overview.idempotent_requests += 1; + } + if !request.record.terminal { + if request.record.lease_expires_at > now { + overview.active_leases += 1; + } else { + overview.expired_leases += 1; + } + } + if request.record.chargeable { + overview.chargeable_requests += 1; + if request.record.terminal + && request.record.billable_cost_microunits.is_none() + { + overview.estimate_only_requests += 1; + } + } + overview.total_cost_microunits = overview + .total_cost_microunits + .saturating_add(request.record.cost_amount_microunits); + overview.total_billable_cost_microunits = overview + .total_billable_cost_microunits + .saturating_add(request.record.billable_cost_microunits.unwrap_or(0)); + organizations.insert(request.record.tenant.organization_id.clone()); + projects.insert(( + request.record.tenant.organization_id.clone(), + request.record.tenant.project_id.clone(), + )); + environments.insert(( + request.record.tenant.organization_id.clone(), + request.record.tenant.project_id.clone(), + request.record.tenant.environment_id.clone(), + )); + } + overview.organization_count = usize_to_i64(organizations.len()); + overview.project_count = usize_to_i64(projects.len()); + overview.environment_count = usize_to_i64(environments.len()); + } + LedgerBackend::Postgres(pool) => { + let row = sqlx::query( + "SELECT + count(*)::bigint AS total_requests, + count(*) FILTER (WHERE state = 'started')::bigint AS started_requests, + count(*) FILTER (WHERE state = 'completed')::bigint AS completed_requests, + count(*) FILTER (WHERE state = 'failed')::bigint AS failed_requests, + count(*) FILTER (WHERE state = 'cancelled')::bigint AS cancelled_requests, + count(*) FILTER (WHERE terminal_reason = 'lease_expired_unreconciled')::bigint AS unreconciled_requests, + count(*) FILTER (WHERE idempotency_key_hash IS NOT NULL)::bigint AS idempotent_requests, + count(*) FILTER (WHERE state = 'started' AND lease_expires_at > now())::bigint AS active_leases, + count(*) FILTER (WHERE state = 'started' AND lease_expires_at <= now())::bigint AS expired_leases, + count(*) FILTER (WHERE chargeable)::bigint AS chargeable_requests, + count(*) FILTER ( + WHERE state <> 'started' AND chargeable + AND billable_cost_microunits IS NULL + )::bigint AS estimate_only_requests, + COALESCE(sum(cost_amount_microunits), 0)::bigint AS total_cost_microunits, + COALESCE(sum(billable_cost_microunits), 0)::bigint AS total_billable_cost_microunits, + count(DISTINCT organization_id)::bigint AS organization_count, + count(DISTINCT (organization_id, project_id))::bigint AS project_count, + count(DISTINCT (organization_id, project_id, environment_id))::bigint AS environment_count + FROM modelport_gateway_requests", + ) + .fetch_one(pool) + .await?; + overview.total_requests = row.try_get("total_requests")?; + overview.started_requests = row.try_get("started_requests")?; + overview.completed_requests = row.try_get("completed_requests")?; + overview.failed_requests = row.try_get("failed_requests")?; + overview.cancelled_requests = row.try_get("cancelled_requests")?; + overview.unreconciled_requests = row.try_get("unreconciled_requests")?; + overview.idempotent_requests = row.try_get("idempotent_requests")?; + overview.active_leases = row.try_get("active_leases")?; + overview.expired_leases = row.try_get("expired_leases")?; + overview.chargeable_requests = row.try_get("chargeable_requests")?; + overview.estimate_only_requests = row.try_get("estimate_only_requests")?; + overview.total_cost_microunits = row.try_get("total_cost_microunits")?; + overview.total_billable_cost_microunits = + row.try_get("total_billable_cost_microunits")?; + overview.organization_count = row.try_get("organization_count")?; + overview.project_count = row.try_get("project_count")?; + overview.environment_count = row.try_get("environment_count")?; + } + } + Ok(overview) + } + + pub(crate) async fn list_requests( + &self, + query: &EnterpriseLedgerQuery, + ) -> Result { + let query = query.normalized()?; + match self.backend.as_ref() { + LedgerBackend::Memory(ledger) => { + let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); + let mut requests = ledger + .requests + .iter() + .filter(|(_, request)| query.matches_memory(request)) + .map(|(ledger_id, request)| { + memory_request_row( + ledger_id, + request, + usize_to_i64( + ledger + .attempts + .values() + .filter(|attempt| attempt.request_ledger_id == *ledger_id) + .count(), + ), + ) + }) + .collect::>(); + requests.sort_by(|left, right| { + right + .created_at_ms + .cmp(&left.created_at_ms) + .then_with(|| right.ledger_id.cmp(&left.ledger_id)) + }); + let total = usize_to_i64(requests.len()); + let start = query.offset().min(requests.len()); + let end = start.saturating_add(query.page_size).min(requests.len()); + Ok(EnterpriseRequestPage { + requests: requests[start..end].to_vec(), + total, + page: query.page, + page_size: query.page_size, + }) + } + LedgerBackend::Postgres(pool) => { + let count = sqlx::query_scalar::<_, i64>(REQUEST_COUNT_SQL) + .bind(query.state.as_deref()) + .bind(query.protocol.as_deref()) + .bind(query.organization_id.as_deref()) + .bind(query.project_id.as_deref()) + .bind(query.environment_id.as_deref()) + .bind(query.search.as_deref()) + .bind(query.traffic_class.as_deref()) + .fetch_one(pool) + .await?; + let rows = sqlx::query(REQUEST_LIST_SQL) + .bind(query.state.as_deref()) + .bind(query.protocol.as_deref()) + .bind(query.organization_id.as_deref()) + .bind(query.project_id.as_deref()) + .bind(query.environment_id.as_deref()) + .bind(query.search.as_deref()) + .bind(query.traffic_class.as_deref()) + .bind(usize_to_i64(query.page_size)) + .bind(usize_to_i64(query.offset())) + .bind(None::) + .fetch_all(pool) + .await?; + Ok(EnterpriseRequestPage { + requests: rows + .iter() + .map(request_row_from_pg) + .collect::>()?, + total: count, + page: query.page, + page_size: query.page_size, + }) + } + } + } + + #[cfg(test)] + pub(crate) async fn usage_rows(&self) -> Result, AppError> { + self.usage_rows_since(None).await + } + + pub(crate) async fn usage_rows_since( + &self, + since_ms: Option, + ) -> Result, AppError> { + let since_ms_i64 = since_ms.map(|value| i64::try_from(value).unwrap_or(i64::MAX)); + let mut requests = match self.backend.as_ref() { + LedgerBackend::Memory(ledger) => { + let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); + ledger + .requests + .iter() + .filter(|(_, request)| { + since_ms_i64.is_none_or(|since| request.record.created_at_ms >= since) + }) + .map(|(ledger_id, request)| { + memory_request_row( + ledger_id, + request, + usize_to_i64( + ledger + .attempts + .values() + .filter(|attempt| attempt.request_ledger_id == *ledger_id) + .count(), + ), + ) + }) + .collect::>() + } + LedgerBackend::Postgres(pool) => sqlx::query(REQUEST_LIST_SQL) + .bind(None::<&str>) + .bind(None::<&str>) + .bind(None::<&str>) + .bind(None::<&str>) + .bind(None::<&str>) + .bind(None::<&str>) + .bind(None::<&str>) + .bind(i64::MAX) + .bind(0_i64) + .bind(since_ms_i64) + .fetch_all(pool) + .await? + .iter() + .map(request_row_from_pg) + .collect::, _>>()?, + }; + requests.retain(|request| request.state != "started"); + requests.sort_by(|left, right| { + right + .created_at_ms + .cmp(&left.created_at_ms) + .then_with(|| right.ledger_id.cmp(&left.ledger_id)) + }); + Ok(requests.iter().map(operational_log_row).collect()) + } + + pub(crate) async fn operational_logs( + &self, + query: &OperationalLogQuery, + ) -> Result, AppError> { + let LedgerBackend::Postgres(pool) = self.backend.as_ref() else { + return Ok(None); + }; + + let mut summary_query = QueryBuilder::::new( + "SELECT + count(*)::bigint AS total_requests, + count(*) FILTER (WHERE r.state = 'completed')::bigint AS success_requests, + count(*) FILTER (WHERE r.tool_use_requested)::bigint AS tool_use_requests, + count(*) FILTER ( + WHERE r.tool_use_requested AND r.state = 'completed' + )::bigint AS tool_use_success_requests, + COALESCE(sum(r.input_tokens), 0)::bigint AS total_input_tokens, + COALESCE(sum(r.output_tokens), 0)::bigint AS total_output_tokens, + COALESCE(sum(r.cache_write_tokens), 0)::bigint AS total_cache_write_tokens, + COALESCE(sum(r.cache_read_tokens), 0)::bigint AS total_cache_read_tokens, + COALESCE(sum(r.cost_amount_microunits), 0)::bigint AS total_cost_microunits, + COALESCE(sum(r.actual_cost_microunits), 0)::bigint AS total_actual_cost_microunits, + COALESCE(sum(r.billable_cost_microunits), 0)::bigint AS total_billable_cost_microunits, + count(*) FILTER ( + WHERE r.billable_cost_microunits IS NOT NULL + )::bigint AS billable_requests, + count(*) FILTER ( + WHERE r.billable_cost_microunits IS NULL + )::bigint AS estimate_only_requests, + percentile_disc(0.95) WITHIN GROUP (ORDER BY r.latency_ms) + FILTER (WHERE r.latency_ms IS NOT NULL) AS latency_p95_ms, + count(r.latency_ms)::bigint AS latency_sample_count, + percentile_disc(0.95) WITHIN GROUP (ORDER BY r.first_byte_latency_ms) + FILTER (WHERE r.first_byte_latency_ms IS NOT NULL) + AS first_byte_latency_p95_ms, + count(r.first_byte_latency_ms)::bigint AS first_byte_latency_sample_count, + (EXTRACT(EPOCH FROM min(r.created_at)) * 1000)::bigint AS first_timestamp_ms, + (EXTRACT(EPOCH FROM max(r.created_at)) * 1000)::bigint AS last_timestamp_ms + FROM modelport_gateway_requests r", + ); + push_operational_log_filters(&mut summary_query, query); + let summary_row = summary_query.build().fetch_one(pool).await?; + let total: i64 = summary_row.try_get("total_requests")?; + let total_input_tokens: i64 = summary_row.try_get("total_input_tokens")?; + let total_output_tokens: i64 = summary_row.try_get("total_output_tokens")?; + let total_cache_write_tokens: i64 = summary_row.try_get("total_cache_write_tokens")?; + let total_cache_read_tokens: i64 = summary_row.try_get("total_cache_read_tokens")?; + let total_tokens = total_input_tokens + .saturating_add(total_output_tokens) + .saturating_add(total_cache_write_tokens) + .saturating_add(total_cache_read_tokens); + let first_timestamp: Option = summary_row.try_get("first_timestamp_ms")?; + let last_timestamp: Option = summary_row.try_get("last_timestamp_ms")?; + let minutes = match (first_timestamp, last_timestamp) { + (Some(first), Some(last)) if last > first => { + ((last - first) as f64 / 60_000.0).max(1.0) + } + _ => 1.0, + }; + let summary = json!({ + "totalRequests": nonnegative_u64(total), + "successRequests": nonnegative_u64(summary_row.try_get("success_requests")?), + "toolUseRequests": nonnegative_u64(summary_row.try_get("tool_use_requests")?), + "toolUseSuccessRequests": nonnegative_u64( + summary_row.try_get("tool_use_success_requests")? + ), + "totalInputTokens": nonnegative_u64(total_input_tokens), + "totalOutputTokens": nonnegative_u64(total_output_tokens), + "totalCacheWriteTokens": nonnegative_u64(total_cache_write_tokens), + "totalCacheReadTokens": nonnegative_u64(total_cache_read_tokens), + "totalTokens": nonnegative_u64(total_tokens), + "totalCostEstimate": microunits_usd( + summary_row.try_get("total_cost_microunits")? + ), + "totalActualCost": microunits_usd( + summary_row.try_get("total_actual_cost_microunits")? + ), + "totalBillableCost": microunits_usd( + summary_row.try_get("total_billable_cost_microunits")? + ), + "billableRequests": nonnegative_u64( + summary_row.try_get("billable_requests")? + ), + "estimateOnlyRequests": nonnegative_u64( + summary_row.try_get("estimate_only_requests")? + ), + "latencyP95Ms": summary_row + .try_get::, _>("latency_p95_ms")? + .map(nonnegative_u64) + .unwrap_or(0), + "latencySampleCount": nonnegative_u64( + summary_row.try_get("latency_sample_count")? + ), + "firstByteLatencyP95Ms": summary_row + .try_get::, _>("first_byte_latency_p95_ms")? + .map(nonnegative_u64) + .unwrap_or(0), + "firstByteLatencySampleCount": nonnegative_u64( + summary_row.try_get("first_byte_latency_sample_count")? + ), + "rpm": total.max(0) as f64 / minutes, + "tpm": total_tokens.max(0) as f64 / minutes, + }); + + let mut rows_query = QueryBuilder::::new(OPERATIONAL_LOG_SELECT_SQL); + push_operational_log_filters(&mut rows_query, query); + rows_query + .push(" ORDER BY r.created_at DESC, r.ledger_id DESC LIMIT ") + .push_bind(usize_to_i64(query.page_size)) + .push(" OFFSET ") + .push_bind(usize_to_i64( + query.page.saturating_sub(1).saturating_mul(query.page_size), + )); + let rows = rows_query.build().fetch_all(pool).await?; + let logs = rows + .iter() + .map(request_row_from_pg) + .collect::, _>>()? + .iter() + .map(operational_log_row) + .collect(); + + Ok(Some(OperationalLogPage { + logs, + total, + summary, + })) + } + + pub(crate) async fn dashboard_snapshot( + &self, + start_ms: u64, + end_ms: u64, + bucket_ms: u64, + today_start_ms: u64, + api_keys: (u64, u64), + ) -> Result, AppError> { + let LedgerBackend::Postgres(pool) = self.backend.as_ref() else { + return Ok(None); + }; + let start_ms = i64::try_from(start_ms).unwrap_or(i64::MAX); + let end_ms = i64::try_from(end_ms).unwrap_or(i64::MAX); + let bucket_ms = i64::try_from(bucket_ms.max(1)).unwrap_or(i64::MAX); + let today_start_ms = i64::try_from(today_start_ms).unwrap_or(i64::MAX); + + let provider_rows = sqlx::query( + "SELECT + COALESCE(provider_id, 'unrouted') AS provider_id, + count(*)::bigint AS requests, + count(*) FILTER (WHERE state = 'completed')::bigint AS successes, + COALESCE(sum(latency_ms), 0)::bigint AS duration_ms, + COALESCE(sum(input_tokens), 0)::bigint AS input_tokens, + COALESCE(sum(output_tokens), 0)::bigint AS output_tokens, + COALESCE(sum(cache_write_tokens), 0)::bigint AS cache_write_tokens, + COALESCE(sum(cache_read_tokens), 0)::bigint AS cache_read_tokens, + COALESCE(sum(cost_amount_microunits), 0)::bigint AS cost_microunits + FROM modelport_gateway_requests + WHERE state <> 'started' + AND traffic_class = 'business' + AND created_at >= to_timestamp($1::double precision / 1000.0) + GROUP BY COALESCE(provider_id, 'unrouted')", + ) + .bind(today_start_ms) + .fetch_all(pool) + .await?; + let mut usage_summary = UsageSummary { + api_keys_total: api_keys.0, + api_keys_active: api_keys.1, + ..UsageSummary::default() + }; + let mut provider_usage = BTreeMap::new(); + let mut total_duration_ms = 0u64; + for row in provider_rows { + let requests = nonnegative_u64(row.try_get("requests")?); + let successes = nonnegative_u64(row.try_get("successes")?); + let duration_ms = nonnegative_u64(row.try_get("duration_ms")?); + let input_tokens = nonnegative_u64(row.try_get("input_tokens")?); + let output_tokens = nonnegative_u64(row.try_get("output_tokens")?); + let cache_write_tokens = nonnegative_u64(row.try_get("cache_write_tokens")?); + let cache_read_tokens = nonnegative_u64(row.try_get("cache_read_tokens")?); + let cost_microunits: i64 = row.try_get("cost_microunits")?; + usage_summary.total_requests = usage_summary.total_requests.saturating_add(requests); + usage_summary.total_successes = usage_summary.total_successes.saturating_add(successes); + usage_summary.total_input_tokens = usage_summary + .total_input_tokens + .saturating_add(input_tokens); + usage_summary.total_output_tokens = usage_summary + .total_output_tokens + .saturating_add(output_tokens); + usage_summary.total_cache_write_tokens = usage_summary + .total_cache_write_tokens + .saturating_add(cache_write_tokens); + usage_summary.total_cache_read_tokens = usage_summary + .total_cache_read_tokens + .saturating_add(cache_read_tokens); + usage_summary.total_cost_estimate += microunits_usd(cost_microunits); + total_duration_ms = total_duration_ms.saturating_add(duration_ms); + provider_usage.insert( + row.try_get("provider_id")?, + ProviderUsageStats { + requests_total: requests, + successes_total: successes, + duration_ms_total: duration_ms, + input_tokens_total: input_tokens, + output_tokens_total: output_tokens, + cache_write_tokens_total: cache_write_tokens, + cache_read_tokens_total: cache_read_tokens, + cost_estimate_usd_total: microunits_usd(cost_microunits), + }, + ); + } + usage_summary.average_latency_ms = total_duration_ms + .checked_div(usage_summary.total_requests) + .unwrap_or(0); + + let bucket_count = + usize::try_from((end_ms.saturating_sub(start_ms) / bucket_ms).saturating_add(1)) + .unwrap_or(1) + .max(1); + let mut requests = vec![0u64; bucket_count]; + let mut errors = vec![0u64; bucket_count]; + let mut input_tokens = vec![0u64; bucket_count]; + let mut output_tokens = vec![0u64; bucket_count]; + let mut cache_write_tokens = vec![0u64; bucket_count]; + let mut cache_read_tokens = vec![0u64; bucket_count]; + let bucket_rows = sqlx::query( + "SELECT + floor( + ((EXTRACT(EPOCH FROM created_at) * 1000) - $1::double precision) + / $3::double precision + )::bigint AS bucket_index, + count(*)::bigint AS requests, + count(*) FILTER (WHERE state <> 'completed')::bigint AS errors, + COALESCE(sum(input_tokens), 0)::bigint AS input_tokens, + COALESCE(sum(output_tokens), 0)::bigint AS output_tokens, + COALESCE(sum(cache_write_tokens), 0)::bigint AS cache_write_tokens, + COALESCE(sum(cache_read_tokens), 0)::bigint AS cache_read_tokens, + COALESCE(sum(cost_amount_microunits), 0)::bigint AS cost_microunits + FROM modelport_gateway_requests + WHERE state <> 'started' + AND traffic_class = 'business' + AND created_at >= to_timestamp($1::double precision / 1000.0) + AND created_at <= to_timestamp($2::double precision / 1000.0) + GROUP BY bucket_index + ORDER BY bucket_index", + ) + .bind(start_ms) + .bind(end_ms) + .bind(bucket_ms) + .fetch_all(pool) + .await?; + let mut matched_requests = 0u64; + let mut success_requests = 0u64; + let mut total_input_tokens = 0u64; + let mut total_output_tokens = 0u64; + let mut total_cache_write_tokens = 0u64; + let mut total_cache_read_tokens = 0u64; + let mut total_cost_microunits = 0i64; + for row in bucket_rows { + let index = usize::try_from(row.try_get::("bucket_index")?) + .unwrap_or(bucket_count.saturating_sub(1)) + .min(bucket_count.saturating_sub(1)); + let row_requests = nonnegative_u64(row.try_get("requests")?); + let row_errors = nonnegative_u64(row.try_get("errors")?); + let row_input_tokens = nonnegative_u64(row.try_get("input_tokens")?); + let row_output_tokens = nonnegative_u64(row.try_get("output_tokens")?); + let row_cache_write_tokens = nonnegative_u64(row.try_get("cache_write_tokens")?); + let row_cache_read_tokens = nonnegative_u64(row.try_get("cache_read_tokens")?); + let row_cost_microunits: i64 = row.try_get("cost_microunits")?; + requests[index] = row_requests; + errors[index] = row_errors; + input_tokens[index] = row_input_tokens; + output_tokens[index] = row_output_tokens; + cache_write_tokens[index] = row_cache_write_tokens; + cache_read_tokens[index] = row_cache_read_tokens; + matched_requests = matched_requests.saturating_add(row_requests); + success_requests = + success_requests.saturating_add(row_requests.saturating_sub(row_errors)); + total_input_tokens = total_input_tokens.saturating_add(row_input_tokens); + total_output_tokens = total_output_tokens.saturating_add(row_output_tokens); + total_cache_write_tokens = + total_cache_write_tokens.saturating_add(row_cache_write_tokens); + total_cache_read_tokens = total_cache_read_tokens.saturating_add(row_cache_read_tokens); + total_cost_microunits = total_cost_microunits.saturating_add(row_cost_microunits); + } + + let model_rows = sqlx::query( + "SELECT + COALESCE(resolved_model, requested_model, 'unknown') AS model, + COALESCE(provider_id, 'unknown') AS provider, + count(*)::bigint AS requests, + COALESCE(sum( + input_tokens + output_tokens + cache_write_tokens + cache_read_tokens + ), 0)::bigint AS tokens, + COALESCE(sum(cost_amount_microunits), 0)::bigint AS cost_microunits + FROM modelport_gateway_requests + WHERE state <> 'started' + AND traffic_class = 'business' + AND created_at >= to_timestamp($1::double precision / 1000.0) + AND created_at <= to_timestamp($2::double precision / 1000.0) + GROUP BY + COALESCE(resolved_model, requested_model, 'unknown'), + COALESCE(provider_id, 'unknown') + ORDER BY tokens DESC, requests DESC, model ASC + LIMIT 200", + ) + .bind(start_ms) + .bind(end_ms) + .fetch_all(pool) + .await?; + let model_usage = model_rows + .iter() + .map(|row| { + Ok(json!({ + "model": row.try_get::("model")?, + "provider": row.try_get::("provider")?, + "requests": nonnegative_u64(row.try_get("requests")?), + "tokens": nonnegative_u64(row.try_get("tokens")?), + "cost": microunits_usd(row.try_get("cost_microunits")?), + })) + }) + .collect::, sqlx::Error>>()?; + let request_time_series = dashboard_value_series(&requests, start_ms, bucket_ms); + let error_time_series = dashboard_value_series(&errors, start_ms, bucket_ms); + let token_time_series = (0..bucket_count) + .map(|index| { + let billed_input = input_tokens[index] + .saturating_add(cache_write_tokens[index]) + .saturating_add(cache_read_tokens[index]); + json!({ + "timestamp": dashboard_bucket_timestamp(start_ms, bucket_ms, index), + "inputTokens": input_tokens[index], + "outputTokens": output_tokens[index], + "cacheWriteTokens": cache_write_tokens[index], + "cacheReadTokens": cache_read_tokens[index], + "cacheHitRate": if billed_input == 0 { + 0.0 + } else { + cache_read_tokens[index] as f64 / billed_input as f64 * 100.0 + }, + }) + }) + .collect(); + let total_tokens = total_input_tokens + .saturating_add(total_output_tokens) + .saturating_add(total_cache_write_tokens) + .saturating_add(total_cache_read_tokens); + let minutes = (end_ms.saturating_sub(start_ms) as f64 / 60_000.0).max(1.0); + + Ok(Some(DashboardLedgerSnapshot { + usage_summary, + provider_usage, + matched_requests, + request_time_series, + error_time_series, + token_time_series, + model_usage, + summary: json!({ + "totalRequests": matched_requests, + "successRequests": success_requests, + "totalInputTokens": total_input_tokens, + "totalOutputTokens": total_output_tokens, + "totalCacheWriteTokens": total_cache_write_tokens, + "totalCacheReadTokens": total_cache_read_tokens, + "totalTokens": total_tokens, + "totalCostEstimate": microunits_usd(total_cost_microunits), + "rpm": matched_requests as f64 / minutes, + "tpm": total_tokens as f64 / minutes, + }), + })) + } + + pub(crate) async fn latency_stats_since( + &self, + since_ms: u64, + ) -> Result, AppError> { + let LedgerBackend::Postgres(pool) = self.backend.as_ref() else { + return Ok(None); + }; + let since_ms = i64::try_from(since_ms).unwrap_or(i64::MAX); + let overall = sqlx::query( + "SELECT + percentile_disc(0.50) WITHIN GROUP (ORDER BY latency_ms) AS p50, + percentile_disc(0.90) WITHIN GROUP (ORDER BY latency_ms) AS p90, + percentile_disc(0.95) WITHIN GROUP (ORDER BY latency_ms) AS p95, + percentile_disc(0.99) WITHIN GROUP (ORDER BY latency_ms) AS p99, + floor(COALESCE(avg(latency_ms), 0))::bigint AS avg, + COALESCE(max(latency_ms), 0)::bigint AS max, + count(*)::bigint AS count + FROM modelport_gateway_requests + WHERE state <> 'started' + AND created_at >= to_timestamp($1::double precision / 1000.0)", + ) + .bind(since_ms) + .fetch_one(pool) + .await?; + let by_model_rows = sqlx::query( + "SELECT + COALESCE(resolved_model, requested_model, 'unknown') AS name, + percentile_disc(0.50) WITHIN GROUP (ORDER BY latency_ms) AS p50, + percentile_disc(0.90) WITHIN GROUP (ORDER BY latency_ms) AS p90, + percentile_disc(0.95) WITHIN GROUP (ORDER BY latency_ms) AS p95, + percentile_disc(0.99) WITHIN GROUP (ORDER BY latency_ms) AS p99, + floor(COALESCE(avg(latency_ms), 0))::bigint AS avg, + COALESCE(max(latency_ms), 0)::bigint AS max, + count(*)::bigint AS count + FROM modelport_gateway_requests + WHERE state <> 'started' + AND created_at >= to_timestamp($1::double precision / 1000.0) + GROUP BY COALESCE(resolved_model, requested_model, 'unknown') + ORDER BY count DESC + LIMIT 200", + ) + .bind(since_ms) + .fetch_all(pool) + .await?; + let by_provider_rows = sqlx::query( + "SELECT + COALESCE(provider_id, 'unrouted') AS name, + percentile_disc(0.50) WITHIN GROUP (ORDER BY latency_ms) AS p50, + percentile_disc(0.90) WITHIN GROUP (ORDER BY latency_ms) AS p90, + percentile_disc(0.95) WITHIN GROUP (ORDER BY latency_ms) AS p95, + percentile_disc(0.99) WITHIN GROUP (ORDER BY latency_ms) AS p99, + floor(COALESCE(avg(latency_ms), 0))::bigint AS avg, + COALESCE(max(latency_ms), 0)::bigint AS max, + count(*)::bigint AS count + FROM modelport_gateway_requests + WHERE state <> 'started' + AND created_at >= to_timestamp($1::double precision / 1000.0) + GROUP BY COALESCE(provider_id, 'unrouted') + ORDER BY count DESC + LIMIT 200", + ) + .bind(since_ms) + .fetch_all(pool) + .await?; + let grouped = |rows: Vec| -> Result { + let mut values = serde_json::Map::new(); + for row in rows { + values.insert(row.try_get("name")?, latency_stats_from_pg(&row)?); + } + Ok(Value::Object(values)) + }; + + Ok(Some(json!({ + "p50": optional_nonnegative_u64(&overall, "p50")?, + "p90": optional_nonnegative_u64(&overall, "p90")?, + "p95": optional_nonnegative_u64(&overall, "p95")?, + "p99": optional_nonnegative_u64(&overall, "p99")?, + "avg": nonnegative_u64(overall.try_get("avg")?), + "max": nonnegative_u64(overall.try_get("max")?), + "byModel": grouped(by_model_rows)?, + "byProvider": grouped(by_provider_rows)?, + "sampleCount": nonnegative_u64(overall.try_get("count")?), + "percentilesEstimated": false, + }))) + } + + pub(crate) async fn usage_row(&self, ledger_id: &str) -> Result, AppError> { + let request = match self.backend.as_ref() { + LedgerBackend::Memory(ledger) => { + let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); + ledger.requests.get(ledger_id).map(|request| { + memory_request_row( + ledger_id, + request, + usize_to_i64( + ledger + .attempts + .values() + .filter(|attempt| attempt.request_ledger_id == ledger_id) + .count(), + ), + ) + }) + } + LedgerBackend::Postgres(pool) => sqlx::query(REQUEST_DETAIL_SQL) + .bind(ledger_id) + .fetch_optional(pool) + .await? + .as_ref() + .map(request_row_from_pg) + .transpose()?, + }; + Ok(request + .filter(|request| request.state != "started") + .as_ref() + .map(operational_log_row)) + } + + pub(crate) async fn management_usage(&self) -> Result { + match self.backend.as_ref() { + LedgerBackend::Memory(ledger) => { + let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); + let now = u64::try_from(now_millis()).unwrap_or(u64::MAX); + let day_start = current_period("daily", now).0; + let month_start = current_period("monthly", now).0; + let rolling_day_start = now.saturating_sub(24 * 60 * 60 * 1_000); + let mut stats = ManagementUsageStats::default(); + for request in ledger + .requests + .values() + .filter(|request| request.record.terminal) + { + let created_at = u64::try_from(request.record.created_at_ms).unwrap_or(0); + if created_at >= rolling_day_start { + let requests = stats + .users_24h + .entry(request.principal_id.clone()) + .or_default(); + *requests = requests.saturating_add(1); + } + if let Some(api_key_id) = request.api_key_id.as_deref() + && created_at >= day_start + { + let row = stats.api_keys.entry(api_key_id.to_owned()).or_default(); + row.requests_today = row.requests_today.saturating_add(1); + row.tokens_today = row + .tokens_today + .saturating_add(request_total_tokens(&request.record)); + } + if let Some(team_id) = request.team_id.as_deref() + && created_at >= month_start + { + let row = stats.teams.entry(team_id.to_owned()).or_default(); + let cost = request + .record + .billable_cost_microunits + .map_or(0.0, microunits_usd); + row.monthly_spend_usd += cost; + if created_at >= day_start { + row.requests_today = row.requests_today.saturating_add(1); + row.daily_spend_usd += cost; + } + } + } + Ok(stats) + } + LedgerBackend::Postgres(pool) => { + let api_key_rows = sqlx::query( + "SELECT + api_key_id, + count(*)::bigint AS requests_today, + COALESCE(sum( + input_tokens + output_tokens + + cache_write_tokens + cache_read_tokens + ), 0)::bigint AS tokens_today + FROM modelport_gateway_requests + WHERE state <> 'started' + AND api_key_id IS NOT NULL + AND created_at >= ( + date_trunc('day', now() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' + ) + GROUP BY api_key_id", + ) + .fetch_all(pool) + .await?; + let team_rows = sqlx::query( + "SELECT + team_id, + count(*) FILTER ( + WHERE created_at >= ( + date_trunc('day', now() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' + ) + )::bigint AS requests_today, + COALESCE(sum(billable_cost_microunits) FILTER ( + WHERE chargeable + AND created_at >= ( + date_trunc('day', now() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' + ) + ), 0)::bigint AS daily_spend_microunits, + COALESCE(sum(billable_cost_microunits) FILTER ( + WHERE chargeable + ), 0)::bigint AS monthly_spend_microunits + FROM modelport_gateway_requests + WHERE state <> 'started' + AND team_id IS NOT NULL + AND created_at >= ( + date_trunc('month', now() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' + ) + GROUP BY team_id", + ) + .fetch_all(pool) + .await?; + let user_rows = sqlx::query( + "SELECT principal_id, count(*)::bigint AS requests_24h + FROM modelport_gateway_requests + WHERE state <> 'started' + AND created_at >= now() - interval '24 hours' + GROUP BY principal_id", + ) + .fetch_all(pool) + .await?; + let mut stats = ManagementUsageStats::default(); + for row in api_key_rows { + stats.api_keys.insert( + row.try_get("api_key_id")?, + ApiKeyUsageStats { + requests_today: nonnegative_u64(row.try_get("requests_today")?), + tokens_today: nonnegative_u64(row.try_get("tokens_today")?), + }, + ); + } + for row in team_rows { + stats.teams.insert( + row.try_get("team_id")?, + TeamUsageStats { + requests_today: nonnegative_u64(row.try_get("requests_today")?), + daily_spend_usd: microunits_usd(row.try_get("daily_spend_microunits")?), + monthly_spend_usd: microunits_usd( + row.try_get("monthly_spend_microunits")?, + ), + }, + ); + } + for row in user_rows { + stats.users_24h.insert( + row.try_get("principal_id")?, + nonnegative_u64(row.try_get("requests_24h")?), + ); + } + Ok(stats) + } + } + } + + pub(crate) async fn request_detail( + &self, + ledger_id: &str, + ) -> Result, AppError> { + match self.backend.as_ref() { + LedgerBackend::Memory(ledger) => { + let ledger = ledger.lock().expect("enterprise ledger lock poisoned"); + let Some(request) = ledger.requests.get(ledger_id) else { + return Ok(None); + }; + let mut attempts = ledger + .attempts + .iter() + .filter(|(_, attempt)| attempt.request_ledger_id == ledger_id) + .map(|(attempt_id, attempt)| memory_attempt_row(attempt_id, attempt)) + .collect::>(); + attempts.sort_by_key(|attempt| attempt.created_at_ms); + Ok(Some(EnterpriseRequestDetail { + request: memory_request_row(ledger_id, request, usize_to_i64(attempts.len())), + attempts, + })) + } + LedgerBackend::Postgres(pool) => { + let Some(row) = sqlx::query(REQUEST_DETAIL_SQL) + .bind(ledger_id) + .fetch_optional(pool) + .await? + else { + return Ok(None); + }; + let request = request_row_from_pg(&row)?; + let attempt_rows = sqlx::query(ATTEMPT_LIST_SQL) + .bind(ledger_id) + .fetch_all(pool) + .await?; + Ok(Some(EnterpriseRequestDetail { + request, + attempts: attempt_rows + .iter() + .map(attempt_row_from_pg) + .collect::>()?, + })) + } + } + } +} From 13392b12d65ad1ddb8f245108c4428d260be8080 Mon Sep 17 00:00:00 2001 From: tiammomo <26957354+tiammomo@users.noreply.github.com> Date: Mon, 7 Sep 2026 12:39:51 +0800 Subject: [PATCH 2/4] refactor(dashboard): encapsulate provider credential management --- dashboard/e2e/provider-management.spec.ts | 19 +- .../src/features/models/ModelFormField.tsx | 19 + .../features/models/ProviderCredentials.tsx | 426 ++++++++++++++++ .../src/features/models/model-data.test.ts | 23 + dashboard/src/features/models/model-data.ts | 13 + .../src/features/models/operator-state.ts | 6 +- dashboard/src/lib/utils.ts | 6 + dashboard/src/pages/ModelsPage.tsx | 464 +----------------- 8 files changed, 520 insertions(+), 456 deletions(-) create mode 100644 dashboard/src/features/models/ProviderCredentials.tsx diff --git a/dashboard/e2e/provider-management.spec.ts b/dashboard/e2e/provider-management.spec.ts index b4fbb34..900f222 100644 --- a/dashboard/e2e/provider-management.spec.ts +++ b/dashboard/e2e/provider-management.spec.ts @@ -55,7 +55,7 @@ test.describe('provider management', () => { await expect(card.getByText('已禁用')).toHaveCount(0) }) - test('exposes credential pool controls on provider cards', async ({ page }) => { + test('validates, creates, edits and deletes credentials on provider cards', async ({ page }) => { const suffix = Date.now().toString(36) const providerId = 'deepseek' const credentialId = `e2e_pool_${suffix}` @@ -70,6 +70,9 @@ test.describe('provider management', () => { try { await card.getByRole('button', { name: '新增' }).click() const credentialDialog = page.getByRole('dialog') + await credentialDialog.getByRole('button', { name: '新增账号' }).click() + await expect(credentialDialog.locator('#credential-id')).toHaveAttribute('aria-invalid', 'true') + await expect(credentialDialog.locator('#credential-id')).toBeFocused() await credentialDialog.getByPlaceholder('例如: account-a').fill(credentialId) await credentialDialog.getByPlaceholder('例如: Mimo 主账号').fill('Pool Account A') await credentialDialog.getByPlaceholder('例如: MIMO_OPENAI_API_KEY_ALT').fill(`E2E_POOL_KEY_${suffix.toUpperCase()}`) @@ -83,6 +86,20 @@ test.describe('provider management', () => { await card.getByRole('combobox').first().click() await page.getByRole('option', { name: '轮询' }).click() await expect(card).toContainText('轮询') + + await card.getByRole('button', { name: '编辑上游账号 Pool Account A' }).click() + const editDialog = page.getByRole('dialog') + await expect(editDialog.locator('#credential-id')).toHaveCount(0) + await editDialog.locator('#credential-name').fill('Pool Account Updated') + await editDialog.getByRole('button', { name: '保存账号' }).click() + await expect(editDialog).not.toBeVisible() + await expect(card).toContainText('Pool Account Updated') + + await card.getByRole('button', { name: '删除上游账号 Pool Account Updated' }).click() + await page.getByRole('dialog').getByRole('button', { name: '删除账号' }).click() + await expect(page.getByRole('dialog')).not.toBeVisible() + await expect(card).not.toContainText('Pool Account Updated') + await expect(card).toContainText('默认凭证') } finally { await page.request.delete( `/admin/providers/${providerId}/credentials/${encodeURIComponent(credentialId)}`, diff --git a/dashboard/src/features/models/ModelFormField.tsx b/dashboard/src/features/models/ModelFormField.tsx index c6ab215..6cb0b58 100644 --- a/dashboard/src/features/models/ModelFormField.tsx +++ b/dashboard/src/features/models/ModelFormField.tsx @@ -1,3 +1,4 @@ +import { Switch } from '@/components/ui/switch' import type { ReactNode } from 'react' import { CircleAlert } from 'lucide-react' import { Label } from '@/components/ui/label' @@ -39,3 +40,21 @@ export function Field({ ) } +export function SwitchRow({ + label, + checked, + disabled, + onCheckedChange, +}: { + label: string + checked: boolean + disabled?: boolean + onCheckedChange: (checked: boolean) => void +}) { + return ( +
+ + +
+ ) +} diff --git a/dashboard/src/features/models/ProviderCredentials.tsx b/dashboard/src/features/models/ProviderCredentials.tsx new file mode 100644 index 0000000..e0aee4b --- /dev/null +++ b/dashboard/src/features/models/ProviderCredentials.tsx @@ -0,0 +1,426 @@ +import { Badge } from '@/components/ui/badge' +import { Button } from '@/components/ui/button' +import { Dialog, DialogContent, DialogDescription, DialogFooter, DialogHeader, DialogTitle } from '@/components/ui/dialog' +import { Input } from '@/components/ui/input' +import { Select, SelectContent, SelectItem, SelectTrigger, SelectValue } from '@/components/ui/select' +import { + useCreateProviderCredential, + useDeleteProviderCredential, + useSelectProviderCredential, + useUpdateProviderCredential, + useUpdateProviderCredentialPoolMode, +} from '@/hooks' +import { focusFirstInvalidDialogField, formatNumber, formatRelativeTime } from '@/lib/utils' +import type { Provider, ProviderCredential, ProviderCredentialPoolMode, ProviderOnlineBalance } from '@/types' +import { AlertTriangle, Loader2, Pencil, Plus, Trash2, WalletCards } from 'lucide-react' +import { useMemo, useState } from 'react' +import { toast } from 'sonner' +import { + CREDENTIAL_POOL_MODE_LABELS, + DEFAULT_CREDENTIAL_FORM, + credentialPayloadFromForm, + credentialToForm, + providerCredentialState, + providerDisplayTitle, + type ProviderCredentialFormState, +} from './model-data' +import { Field, SwitchRow } from './ModelFormField' +import { validateCredentialForm } from './operator-state' + +export function ProviderCredentials({ provider, canManage, onlineBalance, checkingBalance, onCheckBalance }: { + provider: Provider + canManage: boolean + onlineBalance?: ProviderOnlineBalance + checkingBalance: boolean + onCheckBalance: () => void +}) { + const createProviderCredential = useCreateProviderCredential() + const updateProviderCredential = useUpdateProviderCredential() + const selectProviderCredential = useSelectProviderCredential() + const updateProviderCredentialPoolMode = useUpdateProviderCredentialPoolMode() + const deleteProviderCredential = useDeleteProviderCredential() + const [credentialSubmitAttempted, setCredentialSubmitAttempted] = useState(false) + const [credentialDialogOpen, setCredentialDialogOpen] = useState(false) + const [editingCredential, setEditingCredential] = useState(null) + const [credentialForm, setCredentialForm] = useState(DEFAULT_CREDENTIAL_FORM) + const [credentialDeleteTarget, setCredentialDeleteTarget] = useState(null) + + const { credentials, activeCredential, credentialReady, credentialPoolMode } = providerCredentialState(provider) + const displayTitle = providerDisplayTitle(provider) + const credentialBusy = selectProviderCredential.isPending || updateProviderCredentialPoolMode.isPending || deleteProviderCredential.isPending + const credentialValidation = useMemo( + () => validateCredentialForm(credentialForm, !editingCredential), + [credentialForm, editingCredential], + ) + const openCredentialDialog = (credential?: ProviderCredential) => { + setCredentialDialogOpen(true) + setEditingCredential(credential ?? null) + setCredentialForm(credentialToForm(provider, credential)) + setCredentialSubmitAttempted(false) + } + + const closeCredentialDialog = () => { + setCredentialDialogOpen(false) + setEditingCredential(null) + setCredentialForm(DEFAULT_CREDENTIAL_FORM) + setCredentialSubmitAttempted(false) + } + + const handleSubmitCredential = () => { + if (!credentialDialogOpen) return + setCredentialSubmitAttempted(true) + if (!credentialValidation.valid) { + toast.error('请先修正账号表单中的错误') + focusFirstInvalidDialogField() + return + } + const data = credentialPayloadFromForm(credentialForm, !editingCredential) + const options = { + onSuccess: () => { + toast.success(editingCredential + ? '账号引用已更新;如环境变量值有变化,请重启进程并重新测试' + : '账号引用已新增;注入环境变量、重启进程并重新测试后才会生效') + closeCredentialDialog() + }, + onError: (error: unknown) => toast.error(error instanceof Error ? error.message : '保存账号失败'), + } + + if (editingCredential) { + updateProviderCredential.mutate({ + providerId: provider.id, + credentialId: editingCredential.id, + data, + }, options) + } else { + createProviderCredential.mutate({ + providerId: provider.id, + data, + }, options) + } + } + + const handleSelectProviderCredential = (credentialId: string) => { + selectProviderCredential.mutate({ providerId: provider.id, credentialId }, { + onSuccess: () => toast.success(`已切换 ${provider.displayName} 账号`), + onError: (error) => toast.error(error instanceof Error ? error.message : '切换账号失败'), + }) + } + + const handleUpdateProviderCredentialPoolMode = (mode: ProviderCredentialPoolMode) => { + updateProviderCredentialPoolMode.mutate({ providerId: provider.id, mode }, { + onSuccess: () => toast.success(`已更新 ${provider.displayName} 号池策略`), + onError: (error) => toast.error(error instanceof Error ? error.message : '更新号池策略失败'), + }) + } + + const handleDeleteProviderCredential = () => { + if (!credentialDeleteTarget) return + const credential = credentialDeleteTarget + deleteProviderCredential.mutate({ providerId: provider.id, credentialId: credential.id }, { + onSuccess: () => { + toast.success(`已删除账号 ${credential.name}`) + setCredentialDeleteTarget(null) + }, + onError: (error) => toast.error(error instanceof Error ? error.message : '删除账号失败'), + }) + } + + return ( + <> +
+
+
+

上游账号

+

+ {credentials.length > 0 ? `${credentials.length} 个账号 · ${CREDENTIAL_POOL_MODE_LABELS[credentialPoolMode]}` : '默认凭证'} +

+
+
+ + {canManage && } +
+
+ {provider.id === 'deepseek' && canManage && ( +
+
+
+
+

DeepSeek 线上余额

+ {onlineBalance && ( + + {onlineBalance.isAvailable ? '可调用' : '余额不足'} + + )} +
+

+ 实时只读查询;充值、退款与账单以 DeepSeek 控制台为准。 +

+
+ +
+ {onlineBalance && ( +
+ {onlineBalance.balanceInfos.map((balance) => ( +
+

{balance.currency} 可用总额

+

+ {balance.totalBalance} {balance.currency} +

+

+ 赠金 {balance.grantedBalance} · 充值 {balance.toppedUpBalance} +

+
+ ))} +

+ 最近查询:{formatRelativeTime(onlineBalance.checkedAt)} +

+
+ )} +
+ )} + {credentials.length === 0 ? ( +
+ + {credentialReady ? '默认环境变量可用' : '缺少默认密钥'} + + {provider.apiKeyEnv || '无需 API Key'} +
+ ) : ( +
+ +
+ {activeCredential && ( + <> + + {activeCredential.hasApiKey ? 'Key 可用' : 'Key 缺失'} + + {canManage && } + {canManage && } + + )} +
+ {activeCredential && ( +
+

+ 环境变量:{activeCredential.apiKeyEnv} +

+ {activeCredential.baseUrl && ( +

+ Base URL:{activeCredential.baseUrl} +

+ )} +
+ )} +
+ {credentials.map((credential) => { + const health = credential.health + const healthStatus = health?.status ?? (credential.hasApiKey ? 'healthy' : 'degraded') + const credentialRechargeBadge = health?.rechargeRequired ? '等待充值' : null + return ( +
+
+
+ {credential.name} + {credential.active && 当前} + {credential.status === 'disabled' && 禁用} +
+
+ {credential.apiKeyEnv} + {health?.lastUsedAt && 最近 {formatRelativeTime(health.lastUsedAt)}} +
+
+
+ + {credential.hasApiKey ? 'Key 可用' : 'Key 缺失'} + + + {credentialHealthLabel(healthStatus)} + + {credentialRechargeBadge && {credentialRechargeBadge}} + + {health?.requestsTotal ? `${formatNumber(health.requestsTotal)} 次 · ${Math.round(health.successRate)}%` : '暂无请求'} + +
+ {health?.lastError && ( +

{health.lastError}

+ )} +
+ ) + })} +
+
+ )} +
+ + { if (!open) closeCredentialDialog() }}> + + + {editingCredential ? '编辑上游账号' : '新增上游账号'} + + 账号只保存环境变量名;真实 API Key 仍放在 .env、容器 Secret 或系统环境变量中。保存后需重启并重新运行 Provider 连接测试。 + + +
+ {!editingCredential && ( + + setCredentialForm({ ...credentialForm, id: event.target.value.toLowerCase() })} + placeholder="例如: account-a" + aria-invalid={credentialSubmitAttempted && Boolean(credentialValidation.errors.id)} + aria-required="true" + /> + + )} + + setCredentialForm({ ...credentialForm, name: event.target.value })} + placeholder="例如: Mimo 主账号" + aria-invalid={credentialSubmitAttempted && Boolean(credentialValidation.errors.name)} + aria-required="true" + /> + + + setCredentialForm({ ...credentialForm, apiKeyEnv: event.target.value })} + placeholder="例如: MIMO_OPENAI_API_KEY_ALT" + aria-invalid={credentialSubmitAttempted && Boolean(credentialValidation.errors.apiKeyEnv)} + aria-required="true" + /> + + + setCredentialForm({ ...credentialForm, baseUrl: event.target.value })} + placeholder="可选,不填则沿用供应商 Base URL" + aria-invalid={credentialSubmitAttempted && Boolean(credentialValidation.errors.baseUrl)} + /> + +
+ setCredentialForm({ ...credentialForm, status: checked ? 'active' : 'disabled' })} + /> +
+ {credentialValidation.warnings.length > 0 && ( +
+ + {credentialValidation.warnings.join(' ')} +
+ )} +
+ + + + +
+
+ + { if (!open) setCredentialDeleteTarget(null) }}> + + + 删除上游账号 + 账号配置和健康记录会删除;真实环境变量不会被修改。 + +
+

{credentialDeleteTarget?.name}

+

{credentialDeleteTarget?.apiKeyEnv}

+ {credentialDeleteTarget?.active && ( +

这是当前账号;删除后系统会选择其他可用账号,若没有候选则 Provider 可能不可路由。

+ )} +
+ + + + +
+
+ + ) +} + +function credentialHealthLabel(status: string) { + if (status === 'cooldown') return '冷却' + if (status === 'degraded') return '降级' + return '健康' +} + +function credentialHealthVariant(status: string): 'success' | 'warning' { + if (status === 'cooldown' || status === 'degraded') return 'warning' + return 'success' +} + diff --git a/dashboard/src/features/models/model-data.test.ts b/dashboard/src/features/models/model-data.test.ts index 3dcffd7..e581c82 100644 --- a/dashboard/src/features/models/model-data.test.ts +++ b/dashboard/src/features/models/model-data.test.ts @@ -5,6 +5,7 @@ import { defaultToolStreamingArguments, parseList, providerInventoryItems, + providerCredentialState, providerOrigin, providerPayloadFromForm, providerToForm, @@ -128,3 +129,25 @@ describe('model feature data', () => { expect(dependencyLabel('apiKey')).toBe('API 密钥') }) }) + +describe('provider credential availability', () => { + const credential = { id: 'pool-a', providerId: 'openai', name: 'Account A', apiKeyEnv: 'POOL_A_KEY', status: 'active' as const, active: false, hasApiKey: true } + + it('recognizes a resolved pool credential when the default environment key is missing', () => { + const source = provider({ hasApiKey: false, credentials: [credential], activeCredentialId: credential.id }) + expect(providerCredentialState(source).credentialReady).toBe(true) + expect(providerCredentialState(source).activeCredential?.id).toBe(credential.id) + }) + + it('ignores disabled or unresolved pool credentials for readiness', () => { + expect(providerCredentialState(provider({ hasApiKey: false, credentials: [{ ...credential, status: 'disabled' }] })).credentialReady).toBe(false) + expect(providerCredentialState(provider({ hasApiKey: false, credentials: [{ ...credential, hasApiKey: false }] })).credentialReady).toBe(false) + expect(providerCredentialState(provider({ hasApiKey: false, apiKeyRequired: false })).credentialReady).toBe(true) + }) + + it('keeps the runtime active credential ahead of the configured fallback', () => { + const active = { ...credential, id: 'pool-b', active: true } + const source = provider({ credentials: [credential, active], activeCredentialId: credential.id }) + expect(providerCredentialState(source).activeCredential?.id).toBe(active.id) + }) +}) diff --git a/dashboard/src/features/models/model-data.ts b/dashboard/src/features/models/model-data.ts index 3d1348e..e4d4d82 100644 --- a/dashboard/src/features/models/model-data.ts +++ b/dashboard/src/features/models/model-data.ts @@ -506,3 +506,16 @@ export function dependencyLabel(type: string): string { if (type === 'route') return '路由配置' return type } + +export function providerCredentialState(provider: Provider) { + const credentials = provider.credentials ?? [] + return { + credentials, + credentialReady: provider.hasApiKey || !provider.apiKeyRequired + || credentials.some((credential) => credential.status === 'active' && credential.hasApiKey), + activeCredential: credentials.find((credential) => credential.active) + ?? credentials.find((credential) => credential.id === provider.activeCredentialId) + ?? null, + credentialPoolMode: provider.credentialPoolMode ?? 'failover', + } +} diff --git a/dashboard/src/features/models/operator-state.ts b/dashboard/src/features/models/operator-state.ts index bbdd157..1f39963 100644 --- a/dashboard/src/features/models/operator-state.ts +++ b/dashboard/src/features/models/operator-state.ts @@ -1,5 +1,5 @@ import type { Provider } from '@/types' -import { parseList, type ProviderCredentialFormState, type ProviderFormState } from './model-data' +import { parseList, providerCredentialState, type ProviderCredentialFormState, type ProviderFormState } from './model-data' export type ProviderReadinessLevel = 'ready' | 'attention' | 'blocked' | 'disabled' @@ -32,9 +32,7 @@ export function providerReadiness(provider: Provider, isDefault = false): Provid } } - const credentialReady = provider.hasApiKey - || !provider.apiKeyRequired - || Boolean(provider.credentials?.some((credential) => credential.status === 'active' && credential.hasApiKey)) + const { credentialReady } = providerCredentialState(provider) if (!credentialReady) { return { level: 'blocked', diff --git a/dashboard/src/lib/utils.ts b/dashboard/src/lib/utils.ts index 29788a7..e0f0f6c 100644 --- a/dashboard/src/lib/utils.ts +++ b/dashboard/src/lib/utils.ts @@ -106,3 +106,9 @@ export function debounce unknown>( timer = setTimeout(() => fn(...args), ms) } } + +export function focusFirstInvalidDialogField() { + window.requestAnimationFrame(() => { + document.querySelector('[role="dialog"] [aria-invalid="true"]')?.focus() + }) +} diff --git a/dashboard/src/pages/ModelsPage.tsx b/dashboard/src/pages/ModelsPage.tsx index 58e969d..8b95485 100644 --- a/dashboard/src/pages/ModelsPage.tsx +++ b/dashboard/src/pages/ModelsPage.tsx @@ -18,18 +18,15 @@ import { Table, TableBody, TableCell, TableHead, TableHeader, TableRow } from '@ import { Tabs, TabsContent, TabsList, TabsTrigger } from '@/components/ui/tabs' import { apiKeyExpiryState } from '@/features/api-keys/api-key-view' import { - CREDENTIAL_POOL_MODE_LABELS, - DEFAULT_CREDENTIAL_FORM, DEFAULT_PROVIDER_FORM, PROVIDER_OPERATIONAL_FILTERS, - credentialPayloadFromForm, - credentialToForm, defaultToolStreamingArguments, defaultToolUseForProviderForm, dependencyLabel, modelRouteTitle, providerDeleteBlockedFromError, providerDisplayTitle, + providerCredentialState, providerFilterCount, providerIdentity, providerInventoryGroups, @@ -41,17 +38,16 @@ import { providerPayloadFromForm, providerRuntimeState, providerToForm, - type ProviderCredentialFormState, type ProviderFormState, type ProviderInventoryGroup, type ProviderOperationalFilter, } from '@/features/models/model-data' import { ModelAdaptationDialog } from '@/features/models/ModelAdaptationDialog' -import { Field } from '@/features/models/ModelFormField' +import { Field, SwitchRow } from '@/features/models/ModelFormField' +import { ProviderCredentials } from '@/features/models/ProviderCredentials' import { providerReadiness, validateAliasForm, - validateCredentialForm, validateProviderForm, type ProviderReadinessLevel, } from '@/features/models/operator-state' @@ -63,22 +59,17 @@ import { useCheckProviderBalance, useCreateAlias, useCreateProvider, - useCreateProviderCredential, useDeleteAlias, useDeleteProvider, - useDeleteProviderCredential, useDiscoverProviderModels, useNow, useProviders, - useSelectProviderCredential, useSetProviderDisabled, useSettings, useToggleModel, useUpdateDefaultModel, useUpdateDefaultProvider, useUpdateProvider, - useUpdateProviderCredential, - useUpdateProviderCredentialPoolMode, useUpdateProviderOrder, } from '@/hooks' import { PROVIDER_PROTOCOL_LABELS } from '@/lib/constants' @@ -91,14 +82,12 @@ import { type ProviderTemplate, } from '@/lib/model-catalog' import { paginateItems } from '@/lib/pagination' -import { cn, copyToClipboard, formatNumber, formatRelativeTime } from '@/lib/utils' +import { cn, copyToClipboard, focusFirstInvalidDialogField, formatNumber, formatRelativeTime } from '@/lib/utils' import { useAuthStore } from '@/stores' import type { FidelityMode, MaxTokensField, Provider, -ProviderCredential, -ProviderCredentialPoolMode, ProviderDeleteBlocked, ProviderModelInventory, ProviderOnlineBalance, @@ -129,7 +118,6 @@ import { Search, Settings, Trash2, - WalletCards, } from 'lucide-react' import { Fragment, useMemo, useState } from 'react' import { Link } from 'react-router-dom' @@ -199,11 +187,6 @@ export function ModelsPage() { const createProvider = useCreateProvider() const updateProvider = useUpdateProvider() const setProviderDisabled = useSetProviderDisabled() - const createProviderCredential = useCreateProviderCredential() - const updateProviderCredential = useUpdateProviderCredential() - const selectProviderCredential = useSelectProviderCredential() - const updateProviderCredentialPoolMode = useUpdateProviderCredentialPoolMode() - const deleteProviderCredential = useDeleteProviderCredential() const deleteProvider = useDeleteProvider() const toggleModel = useToggleModel() const bulkToggleModels = useBulkToggleModels() @@ -220,17 +203,13 @@ export function ModelsPage() { const [showProviderDialog, setShowProviderDialog] = useState(false) const [aliasSubmitAttempted, setAliasSubmitAttempted] = useState(false) const [providerSubmitAttempted, setProviderSubmitAttempted] = useState(false) - const [credentialSubmitAttempted, setCredentialSubmitAttempted] = useState(false) - const [credentialDialogProvider, setCredentialDialogProvider] = useState(null) const [selectedTemplate, setSelectedTemplate] = useState(null) const [editingProvider, setEditingProvider] = useState(null) - const [editingCredential, setEditingCredential] = useState(null) const [editingModelAdaptation, setEditingModelAdaptation] = useState<{ provider: Provider item: ProviderModelInventory } | null>(null) const [providerForm, setProviderForm] = useState(DEFAULT_PROVIDER_FORM) - const [credentialForm, setCredentialForm] = useState(DEFAULT_CREDENTIAL_FORM) const [deleteTarget, setDeleteTarget] = useState(null) const [deleteBlock, setDeleteBlock] = useState(null) const [deleteConfirmation, setDeleteConfirmation] = useState('') @@ -244,11 +223,6 @@ export function ModelsPage() { const [aliasPageSize, setAliasPageSize] = useState(20) const [activeTab, setActiveTab] = useState('library') const [aliasDeleteTarget, setAliasDeleteTarget] = useState(null) - const [credentialDeleteTarget, setCredentialDeleteTarget] = useState<{ - provider: Provider - credential: ProviderCredential - } | null>(null) - const configuredProviderIds = useMemo(() => new Set(providers.map((provider) => provider.id)), [providers]) const defaultProvider = settings?.gateway.defaultProvider.trim() ?? '' const providerOrder = useMemo( @@ -354,10 +328,6 @@ export function ModelsPage() { } : null const providerValidation = useMemo(() => validateProviderForm(providerForm), [providerForm]) - const credentialValidation = useMemo( - () => validateCredentialForm(credentialForm, !editingCredential), - [credentialForm, editingCredential], - ) const aliasValidation = useMemo( () => validateAliasForm(aliasForm.alias, aliasForm.target), [aliasForm.alias, aliasForm.target], @@ -413,20 +383,6 @@ export function ModelsPage() { setProviderSubmitAttempted(false) } - const openCredentialDialog = (provider: Provider, credential?: ProviderCredential) => { - setCredentialDialogProvider(provider) - setEditingCredential(credential ?? null) - setCredentialForm(credentialToForm(provider, credential)) - setCredentialSubmitAttempted(false) - } - - const closeCredentialDialog = () => { - setCredentialDialogProvider(null) - setEditingCredential(null) - setCredentialForm(DEFAULT_CREDENTIAL_FORM) - setCredentialSubmitAttempted(false) - } - const handleSubmitProvider = () => { setProviderSubmitAttempted(true) if (!providerValidation.valid) { @@ -452,39 +408,6 @@ export function ModelsPage() { } } - const handleSubmitCredential = () => { - if (!credentialDialogProvider) return - setCredentialSubmitAttempted(true) - if (!credentialValidation.valid) { - toast.error('请先修正账号表单中的错误') - focusFirstInvalidDialogField() - return - } - const data = credentialPayloadFromForm(credentialForm, !editingCredential) - const options = { - onSuccess: () => { - toast.success(editingCredential - ? '账号引用已更新;如环境变量值有变化,请重启进程并重新测试' - : '账号引用已新增;注入环境变量、重启进程并重新测试后才会生效') - closeCredentialDialog() - }, - onError: (error: unknown) => toast.error(error instanceof Error ? error.message : '保存账号失败'), - } - - if (editingCredential) { - updateProviderCredential.mutate({ - providerId: credentialDialogProvider.id, - credentialId: editingCredential.id, - data, - }, options) - } else { - createProviderCredential.mutate({ - providerId: credentialDialogProvider.id, - data, - }, options) - } - } - const handleSetProviderDisabled = (provider: Provider) => { const disabled = provider.status !== 'disabled' setProviderDisabled.mutate({ providerId: provider.id, disabled }, { @@ -493,32 +416,6 @@ export function ModelsPage() { }) } - const handleSelectProviderCredential = (provider: Provider, credentialId: string) => { - selectProviderCredential.mutate({ providerId: provider.id, credentialId }, { - onSuccess: () => toast.success(`已切换 ${provider.displayName} 账号`), - onError: (error) => toast.error(error instanceof Error ? error.message : '切换账号失败'), - }) - } - - const handleUpdateProviderCredentialPoolMode = (provider: Provider, mode: ProviderCredentialPoolMode) => { - updateProviderCredentialPoolMode.mutate({ providerId: provider.id, mode }, { - onSuccess: () => toast.success(`已更新 ${provider.displayName} 号池策略`), - onError: (error) => toast.error(error instanceof Error ? error.message : '更新号池策略失败'), - }) - } - - const handleDeleteProviderCredential = () => { - if (!credentialDeleteTarget) return - const { provider, credential } = credentialDeleteTarget - deleteProviderCredential.mutate({ providerId: provider.id, credentialId: credential.id }, { - onSuccess: () => { - toast.success(`已删除账号 ${credential.name}`) - setCredentialDeleteTarget(null) - }, - onError: (error) => toast.error(error instanceof Error ? error.message : '删除账号失败'), - }) - } - const handleDeleteProvider = (force = false) => { if (!deleteTarget) return deleteProvider.mutate({ providerId: deleteTarget.id, force }, { @@ -1081,18 +978,12 @@ export function ModelsPage() { }} onCopy={copyText} onAlias={openAliasDialog} - onCreateCredential={() => openCredentialDialog(provider)} - onEditCredential={(credential) => openCredentialDialog(provider, credential)} - onSelectCredential={(credentialId) => handleSelectProviderCredential(provider, credentialId)} - onUpdateCredentialPoolMode={(mode) => handleUpdateProviderCredentialPoolMode(provider, mode)} - onDeleteCredential={(credential) => setCredentialDeleteTarget({ provider, credential })} onToggleModel={(model, enabled) => handleToggleProviderModel(provider, model, enabled)} onBulkToggleModels={(enabled) => handleBulkToggleProviderModels(provider, enabled)} onSetDefaultModel={(model) => handleSetDefaultModel(provider, model)} onEditModel={(item) => setEditingModelAdaptation({ provider, item })} modelMutationKey={modelMutationKey} bulkModelMutation={bulkModelMutation} - credentialBusy={selectProviderCredential.isPending || updateProviderCredentialPoolMode.isPending || deleteProviderCredential.isPending} defaultModelMutationKey={defaultModelMutationKey} /> ))} @@ -1712,87 +1603,6 @@ export function ModelsPage() { onClose={() => setEditingModelAdaptation(null)} />} - { if (!open) closeCredentialDialog() }}> - - - {editingCredential ? '编辑上游账号' : '新增上游账号'} - - 账号只保存环境变量名;真实 API Key 仍放在 .env、容器 Secret 或系统环境变量中。保存后需重启并重新运行 Provider 连接测试。 - - -
- {!editingCredential && ( - - setCredentialForm({ ...credentialForm, id: event.target.value.toLowerCase() })} - placeholder="例如: account-a" - aria-invalid={credentialSubmitAttempted && Boolean(credentialValidation.errors.id)} - aria-required="true" - /> - - )} - - setCredentialForm({ ...credentialForm, name: event.target.value })} - placeholder="例如: Mimo 主账号" - aria-invalid={credentialSubmitAttempted && Boolean(credentialValidation.errors.name)} - aria-required="true" - /> - - - setCredentialForm({ ...credentialForm, apiKeyEnv: event.target.value })} - placeholder="例如: MIMO_OPENAI_API_KEY_ALT" - aria-invalid={credentialSubmitAttempted && Boolean(credentialValidation.errors.apiKeyEnv)} - aria-required="true" - /> - - - setCredentialForm({ ...credentialForm, baseUrl: event.target.value })} - placeholder="可选,不填则沿用供应商 Base URL" - aria-invalid={credentialSubmitAttempted && Boolean(credentialValidation.errors.baseUrl)} - /> - -
- setCredentialForm({ ...credentialForm, status: checked ? 'active' : 'disabled' })} - /> -
- {credentialValidation.warnings.length > 0 && ( -
- - {credentialValidation.warnings.join(' ')} -
- )} -
- - - - -
-
- { if (!open) { setDeleteTarget(null) @@ -1869,29 +1679,6 @@ export function ModelsPage() { - { if (!open) setCredentialDeleteTarget(null) }}> - - - 删除上游账号 - 账号配置和健康记录会删除;真实环境变量不会被修改。 - -
-

{credentialDeleteTarget?.credential.name}

-

{credentialDeleteTarget?.credential.apiKeyEnv}

- {credentialDeleteTarget?.credential.active && ( -

这是当前账号;删除后系统会选择其他可用账号,若没有候选则 Provider 可能不可路由。

- )} -
- - - - -
-
- { if (!open) setAliasDeleteTarget(null) }}> @@ -1991,7 +1778,7 @@ function ProviderRoutingOverview({ onOpenRouting: () => void }) { const credentialReady = defaultProvider - ? defaultProvider.hasApiKey || !defaultProvider.apiKeyRequired + ? providerCredentialState(defaultProvider).credentialReady : false return ( @@ -2126,18 +1913,12 @@ function ProviderCard({ onDelete, onCopy, onAlias, - onCreateCredential, - onEditCredential, - onSelectCredential, - onUpdateCredentialPoolMode, - onDeleteCredential, onToggleModel, onBulkToggleModels, onSetDefaultModel, onEditModel, modelMutationKey, bulkModelMutation, - credentialBusy, defaultModelMutationKey, }: { provider: Provider @@ -2156,24 +1937,15 @@ function ProviderCard({ onDelete: () => void onCopy: (value: string) => Promise onAlias: (alias?: string, target?: string) => void - onCreateCredential: () => void - onEditCredential: (credential: ProviderCredential) => void - onSelectCredential: (credentialId: string) => void - onUpdateCredentialPoolMode: (mode: ProviderCredentialPoolMode) => void - onDeleteCredential: (credential: ProviderCredential) => void onToggleModel: (model: string, enabled: boolean) => void onBulkToggleModels: (enabled: boolean) => void onSetDefaultModel: (model: string) => void onEditModel: (item: ProviderModelInventory) => void modelMutationKey: string | null bulkModelMutation: { providerId: string; enabled: boolean } | null - credentialBusy: boolean defaultModelMutationKey: string | null }) { - const credentials = provider.credentials ?? [] - const credentialReady = provider.hasApiKey - || !provider.apiKeyRequired - || credentials.some((credential) => credential.status === 'active' && credential.hasApiKey) + const { credentials, credentialReady, activeCredential } = providerCredentialState(provider) const lastTest = provider.lastTest const connectionVerified = lastTest?.success === true const routeReady = provider.status === 'active' @@ -2186,10 +1958,6 @@ function ProviderCard({ const modelListId = `provider-models-${provider.id}` const identity = providerIdentity(provider) const displayTitle = providerDisplayTitle(provider) - const activeCredential = credentials.find((credential) => credential.active) - ?? credentials.find((credential) => credential.id === provider.activeCredentialId) - ?? null - const credentialPoolMode = provider.credentialPoolMode ?? 'failover' const modelGroups = providerModelGroups(provider) const inventoryGroups = providerInventoryGroups(provider) const inventoryItems = providerInventoryItems(provider) @@ -2290,183 +2058,13 @@ function ProviderCard({ )} -
-
-
-

上游账号

-

- {credentials.length > 0 ? `${credentials.length} 个账号 · ${CREDENTIAL_POOL_MODE_LABELS[credentialPoolMode]}` : '默认凭证'} -

-
-
- - {canManage && } -
-
- {provider.id === 'deepseek' && canManage && ( -
-
-
-
-

DeepSeek 线上余额

- {onlineBalance && ( - - {onlineBalance.isAvailable ? '可调用' : '余额不足'} - - )} -
-

- 实时只读查询;充值、退款与账单以 DeepSeek 控制台为准。 -

-
- -
- {onlineBalance && ( -
- {onlineBalance.balanceInfos.map((balance) => ( -
-

{balance.currency} 可用总额

-

- {balance.totalBalance} {balance.currency} -

-

- 赠金 {balance.grantedBalance} · 充值 {balance.toppedUpBalance} -

-
- ))} -

- 最近查询:{formatRelativeTime(onlineBalance.checkedAt)} -

-
- )} -
- )} - {credentials.length === 0 ? ( -
- - {credentialReady ? '默认环境变量可用' : '缺少默认密钥'} - - {provider.apiKeyEnv || '无需 API Key'} -
- ) : ( -
- -
- {activeCredential && ( - <> - - {activeCredential.hasApiKey ? 'Key 可用' : 'Key 缺失'} - - {canManage && } - {canManage && } - - )} -
- {activeCredential && ( -
-

- 环境变量:{activeCredential.apiKeyEnv} -

- {activeCredential.baseUrl && ( -

- Base URL:{activeCredential.baseUrl} -

- )} -
- )} -
- {credentials.map((credential) => { - const health = credential.health - const healthStatus = health?.status ?? (credential.hasApiKey ? 'healthy' : 'degraded') - const credentialRechargeBadge = health?.rechargeRequired ? '等待充值' : null - return ( -
-
-
- {credential.name} - {credential.active && 当前} - {credential.status === 'disabled' && 禁用} -
-
- {credential.apiKeyEnv} - {health?.lastUsedAt && 最近 {formatRelativeTime(health.lastUsedAt)}} -
-
-
- - {credential.hasApiKey ? 'Key 可用' : 'Key 缺失'} - - - {credentialHealthLabel(healthStatus)} - - {credentialRechargeBadge && {credentialRechargeBadge}} - - {health?.requestsTotal ? `${formatNumber(health.requestsTotal)} 次 · ${Math.round(health.successRate)}%` : '暂无请求'} - -
- {health?.lastError && ( -

{health.lastError}

- )} -
- ) - })} -
-
- )} -
+
{canManage &&