diff --git a/livekit-telemetry/Cargo.toml b/livekit-telemetry/Cargo.toml index 0943e1689..c36394dc5 100644 --- a/livekit-telemetry/Cargo.toml +++ b/livekit-telemetry/Cargo.toml @@ -37,6 +37,8 @@ uniffi = ["dep:uniffi"] [dev-dependencies] tokio = { workspace = true, default-features = false, features = ["rt", "rt-multi-thread", "macros", "sync", "time", "test-util", "net", "io-util"] } +# The mock collector test drives the real HTTP stack. +livekit-net = { workspace = true, features = ["native"] } # How CI checks this crate's features, read by # `.github/workflows/feature-combinations-curated.yml` via `cargo metadata`. diff --git a/livekit-telemetry/src/backend_tests.rs b/livekit-telemetry/src/backend_tests.rs new file mode 100644 index 000000000..db8a55c5f --- /dev/null +++ b/livekit-telemetry/src/backend_tests.rs @@ -0,0 +1,978 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! The backend contract, end to end through the pipeline: what the platform hands over (server +//! URL, token) and what the collector answers, and what the core does about each. One test per +//! row of the contract table in the PR description. + +use std::{sync::Arc, time::Duration}; + +use prost::Message; + +use crate::{ + destination::tests::{granted, grantless}, + telemetry::tests::{ + answer, bounded, event_names, offline, start, start_cloud, test_config, FakeTransport, + }, + ExportError, Scope, Telemetry, TelemetryEvent, TelemetryStatus, +}; + +const PROJECT: &str = "wss://p.livekit.cloud"; + +/// A Cloud pipeline with one room connected to [`PROJECT`] with a granted token. +fn connected(transport: &Arc) -> (Telemetry, Scope) { + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + room.set_server(PROJECT, &granted(3600)); + (telemetry, room) +} + +async fn ping(telemetry: &Telemetry, room: &Scope) { + room.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; +} + +fn rpc_status(message: &str, retry_secs: Option) -> Vec { + #[derive(Clone, PartialEq, prost::Message)] + struct Status { + #[prost(int32, tag = "1")] + code: i32, + #[prost(string, tag = "2")] + message: String, + #[prost(message, repeated, tag = "3")] + details: Vec, + } + #[derive(Clone, PartialEq, prost::Message)] + struct RetryInfo { + #[prost(message, optional, tag = "1")] + retry_delay: Option, + } + let details = retry_secs + .map(|seconds| prost_types::Any { + type_url: "type.googleapis.com/google.rpc.RetryInfo".into(), + value: RetryInfo { retry_delay: Some(prost_types::Duration { seconds, nanos: 0 }) } + .encode_to_vec(), + }) + .into_iter() + .collect(); + Status { code: 8, message: message.into(), details }.encode_to_vec() +} + +#[tokio::test(start_paused = true)] +async fn the_ingest_url_and_token_come_from_the_room() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + let token = granted(3600); + room.set_server("wss://p.livekit.cloud/rtc?access_token=secret", &token); + ping(&telemetry, &room).await; + let sent = transport.sent(); + assert_eq!(sent[0].url, "https://p.livekit.cloud/observability/client/logs/otlp/v0"); + assert_eq!(sent[0].headers["Authorization"], format!("Bearer {token}")); + assert!(!sent[0].url.contains("secret"), "path and query of the server URL never leak"); +} + +#[tokio::test(start_paused = true)] +async fn handing_over_the_same_token_again_is_free() { + let transport = FakeTransport::scripted([]); + let (telemetry, room) = connected(&transport); + let token = granted(3600); + room.set_server(PROJECT, &token); + for _ in 0..100 { + room.set_server(PROJECT, &token); + } + ping(&telemetry, &room).await; + assert_eq!(transport.sent().len(), 1, "one upload, no extra wake-ups"); +} + +#[tokio::test(start_paused = true)] +async fn self_hosted_servers_get_nothing_and_nothing_is_kept() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + room.set_server("wss://livekit.example.com", &granted(3600)); + ping(&telemetry, &room).await; + telemetry.emit(TelemetryEvent::new("lk.ping")); // process-level: nobody would take it either + telemetry.flush().await; + assert!(transport.sent().is_empty()); + let stats = telemetry.stats(); + assert_eq!(stats.cached_batches, 0, "not collected, so nothing sits on disk"); + assert_eq!(stats.status, TelemetryStatus::Off); +} + +#[tokio::test(start_paused = true)] +async fn a_room_without_the_observability_grant_sends_nothing() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + room.set_server(PROJECT, &grantless(3600)); + ping(&telemetry, &room).await; + assert!(transport.sent().is_empty(), "no consent, nothing leaves the device"); + assert_eq!(telemetry.stats().cached_batches, 0); +} + +#[tokio::test(start_paused = true)] +async fn an_expired_token_holds_uploads_until_a_fresh_one_arrives() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + room.set_server(PROJECT, &granted(60)); + tokio::time::sleep(Duration::from_secs(61)).await; + ping(&telemetry, &room).await; + assert!(transport.sent().is_empty(), "never sent with a token known to be expired"); + assert_eq!(telemetry.stats().status, TelemetryStatus::Waiting); + assert_eq!(telemetry.stats().dropped, 0, "held, not dropped"); + + let fresh = granted(3600); + room.set_server(PROJECT, &fresh); + tokio::time::sleep(Duration::from_millis(1)).await; + let sent = transport.sent(); + assert_eq!(sent.len(), 1, "the refresh releases the backlog right away"); + assert_eq!(sent[0].headers["Authorization"], format!("Bearer {fresh}")); +} + +#[tokio::test(start_paused = true)] +async fn a_refresh_that_drops_the_grant_keeps_uploading_with_the_granted_token() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + let join = granted(3600); + room.set_server(PROJECT, &join); + room.set_server(PROJECT, &grantless(7200)); // today's server-side refresh + ping(&telemetry, &room).await; + assert_eq!(transport.sent()[0].headers["Authorization"], format!("Bearer {join}")); +} + +#[tokio::test(start_paused = true)] +async fn two_rooms_on_two_projects_never_share_a_token_or_a_destination() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let (a, b) = (telemetry.begin_scope(), telemetry.begin_scope()); + let (token_a, token_b) = (granted(3600), granted(3601)); + a.set_server("wss://a.livekit.cloud", &token_a); + b.set_server("wss://b.livekit.cloud", &token_b); + a.emit(TelemetryEvent::new("custom.a")); + b.emit(TelemetryEvent::new("custom.b")); + telemetry.flush().await; + let sent = transport.sent(); + assert_eq!(sent.len(), 2, "one batch per project"); + for request in &sent { + let names = event_names(request); + if request.url.starts_with("https://a.") { + assert_eq!(names, ["custom.a"]); + assert_eq!(request.headers["Authorization"], format!("Bearer {token_a}")); + } else { + assert!(request.url.starts_with("https://b.")); + assert_eq!(names, ["custom.b"]); + assert_eq!(request.headers["Authorization"], format!("Bearer {token_b}")); + } + } +} + +#[tokio::test(start_paused = true)] +async fn records_before_the_first_connect_wait_and_go_to_that_project() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + telemetry.emit(TelemetryEvent::new("lk.device.capture.failed")); // before any room + telemetry.flush().await; + assert!(transport.sent().is_empty()); + let room = telemetry.begin_scope(); + room.set_server(PROJECT, &granted(3600)); + tokio::time::sleep(Duration::from_millis(1)).await; + assert_eq!(event_names(&transport.sent()[0]), ["lk.device.capture.failed"]); +} + +#[tokio::test(start_paused = true)] +async fn partial_success_counts_the_refused_records_and_never_retries() { + #[derive(Clone, PartialEq, prost::Message)] + struct Partial { + #[prost(int64, tag = "1")] + rejected: i64, + #[prost(string, tag = "2")] + error_message: String, + } + #[derive(Clone, PartialEq, prost::Message)] + struct Response { + #[prost(message, optional, tag = "1")] + partial_success: Option, + } + let body = Response { + partial_success: Some(Partial { rejected: 2, error_message: "too old".into() }), + } + .encode_to_vec(); + let transport = FakeTransport::scripted([answer(200, &[], &body)]); + let (telemetry, room) = connected(&transport); + for _ in 0..3 { + room.emit(TelemetryEvent::new("lk.ping")); + } + telemetry.flush().await; + let stats = telemetry.stats(); + assert_eq!((stats.uploads_sent, stats.dropped_rejected), (1, 2)); + assert_eq!(stats.cached_batches, 0, "accepted: never sent again"); +} + +#[tokio::test(start_paused = true)] +async fn a_bad_request_drops_the_batch() { + for status in [400, 422] { + let transport = FakeTransport::scripted([answer(status, &[], b"nope")]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + ping(&telemetry, &room).await; + let stats = telemetry.stats(); + assert_eq!(stats.dropped_rejected, 1, "{status}: dropped, counted"); + // The second flush: the ping, and (its own owner) the report of the loss. + assert_eq!(transport.sent().len(), 3, "{status}: not retried, uploads carry on"); + } +} + +#[tokio::test(start_paused = true)] +async fn disabled_project_goes_silent_and_purges_its_cache() { + let disabled = answer(401, &[], b"project data recording is disabled by owner"); + let transport = FakeTransport::scripted([disabled]); + let telemetry = start_cloud(test_config(), transport.clone()); + let (a, b) = (telemetry.begin_scope(), telemetry.begin_scope()); + a.set_server("wss://a.livekit.cloud", &granted(3600)); + a.emit(TelemetryEvent::new("custom.a")); + telemetry.flush().await; + a.emit(TelemetryEvent::new("custom.a")); + telemetry.flush().await; + assert_eq!(transport.sent().len(), 1, "never sent again"); + assert_eq!(telemetry.stats().dropped_disabled, 2, "the batch and what followed"); + assert_eq!(telemetry.stats().cached_batches, 0); + + b.set_server("wss://b.livekit.cloud", &granted(3600)); + b.emit(TelemetryEvent::new("custom.b")); + telemetry.flush().await; + assert_eq!(event_names(&transport.sent()[1]), ["custom.b"], "another project is unaffected"); +} + +#[tokio::test(start_paused = true)] +async fn unauthorized_holds_the_batch_until_a_new_token() { + let transport = FakeTransport::scripted([answer(401, &[], b"invalid token")]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + ping(&telemetry, &room).await; + assert_eq!(transport.sent().len(), 1, "the refused token is not tried again"); + let stats = telemetry.stats(); + // Two pings, plus the report of the refusal (process-level: its own batch, same token). + assert_eq!((stats.dropped, stats.cached_batches), (0, 3), "a credential problem loses nothing"); + assert_eq!(stats.status, TelemetryStatus::Waiting); + + room.set_server(PROJECT, &granted(7200)); + tokio::time::sleep(secs(5)).await; + assert_eq!(telemetry.stats().cached_batches, 0, "the new token ships the backlog"); +} + +#[tokio::test(start_paused = true)] +async fn not_found_on_the_derived_endpoint_goes_silent() { + let transport = FakeTransport::scripted([answer(404, &[], b"")]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + ping(&telemetry, &room).await; + assert_eq!(transport.sent().len(), 1); + assert_eq!(telemetry.stats().status, TelemetryStatus::Off); + assert_eq!(telemetry.stats().cached_batches, 0); +} + +#[tokio::test(start_paused = true)] +async fn throttling_honors_retry_after_then_retry_info_then_a_minute() { + let cases: [(crate::ExportResponse, Duration); 4] = [ + (answer(429, &[("Retry-After", "7")], &rpc_status("q", Some(30))).unwrap(), secs(7)), + (answer(429, &[], &rpc_status("q", Some(30))).unwrap(), secs(30)), + (answer(429, &[], &rpc_status("QuotaStatusExceeded", None)).unwrap(), secs(60)), + (answer(503, &[("Retry-After", "12")], b"").unwrap(), secs(12)), + ]; + for (response, wait) in cases { + let transport = FakeTransport::scripted([Ok(response.clone())]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + assert_eq!(telemetry.stats().status, TelemetryStatus::Throttled); + ping(&telemetry, &room).await; // collection goes on during the pause + tokio::time::sleep(wait - Duration::from_millis(10)).await; + assert_eq!(transport.sent().len(), 1, "{}: quiet for {wait:?}", response.status); + tokio::time::sleep(Duration::from_millis(20)).await; + assert!(transport.sent().len() >= 3, "{}: both batches ship after", response.status); + assert_eq!(telemetry.stats().dropped, 0); + } +} + +#[tokio::test(start_paused = true)] +async fn a_retryable_500_waits_for_its_retry_info() { + let transport = FakeTransport::scripted([answer(500, &[], &rpc_status("later", Some(9)))]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + tokio::time::sleep(secs(8)).await; + assert_eq!(transport.sent().len(), 1); + tokio::time::sleep(secs(2)).await; + assert_eq!(transport.sent().len(), 2, "retried after the delay the server named"); + assert_eq!(telemetry.stats().cached_batches, 0); +} + +#[tokio::test(start_paused = true)] +async fn other_server_errors_drop_the_batch() { + for status in [500, 501, 505] { + let transport = FakeTransport::scripted([answer(status, &[], b"boom")]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + assert_eq!(telemetry.stats().dropped_rejected, 1, "{status}: final per OTLP/HTTP"); + } +} + +/// 502/503/504 without a delay: jittered backoff, and running out of patience pauses — it never +/// deletes. Only the cache's age and size bound what is kept. +#[tokio::test(start_paused = true)] +async fn a_failing_server_never_costs_a_batch() { + let failures = [502, 503, 504].into_iter().cycle().take(12).map(|s| answer(s, &[], b"")); + let transport = FakeTransport::scripted(failures); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + tokio::time::sleep(secs(20 * 60)).await; + let stats = telemetry.stats(); + assert_eq!(transport.sent().len(), 13, "twelve failures, then delivered"); + assert_eq!((stats.dropped, stats.cached_batches, stats.upload_failures), (0, 0, 12)); +} + +/// No answer at all (offline, DNS, TLS, connection reset, timeout): 1 s doubling to a 60 s cap, +/// fully jittered, never dropped. +#[tokio::test(start_paused = true)] +async fn no_answer_backs_off_exponentially_with_full_jitter() { + let transport = FakeTransport::scripted(std::iter::repeat_with(offline).take(8)); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + let backoffs = [1.0_f64, 2.0, 4.0, 8.0, 16.0, 32.0, 60.0, 60.0]; + let mut at = Vec::new(); + let start = tokio::time::Instant::now(); + // Twice all eight full backoffs: every retry is due well before, so past it they stopped. + let deadline = Duration::from_secs_f64(2.0 * backoffs.iter().sum::()); + tokio::time::timeout(deadline, async { + while transport.sent().len() < 9 { + let before = transport.sent().len(); + tokio::time::sleep(Duration::from_millis(10)).await; + if transport.sent().len() > before { + at.push(start.elapsed().as_secs_f64()); + } + } + }) + .await + .unwrap_or_else(|_| { + panic!("retry loop: {} of 9 requests by {deadline:?}", transport.sent().len()) + }); + // The flush tick (1 s here) retries too once a wait is over; waits never exceed the backoff. + let gaps: Vec = std::iter::once(at[0]).chain(at.windows(2).map(|w| w[1] - w[0])).collect(); + for (gap, full) in gaps.iter().zip(backoffs) { + assert!(*gap <= full.max(1.0) + 1.05, "gap {gap} over {full} ({gaps:?})"); + } + let stats = telemetry.stats(); + assert_eq!((stats.dropped, stats.cached_batches), (0, 0), "delivered in the end, none lost"); +} + +#[tokio::test(start_paused = true)] +async fn an_invalid_request_is_dropped() { + let invalid = Err(ExportError::Rejected { reason: "invalid URL".into() }); + let transport = FakeTransport::scripted([invalid]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + assert_eq!(telemetry.stats().dropped_rejected, 1); +} + +#[tokio::test(start_paused = true)] +async fn the_override_reaches_a_local_collector_without_cloud_rules() { + let transport = FakeTransport::scripted([]); + let telemetry = start(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + room.set_server("ws://localhost:7880", "dev-token-without-grant"); + ping(&telemetry, &room).await; + let sent = transport.sent(); + assert_eq!(sent[0].url, "http://collector/v1/logs"); + assert!(!sent[0].headers.contains_key("Authorization"), "tokens never go to an override"); +} + +/// A real HTTP exchange: a mock OTLP collector on a local socket answers 429 with `Retry-After`, +/// then 200, through the `livekit-net` HTTP stack the SDKs can use. +#[cfg(feature = "net")] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_real_http_collector_throttles_and_recovers() { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let addr = listener.local_addr().expect("addr"); + let (seen_tx, mut seen) = tokio::sync::mpsc::unbounded_channel(); + tokio::spawn(async move { + let answers = [ + "HTTP/1.1 429 Too Many Requests\r\nRetry-After: 1\r\nConnection: close\r\nContent-Length: 0\r\n\r\n", + "HTTP/1.1 200 OK\r\nConnection: close\r\nContent-Length: 0\r\n\r\n", + ]; + for answer in answers.iter().cycle() { + let Ok((mut socket, _)) = listener.accept().await else { return }; + let mut request = Vec::new(); + let mut chunk = [0u8; 4096]; + // Read the head, then exactly Content-Length bytes of body. + loop { + let n = socket.read(&mut chunk).await.unwrap_or(0); + request.extend_from_slice(&chunk[..n]); + let text = String::from_utf8_lossy(&request).to_string(); + if let Some(head_end) = text.find("\r\n\r\n") { + let length = text[..head_end] + .lines() + .find_map(|l| { + l.to_ascii_lowercase() + .strip_prefix("content-length:") + .map(|v| v.trim().parse::().unwrap_or(0)) + }) + .unwrap_or(0); + if request.len() >= head_end + 4 + length || n == 0 { + let _ = seen_tx.send(text[..head_end].to_owned()); + break; + } + } + if n == 0 { + break; + } + } + let _ = socket.write_all(answer.as_bytes()).await; + let _ = socket.shutdown().await; + } + }); + + let transport = Arc::new(crate::NetTransport::new(livekit_net::testing::native_http_client())); + let (telemetry, exporter) = Telemetry::new(test_config(), transport); + telemetry.override_endpoint(&format!("http://{addr}")); + tokio::spawn(exporter.run()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + let head = bounded(seen.recv()).await.expect("first request"); + assert!(head.starts_with("POST /v1/logs"), "{head}"); + assert!(head.to_ascii_lowercase().contains("content-encoding: gzip")); + assert_eq!(telemetry.stats().status, TelemetryStatus::Throttled); + + bounded(seen.recv()).await.expect("retried request"); + tokio::time::sleep(Duration::from_millis(200)).await; + let stats = telemetry.stats(); + assert_eq!((stats.uploads_sent, stats.dropped, stats.cached_batches), (1, 0, 0)); +} + +fn secs(n: u64) -> Duration { + Duration::from_secs(n) +} + +/// 413: the batch is split in half and retried at once, down to single records; a record too +/// big on its own is dropped and counted. The halves are committed before the original goes. +#[tokio::test(start_paused = true)] +async fn payload_too_large_splits_down_to_single_records() { + let too_large = || answer(413, &[], b""); + // 4 records: 413 → [2, 2]; first half 413 → [1, 1]; the first single 413 → oversized. + let transport = FakeTransport::scripted([too_large(), too_large(), too_large()]); + let (telemetry, room) = connected(&transport); + for n in 0..4 { + room.emit(TelemetryEvent::new(format!("custom.e{n}"))); + } + telemetry.flush().await; + let delivered: Vec = transport.sent()[3..].iter().flat_map(event_names).collect(); + assert_eq!(delivered, ["custom.e1", "custom.e2", "custom.e3"], "in order, e0 alone too big"); + let stats = telemetry.stats(); + assert_eq!((stats.dropped_oversized, stats.dropped, stats.cached_batches), (1, 1, 0)); +} + +/// Repeated 401s while the renewal is pending: every refused token is tried once, never again, +/// and the batch waits through all of them. +#[tokio::test(start_paused = true)] +async fn repeated_unauthorized_answers_wait_for_renewal_without_loss() { + let denied = || answer(401, &[], b"invalid token"); + let transport = FakeTransport::scripted([denied(), denied()]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + room.set_server(PROJECT, &granted(3601)); + tokio::time::sleep(Duration::from_millis(1)).await; + assert_eq!(transport.sent().len(), 2, "the second token is refused too"); + for _ in 0..5 { + ping(&telemetry, &room).await; + } + assert_eq!(transport.sent().len(), 2, "no token left to try: nothing is sent"); + room.set_server(PROJECT, &granted(3602)); + tokio::time::sleep(secs(5)).await; + let stats = telemetry.stats(); + assert_eq!((stats.dropped, stats.cached_batches, stats.uploads_unauthorized), (0, 0, 2)); +} + +/// A server-directed delay is honored in full: a shutdown drains everything else, but not +/// through a `Retry-After`. +#[tokio::test(start_paused = true)] +async fn shutdown_never_cuts_a_server_delay_short() { + let transport = FakeTransport::scripted([answer(503, &[("Retry-After", "600")], b"")]); + let (telemetry, room) = connected(&transport); + ping(&telemetry, &room).await; + ping(&telemetry, &room).await; + telemetry.shutdown().await; + assert_eq!(transport.sent().len(), 1, "no request before the ten minutes are up"); +} + +/// An expired token is a hard hold: the soft holds' one-minute escape never sends with it. +#[tokio::test(start_paused = true)] +async fn hard_holds_have_no_escape_hatch() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + room.set_server(PROJECT, &granted(10)); + let _connecting = room.start(crate::SpanName::Connect, None); // a soft hold on top + tokio::time::sleep(secs(11)).await; + for _ in 0..10 { + ping(&telemetry, &room).await; + tokio::time::sleep(secs(61)).await; + } + assert!(transport.sent().is_empty(), "ten minutes, not a single request"); + assert_eq!(telemetry.stats().dropped, 0); +} + +/// After a restart the cached batches wait for a token of their own project — never another +/// project's — bounded by the cache's 24 h age. +#[tokio::test(start_paused = true)] +async fn a_restart_with_cached_data_and_no_token_waits_for_the_same_project() { + let dir = crate::cache::temp_dir("restart"); + let config = || crate::TelemetryConfig { + storage_dir: Some(dir.to_string_lossy().into_owned()), + ..test_config() + }; + let first = + start_cloud(config(), FakeTransport::scripted(std::iter::repeat_with(offline).take(64))); + let room = first.begin_scope(); + room.set_server("wss://a.livekit.cloud", &granted(3600)); + room.emit(TelemetryEvent::new("custom.yesterday")); + first.flush().await; // offline: stays on disk + first.shutdown().await; // the first launch is over: its exporter has stopped + drop((room, first)); + + let transport = FakeTransport::scripted([]); + let second = start_cloud(config(), transport.clone()); + second.flush().await; + assert!(transport.sent().is_empty(), "no token yet"); + let other = second.begin_scope(); + other.set_server("wss://b.livekit.cloud", &granted(3600)); + other.emit(TelemetryEvent::new("custom.today")); + second.flush().await; + let urls: Vec = transport.sent().iter().map(|r| r.url.clone()).collect(); + assert!(urls.iter().all(|u| u.starts_with("https://b.")), "b's token never carries a's data"); + assert!(!transport.sent().iter().flat_map(event_names).any(|n| n == "custom.yesterday")); + + let again = second.begin_scope(); + again.set_server("wss://a.livekit.cloud", &granted(3600)); + tokio::time::sleep(Duration::from_millis(1)).await; + let last = transport.sent().pop().expect("sent"); + assert!(last.url.starts_with("https://a.") && event_names(&last) == ["custom.yesterday"]); + let _ = std::fs::remove_dir_all(&dir); +} + +/// Finding r1-2: records keep the project they were captured for. A batch cached for project a +/// still goes to a with a's token after the Room reconnects to b; records queued before the +/// switch do too. +#[tokio::test(start_paused = true)] +async fn a_room_switching_projects_takes_nothing_along() { + let transport = FakeTransport::scripted([offline()]); + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + let (a, b) = (granted(3600), granted(3601)); + room.set_server("wss://a.livekit.cloud", &a); + room.emit(TelemetryEvent::new("custom.cached_for_a")); + telemetry.flush().await; // offline: cached for a + room.emit(TelemetryEvent::new("custom.queued_for_a")); // captured for a, not yet encoded + room.set_server("wss://b.livekit.cloud", &b); + room.emit(TelemetryEvent::new("custom.for_b")); + tokio::time::sleep(secs(70)).await; + let sent = transport.sent(); + for request in &sent[1..] { + let names = event_names(request); + let (url, auth) = (&request.url, &request.headers["Authorization"]); + if names.iter().any(|n| n.ends_with("_for_a")) { + assert!(url.starts_with("https://a.") && *auth == format!("Bearer {a}"), "{names:?}"); + assert!(!names.contains(&"custom.for_b".to_owned())); + } + if names.contains(&"custom.for_b".to_owned()) { + assert!(url.starts_with("https://b.") && *auth == format!("Bearer {b}")); + } + } + let all: Vec = sent[1..].iter().flat_map(event_names).collect(); + for name in ["custom.cached_for_a", "custom.queued_for_a", "custom.for_b"] { + assert!(all.contains(&name.to_owned()), "{name} delivered: {all:?}"); + } +} + +/// Finding r1-2: a Room that has no server yet never has its records sent to another Room's +/// project; they go to its own first project once it connects. +#[tokio::test(start_paused = true)] +async fn an_unconnected_room_never_borrows_another_rooms_project() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let (early, connected) = (telemetry.begin_scope(), telemetry.begin_scope()); + early.emit(TelemetryEvent::new("custom.early")); + connected.set_server("wss://a.livekit.cloud", &granted(3600)); + connected.emit(TelemetryEvent::new("custom.a")); + telemetry.flush().await; + tokio::time::sleep(secs(5)).await; + let names: Vec = transport.sent().iter().flat_map(event_names).collect(); + assert!(!names.contains(&"custom.early".to_owned()), "not on a's token: {names:?}"); + + early.set_server("wss://b.livekit.cloud", &granted(3601)); + tokio::time::sleep(secs(5)).await; + let last = + transport.sent().into_iter().find(|r| event_names(r).contains(&"custom.early".to_owned())); + assert!(last.is_some_and(|r| r.url.starts_with("https://b.")), "its own project"); +} + +/// One project's failures pause only that project: the other keeps uploading. +#[tokio::test(start_paused = true)] +async fn a_failing_project_does_not_pause_the_others() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let (a, b) = (telemetry.begin_scope(), telemetry.begin_scope()); + a.set_server("wss://a.livekit.cloud", &granted(3600)); + a.emit(TelemetryEvent::new("custom.a")); + telemetry.flush().await; // a's first batch goes out and is accepted + transport.then(std::iter::repeat_with(|| answer(503, &[], b"")).take(1)); + a.emit(TelemetryEvent::new("custom.a")); + telemetry.flush().await; // a now backs off + b.set_server("wss://b.livekit.cloud", &granted(3600)); + b.emit(TelemetryEvent::new("custom.b")); + telemetry.flush().await; + let last = transport.sent().pop().expect("sent"); + assert!(last.url.starts_with("https://b."), "b is not held by a's backoff"); +} + +/// A reconnect handing over the same URL and token retries an ingest that answered 404. +#[tokio::test(start_paused = true)] +async fn a_reconnect_with_the_same_token_retries_a_404() { + let transport = FakeTransport::scripted([answer(404, &[], b"")]); + let token = granted(3600); + let telemetry = start_cloud(test_config(), transport.clone()); + let room = telemetry.begin_scope(); + room.set_server(PROJECT, &token); + ping(&telemetry, &room).await; + assert_eq!(telemetry.stats().status, TelemetryStatus::Off); + room.set_server(PROJECT, &token); // the reconnect: same pair + ping(&telemetry, &room).await; + assert_eq!(transport.sent().len(), 2, "tried again"); + assert_eq!(telemetry.stats().status, TelemetryStatus::Ok); +} + +/// A one-connection-at-a-time HTTP server answering every request with `answer`, reporting each +/// request head it saw. +#[cfg(feature = "net")] +async fn http_server( + answer: String, +) -> (std::net::SocketAddr, tokio::sync::mpsc::UnboundedReceiver) { + http_server_by_host(move |_| answer.clone()).await +} + +/// Like [`http_server`], answering each request by its `Host` header. +#[cfg(feature = "net")] +async fn http_server_by_host( + answer: impl Fn(&str) -> String + Send + 'static, +) -> (std::net::SocketAddr, tokio::sync::mpsc::UnboundedReceiver) { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let addr = listener.local_addr().expect("addr"); + let (seen_tx, seen) = tokio::sync::mpsc::unbounded_channel(); + tokio::spawn(async move { + loop { + let Ok((mut socket, _)) = listener.accept().await else { return }; + let mut request = Vec::new(); + let mut chunk = [0u8; 4096]; + loop { + let n = socket.read(&mut chunk).await.unwrap_or(0); + request.extend_from_slice(&chunk[..n]); + let text = String::from_utf8_lossy(&request).to_string(); + if let Some(end) = text.find("\r\n\r\n") { + let length = text[..end] + .lines() + .find_map(|l| { + l.to_ascii_lowercase() + .strip_prefix("content-length:") + .map(|v| v.trim().parse::().unwrap_or(0)) + }) + .unwrap_or(0); + if request.len() >= end + 4 + length { + let _ = seen_tx.send(text[..end].to_owned()); + break; + } + } + if n == 0 { + break; + } + } + let text = String::from_utf8_lossy(&request).to_string(); + let host = text + .lines() + .find_map(|l| l.strip_prefix("host: ").or_else(|| l.strip_prefix("Host: "))) + .unwrap_or_default() + .to_owned(); + let _ = socket.write_all(answer(&host).as_bytes()).await; + let _ = socket.shutdown().await; + } + }); + (addr, seen) +} + +/// The transport rule for credentials: a redirect to another port, or to another host on the +/// same port, never carries the token. (A scheme-only change on the same host and explicit port +/// is not covered here: it needs a TLS origin; the derived Cloud URLs are always `https` on the +/// default port.) +#[cfg(feature = "net")] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn redirects_to_another_host_or_port_never_carry_the_token() { + use crate::TelemetryTransport; + let ok = "HTTP/1.1 200 OK\r\nConnection: close\r\nContent-Length: 0\r\n\r\n".to_owned(); + let (target, mut seen) = http_server(ok).await; + for location in [ + format!("http://127.0.0.1:{}/v1/logs", target.port()), // another port + format!("http://localhost:{}/v1/logs", target.port()), // another host + ] { + let redirect = format!( + "HTTP/1.1 307 Temporary Redirect\r\nLocation: {location}\r\nConnection: close\r\nContent-Length: 0\r\n\r\n" + ); + let (origin, _) = http_server(redirect).await; + let transport = crate::NetTransport::new(livekit_net::testing::native_http_client()); + let request = crate::ExportRequest { + url: format!("http://127.0.0.1:{}/v1/logs", origin.port()), + headers: [("Authorization".to_owned(), "Bearer secret".to_owned())].into(), + body: vec![1, 2, 3], + }; + let _ = bounded(transport.send(request)).await; + let head = bounded(seen.recv()).await.expect("followed request head"); + assert!(!head.to_ascii_lowercase().contains("authorization"), "{location}: {head}"); + } + + // Same port, another host: `localhost:` redirects to `127.0.0.1:`. + let (server, mut heads) = http_server_by_host(|host| { + if host.starts_with("localhost") { + let port = host.rsplit(':').next().unwrap_or_default(); + format!( + "HTTP/1.1 307 Temporary Redirect\r\nLocation: http://127.0.0.1:{port}/v1/logs\r\nConnection: close\r\nContent-Length: 0\r\n\r\n" + ) + } else { + "HTTP/1.1 200 OK\r\nConnection: close\r\nContent-Length: 0\r\n\r\n".to_owned() + } + }) + .await; + let transport = crate::NetTransport::new(livekit_net::testing::native_http_client()); + let request = crate::ExportRequest { + url: format!("http://localhost:{}/v1/logs", server.port()), + headers: [("Authorization".to_owned(), "Bearer secret".to_owned())].into(), + body: vec![1, 2, 3], + }; + let _ = bounded(transport.send(request)).await; + let first = bounded(heads.recv()).await.expect("the origin request"); + assert!(first.to_ascii_lowercase().contains("authorization"), "sent to the origin"); + let second = bounded(heads.recv()).await.expect("followed request head"); + assert!(!second.to_ascii_lowercase().contains("authorization"), "{second}"); +} + +/// Finding r2-1: the answer to a batch a Room captured before it connected is attributed to the +/// project the batch was sent to — that Room's first project — never to the latest project: a +/// 404 or a 429 from a leaves b alone. +#[tokio::test(start_paused = true)] +async fn answers_to_pre_connect_batches_are_attributed_to_their_own_project() { + for (answer_a, what) in [(answer(404, &[], b""), "404"), (answer(429, &[], b""), "429")] { + let transport = FakeTransport::scripted([answer_a]); + let telemetry = start_cloud(test_config(), transport.clone()); + let (early, other) = (telemetry.begin_scope(), telemetry.begin_scope()); + early.emit(TelemetryEvent::new("custom.before_connect")); // no project yet + early.set_server("wss://a.livekit.cloud", &granted(3600)); + other.set_server("wss://b.livekit.cloud", &granted(3601)); // b is now the latest + telemetry.flush().await; + assert!(transport.sent()[0].url.starts_with("https://a."), "{what}: sent to a"); + other.emit(TelemetryEvent::new("custom.b")); + telemetry.flush().await; + let to_b = transport.sent().into_iter().find(|r| r.url.starts_with("https://b.")); + assert!(to_b.is_some_and(|r| event_names(&r) == ["custom.b"]), "{what}: b unaffected"); + } +} + +/// Finding r3-3: a Room that captured records before it connected, then went away before they +/// were sent, still gets them sent with its first project's credential. +#[tokio::test(start_paused = true)] +async fn a_gone_rooms_pre_connect_backlog_keeps_its_credential() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + let offline = + crate::DeviceState { network: crate::NetworkType::Unavailable, ..Default::default() }; + telemetry.set_device_state(offline); + let room = telemetry.begin_scope(); + room.emit(TelemetryEvent::new("custom.before_connect")); // captured without a project + room.set_server("wss://a.livekit.cloud", &granted(3600)); + telemetry.flush().await; // cached, host-less; offline: not sent + drop(room); // the Room goes away (its connect failed) + for _ in 0..3 { + telemetry.flush().await; // passes retire what nothing needs any more + } + telemetry.set_device_state(crate::DeviceState::default()); + tokio::time::sleep(secs(2)).await; + let sent = transport + .sent() + .into_iter() + .find(|r| event_names(r).contains(&"custom.before_connect".to_owned())); + assert!(sent.is_some_and(|r| r.url.starts_with("https://a.")), "sent with a's credential"); +} + +/// Codex final review B2: what a Room captured before its first `set_server` — cached before +/// the connect, or encoded after it — names that Room's project on disk, so it replays after a +/// restart with a token of the same project. +#[tokio::test(start_paused = true)] +async fn pre_connect_records_replay_after_a_restart() { + use crate::{DeviceState, NetworkType, TelemetryConfig}; + let dir = crate::cache::temp_dir("pre-connect-restart"); + let config = || TelemetryConfig { + storage_dir: Some(dir.to_string_lossy().into_owned()), + ..test_config() + }; + let (first, exporter) = Telemetry::new(config(), FakeTransport::scripted([])); + let task = tokio::spawn(exporter.run()); + first.set_device_state(DeviceState { network: NetworkType::Unavailable, ..Default::default() }); + let room = first.begin_scope(); + room.emit(TelemetryEvent::new("custom.cached_before")); + first.flush().await; // cached without a project + room.emit(TelemetryEvent::new("custom.encoded_after")); // captured without a project + room.set_server("wss://a.livekit.cloud", &granted(3600)); + first.flush().await; // offline: nothing sent + task.abort(); // killed before going online + let _ = task.await; + drop((room, first)); + + let transport = FakeTransport::scripted([]); + let second = start_cloud(config(), transport.clone()); + let again = second.begin_scope(); + again.set_server("wss://a.livekit.cloud", &granted(3601)); + second.flush().await; + let to_a: Vec = transport + .sent() + .iter() + .filter(|r| r.url.starts_with("https://a.")) + .flat_map(event_names) + .collect(); + for name in ["custom.cached_before", "custom.encoded_after"] { + assert!(to_a.contains(&name.to_owned()), "{name} replayed to a: {to_a:?}"); + } + let _ = std::fs::remove_dir_all(&dir); +} + +/// Whether a journaled rewrite (`@.pend`) is in progress in `dir`. +fn rewriting(dir: &std::path::Path) -> bool { + std::fs::read_dir(dir) + .map(|d| d.flatten().any(|e| e.path().extension().is_some_and(|x| x == "pend"))) + .unwrap_or(false) +} + +/// Final review r2 B2a: the process dies while a pre-connect batch is being bound to its Room's +/// project — once the bound copy is journaled, before or after the commit (the parent's delete) — +/// and a fresh pipeline still replays it to that project, exactly once. +#[tokio::test(start_paused = true)] +async fn a_crash_while_binding_still_replays_to_the_rooms_project() { + use crate::{ + cache::{FileCache, Step}, + DeviceState, NetworkType, TelemetryConfig, + }; + for crash_at in [Step::Delete, Step::Publish] { + let dir = crate::cache::temp_dir("bind-crash"); + let config = || TelemetryConfig { + storage_dir: Some(dir.to_string_lossy().into_owned()), + ..test_config() + }; + let cache = Arc::new(FileCache::open(&dir, 1 << 20).expect("cache")); + let (first, exporter) = + Telemetry::with_cache(config(), FakeTransport::scripted([]), cache.clone()); + let task = tokio::spawn(exporter.run()); + first.set_device_state(DeviceState { + network: NetworkType::Unavailable, + ..Default::default() + }); + let room = first.begin_scope(); + room.emit(TelemetryEvent::new("custom.cached_before")); + first.flush().await; // cached without a project + let journal = dir.clone(); + cache.inject(move |step, _| { + assert!(!(step == crash_at && rewriting(&journal)), "killed while binding"); + Ok(()) + }); + room.set_server("wss://a.livekit.cloud", &granted(3600)); + bounded(first.flush()).await; + assert!(bounded(task).await.is_err(), "{crash_at:?}: killed mid-bind"); + assert!(rewriting(&dir), "{crash_at:?}: the bound copy is journaled"); + drop((room, first)); + + let transport = FakeTransport::scripted([]); + let second = start_cloud(config(), transport.clone()); + let again = second.begin_scope(); + again.set_server("wss://a.livekit.cloud", &granted(3601)); + second.flush().await; + let to_a: Vec = transport + .sent() + .iter() + .filter(|r| r.url.starts_with("https://a.")) + .flat_map(event_names) + .filter(|name| name == "custom.cached_before") + .collect(); + assert_eq!(to_a.len(), 1, "{crash_at:?}: replayed to a once"); + let _ = std::fs::remove_dir_all(&dir); + } +} + +/// Final review r2 B2b: a pre-connect batch whose binding failed is uploaded through the live +/// Room → project map, accepted, and its delete keeps failing; a later chance to bind it must not +/// give it a new id that is sent again. +#[tokio::test(start_paused = true)] +async fn an_accepted_batch_awaiting_its_delete_is_not_rebound_and_resent() { + use std::sync::atomic::{AtomicBool, Ordering::SeqCst}; + + use crate::{ + cache::{FileCache, Step}, + BatchCache, DeviceState, NetworkType, + }; + let dir = crate::cache::temp_dir("bind-undeletable"); + let cache = Arc::new(FileCache::open(&dir, 1 << 20).expect("cache")); + let transport = FakeTransport::scripted([]); + let (telemetry, exporter) = + Telemetry::with_cache(test_config(), transport.clone(), cache.clone()); + tokio::spawn(exporter.run()); + telemetry + .set_device_state(DeviceState { network: NetworkType::Unavailable, ..Default::default() }); + let room = telemetry.begin_scope(); + room.emit(TelemetryEvent::new("custom.once")); + telemetry.flush().await; // cached without a project + let (bind_fails, delete_fails) = + (Arc::new(AtomicBool::new(true)), Arc::new(AtomicBool::new(true))); + let (binds, deletes, journal) = (bind_fails.clone(), delete_fails.clone(), dir.clone()); + cache.inject(move |step, path| { + let bound = path.to_string_lossy().contains("a.livekit.cloud"); + match step { + Step::Write if bound && binds.load(SeqCst) => Err(std::io::Error::other("bind")), + // A plain delete fails; a rewrite's commit (its journal written) goes through. + Step::Delete if !rewriting(&journal) && deletes.load(SeqCst) => { + Err(std::io::Error::other("delete")) + } + _ => Ok(()), + } + }); + room.set_server("wss://a.livekit.cloud", &granted(3600)); + telemetry.set_device_state(DeviceState::default()); + telemetry.flush().await; // binding fails; sent through the live map, accepted; delete fails + let sends = + || transport.sent().iter().flat_map(event_names).filter(|n| n == "custom.once").count(); + assert_eq!(sends(), 1, "precondition: accepted once"); + bind_fails.store(false, SeqCst); + telemetry.flush().await; // the delete fails again; binding could now succeed + telemetry.flush().await; + assert_eq!(sends(), 1, "never sent again this launch"); + delete_fails.store(false, SeqCst); + telemetry.flush().await; + assert!(cache.pending().is_empty(), "deleted once the storage allows"); + assert_eq!(sends(), 1); + let _ = std::fs::remove_dir_all(&dir); +} diff --git a/livekit-telemetry/src/lib.rs b/livekit-telemetry/src/lib.rs index c4045e13e..4b861b487 100644 --- a/livekit-telemetry/src/lib.rs +++ b/livekit-telemetry/src/lib.rs @@ -53,7 +53,7 @@ mod transport; mod destination; /// Entry point and configuration. -#[allow(dead_code)] // `weak_commands` serves `global`; test hooks serve the pipeline tests +#[allow(dead_code)] // `weak_commands` serves the process-wide pipeline (`global`) mod telemetry; mod trace; @@ -70,5 +70,9 @@ pub use telemetry::*; pub use trace::*; pub use transport::*; +/// The backend and device contracts, one test per row of their tables. +#[cfg(test)] +mod backend_tests; + #[cfg(feature = "uniffi")] uniffi::setup_scaffolding!(); diff --git a/livekit-telemetry/src/telemetry.rs b/livekit-telemetry/src/telemetry.rs index a72359419..7a12baa6c 100644 --- a/livekit-telemetry/src/telemetry.rs +++ b/livekit-telemetry/src/telemetry.rs @@ -890,3 +890,1108 @@ fn add_sdk_resource(resource: &mut Vec, sdk: Option<&TelemetryResourc } } } + +#[cfg(test)] +pub(crate) mod tests { + use crate::span::SpanKind; + use crate::{ReconnectReason, RoomIdentity, SpanName, SpanStep}; + use std::{collections::VecDeque, fs, path::Path, sync::Mutex}; + + #[tokio::test(start_paused = true)] + async fn typed_spans_hold_uploads_while_connecting_and_export_when_ended() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + let session = telemetry.begin_scope(); + let span = session + .start(SpanName::Reconnect { reason: ReconnectReason::SignalDisconnected }, None); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + assert!(transport.sent().is_empty(), "an open reconnect holds uploads"); + span.step(SpanStep::Attempt { number: 1, full: false }); + span.end(SpanOutcome::Ok, None); + assert!(span.context().is_some_and(|c| c.trace_id == session.trace_id())); + telemetry.flush().await; + assert!(transport.sent().iter().any(|r| r.url.contains("traces")), "the span is exported"); + } + + pub(crate) fn exported_spans( + transport: &FakeTransport, + ) -> Vec { + transport + .sent() + .iter() + .filter(|r| r.url.contains("traces")) + .flat_map(|r| { + ExportTraceServiceRequest::decode(&gunzip(&r.body)[..]) + .expect("otlp") + .resource_spans + .into_iter() + .flat_map(|rs| rs.scope_spans.into_iter().flat_map(|ss| ss.spans)) + }) + .collect() + } + + #[tokio::test(start_paused = true)] + async fn the_subscribe_span_is_owned_by_the_core() { + use crate::{SpanTrack, TrackSource}; + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + let session = telemetry.begin_scope(); + let track = |sid: &str| SpanTrack { + sid: Some(sid.into()), + kind: TrackKind::Video, + source: TrackSource::Camera, + remote_identity: Some("bob".into()), + }; + // Intent, confirmation, then the first inbound reading with bytes: ok. + session.subscribe_started(track("TR_a")); + session.subscribed(track("TR_a")); + let mut empty = RtcStatsSample::new("TR_a", TrackKind::Video, StreamDirection::Inbound); + empty.bytes = Some(0); + session.record_stats(empty); + let mut media = RtcStatsSample::new("TR_a", TrackKind::Video, StreamDirection::Inbound); + media.bytes = Some(1_500); + session.record_stats(media); + // A second one nobody hears from — not a single reading, no other activity: the core's own + // clock times it out. A third: unpublished before media. + session.subscribe_started(track("TR_b")); + session.subscribe_started(track("TR_c")); + session.track_ended("TR_c"); + tokio::time::sleep(Scope::SUBSCRIBE_TIMEOUT + Duration::from_secs(1)).await; + telemetry.flush().await; + + let spans = exported_spans(&transport); + let by_sid = |sid: &str| { + spans + .iter() + .find(|s| { + s.attributes.iter().any(|kv| { + kv.key == "lk.track.sid" + && kv.value.as_ref().and_then(|v| v.value.clone()) + == Some(Value::StringValue(sid.into())) + }) + }) + .unwrap_or_else(|| panic!("span for {sid}")) + }; + let ok = by_sid("TR_a"); + assert_eq!(ok.name, "lk.subscribe"); + assert_eq!( + ok.events.iter().map(|e| e.name.as_str()).collect::>(), + ["subscribed", "first_media"] + ); + assert!(ok.attributes.iter().any(|kv| kv.key == "lk.outcome" + && kv.value.as_ref().and_then(|v| v.value.clone()) + == Some(Value::StringValue("ok".into())))); + let timed_out = by_sid("TR_b"); + assert_eq!( + timed_out.status.as_ref().map(|s| s.code), + Some(status::StatusCode::Error as i32) + ); + assert!(timed_out.attributes.iter().any(|kv| kv.key == "error.type" + && kv.value.as_ref().and_then(|v| v.value.clone()) + == Some(Value::StringValue("timed_out".into())))); + assert!( + spans.iter().filter(|s| s.name == "lk.subscribe").count() >= 3, + "cancelled exports too" + ); + assert!( + !spans.iter().any(|s| s.attributes.iter().any(|kv| kv.key == "lk.track.sid" + && kv.value.as_ref().and_then(|v| v.value.clone()) + == Some(Value::StringValue("TR_d".into())))), + "still open" + ); + } + + #[tokio::test(start_paused = true)] + async fn room_identity_and_resource_are_typed() { + let transport = FakeTransport::scripted([]); + let mut config = test_config(); + config.sdk = Some(TelemetryResource { + sdk: Sdk::Swift, + sdk_version: "2.16.0".into(), + os_name: "iOS".into(), + os_version: "19.0".into(), + device_model: Some("iPhone17,1".into()), + }); + let telemetry = start(config, transport.clone()); + let session = telemetry.begin_scope(); + session.set_room(RoomIdentity { + sid: Some("RM_a".into()), + name: Some("telemetry".into()), + ..Default::default() + }); + session.emit(TelemetryEvent::new("lk.ping")); + telemetry.device_event(DeviceEvent::AudioInterruption { began: true }); + telemetry.flush().await; + let sent = transport.sent(); + let logs: Vec = sent.iter().flat_map(records).collect(); + let with_room = + logs.iter().find(|r| attribute(r, "lk.room.sid").is_some()).expect("room record"); + assert_eq!(attribute(with_room, "lk.room.sid"), Some(Value::StringValue("RM_a".into()))); + assert!(logs.iter().any(|r| r.body.as_ref().and_then(|b| b.value.clone()) + == Some(Value::StringValue("audio interruption began".into())))); + let decoded = + ExportLogsServiceRequest::decode(&gunzip(&sent[0].body)[..]).expect("valid OTLP"); + let resource = decoded.resource_logs[0].resource.as_ref().expect("resource"); + let value = |key: &str| { + resource + .attributes + .iter() + .find(|kv| kv.key == key) + .and_then(|kv| kv.value.as_ref()) + .and_then(|v| v.value.clone()) + }; + assert_eq!(value("service.name"), Some(Value::StringValue("livekit-client-swift".into()))); + assert_eq!(value("device.model.identifier"), Some(Value::StringValue("iPhone17,1".into()))); + assert!(value("telemetry.sdk.name").is_some()); + } + + #[tokio::test(start_paused = true)] + async fn log_records_apply_the_source_floor_and_carry_code_attributes() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + let line = |source, severity, logger: &str| crate::LogRecord { + severity, + source, + body: "boom".into(), + logger: Some(logger.into()), + function: Some("connect()".into()), + file: Some("Room.swift".into()), + line: Some(42), + timestamp_ns: None, + span_id: None, + }; + telemetry.log(line(LogSource::WebRtc, Severity::Warn, "sctp.cc")); + telemetry.log(line(LogSource::Sdk, Severity::Info, "Room")); + telemetry.log(line(LogSource::Ffi, Severity::Error, "livekit_telemetry::exporter")); + telemetry.log(line(LogSource::Sdk, Severity::Warn, "Room")); + telemetry.log(line(LogSource::WebRtc, Severity::Error, "sctp.cc")); + telemetry.flush().await; + let sent = transport.sent(); + assert_eq!(sent.len(), 1); + let logs = records(&sent[0]); + assert_eq!(logs.len(), 2, "sdk warn + webrtc error; not webrtc warn, sdk info, own module"); + let sdk = &logs[0]; + assert_eq!(attribute(sdk, "lk.log.source"), Some(Value::StringValue("sdk".into()))); + assert_eq!(attribute(sdk, "lk.log.logger"), Some(Value::StringValue("Room".into()))); + assert_eq!( + attribute(sdk, "code.function.name"), + Some(Value::StringValue("connect()".into())) + ); + assert_eq!(attribute(sdk, "code.line.number"), Some(Value::IntValue(42))); + assert_eq!( + sdk.body.as_ref().and_then(|b| b.value.clone()), + Some(Value::StringValue("boom".into())) + ); + assert_eq!(attribute(&logs[1], "lk.log.source"), Some(Value::StringValue("webrtc".into()))); + } + + use prost::Message; + + use super::*; + use crate::{ + cache::temp_dir, + proto::opentelemetry::proto::{ + collector::{logs::v1::ExportLogsServiceRequest, trace::v1::ExportTraceServiceRequest}, + common::v1::any_value::Value, + logs::v1::LogRecord, + trace::v1::{span, status}, + }, + AppState, ExportError, ExportRequest, ExportResponse, SpanOutcome, StreamDirection, + ThermalState, TrackKind, + }; + + /// Fail a test wait that could otherwise hang, naming the line that waited. The bound + /// outlasts one export timeout, which a flush through a hanging transport waits out. + #[track_caller] + pub(crate) fn bounded( + future: F, + ) -> impl std::future::Future { + let caller = std::panic::Location::caller(); + let bound = Duration::from_millis(test_config().export_timeout_ms) + Duration::from_secs(5); + async move { + timeout(bound, future).await.unwrap_or_else(|_| panic!("{caller}: hung for {bound:?}")) + } + } + + #[derive(Default)] + pub(crate) struct FakeTransport { + requests: Mutex>, + script: Mutex>>, + } + + impl FakeTransport { + pub(crate) fn scripted( + results: impl IntoIterator>, + ) -> Arc { + Arc::new(Self { + script: Mutex::new(results.into_iter().collect()), + ..Default::default() + }) + } + pub(crate) fn sent(&self) -> Vec { + self.requests.lock().expect("lock").clone() + } + + /// Queue more answers. + pub(crate) fn then( + &self, + results: impl IntoIterator>, + ) { + self.script.lock().expect("lock").extend(results); + } + } + + #[async_trait::async_trait] + impl TelemetryTransport for FakeTransport { + async fn send(&self, request: ExportRequest) -> Result { + self.requests.lock().expect("lock").push(request); + let scripted = self.script.lock().expect("lock").pop_front(); + scripted.unwrap_or_else(|| Ok(ExportResponse::accepted())) + } + } + + /// An HTTP answer with this status, headers and body. + pub(crate) fn answer( + status: u16, + headers: &[(&str, &str)], + body: &[u8], + ) -> Result { + Ok(ExportResponse { + status, + headers: headers.iter().map(|(k, v)| (k.to_string(), v.to_string())).collect(), + body: body.to_vec(), + }) + } + + /// The old defaults (1 s export, 15 s windows): the pipeline mechanics tests reason in + /// those; the conservative production defaults have their own tests. + pub(crate) fn test_config() -> TelemetryConfig { + TelemetryConfig { flush_interval_ms: 1000, stats_window_ms: 15_000, ..Default::default() } + } + + pub(crate) fn gunzip(body: &[u8]) -> Vec { + use std::io::Read; + let mut out = Vec::new(); + flate2::read::GzDecoder::new(body).read_to_end(&mut out).expect("gzip body"); + out + } + + pub(crate) fn offline() -> Result { + Err(ExportError::Retryable { reason: "offline".into(), retry_after_ms: None }) + } + + pub(crate) fn offline_forever() -> impl Iterator> { + std::iter::repeat_with(offline).take(64) + } + + pub(crate) fn pipeline(transport: Arc) -> Telemetry { + start(test_config(), transport) + } + + pub(crate) fn persisted_pipeline(transport: Arc, dir: &Path) -> Telemetry { + let mut config = test_config(); + config.storage_dir = Some(dir.to_string_lossy().into_owned()); + start(config, transport) + } + + /// A running pipeline whose uploads all go to `http://collector` (the test override). + pub(crate) fn start(config: TelemetryConfig, transport: Arc) -> Telemetry { + let telemetry = start_cloud(config, transport); + telemetry.override_endpoint("http://collector"); + telemetry + } + + /// A running pipeline with LiveKit Cloud routing: nothing uploads before `set_server`. + pub(crate) fn start_cloud(config: TelemetryConfig, transport: Arc) -> Telemetry { + let (telemetry, exporter) = Telemetry::new(config, transport); + tokio::spawn(exporter.run()); + telemetry + } + + pub(crate) fn files_in(dir: &Path) -> usize { + fs::read_dir(dir).map(|d| d.count()).unwrap_or(0) + } + + pub(crate) fn records(request: &ExportRequest) -> Vec { + let decoded = + ExportLogsServiceRequest::decode(&gunzip(&request.body)[..]).expect("valid OTLP"); + decoded.resource_logs[0].scope_logs[0].log_records.clone() + } + + pub(crate) fn event_names(request: &ExportRequest) -> Vec { + records(request).iter().map(|r| r.event_name.clone()).collect() + } + + pub(crate) fn attribute(record: &LogRecord, key: &str) -> Option { + record.attributes.iter().find(|kv| kv.key == key)?.value.as_ref()?.value.clone() + } + + #[tokio::test(start_paused = true)] + async fn batches_events_into_one_otlp_request() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + for _ in 0..3 { + telemetry.emit(TelemetryEvent::new("lk.ping")); + } + telemetry.flush().await; + + let sent = transport.sent(); + assert_eq!(sent.len(), 1); + assert_eq!(sent[0].url, "http://collector/v1/logs"); + assert_eq!(sent[0].headers["Content-Type"], "application/x-protobuf"); + assert_eq!(event_names(&sent[0]), ["lk.ping"; 3]); + let decoded = + ExportLogsServiceRequest::decode(&gunzip(&sent[0].body)[..]).expect("valid OTLP"); + let resource = decoded.resource_logs[0].resource.as_ref().expect("resource"); + assert!(resource.attributes.iter().any(|kv| kv.key == "telemetry.sdk.name")); + assert_eq!(telemetry.stats().dropped, 0); + assert_eq!(telemetry.stats().uploads_sent, 1); + } + + #[tokio::test(start_paused = true)] + async fn failed_upload_waits_in_memory_and_is_retried_after_backoff() { + let transport = FakeTransport::scripted([offline()]); + let telemetry = pipeline(transport.clone()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + assert_eq!(transport.sent().len(), 1, "one attempt: retrying is the backoff's job"); + assert_eq!(telemetry.stats().dropped, 0, "kept in the memory cache, not dropped"); + assert_eq!(telemetry.stats().upload_failures, 1); + assert_eq!(telemetry.stats().cached_batches, 1); + + telemetry.flush().await; + assert_eq!(transport.sent().len(), 1, "backoff: no upload right away"); + + tokio::time::sleep(Duration::from_secs(5)).await; + assert_eq!(transport.sent().len(), 2, "retried once the backoff elapsed, before the tick"); + assert_eq!(event_names(&transport.sent()[1]), ["lk.ping"]); + assert_eq!(telemetry.stats().cached_batches, 0); + } + + #[tokio::test(start_paused = true)] + async fn self_telemetry_report_rides_along_after_problems() { + let transport = FakeTransport::scripted([offline(), offline(), offline()]); + let telemetry = pipeline(transport.clone()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + tokio::time::sleep(Duration::from_secs(61)).await; // backoff over, batch uploads + + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + let sent = transport.sent(); + let last = &sent[sent.len() - 1]; + assert_eq!(event_names(last), ["lk.ping", "lk.telemetry.report"]); + let report = &records(last)[1]; + assert_eq!(attribute(report, "lk.telemetry.uploads.failed"), Some(Value::IntValue(3))); + assert_eq!(attribute(report, "lk.telemetry.uploads.sent"), Some(Value::IntValue(1))); + assert_eq!(attribute(report, "lk.telemetry.cache.batches"), Some(Value::IntValue(0))); + assert_eq!(attribute(report, "lk.telemetry.dropped.queue_full"), None, "zeros omitted"); + + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + let sent = transport.sent(); + assert_eq!(event_names(&sent[sent.len() - 1]), ["lk.ping"], "nothing new to report"); + } + + #[tokio::test(start_paused = true)] + async fn queue_overflow_is_counted_by_reason() { + let mut config = test_config(); + config.max_queue_size = 1; + let telemetry = start(config, FakeTransport::scripted([])); + for _ in 0..3 { + telemetry.emit(TelemetryEvent::new("lk.ping")); + } + let stats = telemetry.stats(); + assert_eq!(stats.dropped_queue_full, 2); + assert_eq!(stats.dropped, 2); + } + + #[tokio::test(start_paused = true)] + async fn rejected_batch_is_dropped_without_retry() { + let transport = + FakeTransport::scripted([Err(ExportError::Rejected { reason: "400".into() })]); + let telemetry = pipeline(transport.clone()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + + assert_eq!(transport.sent().len(), 1); + assert_eq!(telemetry.stats().dropped_rejected, 1); + } + + #[tokio::test(start_paused = true)] + async fn batch_is_written_before_upload_and_replayed_on_next_start() { + let dir = temp_dir("replay"); + let first_transport = FakeTransport::scripted(offline_forever()); + let first = persisted_pipeline(first_transport.clone(), &dir); + first.emit(TelemetryEvent::new("lk.ping")); + first.flush().await; + assert_eq!(first_transport.sent().len(), 1); + assert_eq!(first.stats().dropped, 0); + assert_eq!(files_in(&dir), 1, "written before the first attempt, kept after failure"); + + let second_transport = FakeTransport::scripted([]); + let second = persisted_pipeline(second_transport.clone(), &dir); + second.flush().await; + let sent = second_transport.sent(); + assert_eq!(sent.len(), 1, "replayed on start"); + assert_eq!(event_names(&sent[0]), ["lk.ping"]); + assert_eq!(files_in(&dir), 0); + let _ = fs::remove_dir_all(&dir); + } + + #[tokio::test(start_paused = true)] + async fn throttling_holds_uploads_and_keeps_collecting() { + let dir = temp_dir("throttle"); + let throttled = + Err(ExportError::Retryable { reason: "429".into(), retry_after_ms: Some(5_000) }); + let transport = FakeTransport::scripted([throttled]); + let telemetry = persisted_pipeline(transport.clone(), &dir); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + assert_eq!(transport.sent().len(), 1, "no retries on Retry-After"); + assert_eq!(files_in(&dir), 1, "the throttled batch stays cached"); + assert_eq!(telemetry.stats().dropped, 0); + + // A hold pauses uploads, not collection: what happens during the quiet window is the + // part an operator most wants afterwards. + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + assert_eq!(files_in(&dir), 2, "events inside the window are cached too"); + assert_eq!(telemetry.stats().dropped, 0, "and nothing is thrown away"); + + tokio::time::sleep(Duration::from_secs(6)).await; + assert_eq!(transport.sent().len(), 3, "both cached batches upload after Retry-After"); + assert_eq!(files_in(&dir), 0); + let _ = fs::remove_dir_all(&dir); + } + + /// The cache is the floor under a hold: outlast it and the oldest batches go, counted apart + /// from an ordinary overflow so the report says *why* the session has a hole. + #[tokio::test(start_paused = true)] + async fn a_hold_longer_than_the_cache_reports_what_it_cost() { + let dir = temp_dir("throttle-overflow"); + let throttled = + Err(ExportError::Retryable { reason: "429".into(), retry_after_ms: Some(60_000) }); + let transport = FakeTransport::scripted([throttled]); + let mut config = test_config(); + config.storage_dir = Some(dir.to_string_lossy().into_owned()); + config.max_cache_bytes = 700; // a couple of batches, so the hold overruns it quickly + let telemetry = start(config, transport.clone()); + + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + assert_eq!(telemetry.stats().dropped, 0, "the first batch is cached, not dropped"); + + for _ in 0..8 { + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + } + let stats = telemetry.stats(); + assert!(stats.dropped_throttled > 0, "evictions inside a hold are attributed to it"); + assert_eq!(stats.dropped, stats.dropped_throttled, "and to nothing else"); + let _ = fs::remove_dir_all(&dir); + } + + #[tokio::test(start_paused = true)] + async fn shutdown_offline_keeps_queue_on_disk() { + let dir = temp_dir("spill"); + let transport = FakeTransport::scripted(offline_forever()); + let telemetry = persisted_pipeline(transport.clone(), &dir); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.shutdown().await; + + // One file, or two when the first tick shipped the ping before shutdown added the summary. + let files = files_in(&dir); + assert!( + (1..=2).contains(&files), + "cached before the network was tried, kept after: {files}" + ); + assert_eq!(telemetry.stats().dropped, 0); + let _ = fs::remove_dir_all(&dir); + } + + #[tokio::test(start_paused = true)] + async fn custom_cache_is_used_as_is() { + let cache = Arc::new(MemoryCache::new(1 << 20)); + let transport = FakeTransport::scripted(offline_forever()); + let (telemetry, exporter) = Telemetry::with_cache(test_config(), transport, cache.clone()); + tokio::spawn(exporter.run()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + let pending = cache.pending(); + assert_eq!(pending.len(), 1); + assert!(pending[0].contains("-1-l-"), "id carries count and signal: {}", pending[0]); + } + + #[tokio::test(start_paused = true)] + async fn device_state_emits_change_events_and_stretches_cadence() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + telemetry.set_device_state(DeviceState { + thermal: ThermalState::Critical, + ..DeviceState::default() + }); + telemetry.flush().await; + let sent = transport.sent(); + assert_eq!(sent.len(), 1); + let names = event_names(&sent[0]); + assert!(names.contains(&"lk.device.thermal.changed".to_owned()), "{names:?}"); + assert_eq!( + names.len(), + 4, + "initial value for every known field (battery, low power unknown)" + ); + let thermal = records(&sent[0]) + .into_iter() + .find(|r| r.event_name == "lk.device.thermal.changed") + .expect("thermal event"); + assert_eq!( + attribute(&thermal, "lk.device.thermal.state"), + Some(Value::StringValue("critical".into())) + ); + + // 1 s base interval × 4 under critical thermal pressure. + telemetry.emit(TelemetryEvent::new("lk.ping")); + tokio::time::sleep(Duration::from_secs(2)).await; + assert_eq!(transport.sent().len(), 1, "not yet: cadence stretched to 4 s"); + tokio::time::sleep(Duration::from_millis(2_500)).await; + assert_eq!(transport.sent().len(), 2, "exported on the stretched tick"); + } + + #[tokio::test(start_paused = true)] + async fn requests_are_gzipped_and_low_priority() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + let sent = transport.sent(); + assert_eq!(sent[0].headers["Content-Encoding"], "gzip"); + assert_eq!(sent[0].headers["Priority"], "u=7"); + assert_eq!(event_names(&sent[0]), ["lk.ping"]); + } + + #[tokio::test(start_paused = true)] + async fn uploads_hold_while_connecting_but_never_beyond_the_cap() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + let connect = telemetry.begin_span("lk.connect", SpanKind::Client, None); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + assert!(transport.sent().is_empty(), "the uplink belongs to signaling and ICE"); + assert_eq!(telemetry.stats().cached_batches, 1, "…but the batch is safely cached"); + + tokio::time::sleep(Duration::from_secs(301)).await; + telemetry.flush().await; + assert_eq!(transport.sent().len(), 1, "held 5 min: one batch goes out regardless"); + assert_eq!(telemetry.stats().holds_capped, 1, "…and the starvation is counted"); + + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.end_span(connect, SpanOutcome::Ok, None, Vec::new()); + telemetry.flush().await; + assert_eq!(transport.sent().len(), 3, "connected: the ping and the connect span ship"); + } + + #[tokio::test(start_paused = true)] + async fn bandwidth_limitation_does_not_hold_uploads() { + // WebRTC reports `bandwidth` for minutes during ramp-up and for as long as an encoder + // stalls; holding on it starved a real device of uploads for 8 minutes. + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + let limited = |ms| RtcStatsSample { + quality_limitation_bandwidth_ms: Some(ms), + ..RtcStatsSample::new("TR_1", TrackKind::Video, StreamDirection::Outbound) + }; + telemetry.record_stats(limited(0)); + telemetry.record_stats(limited(800)); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + assert_eq!(transport.sent().len(), 1, "yielding to media is the transport's job"); + } + + #[tokio::test(start_paused = true)] + async fn device_holds_uploads_and_a_backlog_replays_within_the_budget() { + let transport = FakeTransport::scripted([]); + let mut config = test_config(); + config.max_batches_per_upload = 2; + let telemetry = start(config, transport.clone()); + let call = telemetry.begin_scope(); + call.set_server("wss://p.livekit.cloud", "token"); // the budget applies next to a call + telemetry.set_device_state(DeviceState { network_constrained: true, ..Default::default() }); + for _ in 0..5 { + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + } + assert!(transport.sent().is_empty(), "Low Data Mode: record, do not upload"); + assert_eq!(telemetry.stats().cached_batches, 5); + + // Back to normal: the change wakes the exporter, which replays within the budget. + telemetry.set_device_state(DeviceState::default()); + tokio::time::sleep(Duration::from_millis(1)).await; + assert_eq!(transport.sent().len(), 2, "two batches per pass next to a live call"); + tokio::time::sleep(Duration::from_millis(1000)).await; // the next tick + assert_eq!(transport.sent().len(), 4); + // An explicit flush drains the rest (the change event made a sixth batch). + telemetry.flush().await; + assert_eq!(transport.sent().len(), 6, "flush has no per-pass budget"); + telemetry.shutdown().await; + assert_eq!(transport.sent().len(), 7, "shutdown drains too (+ the summary)"); + } + + struct Hanging; + + #[async_trait::async_trait] + impl TelemetryTransport for Hanging { + async fn send(&self, _: ExportRequest) -> Result { + std::future::pending().await + } + } + + #[tokio::test(start_paused = true)] + async fn a_full_queue_flushes_before_the_tick() { + let transport = FakeTransport::scripted([]); + let mut config = test_config(); + config.flush_interval_ms = 60_000; + config.flush_threshold_bytes = 10_000; + let telemetry = start(config, transport.clone()); + tokio::time::sleep(Duration::from_millis(1)).await; // the immediate first tick passes + for _ in 0..3 { + telemetry.emit(TelemetryEvent::new("big").with_body("x".repeat(4_000))); + } + tokio::time::sleep(Duration::from_millis(1)).await; // the wake-up is processed + let sent = transport.sent(); + assert_eq!(sent.len(), 1, "exported on crossing the byte threshold, a minute early"); + assert_eq!(records(&sent[0]).len(), 3); + } + + #[tokio::test(start_paused = true)] + async fn requests_stay_under_the_byte_cap() { + let transport = FakeTransport::scripted([]); + let mut config = test_config(); + config.max_batch_bytes = 10_000; + config.max_batches_per_upload = 10; + let telemetry = start(config, transport.clone()); + for _ in 0..5 { + telemetry.emit(TelemetryEvent::new("big").with_body("x".repeat(4_000))); + } + telemetry.flush().await; + let sent = transport.sent(); + assert_eq!(sent.len(), 3, "5 × ~4 KB under a 10 KB cap: 2 + 2 + 1"); + assert!(sent.iter().all(|request| records(request).len() <= 2)); + } + + #[tokio::test(start_paused = true)] + async fn custom_events_are_namespaced() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + telemetry.emit_custom("acme.checkout", vec![Attribute::new("acme.step", 3i64)]); + telemetry.emit_custom("custom.already", Vec::new()); + telemetry.flush().await; + let sent = transport.sent(); + assert_eq!(event_names(&sent[0]), ["custom.acme.checkout", "custom.already"]); + let record = records(&sent[0]).remove(0); + assert_eq!(attribute(&record, "acme.step"), Some(Value::IntValue(3))); + } + + #[tokio::test(start_paused = true)] + async fn shutdown_leaves_a_session_summary() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + telemetry.shutdown().await; + let sent = transport.sent(); + assert_eq!(sent.len(), 2, "the summary is its own batch when nothing else is queued"); + let report = records(&sent[1]).remove(0); + assert_eq!(report.event_name, "lk.telemetry.report"); + assert_eq!(attribute(&report, "lk.telemetry.uploads.sent"), Some(Value::IntValue(1))); + assert!( + matches!(attribute(&report, "lk.telemetry.uploads.bytes"), Some(Value::IntValue(n)) if n > 0), + "bytes on the wire are reported" + ); + } + + #[tokio::test(start_paused = true)] + async fn cache_eviction_is_counted_as_a_drop() { + let transport = FakeTransport::scripted([offline(), offline(), offline()]); + // Room for exactly one batch: a second push evicts the first. + let (telemetry, exporter) = + Telemetry::with_cache(test_config(), transport.clone(), Arc::new(MemoryCache::new(1))); + tokio::spawn(exporter.run()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; // fails: cached, upload paused + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + let stats = telemetry.stats(); + assert_eq!(stats.cached_batches, 1); + assert_eq!(stats.dropped_cache_full, 1, "the evicted ping is counted, not silently lost"); + } + + #[tokio::test(start_paused = true)] + async fn timeouts_are_counted_apart_from_failures() { + let (telemetry, exporter) = Telemetry::new(test_config(), Arc::new(Hanging)); + telemetry.override_endpoint("http://collector"); + tokio::spawn(exporter.run()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + bounded(telemetry.flush()).await; // one attempt, bounded by export_timeout in paused time + let stats = telemetry.stats(); + assert_eq!(stats.upload_timeouts, 1); + assert_eq!(stats.status, TelemetryStatus::Paused, "backing off"); + assert_eq!(stats.upload_failures, 0); + assert_eq!(stats.cached_batches, 1, "kept for the next attempt"); + } + + #[tokio::test(start_paused = true)] + async fn sessions_have_their_own_trace_and_attributes() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + telemetry.set_attribute("acme.tenant", Some("t1".into())); + let a = telemetry.begin_scope(); + let b = telemetry.begin_scope(); + assert_ne!(a.trace_id(), b.trace_id()); + assert_ne!(a.trace_id(), telemetry.trace_id(), "the process has its own session"); + a.set_room(RoomIdentity { sid: Some("RM_a".into()), ..Default::default() }); + b.set_room(RoomIdentity { sid: Some("RM_b".into()), ..Default::default() }); + let span = a.start(SpanName::Connect, None); + let span_id = span.context().expect("bound").span_id; + a.emit(TelemetryEvent::new("lk.ping")); + b.emit(TelemetryEvent::new("lk.ping")); + // A warn record from the SDK logger, inside room A's connect: no session handle, just + // the ambient span id — the core files it under A. + telemetry.emit( + TelemetryEvent::new("").with_severity(Severity::Warn).with_body("hmm").in_span(span_id), + ); + telemetry.emit(TelemetryEvent::new("lk.device.thermal.changed")); + span.end(SpanOutcome::Ok, None); + telemetry.flush().await; + + let sent = transport.sent(); + assert_eq!(sent.len(), 4, "a batch per owner: A's span, A's logs, B's, the process's"); + let logs: Vec = + sent.iter().filter(|r| r.url.ends_with("logs")).flat_map(records).collect(); + assert_eq!(logs.len(), 4); + assert_eq!(hex(&logs[0].trace_id), a.trace_id()); + assert_eq!(attribute(&logs[0], "lk.room.sid"), Some(Value::StringValue("RM_a".into()))); + assert_eq!(attribute(&logs[0], "session.id"), Some(Value::StringValue(a.trace_id()))); + assert_eq!(hex(&logs[1].trace_id), a.trace_id(), "resolved through the span"); + assert_eq!(hex(&logs[2].trace_id), b.trace_id()); + assert_eq!(attribute(&logs[2], "lk.room.sid"), Some(Value::StringValue("RM_b".into()))); + assert_eq!(hex(&logs[3].trace_id), telemetry.trace_id(), "device state: process session"); + assert_eq!(attribute(&logs[3], "lk.room.sid"), None); + assert!( + logs.iter() + .all(|r| attribute(r, "acme.tenant") == Some(Value::StringValue("t1".into()))), + "a pipeline-wide attribute reaches every session" + ); + let traces_request = sent.iter().find(|r| r.url.ends_with("traces")).expect("traces"); + let traces = + ExportTraceServiceRequest::decode(&gunzip(&traces_request.body)[..]).expect("otlp"); + let otlp_span = &traces.resource_spans[0].scope_spans[0].spans[0]; + assert_eq!(hex(&otlp_span.trace_id), a.trace_id()); + assert!(otlp_span.attributes.iter().any(|a| a.key == "lk.room.sid")); + + // A record that names a span which has already ended — and been exported — is still that + // session's: the SDK's log path hops threads, the span does not wait for it. + telemetry.emit( + TelemetryEvent::new("") + .with_severity(Severity::Error) + .with_body("late") + .in_span(span_id), + ); + telemetry.flush().await; + let late = &records(transport.sent().last().expect("sent"))[0]; + assert_eq!(hex(&late.trace_id), a.trace_id(), "filed under the ended span's session"); + assert_eq!(late.span_id, span_id.to_be_bytes().to_vec()); + } + + #[tokio::test(start_paused = true)] + async fn uploads_wait_for_a_destination() { + let transport = FakeTransport::scripted([]); + let telemetry = start_cloud(test_config(), transport.clone()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + tokio::time::sleep(Duration::from_secs(600)).await; + assert!(transport.sent().is_empty(), "no destination: nothing leaves, no hold cap either"); + assert_eq!(telemetry.stats().cached_batches, 1); + assert_eq!(telemetry.stats().status, TelemetryStatus::Waiting); + + let token = crate::destination::tests::granted(3600); + let session = telemetry.begin_scope(); + session.set_server("wss://x.livekit.cloud", &token); + tokio::time::sleep(Duration::from_millis(1)).await; + let sent = transport.sent(); + assert_eq!(sent.len(), 1, "cached batches ship as soon as the destination is known"); + assert_eq!(sent[0].url, "https://x.livekit.cloud/observability/client/logs/otlp/v0"); + assert_eq!(sent[0].headers["Authorization"], format!("Bearer {token}")); + session.start(SpanName::Publish, None).end(SpanOutcome::Ok, None); + telemetry.flush().await; + assert_eq!( + transport.sent()[1].url, + "https://x.livekit.cloud/observability/client/traces/otlp/v0" + ); + } + + #[tokio::test(start_paused = true)] + async fn debug_and_info_logs_never_leave_the_device() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + telemetry.emit(TelemetryEvent::new("").with_severity(Severity::Info).with_body("noise")); + telemetry.emit(TelemetryEvent::new("").with_severity(Severity::Error).with_body("boom")); + telemetry.flush().await; + + let sent = transport.sent(); + let records = records(&sent[0]); + assert_eq!(records.len(), 1); + assert_eq!(records[0].event_name, "", "a log line, not an event"); + assert_eq!(records[0].severity_text, "ERROR"); + assert_eq!( + records[0].body.as_ref().and_then(|b| b.value.clone()), + Some(Value::StringValue("boom".into())) + ); + } + + #[tokio::test(start_paused = true)] + async fn flood_guard_caps_events_but_not_stats_windows() { + let mut config = test_config(); + config.max_events_per_10min = 1; + config.stats_window_ms = 1_000; + let transport = FakeTransport::scripted([]); + let telemetry = start(config, transport.clone()); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.record_stats(RtcStatsSample::new( + "TR_1", + TrackKind::Audio, + StreamDirection::Inbound, + )); + assert_eq!(telemetry.stats().dropped_rate_limited, 1); + + tokio::time::sleep(Duration::from_millis(1_100)).await; // stats window closes + telemetry.flush().await; + let names: Vec = transport.sent().iter().flat_map(event_names).collect(); + assert_eq!(names.iter().filter(|n| *n == "lk.ping").count(), 1); + assert!(names.contains(&"lk.rtc.stats.sample".to_owned()), "{names:?}"); + assert!(names.contains(&"lk.telemetry.report".to_owned()), "rate limiting is reported"); + } + + #[tokio::test(start_paused = true)] + async fn stats_readings_are_windowed_into_one_event() { + let mut config = test_config(); + config.stats_window_ms = 2_000; + let transport = FakeTransport::scripted([]); + let telemetry = start(config, transport.clone()); + for (bytes, jitter) in [(100, 1.0), (200, 3.0), (300, 2.0)] { + let mut sample = + RtcStatsSample::new("TR_1", TrackKind::Video, StreamDirection::Inbound); + sample.bytes = Some(bytes); + sample.jitter_ms = Some(jitter); + sample.codec = Some("video/VP8".into()); + telemetry.record_stats(sample); + } + telemetry.flush().await; + assert!(transport.sent().is_empty(), "windows do not flush early"); + + tokio::time::sleep(Duration::from_millis(2_100)).await; + telemetry.flush().await; + let sent = transport.sent(); + let window = records(&sent[0]) + .into_iter() + .find(|r| r.event_name == "lk.rtc.stats.sample") + .expect("window"); + assert_eq!(attribute(&window, "lk.track.kind"), Some(Value::StringValue("video".into()))); + assert_eq!( + attribute(&window, "lk.rtc.codec"), + Some(Value::StringValue("video/VP8".into())) + ); + assert_eq!(attribute(&window, "lk.rtc.bytes"), Some(Value::IntValue(300))); + assert_eq!(attribute(&window, "lk.rtc.samples"), Some(Value::IntValue(3))); + assert_eq!(attribute(&window, "lk.rtc.jitter_ms.avg"), Some(Value::DoubleValue(2.0))); + } + + #[tokio::test(start_paused = true)] + async fn shutdown_closes_open_stats_windows() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + telemetry.record_stats(RtcStatsSample::new( + "TR_1", + TrackKind::Audio, + StreamDirection::Outbound, + )); + telemetry.shutdown().await; + let names: Vec = transport.sent().iter().flat_map(event_names).collect(); + assert_eq!( + names, + ["lk.rtc.stats.sample", "lk.telemetry.report"], + "window + shutdown summary" + ); + } + + #[tokio::test(start_paused = true)] + async fn session_attributes_are_attached_to_every_record() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + telemetry.set_attribute("lk.room.sid", Some("RM_1".into())); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.emit(TelemetryEvent::new("lk.ping").with_attribute("lk.room.sid", "RM_override")); + telemetry.flush().await; + let first = records(&transport.sent()[0]); + assert_eq!(attribute(&first[0], "lk.room.sid"), Some(Value::StringValue("RM_1".into()))); + assert_eq!( + attribute(&first[1], "lk.room.sid"), + Some(Value::StringValue("RM_override".into())), + "an explicit attribute wins" + ); + telemetry.set_attribute("lk.room.sid", None); + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.flush().await; + assert_eq!(attribute(&records(&transport.sent()[1])[0], "lk.room.sid"), None); + } + + #[tokio::test(start_paused = true)] + async fn spans_export_as_traces_under_the_session_trace_id() { + let transport = FakeTransport::scripted([]); + let telemetry = start(test_config(), transport.clone()); + telemetry.override_endpoint("http://c/observability/client/logs/otlp/v0"); + let connect = telemetry.begin_span("lk.connect", SpanKind::Client, None); + telemetry.add_span_event(connect, "ws_open", vec![]); + telemetry.emit( + TelemetryEvent::new("") + .with_severity(Severity::Error) + .with_body("boom") + .in_span(connect), + ); + telemetry.end_span( + connect, + SpanOutcome::Error, + Some("timeout".into()), + vec![Attribute::new("lk.connect.attempt", 1i64)], + ); + telemetry.flush().await; + + let sent = transport.sent(); + assert_eq!(sent.len(), 2, "one logs batch, one traces batch"); + let traces = + sent.iter().find(|r| r.url.ends_with("/traces/otlp/v0")).expect("traces request"); + assert_eq!( + traces.url, "http://c/observability/client/traces/otlp/v0", + "derived from logs endpoint" + ); + let decoded = + ExportTraceServiceRequest::decode(&gunzip(&traces.body)[..]).expect("valid OTLP"); + let otlp_span = &decoded.resource_spans[0].scope_spans[0].spans[0]; + assert_eq!(otlp_span.name, "lk.connect"); + assert_eq!(otlp_span.kind, span::SpanKind::Client as i32); + assert_eq!(hex(&otlp_span.trace_id), telemetry.trace_id()); + assert_eq!(otlp_span.span_id, connect.to_be_bytes().to_vec()); + assert!(otlp_span.parent_span_id.is_empty()); + assert!(otlp_span.end_time_unix_nano >= otlp_span.start_time_unix_nano); + assert_eq!(otlp_span.events[0].name, "ws_open"); + assert_eq!( + otlp_span.status.as_ref().map(|s| s.code), + Some(status::StatusCode::Error as i32) + ); + assert_eq!(otlp_span.status.as_ref().map(|s| s.message.as_str()), Some("timeout")); + let attr = |key: &str| { + otlp_span + .attributes + .iter() + .find(|kv| kv.key == key) + .and_then(|kv| kv.value.as_ref()?.value.clone()) + }; + assert_eq!(attr("lk.outcome"), Some(Value::StringValue("error".into()))); + assert_eq!(attr("error.type"), Some(Value::StringValue("timeout".into()))); + assert_eq!(attr("lk.connect.attempt"), Some(Value::IntValue(1))); + + let logs = sent.iter().find(|r| r.url.ends_with("/logs/otlp/v0")).expect("logs request"); + let record = &records(logs)[0]; + assert_eq!(hex(&record.trace_id), telemetry.trace_id(), "every record carries the trace"); + assert_eq!( + record.span_id, + connect.to_be_bytes().to_vec(), + "and the span it was emitted in" + ); + } + + #[tokio::test(start_paused = true)] + async fn cancelled_spans_keep_status_unset() { + let transport = FakeTransport::scripted([]); + let telemetry = pipeline(transport.clone()); + let publish = telemetry.begin_span("lk.publish", SpanKind::Internal, None); + telemetry.end_span(publish, SpanOutcome::Cancelled, None, vec![]); + telemetry.flush().await; + let sent = transport.sent(); + assert_eq!(sent[0].url, "http://collector/v1/traces"); + let decoded = + ExportTraceServiceRequest::decode(&gunzip(&sent[0].body)[..]).expect("valid OTLP"); + let otlp_span = &decoded.resource_spans[0].scope_spans[0].spans[0]; + assert_eq!( + otlp_span.status.as_ref().map(|s| s.code), + Some(status::StatusCode::Unset as i32) + ); + let outcome = otlp_span + .attributes + .iter() + .find(|kv| kv.key == "lk.outcome") + .and_then(|kv| kv.value.as_ref()?.value.clone()); + assert_eq!(outcome, Some(Value::StringValue("cancelled".into()))); + } + + pub(crate) fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|b| format!("{b:02x}")).collect() + } + + #[tokio::test(start_paused = true)] + async fn background_transition_flushes_while_foreground_waits_for_the_tick() { + assert_eq!(DeviceState::default().app_state, AppState::Foreground); + assert_eq!(TelemetryConfig::default().flush_interval_ms, 60_000); + for app_state in [AppState::Foreground, AppState::Background] { + let transport = FakeTransport::scripted([]); + let telemetry = start(TelemetryConfig::default(), transport.clone()); + telemetry.set_device_state(DeviceState::default()); + telemetry.flush().await; + tokio::time::sleep(Duration::from_millis(10)).await; + let initial = transport.sent().len(); + let background = app_state == AppState::Background; + let names = || transport.sent().iter().flat_map(event_names).collect::>(); + + telemetry.emit(TelemetryEvent::new("lk.ping")); + telemetry.set_device_state(DeviceState { app_state, ..Default::default() }); + tokio::time::sleep(Duration::from_millis(10)).await; + assert_eq!(transport.sent().len(), initial + usize::from(background), "{app_state:?}"); + assert_eq!(names().contains(&"lk.ping".to_owned()), background, "{app_state:?}"); + + telemetry.emit(TelemetryEvent::new("lk.next")); + telemetry.set_device_state(DeviceState { app_state, ..Default::default() }); + tokio::time::sleep(Duration::from_secs(59)).await; + assert!(!names().contains(&"lk.next".to_owned()), "{app_state:?}: before the tick"); + tokio::time::sleep(Duration::from_secs(2)).await; + assert!(names().contains(&"lk.ping".to_owned()), "{app_state:?}"); + assert_eq!( + names().contains(&"lk.next".to_owned()), + !background, + "{app_state:?}: the background cadence is doubled" + ); + tokio::time::sleep(Duration::from_secs(60)).await; + assert!(names().contains(&"lk.next".to_owned()), "{app_state:?}: periodic flush"); + telemetry.shutdown().await; + } + } + + /// Finding 12: more finished spans than one batch holds all reach the cache (and the wire) + /// at shutdown; none stay behind in memory uncounted. + #[tokio::test(start_paused = true)] + async fn shutdown_drains_every_finished_span_batch() { + let transport = FakeTransport::scripted([]); + let mut config = test_config(); + config.max_batch_size = 2; + let telemetry = start(config, transport.clone()); + tokio::time::sleep(Duration::from_millis(10)).await; // the start-up tick is behind us + let session = telemetry.begin_scope(); + for _ in 0..5 { + session.start(SpanName::Publish, None).end(SpanOutcome::Ok, None); + } + telemetry.shutdown().await; + assert_eq!(exported_spans(&transport).len(), 5, "three batches: 2 + 2 + 1"); + assert_eq!(telemetry.stats().dropped, 0); + } +}