1mod 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#[derive(Debug)]
27pub enum DataSrcError {
28 FailToRegisterGlobalDataSrc {
31 name: Arc<str>,
33 },
34
35 FailToSetupGlobalDataSrcs {
37 errors: Vec<ErrEntry>,
39 },
40
41 DuringSetupGlobalDataSrcs,
43
44 AlreadySetupGlobalDataSrcs,
46
47 FailToCastDataConn {
49 name: Arc<str>,
51
52 target_type: &'static str,
54 },
55
56 FailToCreateDataConn {
58 name: Arc<str>,
60
61 data_conn_type: &'static str,
63 },
64
65 NotFoundDataSrcToCreateDataConn {
67 name: Arc<str>,
69
70 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 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 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 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}