Skip to main content

sabi/tokio/
data_hub.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
5use 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/// Represents errors that can occur within the `DataHub`.
18#[derive(Debug)]
19pub enum DataHubError {
20    /// An error indicating that one or more local data sources failed during their setup processes.
21    FailToSetupLocalDataSrcs {
22        /// A vector of errors, each containing the name of the data source and the error itself.
23        errors: Vec<ErrEntry>,
24    },
25
26    /// An error indicating that no suitable data source was found to create a data connection
27    /// with the specified name and type.
28    NoDataSrcToCreateDataConn {
29        /// The name of the data connection that could not be created.
30        name: Arc<str>,
31
32        /// The string representation of the data connection type that was requested.
33        data_conn_type: &'static str,
34    },
35}
36
37impl DataHub {
38    /// Creates a new `DataHub` instance.
39    ///
40    /// This initializes the `DataHub` with no local data sources and an empty data connection manager.
41    /// Global data sources, if any, are copied into the `data_src_map`.
42    #[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    /// Creates a new `DataHub` instance with a specified commit order for data connections.
56    ///
57    /// This allows defining the order in which data connections will be committed. Connections
58    /// not specified in `names` will be committed after the specified ones, in an undefined order.
59    /// Global data sources are copied into the `data_src_map`.
60    ///
61    /// # Parameters
62    ///
63    /// * `names` - An array of string slices specifying the desired commit order by data connection
64    ///   name.
65    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    /// Registers a local data source with the `DataHub`.
78    ///
79    /// This method allows adding a custom data source which can provide data connections.
80    /// Data sources can only be added before `run_async` or `txn_async` are called.
81    ///
82    /// # Parameters
83    ///
84    /// * `name` - The name to associate with this data source.
85    /// * `ds` - The data source instance, which must implement `DataSrc` and have a `'static`
86    ///   lifetime. If this `DataHub` is moved between threads, `ds` must also implement `Send`.
87    ///
88    /// # Type Parameters
89    ///
90    /// * `S` - The type of the data source.
91    /// * `C` - The type of the data connection provided by the data source.
92    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    /// Deregisters a local data source from the `DataHub`.
104    ///
105    /// This removes a data source previously added with `uses`. Data sources can only be
106    /// removed before `run_async` or `txn_async` are called.
107    ///
108    /// # Parameters
109    ///
110    /// * `name` - The name of the data source to remove.
111    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    /// Executes an asynchronous logic function with the `DataHub` and handles setup and cleanup.
144    ///
145    /// This method sets up local data sources, runs the provided `logic_fn`, and then
146    /// cleans up all data connections and sources. It does *not* automatically commit
147    /// or rollback any transactions.
148    ///
149    /// # Parameters
150    ///
151    /// * `logic_fn` - An asynchronous function that takes a mutable reference to `DataHub`
152    ///                and returns a `Result`. This function contains the application's logic.
153    ///                The returned `Future` must implement `Send`.
154    ///
155    /// # Type Parameters
156    ///
157    /// * `F` - The type of the asynchronous logic function.
158    ///
159    /// # Returns
160    ///
161    /// A `Result` indicating the success or failure of the `logic_fn` execution or
162    /// the setup of data sources.
163    #[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    /// Executes a given asynchronous logic function within a managed transaction.
178    ///
179    /// This method starts by asynchronously setting up local data sources, runs the provided closure,
180    /// and then attempts to asynchronously commit all open data connections in the session.
181    ///
182    /// If any error occurs during the execution of the closure or during the commit phase,
183    /// it initiates an asynchronous rollback on all data connections and reports the transaction
184    /// failure details.
185    /// Finally, it cleans up session resources.
186    ///
187    /// # Parameters
188    ///
189    /// * `logic_fn`: An asynchronous closure that encapsulates the business logic to be executed.
190    ///   It takes a mutable reference to [`DataHub`] as an argument and returns a pinned, boxed
191    ///   future.
192    ///
193    /// # Type Parameters
194    ///
195    /// * `F` - The type of the asynchronous transactional logic function.
196    ///
197    /// # Returns
198    ///
199    /// * `errs::Result<()>`: `Ok(())` if the closure and the commit phase succeed,
200    ///   or an [`errs::Err`] if any phase fails.
201    #[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    /// Retrieves an existing data connection or creates a new one if it doesn't exist.
226    ///
227    /// This asynchronous method first checks if a data connection with the given `name`
228    /// and type `C` already exists. If not, it attempts to find a suitable data source
229    /// (local or global) to create a new data connection.
230    ///
231    /// # Parameters
232    ///
233    /// * `name` - The name of the data connection to retrieve or create.
234    ///
235    /// # Type Parameters
236    ///
237    /// * `C` - The expected type of the data connection, which must implement `DataConn` and have
238    ///   a `'static` lifetime.
239    ///
240    /// # Returns
241    ///
242    /// A `Result` which is `Ok` containing a mutable reference to the data connection
243    /// if found or successfully created, or an `Err` if no suitable data source is found
244    /// or connection creation fails.
245    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                // impossible case.
272            }
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}