fix(usage): bound terminal event persistence

Add end-to-end terminal admission, bounded database fallback, and observable overload handling. Preserve first-byte lifecycle state across asynchronous runtime and frontend updates.
This commit is contained in:
elky
2026-07-16 13:12:30 +08:00
parent c6d373e6aa
commit 9a47267545
10 changed files with 1952 additions and 447 deletions
+24
View File
@@ -332,10 +332,22 @@ export interface GatewayUsageRuntimeMetrics {
workerProcessFailuresTotal: number | null
workerReadFailuresTotal: number | null
workerReclaimFailuresTotal: number | null
terminalSubmissionLimit: number | null
terminalSubmissionInFlight: number | null
terminalSubmissionMaxInFlight: number | null
terminalSubmissionRejectedTotal: number | null
terminalEnqueueInFlight: number | null
terminalEnqueueDeferredTotal: number | null
terminalEnqueueDeferredDirectWriteTotal: number | null
terminalEnqueueDeferredDroppedTotal: number | null
terminalEnqueueDeferredRetryTotal: number | null
terminalEnqueueFailedTotal: number | null
terminalDirectFallbackLimit: number | null
terminalDirectFallbackInFlight: number | null
terminalDirectFallbackMaxInFlight: number | null
terminalDirectFallbackSucceededTotal: number | null
terminalDirectFallbackFailedTotal: number | null
terminalDirectFallbackRejectedTotal: number | null
lifecycleEnqueueInFlight: number | null
lifecycleEnqueueDeferredTotal: number | null
lifecycleEnqueueDeferredDroppedTotal: number | null
@@ -673,10 +685,22 @@ function buildUsageRuntimeMetrics(samples: PrometheusSample[]): GatewayUsageRunt
workerProcessFailuresTotal: findMetricValueNumber(samples, 'usage_runtime_queue_worker_process_failures_total'),
workerReadFailuresTotal: findMetricValueNumber(samples, 'usage_runtime_queue_worker_read_failures_total'),
workerReclaimFailuresTotal: findMetricValueNumber(samples, 'usage_runtime_queue_worker_reclaim_failures_total'),
terminalSubmissionLimit: findMetricValueNumber(samples, 'usage_runtime_terminal_submission_limit'),
terminalSubmissionInFlight: findMetricValueNumber(samples, 'usage_runtime_terminal_submission_in_flight'),
terminalSubmissionMaxInFlight: findMetricValueNumber(samples, 'usage_runtime_terminal_submission_max_in_flight'),
terminalSubmissionRejectedTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_submission_rejected_total'),
terminalEnqueueInFlight: findMetricValueNumber(samples, 'usage_runtime_terminal_enqueue_in_flight'),
terminalEnqueueDeferredTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_enqueue_deferred_total'),
terminalEnqueueDeferredDirectWriteTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_enqueue_deferred_direct_write_total'),
terminalEnqueueDeferredDroppedTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_enqueue_deferred_dropped_total'),
terminalEnqueueDeferredRetryTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_enqueue_deferred_retry_total'),
terminalEnqueueFailedTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_enqueue_failed_total'),
terminalDirectFallbackLimit: findMetricValueNumber(samples, 'usage_runtime_terminal_direct_fallback_limit'),
terminalDirectFallbackInFlight: findMetricValueNumber(samples, 'usage_runtime_terminal_direct_fallback_in_flight'),
terminalDirectFallbackMaxInFlight: findMetricValueNumber(samples, 'usage_runtime_terminal_direct_fallback_max_in_flight'),
terminalDirectFallbackSucceededTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_direct_fallback_succeeded_total'),
terminalDirectFallbackFailedTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_direct_fallback_failed_total'),
terminalDirectFallbackRejectedTotal: findMetricValueNumber(samples, 'usage_runtime_terminal_direct_fallback_rejected_total'),
lifecycleEnqueueInFlight: findMetricValueNumber(samples, 'usage_runtime_lifecycle_enqueue_in_flight'),
lifecycleEnqueueDeferredTotal: findMetricValueNumber(samples, 'usage_runtime_lifecycle_enqueue_deferred_total'),
lifecycleEnqueueDeferredDroppedTotal: findMetricValueNumber(samples, 'usage_runtime_lifecycle_enqueue_deferred_dropped_total'),
@@ -1,7 +1,10 @@
import { describe, expect, it } from 'vitest'
import type { UsageRecord } from '../../types'
import { syncUsageRecordStreamResolution } from '../recordSync'
import {
mergeUsageRecordFirstByteTimeMs,
syncUsageRecordStreamResolution,
} from '../recordSync'
function buildUsageRecord(overrides: Partial<UsageRecord> = {}): UsageRecord {
return {
@@ -59,3 +62,16 @@ describe('syncUsageRecordStreamResolution', () => {
expect(nextRecords).toBe(records)
})
})
describe('mergeUsageRecordFirstByteTimeMs', () => {
it('does not let a stale active update clear or reduce a resolved first-byte time', () => {
expect(mergeUsageRecordFirstByteTimeMs(500, null)).toBe(500)
expect(mergeUsageRecordFirstByteTimeMs(500, 320)).toBe(500)
expect(mergeUsageRecordFirstByteTimeMs(500, 640)).toBe(640)
expect(mergeUsageRecordFirstByteTimeMs(undefined, 320)).toBe(320)
expect(mergeUsageRecordFirstByteTimeMs(undefined, 0)).toBe(0)
expect(mergeUsageRecordFirstByteTimeMs(0, null)).toBe(0)
expect(mergeUsageRecordFirstByteTimeMs(null, -1)).toBeNull()
expect(mergeUsageRecordFirstByteTimeMs(-1, null)).toBeUndefined()
})
})
@@ -78,6 +78,11 @@ describe('usage status helpers', () => {
status: 'streaming',
first_byte_time_ms: 320,
}))).toBe('streaming')
expect(resolveDisplayRequestStatus(buildUsageRecord({
status: 'streaming',
first_byte_time_ms: 0,
}))).toBe('streaming')
})
it('treats active lifecycle records with failure signals as failed for display', () => {
@@ -5,6 +5,30 @@ export type UsageRecordStreamResolution = Pick<
'id' | 'is_stream' | 'upstream_is_stream' | 'client_requested_stream' | 'client_is_stream'
>
export function mergeUsageRecordFirstByteTimeMs(
existingValue: number | null | undefined,
nextValue: number | null | undefined
): number | null | undefined {
// Millisecond timing is floored, so 0 still means the first byte was observed.
const existingIsResolved = typeof existingValue === 'number' &&
Number.isFinite(existingValue) &&
existingValue >= 0
const nextIsResolved = typeof nextValue === 'number' &&
Number.isFinite(nextValue) &&
nextValue >= 0
if (existingIsResolved && nextIsResolved) {
return Math.max(existingValue, nextValue)
}
if (existingIsResolved) {
return existingValue
}
if (nextIsResolved) {
return nextValue
}
return existingValue == null ? existingValue : undefined
}
export function syncUsageRecordStreamResolution(
records: UsageRecord[],
resolved: UsageRecordStreamResolution
+9 -2
View File
@@ -154,6 +154,7 @@ import {
getDateRangeFromPeriod
} from '@/features/usage/composables'
import { reconcileActiveRequestDiscovery } from '@/features/usage/utils/activeRequestDiscovery'
import { mergeUsageRecordFirstByteTimeMs } from '@/features/usage/utils/recordSync'
import {
hasUsageFallback,
isUsageRecordFailed,
@@ -521,7 +522,10 @@ async function pollActiveRequests() {
record.actual_cost = update.actual_cost ?? undefined
record.rate_multiplier = update.rate_multiplier ?? undefined
record.response_time_ms = update.response_time_ms ?? undefined
record.first_byte_time_ms = update.first_byte_time_ms ?? undefined
record.first_byte_time_ms = mergeUsageRecordFirstByteTimeMs(
record.first_byte_time_ms,
update.first_byte_time_ms
)
if ('updated_at' in update) {
record.updated_at = typeof update.updated_at === 'string' ? update.updated_at : null
}
@@ -1104,7 +1108,10 @@ function handleDetailRequestState(update: {
record.response_time_ms = update.responseTimeMs
}
if ('firstByteTimeMs' in update) {
record.first_byte_time_ms = update.firstByteTimeMs ?? undefined
record.first_byte_time_ms = mergeUsageRecordFirstByteTimeMs(
record.first_byte_time_ms,
update.firstByteTimeMs
)
}
if ('isStream' in update && typeof update.isStream === 'boolean') {
record.is_stream = update.isStream