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

use crate::event_emitter::EventEmitter;
use async_trait::async_trait;

use rusoto_s3::{PutObjectRequest, S3};
use std::future::Future;

#[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>;

    #[tracing::instrument(skip(self, events))]
    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(())
    }
}