1use super::data_src::{copy_global_data_srcs_to_map, create_data_conn_from_global_data_src_async};
6use super::{
7 DataConn, DataConnContainer, DataConnManager, DataHub, DataSrc, DataSrcManager, ErrEntry,
8 SendSyncNonNull,
9};
10
11use std::collections::HashMap;
12use std::future::Future;
13use std::pin::Pin;
14use std::sync::Arc;
15use std::{any, ptr};
16
17#[derive(Debug)]
19pub enum DataHubError {
20 FailToSetupLocalDataSrcs {
22 errors: Vec<ErrEntry>,
24 },
25
26 NoDataSrcToCreateDataConn {
29 name: Arc<str>,
31
32 data_conn_type: &'static str,
34 },
35}
36
37impl DataHub {
38 #[allow(clippy::new_without_default)]
43 pub fn new() -> Self {
44 let mut data_src_map = HashMap::new();
45 copy_global_data_srcs_to_map(&mut data_src_map);
46
47 Self {
48 local_data_src_manager: DataSrcManager::new(true),
49 data_src_map,
50 data_conn_manager: DataConnManager::new(),
51 fixed: false,
52 }
53 }
54
55 pub fn with_commit_order(names: &[&str]) -> Self {
66 let mut data_src_map = HashMap::new();
67 copy_global_data_srcs_to_map(&mut data_src_map);
68
69 Self {
70 local_data_src_manager: DataSrcManager::new(true),
71 data_src_map,
72 data_conn_manager: DataConnManager::with_commit_order(names),
73 fixed: false,
74 }
75 }
76
77 pub fn uses<S, C>(&mut self, name: impl Into<Arc<str>>, ds: S)
93 where
94 S: DataSrc<C> + 'static,
95 C: DataConn + 'static,
96 {
97 if self.fixed {
98 return;
99 }
100 self.local_data_src_manager.add(name, ds);
101 }
102
103 pub fn disuses(&mut self, name: impl AsRef<str>) {
112 if self.fixed {
113 return;
114 }
115 self.data_src_map.remove(name.as_ref());
116 self.local_data_src_manager.remove(name);
117 }
118
119 #[inline]
120 async fn begin_async(&mut self) -> errs::Result<()> {
121 self.fixed = true;
122
123 let mut errors = Vec::new();
124
125 self.local_data_src_manager.setup_async(&mut errors).await;
126 if errors.is_empty() {
127 self.local_data_src_manager
128 .copy_ds_ready_to_map(&mut self.data_src_map);
129 Ok(())
130 } else {
131 Err(errs::Err::new(DataHubError::FailToSetupLocalDataSrcs {
132 errors,
133 }))
134 }
135 }
136
137 #[inline]
138 fn end(&mut self) {
139 self.data_conn_manager.close();
140 self.fixed = false;
141 }
142
143 #[allow(clippy::doc_overindented_list_items)]
164 pub async fn run_async<F>(&mut self, mut logic_fn: F) -> errs::Result<()>
165 where
166 for<'a> F:
167 FnMut(&'a mut DataHub) -> Pin<Box<dyn Future<Output = errs::Result<()>> + Send + 'a>>,
168 {
169 let mut r = self.begin_async().await;
170 if r.is_ok() {
171 r = logic_fn(self).await;
172 }
173 self.end();
174 r
175 }
176
177 #[allow(clippy::doc_overindented_list_items)]
202 pub async fn txn_async<F>(&mut self, mut logic_fn: F) -> errs::Result<()>
203 where
204 for<'a> F:
205 FnMut(&'a mut DataHub) -> Pin<Box<dyn Future<Output = errs::Result<()>> + Send + 'a>>,
206 {
207 let mut r = self.begin_async().await;
208 if r.is_ok() {
209 r = logic_fn(self).await;
210 }
211
212 let mut reports = self.data_conn_manager.new_failure_reports();
213
214 if r.is_ok() {
215 r = self.data_conn_manager.commit_async(&mut reports).await;
216 }
217 if r.is_err() {
218 self.data_conn_manager.rollback_async(reports).await;
219 }
220
221 self.end();
222 r
223 }
224
225 pub async fn get_data_conn_async<C>(&mut self, name: &str) -> errs::Result<&mut C>
246 where
247 C: DataConn + 'static,
248 {
249 if let Some(nnptr) = self.data_conn_manager.find_by_name(name) {
250 let typed_nnptr = DataConnManager::to_typed_ptr::<C>(&nnptr)?;
251 return Ok(unsafe { &mut (*typed_nnptr).data_conn });
252 }
253
254 if let Some((local, index)) = self.data_src_map.get(name) {
255 let boxed = if *local {
256 self.local_data_src_manager
257 .create_data_conn_async::<C>(*index, name)
258 .await?
259 } else {
260 create_data_conn_from_global_data_src_async::<C>(*index, name).await?
261 };
262
263 let ptr = Box::into_raw(boxed);
264 if let Some(nnptr) = ptr::NonNull::new(ptr) {
265 let ssnnptr = SendSyncNonNull::new(nnptr);
266 self.data_conn_manager.add(ssnnptr);
267
268 let typed_ptr = ptr.cast::<DataConnContainer<C>>();
269 return Ok(unsafe { &mut (*typed_ptr).data_conn });
270 } else {
271 }
273 }
274
275 Err(errs::Err::new(DataHubError::NoDataSrcToCreateDataConn {
276 name: name.into(),
277 data_conn_type: any::type_name::<C>(),
278 }))
279 }
280}
281
282#[macro_export]
283#[doc(hidden)]
284macro_rules! _logic {
285 ($f:expr) => {
286 |data| {
287 let fut: std::pin::Pin<Box<dyn std::future::Future<Output = errs::Result<()>> + Send>> =
288 Box::pin(async move { $f(data).await });
289 fut
290 }
291 };
292}
293
294#[cfg_attr(coverage_nightly, coverage(off))]
295#[cfg(test)]
296mod tests_of_data_hub {
297 use super::*;
298 use crate::tokio::{AsyncGroup, DataConnError, DataSrcError};
299 use crate::TxnFailureReport;
300 use std::sync::Mutex;
301
302 #[derive(Clone, Copy, PartialEq)]
303 enum Failure {
304 None,
305 FailToPreCommit,
306 FailToCommit,
307 FailToPostCommit,
308 FailToRollback,
309 FailToSetup,
310 FailToCreateDataConn,
311 }
312
313 struct MyDataConn {
314 id: i8,
315 failure: Failure,
316 committed: bool,
317 logger: Arc<Mutex<Vec<String>>>,
318 }
319 impl MyDataConn {
320 fn new(id: i8, logger: Arc<Mutex<Vec<String>>>, failure: Failure) -> Self {
321 logger
322 .lock()
323 .unwrap()
324 .push(format!("MyDataConn::new {}", id));
325 Self {
326 id,
327 failure,
328 committed: false,
329 logger,
330 }
331 }
332 }
333 impl Drop for MyDataConn {
334 fn drop(&mut self) {
335 self.logger
336 .lock()
337 .unwrap()
338 .push(format!("MyDataConn::drop {}", self.id));
339 }
340 }
341 impl DataConn for MyDataConn {
342 async fn pre_commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
343 if self.failure == Failure::FailToPreCommit {
344 self.logger
345 .lock()
346 .unwrap()
347 .push(format!("MyDataConn::pre_commit_async {} failed", self.id));
348 Err(errs::Err::new("pre commit error"))
349 } else {
350 self.logger
351 .lock()
352 .unwrap()
353 .push(format!("MyDataConn::pre_commit_async {}", self.id));
354 Ok(())
355 }
356 }
357 async fn commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
358 if self.failure == Failure::FailToCommit {
359 self.logger
360 .lock()
361 .unwrap()
362 .push(format!("MyDataConn::commit_async {} failed", self.id));
363 Err(errs::Err::new("commit error"))
364 } else {
365 self.logger
366 .lock()
367 .unwrap()
368 .push(format!("MyDataConn::commit_async {}", self.id));
369 self.committed = true;
370 Ok(())
371 }
372 }
373 async fn post_commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
374 if self.failure == Failure::FailToPostCommit {
375 self.logger
376 .lock()
377 .unwrap()
378 .push(format!("MyDataConn::post_commit_async {} failed", self.id));
379 Err(errs::Err::new("post commit error"))
380 } else {
381 self.logger
382 .lock()
383 .unwrap()
384 .push(format!("MyDataConn::post_commit_async {}", self.id));
385 Ok(())
386 }
387 }
388 fn is_committed(&self) -> bool {
389 self.committed
390 }
391 async fn rollback_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
392 if self.failure == Failure::FailToRollback {
393 self.logger
394 .lock()
395 .unwrap()
396 .push(format!("MyDataConn::rollback_async {} failed", self.id));
397 Ok(())
398 } else {
399 self.logger
400 .lock()
401 .unwrap()
402 .push(format!("MyDataConn::rollback_async {}", self.id));
403 Err(errs::Err::new("rollback error"))
404 }
405 }
406 async fn on_txn_failure_async(
407 &mut self,
408 _ag: &mut AsyncGroup,
409 _reports: Arc<[TxnFailureReport]>,
410 ) {
411 self.logger
412 .lock()
413 .unwrap()
414 .push(format!("MyDataConn::on_txn_failure_async {}", self.id));
415 }
416 fn close(&mut self) {
417 self.logger
418 .lock()
419 .unwrap()
420 .push(format!("MyDataConn::close {}", self.id));
421 }
422 }
423
424 struct MyDataSrc {
425 id: i8,
426 failure: Failure,
427 logger: Arc<Mutex<Vec<String>>>,
428 }
429 impl MyDataSrc {
430 fn new(id: i8, logger: Arc<Mutex<Vec<String>>>, failure: Failure) -> Self {
431 logger
432 .lock()
433 .unwrap()
434 .push(format!("MyDataSrc::new {}", id));
435 Self {
436 id,
437 failure,
438 logger,
439 }
440 }
441 }
442 impl Drop for MyDataSrc {
443 fn drop(&mut self) {
444 self.logger
445 .lock()
446 .unwrap()
447 .push(format!("MyDataSrc::drop {}", self.id));
448 }
449 }
450 impl DataSrc<MyDataConn> for MyDataSrc {
451 async fn setup_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
452 if self.failure == Failure::FailToSetup {
453 self.logger
454 .lock()
455 .unwrap()
456 .push(format!("MyDataSrc::setup_async {} failed", self.id));
457 Err(errs::Err::new("setup error".to_string()))
458 } else {
459 self.logger
460 .lock()
461 .unwrap()
462 .push(format!("MyDataSrc::setup_async {}", self.id));
463 Ok(())
464 }
465 }
466 fn close(&mut self) {
467 self.logger
468 .lock()
469 .unwrap()
470 .push(format!("MyDataSrc::close {}", self.id));
471 }
472 async fn create_data_conn_async(&mut self) -> errs::Result<Box<MyDataConn>> {
473 if self.failure == Failure::FailToCreateDataConn {
474 self.logger.lock().unwrap().push(format!(
475 "MyDataSrc::create_data_conn_async {} failed",
476 self.id
477 ));
478 return Err(errs::Err::new("eeee".to_string()));
479 }
480 {
481 self.logger
482 .lock()
483 .unwrap()
484 .push(format!("MyDataSrc::create_data_conn_async {}", self.id));
485 }
486 let conn = MyDataConn::new(self.id, self.logger.clone(), self.failure);
487 Ok(Box::new(conn))
488 }
489 }
490
491 struct AnotherDataConn {}
492 impl DataConn for AnotherDataConn {
493 async fn pre_commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
494 Ok(())
495 }
496 async fn commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
497 Ok(())
498 }
499 async fn post_commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
500 Ok(())
501 }
502 fn is_committed(&self) -> bool {
503 false
504 }
505 async fn rollback_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
506 Ok(())
507 }
508 async fn on_txn_failure_async(
509 &mut self,
510 _ag: &mut AsyncGroup,
511 _reports: Arc<[TxnFailureReport]>,
512 ) {
513 }
514 fn close(&mut self) {}
515 }
516
517 #[test]
518 fn test_new() {
519 let hub = DataHub::new();
520 assert!(hub.local_data_src_manager.vec_unready.is_empty());
521 assert!(hub.local_data_src_manager.vec_ready.is_empty());
522 assert!(hub.local_data_src_manager.local);
523 assert!(hub.data_src_map.is_empty());
524 assert!(hub.data_conn_manager.vec.is_empty());
525 assert!(hub.data_conn_manager.index_map.is_empty());
526 assert!(!hub.fixed);
527 }
528
529 #[test]
530 fn test_with_commit_order() {
531 let hub = DataHub::with_commit_order(&["bar", "qux", "foo"]);
532 assert!(hub.local_data_src_manager.vec_unready.is_empty());
533 assert!(hub.local_data_src_manager.vec_ready.is_empty());
534 assert!(hub.local_data_src_manager.local);
535 assert!(hub.data_src_map.is_empty());
536 assert_eq!(hub.data_conn_manager.vec.len(), 3);
537 assert_eq!(hub.data_conn_manager.index_map.len(), 3);
538 assert!(!hub.fixed);
539 }
540
541 #[tokio::test]
542 async fn test_uses_and_ok() {
543 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
544
545 let mut hub = DataHub::new();
546 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
547 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
548
549 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 2);
550 assert!(hub.local_data_src_manager.vec_ready.is_empty());
551 assert!(hub.local_data_src_manager.local);
552 assert!(hub.data_src_map.is_empty());
553 assert_eq!(hub.data_conn_manager.vec.len(), 0);
554 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
555 assert!(!hub.fixed);
556
557 assert!(hub.begin_async().await.is_ok());
558
559 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 0);
560 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 2);
561 assert!(hub.local_data_src_manager.local);
562 assert_eq!(hub.data_src_map.len(), 2);
563 assert_eq!(hub.data_conn_manager.vec.len(), 0);
564 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
565 assert!(hub.fixed);
566 }
567
568 #[tokio::test]
569 async fn test_uses_but_already_fixed() {
570 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
571
572 let mut hub = DataHub::new();
573 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
574
575 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 1);
576 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 0);
577 assert!(hub.local_data_src_manager.local);
578 assert_eq!(hub.data_src_map.len(), 0);
579 assert_eq!(hub.data_conn_manager.vec.len(), 0);
580 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
581 assert!(!hub.fixed);
582
583 assert!(hub.begin_async().await.is_ok());
584
585 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 0);
586 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 1);
587 assert!(hub.local_data_src_manager.local);
588 assert_eq!(hub.data_src_map.len(), 1);
589 assert_eq!(hub.data_conn_manager.vec.len(), 0);
590 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
591 assert!(hub.fixed);
592
593 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
594
595 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 0);
596 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 1);
597 assert!(hub.local_data_src_manager.local);
598 assert_eq!(hub.data_src_map.len(), 1);
599 assert_eq!(hub.data_conn_manager.vec.len(), 0);
600 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
601 assert!(hub.fixed);
602 }
603
604 #[test]
605 fn test_disuses_and_ok() {
606 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
607
608 let mut hub = DataHub::new();
609 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
610 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
611
612 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 2);
613 assert!(hub.local_data_src_manager.vec_ready.is_empty());
614 assert!(hub.local_data_src_manager.local);
615 assert!(hub.data_src_map.is_empty());
616 assert_eq!(hub.data_conn_manager.vec.len(), 0);
617 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
618 assert!(!hub.fixed);
619
620 hub.disuses("foo");
621
622 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 1);
623 assert!(hub.local_data_src_manager.vec_ready.is_empty());
624 assert!(hub.local_data_src_manager.local);
625 assert!(hub.data_src_map.is_empty());
626 assert_eq!(hub.data_conn_manager.vec.len(), 0);
627 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
628 assert!(!hub.fixed);
629
630 hub.disuses("bar");
631
632 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 0);
633 assert!(hub.local_data_src_manager.vec_ready.is_empty());
634 assert!(hub.local_data_src_manager.local);
635 assert!(hub.data_src_map.is_empty());
636 assert_eq!(hub.data_conn_manager.vec.len(), 0);
637 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
638 assert!(!hub.fixed);
639 }
640
641 #[tokio::test]
642 async fn test_disuses_and_fix() {
643 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
644
645 let mut hub = DataHub::new();
646 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
647 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
648
649 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 2);
650 assert!(hub.local_data_src_manager.vec_ready.is_empty());
651 assert!(hub.local_data_src_manager.local);
652 assert!(hub.data_src_map.is_empty());
653 assert_eq!(hub.data_conn_manager.vec.len(), 0);
654 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
655 assert!(!hub.fixed);
656
657 hub.disuses("foo");
658
659 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 1);
660 assert!(hub.local_data_src_manager.vec_ready.is_empty());
661 assert!(hub.local_data_src_manager.local);
662 assert!(hub.data_src_map.is_empty());
663 assert_eq!(hub.data_conn_manager.vec.len(), 0);
664 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
665 assert!(!hub.fixed);
666
667 hub.disuses("bar");
668
669 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 0);
670 assert!(hub.local_data_src_manager.vec_ready.is_empty());
671 assert!(hub.local_data_src_manager.local);
672 assert!(hub.data_src_map.is_empty());
673 assert_eq!(hub.data_conn_manager.vec.len(), 0);
674 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
675 assert!(!hub.fixed);
676
677 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
678 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
679
680 assert!(hub.begin_async().await.is_ok());
681
682 assert!(hub.local_data_src_manager.vec_unready.is_empty());
683 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 2);
684 assert!(hub.local_data_src_manager.local);
685 assert_eq!(hub.data_src_map.len(), 2);
686 assert_eq!(hub.data_conn_manager.vec.len(), 0);
687 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
688 assert!(hub.fixed);
689
690 hub.uses("baz", MyDataSrc::new(3, logger.clone(), Failure::None));
691
692 assert!(hub.local_data_src_manager.vec_unready.is_empty());
693 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 2);
694 assert!(hub.local_data_src_manager.local);
695 assert_eq!(hub.data_src_map.len(), 2);
696 assert_eq!(hub.data_conn_manager.vec.len(), 0);
697 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
698 assert!(hub.fixed);
699
700 hub.disuses("bar");
701
702 assert!(hub.local_data_src_manager.vec_unready.is_empty());
703 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 2);
704 assert!(hub.local_data_src_manager.local);
705 assert_eq!(hub.data_src_map.len(), 2);
706 assert_eq!(hub.data_conn_manager.vec.len(), 0);
707 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
708 assert!(hub.fixed);
709
710 hub.end();
711
712 assert!(hub.local_data_src_manager.vec_unready.is_empty());
713 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 2);
714 assert!(hub.local_data_src_manager.local);
715 assert_eq!(hub.data_src_map.len(), 2);
716 assert_eq!(hub.data_conn_manager.vec.len(), 0);
717 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
718 assert!(!hub.fixed);
719
720 hub.disuses("bar");
721
722 assert!(hub.local_data_src_manager.vec_unready.is_empty());
723 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 1);
724 assert!(hub.local_data_src_manager.local);
725 assert_eq!(hub.data_src_map.len(), 1);
726 assert_eq!(hub.data_conn_manager.vec.len(), 0);
727 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
728 assert!(!hub.fixed);
729
730 hub.disuses("foo");
731
732 assert!(hub.local_data_src_manager.vec_unready.is_empty());
733 assert!(hub.local_data_src_manager.vec_ready.is_empty());
734 assert!(hub.local_data_src_manager.local);
735 assert_eq!(hub.data_src_map.len(), 0);
736 assert_eq!(hub.data_conn_manager.vec.len(), 0);
737 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
738 assert!(!hub.fixed);
739 }
740
741 #[tokio::test]
742 async fn test_begin_if_empty() {
743 let mut hub = DataHub::new();
744 assert!(hub.begin_async().await.is_ok());
745
746 assert!(hub.local_data_src_manager.vec_unready.is_empty());
747 assert!(hub.local_data_src_manager.vec_ready.is_empty());
748 assert!(hub.local_data_src_manager.local);
749 assert_eq!(hub.data_src_map.len(), 0);
750 assert_eq!(hub.data_conn_manager.vec.len(), 0);
751 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
752 assert!(hub.fixed);
753
754 hub.end();
755
756 assert!(hub.local_data_src_manager.vec_unready.is_empty());
757 assert!(hub.local_data_src_manager.vec_ready.is_empty());
758 assert!(hub.local_data_src_manager.local);
759 assert_eq!(hub.data_src_map.len(), 0);
760 assert_eq!(hub.data_conn_manager.vec.len(), 0);
761 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
762 assert!(!hub.fixed);
763 }
764
765 #[tokio::test]
766 async fn test_begin_and_ok() {
767 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
768
769 {
770 let mut hub = DataHub::new();
771
772 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
773 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
774
775 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 2);
776 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 0);
777 assert_eq!(hub.local_data_src_manager.local, true);
778 assert_eq!(hub.data_src_map.len(), 0);
779 assert_eq!(hub.data_conn_manager.vec.len(), 0);
780 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
781 assert_eq!(hub.fixed, false);
782
783 assert_eq!(hub.begin_async().await.is_ok(), true);
784
785 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 0);
786 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 2);
787 assert_eq!(hub.local_data_src_manager.local, true);
788 assert_eq!(hub.data_src_map.len(), 2);
789 assert_eq!(hub.data_conn_manager.vec.len(), 0);
790 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
791 assert_eq!(hub.fixed, true);
792
793 hub.end();
794
795 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 0);
796 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 2);
797 assert_eq!(hub.local_data_src_manager.local, true);
798 assert_eq!(hub.data_src_map.len(), 2);
799 assert_eq!(hub.data_conn_manager.vec.len(), 0);
800 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
801 assert_eq!(hub.fixed, false);
802 }
803
804 assert_eq!(
805 *logger.lock().unwrap(),
806 &[
807 "MyDataSrc::new 1",
808 "MyDataSrc::new 2",
809 "MyDataSrc::setup_async 1",
810 "MyDataSrc::setup_async 2",
811 "MyDataSrc::close 2",
812 "MyDataSrc::drop 2",
813 "MyDataSrc::close 1",
814 "MyDataSrc::drop 1",
815 ]
816 );
817 }
818
819 #[tokio::test]
820 async fn test_begin_but_failed() {
821 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
822
823 {
824 let mut hub = DataHub::new();
825
826 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
827 hub.uses(
828 "bar",
829 MyDataSrc::new(2, logger.clone(), Failure::FailToSetup),
830 );
831 hub.uses("baz", MyDataSrc::new(3, logger.clone(), Failure::None));
832
833 assert_eq!(hub.local_data_src_manager.vec_unready.len(), 3);
834 assert_eq!(hub.local_data_src_manager.vec_ready.len(), 0);
835 assert_eq!(hub.local_data_src_manager.local, true);
836 assert_eq!(hub.data_src_map.len(), 0);
837 assert_eq!(hub.data_conn_manager.vec.len(), 0);
838 assert_eq!(hub.data_conn_manager.index_map.len(), 0);
839 assert_eq!(hub.fixed, false);
840
841 if let Err(err) = hub.begin_async().await {
842 match err.reason::<DataHubError>() {
843 Ok(DataHubError::FailToSetupLocalDataSrcs { errors }) => {
844 assert_eq!(errors.len(), 1);
845 assert_eq!(errors[0].index, 1);
846 assert_eq!(errors[0].name, "bar".into());
847 assert_eq!(errors[0].err.reason::<String>().unwrap(), "setup error");
848 }
849 _ => panic!(),
850 }
851 } else {
852 panic!();
853 }
854
855 hub.end();
856 }
857
858 assert_eq!(
859 *logger.lock().unwrap(),
860 &[
861 "MyDataSrc::new 1",
862 "MyDataSrc::new 2",
863 "MyDataSrc::new 3",
864 "MyDataSrc::setup_async 1",
865 "MyDataSrc::setup_async 2 failed",
866 "MyDataSrc::close 1",
867 "MyDataSrc::drop 3",
868 "MyDataSrc::drop 2",
869 "MyDataSrc::drop 1",
870 ]
871 );
872 }
873
874 #[tokio::test]
875 async fn test_run_and_ok() {
876 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
877 {
878 let mut hub = DataHub::new();
879
880 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
881 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
882
883 let logger_clone = logger.clone();
884 assert!(hub
885 .run_async(|_data| {
886 let logger_clone2 = logger_clone.clone();
887 Box::pin(async move {
888 logger_clone2
889 .lock()
890 .unwrap()
891 .push("execute logic".to_string());
892 Ok(())
893 })
894 })
895 .await
896 .is_ok());
897 }
898
899 assert_eq!(
900 *logger.lock().unwrap(),
901 &[
902 "MyDataSrc::new 1",
903 "MyDataSrc::new 2",
904 "MyDataSrc::setup_async 1",
905 "MyDataSrc::setup_async 2",
906 "execute logic",
907 "MyDataSrc::close 2",
908 "MyDataSrc::drop 2",
909 "MyDataSrc::close 1",
910 "MyDataSrc::drop 1",
911 ]
912 );
913 }
914
915 #[tokio::test]
916 async fn test_run_but_failed() {
917 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
918 {
919 let mut hub = DataHub::new();
920
921 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
922 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
923
924 let logger_clone = logger.clone();
925 if let Err(err) = hub
926 .run_async(|_data| {
927 let logger_clone2 = logger_clone.clone();
928 Box::pin(async move {
929 logger_clone2
930 .lock()
931 .unwrap()
932 .push("execute logic but fail".to_string());
933 Err(errs::Err::new("logic error".to_string()))
934 })
935 })
936 .await
937 {
938 match err.reason::<String>() {
939 Ok(s) => assert_eq!(s, "logic error"),
940 _ => panic!(),
941 }
942 } else {
943 panic!();
944 }
945 }
946
947 assert_eq!(
948 *logger.lock().unwrap(),
949 &[
950 "MyDataSrc::new 1",
951 "MyDataSrc::new 2",
952 "MyDataSrc::setup_async 1",
953 "MyDataSrc::setup_async 2",
954 "execute logic but fail",
955 "MyDataSrc::close 2",
956 "MyDataSrc::drop 2",
957 "MyDataSrc::close 1",
958 "MyDataSrc::drop 1",
959 ]
960 );
961 }
962
963 #[tokio::test]
964 async fn test_txn_and_no_data_access_and_ok() {
965 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
966 {
967 let mut hub = DataHub::new();
968
969 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
970 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
971
972 let logger_clone = logger.clone();
973 assert!(hub
974 .txn_async(|_data| {
975 let logger_clone2 = logger_clone.clone();
976 Box::pin(async move {
977 logger_clone2
978 .lock()
979 .unwrap()
980 .push("execute logic".to_string());
981 Ok(())
982 })
983 })
984 .await
985 .is_ok());
986 }
987
988 assert_eq!(
989 *logger.lock().unwrap(),
990 &[
991 "MyDataSrc::new 1",
992 "MyDataSrc::new 2",
993 "MyDataSrc::setup_async 1",
994 "MyDataSrc::setup_async 2",
995 "execute logic",
996 "MyDataSrc::close 2",
997 "MyDataSrc::drop 2",
998 "MyDataSrc::close 1",
999 "MyDataSrc::drop 1",
1000 ]
1001 );
1002 }
1003
1004 #[tokio::test]
1005 async fn test_txn_and_has_data_access_and_ok() {
1006 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1007 {
1008 let mut hub = DataHub::new();
1009
1010 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
1011 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
1012
1013 let logger_clone = logger.clone();
1014 hub.txn_async(move |data| {
1015 let logger_clone2 = logger_clone.clone();
1016 Box::pin(async move {
1017 logger_clone2
1018 .lock()
1019 .unwrap()
1020 .push("execute logic".to_string());
1021 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1022 let _conn2 = data.get_data_conn_async::<MyDataConn>("bar").await?;
1023 Ok(())
1024 })
1025 })
1026 .await
1027 .unwrap()
1028 }
1029
1030 assert_eq!(
1031 *logger.lock().unwrap(),
1032 &[
1033 "MyDataSrc::new 1",
1034 "MyDataSrc::new 2",
1035 "MyDataSrc::setup_async 1",
1036 "MyDataSrc::setup_async 2",
1037 "execute logic",
1038 "MyDataSrc::create_data_conn_async 1",
1039 "MyDataConn::new 1",
1040 "MyDataSrc::create_data_conn_async 2",
1041 "MyDataConn::new 2",
1042 "MyDataConn::pre_commit_async 1",
1043 "MyDataConn::pre_commit_async 2",
1044 "MyDataConn::commit_async 1",
1045 "MyDataConn::commit_async 2",
1046 "MyDataConn::post_commit_async 1",
1047 "MyDataConn::post_commit_async 2",
1048 "MyDataConn::close 2",
1049 "MyDataConn::drop 2",
1050 "MyDataConn::close 1",
1051 "MyDataConn::drop 1",
1052 "MyDataSrc::close 2",
1053 "MyDataSrc::drop 2",
1054 "MyDataSrc::close 1",
1055 "MyDataSrc::drop 1",
1056 ]
1057 );
1058 }
1059
1060 #[tokio::test]
1061 async fn test_txn_but_failed_to_run_logic() {
1062 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1063 {
1064 let mut hub = DataHub::new();
1065
1066 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
1067 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
1068
1069 let logger_clone = logger.clone();
1070 if let Err(e) = hub
1071 .txn_async(move |data| {
1072 let logger_clone2 = logger_clone.clone();
1073 Box::pin(async move {
1074 logger_clone2
1075 .lock()
1076 .unwrap()
1077 .push("execute logic".to_string());
1078 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1079 let _conn2 = data.get_data_conn_async::<MyDataConn>("bar").await?;
1080 Err(errs::Err::new("logic error"))
1081 })
1082 })
1083 .await
1084 {
1085 match e.reason::<&str>() {
1086 Ok(s) => assert_eq!(s, &"logic error"),
1087 _ => panic!(),
1088 }
1089 }
1090 }
1091
1092 assert_eq!(
1093 *logger.lock().unwrap(),
1094 &[
1095 "MyDataSrc::new 1",
1096 "MyDataSrc::new 2",
1097 "MyDataSrc::setup_async 1",
1098 "MyDataSrc::setup_async 2",
1099 "execute logic",
1100 "MyDataSrc::create_data_conn_async 1",
1101 "MyDataConn::new 1",
1102 "MyDataSrc::create_data_conn_async 2",
1103 "MyDataConn::new 2",
1104 "MyDataConn::rollback_async 1",
1105 "MyDataConn::rollback_async 2",
1106 "MyDataConn::on_txn_failure_async 1",
1107 "MyDataConn::on_txn_failure_async 2",
1108 "MyDataConn::close 2",
1109 "MyDataConn::drop 2",
1110 "MyDataConn::close 1",
1111 "MyDataConn::drop 1",
1112 "MyDataSrc::close 2",
1113 "MyDataSrc::drop 2",
1114 "MyDataSrc::close 1",
1115 "MyDataSrc::drop 1",
1116 ]
1117 );
1118 }
1119
1120 #[tokio::test]
1121 async fn test_txn_but_failed_to_pre_commit() {
1122 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1123 {
1124 let mut hub = DataHub::new();
1125
1126 hub.uses(
1127 "foo",
1128 MyDataSrc::new(1, logger.clone(), Failure::FailToPreCommit),
1129 );
1130 hub.uses(
1131 "bar",
1132 MyDataSrc::new(2, logger.clone(), Failure::FailToPreCommit),
1133 );
1134
1135 let logger_clone = logger.clone();
1136 if let Err(e) = hub
1137 .txn_async(move |data| {
1138 let logger_clone2 = logger_clone.clone();
1139 Box::pin(async move {
1140 logger_clone2
1141 .lock()
1142 .unwrap()
1143 .push("execute logic".to_string());
1144 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1145 let _conn2 = data.get_data_conn_async::<MyDataConn>("bar").await?;
1146 Ok(())
1147 })
1148 })
1149 .await
1150 {
1151 match e.reason::<DataConnError>() {
1152 Ok(DataConnError::FailToPreCommitDataConn { errors }) => {
1153 assert_eq!(errors.len(), 1);
1154 assert_eq!(errors[0].index, 0);
1155 assert_eq!(errors[0].name, "foo".into());
1156 assert_eq!(errors[0].err.reason::<&str>().unwrap(), &"pre commit error");
1157 }
1158 _ => panic!("{e:?}"),
1159 }
1160 }
1161 }
1162
1163 assert_eq!(
1164 *logger.lock().unwrap(),
1165 &[
1166 "MyDataSrc::new 1",
1167 "MyDataSrc::new 2",
1168 "MyDataSrc::setup_async 1",
1169 "MyDataSrc::setup_async 2",
1170 "execute logic",
1171 "MyDataSrc::create_data_conn_async 1",
1172 "MyDataConn::new 1",
1173 "MyDataSrc::create_data_conn_async 2",
1174 "MyDataConn::new 2",
1175 "MyDataConn::pre_commit_async 1 failed",
1176 "MyDataConn::rollback_async 1",
1177 "MyDataConn::rollback_async 2",
1178 "MyDataConn::on_txn_failure_async 1",
1179 "MyDataConn::on_txn_failure_async 2",
1180 "MyDataConn::close 2",
1181 "MyDataConn::drop 2",
1182 "MyDataConn::close 1",
1183 "MyDataConn::drop 1",
1184 "MyDataSrc::close 2",
1185 "MyDataSrc::drop 2",
1186 "MyDataSrc::close 1",
1187 "MyDataSrc::drop 1",
1188 ]
1189 );
1190 }
1191
1192 #[tokio::test]
1193 async fn test_txn_but_failed_to_commit() {
1194 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1195 {
1196 let mut hub = DataHub::new();
1197
1198 hub.uses(
1199 "foo",
1200 MyDataSrc::new(1, logger.clone(), Failure::FailToCommit),
1201 );
1202 hub.uses(
1203 "bar",
1204 MyDataSrc::new(2, logger.clone(), Failure::FailToCommit),
1205 );
1206
1207 let logger_clone = logger.clone();
1208 if let Err(e) = hub
1209 .txn_async(move |data| {
1210 let logger_clone2 = logger_clone.clone();
1211 Box::pin(async move {
1212 logger_clone2
1213 .lock()
1214 .unwrap()
1215 .push("execute logic".to_string());
1216 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1217 let _conn2 = data.get_data_conn_async::<MyDataConn>("bar").await?;
1218 Ok(())
1219 })
1220 })
1221 .await
1222 {
1223 match e.reason::<DataConnError>() {
1224 Ok(DataConnError::FailToCommitDataConn { errors }) => {
1225 assert_eq!(errors.len(), 1);
1226 assert_eq!(errors[0].index, 0);
1227 assert_eq!(errors[0].name, "foo".into());
1228 assert_eq!(errors[0].err.reason::<&str>().unwrap(), &"commit error");
1229 }
1230 _ => panic!("{e:?}"),
1231 }
1232 }
1233 }
1234
1235 assert_eq!(
1236 *logger.lock().unwrap(),
1237 &[
1238 "MyDataSrc::new 1",
1239 "MyDataSrc::new 2",
1240 "MyDataSrc::setup_async 1",
1241 "MyDataSrc::setup_async 2",
1242 "execute logic",
1243 "MyDataSrc::create_data_conn_async 1",
1244 "MyDataConn::new 1",
1245 "MyDataSrc::create_data_conn_async 2",
1246 "MyDataConn::new 2",
1247 "MyDataConn::pre_commit_async 1",
1248 "MyDataConn::pre_commit_async 2",
1249 "MyDataConn::commit_async 1 failed",
1250 "MyDataConn::rollback_async 1",
1251 "MyDataConn::rollback_async 2",
1252 "MyDataConn::on_txn_failure_async 1",
1253 "MyDataConn::on_txn_failure_async 2",
1254 "MyDataConn::close 2",
1255 "MyDataConn::drop 2",
1256 "MyDataConn::close 1",
1257 "MyDataConn::drop 1",
1258 "MyDataSrc::close 2",
1259 "MyDataSrc::drop 2",
1260 "MyDataSrc::close 1",
1261 "MyDataSrc::drop 1",
1262 ]
1263 );
1264 }
1265
1266 #[tokio::test]
1267 async fn test_txn_but_failed_to_post_commit() {
1268 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1269 {
1270 let mut hub = DataHub::new();
1271
1272 hub.uses(
1273 "foo",
1274 MyDataSrc::new(1, logger.clone(), Failure::FailToPostCommit),
1275 );
1276 hub.uses(
1277 "bar",
1278 MyDataSrc::new(2, logger.clone(), Failure::FailToPostCommit),
1279 );
1280
1281 let logger_clone = logger.clone();
1282 if let Err(e) = hub
1283 .txn_async(move |data| {
1284 let logger_clone2 = logger_clone.clone();
1285 Box::pin(async move {
1286 logger_clone2
1287 .lock()
1288 .unwrap()
1289 .push("execute logic".to_string());
1290 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1291 let _conn2 = data.get_data_conn_async::<MyDataConn>("bar").await?;
1292 Ok(())
1293 })
1294 })
1295 .await
1296 {
1297 match e.reason::<DataConnError>() {
1298 Ok(DataConnError::FailToPostCommitDataConn { errors }) => {
1299 assert_eq!(errors.len(), 2);
1300 assert_eq!(errors[0].index, 0);
1301 assert_eq!(errors[0].name, "foo".into());
1302 assert_eq!(
1303 errors[0].err.reason::<&str>().unwrap(),
1304 &"post commit error"
1305 );
1306 assert_eq!(errors[1].index, 1);
1307 assert_eq!(errors[1].name, "bar".into());
1308 assert_eq!(
1309 errors[1].err.reason::<&str>().unwrap(),
1310 &"post commit error"
1311 );
1312 }
1313 _ => panic!("{e:?}"),
1314 }
1315 }
1316 }
1317
1318 assert_eq!(
1319 *logger.lock().unwrap(),
1320 &[
1321 "MyDataSrc::new 1",
1322 "MyDataSrc::new 2",
1323 "MyDataSrc::setup_async 1",
1324 "MyDataSrc::setup_async 2",
1325 "execute logic",
1326 "MyDataSrc::create_data_conn_async 1",
1327 "MyDataConn::new 1",
1328 "MyDataSrc::create_data_conn_async 2",
1329 "MyDataConn::new 2",
1330 "MyDataConn::pre_commit_async 1",
1331 "MyDataConn::pre_commit_async 2",
1332 "MyDataConn::commit_async 1",
1333 "MyDataConn::commit_async 2",
1334 "MyDataConn::post_commit_async 1 failed",
1335 "MyDataConn::post_commit_async 2 failed",
1336 "MyDataConn::on_txn_failure_async 1",
1337 "MyDataConn::on_txn_failure_async 2",
1338 "MyDataConn::close 2",
1339 "MyDataConn::drop 2",
1340 "MyDataConn::close 1",
1341 "MyDataConn::drop 1",
1342 "MyDataSrc::close 2",
1343 "MyDataSrc::drop 2",
1344 "MyDataSrc::close 1",
1345 "MyDataSrc::drop 1",
1346 ]
1347 );
1348 }
1349
1350 #[tokio::test]
1351 async fn test_txn_but_failed_to_rollback() {
1352 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1353 {
1354 let mut hub = DataHub::new();
1355
1356 hub.uses(
1357 "foo",
1358 MyDataSrc::new(1, logger.clone(), Failure::FailToRollback),
1359 );
1360 hub.uses(
1361 "bar",
1362 MyDataSrc::new(2, logger.clone(), Failure::FailToRollback),
1363 );
1364
1365 let logger_clone = logger.clone();
1366 if let Err(e) = hub
1367 .txn_async(move |data| {
1368 let logger_clone2 = logger_clone.clone();
1369 Box::pin(async move {
1370 logger_clone2
1371 .lock()
1372 .unwrap()
1373 .push("execute logic".to_string());
1374 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1375 let _conn2 = data.get_data_conn_async::<MyDataConn>("bar").await?;
1376 Err(errs::Err::new("logic error"))
1377 })
1378 })
1379 .await
1380 {
1381 match e.reason::<&str>() {
1382 Ok(s) => assert_eq!(s, &"logic error"),
1383 _ => panic!("{e:?}"),
1384 }
1385 }
1386 }
1387
1388 assert_eq!(
1389 *logger.lock().unwrap(),
1390 &[
1391 "MyDataSrc::new 1",
1392 "MyDataSrc::new 2",
1393 "MyDataSrc::setup_async 1",
1394 "MyDataSrc::setup_async 2",
1395 "execute logic",
1396 "MyDataSrc::create_data_conn_async 1",
1397 "MyDataConn::new 1",
1398 "MyDataSrc::create_data_conn_async 2",
1399 "MyDataConn::new 2",
1400 "MyDataConn::rollback_async 1 failed",
1401 "MyDataConn::rollback_async 2 failed",
1402 "MyDataConn::on_txn_failure_async 1",
1403 "MyDataConn::on_txn_failure_async 2",
1404 "MyDataConn::close 2",
1405 "MyDataConn::drop 2",
1406 "MyDataConn::close 1",
1407 "MyDataConn::drop 1",
1408 "MyDataSrc::close 2",
1409 "MyDataSrc::drop 2",
1410 "MyDataSrc::close 1",
1411 "MyDataSrc::drop 1",
1412 ]
1413 );
1414 }
1415
1416 #[tokio::test]
1417 async fn test_txn_with_commit_order() {
1418 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1419 {
1420 let mut hub = DataHub::with_commit_order(&["bar", "foo"]);
1421
1422 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
1423 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
1424
1425 let logger_clone = logger.clone();
1426 hub.txn_async(move |data| {
1427 let logger_clone2 = logger_clone.clone();
1428 Box::pin(async move {
1429 logger_clone2
1430 .lock()
1431 .unwrap()
1432 .push("execute logic".to_string());
1433 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1434 let _conn2 = data.get_data_conn_async::<MyDataConn>("bar").await?;
1435 Ok(())
1436 })
1437 })
1438 .await
1439 .unwrap();
1440 }
1441
1442 assert_eq!(
1443 *logger.lock().unwrap(),
1444 &[
1445 "MyDataSrc::new 1",
1446 "MyDataSrc::new 2",
1447 "MyDataSrc::setup_async 1",
1448 "MyDataSrc::setup_async 2",
1449 "execute logic",
1450 "MyDataSrc::create_data_conn_async 1",
1451 "MyDataConn::new 1",
1452 "MyDataSrc::create_data_conn_async 2",
1453 "MyDataConn::new 2",
1454 "MyDataConn::pre_commit_async 2",
1455 "MyDataConn::pre_commit_async 1",
1456 "MyDataConn::commit_async 2",
1457 "MyDataConn::commit_async 1",
1458 "MyDataConn::post_commit_async 2",
1459 "MyDataConn::post_commit_async 1",
1460 "MyDataConn::close 1",
1461 "MyDataConn::drop 1",
1462 "MyDataConn::close 2",
1463 "MyDataConn::drop 2",
1464 "MyDataSrc::close 2",
1465 "MyDataSrc::drop 2",
1466 "MyDataSrc::close 1",
1467 "MyDataSrc::drop 1",
1468 ]
1469 );
1470 }
1471
1472 #[tokio::test]
1473 async fn test_txn_but_fail_to_setup() {
1474 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1475 {
1476 let mut hub = DataHub::new();
1477
1478 hub.uses(
1479 "foo",
1480 MyDataSrc::new(1, logger.clone(), Failure::FailToSetup),
1481 );
1482
1483 let logger_clone = logger.clone();
1484
1485 if let Err(e) = hub
1486 .txn_async(move |_data| {
1487 let logger_clone2 = logger_clone.clone();
1488 Box::pin(async move {
1489 logger_clone2
1490 .lock()
1491 .unwrap()
1492 .push("execute logic".to_string());
1493 Ok(())
1494 })
1495 })
1496 .await
1497 {
1498 match e.reason::<DataHubError>() {
1499 Ok(DataHubError::FailToSetupLocalDataSrcs { errors }) => {
1500 assert_eq!(errors.len(), 1);
1501 assert_eq!(errors[0].index, 0);
1502 assert_eq!(errors[0].name, "foo".into());
1503 assert_eq!(errors[0].err.reason::<String>().unwrap(), "setup error");
1504 }
1505 _ => panic!(),
1506 }
1507 }
1508 }
1509
1510 assert_eq!(
1511 *logger.lock().unwrap(),
1512 &[
1513 "MyDataSrc::new 1",
1514 "MyDataSrc::setup_async 1 failed",
1515 "MyDataSrc::drop 1",
1516 ]
1517 );
1518 }
1519
1520 #[tokio::test]
1521 async fn test_get_data_conn_cached() {
1522 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1523 {
1524 let mut hub = DataHub::new();
1525
1526 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
1527
1528 let logger_clone = logger.clone();
1529
1530 if let Err(e) = hub
1531 .txn_async(move |data| {
1532 let logger_clone2 = logger_clone.clone();
1533 Box::pin(async move {
1534 logger_clone2
1535 .lock()
1536 .unwrap()
1537 .push("execute logic".to_string());
1538 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1539 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1540 Ok(())
1541 })
1542 })
1543 .await
1544 {
1545 panic!("{e:?}");
1546 }
1547 }
1548
1549 assert_eq!(
1550 *logger.lock().unwrap(),
1551 &[
1552 "MyDataSrc::new 1",
1553 "MyDataSrc::setup_async 1",
1554 "execute logic",
1555 "MyDataSrc::create_data_conn_async 1",
1556 "MyDataConn::new 1",
1557 "MyDataConn::pre_commit_async 1",
1558 "MyDataConn::commit_async 1",
1559 "MyDataConn::post_commit_async 1",
1560 "MyDataConn::close 1",
1561 "MyDataConn::drop 1",
1562 "MyDataSrc::close 1",
1563 "MyDataSrc::drop 1",
1564 ]
1565 );
1566 }
1567
1568 #[tokio::test]
1569 async fn test_get_data_conn_and_no_data_src_to_create_data_conn() {
1570 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1571 {
1572 let mut hub = DataHub::new();
1573
1574 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
1575 hub.uses("bar", MyDataSrc::new(2, logger.clone(), Failure::None));
1576
1577 let logger_clone = logger.clone();
1578 let err = hub
1579 .txn_async(move |data| {
1580 let logger_clone2 = logger_clone.clone();
1581 Box::pin(async move {
1582 logger_clone2
1583 .lock()
1584 .unwrap()
1585 .push("execute logic".to_string());
1586 let _conn1 = data.get_data_conn_async::<MyDataConn>("fxx").await?;
1587 Ok(())
1588 })
1589 })
1590 .await
1591 .unwrap_err();
1592
1593 match err.reason::<DataHubError>() {
1594 Ok(r) => match r {
1595 DataHubError::NoDataSrcToCreateDataConn {
1596 name,
1597 data_conn_type,
1598 } => {
1599 assert_eq!(name.as_ref(), "fxx");
1600 assert_eq!(
1601 data_conn_type,
1602 &"sabi::tokio::data_hub::tests_of_data_hub::MyDataConn"
1603 );
1604 }
1605 _ => panic!(),
1606 },
1607 _ => panic!(),
1608 }
1609 }
1610
1611 assert_eq!(
1612 *logger.lock().unwrap(),
1613 &[
1614 "MyDataSrc::new 1",
1615 "MyDataSrc::new 2",
1616 "MyDataSrc::setup_async 1",
1617 "MyDataSrc::setup_async 2",
1618 "execute logic",
1619 "MyDataSrc::close 2",
1620 "MyDataSrc::drop 2",
1621 "MyDataSrc::close 1",
1622 "MyDataSrc::drop 1",
1623 ]
1624 );
1625 }
1626
1627 #[tokio::test]
1628 async fn test_get_data_conn_and_failed_to_creata_data_conn() {
1629 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1630 {
1631 let mut hub = DataHub::new();
1632
1633 hub.uses(
1634 "foo",
1635 MyDataSrc::new(1, logger.clone(), Failure::FailToCreateDataConn),
1636 );
1637
1638 let logger_clone = logger.clone();
1639
1640 let err = hub
1641 .txn_async(move |data| {
1642 let logger_clone2 = logger_clone.clone();
1643 Box::pin(async move {
1644 logger_clone2
1645 .lock()
1646 .unwrap()
1647 .push("execute logic".to_string());
1648 data.get_data_conn_async::<MyDataConn>("foo").await?;
1649 Ok(())
1650 })
1651 })
1652 .await
1653 .unwrap_err();
1654
1655 match err.reason::<DataSrcError>() {
1656 Ok(DataSrcError::FailToCreateDataConn {
1657 name,
1658 data_conn_type,
1659 }) => {
1660 assert_eq!(name.as_ref(), "foo");
1661 assert_eq!(
1662 data_conn_type,
1663 &"sabi::tokio::data_hub::tests_of_data_hub::MyDataConn"
1664 );
1665 }
1666 _ => panic!(),
1667 }
1668 }
1669
1670 assert_eq!(
1671 *logger.lock().unwrap(),
1672 &[
1673 "MyDataSrc::new 1",
1674 "MyDataSrc::setup_async 1",
1675 "execute logic",
1676 "MyDataSrc::create_data_conn_async 1 failed",
1677 "MyDataSrc::close 1",
1678 "MyDataSrc::drop 1",
1679 ]
1680 );
1681 }
1682
1683 #[tokio::test]
1684 async fn test_get_data_conn_and_failed_to_cast_data_conn() {
1685 let logger = Arc::new(Mutex::new(Vec::<String>::new()));
1686 {
1687 let mut hub = DataHub::new();
1688
1689 hub.uses("foo", MyDataSrc::new(1, logger.clone(), Failure::None));
1690
1691 let logger_clone = logger.clone();
1692
1693 let err = hub
1694 .txn_async(move |data| {
1695 let logger_clone2 = logger_clone.clone();
1696 Box::pin(async move {
1697 logger_clone2
1698 .lock()
1699 .unwrap()
1700 .push("execute logic".to_string());
1701 if let Err(e) = data.get_data_conn_async::<AnotherDataConn>("foo").await {
1702 match e.reason::<DataSrcError>() {
1703 Ok(DataSrcError::FailToCastDataConn { name, target_type }) => {
1704 assert_eq!(name.as_ref(), "foo");
1705 assert_eq!(
1706 target_type,
1707 &"sabi::tokio::data_hub::tests_of_data_hub::AnotherDataConn"
1708 );
1709 }
1710 _ => panic!("{e:?}"),
1711 }
1712 } else {
1713 panic!();
1714 }
1715
1716 let _conn1 = data.get_data_conn_async::<MyDataConn>("foo").await?;
1717
1718 if let Err(e) = data.get_data_conn_async::<AnotherDataConn>("foo").await {
1719 match e.reason::<DataConnError>() {
1720 Ok(DataConnError::FailToCastDataConn { name, target_type }) => {
1721 assert_eq!(name.as_ref(), "foo");
1722 assert_eq!(
1723 target_type,
1724 &"sabi::tokio::data_hub::tests_of_data_hub::AnotherDataConn"
1725 );
1726 Err(e)
1727 }
1728 _ => panic!("{e:?}"),
1729 }
1730 } else {
1731 panic!();
1732 }
1733 })
1734 })
1735 .await
1736 .unwrap_err();
1737
1738 match err.reason::<DataConnError>() {
1739 Ok(DataConnError::FailToCastDataConn { name, target_type }) => {
1740 assert_eq!(name.as_ref(), "foo");
1741 assert_eq!(
1742 target_type,
1743 &"sabi::tokio::data_hub::tests_of_data_hub::AnotherDataConn"
1744 );
1745 }
1746 _ => panic!("{err:?}"),
1747 }
1748 }
1749
1750 assert_eq!(
1751 *logger.lock().unwrap(),
1752 &[
1753 "MyDataSrc::new 1",
1754 "MyDataSrc::setup_async 1",
1755 "execute logic",
1756 "MyDataSrc::create_data_conn_async 1",
1757 "MyDataConn::new 1",
1758 "MyDataConn::rollback_async 1",
1759 "MyDataConn::on_txn_failure_async 1",
1760 "MyDataConn::close 1",
1761 "MyDataConn::drop 1",
1762 "MyDataSrc::close 1",
1763 "MyDataSrc::drop 1",
1764 ]
1765 );
1766 }
1767
1768 trait Data {}
1769 impl Data for DataHub {}
1770
1771 async fn process_async(_data: &mut impl Data) -> errs::Result<()> {
1772 Ok(())
1773 }
1774
1775 #[tokio::test]
1776 async fn data_hub_implements_send_trait() {
1777 let handle = tokio::spawn(async {
1778 let mut data = DataHub::new();
1779 data.run_async(_logic!(process_async)).await.unwrap();
1780 });
1781
1782 handle.await.unwrap();
1783 }
1784
1785 #[tokio::test]
1786 async fn txn_async_in_spawn() {
1787 let handle = tokio::spawn(async {
1788 let mut data = DataHub::new();
1789 data.txn_async(_logic!(process_async)).await.unwrap();
1790 });
1791
1792 handle.await.unwrap();
1793 }
1794}