1use crate::config::S3SinkConfig;
4#[cfg(feature = "arrow")]
5use crate::config::S3SinkFormat;
6use async_trait::async_trait;
7use aws_sdk_s3::Client;
8use faucet_core::FaucetError;
9use futures::stream::{self, StreamExt, TryStreamExt};
10use serde_json::Value;
11
12pub struct S3Sink {
15 config: S3SinkConfig,
16 client: Client,
17}
18
19impl S3Sink {
20 pub async fn new(config: S3SinkConfig) -> Result<Self, FaucetError> {
24 faucet_core::validate_batch_size(config.batch_size)?;
25 let client = Self::build_client(&config).await?;
26 Ok(Self { config, client })
27 }
28
29 async fn build_client(config: &S3SinkConfig) -> Result<Client, FaucetError> {
31 let mut config_loader = aws_config::defaults(aws_config::BehaviorVersion::latest());
32
33 if let Some(ref region) = config.region {
34 config_loader = config_loader.region(aws_config::Region::new(region.clone()));
35 }
36
37 if let Some(ref endpoint) = config.endpoint_url {
38 config_loader = config_loader.endpoint_url(endpoint);
39 }
40
41 let sdk_config = config_loader.load().await;
42 let client = Client::new(&sdk_config);
43 Ok(client)
44 }
45
46 fn serialize_jsonl(records: &[Value]) -> Result<Vec<u8>, FaucetError> {
48 let mut buf: Vec<u8> = Vec::new();
49 for record in records {
50 let line = serde_json::to_vec(record)
51 .map_err(|e| FaucetError::Sink(format!("JSON serialization failed: {e}")))?;
52 buf.extend_from_slice(&line);
53 buf.push(b'\n');
54 }
55 Ok(buf)
56 }
57
58 fn generate_key(&self) -> String {
60 let id = uuid::Uuid::new_v4();
61 format!("{}{}{}", self.config.prefix, id, self.config.file_extension)
62 }
63
64 fn effective_chunk_cap(&self) -> Option<usize> {
68 match (self.config.batch_size, self.config.max_records_per_file) {
69 (0, None) => None,
70 (0, Some(0)) => None,
71 (0, Some(max)) => Some(max),
72 (bs, None) => Some(bs),
73 (bs, Some(0)) => Some(bs),
74 (bs, Some(max)) => Some(bs.min(max)),
75 }
76 }
77
78 async fn upload_file(&self, key: &str, body: Vec<u8>) -> Result<(), FaucetError> {
80 #[cfg(feature = "compression")]
81 let body = {
82 let codec = self.config.compression.resolve(&self.config.file_extension);
83 faucet_core::compression::warn_mismatch(&self.config.file_extension, codec);
84 faucet_core::compression::compress_buf(&body, codec)?
85 };
86
87 self.client
88 .put_object()
89 .bucket(&self.config.bucket)
90 .key(key)
91 .body(body.into())
92 .content_type("application/x-ndjson")
93 .send()
94 .await
95 .map_err(|e| FaucetError::Sink(format!("S3 put object error for key '{key}': {e}")))?;
96
97 tracing::debug!(key = %key, "Uploaded S3 object");
98 Ok(())
99 }
100
101 #[cfg(feature = "arrow")]
105 async fn upload_parquet_objects(
106 &self,
107 prepared: Vec<(String, Vec<u8>)>,
108 ) -> Result<(), FaucetError> {
109 let concurrency = self.config.concurrency.max(1);
110 stream::iter(prepared)
111 .map(|(key, body)| async move {
112 self.client
113 .put_object()
114 .bucket(&self.config.bucket)
115 .key(&key)
116 .body(body.into())
117 .content_type("application/vnd.apache.parquet")
118 .send()
119 .await
120 .map_err(|e| {
121 FaucetError::Sink(format!("S3 put object error for key '{key}': {e}"))
122 })?;
123 tracing::debug!(key = %key, "Uploaded S3 parquet object");
124 Ok::<(), FaucetError>(())
125 })
126 .buffer_unordered(concurrency)
127 .try_collect::<Vec<()>>()
128 .await?;
129 Ok(())
130 }
131}
132
133#[async_trait]
134impl faucet_core::Sink for S3Sink {
135 fn config_schema(&self) -> serde_json::Value {
136 serde_json::to_value(faucet_core::schema_for!(S3SinkConfig)).expect("schema serialization")
137 }
138
139 fn dataset_uri(&self) -> String {
140 format!("s3://{}/{}", self.config.bucket, self.config.prefix)
141 }
142
143 async fn check(
146 &self,
147 ctx: &faucet_core::check::CheckContext,
148 ) -> Result<faucet_core::check::CheckReport, FaucetError> {
149 use faucet_core::check::{CheckReport, Probe};
150
151 let started = std::time::Instant::now();
152 let probe = match tokio::time::timeout(
153 ctx.timeout,
154 self.client.head_bucket().bucket(&self.config.bucket).send(),
155 )
156 .await
157 {
158 Ok(Ok(_)) => Probe::pass("auth", started.elapsed()),
159 Ok(Err(e)) => Probe::fail_hint(
160 "auth",
161 started.elapsed(),
162 e.to_string(),
163 "check bucket name, credentials, and network",
164 ),
165 Err(_) => Probe::fail("network", started.elapsed(), "timed out"),
166 };
167 Ok(CheckReport::single(probe))
168 }
169
170 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
171 if records.is_empty() {
172 return Ok(0);
173 }
174
175 let chunks: Vec<&[Value]> = match self.effective_chunk_cap() {
176 Some(cap) => records.chunks(cap).collect(),
177 None => vec![records],
178 };
179
180 #[cfg(feature = "arrow")]
182 if matches!(self.config.format, S3SinkFormat::Parquet) {
183 let prepared: Vec<(String, Vec<u8>)> = chunks
184 .iter()
185 .map(|chunk| {
186 let batch = faucet_core::columnar::values_to_record_batch_inferred(chunk)?;
187 let body = encode_parquet(&batch)?;
188 Ok((self.generate_key(), body))
189 })
190 .collect::<Result<Vec<_>, FaucetError>>()?;
191 self.upload_parquet_objects(prepared).await?;
192 tracing::info!(
193 records = records.len(),
194 files = chunks.len(),
195 "S3 parquet batch write complete"
196 );
197 return Ok(records.len());
198 }
199
200 let total_files = chunks.len();
201 let concurrency = self.config.concurrency.max(1);
202
203 let prepared: Vec<(String, Vec<u8>)> = chunks
205 .iter()
206 .map(|chunk| {
207 let body = Self::serialize_jsonl(chunk)?;
208 let key = self.generate_key();
209 Ok((key, body))
210 })
211 .collect::<Result<Vec<_>, FaucetError>>()?;
212
213 stream::iter(prepared)
214 .map(|(key, body)| async move { self.upload_file(&key, body).await })
215 .buffer_unordered(concurrency)
216 .try_collect::<Vec<()>>()
217 .await?;
218
219 tracing::info!(
220 records = records.len(),
221 files = total_files,
222 "S3 batch write complete"
223 );
224 Ok(records.len())
225 }
226
227 #[cfg(feature = "arrow")]
232 fn supports_columnar(&self) -> bool {
233 matches!(self.config.format, S3SinkFormat::Parquet)
234 }
235
236 #[cfg(feature = "arrow")]
242 async fn write_batch_columnar(
243 &self,
244 batch: &arrow::array::RecordBatch,
245 ) -> Result<usize, FaucetError> {
246 if batch.num_rows() == 0 {
247 return Ok(0);
248 }
249 if !matches!(self.config.format, S3SinkFormat::Parquet) {
250 let rows = faucet_core::columnar::record_batch_to_values(batch)?;
251 return self.write_batch(&rows).await;
252 }
253
254 let n = batch.num_rows();
255 let cap = self.effective_chunk_cap().unwrap_or(n).max(1);
256 let mut prepared: Vec<(String, Vec<u8>)> = Vec::new();
257 let mut offset = 0usize;
258 while offset < n {
259 let len = cap.min(n - offset);
260 let slice = batch.slice(offset, len);
261 let body = encode_parquet(&slice)?;
262 prepared.push((self.generate_key(), body));
263 offset += len;
264 }
265
266 let files = prepared.len();
267 self.upload_parquet_objects(prepared).await?;
268 tracing::info!(records = n, files, "S3 parquet columnar write complete");
269 Ok(n)
270 }
271}
272
273#[cfg(feature = "arrow")]
276fn encode_parquet(batch: &arrow::array::RecordBatch) -> Result<Vec<u8>, FaucetError> {
277 use parquet::arrow::ArrowWriter;
278 use parquet::basic::{Compression, ZstdLevel};
279 use parquet::file::properties::WriterProperties;
280
281 let props = WriterProperties::builder()
282 .set_compression(Compression::ZSTD(ZstdLevel::default()))
283 .build();
284 let mut buf: Vec<u8> = Vec::new();
285 {
286 let mut writer = ArrowWriter::try_new(&mut buf, batch.schema(), Some(props))
287 .map_err(|e| FaucetError::Sink(format!("parquet writer init failed: {e}")))?;
288 writer
289 .write(batch)
290 .map_err(|e| FaucetError::Sink(format!("parquet write failed: {e}")))?;
291 writer
292 .close()
293 .map_err(|e| FaucetError::Sink(format!("parquet finalize failed: {e}")))?;
294 }
295 Ok(buf)
296}
297
298#[cfg(test)]
299mod tests {
300 use super::*;
301 use crate::config::S3SinkConfig;
302 use faucet_core::Sink as _;
303 use serde_json::json;
304
305 fn test_sink(config: S3SinkConfig) -> S3Sink {
307 let sdk_config = aws_config::SdkConfig::builder()
308 .behavior_version(aws_config::BehaviorVersion::latest())
309 .build();
310 let client = Client::new(&sdk_config);
311 S3Sink { config, client }
312 }
313
314 #[test]
315 fn dataset_uri_includes_bucket_and_prefix() {
316 let sink = test_sink(S3SinkConfig::new("my-bucket").prefix("data/events/"));
317 assert_eq!(sink.dataset_uri(), "s3://my-bucket/data/events/");
318 }
319
320 #[test]
321 fn serialize_jsonl_produces_newline_delimited() {
322 let records = vec![
323 json!({"id": 1, "name": "Alice"}),
324 json!({"id": 2, "name": "Bob"}),
325 ];
326 let result = S3Sink::serialize_jsonl(&records).unwrap();
327 let text = String::from_utf8(result).unwrap();
328 let lines: Vec<&str> = text.trim().split('\n').collect();
329 assert_eq!(lines.len(), 2);
330
331 let first: Value = serde_json::from_str(lines[0]).unwrap();
332 assert_eq!(first["id"], 1);
333 }
334
335 #[test]
336 fn serialize_jsonl_empty() {
337 let result = S3Sink::serialize_jsonl(&[]).unwrap();
338 assert!(result.is_empty());
339 }
340
341 #[test]
342 fn generate_key_uses_prefix_and_extension() {
343 let sink = test_sink(
344 S3SinkConfig::new("bucket")
345 .prefix("data/")
346 .file_extension(".jsonl"),
347 );
348 let key = sink.generate_key();
349 assert!(key.starts_with("data/"));
350 assert!(key.ends_with(".jsonl"));
351 assert!(key.len() > "data/".len() + ".jsonl".len());
353 }
354
355 #[test]
356 fn generate_key_no_prefix() {
357 let sink = test_sink(S3SinkConfig::new("bucket"));
358 let key = sink.generate_key();
359 assert!(key.ends_with(".jsonl"));
360 assert!(!key.starts_with('/'));
362 }
363
364 #[tokio::test]
365 async fn new_rejects_out_of_range_batch_size() {
366 let mut config = S3SinkConfig::new("bucket");
367 config.batch_size = faucet_core::MAX_BATCH_SIZE + 1;
368 match S3Sink::new(config).await {
369 Err(faucet_core::FaucetError::Config(m)) => {
370 assert!(m.contains("batch_size"), "got: {m}")
371 }
372 _ => panic!("expected a batch_size Config error"),
373 }
374 }
375
376 #[cfg(feature = "arrow")]
379 #[test]
380 fn supports_columnar_only_for_parquet_format() {
381 let parquet_sink = test_sink(S3SinkConfig::new("b").format(S3SinkFormat::Parquet));
382 assert!(faucet_core::Sink::supports_columnar(&parquet_sink));
383 let json_sink = test_sink(S3SinkConfig::new("b"));
384 assert!(!faucet_core::Sink::supports_columnar(&json_sink));
385 }
386
387 #[cfg(feature = "arrow")]
388 #[test]
389 fn encode_parquet_round_trips_via_reader() {
390 use arrow::array::{Int32Array, RecordBatch, StringArray};
391 use arrow::datatypes::{DataType, Field, Schema};
392 use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
393 use std::sync::Arc;
394
395 let schema = Arc::new(Schema::new(vec![
396 Field::new("id", DataType::Int32, false),
397 Field::new("name", DataType::Utf8, true),
398 ]));
399 let batch = RecordBatch::try_new(
400 schema,
401 vec![
402 Arc::new(Int32Array::from(vec![1, 2, 3])),
403 Arc::new(StringArray::from(vec![Some("a"), None, Some("c")])),
404 ],
405 )
406 .unwrap();
407
408 let bytes = encode_parquet(&batch).unwrap();
409 assert_eq!(&bytes[..4], b"PAR1");
411
412 let reader = ParquetRecordBatchReaderBuilder::try_new(bytes::Bytes::from(bytes))
413 .unwrap()
414 .build()
415 .unwrap();
416 let total: usize = reader.map(|b| b.unwrap().num_rows()).sum();
417 assert_eq!(total, 3);
418 }
419
420 #[cfg(feature = "compression")]
421 #[test]
422 fn compress_buf_used_for_gzip_extension() {
423 let cfg = S3SinkConfig::new("bucket").file_extension(".jsonl.gz");
426 let codec = cfg.compression.resolve(&cfg.file_extension);
427 assert_eq!(codec, faucet_core::Compression::Gzip);
428 let compressed = faucet_core::compression::compress_buf(b"hello\n", codec).unwrap();
429 assert_eq!(&compressed[..2], b"\x1f\x8b");
431 }
432}