From 8424d01d2eb4a7b14cb309e223ad356646ad073b Mon Sep 17 00:00:00 2001 From: hanafish <1106510024@qq.com> Date: Wed, 22 Jul 2026 21:47:32 +0800 Subject: [PATCH] perf(diagnostics): bound RPC latency aggregation Pre-commit hook ran. Total eslint: 0, total circular: 0 --- src-tauri/src/usage_diagnostics/sanitize.rs | 24 ++++- src-tauri/src/usage_diagnostics/types.rs | 2 +- .../usage_diagnostics_tests.rs | 29 +++++- src/diagnostics/aggregate.ts | 2 +- src/diagnostics/runtimeCounters.test.ts | 68 +++++++++++++ src/diagnostics/runtimeCounters.ts | 95 ++++++++++++++++--- src/diagnostics/types.ts | 3 +- src/diagnostics/useDiagnosticsBootstrap.ts | 8 ++ 8 files changed, 207 insertions(+), 24 deletions(-) create mode 100644 src/diagnostics/runtimeCounters.test.ts diff --git a/src-tauri/src/usage_diagnostics/sanitize.rs b/src-tauri/src/usage_diagnostics/sanitize.rs index ba87284ef..6b6cccc02 100644 --- a/src-tauri/src/usage_diagnostics/sanitize.rs +++ b/src-tauri/src/usage_diagnostics/sanitize.rs @@ -6,7 +6,7 @@ use super::types::{ }; const UNKNOWN_BUCKET: &str = "unknown"; -const MAX_RUNTIME_OPERATIONS: usize = 100; +const MAX_RUNTIME_OPERATIONS: usize = 128; const MAX_LIST_ENTRIES: usize = 25; pub fn bucket_duration_ms(duration_ms: f64) -> &'static str { @@ -105,10 +105,24 @@ fn sanitize_runtime_summary(summary: DiagnosticsRuntimeSummary) -> DiagnosticsRu item.insert("total".into(), Value::from(total)); item.insert("success".into(), Value::from(success)); item.insert("failure".into(), Value::from(failure)); - item.insert( - "durationBucket".into(), - Value::from(bucket_string_field(&value, "durationBucket")), - ); + if value.get("averageDurationBucket").is_some() || value.get("p95DurationBucket").is_some() + { + item.insert( + "averageDurationBucket".into(), + Value::from(bucket_string_field(&value, "averageDurationBucket")), + ); + item.insert( + "p95DurationBucket".into(), + Value::from(bucket_string_field(&value, "p95DurationBucket")), + ); + } else { + // Queue records produced by diagnostics schema v1 used one coarse + // duration bucket. Keep accepting them while emitting v2 snapshots. + item.insert( + "durationBucket".into(), + Value::from(bucket_string_field(&value, "durationBucket")), + ); + } by_operation.insert(sanitize_operation_name(&operation), Value::Object(item)); } diff --git a/src-tauri/src/usage_diagnostics/types.rs b/src-tauri/src/usage_diagnostics/types.rs index cfeb3e701..03103549f 100644 --- a/src-tauri/src/usage_diagnostics/types.rs +++ b/src-tauri/src/usage_diagnostics/types.rs @@ -1,7 +1,7 @@ use serde::{Deserialize, Serialize}; use serde_json::{Map, Value}; -pub const DIAGNOSTICS_SCHEMA_VERSION: u32 = 1; +pub const DIAGNOSTICS_SCHEMA_VERSION: u32 = 2; #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "kebab-case")] diff --git a/src-tauri/src/usage_diagnostics/usage_diagnostics_tests.rs b/src-tauri/src/usage_diagnostics/usage_diagnostics_tests.rs index 491b7c804..5a9db60ec 100644 --- a/src-tauri/src/usage_diagnostics/usage_diagnostics_tests.rs +++ b/src-tauri/src/usage_diagnostics/usage_diagnostics_tests.rs @@ -109,7 +109,7 @@ fn performance_only_strips_usage_aggregates() { snapshot(DiagnosticsLevel::Default), DiagnosticsLevel::PerformanceOnly, ); - assert_eq!(sanitized.schema_version, 1); + assert_eq!(sanitized.schema_version, 2); assert_eq!(sanitized.app_version.as_deref(), Some("1.0.1+test")); assert!(sanitized.app_usage_duration_bucket.is_some()); assert!(sanitized.system_profile.is_some()); @@ -125,7 +125,7 @@ fn performance_only_strips_usage_aggregates() { #[test] fn off_keeps_minimal_existence_usage_only() { let sanitized = sanitize_snapshot(snapshot(DiagnosticsLevel::Default), DiagnosticsLevel::Off); - assert_eq!(sanitized.schema_version, 1); + assert_eq!(sanitized.schema_version, 2); assert_eq!(sanitized.diagnostics_level, DiagnosticsLevel::Off); assert_eq!(sanitized.app_version.as_deref(), Some("1.0.1+test")); assert!(sanitized.app_usage_duration_bucket.is_some()); @@ -158,6 +158,31 @@ fn default_sanitization_drops_raw_or_path_fields() { assert!(serialized.contains("agent_send_message")); } +#[test] +fn runtime_sanitizer_accepts_v1_and_v2_duration_shapes() { + let mut legacy = snapshot(DiagnosticsLevel::PerformanceOnly); + let sanitized_legacy = sanitize_snapshot(legacy.clone(), DiagnosticsLevel::PerformanceOnly); + let legacy_operation = &sanitized_legacy.rpc.unwrap().by_operation["agent_send_message"]; + assert_eq!(legacy_operation["durationBucket"], "lt_1m"); + + legacy.rpc.as_mut().unwrap().by_operation.insert( + "cli_agent_message".to_string(), + json!({ + "total": 2, + "success": 1, + "failure": 1, + "averageDurationBucket": "20_100ms", + "p95DurationBucket": "500ms_2s", + "payload": "drop" + }), + ); + let sanitized_v2 = sanitize_snapshot(legacy, DiagnosticsLevel::PerformanceOnly); + let v2_operation = &sanitized_v2.rpc.unwrap().by_operation["cli_agent_message"]; + assert_eq!(v2_operation["averageDurationBucket"], "20_100ms"); + assert_eq!(v2_operation["p95DurationBucket"], "500ms_2s"); + assert!(v2_operation.get("payload").is_none()); +} + #[test] fn queue_reads_only_unsent_and_marks_sent() { let dir = tempfile::tempdir().unwrap(); diff --git a/src/diagnostics/aggregate.ts b/src/diagnostics/aggregate.ts index 0f64ac250..8e4209f4c 100644 --- a/src/diagnostics/aggregate.ts +++ b/src/diagnostics/aggregate.ts @@ -36,7 +36,7 @@ import type { DiagnosticsUsageSnapshot, } from "./types"; -const SCHEMA_VERSION = 1; +const SCHEMA_VERSION = 2; const MAX_TOP_MODELS = 10; const MAX_RUST_AGENT_TOP_SESSIONS_PER_DAY = 10; const EXTERNAL_HISTORY_LIMIT = 200; diff --git a/src/diagnostics/runtimeCounters.test.ts b/src/diagnostics/runtimeCounters.test.ts new file mode 100644 index 000000000..18caed1ac --- /dev/null +++ b/src/diagnostics/runtimeCounters.test.ts @@ -0,0 +1,68 @@ +import { beforeEach, describe, expect, it } from "vitest"; + +import { + consumeHttpDiagnosticsSummary, + consumeRpcDiagnosticsSummary, + discardRuntimeDiagnosticsCounters, + recordDiagnosticsRpc, +} from "./runtimeCounters"; + +describe("runtime diagnostics counters", () => { + beforeEach(() => { + consumeRpcDiagnosticsSummary(); + consumeHttpDiagnosticsSummary(); + }); + + it("reports average, p95, and failures from a fixed histogram", () => { + for (let index = 0; index < 95; index += 1) { + recordDiagnosticsRpc("fast", 4, true); + } + for (let index = 0; index < 5; index += 1) { + recordDiagnosticsRpc("fast", 600, index !== 4); + } + + const summary = consumeRpcDiagnosticsSummary(); + + expect(summary).toMatchObject({ total: 100, success: 99, failure: 1 }); + expect(summary.byOperation.fast).toEqual({ + total: 100, + success: 99, + failure: 1, + averageDurationBucket: "20_100ms", + p95DurationBucket: "1_5ms", + }); + }); + + it("keeps at most 128 operation entries and merges overflow", () => { + for (let index = 0; index < 1_000; index += 1) { + recordDiagnosticsRpc(`operation-${index}`, index % 10, index % 11 !== 0); + } + + const summary = consumeRpcDiagnosticsSummary(); + const operations = Object.keys(summary.byOperation); + + expect(operations).toHaveLength(128); + expect(summary.byOperation.__other__).toMatchObject({ + total: 873, + failure: 79, + }); + }); + + it("consumes and releases the current interval", () => { + recordDiagnosticsRpc("one", 1, false); + expect(consumeRpcDiagnosticsSummary().total).toBe(1); + expect(consumeRpcDiagnosticsSummary()).toEqual({ + total: 0, + success: 0, + failure: 0, + byOperation: {}, + }); + }); + + it("discards bounded counters while diagnostics cannot upload", () => { + recordDiagnosticsRpc("offline", 2_500, true); + discardRuntimeDiagnosticsCounters(); + + expect(consumeRpcDiagnosticsSummary().total).toBe(0); + }); +}); diff --git a/src/diagnostics/runtimeCounters.ts b/src/diagnostics/runtimeCounters.ts index 747895222..037c8c7df 100644 --- a/src/diagnostics/runtimeCounters.ts +++ b/src/diagnostics/runtimeCounters.ts @@ -1,35 +1,83 @@ -import { bucketDurationMs } from "./buckets"; import type { DiagnosticsRuntimeSummary } from "./types"; interface RuntimeCounter { total: number; failure: number; - durations: number[]; + totalDurationMs: number; + durationHistogram: number[]; } +const MAX_RUNTIME_OPERATIONS = 128; +const OTHER_OPERATION = "__other__"; +const DURATION_BUCKETS = [ + { upperBoundMs: 1, label: "lt_1ms" }, + { upperBoundMs: 5, label: "1_5ms" }, + { upperBoundMs: 20, label: "5_20ms" }, + { upperBoundMs: 100, label: "20_100ms" }, + { upperBoundMs: 500, label: "100_500ms" }, + { upperBoundMs: 2_000, label: "500ms_2s" }, + { upperBoundMs: Number.POSITIVE_INFINITY, label: "2s_plus" }, +] as const; + const rpcCounters = new Map(); const httpCounters = new Map(); +function createCounter(): RuntimeCounter { + return { + total: 0, + failure: 0, + totalDurationMs: 0, + durationHistogram: Array.from({ length: DURATION_BUCKETS.length }, () => 0), + }; +} + function getCounter( counters: Map, operation: string ): RuntimeCounter { const existing = counters.get(operation); if (existing) return existing; - const created: RuntimeCounter = { total: 0, failure: 0, durations: [] }; - counters.set(operation, created); + + const boundedOperation = + counters.size < MAX_RUNTIME_OPERATIONS - 1 ? operation : OTHER_OPERATION; + const overflow = counters.get(boundedOperation); + if (overflow) return overflow; + + const created = createCounter(); + counters.set(boundedOperation, created); return created; } +function normalizeDuration(durationMs: number): number { + return Number.isFinite(durationMs) && durationMs >= 0 ? durationMs : 0; +} + +function durationBucketIndex(durationMs: number): number { + const index = DURATION_BUCKETS.findIndex( + ({ upperBoundMs }) => durationMs < upperBoundMs + ); + return index === -1 ? DURATION_BUCKETS.length - 1 : index; +} + +function recordCounter( + counter: RuntimeCounter, + durationMs: number, + ok: boolean +): void { + const normalizedDurationMs = normalizeDuration(durationMs); + counter.total += 1; + counter.totalDurationMs += normalizedDurationMs; + counter.durationHistogram[durationBucketIndex(normalizedDurationMs)] += 1; + if (!ok) counter.failure += 1; +} + export function recordDiagnosticsRpc( command: string, durationMs: number, ok: boolean ): void { const counter = getCounter(rpcCounters, command); - counter.total += 1; - if (!ok) counter.failure += 1; - counter.durations.push(durationMs); + recordCounter(counter, durationMs, ok); } export function recordDiagnosticsHttp( @@ -38,14 +86,22 @@ export function recordDiagnosticsHttp( ok: boolean ): void { const counter = getCounter(httpCounters, target); - counter.total += 1; - if (!ok) counter.failure += 1; - counter.durations.push(durationMs); + recordCounter(counter, durationMs, ok); +} + +function bucketLabelForDuration(durationMs: number): string { + return DURATION_BUCKETS[durationBucketIndex(durationMs)].label; } -function average(values: number[]): number { - if (values.length === 0) return 0; - return values.reduce((sum, value) => sum + value, 0) / values.length; +function percentileBucket(histogram: number[], total: number): string { + if (total === 0) return DURATION_BUCKETS[0].label; + const target = Math.ceil(total * 0.95); + let cumulative = 0; + for (let index = 0; index < histogram.length; index += 1) { + cumulative += histogram[index] ?? 0; + if (cumulative >= target) return DURATION_BUCKETS[index].label; + } + return DURATION_BUCKETS[DURATION_BUCKETS.length - 1].label; } function consumeDiagnosticsSummary( @@ -62,7 +118,13 @@ function consumeDiagnosticsSummary( total: counter.total, success: counter.total - counter.failure, failure: counter.failure, - durationBucket: bucketDurationMs(average(counter.durations)), + averageDurationBucket: bucketLabelForDuration( + counter.total === 0 ? 0 : counter.totalDurationMs / counter.total + ), + p95DurationBucket: percentileBucket( + counter.durationHistogram, + counter.total + ), }; } @@ -77,3 +139,8 @@ export function consumeRpcDiagnosticsSummary(): DiagnosticsRuntimeSummary { export function consumeHttpDiagnosticsSummary(): DiagnosticsRuntimeSummary { return consumeDiagnosticsSummary(httpCounters); } + +export function discardRuntimeDiagnosticsCounters(): void { + consumeRpcDiagnosticsSummary(); + consumeHttpDiagnosticsSummary(); +} diff --git a/src/diagnostics/types.ts b/src/diagnostics/types.ts index 27727441d..502937efb 100644 --- a/src/diagnostics/types.ts +++ b/src/diagnostics/types.ts @@ -16,7 +16,8 @@ export interface DiagnosticsRuntimeOperationSummary { total: number; success: number; failure: number; - durationBucket: string; + averageDurationBucket: string; + p95DurationBucket: string; } export interface DiagnosticsRuntimeSummary { diff --git a/src/diagnostics/useDiagnosticsBootstrap.ts b/src/diagnostics/useDiagnosticsBootstrap.ts index b6a3faa34..ea90ccab6 100644 --- a/src/diagnostics/useDiagnosticsBootstrap.ts +++ b/src/diagnostics/useDiagnosticsBootstrap.ts @@ -12,6 +12,7 @@ import { workspaceFoldersAtom } from "@src/store/ui/workspaceFoldersAtom"; import type { WorkspaceFolder } from "@src/types/workspace"; import { createDiagnosticsUsageSnapshot } from "./aggregate"; +import { discardRuntimeDiagnosticsCounters } from "./runtimeCounters"; import { diagnosticsInitialize, diagnosticsSubmitUsageSnapshot, @@ -24,6 +25,10 @@ const HOUR_MS = 60 * MINUTE_MS; const LAST_FLUSH_STORAGE_KEY = "orgii:diagnostics:lastFlushAt"; const logger = createLogger("DiagnosticsBootstrap"); +function discardRuntimeDiagnostics(): void { + discardRuntimeDiagnosticsCounters(); +} + function reportDiagnosticsFailure(operation: string, error: unknown): void { logger.warn(`${operation} failed`, error); } @@ -123,6 +128,9 @@ export function useDiagnosticsBootstrap(): void { useEffect(() => { if (!settingsLoaded) return; + if (offlineMode) { + discardRuntimeDiagnostics(); + } const generation = ++schedulerGenerationRef.current; let cancelled = false;