Skip to main content

statsig_rust/specs_adapter/
statsig_http_specs_adapter.rs

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
60// OB client -- START
61// These types are only for the config_sync_overall.latency observability metric added in this change.
62
63enum 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}
117// OB client -- END
118
119impl 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        // only fallback when the spec_url is not the DEFAULT_SPECS_URL
143        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(&current_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) = &current_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        // TODO logging
278
279        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(&current_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                    // Backup request failed, so no successful network payload was returned.
428                    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}