Use header-only usage server timing

This commit is contained in:
Novick Yuan
2026-05-28 14:30:15 +08:00
parent 01c8592ca6
commit 6412294262
14 changed files with 425 additions and 75 deletions
@@ -9,9 +9,9 @@ use aether_admin::observability::usage::{
admin_usage_data_unavailable_response, admin_usage_has_fallback, admin_usage_is_failed,
admin_usage_matches_search, admin_usage_matches_username, admin_usage_parse_ids,
admin_usage_parse_limit, admin_usage_parse_offset, admin_usage_provider_key_name,
admin_usage_record_json, build_admin_usage_active_requests_response,
build_admin_usage_records_response, build_admin_usage_summary_stats_response_from_summary,
usage_server_now_unix_ms, ADMIN_USAGE_DATA_UNAVAILABLE_DETAIL,
admin_usage_record_json, attach_usage_server_now_header,
build_admin_usage_active_requests_response, build_admin_usage_records_response,
build_admin_usage_summary_stats_response_from_summary, ADMIN_USAGE_DATA_UNAVAILABLE_DETAIL,
};
use aether_data::repository::users::StoredUserSummary;
use aether_data_contracts::repository::{
@@ -314,14 +314,15 @@ fn build_admin_usage_records_response_with_attempt_flags(
})
.collect();
Json(json!({
"server_now_unix_ms": usage_server_now_unix_ms(),
"records": records,
"total": total,
"limit": limit,
"offset": offset,
}))
.into_response()
attach_usage_server_now_header(
Json(json!({
"records": records,
"total": total,
"limit": limit,
"offset": offset,
}))
.into_response(),
)
}
fn build_admin_usage_records_query(
@@ -4,6 +4,7 @@ use aether_ai_serving::UPSTREAM_IS_STREAM_KEY;
use aether_billing::{
normalize_input_tokens_for_billing, normalize_total_input_context_for_cache_hit_rate,
};
use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
use aether_data_contracts::repository::usage::{
StoredRequestUsageAudit, StoredUsageBreakdownSummaryRow, StoredUsageDailySummary,
UsageAuditKeywordSearchQuery, UsageAuditListQuery, UsageBreakdownGroupBy,
@@ -33,6 +34,17 @@ fn users_me_usage_server_now_unix_ms() -> u64 {
u64::try_from(Utc::now().timestamp_millis()).unwrap_or_default()
}
fn attach_users_me_usage_server_now_header(mut response: Response<Body>) -> Response<Body> {
if let Ok(value) = http::HeaderValue::from_str(&users_me_usage_server_now_unix_ms().to_string())
{
response.headers_mut().insert(
http::HeaderName::from_static(USAGE_SERVER_NOW_UNIX_MS_HEADER),
value,
);
}
response
}
fn build_users_me_usage_reader_unavailable_response() -> Response<Body> {
build_auth_error_response(
http::StatusCode::SERVICE_UNAVAILABLE,
@@ -1075,7 +1087,6 @@ pub(super) async fn handle_users_me_usage_get(
.flatten();
let mut payload = json!({
"server_now_unix_ms": users_me_usage_server_now_unix_ms(),
"total_requests": total_requests,
"total_input_tokens": total_input_tokens,
"total_output_tokens": total_output_tokens,
@@ -1099,7 +1110,7 @@ pub(super) async fn handle_users_me_usage_get(
&summary_by_provider
));
}
Json(payload).into_response()
attach_users_me_usage_server_now_header(Json(payload).into_response())
}
pub(super) async fn handle_users_me_usage_active_get(
@@ -1174,14 +1185,15 @@ pub(super) async fn handle_users_me_usage_active_get(
.collect::<Vec<_>>()
};
Json(json!({
"server_now_unix_ms": users_me_usage_server_now_unix_ms(),
"requests": items
.iter()
.map(build_users_me_usage_active_payload)
.collect::<Vec<_>>(),
}))
.into_response()
attach_users_me_usage_server_now_header(
Json(json!({
"requests": items
.iter()
.map(build_users_me_usage_active_payload)
.collect::<Vec<_>>(),
}))
.into_response(),
)
}
pub(super) async fn handle_users_me_usage_interval_timeline_get(
@@ -1370,14 +1382,21 @@ async fn build_usage_heatmap_summaries(
mod tests {
use std::collections::BTreeMap;
use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
use aether_data_contracts::repository::usage::StoredRequestUsageAudit;
use axum::{
body::Body,
response::{IntoResponse, Response},
Json,
};
use chrono::Utc;
use serde_json::json;
use super::{
build_users_me_usage_active_payload, build_users_me_usage_record_payload,
users_me_usage_client_is_stream, users_me_usage_is_failed,
users_me_usage_server_now_unix_ms, users_me_usage_upstream_is_stream,
attach_users_me_usage_server_now_header, build_users_me_usage_active_payload,
build_users_me_usage_record_payload, users_me_usage_client_is_stream,
users_me_usage_is_failed, users_me_usage_server_now_unix_ms,
users_me_usage_upstream_is_stream,
};
fn sample_usage(status: &str) -> StoredRequestUsageAudit {
@@ -1433,6 +1452,28 @@ mod tests {
assert!(value > 1_000_000_000_000);
}
fn assert_users_me_usage_server_now_header(response: &Response<Body>) {
let value = response
.headers()
.get(USAGE_SERVER_NOW_UNIX_MS_HEADER)
.expect("user usage response should include server now header")
.to_str()
.expect("server now header should be valid ASCII")
.parse::<u64>()
.expect("server now header should be epoch millis");
assert!(value > 1_000_000_000_000);
}
#[test]
fn users_me_usage_server_now_header_is_added_to_response() {
let response = attach_users_me_usage_server_now_header(
Json(json!({ "requests": [] })).into_response(),
);
assert_users_me_usage_server_now_header(&response);
}
#[test]
fn user_usage_record_payload_rehydrates_cache_creation_total_from_classified_fields() {
let item = StoredRequestUsageAudit {
@@ -1,3 +1,4 @@
use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
use axum::body::Body;
use axum::extract::{Request, State};
use axum::http::{self, HeaderValue, Response};
@@ -6,6 +7,8 @@ use axum::middleware::Next;
use crate::headers::header_value_str;
use crate::state::{AppState, FrontdoorCorsConfig};
const FRONTDOOR_CREDENTIALS_EXPOSE_HEADERS: &str = "*, x-aether-server-now-unix-ms";
fn append_vary(headers: &mut http::HeaderMap, value: &'static str) {
headers.append(http::header::VARY, HeaderValue::from_static(value));
}
@@ -31,7 +34,11 @@ fn apply_frontdoor_cors_headers(
);
headers.insert(
http::header::ACCESS_CONTROL_EXPOSE_HEADERS,
HeaderValue::from_static("*"),
HeaderValue::from_static(if cors.allow_credentials() {
FRONTDOOR_CREDENTIALS_EXPOSE_HEADERS
} else {
"*"
}),
);
if let Some(value) = requested_headers {
if let Ok(value) = HeaderValue::from_str(value) {
@@ -109,3 +116,52 @@ pub(crate) async fn frontdoor_cors_middleware(
);
response
}
#[cfg(test)]
mod tests {
use super::*;
fn assert_exposes_header(value: &HeaderValue, expected: &str) {
let exposed_headers = value
.to_str()
.expect("expose headers should be valid ASCII");
assert!(
exposed_headers
.split(',')
.map(str::trim)
.any(|header| header.eq_ignore_ascii_case(expected)),
"{exposed_headers} should include {expected}"
);
}
#[test]
fn frontdoor_cors_explicitly_exposes_usage_server_time_for_credentials() {
let cors = FrontdoorCorsConfig::new(vec!["http://localhost:5173".to_string()], true)
.expect("cors config should build");
let mut headers = http::HeaderMap::new();
apply_frontdoor_cors_headers(&mut headers, &cors, "http://localhost:5173", None);
let expose_headers = headers
.get(http::header::ACCESS_CONTROL_EXPOSE_HEADERS)
.expect("expose headers should be set");
assert_exposes_header(expose_headers, "*");
assert_exposes_header(expose_headers, USAGE_SERVER_NOW_UNIX_MS_HEADER);
}
#[test]
fn frontdoor_cors_keeps_wildcard_expose_headers_without_credentials() {
let cors = FrontdoorCorsConfig::new(vec!["http://localhost:5173".to_string()], false)
.expect("cors config should build");
let mut headers = http::HeaderMap::new();
apply_frontdoor_cors_headers(&mut headers, &cors, "http://localhost:5173", None);
assert_eq!(
headers
.get(http::header::ACCESS_CONTROL_EXPOSE_HEADERS)
.expect("expose headers should be set"),
"*"
);
}
}
+14 -6
View File
@@ -7,6 +7,7 @@ use crate::tests::{
build_state_with_execution_runtime_override, json, start_server, AppState, Arc, Body,
FrontdoorCorsConfig, Mutex, Request, Router, StatusCode, FRONTDOOR_MANIFEST_PATH, READYZ_PATH,
};
use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
use aether_crypto::DEVELOPMENT_ENCRYPTION_KEY;
use aether_data::repository::auth::InMemoryAuthApiKeySnapshotRepository;
use aether_data::repository::candidate_selection::InMemoryMinimalCandidateSelectionReadRepository;
@@ -440,12 +441,19 @@ async fn gateway_adds_cors_headers_to_proxied_responses() {
.expect("allow origin header"),
"http://localhost:3000"
);
assert_eq!(
response_headers
.get("access-control-expose-headers")
.expect("expose headers header"),
"*"
);
let expose_headers = response_headers
.get("access-control-expose-headers")
.expect("expose headers header")
.to_str()
.expect("expose headers should be valid ASCII");
assert!(expose_headers
.split(',')
.map(str::trim)
.any(|header| header == "*"));
assert!(expose_headers
.split(',')
.map(str::trim)
.any(|header| header.eq_ignore_ascii_case(USAGE_SERVER_NOW_UNIX_MS_HEADER)));
assert_eq!(
*execution_runtime_hits.lock().expect("mutex should lock"),
1
+72 -15
View File
@@ -19,12 +19,24 @@ use serde_json::{json, Value};
use std::collections::{BTreeMap, BTreeSet};
use url::form_urlencoded;
pub use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
pub const ADMIN_USAGE_DATA_UNAVAILABLE_DETAIL: &str = "Admin usage data unavailable";
pub fn usage_server_now_unix_ms() -> u64 {
u64::try_from(chrono::Utc::now().timestamp_millis()).unwrap_or_default()
}
pub fn attach_usage_server_now_header(mut response: Response<Body>) -> Response<Body> {
if let Ok(value) = http::HeaderValue::from_str(&usage_server_now_unix_ms().to_string()) {
response.headers_mut().insert(
http::HeaderName::from_static(USAGE_SERVER_NOW_UNIX_MS_HEADER),
value,
);
}
response
}
pub fn admin_usage_data_unavailable_response(detail: &'static str) -> Response<Body> {
(
http::StatusCode::SERVICE_UNAVAILABLE,
@@ -2279,11 +2291,7 @@ pub fn build_admin_usage_active_requests_response(
})
.collect();
Json(json!({
"server_now_unix_ms": usage_server_now_unix_ms(),
"requests": payload,
}))
.into_response()
attach_usage_server_now_header(Json(json!({ "requests": payload })).into_response())
}
#[allow(clippy::too_many_arguments)]
@@ -2313,14 +2321,15 @@ pub fn build_admin_usage_records_response(
})
.collect();
Json(json!({
"server_now_unix_ms": usage_server_now_unix_ms(),
"records": records,
"total": total,
"limit": limit,
"offset": offset,
}))
.into_response()
attach_usage_server_now_header(
Json(json!({
"records": records,
"total": total,
"limit": limit,
"offset": offset,
}))
.into_response(),
)
}
pub fn build_admin_usage_curl_response(
@@ -2512,6 +2521,7 @@ pub fn build_admin_usage_replay_plan_response(
mod tests {
use std::collections::BTreeMap;
use axum::{body::Body, response::Response};
use serde_json::json;
use super::{
@@ -2519,8 +2529,10 @@ mod tests {
admin_usage_has_fallback, admin_usage_is_failed, admin_usage_is_success,
admin_usage_matches_search, admin_usage_matches_status, admin_usage_matches_username,
admin_usage_record_json, admin_usage_resolve_request_capture_body,
admin_usage_total_tokens, admin_usage_upstream_is_stream, build_admin_usage_detail_payload,
usage_server_now_unix_ms,
admin_usage_total_tokens, admin_usage_upstream_is_stream,
build_admin_usage_active_requests_response, build_admin_usage_detail_payload,
build_admin_usage_records_response, usage_server_now_unix_ms,
USAGE_SERVER_NOW_UNIX_MS_HEADER,
};
use aether_data_contracts::repository::usage::{StoredRequestUsageAudit, UsageBodyField};
@@ -2581,6 +2593,51 @@ mod tests {
assert!(value > 1_000_000_000_000);
}
fn assert_usage_server_now_header(response: &Response<Body>) {
let value = response
.headers()
.get(USAGE_SERVER_NOW_UNIX_MS_HEADER)
.expect("usage response should include server now header")
.to_str()
.expect("server now header should be valid ASCII")
.parse::<u64>()
.expect("server now header should be epoch millis");
assert!(value > 1_000_000_000_000);
}
#[test]
fn admin_usage_records_response_sets_server_now_header() {
let item = sample_usage("completed", Some(200), None);
let response = build_admin_usage_records_response(
&[item],
&BTreeMap::new(),
&BTreeMap::new(),
true,
true,
&BTreeMap::new(),
1,
20,
0,
);
assert_usage_server_now_header(&response);
}
#[test]
fn admin_usage_active_response_sets_server_now_header() {
let item = sample_usage("streaming", Some(200), None);
let response = build_admin_usage_active_requests_response(
&[item],
&BTreeMap::new(),
true,
&BTreeMap::new(),
&BTreeMap::new(),
);
assert_usage_server_now_header(&response);
}
#[test]
fn explicit_completed_status_wins_over_legacy_failure_fields() {
let item = sample_usage(
+3 -1
View File
@@ -16,4 +16,6 @@ pub use plan::{
TRANSPORT_HTTP_MODE_HTTP1_ONLY, TRANSPORT_POOL_SCOPE_KEY,
};
pub use result::{ExecutionResult, ExecutionTelemetry, ResponseBody};
pub use usage::{ExecutionStreamTerminalSummary, StandardizedUsage};
pub use usage::{
ExecutionStreamTerminalSummary, StandardizedUsage, USAGE_SERVER_NOW_UNIX_MS_HEADER,
};
+2
View File
@@ -2,6 +2,8 @@ use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
pub const USAGE_SERVER_NOW_UNIX_MS_HEADER: &str = "x-aether-server-now-unix-ms";
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
pub struct StandardizedUsage {
pub input_tokens: i64,
@@ -0,0 +1,62 @@
import { describe, expect, it } from 'vitest'
import {
SERVER_NOW_UNIX_MS_HEADER,
buildServerTimingMetadata,
readServerNowUnixMsFromHeaders,
withServerTiming,
} from '../serverTiming'
describe('serverTiming', () => {
it('reads server time from response headers', () => {
expect(readServerNowUnixMsFromHeaders({
[SERVER_NOW_UNIX_MS_HEADER]: '1779999000123',
})).toBe(1_779_999_000_123)
expect(readServerNowUnixMsFromHeaders({
'X-Aether-Server-Now-Unix-Ms': '1779999000456',
})).toBe(1_779_999_000_456)
})
it('does not fall back to body fields', () => {
const timing = buildServerTimingMetadata({
headers: {},
data: {
server_now_unix_ms: 1_779_999_000_123,
},
}, 1_000, 1_100)
expect(timing).toBeUndefined()
})
it('builds metadata with round trip duration', () => {
const timing = buildServerTimingMetadata({
headers: {
[SERVER_NOW_UNIX_MS_HEADER]: '1050',
},
}, 1_000, 1_125)
expect(timing).toEqual({
server_now_unix_ms: 1_050,
client_send_unix_ms: 1_000,
client_receive_unix_ms: 1_125,
round_trip_ms: 125,
})
})
it('returns the original payload when the header is missing or invalid', () => {
const payload = { records: [] }
expect(withServerTiming({ data: payload, headers: {} }, 1_000)).toBe(payload)
expect(withServerTiming({
data: payload,
headers: { [SERVER_NOW_UNIX_MS_HEADER]: 'not-a-number' },
}, 1_000)).toBe(payload)
expect(withServerTiming({
data: payload,
headers: { [SERVER_NOW_UNIX_MS_HEADER]: '0' },
}, 1_000)).toBe(payload)
expect(withServerTiming({
data: payload,
headers: { [SERVER_NOW_UNIX_MS_HEADER]: '1050.5' },
}, 1_000)).toBe(payload)
})
})
@@ -132,8 +132,10 @@ describe('usageApi contract alignment', () => {
}
if (url === '/api/admin/usage/records') {
return Promise.resolve({
headers: {
'x-aether-server-now-unix-ms': '10050',
},
data: {
server_now_unix_ms: 10_050,
records: [{ id: 'record-3' }],
total: 1,
limit: 25,
@@ -163,6 +165,7 @@ describe('usageApi contract alignment', () => {
server_now_unix_ms: 10_050,
client_send_unix_ms: 1_000,
client_receive_unix_ms: 1_100,
round_trip_ms: 100,
})
})
+2 -3
View File
@@ -148,7 +148,6 @@ export interface ApiFormatSummary {
// 使用统计响应接口
export interface UsageResponse extends ServerTimedPayload {
server_now_unix_ms?: number
total_requests: number
total_input_tokens: number
total_output_tokens: number
@@ -329,7 +328,7 @@ export const meApi = {
}): Promise<UsageResponse> {
const clientSendUnixMs = beginServerTimingSample()
const response = await apiClient.get<UsageResponse>('/api/users/me/usage', { params })
return withServerTiming(response.data, clientSendUnixMs)
return withServerTiming(response, clientSendUnixMs)
},
// 获取活跃请求状态(用于轮询更新)
@@ -368,7 +367,7 @@ export const meApi = {
const params = ids ? { ids } : {}
const clientSendUnixMs = beginServerTimingSample()
const response = await apiClient.get('/api/users/me/usage/active', { params })
return withServerTiming(response.data, clientSendUnixMs)
return withServerTiming(response, clientSendUnixMs)
},
// 获取可用的提供商
+44 -10
View File
@@ -1,7 +1,12 @@
import type { AxiosResponse } from 'axios'
export const SERVER_NOW_UNIX_MS_HEADER = 'x-aether-server-now-unix-ms'
export interface ServerTimingMetadata {
server_now_unix_ms: number
client_send_unix_ms: number
client_receive_unix_ms: number
round_trip_ms: number
}
export interface ServerTimedPayload {
@@ -12,33 +17,62 @@ export function beginServerTimingSample(): number {
return Date.now()
}
export function readServerNowUnixMs(payload: unknown): number | null {
if (!payload || typeof payload !== 'object') return null
const value = (payload as { server_now_unix_ms?: unknown }).server_now_unix_ms
return typeof value === 'number' && Number.isFinite(value) ? value : null
function readHeaderValue(headers: unknown, name: string): unknown {
if (!headers || typeof headers !== 'object') return undefined
const get = (headers as { get?: unknown }).get
if (typeof get === 'function') {
return get.call(headers, name)
}
const lowerName = name.toLowerCase()
for (const [key, value] of Object.entries(headers as Record<string, unknown>)) {
if (key.toLowerCase() === lowerName) return value
}
return undefined
}
export function readServerNowUnixMsFromHeaders(headers: unknown): number | null {
const value = readHeaderValue(headers, SERVER_NOW_UNIX_MS_HEADER)
const raw = Array.isArray(value) ? value[0] : value
const parsed = typeof raw === 'number'
? raw
: typeof raw === 'string'
? Number(raw.trim())
: Number.NaN
return Number.isSafeInteger(parsed) && parsed > 0 ? parsed : null
}
export function buildServerTimingMetadata(
payload: unknown,
response: Pick<AxiosResponse, 'headers'> | { headers?: unknown } | null | undefined,
clientSendUnixMs: number,
clientReceiveUnixMs = Date.now()
): ServerTimingMetadata | undefined {
const serverNowUnixMs = readServerNowUnixMs(payload)
const serverNowUnixMs = readServerNowUnixMsFromHeaders(response?.headers)
if (serverNowUnixMs == null) return undefined
if (!Number.isFinite(clientSendUnixMs) || !Number.isFinite(clientReceiveUnixMs)) return undefined
if (clientReceiveUnixMs < clientSendUnixMs) return undefined
const roundTripMs = clientReceiveUnixMs - clientSendUnixMs
return {
server_now_unix_ms: serverNowUnixMs,
client_send_unix_ms: clientSendUnixMs,
client_receive_unix_ms: clientReceiveUnixMs,
round_trip_ms: roundTripMs,
}
}
export function withServerTiming<T extends object>(payload: T, clientSendUnixMs: number): T & ServerTimedPayload {
const serverTiming = buildServerTimingMetadata(payload, clientSendUnixMs)
if (!serverTiming) return payload
export function withServerTiming<T extends object>(
response: Pick<AxiosResponse<T>, 'data' | 'headers'>,
clientSendUnixMs: number
): T & ServerTimedPayload {
const serverTiming = buildServerTimingMetadata(response, clientSendUnixMs)
if (!serverTiming) return response.data
return {
...payload,
...response.data,
server_timing: serverTiming,
}
}
+4 -5
View File
@@ -134,7 +134,6 @@ export interface UsageRequestOptions {
type UsageListResponse = ServerTimedPayload & {
records?: unknown
server_now_unix_ms?: unknown
pagination?: {
total?: unknown
limit?: unknown
@@ -379,7 +378,7 @@ export const usageApi = {
const { params, pagination } = buildCurrentUserUsageParams(filters)
const clientSendUnixMs = beginServerTimingSample()
const response = await apiClient.get<UsageListResponse>('/api/users/me/usage', { params })
return normalizeUsageRecordPage(withServerTiming(response.data, clientSendUnixMs), pagination)
return normalizeUsageRecordPage(withServerTiming(response, clientSendUnixMs), pagination)
},
async getUsageStats(filters?: UsageFilters, options?: UsageRequestOptions): Promise<UsageStats> {
@@ -462,7 +461,7 @@ export const usageApi = {
const recordsClientSendUnixMs = beginServerTimingSample()
const recordsRequest = apiClient
.get<UsageListResponse>('/api/admin/usage/records', { params: recordParams })
.then(response => withServerTiming(response.data, recordsClientSendUnixMs))
.then(response => withServerTiming(response, recordsClientSendUnixMs))
const [statsResponse, recordsResponse] = await Promise.all([
statsRequest,
@@ -503,7 +502,7 @@ export const usageApi = {
return dedupedRequest(key, async () => {
const clientSendUnixMs = beginServerTimingSample()
const response = await apiClient.get('/api/admin/usage/records', { params })
return withServerTiming(response.data, clientSendUnixMs)
return withServerTiming(response, clientSendUnixMs)
})
},
@@ -571,7 +570,7 @@ export const usageApi = {
}
const clientSendUnixMs = beginServerTimingSample()
const response = await apiClient.get('/api/admin/usage/active', { params })
return withServerTiming(response.data, clientSendUnixMs)
return withServerTiming(response, clientSendUnixMs)
},
/**
@@ -1,15 +1,20 @@
import { describe, expect, it } from 'vitest'
import { calculateServerClockOffsetMs, useServerClock } from '../useServerClock'
import {
calculateServerClockOffsetMs,
shouldUseServerClockSample,
useServerClock
} from '../useServerClock'
describe('useServerClock', () => {
it('calculates offset from the request midpoint', () => {
it('calculates offset from the response receive time', () => {
const offset = calculateServerClockOffsetMs({
server_now_unix_ms: 10_500,
client_send_unix_ms: 20_000,
client_receive_unix_ms: 20_200,
round_trip_ms: 200,
})
expect(offset).toBe(-9_600)
expect(offset).toBe(-9_700)
})
it('ignores missing or invalid timing samples', () => {
@@ -18,11 +23,19 @@ describe('useServerClock', () => {
server_now_unix_ms: Number.NaN,
client_send_unix_ms: 20_000,
client_receive_unix_ms: 20_200,
round_trip_ms: 200,
})).toBeNull()
expect(calculateServerClockOffsetMs({
server_now_unix_ms: 10_500,
client_send_unix_ms: 20_200,
client_receive_unix_ms: 20_000,
round_trip_ms: 200,
})).toBeNull()
expect(calculateServerClockOffsetMs({
server_now_unix_ms: 10_500,
client_send_unix_ms: 20_000,
client_receive_unix_ms: 20_200,
round_trip_ms: Number.NaN,
})).toBeNull()
})
@@ -33,10 +46,57 @@ describe('useServerClock', () => {
server_now_unix_ms: 10_500,
client_send_unix_ms: 20_000,
client_receive_unix_ms: 20_200,
round_trip_ms: 200,
})
clock.updateServerClockOffset(undefined)
expect(clock.hasServerClockOffset.value).toBe(true)
expect(clock.serverClockOffsetMs.value).toBe(-9_600)
expect(clock.serverClockOffsetMs.value).toBe(-9_700)
expect(clock.serverClockSampleRoundTripMs.value).toBe(200)
})
it('does not let a much slower sample overwrite a better clock offset', () => {
const clock = useServerClock()
clock.updateServerClockOffset({
server_now_unix_ms: 10_500,
client_send_unix_ms: 20_000,
client_receive_unix_ms: 20_050,
round_trip_ms: 50,
})
clock.updateServerClockOffset({
server_now_unix_ms: 20_500,
client_send_unix_ms: 30_000,
client_receive_unix_ms: 30_500,
round_trip_ms: 500,
})
expect(clock.serverClockOffsetMs.value).toBe(-9_550)
expect(clock.serverClockSampleRoundTripMs.value).toBe(50)
})
it('accepts a faster sample after an initial slow sample', () => {
const clock = useServerClock()
clock.updateServerClockOffset({
server_now_unix_ms: 10_500,
client_send_unix_ms: 20_000,
client_receive_unix_ms: 20_500,
round_trip_ms: 500,
})
clock.updateServerClockOffset({
server_now_unix_ms: 20_500,
client_send_unix_ms: 30_000,
client_receive_unix_ms: 30_050,
round_trip_ms: 50,
})
expect(clock.serverClockOffsetMs.value).toBe(-9_550)
expect(clock.serverClockSampleRoundTripMs.value).toBe(50)
})
it('allows small RTT regressions so the offset can stay fresh', () => {
expect(shouldUseServerClockSample(140, 50)).toBe(true)
expect(shouldUseServerClockSample(151, 50)).toBe(false)
})
})
@@ -1,36 +1,62 @@
import { ref } from 'vue'
import type { ServerTimingMetadata } from '@/api/serverTiming'
const SERVER_CLOCK_RTT_REGRESSION_TOLERANCE_MS = 100
export function calculateServerClockOffsetMs(timing: ServerTimingMetadata | null | undefined): number | null {
if (!timing) return null
const { server_now_unix_ms: serverNowUnixMs, client_send_unix_ms: clientSendUnixMs, client_receive_unix_ms: clientReceiveUnixMs } = timing
const {
server_now_unix_ms: serverNowUnixMs,
client_send_unix_ms: clientSendUnixMs,
client_receive_unix_ms: clientReceiveUnixMs,
round_trip_ms: roundTripMs,
} = timing
if (!Number.isFinite(serverNowUnixMs) || !Number.isFinite(clientSendUnixMs) || !Number.isFinite(clientReceiveUnixMs)) {
if (
!Number.isFinite(serverNowUnixMs) ||
!Number.isFinite(clientSendUnixMs) ||
!Number.isFinite(clientReceiveUnixMs) ||
!Number.isFinite(roundTripMs)
) {
return null
}
if (clientReceiveUnixMs < clientSendUnixMs) {
if (clientReceiveUnixMs < clientSendUnixMs || roundTripMs < 0) {
return null
}
const clientMidpointUnixMs = clientSendUnixMs + ((clientReceiveUnixMs - clientSendUnixMs) / 2)
return serverNowUnixMs - clientMidpointUnixMs
return serverNowUnixMs - clientReceiveUnixMs
}
export function shouldUseServerClockSample(
nextRoundTripMs: number,
currentRoundTripMs: number | null | undefined
): boolean {
if (!Number.isFinite(nextRoundTripMs) || nextRoundTripMs < 0) return false
if (currentRoundTripMs == null || !Number.isFinite(currentRoundTripMs)) return true
return nextRoundTripMs <= currentRoundTripMs + SERVER_CLOCK_RTT_REGRESSION_TOLERANCE_MS
}
export function useServerClock() {
const serverClockOffsetMs = ref(0)
const hasServerClockOffset = ref(false)
const serverClockSampleRoundTripMs = ref<number | null>(null)
function updateServerClockOffset(timing: ServerTimingMetadata | null | undefined): void {
const offsetMs = calculateServerClockOffsetMs(timing)
if (offsetMs == null) return
if (!shouldUseServerClockSample(timing?.round_trip_ms ?? Number.NaN, serverClockSampleRoundTripMs.value)) {
return
}
serverClockOffsetMs.value = offsetMs
serverClockSampleRoundTripMs.value = timing?.round_trip_ms ?? null
hasServerClockOffset.value = true
}
return {
serverClockOffsetMs,
hasServerClockOffset,
serverClockSampleRoundTripMs,
updateServerClockOffset,
}
}