From a2cfc50a6b1cbb10559110b54aab697df97886f1 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 4 Sep 2026 17:11:08 +0000 Subject: [PATCH 1/8] feat(sip): add SIP telephony REST endpoints and models Co-authored-by: nash --- src/models/mod.rs | 2 + src/models/sip.rs | 471 ++++++++++++++++++++++++++++++++++++++++++++++ src/video/mod.rs | 1 + src/video/sip.rs | 137 ++++++++++++++ 4 files changed, 611 insertions(+) create mode 100644 src/models/sip.rs create mode 100644 src/video/sip.rs diff --git a/src/models/mod.rs b/src/models/mod.rs index 1a15aa7..e0b98bb 100644 --- a/src/models/mod.rs +++ b/src/models/mod.rs @@ -7,8 +7,10 @@ mod call; mod shared; +mod sip; mod user; pub use call::*; pub use shared::*; +pub use sip::*; pub use user::*; diff --git a/src/models/sip.rs b/src/models/sip.rs new file mode 100644 index 0000000..22b8906 --- /dev/null +++ b/src/models/sip.rs @@ -0,0 +1,471 @@ +//! SIP (telephony) request/response models. +//! +//! Field names track the getstream-go JSON tags so a later OpenAPI codegen pass +//! can replace these transparently. Response types derive `Default` + +//! `#[serde(default)]` so partial payloads deserialize cleanly. + +use std::collections::HashMap; + +use serde::{Deserialize, Serialize}; + +use super::shared::{CustomData, Timestamp}; + +// SIP trunks + +/// `create_sip_trunk` request (`CreateSIPTrunkRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct CreateSipTrunkRequest { + /// Name of the SIP trunk. + pub name: String, + /// Phone numbers associated with this SIP trunk. + pub numbers: Vec, + /// Optional password for SIP trunk authentication. + #[serde(skip_serializing_if = "Option::is_none")] + pub password: Option, + /// Optional list of allowed IPv4/IPv6 addresses or CIDR blocks. + #[serde(skip_serializing_if = "Option::is_none")] + pub allowed_ips: Option>, +} + +impl CreateSipTrunkRequest { + /// Build a create request for a named trunk with its phone numbers. + pub fn new(name: impl Into, numbers: impl IntoIterator) -> Self { + Self { + name: name.into(), + numbers: numbers.into_iter().collect(), + ..Default::default() + } + } +} + +/// `update_sip_trunk` request (`UpdateSIPTrunkRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct UpdateSipTrunkRequest { + /// Name of the SIP trunk. + pub name: String, + /// Phone numbers associated with this SIP trunk. + pub numbers: Vec, + /// Optional password for SIP trunk authentication. + #[serde(skip_serializing_if = "Option::is_none")] + pub password: Option, + /// Optional list of allowed IPv4/IPv6 addresses or CIDR blocks. + #[serde(skip_serializing_if = "Option::is_none")] + pub allowed_ips: Option>, +} + +/// A SIP trunk (`SIPTrunkResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct SipTrunkResponse { + pub id: String, + pub name: String, + /// Password for SIP trunk authentication. + pub password: String, + /// Username for SIP trunk authentication. + pub username: String, + /// The URI for the SIP trunk. + pub uri: String, + pub numbers: Vec, + pub allowed_ips: Vec, + pub created_at: Timestamp, + pub updated_at: Timestamp, +} + +/// `create_sip_trunk` response (`CreateSIPTrunkResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct CreateSipTrunkResponse { + pub duration: String, + pub sip_trunk: Option, +} + +/// `update_sip_trunk` response (`UpdateSIPTrunkResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct UpdateSipTrunkResponse { + pub duration: String, + pub sip_trunk: Option, +} + +/// `delete_sip_trunk` response (`DeleteSIPTrunkResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct DeleteSipTrunkResponse { + pub duration: String, +} + +/// `list_sip_trunks` response (`ListSIPTrunksResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct ListSipTrunksResponse { + pub duration: String, + pub sip_trunks: Vec, +} + +// SIP inbound routing rules + +/// SIP caller settings (`SIPCallerConfigsRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct SipCallerConfigsRequest { + /// Unique identifier for the caller (handlebars template). + pub id: String, + /// Custom data associated with the caller (values are handlebars templates). + #[serde(skip_serializing_if = "Option::is_none")] + pub custom_data: Option, +} + +/// SIP call settings (`SIPCallConfigsRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct SipCallConfigsRequest { + /// Custom data associated with the call. + #[serde(skip_serializing_if = "Option::is_none")] + pub custom_data: Option, +} + +/// Direct routing rule call settings (`SIPDirectRoutingRuleCallConfigsRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct SipDirectRoutingRuleCallConfigsRequest { + /// ID of the call (handlebars template). + pub call_id: String, + /// Type of the call. + pub call_type: String, +} + +/// PIN protection settings (`SIPPinProtectionConfigsRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct SipPinProtectionConfigsRequest { + #[serde(skip_serializing_if = "Option::is_none")] + pub default_pin: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub enabled: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub max_attempts: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub required_pin_digits: Option, +} + +/// PIN routing rule call settings (`SIPInboundRoutingRulePinConfigsRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct SipInboundRoutingRulePinConfigsRequest { + #[serde(skip_serializing_if = "Option::is_none")] + pub custom_webhook_url: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pin_failed_attempt_prompt: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pin_hangup_prompt: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pin_prompt: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pin_success_prompt: Option, +} + +/// `create_sip_inbound_routing_rule` request (`CreateSIPInboundRoutingRuleRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct CreateSipInboundRoutingRuleRequest { + pub name: String, + pub trunk_ids: Vec, + pub caller_configs: SipCallerConfigsRequest, + #[serde(skip_serializing_if = "Option::is_none")] + pub called_numbers: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub caller_numbers: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub call_configs: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub direct_routing_configs: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pin_protection_configs: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pin_routing_configs: Option, +} + +/// `update_sip_inbound_routing_rule` request (`UpdateSIPInboundRoutingRuleRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct UpdateSipInboundRoutingRuleRequest { + pub name: String, + pub trunk_ids: Vec, + pub caller_configs: SipCallerConfigsRequest, + #[serde(skip_serializing_if = "Option::is_none")] + pub called_numbers: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub caller_numbers: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub call_configs: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub direct_routing_configs: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pin_protection_configs: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pin_routing_configs: Option, +} + +/// SIP call settings response (`SIPCallConfigsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct SipCallConfigsResponse { + pub custom_data: CustomData, +} + +/// SIP caller settings response (`SIPCallerConfigsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct SipCallerConfigsResponse { + pub id: String, + pub custom_data: CustomData, +} + +/// Direct routing rule call settings response +/// (`SIPDirectRoutingRuleCallConfigsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct SipDirectRoutingRuleCallConfigsResponse { + pub call_id: String, + pub call_type: String, +} + +/// PIN protection settings response (`SIPPinProtectionConfigsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct SipPinProtectionConfigsResponse { + pub enabled: bool, + pub default_pin: Option, + pub max_attempts: Option, + pub required_pin_digits: Option, +} + +/// PIN routing rule call settings response +/// (`SIPInboundRoutingRulePinConfigsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct SipInboundRoutingRulePinConfigsResponse { + pub custom_webhook_url: Option, + pub pin_failed_attempt_prompt: Option, + pub pin_hangup_prompt: Option, + pub pin_prompt: Option, + pub pin_success_prompt: Option, +} + +/// A SIP inbound routing rule (`SIPInboundRoutingRuleResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct SipInboundRoutingRuleResponse { + pub id: String, + pub name: String, + pub called_numbers: Vec, + pub trunk_ids: Vec, + pub caller_numbers: Vec, + pub created_at: Timestamp, + pub updated_at: Timestamp, + pub call_configs: Option, + pub caller_configs: Option, + pub direct_routing_configs: Option, + pub pin_protection_configs: Option, + pub pin_routing_configs: Option, +} + +/// `create_sip_inbound_routing_rule` response (`SIPInboundRoutingRuleResponse` +/// envelope, returned directly for create). +/// +/// The create endpoint returns the rule fields inline alongside `duration`. +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct CreateSipInboundRoutingRuleResponse { + pub duration: String, + pub id: String, + pub name: String, + pub called_numbers: Vec, + pub trunk_ids: Vec, + pub caller_numbers: Vec, + pub created_at: Timestamp, + pub updated_at: Timestamp, + pub call_configs: Option, + pub caller_configs: Option, + pub direct_routing_configs: Option, + pub pin_protection_configs: Option, + pub pin_routing_configs: Option, +} + +/// `update_sip_inbound_routing_rule` response (`UpdateSIPInboundRoutingRuleResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct UpdateSipInboundRoutingRuleResponse { + pub duration: String, + pub sip_inbound_routing_rule: Option, +} + +/// `delete_sip_inbound_routing_rule` response (`DeleteSIPInboundRoutingRuleResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct DeleteSipInboundRoutingRuleResponse { + pub duration: String, +} + +/// `list_sip_inbound_routing_rules` response (`ListSIPInboundRoutingRuleResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct ListSipInboundRoutingRuleResponse { + pub duration: String, + pub sip_inbound_routing_rules: Vec, +} + +// SIP auth / resolve + +/// `resolve_sip_auth` request (`ResolveSipAuthRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct ResolveSipAuthRequest { + pub sip_caller_number: String, + pub sip_trunk_number: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub from_host: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub source_ip: Option, +} + +/// `resolve_sip_auth` response (`ResolveSipAuthResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct ResolveSipAuthResponse { + /// Authentication result: `password`, `accept`, or `no_trunk_found`. + pub auth_result: String, + pub duration: String, + pub password: Option, + pub trunk_id: Option, + pub username: Option, +} + +/// SIP digest challenge authentication data (`SIPChallengeRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct SipChallengeRequest { + #[serde(skip_serializing_if = "Option::is_none")] + pub algorithm: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub charset: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cnonce: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub method: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub nc: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub nonce: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub opaque: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub realm: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub response: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub stale: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub uri: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub userhash: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub username: Option, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub domain: Vec, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub qop: Vec, +} + +/// `resolve_sip_inbound` request (`ResolveSipInboundRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct ResolveSipInboundRequest { + pub sip_caller_number: String, + pub sip_trunk_number: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub routing_number: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub trunk_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub challenge: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub sip_headers: Option>, +} + +/// Credentials for SIP inbound call authentication (`SipInboundCredentials`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct SipInboundCredentials { + pub api_key: String, + pub call_id: String, + pub call_type: String, + pub token: String, + pub user_id: String, + pub call_custom_data: CustomData, + pub user_custom_data: CustomData, +} + +/// `resolve_sip_inbound` response (`ResolveSipInboundResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct ResolveSipInboundResponse { + pub duration: String, + pub credentials: SipInboundCredentials, + pub sip_routing_rule: Option, + pub sip_trunk: Option, +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use super::*; + + #[test] + fn create_trunk_omits_absent_optionals_and_keeps_required_fields() { + let value = serde_json::to_value(CreateSipTrunkRequest::new( + "primary", + ["+15551230000".to_owned()], + )) + .expect("request should serialize"); + assert_eq!( + value, + json!({ "name": "primary", "numbers": ["+15551230000"] }) + ); + } + + #[test] + fn create_routing_rule_serializes_required_and_nested_configs() { + let request = CreateSipInboundRoutingRuleRequest { + name: "rule-1".to_owned(), + trunk_ids: vec!["trunk-1".to_owned()], + caller_configs: SipCallerConfigsRequest { + id: "{{caller_number}}".to_owned(), + custom_data: None, + }, + direct_routing_configs: Some(SipDirectRoutingRuleCallConfigsRequest { + call_id: "support".to_owned(), + call_type: "default".to_owned(), + }), + ..Default::default() + }; + let value = serde_json::to_value(&request).expect("request should serialize"); + assert_eq!( + value, + json!({ + "name": "rule-1", + "trunk_ids": ["trunk-1"], + "caller_configs": { "id": "{{caller_number}}" }, + "direct_routing_configs": { "call_id": "support", "call_type": "default" } + }) + ); + } + + #[test] + fn trunk_response_deserializes_partial_payload() { + let response: CreateSipTrunkResponse = serde_json::from_value(json!({ + "duration": "1.2ms", + "sip_trunk": { + "id": "trunk-1", + "name": "primary", + "numbers": ["+15551230000"] + } + })) + .expect("response should deserialize"); + let trunk = response.sip_trunk.expect("trunk present"); + assert_eq!(trunk.id, "trunk-1"); + assert_eq!(trunk.numbers, vec!["+15551230000".to_owned()]); + assert!(trunk.allowed_ips.is_empty()); + } +} diff --git a/src/video/mod.rs b/src/video/mod.rs index 8eca82c..ccd85d7 100644 --- a/src/video/mod.rs +++ b/src/video/mod.rs @@ -1,6 +1,7 @@ //! Video coordinator REST: [`VideoClient`] and [`Call`]. mod call; +mod sip; pub use call::Call; diff --git a/src/video/sip.rs b/src/video/sip.rs new file mode 100644 index 0000000..31df9b9 --- /dev/null +++ b/src/video/sip.rs @@ -0,0 +1,137 @@ +//! SIP (telephony) coordinator REST endpoints on [`VideoClient`]. + +use reqwest::Method; + +use super::VideoClient; +use crate::client::Client; +use crate::error::Result; +use crate::models::{ + CreateSipInboundRoutingRuleRequest, CreateSipInboundRoutingRuleResponse, CreateSipTrunkRequest, + CreateSipTrunkResponse, DeleteSipInboundRoutingRuleResponse, DeleteSipTrunkResponse, + ListSipInboundRoutingRuleResponse, ListSipTrunksResponse, ResolveSipAuthRequest, + ResolveSipAuthResponse, ResolveSipInboundRequest, ResolveSipInboundResponse, + UpdateSipInboundRoutingRuleRequest, UpdateSipInboundRoutingRuleResponse, UpdateSipTrunkRequest, + UpdateSipTrunkResponse, +}; + +const TRUNKS: &str = "/api/v2/video/sip/inbound_trunks"; +const TRUNK_BY_ID: &str = "/api/v2/video/sip/inbound_trunks/{id}"; +const RULES: &str = "/api/v2/video/sip/inbound_routing_rules"; +const RULE_BY_ID: &str = "/api/v2/video/sip/inbound_routing_rules/{id}"; + +impl VideoClient { + // SIP inbound trunks + + /// List SIP inbound trunks (`GET /api/v2/video/sip/inbound_trunks`). + pub async fn list_sip_trunks(&self) -> Result { + self.client + .request::<(), _>(Method::GET, TRUNKS, &[], None) + .await + } + + /// Create a SIP inbound trunk (`POST /api/v2/video/sip/inbound_trunks`). + pub async fn create_sip_trunk( + &self, + request: CreateSipTrunkRequest, + ) -> Result { + self.client + .request(Method::POST, TRUNKS, &[], Some(&request)) + .await + } + + /// Update a SIP inbound trunk (`PUT /api/v2/video/sip/inbound_trunks/{id}`). + pub async fn update_sip_trunk( + &self, + id: &str, + request: UpdateSipTrunkRequest, + ) -> Result { + let path = Client::build_path(TRUNK_BY_ID, &[("id", id)]); + self.client + .request(Method::PUT, &path, &[], Some(&request)) + .await + } + + /// Delete a SIP inbound trunk (`DELETE /api/v2/video/sip/inbound_trunks/{id}`). + pub async fn delete_sip_trunk(&self, id: &str) -> Result { + let path = Client::build_path(TRUNK_BY_ID, &[("id", id)]); + self.client + .request::<(), _>(Method::DELETE, &path, &[], None) + .await + } + + // SIP inbound routing rules + + /// List SIP inbound routing rules + /// (`GET /api/v2/video/sip/inbound_routing_rules`). + pub async fn list_sip_inbound_routing_rules( + &self, + ) -> Result { + self.client + .request::<(), _>(Method::GET, RULES, &[], None) + .await + } + + /// Create a SIP inbound routing rule + /// (`POST /api/v2/video/sip/inbound_routing_rules`). + pub async fn create_sip_inbound_routing_rule( + &self, + request: CreateSipInboundRoutingRuleRequest, + ) -> Result { + self.client + .request(Method::POST, RULES, &[], Some(&request)) + .await + } + + /// Update a SIP inbound routing rule + /// (`PUT /api/v2/video/sip/inbound_routing_rules/{id}`). + pub async fn update_sip_inbound_routing_rule( + &self, + id: &str, + request: UpdateSipInboundRoutingRuleRequest, + ) -> Result { + let path = Client::build_path(RULE_BY_ID, &[("id", id)]); + self.client + .request(Method::PUT, &path, &[], Some(&request)) + .await + } + + /// Delete a SIP inbound routing rule + /// (`DELETE /api/v2/video/sip/inbound_routing_rules/{id}`). + pub async fn delete_sip_inbound_routing_rule( + &self, + id: &str, + ) -> Result { + let path = Client::build_path(RULE_BY_ID, &[("id", id)]); + self.client + .request::<(), _>(Method::DELETE, &path, &[], None) + .await + } + + // SIP auth / resolve + + /// Resolve SIP authentication requirements for an inbound call + /// (`POST /api/v2/video/sip/auth`). + pub async fn resolve_sip_auth( + &self, + request: ResolveSipAuthRequest, + ) -> Result { + self.client + .request(Method::POST, "/api/v2/video/sip/auth", &[], Some(&request)) + .await + } + + /// Resolve SIP inbound routing (`POST /api/v2/video/sip/resolve`). + pub async fn resolve_sip_inbound( + &self, + request: ResolveSipInboundRequest, + ) -> Result { + self.client + .request( + Method::POST, + "/api/v2/video/sip/resolve", + &[], + Some(&request), + ) + .await + } +} From a9d57f6488f4f7376559a4ade872fe6e8b8a6d41 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 4 Sep 2026 17:11:46 +0000 Subject: [PATCH 2/8] feat(stats): add advanced call statistics and reporting endpoints Co-authored-by: nash --- src/models/mod.rs | 2 + src/models/stats.rs | 435 ++++++++++++++++++++++++++++++++++++++++++++ src/video/call.rs | 235 ++++++++++++++++++++++++ src/video/mod.rs | 1 + src/video/stats.rs | 107 +++++++++++ 5 files changed, 780 insertions(+) create mode 100644 src/models/stats.rs create mode 100644 src/video/stats.rs diff --git a/src/models/mod.rs b/src/models/mod.rs index e0b98bb..9ee067b 100644 --- a/src/models/mod.rs +++ b/src/models/mod.rs @@ -8,9 +8,11 @@ mod call; mod shared; mod sip; +mod stats; mod user; pub use call::*; pub use shared::*; pub use sip::*; +pub use stats::*; pub use user::*; diff --git a/src/models/stats.rs b/src/models/stats.rs new file mode 100644 index 0000000..af74678 --- /dev/null +++ b/src/models/stats.rs @@ -0,0 +1,435 @@ +//! Advanced call statistics and reporting request/response models. +//! +//! Field names track the getstream-go JSON tags. Deeply nested analytics +//! payloads are kept as [`serde_json::Value`] to stay robust across server-side +//! schema additions, matching the existing stats/report types in +//! [`super::call`]. Response types derive `Default` + `#[serde(default)]`. + +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +use super::shared::{CustomData, SortParamRequest, Timestamp}; + +// Active calls status + +/// Aggregate counts for the current active-calls snapshot (`ActiveCallsSummary`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct ActiveCallsSummary { + pub active_calls: i32, + pub active_publishers: i32, + pub active_subscribers: i32, + pub participants: i32, +} + +/// `get_active_calls_status` response (`GetActiveCallsStatusResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct GetActiveCallsStatusResponse { + pub duration: String, + pub start_time: Timestamp, + pub end_time: Timestamp, + /// Detailed join/publisher/subscriber metrics (opaque, schema-versioned). + pub metrics: Option, + pub summary: Option, +} + +// Aggregate call stats + +/// `query_aggregate_call_stats` request (`QueryAggregateCallStatsRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct QueryAggregateCallStatsRequest { + #[serde(skip_serializing_if = "Option::is_none")] + pub from: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub to: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub report_types: Option>, +} + +/// `query_aggregate_call_stats` response (`QueryAggregateCallStatsResponse`). +/// +/// Each report is an opaque, schema-versioned analytics bundle. +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct QueryAggregateCallStatsResponse { + pub duration: String, + pub call_duration_report: Option, + pub call_participant_count_report: Option, + pub calls_per_day_report: Option, + pub network_metrics_report: Option, + pub quality_score_report: Option, + pub sdk_usage_report: Option, + pub user_feedback_report: Option, +} + +// Call session stats + +/// `query_call_session_stats` request (`QueryCallSessionStatsRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct QueryCallSessionStatsRequest { + #[serde(skip_serializing_if = "Option::is_none")] + pub limit: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub next: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub prev: Option, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub sort: Vec, + #[serde(skip_serializing_if = "std::collections::HashMap::is_empty")] + pub filter_conditions: CustomData, +} + +/// `query_call_session_stats` response (`QueryCallSessionStatsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct QueryCallSessionStatsResponse { + pub duration: String, + /// Per-session stat summaries (opaque, schema-versioned). + pub call_stats: Vec, + pub next: Option, + pub prev: Option, +} + +// Session participant stats (call-scoped) + +/// Query params for `get_call_session_participant_stats_details`. +#[derive(Debug, Clone, Default)] +pub struct GetCallSessionParticipantStatsDetailsRequest { + pub since: Option, + pub until: Option, + pub max_points: Option, +} + +/// `get_call_session_participant_stats_details` response +/// (`GetCallSessionParticipantStatsDetailsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct GetCallSessionParticipantStatsDetailsResponse { + pub duration: String, + pub call_id: String, + pub call_session_id: String, + pub call_type: String, + pub user_id: String, + pub user_session_id: String, + pub publisher: Option, + pub subscriber: Option, + pub timeframe: Option, + pub user: Option, +} + +/// Query params for `query_call_session_participant_stats`. +#[derive(Debug, Clone, Default)] +pub struct QueryCallSessionParticipantStatsRequest { + pub limit: Option, + pub prev: Option, + pub next: Option, + pub sort: Vec, + pub filter_conditions: CustomData, +} + +/// `query_call_session_participant_stats` response +/// (`QueryCallSessionParticipantStatsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct QueryCallSessionParticipantStatsResponse { + pub duration: String, + pub call_id: String, + pub call_session_id: String, + pub call_type: String, + /// Per-participant stat summaries (opaque, schema-versioned). + pub participants: Vec, + pub counts: Value, + pub call_started_at: Option, + pub call_ended_at: Option, + pub next: Option, + pub prev: Option, + pub tmp_data_source: Option, + pub call_events: Vec, +} + +/// Query params for `get_call_session_participant_stats_timeline`. +#[derive(Debug, Clone, Default)] +pub struct GetCallSessionParticipantStatsTimelineRequest { + pub start_time: Option, + pub end_time: Option, + pub severity: Vec, +} + +/// `get_call_session_participant_stats_timeline` response +/// (`QueryCallSessionParticipantStatsTimelineResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct QueryCallSessionParticipantStatsTimelineResponse { + pub duration: String, + pub call_id: String, + pub call_session_id: String, + pub call_type: String, + pub user_id: String, + pub user_session_id: String, + /// Timeline events (opaque, schema-versioned). + pub events: Vec, +} + +// Participant session metrics / sessions (call-scoped) + +/// Query params for `get_call_participant_session_metrics`. +#[derive(Debug, Clone, Default)] +pub struct GetCallParticipantSessionMetricsRequest { + pub since: Option, + pub until: Option, +} + +/// `get_call_participant_session_metrics` response +/// (`GetCallParticipantSessionMetricsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct GetCallParticipantSessionMetricsResponse { + pub duration: String, + pub is_publisher: Option, + pub is_subscriber: Option, + pub joined_at: Option, + pub publisher_type: Option, + pub user_id: Option, + pub user_session_id: Option, + /// Per-track publish metrics (opaque, schema-versioned). + pub published_tracks: Vec, + pub client: Option, +} + +/// Query params for `query_call_participant_sessions`. +#[derive(Debug, Clone, Default)] +pub struct QueryCallParticipantSessionsRequest { + pub limit: Option, + pub prev: Option, + pub next: Option, + pub filter_conditions: CustomData, +} + +/// `query_call_participant_sessions` response +/// (`QueryCallParticipantSessionsResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct QueryCallParticipantSessionsResponse { + pub duration: i64, + pub call_id: String, + pub call_session_id: String, + pub call_type: String, + pub total_participant_duration: i64, + pub total_participant_sessions: i64, + /// Per-participant-session details (opaque, schema-versioned). + pub participants_sessions: Vec, + pub next: Option, + pub prev: Option, + pub session: Option, +} + +// Daily digest + +/// Query params for `get_daily_digest`. +#[derive(Debug, Clone, Default)] +pub struct GetDailyDigestRequest { + pub date: Option, + pub target_app_id: Option, +} + +/// `get_daily_digest` response (`GetDailyDigestResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct GetDailyDigestResponse { + pub duration: String, + pub date: String, + /// Readiness status: `ready`, `pending`, `failed`, `future_date`, `expired`. + pub status: String, + pub generated_at: Option, + pub retry_after: Option, + pub revision: Option, + pub schema_version: Option, + pub digest_kinds: Vec, + /// Per-broadcast digests (opaque, present only when `status` is `ready`). + pub broadcasts: Vec, + /// Per-call-session summaries (opaque, present only when `status` is `ready`). + pub call_sessions: Vec, + pub broadcast_rollup: Option, +} + +// User feedback + +/// `query_user_feedback` request (`QueryUserFeedbackRequest`). +/// +/// `full` is sent as a query parameter; the remaining fields form the JSON body. +#[derive(Debug, Clone, Default, Serialize)] +pub struct QueryUserFeedbackRequest { + /// Return full feedback records. Sent as a query parameter. + #[serde(skip)] + pub full: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub limit: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub next: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub prev: Option, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub sort: Vec, + #[serde(skip_serializing_if = "std::collections::HashMap::is_empty")] + pub filter_conditions: CustomData, +} + +/// A single user feedback record (`UserFeedbackResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct UserFeedbackResponse { + pub cid: String, + pub rating: i32, + pub reason: String, + pub sdk: String, + pub sdk_version: String, + pub session_id: String, + pub user_id: String, + pub platform: Value, + pub custom: CustomData, +} + +/// `query_user_feedback` response (`QueryUserFeedbackResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct QueryUserFeedbackResponse { + pub duration: String, + pub user_feedback: Vec, + pub next: Option, + pub prev: Option, +} + +// Client call events + +/// A single client-side telemetry event (`ClientEvent`). +/// +/// Every field is optional; which fields are required depends on the event's +/// `stage`/`event_type`. See the serverside API reference for the per-stage +/// requirements. +#[derive(Debug, Clone, Default, Serialize)] +pub struct ClientEvent { + #[serde(skip_serializing_if = "Option::is_none")] + pub stage: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub stage_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub event_type: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub outcome: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub id: Option, + #[serde(rename = "type", skip_serializing_if = "Option::is_none")] + pub call_type: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub call_session_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub coordinator_connect_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub join_attempt_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub join_reason: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub elapsed_time: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub retry_count_attempt: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub retry_failure_code: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub retry_failure_reason: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub peer_connection: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub ice_state: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub sfu_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub was_previously_connected: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub previously_connected_timestamp: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub camera_permission_status: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub microphone_permission_status: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub screen_share_status: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub track_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub sdk_version: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub user_agent: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub user_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub timestamp: Option, +} + +/// `report_client_call_event` request (`ReportClientCallEventRequest`). +#[derive(Debug, Clone, Default, Serialize)] +pub struct ReportClientCallEventRequest { + /// Client-side events to report (1–100 per request). + pub events: Vec, +} + +/// `report_client_call_event` response (`ReportClientEventResponse`). +#[derive(Debug, Clone, Default, Deserialize)] +#[serde(default)] +pub struct ReportClientEventResponse { + pub duration: String, +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use super::*; + + #[test] + fn user_feedback_request_excludes_full_from_body() { + let value = serde_json::to_value(QueryUserFeedbackRequest { + full: Some(true), + limit: Some(10), + ..Default::default() + }) + .expect("request should serialize"); + assert_eq!(value, json!({ "limit": 10 })); + } + + #[test] + fn client_event_uses_type_wire_name_and_omits_absent_fields() { + let value = serde_json::to_value(ClientEvent { + stage: Some("JoinInitiated".to_owned()), + call_type: Some("default".to_owned()), + id: Some("call-1".to_owned()), + ..Default::default() + }) + .expect("event should serialize"); + assert_eq!( + value, + json!({ "stage": "JoinInitiated", "type": "default", "id": "call-1" }) + ); + } + + #[test] + fn aggregate_stats_request_omits_absent_optionals() { + let value = serde_json::to_value(QueryAggregateCallStatsRequest { + report_types: Some(vec!["call_duration_report".to_owned()]), + ..Default::default() + }) + .expect("request should serialize"); + assert_eq!(value, json!({ "report_types": ["call_duration_report"] })); + } + + #[test] + fn participant_sessions_response_tolerates_integer_duration() { + let response: QueryCallParticipantSessionsResponse = serde_json::from_value(json!({ + "duration": 12, + "call_id": "call-1", + "total_participant_sessions": 3 + })) + .expect("response should deserialize"); + assert_eq!(response.duration, 12); + assert_eq!(response.total_participant_sessions, 3); + } +} diff --git a/src/video/call.rs b/src/video/call.rs index 4306487..ac9ff60 100644 --- a/src/video/call.rs +++ b/src/video/call.rs @@ -663,6 +663,173 @@ impl Call { .await } + /// Retrieve per-participant session metrics for one participant session. + /// + /// `GET .../call/{type}/{id}/session/{session}/participant/{user}/{user_session}/details/track` + pub async fn get_call_participant_session_metrics( + &self, + session: &str, + user: &str, + user_session: &str, + request: GetCallParticipantSessionMetricsRequest, + ) -> Result { + let mut query = Vec::new(); + if let Some(value) = request.since.as_ref() { + query.push(("since".to_owned(), timestamp_query(value))); + } + if let Some(value) = request.until.as_ref() { + query.push(("until".to_owned(), timestamp_query(value))); + } + let path = self.path( + "/session/{session}/participant/{user}/{user_session}/details/track", + &[ + ("session", session), + ("user", user), + ("user_session", user_session), + ], + ); + self.client + .request::<(), _>(Method::GET, &path, &query, None) + .await + } + + /// List participant sessions for one call session. + /// + /// `GET .../call/{type}/{id}/session/{session}/participant_sessions` + pub async fn query_call_participant_sessions( + &self, + session: &str, + request: QueryCallParticipantSessionsRequest, + ) -> Result { + let mut query = Vec::new(); + if let Some(limit) = request.limit { + query.push(("limit".to_owned(), limit.to_string())); + } + if let Some(prev) = request.prev { + query.push(("prev".to_owned(), prev)); + } + if let Some(next) = request.next { + query.push(("next".to_owned(), next)); + } + if let Some(encoded) = filter_conditions_query(&request.filter_conditions)? { + query.push(("filter_conditions".to_owned(), encoded)); + } + let path = self.path( + "/session/{session}/participant_sessions", + &[("session", session)], + ); + self.client + .request::<(), _>(Method::GET, &path, &query, None) + .await + } + + /// Retrieve detailed participant stats time series for one participant session. + /// + /// `GET .../call_stats/{type}/{id}/{session}/participant/{user}/{user_session}/details` + pub async fn get_call_session_participant_stats_details( + &self, + session: &str, + user: &str, + user_session: &str, + request: GetCallSessionParticipantStatsDetailsRequest, + ) -> Result { + let mut query = Vec::new(); + if let Some(since) = request.since { + query.push(("since".to_owned(), since)); + } + if let Some(until) = request.until { + query.push(("until".to_owned(), until)); + } + if let Some(max_points) = request.max_points { + query.push(("max_points".to_owned(), max_points.to_string())); + } + let path = Client::build_path( + "/api/v2/video/call_stats/{type}/{id}/{session}/participant/{user}/{user_session}/details", + &[ + ("type", &self.call_type), + ("id", &self.call_id), + ("session", session), + ("user", user), + ("user_session", user_session), + ], + ); + self.client + .request::<(), _>(Method::GET, &path, &query, None) + .await + } + + /// Query participant stats for one call session. + /// + /// `GET .../call_stats/{type}/{id}/{session}/participants` + pub async fn query_call_session_participant_stats( + &self, + session: &str, + request: QueryCallSessionParticipantStatsRequest, + ) -> Result { + let mut query = Vec::new(); + if let Some(limit) = request.limit { + query.push(("limit".to_owned(), limit.to_string())); + } + if let Some(prev) = request.prev { + query.push(("prev".to_owned(), prev)); + } + if let Some(next) = request.next { + query.push(("next".to_owned(), next)); + } + if let Some(encoded) = sort_query(&request.sort)? { + query.push(("sort".to_owned(), encoded)); + } + if let Some(encoded) = filter_conditions_query(&request.filter_conditions)? { + query.push(("filter_conditions".to_owned(), encoded)); + } + let path = Client::build_path( + "/api/v2/video/call_stats/{type}/{id}/{session}/participants", + &[ + ("type", &self.call_type), + ("id", &self.call_id), + ("session", session), + ], + ); + self.client + .request::<(), _>(Method::GET, &path, &query, None) + .await + } + + /// Retrieve the participant stats timeline for one participant session. + /// + /// `GET .../call_stats/{type}/{id}/{session}/participants/{user}/{user_session}/timeline` + pub async fn get_call_session_participant_stats_timeline( + &self, + session: &str, + user: &str, + user_session: &str, + request: GetCallSessionParticipantStatsTimelineRequest, + ) -> Result { + let mut query = Vec::new(); + if let Some(start_time) = request.start_time { + query.push(("start_time".to_owned(), start_time)); + } + if let Some(end_time) = request.end_time { + query.push(("end_time".to_owned(), end_time)); + } + if !request.severity.is_empty() { + query.push(("severity".to_owned(), request.severity.join(","))); + } + let path = Client::build_path( + "/api/v2/video/call_stats/{type}/{id}/{session}/participants/{user}/{user_session}/timeline", + &[ + ("type", &self.call_type), + ("id", &self.call_id), + ("session", session), + ("user", user), + ("user_session", user_session), + ], + ); + self.client + .request::<(), _>(Method::GET, &path, &query, None) + .await + } + // participant path (SFU WebRTC) /// Set the maximum reconnect duration. Zero keeps reconnecting indefinitely. @@ -879,3 +1046,71 @@ fn timestamp_query(value: &Timestamp) -> String { .map(str::to_owned) .unwrap_or_else(|| value.to_string()) } + +/// JSON-encode a `filter_conditions` map for a query parameter, or `None` when +/// empty. Matches the getstream-go query encoding for map-valued params. +fn filter_conditions_query(filter: &CustomData) -> Result> { + if filter.is_empty() { + return Ok(None); + } + Ok(Some(serde_json::to_string(filter)?)) +} + +/// Encode a `sort` list for a query parameter, or `None` when empty. Matches the +/// getstream-go query encoding: each entry is JSON-encoded and comma-joined. +fn sort_query(sort: &[SortParamRequest]) -> Result> { + if sort.is_empty() { + return Ok(None); + } + let parts = sort + .iter() + .map(serde_json::to_string) + .collect::, _>>()?; + Ok(Some(parts.join(","))) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn empty_stats_query_params_are_omitted() { + assert!( + filter_conditions_query(&CustomData::new()) + .expect("ok") + .is_none() + ); + assert!(sort_query(&[]).expect("ok").is_none()); + } + + #[test] + fn sort_query_matches_go_comma_joined_json_encoding() { + let sort = vec![ + SortParamRequest { + field: Some("quality_score".to_owned()), + direction: Some(-1), + }, + SortParamRequest { + field: Some("user_id".to_owned()), + direction: Some(1), + }, + ]; + assert_eq!( + sort_query(&sort).expect("ok"), + Some( + "{\"direction\":-1,\"field\":\"quality_score\"},{\"direction\":1,\"field\":\"user_id\"}" + .to_owned() + ) + ); + } + + #[test] + fn filter_conditions_query_json_encodes_map() { + let mut filter = CustomData::new(); + filter.insert("call_cid".to_owned(), serde_json::json!("default:c1")); + assert_eq!( + filter_conditions_query(&filter).expect("ok"), + Some("{\"call_cid\":\"default:c1\"}".to_owned()) + ); + } +} diff --git a/src/video/mod.rs b/src/video/mod.rs index ccd85d7..7855403 100644 --- a/src/video/mod.rs +++ b/src/video/mod.rs @@ -2,6 +2,7 @@ mod call; mod sip; +mod stats; pub use call::Call; diff --git a/src/video/stats.rs b/src/video/stats.rs new file mode 100644 index 0000000..26ea2e1 --- /dev/null +++ b/src/video/stats.rs @@ -0,0 +1,107 @@ +//! Application-level call statistics and reporting endpoints on [`VideoClient`]. + +use reqwest::Method; + +use super::VideoClient; +use crate::error::Result; +use crate::models::{ + GetActiveCallsStatusResponse, GetDailyDigestRequest, GetDailyDigestResponse, + QueryAggregateCallStatsRequest, QueryAggregateCallStatsResponse, QueryCallSessionStatsRequest, + QueryCallSessionStatsResponse, QueryUserFeedbackRequest, QueryUserFeedbackResponse, + ReportClientCallEventRequest, ReportClientEventResponse, +}; + +impl VideoClient { + /// Get the status of all active calls with metrics and summary + /// (`GET /api/v2/video/active_calls_status`). + pub async fn get_active_calls_status(&self) -> Result { + self.client + .request::<(), _>(Method::GET, "/api/v2/video/active_calls_status", &[], None) + .await + } + + /// Query aggregate call stats reports (`POST /api/v2/video/stats`). + pub async fn query_aggregate_call_stats( + &self, + request: QueryAggregateCallStatsRequest, + ) -> Result { + self.client + .request(Method::POST, "/api/v2/video/stats", &[], Some(&request)) + .await + } + + /// Query per-session call stats with filter/sort/pagination + /// (`POST /api/v2/video/call_stats`). + pub async fn query_call_session_stats( + &self, + request: QueryCallSessionStatsRequest, + ) -> Result { + self.client + .request( + Method::POST, + "/api/v2/video/call_stats", + &[], + Some(&request), + ) + .await + } + + /// Get the per-broadcast daily digest bundle for one UTC day + /// (`GET /api/v2/video/stats/daily_digest`). + pub async fn get_daily_digest( + &self, + request: GetDailyDigestRequest, + ) -> Result { + let mut query: Vec<(String, String)> = Vec::new(); + if let Some(date) = request.date { + query.push(("date".to_owned(), date)); + } + if let Some(target_app_id) = request.target_app_id { + query.push(("target_app_id".to_owned(), target_app_id)); + } + self.client + .request::<(), _>( + Method::GET, + "/api/v2/video/stats/daily_digest", + &query, + None, + ) + .await + } + + /// Query user feedback with filter/sort/pagination + /// (`POST /api/v2/video/call/feedback`). + pub async fn query_user_feedback( + &self, + request: QueryUserFeedbackRequest, + ) -> Result { + let query = request + .full + .map(|full| vec![("full".to_owned(), full.to_string())]) + .unwrap_or_default(); + self.client + .request( + Method::POST, + "/api/v2/video/call/feedback", + &query, + Some(&request), + ) + .await + } + + /// Report a batch of client-side telemetry events + /// (`POST /api/v2/video/call_client_event`). + pub async fn report_client_call_event( + &self, + request: ReportClientCallEventRequest, + ) -> Result { + self.client + .request( + Method::POST, + "/api/v2/video/call_client_event", + &[], + Some(&request), + ) + .await + } +} From 1818e4281304000a0a5b5917fb7d7aa2a940aca4 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 4 Sep 2026 17:11:46 +0000 Subject: [PATCH 3/8] test: add live SIP and call-stats integration tests Co-authored-by: nash --- tests/video_sip.rs | 230 +++++++++++++++++++++++++++++++++++++++++++ tests/video_stats.rs | 162 ++++++++++++++++++++++++++++++ 2 files changed, 392 insertions(+) create mode 100644 tests/video_sip.rs create mode 100644 tests/video_stats.rs diff --git a/tests/video_sip.rs b/tests/video_sip.rs new file mode 100644 index 0000000..0e07c80 --- /dev/null +++ b/tests/video_sip.rs @@ -0,0 +1,230 @@ +//! Live integration tests for the SIP (telephony) server REST surface. +//! +//! Run with credentials present (repo `.env`): `cargo test`. Without credentials +//! the tests print a SKIP line and pass without touching the API. Each test +//! creates unique resources and cleans them up on success, failure, and timeout. + +mod common; + +use std::time::Duration; + +use anyhow::{Context, Result, ensure}; +use getstream::Stream; +use getstream::models::{ + CreateSipInboundRoutingRuleRequest, CreateSipTrunkRequest, ResolveSipAuthRequest, + SipCallerConfigsRequest, SipDirectRoutingRuleCallConfigsRequest, UpdateSipTrunkRequest, +}; + +const TEST_TIMEOUT: Duration = Duration::from_secs(60); + +/// A unique, plausible E.164 number derived from a fresh UUID. +fn unique_number() -> String { + let digits = uuid::Uuid::new_v4().as_u128().to_string(); + format!("+1{}", &digits[..10]) +} + +/// Full trunk lifecycle: create → list (present) → update → delete → list (gone). +#[tokio::test] +async fn sip_trunk_crud_lifecycle() { + let Some(client) = common::client_or_skip() else { + return; + }; + + let name = common::unique_id("rust-it-sip-trunk"); + let number = unique_number(); + + let created = client + .video() + .create_sip_trunk(CreateSipTrunkRequest { + password: Some("s3cr3t-pass".to_owned()), + ..CreateSipTrunkRequest::new(&name, [number.clone()]) + }) + .await + .expect("create_sip_trunk failed"); + let trunk_id = created + .sip_trunk + .expect("create response missing sip_trunk") + .id; + assert!(!trunk_id.is_empty(), "created trunk has empty id"); + + let outcome = tokio::time::timeout( + TEST_TIMEOUT, + exercise_trunk(&client, &trunk_id, &name, &number), + ) + .await; + + // Best-effort cleanup regardless of how the assertions above resolved. + let _ = client.video().delete_sip_trunk(&trunk_id).await; + + outcome + .expect("sip trunk lifecycle timed out") + .expect("sip trunk lifecycle assertions failed"); +} + +async fn exercise_trunk(client: &Stream, trunk_id: &str, name: &str, number: &str) -> Result<()> { + let video = client.video(); + + let listed = video.list_sip_trunks().await.context("list_sip_trunks")?; + let found = listed + .sip_trunks + .iter() + .find(|t| t.id == trunk_id) + .context("created trunk not present in list")?; + ensure!(found.name == name, "listed trunk name mismatch"); + ensure!( + found.numbers.iter().any(|n| n == number), + "listed trunk missing its number" + ); + + let updated_name = format!("{name}-updated"); + let updated = video + .update_sip_trunk( + trunk_id, + UpdateSipTrunkRequest { + name: updated_name.clone(), + numbers: vec![number.to_owned()], + ..Default::default() + }, + ) + .await + .context("update_sip_trunk")?; + ensure!( + updated + .sip_trunk + .map(|t| t.name == updated_name) + .unwrap_or(false), + "update did not return the new trunk name" + ); + + // A created trunk with a password authenticates by digest: resolving its + // number returns the trunk id, while an unknown number finds no trunk. + let resolved = video + .resolve_sip_auth(ResolveSipAuthRequest { + sip_caller_number: unique_number(), + sip_trunk_number: number.to_owned(), + ..Default::default() + }) + .await + .context("resolve_sip_auth for known trunk")?; + if resolved.auth_result == "password" { + ensure!( + resolved.trunk_id.as_deref() == Some(trunk_id), + "resolve_sip_auth matched a different trunk" + ); + } + + let unknown = video + .resolve_sip_auth(ResolveSipAuthRequest { + sip_caller_number: unique_number(), + sip_trunk_number: unique_number(), + ..Default::default() + }) + .await + .context("resolve_sip_auth for unknown trunk")?; + ensure!( + unknown.auth_result == "no_trunk_found", + "unknown trunk number should not resolve to a trunk, got {:?}", + unknown.auth_result + ); + + video + .delete_sip_trunk(trunk_id) + .await + .context("delete_sip_trunk")?; + + let after = video + .list_sip_trunks() + .await + .context("list_sip_trunks after delete")?; + ensure!( + after.sip_trunks.iter().all(|t| t.id != trunk_id), + "deleted trunk still present in list" + ); + + Ok(()) +} + +/// Full routing-rule lifecycle against a real trunk, cleaning up both resources. +#[tokio::test] +async fn sip_inbound_routing_rule_crud_lifecycle() { + let Some(client) = common::client_or_skip() else { + return; + }; + + let trunk_name = common::unique_id("rust-it-sip-rule-trunk"); + let created = client + .video() + .create_sip_trunk(CreateSipTrunkRequest::new(&trunk_name, [unique_number()])) + .await + .expect("create_sip_trunk failed"); + let trunk_id = created + .sip_trunk + .expect("create response missing sip_trunk") + .id; + + let outcome = + tokio::time::timeout(TEST_TIMEOUT, exercise_routing_rule(&client, &trunk_id)).await; + + let _ = client.video().delete_sip_trunk(&trunk_id).await; + + outcome + .expect("sip routing rule lifecycle timed out") + .expect("sip routing rule lifecycle assertions failed"); +} + +async fn exercise_routing_rule(client: &Stream, trunk_id: &str) -> Result<()> { + let video = client.video(); + let rule_name = common::unique_id("rust-it-sip-rule"); + + let created = video + .create_sip_inbound_routing_rule(CreateSipInboundRoutingRuleRequest { + name: rule_name.clone(), + trunk_ids: vec![trunk_id.to_owned()], + caller_configs: SipCallerConfigsRequest { + id: "{{caller_number}}".to_owned(), + custom_data: None, + }, + direct_routing_configs: Some(SipDirectRoutingRuleCallConfigsRequest { + call_id: "{{caller_number}}".to_owned(), + call_type: "default".to_owned(), + }), + ..Default::default() + }) + .await + .context("create_sip_inbound_routing_rule")?; + let rule_id = created.id; + ensure!(!rule_id.is_empty(), "created rule has empty id"); + + let listed = video + .list_sip_inbound_routing_rules() + .await + .context("list_sip_inbound_routing_rules")?; + let found = listed + .sip_inbound_routing_rules + .iter() + .find(|r| r.id == rule_id) + .context("created rule not present in list")?; + ensure!( + found.trunk_ids.iter().any(|id| id == trunk_id), + "listed rule missing its trunk id" + ); + + video + .delete_sip_inbound_routing_rule(&rule_id) + .await + .context("delete_sip_inbound_routing_rule")?; + + let after = video + .list_sip_inbound_routing_rules() + .await + .context("list after delete")?; + ensure!( + after + .sip_inbound_routing_rules + .iter() + .all(|r| r.id != rule_id), + "deleted rule still present in list" + ); + + Ok(()) +} diff --git a/tests/video_stats.rs b/tests/video_stats.rs new file mode 100644 index 0000000..f9d8876 --- /dev/null +++ b/tests/video_stats.rs @@ -0,0 +1,162 @@ +//! Live integration tests for the advanced call-stats REST surface. +//! +//! Run with credentials present (repo `.env`): `cargo test`. Without credentials +//! the tests print a SKIP line and pass without touching the API. The call +//! created for the session-scoped queries is deleted on every exit path. +//! +//! Stats are computed asynchronously by the server, so the per-session queries +//! tolerate a `404` (stats not yet available) rather than asserting on analytics +//! values; when a payload is returned, its identity fields must echo the call. + +mod common; + +use getstream::models::{ + CallRequest, DeleteCallRequest, GetOrCreateCallRequest, QueryCallParticipantSessionsRequest, + QueryCallSessionParticipantStatsRequest, UserRequest, +}; +use getstream::rtc::JoinCallData; + +/// Active-calls status is an application-level read with structural invariants +/// that do not depend on analytics timing. +#[tokio::test] +async fn active_calls_status_has_consistent_summary() { + let Some(client) = common::client_or_skip() else { + return; + }; + + let status = client + .video() + .get_active_calls_status() + .await + .expect("get_active_calls_status failed"); + + if let Some(summary) = status.summary { + assert!( + summary.active_calls >= 0, + "active_calls must be non-negative" + ); + assert!( + summary.participants >= 0, + "participants must be non-negative" + ); + assert!( + summary.active_publishers >= 0 && summary.active_subscribers >= 0, + "publisher/subscriber counts must be non-negative" + ); + } +} + +/// Create a call, join a real session, then query its per-session participant +/// stats. Identity fields must echo the call; analytics-not-ready is tolerated. +#[tokio::test] +async fn session_scoped_participant_stats_echo_call_identity() { + let Some(client) = common::client_or_skip() else { + return; + }; + + let user_id = common::unique_id("rust-it-stats-user"); + let call_id = common::unique_id("rust-it-stats-call"); + client + .upsert_users([UserRequest::new(&user_id)]) + .await + .expect("upsert_users failed"); + + let call = client.video().call("default", &call_id); + call.get_or_create(GetOrCreateCallRequest { + data: Some(CallRequest { + created_by_id: Some(user_id.clone()), + ..Default::default() + }), + ..Default::default() + }) + .await + .expect("get_or_create failed"); + + let outcome: Result<(), String> = async { + call.join(JoinCallData::new(&user_id)) + .await + .map_err(|error| format!("join failed: {error}"))?; + let session_id = call + .session_id() + .await + .ok_or_else(|| "joined call did not expose a session id".to_owned())?; + call.leave() + .await + .map_err(|error| format!("leave failed: {error}"))?; + call.end() + .await + .map_err(|error| format!("end failed: {error}"))?; + + if let Some(stats) = allow_stats_pending_skip( + "query_call_session_participant_stats", + call.query_call_session_participant_stats( + &session_id, + QueryCallSessionParticipantStatsRequest::default(), + ) + .await, + )? { + assert_eq!(stats.call_id, call_id, "participant stats call_id mismatch"); + assert_eq!( + stats.call_type, "default", + "participant stats type mismatch" + ); + assert_eq!( + stats.call_session_id, session_id, + "participant stats session mismatch" + ); + } + + if let Some(sessions) = allow_stats_pending_skip( + "query_call_participant_sessions", + call.query_call_participant_sessions( + &session_id, + QueryCallParticipantSessionsRequest::default(), + ) + .await, + )? { + assert_eq!( + sessions.call_id, call_id, + "participant sessions call_id mismatch" + ); + assert_eq!( + sessions.call_type, "default", + "participant sessions type mismatch" + ); + assert_eq!( + sessions.call_session_id, session_id, + "participant sessions session mismatch" + ); + } + + Ok(()) + } + .await; + + let leave_cleanup = call.leave().await; + let delete_cleanup = call.delete(DeleteCallRequest { hard: Some(true) }).await; + if let Err(error) = outcome { + panic!("{error}; leave cleanup: {leave_cleanup:?}; delete cleanup: {delete_cleanup:?}"); + } + let _ = leave_cleanup; + delete_cleanup.expect("delete cleanup failed"); +} + +/// Treat a `404` as "stats not yet computed" (async analytics pipeline), passing +/// the test without asserting on values; surface any other error. +fn allow_stats_pending_skip( + endpoint: &str, + result: Result, +) -> Result, String> { + match result { + Ok(value) => Ok(Some(value)), + Err(error) + if error + .as_api_error() + .is_some_and(|api_error| api_error.status == 404) => + { + eprintln!("SKIP stats pending: {endpoint}: {error}"); + Ok(None) + } + Err(error) => Err(format!("{endpoint} failed: {error}")), + } +} From f817476cabacb03bddb502106c95aeffbe5c5878 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 4 Sep 2026 17:11:46 +0000 Subject: [PATCH 4/8] docs: document SIP and advanced call-stats surface Co-authored-by: nash --- CHANGELOG.md | 23 +++++++++++++++++++++++ README.md | 5 +++++ 2 files changed, 28 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 154791a..259f4d0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,26 @@ +# Unreleased + +## New Features + +### Video REST: SIP telephony + +SIP inbound trunk CRUD (`create_sip_trunk`, `list_sip_trunks`, +`update_sip_trunk`, `delete_sip_trunk`), SIP inbound routing rule CRUD +(`create_sip_inbound_routing_rule`, `list_sip_inbound_routing_rules`, +`update_sip_inbound_routing_rule`, `delete_sip_inbound_routing_rule`), and SIP +resolution (`resolve_sip_auth`, `resolve_sip_inbound`) on `VideoClient`, with +typed request/response models. + +### Video REST: advanced call statistics and reporting + +Application-level stats on `VideoClient` (`get_active_calls_status`, +`query_aggregate_call_stats`, `query_call_session_stats`, `get_daily_digest`, +`query_user_feedback`, `report_client_call_event`) and call-session-scoped stats +on `Call` (`get_call_participant_session_metrics`, +`query_call_participant_sessions`, `get_call_session_participant_stats_details`, +`query_call_session_participant_stats`, +`get_call_session_participant_stats_timeline`). + # v0.1.0-preview.2 docs.rs builds on current nightly. `doc_auto_cfg` was removed in 1.92 and diff --git a/README.md b/README.md index d495210..77ccced 100644 --- a/README.md +++ b/README.md @@ -26,6 +26,11 @@ remote audio and video, transform it, and publish media back into the call. - Create, query, update, end, and delete video calls. - Manage call members, permissions, recording, transcription, captions, livestreaming, custom events, and reactions. +- Configure SIP telephony: inbound trunks, inbound routing rules, and SIP + auth/inbound resolution. +- Query advanced call statistics and reporting: active-calls status, aggregate + and per-session stats, participant stats and metrics, daily digest, user + feedback, and client call-event reporting. - Join a call as a server-side SFU participant with retry, reconnect, and migration handling. - Subscribe globally or by participant session to remote audio, video, and From e771154d0b52f9ee9595caa70abb79834f80a973 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 4 Sep 2026 17:22:51 +0000 Subject: [PATCH 5/8] chore: satisfy clippy 1.98 lints (drain_collect, chunks_exact_to_as_chunks) Co-authored-by: nash --- src/rtc/join/reconnect_runtime.rs | 2 +- src/rtc/pcm/convert.rs | 4 +++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/src/rtc/join/reconnect_runtime.rs b/src/rtc/join/reconnect_runtime.rs index 058505d..9f52a76 100644 --- a/src/rtc/join/reconnect_runtime.rs +++ b/src/rtc/join/reconnect_runtime.rs @@ -845,7 +845,7 @@ impl RtcCore { { return Err(join_cancelled()); } - connection.signal_tasks.drain(..).collect() + std::mem::take(&mut connection.signal_tasks) }; abort_tasks(old_tasks).await; { diff --git a/src/rtc/pcm/convert.rs b/src/rtc/pcm/convert.rs index f68a386..4799b66 100644 --- a/src/rtc/pcm/convert.rs +++ b/src/rtc/pcm/convert.rs @@ -65,7 +65,9 @@ impl PcmFrame { /// one byte. pub fn from_bytes(bytes: &[u8], sample_rate: u32, channels: u16) -> Self { let samples = bytes - .chunks_exact(2) + .as_chunks::<2>() + .0 + .iter() .map(|b| i16::from_le_bytes([b[0], b[1]])) .collect(); Self::new(samples, sample_rate, channels) From c984e0ac3bc9b8df279419e7cdd585b69e528e07 Mon Sep 17 00:00:00 2001 From: "Neevash Ramdial (Nash)" Date: Fri, 4 Sep 2026 12:34:40 -0600 Subject: [PATCH 6/8] fix: key call-stats queries by call session id, trim SIP to CRUD The call-session-scoped stats endpoints are keyed by the coordinator call session id (`get().call.session.id`). The integration test passed `Call::session_id()`, which is the caller's own SFU session and is reported as `user_session_id` within those payloads, so both queries returned 404. The test's unconditional 404 skip then let it pass without asserting anything, leaving the endpoints effectively uncovered. Use the call session id, replace the blanket skip with a bounded retry that still fails on a persistent 404, and assert that the returned participant session carries this join's SFU session id so the two cannot be conflated again. Also drop `resolve_sip_auth` and `resolve_sip_inbound` along with their models; these endpoints are not supported through the server client. The eight trunk and routing-rule CRUD methods are unchanged. Co-Authored-By: Claude Opus 5 (1M context) --- CHANGELOG.md | 7 +-- README.md | 3 +- src/models/sip.rs | 100 ------------------------------ src/video/sip.rs | 34 +--------- tests/video_sip.rs | 35 +---------- tests/video_stats.rs | 144 +++++++++++++++++++++++++++---------------- 6 files changed, 99 insertions(+), 224 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 259f4d0..af54aca 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,11 +5,10 @@ ### Video REST: SIP telephony SIP inbound trunk CRUD (`create_sip_trunk`, `list_sip_trunks`, -`update_sip_trunk`, `delete_sip_trunk`), SIP inbound routing rule CRUD +`update_sip_trunk`, `delete_sip_trunk`) and SIP inbound routing rule CRUD (`create_sip_inbound_routing_rule`, `list_sip_inbound_routing_rules`, -`update_sip_inbound_routing_rule`, `delete_sip_inbound_routing_rule`), and SIP -resolution (`resolve_sip_auth`, `resolve_sip_inbound`) on `VideoClient`, with -typed request/response models. +`update_sip_inbound_routing_rule`, `delete_sip_inbound_routing_rule`) on +`VideoClient`, with typed request/response models. ### Video REST: advanced call statistics and reporting diff --git a/README.md b/README.md index 77ccced..eaf9610 100644 --- a/README.md +++ b/README.md @@ -26,8 +26,7 @@ remote audio and video, transform it, and publish media back into the call. - Create, query, update, end, and delete video calls. - Manage call members, permissions, recording, transcription, captions, livestreaming, custom events, and reactions. -- Configure SIP telephony: inbound trunks, inbound routing rules, and SIP - auth/inbound resolution. +- Configure SIP telephony: inbound trunks and inbound routing rules. - Query advanced call statistics and reporting: active-calls status, aggregate and per-session stats, participant stats and metrics, daily digest, user feedback, and client call-event reporting. diff --git a/src/models/sip.rs b/src/models/sip.rs index 22b8906..7106eaf 100644 --- a/src/models/sip.rs +++ b/src/models/sip.rs @@ -4,8 +4,6 @@ //! can replace these transparently. Response types derive `Default` + //! `#[serde(default)]` so partial payloads deserialize cleanly. -use std::collections::HashMap; - use serde::{Deserialize, Serialize}; use super::shared::{CustomData, Timestamp}; @@ -308,104 +306,6 @@ pub struct ListSipInboundRoutingRuleResponse { pub sip_inbound_routing_rules: Vec, } -// SIP auth / resolve - -/// `resolve_sip_auth` request (`ResolveSipAuthRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct ResolveSipAuthRequest { - pub sip_caller_number: String, - pub sip_trunk_number: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub from_host: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub source_ip: Option, -} - -/// `resolve_sip_auth` response (`ResolveSipAuthResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct ResolveSipAuthResponse { - /// Authentication result: `password`, `accept`, or `no_trunk_found`. - pub auth_result: String, - pub duration: String, - pub password: Option, - pub trunk_id: Option, - pub username: Option, -} - -/// SIP digest challenge authentication data (`SIPChallengeRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct SipChallengeRequest { - #[serde(skip_serializing_if = "Option::is_none")] - pub algorithm: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub charset: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub cnonce: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub method: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub nc: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub nonce: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub opaque: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub realm: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub response: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub stale: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub uri: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub userhash: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub username: Option, - #[serde(skip_serializing_if = "Vec::is_empty")] - pub domain: Vec, - #[serde(skip_serializing_if = "Vec::is_empty")] - pub qop: Vec, -} - -/// `resolve_sip_inbound` request (`ResolveSipInboundRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct ResolveSipInboundRequest { - pub sip_caller_number: String, - pub sip_trunk_number: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub routing_number: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub trunk_id: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub challenge: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub sip_headers: Option>, -} - -/// Credentials for SIP inbound call authentication (`SipInboundCredentials`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct SipInboundCredentials { - pub api_key: String, - pub call_id: String, - pub call_type: String, - pub token: String, - pub user_id: String, - pub call_custom_data: CustomData, - pub user_custom_data: CustomData, -} - -/// `resolve_sip_inbound` response (`ResolveSipInboundResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct ResolveSipInboundResponse { - pub duration: String, - pub credentials: SipInboundCredentials, - pub sip_routing_rule: Option, - pub sip_trunk: Option, -} - #[cfg(test)] mod tests { use serde_json::json; diff --git a/src/video/sip.rs b/src/video/sip.rs index 31df9b9..0c5c5aa 100644 --- a/src/video/sip.rs +++ b/src/video/sip.rs @@ -8,10 +8,8 @@ use crate::error::Result; use crate::models::{ CreateSipInboundRoutingRuleRequest, CreateSipInboundRoutingRuleResponse, CreateSipTrunkRequest, CreateSipTrunkResponse, DeleteSipInboundRoutingRuleResponse, DeleteSipTrunkResponse, - ListSipInboundRoutingRuleResponse, ListSipTrunksResponse, ResolveSipAuthRequest, - ResolveSipAuthResponse, ResolveSipInboundRequest, ResolveSipInboundResponse, - UpdateSipInboundRoutingRuleRequest, UpdateSipInboundRoutingRuleResponse, UpdateSipTrunkRequest, - UpdateSipTrunkResponse, + ListSipInboundRoutingRuleResponse, ListSipTrunksResponse, UpdateSipInboundRoutingRuleRequest, + UpdateSipInboundRoutingRuleResponse, UpdateSipTrunkRequest, UpdateSipTrunkResponse, }; const TRUNKS: &str = "/api/v2/video/sip/inbound_trunks"; @@ -106,32 +104,4 @@ impl VideoClient { .request::<(), _>(Method::DELETE, &path, &[], None) .await } - - // SIP auth / resolve - - /// Resolve SIP authentication requirements for an inbound call - /// (`POST /api/v2/video/sip/auth`). - pub async fn resolve_sip_auth( - &self, - request: ResolveSipAuthRequest, - ) -> Result { - self.client - .request(Method::POST, "/api/v2/video/sip/auth", &[], Some(&request)) - .await - } - - /// Resolve SIP inbound routing (`POST /api/v2/video/sip/resolve`). - pub async fn resolve_sip_inbound( - &self, - request: ResolveSipInboundRequest, - ) -> Result { - self.client - .request( - Method::POST, - "/api/v2/video/sip/resolve", - &[], - Some(&request), - ) - .await - } } diff --git a/tests/video_sip.rs b/tests/video_sip.rs index 0e07c80..b448569 100644 --- a/tests/video_sip.rs +++ b/tests/video_sip.rs @@ -11,8 +11,8 @@ use std::time::Duration; use anyhow::{Context, Result, ensure}; use getstream::Stream; use getstream::models::{ - CreateSipInboundRoutingRuleRequest, CreateSipTrunkRequest, ResolveSipAuthRequest, - SipCallerConfigsRequest, SipDirectRoutingRuleCallConfigsRequest, UpdateSipTrunkRequest, + CreateSipInboundRoutingRuleRequest, CreateSipTrunkRequest, SipCallerConfigsRequest, + SipDirectRoutingRuleCallConfigsRequest, UpdateSipTrunkRequest, }; const TEST_TIMEOUT: Duration = Duration::from_secs(60); @@ -96,37 +96,6 @@ async fn exercise_trunk(client: &Stream, trunk_id: &str, name: &str, number: &st "update did not return the new trunk name" ); - // A created trunk with a password authenticates by digest: resolving its - // number returns the trunk id, while an unknown number finds no trunk. - let resolved = video - .resolve_sip_auth(ResolveSipAuthRequest { - sip_caller_number: unique_number(), - sip_trunk_number: number.to_owned(), - ..Default::default() - }) - .await - .context("resolve_sip_auth for known trunk")?; - if resolved.auth_result == "password" { - ensure!( - resolved.trunk_id.as_deref() == Some(trunk_id), - "resolve_sip_auth matched a different trunk" - ); - } - - let unknown = video - .resolve_sip_auth(ResolveSipAuthRequest { - sip_caller_number: unique_number(), - sip_trunk_number: unique_number(), - ..Default::default() - }) - .await - .context("resolve_sip_auth for unknown trunk")?; - ensure!( - unknown.auth_result == "no_trunk_found", - "unknown trunk number should not resolve to a trunk, got {:?}", - unknown.auth_result - ); - video .delete_sip_trunk(trunk_id) .await diff --git a/tests/video_stats.rs b/tests/video_stats.rs index f9d8876..bfab6ed 100644 --- a/tests/video_stats.rs +++ b/tests/video_stats.rs @@ -4,12 +4,18 @@ //! the tests print a SKIP line and pass without touching the API. The call //! created for the session-scoped queries is deleted on every exit path. //! -//! Stats are computed asynchronously by the server, so the per-session queries -//! tolerate a `404` (stats not yet available) rather than asserting on analytics -//! values; when a payload is returned, its identity fields must echo the call. +//! The per-session queries are keyed by the *coordinator* call session id +//! (`get().call.session.id`), not by `Call::session_id()` -- the latter is this +//! participant's SFU session, which these endpoints report as `user_session_id` +//! nested inside the payload. Passing the wrong one returns `404 call session +//! not found`. Stats can lag the call by a few seconds, so the queries retry on +//! `404` for a bounded window and then fail rather than skipping. mod common; +use std::future::Future; +use std::time::Duration; + use getstream::models::{ CallRequest, DeleteCallRequest, GetOrCreateCallRequest, QueryCallParticipantSessionsRequest, QueryCallSessionParticipantStatsRequest, UserRequest, @@ -47,7 +53,8 @@ async fn active_calls_status_has_consistent_summary() { } /// Create a call, join a real session, then query its per-session participant -/// stats. Identity fields must echo the call; analytics-not-ready is tolerated. +/// stats. Identity fields must echo the call, and the participant session must +/// carry this join's SFU session id. #[tokio::test] async fn session_scoped_participant_stats_echo_call_identity() { let Some(client) = common::client_or_skip() else { @@ -76,10 +83,22 @@ async fn session_scoped_participant_stats_echo_call_identity() { call.join(JoinCallData::new(&user_id)) .await .map_err(|error| format!("join failed: {error}"))?; - let session_id = call + + // `session_id()` is this participant's SFU session; the stats endpoints + // are keyed by the coordinator's call session, which is a different id. + let user_session_id = call .session_id() .await - .ok_or_else(|| "joined call did not expose a session id".to_owned())?; + .ok_or_else(|| "joined call did not expose an SFU session id".to_owned())?; + let session_id = call + .get(Default::default()) + .await + .map_err(|error| format!("get failed: {error}"))? + .call + .session + .map(|session| session.id) + .ok_or_else(|| "joined call did not expose a call session".to_owned())?; + call.leave() .await .map_err(|error| format!("leave failed: {error}"))?; @@ -87,46 +106,55 @@ async fn session_scoped_participant_stats_echo_call_identity() { .await .map_err(|error| format!("end failed: {error}"))?; - if let Some(stats) = allow_stats_pending_skip( - "query_call_session_participant_stats", + let stats = await_stats("query_call_session_participant_stats", || { call.query_call_session_participant_stats( &session_id, QueryCallSessionParticipantStatsRequest::default(), ) - .await, - )? { - assert_eq!(stats.call_id, call_id, "participant stats call_id mismatch"); - assert_eq!( - stats.call_type, "default", - "participant stats type mismatch" - ); - assert_eq!( - stats.call_session_id, session_id, - "participant stats session mismatch" - ); - } + }) + .await?; + assert_eq!(stats.call_id, call_id, "participant stats call_id mismatch"); + assert_eq!( + stats.call_type, "default", + "participant stats type mismatch" + ); + assert_eq!( + stats.call_session_id, session_id, + "participant stats session mismatch" + ); - if let Some(sessions) = allow_stats_pending_skip( - "query_call_participant_sessions", + let sessions = await_stats("query_call_participant_sessions", || { call.query_call_participant_sessions( &session_id, QueryCallParticipantSessionsRequest::default(), ) - .await, - )? { - assert_eq!( - sessions.call_id, call_id, - "participant sessions call_id mismatch" - ); - assert_eq!( - sessions.call_type, "default", - "participant sessions type mismatch" - ); - assert_eq!( - sessions.call_session_id, session_id, - "participant sessions session mismatch" - ); - } + }) + .await?; + assert_eq!( + sessions.call_id, call_id, + "participant sessions call_id mismatch" + ); + assert_eq!( + sessions.call_type, "default", + "participant sessions type mismatch" + ); + assert_eq!( + sessions.call_session_id, session_id, + "participant sessions session mismatch" + ); + + // The join above is the only participant session, and it must be + // reported under the SFU session id -- guarding the two ids from being + // conflated again. + let reported: Vec<&str> = sessions + .participants_sessions + .iter() + .filter_map(|entry| entry.get("user_session_id")?.as_str()) + .collect(); + assert!( + reported.contains(&user_session_id.as_str()), + "participant sessions {reported:?} missing this join's SFU session {user_session_id}" + ); Ok(()) } @@ -141,22 +169,32 @@ async fn session_scoped_participant_stats_echo_call_identity() { delete_cleanup.expect("delete cleanup failed"); } -/// Treat a `404` as "stats not yet computed" (async analytics pipeline), passing -/// the test without asserting on values; surface any other error. -fn allow_stats_pending_skip( - endpoint: &str, - result: Result, -) -> Result, String> { - match result { - Ok(value) => Ok(Some(value)), - Err(error) - if error - .as_api_error() - .is_some_and(|api_error| api_error.status == 404) => - { - eprintln!("SKIP stats pending: {endpoint}: {error}"); - Ok(None) +/// Analytics can trail the call by a few seconds, so retry a `404` for a bounded +/// window. Unlike an unconditional skip, a persistent `404` still fails the test. +async fn await_stats(endpoint: &str, mut query: F) -> Result +where + F: FnMut() -> Fut, + Fut: Future>, +{ + const DEADLINE: Duration = Duration::from_secs(30); + const INTERVAL: Duration = Duration::from_secs(3); + + let mut waited = Duration::ZERO; + loop { + match query().await { + Ok(value) => return Ok(value), + Err(error) + if waited < DEADLINE + && error + .as_api_error() + .is_some_and(|api_error| api_error.status == 404) => + { + tokio::time::sleep(INTERVAL).await; + waited += INTERVAL; + } + Err(error) => { + return Err(format!("{endpoint} failed after {waited:?}: {error}")); + } } - Err(error) => Err(format!("{endpoint} failed: {error}")), } } From 79956cb1ffae72a3eb6b537622b5f7bc25d69081 Mon Sep 17 00:00:00 2001 From: "Neevash Ramdial (Nash)" Date: Fri, 4 Sep 2026 13:03:50 -0600 Subject: [PATCH 7/8] feat: drop SIP telephony from this change Removes the SIP trunk and inbound routing-rule surface, its models, and its integration test, narrowing this change to the advanced call-stats REST coverage. SIP will ship separately. Co-Authored-By: Claude Opus 5 (1M context) --- CHANGELOG.md | 8 - README.md | 1 - src/models/mod.rs | 2 - src/models/sip.rs | 371 --------------------------------------------- src/video/mod.rs | 1 - src/video/sip.rs | 107 ------------- tests/video_sip.rs | 199 ------------------------ 7 files changed, 689 deletions(-) delete mode 100644 src/models/sip.rs delete mode 100644 src/video/sip.rs delete mode 100644 tests/video_sip.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index af54aca..cd25a8e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,14 +2,6 @@ ## New Features -### Video REST: SIP telephony - -SIP inbound trunk CRUD (`create_sip_trunk`, `list_sip_trunks`, -`update_sip_trunk`, `delete_sip_trunk`) and SIP inbound routing rule CRUD -(`create_sip_inbound_routing_rule`, `list_sip_inbound_routing_rules`, -`update_sip_inbound_routing_rule`, `delete_sip_inbound_routing_rule`) on -`VideoClient`, with typed request/response models. - ### Video REST: advanced call statistics and reporting Application-level stats on `VideoClient` (`get_active_calls_status`, diff --git a/README.md b/README.md index eaf9610..12ff8a3 100644 --- a/README.md +++ b/README.md @@ -26,7 +26,6 @@ remote audio and video, transform it, and publish media back into the call. - Create, query, update, end, and delete video calls. - Manage call members, permissions, recording, transcription, captions, livestreaming, custom events, and reactions. -- Configure SIP telephony: inbound trunks and inbound routing rules. - Query advanced call statistics and reporting: active-calls status, aggregate and per-session stats, participant stats and metrics, daily digest, user feedback, and client call-event reporting. diff --git a/src/models/mod.rs b/src/models/mod.rs index 9ee067b..6f42abd 100644 --- a/src/models/mod.rs +++ b/src/models/mod.rs @@ -7,12 +7,10 @@ mod call; mod shared; -mod sip; mod stats; mod user; pub use call::*; pub use shared::*; -pub use sip::*; pub use stats::*; pub use user::*; diff --git a/src/models/sip.rs b/src/models/sip.rs deleted file mode 100644 index 7106eaf..0000000 --- a/src/models/sip.rs +++ /dev/null @@ -1,371 +0,0 @@ -//! SIP (telephony) request/response models. -//! -//! Field names track the getstream-go JSON tags so a later OpenAPI codegen pass -//! can replace these transparently. Response types derive `Default` + -//! `#[serde(default)]` so partial payloads deserialize cleanly. - -use serde::{Deserialize, Serialize}; - -use super::shared::{CustomData, Timestamp}; - -// SIP trunks - -/// `create_sip_trunk` request (`CreateSIPTrunkRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct CreateSipTrunkRequest { - /// Name of the SIP trunk. - pub name: String, - /// Phone numbers associated with this SIP trunk. - pub numbers: Vec, - /// Optional password for SIP trunk authentication. - #[serde(skip_serializing_if = "Option::is_none")] - pub password: Option, - /// Optional list of allowed IPv4/IPv6 addresses or CIDR blocks. - #[serde(skip_serializing_if = "Option::is_none")] - pub allowed_ips: Option>, -} - -impl CreateSipTrunkRequest { - /// Build a create request for a named trunk with its phone numbers. - pub fn new(name: impl Into, numbers: impl IntoIterator) -> Self { - Self { - name: name.into(), - numbers: numbers.into_iter().collect(), - ..Default::default() - } - } -} - -/// `update_sip_trunk` request (`UpdateSIPTrunkRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct UpdateSipTrunkRequest { - /// Name of the SIP trunk. - pub name: String, - /// Phone numbers associated with this SIP trunk. - pub numbers: Vec, - /// Optional password for SIP trunk authentication. - #[serde(skip_serializing_if = "Option::is_none")] - pub password: Option, - /// Optional list of allowed IPv4/IPv6 addresses or CIDR blocks. - #[serde(skip_serializing_if = "Option::is_none")] - pub allowed_ips: Option>, -} - -/// A SIP trunk (`SIPTrunkResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct SipTrunkResponse { - pub id: String, - pub name: String, - /// Password for SIP trunk authentication. - pub password: String, - /// Username for SIP trunk authentication. - pub username: String, - /// The URI for the SIP trunk. - pub uri: String, - pub numbers: Vec, - pub allowed_ips: Vec, - pub created_at: Timestamp, - pub updated_at: Timestamp, -} - -/// `create_sip_trunk` response (`CreateSIPTrunkResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct CreateSipTrunkResponse { - pub duration: String, - pub sip_trunk: Option, -} - -/// `update_sip_trunk` response (`UpdateSIPTrunkResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct UpdateSipTrunkResponse { - pub duration: String, - pub sip_trunk: Option, -} - -/// `delete_sip_trunk` response (`DeleteSIPTrunkResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct DeleteSipTrunkResponse { - pub duration: String, -} - -/// `list_sip_trunks` response (`ListSIPTrunksResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct ListSipTrunksResponse { - pub duration: String, - pub sip_trunks: Vec, -} - -// SIP inbound routing rules - -/// SIP caller settings (`SIPCallerConfigsRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct SipCallerConfigsRequest { - /// Unique identifier for the caller (handlebars template). - pub id: String, - /// Custom data associated with the caller (values are handlebars templates). - #[serde(skip_serializing_if = "Option::is_none")] - pub custom_data: Option, -} - -/// SIP call settings (`SIPCallConfigsRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct SipCallConfigsRequest { - /// Custom data associated with the call. - #[serde(skip_serializing_if = "Option::is_none")] - pub custom_data: Option, -} - -/// Direct routing rule call settings (`SIPDirectRoutingRuleCallConfigsRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct SipDirectRoutingRuleCallConfigsRequest { - /// ID of the call (handlebars template). - pub call_id: String, - /// Type of the call. - pub call_type: String, -} - -/// PIN protection settings (`SIPPinProtectionConfigsRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct SipPinProtectionConfigsRequest { - #[serde(skip_serializing_if = "Option::is_none")] - pub default_pin: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub enabled: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub max_attempts: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub required_pin_digits: Option, -} - -/// PIN routing rule call settings (`SIPInboundRoutingRulePinConfigsRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct SipInboundRoutingRulePinConfigsRequest { - #[serde(skip_serializing_if = "Option::is_none")] - pub custom_webhook_url: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub pin_failed_attempt_prompt: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub pin_hangup_prompt: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub pin_prompt: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub pin_success_prompt: Option, -} - -/// `create_sip_inbound_routing_rule` request (`CreateSIPInboundRoutingRuleRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct CreateSipInboundRoutingRuleRequest { - pub name: String, - pub trunk_ids: Vec, - pub caller_configs: SipCallerConfigsRequest, - #[serde(skip_serializing_if = "Option::is_none")] - pub called_numbers: Option>, - #[serde(skip_serializing_if = "Option::is_none")] - pub caller_numbers: Option>, - #[serde(skip_serializing_if = "Option::is_none")] - pub call_configs: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub direct_routing_configs: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub pin_protection_configs: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub pin_routing_configs: Option, -} - -/// `update_sip_inbound_routing_rule` request (`UpdateSIPInboundRoutingRuleRequest`). -#[derive(Debug, Clone, Default, Serialize)] -pub struct UpdateSipInboundRoutingRuleRequest { - pub name: String, - pub trunk_ids: Vec, - pub caller_configs: SipCallerConfigsRequest, - #[serde(skip_serializing_if = "Option::is_none")] - pub called_numbers: Option>, - #[serde(skip_serializing_if = "Option::is_none")] - pub caller_numbers: Option>, - #[serde(skip_serializing_if = "Option::is_none")] - pub call_configs: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub direct_routing_configs: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub pin_protection_configs: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub pin_routing_configs: Option, -} - -/// SIP call settings response (`SIPCallConfigsResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct SipCallConfigsResponse { - pub custom_data: CustomData, -} - -/// SIP caller settings response (`SIPCallerConfigsResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct SipCallerConfigsResponse { - pub id: String, - pub custom_data: CustomData, -} - -/// Direct routing rule call settings response -/// (`SIPDirectRoutingRuleCallConfigsResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct SipDirectRoutingRuleCallConfigsResponse { - pub call_id: String, - pub call_type: String, -} - -/// PIN protection settings response (`SIPPinProtectionConfigsResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct SipPinProtectionConfigsResponse { - pub enabled: bool, - pub default_pin: Option, - pub max_attempts: Option, - pub required_pin_digits: Option, -} - -/// PIN routing rule call settings response -/// (`SIPInboundRoutingRulePinConfigsResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct SipInboundRoutingRulePinConfigsResponse { - pub custom_webhook_url: Option, - pub pin_failed_attempt_prompt: Option, - pub pin_hangup_prompt: Option, - pub pin_prompt: Option, - pub pin_success_prompt: Option, -} - -/// A SIP inbound routing rule (`SIPInboundRoutingRuleResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct SipInboundRoutingRuleResponse { - pub id: String, - pub name: String, - pub called_numbers: Vec, - pub trunk_ids: Vec, - pub caller_numbers: Vec, - pub created_at: Timestamp, - pub updated_at: Timestamp, - pub call_configs: Option, - pub caller_configs: Option, - pub direct_routing_configs: Option, - pub pin_protection_configs: Option, - pub pin_routing_configs: Option, -} - -/// `create_sip_inbound_routing_rule` response (`SIPInboundRoutingRuleResponse` -/// envelope, returned directly for create). -/// -/// The create endpoint returns the rule fields inline alongside `duration`. -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct CreateSipInboundRoutingRuleResponse { - pub duration: String, - pub id: String, - pub name: String, - pub called_numbers: Vec, - pub trunk_ids: Vec, - pub caller_numbers: Vec, - pub created_at: Timestamp, - pub updated_at: Timestamp, - pub call_configs: Option, - pub caller_configs: Option, - pub direct_routing_configs: Option, - pub pin_protection_configs: Option, - pub pin_routing_configs: Option, -} - -/// `update_sip_inbound_routing_rule` response (`UpdateSIPInboundRoutingRuleResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct UpdateSipInboundRoutingRuleResponse { - pub duration: String, - pub sip_inbound_routing_rule: Option, -} - -/// `delete_sip_inbound_routing_rule` response (`DeleteSIPInboundRoutingRuleResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct DeleteSipInboundRoutingRuleResponse { - pub duration: String, -} - -/// `list_sip_inbound_routing_rules` response (`ListSIPInboundRoutingRuleResponse`). -#[derive(Debug, Clone, Default, Deserialize)] -#[serde(default)] -pub struct ListSipInboundRoutingRuleResponse { - pub duration: String, - pub sip_inbound_routing_rules: Vec, -} - -#[cfg(test)] -mod tests { - use serde_json::json; - - use super::*; - - #[test] - fn create_trunk_omits_absent_optionals_and_keeps_required_fields() { - let value = serde_json::to_value(CreateSipTrunkRequest::new( - "primary", - ["+15551230000".to_owned()], - )) - .expect("request should serialize"); - assert_eq!( - value, - json!({ "name": "primary", "numbers": ["+15551230000"] }) - ); - } - - #[test] - fn create_routing_rule_serializes_required_and_nested_configs() { - let request = CreateSipInboundRoutingRuleRequest { - name: "rule-1".to_owned(), - trunk_ids: vec!["trunk-1".to_owned()], - caller_configs: SipCallerConfigsRequest { - id: "{{caller_number}}".to_owned(), - custom_data: None, - }, - direct_routing_configs: Some(SipDirectRoutingRuleCallConfigsRequest { - call_id: "support".to_owned(), - call_type: "default".to_owned(), - }), - ..Default::default() - }; - let value = serde_json::to_value(&request).expect("request should serialize"); - assert_eq!( - value, - json!({ - "name": "rule-1", - "trunk_ids": ["trunk-1"], - "caller_configs": { "id": "{{caller_number}}" }, - "direct_routing_configs": { "call_id": "support", "call_type": "default" } - }) - ); - } - - #[test] - fn trunk_response_deserializes_partial_payload() { - let response: CreateSipTrunkResponse = serde_json::from_value(json!({ - "duration": "1.2ms", - "sip_trunk": { - "id": "trunk-1", - "name": "primary", - "numbers": ["+15551230000"] - } - })) - .expect("response should deserialize"); - let trunk = response.sip_trunk.expect("trunk present"); - assert_eq!(trunk.id, "trunk-1"); - assert_eq!(trunk.numbers, vec!["+15551230000".to_owned()]); - assert!(trunk.allowed_ips.is_empty()); - } -} diff --git a/src/video/mod.rs b/src/video/mod.rs index 7855403..9704bbb 100644 --- a/src/video/mod.rs +++ b/src/video/mod.rs @@ -1,7 +1,6 @@ //! Video coordinator REST: [`VideoClient`] and [`Call`]. mod call; -mod sip; mod stats; pub use call::Call; diff --git a/src/video/sip.rs b/src/video/sip.rs deleted file mode 100644 index 0c5c5aa..0000000 --- a/src/video/sip.rs +++ /dev/null @@ -1,107 +0,0 @@ -//! SIP (telephony) coordinator REST endpoints on [`VideoClient`]. - -use reqwest::Method; - -use super::VideoClient; -use crate::client::Client; -use crate::error::Result; -use crate::models::{ - CreateSipInboundRoutingRuleRequest, CreateSipInboundRoutingRuleResponse, CreateSipTrunkRequest, - CreateSipTrunkResponse, DeleteSipInboundRoutingRuleResponse, DeleteSipTrunkResponse, - ListSipInboundRoutingRuleResponse, ListSipTrunksResponse, UpdateSipInboundRoutingRuleRequest, - UpdateSipInboundRoutingRuleResponse, UpdateSipTrunkRequest, UpdateSipTrunkResponse, -}; - -const TRUNKS: &str = "/api/v2/video/sip/inbound_trunks"; -const TRUNK_BY_ID: &str = "/api/v2/video/sip/inbound_trunks/{id}"; -const RULES: &str = "/api/v2/video/sip/inbound_routing_rules"; -const RULE_BY_ID: &str = "/api/v2/video/sip/inbound_routing_rules/{id}"; - -impl VideoClient { - // SIP inbound trunks - - /// List SIP inbound trunks (`GET /api/v2/video/sip/inbound_trunks`). - pub async fn list_sip_trunks(&self) -> Result { - self.client - .request::<(), _>(Method::GET, TRUNKS, &[], None) - .await - } - - /// Create a SIP inbound trunk (`POST /api/v2/video/sip/inbound_trunks`). - pub async fn create_sip_trunk( - &self, - request: CreateSipTrunkRequest, - ) -> Result { - self.client - .request(Method::POST, TRUNKS, &[], Some(&request)) - .await - } - - /// Update a SIP inbound trunk (`PUT /api/v2/video/sip/inbound_trunks/{id}`). - pub async fn update_sip_trunk( - &self, - id: &str, - request: UpdateSipTrunkRequest, - ) -> Result { - let path = Client::build_path(TRUNK_BY_ID, &[("id", id)]); - self.client - .request(Method::PUT, &path, &[], Some(&request)) - .await - } - - /// Delete a SIP inbound trunk (`DELETE /api/v2/video/sip/inbound_trunks/{id}`). - pub async fn delete_sip_trunk(&self, id: &str) -> Result { - let path = Client::build_path(TRUNK_BY_ID, &[("id", id)]); - self.client - .request::<(), _>(Method::DELETE, &path, &[], None) - .await - } - - // SIP inbound routing rules - - /// List SIP inbound routing rules - /// (`GET /api/v2/video/sip/inbound_routing_rules`). - pub async fn list_sip_inbound_routing_rules( - &self, - ) -> Result { - self.client - .request::<(), _>(Method::GET, RULES, &[], None) - .await - } - - /// Create a SIP inbound routing rule - /// (`POST /api/v2/video/sip/inbound_routing_rules`). - pub async fn create_sip_inbound_routing_rule( - &self, - request: CreateSipInboundRoutingRuleRequest, - ) -> Result { - self.client - .request(Method::POST, RULES, &[], Some(&request)) - .await - } - - /// Update a SIP inbound routing rule - /// (`PUT /api/v2/video/sip/inbound_routing_rules/{id}`). - pub async fn update_sip_inbound_routing_rule( - &self, - id: &str, - request: UpdateSipInboundRoutingRuleRequest, - ) -> Result { - let path = Client::build_path(RULE_BY_ID, &[("id", id)]); - self.client - .request(Method::PUT, &path, &[], Some(&request)) - .await - } - - /// Delete a SIP inbound routing rule - /// (`DELETE /api/v2/video/sip/inbound_routing_rules/{id}`). - pub async fn delete_sip_inbound_routing_rule( - &self, - id: &str, - ) -> Result { - let path = Client::build_path(RULE_BY_ID, &[("id", id)]); - self.client - .request::<(), _>(Method::DELETE, &path, &[], None) - .await - } -} diff --git a/tests/video_sip.rs b/tests/video_sip.rs deleted file mode 100644 index b448569..0000000 --- a/tests/video_sip.rs +++ /dev/null @@ -1,199 +0,0 @@ -//! Live integration tests for the SIP (telephony) server REST surface. -//! -//! Run with credentials present (repo `.env`): `cargo test`. Without credentials -//! the tests print a SKIP line and pass without touching the API. Each test -//! creates unique resources and cleans them up on success, failure, and timeout. - -mod common; - -use std::time::Duration; - -use anyhow::{Context, Result, ensure}; -use getstream::Stream; -use getstream::models::{ - CreateSipInboundRoutingRuleRequest, CreateSipTrunkRequest, SipCallerConfigsRequest, - SipDirectRoutingRuleCallConfigsRequest, UpdateSipTrunkRequest, -}; - -const TEST_TIMEOUT: Duration = Duration::from_secs(60); - -/// A unique, plausible E.164 number derived from a fresh UUID. -fn unique_number() -> String { - let digits = uuid::Uuid::new_v4().as_u128().to_string(); - format!("+1{}", &digits[..10]) -} - -/// Full trunk lifecycle: create → list (present) → update → delete → list (gone). -#[tokio::test] -async fn sip_trunk_crud_lifecycle() { - let Some(client) = common::client_or_skip() else { - return; - }; - - let name = common::unique_id("rust-it-sip-trunk"); - let number = unique_number(); - - let created = client - .video() - .create_sip_trunk(CreateSipTrunkRequest { - password: Some("s3cr3t-pass".to_owned()), - ..CreateSipTrunkRequest::new(&name, [number.clone()]) - }) - .await - .expect("create_sip_trunk failed"); - let trunk_id = created - .sip_trunk - .expect("create response missing sip_trunk") - .id; - assert!(!trunk_id.is_empty(), "created trunk has empty id"); - - let outcome = tokio::time::timeout( - TEST_TIMEOUT, - exercise_trunk(&client, &trunk_id, &name, &number), - ) - .await; - - // Best-effort cleanup regardless of how the assertions above resolved. - let _ = client.video().delete_sip_trunk(&trunk_id).await; - - outcome - .expect("sip trunk lifecycle timed out") - .expect("sip trunk lifecycle assertions failed"); -} - -async fn exercise_trunk(client: &Stream, trunk_id: &str, name: &str, number: &str) -> Result<()> { - let video = client.video(); - - let listed = video.list_sip_trunks().await.context("list_sip_trunks")?; - let found = listed - .sip_trunks - .iter() - .find(|t| t.id == trunk_id) - .context("created trunk not present in list")?; - ensure!(found.name == name, "listed trunk name mismatch"); - ensure!( - found.numbers.iter().any(|n| n == number), - "listed trunk missing its number" - ); - - let updated_name = format!("{name}-updated"); - let updated = video - .update_sip_trunk( - trunk_id, - UpdateSipTrunkRequest { - name: updated_name.clone(), - numbers: vec![number.to_owned()], - ..Default::default() - }, - ) - .await - .context("update_sip_trunk")?; - ensure!( - updated - .sip_trunk - .map(|t| t.name == updated_name) - .unwrap_or(false), - "update did not return the new trunk name" - ); - - video - .delete_sip_trunk(trunk_id) - .await - .context("delete_sip_trunk")?; - - let after = video - .list_sip_trunks() - .await - .context("list_sip_trunks after delete")?; - ensure!( - after.sip_trunks.iter().all(|t| t.id != trunk_id), - "deleted trunk still present in list" - ); - - Ok(()) -} - -/// Full routing-rule lifecycle against a real trunk, cleaning up both resources. -#[tokio::test] -async fn sip_inbound_routing_rule_crud_lifecycle() { - let Some(client) = common::client_or_skip() else { - return; - }; - - let trunk_name = common::unique_id("rust-it-sip-rule-trunk"); - let created = client - .video() - .create_sip_trunk(CreateSipTrunkRequest::new(&trunk_name, [unique_number()])) - .await - .expect("create_sip_trunk failed"); - let trunk_id = created - .sip_trunk - .expect("create response missing sip_trunk") - .id; - - let outcome = - tokio::time::timeout(TEST_TIMEOUT, exercise_routing_rule(&client, &trunk_id)).await; - - let _ = client.video().delete_sip_trunk(&trunk_id).await; - - outcome - .expect("sip routing rule lifecycle timed out") - .expect("sip routing rule lifecycle assertions failed"); -} - -async fn exercise_routing_rule(client: &Stream, trunk_id: &str) -> Result<()> { - let video = client.video(); - let rule_name = common::unique_id("rust-it-sip-rule"); - - let created = video - .create_sip_inbound_routing_rule(CreateSipInboundRoutingRuleRequest { - name: rule_name.clone(), - trunk_ids: vec![trunk_id.to_owned()], - caller_configs: SipCallerConfigsRequest { - id: "{{caller_number}}".to_owned(), - custom_data: None, - }, - direct_routing_configs: Some(SipDirectRoutingRuleCallConfigsRequest { - call_id: "{{caller_number}}".to_owned(), - call_type: "default".to_owned(), - }), - ..Default::default() - }) - .await - .context("create_sip_inbound_routing_rule")?; - let rule_id = created.id; - ensure!(!rule_id.is_empty(), "created rule has empty id"); - - let listed = video - .list_sip_inbound_routing_rules() - .await - .context("list_sip_inbound_routing_rules")?; - let found = listed - .sip_inbound_routing_rules - .iter() - .find(|r| r.id == rule_id) - .context("created rule not present in list")?; - ensure!( - found.trunk_ids.iter().any(|id| id == trunk_id), - "listed rule missing its trunk id" - ); - - video - .delete_sip_inbound_routing_rule(&rule_id) - .await - .context("delete_sip_inbound_routing_rule")?; - - let after = video - .list_sip_inbound_routing_rules() - .await - .context("list after delete")?; - ensure!( - after - .sip_inbound_routing_rules - .iter() - .all(|r| r.id != rule_id), - "deleted rule still present in list" - ); - - Ok(()) -} From ff569fe6c7d511a78342eb67d2dfc6125ef75f12 Mon Sep 17 00:00:00 2001 From: "Neevash Ramdial (Nash)" Date: Fri, 4 Sep 2026 15:05:36 -0600 Subject: [PATCH 8/8] fix: encode the sort query param as a JSON array `sort_query` comma-joined the encoded entries, which is not valid JSON for either a single entry or several; the coordinator rejected it with "is not a valid JSON for field 'sort'". Encode the list directly instead, and pin the array shape in a test that parses the output rather than comparing it against the encoder's own formatting. Also align the stats surface with the conventions around it: time-range request params take `Timestamp` like their siblings in `models::call`, `call_stats` paths go through a `stats_path` helper alongside `path`, and optional query params are pushed by a shared `push_opt`. The live test now populates `filter_conditions` and `limit` so the encoding is exercised end to end, and runs its two independent reads concurrently. Co-Authored-By: Claude Opus 5 (1M context) --- src/models/stats.rs | 19 +++-- src/rtc/pcm/convert.rs | 2 +- src/video/call.rs | 170 +++++++++++++++++++++-------------------- tests/video_stats.rs | 48 +++++++----- 4 files changed, 132 insertions(+), 107 deletions(-) diff --git a/src/models/stats.rs b/src/models/stats.rs index af74678..dbce3f8 100644 --- a/src/models/stats.rs +++ b/src/models/stats.rs @@ -96,8 +96,8 @@ pub struct QueryCallSessionStatsResponse { /// Query params for `get_call_session_participant_stats_details`. #[derive(Debug, Clone, Default)] pub struct GetCallSessionParticipantStatsDetailsRequest { - pub since: Option, - pub until: Option, + pub since: Option, + pub until: Option, pub max_points: Option, } @@ -124,6 +124,13 @@ pub struct QueryCallSessionParticipantStatsRequest { pub limit: Option, pub prev: Option, pub next: Option, + /// Sort order for the returned participants. + /// + /// The coordinator currently answers any non-empty value on this endpoint + /// with `custom sorting is not supported`; it is accepted here so callers + /// are ready when sorting is enabled server-side. The `sort` on the + /// `query_call_session_stats` / `query_user_feedback` request bodies is + /// supported today. pub sort: Vec, pub filter_conditions: CustomData, } @@ -151,8 +158,8 @@ pub struct QueryCallSessionParticipantStatsResponse { /// Query params for `get_call_session_participant_stats_timeline`. #[derive(Debug, Clone, Default)] pub struct GetCallSessionParticipantStatsTimelineRequest { - pub start_time: Option, - pub end_time: Option, + pub start_time: Option, + pub end_time: Option, pub severity: Vec, } @@ -211,6 +218,8 @@ pub struct QueryCallParticipantSessionsRequest { #[derive(Debug, Clone, Default, Deserialize)] #[serde(default)] pub struct QueryCallParticipantSessionsResponse { + /// Session length in seconds. Unlike the `duration` string other endpoints + /// return (`"23.27ms"`), this endpoint returns an integer. pub duration: i64, pub call_id: String, pub call_session_id: String, @@ -241,7 +250,7 @@ pub struct GetDailyDigestResponse { pub date: String, /// Readiness status: `ready`, `pending`, `failed`, `future_date`, `expired`. pub status: String, - pub generated_at: Option, + pub generated_at: Option, pub retry_after: Option, pub revision: Option, pub schema_version: Option, diff --git a/src/rtc/pcm/convert.rs b/src/rtc/pcm/convert.rs index 4799b66..ab01f6d 100644 --- a/src/rtc/pcm/convert.rs +++ b/src/rtc/pcm/convert.rs @@ -68,7 +68,7 @@ impl PcmFrame { .as_chunks::<2>() .0 .iter() - .map(|b| i16::from_le_bytes([b[0], b[1]])) + .map(|&b| i16::from_le_bytes(b)) .collect(); Self::new(samples, sample_rate, channels) } diff --git a/src/video/call.rs b/src/video/call.rs index ac9ff60..4560543 100644 --- a/src/video/call.rs +++ b/src/video/call.rs @@ -10,6 +10,7 @@ use crate::error::{Error, Result}; use crate::models::*; const CALL_BASE: &str = "/api/v2/video/call/{type}/{id}"; +const CALL_STATS_BASE: &str = "/api/v2/video/call_stats/{type}/{id}"; const INTERNAL_RTC_TOKEN_LIFETIME: Duration = Duration::from_secs(10 * 60); /// A handle to a specific call (`:`). Cheap to construct; no request is @@ -52,7 +53,16 @@ impl Call { } fn path(&self, suffix: &str, extra: &[(&str, &str)]) -> String { - let template = format!("{CALL_BASE}{suffix}"); + self.path_from(CALL_BASE, suffix, extra) + } + + /// Same substitution as [`Self::path`] against the `call_stats` base. + fn stats_path(&self, suffix: &str, extra: &[(&str, &str)]) -> String { + self.path_from(CALL_STATS_BASE, suffix, extra) + } + + fn path_from(&self, base: &str, suffix: &str, extra: &[(&str, &str)]) -> String { + let template = format!("{base}{suffix}"); let mut params: Vec<(&str, &str)> = vec![("type", &self.call_type), ("id", &self.call_id)]; params.extend_from_slice(extra); Client::build_path(&template, ¶ms) @@ -650,14 +660,7 @@ impl Call { if let Some(value) = request.exclude_sfus { query.push(("exclude_sfus".to_owned(), value.to_string())); } - let path = Client::build_path( - "/api/v2/video/call_stats/{type}/{id}/{session_id}/map", - &[ - ("type", &self.call_type), - ("id", &self.call_id), - ("session_id", session_id), - ], - ); + let path = self.stats_path("/{session_id}/map", &[("session_id", session_id)]); self.client .request::<(), _>(Method::GET, &path, &query, None) .await @@ -674,12 +677,16 @@ impl Call { request: GetCallParticipantSessionMetricsRequest, ) -> Result { let mut query = Vec::new(); - if let Some(value) = request.since.as_ref() { - query.push(("since".to_owned(), timestamp_query(value))); - } - if let Some(value) = request.until.as_ref() { - query.push(("until".to_owned(), timestamp_query(value))); - } + push_opt( + &mut query, + "since", + request.since.as_ref().map(timestamp_query), + ); + push_opt( + &mut query, + "until", + request.until.as_ref().map(timestamp_query), + ); let path = self.path( "/session/{session}/participant/{user}/{user_session}/details/track", &[ @@ -702,15 +709,9 @@ impl Call { request: QueryCallParticipantSessionsRequest, ) -> Result { let mut query = Vec::new(); - if let Some(limit) = request.limit { - query.push(("limit".to_owned(), limit.to_string())); - } - if let Some(prev) = request.prev { - query.push(("prev".to_owned(), prev)); - } - if let Some(next) = request.next { - query.push(("next".to_owned(), next)); - } + push_opt(&mut query, "limit", request.limit); + push_opt(&mut query, "prev", request.prev); + push_opt(&mut query, "next", request.next); if let Some(encoded) = filter_conditions_query(&request.filter_conditions)? { query.push(("filter_conditions".to_owned(), encoded)); } @@ -734,20 +735,20 @@ impl Call { request: GetCallSessionParticipantStatsDetailsRequest, ) -> Result { let mut query = Vec::new(); - if let Some(since) = request.since { - query.push(("since".to_owned(), since)); - } - if let Some(until) = request.until { - query.push(("until".to_owned(), until)); - } - if let Some(max_points) = request.max_points { - query.push(("max_points".to_owned(), max_points.to_string())); - } - let path = Client::build_path( - "/api/v2/video/call_stats/{type}/{id}/{session}/participant/{user}/{user_session}/details", + push_opt( + &mut query, + "since", + request.since.as_ref().map(timestamp_query), + ); + push_opt( + &mut query, + "until", + request.until.as_ref().map(timestamp_query), + ); + push_opt(&mut query, "max_points", request.max_points); + let path = self.stats_path( + "/{session}/participant/{user}/{user_session}/details", &[ - ("type", &self.call_type), - ("id", &self.call_id), ("session", session), ("user", user), ("user_session", user_session), @@ -767,29 +768,16 @@ impl Call { request: QueryCallSessionParticipantStatsRequest, ) -> Result { let mut query = Vec::new(); - if let Some(limit) = request.limit { - query.push(("limit".to_owned(), limit.to_string())); - } - if let Some(prev) = request.prev { - query.push(("prev".to_owned(), prev)); - } - if let Some(next) = request.next { - query.push(("next".to_owned(), next)); - } + push_opt(&mut query, "limit", request.limit); + push_opt(&mut query, "prev", request.prev); + push_opt(&mut query, "next", request.next); if let Some(encoded) = sort_query(&request.sort)? { query.push(("sort".to_owned(), encoded)); } if let Some(encoded) = filter_conditions_query(&request.filter_conditions)? { query.push(("filter_conditions".to_owned(), encoded)); } - let path = Client::build_path( - "/api/v2/video/call_stats/{type}/{id}/{session}/participants", - &[ - ("type", &self.call_type), - ("id", &self.call_id), - ("session", session), - ], - ); + let path = self.stats_path("/{session}/participants", &[("session", session)]); self.client .request::<(), _>(Method::GET, &path, &query, None) .await @@ -806,20 +794,22 @@ impl Call { request: GetCallSessionParticipantStatsTimelineRequest, ) -> Result { let mut query = Vec::new(); - if let Some(start_time) = request.start_time { - query.push(("start_time".to_owned(), start_time)); - } - if let Some(end_time) = request.end_time { - query.push(("end_time".to_owned(), end_time)); - } + push_opt( + &mut query, + "start_time", + request.start_time.as_ref().map(timestamp_query), + ); + push_opt( + &mut query, + "end_time", + request.end_time.as_ref().map(timestamp_query), + ); if !request.severity.is_empty() { query.push(("severity".to_owned(), request.severity.join(","))); } - let path = Client::build_path( - "/api/v2/video/call_stats/{type}/{id}/{session}/participants/{user}/{user_session}/timeline", + let path = self.stats_path( + "/{session}/participants/{user}/{user_session}/timeline", &[ - ("type", &self.call_type), - ("id", &self.call_id), ("session", session), ("user", user), ("user_session", user_session), @@ -1047,26 +1037,31 @@ fn timestamp_query(value: &Timestamp) -> String { .unwrap_or_else(|| value.to_string()) } -/// JSON-encode a `filter_conditions` map for a query parameter, or `None` when -/// empty. Matches the getstream-go query encoding for map-valued params. -fn filter_conditions_query(filter: &CustomData) -> Result> { - if filter.is_empty() { - return Ok(None); +/// Push `name=value` when the option is set, stringifying the value. +fn push_opt(query: &mut Vec<(String, String)>, name: &str, value: Option) { + if let Some(value) = value { + query.push((name.to_owned(), value.to_string())); } - Ok(Some(serde_json::to_string(filter)?)) } -/// Encode a `sort` list for a query parameter, or `None` when empty. Matches the -/// getstream-go query encoding: each entry is JSON-encoded and comma-joined. +/// JSON-encode a `sort` list for a query parameter, or `None` when empty. +/// +/// The coordinator parses this parameter as a JSON array; comma-joining the +/// encoded entries instead yields `is not a valid JSON for field 'sort'`. fn sort_query(sort: &[SortParamRequest]) -> Result> { if sort.is_empty() { return Ok(None); } - let parts = sort - .iter() - .map(serde_json::to_string) - .collect::, _>>()?; - Ok(Some(parts.join(","))) + Ok(Some(serde_json::to_string(sort)?)) +} + +/// JSON-encode a `filter_conditions` map for a query parameter, or `None` when +/// empty. Matches the getstream-go query encoding for map-valued params. +fn filter_conditions_query(filter: &CustomData) -> Result> { + if filter.is_empty() { + return Ok(None); + } + Ok(Some(serde_json::to_string(filter)?)) } #[cfg(test)] @@ -1080,11 +1075,12 @@ mod tests { .expect("ok") .is_none() ); - assert!(sort_query(&[]).expect("ok").is_none()); } #[test] - fn sort_query_matches_go_comma_joined_json_encoding() { + fn sort_query_encodes_a_json_array() { + assert!(sort_query(&[]).expect("ok").is_none()); + let sort = vec![ SortParamRequest { field: Some("quality_score".to_owned()), @@ -1095,12 +1091,18 @@ mod tests { direction: Some(1), }, ]; + let encoded = sort_query(&sort).expect("ok").expect("some"); + + // Must parse as a JSON array: the coordinator rejects comma-joined + // objects with "is not a valid JSON for field 'sort'". + let parsed: serde_json::Value = + serde_json::from_str(&encoded).expect("sort query must be valid JSON"); assert_eq!( - sort_query(&sort).expect("ok"), - Some( - "{\"direction\":-1,\"field\":\"quality_score\"},{\"direction\":1,\"field\":\"user_id\"}" - .to_owned() - ) + parsed, + serde_json::json!([ + {"direction": -1, "field": "quality_score"}, + {"direction": 1, "field": "user_id"}, + ]) ); } diff --git a/tests/video_stats.rs b/tests/video_stats.rs index bfab6ed..0d8b491 100644 --- a/tests/video_stats.rs +++ b/tests/video_stats.rs @@ -17,8 +17,8 @@ use std::future::Future; use std::time::Duration; use getstream::models::{ - CallRequest, DeleteCallRequest, GetOrCreateCallRequest, QueryCallParticipantSessionsRequest, - QueryCallSessionParticipantStatsRequest, UserRequest, + CallRequest, CustomData, DeleteCallRequest, GetOrCreateCallRequest, + QueryCallParticipantSessionsRequest, QueryCallSessionParticipantStatsRequest, UserRequest, }; use getstream::rtc::JoinCallData; @@ -106,13 +106,28 @@ async fn session_scoped_participant_stats_echo_call_identity() { .await .map_err(|error| format!("end failed: {error}"))?; - let stats = await_stats("query_call_session_participant_stats", || { - call.query_call_session_participant_stats( - &session_id, - QueryCallSessionParticipantStatsRequest::default(), - ) - }) - .await?; + // Independent reads of the same ended session: no need to serialise + // their retry windows. + let (stats, sessions) = tokio::try_join!( + await_stats("query_call_session_participant_stats", || { + call.query_call_session_participant_stats( + &session_id, + QueryCallSessionParticipantStatsRequest { + // Populated so the query encoding is validated against + // the server, not just against its own unit test. + limit: Some(5), + filter_conditions: participant_filter(&user_id), + ..Default::default() + }, + ) + }), + await_stats("query_call_participant_sessions", || { + call.query_call_participant_sessions( + &session_id, + QueryCallParticipantSessionsRequest::default(), + ) + }), + )?; assert_eq!(stats.call_id, call_id, "participant stats call_id mismatch"); assert_eq!( stats.call_type, "default", @@ -123,13 +138,6 @@ async fn session_scoped_participant_stats_echo_call_identity() { "participant stats session mismatch" ); - let sessions = await_stats("query_call_participant_sessions", || { - call.query_call_participant_sessions( - &session_id, - QueryCallParticipantSessionsRequest::default(), - ) - }) - .await?; assert_eq!( sessions.call_id, call_id, "participant sessions call_id mismatch" @@ -165,10 +173,16 @@ async fn session_scoped_participant_stats_echo_call_identity() { if let Err(error) = outcome { panic!("{error}; leave cleanup: {leave_cleanup:?}; delete cleanup: {delete_cleanup:?}"); } - let _ = leave_cleanup; delete_cleanup.expect("delete cleanup failed"); } +/// Restrict a participant-stats query to one user, exercising the +/// `filter_conditions` query encoding against the server rather than only +/// against the encoder's own unit test. +fn participant_filter(user_id: &str) -> CustomData { + CustomData::from([("user_id".to_owned(), serde_json::Value::from(user_id))]) +} + /// Analytics can trail the call by a few seconds, so retry a `404` for a bounded /// window. Unlike an unconditional skip, a persistent `404` still fails the test. async fn await_stats(endpoint: &str, mut query: F) -> Result