1use std::fs::{self, File, OpenOptions};
4use std::io::{self, BufReader, BufWriter, Write};
5use std::os::unix::fs::OpenOptionsExt;
6use std::path::{Path, PathBuf};
7
8use ahri_tre_tabular::{
9 AhriTreArrowMetadata, AnnotatedArrowTable, ArrowIpcCompression, ArrowIpcWriteOptions,
10 ArrowTable, annotate_arrow_table, write_ipc_file_to_writer, write_parquet_file,
11};
12use ahri_tre_types::{NcName, StudyId};
13use uuid::Uuid;
14
15use crate::{DuckLakeAdapter, DuckLakeOpenedCatalog, LakeError, ScratchAttempt};
16
17#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub enum DatasetFileFormat {
20 ArrowIpc,
21 Csv,
22 Json,
23 Parquet,
24 Xlsx,
25}
26
27impl DatasetFileFormat {
28 pub fn extension(self) -> &'static str {
30 match self {
31 Self::ArrowIpc => "arrow",
32 Self::Csv => "csv",
33 Self::Json => "json",
34 Self::Parquet => "parquet",
35 Self::Xlsx => "xlsx",
36 }
37 }
38}
39
40#[derive(Debug, Clone, PartialEq, Eq)]
42pub struct DatasetFileReadOptions {
43 pub header: bool,
44 pub delimiter: char,
45 pub null_strings: Vec<String>,
46 pub json_format: String,
47 pub sheet: Option<String>,
48}
49
50impl Default for DatasetFileReadOptions {
51 fn default() -> Self {
52 Self {
53 header: true,
54 delimiter: ',',
55 null_strings: Vec::new(),
56 json_format: "auto".to_string(),
57 sheet: None,
58 }
59 }
60}
61
62#[derive(Debug, Clone, PartialEq, Eq)]
64pub struct LoadDatasetTableFromFileRequest {
65 pub study_id: StudyId,
66 pub dataset_name: NcName,
67 pub major: i32,
68 pub minor: i32,
69 pub patch: i32,
70 pub source_path: PathBuf,
71 pub format: DatasetFileFormat,
72 pub options: DatasetFileReadOptions,
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
77pub struct DatasetCsvSeed<'a> {
78 pub dataset_name: NcName,
79 pub contents: &'a [u8],
80}
81
82#[derive(Debug, Clone, PartialEq, Eq)]
84pub struct LoadDatasetTableFromSqlRequest {
85 pub study_id: StudyId,
86 pub dataset_name: NcName,
87 pub major: i32,
88 pub minor: i32,
89 pub patch: i32,
90 pub source_sql: String,
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq)]
95pub enum DatasetExportFormat {
96 Csv,
97 Ndjson,
98 Parquet,
99 ArrowIpc,
100}
101
102#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct ExportDatasetTableRequest {
105 pub study_id: StudyId,
106 pub dataset_name: NcName,
107 pub major: i32,
108 pub minor: i32,
109 pub patch: i32,
110 pub limit: Option<usize>,
111 pub destination_path: PathBuf,
112 pub format: DatasetExportFormat,
113 pub compress: bool,
114 pub arrow_metadata: Option<AhriTreArrowMetadata>,
115}
116
117#[derive(Debug, Clone, PartialEq, Eq)]
119pub struct PreviewDatasetTableRequest {
120 pub study_id: StudyId,
121 pub dataset_name: NcName,
122 pub major: i32,
123 pub minor: i32,
124 pub patch: i32,
125 pub limit: usize,
126}
127
128#[derive(Debug, Clone, PartialEq, Eq)]
130pub struct ExportedDatasetTable {
131 pub relation: DatasetTableRelation,
132 pub destination_path: PathBuf,
133 pub format: DatasetExportFormat,
134 pub compressed: bool,
135 pub row_count: usize,
136 pub column_count: usize,
137 pub exported_size_bytes: u64,
138 pub warnings: Vec<String>,
139}
140
141#[derive(Debug, Clone, PartialEq, Eq)]
143pub struct PreviewedDatasetTable {
144 pub relation: DatasetTableRelation,
145 pub columns: Vec<ImportedColumn>,
146 pub rows: Vec<Vec<Option<String>>>,
147 pub displayed_row_count: usize,
148 pub total_row_count: usize,
149 pub truncated: bool,
150}
151
152#[derive(Debug, Clone, PartialEq, Eq)]
154pub struct LoadRedcapFormDatasetRequest {
155 pub study_id: StudyId,
156 pub dataset_name: NcName,
157 pub major: i32,
158 pub minor: i32,
159 pub patch: i32,
160 pub source_eav_path: PathBuf,
161 pub fields: Vec<RedcapEavField>,
162 pub synthetic_columns: Vec<RedcapSyntheticColumn>,
163}
164
165#[derive(Debug, Clone, PartialEq, Eq)]
167pub struct RedcapEavField {
168 pub name: String,
169 pub value_type_id: ahri_tre_types::ValueTypeId,
170 pub value_format: Option<String>,
171 pub choices: Vec<RedcapEavChoice>,
172}
173
174#[derive(Debug, Clone, PartialEq, Eq)]
176pub struct RedcapEavChoice {
177 pub redcap_value: String,
178 pub value: i32,
179}
180
181#[derive(Debug, Clone, PartialEq, Eq)]
183pub struct RedcapSyntheticColumn {
184 pub source_column: String,
185 pub output_column: String,
186 pub value_type_id: ahri_tre_types::ValueTypeId,
187 pub choices: Vec<RedcapEavChoice>,
188}
189
190#[derive(Debug, Clone, PartialEq, Eq)]
192pub struct DatasetTableRelation {
193 pub schema_name: String,
194 pub table_name: String,
195}
196
197impl DatasetTableRelation {
198 pub fn qualified_name(&self) -> String {
200 format!("{}.{}", self.schema_name, self.table_name)
201 }
202
203 fn to_sql(&self) -> String {
204 format!(
205 "{}.{}",
206 quote_identifier(&self.schema_name),
207 quote_identifier(&self.table_name)
208 )
209 }
210}
211
212#[derive(Debug)]
214pub struct LoadedDatasetTable {
215 pub relation: DatasetTableRelation,
216 pub row_count: usize,
217 pub column_count: usize,
218 pub column_names: Vec<String>,
219 pub columns: Vec<ImportedColumn>,
220}
221
222#[derive(Debug, Clone, PartialEq, Eq)]
224pub struct ImportedColumn {
225 pub name: String,
226 pub duckdb_type: String,
227}
228
229impl DuckLakeAdapter {
230 pub fn validate_dataset_file(
238 &self,
239 connection: &duckdb::Connection,
240 source_path: &Path,
241 format: DatasetFileFormat,
242 options: &DatasetFileReadOptions,
243 ) -> Result<(), LakeError> {
244 if !source_path.is_file() {
245 return Err(LakeError::InvalidFileName {
246 path: source_path.display().to_string(),
247 });
248 }
249 if format == DatasetFileFormat::ArrowIpc {
250 return Err(LakeError::InvalidMetadataValue {
251 key: "format",
252 value: "Arrow IPC validation requires adapter-managed conversion".to_string(),
253 });
254 }
255
256 let relation_sql = file_relation_sql(source_path, format, options)?;
257 connection.query_row(
258 &format!("SELECT count(*) FROM ({relation_sql}) AS validated_dataset_file"),
259 [],
260 |_| Ok(()),
261 )?;
262 Ok(())
263 }
264
265 pub fn replace_server_managed_dataset_tables_from_csv(
274 &self,
275 opened: &DuckLakeOpenedCatalog,
276 scratch: &ScratchAttempt,
277 study_id: StudyId,
278 seeds: &[DatasetCsvSeed<'_>],
279 authority: &mut dyn ahri_tre_core::DatasetMaintenanceAuthority,
280 ) -> Result<Vec<LoadedDatasetTable>, LakeError> {
281 authority.ensure_exclusive()?;
282 if opened.attach_plan.create_if_not_exists
283 || opened.attach_plan.attach_description.contains("ENCRYPTED")
284 || opened.attach_plan.configured_encryption_mode
285 != ahri_tre_types::EncryptionMode::ServerManaged
286 || opened.health.detected_encryption_mode
287 != ahri_tre_types::EncryptionMode::ServerManaged
288 {
289 return Err(LakeError::InvalidMetadataValue {
290 key: "dataset_seed_catalog",
291 value: "existing server-managed catalog required".to_string(),
292 });
293 }
294 self.replace_dataset_tables_from_csv(&opened.connection, scratch, study_id, seeds)
295 }
296
297 fn replace_dataset_tables_from_csv(
298 &self,
299 connection: &duckdb::Connection,
300 scratch: &ScratchAttempt,
301 study_id: StudyId,
302 seeds: &[DatasetCsvSeed<'_>],
303 ) -> Result<Vec<LoadedDatasetTable>, LakeError> {
304 let retained = seeds
305 .iter()
306 .map(|seed| dataset_table_relation(study_id, seed.dataset_name.as_str(), 1, 0, 0))
307 .collect::<Vec<_>>();
308 if seeds.is_empty()
309 || seeds.iter().any(|seed| seed.contents.is_empty())
310 || retained
311 .iter()
312 .enumerate()
313 .any(|(index, relation)| retained[..index].contains(relation))
314 {
315 return Err(LakeError::InvalidMetadataValue {
316 key: "dataset_seed",
317 value: "non-empty unique dataset inputs required".to_string(),
318 });
319 }
320
321 connection.execute_batch("BEGIN TRANSACTION;")?;
322 let replaced = (|| {
323 self.retain_only_dataset_tables(connection, study_id, &[])?;
326 let mut loaded = Vec::with_capacity(seeds.len());
327 for seed in seeds {
328 let relation =
329 dataset_table_relation(study_id, seed.dataset_name.as_str(), 1, 0, 0);
330 let source_path = scratch.path().join(format!("{}.csv", relation.table_name));
331 let mut file = OpenOptions::new()
332 .write(true)
333 .create_new(true)
334 .mode(0o600)
335 .open(&source_path)?;
336 file.write_all(seed.contents)?;
337 file.sync_all()?;
338 loaded.push(self.load_dataset_table_from_file_with_creation_intent(
339 connection,
340 &LoadDatasetTableFromFileRequest {
341 study_id,
342 dataset_name: seed.dataset_name.clone(),
343 major: 1,
344 minor: 0,
345 patch: 0,
346 source_path,
347 format: DatasetFileFormat::Csv,
348 options: DatasetFileReadOptions::default(),
349 },
350 || Ok(()),
351 )?);
352 }
353
354 let mut actual = study_dataset_tables(connection, &study_lake_schema(study_id))?;
355 let mut expected = retained;
356 actual.sort_by(|left, right| left.table_name.cmp(&right.table_name));
357 expected.sort_by(|left, right| left.table_name.cmp(&right.table_name));
358 if actual != expected {
359 return Err(LakeError::InvalidMetadataValue {
360 key: "dataset_seed",
361 value: "exact dataset replacement verification failed".to_string(),
362 });
363 }
364 Ok(loaded)
365 })();
366
367 match replaced {
368 Ok(loaded) => {
369 if let Err(error) = connection.execute_batch("COMMIT;") {
370 let _ = connection.execute_batch("ROLLBACK;");
371 return Err(error.into());
372 }
373 Ok(loaded)
374 }
375 Err(error) => {
376 let _ = connection.execute_batch("ROLLBACK;");
377 Err(error)
378 }
379 }
380 }
381
382 fn retain_only_dataset_tables(
389 &self,
390 connection: &duckdb::Connection,
391 study_id: StudyId,
392 retained: &[DatasetTableRelation],
393 ) -> Result<Vec<DatasetTableRelation>, LakeError> {
394 let schema_name = study_lake_schema(study_id);
395 if retained
396 .iter()
397 .any(|relation| relation.schema_name != schema_name)
398 {
399 return Err(LakeError::InvalidMetadataValue {
400 key: "retained_dataset_relation",
401 value: "relation belongs to another Study namespace".to_string(),
402 });
403 }
404
405 let existing = study_dataset_tables(connection, &schema_name)?;
406 let mut removed = Vec::new();
407 for relation in existing {
408 if !retained.contains(&relation) {
409 self.drop_dataset_table(connection, &relation)?;
410 removed.push(relation);
411 }
412 }
413 Ok(removed)
414 }
415
416 #[cfg(test)]
424 fn load_dataset_table_from_file(
425 &self,
426 connection: &duckdb::Connection,
427 request: &LoadDatasetTableFromFileRequest,
428 ) -> Result<LoadedDatasetTable, LakeError> {
429 self.load_dataset_table_from_file_with_creation_intent(connection, request, || Ok(()))
430 }
431
432 pub fn load_dataset_table_from_file_reserved(
434 &self,
435 connection: &duckdb::Connection,
436 authority: &mut dyn ahri_tre_core::DatasetWriteAuthority,
437 request: &LoadDatasetTableFromFileRequest,
438 ) -> Result<LoadedDatasetTable, LakeError> {
439 authorize_dataset_write(
440 authority,
441 request.study_id,
442 &request.dataset_name,
443 request.major,
444 request.minor,
445 request.patch,
446 )?;
447 self.load_dataset_table_from_file_with_creation_intent(connection, request, || {
448 authority.record_creation_intent().map_err(Into::into)
449 })
450 }
451
452 pub fn load_dataset_table_from_file_reserved_observed(
454 &self,
455 context: Option<ahri_tre_observability::CorrelationContext>,
456 connection: &duckdb::Connection,
457 authority: &mut dyn ahri_tre_core::DatasetWriteAuthority,
458 request: &LoadDatasetTableFromFileRequest,
459 ) -> Result<LoadedDatasetTable, LakeError> {
460 let span = context.map(|context| context.span(ahri_tre_observability::Stage::LakeRead));
461 let result = self.load_dataset_table_from_file_reserved(connection, authority, request);
462 let evidence = result.as_ref().map(|loaded| {
463 [
464 ahri_tre_observability::Measurement::Rows(loaded.row_count as u64),
465 ahri_tre_observability::Measurement::Columns(loaded.column_count as u64),
466 ]
467 });
468 crate::error::finish_dataset_observation(
469 span,
470 &result,
471 evidence
472 .as_ref()
473 .map(|counts| counts.as_slice())
474 .unwrap_or_default(),
475 );
476 result
477 }
478
479 fn load_dataset_table_from_file_with_creation_intent(
480 &self,
481 connection: &duckdb::Connection,
482 request: &LoadDatasetTableFromFileRequest,
483 mut before_create: impl FnMut() -> Result<(), LakeError>,
484 ) -> Result<LoadedDatasetTable, LakeError> {
485 if !request.source_path.is_file() {
486 return Err(LakeError::InvalidFileName {
487 path: request.source_path.display().to_string(),
488 });
489 }
490 validate_version_numbers(request.major, request.minor, request.patch)?;
491
492 let relation = dataset_table_relation(
493 request.study_id,
494 request.dataset_name.as_str(),
495 request.major,
496 request.minor,
497 request.patch,
498 );
499 let stage_table = dataset_stage_table_name();
500 let stage_sql = quote_identifier(&stage_table);
501 let converted_source_path = (request.format == DatasetFileFormat::ArrowIpc)
502 .then(|| temporary_arrow_ipc_parquet_path(&request.source_path));
503 let target_sql = relation.to_sql();
504 let mut target_created = false;
505 let result = (|| {
506 let source_path = if let Some(path) = &converted_source_path {
507 arrow_ipc_file_to_parquet(&request.source_path, path)?;
508 path.as_path()
509 } else {
510 request.source_path.as_path()
511 };
512 let source_format = if request.format == DatasetFileFormat::ArrowIpc {
513 DatasetFileFormat::Parquet
514 } else {
515 request.format
516 };
517 let relation_sql = file_relation_sql(source_path, source_format, &request.options)?;
518 connection.execute_batch(&format!("DROP TABLE IF EXISTS {stage_sql};"))?;
519 connection
520 .execute_batch(&format!("CREATE TEMP TABLE {stage_sql} AS {relation_sql};"))?;
521 connection.execute_batch(&format!(
522 "CREATE SCHEMA IF NOT EXISTS {};",
523 quote_identifier(&relation.schema_name)
524 ))?;
525 ensure_dataset_target_empty(connection, &relation)?;
526 before_create()?;
527 connection.execute_batch(&format!(
528 "CREATE TABLE {target_sql} AS SELECT * FROM {stage_sql};"
529 ))?;
530 target_created = true;
531
532 let row_count = table_row_count(connection, &target_sql)?;
533 let columns = table_columns(connection, &target_sql)?;
534 let column_names: Vec<String> =
535 columns.iter().map(|column| column.name.clone()).collect();
536 let column_count = column_names.len();
537 Ok::<_, LakeError>(LoadedDatasetTable {
538 relation,
539 row_count,
540 column_count,
541 column_names,
542 columns,
543 })
544 })();
545
546 finish_dataset_load(
547 connection,
548 Some(&stage_sql),
549 &target_sql,
550 target_created,
551 converted_source_path.as_deref(),
552 result,
553 )
554 }
555
556 #[cfg(test)]
564 fn load_dataset_table_from_sql(
565 &self,
566 connection: &duckdb::Connection,
567 request: &LoadDatasetTableFromSqlRequest,
568 ) -> Result<LoadedDatasetTable, LakeError> {
569 self.load_dataset_table_from_sql_with_creation_intent(connection, request, || Ok(()))
570 }
571
572 pub fn load_restricted_query_output_reserved(
576 &self,
577 connection: &duckdb::Connection,
578 authority: &mut dyn ahri_tre_core::DatasetWriteAuthority,
579 target: &LoadDatasetTableFromSqlRequest,
580 output: crate::PreparedQueryOutput,
581 ) -> Result<LoadedDatasetTable, LakeError> {
582 self.load_dataset_table_from_file_reserved(
583 connection,
584 authority,
585 &LoadDatasetTableFromFileRequest {
586 study_id: target.study_id,
587 dataset_name: target.dataset_name.clone(),
588 major: target.major,
589 minor: target.minor,
590 patch: target.patch,
591 source_path: output.path.clone(),
592 format: DatasetFileFormat::Parquet,
593 options: DatasetFileReadOptions::default(),
594 },
595 )
596 }
597
598 pub fn load_dataset_table_from_sql_reserved(
600 &self,
601 connection: &duckdb::Connection,
602 authority: &mut dyn ahri_tre_core::DatasetWriteAuthority,
603 request: &LoadDatasetTableFromSqlRequest,
604 ) -> Result<LoadedDatasetTable, LakeError> {
605 authorize_dataset_write(
606 authority,
607 request.study_id,
608 &request.dataset_name,
609 request.major,
610 request.minor,
611 request.patch,
612 )?;
613 self.load_dataset_table_from_sql_with_creation_intent(connection, request, || {
614 authority.record_creation_intent().map_err(Into::into)
615 })
616 }
617
618 fn load_dataset_table_from_sql_with_creation_intent(
619 &self,
620 connection: &duckdb::Connection,
621 request: &LoadDatasetTableFromSqlRequest,
622 mut before_create: impl FnMut() -> Result<(), LakeError>,
623 ) -> Result<LoadedDatasetTable, LakeError> {
624 validate_version_numbers(request.major, request.minor, request.patch)?;
625 let source_sql = trimmed_select_sql(&request.source_sql)?;
626
627 let relation = dataset_table_relation(
628 request.study_id,
629 request.dataset_name.as_str(),
630 request.major,
631 request.minor,
632 request.patch,
633 );
634 let stage_table = dataset_stage_table_name();
635 let stage_sql = quote_identifier(&stage_table);
636 let target_sql = relation.to_sql();
637
638 let mut target_created = false;
639 let result = (|| {
640 connection.execute_batch(&format!("DROP TABLE IF EXISTS {stage_sql};"))?;
641 connection.execute_batch(&format!(
642 "CREATE TEMP TABLE {stage_sql} AS SELECT * FROM ({source_sql}) src;"
643 ))?;
644 connection.execute_batch(&format!(
645 "CREATE SCHEMA IF NOT EXISTS {};",
646 quote_identifier(&relation.schema_name)
647 ))?;
648 ensure_dataset_target_empty(connection, &relation)?;
649 before_create()?;
650 connection.execute_batch(&format!(
651 "CREATE TABLE {target_sql} AS SELECT * FROM {stage_sql};"
652 ))?;
653 target_created = true;
654
655 let row_count = table_row_count(connection, &target_sql)?;
656 let columns = table_columns(connection, &target_sql)?;
657 let column_names: Vec<String> =
658 columns.iter().map(|column| column.name.clone()).collect();
659 let column_count = column_names.len();
660 Ok::<_, LakeError>(LoadedDatasetTable {
661 relation,
662 row_count,
663 column_count,
664 column_names,
665 columns,
666 })
667 })();
668
669 finish_dataset_load(
670 connection,
671 Some(&stage_sql),
672 &target_sql,
673 target_created,
674 None,
675 result,
676 )
677 }
678
679 #[cfg(test)]
687 fn load_redcap_form_dataset(
688 &self,
689 connection: &duckdb::Connection,
690 request: &LoadRedcapFormDatasetRequest,
691 ) -> Result<LoadedDatasetTable, LakeError> {
692 self.load_redcap_form_dataset_with_creation_intent(connection, request, || Ok(()))
693 }
694
695 pub fn load_redcap_form_dataset_reserved(
697 &self,
698 connection: &duckdb::Connection,
699 authority: &mut dyn ahri_tre_core::DatasetWriteAuthority,
700 request: &LoadRedcapFormDatasetRequest,
701 ) -> Result<LoadedDatasetTable, LakeError> {
702 authorize_dataset_write(
703 authority,
704 request.study_id,
705 &request.dataset_name,
706 request.major,
707 request.minor,
708 request.patch,
709 )?;
710 self.load_redcap_form_dataset_with_creation_intent(connection, request, || {
711 authority.record_creation_intent().map_err(Into::into)
712 })
713 }
714
715 fn load_redcap_form_dataset_with_creation_intent(
716 &self,
717 connection: &duckdb::Connection,
718 request: &LoadRedcapFormDatasetRequest,
719 mut before_create: impl FnMut() -> Result<(), LakeError>,
720 ) -> Result<LoadedDatasetTable, LakeError> {
721 if !request.source_eav_path.is_file() {
722 return Err(LakeError::InvalidFileName {
723 path: request.source_eav_path.display().to_string(),
724 });
725 }
726 validate_version_numbers(request.major, request.minor, request.patch)?;
727
728 let relation = dataset_table_relation(
729 request.study_id,
730 request.dataset_name.as_str(),
731 request.major,
732 request.minor,
733 request.patch,
734 );
735 let target_sql = relation.to_sql();
736 let eav_sql = sql_literal(&request.source_eav_path.display().to_string());
737 let field_list = sql_string_list(request.fields.iter().map(|field| field.name.as_str()));
738 ensure_redcap_eav_columns(
739 connection,
740 &request.source_eav_path,
741 &["field_name", "value"],
742 )?;
743 let has_event = csv_has_column(connection, &request.source_eav_path, "redcap_event_name")?;
744 let has_repeat_instrument = csv_has_column(
745 connection,
746 &request.source_eav_path,
747 "redcap_repeat_instrument",
748 )?;
749 let has_repeat_instance = csv_has_column(
750 connection,
751 &request.source_eav_path,
752 "redcap_repeat_instance",
753 )?;
754
755 let mut key_columns = vec!["record".to_string()];
756 if has_event {
757 key_columns.push("redcap_event_name".to_string());
758 }
759 if has_repeat_instrument {
760 key_columns.push("redcap_repeat_instrument".to_string());
761 }
762 if has_repeat_instance {
763 key_columns.push("redcap_repeat_instance".to_string());
764 }
765 let key_sql = key_columns
766 .iter()
767 .map(|column| quote_identifier(column))
768 .collect::<Vec<_>>()
769 .join(", ");
770
771 let projections = redcap_projection_sql(&request.fields, &request.synthetic_columns);
772 let mut target_created = false;
773 let result = (|| {
774 connection.execute_batch(&format!(
775 "CREATE SCHEMA IF NOT EXISTS {};",
776 quote_identifier(&relation.schema_name)
777 ))?;
778 ensure_dataset_target_empty(connection, &relation)?;
779 before_create()?;
780 connection.execute_batch(&format!(
781 r#"
782 CREATE TABLE {target_sql} AS
783 WITH src AS (
784 SELECT * FROM {}
785 WHERE field_name IN ({field_list})
786 ),
787 grouped AS (
788 SELECT {key_sql}, field_name, string_agg(CAST(value AS VARCHAR), ', ') AS value
789 FROM src
790 GROUP BY {key_sql}, field_name
791 ),
792 pivoted AS (
793 PIVOT grouped
794 ON field_name
795 USING any_value(value)
796 GROUP BY {key_sql}
797 )
798 SELECT {projections}
799 FROM pivoted
800 ORDER BY {key_sql};
801 "#,
802 redcap_eav_scan_sql(&eav_sql)
803 ))?;
804 target_created = true;
805
806 let row_count = table_row_count(connection, &target_sql)?;
807 let columns = table_columns(connection, &target_sql)?;
808 let column_names: Vec<String> =
809 columns.iter().map(|column| column.name.clone()).collect();
810 let column_count = column_names.len();
811 Ok::<_, LakeError>(LoadedDatasetTable {
812 relation,
813 row_count,
814 column_count,
815 column_names,
816 columns,
817 })
818 })();
819
820 finish_dataset_load(connection, None, &target_sql, target_created, None, result)
821 }
822
823 pub fn cleanup_reserved_dataset_output(
825 &self,
826 connection: &duckdb::Connection,
827 authority: &mut dyn ahri_tre_core::DatasetOutputAuthority,
828 ) -> Result<(), LakeError> {
829 authority.ensure_current()?;
830 let version = authority.output_version()?;
831 let relation = dataset_table_relation(
832 version.output.study_id,
833 &version.output.lake_name,
834 version.major,
835 version.minor,
836 version.patch,
837 );
838 self.drop_dataset_table(connection, &relation)
839 }
840
841 pub fn drop_dataset_table(
847 &self,
848 connection: &duckdb::Connection,
849 relation: &DatasetTableRelation,
850 ) -> Result<(), LakeError> {
851 connection.execute_batch(&format!("DROP TABLE IF EXISTS {};", relation.to_sql()))?;
852 Ok(())
853 }
854
855 pub fn export_dataset_table(
863 &self,
864 connection: &duckdb::Connection,
865 request: &ExportDatasetTableRequest,
866 ) -> Result<ExportedDatasetTable, LakeError> {
867 validate_version_numbers(request.major, request.minor, request.patch)?;
868 let relation = dataset_table_relation(
869 request.study_id,
870 request.dataset_name.as_str(),
871 request.major,
872 request.minor,
873 request.patch,
874 );
875 let source_sql = limited_relation_sql(&relation.to_sql(), request.limit);
876 let row_count = table_row_count(connection, &source_sql)?;
877 let column_names = table_column_names(connection, &source_sql)?;
878 let column_count = column_names.len();
879
880 if let Some(parent) = request.destination_path.parent() {
881 fs::create_dir_all(parent)?;
882 }
883
884 let mut warnings = Vec::new();
885 match request.format {
886 DatasetExportFormat::Csv => {
887 export_relation_with_copy(
888 connection,
889 &source_sql,
890 &request.destination_path,
891 "FORMAT CSV, HEADER TRUE",
892 request.compress,
893 )?;
894 }
895 DatasetExportFormat::Ndjson => {
896 export_relation_with_copy(
897 connection,
898 &source_sql,
899 &request.destination_path,
900 "FORMAT JSON, ARRAY FALSE",
901 request.compress,
902 )?;
903 }
904 DatasetExportFormat::Parquet => {
905 let exported = export_relation_to_arrow_table(
906 connection,
907 &source_sql,
908 request.arrow_metadata.as_ref(),
909 )?;
910 warnings.extend(exported.warnings);
911 write_parquet_file(&request.destination_path, &exported.table)?;
912 }
913 DatasetExportFormat::ArrowIpc => {
914 let exported = export_relation_to_arrow_table(
915 connection,
916 &source_sql,
917 request.arrow_metadata.as_ref(),
918 )?;
919 warnings.extend(exported.warnings);
920 write_ipc_file_to_writer(
921 BufWriter::new(File::create(&request.destination_path)?),
922 &exported.table,
923 ArrowIpcWriteOptions {
924 compression: Some(ArrowIpcCompression::Zstd),
925 },
926 )?;
927 }
928 }
929
930 Ok(ExportedDatasetTable {
931 relation,
932 destination_path: request.destination_path.clone(),
933 format: request.format,
934 compressed: request.compress
935 || matches!(
936 request.format,
937 DatasetExportFormat::Parquet | DatasetExportFormat::ArrowIpc
938 ),
939 row_count,
940 column_count,
941 exported_size_bytes: fs::metadata(&request.destination_path)?.len(),
942 warnings,
943 })
944 }
945
946 pub fn preview_dataset_table(
952 &self,
953 connection: &duckdb::Connection,
954 request: &PreviewDatasetTableRequest,
955 ) -> Result<PreviewedDatasetTable, LakeError> {
956 validate_version_numbers(request.major, request.minor, request.patch)?;
957 let relation = dataset_table_relation(
958 request.study_id,
959 request.dataset_name.as_str(),
960 request.major,
961 request.minor,
962 request.patch,
963 );
964 let relation_sql = relation.to_sql();
965 let total_row_count = table_row_count(connection, &relation_sql)?;
966 let columns = table_columns(connection, &relation_sql)?;
967 let rows = preview_rows(connection, &relation_sql, &columns, request.limit)?;
968 let displayed_row_count = rows.len();
969 Ok(PreviewedDatasetTable {
970 relation,
971 columns,
972 rows,
973 displayed_row_count,
974 total_row_count,
975 truncated: displayed_row_count < total_row_count,
976 })
977 }
978}
979
980fn study_dataset_tables(
981 connection: &duckdb::Connection,
982 schema_name: &str,
983) -> Result<Vec<DatasetTableRelation>, LakeError> {
984 let mut statement = connection.prepare(
985 "SELECT table_name
986 FROM information_schema.tables
987 WHERE table_catalog = current_database()
988 AND table_schema = ?
989 AND table_type = 'BASE TABLE'
990 ORDER BY table_name",
991 )?;
992 let rows = statement.query_map([schema_name], |row| row.get::<_, String>(0))?;
993 let mut relations = Vec::new();
994 for row in rows {
995 relations.push(DatasetTableRelation {
996 schema_name: schema_name.to_string(),
997 table_name: row?,
998 });
999 }
1000 Ok(relations)
1001}
1002
1003fn redcap_projection_sql(
1004 fields: &[RedcapEavField],
1005 synthetic_columns: &[RedcapSyntheticColumn],
1006) -> String {
1007 let mut projections = vec!["TRY_CAST(record AS INTEGER) AS record".to_string()];
1008 for synthetic in synthetic_columns {
1009 projections.push(redcap_synthetic_projection_sql(synthetic));
1010 }
1011 for field in fields {
1012 projections.push(redcap_field_projection_sql(field));
1013 }
1014 projections.join(", ")
1015}
1016
1017fn redcap_synthetic_projection_sql(column: &RedcapSyntheticColumn) -> String {
1018 let source = quote_identifier(&column.source_column);
1019 let target = quote_identifier(&column.output_column);
1020 match column.value_type_id.0 {
1021 1 => format!("TRY_CAST({source} AS INTEGER) AS {target}"),
1022 7 => redcap_category_case_sql(&source, &target, &column.choices),
1023 _ => format!("CAST({source} AS VARCHAR) AS {target}"),
1024 }
1025}
1026
1027fn redcap_field_projection_sql(field: &RedcapEavField) -> String {
1028 let column = quote_identifier(&field.name);
1029 match field.value_type_id.0 {
1030 1 => format!("TRY_CAST({column} AS INTEGER) AS {column}"),
1031 2 => format!("TRY_CAST({column} AS DOUBLE) AS {column}"),
1032 4 => redcap_temporal_projection_sql(field, &column, "DATE"),
1033 5 => redcap_temporal_projection_sql(field, &column, "TIMESTAMP"),
1034 6 => redcap_temporal_projection_sql(field, &column, "TIME"),
1035 7 => redcap_category_case_sql(&column, &column, &field.choices),
1036 8 => format!("CAST({column} AS VARCHAR) AS {column}"),
1037 _ => format!("CAST({column} AS VARCHAR) AS {column}"),
1038 }
1039}
1040
1041fn redcap_category_case_sql(source: &str, target: &str, choices: &[RedcapEavChoice]) -> String {
1042 let clauses = choices
1043 .iter()
1044 .map(|choice| {
1045 format!(
1046 "WHEN CAST({source} AS VARCHAR) = {} THEN {}",
1047 sql_literal(&choice.redcap_value),
1048 choice.value
1049 )
1050 })
1051 .collect::<Vec<_>>()
1052 .join(" ");
1053 format!("CASE WHEN {source} IS NULL THEN NULL {clauses} ELSE NULL END AS {target}")
1054}
1055
1056fn redcap_temporal_projection_sql(
1057 field: &RedcapEavField,
1058 column: &str,
1059 target_type: &str,
1060) -> String {
1061 if let Some(format) = field
1062 .value_format
1063 .as_deref()
1064 .map(str::trim)
1065 .filter(|value| !value.is_empty())
1066 {
1067 let parsed = format!("TRY_STRPTIME({column}, {})", sql_literal(format));
1068 if target_type == "TIMESTAMP" {
1069 format!("{parsed} AS {column}")
1070 } else {
1071 format!("CAST({parsed} AS {target_type}) AS {column}")
1072 }
1073 } else {
1074 format!("TRY_CAST({column} AS {target_type}) AS {column}")
1075 }
1076}
1077
1078fn csv_has_column(
1079 connection: &duckdb::Connection,
1080 path: &Path,
1081 column_name: &str,
1082) -> Result<bool, LakeError> {
1083 let columns = redcap_eav_columns(connection, path)?;
1084 Ok(columns.iter().any(|value| value == column_name))
1085}
1086
1087fn ensure_redcap_eav_columns(
1088 connection: &duckdb::Connection,
1089 path: &Path,
1090 required: &[&str],
1091) -> Result<(), LakeError> {
1092 let columns = redcap_eav_columns(connection, path)?;
1093 let missing = required
1094 .iter()
1095 .filter(|column| !columns.iter().any(|value| value == **column))
1096 .copied()
1097 .collect::<Vec<_>>();
1098 if !missing.is_empty() {
1099 return Err(LakeError::InvalidDataFileMetadata {
1100 field: "redcap_eav_columns",
1101 message: format!(
1102 "expected REDCap EAV CSV columns {}, got {}",
1103 missing.join(", "),
1104 if columns.is_empty() {
1105 "<none>".to_string()
1106 } else {
1107 columns.join(", ")
1108 }
1109 ),
1110 });
1111 }
1112 Ok(())
1113}
1114
1115fn redcap_eav_columns(
1116 connection: &duckdb::Connection,
1117 path: &Path,
1118) -> Result<Vec<String>, LakeError> {
1119 let path = sql_literal(&path.display().to_string());
1120 let mut statement = connection.prepare(&format!(
1121 "DESCRIBE SELECT * FROM {};",
1122 redcap_eav_scan_sql(&path)
1123 ))?;
1124 let rows = statement.query_map([], |row| row.get::<_, String>(0))?;
1125 let mut columns = Vec::new();
1126 for row in rows {
1127 columns.push(row?);
1128 }
1129 Ok(columns)
1130}
1131
1132fn redcap_eav_scan_sql(path: &str) -> String {
1133 format!("read_csv_auto({path}, header=true, delim=',', all_varchar=true)")
1134}
1135
1136fn sql_string_list<'a>(values: impl Iterator<Item = &'a str>) -> String {
1137 values.map(sql_literal).collect::<Vec<_>>().join(", ")
1138}
1139
1140pub fn dataset_output_key(
1143 study_id: StudyId,
1144 dataset_name: &ahri_tre_types::NcName,
1145) -> ahri_tre_core::DatasetOutputKey {
1146 ahri_tre_core::DatasetOutputKey {
1147 study_id,
1148 lake_name: strict_ncname(dataset_name.as_str()).to_ascii_lowercase(),
1149 }
1150}
1151
1152pub fn dataset_table_relation(
1154 study_id: StudyId,
1155 dataset_name: &str,
1156 major: i32,
1157 minor: i32,
1158 patch: i32,
1159) -> DatasetTableRelation {
1160 DatasetTableRelation {
1161 schema_name: study_lake_schema(study_id),
1162 table_name: format!(
1163 "{}{}",
1164 strict_ncname(dataset_name),
1165 version_ncname(major, minor, patch)
1166 ),
1167 }
1168}
1169
1170fn study_lake_schema(study_id: StudyId) -> String {
1171 format!(
1172 "study_{}",
1173 study_id
1174 .0
1175 .to_string()
1176 .to_ascii_lowercase()
1177 .replace('-', "_")
1178 )
1179}
1180
1181fn authorize_dataset_write(
1184 authority: &mut dyn ahri_tre_core::DatasetWriteAuthority,
1185 study_id: StudyId,
1186 dataset_name: &ahri_tre_types::NcName,
1187 major: i32,
1188 minor: i32,
1189 patch: i32,
1190) -> Result<(), LakeError> {
1191 let owned = authority.output_version()?;
1192 if owned.output != dataset_output_key(study_id, dataset_name)
1193 || (owned.major, owned.minor, owned.patch) != (major, minor, patch)
1194 {
1195 return Err(ahri_tre_core::CoreError::Conflict(
1196 "Dataset reservation does not own this output".into(),
1197 )
1198 .into());
1199 }
1200 authority.ensure_current().map_err(Into::into)
1201}
1202
1203fn ensure_dataset_target_empty(
1204 connection: &duckdb::Connection,
1205 relation: &DatasetTableRelation,
1206) -> Result<(), LakeError> {
1207 let occupied: bool = connection.query_row("SELECT EXISTS (SELECT 1 FROM information_schema.tables WHERE lower(table_schema) = lower(?) AND lower(table_name) = lower(?))",
1208 [&relation.schema_name, &relation.table_name], |row| row.get(0))?;
1209 if occupied {
1210 return Err(LakeError::DatasetOutputOccupied);
1211 }
1212 Ok(())
1213}
1214
1215fn finish_dataset_load(
1216 connection: &duckdb::Connection,
1217 stage_sql: Option<&str>,
1218 target_sql: &str,
1219 target_created: bool,
1220 converted_source: Option<&Path>,
1221 result: Result<LoadedDatasetTable, LakeError>,
1222) -> Result<LoadedDatasetTable, LakeError> {
1223 let mut cleanup_failed = false;
1224 if let Some(stage_sql) = stage_sql {
1225 cleanup_failed |= connection
1226 .execute_batch(&format!("DROP TABLE IF EXISTS {stage_sql};"))
1227 .is_err();
1228 }
1229 if let Some(path) = converted_source {
1230 cleanup_failed |=
1231 fs::remove_file(path).is_err_and(|error| error.kind() != std::io::ErrorKind::NotFound);
1232 }
1233 if target_created && (result.is_err() || cleanup_failed) {
1234 cleanup_failed |= connection
1235 .execute_batch(&format!("DROP TABLE IF EXISTS {target_sql};"))
1236 .is_err();
1237 }
1238 if cleanup_failed {
1239 Err(LakeError::PhysicalRollback {
1240 operation: "Dataset loading did not complete".into(),
1241 rollback: "Dataset staging or output cleanup could not be confirmed".into(),
1242 })
1243 } else {
1244 result
1245 }
1246}
1247
1248fn dataset_stage_table_name() -> String {
1249 format!("__tre_stage_{}", Uuid::new_v4().simple())
1250}
1251
1252fn validate_version_numbers(major: i32, minor: i32, patch: i32) -> Result<(), LakeError> {
1253 if major < 0 || minor < 0 || patch < 0 {
1254 return Err(LakeError::InvalidMetadataValue {
1255 key: "dataset_version",
1256 value: format!("{major}.{minor}.{patch}"),
1257 });
1258 }
1259 Ok(())
1260}
1261
1262fn version_ncname(major: i32, minor: i32, patch: i32) -> String {
1263 strict_ncname(&format!("{major}.{minor}.{patch}"))
1264}
1265
1266fn strict_ncname(value: &str) -> String {
1267 let trimmed = value.trim();
1268 if trimmed.is_empty() {
1269 return "_x".to_string();
1270 }
1271
1272 let mut output = String::new();
1273 for (index, character) in trimmed.chars().enumerate() {
1274 let valid = character.is_ascii_alphanumeric() || character == '_';
1275 if index == 0 && !(character.is_ascii_alphabetic() || character == '_') {
1276 output.push('_');
1277 }
1278 output.push(if valid { character } else { '_' });
1279 }
1280
1281 let mut collapsed = String::new();
1282 let mut previous_was_replacement = false;
1283 for character in output.chars() {
1284 if character == '_' {
1285 if !previous_was_replacement {
1286 collapsed.push(character);
1287 }
1288 previous_was_replacement = true;
1289 } else {
1290 collapsed.push(character);
1291 previous_was_replacement = false;
1292 }
1293 }
1294 if collapsed.to_ascii_lowercase().starts_with("xml") {
1295 collapsed.insert(0, '_');
1296 }
1297 collapsed
1298}
1299
1300pub(crate) fn file_relation_sql(
1301 source_path: &Path,
1302 format: DatasetFileFormat,
1303 options: &DatasetFileReadOptions,
1304) -> Result<String, LakeError> {
1305 let canonical_path =
1306 fs::canonicalize(source_path).unwrap_or_else(|_| source_path.to_path_buf());
1307 let path_sql = sql_literal(&canonical_path.display().to_string());
1308 let relation = match format {
1309 DatasetFileFormat::ArrowIpc => {
1310 return Err(LakeError::InvalidMetadataValue {
1311 key: "format",
1312 value: "Arrow IPC must be converted inside the lake adapter before DuckDB import"
1313 .to_string(),
1314 });
1315 }
1316 DatasetFileFormat::Csv => {
1317 let mut args = vec![
1318 path_sql,
1319 format!("header={}", sql_bool(options.header)),
1320 format!("delim={}", sql_literal(&options.delimiter.to_string())),
1321 ];
1322 if !options.null_strings.is_empty() {
1323 args.push(format!(
1324 "nullstr={}",
1325 string_list_sql(&options.null_strings)
1326 ));
1327 }
1328 format!("SELECT * FROM read_csv_auto({})", args.join(", "))
1329 }
1330 DatasetFileFormat::Json => format!(
1331 "SELECT * FROM read_json_auto({}, format={})",
1332 path_sql,
1333 sql_literal(&options.json_format)
1334 ),
1335 DatasetFileFormat::Parquet => format!("SELECT * FROM read_parquet({path_sql})"),
1336 DatasetFileFormat::Xlsx => {
1337 let mut args = vec![path_sql];
1338 if let Some(sheet) = options
1339 .sheet
1340 .as_deref()
1341 .filter(|value| !value.trim().is_empty())
1342 {
1343 args.push(format!("sheet={}", sql_literal(sheet)));
1344 }
1345 format!("SELECT * FROM read_xlsx({})", args.join(", "))
1346 }
1347 };
1348 Ok(relation)
1349}
1350
1351fn arrow_ipc_file_to_parquet(source_path: &Path, parquet_path: &Path) -> Result<(), LakeError> {
1352 let result = (|| {
1353 let reader = arrow_ipc::reader::StreamReader::try_new(
1354 ahri_tre_tabular::open_ipc_messages(source_path)?,
1355 None,
1356 )?;
1357 let mut writer = parquet::arrow::ArrowWriter::try_new(
1358 File::create(parquet_path)?,
1359 reader.schema(),
1360 None,
1361 )
1362 .map_err(ahri_tre_tabular::TabularError::from)?;
1363 for batch in reader {
1364 let batch = batch?;
1365 if batch.get_array_memory_size() > ahri_tre_tabular::streaming::MAX_BATCH_BYTES {
1366 return Err(
1367 std::io::Error::other("Arrow batch exceeds upload memory ceiling").into(),
1368 );
1369 }
1370 writer
1371 .write(&batch)
1372 .map_err(ahri_tre_tabular::TabularError::from)?;
1373 writer
1374 .flush()
1375 .map_err(ahri_tre_tabular::TabularError::from)?;
1376 }
1377 writer
1378 .close()
1379 .map_err(ahri_tre_tabular::TabularError::from)?;
1380 Ok::<_, LakeError>(())
1381 })();
1382 if result.is_err() {
1383 let _ = fs::remove_file(parquet_path);
1384 }
1385 result
1386}
1387
1388fn temporary_arrow_ipc_parquet_path(source_path: &Path) -> PathBuf {
1389 source_path.with_extension(format!("arrow-ipc.{}.parquet", Uuid::new_v4().simple()))
1390}
1391
1392fn trimmed_select_sql(sql: &str) -> Result<String, LakeError> {
1393 let trimmed = sql.trim().trim_end_matches(';').trim();
1394 if trimmed.is_empty()
1395 || trimmed.contains(';')
1396 || !trimmed.to_ascii_lowercase().starts_with("select ")
1397 {
1398 return Err(LakeError::InvalidMetadataValue {
1399 key: "source_sql",
1400 value: "expected a single read-only SELECT statement".to_string(),
1401 });
1402 }
1403 Ok(trimmed.to_string())
1404}
1405
1406fn table_row_count(connection: &duckdb::Connection, table_sql: &str) -> Result<usize, LakeError> {
1407 let count: i64 =
1408 connection.query_row(&format!("SELECT count(*) FROM {table_sql}"), [], |row| {
1409 row.get(0)
1410 })?;
1411 usize::try_from(count).map_err(|_| LakeError::InvalidMetadataValue {
1412 key: "row_count",
1413 value: count.to_string(),
1414 })
1415}
1416
1417fn limited_relation_sql(relation_sql: &str, limit: Option<usize>) -> String {
1418 match limit {
1419 Some(limit) => format!("(SELECT * FROM {relation_sql} LIMIT {limit}) AS limited_dataset"),
1420 None => relation_sql.to_string(),
1421 }
1422}
1423
1424fn table_column_names(
1425 connection: &duckdb::Connection,
1426 table_sql: &str,
1427) -> Result<Vec<String>, LakeError> {
1428 Ok(table_columns(connection, table_sql)?
1429 .into_iter()
1430 .map(|column| column.name)
1431 .collect())
1432}
1433
1434fn table_columns(
1435 connection: &duckdb::Connection,
1436 table_sql: &str,
1437) -> Result<Vec<ImportedColumn>, LakeError> {
1438 let mut statement = connection.prepare(&format!("DESCRIBE SELECT * FROM {table_sql}"))?;
1439 let columns = statement.query_map([], |row| {
1440 Ok(ImportedColumn {
1441 name: row.get::<_, String>(0)?,
1442 duckdb_type: row.get::<_, String>(1)?,
1443 })
1444 })?;
1445 let mut imported = Vec::new();
1446 for column in columns {
1447 imported.push(column?);
1448 }
1449 Ok(imported)
1450}
1451
1452fn preview_rows(
1453 connection: &duckdb::Connection,
1454 relation_sql: &str,
1455 columns: &[ImportedColumn],
1456 limit: usize,
1457) -> Result<Vec<Vec<Option<String>>>, LakeError> {
1458 if columns.is_empty() || limit == 0 {
1459 return Ok(Vec::new());
1460 }
1461 let projection = columns
1462 .iter()
1463 .map(|column| {
1464 let name = quote_identifier(&column.name);
1465 format!("CAST({name} AS VARCHAR)")
1466 })
1467 .collect::<Vec<_>>()
1468 .join(", ");
1469 let mut statement = connection.prepare(&format!(
1470 "SELECT {projection} FROM {relation_sql} LIMIT {limit}"
1471 ))?;
1472 let rows = statement.query_map([], |row| {
1473 let mut values = Vec::with_capacity(columns.len());
1474 for index in 0..columns.len() {
1475 values.push(row.get::<_, Option<String>>(index)?);
1476 }
1477 Ok(values)
1478 })?;
1479 let mut preview = Vec::new();
1480 for row in rows {
1481 preview.push(row?);
1482 }
1483 Ok(preview)
1484}
1485
1486fn export_relation_with_copy(
1487 connection: &duckdb::Connection,
1488 source_sql: &str,
1489 destination_path: &Path,
1490 copy_options: &str,
1491 compress: bool,
1492) -> Result<(), LakeError> {
1493 if compress {
1494 let temporary_path = temporary_export_path(destination_path);
1495 let result = (|| {
1496 copy_relation_to_path(connection, source_sql, &temporary_path, copy_options)?;
1497 compress_zstd_file(&temporary_path, destination_path)
1498 })();
1499 let _ = fs::remove_file(&temporary_path);
1500 result
1501 } else {
1502 copy_relation_to_path(connection, source_sql, destination_path, copy_options)
1503 }
1504}
1505
1506fn copy_relation_to_path(
1507 connection: &duckdb::Connection,
1508 source_sql: &str,
1509 destination_path: &Path,
1510 copy_options: &str,
1511) -> Result<(), LakeError> {
1512 let destination_sql = sql_literal(&destination_path.display().to_string());
1513 connection.execute_batch(&format!(
1514 "COPY (SELECT * FROM {source_sql}) TO {destination_sql} ({copy_options});"
1515 ))?;
1516 Ok(())
1517}
1518
1519fn export_relation_to_arrow_table(
1520 connection: &duckdb::Connection,
1521 source_sql: &str,
1522 metadata: Option<&AhriTreArrowMetadata>,
1523) -> Result<AnnotatedArrowTable, LakeError> {
1524 let mut statement = connection.prepare(&format!("SELECT * FROM {source_sql}"))?;
1525 let mut arrow = statement.query_arrow([])?;
1526 let schema = arrow.get_schema();
1527 let table = ArrowTable::try_new(schema, arrow.by_ref().collect())?;
1528 if let Some(metadata) = metadata {
1529 Ok(annotate_arrow_table(&table, metadata)?)
1530 } else {
1531 Ok(AnnotatedArrowTable {
1532 table,
1533 warnings: Vec::new(),
1534 })
1535 }
1536}
1537
1538fn compress_zstd_file(source_path: &Path, destination_path: &Path) -> Result<(), LakeError> {
1539 let source = BufReader::new(File::open(source_path)?);
1540 let mut encoder =
1541 zstd::stream::write::Encoder::new(BufWriter::new(File::create(destination_path)?), 3)?;
1542 let mut reader = source;
1543 io::copy(&mut reader, &mut encoder)?;
1544 encoder.finish()?.flush()?;
1545 Ok(())
1546}
1547
1548fn temporary_export_path(destination_path: &Path) -> PathBuf {
1549 let extension = destination_path
1550 .extension()
1551 .and_then(|value| value.to_str())
1552 .filter(|value| !value.trim().is_empty())
1553 .unwrap_or("tmp");
1554 destination_path.with_extension(format!("{}.{}.tmp", extension, Uuid::new_v4().simple()))
1555}
1556
1557fn quote_identifier(value: &str) -> String {
1558 format!("\"{}\"", value.replace('"', "\"\""))
1559}
1560
1561fn sql_literal(value: &str) -> String {
1562 format!("'{}'", value.replace('\'', "''"))
1563}
1564
1565fn string_list_sql(values: &[String]) -> String {
1566 format!(
1567 "[{}]",
1568 values
1569 .iter()
1570 .map(|value| sql_literal(value))
1571 .collect::<Vec<_>>()
1572 .join(", ")
1573 )
1574}
1575
1576fn sql_bool(value: bool) -> &'static str {
1577 if value { "true" } else { "false" }
1578}
1579
1580#[cfg(test)]
1581mod tests {
1582 use std::fs;
1583 use std::os::unix::fs::PermissionsExt;
1584 use std::path::{Path, PathBuf};
1585
1586 use ahri_tre_types::{NcName, StudyId};
1587 use uuid::Uuid;
1588
1589 use super::*;
1590 use crate::{ScratchAttemptId, TrustedScratch};
1591
1592 #[test]
1593 fn rejected_sql_load_preserves_existing_target_and_removes_staging() {
1594 let connection = duckdb::Connection::open_in_memory().unwrap();
1595 let lake = DuckLakeAdapter::new("/unused");
1596 let mut request = LoadDatasetTableFromSqlRequest {
1597 study_id: StudyId(Uuid::new_v4()),
1598 dataset_name: NcName::parse("preserved").unwrap(),
1599 major: 1,
1600 minor: 0,
1601 patch: 0,
1602 source_sql: "SELECT 42 AS answer".into(),
1603 };
1604 let loaded = lake
1605 .load_dataset_table_from_sql(&connection, &request)
1606 .unwrap();
1607 request.source_sql = "SELECT * FROM missing_source".into();
1608 assert!(
1609 lake.load_dataset_table_from_sql(&connection, &request)
1610 .is_err()
1611 );
1612 let answer: i64 = connection
1613 .query_row(
1614 &format!("SELECT answer FROM {}", loaded.relation.to_sql()),
1615 [],
1616 |row| row.get(0),
1617 )
1618 .expect("rejected staging must preserve existing target");
1619 assert_eq!(answer, 42);
1620 request.source_sql = "SELECT 99 AS answer".into();
1621 assert!(
1622 lake.load_dataset_table_from_sql(&connection, &request)
1623 .is_err(),
1624 "an occupied version must not be overwritten"
1625 );
1626 let stages: i64 = connection
1627 .query_row(
1628 "SELECT count(*) FROM duckdb_tables() WHERE table_name LIKE '__tre_stage_%'",
1629 [],
1630 |row| row.get(0),
1631 )
1632 .unwrap();
1633 assert_eq!(stages, 0);
1634 }
1635
1636 #[test]
1637 fn loads_csv_file_into_study_scoped_dataset_table() {
1638 let root = temp_lake_root();
1639 let source_path = root.join("ingest.csv");
1640 fs::write(&source_path, b"category,country\nA,USA\nB,South Africa\n")
1641 .expect("source fixture should write");
1642 let adapter = DuckLakeAdapter::new(path_string(&root));
1643 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
1644 let study_id = StudyId(Uuid::new_v4());
1645
1646 let loaded = adapter
1647 .load_dataset_table_from_file(
1648 &connection,
1649 &LoadDatasetTableFromFileRequest {
1650 study_id,
1651 dataset_name: NcName::parse("ingested_dataset")
1652 .expect("dataset name should parse"),
1653 major: 1,
1654 minor: 0,
1655 patch: 0,
1656 source_path,
1657 format: DatasetFileFormat::Csv,
1658 options: DatasetFileReadOptions::default(),
1659 },
1660 )
1661 .expect("dataset table should load");
1662
1663 assert_eq!(
1664 loaded.relation.schema_name,
1665 format!(
1666 "study_{}",
1667 study_id
1668 .0
1669 .to_string()
1670 .to_ascii_lowercase()
1671 .replace('-', "_")
1672 )
1673 );
1674 assert_eq!(loaded.relation.table_name, "ingested_dataset_1_0_0");
1675 assert_eq!(loaded.row_count, 2);
1676 assert_eq!(loaded.column_count, 2);
1677 assert_eq!(
1678 loaded.columns,
1679 vec![
1680 ImportedColumn {
1681 name: "category".to_string(),
1682 duckdb_type: "VARCHAR".to_string(),
1683 },
1684 ImportedColumn {
1685 name: "country".to_string(),
1686 duckdb_type: "VARCHAR".to_string(),
1687 },
1688 ]
1689 );
1690 let count: i64 = connection
1691 .query_row(
1692 &format!("SELECT count(*) FROM {}", loaded.relation.to_sql()),
1693 [],
1694 |row| row.get(0),
1695 )
1696 .expect("loaded dataset should query");
1697 assert_eq!(count, 2);
1698
1699 let _ = fs::remove_dir_all(&root);
1700 }
1701
1702 #[test]
1703 fn validates_csv_with_the_same_duckdb_reader_without_persisting_a_relation() {
1704 let root = temp_lake_root();
1705 let valid_path = root.join("valid.csv");
1706 let invalid_path = root.join("invalid.csv");
1707 fs::write(&valid_path, b"id,notes\n1,\"first line\nsecond line\"\n")
1708 .expect("valid source fixture should write");
1709 fs::write(&invalid_path, b"id,value\n1\n").expect("invalid source fixture should write");
1710 let adapter = DuckLakeAdapter::new(path_string(&root));
1711 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
1712
1713 adapter
1714 .validate_dataset_file(
1715 &connection,
1716 &valid_path,
1717 DatasetFileFormat::Csv,
1718 &DatasetFileReadOptions::default(),
1719 )
1720 .expect("DuckDB should accept the valid CSV");
1721 assert!(
1722 adapter
1723 .validate_dataset_file(
1724 &connection,
1725 &invalid_path,
1726 DatasetFileFormat::Csv,
1727 &DatasetFileReadOptions::default(),
1728 )
1729 .is_err()
1730 );
1731
1732 let _ = fs::remove_dir_all(&root);
1733 }
1734
1735 #[test]
1736 fn retaining_an_exact_dataset_set_removes_stale_study_tables() {
1737 let adapter = DuckLakeAdapter::new("/lake");
1738 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
1739 let study_id = StudyId(Uuid::new_v4());
1740 let retained = dataset_table_relation(study_id, "retained", 1, 0, 0);
1741 let stale = dataset_table_relation(study_id, "stale", 1, 0, 0);
1742 connection
1743 .execute_batch(&format!(
1744 "CREATE SCHEMA {}; CREATE TABLE {} (id INTEGER); CREATE TABLE {} (id INTEGER);",
1745 quote_identifier(&retained.schema_name),
1746 retained.to_sql(),
1747 stale.to_sql()
1748 ))
1749 .expect("study fixtures should be created");
1750
1751 let removed = adapter
1752 .retain_only_dataset_tables(&connection, study_id, std::slice::from_ref(&retained))
1753 .expect("stale study tables should be removed");
1754
1755 assert_eq!(removed, vec![stale]);
1756 connection
1757 .query_row(
1758 &format!("SELECT count(*) FROM {}", retained.to_sql()),
1759 [],
1760 |row| row.get::<_, i64>(0),
1761 )
1762 .expect("retained table should remain queryable");
1763 }
1764
1765 #[test]
1766 fn exact_csv_replacement_removes_contamination_and_reloads_expected_tables() {
1767 let root = temp_lake_root();
1768 let scratch_root = root.join("scratch");
1769 fs::create_dir(&scratch_root).expect("scratch root should be created");
1770 fs::set_permissions(&scratch_root, fs::Permissions::from_mode(0o700))
1771 .expect("scratch root should be private");
1772 let scratch = TrustedScratch::open(&scratch_root, [Path::new("/lake")])
1773 .expect("trusted scratch should open")
1774 .create_attempt(
1775 ScratchAttemptId::new("00000000000000000000000000000015")
1776 .expect("attempt identifier should be valid"),
1777 )
1778 .expect("scratch attempt should be created");
1779 let adapter = DuckLakeAdapter::new("/lake");
1780 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
1781 let study_id = StudyId(Uuid::new_v4());
1782 let stale = dataset_table_relation(study_id, "stale", 1, 0, 0);
1783 connection
1784 .execute_batch(&format!(
1785 "CREATE SCHEMA {}; CREATE TABLE {} (id INTEGER);",
1786 quote_identifier(&stale.schema_name),
1787 stale.to_sql()
1788 ))
1789 .expect("contaminating table should be created");
1790 let prior = dataset_table_relation(study_id, "alpha", 1, 0, 0);
1791 connection
1792 .execute_batch(&format!(
1793 "CREATE TABLE {} (old_value INTEGER); INSERT INTO {} VALUES (99);",
1794 prior.to_sql(),
1795 prior.to_sql()
1796 ))
1797 .unwrap();
1798 let seeds = [
1799 DatasetCsvSeed {
1800 dataset_name: NcName::parse("alpha").unwrap(),
1801 contents: b"id,label\n1,first\n",
1802 },
1803 DatasetCsvSeed {
1804 dataset_name: NcName::parse("beta").unwrap(),
1805 contents: b"id,label\n2,second\n",
1806 },
1807 ];
1808
1809 let loaded = adapter
1810 .replace_dataset_tables_from_csv(&connection, &scratch, study_id, &seeds)
1811 .expect("exact seed replacement should succeed");
1812
1813 assert_eq!(loaded.len(), 2);
1814 assert_eq!(
1815 connection
1816 .query_row(
1817 &format!("SELECT label FROM {}", prior.to_sql()),
1818 [],
1819 |row| row.get::<_, String>(0)
1820 )
1821 .unwrap(),
1822 "first"
1823 );
1824 assert_eq!(
1825 study_dataset_tables(&connection, &study_lake_schema(study_id)).unwrap(),
1826 vec![
1827 dataset_table_relation(study_id, "alpha", 1, 0, 0),
1828 dataset_table_relation(study_id, "beta", 1, 0, 0),
1829 ]
1830 );
1831 drop(scratch);
1832 let _ = fs::remove_dir_all(root);
1833 }
1834
1835 #[test]
1836 fn exact_csv_replacement_rolls_back_after_a_mid_seed_failure() {
1837 let root = temp_lake_root();
1838 let scratch_root = root.join("scratch");
1839 fs::create_dir(&scratch_root).expect("scratch root should be created");
1840 fs::set_permissions(&scratch_root, fs::Permissions::from_mode(0o700))
1841 .expect("scratch root should be private");
1842 let scratch = TrustedScratch::open(&scratch_root, [Path::new("/lake")])
1843 .expect("trusted scratch should open")
1844 .create_attempt(
1845 ScratchAttemptId::new("00000000000000000000000000000016")
1846 .expect("attempt identifier should be valid"),
1847 )
1848 .expect("scratch attempt should be created");
1849 let adapter = DuckLakeAdapter::new("/lake");
1850 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
1851 let study_id = StudyId(Uuid::new_v4());
1852 let stale = dataset_table_relation(study_id, "stale", 1, 0, 0);
1853 connection
1854 .execute_batch(&format!(
1855 "CREATE SCHEMA {}; CREATE TABLE {} (id INTEGER); INSERT INTO {} VALUES (7);",
1856 quote_identifier(&stale.schema_name),
1857 stale.to_sql(),
1858 stale.to_sql()
1859 ))
1860 .expect("pre-seed state should be created");
1861 fs::write(scratch.path().join("beta_1_0_0.csv"), b"occupied")
1862 .expect("later seed path should be occupied");
1863 let seeds = [
1864 DatasetCsvSeed {
1865 dataset_name: NcName::parse("alpha").unwrap(),
1866 contents: b"id\n1\n",
1867 },
1868 DatasetCsvSeed {
1869 dataset_name: NcName::parse("beta").unwrap(),
1870 contents: b"id\n2\n",
1871 },
1872 ];
1873
1874 adapter
1875 .replace_dataset_tables_from_csv(&connection, &scratch, study_id, &seeds)
1876 .expect_err("occupied second seed path should fail after the first replacement");
1877
1878 assert_eq!(
1879 study_dataset_tables(&connection, &study_lake_schema(study_id)).unwrap(),
1880 vec![stale.clone()]
1881 );
1882 assert_eq!(
1883 connection
1884 .query_row(&format!("SELECT id FROM {}", stale.to_sql()), [], |row| {
1885 row.get::<_, i64>(0)
1886 })
1887 .unwrap(),
1888 7
1889 );
1890 drop(scratch);
1891 let _ = fs::remove_dir_all(root);
1892 }
1893
1894 #[test]
1895 fn loads_arrow_ipc_file_into_study_scoped_dataset_table() {
1896 let root = temp_lake_root();
1897 let source_path = root.join("ingest.arrow");
1898 let batch = ahri_tre_tabular::columns_to_record_batch(vec![
1899 ahri_tre_tabular::TabularColumn::Utf8 {
1900 name: "status".to_string(),
1901 values: vec![Some("A".to_string()), Some("D".to_string())],
1902 },
1903 ahri_tre_tabular::TabularColumn::Int64 {
1904 name: "age".to_string(),
1905 values: vec![Some(34), Some(41)],
1906 },
1907 ])
1908 .expect("Arrow IPC fixture batch should build");
1909 let table = ArrowTable::from_batches(vec![batch]).expect("Arrow IPC table should build");
1910 let bytes =
1911 ahri_tre_tabular::write_ipc_file(&table).expect("Arrow IPC fixture should serialize");
1912 fs::write(&source_path, bytes).expect("Arrow IPC fixture should write");
1913 let adapter = DuckLakeAdapter::new(path_string(&root));
1914 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
1915 let study_id = StudyId(Uuid::new_v4());
1916
1917 let loaded = adapter
1918 .load_dataset_table_from_file(
1919 &connection,
1920 &LoadDatasetTableFromFileRequest {
1921 study_id,
1922 dataset_name: NcName::parse("ingested_arrow_dataset")
1923 .expect("dataset name should parse"),
1924 major: 1,
1925 minor: 0,
1926 patch: 0,
1927 source_path,
1928 format: DatasetFileFormat::ArrowIpc,
1929 options: DatasetFileReadOptions {
1930 null_strings: vec!["".to_string()],
1931 ..DatasetFileReadOptions::default()
1932 },
1933 },
1934 )
1935 .expect("Arrow IPC dataset table should load");
1936
1937 assert_eq!(loaded.row_count, 2);
1938 assert_eq!(loaded.column_count, 2);
1939 assert_eq!(
1940 loaded.columns,
1941 vec![
1942 ImportedColumn {
1943 name: "status".to_string(),
1944 duckdb_type: "VARCHAR".to_string(),
1945 },
1946 ImportedColumn {
1947 name: "age".to_string(),
1948 duckdb_type: "BIGINT".to_string(),
1949 },
1950 ]
1951 );
1952 let count: i64 = connection
1953 .query_row(
1954 &format!("SELECT count(*) FROM {}", loaded.relation.to_sql()),
1955 [],
1956 |row| row.get(0),
1957 )
1958 .expect("loaded Arrow IPC dataset should query");
1959 assert_eq!(count, 2);
1960 assert!(
1961 fs::read_dir(&root)
1962 .expect("lake root should list")
1963 .all(|entry| entry
1964 .expect("lake entry should read")
1965 .path()
1966 .extension()
1967 .and_then(|extension| extension.to_str())
1968 .is_none_or(|extension| extension != "parquet"))
1969 );
1970
1971 let _ = fs::remove_dir_all(&root);
1972 }
1973
1974 #[test]
1975 fn exports_dataset_table_to_supported_local_formats() {
1976 let root = temp_lake_root();
1977 let source_path = root.join("ingest.csv");
1978 fs::write(&source_path, b"category,country\nA,USA\nB,South Africa\n")
1979 .expect("source fixture should write");
1980 let adapter = DuckLakeAdapter::new(path_string(&root));
1981 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
1982 let study_id = StudyId(Uuid::new_v4());
1983 let dataset_name = NcName::parse("exportable_dataset").unwrap();
1984 adapter
1985 .load_dataset_table_from_file(
1986 &connection,
1987 &LoadDatasetTableFromFileRequest {
1988 study_id,
1989 dataset_name: dataset_name.clone(),
1990 major: 1,
1991 minor: 0,
1992 patch: 0,
1993 source_path,
1994 format: DatasetFileFormat::Csv,
1995 options: DatasetFileReadOptions::default(),
1996 },
1997 )
1998 .expect("dataset table should load");
1999
2000 let parquet = adapter
2001 .export_dataset_table(
2002 &connection,
2003 &ExportDatasetTableRequest {
2004 study_id,
2005 dataset_name: dataset_name.clone(),
2006 major: 1,
2007 minor: 0,
2008 patch: 0,
2009 limit: None,
2010 destination_path: root.join("export.parquet"),
2011 format: DatasetExportFormat::Parquet,
2012 compress: false,
2013 arrow_metadata: None,
2014 },
2015 )
2016 .expect("parquet export should succeed");
2017 assert_eq!(parquet.row_count, 2);
2018 assert_eq!(parquet.column_count, 2);
2019 assert!(parquet.compressed);
2020 assert!(parquet.exported_size_bytes > 0);
2021
2022 let ndjson = adapter
2023 .export_dataset_table(
2024 &connection,
2025 &ExportDatasetTableRequest {
2026 study_id,
2027 dataset_name: dataset_name.clone(),
2028 major: 1,
2029 minor: 0,
2030 patch: 0,
2031 limit: None,
2032 destination_path: root.join("export.ndjson.zst"),
2033 format: DatasetExportFormat::Ndjson,
2034 compress: true,
2035 arrow_metadata: None,
2036 },
2037 )
2038 .expect("compressed ndjson export should succeed");
2039 assert_eq!(ndjson.row_count, 2);
2040 assert!(ndjson.compressed);
2041 let compressed = fs::read(root.join("export.ndjson.zst")).unwrap();
2042 let decoded =
2043 zstd::stream::decode_all(compressed.as_slice()).expect("ndjson should decompress");
2044 let text = String::from_utf8(decoded).expect("ndjson should be UTF-8");
2045 assert!(text.contains("\"category\""));
2046 assert!(text.lines().count() >= 2);
2047
2048 let arrow = adapter
2049 .export_dataset_table(
2050 &connection,
2051 &ExportDatasetTableRequest {
2052 study_id,
2053 dataset_name,
2054 major: 1,
2055 minor: 0,
2056 patch: 0,
2057 limit: None,
2058 destination_path: root.join("export.arrow"),
2059 format: DatasetExportFormat::ArrowIpc,
2060 compress: false,
2061 arrow_metadata: None,
2062 },
2063 )
2064 .expect("compressed Arrow IPC export should succeed");
2065 assert!(arrow.compressed);
2066 assert!(
2067 fs::read(root.join("export.arrow"))
2068 .expect("Arrow IPC file should read")
2069 .starts_with(b"ARROW1")
2070 );
2071
2072 let limited_csv = adapter
2073 .export_dataset_table(
2074 &connection,
2075 &ExportDatasetTableRequest {
2076 study_id,
2077 dataset_name: NcName::parse("exportable_dataset").unwrap(),
2078 major: 1,
2079 minor: 0,
2080 patch: 0,
2081 limit: Some(1),
2082 destination_path: root.join("limited.csv"),
2083 format: DatasetExportFormat::Csv,
2084 compress: false,
2085 arrow_metadata: None,
2086 },
2087 )
2088 .expect("limited CSV export should succeed");
2089 assert_eq!(limited_csv.row_count, 1);
2090 assert_eq!(
2091 fs::read_to_string(root.join("limited.csv")).expect("limited CSV should read"),
2092 "category,country\nA,USA\n"
2093 );
2094
2095 let _ = fs::remove_dir_all(&root);
2096 }
2097
2098 #[test]
2099 fn previews_dataset_table_with_limit_and_nulls_without_exporting() {
2100 let root = temp_lake_root();
2101 let source_path = root.join("preview.csv");
2102 fs::write(
2103 &source_path,
2104 b"id,label,note\n1,alpha,short\n2,,a very long note\n3,gamma,\n",
2105 )
2106 .expect("source fixture should write");
2107 let adapter = DuckLakeAdapter::new(path_string(&root));
2108 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
2109 let study_id = StudyId(Uuid::new_v4());
2110 let dataset_name = NcName::parse("preview_dataset").unwrap();
2111 adapter
2112 .load_dataset_table_from_file(
2113 &connection,
2114 &LoadDatasetTableFromFileRequest {
2115 study_id,
2116 dataset_name: dataset_name.clone(),
2117 major: 1,
2118 minor: 0,
2119 patch: 0,
2120 source_path,
2121 format: DatasetFileFormat::Csv,
2122 options: DatasetFileReadOptions::default(),
2123 },
2124 )
2125 .expect("dataset table should load");
2126
2127 let preview = adapter
2128 .preview_dataset_table(
2129 &connection,
2130 &PreviewDatasetTableRequest {
2131 study_id,
2132 dataset_name,
2133 major: 1,
2134 minor: 0,
2135 patch: 0,
2136 limit: 2,
2137 },
2138 )
2139 .expect("dataset preview should read");
2140
2141 assert_eq!(
2142 preview
2143 .columns
2144 .iter()
2145 .map(|column| column.name.as_str())
2146 .collect::<Vec<_>>(),
2147 vec!["id", "label", "note"]
2148 );
2149 assert_eq!(preview.displayed_row_count, 2);
2150 assert_eq!(preview.total_row_count, 3);
2151 assert!(preview.truncated);
2152 assert_eq!(
2153 preview.rows,
2154 vec![
2155 vec![
2156 Some("1".to_string()),
2157 Some("alpha".to_string()),
2158 Some("short".to_string())
2159 ],
2160 vec![
2161 Some("2".to_string()),
2162 None,
2163 Some("a very long note".to_string())
2164 ],
2165 ]
2166 );
2167 assert!(!root.join("preview_dataset.csv").exists());
2168
2169 let _ = fs::remove_dir_all(&root);
2170 }
2171
2172 #[test]
2173 fn loads_redcap_eav_with_event_and_repeat_metadata() {
2174 let root = temp_lake_root();
2175 let source_path = root.join("redcap_eav.csv");
2176 fs::write(
2177 &source_path,
2178 b"record,redcap_event_name,redcap_repeat_instrument,redcap_repeat_instance,field_name,value\n\
21791,baseline,visit,1,age,40\n\
21801,baseline,visit,1,status,1\n\
21811,followup,visit,2,age,41\n\
21821,followup,visit,2,status,2\n",
2183 )
2184 .expect("source fixture should write");
2185 let adapter = DuckLakeAdapter::new(path_string(&root));
2186 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
2187 let study_id = StudyId(Uuid::new_v4());
2188
2189 let loaded = adapter
2190 .load_redcap_form_dataset(
2191 &connection,
2192 &LoadRedcapFormDatasetRequest {
2193 study_id,
2194 dataset_name: NcName::parse("redcap_1_visit").unwrap(),
2195 major: 1,
2196 minor: 0,
2197 patch: 0,
2198 source_eav_path: source_path,
2199 fields: vec![
2200 RedcapEavField {
2201 name: "age".to_string(),
2202 value_type_id: ahri_tre_types::ValueTypeId(1),
2203 value_format: None,
2204 choices: Vec::new(),
2205 },
2206 RedcapEavField {
2207 name: "status".to_string(),
2208 value_type_id: ahri_tre_types::ValueTypeId(7),
2209 value_format: None,
2210 choices: vec![
2211 RedcapEavChoice {
2212 redcap_value: "1".to_string(),
2213 value: 1,
2214 },
2215 RedcapEavChoice {
2216 redcap_value: "2".to_string(),
2217 value: 2,
2218 },
2219 ],
2220 },
2221 ],
2222 synthetic_columns: vec![
2223 RedcapSyntheticColumn {
2224 source_column: "redcap_event_name".to_string(),
2225 output_column: "redcap_1_visit_redcap_event".to_string(),
2226 value_type_id: ahri_tre_types::ValueTypeId(7),
2227 choices: vec![
2228 RedcapEavChoice {
2229 redcap_value: "baseline".to_string(),
2230 value: 1,
2231 },
2232 RedcapEavChoice {
2233 redcap_value: "followup".to_string(),
2234 value: 2,
2235 },
2236 ],
2237 },
2238 RedcapSyntheticColumn {
2239 source_column: "redcap_repeat_instance".to_string(),
2240 output_column: "redcap_1_visit_redcap_repeat_instance".to_string(),
2241 value_type_id: ahri_tre_types::ValueTypeId(1),
2242 choices: Vec::new(),
2243 },
2244 ],
2245 },
2246 )
2247 .expect("REDCap EAV dataset should load");
2248
2249 assert_eq!(loaded.row_count, 2);
2250 assert!(
2251 loaded
2252 .column_names
2253 .contains(&"redcap_1_visit_redcap_event".to_string())
2254 );
2255 let sum_age: i64 = connection
2256 .query_row(
2257 &format!("SELECT sum(age) FROM {}", loaded.relation.to_sql()),
2258 [],
2259 |row| row.get(0),
2260 )
2261 .expect("loaded REDCap dataset should query");
2262 assert_eq!(sum_age, 81);
2263 let repeat_type: String = connection
2264 .query_row(
2265 &format!(
2266 "SELECT typeof(redcap_1_visit_redcap_repeat_instance) FROM {} LIMIT 1",
2267 loaded.relation.to_sql()
2268 ),
2269 [],
2270 |row| row.get(0),
2271 )
2272 .expect("repeat instance type should be readable");
2273 assert_eq!(repeat_type, "INTEGER");
2274
2275 let _ = fs::remove_dir_all(&root);
2276 }
2277
2278 #[test]
2279 fn redcap_eav_loader_reports_missing_required_eav_columns() {
2280 let root = temp_lake_root();
2281 let source_path = root.join("redcap_eav_error.csv");
2282 fs::write(&source_path, b"{\"error\":\"export too large\"}\n")
2283 .expect("source fixture should write");
2284 let adapter = DuckLakeAdapter::new(path_string(&root));
2285 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
2286 let study_id = StudyId(Uuid::new_v4());
2287
2288 let error = adapter
2289 .load_redcap_form_dataset(
2290 &connection,
2291 &LoadRedcapFormDatasetRequest {
2292 study_id,
2293 dataset_name: NcName::parse("redcap_1_visit").unwrap(),
2294 major: 1,
2295 minor: 0,
2296 patch: 0,
2297 source_eav_path: source_path,
2298 fields: vec![RedcapEavField {
2299 name: "age".to_string(),
2300 value_type_id: ahri_tre_types::ValueTypeId(1),
2301 value_format: None,
2302 choices: Vec::new(),
2303 }],
2304 synthetic_columns: Vec::new(),
2305 },
2306 )
2307 .expect_err("malformed REDCap EAV export should fail before binding field_name");
2308
2309 assert!(
2310 error
2311 .to_string()
2312 .contains("expected REDCap EAV CSV columns field_name, value")
2313 );
2314
2315 let _ = fs::remove_dir_all(&root);
2316 }
2317
2318 #[test]
2319 fn loads_redcap_eav_with_event_column_after_value() {
2320 let root = temp_lake_root();
2321 let source_path = root.join("redcap_eav_event_after_value.csv");
2322 fs::write(
2323 &source_path,
2324 b"record,field_name,value,redcap_event_name\n\
23252,ls_clinic_attended,1,initial_visit_arm_1\n\
23262,ls_services_requested,3,initial_visit_arm_1\n\
23272,attendance_and_services_complete,2,initial_visit_arm_1\n\
23282,ls_clinic_attended,1,follow_up_visit__1_arm_1\n\
23292,ls_services_requested,1,follow_up_visit__1_arm_1\n\
23302,ls_services_requested,3,follow_up_visit__1_arm_1\n",
2331 )
2332 .expect("source fixture should write");
2333 let adapter = DuckLakeAdapter::new(path_string(&root));
2334 let connection = duckdb::Connection::open_in_memory().expect("DuckDB should open");
2335 let study_id = StudyId(Uuid::new_v4());
2336
2337 let loaded = adapter
2338 .load_redcap_form_dataset(
2339 &connection,
2340 &LoadRedcapFormDatasetRequest {
2341 study_id,
2342 dataset_name: NcName::parse("redcap_797_attendance_and_services").unwrap(),
2343 major: 1,
2344 minor: 0,
2345 patch: 0,
2346 source_eav_path: source_path,
2347 fields: vec![
2348 RedcapEavField {
2349 name: "ls_clinic_attended".to_string(),
2350 value_type_id: ahri_tre_types::ValueTypeId(7),
2351 value_format: None,
2352 choices: vec![RedcapEavChoice {
2353 redcap_value: "1".to_string(),
2354 value: 1,
2355 }],
2356 },
2357 RedcapEavField {
2358 name: "ls_services_requested".to_string(),
2359 value_type_id: ahri_tre_types::ValueTypeId(8),
2360 value_format: None,
2361 choices: Vec::new(),
2362 },
2363 RedcapEavField {
2364 name: "attendance_and_services_complete".to_string(),
2365 value_type_id: ahri_tre_types::ValueTypeId(7),
2366 value_format: None,
2367 choices: vec![RedcapEavChoice {
2368 redcap_value: "2".to_string(),
2369 value: 2,
2370 }],
2371 },
2372 ],
2373 synthetic_columns: vec![RedcapSyntheticColumn {
2374 source_column: "redcap_event_name".to_string(),
2375 output_column: "redcap_797_attendance_and_services_redcap_event"
2376 .to_string(),
2377 value_type_id: ahri_tre_types::ValueTypeId(7),
2378 choices: vec![
2379 RedcapEavChoice {
2380 redcap_value: "initial_visit_arm_1".to_string(),
2381 value: 1,
2382 },
2383 RedcapEavChoice {
2384 redcap_value: "follow_up_visit__1_arm_1".to_string(),
2385 value: 2,
2386 },
2387 ],
2388 }],
2389 },
2390 )
2391 .expect("REDCap EAV dataset should load when event appears after value");
2392
2393 assert_eq!(loaded.row_count, 2);
2394 assert!(
2395 loaded
2396 .column_names
2397 .contains(&"redcap_797_attendance_and_services_redcap_event".to_string())
2398 );
2399 let services: String = connection
2400 .query_row(
2401 &format!(
2402 "SELECT ls_services_requested FROM {} WHERE redcap_797_attendance_and_services_redcap_event = 2",
2403 loaded.relation.to_sql()
2404 ),
2405 [],
2406 |row| row.get(0),
2407 )
2408 .expect("loaded REDCap follow-up row should query");
2409 assert_eq!(services, "1, 3");
2410
2411 let _ = fs::remove_dir_all(&root);
2412 }
2413
2414 fn temp_lake_root() -> PathBuf {
2415 let root = std::env::temp_dir().join(format!(
2416 "ahri_tre_lake_dataset_{}_{}",
2417 std::process::id(),
2418 Uuid::new_v4().simple()
2419 ));
2420 fs::create_dir_all(&root).expect("temp lake root should be created");
2421 root
2422 }
2423
2424 fn path_string(path: &Path) -> String {
2425 path.display().to_string()
2426 }
2427}