Skip to main content

ahri_tre_pgmeta/
write_models.rs

1//! PostgreSQL write models and insert/update helpers for metadata tables.
2//!
3//! Request structs in this module represent persistence-facing inputs. Domain
4//! repositories map validated domain records into these shapes before executing
5//! PostgreSQL writes.
6
7use crate::{PgMetaError, PgMetadataConnection};
8use ahri_tre_types::{
9    AssetId, AssetType, AssetVersionDuoRestrictionId, DomainId, EntityId, EntityRelationId,
10    KeyRole, StudyDuoRestrictionId, StudyId, StudyTypeId, TransformationId, TransformationType,
11    ValueTypeId, VariableId, VersionId, VocabularyId, VocabularyItemId,
12};
13use chrono::{DateTime, NaiveDate, Utc};
14use uuid::Uuid;
15
16/// Request to create a study row.
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct PgCreateStudy {
19    pub study_id: Option<StudyId>,
20    pub name: String,
21    pub description: Option<String>,
22    pub documentation: Option<String>,
23    pub agent_instructions: Option<String>,
24    pub external_id: Option<String>,
25    pub study_type_id: Option<StudyTypeId>,
26    pub date_created: Option<DateTime<Utc>>,
27    pub created_by: Option<String>,
28}
29
30/// Request to create a domain row.
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct PgCreateDomain {
33    pub domain_id: Option<DomainId>,
34    pub name: String,
35    pub uri: Option<String>,
36    pub description: Option<String>,
37}
38
39/// Request to create a vocabulary row.
40#[derive(Debug, Clone, PartialEq, Eq)]
41pub struct PgCreateVocabulary {
42    pub vocabulary_id: Option<VocabularyId>,
43    pub domain_id: DomainId,
44    pub name: String,
45    pub description: Option<String>,
46}
47
48/// Request to create a vocabulary item row.
49#[derive(Debug, Clone, PartialEq, Eq)]
50pub struct PgCreateVocabularyItem {
51    pub vocabulary_item_id: Option<VocabularyItemId>,
52    pub vocabulary_id: VocabularyId,
53    pub value: i32,
54    pub code: String,
55    pub description: Option<String>,
56}
57
58/// Request to create a vocabulary mapping row.
59#[derive(Debug, Clone, PartialEq, Eq)]
60pub struct PgCreateVocabularyMapping {
61    pub vocabulary_mapping_id: Option<i64>,
62    pub from_vocabulary_item: VocabularyItemId,
63    pub to_vocabulary_item: VocabularyItemId,
64}
65
66#[derive(Debug, Clone, PartialEq, Eq)]
67pub struct PgCreateVariable {
68    pub variable_id: Option<VariableId>,
69    pub domain_id: DomainId,
70    pub name: String,
71    pub value_type_id: ValueTypeId,
72    pub value_format: Option<String>,
73    pub vocabulary_id: Option<VocabularyId>,
74    pub key_role: KeyRole,
75    pub description: Option<String>,
76    pub note: Option<String>,
77    pub ontology_namespace: Option<String>,
78    pub ontology_class: Option<String>,
79}
80
81#[derive(Debug, Clone, PartialEq, Eq)]
82pub struct PgCreateAsset {
83    pub asset_id: Option<AssetId>,
84    pub study_id: StudyId,
85    pub name: String,
86    pub description: Option<String>,
87    pub agent_instructions: Option<String>,
88    pub asset_type: AssetType,
89    pub date_created: Option<DateTime<Utc>>,
90    pub created_by: Option<String>,
91}
92
93#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct PgCreateAssetVersion {
95    pub classification: ahri_tre_types::IngestClassification,
96    pub version_id: Option<VersionId>,
97    pub asset_id: AssetId,
98    pub major: Option<i32>,
99    pub minor: Option<i32>,
100    pub patch: Option<i32>,
101    pub version_note: Option<String>,
102    pub is_latest: Option<bool>,
103    pub doi: Option<String>,
104    pub created_at: Option<DateTime<Utc>>,
105    pub created_by: Option<String>,
106}
107
108#[derive(Debug, Clone, PartialEq, Eq)]
109pub struct PgCreateDataFile {
110    pub datafile_id: VersionId,
111    /// Size in bytes of the managed stored artifact, when known.
112    pub size_bytes: Option<i64>,
113    pub compressed: Option<bool>,
114    pub encrypted: Option<bool>,
115    pub compression_algorithm: Option<String>,
116    pub encryption_algorithm: Option<String>,
117    pub encryption_key: Option<Vec<u8>>,
118    pub storage_uri: String,
119    pub edam_format: String,
120    pub digest: String,
121}
122
123#[derive(Debug, Clone, PartialEq, Eq)]
124pub struct PgCreateTransformation {
125    pub transformation_id: Option<TransformationId>,
126    pub transformation_type: TransformationType,
127    pub description: String,
128    pub repository_url: Option<String>,
129    pub commit_hash: Option<String>,
130    pub file_path: Option<String>,
131    pub date_created: Option<DateTime<Utc>>,
132    pub created_by: Option<String>,
133}
134
135#[derive(Debug, Clone, PartialEq, Eq)]
136pub struct PgCreateDuoCode {
137    pub duo_id: String,
138    pub duo_code: String,
139    pub label: String,
140    pub definition: String,
141    pub category: String,
142    pub qualifier_type: String,
143    pub qualifier_required: bool,
144    pub qualifier_note: Option<String>,
145}
146
147#[derive(Debug, Clone, PartialEq, Eq, Default)]
148pub struct PgDuoQualifiers {
149    pub qualifier_ontology_namespace: Option<String>,
150    pub qualifier_ontology_id: Option<String>,
151    pub qualifier_text: Option<String>,
152    pub qualifier_date: Option<NaiveDate>,
153    pub qualifier_integer: Option<i32>,
154}
155
156#[derive(Debug, Clone, PartialEq, Eq)]
157pub struct PgCreateStudyDuoRestriction {
158    pub study_duo_restriction_id: Option<StudyDuoRestrictionId>,
159    pub study_id: StudyId,
160    pub duo_id: String,
161    pub qualifiers: PgDuoQualifiers,
162}
163
164#[derive(Debug, Clone, PartialEq, Eq)]
165pub struct PgCreateAssetVersionDuoRestriction {
166    pub asset_version_duo_restriction_id: Option<AssetVersionDuoRestrictionId>,
167    pub version_id: VersionId,
168    pub duo_id: String,
169    pub qualifiers: PgDuoQualifiers,
170}
171
172#[derive(Debug, Clone, PartialEq, Eq)]
173pub struct PgCreateEntity {
174    pub entity_id: Option<EntityId>,
175    pub domain_id: DomainId,
176    pub uuid: Option<Uuid>,
177    pub name: String,
178    pub description: Option<String>,
179    pub agent_instructions: Option<String>,
180    pub ontology_namespace: Option<String>,
181    pub ontology_class: Option<String>,
182}
183
184#[derive(Debug, Clone, PartialEq, Eq)]
185pub struct PgCreateEntityRelation {
186    pub entity_relation_id: Option<EntityRelationId>,
187    pub subject_entity_id: EntityId,
188    pub object_entity_id: EntityId,
189    pub domain_id: DomainId,
190    pub uuid: Option<Uuid>,
191    pub name: String,
192    pub description: Option<String>,
193    pub agent_instructions: Option<String>,
194    pub ontology_namespace: Option<String>,
195    pub ontology_class: Option<String>,
196}
197
198#[derive(Debug, Clone, PartialEq, Eq)]
199pub struct PgCreateEntityInstanceInStudy {
200    pub instance_id: Option<i64>,
201    pub uuid: Option<Uuid>,
202    pub study_id: StudyId,
203    pub entity_id: EntityId,
204    pub external_id: String,
205    pub label: Option<String>,
206    pub note: Option<String>,
207    pub transformation_id: TransformationId,
208}
209
210#[derive(Debug, Clone, PartialEq, Eq)]
211pub struct PgWriteEntityInstance {
212    pub instance_id: Option<i64>,
213    pub uuid: Option<Uuid>,
214    pub entity_id: EntityId,
215    pub label: Option<String>,
216    pub note: Option<String>,
217    pub transformation_id: TransformationId,
218}
219
220#[derive(Debug, Clone, PartialEq, Eq)]
221pub struct PgWriteStudyEntityInstance {
222    pub study_id: StudyId,
223    pub entity_instance_id: i64,
224    pub entity_id: EntityId,
225    pub external_id: String,
226    pub transformation_id: TransformationId,
227}
228
229#[derive(Debug, Clone, PartialEq, Eq)]
230pub struct PgCreateRelationInstanceInStudy {
231    pub relation_instance_id: Option<i64>,
232    pub uuid: Option<Uuid>,
233    pub study_id: StudyId,
234    pub entity_relation_id: EntityRelationId,
235    pub entity_instance_id_1: i64,
236    pub entity_instance_id_2: i64,
237    pub external_id: String,
238    pub valid_from: Option<NaiveDate>,
239    pub valid_to: Option<NaiveDate>,
240    pub note: Option<String>,
241    pub transformation_id: TransformationId,
242}
243
244#[derive(Debug, Clone, PartialEq, Eq)]
245pub struct PgWriteRelationInstance {
246    pub relation_instance_id: Option<i64>,
247    pub uuid: Option<Uuid>,
248    pub entity_relation_id: EntityRelationId,
249    pub entity_instance_id_1: i64,
250    pub entity_instance_id_2: i64,
251    pub valid_from: Option<NaiveDate>,
252    pub valid_to: Option<NaiveDate>,
253    pub note: Option<String>,
254    pub transformation_id: TransformationId,
255}
256
257#[derive(Debug, Clone, PartialEq, Eq)]
258pub struct PgWriteStudyRelationInstance {
259    pub study_id: StudyId,
260    pub relation_instance_id: i64,
261    pub entity_relation_id: EntityRelationId,
262    pub external_id: String,
263    pub transformation_id: TransformationId,
264}
265
266type DatasetVersionEntityLinkInput = (VersionId, i64, VariableId, TransformationId);
267type DatasetVersionRelationLinkInput = (
268    VersionId,
269    i64,
270    VariableId,
271    VariableId,
272    Option<VariableId>,
273    TransformationId,
274);
275
276impl PgMetadataConnection {
277    pub fn create_study(&mut self, input: PgCreateStudy) -> Result<StudyId, PgMetaError> {
278        let study_type_id = optional_i32(input.study_type_id.map(|id| id.0), "study_type_id")?;
279        let date_created = input.date_created;
280        let row = match input.study_id {
281            Some(study_id) => self.client.query_one(
282                "
283                INSERT INTO public.studies (
284                    study_id, name, description, documentation, agent_instructions,
285                    external_id, study_type_id, date_created, created_by
286                )
287                VALUES ($1, $2, $3, $4, $5, $6, $7, COALESCE($8, CURRENT_TIMESTAMP),
288                        COALESCE($9, CURRENT_USER))
289                RETURNING study_id
290                ",
291                &[
292                    &study_id.0,
293                    &input.name,
294                    &input.description,
295                    &input.documentation,
296                    &input.agent_instructions,
297                    &input.external_id,
298                    &study_type_id,
299                    &date_created,
300                    &input.created_by,
301                ],
302            ),
303            None => self.client.query_one(
304                "
305                INSERT INTO public.studies (
306                    name, description, documentation, agent_instructions,
307                    external_id, study_type_id, date_created, created_by
308                )
309                VALUES ($1, $2, $3, $4, $5, $6, COALESCE($7, CURRENT_TIMESTAMP),
310                        COALESCE($8, CURRENT_USER))
311                RETURNING study_id
312                ",
313                &[
314                    &input.name,
315                    &input.description,
316                    &input.documentation,
317                    &input.agent_instructions,
318                    &input.external_id,
319                    &study_type_id,
320                    &date_created,
321                    &input.created_by,
322                ],
323            ),
324        }
325        .map_err(|source| write_error("create study", source))?;
326        Ok(StudyId(row.get("study_id")))
327    }
328
329    pub fn create_domain(&mut self, input: PgCreateDomain) -> Result<DomainId, PgMetaError> {
330        let row = match input.domain_id {
331            Some(domain_id) => self.client.query_one(
332                "
333                INSERT INTO public.domains (domain_id, name, uri, description)
334                VALUES ($1, $2, $3, $4)
335                RETURNING domain_id
336                ",
337                &[
338                    &id_i32(domain_id.0, "domain_id")?,
339                    &input.name,
340                    &input.uri,
341                    &input.description,
342                ],
343            ),
344            None => self.client.query_one(
345                "
346                INSERT INTO public.domains (name, uri, description)
347                VALUES ($1, $2, $3)
348                RETURNING domain_id
349                ",
350                &[&input.name, &input.uri, &input.description],
351            ),
352        }
353        .map_err(|source| write_error("create domain", source))?;
354        Ok(DomainId(i64::from(row.get::<_, i32>("domain_id"))))
355    }
356
357    pub fn link_study_domain(
358        &mut self,
359        study_id: StudyId,
360        domain_id: DomainId,
361    ) -> Result<(), PgMetaError> {
362        self.client
363            .execute(
364                "
365                INSERT INTO public.study_domains (study_id, domain_id)
366                VALUES ($1, $2)
367                ON CONFLICT DO NOTHING
368                ",
369                &[&study_id.0, &id_i32(domain_id.0, "domain_id")?],
370            )
371            .map_err(|source| write_error("link study domain", source))?;
372        Ok(())
373    }
374
375    pub fn grant_study_access(
376        &mut self,
377        study_id: StudyId,
378        user_id: &str,
379    ) -> Result<(), PgMetaError> {
380        self.client
381            .query(
382                "SELECT public.tre_grant_study_access($1, $2)",
383                &[&study_id.0, &user_id],
384            )
385            .map_err(|source| write_error("grant study access", source))?;
386        Ok(())
387    }
388
389    pub fn add_study_delegate_custodian(
390        &mut self,
391        study_id: StudyId,
392        user_id: &str,
393    ) -> Result<(), PgMetaError> {
394        self.client
395            .query(
396                "SELECT public.tre_add_study_delegate_custodian($1, $2)",
397                &[&study_id.0, &user_id],
398            )
399            .map_err(|source| write_error("add study delegate custodian", source))?;
400        Ok(())
401    }
402
403    pub fn transfer_study_primary_custodianship(
404        &mut self,
405        study_id: StudyId,
406        user_id: &str,
407    ) -> Result<(), PgMetaError> {
408        self.client
409            .query(
410                "SELECT public.tre_transfer_study_primary_custodianship($1, $2)",
411                &[&study_id.0, &user_id],
412            )
413            .map_err(|source| write_error("transfer study primary custodianship", source))?;
414        Ok(())
415    }
416
417    pub fn remove_study_delegate_custodian(
418        &mut self,
419        study_id: StudyId,
420        user_id: &str,
421    ) -> Result<(), PgMetaError> {
422        self.client
423            .query(
424                "SELECT public.tre_remove_study_delegate_custodian($1, $2)",
425                &[&study_id.0, &user_id],
426            )
427            .map_err(|source| write_error("remove study delegate custodian", source))?;
428        Ok(())
429    }
430
431    pub fn revoke_study_access(
432        &mut self,
433        study_id: StudyId,
434        user_id: &str,
435    ) -> Result<(), PgMetaError> {
436        self.client
437            .query(
438                "SELECT public.tre_revoke_study_access($1, $2)",
439                &[&study_id.0, &user_id],
440            )
441            .map_err(|source| write_error("revoke study access", source))?;
442        Ok(())
443    }
444
445    pub fn grant_study_access_with_metadata(
446        &mut self,
447        study_id: StudyId,
448        user_id: &str,
449        date_granted: Option<DateTime<Utc>>,
450        granted_by: Option<&str>,
451    ) -> Result<(), PgMetaError> {
452        let date_granted = date_granted.map(|value| value.naive_utc());
453        self.client
454            .execute(
455                "
456                INSERT INTO public.study_access (study_id, user_id, date_granted, granted_by)
457                VALUES ($1, $2, COALESCE($3, CURRENT_TIMESTAMP), COALESCE($4, CURRENT_USER))
458                ON CONFLICT DO NOTHING
459                ",
460                &[&study_id.0, &user_id, &date_granted, &granted_by],
461            )
462            .map_err(|source| write_error("grant study access", source))?;
463        Ok(())
464    }
465
466    pub fn add_study_custodian(
467        &mut self,
468        study_id: StudyId,
469        user_id: &str,
470        is_primary: bool,
471    ) -> Result<(), PgMetaError> {
472        self.add_study_custodian_with_metadata(study_id, user_id, is_primary, None, None)
473    }
474
475    pub fn add_study_custodian_with_metadata(
476        &mut self,
477        study_id: StudyId,
478        user_id: &str,
479        is_primary: bool,
480        date_granted: Option<DateTime<Utc>>,
481        granted_by: Option<&str>,
482    ) -> Result<(), PgMetaError> {
483        let date_granted = date_granted.map(|value| value.naive_utc());
484        let mut transaction = self
485            .client
486            .transaction()
487            .map_err(PgMetaError::Transaction)?;
488        if is_primary {
489            transaction
490                .execute(
491                    "UPDATE public.study_custodians SET is_primary = FALSE WHERE study_id = $1",
492                    &[&study_id.0],
493                )
494                .map_err(|source| write_error("unset primary study custodian", source))?;
495        }
496        transaction
497            .execute(
498                "
499                INSERT INTO public.study_custodians (
500                    study_id, user_id, is_primary, date_granted, granted_by
501                )
502                VALUES ($1, $2, $3, COALESCE($4, CURRENT_TIMESTAMP), COALESCE($5, CURRENT_USER))
503                ON CONFLICT (study_id, user_id) DO UPDATE SET
504                    is_primary = EXCLUDED.is_primary
505                ",
506                &[
507                    &study_id.0,
508                    &user_id,
509                    &is_primary,
510                    &date_granted,
511                    &granted_by,
512                ],
513            )
514            .map_err(|source| write_error("add study custodian", source))?;
515        transaction.commit().map_err(PgMetaError::Transaction)?;
516        Ok(())
517    }
518
519    pub fn create_vocabulary(
520        &mut self,
521        input: PgCreateVocabulary,
522    ) -> Result<VocabularyId, PgMetaError> {
523        let row = match input.vocabulary_id {
524            Some(vocabulary_id) => self.client.query_one(
525                "
526                INSERT INTO public.vocabularies (vocabulary_id, domain_id, name, description)
527                VALUES ($1, $2, $3, $4)
528                RETURNING vocabulary_id
529                ",
530                &[
531                    &id_i32(vocabulary_id.0, "vocabulary_id")?,
532                    &id_i32(input.domain_id.0, "domain_id")?,
533                    &input.name,
534                    &input.description,
535                ],
536            ),
537            None => self.client.query_one(
538                "
539                INSERT INTO public.vocabularies (domain_id, name, description)
540                VALUES ($1, $2, $3)
541                RETURNING vocabulary_id
542                ",
543                &[
544                    &id_i32(input.domain_id.0, "domain_id")?,
545                    &input.name,
546                    &input.description,
547                ],
548            ),
549        }
550        .map_err(|source| write_error("create vocabulary", source))?;
551        Ok(VocabularyId(i64::from(row.get::<_, i32>("vocabulary_id"))))
552    }
553
554    pub fn create_vocabulary_item(
555        &mut self,
556        input: PgCreateVocabularyItem,
557    ) -> Result<VocabularyItemId, PgMetaError> {
558        let row = match input.vocabulary_item_id {
559            Some(vocabulary_item_id) => self.client.query_one(
560                "
561                INSERT INTO public.vocabulary_items (
562                    vocabulary_item_id, vocabulary_id, value, code, description
563                )
564                VALUES ($1, $2, $3, $4, $5)
565                RETURNING vocabulary_item_id
566                ",
567                &[
568                    &id_i32(vocabulary_item_id.0, "vocabulary_item_id")?,
569                    &id_i32(input.vocabulary_id.0, "vocabulary_id")?,
570                    &input.value,
571                    &input.code,
572                    &input.description,
573                ],
574            ),
575            None => self.client.query_one(
576                "
577                INSERT INTO public.vocabulary_items (vocabulary_id, value, code, description)
578                VALUES ($1, $2, $3, $4)
579                RETURNING vocabulary_item_id
580                ",
581                &[
582                    &id_i32(input.vocabulary_id.0, "vocabulary_id")?,
583                    &input.value,
584                    &input.code,
585                    &input.description,
586                ],
587            ),
588        }
589        .map_err(|source| write_error("create vocabulary item", source))?;
590        Ok(VocabularyItemId(i64::from(
591            row.get::<_, i32>("vocabulary_item_id"),
592        )))
593    }
594
595    pub fn create_vocabulary_mapping(
596        &mut self,
597        input: PgCreateVocabularyMapping,
598    ) -> Result<i64, PgMetaError> {
599        let row = match input.vocabulary_mapping_id {
600            Some(vocabulary_mapping_id) => self.client.query_one(
601                "
602                INSERT INTO public.vocabulary_mapping (
603                    vocabulary_mapping_id, from_vocabulary_item, to_vocabulary_item
604                )
605                VALUES ($1, $2, $3)
606                RETURNING vocabulary_mapping_id
607                ",
608                &[
609                    &id_i32(vocabulary_mapping_id, "vocabulary_mapping_id")?,
610                    &id_i32(input.from_vocabulary_item.0, "from_vocabulary_item")?,
611                    &id_i32(input.to_vocabulary_item.0, "to_vocabulary_item")?,
612                ],
613            ),
614            None => self.client.query_one(
615                "
616                INSERT INTO public.vocabulary_mapping (
617                    from_vocabulary_item, to_vocabulary_item
618                )
619                VALUES ($1, $2)
620                RETURNING vocabulary_mapping_id
621                ",
622                &[
623                    &id_i32(input.from_vocabulary_item.0, "from_vocabulary_item")?,
624                    &id_i32(input.to_vocabulary_item.0, "to_vocabulary_item")?,
625                ],
626            ),
627        }
628        .map_err(|source| write_error("create vocabulary mapping", source))?;
629        Ok(i64::from(row.get::<_, i32>("vocabulary_mapping_id")))
630    }
631
632    pub fn create_variable(&mut self, input: PgCreateVariable) -> Result<VariableId, PgMetaError> {
633        let value_type_id = id_i32(input.value_type_id.0, "value_type_id")?;
634        let vocabulary_id = optional_i32(input.vocabulary_id.map(|id| id.0), "vocabulary_id")?;
635        let key_role = key_role_sql(input.key_role)?;
636        let row = match input.variable_id {
637            Some(variable_id) => self.client.query_one(
638                "
639                INSERT INTO public.variables (
640                    variable_id, domain_id, name, value_type_id, value_format,
641                    vocabulary_id, keyrole, description, note,
642                    ontology_namespace, ontology_class
643                )
644                VALUES ($1, $2, $3, $4, $5, $6, CAST($7 AS text)::public.variable_keyrole_enum,
645                        $8, $9, $10, $11)
646                RETURNING variable_id
647                ",
648                &[
649                    &id_i32(variable_id.0, "variable_id")?,
650                    &id_i32(input.domain_id.0, "domain_id")?,
651                    &input.name,
652                    &value_type_id,
653                    &input.value_format,
654                    &vocabulary_id,
655                    &key_role,
656                    &input.description,
657                    &input.note,
658                    &input.ontology_namespace,
659                    &input.ontology_class,
660                ],
661            ),
662            None => self.client.query_one(
663                "
664                INSERT INTO public.variables (
665                    domain_id, name, value_type_id, value_format,
666                    vocabulary_id, keyrole, description, note,
667                    ontology_namespace, ontology_class
668                )
669                VALUES ($1, $2, $3, $4, $5, CAST($6 AS text)::public.variable_keyrole_enum,
670                        $7, $8, $9, $10)
671                RETURNING variable_id
672                ",
673                &[
674                    &id_i32(input.domain_id.0, "domain_id")?,
675                    &input.name,
676                    &value_type_id,
677                    &input.value_format,
678                    &vocabulary_id,
679                    &key_role,
680                    &input.description,
681                    &input.note,
682                    &input.ontology_namespace,
683                    &input.ontology_class,
684                ],
685            ),
686        }
687        .map_err(|source| write_error("create variable", source))?;
688        Ok(VariableId(i64::from(row.get::<_, i32>("variable_id"))))
689    }
690
691    pub fn update_variable(&mut self, input: PgCreateVariable) -> Result<(), PgMetaError> {
692        let variable_id = input.variable_id.ok_or_else(|| PgMetaError::Decode {
693            field: "variable_id",
694            value: "NULL".to_string(),
695            message: "variable update requires a variable_id".to_string(),
696        })?;
697        let value_type_id = id_i32(input.value_type_id.0, "value_type_id")?;
698        let vocabulary_id = optional_i32(input.vocabulary_id.map(|id| id.0), "vocabulary_id")?;
699        let key_role = key_role_sql(input.key_role)?;
700        self.client
701            .execute(
702                "
703                UPDATE public.variables
704                   SET value_type_id = $2,
705                       value_format = $3,
706                       vocabulary_id = $4,
707                       keyrole = CAST($5 AS text)::public.variable_keyrole_enum,
708                       description = $6,
709                       note = $7,
710                       ontology_namespace = $8,
711                       ontology_class = $9
712                 WHERE variable_id = $1
713                ",
714                &[
715                    &id_i32(variable_id.0, "variable_id")?,
716                    &value_type_id,
717                    &input.value_format,
718                    &vocabulary_id,
719                    &key_role,
720                    &input.description,
721                    &input.note,
722                    &input.ontology_namespace,
723                    &input.ontology_class,
724                ],
725            )
726            .map_err(|source| write_error("update variable", source))?;
727        Ok(())
728    }
729
730    pub fn update_vocabulary(&mut self, input: PgCreateVocabulary) -> Result<(), PgMetaError> {
731        let vocabulary_id = input.vocabulary_id.ok_or_else(|| PgMetaError::Decode {
732            field: "vocabulary_id",
733            value: "NULL".to_string(),
734            message: "vocabulary update requires a vocabulary_id".to_string(),
735        })?;
736        self.client
737            .execute(
738                "
739                UPDATE public.vocabularies
740                   SET name = $2,
741                       description = $3
742                 WHERE vocabulary_id = $1
743                ",
744                &[
745                    &id_i32(vocabulary_id.0, "vocabulary_id")?,
746                    &input.name,
747                    &input.description,
748                ],
749            )
750            .map_err(|source| write_error("update vocabulary", source))?;
751        Ok(())
752    }
753
754    pub fn update_vocabulary_item(
755        &mut self,
756        input: PgCreateVocabularyItem,
757    ) -> Result<(), PgMetaError> {
758        let vocabulary_item_id = input
759            .vocabulary_item_id
760            .ok_or_else(|| PgMetaError::Decode {
761                field: "vocabulary_item_id",
762                value: "NULL".to_string(),
763                message: "vocabulary item update requires a vocabulary_item_id".to_string(),
764            })?;
765        self.client
766            .execute(
767                "
768                UPDATE public.vocabulary_items
769                   SET vocabulary_id = $2,
770                       value = $3,
771                       code = $4,
772                       description = $5
773                 WHERE vocabulary_item_id = $1
774                ",
775                &[
776                    &id_i32(vocabulary_item_id.0, "vocabulary_item_id")?,
777                    &id_i32(input.vocabulary_id.0, "vocabulary_id")?,
778                    &input.value,
779                    &input.code,
780                    &input.description,
781                ],
782            )
783            .map_err(|source| write_error("update vocabulary item", source))?;
784        Ok(())
785    }
786
787    pub fn link_dataset_variable(
788        &mut self,
789        dataset_id: VersionId,
790        variable_id: VariableId,
791        row_role: KeyRole,
792    ) -> Result<(), PgMetaError> {
793        let row_role = key_role_sql(row_role)?;
794        self.client
795            .execute(
796                "
797                INSERT INTO public.dataset_variables (dataset_id, variable_id, row_role)
798                VALUES ($1, $2, CAST($3 AS text)::public.variable_keyrole_enum)
799                ON CONFLICT (dataset_id, variable_id) DO UPDATE SET row_role = EXCLUDED.row_role
800                ",
801                &[
802                    &dataset_id.0,
803                    &id_i32(variable_id.0, "variable_id")?,
804                    &row_role,
805                ],
806            )
807            .map_err(|source| write_error("link dataset variable", source))?;
808        Ok(())
809    }
810
811    pub fn create_transformation(
812        &mut self,
813        input: PgCreateTransformation,
814    ) -> Result<TransformationId, PgMetaError> {
815        let transformation_type = transformation_type_sql(input.transformation_type);
816        let date_created = input.date_created;
817        let row = match input.transformation_id {
818            Some(transformation_id) => self.client.query_one(
819                "
820                INSERT INTO public.transformations (
821                    transformation_id, transformation_type, description,
822                    repository_url, commit_hash, file_path, date_created, created_by
823                )
824                VALUES (
825                    $1, CAST($2 AS text)::public.transformation_type_enum, $3, $4, $5, $6,
826                    COALESCE($7, CURRENT_TIMESTAMP), COALESCE($8, CURRENT_USER)
827                )
828                RETURNING transformation_id
829                ",
830                &[
831                    &id_i32(transformation_id.0, "transformation_id")?,
832                    &transformation_type,
833                    &input.description,
834                    &input.repository_url,
835                    &input.commit_hash,
836                    &input.file_path,
837                    &date_created,
838                    &input.created_by,
839                ],
840            ),
841            None => self.client.query_one(
842                "
843                INSERT INTO public.transformations (
844                    transformation_type, description, repository_url, commit_hash, file_path,
845                    date_created, created_by
846                )
847                VALUES (
848                    CAST($1 AS text)::public.transformation_type_enum, $2, $3, $4, $5,
849                    COALESCE($6, CURRENT_TIMESTAMP), COALESCE($7, CURRENT_USER)
850                )
851                RETURNING transformation_id
852                ",
853                &[
854                    &transformation_type,
855                    &input.description,
856                    &input.repository_url,
857                    &input.commit_hash,
858                    &input.file_path,
859                    &date_created,
860                    &input.created_by,
861                ],
862            ),
863        }
864        .map_err(|source| write_error("create transformation", source))?;
865        Ok(TransformationId(i64::from(
866            row.get::<_, i32>("transformation_id"),
867        )))
868    }
869
870    pub fn delete_transformation(
871        &mut self,
872        transformation_id: TransformationId,
873    ) -> Result<bool, PgMetaError> {
874        let deleted = self
875            .client
876            .execute(
877                "DELETE FROM public.transformations WHERE transformation_id = $1",
878                &[&id_i32(transformation_id.0, "transformation_id")?],
879            )
880            .map_err(|source| write_error("delete transformation", source))?;
881        Ok(deleted > 0)
882    }
883
884    pub fn create_asset(&mut self, input: PgCreateAsset) -> Result<AssetId, PgMetaError> {
885        let asset_type = asset_type_sql(input.asset_type);
886        let date_created = input.date_created;
887        let row = match input.asset_id {
888            Some(asset_id) => self.client.query_one(
889                "
890                INSERT INTO public.assets (
891                    asset_id, study_id, name, description, agent_instructions, asset_type,
892                    date_created, created_by
893                )
894                VALUES ($1, $2, $3, $4, $5, CAST($6 AS text)::public.asset_type_enum,
895                        COALESCE($7, CURRENT_TIMESTAMP), COALESCE($8, CURRENT_USER))
896                RETURNING asset_id
897                ",
898                &[
899                    &asset_id.0,
900                    &input.study_id.0,
901                    &input.name,
902                    &input.description,
903                    &input.agent_instructions,
904                    &asset_type,
905                    &date_created,
906                    &input.created_by,
907                ],
908            ),
909            None => self.client.query_one(
910                "
911                INSERT INTO public.assets (
912                    study_id, name, description, agent_instructions, asset_type,
913                    date_created, created_by
914                )
915                VALUES ($1, $2, $3, $4, CAST($5 AS text)::public.asset_type_enum,
916                        COALESCE($6, CURRENT_TIMESTAMP), COALESCE($7, CURRENT_USER))
917                RETURNING asset_id
918                ",
919                &[
920                    &input.study_id.0,
921                    &input.name,
922                    &input.description,
923                    &input.agent_instructions,
924                    &asset_type,
925                    &date_created,
926                    &input.created_by,
927                ],
928            ),
929        }
930        .map_err(|source| write_error("create asset", source))?;
931        Ok(AssetId(row.get("asset_id")))
932    }
933
934    pub fn create_asset_version(
935        &mut self,
936        input: PgCreateAssetVersion,
937    ) -> Result<VersionId, PgMetaError> {
938        if self.admission_active {
939            return insert_asset_version(&mut self.client, input);
940        }
941        let mut transaction = self
942            .client
943            .transaction()
944            .map_err(PgMetaError::Transaction)?;
945        let version_id = insert_asset_version(&mut transaction, input)?;
946        transaction.commit().map_err(PgMetaError::Transaction)?;
947        Ok(version_id)
948    }
949
950    pub fn create_dataset(&mut self, dataset_id: VersionId) -> Result<(), PgMetaError> {
951        self.client
952            .execute(
953                "INSERT INTO public.datasets (dataset_id) VALUES ($1)",
954                &[&dataset_id.0],
955            )
956            .map_err(|source| write_error("create dataset", source))?;
957        Ok(())
958    }
959
960    pub fn create_datafile(&mut self, input: PgCreateDataFile) -> Result<(), PgMetaError> {
961        self.client
962            .execute(
963                "
964                INSERT INTO public.datafiles (
965                    datafile_id, size_bytes, compressed, encrypted, compression_algorithm,
966                    encryption_algorithm, encryption_key, storage_uri, edam_format, digest
967                )
968                VALUES (
969                    $1, $2, COALESCE($3, FALSE), COALESCE($4, FALSE),
970                    COALESCE($5, 'zstd'), COALESCE($6, 'AES-256-CBC with PBKDF2'),
971                    $7, $8, $9, $10
972                )
973                ",
974                &[
975                    &input.datafile_id.0,
976                    &input.size_bytes,
977                    &input.compressed,
978                    &input.encrypted,
979                    &input.compression_algorithm,
980                    &input.encryption_algorithm,
981                    &input.encryption_key,
982                    &input.storage_uri,
983                    &input.edam_format,
984                    &input.digest,
985                ],
986            )
987            .map_err(|source| write_error("create datafile", source))?;
988        Ok(())
989    }
990
991    pub fn create_dataset_version(
992        &mut self,
993        asset: PgCreateAsset,
994        version: PgCreateAssetVersion,
995    ) -> Result<VersionId, PgMetaError> {
996        let mut transaction = self
997            .client
998            .transaction()
999            .map_err(PgMetaError::Transaction)?;
1000        let asset_id = insert_asset(&mut transaction, asset)?;
1001        let mut version = version;
1002        version.asset_id = asset_id;
1003        let version_id = insert_asset_version(&mut transaction, version)?;
1004        transaction
1005            .execute(
1006                "INSERT INTO public.datasets (dataset_id) VALUES ($1)",
1007                &[&version_id.0],
1008            )
1009            .map_err(|source| write_error("create dataset", source))?;
1010        transaction.commit().map_err(PgMetaError::Transaction)?;
1011        Ok(version_id)
1012    }
1013
1014    pub fn create_datafile_version(
1015        &mut self,
1016        asset: PgCreateAsset,
1017        version: PgCreateAssetVersion,
1018        datafile: PgCreateDataFile,
1019    ) -> Result<VersionId, PgMetaError> {
1020        let mut transaction = self
1021            .client
1022            .transaction()
1023            .map_err(PgMetaError::Transaction)?;
1024        let asset_id = insert_asset(&mut transaction, asset)?;
1025        let mut version = version;
1026        version.asset_id = asset_id;
1027        let version_id = insert_asset_version(&mut transaction, version)?;
1028        let mut datafile = datafile;
1029        datafile.datafile_id = version_id;
1030        insert_datafile(&mut transaction, datafile)?;
1031        transaction.commit().map_err(PgMetaError::Transaction)?;
1032        Ok(version_id)
1033    }
1034
1035    pub fn create_duo_code(&mut self, input: PgCreateDuoCode) -> Result<(), PgMetaError> {
1036        self.client
1037            .execute(
1038                "
1039                INSERT INTO public.duo_codes (
1040                    duo_id, duo_code, label, definition, category,
1041                    qualifier_type, qualifier_required, qualifier_note
1042                )
1043                VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
1044                ",
1045                &[
1046                    &input.duo_id,
1047                    &input.duo_code,
1048                    &input.label,
1049                    &input.definition,
1050                    &input.category,
1051                    &input.qualifier_type,
1052                    &input.qualifier_required,
1053                    &input.qualifier_note,
1054                ],
1055            )
1056            .map_err(|source| write_error("create DUO code", source))?;
1057        Ok(())
1058    }
1059
1060    pub fn add_study_duo_restriction(
1061        &mut self,
1062        study_id: StudyId,
1063        duo_id: &str,
1064        qualifiers: PgDuoQualifiers,
1065    ) -> Result<i64, PgMetaError> {
1066        self.create_study_duo_restriction(PgCreateStudyDuoRestriction {
1067            study_duo_restriction_id: None,
1068            study_id,
1069            duo_id: duo_id.to_string(),
1070            qualifiers,
1071        })
1072    }
1073
1074    pub fn create_study_duo_restriction(
1075        &mut self,
1076        input: PgCreateStudyDuoRestriction,
1077    ) -> Result<i64, PgMetaError> {
1078        let row = match input.study_duo_restriction_id {
1079            Some(restriction_id) => self.client.query_one(
1080                "
1081                INSERT INTO public.study_duo_restrictions (
1082                    study_duo_restriction_id, study_id, duo_id,
1083                    qualifier_ontology_namespace, qualifier_ontology_id,
1084                    qualifier_text, qualifier_date, qualifier_integer
1085                )
1086                VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
1087                RETURNING study_duo_restriction_id
1088                ",
1089                &[
1090                    &id_i32(restriction_id.0, "study_duo_restriction_id")?,
1091                    &input.study_id.0,
1092                    &input.duo_id,
1093                    &input.qualifiers.qualifier_ontology_namespace,
1094                    &input.qualifiers.qualifier_ontology_id,
1095                    &input.qualifiers.qualifier_text,
1096                    &input.qualifiers.qualifier_date,
1097                    &input.qualifiers.qualifier_integer,
1098                ],
1099            ),
1100            None => self.client.query_one(
1101                "
1102                INSERT INTO public.study_duo_restrictions (
1103                    study_id, duo_id, qualifier_ontology_namespace, qualifier_ontology_id,
1104                    qualifier_text, qualifier_date, qualifier_integer
1105                )
1106                VALUES ($1, $2, $3, $4, $5, $6, $7)
1107                RETURNING study_duo_restriction_id
1108                ",
1109                &[
1110                    &input.study_id.0,
1111                    &input.duo_id,
1112                    &input.qualifiers.qualifier_ontology_namespace,
1113                    &input.qualifiers.qualifier_ontology_id,
1114                    &input.qualifiers.qualifier_text,
1115                    &input.qualifiers.qualifier_date,
1116                    &input.qualifiers.qualifier_integer,
1117                ],
1118            ),
1119        }
1120        .map_err(|source| write_error("create study DUO restriction", source))?;
1121        Ok(i64::from(row.get::<_, i32>("study_duo_restriction_id")))
1122    }
1123
1124    pub fn add_asset_version_duo_restriction(
1125        &mut self,
1126        version_id: VersionId,
1127        duo_id: &str,
1128        qualifiers: PgDuoQualifiers,
1129    ) -> Result<i64, PgMetaError> {
1130        self.create_asset_version_duo_restriction(PgCreateAssetVersionDuoRestriction {
1131            asset_version_duo_restriction_id: None,
1132            version_id,
1133            duo_id: duo_id.to_string(),
1134            qualifiers,
1135        })
1136    }
1137
1138    pub fn create_asset_version_duo_restriction(
1139        &mut self,
1140        input: PgCreateAssetVersionDuoRestriction,
1141    ) -> Result<i64, PgMetaError> {
1142        let row = match input.asset_version_duo_restriction_id {
1143            Some(restriction_id) => self.client.query_one(
1144                "
1145                INSERT INTO public.asset_version_duo_restrictions (
1146                    asset_version_duo_restriction_id, version_id, duo_id,
1147                    qualifier_ontology_namespace, qualifier_ontology_id,
1148                    qualifier_text, qualifier_date, qualifier_integer
1149                )
1150                VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
1151                RETURNING asset_version_duo_restriction_id
1152                ",
1153                &[
1154                    &id_i32(restriction_id.0, "asset_version_duo_restriction_id")?,
1155                    &input.version_id.0,
1156                    &input.duo_id,
1157                    &input.qualifiers.qualifier_ontology_namespace,
1158                    &input.qualifiers.qualifier_ontology_id,
1159                    &input.qualifiers.qualifier_text,
1160                    &input.qualifiers.qualifier_date,
1161                    &input.qualifiers.qualifier_integer,
1162                ],
1163            ),
1164            None => self.client.query_one(
1165                "
1166                INSERT INTO public.asset_version_duo_restrictions (
1167                    version_id, duo_id, qualifier_ontology_namespace, qualifier_ontology_id,
1168                    qualifier_text, qualifier_date, qualifier_integer
1169                )
1170                VALUES ($1, $2, $3, $4, $5, $6, $7)
1171                RETURNING asset_version_duo_restriction_id
1172                ",
1173                &[
1174                    &input.version_id.0,
1175                    &input.duo_id,
1176                    &input.qualifiers.qualifier_ontology_namespace,
1177                    &input.qualifiers.qualifier_ontology_id,
1178                    &input.qualifiers.qualifier_text,
1179                    &input.qualifiers.qualifier_date,
1180                    &input.qualifiers.qualifier_integer,
1181                ],
1182            ),
1183        }
1184        .map_err(|source| write_error("create asset version DUO restriction", source))?;
1185        Ok(i64::from(
1186            row.get::<_, i32>("asset_version_duo_restriction_id"),
1187        ))
1188    }
1189
1190    pub fn create_entity(&mut self, input: PgCreateEntity) -> Result<EntityId, PgMetaError> {
1191        let row = match input.entity_id {
1192            Some(entity_id) => self.client.query_one(
1193                "
1194                INSERT INTO public.entities (
1195                    entity_id, domain_id, uuid, name, description,
1196                    agent_instructions, ontology_namespace, ontology_class
1197                )
1198                VALUES ($1, $2, COALESCE($3, uuidv7()), $4, $5, $6, $7, $8)
1199                RETURNING entity_id
1200                ",
1201                &[
1202                    &id_i32(entity_id.0, "entity_id")?,
1203                    &id_i32(input.domain_id.0, "domain_id")?,
1204                    &input.uuid,
1205                    &input.name,
1206                    &input.description,
1207                    &input.agent_instructions,
1208                    &input.ontology_namespace,
1209                    &input.ontology_class,
1210                ],
1211            ),
1212            None => self.client.query_one(
1213                "
1214                INSERT INTO public.entities (
1215                    domain_id, uuid, name, description, agent_instructions,
1216                    ontology_namespace, ontology_class
1217                )
1218                VALUES ($1, COALESCE($2, uuidv7()), $3, $4, $5, $6, $7)
1219                RETURNING entity_id
1220                ",
1221                &[
1222                    &id_i32(input.domain_id.0, "domain_id")?,
1223                    &input.uuid,
1224                    &input.name,
1225                    &input.description,
1226                    &input.agent_instructions,
1227                    &input.ontology_namespace,
1228                    &input.ontology_class,
1229                ],
1230            ),
1231        }
1232        .map_err(|source| write_error("create entity", source))?;
1233        Ok(EntityId(i64::from(row.get::<_, i32>("entity_id"))))
1234    }
1235
1236    pub fn create_entity_relation(
1237        &mut self,
1238        input: PgCreateEntityRelation,
1239    ) -> Result<EntityRelationId, PgMetaError> {
1240        let row = match input.entity_relation_id {
1241            Some(entity_relation_id) => self.client.query_one(
1242                "
1243                INSERT INTO public.entityrelations (
1244                    entityrelation_id, subject_entity_id, object_entity_id,
1245                    domain_id, uuid, name, description, agent_instructions,
1246                    ontology_namespace, ontology_class
1247                )
1248                VALUES ($1, $2, $3, $4, COALESCE($5, uuidv7()), $6, $7, $8, $9, $10)
1249                RETURNING entityrelation_id
1250                ",
1251                &[
1252                    &id_i32(entity_relation_id.0, "entityrelation_id")?,
1253                    &id_i32(input.subject_entity_id.0, "subject_entity_id")?,
1254                    &id_i32(input.object_entity_id.0, "object_entity_id")?,
1255                    &id_i32(input.domain_id.0, "domain_id")?,
1256                    &input.uuid,
1257                    &input.name,
1258                    &input.description,
1259                    &input.agent_instructions,
1260                    &input.ontology_namespace,
1261                    &input.ontology_class,
1262                ],
1263            ),
1264            None => self.client.query_one(
1265                "
1266                INSERT INTO public.entityrelations (
1267                    subject_entity_id, object_entity_id, domain_id, uuid, name, description,
1268                    agent_instructions, ontology_namespace, ontology_class
1269                )
1270                VALUES ($1, $2, $3, COALESCE($4, uuidv7()), $5, $6, $7, $8, $9)
1271                RETURNING entityrelation_id
1272                ",
1273                &[
1274                    &id_i32(input.subject_entity_id.0, "subject_entity_id")?,
1275                    &id_i32(input.object_entity_id.0, "object_entity_id")?,
1276                    &id_i32(input.domain_id.0, "domain_id")?,
1277                    &input.uuid,
1278                    &input.name,
1279                    &input.description,
1280                    &input.agent_instructions,
1281                    &input.ontology_namespace,
1282                    &input.ontology_class,
1283                ],
1284            ),
1285        }
1286        .map_err(|source| write_error("create entity relation", source))?;
1287        Ok(EntityRelationId(i64::from(
1288            row.get::<_, i32>("entityrelation_id"),
1289        )))
1290    }
1291
1292    pub fn create_entity_instance_in_study(
1293        &mut self,
1294        input: PgCreateEntityInstanceInStudy,
1295    ) -> Result<i64, PgMetaError> {
1296        let mut transaction = self
1297            .client
1298            .transaction()
1299            .map_err(PgMetaError::Transaction)?;
1300        let row = match input.instance_id {
1301            Some(instance_id) => transaction.query_one(
1302                "
1303                INSERT INTO public.entity_instances (
1304                    instance_id, uuid, entity_id, label, note, transformation_id
1305                )
1306                VALUES ($1, COALESCE($2, uuidv7()), $3, $4, $5, $6)
1307                RETURNING instance_id
1308                ",
1309                &[
1310                    &instance_id,
1311                    &input.uuid,
1312                    &id_i32(input.entity_id.0, "entity_id")?,
1313                    &input.label,
1314                    &input.note,
1315                    &id_i32(input.transformation_id.0, "transformation_id")?,
1316                ],
1317            ),
1318            None => transaction.query_one(
1319                "
1320                INSERT INTO public.entity_instances (uuid, entity_id, label, note, transformation_id)
1321                VALUES (COALESCE($1, uuidv7()), $2, $3, $4, $5)
1322                RETURNING instance_id
1323                ",
1324                &[
1325                    &input.uuid,
1326                    &id_i32(input.entity_id.0, "entity_id")?,
1327                    &input.label,
1328                    &input.note,
1329                    &id_i32(input.transformation_id.0, "transformation_id")?,
1330                ],
1331            ),
1332        }
1333        .map_err(|source| write_error("create entity instance", source))?;
1334        let instance_id = row.get("instance_id");
1335        transaction
1336            .execute(
1337                "
1338                INSERT INTO public.study_entity_instances (
1339                    study_id, entity_instance_id, entity_id, external_id, transformation_id
1340                )
1341                VALUES ($1, $2, $3, $4, $5)
1342                ",
1343                &[
1344                    &input.study_id.0,
1345                    &instance_id,
1346                    &id_i32(input.entity_id.0, "entity_id")?,
1347                    &input.external_id,
1348                    &id_i32(input.transformation_id.0, "transformation_id")?,
1349                ],
1350            )
1351            .map_err(|source| write_error("link study entity instance", source))?;
1352        transaction.commit().map_err(PgMetaError::Transaction)?;
1353        Ok(instance_id)
1354    }
1355
1356    pub fn create_relation_instance_in_study(
1357        &mut self,
1358        input: PgCreateRelationInstanceInStudy,
1359    ) -> Result<i64, PgMetaError> {
1360        let mut transaction = self
1361            .client
1362            .transaction()
1363            .map_err(PgMetaError::Transaction)?;
1364        let row = match input.relation_instance_id {
1365            Some(relation_instance_id) => transaction.query_one(
1366                "
1367                INSERT INTO public.relation_instances (
1368                    relation_instance_id, uuid, entityrelation_id,
1369                    entity_instance_id_1, entity_instance_id_2,
1370                    valid_from, valid_to, note, transformation_id
1371                )
1372                VALUES ($1, COALESCE($2, uuidv7()), $3, $4, $5, $6, $7, $8, $9)
1373                RETURNING relation_instance_id
1374                ",
1375                &[
1376                    &relation_instance_id,
1377                    &input.uuid,
1378                    &id_i32(input.entity_relation_id.0, "entityrelation_id")?,
1379                    &input.entity_instance_id_1,
1380                    &input.entity_instance_id_2,
1381                    &input.valid_from,
1382                    &input.valid_to,
1383                    &input.note,
1384                    &id_i32(input.transformation_id.0, "transformation_id")?,
1385                ],
1386            ),
1387            None => transaction.query_one(
1388                "
1389                INSERT INTO public.relation_instances (
1390                    uuid, entityrelation_id, entity_instance_id_1, entity_instance_id_2,
1391                    valid_from, valid_to, note, transformation_id
1392                )
1393                VALUES (COALESCE($1, uuidv7()), $2, $3, $4, $5, $6, $7, $8)
1394                RETURNING relation_instance_id
1395                ",
1396                &[
1397                    &input.uuid,
1398                    &id_i32(input.entity_relation_id.0, "entityrelation_id")?,
1399                    &input.entity_instance_id_1,
1400                    &input.entity_instance_id_2,
1401                    &input.valid_from,
1402                    &input.valid_to,
1403                    &input.note,
1404                    &id_i32(input.transformation_id.0, "transformation_id")?,
1405                ],
1406            ),
1407        }
1408        .map_err(|source| write_error("create relation instance", source))?;
1409        let relation_instance_id = row.get("relation_instance_id");
1410        transaction
1411            .execute(
1412                "
1413                INSERT INTO public.study_relation_instances (
1414                    study_id, relation_instance_id, entityrelation_id,
1415                    external_id, transformation_id
1416                )
1417                VALUES ($1, $2, $3, $4, $5)
1418                ",
1419                &[
1420                    &input.study_id.0,
1421                    &relation_instance_id,
1422                    &id_i32(input.entity_relation_id.0, "entityrelation_id")?,
1423                    &input.external_id,
1424                    &id_i32(input.transformation_id.0, "transformation_id")?,
1425                ],
1426            )
1427            .map_err(|source| write_error("link study relation instance", source))?;
1428        transaction.commit().map_err(PgMetaError::Transaction)?;
1429        Ok(relation_instance_id)
1430    }
1431
1432    pub fn save_entity_instance(
1433        &mut self,
1434        input: PgWriteEntityInstance,
1435    ) -> Result<i64, PgMetaError> {
1436        let row = match input.instance_id {
1437            Some(instance_id) => self.client.query_one(
1438                "
1439                UPDATE public.entity_instances
1440                   SET entity_id = $2,
1441                       label = $3,
1442                       note = $4,
1443                       transformation_id = $5
1444                 WHERE instance_id = $1
1445                   AND ($6::uuid IS NULL OR uuid = $6)
1446                RETURNING instance_id
1447                ",
1448                &[
1449                    &instance_id,
1450                    &id_i32(input.entity_id.0, "entity_id")?,
1451                    &input.label,
1452                    &input.note,
1453                    &id_i32(input.transformation_id.0, "transformation_id")?,
1454                    &input.uuid,
1455                ],
1456            ),
1457            None => self.client.query_one(
1458                "
1459                INSERT INTO public.entity_instances (uuid, entity_id, label, note, transformation_id)
1460                VALUES (COALESCE($1, uuidv7()), $2, $3, $4, $5)
1461                RETURNING instance_id
1462                ",
1463                &[
1464                    &input.uuid,
1465                    &id_i32(input.entity_id.0, "entity_id")?,
1466                    &input.label,
1467                    &input.note,
1468                    &id_i32(input.transformation_id.0, "transformation_id")?,
1469                ],
1470            ),
1471        }
1472        .map_err(|source| write_error("save entity instance", source))?;
1473        Ok(row.get("instance_id"))
1474    }
1475
1476    pub fn save_study_entity_instance(
1477        &mut self,
1478        input: PgWriteStudyEntityInstance,
1479    ) -> Result<(), PgMetaError> {
1480        self.client
1481            .execute(
1482                "
1483                INSERT INTO public.study_entity_instances (
1484                    study_id, entity_instance_id, entity_id, external_id, transformation_id
1485                )
1486                VALUES ($1, $2, $3, $4, $5)
1487                ON CONFLICT (study_id, entity_instance_id) DO UPDATE
1488                  SET entity_id = EXCLUDED.entity_id,
1489                      external_id = EXCLUDED.external_id,
1490                      transformation_id = EXCLUDED.transformation_id
1491                ",
1492                &[
1493                    &input.study_id.0,
1494                    &input.entity_instance_id,
1495                    &id_i32(input.entity_id.0, "entity_id")?,
1496                    &input.external_id,
1497                    &id_i32(input.transformation_id.0, "transformation_id")?,
1498                ],
1499            )
1500            .map_err(|source| write_error("save study entity instance", source))?;
1501        Ok(())
1502    }
1503
1504    pub fn save_relation_instance(
1505        &mut self,
1506        input: PgWriteRelationInstance,
1507    ) -> Result<i64, PgMetaError> {
1508        let row = match input.relation_instance_id {
1509            Some(relation_instance_id) => self.client.query_one(
1510                "
1511                UPDATE public.relation_instances
1512                   SET entityrelation_id = $2,
1513                       entity_instance_id_1 = $3,
1514                       entity_instance_id_2 = $4,
1515                       valid_from = $5,
1516                       valid_to = $6,
1517                       note = $7,
1518                       transformation_id = $8
1519                 WHERE relation_instance_id = $1
1520                   AND ($9::uuid IS NULL OR uuid = $9)
1521                RETURNING relation_instance_id
1522                ",
1523                &[
1524                    &relation_instance_id,
1525                    &id_i32(input.entity_relation_id.0, "entityrelation_id")?,
1526                    &input.entity_instance_id_1,
1527                    &input.entity_instance_id_2,
1528                    &input.valid_from,
1529                    &input.valid_to,
1530                    &input.note,
1531                    &id_i32(input.transformation_id.0, "transformation_id")?,
1532                    &input.uuid,
1533                ],
1534            ),
1535            None => self.client.query_one(
1536                "
1537                INSERT INTO public.relation_instances (
1538                    uuid, entityrelation_id, entity_instance_id_1, entity_instance_id_2,
1539                    valid_from, valid_to, note, transformation_id
1540                )
1541                VALUES (COALESCE($1, uuidv7()), $2, $3, $4, $5, $6, $7, $8)
1542                RETURNING relation_instance_id
1543                ",
1544                &[
1545                    &input.uuid,
1546                    &id_i32(input.entity_relation_id.0, "entityrelation_id")?,
1547                    &input.entity_instance_id_1,
1548                    &input.entity_instance_id_2,
1549                    &input.valid_from,
1550                    &input.valid_to,
1551                    &input.note,
1552                    &id_i32(input.transformation_id.0, "transformation_id")?,
1553                ],
1554            ),
1555        }
1556        .map_err(|source| write_error("save relation instance", source))?;
1557        Ok(row.get("relation_instance_id"))
1558    }
1559
1560    pub fn save_study_relation_instance(
1561        &mut self,
1562        input: PgWriteStudyRelationInstance,
1563    ) -> Result<(), PgMetaError> {
1564        self.client
1565            .execute(
1566                "
1567                INSERT INTO public.study_relation_instances (
1568                    study_id, relation_instance_id, entityrelation_id,
1569                    external_id, transformation_id
1570                )
1571                VALUES ($1, $2, $3, $4, $5)
1572                ON CONFLICT (study_id, relation_instance_id) DO UPDATE
1573                  SET entityrelation_id = EXCLUDED.entityrelation_id,
1574                      external_id = EXCLUDED.external_id,
1575                      transformation_id = EXCLUDED.transformation_id
1576                ",
1577                &[
1578                    &input.study_id.0,
1579                    &input.relation_instance_id,
1580                    &id_i32(input.entity_relation_id.0, "entityrelation_id")?,
1581                    &input.external_id,
1582                    &id_i32(input.transformation_id.0, "transformation_id")?,
1583                ],
1584            )
1585            .map_err(|source| write_error("save study relation instance", source))?;
1586        Ok(())
1587    }
1588
1589    pub fn link_asset_version_entity(
1590        &mut self,
1591        version_id: VersionId,
1592        entity_instance_id: i64,
1593        transformation_id: TransformationId,
1594    ) -> Result<(), PgMetaError> {
1595        self.client
1596            .execute(
1597                "
1598                INSERT INTO public.asset_version_entities (
1599                    version_id, entity_instance_id, transformation_id
1600                )
1601                VALUES ($1, $2, $3)
1602                ON CONFLICT (version_id, entity_instance_id) DO UPDATE
1603                  SET transformation_id = EXCLUDED.transformation_id
1604                ",
1605                &[
1606                    &version_id.0,
1607                    &entity_instance_id,
1608                    &id_i32(transformation_id.0, "transformation_id")?,
1609                ],
1610            )
1611            .map_err(|source| write_error("link asset version entity", source))?;
1612        Ok(())
1613    }
1614
1615    pub fn add_dataset_version_entity(
1616        &mut self,
1617        version_id: VersionId,
1618        entity_instance_id: i64,
1619        entity_variable_id: VariableId,
1620        transformation_id: TransformationId,
1621    ) -> Result<(), PgMetaError> {
1622        let mut transaction = self
1623            .client
1624            .transaction()
1625            .map_err(PgMetaError::Transaction)?;
1626        transaction
1627            .execute(
1628                "
1629                INSERT INTO public.asset_version_entities (
1630                    version_id, entity_instance_id, transformation_id
1631                )
1632                VALUES ($1, $2, $3)
1633                ON CONFLICT (version_id, entity_instance_id) DO UPDATE
1634                  SET transformation_id = EXCLUDED.transformation_id
1635                ",
1636                &[
1637                    &version_id.0,
1638                    &entity_instance_id,
1639                    &id_i32(transformation_id.0, "transformation_id")?,
1640                ],
1641            )
1642            .map_err(|source| write_error("link asset version entity", source))?;
1643        transaction
1644            .execute(
1645                "
1646                INSERT INTO public.dataset_version_entities (
1647                    version_id, entity_instance_id, entity_variable_id, transformation_id
1648                )
1649                VALUES ($1, $2, $3, $4)
1650                ON CONFLICT (version_id, entity_instance_id) DO UPDATE
1651                  SET entity_variable_id = EXCLUDED.entity_variable_id,
1652                      transformation_id = EXCLUDED.transformation_id
1653                ",
1654                &[
1655                    &version_id.0,
1656                    &entity_instance_id,
1657                    &id_i32(entity_variable_id.0, "entity_variable_id")?,
1658                    &id_i32(transformation_id.0, "transformation_id")?,
1659                ],
1660            )
1661            .map_err(|source| write_error("link dataset version entity", source))?;
1662        transaction.commit().map_err(PgMetaError::Transaction)?;
1663        Ok(())
1664    }
1665
1666    pub fn link_dataset_version_entities(
1667        &mut self,
1668        links: &[DatasetVersionEntityLinkInput],
1669    ) -> Result<(), PgMetaError> {
1670        let mut transaction = self
1671            .client
1672            .transaction()
1673            .map_err(PgMetaError::Transaction)?;
1674        for (version_id, entity_instance_id, entity_variable_id, transformation_id) in links {
1675            transaction
1676                .execute(
1677                    "
1678                    INSERT INTO public.asset_version_entities (
1679                        version_id, entity_instance_id, transformation_id
1680                    )
1681                    VALUES ($1, $2, $3)
1682                    ON CONFLICT (version_id, entity_instance_id) DO UPDATE
1683                      SET transformation_id = EXCLUDED.transformation_id
1684                    ",
1685                    &[
1686                        &version_id.0,
1687                        entity_instance_id,
1688                        &id_i32(transformation_id.0, "transformation_id")?,
1689                    ],
1690                )
1691                .map_err(|source| write_error("link asset version entity", source))?;
1692            transaction
1693                .execute(
1694                    "
1695                    INSERT INTO public.dataset_version_entities (
1696                        version_id, entity_instance_id, entity_variable_id, transformation_id
1697                    )
1698                    VALUES ($1, $2, $3, $4)
1699                    ON CONFLICT (version_id, entity_instance_id) DO UPDATE
1700                      SET entity_variable_id = EXCLUDED.entity_variable_id,
1701                          transformation_id = EXCLUDED.transformation_id
1702                    ",
1703                    &[
1704                        &version_id.0,
1705                        entity_instance_id,
1706                        &id_i32(entity_variable_id.0, "entity_variable_id")?,
1707                        &id_i32(transformation_id.0, "transformation_id")?,
1708                    ],
1709                )
1710                .map_err(|source| write_error("link dataset version entity", source))?;
1711        }
1712        transaction.commit().map_err(PgMetaError::Transaction)?;
1713        Ok(())
1714    }
1715
1716    pub fn link_asset_version_relation_instance(
1717        &mut self,
1718        version_id: VersionId,
1719        relation_instance_id: i64,
1720        transformation_id: TransformationId,
1721    ) -> Result<(), PgMetaError> {
1722        self.client
1723            .execute(
1724                "
1725                INSERT INTO public.asset_version_relation_instances (
1726                    version_id, relation_instance_id, transformation_id
1727                )
1728                VALUES ($1, $2, $3)
1729                ON CONFLICT (version_id, relation_instance_id) DO UPDATE
1730                  SET transformation_id = EXCLUDED.transformation_id
1731                ",
1732                &[
1733                    &version_id.0,
1734                    &relation_instance_id,
1735                    &id_i32(transformation_id.0, "transformation_id")?,
1736                ],
1737            )
1738            .map_err(|source| write_error("link asset version relation", source))?;
1739        Ok(())
1740    }
1741
1742    pub fn add_dataset_version_relation(
1743        &mut self,
1744        version_id: VersionId,
1745        relation_instance_id: i64,
1746        subject_variable_id: VariableId,
1747        object_variable_id: VariableId,
1748        relation_variable_id: Option<VariableId>,
1749        transformation_id: TransformationId,
1750    ) -> Result<(), PgMetaError> {
1751        let mut transaction = self
1752            .client
1753            .transaction()
1754            .map_err(PgMetaError::Transaction)?;
1755        transaction
1756            .execute(
1757                "
1758                INSERT INTO public.asset_version_relation_instances (
1759                    version_id, relation_instance_id, transformation_id
1760                )
1761                VALUES ($1, $2, $3)
1762                ON CONFLICT (version_id, relation_instance_id) DO UPDATE
1763                  SET transformation_id = EXCLUDED.transformation_id
1764                ",
1765                &[
1766                    &version_id.0,
1767                    &relation_instance_id,
1768                    &id_i32(transformation_id.0, "transformation_id")?,
1769                ],
1770            )
1771            .map_err(|source| write_error("link asset version relation", source))?;
1772        transaction
1773            .execute(
1774                "
1775                INSERT INTO public.dataset_version_relation_instances (
1776                    version_id, relation_instance_id, subject_variable_id,
1777                    object_variable_id, relation_variable_id, transformation_id
1778                )
1779                VALUES ($1, $2, $3, $4, $5, $6)
1780                ON CONFLICT (version_id, relation_instance_id) DO UPDATE
1781                  SET subject_variable_id = EXCLUDED.subject_variable_id,
1782                      object_variable_id = EXCLUDED.object_variable_id,
1783                      relation_variable_id = EXCLUDED.relation_variable_id,
1784                      transformation_id = EXCLUDED.transformation_id
1785                ",
1786                &[
1787                    &version_id.0,
1788                    &relation_instance_id,
1789                    &id_i32(subject_variable_id.0, "subject_variable_id")?,
1790                    &id_i32(object_variable_id.0, "object_variable_id")?,
1791                    &optional_i32(relation_variable_id.map(|id| id.0), "relation_variable_id")?,
1792                    &id_i32(transformation_id.0, "transformation_id")?,
1793                ],
1794            )
1795            .map_err(|source| write_error("link dataset version relation", source))?;
1796        transaction.commit().map_err(PgMetaError::Transaction)?;
1797        Ok(())
1798    }
1799
1800    pub fn link_dataset_version_relation_instances(
1801        &mut self,
1802        links: &[DatasetVersionRelationLinkInput],
1803    ) -> Result<(), PgMetaError> {
1804        let mut transaction = self
1805            .client
1806            .transaction()
1807            .map_err(PgMetaError::Transaction)?;
1808        for (
1809            version_id,
1810            relation_instance_id,
1811            subject_variable_id,
1812            object_variable_id,
1813            relation_variable_id,
1814            transformation_id,
1815        ) in links
1816        {
1817            transaction
1818                .execute(
1819                    "
1820                    INSERT INTO public.asset_version_relation_instances (
1821                        version_id, relation_instance_id, transformation_id
1822                    )
1823                    VALUES ($1, $2, $3)
1824                    ON CONFLICT (version_id, relation_instance_id) DO UPDATE
1825                      SET transformation_id = EXCLUDED.transformation_id
1826                    ",
1827                    &[
1828                        &version_id.0,
1829                        relation_instance_id,
1830                        &id_i32(transformation_id.0, "transformation_id")?,
1831                    ],
1832                )
1833                .map_err(|source| write_error("link asset version relation", source))?;
1834            transaction
1835                .execute(
1836                    "
1837                    INSERT INTO public.dataset_version_relation_instances (
1838                        version_id, relation_instance_id, subject_variable_id,
1839                        object_variable_id, relation_variable_id, transformation_id
1840                    )
1841                    VALUES ($1, $2, $3, $4, $5, $6)
1842                    ON CONFLICT (version_id, relation_instance_id) DO UPDATE
1843                      SET subject_variable_id = EXCLUDED.subject_variable_id,
1844                          object_variable_id = EXCLUDED.object_variable_id,
1845                          relation_variable_id = EXCLUDED.relation_variable_id,
1846                          transformation_id = EXCLUDED.transformation_id
1847                    ",
1848                    &[
1849                        &version_id.0,
1850                        relation_instance_id,
1851                        &id_i32(subject_variable_id.0, "subject_variable_id")?,
1852                        &id_i32(object_variable_id.0, "object_variable_id")?,
1853                        &optional_i32(relation_variable_id.map(|id| id.0), "relation_variable_id")?,
1854                        &id_i32(transformation_id.0, "transformation_id")?,
1855                    ],
1856                )
1857                .map_err(|source| write_error("link dataset version relation", source))?;
1858        }
1859        transaction.commit().map_err(PgMetaError::Transaction)?;
1860        Ok(())
1861    }
1862
1863    pub fn link_transformation_input(
1864        &mut self,
1865        transformation_id: TransformationId,
1866        version_id: VersionId,
1867    ) -> Result<(), PgMetaError> {
1868        self.client
1869            .execute(
1870                "
1871                INSERT INTO public.transformation_inputs (transformation_id, version_id)
1872                VALUES ($1, $2)
1873                ON CONFLICT DO NOTHING
1874                ",
1875                &[
1876                    &id_i32(transformation_id.0, "transformation_id")?,
1877                    &version_id.0,
1878                ],
1879            )
1880            .map_err(|source| write_error("link transformation input", source))?;
1881        Ok(())
1882    }
1883
1884    pub fn link_transformation_output(
1885        &mut self,
1886        transformation_id: TransformationId,
1887        version_id: VersionId,
1888    ) -> Result<(), PgMetaError> {
1889        self.client
1890            .execute(
1891                "
1892                INSERT INTO public.transformation_outputs (transformation_id, version_id)
1893                VALUES ($1, $2)
1894                ON CONFLICT DO NOTHING
1895                ",
1896                &[
1897                    &id_i32(transformation_id.0, "transformation_id")?,
1898                    &version_id.0,
1899                ],
1900            )
1901            .map_err(|source| write_error("link transformation output", source))?;
1902        Ok(())
1903    }
1904
1905    pub fn upsert_tag(&mut self, name: &str) -> Result<i64, PgMetaError> {
1906        let row = self
1907            .client
1908            .query_one(
1909                "
1910                INSERT INTO public.tags (name)
1911                VALUES ($1)
1912                ON CONFLICT (name) DO UPDATE SET name = EXCLUDED.name
1913                RETURNING tag_id
1914                ",
1915                &[&name],
1916            )
1917            .map_err(|source| write_error("upsert tag", source))?;
1918        Ok(i64::from(row.get::<_, i32>("tag_id")))
1919    }
1920
1921    pub fn attach_tag(&mut self, target: PgWriteTagTarget, tag_id: i64) -> Result<(), PgMetaError> {
1922        let tag_id = id_i32(tag_id, "tag_id")?;
1923        match target {
1924            PgWriteTagTarget::Study(id) => {
1925                self.insert_tag_link("study_tags", "study_id", &id.0, &tag_id)
1926            }
1927            PgWriteTagTarget::Domain(id) => {
1928                let id = id_i32(id.0, "domain_id")?;
1929                self.insert_tag_link("domain_tags", "domain_id", &id, &tag_id)
1930            }
1931            PgWriteTagTarget::Variable(id) => {
1932                let id = id_i32(id.0, "variable_id")?;
1933                self.insert_tag_link("variable_tags", "variable_id", &id, &tag_id)
1934            }
1935            PgWriteTagTarget::Asset(id) => {
1936                self.insert_tag_link("asset_tags", "asset_id", &id.0, &tag_id)
1937            }
1938            PgWriteTagTarget::AssetVersion(id) => {
1939                self.insert_tag_link("asset_version_tags", "version_id", &id.0, &tag_id)
1940            }
1941            PgWriteTagTarget::Entity(id) => {
1942                let id = id_i32(id.0, "entity_id")?;
1943                self.insert_tag_link("entity_tags", "entity_id", &id, &tag_id)
1944            }
1945            PgWriteTagTarget::EntityRelation(id) => {
1946                let id = id_i32(id.0, "entityrelation_id")?;
1947                self.insert_tag_link("entityrelation_tags", "entityrelation_id", &id, &tag_id)
1948            }
1949        }
1950    }
1951
1952    pub fn detach_tag(&mut self, target: PgWriteTagTarget, tag_id: i64) -> Result<(), PgMetaError> {
1953        let tag_id = id_i32(tag_id, "tag_id")?;
1954        match target {
1955            PgWriteTagTarget::Study(id) => {
1956                self.delete_tag_link("study_tags", "study_id", &id.0, &tag_id)
1957            }
1958            PgWriteTagTarget::Domain(id) => {
1959                let id = id_i32(id.0, "domain_id")?;
1960                self.delete_tag_link("domain_tags", "domain_id", &id, &tag_id)
1961            }
1962            PgWriteTagTarget::Variable(id) => {
1963                let id = id_i32(id.0, "variable_id")?;
1964                self.delete_tag_link("variable_tags", "variable_id", &id, &tag_id)
1965            }
1966            PgWriteTagTarget::Asset(id) => {
1967                self.delete_tag_link("asset_tags", "asset_id", &id.0, &tag_id)
1968            }
1969            PgWriteTagTarget::AssetVersion(id) => {
1970                self.delete_tag_link("asset_version_tags", "version_id", &id.0, &tag_id)
1971            }
1972            PgWriteTagTarget::Entity(id) => {
1973                let id = id_i32(id.0, "entity_id")?;
1974                self.delete_tag_link("entity_tags", "entity_id", &id, &tag_id)
1975            }
1976            PgWriteTagTarget::EntityRelation(id) => {
1977                let id = id_i32(id.0, "entityrelation_id")?;
1978                self.delete_tag_link("entityrelation_tags", "entityrelation_id", &id, &tag_id)
1979            }
1980        }
1981    }
1982
1983    pub fn delete_tag_if_unused(&mut self, tag_id: i64) -> Result<bool, PgMetaError> {
1984        let tag_id = id_i32(tag_id, "tag_id")?;
1985        let deleted = self
1986            .client
1987            .execute(
1988                "
1989                DELETE FROM public.tags t
1990                 WHERE t.tag_id = $1
1991                   AND NOT EXISTS (SELECT 1 FROM public.study_tags WHERE tag_id = t.tag_id)
1992                   AND NOT EXISTS (SELECT 1 FROM public.domain_tags WHERE tag_id = t.tag_id)
1993                   AND NOT EXISTS (SELECT 1 FROM public.variable_tags WHERE tag_id = t.tag_id)
1994                   AND NOT EXISTS (SELECT 1 FROM public.asset_tags WHERE tag_id = t.tag_id)
1995                   AND NOT EXISTS (SELECT 1 FROM public.asset_version_tags WHERE tag_id = t.tag_id)
1996                   AND NOT EXISTS (SELECT 1 FROM public.entity_tags WHERE tag_id = t.tag_id)
1997                   AND NOT EXISTS (SELECT 1 FROM public.entityrelation_tags WHERE tag_id = t.tag_id)
1998                ",
1999                &[&tag_id],
2000            )
2001            .map_err(|source| write_error("delete unused tag", source))?;
2002        Ok(deleted == 1)
2003    }
2004
2005    fn insert_tag_link<T>(
2006        &mut self,
2007        table: &str,
2008        target_column: &str,
2009        target_id: &T,
2010        tag_id: &i32,
2011    ) -> Result<(), PgMetaError>
2012    where
2013        T: postgres::types::ToSql + Sync,
2014    {
2015        let sql = format!(
2016            "
2017            INSERT INTO public.{table} ({target_column}, tag_id)
2018            VALUES ($1, $2)
2019            ON CONFLICT DO NOTHING
2020            "
2021        );
2022        self.client
2023            .execute(&sql, &[target_id, tag_id])
2024            .map_err(|source| write_error("attach tag", source))?;
2025        Ok(())
2026    }
2027
2028    fn delete_tag_link<T>(
2029        &mut self,
2030        table: &str,
2031        target_column: &str,
2032        target_id: &T,
2033        tag_id: &i32,
2034    ) -> Result<(), PgMetaError>
2035    where
2036        T: postgres::types::ToSql + Sync,
2037    {
2038        let sql = format!(
2039            "
2040            DELETE FROM public.{table}
2041             WHERE {target_column} = $1
2042               AND tag_id = $2
2043            "
2044        );
2045        self.client
2046            .execute(&sql, &[target_id, tag_id])
2047            .map_err(|source| write_error("detach tag", source))?;
2048        Ok(())
2049    }
2050}
2051
2052/// Target that can receive tag writes through the metadata write API.
2053#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2054pub enum PgWriteTagTarget {
2055    Study(StudyId),
2056    Domain(DomainId),
2057    Variable(VariableId),
2058    Asset(AssetId),
2059    AssetVersion(VersionId),
2060    Entity(EntityId),
2061    EntityRelation(EntityRelationId),
2062}
2063
2064fn insert_asset(
2065    transaction: &mut postgres::Transaction<'_>,
2066    input: PgCreateAsset,
2067) -> Result<AssetId, PgMetaError> {
2068    let asset_type = asset_type_sql(input.asset_type);
2069    let date_created = input.date_created;
2070    let row = match input.asset_id {
2071        Some(asset_id) => transaction.query_one(
2072            "
2073            INSERT INTO public.assets (
2074                asset_id, study_id, name, description, agent_instructions, asset_type,
2075                date_created, created_by
2076            )
2077            VALUES ($1, $2, $3, $4, $5, CAST($6 AS text)::public.asset_type_enum,
2078                    COALESCE($7, CURRENT_TIMESTAMP), COALESCE($8, CURRENT_USER))
2079            RETURNING asset_id
2080            ",
2081            &[
2082                &asset_id.0,
2083                &input.study_id.0,
2084                &input.name,
2085                &input.description,
2086                &input.agent_instructions,
2087                &asset_type,
2088                &date_created,
2089                &input.created_by,
2090            ],
2091        ),
2092        None => transaction.query_one(
2093            "
2094            INSERT INTO public.assets (
2095                study_id, name, description, agent_instructions, asset_type,
2096                date_created, created_by
2097            )
2098            VALUES ($1, $2, $3, $4, CAST($5 AS text)::public.asset_type_enum,
2099                    COALESCE($6, CURRENT_TIMESTAMP), COALESCE($7, CURRENT_USER))
2100            RETURNING asset_id
2101            ",
2102            &[
2103                &input.study_id.0,
2104                &input.name,
2105                &input.description,
2106                &input.agent_instructions,
2107                &asset_type,
2108                &date_created,
2109                &input.created_by,
2110            ],
2111        ),
2112    }
2113    .map_err(|source| write_error("create asset", source))?;
2114    Ok(AssetId(row.get("asset_id")))
2115}
2116
2117fn insert_asset_version(
2118    transaction: &mut impl postgres::GenericClient,
2119    input: PgCreateAssetVersion,
2120) -> Result<VersionId, PgMetaError> {
2121    let created_at = input.created_at;
2122    transaction
2123        .query_one(
2124            "SELECT public.tre_validate_ingest($1,$2,$3)",
2125            &[
2126                &input.asset_id.0,
2127                &input.classification.risk.as_str(),
2128                &input.classification.justification,
2129            ],
2130        )
2131        .map_err(|source| write_error("validate ingest classification", source))?;
2132    clear_latest_asset_version_if_needed(transaction, &input)?;
2133    let row = match input.version_id {
2134        Some(version_id) => transaction.query_one(
2135            "
2136            INSERT INTO public.asset_versions (
2137                version_id, asset_id, major, minor, patch, version_note, is_latest, doi,
2138                created_at, created_by
2139            )
2140            VALUES ($1, $2, COALESCE($3, 1), COALESCE($4, 0), COALESCE($5, 0),
2141                    COALESCE($6, 'Original version'), COALESCE($7, TRUE), $8,
2142                    COALESCE($9, CURRENT_TIMESTAMP), COALESCE($10, CURRENT_USER))
2143            RETURNING version_id
2144            ",
2145            &[
2146                &version_id.0,
2147                &input.asset_id.0,
2148                &input.major,
2149                &input.minor,
2150                &input.patch,
2151                &input.version_note,
2152                &input.is_latest,
2153                &input.doi,
2154                &created_at,
2155                &input.created_by,
2156            ],
2157        ),
2158        None => transaction.query_one(
2159            "
2160            INSERT INTO public.asset_versions (
2161                asset_id, major, minor, patch, version_note, is_latest, doi,
2162                created_at, created_by
2163            )
2164            VALUES ($1, COALESCE($2, 1), COALESCE($3, 0), COALESCE($4, 0),
2165                    COALESCE($5, 'Original version'), COALESCE($6, TRUE), $7,
2166                    COALESCE($8, CURRENT_TIMESTAMP), COALESCE($9, CURRENT_USER))
2167            RETURNING version_id
2168            ",
2169            &[
2170                &input.asset_id.0,
2171                &input.major,
2172                &input.minor,
2173                &input.patch,
2174                &input.version_note,
2175                &input.is_latest,
2176                &input.doi,
2177                &created_at,
2178                &input.created_by,
2179            ],
2180        ),
2181    }
2182    .map_err(|source| write_error("create asset version", source))?;
2183    Ok(VersionId(row.get("version_id")))
2184}
2185
2186fn clear_latest_asset_version_if_needed(
2187    transaction: &mut impl postgres::GenericClient,
2188    input: &PgCreateAssetVersion,
2189) -> Result<(), PgMetaError> {
2190    if !input.is_latest.unwrap_or(true) {
2191        return Ok(());
2192    }
2193
2194    transaction
2195        .execute(
2196            "UPDATE public.asset_versions SET is_latest = FALSE WHERE asset_id = $1 AND is_latest = TRUE",
2197            &[&input.asset_id.0],
2198        )
2199        .map_err(|source| write_error("clear latest asset version", source))?;
2200    Ok(())
2201}
2202
2203fn insert_datafile(
2204    transaction: &mut postgres::Transaction<'_>,
2205    input: PgCreateDataFile,
2206) -> Result<(), PgMetaError> {
2207    transaction
2208        .execute(
2209            "
2210            INSERT INTO public.datafiles (
2211                datafile_id, size_bytes, compressed, encrypted, compression_algorithm,
2212                encryption_algorithm, encryption_key, storage_uri, edam_format, digest
2213            )
2214            VALUES (
2215                $1, $2, COALESCE($3, FALSE), COALESCE($4, FALSE),
2216                COALESCE($5, 'zstd'), COALESCE($6, 'AES-256-CBC with PBKDF2'),
2217                $7, $8, $9, $10
2218            )
2219            ",
2220            &[
2221                &input.datafile_id.0,
2222                &input.size_bytes,
2223                &input.compressed,
2224                &input.encrypted,
2225                &input.compression_algorithm,
2226                &input.encryption_algorithm,
2227                &input.encryption_key,
2228                &input.storage_uri,
2229                &input.edam_format,
2230                &input.digest,
2231            ],
2232        )
2233        .map_err(|source| write_error("create datafile", source))?;
2234    Ok(())
2235}
2236
2237fn asset_type_sql(value: AssetType) -> &'static str {
2238    match value {
2239        AssetType::Dataset => "dataset",
2240        AssetType::File => "file",
2241    }
2242}
2243
2244fn key_role_sql(value: KeyRole) -> Result<&'static str, PgMetaError> {
2245    match value {
2246        KeyRole::None => Ok("none"),
2247        KeyRole::Record => Ok("record"),
2248        KeyRole::External => Ok("external"),
2249    }
2250}
2251
2252fn transformation_type_sql(value: TransformationType) -> &'static str {
2253    match value {
2254        TransformationType::Ingest => "ingest",
2255        TransformationType::Transform => "transform",
2256        TransformationType::Entity => "entity",
2257        TransformationType::Export => "export",
2258        TransformationType::Repository => "repository",
2259    }
2260}
2261
2262fn optional_i32(value: Option<i64>, field: &'static str) -> Result<Option<i32>, PgMetaError> {
2263    value.map(|value| id_i32(value, field)).transpose()
2264}
2265
2266fn id_i32(value: i64, field: &'static str) -> Result<i32, PgMetaError> {
2267    i32::try_from(value).map_err(|_| PgMetaError::Decode {
2268        field,
2269        value: value.to_string(),
2270        message: "integer identifier is outside PostgreSQL INTEGER range".to_string(),
2271    })
2272}
2273
2274fn write_error(operation: &'static str, source: postgres::Error) -> PgMetaError {
2275    PgMetaError::Write { operation, source }
2276}
2277
2278#[cfg(test)]
2279mod tests {
2280    use super::*;
2281
2282    #[test]
2283    fn key_role_sql_encodes_external() {
2284        assert_eq!(
2285            key_role_sql(KeyRole::None).expect("none should encode"),
2286            "none"
2287        );
2288        assert_eq!(
2289            key_role_sql(KeyRole::Record).expect("record should encode"),
2290            "record"
2291        );
2292        assert_eq!(
2293            key_role_sql(KeyRole::External).expect("external should encode"),
2294            "external"
2295        );
2296    }
2297}