Skip to main content

statsig_rust/event_logging_adapter/
statsig_http_event_logging_adapter.rs

1use crate::compression::compression_helper::{compress_data, get_compression_format};
2use crate::event_logging_adapter::EventLoggingAdapter;
3use crate::log_event_payload::LogEventRequest;
4use crate::networking::{NetworkClient, NetworkError, RequestArgs};
5use crate::observability::ops_stats::{OpsStatsForInstance, OPS_STATS};
6use crate::statsig_metadata::StatsigMetadata;
7use crate::{log_d, StatsigErr, StatsigOptions, StatsigRuntime};
8use async_trait::async_trait;
9use serde::Deserialize;
10use serde_json::json;
11use std::collections::HashMap;
12use std::sync::Arc;
13
14const DEFAULT_LOG_EVENT_URL: &str = "https://prodregistryv2.org/v1/log_event";
15
16#[derive(Deserialize)]
17struct LogEventResult {
18    success: Option<bool>,
19}
20
21const TAG: &str = stringify!(StatsigHttpEventLoggingAdapter);
22
23pub struct StatsigHttpEventLoggingAdapter {
24    log_event_url: String,
25    network: NetworkClient,
26    ops_stats: Arc<OpsStatsForInstance>,
27}
28
29impl StatsigHttpEventLoggingAdapter {
30    #[must_use]
31    pub fn new(sdk_key: &str, options: Option<&StatsigOptions>) -> Self {
32        let headers = StatsigMetadata::get_constant_request_headers(
33            sdk_key,
34            options.and_then(|opts| opts.service_name.as_deref()),
35        );
36
37        let log_event_url = options
38            .and_then(|opts| opts.log_event_url.as_ref())
39            .map(|u| u.to_string())
40            .unwrap_or_else(|| DEFAULT_LOG_EVENT_URL.to_string());
41
42        let sdk_instance_id = options
43            .map(|opts| opts.get_sdk_instance_id(sdk_key))
44            .unwrap_or(sdk_key);
45
46        Self {
47            log_event_url,
48            network: NetworkClient::new(sdk_key, Some(headers), options),
49            ops_stats: OPS_STATS.get_for_instance(sdk_instance_id),
50        }
51    }
52
53    pub async fn send_events_over_http(&self, request: &LogEventRequest) -> Result<(), StatsigErr> {
54        log_d!(
55            TAG,
56            "Logging Events ({}): {}",
57            &request.event_count,
58            json!(&request.payload).to_string()
59        );
60
61        let compression_format = get_compression_format();
62
63        // Set headers
64        let headers = HashMap::from([
65            (
66                "statsig-event-count".to_string(),
67                request.event_count.to_string(),
68            ),
69            (
70                "statsig-retry-count".to_string(),
71                request.retries.to_string(),
72            ),
73            ("Content-Encoding".to_owned(), compression_format.clone()),
74            ("Content-Type".to_owned(), "application/json".to_owned()),
75        ]);
76
77        // Compress data before sending it
78        let bytes = serde_json::to_vec(&request.payload)
79            .map_err(|e| StatsigErr::SerializationError(e.to_string()))?;
80        self.ops_stats
81            .log_event_request_uncompressed_body_size_bytes(
82                bytes.len(),
83                get_request_flush_type(request),
84                self.get_observability_tags(),
85            );
86
87        let compressed = match compress_data(&bytes) {
88            Ok(c) => c,
89            Err(e) => return Err(e),
90        };
91
92        // Make request
93        let response = self
94            .network
95            .post(
96                RequestArgs {
97                    url: self.log_event_url.clone(),
98                    headers: Some(headers),
99                    accept_gzip_response: true,
100                    ..RequestArgs::new()
101                },
102                Some(compressed),
103            )
104            .await
105            .map_err(StatsigErr::NetworkError)?;
106
107        let mut res_data = match response.data {
108            Some(data) => data,
109            None => {
110                return Err(StatsigErr::NetworkError(NetworkError::RequestFailed(
111                    self.log_event_url.clone(),
112                    response.status_code,
113                    "Empty response from network".to_string(),
114                )));
115            }
116        };
117
118        let result = res_data
119            .deserialize_into::<LogEventResult>()
120            .map(|result| result.success != Some(false))
121            .map_err(|e| {
122                StatsigErr::JsonParseError(stringify!(LogEventResult).to_string(), e.to_string())
123            })?;
124
125        if result {
126            Ok(())
127        } else {
128            Err(StatsigErr::LogEventError(
129                "Unsuccessful response from network".into(),
130            ))
131        }
132    }
133}
134
135fn get_request_flush_type(request: &LogEventRequest) -> String {
136    request
137        .payload
138        .statsig_metadata
139        .get("flushType")
140        .and_then(|value| value.as_str())
141        .unwrap_or("unknown")
142        .to_string()
143}
144
145#[async_trait]
146impl EventLoggingAdapter for StatsigHttpEventLoggingAdapter {
147    async fn start(&self, _statsig_runtime: &Arc<StatsigRuntime>) -> Result<(), StatsigErr> {
148        Ok(())
149    }
150
151    async fn log_events(&self, request: LogEventRequest) -> Result<bool, StatsigErr> {
152        match self.send_events_over_http(&request).await {
153            Ok(()) => Ok(true),
154            Err(e) => Err(e),
155        }
156    }
157
158    async fn shutdown(&self) -> Result<(), StatsigErr> {
159        Ok(())
160    }
161
162    fn should_schedule_background_flush(&self) -> bool {
163        true
164    }
165
166    fn get_observability_tags(&self) -> Option<HashMap<String, String>> {
167        Some(HashMap::from([(
168            "source_api".to_string(),
169            self.log_event_url.clone(),
170        )]))
171    }
172}
173
174#[cfg(not(feature = "with_zstd"))]
175#[tokio::test]
176async fn test_event_logging() {
177    use crate::log_event_payload::{LogEventPayload, LogEventRequest};
178    use std::env;
179
180    let sdk_key = env::var("test_api_key").expect("test_api_key environment variable not set");
181
182    let adapter = StatsigHttpEventLoggingAdapter::new(&sdk_key, None);
183
184    let payload_str = r#"{"events":[{"eventName":"statsig::config_exposure","metadata":{"config":"running_exp_in_unlayered_with_holdout","ruleID":"5suobe8yyvznqasn9Ph1dI"},"secondaryExposures":[{"gate":"global_holdout","gateValue":"false","ruleID":"3QoA4ncNdVGBaMt3N1KYjz:0.50:1"},{"gate":"exp_holdout","gateValue":"false","ruleID":"1rEqLOpCROaRafv7ubGgax"}],"time":1722386636538,"user":{"appVersion":null,"country":null,"custom":null,"customIDs":null,"email":"daniel@statsig.com","ip":null,"locale":null,"privateAttributes":null,"statsigEnvironment":null,"userAgent":null,"userID":"a-user"},"value":null}],"statsigMetadata":{"sdk_type":"statsig-server-core","sdk_version":"0.0.1"}}"#;
185    let payload = serde_json::from_str::<LogEventPayload>(payload_str).unwrap();
186
187    let request = LogEventRequest {
188        payload,
189        event_count: 1,
190        retries: 0,
191    };
192
193    let result = adapter.log_events(request).await;
194
195    assert!(result.is_ok(), "Error logging events: {:?}", result.err());
196}