|
| 1 | +use anyhow::{anyhow, Result}; |
| 2 | +use aws_config::meta::region::RegionProviderChain; |
| 3 | +use aws_sdk_s3::{Client, Endpoint, Region}; |
| 4 | +use chrono::Utc; |
| 5 | +use file_store::traits::MsgBytes; |
| 6 | +use file_store::{file_sink, file_upload, FileStore, FileType, Settings}; |
| 7 | +use std::path::Path; |
| 8 | +use std::{str::FromStr, sync::Arc}; |
| 9 | +use tempfile::TempDir; |
| 10 | +use tokio::sync::Mutex; |
| 11 | +use tonic::transport::Uri; |
| 12 | +use uuid::Uuid; |
| 13 | + |
| 14 | +pub const AWSLOCAL_DEFAULT_ENDPOINT: &str = "http://127.0.0.1:4566"; |
| 15 | + |
| 16 | +pub fn gen_bucket_name() -> String { |
| 17 | + format!("mvr-{}-{}", Uuid::new_v4(), Utc::now().timestamp_millis()) |
| 18 | +} |
| 19 | + |
| 20 | +// Interacts with the locastack. |
| 21 | +// Used to create mocked aws buckets and files. |
| 22 | +pub struct AwsLocal { |
| 23 | + pub fs_settings: Settings, |
| 24 | + pub file_store: FileStore, |
| 25 | + pub aws_client: aws_sdk_s3::Client, |
| 26 | +} |
| 27 | + |
| 28 | +impl AwsLocal { |
| 29 | + async fn create_aws_client(settings: &Settings) -> aws_sdk_s3::Client { |
| 30 | + let endpoint: Option<Endpoint> = match &settings.endpoint { |
| 31 | + Some(endpoint) => Uri::from_str(endpoint) |
| 32 | + .map(Endpoint::immutable) |
| 33 | + .map(Some) |
| 34 | + .unwrap(), |
| 35 | + _ => None, |
| 36 | + }; |
| 37 | + let region = Region::new(settings.region.clone()); |
| 38 | + let region_provider = RegionProviderChain::first_try(region).or_default_provider(); |
| 39 | + |
| 40 | + let mut config = aws_config::from_env().region(region_provider); |
| 41 | + config = config.endpoint_resolver(endpoint.unwrap()); |
| 42 | + |
| 43 | + let creds = aws_types::credentials::Credentials::from_keys( |
| 44 | + settings.access_key_id.as_ref().unwrap(), |
| 45 | + settings.secret_access_key.as_ref().unwrap(), |
| 46 | + None, |
| 47 | + ); |
| 48 | + config = config.credentials_provider(creds); |
| 49 | + |
| 50 | + let config = config.load().await; |
| 51 | + |
| 52 | + Client::new(&config) |
| 53 | + } |
| 54 | + |
| 55 | + pub async fn new(endpoint: &str, bucket: &str) -> AwsLocal { |
| 56 | + let settings = Settings { |
| 57 | + bucket: bucket.into(), |
| 58 | + endpoint: Some(endpoint.into()), |
| 59 | + region: "us-east-1".into(), |
| 60 | + access_key_id: Some("random".into()), |
| 61 | + secret_access_key: Some("random2".into()), |
| 62 | + }; |
| 63 | + let client = Self::create_aws_client(&settings).await; |
| 64 | + client.create_bucket().bucket(bucket).send().await.unwrap(); |
| 65 | + AwsLocal { |
| 66 | + aws_client: client, |
| 67 | + fs_settings: settings.clone(), |
| 68 | + file_store: file_store::FileStore::from_settings(&settings) |
| 69 | + .await |
| 70 | + .unwrap(), |
| 71 | + } |
| 72 | + } |
| 73 | + |
| 74 | + pub fn fs_settings(&self) -> Settings { |
| 75 | + self.fs_settings.clone() |
| 76 | + } |
| 77 | + |
| 78 | + pub async fn put_proto_to_aws<T: prost::Message + MsgBytes>( |
| 79 | + &self, |
| 80 | + items: Vec<T>, |
| 81 | + file_type: FileType, |
| 82 | + metric_name: &'static str, |
| 83 | + ) -> Result<String> { |
| 84 | + let tmp_dir = TempDir::new()?; |
| 85 | + let tmp_dir_path = tmp_dir.path().to_owned(); |
| 86 | + |
| 87 | + let (shutdown_trigger, shutdown_listener) = triggered::trigger(); |
| 88 | + |
| 89 | + let (file_upload, file_upload_server) = |
| 90 | + file_upload::FileUpload::from_settings_tm(&self.fs_settings) |
| 91 | + .await |
| 92 | + .unwrap(); |
| 93 | + |
| 94 | + let (item_sink, item_server) = |
| 95 | + file_sink::FileSinkBuilder::new(file_type, &tmp_dir_path, file_upload, metric_name) |
| 96 | + .auto_commit(false) |
| 97 | + .roll_time(std::time::Duration::new(15, 0)) |
| 98 | + .create::<T>() |
| 99 | + .await |
| 100 | + .unwrap(); |
| 101 | + |
| 102 | + for item in items { |
| 103 | + item_sink.write(item, &[]).await.unwrap(); |
| 104 | + } |
| 105 | + let item_recv = item_sink.commit().await.unwrap(); |
| 106 | + |
| 107 | + let uploaded_file = Arc::new(Mutex::new(String::default())); |
| 108 | + let up_2 = uploaded_file.clone(); |
| 109 | + let mut timeout = std::time::Duration::new(5, 0); |
| 110 | + |
| 111 | + tokio::spawn(async move { |
| 112 | + let uploaded_files = item_recv.await.unwrap().unwrap(); |
| 113 | + assert!(uploaded_files.len() == 1); |
| 114 | + let mut val = up_2.lock().await; |
| 115 | + *val = uploaded_files.first().unwrap().to_string(); |
| 116 | + |
| 117 | + // After files uploaded to aws the must be removed. |
| 118 | + // So we wait when dir will be empty. |
| 119 | + // It means all files are uploaded to aws |
| 120 | + loop { |
| 121 | + if is_dir_has_files(&tmp_dir_path) { |
| 122 | + let dur = std::time::Duration::from_millis(10); |
| 123 | + tokio::time::sleep(dur).await; |
| 124 | + timeout -= dur; |
| 125 | + continue; |
| 126 | + } |
| 127 | + break; |
| 128 | + } |
| 129 | + |
| 130 | + shutdown_trigger.trigger(); |
| 131 | + }); |
| 132 | + |
| 133 | + tokio::try_join!( |
| 134 | + file_upload_server.run(shutdown_listener.clone()), |
| 135 | + item_server.run(shutdown_listener.clone()) |
| 136 | + ) |
| 137 | + .unwrap(); |
| 138 | + |
| 139 | + tmp_dir.close()?; |
| 140 | + |
| 141 | + let res = uploaded_file.lock().await; |
| 142 | + Ok(res.clone()) |
| 143 | + } |
| 144 | + |
| 145 | + pub async fn put_file_to_aws(&self, file_path: &Path) -> Result<()> { |
| 146 | + let path_str = file_path.display(); |
| 147 | + if !file_path.exists() { |
| 148 | + return Err(anyhow!("File {path_str} is absent")); |
| 149 | + } |
| 150 | + if !file_path.is_file() { |
| 151 | + return Err(anyhow!("File {path_str} is not a file")); |
| 152 | + } |
| 153 | + self.file_store.put(file_path).await?; |
| 154 | + |
| 155 | + Ok(()) |
| 156 | + } |
| 157 | +} |
| 158 | + |
| 159 | +fn is_dir_has_files(dir_path: &Path) -> bool { |
| 160 | + let entries = std::fs::read_dir(dir_path) |
| 161 | + .unwrap() |
| 162 | + .map(|res| res.map(|e| e.path().is_dir())) |
| 163 | + .collect::<Result<Vec<_>, std::io::Error>>() |
| 164 | + .unwrap(); |
| 165 | + entries.contains(&false) |
| 166 | +} |
0 commit comments