Skip to main content

ahri_tre_lake/
dataset_loading.rs

1//! Dataset table loading and export helpers for DuckLake relations.
2
3use 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/// File formats accepted when loading a dataset table from disk.
18#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub enum DatasetFileFormat {
20    ArrowIpc,
21    Csv,
22    Json,
23    Parquet,
24    Xlsx,
25}
26
27impl DatasetFileFormat {
28    /// Returns the conventional extension for this dataset file format.
29    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/// Options passed to DuckDB file readers when importing a dataset table.
41#[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/// Request to load one versioned dataset table from a source file.
63#[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/// One trusted in-memory CSV selected for atomic dataset replacement.
76#[derive(Debug, Clone, PartialEq, Eq)]
77pub struct DatasetCsvSeed<'a> {
78    pub dataset_name: NcName,
79    pub contents: &'a [u8],
80}
81
82/// Request to materialize one versioned dataset table from a SQL query.
83#[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/// Output format for exporting a dataset table from DuckLake.
94#[derive(Debug, Clone, Copy, PartialEq, Eq)]
95pub enum DatasetExportFormat {
96    Csv,
97    Ndjson,
98    Parquet,
99    ArrowIpc,
100}
101
102/// Request to export one versioned dataset table to a local file.
103#[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/// Request to read a bounded preview of one versioned dataset table.
118#[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/// Summary of a dataset table export operation.
129#[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/// Bounded row preview for one dataset table.
142#[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/// Request to reshape a REDCap EAV export into a versioned dataset table.
153#[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/// REDCap field projection metadata used when pivoting EAV data.
166#[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/// Mapping from a REDCap raw choice value to the TRE categorical integer value.
175#[derive(Debug, Clone, PartialEq, Eq)]
176pub struct RedcapEavChoice {
177    pub redcap_value: String,
178    pub value: i32,
179}
180
181/// Derived REDCap output column projected from an existing source column.
182#[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/// DuckDB schema and table name for a versioned dataset table.
191#[derive(Debug, Clone, PartialEq, Eq)]
192pub struct DatasetTableRelation {
193    pub schema_name: String,
194    pub table_name: String,
195}
196
197impl DatasetTableRelation {
198    /// Returns the unquoted `schema.table` relation name for display.
199    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/// Summary of a dataset table loaded into DuckLake.
213#[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/// DuckDB-reported shape for one imported dataset column.
223#[derive(Debug, Clone, PartialEq, Eq)]
224pub struct ImportedColumn {
225    pub name: String,
226    pub duckdb_type: String,
227}
228
229impl DuckLakeAdapter {
230    /// Validates a dataset file by scanning it with the same DuckDB reader used
231    /// by governed dataset imports, without creating a persistent relation.
232    ///
233    /// # Errors
234    ///
235    /// Returns an error if the source file is missing or DuckDB cannot parse
236    /// and scan the complete file with the requested format and options.
237    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    /// Atomically replaces the exact Study dataset set in an existing
266    /// server-managed DuckLake catalog from trusted CSV inputs.
267    ///
268    /// # Errors
269    ///
270    /// Returns an error unless the catalog was opened as an existing
271    /// server-managed catalog without an encryption override, or if private
272    /// staging, dataset replacement, or exact-set verification fails.
273    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            // Exclusive maintenance replaces the contents of retained names too.
324            // Ordinary writers still require an absent target before creation.
325            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    /// Removes versioned dataset tables in one Study namespace unless they are
383    /// members of the supplied exact retained set.
384    ///
385    /// # Errors
386    ///
387    /// Returns an error if catalog inspection or a table drop fails.
388    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    /// Loads a versioned dataset table from a CSV, JSON, Parquet, Arrow IPC, or XLSX file.
417    ///
418    /// # Errors
419    ///
420    /// Returns an error if the source file is missing, a version component is
421    /// negative, DuckDB cannot read the file, or the target relation cannot be
422    /// created.
423    #[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    /// Loads only the exact output held by an active shared reservation.
433    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    /// Observes one reserved load, including its existing rollback outcome.
453    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    /// Loads a versioned dataset table from a caller-supplied SQL `SELECT`.
557    ///
558    /// # Errors
559    ///
560    /// Returns an error if a version component is negative, the source SQL is
561    /// empty or not a `SELECT`, or DuckDB cannot execute the staging or create
562    /// statements.
563    #[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    /// The output file is produced by the restricted executor, never supplied
573    /// as a client path or SQL document. Existing reservation/cleanup semantics
574    /// remain the sole authority for publishing its Dataset table.
575    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    /// Loads only the exact output held by an active shared reservation.
599    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    /// Pivots a REDCap EAV CSV export into a versioned dataset table.
680    ///
681    /// # Errors
682    ///
683    /// Returns an error if the source file is missing, a version component is
684    /// negative, the REDCap CSV cannot be inspected, or DuckDB cannot create the
685    /// target relation.
686    #[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    /// Loads only the exact output held by an active shared reservation.
696    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    /// Compensates only a still-owned, non-admitted output of a known attempt.
824    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    /// Drops a dataset table if it exists.
842    ///
843    /// # Errors
844    ///
845    /// Returns an error if DuckDB rejects the `DROP TABLE` command.
846    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    /// Exports a versioned dataset table to CSV, NDJSON, Parquet, or Arrow IPC.
856    ///
857    /// # Errors
858    ///
859    /// Returns an error if the relation cannot be read, the destination
860    /// directory cannot be created, DuckDB export fails, or the Arrow/Parquet
861    /// writer fails.
862    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    /// Reads a bounded, stringified row preview for a versioned dataset table.
947    ///
948    /// # Errors
949    ///
950    /// Returns an error if the relation cannot be inspected or queried.
951    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
1140/// Canonical reservation key shared by new-name and existing-Asset selectors.
1141/// Version allocation occurs after this logical Dataset key is reserved.
1142pub 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
1152/// Builds the DuckDB relation name used for a versioned dataset table.
1153pub 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
1181/// Cleanup includes partial conversion files and staging tables on every exit.
1182/// Only a target successfully created by this call can be compensated.
1183fn 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}