mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-12 14:10:19 +08:00
Merge remote-tracking branch 'origin/pr/597'
This commit is contained in:
@@ -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)
|
||||
})
|
||||
})
|
||||
@@ -1,4 +1,4 @@
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
const { getMock, cachedRequestMock, dedupedRequestMock, buildCacheKeyMock } = vi.hoisted(() => ({
|
||||
getMock: vi.fn(),
|
||||
@@ -29,6 +29,10 @@ describe('usageApi contract alignment', () => {
|
||||
buildCacheKeyMock.mockClear()
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks()
|
||||
})
|
||||
|
||||
it('loads current-user usage records from the Rust usage endpoint and normalizes pagination', async () => {
|
||||
getMock.mockResolvedValueOnce({
|
||||
data: {
|
||||
@@ -115,6 +119,56 @@ describe('usageApi contract alignment', () => {
|
||||
})
|
||||
})
|
||||
|
||||
it('captures admin user usage server timing when the records request resolves', async () => {
|
||||
let now = 1_000
|
||||
vi.spyOn(Date, 'now').mockImplementation(() => now)
|
||||
|
||||
let resolveStats: ((value: unknown) => void) | null = null
|
||||
getMock.mockImplementation((url: string) => {
|
||||
if (url === '/api/admin/usage/stats') {
|
||||
return new Promise(resolve => {
|
||||
resolveStats = resolve
|
||||
})
|
||||
}
|
||||
if (url === '/api/admin/usage/records') {
|
||||
return Promise.resolve({
|
||||
headers: {
|
||||
'x-aether-server-now-unix-ms': '10050',
|
||||
},
|
||||
data: {
|
||||
records: [{ id: 'record-3' }],
|
||||
total: 1,
|
||||
limit: 25,
|
||||
offset: 0,
|
||||
},
|
||||
})
|
||||
}
|
||||
return Promise.reject(new Error(`unexpected url: ${url}`))
|
||||
})
|
||||
|
||||
const resultPromise = usageApi.getUserUsage('user-123', { page: 1, page_size: 25 })
|
||||
now = 1_100
|
||||
await Promise.resolve()
|
||||
now = 20_000
|
||||
resolveStats?.({
|
||||
data: {
|
||||
total_requests: 1,
|
||||
total_tokens: 10,
|
||||
total_cost: 0.1,
|
||||
avg_response_time: 500,
|
||||
},
|
||||
})
|
||||
|
||||
const result = await resultPromise
|
||||
|
||||
expect(result.server_timing).toEqual({
|
||||
server_now_unix_ms: 10_050,
|
||||
client_send_unix_ms: 1_000,
|
||||
client_receive_unix_ms: 1_100,
|
||||
round_trip_ms: 100,
|
||||
})
|
||||
})
|
||||
|
||||
it('uses an extended timeout and cache bypass option for admin analytics', async () => {
|
||||
getMock
|
||||
.mockResolvedValueOnce({
|
||||
|
||||
+11
-3
@@ -5,6 +5,11 @@ import { cachedRequest, buildCacheKey } from '@/utils/cache'
|
||||
import type { BillingSummary } from './auth'
|
||||
import type { UserSession } from '@/types/session'
|
||||
import type { FeatureSettingsMap } from '@/utils/featureSettings'
|
||||
import {
|
||||
beginServerTimingSample,
|
||||
withServerTiming,
|
||||
type ServerTimedPayload,
|
||||
} from './serverTiming'
|
||||
|
||||
const ACTIVITY_HEATMAP_CACHE_TTL_MS = 30 * 60 * 1000
|
||||
|
||||
@@ -142,7 +147,7 @@ export interface ApiFormatSummary {
|
||||
}
|
||||
|
||||
// 使用统计响应接口
|
||||
export interface UsageResponse {
|
||||
export interface UsageResponse extends ServerTimedPayload {
|
||||
total_requests: number
|
||||
total_input_tokens: number
|
||||
total_output_tokens: number
|
||||
@@ -321,12 +326,14 @@ export const meApi = {
|
||||
limit?: number
|
||||
offset?: number
|
||||
}): Promise<UsageResponse> {
|
||||
const clientSendUnixMs = beginServerTimingSample()
|
||||
const response = await apiClient.get<UsageResponse>('/api/users/me/usage', { params })
|
||||
return response.data
|
||||
return withServerTiming(response, clientSendUnixMs)
|
||||
},
|
||||
|
||||
// 获取活跃请求状态(用于轮询更新)
|
||||
async getActiveRequests(ids?: string): Promise<{
|
||||
server_timing?: ServerTimedPayload['server_timing']
|
||||
requests: Array<{
|
||||
id: string
|
||||
status: 'pending' | 'streaming' | 'completed' | 'failed' | 'cancelled'
|
||||
@@ -358,8 +365,9 @@ export const meApi = {
|
||||
}>
|
||||
}> {
|
||||
const params = ids ? { ids } : {}
|
||||
const clientSendUnixMs = beginServerTimingSample()
|
||||
const response = await apiClient.get('/api/users/me/usage/active', { params })
|
||||
return response.data
|
||||
return withServerTiming(response, clientSendUnixMs)
|
||||
},
|
||||
|
||||
// 获取可用的提供商
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
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 {
|
||||
server_timing?: ServerTimingMetadata
|
||||
}
|
||||
|
||||
export function beginServerTimingSample(): number {
|
||||
return Date.now()
|
||||
}
|
||||
|
||||
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(
|
||||
response: Pick<AxiosResponse, 'headers'> | { headers?: unknown } | null | undefined,
|
||||
clientSendUnixMs: number,
|
||||
clientReceiveUnixMs = Date.now()
|
||||
): ServerTimingMetadata | undefined {
|
||||
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>(
|
||||
response: Pick<AxiosResponse<T>, 'data' | 'headers'>,
|
||||
clientSendUnixMs: number
|
||||
): T & ServerTimedPayload {
|
||||
const serverTiming = buildServerTimingMetadata(response, clientSendUnixMs)
|
||||
if (!serverTiming) return response.data
|
||||
return {
|
||||
...response.data,
|
||||
server_timing: serverTiming,
|
||||
}
|
||||
}
|
||||
@@ -2,6 +2,11 @@ import apiClient from './client'
|
||||
import { cachedRequest, dedupedRequest, buildCacheKey } from '@/utils/cache'
|
||||
import type { ActivityHeatmap } from '@/types/activity'
|
||||
import type { ImageProgress } from './requestTrace'
|
||||
import {
|
||||
beginServerTimingSample,
|
||||
withServerTiming,
|
||||
type ServerTimedPayload,
|
||||
} from './serverTiming'
|
||||
|
||||
const ACTIVITY_HEATMAP_CACHE_TTL_MS = 30 * 60 * 1000
|
||||
const USAGE_ANALYTICS_CACHE_TTL_MS = 30 * 1000
|
||||
@@ -127,7 +132,7 @@ export interface UsageRequestOptions {
|
||||
skipCache?: boolean
|
||||
}
|
||||
|
||||
type UsageListResponse = {
|
||||
type UsageListResponse = ServerTimedPayload & {
|
||||
records?: unknown
|
||||
pagination?: {
|
||||
total?: unknown
|
||||
@@ -199,6 +204,7 @@ function normalizeUsageRecordPage(
|
||||
total: number
|
||||
page: number
|
||||
page_size: number
|
||||
server_timing?: ServerTimedPayload['server_timing']
|
||||
} {
|
||||
const records = assertUsageRecords(payload.records)
|
||||
const pagination = payload.pagination
|
||||
@@ -220,6 +226,7 @@ function normalizeUsageRecordPage(
|
||||
total,
|
||||
page: resolvedPage,
|
||||
page_size: limit,
|
||||
...(payload.server_timing ? { server_timing: payload.server_timing } : {}),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -366,10 +373,12 @@ export const usageApi = {
|
||||
total: number
|
||||
page: number
|
||||
page_size: number
|
||||
server_timing?: ServerTimedPayload['server_timing']
|
||||
}> {
|
||||
const { params, pagination } = buildCurrentUserUsageParams(filters)
|
||||
const clientSendUnixMs = beginServerTimingSample()
|
||||
const response = await apiClient.get<UsageListResponse>('/api/users/me/usage', { params })
|
||||
return normalizeUsageRecordPage(response.data, pagination)
|
||||
return normalizeUsageRecordPage(withServerTiming(response, clientSendUnixMs), pagination)
|
||||
},
|
||||
|
||||
async getUsageStats(filters?: UsageFilters, options?: UsageRequestOptions): Promise<UsageStats> {
|
||||
@@ -444,17 +453,25 @@ export const usageApi = {
|
||||
async getUserUsage(userId: string, filters?: UsageFilters): Promise<{
|
||||
records: UsageRecord[]
|
||||
stats: UsageStats
|
||||
server_timing?: ServerTimedPayload['server_timing']
|
||||
}> {
|
||||
const statsParams = buildAdminUsageStatsParams(userId, filters)
|
||||
const { params: recordParams } = buildAdminUsageRecordParams(userId, filters)
|
||||
const statsRequest = apiClient.get<UsageStats>('/api/admin/usage/stats', { params: statsParams })
|
||||
const recordsClientSendUnixMs = beginServerTimingSample()
|
||||
const recordsRequest = apiClient
|
||||
.get<UsageListResponse>('/api/admin/usage/records', { params: recordParams })
|
||||
.then(response => withServerTiming(response, recordsClientSendUnixMs))
|
||||
|
||||
const [statsResponse, recordsResponse] = await Promise.all([
|
||||
apiClient.get<UsageStats>('/api/admin/usage/stats', { params: statsParams }),
|
||||
apiClient.get<UsageListResponse>('/api/admin/usage/records', { params: recordParams }),
|
||||
statsRequest,
|
||||
recordsRequest,
|
||||
])
|
||||
|
||||
return {
|
||||
records: assertUsageRecords(recordsResponse.data.records),
|
||||
records: assertUsageRecords(recordsResponse.records),
|
||||
stats: statsResponse.data,
|
||||
...(recordsResponse.server_timing ? { server_timing: recordsResponse.server_timing } : {}),
|
||||
}
|
||||
},
|
||||
|
||||
@@ -480,11 +497,13 @@ export const usageApi = {
|
||||
total: number
|
||||
limit: number
|
||||
offset: number
|
||||
server_timing?: ServerTimedPayload['server_timing']
|
||||
}> {
|
||||
const key = buildCacheKey('usage:records', params as Record<string, unknown> | undefined)
|
||||
return dedupedRequest(key, async () => {
|
||||
const clientSendUnixMs = beginServerTimingSample()
|
||||
const response = await apiClient.get('/api/admin/usage/records', { params })
|
||||
return response.data
|
||||
return withServerTiming(response, clientSendUnixMs)
|
||||
})
|
||||
},
|
||||
|
||||
@@ -496,6 +515,7 @@ export const usageApi = {
|
||||
ids?: string[],
|
||||
timeRange?: Pick<UsageFilters, 'start_date' | 'end_date' | 'preset' | 'timezone' | 'tz_offset_minutes'>
|
||||
): Promise<{
|
||||
server_timing?: ServerTimedPayload['server_timing']
|
||||
requests: Array<{
|
||||
id: string
|
||||
status: 'pending' | 'streaming' | 'completed' | 'failed' | 'cancelled'
|
||||
@@ -549,8 +569,9 @@ export const usageApi = {
|
||||
if (typeof timeRange?.tz_offset_minutes === 'number') {
|
||||
params.tz_offset_minutes = timeRange.tz_offset_minutes
|
||||
}
|
||||
const clientSendUnixMs = beginServerTimingSample()
|
||||
const response = await apiClient.get('/api/admin/usage/active', { params })
|
||||
return response.data
|
||||
return withServerTiming(response, clientSendUnixMs)
|
||||
},
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user