Files
Genarrative/server-rs/crates/api-server/src/telemetry.rs
T
kdletters 5bf036bb81
Project CI / AI game creator shell Rust shard 4/4 (push) Has been cancelled
Project CI / AI game creator shell Rust shard 3/4 (push) Has been cancelled
Project CI / AI game creator shell Rust shard 2/4 (push) Has been cancelled
Project CI / AI game creator shell Rust smoke (push) Has been cancelled
Project CI / AI game creator shell Rust crates (push) Has been cancelled
Project CI / Backend tests (push) Has been cancelled
Project CI / Native shell tests (push) Has been cancelled
Project CI / Frontend tests (push) Has been cancelled
Project CI / Repository checks (push) Has been cancelled
Project CI / AI game creator shell web tests (push) Has been cancelled
Project CI / AI game creator shell Rust shard 1/4 (push) Has been cancelled
收敛后端依赖装配与鉴权并完善异步追踪 (#425)
HTTP 请求取消后,在途计数原先无法释放;项目元数据和 External API 鉴权依赖完整 AppState,相同鉴权和追踪配置也散落在多个入口。本次集中装配这些依赖和横切能力,保持现有公开 API、权限、计费、幂等和事务规则。

## 修改

- 91 个受保护路由集中应用鉴权;保留方法级 404/405/HEAD/Allow、公开入口、MCP、精选缓存头及 2/4 MiB 请求限制。
- 七个项目元数据入口改用缓存的 EditorProjectState,External/MCP 鉴权改用 ExternalApiAuthState;生产实现复用 SpacetimeClient,媒体修复维持原有上传和登记顺序。
- RAII 覆盖请求 Future 取消和 panic unwind 的计数清理;正常与降级服务复用 TraceLayer,指标采用 MatchedPath 模板及固定兜底。
- LLM 普通与流式调用增加跳过参数的异步 span,保持父上下文、流式回调、错误与重试行为;补充替代依赖测试并同步锁文件和文档。

## 验证

| 验证面 | 结果 |
| --- | --- |
| platform-llm 完整本地回归 | 161 个单元测试、3 个集成测试通过;1 个真实 Provider 用例按原配置忽略 |
| api-server 完整回归 | 执行时 1075 通过、12 失败、6 忽略;其中 1 个新增公开读取 fixture 断言已修正,14 个路由契约回归随后全部通过;剩余 11 个是下述既有 Windows 失败 |
| 窄依赖、取消与追踪 | 元数据 owner/幂等/revision、鉴权及 MCP 错误传播、取消/panic/流式响应、追踪父子关系与敏感参数省略均通过 |
| 实际本地服务 | 独立 SpacetimeDB 上 102/102 检查通过,两个动态项目 ID 的路由模板及请求 ID 日志核验 3/3 通过 |
| 编译与边界 | api-server cargo check、AGC 锁文件下 platform-llm cargo check、rustfmt、编码、文档索引、DDD 与 diff 检查通过 |
| 合入最新 master 后 | 后端源码及锁文件保持已测内容;再次通过 14 个路由契约测试、3 个 Provider 追踪测试及编码/文档/DDD/diff 检查 |

实际服务检查覆盖 health/ready、两账号登录、项目 CRUD、幂等重复、跨 owner 拒绝、revision 冲突、External/MCP 读取、Key 撤销及 404/405。使用既有 test 环境的本地 Router 拒绝 fixture,未调用真实付费 Provider;自建服务已关闭,原开发实例保留。

## 已知测试限制

API 全量测试尚未全绿:11 个 wallet_refund_outbox 用例在 Windows 的目录同步处失败。其生产文件与变更前内容一致;标准库隔离复现确认 File::open(目录) 返回 OS 5,而普通文件写入、同步及 hard_link 正常。这个已有的目录持久化问题未混入本次重构,也未通过跳过或弱化相关断言掩盖。

---------

Co-authored-by: kdletters <61648117+kdletters@users.noreply.github.com>
Reviewed-on: http://192.168.35.82/git/GenarrativeAI/Genarrative/pulls/425
2026-09-19 12:29:19 +08:00

767 lines
26 KiB
Rust

use axum::{
body::Body,
extract::{MatchedPath, State},
http::{HeaderMap, Request, Response},
middleware::Next,
};
use http_body_util::BodyExt;
use opentelemetry::{KeyValue, global, metrics::Counter};
use std::sync::{
Arc, OnceLock,
atomic::{AtomicI64, Ordering},
};
use tracing::{info, warn};
use crate::{
request_context::resolve_request_id,
state::{AppState, HttpRequestPermitPoolKind},
};
static HTTP_RESPONSE_BODY_IN_FLIGHT: AtomicI64 = AtomicI64::new(0);
static TRACKING_OUTBOX_PENDING_BYTES: AtomicI64 = AtomicI64::new(0);
static TRACKING_OUTBOX_PENDING_FILES: AtomicI64 = AtomicI64::new(0);
static HTTP_REQUEST_PERMITS_AVAILABLE: OnceLock<HttpRequestPermitsAvailableGauges> =
OnceLock::new();
// 集中维护 api-server HTTP 观测,避免在 handler 中散落高基数字段或重复创建 instrument。
pub async fn record_http_observability(
State(state): State<AppState>,
request: Request<Body>,
next: Next,
) -> Response<Body> {
observe_http_request(
http_metrics(),
state.config.slow_request_threshold_ms,
request,
next,
)
.await
}
async fn observe_http_request(
metrics: &HttpMetrics,
slow_request_threshold_ms: u64,
request: Request<Body>,
next: Next,
) -> Response<Body> {
let method = request.method().as_str().to_string();
let route = observability_route(&request);
let scheme = resolve_request_scheme(request.headers());
let path = request.uri().path().to_string();
let request_id = resolve_request_id(&request).unwrap_or_else(|| "unknown".to_string());
let base_labels = http_base_labels(method.clone(), route.clone());
let in_flight = RequestInFlightGuard::new(&metrics.in_flight, base_labels.clone());
let started_at = std::time::Instant::now();
let response = next.run(request).await;
let status = response.status().as_u16();
let status_class = status_class(status);
let latency_ms = started_at.elapsed().as_millis().min(u64::MAX as u128) as u64;
let slow_request = latency_ms >= slow_request_threshold_ms;
let labels = http_response_labels(base_labels, status);
metrics.requests.add(1, &labels);
metrics
.duration
.record(started_at.elapsed().as_secs_f64(), &labels);
drop(in_flight);
if slow_request {
warn!(
request_id = %request_id,
http.request.method = %method,
http.route = %route,
url.scheme = %scheme,
url.path = %path,
http.response.status_code = status,
status,
status_class,
latency_ms,
slow_request = true,
"http request completed slowly"
);
} else {
info!(
request_id = %request_id,
http.request.method = %method,
http.route = %route,
url.scheme = %scheme,
url.path = %path,
http.response.status_code = status,
status,
status_class,
latency_ms,
slow_request = false,
"http request completed"
);
}
track_response_body_in_flight(response)
}
pub(crate) fn update_http_request_permits_available(
pool: HttpRequestPermitPoolKind,
available: usize,
) {
HTTP_REQUEST_PERMITS_AVAILABLE
.get_or_init(register_http_request_permits_available_metric)
.store(pool, available);
}
#[cfg(any())]
pub(crate) fn record_puzzle_gallery_cache_hit() {
puzzle_gallery_cache_metrics().hits.add(1, &[]);
}
#[cfg(any())]
pub(crate) fn record_puzzle_gallery_cache_miss() {
puzzle_gallery_cache_metrics().misses.add(1, &[]);
}
#[cfg(any())]
pub(crate) fn record_puzzle_gallery_cache_rebuild(
duration: std::time::Duration,
data_bytes: usize,
) {
let metrics = puzzle_gallery_cache_metrics();
metrics.rebuilds.add(1, &[]);
metrics.rebuild_duration.record(duration.as_secs_f64(), &[]);
metrics
.data_json_bytes
.record(data_bytes.min(u64::MAX as usize) as u64, &[]);
}
pub(crate) fn record_tracking_outbox_enqueued() {
tracking_outbox_metrics().enqueued.add(1, &[]);
}
pub(crate) fn record_tracking_outbox_dropped(reason: &'static str) {
tracking_outbox_metrics()
.dropped
.add(1, &[KeyValue::new("reason", reason)]);
}
pub(crate) fn record_tracking_outbox_sealed(reason: &'static str) {
tracking_outbox_metrics()
.sealed_files
.add(1, &[KeyValue::new("reason", reason)]);
}
pub(crate) fn record_tracking_outbox_corrupt_file() {
tracking_outbox_metrics().corrupt_files.add(1, &[]);
}
pub(crate) fn record_tracking_outbox_flush(
duration: std::time::Duration,
accepted_count: u32,
file_bytes: u64,
failed: bool,
) {
let status_class = if failed { "error" } else { "ok" };
let labels = [KeyValue::new("status_class", status_class)];
let metrics = tracking_outbox_metrics();
metrics.flushes.add(1, &labels);
metrics
.flush_duration
.record(duration.as_secs_f64(), &labels);
metrics
.flushed_events
.add(u64::from(accepted_count), &labels);
metrics.flushed_bytes.add(file_bytes, &labels);
}
pub(crate) fn update_tracking_outbox_pending_bytes(bytes: u64) {
TRACKING_OUTBOX_PENDING_BYTES.store(bytes.min(i64::MAX as u64) as i64, Ordering::Relaxed);
}
pub(crate) fn update_tracking_outbox_pending_files(files: usize) {
TRACKING_OUTBOX_PENDING_FILES.store(files.min(i64::MAX as usize) as i64, Ordering::Relaxed);
}
pub(crate) fn record_external_api_failure(
provider: &'static str,
failure_stage: &'static str,
status_class: &'static str,
retryable: bool,
) {
external_api_metrics().failures.add(
1,
&[
KeyValue::new("provider", provider),
KeyValue::new("failure_stage", failure_stage),
KeyValue::new("status_class", status_class),
KeyValue::new("retryable", retryable),
],
);
}
pub(crate) fn record_external_api_audit_dropped(provider: &'static str, reason: &'static str) {
external_api_metrics().audit_dropped.add(
1,
&[
KeyValue::new("provider", provider),
KeyValue::new("reason", reason),
],
);
}
fn track_response_body_in_flight(response: Response<Body>) -> Response<Body> {
response.map(|body| {
HTTP_RESPONSE_BODY_IN_FLIGHT.fetch_add(1, Ordering::Relaxed);
let guard = ResponseBodyInFlightGuard;
Body::new(body.map_frame(move |frame| {
let _guard = &guard;
frame
}))
})
}
#[derive(Clone)]
struct HttpMetrics {
requests: Counter<u64>,
in_flight: opentelemetry::metrics::UpDownCounter<i64>,
duration: opentelemetry::metrics::Histogram<f64>,
}
// 请求 Future 被取消或 panic unwind 时也必须释放计数;响应体存活由另一 guard 统计。
struct RequestInFlightGuard<'a> {
counter: &'a opentelemetry::metrics::UpDownCounter<i64>,
labels: Vec<KeyValue>,
}
impl<'a> RequestInFlightGuard<'a> {
fn new(counter: &'a opentelemetry::metrics::UpDownCounter<i64>, labels: Vec<KeyValue>) -> Self {
counter.add(1, &labels);
Self { counter, labels }
}
}
impl Drop for RequestInFlightGuard<'_> {
fn drop(&mut self) {
self.counter.add(-1, &self.labels);
}
}
#[cfg(any())]
struct PuzzleGalleryCacheMetrics {
hits: Counter<u64>,
misses: Counter<u64>,
rebuilds: Counter<u64>,
rebuild_duration: opentelemetry::metrics::Histogram<f64>,
data_json_bytes: opentelemetry::metrics::Histogram<u64>,
}
struct TrackingOutboxMetrics {
enqueued: Counter<u64>,
dropped: Counter<u64>,
sealed_files: Counter<u64>,
corrupt_files: Counter<u64>,
flushes: Counter<u64>,
flush_duration: opentelemetry::metrics::Histogram<f64>,
flushed_events: Counter<u64>,
flushed_bytes: Counter<u64>,
}
struct ExternalApiMetrics {
failures: Counter<u64>,
audit_dropped: Counter<u64>,
}
struct HttpRequestPermitsAvailableGauges {
default: Arc<AtomicI64>,
admin: Arc<AtomicI64>,
}
impl HttpRequestPermitsAvailableGauges {
fn new() -> Self {
Self {
default: Arc::new(AtomicI64::new(0)),
admin: Arc::new(AtomicI64::new(0)),
}
}
fn store(&self, pool: HttpRequestPermitPoolKind, available: usize) {
let value = available.min(i64::MAX as usize) as i64;
match pool {
HttpRequestPermitPoolKind::Default => &self.default,
HttpRequestPermitPoolKind::Admin => &self.admin,
}
.store(value, Ordering::Relaxed);
}
}
struct ResponseBodyInFlightGuard;
impl Drop for ResponseBodyInFlightGuard {
fn drop(&mut self) {
HTTP_RESPONSE_BODY_IN_FLIGHT.fetch_sub(1, Ordering::Relaxed);
}
}
fn http_metrics() -> &'static HttpMetrics {
static METRICS: std::sync::OnceLock<HttpMetrics> = std::sync::OnceLock::new();
METRICS.get_or_init(|| {
let meter = global::meter("genarrative-api");
HttpMetrics {
requests: meter
.u64_counter("genarrative.http.server.requests")
.with_description("HTTP request count grouped by route and status class")
.build(),
in_flight: meter
.i64_up_down_counter("http.server.active_requests")
.with_unit("{request}")
.with_description("Number of active HTTP server requests")
.build(),
duration: meter
.f64_histogram("http.server.request.duration")
.with_unit("s")
.with_description("Duration of HTTP server requests")
.build(),
}
})
}
#[cfg(any())]
fn puzzle_gallery_cache_metrics() -> &'static PuzzleGalleryCacheMetrics {
static METRICS: std::sync::OnceLock<PuzzleGalleryCacheMetrics> = std::sync::OnceLock::new();
METRICS.get_or_init(|| {
let meter = global::meter("genarrative-api");
PuzzleGalleryCacheMetrics {
hits: meter
.u64_counter("genarrative.puzzle_gallery.cache.hits")
.with_description("Puzzle gallery response cache hits")
.build(),
misses: meter
.u64_counter("genarrative.puzzle_gallery.cache.misses")
.with_description("Puzzle gallery response cache misses")
.build(),
rebuilds: meter
.u64_counter("genarrative.puzzle_gallery.cache.rebuilds")
.with_description("Puzzle gallery response cache rebuild count")
.build(),
rebuild_duration: meter
.f64_histogram("genarrative.puzzle_gallery.cache.rebuild.duration")
.with_unit("s")
.with_description("Puzzle gallery response cache rebuild duration")
.build(),
data_json_bytes: meter
.u64_histogram("genarrative.puzzle_gallery.cache.data_json_bytes")
.with_unit("By")
.with_description("Serialized puzzle gallery data JSON size")
.build(),
}
})
}
fn tracking_outbox_metrics() -> &'static TrackingOutboxMetrics {
static METRICS: std::sync::OnceLock<TrackingOutboxMetrics> = std::sync::OnceLock::new();
METRICS.get_or_init(|| {
let meter = global::meter("genarrative-api");
TrackingOutboxMetrics {
enqueued: meter
.u64_counter("genarrative.tracking_outbox.events.enqueued")
.with_description("Tracking events appended to the local outbox")
.build(),
dropped: meter
.u64_counter("genarrative.tracking_outbox.events.dropped")
.with_description("Tracking events dropped by local outbox protection")
.build(),
sealed_files: meter
.u64_counter("genarrative.tracking_outbox.files.sealed")
.with_description("Tracking outbox active files sealed for flushing")
.build(),
corrupt_files: meter
.u64_counter("genarrative.tracking_outbox.files.corrupt")
.with_description(
"Tracking outbox sealed files quarantined because they could not be parsed",
)
.build(),
flushes: meter
.u64_counter("genarrative.tracking_outbox.flushes")
.with_description("Tracking outbox sealed file flush attempts")
.build(),
flush_duration: meter
.f64_histogram("genarrative.tracking_outbox.flush.duration")
.with_unit("s")
.with_description("Tracking outbox sealed file flush duration")
.build(),
flushed_events: meter
.u64_counter("genarrative.tracking_outbox.events.flushed")
.with_description("Tracking events accepted by SpacetimeDB batch procedure")
.build(),
flushed_bytes: meter
.u64_counter("genarrative.tracking_outbox.bytes.flushed")
.with_unit("By")
.with_description("Tracking outbox bytes removed after successful flush")
.build(),
}
})
}
fn external_api_metrics() -> &'static ExternalApiMetrics {
static METRICS: std::sync::OnceLock<ExternalApiMetrics> = std::sync::OnceLock::new();
METRICS.get_or_init(|| {
let meter = global::meter("genarrative-api");
ExternalApiMetrics {
failures: meter
.u64_counter("genarrative.external_api.failures")
.with_description(
"External API call failures grouped by provider and failure stage",
)
.build(),
audit_dropped: meter
.u64_counter("genarrative.external_api.audit.dropped")
.with_description(
"External API failure audit records dropped when synchronous fallback is disabled",
)
.build(),
}
})
}
fn register_http_request_permits_available_metric() -> HttpRequestPermitsAvailableGauges {
let gauges = HttpRequestPermitsAvailableGauges::new();
let meter = global::meter("genarrative-api");
let default_gauge = gauges.default.clone();
let admin_gauge = gauges.admin.clone();
meter
.i64_observable_up_down_counter("genarrative.http.server.request_permits.available")
.with_unit("{permit}")
.with_description("Available api-server HTTP backpressure permits")
.with_callback(move |observer| {
observer.observe(
default_gauge.load(Ordering::Relaxed),
&[KeyValue::new(
"pool",
HttpRequestPermitPoolKind::Default.as_str(),
)],
);
observer.observe(
admin_gauge.load(Ordering::Relaxed),
&[KeyValue::new(
"pool",
HttpRequestPermitPoolKind::Admin.as_str(),
)],
);
})
.build();
gauges
}
pub(crate) fn register_http_runtime_metrics() {
static REGISTERED: OnceLock<()> = OnceLock::new();
REGISTERED.get_or_init(|| {
let meter = global::meter("genarrative-api");
meter
.i64_observable_up_down_counter("genarrative.http.server.response_bodies.in_flight")
.with_unit("{response}")
.with_description("HTTP response bodies still owned by Axum/Hyper")
.with_callback(|observer| {
observer.observe(HTTP_RESPONSE_BODY_IN_FLIGHT.load(Ordering::Relaxed), &[]);
})
.build();
meter
.i64_observable_up_down_counter("genarrative.tracking_outbox.pending.bytes")
.with_unit("By")
.with_description("Tracking outbox bytes waiting on local disk")
.with_callback(|observer| {
observer.observe(TRACKING_OUTBOX_PENDING_BYTES.load(Ordering::Relaxed), &[]);
})
.build();
meter
.i64_observable_up_down_counter("genarrative.tracking_outbox.pending.files")
.with_unit("{file}")
.with_description("Tracking outbox sealed files waiting for flush")
.with_callback(|observer| {
observer.observe(TRACKING_OUTBOX_PENDING_FILES.load(Ordering::Relaxed), &[]);
})
.build();
});
}
fn http_base_labels(method: String, route: String) -> Vec<KeyValue> {
vec![
KeyValue::new("http.request.method", method),
KeyValue::new("http.route", route),
]
}
fn http_response_labels(mut labels: Vec<KeyValue>, status: u16) -> Vec<KeyValue> {
labels.push(KeyValue::new("status_class", status_class(status)));
labels
}
fn status_class(status: u16) -> &'static str {
match status {
100..=199 => "1xx",
200..=299 => "2xx",
300..=399 => "3xx",
400..=499 => "4xx",
500..=599 => "5xx",
_ => "unknown",
}
}
pub(crate) fn observability_route<B>(request: &Request<B>) -> String {
if let Some(path) = request.extensions().get::<MatchedPath>() {
return path.as_str().to_string();
}
let path = request.uri().path();
if path.starts_with("/admin/api/") {
"/admin/api/*".to_string()
} else if path.starts_with("/api/") {
"/api/*".to_string()
} else {
"other".to_string()
}
}
pub(crate) fn resolve_request_scheme(headers: &HeaderMap) -> String {
headers
.get("x-forwarded-proto")
.and_then(|value| value.to_str().ok())
.and_then(|value| value.split(',').next())
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or("http")
.to_string()
}
#[cfg(test)]
mod tests {
use axum::{
Router,
body::{Body, Bytes},
http::{HeaderMap, HeaderValue, Request, StatusCode},
middleware,
routing::get,
};
use http_body_util::BodyExt;
use opentelemetry::{
KeyValue,
metrics::{Counter, Histogram, SyncInstrument, UpDownCounter},
};
use std::{
convert::Infallible,
sync::{Arc, Mutex},
time::Duration,
};
use tokio::sync::Notify;
use tower::ServiceExt;
use super::{HttpMetrics, observability_route, observe_http_request, resolve_request_scheme};
#[derive(Default)]
struct Measurements<T> {
values: Mutex<Vec<(T, Vec<KeyValue>)>>,
}
impl<T: Send> SyncInstrument<T> for Measurements<T> {
fn measure(&self, value: T, attributes: &[KeyValue]) {
self.values
.lock()
.expect("measurement lock")
.push((value, attributes.to_vec()));
}
}
fn observed_router(router: Router) -> (Router, Arc<Measurements<i64>>) {
let in_flight = Arc::new(Measurements::default());
let metrics = HttpMetrics {
requests: Counter::new(Arc::new(Measurements::<u64>::default())),
in_flight: UpDownCounter::new(in_flight.clone()),
duration: Histogram::new(Arc::new(Measurements::<f64>::default())),
};
let router = router.layer(middleware::from_fn(move |request, next| {
let metrics = metrics.clone();
async move { observe_http_request(&metrics, u64::MAX, request, next).await }
}));
(router, in_flight)
}
fn request(uri: &str) -> Request<Body> {
Request::builder()
.uri(uri)
.body(Body::empty())
.expect("request")
}
fn assert_request_finished(measurements: &Measurements<i64>) {
let values = measurements.values.lock().expect("measurement lock");
assert_eq!(
values.iter().map(|(value, _)| *value).collect::<Vec<_>>(),
vec![1, -1]
);
assert_eq!(
values[0].1, values[1].1,
"decrement must use the original labels"
);
}
#[tokio::test]
async fn cancelled_request_releases_in_flight_measurement() {
let entered = Arc::new(Notify::new());
let handler_entered = entered.clone();
let (router, measurements) = observed_router(Router::new().route(
"/api/pending/{id}",
get(move || {
let entered = handler_entered.clone();
async move {
entered.notify_one();
std::future::pending::<StatusCode>().await
}
}),
));
let task = tokio::spawn(router.oneshot(request("/api/pending/123")));
tokio::time::timeout(Duration::from_secs(5), entered.notified())
.await
.expect("handler entered");
assert_eq!(
measurements.values.lock().expect("measurement lock")[0].0,
1
);
task.abort();
assert!(
task.await
.expect_err("request should be cancelled")
.is_cancelled()
);
assert_request_finished(&measurements);
}
#[tokio::test]
async fn panicking_request_releases_in_flight_measurement() {
async fn panic_handler() -> StatusCode {
panic!("handler panic for cancellation cleanup test");
}
let (router, measurements) =
observed_router(Router::new().route("/panic", get(panic_handler)));
let task = tokio::spawn(router.oneshot(request("/panic")));
assert!(task.await.expect_err("handler should panic").is_panic());
assert_request_finished(&measurements);
}
#[tokio::test]
async fn responses_release_in_flight_once_for_success_and_errors() {
for status in [
StatusCode::OK,
StatusCode::UNAUTHORIZED,
StatusCode::INTERNAL_SERVER_ERROR,
] {
let (router, measurements) = observed_router(
Router::new().route("/response", get(move || async move { status })),
);
let response = router
.oneshot(request("/response"))
.await
.expect("response");
assert_eq!(response.status(), status);
assert_request_finished(&measurements);
drop(response);
assert_request_finished(&measurements);
}
}
#[tokio::test]
async fn streaming_response_releases_request_before_body_completion() {
use futures_util::StreamExt;
let (router, measurements) = observed_router(Router::new().route(
"/events",
get(|| async {
let stream = futures_util::stream::iter([Ok::<_, Infallible>(Bytes::from_static(
b"data: ready\n\n",
))])
.chain(futures_util::stream::pending());
(
[("content-type", "text/event-stream")],
Body::from_stream(stream),
)
}),
));
let mut response = router
.oneshot(request("/events"))
.await
.expect("streaming response");
assert_request_finished(&measurements);
let frame = response
.body_mut()
.frame()
.await
.expect("first frame")
.expect("body frame");
assert_eq!(frame.into_data().expect("data frame"), "data: ready\n\n");
drop(response);
assert_request_finished(&measurements);
}
#[tokio::test]
async fn matched_routes_share_templates_and_preserve_distinct_endpoints() {
let routes = Router::new().nest(
"/api",
Router::new()
.route("/projects/{project_id}", get(|| async {}))
.route("/assets/{asset_id}", get(|| async {})),
);
let (router, measurements) = observed_router(routes);
for (uri, template) in [
(
"/api/projects/project-1?cursor=private",
"/api/projects/{project_id}",
),
("/api/projects/project-2", "/api/projects/{project_id}"),
("/api/assets/asset-1", "/api/assets/{asset_id}"),
("/api/missing/private-id?token=private", "/api/*"),
] {
let response = router
.clone()
.oneshot(request(uri))
.await
.expect("response");
drop(response);
assert_request_finished(&measurements);
let mut values = measurements.values.lock().expect("measurement lock");
assert_eq!(
values[0].1,
vec![
KeyValue::new("http.request.method", "GET"),
KeyValue::new("http.route", template)
]
);
values.clear();
}
}
#[test]
fn observability_route_keeps_metrics_labels_low_cardinality() {
assert_eq!(
observability_route(&request("/api/editor/showcase/resources?cursor=abc")),
"/api/*"
);
assert_eq!(
observability_route(&request("/api/editor/projects/project-1")),
"/api/*"
);
assert_eq!(
observability_route(&request("/api/runtime/settings")),
"/api/*"
);
assert_eq!(
observability_route(&request("/admin/api/debug/http")),
"/admin/api/*"
);
assert_eq!(
observability_route(&request("/missing/private-id")),
"other"
);
}
#[test]
fn resolve_request_scheme_uses_forwarded_proto_first_value() {
let mut headers = HeaderMap::new();
headers.insert("x-forwarded-proto", HeaderValue::from_static("https, http"));
assert_eq!(resolve_request_scheme(&headers), "https");
}
}