Skip to main content

mz_testdrive/action/
s3.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10// The metadata for Arrow `Field` type requires `std::collections::HashMap`, which is disallowed.
11#[allow(clippy::disallowed_types)]
12use std::collections::HashMap;
13use std::pin::Pin;
14use std::str;
15use std::sync::Arc;
16use std::thread;
17use std::time::Duration;
18
19use anyhow::Context;
20use anyhow::bail;
21use arrow::array::{
22    ArrayRef, BinaryBuilder, BooleanArray, Date32Array, Decimal128Array, FixedSizeBinaryBuilder,
23    Float32Array, Float64Array, Int8Array, Int16Array, Int32Array, Int64Array, Int64Builder,
24    IntervalDayTimeArray, IntervalYearMonthArray, ListBuilder, MapBuilder, StringArray,
25    StringBuilder, StructArray, Time32SecondArray, TimestampMillisecondArray, UInt8Array,
26    UInt16Array, UInt32Array, UInt64Array,
27};
28use arrow::datatypes::{DataType, Field, IntervalDayTime, IntervalUnit, Schema, TimeUnit};
29use arrow::record_batch::RecordBatch;
30use arrow::util::display::ArrayFormatter;
31use arrow::util::display::FormatOptions;
32use async_compression::tokio::bufread::{BzEncoder, GzipEncoder, XzEncoder, ZstdEncoder};
33use chrono::{NaiveDate, NaiveDateTime, NaiveTime, Timelike};
34use parquet::arrow::ArrowWriter;
35use parquet::basic::{BrotliLevel, Compression as ParquetCompression, GzipLevel, ZstdLevel};
36use parquet::file::properties::WriterProperties;
37use regex::Regex;
38use tokio::io::{AsyncRead, AsyncReadExt};
39
40use crate::action::file::Compression;
41use crate::action::file::Contents;
42use crate::action::file::build_compression;
43use crate::action::{ControlFlow, State};
44use crate::parser::BuiltinCommand;
45
46pub async fn run_verify_data(
47    mut cmd: BuiltinCommand,
48    state: &State,
49) -> Result<ControlFlow, anyhow::Error> {
50    let mut expected_body = cmd
51        .input
52        .into_iter()
53        // Strip suffix to allow lines with trailing whitespace
54        .map(|line| {
55            line.trim_end_matches("// allow-trailing-whitespace")
56                .to_string()
57        })
58        .collect::<Vec<String>>();
59    let bucket: String = cmd.args.parse("bucket")?;
60    let key: String = cmd.args.parse("key")?;
61    let sort_rows = cmd.args.opt_bool("sort-rows")?.unwrap_or(false);
62    cmd.args.done()?;
63
64    println!("Verifying contents of S3 bucket {bucket} key {key}...");
65
66    let client = mz_aws_util::s3::new_client(&state.aws_config);
67
68    // List the path until the INCOMPLETE sentinel file disappears so we know the
69    // data is complete.
70    let mut attempts = 0;
71    let all_files;
72    loop {
73        attempts += 1;
74        if attempts > 10 {
75            bail!("S3 path {key} was still incomplete or empty after 10 attempts")
76        }
77
78        let files = client
79            .list_objects_v2()
80            .bucket(&bucket)
81            .prefix(&format!("{}/", key))
82            .send()
83            .await?;
84        match files.contents {
85            Some(files)
86                if files
87                    .iter()
88                    .any(|obj| obj.key().map_or(false, |key| key.contains("INCOMPLETE"))) =>
89            {
90                thread::sleep(Duration::from_secs(1))
91            }
92            // An empty prefix means the writer has not uploaded anything yet,
93            // so retry just like for the INCOMPLETE sentinel.
94            None => thread::sleep(Duration::from_secs(1)),
95            Some(files) => {
96                all_files = files;
97                break;
98            }
99        }
100    }
101
102    let mut rows = vec![];
103    for obj in all_files.iter() {
104        let file = client
105            .get_object()
106            .bucket(&bucket)
107            .key(obj.key().unwrap())
108            .send()
109            .await?;
110        let bytes = file.body.collect().await?.into_bytes();
111
112        let new_rows = match obj.key().unwrap() {
113            key if key.ends_with(".csv") => {
114                let actual_body = str::from_utf8(bytes.as_ref())?;
115                actual_body.lines().map(|l| l.to_string()).collect()
116            }
117            key if key.ends_with(".parquet") => rows_from_parquet(bytes),
118            key => bail!("unexpected file type: {key}"),
119        };
120        rows.extend(new_rows);
121    }
122    if sort_rows {
123        expected_body.sort();
124        rows.sort();
125    }
126    if rows != expected_body {
127        bail!(
128            "content did not match\nexpected:\n{:?}\n\nactual:\n{:?}",
129            expected_body,
130            rows
131        );
132    }
133
134    Ok(ControlFlow::Continue)
135}
136
137pub async fn run_verify_keys(
138    mut cmd: BuiltinCommand,
139    state: &State,
140) -> Result<ControlFlow, anyhow::Error> {
141    let bucket: String = cmd.args.parse("bucket")?;
142    let prefix_path: String = cmd.args.parse("prefix-path")?;
143    let key_pattern: Regex = cmd.args.parse("key-pattern")?;
144    let num_attempts = cmd.args.opt_parse("num-attempts")?.unwrap_or(30);
145    cmd.args.done()?;
146
147    println!("Verifying {key_pattern} in S3 bucket {bucket} path {prefix_path}...");
148
149    let client = mz_aws_util::s3::new_client(&state.aws_config);
150
151    let mut attempts = 0;
152    while attempts <= num_attempts {
153        attempts += 1;
154        let files = client
155            .list_objects_v2()
156            .bucket(&bucket)
157            .prefix(&format!("{}/", prefix_path))
158            .send()
159            .await?;
160        if let Some(files) = files.contents {
161            let files: Vec<_> = files
162                .iter()
163                .filter(|obj| key_pattern.is_match(obj.key().unwrap()))
164                .map(|obj| obj.key().unwrap())
165                .collect();
166            if !files.is_empty() {
167                println!("Found matching files: {files:?}");
168                return Ok(ControlFlow::Continue);
169            }
170        }
171        // Sleep whenever no attempt matched, not only when the prefix was
172        // empty. Otherwise the presence of any non-matching object burns all
173        // attempts back-to-back without giving the writer time to catch up.
174        thread::sleep(Duration::from_secs(1));
175    }
176
177    bail!("Did not find matching files in bucket {bucket} prefix {prefix_path}");
178}
179
180fn rows_from_parquet(bytes: bytes::Bytes) -> Vec<String> {
181    let reader =
182        parquet::arrow::arrow_reader::ParquetRecordBatchReader::try_new(bytes, 1_000_000).unwrap();
183
184    let mut ret = vec![];
185    let format_options = FormatOptions::default();
186    for batch in reader {
187        let batch = batch.unwrap();
188        let converters = batch
189            .columns()
190            .iter()
191            .map(|a| ArrayFormatter::try_new(a.as_ref(), &format_options).unwrap())
192            .collect::<Vec<_>>();
193
194        for row_idx in 0..batch.num_rows() {
195            let mut buf = String::new();
196            for (col_idx, converter) in converters.iter().enumerate() {
197                if col_idx > 0 {
198                    buf.push_str(" ");
199                }
200                converter.value(row_idx).write(&mut buf).unwrap();
201            }
202            ret.push(buf);
203        }
204    }
205    ret
206}
207
208pub async fn run_upload(
209    mut cmd: BuiltinCommand,
210    state: &State,
211) -> Result<ControlFlow, anyhow::Error> {
212    let bucket = cmd.args.string("bucket")?;
213    let count: Option<usize> = cmd.args.opt_parse("count")?;
214
215    let keys: Vec<String> = if let Some(count) = count {
216        // Bulk mode uses `key-prefix` + `i` + optional `key-suffix`,
217        let prefix = cmd.args.string("key-prefix")?;
218        let suffix = cmd.args.opt_string("key-suffix").unwrap_or_default();
219        (0..count).map(|i| format!("{prefix}{i}{suffix}")).collect()
220    } else {
221        // Single-file mode uses `key`.
222        vec![cmd.args.string("key")?]
223    };
224
225    let compression = build_compression(&mut cmd)?;
226    let contents = Contents::parse(&mut cmd)?;
227    cmd.args.done()?;
228
229    // S3 puts need the full body in memory, so materialize the streamed
230    // contents into a buffer before compressing and uploading.
231    let mut body = Vec::new();
232    contents.write_to(&mut body).await?;
233
234    let aws_client = mz_aws_util::s3::new_client(&state.aws_config);
235
236    // TODO(parkmycar): Stream data to S3. The ByteStream type from the AWS config is a bit
237    // cumbersome to work with, so for now just stick with this.
238    let mut reader: Pin<Box<dyn AsyncRead + Send + Sync>> = match compression {
239        Compression::None => Box::pin(&body[..]),
240        Compression::Gzip => Box::pin(GzipEncoder::new(&body[..])),
241        Compression::Bzip2 => Box::pin(BzEncoder::new(&body[..])),
242        Compression::Xz => Box::pin(XzEncoder::new(&body[..])),
243        Compression::Zstd => Box::pin(ZstdEncoder::new(&body[..])),
244    };
245    let mut content = vec![];
246    reader
247        .read_to_end(&mut content)
248        .await
249        .context("compressing")?;
250
251    // Upload the file(s) to S3.
252    println!(
253        "Uploading {} files to S3 bucket, starting with '{bucket}/{}'",
254        keys.len(),
255        keys.first().map(String::as_str).unwrap_or("<none>")
256    );
257    for key in &keys {
258        aws_client
259            .put_object()
260            .bucket(&bucket)
261            .key(key)
262            .body(content.clone().into())
263            .send()
264            .await
265            .context("s3 put")?;
266    }
267
268    Ok(ControlFlow::Continue)
269}
270
271pub async fn run_set_presigned_url(
272    mut cmd: BuiltinCommand,
273    state: &mut State,
274) -> Result<ControlFlow, anyhow::Error> {
275    let key = cmd.args.string("key")?;
276    let bucket = cmd.args.string("bucket")?;
277    let var_name = cmd.args.string("var-name")?;
278    cmd.args.done()?;
279
280    let aws_client = mz_aws_util::s3::new_client(&state.aws_config);
281    let presign_config = mz_aws_util::s3::new_presigned_config();
282    let request = aws_client
283        .get_object()
284        .bucket(&bucket)
285        .key(&key)
286        .presigned(presign_config)
287        .await
288        .context("s3 presign")?;
289
290    println!("Setting '{var_name}' to presigned URL for {bucket}/{key}");
291    state.cmd_vars.insert(var_name, request.uri().to_string());
292
293    Ok(ControlFlow::Continue)
294}
295
296/// Generates parquet files covering a wide range of Arrow types and uploads them to S3 with
297/// multiple compression variants. This is the Rust equivalent of the Python
298/// `generate_parquet_files()` function.
299///
300/// Uploads:
301/// - `{key-prefix}` (uncompressed)
302/// - `{key-prefix}.snappy`
303/// - `{key-prefix}.gzip`
304/// - `{key-prefix}.brotli`
305/// - `{key-prefix}.zstd`
306/// - `{key-prefix}.lz4`
307pub async fn run_upload_parquet_types(
308    mut cmd: BuiltinCommand,
309    state: &State,
310) -> Result<ControlFlow, anyhow::Error> {
311    let bucket = cmd.args.string("bucket")?;
312    let key_prefix = cmd.args.string("key-prefix")?;
313    cmd.args.done()?;
314
315    let batch = build_parquet_types_batch().context("building parquet types batch")?;
316
317    let compressions = vec![
318        ("".to_string(), ParquetCompression::UNCOMPRESSED),
319        (".snappy".to_string(), ParquetCompression::SNAPPY),
320        (
321            ".gzip".to_string(),
322            ParquetCompression::GZIP(GzipLevel::default()),
323        ),
324        (
325            ".brotli".to_string(),
326            ParquetCompression::BROTLI(BrotliLevel::default()),
327        ),
328        (
329            ".zstd".to_string(),
330            ParquetCompression::ZSTD(ZstdLevel::default()),
331        ),
332        (".lz4".to_string(), ParquetCompression::LZ4_RAW),
333    ];
334
335    let client = mz_aws_util::s3::new_client(&state.aws_config);
336
337    for (suffix, compression) in compressions {
338        let key = format!("{key_prefix}{suffix}");
339        println!("Uploading parquet types file to S3 bucket {bucket}/{key}");
340
341        let props = WriterProperties::builder()
342            .set_compression(compression)
343            .build();
344        let mut buf = Vec::new();
345        {
346            let mut writer = ArrowWriter::try_new(&mut buf, batch.schema(), Some(props))
347                .context("creating parquet writer")?;
348            writer.write(&batch).context("writing parquet batch")?;
349            writer.close().context("closing parquet writer")?;
350        }
351
352        client
353            .put_object()
354            .bucket(&bucket)
355            .key(&key)
356            .body(buf.into())
357            .send()
358            .await
359            .context("s3 put")?;
360    }
361
362    Ok(ControlFlow::Continue)
363}
364
365// Using `as ArrayRef` is necessary when creating the struct array because the inner arrays have different types.
366// The metadata for Arrow `Field` type requires `std::collections::HashMap`, which is disallowed.
367#[allow(clippy::as_conversions, clippy::disallowed_types)]
368fn build_parquet_types_batch() -> Result<RecordBatch, anyhow::Error> {
369    let epoch = NaiveDate::from_ymd_opt(1970, 1, 1).unwrap();
370
371    // date32: days since epoch
372    let date_values: Vec<i32> = [
373        NaiveDate::from_ymd_opt(2025, 11, 1).unwrap(),
374        NaiveDate::from_ymd_opt(2025, 11, 2).unwrap(),
375        NaiveDate::from_ymd_opt(2025, 11, 3).unwrap(),
376    ]
377    .into_iter()
378    .map(|d| d.signed_duration_since(epoch).num_days() as i32)
379    .collect();
380
381    // timestamp(ms): ms since epoch (no timezone)
382    let datetime_values: Vec<i64> = [
383        NaiveDateTime::new(
384            NaiveDate::from_ymd_opt(2025, 11, 1).unwrap(),
385            NaiveTime::from_hms_opt(10, 0, 0).unwrap(),
386        ),
387        NaiveDateTime::new(
388            NaiveDate::from_ymd_opt(2025, 11, 1).unwrap(),
389            NaiveTime::from_hms_opt(11, 30, 0).unwrap(),
390        ),
391        NaiveDateTime::new(
392            NaiveDate::from_ymd_opt(2025, 11, 1).unwrap(),
393            NaiveTime::from_hms_opt(12, 0, 0).unwrap(),
394        ),
395    ]
396    .into_iter()
397    .map(|dt| dt.and_utc().timestamp_millis())
398    .collect();
399
400    // time32(s): seconds since midnight
401    let time_values: Vec<i32> = [
402        NaiveTime::from_hms_opt(9, 0, 0).unwrap(),
403        NaiveTime::from_hms_opt(10, 30, 15).unwrap(),
404        NaiveTime::from_hms_opt(11, 45, 30).unwrap(),
405    ]
406    .into_iter()
407    .map(|t| t.num_seconds_from_midnight() as i32)
408    .collect();
409
410    // list<int64>: [-1, 2], [3, 4, 5], []
411    let mut list_builder = ListBuilder::new(Int64Builder::new());
412    for &val in &[-1i64, 2] {
413        list_builder.values().append_value(val);
414    }
415    list_builder.append(true);
416    for &val in &[3i64, 4, 5] {
417        list_builder.values().append_value(val);
418    }
419    list_builder.append(true);
420    list_builder.append(true); // empty list
421    let list_array = Arc::new(list_builder.finish());
422
423    // decimal128(precision=10, scale=5): -54.321, 123.45, null
424    let decimal_array = Arc::new(
425        Decimal128Array::from(vec![Some(-5_432_100i128), Some(12_345_000i128), None])
426            .with_precision_and_scale(10, 5)
427            .context("setting decimal precision/scale")?,
428    );
429
430    // struct/record: (name text, age int32, avg float64)
431    let struct_array = Arc::new(StructArray::from(vec![
432        (
433            Arc::new(Field::new("name", DataType::Utf8, true)),
434            Arc::new(StringArray::from(vec!["Taco", "Burger", "SlimJim"])) as ArrayRef,
435        ),
436        (
437            Arc::new(Field::new("age", DataType::Int32, true)),
438            Arc::new(Int32Array::from(vec![3, 2, 1])) as ArrayRef,
439        ),
440        (
441            Arc::new(Field::new("avg", DataType::Float64, true)),
442            Arc::new(Float64Array::from(vec![2.2, 4.5, 1.14])) as ArrayRef,
443        ),
444    ]));
445
446    // uuid: FixedSizeBinary(16) with arrow.uuid extension metadata
447    let mut uuid_builder = FixedSizeBinaryBuilder::with_capacity(3, 16);
448    for uuid_str in &[
449        "badc0deb-adc0-deba-dc0d-ebadc0debadc",
450        "deadbeef-dead-4eef-8eef-deaddeadbeef",
451        "00000000-0000-0000-0000-000000000000",
452    ] {
453        let uuid_val = uuid::Uuid::parse_str(uuid_str).context("parsing uuid")?;
454        uuid_builder
455            .append_value(uuid_val.as_bytes())
456            .context("appending uuid bytes")?;
457    }
458    let uuid_array = Arc::new(uuid_builder.finish());
459
460    // variable-length binary
461    let mut binary_builder = BinaryBuilder::new();
462    binary_builder.append_value(b"raw1");
463    binary_builder.append_value(b"raw2");
464    binary_builder.append_value(b"raw3");
465    let binary_array = Arc::new(binary_builder.finish());
466
467    // interval(year-month): stored as i32 months
468    let interval_ym_array = Arc::new(IntervalYearMonthArray::from(vec![1i32, 13, -2]));
469
470    // interval(day-time): stored as (days: i32, milliseconds: i32)
471    let interval_dt_array = Arc::new(IntervalDayTimeArray::from(vec![
472        IntervalDayTime {
473            days: 1,
474            milliseconds: 500,
475        },
476        IntervalDayTime {
477            days: 30,
478            milliseconds: 0,
479        },
480        IntervalDayTime {
481            days: -1,
482            milliseconds: -1000,
483        },
484    ]));
485
486    // map<utf8, utf8>: {"k1": "v1", "k2": "v2"}, {"k3": "v3"}, {}
487    let mut map_builder = MapBuilder::new(None, StringBuilder::new(), StringBuilder::new());
488    map_builder.keys().append_value("k1");
489    map_builder.values().append_value("v1");
490    map_builder.keys().append_value("k2");
491    map_builder.values().append_value("v2");
492    map_builder.append(true).context("appending map row 0")?;
493    map_builder.keys().append_value("k3");
494    map_builder.values().append_value("v3");
495    map_builder.append(true).context("appending map row 1")?;
496    map_builder.append(true).context("appending map row 2")?; // empty map
497    let map_array = Arc::new(map_builder.finish());
498
499    let mut uuid_metadata = HashMap::new();
500    uuid_metadata.insert("ARROW:extension:name".to_string(), "arrow.uuid".to_string());
501
502    let schema = Arc::new(Schema::new(vec![
503        Field::new("int8_col", DataType::Int8, true),
504        Field::new("uint8_col", DataType::UInt8, true),
505        Field::new("int16_col", DataType::Int16, true),
506        Field::new("uint16_col", DataType::UInt16, true),
507        Field::new("int32_col", DataType::Int32, true),
508        Field::new("uint32_col", DataType::UInt32, true),
509        Field::new("int64_col", DataType::Int64, true),
510        Field::new("uint64_col", DataType::UInt64, true),
511        Field::new("float32_col", DataType::Float32, true),
512        Field::new("float64_col", DataType::Float64, true),
513        Field::new("bool_col", DataType::Boolean, true),
514        Field::new("string_col", DataType::Utf8, true),
515        Field::new("binary_col", DataType::Binary, true),
516        Field::new("date32_col", DataType::Date32, true),
517        Field::new(
518            "timestamp_ms_col",
519            DataType::Timestamp(TimeUnit::Millisecond, None),
520            true,
521        ),
522        Field::new("time32_col", DataType::Time32(TimeUnit::Second), true),
523        Field::new(
524            "list_col",
525            DataType::List(Arc::new(Field::new("item", DataType::Int64, true))),
526            true,
527        ),
528        Field::new("decimal_col", DataType::Decimal128(10, 5), true),
529        Field::new("json_col", DataType::Utf8, true),
530        Field::new(
531            "record_col",
532            DataType::Struct(
533                vec![
534                    Field::new("name", DataType::Utf8, true),
535                    Field::new("age", DataType::Int32, true),
536                    Field::new("avg", DataType::Float64, true),
537                ]
538                .into(),
539            ),
540            true,
541        ),
542        Field::new("uuid_col", DataType::FixedSizeBinary(16), false).with_metadata(uuid_metadata),
543        Field::new(
544            "interval_ym_col",
545            DataType::Interval(IntervalUnit::YearMonth),
546            true,
547        ),
548        Field::new(
549            "interval_dt_col",
550            DataType::Interval(IntervalUnit::DayTime),
551            true,
552        ),
553        Field::new(
554            "map_col",
555            DataType::Map(
556                Arc::new(Field::new(
557                    "entries",
558                    DataType::Struct(
559                        vec![
560                            Field::new("keys", DataType::Utf8, false),
561                            Field::new("values", DataType::Utf8, true),
562                        ]
563                        .into(),
564                    ),
565                    false,
566                )),
567                false,
568            ),
569            true,
570        ),
571    ]));
572
573    let batch = RecordBatch::try_new(
574        schema,
575        vec![
576            Arc::new(Int8Array::from(vec![-1i8, 2, 3])),
577            Arc::new(UInt8Array::from(vec![10u8, 20, 30])),
578            Arc::new(Int16Array::from(vec![-1000i16, 2000, 3000])),
579            Arc::new(UInt16Array::from(vec![10000u16, 20000, 30000])),
580            Arc::new(Int32Array::from(vec![-100000i32, 200000, 300000])),
581            Arc::new(UInt32Array::from(vec![1000000u32, 2000000, 3000000])),
582            Arc::new(Int64Array::from(vec![
583                -1_000_000_000i64,
584                2_000_000_000,
585                3_000_000_000,
586            ])),
587            Arc::new(UInt64Array::from(vec![
588                1_000_000_000_000_000_000u64,
589                2_000_000_000_000_000_000,
590                3_000_000_000_000_000_000,
591            ])),
592            Arc::new(Float32Array::from(vec![-1.0f32, 2.5, 3.7])),
593            Arc::new(Float64Array::from(vec![-1.0f64, 2.5, 3.7])),
594            Arc::new(BooleanArray::from(vec![true, false, true])),
595            Arc::new(StringArray::from(vec!["apple", "banana", "cherry"])),
596            binary_array,
597            Arc::new(Date32Array::from(date_values)),
598            Arc::new(TimestampMillisecondArray::from(datetime_values)),
599            Arc::new(Time32SecondArray::from(time_values)),
600            list_array,
601            decimal_array,
602            Arc::new(StringArray::from(vec![
603                r#"{"a": 5, "b": { "c": 1.1 } }"#,
604                r#"{ "d": "str", "e" : [1,2,3] }"#,
605                "{}",
606            ])),
607            struct_array,
608            uuid_array,
609            interval_ym_array,
610            interval_dt_array,
611            map_array,
612        ],
613    )
614    .context("building record batch")?;
615
616    Ok(batch)
617}
618
619/// Upload a parquet file with map columns whose keys are deliberately unsorted.
620/// This is used to test that COPY FROM correctly handles unsorted map keys.
621pub async fn run_upload_parquet_unsorted_map(
622    mut cmd: BuiltinCommand,
623    state: &State,
624) -> Result<ControlFlow, anyhow::Error> {
625    let bucket = cmd.args.string("bucket")?;
626    let key = cmd.args.string("key")?;
627    cmd.args.done()?;
628
629    let batch =
630        build_parquet_unsorted_map_batch().context("building parquet unsorted map batch")?;
631
632    let client = mz_aws_util::s3::new_client(&state.aws_config);
633
634    println!("Uploading parquet unsorted-map file to S3 bucket {bucket}/{key}");
635
636    let props = WriterProperties::builder().build();
637    let mut buf = Vec::new();
638    {
639        let mut writer = ArrowWriter::try_new(&mut buf, batch.schema(), Some(props))
640            .context("creating parquet writer")?;
641        writer.write(&batch).context("writing parquet batch")?;
642        writer.close().context("closing parquet writer")?;
643    }
644
645    client
646        .put_object()
647        .bucket(&bucket)
648        .key(&key)
649        .body(buf.into())
650        .send()
651        .await
652        .context("s3 put")?;
653
654    Ok(ControlFlow::Continue)
655}
656
657/// Build a parquet record batch with map columns whose keys are deliberately
658/// unsorted, to test that the COPY FROM reader sorts them correctly.
659#[allow(clippy::as_conversions, clippy::disallowed_types)]
660fn build_parquet_unsorted_map_batch() -> Result<RecordBatch, anyhow::Error> {
661    // Row 0: id=1, map keys in reverse order: {"z": "val_z", "a": "val_a", "m": "val_m"}
662    // Row 1: id=2, map keys unsorted: {"b": "val_b", "a": "val_a"}
663    // Row 2: id=3, single key (trivially sorted): {"x": "val_x"}
664    // Row 3: id=4, empty map: {}
665
666    let mut map_builder = MapBuilder::new(None, StringBuilder::new(), StringBuilder::new());
667
668    // Row 0: keys deliberately in reverse order
669    map_builder.keys().append_value("z");
670    map_builder.values().append_value("val_z");
671    map_builder.keys().append_value("a");
672    map_builder.values().append_value("val_a");
673    map_builder.keys().append_value("m");
674    map_builder.values().append_value("val_m");
675    map_builder.append(true).context("appending map row 0")?;
676
677    // Row 1: keys unsorted
678    map_builder.keys().append_value("b");
679    map_builder.values().append_value("val_b");
680    map_builder.keys().append_value("a");
681    map_builder.values().append_value("val_a");
682    map_builder.append(true).context("appending map row 1")?;
683
684    // Row 2: single key (trivially sorted)
685    map_builder.keys().append_value("x");
686    map_builder.values().append_value("val_x");
687    map_builder.append(true).context("appending map row 2")?;
688
689    // Row 3: empty map
690    map_builder.append(true).context("appending map row 3")?;
691
692    // Row 4: duplicate keys - should be deduped and the last one chosen by the reader
693    map_builder.keys().append_value("y");
694    map_builder.values().append_value("val_y");
695    map_builder.keys().append_value("y");
696    map_builder.values().append_value("val_y2");
697    map_builder.keys().append_value("y");
698    map_builder.values().append_value("val_y3");
699    map_builder.append(true).context("appending map row 4")?;
700
701    let map_array = Arc::new(map_builder.finish());
702
703    let schema = Arc::new(Schema::new(vec![
704        Field::new("id", DataType::Int32, false),
705        Field::new(
706            "map_col",
707            DataType::Map(
708                Arc::new(Field::new(
709                    "entries",
710                    DataType::Struct(
711                        vec![
712                            Field::new("keys", DataType::Utf8, false),
713                            Field::new("values", DataType::Utf8, true),
714                        ]
715                        .into(),
716                    ),
717                    false,
718                )),
719                false, // keys_sorted = false
720            ),
721            true,
722        ),
723    ]));
724
725    let batch = RecordBatch::try_new(
726        schema,
727        vec![Arc::new(Int32Array::from(vec![1, 2, 3, 4, 5])), map_array],
728    )
729    .context("building record batch")?;
730
731    Ok(batch)
732}