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
use std::error::Error;
use std::time::Duration;

use futures::compat::Future01CompatExt;
use log::*;
use rusoto_s3::{PutObjectRequest, S3};
use async_trait::async_trait;
use crate::event_emitter::EventEmitter;
use std::future::Future;
use futures::future::FutureExt;
use futures::TryStreamExt;


#[derive(Clone)]
pub struct S3EventEmitter<S, F, OnEmission, EmissionResult>
    where
        S: S3 + Send + 'static,
        F: Fn(&[u8]) -> String,
        EmissionResult: Future<Output=Result<(), Box<dyn Error + Send + Sync + 'static>>>  + Send + 'static,
        OnEmission: Fn(String, String) -> EmissionResult  + Send + Sync + 'static,
{
    s3: S,
    output_bucket: String,
    key_fn: F,
    on_emission: OnEmission,
}

impl<S, F, OnEmission, EmissionResult> S3EventEmitter<S, F, OnEmission, EmissionResult>
    where
        S: S3 + Send + 'static,
        F: Fn(&[u8]) -> String,
        EmissionResult: Future<Output=Result<(), Box<dyn Error + Send + Sync + 'static>>>  + Send + 'static,
        OnEmission: Fn(String, String) -> EmissionResult  + Send + Sync + 'static,

{
    pub fn new(
        s3: S,
        output_bucket: impl Into<String>,
        key_fn: F,
        on_emission: OnEmission,
    ) -> Self
    {
        let output_bucket = output_bucket.into();
        Self {
            s3,
            output_bucket,
            key_fn,
            on_emission
        }
    }
}

#[async_trait]
impl<S, F, OnEmission, EmissionResult> EventEmitter for S3EventEmitter<S, F, OnEmission, EmissionResult>
    where
        S: S3 + Send + Sync + 'static,
        F: Fn(&[u8]) -> String + Send + Sync,
        EmissionResult: Future<Output=Result<(), Box<dyn Error + Send + Sync + 'static>>>  + Send + 'static,
        OnEmission: Fn(String, String) -> EmissionResult  + Send + Sync + 'static,
{
    type Event = Vec<u8>;
    type Error = Box<dyn Error>;

    async fn emit_event(&mut self, events: Vec<Self::Event>) -> Result<(), Self::Error> {

        for event in events {
            let key = (self.key_fn)(&event);
            self.s3.put_object(
                PutObjectRequest {
                    body: Some(event.into()),
                    bucket: self.output_bucket.clone(),
                    key: key.clone(),
                    ..Default::default()
                }
            ).await?;

            // TODO: We shouldn't panic when this happens, we should retry or move on to the next event
            (self.on_emission)(self.output_bucket.clone(), key.clone()).await
                .expect("on_emission failed");
        }


        // let event_uploads = tokio::time::timeout(
        //     Duration::from_secs(5),
        //     futures::future::join_all(event_uploads)
        // ).await?;

        // let mut err = None;
        // for upload in event_uploads {
        //     // upload?;
        // }


        Ok(())
    }
}