1
  2
  3
  4
  5
  6
  7
  8
  9
 10
 11
 12
 13
 14
 15
 16
 17
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
use async_trait::async_trait;
use tokio::sync::mpsc;

use crate::event::*;
use crate::error::*;

#[derive(Debug, Clone)]
pub struct LoadData
{
    pub(crate) index: u32,
    pub(crate) offset: u64,
    pub header: EventHeaderRaw,
    pub data: EventData,
}

#[async_trait]
pub trait Loader: Send + Sync + 'static
{
    /// Function invoked when the start of the history is being loaded
    async fn start_of_history(&mut self, _size: usize) { }

    /// Events are being processed
    async fn feed_events(&mut self, _evts: &Vec<EventData>) { }

    /// Load data is being processed
    async fn feed_load_data(&mut self, _data: LoadData) { }

    /// The last event is now received
    async fn end_of_history(&mut self) { }

    /// Callback when the load has failed
    async fn failed(&mut self, err: ChainCreationError) -> Option<ChainCreationError>
    {
        Some(err)
    }
}

#[derive(Debug, Clone, Default)]
pub struct DummyLoader { }

impl Loader
for DummyLoader { }

#[derive(Default)]
pub struct CompositionLoader
{
    pub loaders: Vec<Box<dyn Loader>>,
}

#[async_trait]
impl Loader
for CompositionLoader
{
    async fn start_of_history(&mut self, size: usize)
    {
        for loader in self.loaders.iter_mut() {
            loader.start_of_history(size).await;
        }
    }

    async fn feed_events(&mut self, evts: &Vec<EventData>)
    {
        for loader in self.loaders.iter_mut() {
            loader.feed_events(evts).await;
        }
    }

    async fn feed_load_data(&mut self, data: LoadData)
    {
        for loader in self.loaders.iter_mut() {
            loader.feed_load_data(data.clone()).await;
        }
    }

    async fn end_of_history(&mut self)
    {
        for loader in self.loaders.iter_mut() {
            loader.end_of_history().await;
        }
    }

    async fn failed(&mut self, mut err: ChainCreationError) -> Option<ChainCreationError>
    {
        let err_msg = err.to_string();
        for loader in self.loaders.iter_mut() {
            err = match loader.failed(err).await {
                Some(a) => a,
                None => {
                    ChainCreationError::InternalError(err_msg.clone())
                }
            };
        }
        Some(err)
    }
}

pub struct NotificationLoader
{
    notify: mpsc::Sender<Result<(), ChainCreationError>>
}

impl NotificationLoader {
    pub fn new(notify: mpsc::Sender<Result<(), ChainCreationError>>) -> NotificationLoader {
        NotificationLoader {
            notify
        }
    }
}

#[async_trait]
impl Loader
for NotificationLoader
{
    async fn end_of_history(&mut self)
    {
        let _ = self.notify.send(Ok(())).await;
    }

    async fn failed(&mut self, err: ChainCreationError) -> Option<ChainCreationError>
    {
        let _ = self.notify.send(Err(err)).await;
        None
    }
}