diff --git a/bun.lock b/bun.lock index 4e1b8f0c446..a851742abb0 100644 --- a/bun.lock +++ b/bun.lock @@ -178,7 +178,7 @@ "@fontsource/dm-mono": "5.2.7", "@fortawesome/fontawesome-free": "7.2.0", "@hey-api/client-fetch": "0.13.1", - "@hey-api/openapi-ts": "0.97.3", + "@hey-api/openapi-ts": "0.99.0", "@monaco-editor/loader": "1.7.0", "@playwright/test": "1.58.2", "@poppanator/sveltekit-svg": "6.0.1", @@ -364,13 +364,13 @@ "@hey-api/client-fetch": ["@hey-api/client-fetch@0.13.1", "", { "peerDependencies": { "@hey-api/openapi-ts": "< 2" } }, "sha512-29jBRYNdxVGlx5oewFgOrkulZckpIpBIRHth3uHFn1PrL2ucMy52FvWOY3U3dVx2go1Z3kUmMi6lr07iOpUqqA=="], - "@hey-api/codegen-core": ["@hey-api/codegen-core@0.8.2", "", { "dependencies": { "@hey-api/types": "0.1.4", "ansi-colors": "4.1.3", "c12": "3.3.4", "color-support": "1.1.3" } }, "sha512-R2NMf3wq97rh1mjz33WJQU8svz3F0RYUjvx/QzXucjpSqQ3O5huTdDjErG4fMxSr1X+X56NuDrqtGfHmo1TRUQ=="], + "@hey-api/codegen-core": ["@hey-api/codegen-core@0.9.1", "", { "dependencies": { "@hey-api/types": "0.1.4", "ansi-colors": "4.1.3", "c12": "3.3.4", "color-support": "1.1.3" } }, "sha512-s97jL1dgTMuiMHv2BZ1X4Tgd99Mf9GOvGdNqNcGwIMmnR+PgYNoraj4Zvp134MKsNCap/m7k0r0vKKnl56pj4w=="], - "@hey-api/json-schema-ref-parser": ["@hey-api/json-schema-ref-parser@1.4.2", "", { "dependencies": { "@jsdevtools/ono": "7.1.3", "@types/json-schema": "7.0.15", "js-yaml": "4.1.1" } }, "sha512-ZhCFSKI2ipZHEbgmtUHdyddvRU3wJ4elgCfYUC7T7hZa4EivSrVflTQf2w+v3TuaYxR1Y2V2kq3otqTttrrK8Q=="], + "@hey-api/json-schema-ref-parser": ["@hey-api/json-schema-ref-parser@1.4.4", "", { "dependencies": { "@jsdevtools/ono": "7.1.3", "@types/json-schema": "7.0.15", "js-yaml": "4.2.0" } }, "sha512-otmd+zCxbYVBIp/mlMTnGkvlNYLkVKgs3VOIq0kSnenhB1+fRwLPQIeSwyWM6E51oXhUedkYjVsVpkVexeuJOA=="], - "@hey-api/openapi-ts": ["@hey-api/openapi-ts@0.97.3", "", { "dependencies": { "@hey-api/codegen-core": "0.8.2", "@hey-api/json-schema-ref-parser": "1.4.2", "@hey-api/shared": "0.4.5", "@hey-api/spec-types": "0.2.0", "@hey-api/types": "0.1.4", "@lukeed/ms": "2.0.2", "ansi-colors": "4.1.3", "color-support": "1.1.3", "commander": "14.0.3", "get-tsconfig": "4.14.0" }, "peerDependencies": { "typescript": ">=5.5.3 || >=6.0.0 || 6.0.1-rc" }, "bin": { "openapi-ts": "./bin/run.js" } }, "sha512-4sR6/E/POuy7aPZW9DDjhObzZCq7eSJWiW0+epXeKNczoTWEwdOyWFy9Ca/CnXYlZ3oJsrv0ZD0OO+YuczT7CA=="], + "@hey-api/openapi-ts": ["@hey-api/openapi-ts@0.99.0", "", { "dependencies": { "@hey-api/codegen-core": "0.9.1", "@hey-api/json-schema-ref-parser": "1.4.4", "@hey-api/shared": "0.5.0", "@hey-api/spec-types": "0.2.0", "@hey-api/types": "0.1.4", "@lukeed/ms": "2.0.2", "ansi-colors": "4.1.3", "color-support": "1.1.3", "commander": "15.0.0", "get-tsconfig": "4.14.0" }, "peerDependencies": { "typescript": ">=5.5.3 || >=6.0.0 || 6.0.1-rc" }, "bin": { "openapi-ts": "./bin/run.js" } }, "sha512-SePU/5oEWWkvUBYmvzdYRctseoLuskyhs4ET0RvLIcmzc8yLQoA2R+KtBIQ8bPsoSUB0m4E5SmBnl6aGSA0szQ=="], - "@hey-api/shared": ["@hey-api/shared@0.4.5", "", { "dependencies": { "@hey-api/codegen-core": "0.8.2", "@hey-api/json-schema-ref-parser": "1.4.2", "@hey-api/spec-types": "0.2.0", "@hey-api/types": "0.1.4", "ansi-colors": "4.1.3", "cross-spawn": "7.0.6", "open": "11.0.0", "semver": "7.7.4" } }, "sha512-au4eHpBXAe1du0iMp6ESYuEaMS2jsoEyrbcT246btRhI9rMeQFEs7ZjtcMGXGsxhpaR38A8cPGNHx7QOrWAdMw=="], + "@hey-api/shared": ["@hey-api/shared@0.5.0", "", { "dependencies": { "@hey-api/codegen-core": "0.9.1", "@hey-api/json-schema-ref-parser": "1.4.4", "@hey-api/spec-types": "0.2.0", "@hey-api/types": "0.1.4", "ansi-colors": "4.1.3", "cross-spawn": "7.0.6", "open": "11.0.0", "semver": "7.8.4" } }, "sha512-JN/j4Ebh4cJGYIQ5cwWuqe7GeSUyQoz7oC51WqyhKOcrejK6DKZMDkshc5d1eKTRuRL+rjozuRcoUaZZn2DGPw=="], "@hey-api/spec-types": ["@hey-api/spec-types@0.2.0", "", { "dependencies": { "@hey-api/types": "0.1.4" } }, "sha512-ibQ8Is7evMavzr8GNyJCcTg975d8DpaMUyLmOrQ85UBdy1l6t1KuRAwgChAbesJsIlNV6gjmlXruWyegDX18Fg=="], @@ -1014,7 +1014,7 @@ "any-base": ["any-base@1.1.0", "", {}, "sha512-uMgjozySS8adZZYePpaWs8cxB9/kdzmpX6SgJZ+wbz1K5eYk5QMYDVJaZKhxyIHUdnnJkfR7SVgStgH7LkGUyg=="], - "apache-arrow": ["apache-arrow@github:Karakatiza666/arrow-js#ddba834", { "dependencies": { "@swc/helpers": "^0.5.11", "@types/command-line-args": "^5.2.3", "@types/command-line-usage": "^5.0.4", "@types/node": "^25.2.0", "command-line-args": "^6.0.1", "command-line-usage": "^7.0.1", "flatbuffers": "^25.1.24", "json-with-bigint": "^3.5.3", "tslib": "^2.6.2" }, "bin": { "arrow2csv": "bin/arrow2csv.cjs" } }, "Karakatiza666-arrow-js-ddba834", "sha512-N2H0djCXEuIXEXsZyPKGI4NT0xBhH3DdsNNP3Z273ym1h0OKHR6qrGld4l8RnxGjDwFzXhMZDEqEfN0NMFZUTA=="], + "apache-arrow": ["apache-arrow@github:Karakatiza666/arrow-js#ddba834", { "dependencies": { "@swc/helpers": "^0.5.11", "@types/command-line-args": "^5.2.3", "@types/command-line-usage": "^5.0.4", "@types/node": "^25.2.0", "command-line-args": "^6.0.1", "command-line-usage": "^7.0.1", "flatbuffers": "^25.1.24", "json-with-bigint": "^3.5.3", "tslib": "^2.6.2" }, "bin": { "arrow2csv": "bin/arrow2csv.cjs" } }, "Karakatiza666-arrow-js-ddba834"], "apexcharts": ["apexcharts@5.6.0", "", { "dependencies": { "@yr/monotone-cubic-spline": "^1.0.3" } }, "sha512-BZua59yedRsaDfnxkzNrkyLCvluq2c3ZDBIz4joxSKtgr0xDQXQ5dzceMhf/TpTbAjaF+2NYIpLP3BEEIG2s/w=="], @@ -1120,7 +1120,7 @@ "command-line-usage": ["command-line-usage@7.0.4", "", { "dependencies": { "array-back": "^6.2.2", "chalk-template": "^0.4.0", "table-layout": "^4.1.1", "typical": "^7.3.0" } }, "sha512-85UdvzTNx/+s5CkSgBm/0hzP80RFHAa7PsfeADE5ezZF3uHz3/Tqj9gIKGT9PTtpycc3Ua64T0oVulGfKxzfqg=="], - "commander": ["commander@14.0.3", "", {}, "sha512-H+y0Jo/T1RZ9qPP4Eh1pkcQcLRglraJaSLoyOtHxu6AapkjWVCy2Sit1QQ4x3Dng8qDlSsZEet7g5Pq06MvTgw=="], + "commander": ["commander@15.0.0", "", {}, "sha512-z67u4ZhzCL/Tydu1lJARtEZYWbWaN7oYLHbsuzocr6y4N6WZAagG3RQ4FW61V1/0+jImpj293XfrcYnd1qxtPg=="], "common-ui": ["common-ui@workspace:js-packages/common-ui"], @@ -1446,7 +1446,7 @@ "jpeg-js": ["jpeg-js@0.4.4", "", {}, "sha512-WZzeDOEtTOBK4Mdsar0IqEU5sMr3vSV2RqkAIzUEV2BHnUfKGyswWFPFwK5EeDo93K3FohSHbLAjj0s1Wzd+dg=="], - "js-yaml": ["js-yaml@4.1.1", "", { "dependencies": { "argparse": "^2.0.1" }, "bin": { "js-yaml": "bin/js-yaml.js" } }, "sha512-qQKT4zQxXl8lLwBtHMWwaTcGfFOZviOJet3Oy/xmGk2gZH677CJM9EvtfdSkgWcATZhj/55JZ0rmy3myCT5lsA=="], + "js-yaml": ["js-yaml@4.2.0", "", { "dependencies": { "argparse": "^2.0.1" }, "bin": { "js-yaml": "bin/js-yaml.js" } }, "sha512-ePWsvanv0DWuDRsW8dnt+R4jQ31SCRCQ7hhNcPXZPsoBZiemuZNYGf7adZdqX2D86j6rvKp3RpCxVTSb8WQlOw=="], "json-buffer": ["json-buffer@3.0.1", "", {}, "sha512-4bV5BfR2mqfQTJm+V5tPPdf+ZpuhiIvTuAB5g8kcrXOZpTT/QwwVRWBywX1ozr6lEuPdbHxwaJlm9G6mI2sfSQ=="], @@ -2108,6 +2108,8 @@ "@eslint-community/eslint-utils/eslint-visitor-keys": ["eslint-visitor-keys@3.4.3", "", {}, "sha512-wpc+LXeiyiisxPlEkUzU6svyS1frIO3Mgxj1fdy7Pm8Ygzguax2N3Fa/D/ag1WqbOprdI+uY6wMUl8/a2G+iag=="], + "@hey-api/shared/semver": ["semver@7.8.4", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-rUCObTnP32Q08R2uuIrt7r9PlEonuTmtuXYcW6s5kjdlj3xbnwe+21yXptAUYcMAABLkYYTtnmzb3w3EDZfueA=="], + "@jimp/png/pngjs": ["pngjs@6.0.0", "", {}, "sha512-TRzzuFRRmEoSW/p1KVAmiOgPco2Irlah+bGFCeNfJXxxYGwSw7YwAOAcd7X28K/m5bjBWKsC29KyoMfHbypayg=="], "@opentelemetry/otlp-transformer/@opentelemetry/resources": ["@opentelemetry/resources@2.2.0", "", { "dependencies": { "@opentelemetry/core": "2.2.0", "@opentelemetry/semantic-conventions": "^1.29.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.3.0 <1.10.0" } }, "sha512-1pNQf/JazQTMA0BiO5NINUzH0cbLbbl7mntLa4aJNmCCXSj0q03T5ZXXL0zw4G55TjdL9Tz32cznGClf+8zr5A=="], diff --git a/crates/adapters/src/controller/stats.rs b/crates/adapters/src/controller/stats.rs index b9cf8f068ab..8f28f958e29 100644 --- a/crates/adapters/src/controller/stats.rs +++ b/crates/adapters/src/controller/stats.rs @@ -1641,8 +1641,21 @@ impl InputEndpointMetrics { num_transport_errors: self.num_transport_errors.load(Ordering::Relaxed), num_parse_errors: self.num_parse_errors.load(Ordering::Relaxed), end_of_input: self.end_of_input.load(Ordering::Relaxed), + processing_latency_p99_micros: self.processing_latency_p99_micros(), } } + + /// 99th percentile processing latency in microseconds over the sliding + /// histogram's window, or `None` if the endpoint has no samples. See + /// [`ExternalInputEndpointMetrics::processing_latency_p99_micros`] for what + /// the window covers. + fn processing_latency_p99_micros(&self) -> Option { + self.processing_latency_micros_histogram + .lock() + .ok()? + .snapshot() + .quantile(0.99) + } } // Latency histogram creation functions. @@ -3026,3 +3039,85 @@ impl OutputEndpointStatus { self.metrics.total_processed_steps.load(Ordering::Acquire) } } + +#[cfg(test)] +mod test { + use super::InputEndpointMetrics; + + #[test] + fn latency_p99_absent_without_samples() { + let metrics = InputEndpointMetrics::default(); + assert_eq!(metrics.to_api_type().processing_latency_p99_micros, None); + } + + #[test] + fn latency_p99_reports_uniform_bucket() { + let metrics = InputEndpointMetrics::default(); + { + let mut histogram = metrics.processing_latency_micros_histogram.lock().unwrap(); + for _ in 0..100 { + histogram.record(1_500u64); + } + } + // Bucket [1000, 1999] holds 1_500, and `quantile` reports its lower bound. + assert_eq!( + metrics.to_api_type().processing_latency_p99_micros, + Some(1_000) + ); + } + + /// Guards against reporting the median, which hides the slow tail the + /// column exists to surface. + #[test] + fn latency_p99_reports_slow_tail_not_median() { + let metrics = InputEndpointMetrics::default(); + { + let mut histogram = metrics.processing_latency_micros_histogram.lock().unwrap(); + // 90 fast samples against 10 slow ones: the median sits in the fast + // band, rank 99 (ceil(0.99 * 100)) lands in the slow one. + for _ in 0..90 { + histogram.record(1_500u64); + } + for _ in 0..10 { + histogram.record(250_000u64); + } + } + // Bucket [200000, 299999] holds 250_000, and `quantile` reports its + // lower bound. A median would report 1_000 here. + assert_eq!( + metrics.to_api_type().processing_latency_p99_micros, + Some(200_000) + ); + } + + /// Documents the known limit of a percentile over few samples: one slow + /// batch in a hundred stays below rank 99. + #[test] + fn latency_p99_misses_lone_outlier() { + let metrics = InputEndpointMetrics::default(); + { + let mut histogram = metrics.processing_latency_micros_histogram.lock().unwrap(); + for _ in 0..99 { + histogram.record(1_500u64); + } + histogram.record(250_000u64); + } + assert_eq!( + metrics.to_api_type().processing_latency_p99_micros, + Some(1_000) + ); + } + + /// Guards against reporting the completion histogram, which tracks a + /// different span and would silently change the column's meaning. + #[test] + fn latency_p99_ignores_completion_histogram() { + let metrics = InputEndpointMetrics::default(); + metrics + .completion_latency_micros_histogram + .lock() + .unwrap() + .record(750_000u64); + assert_eq!(metrics.to_api_type().processing_latency_p99_micros, None); + } +} diff --git a/crates/feldera-types/src/adapter_stats.rs b/crates/feldera-types/src/adapter_stats.rs index fad7007c8c8..107c5557348 100644 --- a/crates/feldera-types/src/adapter_stats.rs +++ b/crates/feldera-types/src/adapter_stats.rs @@ -215,6 +215,23 @@ pub struct ExternalInputEndpointMetrics { pub num_parse_errors: u64, /// True if end-of-input has been signaled. pub end_of_input: bool, + /// 99th percentile processing latency (from ingesting a batch to finishing processing it) in microseconds. + /// + /// The time from ingesting a batch of records off the wire + /// to the circuit finishing processing them, + /// covering parsing, queuing, and the circuit step. + /// + /// Does not account for completion latency. + /// + /// Taken over the endpoint's sliding histogram, which holds the 10,000 most + /// recent samples spanning at most 10 minutes, whichever bound is reached + /// first. One sample is recorded per completed batch. An endpoint that stops + /// ingesting keeps reporting its last known latency instead of dropping to + /// `None`. + /// + /// `None` until the endpoint records its first sample. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub processing_latency_p99_micros: Option, } /// Input endpoint status information. diff --git a/crates/storage/src/histogram.rs b/crates/storage/src/histogram.rs index f4de7018bd7..fc777592c9d 100644 --- a/crates/storage/src/histogram.rs +++ b/crates/storage/src/histogram.rs @@ -115,6 +115,25 @@ impl ExponentialHistogramSnapshot { pub fn sum(&self) -> u64 { self.sum } + + /// Approximate `q`-quantile (`q` in `0.0..=1.0`), as the lower bound of the + /// containing bucket. `None` if empty. + pub fn quantile(&self, q: f64) -> Option { + let total: u64 = self.buckets.iter().sum(); + if total == 0 { + return None; + } + let rank = ((q.clamp(0.0, 1.0) * total as f64).ceil() as u64).clamp(1, total); + let mut cumulative = 0u64; + for (index, count) in self.buckets.iter().enumerate() { + cumulative += *count; + if cumulative >= rank { + return Some(*bucket_to_range(index).start()); + } + } + // Unreachable: `cumulative` reaches `total >= rank` in the loop. + Some(*bucket_to_range(N_BUCKETS - 1).start()) + } } pub struct Bucket { @@ -278,7 +297,52 @@ impl SlidingHistogram { #[cfg(test)] mod test { - use crate::histogram::{N_BUCKETS, bucket_to_range, number_to_bucket}; + use crate::histogram::{ExponentialHistogram, N_BUCKETS, bucket_to_range, number_to_bucket}; + + #[test] + fn quantile_empty() { + let hist = ExponentialHistogram::new(); + assert_eq!(hist.snapshot().quantile(0.5), None); + assert_eq!(hist.snapshot().quantile(0.99), None); + } + + #[test] + fn quantile_single_sample() { + let hist = ExponentialHistogram::new(); + hist.record(42u64); + // 42 lands in the bucket for [40, 49]; the lower bound is 40. + assert_eq!(hist.snapshot().quantile(0.0), Some(40)); + assert_eq!(hist.snapshot().quantile(0.5), Some(40)); + assert_eq!(hist.snapshot().quantile(1.0), Some(40)); + } + + #[test] + fn quantile_all_same_bucket() { + let hist = ExponentialHistogram::new(); + for _ in 0..100 { + hist.record(5u64); + } + // Values 0..=9 each occupy their own bucket, so 5 is exact. + assert_eq!(hist.snapshot().quantile(0.5), Some(5)); + assert_eq!(hist.snapshot().quantile(0.99), Some(5)); + } + + #[test] + fn quantile_known_distribution() { + let hist = ExponentialHistogram::new(); + // 99 fast samples (bucket [0,0]) and 1 slow sample (bucket [900,999]). + for _ in 0..99 { + hist.record(0u64); + } + hist.record(950u64); + let snap = hist.snapshot(); + // p50 stays with the fast majority. + assert_eq!(snap.quantile(0.5), Some(0)); + // p99 (ceil(0.99*100)=99th value) is still the last fast sample. + assert_eq!(snap.quantile(0.99), Some(0)); + // p100 reaches the lone slow sample: bucket [900, 999], lower bound 900. + assert_eq!(snap.quantile(1.0), Some(900)); + } #[test] fn buckets() { diff --git a/js-packages/web-console/openapi-ts.config.ts b/js-packages/web-console/openapi-ts.config.ts index 1f8e287e497..64a8995d5a3 100644 --- a/js-packages/web-console/openapi-ts.config.ts +++ b/js-packages/web-console/openapi-ts.config.ts @@ -1,8 +1,30 @@ import { defineConfig } from '@hey-api/openapi-ts' +import { + type CustomTypes, + customTypesPlugin, + overrideOpenapiType +} from './src/lib/functions/common/openapi-ts' + +/** + * Hand-written types substituted for the generated ones, keyed by the custom + * `format` that marks a schema. + */ +const customTypes = { + microseconds: { module: '$lib/functions/common/duration', name: 'Microseconds' } +} as const satisfies CustomTypes export default defineConfig({ input: '../../openapi.json', output: './src/lib/services/manager', + parser: { + patch: { + schemas: { + InputEndpointMetrics: (schema) => { + overrideOpenapiType(schema, 'processing_latency_p99_micros', 'microseconds') + } + } + } + }, plugins: [ { name: '@hey-api/client-fetch', @@ -12,6 +34,7 @@ export default defineConfig({ { name: '@hey-api/sdk', responseStyle: 'data' - } + }, + customTypesPlugin(customTypes) ] }) diff --git a/js-packages/web-console/package.json b/js-packages/web-console/package.json index fce7a44b27f..b5c975d993e 100644 --- a/js-packages/web-console/package.json +++ b/js-packages/web-console/package.json @@ -10,7 +10,7 @@ "@fontsource/dm-mono": "5.2.7", "@fortawesome/fontawesome-free": "7.2.0", "@hey-api/client-fetch": "0.13.1", - "@hey-api/openapi-ts": "0.97.3", + "@hey-api/openapi-ts": "0.99.0", "@monaco-editor/loader": "1.7.0", "@playwright/test": "1.58.2", "@poppanator/sveltekit-svg": "6.0.1", diff --git a/js-packages/web-console/src/assets/icons/feldera-material-icons/circle-help.svg b/js-packages/web-console/src/assets/icons/feldera-material-icons/circle-help.svg new file mode 100755 index 00000000000..4d290c0da59 --- /dev/null +++ b/js-packages/web-console/src/assets/icons/feldera-material-icons/circle-help.svg @@ -0,0 +1 @@ + \ No newline at end of file diff --git a/js-packages/web-console/src/lib/components/pipelines/editor/performance/MetricsTables.svelte b/js-packages/web-console/src/lib/components/pipelines/editor/performance/MetricsTables.svelte index 62aca2c1cc2..b758c9c045f 100644 --- a/js-packages/web-console/src/lib/components/pipelines/editor/performance/MetricsTables.svelte +++ b/js-packages/web-console/src/lib/components/pipelines/editor/performance/MetricsTables.svelte @@ -4,7 +4,12 @@ import ClipboardCopyButton from '$lib/components/other/ClipboardCopyButton.svelte' import { count } from '$lib/functions/common/array' import { humanSize } from '$lib/functions/common/string' - import { formatQty } from '$lib/functions/format' + import { formatDuration, formatQty } from '$lib/functions/format' + import { + defaultLatencyColorSpread, + latencyColor, + latencyColorScale + } from '$lib/functions/latencyColor' import type { AggregatedInputEndpointMetrics, AggregatedMetrics, @@ -28,6 +33,17 @@ ) => void } = $props() + // Calculated from individual connectors, so that a relation's color + // stays comparable to the connectors it summarizes. + const inputLatencyScale = $derived( + latencyColorScale( + [...metrics.current.tables.values()].flatMap((data) => + data.connectors.map((connector) => connector.metrics.processing_latency_p99_micros) + ), + defaultLatencyColorSpread + ) + ) + type HealthFilter = 'all' | 'unhealthy' let tableHealthFilter = $state('all') let viewHealthFilter = $state('all') @@ -352,6 +368,20 @@ Transaction Ingested Buffered + + + Latency p99 + + + 99th percentile time from ingesting a batch of records to the circuit finishing processing + it, taken over the last 10 minutes. Colored relative to the other connectors — red marks the + slowest. + Parse errors Transport errors @@ -377,6 +407,15 @@ > {formatQty(m.buffered_records)} {humanSize(m.buffered_bytes)} + + {#if typeof m.processing_latency_p99_micros === 'number'} + + {formatDuration(m.processing_latency_p99_micros)} + + {:else} + + {/if} + {#if m.num_parse_errors > 0 && relation && connectorEndpointName}