1#[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 .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 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 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 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 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 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 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 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 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
296pub 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#[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 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 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 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 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); let list_array = Arc::new(list_builder.finish());
422
423 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 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 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 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 let interval_ym_array = Arc::new(IntervalYearMonthArray::from(vec![1i32, 13, -2]));
469
470 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 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")?; 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
619pub 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#[allow(clippy::as_conversions, clippy::disallowed_types)]
660fn build_parquet_unsorted_map_batch() -> Result<RecordBatch, anyhow::Error> {
661 let mut map_builder = MapBuilder::new(None, StringBuilder::new(), StringBuilder::new());
667
668 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 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 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 map_builder.append(true).context("appending map row 3")?;
691
692 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, ),
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}