1use super::config_spec_background_sync_metrics::log_config_sync_overall_latency;
2use super::response_format::{get_specs_response_format, SpecsResponseFormat};
3use crate::networking::{NetworkClient, NetworkError, RequestArgs, ResponseData};
4use crate::observability::ops_stats::{OpsStatsForInstance, OPS_STATS};
5use crate::observability::sdk_errors_observer::ErrorBoundaryEvent;
6use crate::sdk_diagnostics::diagnostics::ContextType;
7use crate::sdk_diagnostics::marker::{ActionType, KeyType, Marker, StepType};
8use crate::specs_adapter::{SpecsAdapter, SpecsUpdate, SpecsUpdateListener};
9use crate::specs_response::spec_types::SpecsResponseNoUpdates;
10use crate::statsig_err::StatsigErr;
11use crate::statsig_metadata::StatsigMetadata;
12use crate::utils::get_api_from_url;
13use crate::DEFAULT_INIT_TIMEOUT_MS;
14use crate::{
15 log_d, log_e, log_error_to_statsig_and_console, SpecsSource, StatsigOptions, StatsigRuntime,
16};
17use async_trait::async_trait;
18use chrono::Utc;
19use parking_lot::RwLock;
20use percent_encoding::percent_encode;
21use std::collections::HashMap;
22use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
23use std::sync::{Arc, Weak};
24use std::time::Duration;
25use tokio::sync::Notify;
26use tokio::time::sleep;
27
28use super::SpecsInfo;
29
30pub struct NetworkResponse {
31 pub data: ResponseData,
32 pub loggable_api: String,
33 pub requested_deltas: bool,
34}
35
36pub const DEFAULT_SPECS_URL: &str = "https://statsigcdn.openai.com/v2/download_config_specs";
37pub const DEFAULT_SYNC_INTERVAL_MS: u32 = 10_000;
38
39#[allow(unused)]
40pub const INIT_DICT_ID: &str = "null";
41
42const TAG: &str = stringify!(StatsigHttpSpecsAdapter);
43const STATSIG_NETWORK_FALLBACK_THRESHOLD: u32 = 5;
44
45pub struct StatsigHttpSpecsAdapter {
46 listener: RwLock<Option<Arc<dyn SpecsUpdateListener>>>,
47 network: NetworkClient,
48 sdk_key: String,
49 specs_url: String,
50 fallback_url: Option<String>,
51 init_timeout_ms: u64,
52 sync_interval_duration: Duration,
53 ops_stats: Arc<OpsStatsForInstance>,
54 shutdown_notify: Arc<Notify>,
55 allow_dcs_deltas: bool,
56 use_deltas_next_request: AtomicBool,
57 background_sync_failure_count: AtomicU32,
58}
59
60enum NetworkSyncOutcome {
64 Success,
65 Failure,
66}
67
68impl NetworkSyncOutcome {
69 fn as_bool(&self) -> bool {
70 matches!(self, Self::Success)
71 }
72}
73
74enum ConfigSyncResponseType {
75 Delta,
76 Full,
77 NoUpdate,
78 NetworkError,
79}
80
81impl ConfigSyncResponseType {
82 fn from_response_data(data: &mut ResponseData) -> Self {
83 if is_true_header(data, "x-cache-hit") {
84 return Self::NoUpdate;
85 }
86
87 if data.get_header_ref("x-deltas-used").is_some() {
88 return Self::Delta;
89 }
90
91 let response_type = match data.deserialize_into::<SpecsResponseNoUpdates>() {
92 Ok(response) if !response.has_updates => Self::NoUpdate,
93 _ => Self::Full,
94 };
95 let _ = data.rewind();
96 response_type
97 }
98
99 fn as_str(&self) -> &'static str {
100 match self {
101 Self::Delta => "delta",
102 Self::Full => "full",
103 Self::NoUpdate => "no_update",
104 Self::NetworkError => "network_error",
105 }
106 }
107}
108
109fn is_true_header(data: &ResponseData, key: &str) -> bool {
110 data.get_header_ref(key)
111 .is_some_and(|value| value.eq_ignore_ascii_case("true"))
112}
113
114fn is_process_success(result: &Result<(), StatsigErr>) -> bool {
115 result.is_ok()
116}
117impl StatsigHttpSpecsAdapter {
120 #[must_use]
121 pub fn new(
122 sdk_key: &str,
123 options: Option<&StatsigOptions>,
124 override_url: Option<String>,
125 ) -> Self {
126 let default_options = StatsigOptions::default();
127 let options_ref = options.unwrap_or(&default_options);
128
129 let init_timeout_ms = options_ref
130 .init_timeout_ms
131 .unwrap_or(DEFAULT_INIT_TIMEOUT_MS);
132
133 let specs_url = match override_url {
134 Some(url) => url,
135 None => options_ref
136 .specs_url
137 .as_ref()
138 .map(|u| u.to_string())
139 .unwrap_or(DEFAULT_SPECS_URL.to_string()),
140 };
141
142 let fallback_url = if options_ref.fallback_to_statsig_api.unwrap_or(false)
144 && specs_url != DEFAULT_SPECS_URL
145 {
146 Some(DEFAULT_SPECS_URL.to_string())
147 } else {
148 None
149 };
150
151 let headers = StatsigMetadata::get_constant_request_headers(
152 sdk_key,
153 options_ref.service_name.as_deref(),
154 );
155 let enable_dcs_deltas = options_ref.enable_dcs_deltas.unwrap_or(false);
156
157 let sdk_instance_id = options_ref.get_sdk_instance_id(sdk_key);
158
159 Self {
160 listener: RwLock::new(None),
161 network: NetworkClient::new(sdk_key, Some(headers), Some(options_ref)),
162 sdk_key: sdk_key.to_string(),
163 specs_url,
164 fallback_url,
165 init_timeout_ms,
166 sync_interval_duration: Duration::from_millis(u64::from(
167 options_ref
168 .specs_sync_interval_ms
169 .unwrap_or(DEFAULT_SYNC_INTERVAL_MS),
170 )),
171 ops_stats: OPS_STATS.get_for_instance(sdk_instance_id),
172 shutdown_notify: Arc::new(Notify::new()),
173 allow_dcs_deltas: enable_dcs_deltas,
174 use_deltas_next_request: AtomicBool::new(enable_dcs_deltas),
175 background_sync_failure_count: AtomicU32::new(0),
176 }
177 }
178
179 pub fn force_shutdown(&self) {
180 self.shutdown_notify.notify_one();
181 }
182
183 pub async fn fetch_specs_from_network(
184 &self,
185 current_specs_info: SpecsInfo,
186 trigger: SpecsSyncTrigger,
187 ) -> Result<NetworkResponse, NetworkError> {
188 let request_args = self.get_request_args(¤t_specs_info, trigger);
189 let url = request_args.url.clone();
190 let requested_deltas = request_args.deltas_enabled;
191 match self.handle_specs_request(request_args).await {
192 Ok(response) => Ok(NetworkResponse {
193 data: response,
194 loggable_api: get_api_from_url(&url),
195 requested_deltas,
196 }),
197 Err(e) => Err(e),
198 }
199 }
200
201 fn get_request_args(
202 &self,
203 current_specs_info: &SpecsInfo,
204 trigger: SpecsSyncTrigger,
205 ) -> RequestArgs {
206 let mut params = HashMap::new();
207
208 params.insert("supports_proto".to_string(), "true".to_string());
209 let headers = Some(HashMap::from([
210 ("statsig-supports-proto".to_string(), "true".to_string()),
211 (
212 "accept-encoding".to_string(),
213 "statsig-br, gzip, deflate, br".to_string(),
214 ),
215 ]));
216
217 if let Some(lcut) = current_specs_info.lcut {
218 if lcut > 0 {
219 params.insert("sinceTime".to_string(), lcut.to_string());
220 }
221 }
222
223 let is_init_request = trigger == SpecsSyncTrigger::Initial;
224
225 let timeout_ms = if is_init_request && self.init_timeout_ms > 0 {
226 self.init_timeout_ms
227 } else {
228 0
229 };
230
231 if let Some(cs) = ¤t_specs_info.checksum {
232 params.insert(
233 "checksum".to_string(),
234 percent_encode(cs.as_bytes(), percent_encoding::NON_ALPHANUMERIC).to_string(),
235 );
236 }
237
238 let use_deltas_next_req = self.use_deltas_next_request.load(Ordering::SeqCst);
239 if use_deltas_next_req {
240 params.insert("accept_deltas".to_string(), "true".to_string());
241 }
242
243 RequestArgs {
244 url: construct_specs_url(self.specs_url.as_str(), self.sdk_key.as_str()),
245 retries: match trigger {
246 SpecsSyncTrigger::Initial | SpecsSyncTrigger::Manual => 0,
247 SpecsSyncTrigger::Background => 3,
248 },
249 query_params: Some(params),
250 deltas_enabled: use_deltas_next_req,
251 accept_gzip_response: true,
252 diagnostics_key: Some(KeyType::DownloadConfigSpecs),
253 timeout_ms,
254 headers,
255 ..RequestArgs::new()
256 }
257 }
258
259 async fn handle_fallback_request(
260 &self,
261 mut request_args: RequestArgs,
262 ) -> Result<NetworkResponse, NetworkError> {
263 let requested_deltas = request_args.deltas_enabled;
264 let fallback_url = match &self.fallback_url {
265 Some(url) => construct_specs_url(url.as_str(), &self.sdk_key),
266 None => {
267 return Err(NetworkError::RequestFailed(
268 request_args.url.clone(),
269 None,
270 "No fallback URL".to_string(),
271 ))
272 }
273 };
274
275 request_args.url = fallback_url.clone();
276
277 let response = self.handle_specs_request(request_args).await?;
280 Ok(NetworkResponse {
281 data: response,
282 loggable_api: get_api_from_url(&fallback_url),
283 requested_deltas,
284 })
285 }
286
287 async fn handle_specs_request(
288 &self,
289 request_args: RequestArgs,
290 ) -> Result<ResponseData, NetworkError> {
291 let url = request_args.url.clone();
292 let response = self.network.get(request_args).await?;
293 match response.data {
294 Some(data) => Ok(data),
295 None => Err(NetworkError::RequestFailed(
296 url,
297 None,
298 response.error.unwrap_or("No data in response".to_string()),
299 )),
300 }
301 }
302
303 fn should_attempt_fallback(
304 &self,
305 trigger: SpecsSyncTrigger,
306 result: &Result<(), StatsigErr>,
307 ) -> bool {
308 if result.is_ok() || self.fallback_url.is_none() {
309 return false;
310 }
311
312 if trigger != SpecsSyncTrigger::Background {
313 return true;
314 }
315
316 let failure_count = self
317 .background_sync_failure_count
318 .fetch_add(1, Ordering::SeqCst)
319 + 1;
320
321 if failure_count.is_multiple_of(STATSIG_NETWORK_FALLBACK_THRESHOLD) {
322 return true;
323 }
324
325 log_d!(
326 TAG,
327 "Skipping fallback on background sync failure {}. Retrying fallback every {} failures.",
328 failure_count,
329 STATSIG_NETWORK_FALLBACK_THRESHOLD
330 );
331
332 false
333 }
334
335 pub async fn run_background_sync(self: Arc<Self>) {
336 let specs_info = match self
337 .listener
338 .try_read_for(std::time::Duration::from_secs(5))
339 {
340 Some(lock) => match lock.as_ref() {
341 Some(listener) => listener.get_current_specs_info(),
342 None => SpecsInfo::empty(),
343 },
344 None => SpecsInfo::error(),
345 };
346
347 self.ops_stats
348 .set_diagnostics_context(ContextType::ConfigSync);
349 if let Err(e) = self
350 .manually_sync_specs(specs_info, SpecsSyncTrigger::Background)
351 .await
352 {
353 if let StatsigErr::NetworkError(NetworkError::DisableNetworkOn(_)) = e {
354 return;
355 }
356 log_e!(TAG, "Background specs sync failed: {}", e);
357 }
358 self.ops_stats.enqueue_diagnostics_event(
359 Some(KeyType::DownloadConfigSpecs),
360 Some(ContextType::ConfigSync),
361 );
362 }
363
364 async fn manually_sync_specs(
365 &self,
366 current_specs_info: SpecsInfo,
367 trigger: SpecsSyncTrigger,
368 ) -> Result<(), StatsigErr> {
369 if let Some(lock) = self
370 .listener
371 .try_read_for(std::time::Duration::from_secs(5))
372 {
373 if lock.is_none() {
374 return Err(StatsigErr::UnstartedAdapter("Listener not set".to_string()));
375 }
376 }
377
378 let sync_start_ms = Utc::now().timestamp_millis() as u64;
379 let mut deltas_used = self.use_deltas_next_request.load(Ordering::SeqCst);
380 let mut response = self
381 .fetch_specs_from_network(current_specs_info.clone(), trigger)
382 .await;
383 let mut response_type = response
384 .as_mut()
385 .map_or(ConfigSyncResponseType::NetworkError, |response| {
386 ConfigSyncResponseType::from_response_data(&mut response.data)
387 });
388 let (mut source_api, mut response_format, mut network_success) = match &response {
389 Ok(response) => (
390 response.loggable_api.clone(),
391 get_specs_response_format(&response.data),
392 NetworkSyncOutcome::Success,
393 ),
394 Err(_) => (
395 get_api_from_url(&construct_specs_url(
396 self.specs_url.as_str(),
397 self.sdk_key.as_str(),
398 )),
399 SpecsResponseFormat::Unknown,
400 NetworkSyncOutcome::Failure,
401 ),
402 };
403 if let Ok(response) = &response {
404 deltas_used = response.requested_deltas;
405 }
406
407 let mut result = self.process_spec_data(response).await;
408
409 if self.should_attempt_fallback(trigger, &result) {
410 log_d!(TAG, "Falling back to DCS CDN");
411 let fallback_args = self.get_request_args(¤t_specs_info, trigger);
412 deltas_used = fallback_args.deltas_enabled;
413 let mut response = self.handle_fallback_request(fallback_args).await;
414 response_type = response
415 .as_mut()
416 .map_or(ConfigSyncResponseType::NetworkError, |response| {
417 ConfigSyncResponseType::from_response_data(&mut response.data)
418 });
419 match &response {
420 Ok(response) => {
421 source_api = response.loggable_api.clone();
422 response_format = get_specs_response_format(&response.data);
423 network_success = NetworkSyncOutcome::Success;
424 deltas_used = response.requested_deltas;
425 }
426 Err(_) => {
427 if let Some(fallback_url) = self.fallback_url.as_ref() {
429 source_api = get_api_from_url(&construct_specs_url(
430 fallback_url.as_str(),
431 self.sdk_key.as_str(),
432 ));
433 }
434 network_success = NetworkSyncOutcome::Failure;
435 }
436 }
437 result = self.process_spec_data(response).await;
438 }
439
440 let process_success = is_process_success(&result);
441 log_config_sync_overall_latency(
442 &self.ops_stats,
443 sync_start_ms,
444 source_api.as_str(),
445 response_format.as_str(),
446 network_success.as_bool(),
447 process_success,
448 result
449 .as_ref()
450 .err()
451 .map_or_else(String::new, |e| e.to_string()),
452 deltas_used,
453 response_type.as_str(),
454 );
455
456 result
457 }
458
459 async fn process_spec_data(
460 &self,
461 response: Result<NetworkResponse, NetworkError>,
462 ) -> Result<(), StatsigErr> {
463 let resp = response.map_err(StatsigErr::NetworkError)?;
464 let requested_deltas = resp.requested_deltas;
465
466 let update = SpecsUpdate {
467 data: resp.data,
468 source: SpecsSource::Network,
469 received_at: Utc::now().timestamp_millis() as u64,
470 source_api: Some(resp.loggable_api),
471 has_updates: None,
472 };
473
474 self.ops_stats.add_marker(
475 Marker::new(
476 KeyType::DownloadConfigSpecs,
477 ActionType::Start,
478 Some(StepType::Process),
479 ),
480 None,
481 );
482
483 let result = match self
484 .listener
485 .try_read_for(std::time::Duration::from_secs(5))
486 {
487 Some(lock) => match lock.as_ref() {
488 Some(listener) => listener.did_receive_specs_update(update),
489 None => Err(StatsigErr::UnstartedAdapter("Listener not set".to_string())),
490 },
491 None => {
492 let err =
493 StatsigErr::LockFailure("Failed to acquire read lock on listener".to_string());
494 log_error_to_statsig_and_console!(&self.ops_stats, TAG, err.clone());
495 Err(err)
496 }
497 };
498
499 if matches!(&result, Err(StatsigErr::ChecksumFailure(_))) {
500 let was_deltas_used = self.use_deltas_next_request.swap(false, Ordering::SeqCst);
501 if was_deltas_used {
502 log_d!(TAG, "Disabling delta requests after checksum failure");
503 }
504 } else if result.is_ok() && !requested_deltas && self.allow_dcs_deltas {
505 let was_deltas_used = self.use_deltas_next_request.swap(true, Ordering::SeqCst);
506 if !was_deltas_used {
507 log_d!(
508 TAG,
509 "Re-enabling delta requests after successful non-delta specs update"
510 );
511 }
512 }
513
514 self.ops_stats.add_marker(
515 Marker::new(
516 KeyType::DownloadConfigSpecs,
517 ActionType::End,
518 Some(StepType::Process),
519 )
520 .with_is_success(result.is_ok()),
521 None,
522 );
523
524 result
525 }
526}
527
528#[async_trait]
529impl SpecsAdapter for StatsigHttpSpecsAdapter {
530 async fn start(
531 self: Arc<Self>,
532 _statsig_runtime: &Arc<StatsigRuntime>,
533 ) -> Result<(), StatsigErr> {
534 let specs_info = match self
535 .listener
536 .try_read_for(std::time::Duration::from_secs(5))
537 {
538 Some(lock) => match lock.as_ref() {
539 Some(listener) => listener.get_current_specs_info(),
540 None => SpecsInfo::empty(),
541 },
542 None => SpecsInfo::error(),
543 };
544 self.manually_sync_specs(specs_info, SpecsSyncTrigger::Initial)
545 .await
546 }
547
548 fn initialize(&self, listener: Arc<dyn SpecsUpdateListener>) {
549 match self
550 .listener
551 .try_write_for(std::time::Duration::from_secs(5))
552 {
553 Some(mut lock) => *lock = Some(listener),
554 None => {
555 log_e!(TAG, "Failed to acquire write lock on listener");
556 }
557 }
558 }
559
560 async fn schedule_background_sync(
561 self: Arc<Self>,
562 statsig_runtime: &Arc<StatsigRuntime>,
563 ) -> Result<(), StatsigErr> {
564 let weak_self: Weak<StatsigHttpSpecsAdapter> = Arc::downgrade(&self);
565 let interval_duration = self.sync_interval_duration;
566 let shutdown_notify = self.shutdown_notify.clone();
567
568 statsig_runtime.spawn("http_specs_bg_sync", move |rt_shutdown_notify| async move {
569 loop {
570 tokio::select! {
571 () = sleep(interval_duration) => {
572 if let Some(strong_self) = weak_self.upgrade() {
573 Self::run_background_sync(strong_self).await;
574 } else {
575 log_e!(TAG, "Strong reference to StatsigHttpSpecsAdapter lost. Stopping background sync");
576 break;
577 }
578 }
579 () = rt_shutdown_notify.notified() => {
580 log_d!(TAG, "Runtime shutdown. Shutting down specs background sync");
581 break;
582 },
583 () = shutdown_notify.notified() => {
584 log_d!(TAG, "Shutting down specs background sync");
585 break;
586 }
587 }
588 }
589 })?;
590
591 Ok(())
592 }
593
594 async fn shutdown(
595 &self,
596 _timeout: Duration,
597 _statsig_runtime: &Arc<StatsigRuntime>,
598 ) -> Result<(), StatsigErr> {
599 self.shutdown_notify.notify_one();
600 Ok(())
601 }
602
603 fn get_type_name(&self) -> String {
604 stringify!(StatsigHttpSpecsAdapter).to_string()
605 }
606}
607
608#[allow(unused)]
609fn construct_specs_url(spec_url: &str, sdk_key: &str) -> String {
610 format!("{spec_url}/{sdk_key}.json")
611}
612
613#[derive(Debug, Clone, Copy, PartialEq, Eq)]
614pub enum SpecsSyncTrigger {
615 Initial,
616 Background,
617 Manual,
618}
619
620#[cfg(test)]
621mod tests {
622 use super::*;
623 use crate::{networking::ResponseData, specs_adapter::SpecsUpdate, StatsigOptions};
624 use std::collections::HashMap;
625 use std::sync::atomic::AtomicUsize;
626
627 struct ChecksumFailingListener;
628
629 impl SpecsUpdateListener for ChecksumFailingListener {
630 fn did_receive_specs_update(&self, _update: SpecsUpdate) -> Result<(), StatsigErr> {
631 Err(StatsigErr::ChecksumFailure(
632 "simulated checksum failure".to_string(),
633 ))
634 }
635
636 fn get_current_specs_info(&self) -> SpecsInfo {
637 SpecsInfo::empty()
638 }
639 }
640
641 struct ChecksumFailingThenSuccessListener {
642 calls: AtomicUsize,
643 }
644
645 impl SpecsUpdateListener for ChecksumFailingThenSuccessListener {
646 fn did_receive_specs_update(&self, _update: SpecsUpdate) -> Result<(), StatsigErr> {
647 let curr = self.calls.fetch_add(1, Ordering::SeqCst);
648 if curr == 0 {
649 Err(StatsigErr::ChecksumFailure(
650 "simulated checksum failure".to_string(),
651 ))
652 } else {
653 Ok(())
654 }
655 }
656
657 fn get_current_specs_info(&self) -> SpecsInfo {
658 SpecsInfo::empty()
659 }
660 }
661
662 #[tokio::test]
663 async fn test_disable_accept_deltas_after_checksum_failure() {
664 let options = StatsigOptions {
665 enable_dcs_deltas: Some(true),
666 ..StatsigOptions::default()
667 };
668 let adapter = StatsigHttpSpecsAdapter::new(
669 "secret-key",
670 Some(&options),
671 Some("https://example.com/v2/download_config_specs".to_string()),
672 );
673 let specs_info = SpecsInfo::empty();
674
675 let request_before = adapter.get_request_args(&specs_info, SpecsSyncTrigger::Manual);
676 assert_eq!(
677 request_before
678 .query_params
679 .as_ref()
680 .and_then(|p| p.get("accept_deltas"))
681 .map(String::as_str),
682 Some("true")
683 );
684
685 adapter.initialize(Arc::new(ChecksumFailingListener));
686 let result = adapter
687 .process_spec_data(Ok(NetworkResponse {
688 data: ResponseData::from_bytes(vec![]),
689 loggable_api: "test-api".to_string(),
690 requested_deltas: true,
691 }))
692 .await;
693
694 assert!(matches!(result, Err(StatsigErr::ChecksumFailure(_))));
695
696 let request_after = adapter.get_request_args(&specs_info, SpecsSyncTrigger::Manual);
697 assert!(request_after
698 .query_params
699 .as_ref()
700 .is_none_or(|p| !p.contains_key("accept_deltas")));
701 }
702
703 #[tokio::test]
704 async fn test_reenable_accept_deltas_after_successful_non_delta_update() {
705 let options = StatsigOptions {
706 enable_dcs_deltas: Some(true),
707 ..StatsigOptions::default()
708 };
709 let adapter = StatsigHttpSpecsAdapter::new(
710 "secret-key",
711 Some(&options),
712 Some("https://example.com/v2/download_config_specs".to_string()),
713 );
714 let specs_info = SpecsInfo::empty();
715
716 adapter.initialize(Arc::new(ChecksumFailingThenSuccessListener {
717 calls: AtomicUsize::new(0),
718 }));
719
720 let first_result = adapter
721 .process_spec_data(Ok(NetworkResponse {
722 data: ResponseData::from_bytes(vec![]),
723 loggable_api: "test-api".to_string(),
724 requested_deltas: true,
725 }))
726 .await;
727
728 assert!(matches!(first_result, Err(StatsigErr::ChecksumFailure(_))));
729
730 let request_after_failure = adapter.get_request_args(&specs_info, SpecsSyncTrigger::Manual);
731 assert!(request_after_failure
732 .query_params
733 .as_ref()
734 .is_none_or(|p| !p.contains_key("accept_deltas")));
735
736 let second_result = adapter
737 .process_spec_data(Ok(NetworkResponse {
738 data: ResponseData::from_bytes(vec![]),
739 loggable_api: "test-api".to_string(),
740 requested_deltas: false,
741 }))
742 .await;
743
744 assert!(second_result.is_ok());
745
746 let request_after_success = adapter.get_request_args(&specs_info, SpecsSyncTrigger::Manual);
747 assert_eq!(
748 request_after_success
749 .query_params
750 .as_ref()
751 .and_then(|p| p.get("accept_deltas"))
752 .map(String::as_str),
753 Some("true")
754 );
755 }
756
757 #[test]
758 fn test_checksum_failure_is_not_process_success() {
759 let result = Err(StatsigErr::ChecksumFailure(
760 "simulated checksum failure".to_string(),
761 ));
762
763 assert!(!is_process_success(&result));
764 assert!(is_process_success(&Ok(())));
765 }
766
767 #[test]
768 fn test_fallback_uses_openai_cdn_for_non_default_specs_url() {
769 let options = StatsigOptions {
770 fallback_to_statsig_api: Some(true),
771 ..StatsigOptions::default()
772 };
773 let adapter = StatsigHttpSpecsAdapter::new(
774 "secret-key",
775 Some(&options),
776 Some("https://example.com/v2/download_config_specs".to_string()),
777 );
778
779 assert_eq!(adapter.fallback_url.as_deref(), Some(DEFAULT_SPECS_URL));
780 }
781
782 #[test]
783 fn test_config_sync_response_type_delta() {
784 let mut headers = HashMap::new();
785 headers.insert("x-deltas-used".to_string(), "true".to_string());
786 let mut data = ResponseData::from_bytes_with_headers(
787 b"{\"has_updates\": false}".to_vec(),
788 Some(headers),
789 );
790
791 let response_type = ConfigSyncResponseType::from_response_data(&mut data);
792
793 assert_eq!(response_type.as_str(), "delta");
794 }
795
796 #[test]
797 fn test_config_sync_response_type_no_update() {
798 let mut data = ResponseData::from_bytes(b"{\"has_updates\": false}".to_vec());
799
800 let response_type = ConfigSyncResponseType::from_response_data(&mut data);
801
802 assert_eq!(response_type.as_str(), "no_update");
803 let response = data.deserialize_into::<SpecsResponseNoUpdates>().unwrap();
804 assert!(!response.has_updates);
805 }
806
807 #[test]
808 fn test_config_sync_response_type_no_update_with_delta_header() {
809 let mut headers = HashMap::new();
810 headers.insert("x-cache-hit".to_string(), "true".to_string());
811 headers.insert("x-deltas-used".to_string(), "true".to_string());
812 let mut data = ResponseData::from_bytes_with_headers(
813 b"{\"has_updates\": false}".to_vec(),
814 Some(headers),
815 );
816
817 let response_type = ConfigSyncResponseType::from_response_data(&mut data);
818
819 assert_eq!(response_type.as_str(), "no_update");
820 }
821
822 #[test]
823 fn test_config_sync_response_type_full() {
824 let mut data = ResponseData::from_bytes(b"{\"has_updates\": true}".to_vec());
825
826 let response_type = ConfigSyncResponseType::from_response_data(&mut data);
827
828 assert_eq!(response_type.as_str(), "full");
829 }
830
831 #[test]
832 fn test_config_sync_response_type_full_for_non_json_payload() {
833 let mut data = ResponseData::from_bytes(vec![0, 1, 2]);
834
835 let response_type = ConfigSyncResponseType::from_response_data(&mut data);
836
837 assert_eq!(response_type.as_str(), "full");
838 let mut first_byte = [1];
839 data.get_stream_mut().read_exact(&mut first_byte).unwrap();
840 assert_eq!(first_byte, [0]);
841 }
842
843 #[test]
844 fn test_get_response_format_json() {
845 let mut headers = HashMap::new();
846 headers.insert("content-type".to_string(), "application/json".to_string());
847 let data = ResponseData::from_bytes_with_headers(vec![], Some(headers));
848 assert!(matches!(
849 get_specs_response_format(&data),
850 SpecsResponseFormat::Json
851 ));
852 }
853
854 #[test]
855 fn test_get_response_format_plain_text() {
856 let mut headers = HashMap::new();
857 headers.insert(
858 "content-type".to_string(),
859 "text/plain; charset=utf-8".to_string(),
860 );
861 let data = ResponseData::from_bytes_with_headers(vec![], Some(headers));
862 assert!(matches!(
863 get_specs_response_format(&data),
864 SpecsResponseFormat::PlainText
865 ));
866 }
867
868 #[test]
869 fn test_get_response_format_protobuf() {
870 let mut headers = HashMap::new();
871 headers.insert(
872 "content-type".to_string(),
873 "application/octet-stream".to_string(),
874 );
875 headers.insert("content-encoding".to_string(), "statsig-br".to_string());
876 let data = ResponseData::from_bytes_with_headers(vec![], Some(headers));
877 assert!(matches!(
878 get_specs_response_format(&data),
879 SpecsResponseFormat::Protobuf
880 ));
881 }
882
883 #[test]
884 fn test_get_response_format_unknown_without_content_type() {
885 let data = ResponseData::from_bytes(vec![]);
886 assert!(matches!(
887 get_specs_response_format(&data),
888 SpecsResponseFormat::Unknown
889 ));
890 }
891}