Skip to main content

ahri_tre_app/service/
dataset_transform.rs

1use super::*;
2use ahri_tre_protocol::dataset::DatasetTransformRequest;
3use ahri_tre_types::NcName;
4
5impl AppService {
6    /// Resolve every supplied constraint before reserving the independent Output Study.
7    pub async fn transform_catalogue_dataset(
8        &self,
9        session: &mut DataStoreSession,
10        request: DatasetTransformRequest,
11    ) -> Result<DatasetMaterialization, AppError> {
12        if request.inputs.is_empty() || request.inputs.len() > 128 {
13            return Err(AppError::Validation("Invalid transform input count".into()));
14        }
15        let study = self
16            .resolve_study_selector(
17                self.catalogue_study_selector(request.study)?,
18                "Output Study",
19            )
20            .await?;
21        let mut inputs = std::collections::BTreeMap::new();
22        for input in request.inputs {
23            let (_, catalog, pinned) = self.resolve_catalogue_asset(input.asset).await?;
24            if catalog.asset.asset_type != ahri_tre_types::AssetType::Dataset {
25                return Err(AppError::Conflict(
26                    "Transform input must be a Dataset".into(),
27                ));
28            }
29            let version = self.catalogue_version(&catalog, pinned, input.version.as_deref())?;
30            if inputs.insert(input.alias, version.version_id).is_some() {
31                return Err(AppError::Conflict("Duplicate transform input alias".into()));
32            }
33        }
34        let actor = session
35            .authenticated_actor()
36            .ok_or_else(|| {
37                AppError::Validation("Authenticated transform actor is unavailable".into())
38            })?
39            .to_owned();
40        self.transform_dataset_with_lake_sql(
41            session,
42            DatasetSqlTransformRequest {
43                study_id: study.study_id,
44                dataset_asset_id: None,
45                dataset_name: NcName::parse(request.dataset.as_str().to_string())
46                    .map_err(|_| AppError::Validation("Invalid output Dataset name".into()))?,
47                dataset_version_id: None,
48                inputs,
49                views: request.views,
50                budgets: request.budgets,
51                sql: request.sql,
52                description: Some(request.description),
53                agent_instructions: None,
54                version_note: request.version_note,
55                transformation: ahri_tre_types::NewTransformationRecord {
56                    transformation_type: ahri_tre_types::TransformationType::Transform,
57                    description: "Arbitrary governed Dataset SQL transform; High output".into(),
58                    repository_url: None,
59                    commit_hash: None,
60                    file_path: None,
61                    date_created: None,
62                    created_by: Some(actor),
63                },
64            },
65        )
66        .await
67    }
68}