use aws_config::BehaviorVersion;
use aws_sdk_s3::config::Credentials;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::{Client, Config as S3Config};
use faucet_conformance::{assert_config_schema_valid_value, assert_errors_not_panics};
use faucet_source_s3::{S3Source, S3SourceConfig};
use testcontainers::{ContainerAsync, runners::AsyncRunner};
use testcontainers_modules::minio::MinIO;
const ACCESS_KEY: &str = "minioadmin";
const SECRET_KEY: &str = "minioadmin";
const REGION: &str = "us-east-1";
const TEST_BUCKET: &str = "faucet-stream-conformance";
#[test]
fn conformance_config_schema_valid() {
let schema = serde_json::to_value(schemars::schema_for!(S3SourceConfig)).unwrap();
assert_config_schema_valid_value(&schema, "faucet-source-s3");
}
async fn start_minio() -> (ContainerAsync<MinIO>, String) {
let container: ContainerAsync<MinIO> = MinIO::default()
.start()
.await
.expect("minio container start");
let port = container
.get_host_port_ipv4(9000)
.await
.expect("minio port");
let endpoint = format!("http://127.0.0.1:{port}");
(container, endpoint)
}
async fn build_admin_client(endpoint: &str) -> Client {
let creds = Credentials::new(ACCESS_KEY, SECRET_KEY, None, None, "test");
let sdk_config = aws_config::defaults(BehaviorVersion::latest())
.region(aws_config::Region::new(REGION))
.endpoint_url(endpoint)
.credentials_provider(creds)
.load()
.await;
let s3_config = S3Config::from(&sdk_config)
.to_builder()
.force_path_style(true)
.build();
Client::from_conf(s3_config)
}
async fn seed_bucket(endpoint: &str, objects: &[(String, String)]) {
let client = build_admin_client(endpoint).await;
client
.create_bucket()
.bucket(TEST_BUCKET)
.send()
.await
.expect("create bucket");
for (key, body) in objects {
client
.put_object()
.bucket(TEST_BUCKET)
.key(key)
.body(ByteStream::from(body.clone().into_bytes()))
.send()
.await
.expect("put object");
}
}
async fn build_source(endpoint: &str, config: S3SourceConfig) -> S3Source {
unsafe {
std::env::set_var("AWS_ACCESS_KEY_ID", ACCESS_KEY);
std::env::set_var("AWS_SECRET_ACCESS_KEY", SECRET_KEY);
std::env::set_var("AWS_DEFAULT_REGION", REGION);
}
let config = config
.endpoint_url(endpoint.to_string())
.region(REGION.to_string());
S3Source::new(config).await.expect("S3Source::new")
}
fn jsonl_body(start: i64, end_inclusive: i64) -> String {
let mut out = String::new();
for i in start..=end_inclusive {
out.push_str(&format!("{{\"id\":{i}}}\n"));
}
out
}
#[tokio::test(flavor = "multi_thread")]
async fn conformance_bounded_memory() {
let (_container, endpoint) = start_minio().await;
seed_bucket(
&endpoint,
&[("data.jsonl".to_string(), jsonl_body(1, 5_000))],
)
.await;
let config = S3SourceConfig::new(TEST_BUCKET).with_batch_size(250);
let source = build_source(&endpoint, config).await;
faucet_conformance::assert_bounded_memory(&source, 250, 5_000).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn conformance_errors_not_panics() {
unsafe {
std::env::set_var("AWS_ACCESS_KEY_ID", "bogus");
std::env::set_var("AWS_SECRET_ACCESS_KEY", "bogus");
std::env::set_var("AWS_DEFAULT_REGION", REGION);
}
let config = S3SourceConfig::new("does-not-exist")
.endpoint_url("http://127.0.0.1:1".to_string())
.region(REGION.to_string());
let source = S3Source::new(config).await.expect("source builds lazily");
assert_errors_not_panics(&source).await;
}