Skip to main content

sabi/tokio/data_src/
mod.rs

1// Copyright (C) 2024-2026 Takayuki Sato. All Rights Reserved.
2// This program is free software under MIT License.
3// See the file LICENSE in this distribution for more details.
4
5mod global_setup;
6
7pub(crate) use global_setup::{
8    copy_global_data_srcs_to_map, create_data_conn_from_global_data_src_async,
9};
10pub use global_setup::{
11    create_static_data_src_container, setup_async, setup_with_order_async, uses, uses_async,
12};
13
14use crate::tokio::{
15    AsyncGroup, DataConn, DataConnContainer, DataSrc, DataSrcContainer, DataSrcManager, ErrEntry,
16    SendSyncNonNull,
17};
18
19use std::collections::HashMap;
20use std::future::Future;
21use std::pin::Pin;
22use std::sync::Arc;
23use std::{any, mem, ptr};
24
25/// Represents errors that can occur during data source operations.
26#[derive(Debug)]
27pub enum DataSrcError {
28    /// An error indicating a failure to register a global data source.
29    /// This can happen if the global data source manager is in an invalid state.
30    FailToRegisterGlobalDataSrc {
31        /// The name of the data source that failed to register.
32        name: Arc<str>,
33    },
34
35    /// An error indicating that one or more global data sources failed during their setup process.
36    FailToSetupGlobalDataSrcs {
37        /// A vector of errors, each containing the name of the data source and the error itself.
38        errors: Vec<ErrEntry>,
39    },
40
41    /// An error indicating that a global data source setup is currently in progress.
42    DuringSetupGlobalDataSrcs,
43
44    /// An error indicating that global data sources have already been set up.
45    AlreadySetupGlobalDataSrcs,
46
47    /// An error indicating that a data connection could not be cast to the target type.
48    FailToCastDataConn {
49        /// The name of the data source that failed to provide the correct connection type.
50        name: Arc<str>,
51
52        /// The string representation of the target data connection type that was requested.
53        target_type: &'static str,
54    },
55
56    /// An error indicating that a data connection could not be created by its data source.
57    FailToCreateDataConn {
58        /// The name of the data source that failed to create a data connection.
59        name: Arc<str>,
60
61        /// The string representation of the data connection type that was requested.
62        data_conn_type: &'static str,
63    },
64
65    /// An error indicating that no data source was found for the requested data connection.
66    NotFoundDataSrcToCreateDataConn {
67        /// The name of the data source that was not found.
68        name: Arc<str>,
69
70        /// The string representation of the data connection type that was requested.
71        data_conn_type: &'static str,
72    },
73}
74
75impl<S, C> DataSrcContainer<S, C>
76where
77    S: DataSrc<C> + 'static,
78    C: DataConn + 'static,
79{
80    pub(crate) fn new(name: impl Into<Arc<str>>, data_src: S, local: bool) -> Self {
81        Self {
82            drop_fn: drop_data_src::<S, C>,
83            close_fn: close_data_src::<S, C>,
84            is_data_conn_fn: is_data_conn::<C>,
85
86            setup_fn: setup_data_src_async::<S, C>,
87            create_data_conn_fn: create_data_conn_async::<S, C>,
88
89            local,
90            name: name.into(),
91            data_src,
92        }
93    }
94}
95
96fn drop_data_src<S, C>(ptr: *const DataSrcContainer)
97where
98    S: DataSrc<C>,
99    C: DataConn + 'static,
100{
101    let typed_ptr = ptr as *mut DataSrcContainer<S, C>;
102    drop(unsafe { Box::from_raw(typed_ptr) });
103}
104
105fn close_data_src<S, C>(ptr: *const DataSrcContainer)
106where
107    S: DataSrc<C>,
108    C: DataConn + 'static,
109{
110    let typed_ptr = ptr as *mut DataSrcContainer<S, C>;
111    unsafe { (*typed_ptr).data_src.close() };
112}
113
114fn is_data_conn<C>(type_id: any::TypeId) -> bool
115where
116    C: DataConn + 'static,
117{
118    any::TypeId::of::<C>() == type_id
119}
120
121fn setup_data_src_async<S, C>(
122    ptr: *const DataSrcContainer,
123    ag: &mut AsyncGroup,
124) -> Pin<Box<dyn Future<Output = errs::Result<()>> + Send + '_>>
125where
126    S: DataSrc<C> + 'static,
127    C: DataConn + 'static,
128{
129    let typed_ptr = ptr as *mut DataSrcContainer<S, C>;
130    let data_src = unsafe { &mut (*typed_ptr).data_src };
131    Box::pin(data_src.setup_async(ag))
132}
133
134#[allow(clippy::type_complexity)]
135fn create_data_conn_async<'a, S, C>(
136    ptr: *const DataSrcContainer,
137) -> Pin<Box<dyn Future<Output = errs::Result<Box<DataConnContainer<C>>>> + Send + 'a>>
138where
139    S: DataSrc<C> + 'a,
140    C: DataConn + 'static,
141{
142    let typed_ptr = ptr as *mut DataSrcContainer<S, C>;
143    let data_src = unsafe { &mut (*typed_ptr).data_src };
144    let name = unsafe { &(*typed_ptr).name };
145    Box::pin(async move {
146        let conn: Box<C> = data_src.create_data_conn_async().await?;
147        Ok(Box::new(DataConnContainer::<C>::new(
148            name.to_string(),
149            conn,
150        )))
151    })
152}
153
154impl DataSrcManager {
155    pub(crate) const fn new(local: bool) -> Self {
156        Self {
157            vec_unready: Vec::new(),
158            vec_ready: Vec::new(),
159            local,
160        }
161    }
162
163    pub(crate) fn prepend(&mut self, vec: Vec<SendSyncNonNull<DataSrcContainer>>) {
164        self.vec_unready.splice(0..0, vec);
165    }
166
167    pub(crate) fn add<S, C>(&mut self, name: impl Into<Arc<str>>, ds: S)
168    where
169        S: DataSrc<C> + 'static,
170        C: DataConn + 'static,
171    {
172        let boxed = Box::new(DataSrcContainer::<S, C>::new(name, ds, self.local));
173        let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
174        self.vec_unready.push(SendSyncNonNull::new(ptr));
175    }
176
177    pub(crate) fn remove(&mut self, name: impl AsRef<str>) {
178        let extracted_vec: Vec<_> = self
179            .vec_ready
180            .extract_if(.., |ssnnptr| {
181                unsafe { &(*ssnnptr.non_null_ptr.as_ptr()).name }.as_ref() == name.as_ref()
182            })
183            .collect();
184
185        for ssnnptr in extracted_vec.iter().rev() {
186            let ptr = ssnnptr.non_null_ptr.as_ptr();
187            let close_fn = unsafe { (*ptr).close_fn };
188            let drop_fn = unsafe { (*ptr).drop_fn };
189            close_fn(ptr);
190            drop_fn(ptr);
191        }
192
193        let extracted_vec: Vec<_> = self
194            .vec_unready
195            .extract_if(.., |ssnnptr| {
196                unsafe { &(*ssnnptr.non_null_ptr.as_ptr()).name }.as_ref() == name.as_ref()
197            })
198            .collect();
199
200        for ssnnptr in extracted_vec.iter().rev() {
201            let ptr = ssnnptr.non_null_ptr.as_ptr();
202            let drop_fn = unsafe { (*ptr).drop_fn };
203            drop_fn(ptr);
204        }
205    }
206
207    pub(crate) fn close(&mut self) {
208        let vec = mem::take(&mut self.vec_ready);
209        for ssnnptr in vec.into_iter().rev() {
210            let ptr = ssnnptr.non_null_ptr.as_ptr();
211            let close_fn = unsafe { (*ptr).close_fn };
212            let drop_fn = unsafe { (*ptr).drop_fn };
213            close_fn(ptr);
214            drop_fn(ptr);
215        }
216        let vec = mem::take(&mut self.vec_unready);
217        for ssnnptr in vec.into_iter().rev() {
218            let ptr = ssnnptr.non_null_ptr.as_ptr();
219            let drop_fn = unsafe { (*ptr).drop_fn };
220            drop_fn(ptr);
221        }
222    }
223
224    pub(crate) async fn setup_async(&mut self, errors: &mut Vec<ErrEntry>) {
225        if self.vec_unready.is_empty() {
226            return;
227        }
228
229        let mut ag = AsyncGroup::new();
230        for (i, ssnnptr) in self.vec_unready.iter().enumerate() {
231            let ptr = ssnnptr.non_null_ptr.as_ptr();
232            let setup_fn = unsafe { (*ptr).setup_fn };
233            let name = unsafe { &(*ptr).name };
234            ag._index = i;
235            ag._name = name.clone();
236            if let Err(err) = setup_fn(ptr, &mut ag).await {
237                errors.push(ErrEntry {
238                    index: i,
239                    name: name.clone(),
240                    err,
241                });
242                break;
243            }
244        }
245        let n_done = ag._index;
246        ag.join_and_collect_errors_async(errors).await;
247
248        if errors.is_empty() {
249            self.vec_ready.append(&mut self.vec_unready);
250        } else {
251            for ssnnptr in self.vec_unready[0..n_done].iter().rev() {
252                let ptr = ssnnptr.non_null_ptr.as_ptr();
253                let close_fn = unsafe { (*ptr).close_fn };
254                close_fn(ptr);
255            }
256        }
257    }
258
259    pub(crate) async fn setup_with_order_async(
260        &mut self,
261        names: &[&str],
262        errors: &mut Vec<ErrEntry>,
263    ) {
264        if self.vec_unready.is_empty() {
265            return;
266        }
267
268        let mut index_map: HashMap<&str, usize> = HashMap::with_capacity(names.len());
269        // Using rev because earlier ones take precedence when names overlap
270        for (i, nm) in names.iter().rev().enumerate() {
271            index_map.insert(*nm, names.len() - 1 - i);
272        }
273
274        let mut ordered_indexes = Vec::<Option<usize>>::with_capacity(self.vec_unready.len());
275        ordered_indexes.resize(names.len(), None);
276
277        for vec_index in 0..self.vec_unready.len() {
278            let ssnnptr = &self.vec_unready[vec_index];
279            let ptr = ssnnptr.non_null_ptr.as_ptr();
280            let name = unsafe { (*ptr).name.clone() };
281            if let Some(order_index) = index_map.remove(name.as_ref()) {
282                ordered_indexes[order_index] = Some(vec_index);
283            } else {
284                ordered_indexes.push(Some(vec_index));
285            }
286        }
287
288        let mut ag = AsyncGroup::new();
289        let mut n_done = 0;
290        for vec_index_opt in ordered_indexes.iter() {
291            if let Some(vec_index) = vec_index_opt {
292                let ssnnptr = &self.vec_unready[*vec_index];
293                let ptr = ssnnptr.non_null_ptr.as_ptr();
294                let setup_fn = unsafe { (*ptr).setup_fn };
295                let name = unsafe { &(*ptr).name };
296                ag._index = *vec_index;
297                ag._name = name.clone();
298                if let Err(err) = setup_fn(ptr, &mut ag).await {
299                    errors.push(ErrEntry {
300                        index: *vec_index,
301                        name: name.clone(),
302                        err,
303                    });
304                    break;
305                }
306            }
307            n_done += 1;
308        }
309        ag.join_and_collect_errors_async(errors).await;
310
311        if errors.is_empty() {
312            let old_unready = mem::take(&mut self.vec_unready);
313            // Maximizing performance by pre-allocating the required capacity.
314            self.vec_ready
315                .reserve(ordered_indexes.iter().flatten().count());
316            for vec_index in ordered_indexes.iter().flatten() {
317                self.vec_ready.push(old_unready[*vec_index]);
318            }
319        } else {
320            for vec_index in ordered_indexes.iter().take(n_done).flatten() {
321                let ssnnptr = &self.vec_unready[*vec_index];
322                let ptr = ssnnptr.non_null_ptr.as_ptr();
323                let close_fn = unsafe { (*ptr).close_fn };
324                close_fn(ptr);
325            }
326        }
327    }
328
329    pub(crate) fn copy_ds_ready_to_map(&self, index_map: &mut HashMap<Arc<str>, (bool, usize)>) {
330        for (i, ssnnptr) in self.vec_ready.iter().enumerate() {
331            let ptr = ssnnptr.non_null_ptr.as_ptr();
332            let name = unsafe { (*ptr).name.clone() };
333            index_map.insert(name, (self.local, i));
334        }
335    }
336
337    pub(crate) async fn create_data_conn_async<C>(
338        &self,
339        index: usize,
340        name: impl AsRef<str>,
341    ) -> errs::Result<Box<DataConnContainer>>
342    where
343        C: DataConn + 'static,
344    {
345        if let Some(ssnnptr) = self.vec_ready.get(index) {
346            let ptr = ssnnptr.non_null_ptr.as_ptr();
347            let type_id = any::TypeId::of::<C>();
348            let is_fn = unsafe { (*ptr).is_data_conn_fn };
349            let create_data_conn_fn = unsafe { (*ptr).create_data_conn_fn };
350            if !is_fn(type_id) {
351                Err(errs::Err::new(DataSrcError::FailToCastDataConn {
352                    name: name.as_ref().into(),
353                    target_type: any::type_name::<C>(),
354                }))
355            } else {
356                match create_data_conn_fn(ptr).await {
357                    Ok(boxed) => Ok(boxed),
358                    Err(err) => Err(errs::Err::with_source(
359                        DataSrcError::FailToCreateDataConn {
360                            name: name.as_ref().into(),
361                            data_conn_type: any::type_name::<C>(),
362                        },
363                        err,
364                    )),
365                }
366            }
367        } else {
368            Err(errs::Err::new(
369                DataSrcError::NotFoundDataSrcToCreateDataConn {
370                    name: name.as_ref().into(),
371                    data_conn_type: any::type_name::<C>(),
372                },
373            ))
374        }
375    }
376}
377
378impl Drop for DataSrcManager {
379    fn drop(&mut self) {
380        self.close();
381    }
382}
383
384#[cfg_attr(coverage_nightly, coverage(off))]
385#[cfg(test)]
386mod tests_of_data_src {
387    use super::*;
388    use std::sync::Arc;
389    use tokio::sync::Mutex;
390
391    struct SyncDataConn {}
392    impl SyncDataConn {
393        fn new() -> Self {
394            Self {}
395        }
396    }
397    impl DataConn for SyncDataConn {
398        async fn commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
399            Ok(())
400        }
401        fn is_committed(&self) -> bool {
402            false
403        }
404        async fn rollback_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
405            Ok(())
406        }
407        fn close(&mut self) {}
408    }
409
410    struct AsyncDataConn {}
411    impl AsyncDataConn {
412        fn new() -> Self {
413            Self {}
414        }
415    }
416    impl DataConn for AsyncDataConn {
417        async fn commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
418            Ok(())
419        }
420        fn is_committed(&self) -> bool {
421            false
422        }
423        async fn rollback_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
424            Ok(())
425        }
426        fn close(&mut self) {}
427    }
428
429    struct SyncDataSrc {
430        id: i8,
431        logger: Arc<Mutex<Vec<String>>>,
432        fail_to_setup: bool,
433        fail_to_create_data_conn: bool,
434    }
435    impl SyncDataSrc {
436        fn new(id: i8, logger: Arc<Mutex<Vec<String>>>, fail_to_setup: bool) -> Self {
437            let logger_clone = logger.clone();
438            tokio::spawn(async move {
439                logger_clone
440                    .lock()
441                    .await
442                    .push(format!("SyncDataSrc::new {}", id));
443            });
444            Self {
445                id,
446                logger: logger,
447                fail_to_setup,
448                fail_to_create_data_conn: false,
449            }
450        }
451        fn new_for_fail_to_create_data_conn(id: i8, logger: Arc<Mutex<Vec<String>>>) -> Self {
452            Self {
453                id,
454                logger: logger,
455                fail_to_setup: false,
456                fail_to_create_data_conn: true,
457            }
458        }
459    }
460    impl Drop for SyncDataSrc {
461        fn drop(&mut self) {
462            let logger = self.logger.clone();
463            let id = self.id;
464            tokio::spawn(async move {
465                logger
466                    .lock()
467                    .await
468                    .push(format!("SyncDataSrc::drop {}", id));
469            });
470        }
471    }
472    impl DataSrc<SyncDataConn> for SyncDataSrc {
473        async fn setup_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
474            let fail = self.fail_to_setup;
475            let id = self.id;
476            let logger = self.logger.clone();
477
478            if fail {
479                logger
480                    .lock()
481                    .await
482                    .push(format!("SyncDataSrc::setup {} failed", id));
483                return Err(errs::Err::new("XXX".to_string()));
484            }
485            logger
486                .lock()
487                .await
488                .push(format!("SyncDataSrc::setup {}", id));
489            Ok(())
490        }
491        fn close(&mut self) {
492            let logger = self.logger.clone();
493            let id = self.id;
494            tokio::spawn(async move {
495                logger
496                    .lock()
497                    .await
498                    .push(format!("SyncDataSrc::close {}", id));
499            });
500        }
501        async fn create_data_conn_async(&mut self) -> errs::Result<Box<SyncDataConn>> {
502            let id = self.id;
503            let logger = self.logger.clone();
504            {
505                logger
506                    .lock()
507                    .await
508                    .push(format!("SyncDataSrc::create_data_conn {}", id));
509            }
510            if self.fail_to_create_data_conn {
511                return Err(errs::Err::new("eeee".to_string()));
512            }
513            let conn = SyncDataConn::new();
514            Ok(Box::new(conn))
515        }
516    }
517
518    struct AsyncDataSrc {
519        id: i8,
520        fail: bool,
521        logger: Arc<Mutex<Vec<String>>>,
522        wait: u64,
523    }
524    impl AsyncDataSrc {
525        fn new(id: i8, logger: Arc<Mutex<Vec<String>>>, fail: bool, wait: u64) -> Self {
526            let logger_clone = logger.clone();
527            tokio::spawn(async move {
528                logger_clone
529                    .lock()
530                    .await
531                    .push(format!("AsyncDataSrc::new {}", id));
532            });
533            Self {
534                id,
535                fail,
536                logger,
537                wait,
538            }
539        }
540    }
541    impl Drop for AsyncDataSrc {
542        fn drop(&mut self) {
543            let logger = self.logger.clone();
544            let id = self.id;
545            tokio::spawn(async move {
546                logger
547                    .lock()
548                    .await
549                    .push(format!("AsyncDataSrc::drop {}", id));
550            });
551        }
552    }
553    impl DataSrc<AsyncDataConn> for AsyncDataSrc {
554        async fn setup_async(&mut self, ag: &mut AsyncGroup) -> errs::Result<()> {
555            let logger = self.logger.clone();
556            let fail = self.fail;
557            let id = self.id;
558            let wait = self.wait;
559
560            ag.add(async move {
561                tokio::time::sleep(std::time::Duration::from_millis(wait)).await;
562                let mut logger = logger.lock().await;
563                if fail {
564                    logger.push(format!("AsyncDataSrc::setup {} failed to setup", id));
565                    return Err(errs::Err::new("XXX".to_string()));
566                }
567                logger.push(format!("AsyncDataSrc::setup {}", id));
568                Ok(())
569            });
570            Ok(())
571        }
572        fn close(&mut self) {
573            let logger = self.logger.clone();
574            let id = self.id;
575            tokio::spawn(async move {
576                logger
577                    .lock()
578                    .await
579                    .push(format!("AsyncDataSrc::close {}", id));
580            });
581        }
582        async fn create_data_conn_async(&mut self) -> errs::Result<Box<AsyncDataConn>> {
583            let logger = self.logger.clone();
584            {
585                logger
586                    .lock()
587                    .await
588                    .push(format!("AsyncDataSrc::create_data_conn {}", self.id));
589            }
590            let conn = AsyncDataConn::new();
591            Ok(Box::new(conn))
592        }
593    }
594
595    #[tokio::test]
596    async fn test_of_new() {
597        let manager = DataSrcManager::new(true);
598        assert!(manager.local);
599        assert_eq!(manager.vec_unready.len(), 0);
600        assert_eq!(manager.vec_ready.len(), 0);
601
602        let manager = DataSrcManager::new(false);
603        assert!(!manager.local);
604        assert_eq!(manager.vec_unready.len(), 0);
605        assert_eq!(manager.vec_ready.len(), 0);
606    }
607
608    #[tokio::test]
609    async fn test_of_prepend() {
610        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
611
612        {
613            let mut vec = Vec::<SendSyncNonNull<DataSrcContainer>>::new();
614
615            let ds = SyncDataSrc::new(1, logger.clone(), false);
616            let boxed = Box::new(DataSrcContainer::new("foo", ds, true));
617            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
618            vec.push(SendSyncNonNull::new(ptr));
619
620            let ds = AsyncDataSrc::new(2, logger.clone(), false, 0);
621            let boxed = Box::new(DataSrcContainer::new("bar", ds, true));
622            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
623            vec.push(SendSyncNonNull::new(ptr));
624
625            let mut manager = DataSrcManager::new(true);
626            manager.prepend(vec);
627
628            assert!(manager.local);
629            assert_eq!(manager.vec_unready.len(), 2);
630            assert_eq!(manager.vec_ready.len(), 0);
631
632            assert_eq!(
633                unsafe { manager.vec_unready[0].non_null_ptr.as_ref().name.clone() },
634                "foo".into()
635            );
636            assert_eq!(
637                unsafe { manager.vec_unready[1].non_null_ptr.as_ref().name.clone() },
638                "bar".into()
639            );
640
641            let mut vec = Vec::<SendSyncNonNull<DataSrcContainer>>::new();
642
643            let ds = SyncDataSrc::new(3, logger.clone(), false);
644            let boxed = Box::new(DataSrcContainer::new("baz", ds, true));
645            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
646            vec.push(SendSyncNonNull::new(ptr));
647
648            let ds = AsyncDataSrc::new(4, logger.clone(), false, 0);
649            let boxed = Box::new(DataSrcContainer::new("qux", ds, true));
650            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
651            vec.push(SendSyncNonNull::new(ptr));
652
653            manager.prepend(vec);
654
655            assert!(manager.local);
656            assert_eq!(manager.vec_unready.len(), 4);
657            assert_eq!(manager.vec_ready.len(), 0);
658
659            assert_eq!(
660                unsafe { manager.vec_unready[0].non_null_ptr.as_ref().name.clone() },
661                "baz".into()
662            );
663            assert_eq!(
664                unsafe { manager.vec_unready[1].non_null_ptr.as_ref().name.clone() },
665                "qux".into()
666            );
667            assert_eq!(
668                unsafe { manager.vec_unready[2].non_null_ptr.as_ref().name.clone() },
669                "foo".into()
670            );
671            assert_eq!(
672                unsafe { manager.vec_unready[3].non_null_ptr.as_ref().name.clone() },
673                "bar".into()
674            );
675        }
676
677        // Give some time for drop to be called
678        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
679
680        let locked = logger.lock().await;
681        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
682        assert!(locked.contains(&"AsyncDataSrc::new 2".to_string()));
683        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
684        assert!(locked.contains(&"AsyncDataSrc::new 4".to_string()));
685        assert!(locked.contains(&"AsyncDataSrc::drop 2".to_string()));
686        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
687        assert!(locked.contains(&"AsyncDataSrc::drop 4".to_string()));
688        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
689    }
690
691    #[tokio::test]
692    async fn test_of_add() {
693        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
694
695        {
696            let mut manager = DataSrcManager::new(true);
697
698            let ds = SyncDataSrc::new(1, logger.clone(), false);
699            manager.add("foo", ds);
700
701            assert!(manager.local);
702            assert_eq!(manager.vec_unready.len(), 1);
703            assert_eq!(manager.vec_ready.len(), 0);
704
705            assert_eq!(
706                unsafe { manager.vec_unready[0].non_null_ptr.as_ref().name.clone() },
707                "foo".into()
708            );
709
710            let ds = AsyncDataSrc::new(2, logger.clone(), false, 0);
711            manager.add("bar", ds);
712
713            assert!(manager.local);
714            assert_eq!(manager.vec_unready.len(), 2);
715            assert_eq!(manager.vec_ready.len(), 0);
716
717            assert_eq!(
718                unsafe { manager.vec_unready[0].non_null_ptr.as_ref().name.clone() },
719                "foo".into()
720            );
721            assert_eq!(
722                unsafe { manager.vec_unready[1].non_null_ptr.as_ref().name.clone() },
723                "bar".into()
724            );
725        }
726
727        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
728
729        let locked = logger.lock().await;
730        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
731        assert!(locked.contains(&"AsyncDataSrc::new 2".to_string()));
732        assert!(locked.contains(&"AsyncDataSrc::drop 2".to_string()));
733        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
734    }
735
736    #[tokio::test]
737    async fn test_of_remove() {
738        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
739
740        {
741            let mut manager = DataSrcManager::new(true);
742
743            let ds1 = SyncDataSrc::new(1, logger.clone(), false);
744            let boxed = Box::new(DataSrcContainer::new("foo", ds1, true));
745            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
746            manager.vec_unready.push(SendSyncNonNull::new(ptr));
747
748            let ds2 = AsyncDataSrc::new(2, logger.clone(), false, 0);
749            let boxed = Box::new(DataSrcContainer::new("bar", ds2, true));
750            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
751            manager.vec_unready.push(SendSyncNonNull::new(ptr));
752
753            let ds3 = SyncDataSrc::new(3, logger.clone(), false);
754            let boxed = Box::new(DataSrcContainer::new("baz", ds3, true));
755            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
756            manager.vec_ready.push(SendSyncNonNull::new(ptr));
757
758            let ds4 = AsyncDataSrc::new(4, logger.clone(), false, 0);
759            let boxed = Box::new(DataSrcContainer::new("qux", ds4, true));
760            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
761            manager.vec_ready.push(SendSyncNonNull::new(ptr));
762
763            assert!(manager.local);
764            assert_eq!(manager.vec_unready.len(), 2);
765            assert_eq!(manager.vec_ready.len(), 2);
766
767            manager.remove("baz");
768            manager.remove("foo");
769            manager.remove("qux");
770            manager.remove("bar");
771        }
772
773        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
774
775        let locked = logger.lock().await;
776        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
777        assert!(locked.contains(&"AsyncDataSrc::new 2".to_string()));
778        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
779        assert!(locked.contains(&"AsyncDataSrc::new 4".to_string()));
780        assert!(locked.contains(&"SyncDataSrc::close 3".to_string()));
781        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
782        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
783        assert!(locked.contains(&"AsyncDataSrc::close 4".to_string()));
784        assert!(locked.contains(&"AsyncDataSrc::drop 4".to_string()));
785        assert!(locked.contains(&"AsyncDataSrc::drop 2".to_string()));
786    }
787
788    #[tokio::test]
789    async fn test_of_close() {
790        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
791
792        {
793            let mut manager = DataSrcManager::new(true);
794
795            let ds1 = SyncDataSrc::new(1, logger.clone(), false);
796            let boxed = Box::new(DataSrcContainer::new("foo", ds1, true));
797            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
798            manager.vec_unready.push(SendSyncNonNull::new(ptr));
799
800            let ds2 = AsyncDataSrc::new(2, logger.clone(), false, 0);
801            let boxed = Box::new(DataSrcContainer::new("bar", ds2, true));
802            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
803            manager.vec_unready.push(SendSyncNonNull::new(ptr));
804
805            let ds3 = SyncDataSrc::new(3, logger.clone(), false);
806            let boxed = Box::new(DataSrcContainer::new("baz", ds3, true));
807            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
808            manager.vec_ready.push(SendSyncNonNull::new(ptr));
809
810            let ds4 = AsyncDataSrc::new(4, logger.clone(), false, 0);
811            let boxed = Box::new(DataSrcContainer::new("qux", ds4, true));
812            let ptr = ptr::NonNull::from(Box::leak(boxed)).cast::<DataSrcContainer>();
813            manager.vec_ready.push(SendSyncNonNull::new(ptr));
814
815            assert!(manager.local);
816            assert_eq!(manager.vec_unready.len(), 2);
817            assert_eq!(manager.vec_ready.len(), 2);
818
819            manager.close();
820        }
821
822        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
823
824        let locked = logger.lock().await;
825        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
826        assert!(locked.contains(&"AsyncDataSrc::new 2".to_string()));
827        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
828        assert!(locked.contains(&"AsyncDataSrc::new 4".to_string()));
829        assert!(locked.contains(&"AsyncDataSrc::close 4".to_string()));
830        assert!(locked.contains(&"AsyncDataSrc::drop 4".to_string()));
831        assert!(locked.contains(&"SyncDataSrc::close 3".to_string()));
832        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
833        assert!(locked.contains(&"AsyncDataSrc::drop 2".to_string()));
834        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
835    }
836
837    #[tokio::test]
838    async fn test_of_setup_async_and_ok() {
839        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
840
841        {
842            let mut manager = DataSrcManager::new(true);
843
844            let ds1 = SyncDataSrc::new(1, logger.clone(), false);
845            manager.add("foo", ds1);
846
847            let ds2 = SyncDataSrc::new(2, logger.clone(), false);
848            manager.add("bar", ds2);
849
850            assert!(manager.local);
851            assert_eq!(manager.vec_unready.len(), 2);
852            assert_eq!(manager.vec_ready.len(), 0);
853
854            let mut vec = Vec::new();
855            manager.setup_async(&mut vec).await;
856
857            assert!(manager.local);
858            assert_eq!(manager.vec_unready.len(), 0);
859            assert_eq!(manager.vec_ready.len(), 2);
860        }
861
862        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
863
864        let locked = logger.lock().await;
865        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
866        assert!(locked.contains(&"SyncDataSrc::new 2".to_string()));
867        assert!(locked.contains(&"SyncDataSrc::setup 1".to_string()));
868        assert!(locked.contains(&"SyncDataSrc::setup 2".to_string()));
869        assert!(locked.contains(&"SyncDataSrc::close 2".to_string()));
870        assert!(locked.contains(&"SyncDataSrc::drop 2".to_string()));
871        assert!(locked.contains(&"SyncDataSrc::close 1".to_string()));
872        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
873    }
874
875    #[tokio::test]
876    async fn test_of_setup_but_error() {
877        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
878
879        {
880            let mut manager = DataSrcManager::new(true);
881
882            let ds1 = SyncDataSrc::new(1, logger.clone(), false);
883            manager.add("foo", ds1);
884
885            let ds2 = SyncDataSrc::new(2, logger.clone(), true);
886            manager.add("bar", ds2);
887
888            let ds3 = SyncDataSrc::new(3, logger.clone(), true);
889            manager.add("bar", ds3);
890
891            assert!(manager.local);
892            assert_eq!(manager.vec_unready.len(), 3);
893            assert_eq!(manager.vec_ready.len(), 0);
894
895            let mut vec = Vec::new();
896            manager.setup_async(&mut vec).await;
897
898            assert!(manager.local);
899            assert_eq!(manager.vec_unready.len(), 3);
900            assert_eq!(manager.vec_ready.len(), 0);
901        }
902
903        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
904
905        let locked = logger.lock().await;
906        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
907        assert!(locked.contains(&"SyncDataSrc::new 2".to_string()));
908        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
909        assert!(locked.contains(&"SyncDataSrc::setup 1".to_string()));
910        assert!(locked.contains(&"SyncDataSrc::setup 2 failed".to_string()));
911        assert!(locked.contains(&"SyncDataSrc::close 1".to_string()));
912        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
913        assert!(locked.contains(&"SyncDataSrc::drop 2".to_string()));
914        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
915    }
916
917    #[tokio::test]
918    async fn test_of_setup_with_order_and_ok() {
919        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
920
921        {
922            let mut manager = DataSrcManager::new(true);
923
924            let ds1 = SyncDataSrc::new(1, logger.clone(), false);
925            manager.add("foo", ds1);
926
927            let ds2 = SyncDataSrc::new(2, logger.clone(), false);
928            manager.add("bar", ds2);
929
930            let ds3 = SyncDataSrc::new(3, logger.clone(), false);
931            manager.add("baz", ds3);
932
933            assert!(manager.local);
934            assert_eq!(manager.vec_unready.len(), 3);
935            assert_eq!(manager.vec_ready.len(), 0);
936
937            let mut errors = Vec::new();
938            manager
939                .setup_with_order_async(&["baz", "foo"], &mut errors)
940                .await;
941
942            assert!(manager.local);
943            assert_eq!(manager.vec_unready.len(), 0);
944            assert_eq!(manager.vec_ready.len(), 3);
945
946            assert_eq!(errors.len(), 0);
947        }
948
949        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
950
951        let locked = logger.lock().await;
952        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
953        assert!(locked.contains(&"SyncDataSrc::new 2".to_string()));
954        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
955        assert!(locked.contains(&"SyncDataSrc::setup 3".to_string()));
956        assert!(locked.contains(&"SyncDataSrc::setup 1".to_string()));
957        assert!(locked.contains(&"SyncDataSrc::setup 2".to_string()));
958        assert!(locked.contains(&"SyncDataSrc::close 2".to_string()));
959        assert!(locked.contains(&"SyncDataSrc::drop 2".to_string()));
960        assert!(locked.contains(&"SyncDataSrc::close 1".to_string()));
961        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
962        assert!(locked.contains(&"SyncDataSrc::close 3".to_string()));
963        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
964    }
965
966    #[tokio::test]
967    async fn test_of_setup_with_order_but_fail() {
968        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
969
970        {
971            let mut manager = DataSrcManager::new(true);
972
973            let ds1 = SyncDataSrc::new(1, logger.clone(), true);
974            manager.add("foo", ds1);
975
976            let ds2 = SyncDataSrc::new(2, logger.clone(), true);
977            manager.add("bar", ds2);
978
979            let ds3 = SyncDataSrc::new(3, logger.clone(), false);
980            manager.add("baz", ds3);
981
982            let ds4 = SyncDataSrc::new(4, logger.clone(), false);
983            manager.add("qux", ds4);
984
985            assert!(manager.local);
986            assert_eq!(manager.vec_unready.len(), 4);
987            assert_eq!(manager.vec_ready.len(), 0);
988
989            let mut errors = Vec::new();
990            manager
991                .setup_with_order_async(&["qux", "baz", "foo"], &mut errors)
992                .await;
993
994            assert!(manager.local);
995            assert_eq!(manager.vec_unready.len(), 4);
996            assert_eq!(manager.vec_ready.len(), 0);
997
998            assert_eq!(errors.len(), 1);
999            assert_eq!(errors[0].index, 0);
1000            assert_eq!(errors[0].name, "foo".into());
1001            #[cfg(unix)]
1002            assert_eq!(format!("{:?}", errors[0].err), "errs::Err { reason = alloc::string::String \"XXX\", file = src/tokio/data_src/mod.rs, line = 483 }");
1003            #[cfg(windows)]
1004            assert_eq!(format!("{:?}", errors[0].err), "errs::Err { reason = alloc::string::String \"XXX\", file = src\\tokio\\data_src\\mod.rs, line = 483 }");
1005        }
1006
1007        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
1008
1009        let locked = logger.lock().await;
1010        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
1011        assert!(locked.contains(&"SyncDataSrc::new 2".to_string()));
1012        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
1013        assert!(locked.contains(&"SyncDataSrc::new 4".to_string()));
1014        assert!(locked.contains(&"SyncDataSrc::setup 4".to_string()));
1015        assert!(locked.contains(&"SyncDataSrc::setup 3".to_string()));
1016        assert!(locked.contains(&"SyncDataSrc::setup 1 failed".to_string()));
1017        assert!(locked.contains(&"SyncDataSrc::close 4".to_string()));
1018        assert!(locked.contains(&"SyncDataSrc::close 3".to_string()));
1019        assert!(locked.contains(&"SyncDataSrc::drop 2".to_string()));
1020        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
1021        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
1022        assert!(locked.contains(&"SyncDataSrc::drop 4".to_string()));
1023    }
1024
1025    #[tokio::test]
1026    async fn test_of_setup_with_order_containing_duplicated_name_and_ok() {
1027        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1028
1029        {
1030            let mut manager = DataSrcManager::new(true);
1031
1032            let ds1 = SyncDataSrc::new(1, logger.clone(), false);
1033            manager.add("foo", ds1);
1034
1035            let ds2 = SyncDataSrc::new(2, logger.clone(), false);
1036            manager.add("bar", ds2);
1037
1038            let ds3 = SyncDataSrc::new(3, logger.clone(), false);
1039            manager.add("baz", ds3);
1040
1041            assert!(manager.local);
1042            assert_eq!(manager.vec_unready.len(), 3);
1043            assert_eq!(manager.vec_ready.len(), 0);
1044
1045            let mut vec = Vec::new();
1046            manager
1047                .setup_with_order_async(&["baz", "baz", "foo"], &mut vec)
1048                .await;
1049
1050            assert!(manager.local);
1051            assert_eq!(manager.vec_unready.len(), 0);
1052            assert_eq!(manager.vec_ready.len(), 3);
1053        }
1054
1055        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
1056
1057        let locked = logger.lock().await;
1058        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
1059        assert!(locked.contains(&"SyncDataSrc::new 2".to_string()));
1060        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
1061        assert!(locked.contains(&"SyncDataSrc::setup 3".to_string()));
1062        assert!(locked.contains(&"SyncDataSrc::setup 1".to_string()));
1063        assert!(locked.contains(&"SyncDataSrc::setup 2".to_string()));
1064        assert!(locked.contains(&"SyncDataSrc::close 2".to_string()));
1065        assert!(locked.contains(&"SyncDataSrc::drop 2".to_string()));
1066        assert!(locked.contains(&"SyncDataSrc::close 1".to_string()));
1067        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
1068        assert!(locked.contains(&"SyncDataSrc::close 3".to_string()));
1069        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
1070    }
1071
1072    #[tokio::test]
1073    async fn test_of_setup_with_order_containing_duplicated_name_and_ok_2() {
1074        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1075
1076        {
1077            let mut manager = DataSrcManager::new(true);
1078
1079            let ds1 = SyncDataSrc::new(1, logger.clone(), false);
1080            manager.add("foo", ds1);
1081
1082            let ds2 = SyncDataSrc::new(2, logger.clone(), false);
1083            manager.add("bar", ds2);
1084
1085            let ds3 = SyncDataSrc::new(3, logger.clone(), false);
1086            manager.add("baz", ds3);
1087
1088            let ds4 = SyncDataSrc::new(4, logger.clone(), false);
1089            manager.add("qux", ds4);
1090
1091            assert!(manager.local);
1092            assert_eq!(manager.vec_unready.len(), 4);
1093            assert_eq!(manager.vec_ready.len(), 0);
1094
1095            let mut vec = Vec::new();
1096            manager
1097                .setup_with_order_async(&["baz", "foo", "baz", "qux"], &mut vec)
1098                .await;
1099
1100            assert!(manager.local);
1101            assert_eq!(manager.vec_unready.len(), 0);
1102            assert_eq!(manager.vec_ready.len(), 4);
1103        }
1104
1105        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
1106
1107        let locked = logger.lock().await;
1108        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
1109        assert!(locked.contains(&"SyncDataSrc::new 2".to_string()));
1110        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
1111        assert!(locked.contains(&"SyncDataSrc::new 4".to_string()));
1112        assert!(locked.contains(&"SyncDataSrc::setup 3".to_string()));
1113        assert!(locked.contains(&"SyncDataSrc::setup 1".to_string()));
1114        assert!(locked.contains(&"SyncDataSrc::setup 4".to_string()));
1115        assert!(locked.contains(&"SyncDataSrc::setup 2".to_string()));
1116        assert!(locked.contains(&"SyncDataSrc::close 2".to_string()));
1117        assert!(locked.contains(&"SyncDataSrc::drop 2".to_string()));
1118        assert!(locked.contains(&"SyncDataSrc::close 4".to_string()));
1119        assert!(locked.contains(&"SyncDataSrc::drop 4".to_string()));
1120        assert!(locked.contains(&"SyncDataSrc::close 1".to_string()));
1121        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
1122        assert!(locked.contains(&"SyncDataSrc::close 3".to_string()));
1123        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
1124    }
1125
1126    #[tokio::test]
1127    async fn test_of_setup_with_order_buf_one_of_names_is_not_used() {
1128        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1129
1130        {
1131            let mut manager = DataSrcManager::new(true);
1132
1133            let ds1 = SyncDataSrc::new(1, logger.clone(), false);
1134            manager.add("foo", ds1);
1135
1136            let ds2 = SyncDataSrc::new(2, logger.clone(), false);
1137            manager.add("bar", ds2);
1138
1139            let ds3 = SyncDataSrc::new(3, logger.clone(), false);
1140            manager.add("baz", ds3);
1141
1142            assert!(manager.local);
1143            assert_eq!(manager.vec_unready.len(), 3);
1144            assert_eq!(manager.vec_ready.len(), 0);
1145
1146            let mut vec = Vec::new();
1147            manager
1148                .setup_with_order_async(&["baz", "foo", "xxx"], &mut vec)
1149                .await;
1150
1151            assert!(manager.local);
1152            assert_eq!(manager.vec_unready.len(), 0);
1153            assert_eq!(manager.vec_ready.len(), 3);
1154        }
1155
1156        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
1157
1158        let locked = logger.lock().await;
1159        assert!(locked.contains(&"SyncDataSrc::new 1".to_string()));
1160        assert!(locked.contains(&"SyncDataSrc::new 2".to_string()));
1161        assert!(locked.contains(&"SyncDataSrc::new 3".to_string()));
1162        assert!(locked.contains(&"SyncDataSrc::setup 3".to_string()));
1163        assert!(locked.contains(&"SyncDataSrc::setup 1".to_string()));
1164        assert!(locked.contains(&"SyncDataSrc::setup 2".to_string()));
1165        assert!(locked.contains(&"SyncDataSrc::close 2".to_string()));
1166        assert!(locked.contains(&"SyncDataSrc::drop 2".to_string()));
1167        assert!(locked.contains(&"SyncDataSrc::close 1".to_string()));
1168        assert!(locked.contains(&"SyncDataSrc::drop 1".to_string()));
1169        assert!(locked.contains(&"SyncDataSrc::close 3".to_string()));
1170        assert!(locked.contains(&"SyncDataSrc::drop 3".to_string()));
1171    }
1172
1173    #[tokio::test]
1174    async fn test_of_copy_ds_ready_to_map() {
1175        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1176        let mut errors = Vec::new();
1177
1178        let mut index_map = HashMap::<Arc<str>, (bool, usize)>::new();
1179
1180        let manager = DataSrcManager::new(true);
1181        manager.copy_ds_ready_to_map(&mut index_map);
1182        assert!(index_map.is_empty());
1183
1184        let mut manager = DataSrcManager::new(true);
1185        let ds1 = SyncDataSrc::new(1, logger.clone(), false);
1186        manager.add("foo", ds1);
1187        manager.setup_async(&mut errors).await;
1188        assert!(errors.is_empty());
1189        manager.copy_ds_ready_to_map(&mut index_map);
1190        assert_eq!(index_map.len(), 1);
1191        assert_eq!(index_map.get("foo").unwrap(), &(true, 0));
1192
1193        let mut manager = DataSrcManager::new(false);
1194        let ds2 = AsyncDataSrc::new(2, logger.clone(), false, 0);
1195        let ds3 = SyncDataSrc::new(3, logger.clone(), false);
1196        manager.add("bar", ds2);
1197        manager.add("baz", ds3);
1198        manager.setup_async(&mut errors).await;
1199        assert!(errors.is_empty());
1200        manager.copy_ds_ready_to_map(&mut index_map);
1201        assert_eq!(index_map.len(), 3);
1202        assert_eq!(index_map.get("foo").unwrap(), &(true, 0));
1203        assert_eq!(index_map.get("bar").unwrap(), &(false, 0));
1204        assert_eq!(index_map.get("baz").unwrap(), &(false, 1));
1205    }
1206
1207    #[tokio::test]
1208    async fn test_of_create_data_conn_and_ok() {
1209        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1210        let mut errors = Vec::new();
1211
1212        let mut manager = DataSrcManager::new(true);
1213        let ds1 = SyncDataSrc::new(1, logger.clone(), false);
1214        manager.add("foo", ds1);
1215        manager.setup_async(&mut errors).await;
1216
1217        if let Ok(boxed) = manager
1218            .create_data_conn_async::<SyncDataConn>(0, "foo")
1219            .await
1220        {
1221            assert_eq!(boxed.name.clone(), "foo".into());
1222        } else {
1223            panic!();
1224        }
1225    }
1226
1227    #[tokio::test]
1228    async fn test_of_create_data_conn_but_not_found() {
1229        let mut errors = Vec::new();
1230
1231        let mut manager = DataSrcManager::new(true);
1232        manager.setup_async(&mut errors).await;
1233
1234        if let Err(err) = manager
1235            .create_data_conn_async::<SyncDataConn>(0, "foo")
1236            .await
1237        {
1238            match err.reason::<DataSrcError>() {
1239                Ok(DataSrcError::NotFoundDataSrcToCreateDataConn {
1240                    name,
1241                    data_conn_type,
1242                }) => {
1243                    assert_eq!(*name, "foo".into());
1244                    assert_eq!(
1245                        *data_conn_type,
1246                        "sabi::tokio::data_src::tests_of_data_src::SyncDataConn"
1247                    );
1248                }
1249                _ => panic!(),
1250            }
1251        } else {
1252            panic!();
1253        }
1254    }
1255
1256    #[tokio::test]
1257    async fn test_of_create_data_conn_but_fail_to_cast() {
1258        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1259        let mut errors = Vec::new();
1260
1261        let mut manager = DataSrcManager::new(true);
1262        let ds1 = SyncDataSrc::new(1, logger.clone(), false);
1263        manager.add("foo", ds1);
1264        manager.setup_async(&mut errors).await;
1265
1266        if let Err(err) = manager
1267            .create_data_conn_async::<AsyncDataConn>(0, "foo")
1268            .await
1269        {
1270            match err.reason::<DataSrcError>() {
1271                Ok(DataSrcError::FailToCastDataConn { name, target_type }) => {
1272                    assert_eq!(*name, "foo".into());
1273                    assert_eq!(
1274                        *target_type,
1275                        "sabi::tokio::data_src::tests_of_data_src::AsyncDataConn"
1276                    );
1277                }
1278                _ => panic!(),
1279            }
1280        } else {
1281            panic!();
1282        }
1283    }
1284
1285    #[tokio::test]
1286    async fn test_of_create_data_conn_but_fail_to_create() {
1287        let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1288        let mut errors = Vec::new();
1289
1290        let mut manager = DataSrcManager::new(true);
1291        let ds1 = SyncDataSrc::new_for_fail_to_create_data_conn(1, logger.clone());
1292        manager.add("foo", ds1);
1293        manager.setup_async(&mut errors).await;
1294
1295        if let Err(err) = manager
1296            .create_data_conn_async::<SyncDataConn>(0, "foo")
1297            .await
1298        {
1299            match err.reason::<DataSrcError>() {
1300                Ok(DataSrcError::FailToCreateDataConn {
1301                    name,
1302                    data_conn_type,
1303                }) => {
1304                    assert_eq!(*name, "foo".into());
1305                    assert_eq!(
1306                        *data_conn_type,
1307                        "sabi::tokio::data_src::tests_of_data_src::SyncDataConn"
1308                    );
1309                }
1310                _ => panic!(),
1311            }
1312        } else {
1313            panic!();
1314        }
1315    }
1316}