1mod semantic;
2pub use semantic::*;
3mod lifecycle;
4pub use lifecycle::*;
5mod governance;
6use ahri_tre_protocol::{
7 datastore::{
8 ComponentStatus, DatastoreMigrationAction, DatastoreMigrationApplyResult,
9 DatastoreMigrationPlanResult, DatastoreMigrationScope, DatastoreMigrationStepResult,
10 DatastorePingResult, DatastoreSchemaCompatibilityStatus, DatastoreSchemaStatusResult,
11 DatastoreSummary, DuckLakeCatalogMigrationOwnership,
12 },
13 public_error::{ProtocolError, ProtocolErrorCode},
14 refs::{
15 DatastoreRef, DomainRef, StudyRef, TagRef, TagScope, TagSummary, encode_scoped_integer_ref,
16 },
17 study::{StudyRegistrationPayload, StudySummary},
18 warning::{ProtocolWarning, ProtocolWarningCode},
19};
20pub use governance::*;
21use serde_json::json;
22
23use crate::{AppError, RegistrationResource};
24
25pub fn configured_datastore_open_protocol_error(
28 error: crate::ConfiguredDatastoreOpenError,
29) -> ProtocolError {
30 match error {
31 crate::ConfiguredDatastoreOpenError::AuthenticationRejected => ProtocolError::new(
32 ProtocolErrorCode::Unauthorized,
33 "Authenticated TRE user is not admitted or is no longer authorized",
34 ),
35 _ => ProtocolError::new(
36 ProtocolErrorCode::OperationFailed,
37 "Authenticated Session could not be opened",
38 ),
39 }
40}
41
42pub fn domain_ref(
44 datastore_id: ahri_tre_protocol::PublicUuid,
45 id: ahri_tre_types::DomainId,
46) -> DomainRef {
47 DomainRef {
48 datastore_id,
49 kind: ahri_tre_protocol::refs::ObjectKind::Domain,
50 id: encode_scoped_integer_ref("domain", id.0),
51 }
52}
53
54pub fn domain_summary(
56 datastore_id: ahri_tre_protocol::PublicUuid,
57 domain: ahri_tre_types::DomainRecord,
58) -> ahri_tre_protocol::domain::DomainSummary {
59 ahri_tre_protocol::domain::DomainSummary {
60 domain: domain_ref(datastore_id, domain.domain_id),
61 name: domain.name.as_str().to_string(),
62 description: domain.description,
63 uri: domain.uri,
64 }
65}
66
67pub fn domain_list_response(
69 session: ahri_tre_protocol::session::SessionStatusPayload,
70 result: crate::DomainListResult,
71) -> ahri_tre_protocol::domain::DomainListResponse {
72 let datastore_id = session
73 .datastore_id
74 .expect("catalogue Session retains its Datastore binding");
75 ahri_tre_protocol::domain::DomainListResponse {
76 session,
77 domains: result
78 .domains
79 .into_iter()
80 .map(|domain| domain_summary(datastore_id, domain))
81 .collect(),
82 }
83}
84
85pub fn domain_add_response(
87 session: ahri_tre_protocol::session::SessionStatusPayload,
88 result: crate::DomainAddResult,
89) -> ahri_tre_protocol::domain::DomainAddResponse {
90 let datastore_id = session
91 .datastore_id
92 .expect("catalogue Session retains its Datastore binding");
93 ahri_tre_protocol::domain::DomainAddResponse {
94 session,
95 domain: ahri_tre_protocol::domain::DomainDetail {
96 summary: domain_summary(datastore_id, result.domain),
97 },
98 }
99}
100
101pub fn study_context_protocol_error(error: AppError) -> ProtocolError {
103 let code = match error {
104 AppError::NotFound(_) => ProtocolErrorCode::NotFound,
105 AppError::Conflict(_) => ProtocolErrorCode::Conflict,
106 AppError::Validation(_) => ProtocolErrorCode::ValidationFailed,
107 _ => ProtocolErrorCode::OperationFailed,
108 };
109 ProtocolError::new(code, "Current study could not be resolved or authorized")
110 .with_target("current_study")
111}
112
113pub fn domain_list_protocol_error(error: AppError) -> ProtocolError {
115 catalogue_list_protocol_error(error)
116}
117
118fn catalogue_list_protocol_error(error: AppError) -> ProtocolError {
119 match error {
120 AppError::NotFound(_) => ProtocolError::new(
121 ProtocolErrorCode::NotFound,
122 "Catalogue scope was not found or is not authorized",
123 ),
124 _ => ProtocolError::new(
125 ProtocolErrorCode::OperationFailed,
126 "Catalogue metadata could not be listed",
127 ),
128 }
129}
130
131pub fn domain_registration_protocol_error(error: AppError) -> ProtocolError {
133 registration_protocol_error(error, RegistrationResource::Domain)
134}
135
136fn registration_protocol_error(error: AppError, resource: RegistrationResource) -> ProtocolError {
137 match error {
138 AppError::Conflict(_) => ProtocolError::new(
139 ProtocolErrorCode::Conflict,
140 format!("{resource} registration conflicts with existing catalogue metadata"),
141 ),
142 AppError::NotFound(_) => ProtocolError::new(
143 ProtocolErrorCode::NotFound,
144 format!("{resource} registration scope was not found or is not authorized"),
145 ),
146 AppError::Validation(_) => ProtocolError::new(
147 ProtocolErrorCode::ValidationFailed,
148 format!("{resource} registration request is invalid"),
149 ),
150 AppError::NotImplemented(_) => ProtocolError::new(
151 ProtocolErrorCode::UnsupportedOperation,
152 format!("{resource} registration is not available"),
153 ),
154 AppError::Infrastructure(_)
155 | AppError::DatasetAdmissionOutcomeUnknown(_)
156 | AppError::RegistrationRollback { .. } => ProtocolError::new(
157 ProtocolErrorCode::OperationFailed,
158 format!("{resource} could not be registered"),
159 ),
160 }
161}
162
163pub fn study_registration_payload(
165 datastore_id: ahri_tre_protocol::PublicUuid,
166 registration: crate::StudyRegistration,
167) -> StudyRegistrationPayload {
168 let domains = registration
169 .domains
170 .into_iter()
171 .map(|domain| ahri_tre_protocol::domain::DomainSummary {
172 domain: DomainRef {
173 datastore_id,
174 kind: ahri_tre_protocol::refs::ObjectKind::Domain,
175 id: encode_scoped_integer_ref("domain", domain.domain_id.0),
176 },
177 name: domain.name.as_str().to_string(),
178 description: domain.description,
179 uri: domain.uri,
180 })
181 .collect::<Vec<_>>();
182 let tags = registration
183 .tags
184 .into_iter()
185 .map(|tag| TagSummary {
186 tag: TagRef {
187 datastore_id,
188 kind: ahri_tre_protocol::refs::ObjectKind::Tag,
189 id: encode_scoped_integer_ref("tag", tag.tag_id.0),
190 },
191 scope: TagScope::Global,
192 label: tag.name,
193 })
194 .collect();
195
196 StudyRegistrationPayload {
197 study: StudySummary {
198 study_type: registration.study.study_type_id.map(|id| {
199 ahri_tre_protocol::refs::ObjectRef {
200 datastore_id,
201 kind: ahri_tre_protocol::refs::ObjectKind::StudyType,
202 id: ahri_tre_protocol::refs::encode_scoped_integer_ref("study_type", id.0),
203 }
204 }),
205 study: StudyRef {
206 datastore_id,
207 kind: ahri_tre_protocol::refs::ObjectKind::Study,
208 id: ahri_tre_protocol::PublicUuid::from_uuid(registration.study.study_id.0),
209 },
210 name: registration.study.name.as_str().to_string(),
211 external_id: registration.study.external_id,
212 title: None,
213 description: registration.study.description,
214 domains: domains.iter().map(|domain| domain.domain).collect(),
215 tags,
216 },
217 domains,
218 }
219}
220
221pub fn study_list_response(
223 session: ahri_tre_protocol::session::SessionStatusPayload,
224 result: crate::StudyListResult,
225) -> ahri_tre_protocol::study::StudyListResponse {
226 let datastore_id = session
227 .datastore_id
228 .expect("catalogue Session retains its Datastore binding");
229 ahri_tre_protocol::study::StudyListResponse {
230 session,
231 studies: result
232 .studies
233 .into_iter()
234 .map(|registration| study_registration_payload(datastore_id, registration))
235 .collect(),
236 }
237}
238
239pub fn study_add_response(
241 session: ahri_tre_protocol::session::SessionStatusPayload,
242 result: crate::StudyAddResult,
243) -> ahri_tre_protocol::study::StudyAddResponse {
244 let datastore_id = session
245 .datastore_id
246 .expect("catalogue Session retains its Datastore binding");
247 ahri_tre_protocol::study::StudyAddResponse {
248 session,
249 registration: study_registration_payload(datastore_id, result.registration),
250 }
251}
252
253pub fn study_list_protocol_error(error: AppError) -> ProtocolError {
255 catalogue_list_protocol_error(error)
256}
257
258pub fn study_registration_protocol_error(error: AppError) -> ProtocolError {
260 registration_protocol_error(error, RegistrationResource::Study)
261}
262
263pub fn datafile_list_response(
265 session: ahri_tre_protocol::session::SessionStatusPayload,
266 result: crate::DataFileListResult,
267) -> ahri_tre_protocol::datafile::DataFileListResponse {
268 let datastore_id = session
269 .datastore_id
270 .expect("catalogue Session retains its Datastore binding");
271 ahri_tre_protocol::datafile::DataFileListResponse {
272 session,
273 study: result.study_name.as_str().to_string(),
274 study_id: object_ref(
275 datastore_id,
276 ahri_tre_protocol::refs::ObjectKind::Study,
277 result.study_id.0,
278 ),
279 datafiles: result
280 .datafiles
281 .into_iter()
282 .map(|entry| ahri_tre_protocol::datafile::DataFileEntryPayload {
283 catalog: datafile_list_asset_catalog(datastore_id, entry.catalog),
284 datafiles: entry
285 .datafiles
286 .into_iter()
287 .map(
288 |datafile| ahri_tre_protocol::datafile::SafeDataFileVersionSummary {
289 datafile_id: object_ref(
290 datastore_id,
291 ahri_tre_protocol::refs::ObjectKind::AssetVersion,
292 datafile.datafile_id.0,
293 ),
294 compressed: datafile.compressed,
295 encrypted: datafile.encrypted,
296 compression_algorithm: datafile.compression_algorithm,
297 encryption_algorithm: datafile.encryption_algorithm,
298 edam_format: datafile.edam_format,
299 },
300 )
301 .collect(),
302 })
303 .collect(),
304 }
305}
306
307fn datafile_list_asset_catalog(
308 datastore_id: ahri_tre_protocol::PublicUuid,
309 catalog: crate::DataFileListAssetCatalog,
310) -> ahri_tre_protocol::asset::SafeAssetCatalog {
311 ahri_tre_protocol::asset::SafeAssetCatalog {
312 asset: ahri_tre_protocol::asset::SafeAssetSummary {
313 asset_id: object_ref(
314 datastore_id,
315 ahri_tre_protocol::refs::ObjectKind::Asset,
316 catalog.asset.asset_id.0,
317 ),
318 study_id: object_ref(
319 datastore_id,
320 ahri_tre_protocol::refs::ObjectKind::Study,
321 catalog.asset.study_id.0,
322 ),
323 name: catalog.asset.name.as_str().to_string(),
324 description: catalog.asset.description,
325 asset_type: catalog.asset.asset_type,
326 date_created: catalog.asset.date_created,
327 created_by: catalog.asset.created_by,
328 },
329 versions: catalog
330 .versions
331 .into_iter()
332 .map(|value| safe_asset_version_summary(datastore_id, value))
333 .collect(),
334 latest_version: catalog
335 .latest_version
336 .map(|value| safe_asset_version_summary(datastore_id, value)),
337 }
338}
339
340pub fn safe_asset_catalog(
342 datastore_id: ahri_tre_protocol::PublicUuid,
343 catalog: crate::AssetVersionCatalog,
344) -> ahri_tre_protocol::asset::SafeAssetCatalog {
345 datafile_list_asset_catalog(
346 datastore_id,
347 crate::DataFileListAssetCatalog::from_catalogue(catalog),
348 )
349}
350
351pub fn safe_asset_version_summary(
352 datastore_id: ahri_tre_protocol::PublicUuid,
353 version: ahri_tre_types::AssetVersionRecord,
354) -> ahri_tre_protocol::asset::SafeAssetVersionSummary {
355 ahri_tre_protocol::asset::SafeAssetVersionSummary {
356 version_id: object_ref(
357 datastore_id,
358 ahri_tre_protocol::refs::ObjectKind::AssetVersion,
359 version.version_id.0,
360 ),
361 asset_id: object_ref(
362 datastore_id,
363 ahri_tre_protocol::refs::ObjectKind::Asset,
364 version.asset_id.0,
365 ),
366 major: version.major,
367 minor: version.minor,
368 patch: version.patch,
369 version_label: version.version_label,
370 version_note: version.version_note,
371 is_latest: version.is_latest,
372 doi: version.doi,
373 created_at: version.created_at,
374 created_by: version.created_by,
375 }
376}
377
378pub fn datafile_list_protocol_error(error: AppError) -> ProtocolError {
380 catalogue_list_protocol_error(error)
381}
382
383pub fn datafile_ingest_response(
385 session: ahri_tre_protocol::session::SessionStatusPayload,
386 result: crate::DataFileIngestResult,
387 source: ahri_tre_protocol::ingest::IngestSourceSummary,
388) -> ahri_tre_protocol::ingest::IngestFileResponse {
389 let datastore_id = session
390 .datastore_id
391 .expect("catalogue Session retains its Datastore binding");
392 let response_version = result.catalog.latest_version.clone();
393 let catalog = datafile_list_asset_catalog(datastore_id, result.catalog);
394 let asset = catalog.asset;
395 let study = StudyRef {
396 datastore_id,
397 kind: ahri_tre_protocol::refs::ObjectKind::Study,
398 id: ahri_tre_protocol::PublicUuid::from_uuid(result.study_id.0),
399 };
400 let version = response_version
401 .as_ref()
402 .map(|version| ahri_tre_protocol::ingest::version_ref(datastore_id, version.version_id));
403 ahri_tre_protocol::ingest::IngestFileResponse {
404 session,
405 study: result.study_name.as_str().to_string(),
406 asset: ahri_tre_protocol::ingest::IngestedAssetSummary {
407 name: asset.name.clone(),
408 asset_type: asset.asset_type,
409 study,
410 description: asset.description.clone(),
411 version,
412 },
413 datafile: ahri_tre_protocol::datafile::DataFileSummary {
414 version,
415 datafile: asset.asset_id,
416 study,
417 name: asset.name,
418 format: result.datafile.edam_format,
419 description: asset.description,
420 size_bytes: Some(result.stored_size_bytes),
421 digest_ref: Some(result.datafile.digest),
422 tags: Vec::new(),
423 },
424 version: response_version.map(|version| safe_asset_version_summary(datastore_id, version)),
425 transformation_id: result.transformation_id,
426 transformation: Some(ahri_tre_protocol::refs::ObjectRef {
427 datastore_id,
428 kind: ahri_tre_protocol::refs::ObjectKind::Transformation,
429 id: ahri_tre_protocol::refs::encode_scoped_integer_ref(
430 "transformation",
431 result.transformation_id.0,
432 ),
433 }),
434 source_size_bytes: result.source_size_bytes,
435 stored_size_bytes: result.stored_size_bytes,
436 source,
437 }
438}
439
440pub fn datafile_ingest_protocol_error(error: AppError) -> ProtocolError {
442 match error {
443 AppError::Validation(_) => ProtocolError::new(
444 ProtocolErrorCode::ValidationFailed,
445 "Remote CSV upload request is invalid",
446 ),
447 AppError::NotFound(_) => ProtocolError::new(
448 ProtocolErrorCode::NotFound,
449 "Remote CSV upload Study was not found or is not authorized",
450 ),
451 AppError::Conflict(_) => ProtocolError::new(
452 ProtocolErrorCode::Conflict,
453 "Remote CSV upload conflicts with existing Datafile metadata",
454 ),
455 AppError::NotImplemented(_) => ProtocolError::new(
456 ProtocolErrorCode::UnsupportedOperation,
457 "Remote CSV upload is not available",
458 ),
459 AppError::Infrastructure(_)
460 | AppError::DatasetAdmissionOutcomeUnknown(_)
461 | AppError::RegistrationRollback { .. } => ProtocolError::new(
462 ProtocolErrorCode::OperationFailed,
463 "Remote CSV upload could not be completed",
464 ),
465 }
466}
467
468pub fn dataset_materialization_response(
471 session: ahri_tre_protocol::session::SessionStatusPayload,
472 study: String,
473 result: crate::ManagedDataFileDatasetMaterialization,
474) -> ahri_tre_protocol::ingest::DatasetMaterializationResponse {
475 let datastore_id = session
476 .datastore_id
477 .expect("catalogue Session retains its Datastore binding");
478 let source_study = StudyRef {
479 datastore_id,
480 kind: ahri_tre_protocol::refs::ObjectKind::Study,
481 id: ahri_tre_protocol::PublicUuid::from_uuid(result.source.catalog.asset.study_id.0),
482 };
483 let source_name = result.source.catalog.asset.name.as_str().to_string();
484 let source_description = result.source.catalog.asset.description.clone();
485 let source_datafile = ahri_tre_protocol::datafile::DataFileSummary {
486 version: Some(ahri_tre_protocol::ingest::version_ref(
487 datastore_id,
488 result.source.version.version_id,
489 )),
490 datafile: ahri_tre_protocol::ingest::datafile_ref(
491 datastore_id,
492 result.source.catalog.asset.asset_id,
493 ),
494 study: source_study,
495 name: source_name,
496 format: result.source.datafile.edam_format.clone(),
497 description: source_description,
498 size_bytes: None,
499 digest_ref: Some(result.source.datafile.digest.as_str().to_string()),
500 tags: Vec::new(),
501 };
502 let source_catalog = public_asset_catalog(result.source.catalog);
503 let dataset_catalog = public_asset_catalog(result.dataset.catalog);
504 ahri_tre_protocol::ingest::DatasetMaterializationResponse {
505 session,
506 study,
507 catalog: dataset_catalog,
508 dataset: result.dataset.dataset,
509 lineage: workflow_lineage(datastore_id, result.dataset.lineage),
510 row_count: result.dataset.row_count,
511 column_count: result.dataset.column_count,
512 variables: result
513 .dataset
514 .variables
515 .into_iter()
516 .map(
517 |variable| ahri_tre_protocol::dataset::RegisteredDatasetVariablePayload {
518 variable: variable_summary(datastore_id, variable.variable),
519 vocabulary_items_withheld: variable.vocabulary.is_some(),
520 vocabulary: variable
521 .vocabulary
522 .map(|value| vocabulary_summary(datastore_id, value)),
523 vocabulary_items: Vec::new(),
524 link: dataset_variable_link_summary(datastore_id, variable.link),
525 },
526 )
527 .collect(),
528 metadata_warnings: result.dataset.metadata_warnings,
529 source: ahri_tre_protocol::ingest::DatasetMaterializationSourceSummary {
530 upload: None,
531 kind: ahri_tre_protocol::ingest::DatasetMaterializationSourceKind::ManagedDataFile,
532 request_only: false,
533 },
534 source_catalog: Some(source_catalog),
535 source_version: Some(result.source.version),
536 source_datafile: Some(source_datafile),
537 }
538}
539
540fn public_asset_catalog(
541 catalog: crate::AssetVersionCatalog,
542) -> ahri_tre_protocol::asset::AssetCatalog {
543 let mut asset = catalog.asset;
544 asset.agent_instructions = None;
545 ahri_tre_protocol::asset::AssetCatalog {
546 asset,
547 versions: catalog.versions,
548 latest_version: catalog.latest_version,
549 }
550}
551
552pub fn dataset_materialization_protocol_error(error: AppError) -> ProtocolError {
554 match error {
555 AppError::Validation(_) => ProtocolError::new(
556 ProtocolErrorCode::ValidationFailed,
557 "Dataset derivation request is invalid",
558 ),
559 AppError::NotFound(_) => ProtocolError::new(
560 ProtocolErrorCode::NotFound,
561 "Dataset derivation scope was not found or is not authorized",
562 ),
563 AppError::Conflict(_) => ProtocolError::new(
564 ProtocolErrorCode::Conflict,
565 "Dataset derivation conflicts with existing catalogue metadata",
566 ),
567 AppError::NotImplemented(_) => ProtocolError::new(
568 ProtocolErrorCode::UnsupportedOperation,
569 "Dataset derivation is not available",
570 ),
571 AppError::Infrastructure(_)
572 | AppError::DatasetAdmissionOutcomeUnknown(_)
573 | AppError::RegistrationRollback { .. } => ProtocolError::new(
574 ProtocolErrorCode::OperationFailed,
575 "Dataset derivation could not be completed",
576 ),
577 }
578}
579
580pub fn dataset_metadata_response(
582 session: ahri_tre_protocol::session::SessionStatusPayload,
583 result: crate::DatasetMetadataResult,
584) -> ahri_tre_protocol::dataset::DatasetMetadataResponse {
585 let datastore_id = session
586 .datastore_id
587 .expect("catalogue Session retains its Datastore binding");
588 ahri_tre_protocol::dataset::DatasetMetadataResponse {
589 session,
590 study: result.study.name.as_str().to_string(),
591 metadata: ahri_tre_protocol::dataset::DatasetMetadataPayload {
592 catalog: safe_asset_catalog(datastore_id, result.metadata.catalog),
593 dataset: dataset_summary(datastore_id, result.metadata.dataset),
594 variables_included: result.metadata.variables_included,
595 variable_count: result.metadata.variable_count,
596 variables: result
597 .metadata
598 .variables
599 .into_iter()
600 .map(
601 |variable| ahri_tre_protocol::dataset::DatasetVariableMetadataPayload {
602 link: dataset_variable_link_summary(datastore_id, variable.link),
603 variable: variable_summary(datastore_id, variable.variable),
604 value_type: value_type_summary(datastore_id, variable.value_type),
605 vocabulary_items_withheld: variable.vocabulary.is_some(),
606 vocabulary: variable
607 .vocabulary
608 .map(|value| vocabulary_summary(datastore_id, value)),
609 vocabulary_items: Vec::new(),
610 },
611 )
612 .collect(),
613 },
614 }
615}
616
617pub fn dataset_metadata_protocol_error(error: AppError) -> ProtocolError {
619 catalogue_protocol_error(error)
620}
621
622pub fn catalogue_protocol_error(error: AppError) -> ProtocolError {
623 let (code, message) = match error {
624 AppError::Validation(_) => (
625 ProtocolErrorCode::ValidationFailed,
626 "Catalogue selector is invalid",
627 ),
628 AppError::Conflict(_) => (
629 ProtocolErrorCode::Conflict,
630 "Catalogue selectors are ambiguous or disagree",
631 ),
632 AppError::NotFound(_) => (
633 ProtocolErrorCode::NotFound,
634 "Catalogue scope was not found or is not authorized",
635 ),
636 AppError::NotImplemented(_) => (
637 ProtocolErrorCode::UnsupportedOperation,
638 "Catalogue workflow is unavailable",
639 ),
640 _ => (
641 ProtocolErrorCode::OperationFailed,
642 "Catalogue metadata could not be read",
643 ),
644 };
645 ProtocolError::new(code, message)
646}
647
648pub fn app_error_to_protocol_error(error: &AppError) -> ProtocolError {
649 match error {
650 AppError::Validation(message) => {
651 ProtocolError::new(ProtocolErrorCode::ValidationFailed, safe_message(message))
652 }
653 AppError::NotFound(message) => {
654 ProtocolError::new(ProtocolErrorCode::NotFound, safe_message(message))
655 }
656 AppError::Conflict(message) => {
657 ProtocolError::new(ProtocolErrorCode::Conflict, safe_message(message))
658 }
659 AppError::NotImplemented(message) => ProtocolError::new(
660 ProtocolErrorCode::UnsupportedOperation,
661 format!("unsupported operation: {}", safe_message(message)),
662 ),
663 _ => ProtocolError::new(
664 ProtocolErrorCode::Internal,
665 "The operation failed before a public result was available.",
666 ),
667 }
668}
669
670pub fn safe_warning(code: ProtocolWarningCode, message: impl AsRef<str>) -> ProtocolWarning {
671 ProtocolWarning::new(code, safe_message(message.as_ref()))
672}
673
674pub fn datastore_summary(
675 datastore: ahri_tre_protocol::refs::DatastoreRef,
676 name: impl Into<String>,
677 ready: bool,
678) -> DatastoreSummary {
679 DatastoreSummary {
680 datastore,
681 name: name.into(),
682 ready,
683 }
684}
685
686pub fn datastore_ping_result(
687 datastore: ahri_tre_protocol::refs::DatastoreRef,
688 metadata_ready: bool,
689 lake_ready: bool,
690) -> DatastorePingResult {
691 DatastorePingResult {
692 datastore,
693 metadata: ComponentStatus {
694 ready: metadata_ready,
695 message: None,
696 },
697 lake: ComponentStatus {
698 ready: lake_ready,
699 message: None,
700 },
701 }
702}
703
704pub fn datastore_schema_status_result(
705 datastore: DatastoreRef,
706 runtime: &ahri_tre_runtime::DataStoreRuntime,
707 status: ahri_tre_pgmeta::DatastoreSchemaStatus,
708) -> DatastoreSchemaStatusResult {
709 DatastoreSchemaStatusResult {
710 datastore,
711 database: status.database,
712 status: datastore_schema_compatibility_status(status.status),
713 current_version: status.current_version,
714 target_version: status.target_version,
715 migration_history_present: status.migration_history_present,
716 applied_versions: status.applied_versions,
717 missing_tables: status.missing_tables,
718 ducklake_catalog_migrations: ducklake_catalog_migration_ownership(runtime),
719 }
720}
721
722pub fn datastore_migration_plan_result(
723 runtime: &ahri_tre_runtime::DataStoreRuntime,
724 plan: ahri_tre_pgmeta::DatastoreMigrationPlan,
725) -> DatastoreMigrationPlanResult {
726 let blocked_reason = plan.blocked_reason;
727 let blocked = blocked_reason
728 .as_deref()
729 .is_some_and(|reason| !reason.trim().is_empty());
730 let target_version = plan.target_version;
731 DatastoreMigrationPlanResult {
732 status: datastore_schema_compatibility_status(plan.status),
733 current_version: plan.current_version,
734 target_version: target_version.clone(),
735 steps: plan
736 .steps
737 .into_iter()
738 .map(datastore_migration_step_result)
739 .collect(),
740 blocked,
741 blocked_reason: blocked_reason.clone(),
742 dependency_summary: blocked.then(|| {
743 json!({
744 "status": "blocked",
745 "items": [{
746 "kind": "datastore_schema",
747 "identifier": target_version,
748 "action": "block"
749 }]
750 })
751 }),
752 next_step: blocked.then(|| {
753 "resolve the datastore schema compatibility issue, then rerun schema-plan or schema-migrate"
754 .to_string()
755 }),
756 ducklake_catalog_migrations: ducklake_catalog_migration_ownership(runtime),
757 }
758}
759
760pub fn datastore_migration_apply_result(
761 datastore: DatastoreRef,
762 runtime: &ahri_tre_runtime::DataStoreRuntime,
763 report: ahri_tre_pgmeta::DatastoreMigrationApplyReport,
764) -> DatastoreMigrationApplyResult {
765 let blocked_reason = report.blocked_reason;
766 let blocked = blocked_reason
767 .as_deref()
768 .is_some_and(|reason| !reason.trim().is_empty());
769 let before = datastore_migration_plan_result(runtime, report.before);
770 let target_version = before.target_version.clone();
771 DatastoreMigrationApplyResult {
772 before,
773 after: datastore_schema_status_result(datastore, runtime, report.after),
774 applied_steps: report
775 .applied_steps
776 .into_iter()
777 .map(datastore_migration_step_result)
778 .collect(),
779 blocked,
780 blocked_reason: blocked_reason.clone(),
781 dependency_summary: blocked.then(|| {
782 json!({
783 "status": "blocked",
784 "items": [{
785 "kind": "datastore_schema",
786 "identifier": target_version,
787 "action": "block"
788 }]
789 })
790 }),
791 next_step: blocked.then(|| {
792 "resolve the datastore schema compatibility issue, then rerun schema-plan or schema-migrate"
793 .to_string()
794 }),
795 }
796}
797
798pub fn datastore_schema_compatibility_status(
799 status: ahri_tre_pgmeta::DatastoreSchemaCompatibility,
800) -> DatastoreSchemaCompatibilityStatus {
801 match status {
802 ahri_tre_pgmeta::DatastoreSchemaCompatibility::Current => {
803 DatastoreSchemaCompatibilityStatus::Current
804 }
805 ahri_tre_pgmeta::DatastoreSchemaCompatibility::Pending => {
806 DatastoreSchemaCompatibilityStatus::Pending
807 }
808 ahri_tre_pgmeta::DatastoreSchemaCompatibility::MissingHistory => {
809 DatastoreSchemaCompatibilityStatus::MissingHistory
810 }
811 ahri_tre_pgmeta::DatastoreSchemaCompatibility::Dirty => {
812 DatastoreSchemaCompatibilityStatus::Dirty
813 }
814 ahri_tre_pgmeta::DatastoreSchemaCompatibility::TooNew => {
815 DatastoreSchemaCompatibilityStatus::TooNew
816 }
817 ahri_tre_pgmeta::DatastoreSchemaCompatibility::Unsupported => {
818 DatastoreSchemaCompatibilityStatus::Unsupported
819 }
820 }
821}
822
823fn ducklake_catalog_migration_ownership(
824 runtime: &ahri_tre_runtime::DataStoreRuntime,
825) -> DuckLakeCatalogMigrationOwnership {
826 DuckLakeCatalogMigrationOwnership {
827 database: runtime.lake.catalog_database.clone(),
828 schema: runtime.lake.catalog_schema.clone(),
829 managed_by: "ducklake".to_string(),
830 included_in_datastore_schema_migration: false,
831 }
832}
833
834fn datastore_migration_step_result(
835 step: ahri_tre_pgmeta::DatastoreMigrationStep,
836) -> DatastoreMigrationStepResult {
837 DatastoreMigrationStepResult {
838 version: step.version,
839 description: step.description,
840 scope: match step.scope {
841 ahri_tre_pgmeta::DatastoreMigrationScope::Metadata => DatastoreMigrationScope::Metadata,
842 },
843 transactional: step.transactional,
844 action: match step.action {
845 ahri_tre_pgmeta::DatastoreMigrationAction::Apply => DatastoreMigrationAction::Apply,
846 ahri_tre_pgmeta::DatastoreMigrationAction::Skipped => DatastoreMigrationAction::Skipped,
847 ahri_tre_pgmeta::DatastoreMigrationAction::Failed => DatastoreMigrationAction::Failed,
848 ahri_tre_pgmeta::DatastoreMigrationAction::Blocked => DatastoreMigrationAction::Blocked,
849 },
850 }
851}
852
853fn safe_message(message: &str) -> String {
854 let lower = message.to_ascii_lowercase();
855 if lower.contains("password")
856 || lower.contains("token")
857 || lower.contains("secret")
858 || lower.contains("postgres://")
859 || lower.contains("select ")
860 || lower.contains("/tmp/")
861 || lower.contains("/workspaces/")
862 {
863 "The operation failed with details omitted from the public protocol response.".to_string()
864 } else {
865 message.to_string()
866 }
867}
868
869pub fn object_ref(
870 datastore_id: ahri_tre_protocol::PublicUuid,
871 kind: ahri_tre_protocol::refs::ObjectKind,
872 id: uuid::Uuid,
873) -> ahri_tre_protocol::refs::ObjectRef {
874 ahri_tre_protocol::refs::ObjectRef {
875 datastore_id,
876 kind,
877 id: ahri_tre_protocol::PublicUuid::from_uuid(id),
878 }
879}
880
881pub fn safe_datafile_version_summary(
882 datastore_id: ahri_tre_protocol::PublicUuid,
883 datafile: ahri_tre_types::DataFileRecord,
884) -> ahri_tre_protocol::datafile::SafeDataFileVersionSummary {
885 ahri_tre_protocol::datafile::SafeDataFileVersionSummary {
886 datafile_id: object_ref(
887 datastore_id,
888 ahri_tre_protocol::refs::ObjectKind::AssetVersion,
889 datafile.datafile_id.0,
890 ),
891 compressed: datafile.compressed,
892 encrypted: datafile.encrypted,
893 compression_algorithm: datafile.compression_algorithm,
894 encryption_algorithm: datafile.encryption_algorithm,
895 edam_format: datafile.edam_format,
896 }
897}
898
899pub fn dataset_summary(
900 datastore_id: ahri_tre_protocol::PublicUuid,
901 value: ahri_tre_types::DatasetRecord,
902) -> ahri_tre_protocol::catalogue::DatasetSummary {
903 use ahri_tre_protocol::refs::ObjectKind;
904 ahri_tre_protocol::catalogue::DatasetSummary {
905 dataset_id: object_ref(datastore_id, ObjectKind::AssetVersion, value.dataset_id.0),
906 }
907}
908
909pub fn dataset_variable_link_summary(
910 datastore_id: ahri_tre_protocol::PublicUuid,
911 value: ahri_tre_types::DatasetVariableLinkRecord,
912) -> ahri_tre_protocol::catalogue::DatasetVariableLinkSummary {
913 use ahri_tre_protocol::refs::ObjectKind;
914 ahri_tre_protocol::catalogue::DatasetVariableLinkSummary {
915 dataset_id: object_ref(datastore_id, ObjectKind::AssetVersion, value.dataset_id.0),
916 variable_id: object_ref(
917 datastore_id,
918 ObjectKind::Variable,
919 ahri_tre_protocol::refs::encode_scoped_integer_ref("variable", value.variable_id.0)
920 .as_uuid(),
921 ),
922 row_role: value.row_role,
923 }
924}
925
926pub fn variable_summary(
927 datastore_id: ahri_tre_protocol::PublicUuid,
928 value: ahri_tre_types::VariableRecord,
929) -> ahri_tre_protocol::catalogue::VariableSummary {
930 use ahri_tre_protocol::refs::ObjectKind;
931 ahri_tre_protocol::catalogue::VariableSummary {
932 variable_id: object_ref(
933 datastore_id,
934 ObjectKind::Variable,
935 ahri_tre_protocol::refs::encode_scoped_integer_ref("variable", value.variable_id.0)
936 .as_uuid(),
937 ),
938 domain_id: object_ref(
939 datastore_id,
940 ObjectKind::Domain,
941 ahri_tre_protocol::refs::encode_scoped_integer_ref("domain", value.domain_id.0)
942 .as_uuid(),
943 ),
944 name: value.name,
945 value_type_id: object_ref(
946 datastore_id,
947 ObjectKind::ValueType,
948 ahri_tre_protocol::refs::encode_scoped_integer_ref("value_type", value.value_type_id.0)
949 .as_uuid(),
950 ),
951 value_format: value.value_format,
952 vocabulary_id: value.vocabulary_id.map(|id| {
953 object_ref(
954 datastore_id,
955 ObjectKind::Vocabulary,
956 ahri_tre_protocol::refs::encode_scoped_integer_ref("vocabulary", id.0).as_uuid(),
957 )
958 }),
959 key_role: value.key_role,
960 description: value.description,
961 note: value.note,
962 ontology_namespace: value.ontology_namespace,
963 ontology_class: value.ontology_class,
964 }
965}
966
967pub fn value_type_summary(
968 datastore_id: ahri_tre_protocol::PublicUuid,
969 value: ahri_tre_types::ValueTypeRecord,
970) -> ahri_tre_protocol::catalogue::ValueTypeSummary {
971 use ahri_tre_protocol::refs::ObjectKind;
972 ahri_tre_protocol::catalogue::ValueTypeSummary {
973 value_type_id: object_ref(
974 datastore_id,
975 ObjectKind::ValueType,
976 ahri_tre_protocol::refs::encode_scoped_integer_ref("value_type", value.value_type_id.0)
977 .as_uuid(),
978 ),
979 value_type: value.value_type,
980 description: value.description,
981 }
982}
983
984pub fn vocabulary_summary(
985 datastore_id: ahri_tre_protocol::PublicUuid,
986 value: ahri_tre_types::VocabularyRecord,
987) -> ahri_tre_protocol::catalogue::VocabularySummary {
988 use ahri_tre_protocol::refs::ObjectKind;
989 ahri_tre_protocol::catalogue::VocabularySummary {
990 vocabulary_id: object_ref(
991 datastore_id,
992 ObjectKind::Vocabulary,
993 ahri_tre_protocol::refs::encode_scoped_integer_ref("vocabulary", value.vocabulary_id.0)
994 .as_uuid(),
995 ),
996 domain_id: object_ref(
997 datastore_id,
998 ObjectKind::Domain,
999 ahri_tre_protocol::refs::encode_scoped_integer_ref("domain", value.domain_id.0)
1000 .as_uuid(),
1001 ),
1002 name: value.name,
1003 description: value.description,
1004 }
1005}
1006
1007pub fn vocabulary_item_summary(
1008 datastore_id: ahri_tre_protocol::PublicUuid,
1009 value: ahri_tre_types::VocabularyItemRecord,
1010) -> ahri_tre_protocol::catalogue::VocabularyItemSummary {
1011 use ahri_tre_protocol::refs::ObjectKind;
1012 ahri_tre_protocol::catalogue::VocabularyItemSummary {
1013 vocabulary_item_id: object_ref(
1014 datastore_id,
1015 ObjectKind::VocabularyItem,
1016 ahri_tre_protocol::refs::encode_scoped_integer_ref(
1017 "vocabulary_item",
1018 value.vocabulary_item_id.0,
1019 )
1020 .as_uuid(),
1021 ),
1022 vocabulary_id: object_ref(
1023 datastore_id,
1024 ObjectKind::Vocabulary,
1025 ahri_tre_protocol::refs::encode_scoped_integer_ref("vocabulary", value.vocabulary_id.0)
1026 .as_uuid(),
1027 ),
1028 value: value.value,
1029 code: value.code,
1030 description: value.description,
1031 }
1032}
1033
1034pub fn dataset_list_response(
1036 datastore_id: ahri_tre_protocol::PublicUuid,
1037 session: ahri_tre_protocol::session::SessionStatusPayload,
1038 study: ahri_tre_types::StudyRecord,
1039 listing: crate::StudyDatasetList,
1040) -> ahri_tre_protocol::dataset::DatasetListResponse {
1041 ahri_tre_protocol::dataset::DatasetListResponse {
1042 session,
1043 study: study.name.as_str().to_string(),
1044 study_id: object_ref(
1045 datastore_id,
1046 ahri_tre_protocol::refs::ObjectKind::Study,
1047 study.study_id.0,
1048 ),
1049 datasets: listing
1050 .datasets
1051 .into_iter()
1052 .map(
1053 |entry| ahri_tre_protocol::dataset::DatasetListEntryPayload {
1054 catalog: safe_asset_catalog(datastore_id, entry.catalog),
1055 latest_dataset: entry
1056 .latest_dataset
1057 .map(|value| dataset_summary(datastore_id, value)),
1058 },
1059 )
1060 .collect(),
1061 }
1062}
1063
1064#[cfg(test)]
1065mod tests {
1066 use super::*;
1067 use ahri_tre_protocol::{
1068 ProtocolErrorCode, PublicName,
1069 refs::SessionSelector,
1070 session::{SessionAvailability, SessionStatusPayload},
1071 };
1072 use uuid::Uuid;
1073
1074 fn test_session_status() -> SessionStatusPayload {
1075 SessionStatusPayload {
1076 name: (SessionSelector::Name {
1077 name: PublicName::new("analysis").unwrap(),
1078 })
1079 .name()
1080 .map(|name| name.as_str().to_owned())
1081 .unwrap_or_default(),
1082 session: None,
1083 adapter_observations: Vec::new(),
1084 execution_profile: Some("research".to_string()),
1085 datastore_key: Some("test-datastore".to_string()),
1086 datastore_id: Some(ahri_tre_protocol::PublicUuid::from_uuid(
1087 uuid::Uuid::from_u128(5),
1088 )),
1089 availability: SessionAvailability::Live,
1090 authenticated_tre_user: Some("0000-0001-2345-6789".to_string()),
1091 current_study: None,
1092 current_study_name: None,
1093 created_at: None,
1094 updated_at: None,
1095 unavailable_reason: None,
1096 }
1097 }
1098
1099 #[test]
1100 fn domain_list_workflow_has_one_safe_public_result() {
1101 let result = crate::DomainListResult {
1102 domains: vec![ahri_tre_types::DomainRecord {
1103 domain_id: ahri_tre_types::DomainId(10),
1104 name: ahri_tre_types::NcName::parse("public_health").unwrap(),
1105 description: Some("Public health".to_string()),
1106 uri: Some("https://example.test/public-health".to_string()),
1107 }],
1108 };
1109
1110 let projected = domain_list_response(test_session_status(), result);
1111
1112 assert_eq!(projected.session.name, "analysis");
1113 assert_eq!(projected.domains.len(), 1);
1114 assert_eq!(projected.domains[0].name, "public_health");
1115 assert_eq!(
1116 projected.domains[0].domain.id,
1117 encode_scoped_integer_ref("domain", 10)
1118 );
1119 }
1120
1121 #[test]
1122 fn domain_add_workflow_has_one_safe_public_result_and_error() {
1123 let result = crate::DomainAddResult {
1124 domain: ahri_tre_types::DomainRecord {
1125 domain_id: ahri_tre_types::DomainId(11),
1126 name: ahri_tre_types::NcName::parse("population_health").unwrap(),
1127 description: Some("Population health".to_string()),
1128 uri: None,
1129 },
1130 };
1131
1132 let projected = domain_add_response(test_session_status(), result);
1133 let private_error =
1134 crate::AppError::Infrastructure("database password and /tmp/private-path".to_string());
1135 let error = domain_registration_protocol_error(private_error);
1136
1137 assert_eq!(projected.domain.summary.name, "population_health");
1138 assert_eq!(error.code, ProtocolErrorCode::OperationFailed);
1139 assert_eq!(error.message, "Domain could not be registered");
1140 }
1141
1142 #[test]
1143 fn domain_workflow_errors_have_one_safe_public_meaning() {
1144 let list_missing = domain_list_protocol_error(crate::AppError::NotFound(
1145 "private catalogue detail".to_string(),
1146 ));
1147 let list_failed = domain_list_protocol_error(crate::AppError::Infrastructure(
1148 "postgres://private.example/tre".to_string(),
1149 ));
1150 let registration_cases = [
1151 (
1152 crate::AppError::Validation("private bad name".to_string()),
1153 ProtocolErrorCode::ValidationFailed,
1154 "Domain registration request is invalid",
1155 ),
1156 (
1157 crate::AppError::NotFound("private scope".to_string()),
1158 ProtocolErrorCode::NotFound,
1159 "Domain registration scope was not found or is not authorized",
1160 ),
1161 (
1162 crate::AppError::Conflict("private metadata".to_string()),
1163 ProtocolErrorCode::Conflict,
1164 "Domain registration conflicts with existing catalogue metadata",
1165 ),
1166 ];
1167
1168 assert_eq!(list_missing.code, ProtocolErrorCode::NotFound);
1169 assert_eq!(
1170 list_missing.message,
1171 "Catalogue scope was not found or is not authorized"
1172 );
1173 assert_eq!(list_failed.code, ProtocolErrorCode::OperationFailed);
1174 assert_eq!(
1175 list_failed.message,
1176 "Catalogue metadata could not be listed"
1177 );
1178 for (app_error, code, message) in registration_cases {
1179 let projected = domain_registration_protocol_error(app_error);
1180 assert_eq!(projected.code, code);
1181 assert_eq!(projected.message, message);
1182 assert!(!projected.message.contains("private"));
1183 }
1184 }
1185
1186 #[test]
1187 fn study_registration_projection_keeps_public_domains_and_tags() {
1188 let registration = crate::StudyRegistration {
1189 study: ahri_tre_types::StudyRecord {
1190 study_id: ahri_tre_types::StudyId(
1191 Uuid::parse_str("33333333-3333-4333-8333-333333333333").unwrap(),
1192 ),
1193 name: ahri_tre_types::NcName::parse("cohort").unwrap(),
1194 description: Some("Safe catalogue metadata".to_string()),
1195 documentation: None,
1196 agent_instructions: None,
1197 external_id: Some("EXT-001".to_string()),
1198 study_type_id: None,
1199 date_created: None,
1200 created_by: Some("0000-0000-0000-0001".to_string()),
1201 },
1202 domains: vec![ahri_tre_types::DomainRecord {
1203 domain_id: ahri_tre_types::DomainId(10),
1204 name: ahri_tre_types::NcName::parse("public_health").unwrap(),
1205 description: Some("Public health".to_string()),
1206 uri: None,
1207 }],
1208 tags: vec![ahri_tre_types::TagRecord {
1209 tag_id: ahri_tre_types::TagId(1),
1210 name: "curated".to_string(),
1211 }],
1212 };
1213
1214 let projected = study_registration_payload(
1215 ahri_tre_protocol::PublicUuid::from_uuid(uuid::Uuid::from_u128(5)),
1216 registration,
1217 );
1218
1219 assert_eq!(projected.domains.len(), 1);
1220 assert_eq!(projected.study.domains, vec![projected.domains[0].domain]);
1221 assert_eq!(projected.study.tags.len(), 1);
1222 assert_eq!(projected.study.tags[0].label, "curated");
1223 }
1224
1225 #[test]
1226 fn study_workflows_have_one_safe_public_result_path() {
1227 let registration = crate::StudyRegistration {
1228 study: ahri_tre_types::StudyRecord {
1229 study_id: ahri_tre_types::StudyId(
1230 Uuid::parse_str("44444444-4444-4444-8444-444444444444").unwrap(),
1231 ),
1232 name: ahri_tre_types::NcName::parse("cohort").unwrap(),
1233 description: Some("Safe catalogue metadata".to_string()),
1234 documentation: None,
1235 agent_instructions: None,
1236 external_id: Some("EXT-002".to_string()),
1237 study_type_id: Some(ahri_tre_types::StudyTypeId(1)),
1238 date_created: None,
1239 created_by: Some("0000-0000-0000-0002".to_string()),
1240 },
1241 domains: vec![ahri_tre_types::DomainRecord {
1242 domain_id: ahri_tre_types::DomainId(12),
1243 name: ahri_tre_types::NcName::parse("population_health").unwrap(),
1244 description: Some("Population health".to_string()),
1245 uri: None,
1246 }],
1247 tags: vec![ahri_tre_types::TagRecord {
1248 tag_id: ahri_tre_types::TagId(2),
1249 name: "approved".to_string(),
1250 }],
1251 };
1252
1253 let listed = study_list_response(
1254 test_session_status(),
1255 crate::StudyListResult {
1256 studies: vec![registration.clone()],
1257 },
1258 );
1259 let added = study_add_response(
1260 test_session_status(),
1261 crate::StudyAddResult { registration },
1262 );
1263
1264 assert_eq!(listed.studies.len(), 1);
1265 assert_eq!(listed.studies[0].study.name, "cohort");
1266 assert_eq!(listed.studies[0].domains.len(), 1);
1267 assert_eq!(listed.studies[0].study.tags[0].label, "approved");
1268 assert_eq!(added.registration.study.name, "cohort");
1269 assert_eq!(added.registration.domains.len(), 1);
1270 }
1271
1272 #[test]
1273 fn study_workflow_errors_have_one_safe_public_meaning() {
1274 let list_missing =
1275 study_list_protocol_error(crate::AppError::NotFound("private Study scope".to_string()));
1276 let list_failed = study_list_protocol_error(crate::AppError::Infrastructure(
1277 "postgres://private.example/tre".to_string(),
1278 ));
1279 let registration_cases = [
1280 (
1281 crate::AppError::Validation("private bad Study".to_string()),
1282 ProtocolErrorCode::ValidationFailed,
1283 "Study registration request is invalid",
1284 ),
1285 (
1286 crate::AppError::NotFound("private Domain scope".to_string()),
1287 ProtocolErrorCode::NotFound,
1288 "Study registration scope was not found or is not authorized",
1289 ),
1290 (
1291 crate::AppError::Conflict("private metadata".to_string()),
1292 ProtocolErrorCode::Conflict,
1293 "Study registration conflicts with existing catalogue metadata",
1294 ),
1295 ];
1296
1297 assert_eq!(list_missing.code, ProtocolErrorCode::NotFound);
1298 assert_eq!(
1299 list_missing.message,
1300 "Catalogue scope was not found or is not authorized"
1301 );
1302 assert_eq!(list_failed.code, ProtocolErrorCode::OperationFailed);
1303 assert_eq!(
1304 list_failed.message,
1305 "Catalogue metadata could not be listed"
1306 );
1307 for (app_error, code, message) in registration_cases {
1308 let projected = study_registration_protocol_error(app_error);
1309 assert_eq!(projected.code, code);
1310 assert_eq!(projected.message, message);
1311 assert!(!projected.message.contains("private"));
1312 }
1313 }
1314
1315 #[test]
1316 fn datafile_list_workflow_has_one_safe_public_result() {
1317 let study_id = ahri_tre_types::StudyId(
1318 Uuid::parse_str("55555555-5555-4555-8555-555555555555").unwrap(),
1319 );
1320 let asset_id = ahri_tre_types::AssetId(
1321 Uuid::parse_str("66666666-6666-4666-8666-666666666666").unwrap(),
1322 );
1323 let version_id = ahri_tre_types::VersionId(
1324 Uuid::parse_str("77777777-7777-4777-8777-777777777777").unwrap(),
1325 );
1326 let result = crate::DataFileListResult::from_catalogue(
1327 ahri_tre_types::StudyRecord {
1328 study_id,
1329 name: ahri_tre_types::NcName::parse("cohort").unwrap(),
1330 description: None,
1331 documentation: None,
1332 agent_instructions: None,
1333 external_id: Some("EXT-003".to_string()),
1334 study_type_id: Some(ahri_tre_types::StudyTypeId(1)),
1335 date_created: None,
1336 created_by: None,
1337 },
1338 vec![crate::StudyDataFileEntry {
1339 catalog: crate::AssetVersionCatalog {
1340 asset: ahri_tre_types::AssetRecord {
1341 asset_id,
1342 study_id,
1343 name: ahri_tre_types::NcName::parse("baseline_csv").unwrap(),
1344 description: Some("Baseline data".to_string()),
1345 agent_instructions: Some("private analyst note".to_string()),
1346 asset_type: ahri_tre_types::AssetType::File,
1347 date_created: None,
1348 created_by: None,
1349 },
1350 versions: vec![ahri_tre_types::AssetVersionRecord {
1351 version_id,
1352 asset_id,
1353 major: 1,
1354 minor: 0,
1355 patch: 0,
1356 version_label: Some("v1.0.0".to_string()),
1357 version_note: None,
1358 is_latest: Some(true),
1359 doi: None,
1360 created_at: None,
1361 created_by: None,
1362 }],
1363 latest_version: None,
1364 },
1365 datafiles: vec![ahri_tre_types::DataFileRecord {
1366 datafile_id: version_id,
1367 size_bytes: Some(1_024),
1368 storage_uri: "lake://private/baseline.csv".to_string(),
1369 digest: ahri_tre_types::Sha256Hex::parse(
1370 "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
1371 )
1372 .unwrap(),
1373 encryption_key: Some(vec![1, 2, 3, 4]),
1374 compressed: Some(true),
1375 encrypted: Some(true),
1376 compression_algorithm: Some("zstd".to_string()),
1377 encryption_algorithm: Some("aes-256-gcm".to_string()),
1378 edam_format: "EDAM:format_3752".to_string(),
1379 }],
1380 }],
1381 );
1382 let app_json = serde_json::to_string(&result).unwrap();
1383
1384 let projected = datafile_list_response(test_session_status(), result);
1385 let json = serde_json::to_string(&projected).unwrap();
1386
1387 assert_eq!(projected.study, "cohort");
1388 assert_eq!(projected.study_id.id.as_uuid(), study_id.0);
1389 assert_eq!(projected.datafiles.len(), 1);
1390 assert_eq!(
1391 projected.datafiles[0].datafiles[0].edam_format,
1392 "EDAM:format_3752"
1393 );
1394 for forbidden in [
1395 "lake://",
1396 "private/baseline.csv",
1397 "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
1398 "storage_uri",
1399 "digest",
1400 "private analyst note",
1401 "encryption_key",
1402 ] {
1403 assert!(
1404 !app_json.contains(forbidden),
1405 "app result contains {forbidden}"
1406 );
1407 assert!(
1408 !json.contains(forbidden),
1409 "public result contains {forbidden}"
1410 );
1411 }
1412 }
1413
1414 #[test]
1415 fn datafile_list_errors_hide_missing_and_unauthorized_scope() {
1416 let missing =
1417 datafile_list_protocol_error(crate::AppError::NotFound("missing Study".to_string()));
1418 let unauthorized = datafile_list_protocol_error(crate::AppError::NotFound(
1419 "Study hidden by PostgreSQL policy".to_string(),
1420 ));
1421
1422 assert_eq!(missing, unauthorized);
1423 assert_eq!(missing.code, ProtocolErrorCode::NotFound);
1424 assert_eq!(
1425 missing.message,
1426 "Catalogue scope was not found or is not authorized"
1427 );
1428 }
1429
1430 #[test]
1431 fn dataset_materialization_errors_are_stable_and_redacted() {
1432 let cases = [
1433 (
1434 crate::AppError::Validation("private source path".to_string()),
1435 ProtocolErrorCode::ValidationFailed,
1436 ),
1437 (
1438 crate::AppError::NotFound("private Study identity".to_string()),
1439 ProtocolErrorCode::NotFound,
1440 ),
1441 (
1442 crate::AppError::Conflict("private version identity".to_string()),
1443 ProtocolErrorCode::Conflict,
1444 ),
1445 (
1446 crate::AppError::Infrastructure("postgres password".to_string()),
1447 ProtocolErrorCode::OperationFailed,
1448 ),
1449 ];
1450
1451 for (error, expected_code) in cases {
1452 let projected = dataset_materialization_protocol_error(error);
1453 assert_eq!(projected.code, expected_code);
1454 for forbidden in ["private", "path", "postgres", "password"] {
1455 assert!(!projected.message.to_ascii_lowercase().contains(forbidden));
1456 }
1457 }
1458 }
1459
1460 #[test]
1461 fn datafile_ingest_projection_has_one_safe_meaning_for_all_adapters() {
1462 let study_id = ahri_tre_types::StudyId(uuid::Uuid::new_v4());
1463 let asset_id = ahri_tre_types::AssetId(uuid::Uuid::new_v4());
1464 let version_id = ahri_tre_types::VersionId(uuid::Uuid::new_v4());
1465 let result = crate::DataFileIngestResult {
1466 study_id,
1467 study_name: ahri_tre_types::NcName::parse("cohort").unwrap(),
1468 catalog: crate::DataFileListAssetCatalog {
1469 asset: crate::DataFileListAssetSummary {
1470 asset_id,
1471 study_id,
1472 name: ahri_tre_types::NcName::parse("visits").unwrap(),
1473 description: Some("Visit rows".to_string()),
1474 asset_type: ahri_tre_types::AssetType::File,
1475 date_created: None,
1476 created_by: None,
1477 },
1478 versions: Vec::new(),
1479 latest_version: Some(ahri_tre_types::AssetVersionRecord {
1480 version_id,
1481 asset_id,
1482 major: 1,
1483 minor: 0,
1484 patch: 0,
1485 version_label: Some("v1.0.0".to_string()),
1486 version_note: None,
1487 is_latest: Some(true),
1488 doi: None,
1489 created_at: None,
1490 created_by: None,
1491 }),
1492 },
1493 datafile: crate::DataFileIngestVersion {
1494 datafile_id: version_id,
1495 edam_format: "EDAM:format_3752".to_string(),
1496 digest: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
1497 .to_string(),
1498 },
1499 transformation_id: ahri_tre_types::TransformationId(7),
1500 source_size_bytes: 18,
1501 stored_size_bytes: 18,
1502 };
1503
1504 let projected = datafile_ingest_response(
1505 test_session_status(),
1506 result,
1507 ahri_tre_protocol::ingest::IngestSourceSummary {
1508 kind: ahri_tre_protocol::ingest::IngestSourceKind::RestrictedLocalPath,
1509 request_only: true,
1510 upload: None,
1511 },
1512 );
1513 let json = serde_json::to_string(&projected).unwrap();
1514
1515 assert_eq!(projected.study, "cohort");
1516 assert_eq!(projected.asset.name, "visits");
1517 assert_eq!(projected.datafile.format, "EDAM:format_3752");
1518 assert_eq!(
1519 projected.datafile.digest_ref.as_deref(),
1520 Some("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa")
1521 );
1522 for forbidden in ["lake://", "storage_uri", "encryption_key", "/workspaces/"] {
1523 assert!(!json.contains(forbidden));
1524 }
1525 }
1526
1527 #[test]
1528 fn datafile_ingest_errors_have_one_safe_public_meaning() {
1529 let missing = datafile_ingest_protocol_error(crate::AppError::NotFound(
1530 "missing private Study".to_string(),
1531 ));
1532 let unauthorized = datafile_ingest_protocol_error(crate::AppError::NotFound(
1533 "PostgreSQL policy denied private Study".to_string(),
1534 ));
1535 let failed = datafile_ingest_protocol_error(crate::AppError::Infrastructure(
1536 "Lake location /private/lake and password leaked".to_string(),
1537 ));
1538
1539 assert_eq!(missing, unauthorized);
1540 assert_eq!(missing.code, ProtocolErrorCode::NotFound);
1541 assert_eq!(failed.code, ProtocolErrorCode::OperationFailed);
1542 for forbidden in ["private", "PostgreSQL", "/private/lake", "password"] {
1543 assert!(!missing.message.contains(forbidden));
1544 assert!(!failed.message.contains(forbidden));
1545 }
1546 }
1547
1548 #[test]
1549 fn projection_redacts_sensitive_public_messages() {
1550 let warning = safe_warning(
1551 ProtocolWarningCode::UnsupportedDetailOmitted,
1552 "token leaked at /workspaces/secret.sql select * from users",
1553 );
1554 assert!(!warning.message.contains("token"));
1555 assert!(!warning.message.contains("/workspaces"));
1556 assert!(!warning.message.contains("select"));
1557 }
1558
1559 #[test]
1560 fn app_error_projection_uses_public_error_codes() {
1561 let error = app_error_to_protocol_error(&AppError::Validation("bad request".to_string()));
1562 assert_eq!(error.code, ProtocolErrorCode::ValidationFailed);
1563 assert_eq!(error.message, "bad request");
1564
1565 let error =
1566 app_error_to_protocol_error(&AppError::NotImplemented("datastore.lake.move.resume"));
1567 assert_eq!(error.code, ProtocolErrorCode::UnsupportedOperation);
1568 assert_eq!(
1569 error.message,
1570 "unsupported operation: datastore.lake.move.resume"
1571 );
1572 }
1573}