Skip to main content

faucet_sink_s3/
sink.rs

1//! S3 sink executor.
2
3use 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
12/// A sink that writes JSON records to S3 as JSON Lines files (or, with the
13/// `arrow` feature, self-contained Parquet objects).
14pub struct S3Sink {
15    config: S3SinkConfig,
16    client: Client,
17}
18
19impl S3Sink {
20    /// Create a new S3 sink from the given configuration.
21    ///
22    /// Builds the S3 client eagerly so it is reused across calls.
23    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    /// Build an S3 client from the configuration.
30    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    /// Serialize a slice of records as JSON Lines bytes.
47    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    /// Generate a unique S3 key for a file.
59    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    /// The effective per-object record cap, combining `batch_size` (write-side
65    /// re-chunking) and `max_records_per_file`. `None` means "one object for
66    /// the whole call". Shared by the JSONL and Parquet write paths.
67    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    /// Upload a single JSONL file to S3.
79    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    /// Upload pre-encoded Parquet objects concurrently. Parquet carries its own
102    /// internal compression, so the crate-local `compression` wrapper is
103    /// deliberately **not** applied here; the content type advertises Parquet.
104    #[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    /// Preflight probe: confirm the configured bucket is reachable and the
144    /// credentials work via a non-mutating `HeadBucket` call. Uploads nothing.
145    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        // Parquet path: encode each chunk as a self-contained Parquet object.
181        #[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        // Pre-serialize each chunk and generate keys before uploading.
204        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    /// The S3 sink consumes Arrow `RecordBatch`es natively **only** when
228    /// configured for the [`Parquet`](S3SinkFormat::Parquet) format; the JSONL
229    /// format has no columnar representation and stays on the row path
230    /// (RFC 0002 / #375).
231    #[cfg(feature = "arrow")]
232    fn supports_columnar(&self) -> bool {
233        matches!(self.config.format, S3SinkFormat::Parquet)
234    }
235
236    /// Write an Arrow `RecordBatch` as one or more self-contained Parquet
237    /// objects (sliced by the effective per-object cap), skipping the
238    /// `Value` round-trip. Falls back to the row path for a non-Parquet
239    /// format (which `supports_columnar` prevents the pipeline from reaching,
240    /// but a direct caller might).
241    #[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/// Encode an Arrow `RecordBatch` into a complete, self-contained Parquet file
274/// (ZSTD-compressed) in memory.
275#[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    /// Helper to build an S3Sink synchronously for tests that never make network calls.
306    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        // UUID is 36 chars
352        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        // No prefix means key starts with UUID
361        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    // ── Parquet columnar path (feature `arrow`) ──────────────────────────────
377
378    #[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        // Parquet magic header/footer.
410        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        // White-box: confirm the codec resolved from file_extension is Gzip
424        // and that compress_buf produces a gzip-magic-prefixed buffer.
425        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        // gzip magic bytes.
430        assert_eq!(&compressed[..2], b"\x1f\x8b");
431    }
432}