1use std::collections::{BTreeMap, BTreeSet, HashMap};
2
3use arrow_schema::{Field, Schema};
4use serde_json::Value;
5
6use crate::{ArrowTable, Result};
7
8pub const AHRI_TRE_METADATA_PREFIX: &str = "ahri_tre.";
9
10pub mod keys {
11 pub const DATASET_ID: &str = "ahri_tre.dataset_id";
12 pub const ASSET_ID: &str = "ahri_tre.asset_id";
13 pub const ASSET_NAME: &str = "ahri_tre.asset_name";
14 pub const STUDY_ID: &str = "ahri_tre.study_id";
15 pub const VERSION_ID: &str = "ahri_tre.version_id";
16 pub const VERSION: &str = "ahri_tre.version";
17 pub const DESCRIPTION: &str = "ahri_tre.description";
18 pub const VERSION_NOTE: &str = "ahri_tre.version_note";
19
20 pub const VARIABLE_ID: &str = "ahri_tre.variable_id";
21 pub const VARIABLE_NAME: &str = "ahri_tre.variable_name";
22 pub const VALUE_TYPE_ID: &str = "ahri_tre.value_type_id";
23 pub const VALUE_TYPE: &str = "ahri_tre.value_type";
24 pub const VALUE_FORMAT: &str = "ahri_tre.value_format";
25 pub const ROW_ROLE: &str = "ahri_tre.row_role";
26 pub const NOTE: &str = "ahri_tre.note";
27 pub const ONTOLOGY_NAMESPACE: &str = "ahri_tre.ontology_namespace";
28 pub const ONTOLOGY_CLASS: &str = "ahri_tre.ontology_class";
29 pub const VOCABULARY_ID: &str = "ahri_tre.vocabulary_id";
30 pub const VOCABULARY_NAME: &str = "ahri_tre.vocabulary_name";
31 pub const VOCABULARY_DESCRIPTION: &str = "ahri_tre.vocabulary_description";
32 pub const VOCABULARY_ITEMS: &str = "ahri_tre.vocabulary_items";
33}
34
35const CONVENTIONAL_DESCRIPTION_KEYS: [&str; 4] = ["description", "title", "label", "comment"];
36const CONVENTIONAL_VOCABULARY_KEYS: [&str; 6] = [
37 "value_labels",
38 "value_descriptions",
39 "code_labels",
40 "code_descriptions",
41 "codes",
42 "categories",
43];
44
45#[derive(Debug, Clone, PartialEq, Eq, Default)]
46pub struct AhriTreArrowMetadata {
47 pub schema: BTreeMap<String, String>,
48 pub fields: Vec<AhriTreFieldMetadata>,
49}
50
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct AhriTreFieldMetadata {
53 pub field_name: String,
54 pub metadata: BTreeMap<String, String>,
55}
56
57#[derive(Debug, Clone)]
58pub struct AnnotatedArrowTable {
59 pub table: ArrowTable,
60 pub warnings: Vec<String>,
61}
62
63#[derive(Debug, Clone, PartialEq, Eq)]
64pub struct FieldMetadataProposal {
65 pub field_name: String,
66 pub variable_name: Option<String>,
67 pub value_type_id: Option<i64>,
68 pub value_type: Option<String>,
69 pub description: Option<String>,
70 pub vocabulary_name: Option<String>,
71 pub vocabulary_description: Option<String>,
72 pub vocabulary_items: Vec<VocabularyItemProposal>,
73 pub warnings: Vec<String>,
74}
75
76#[derive(Debug, Clone, PartialEq, Eq)]
77pub struct VocabularyItemProposal {
78 pub code: String,
79 pub description: Option<String>,
80}
81
82#[derive(Debug, Clone, PartialEq, Eq, Default)]
83pub struct JsonMetadataExtraction {
84 pub proposals: Vec<FieldMetadataProposal>,
85 pub warnings: Vec<String>,
86}
87
88pub fn annotate_arrow_table(
89 table: &ArrowTable,
90 metadata: &AhriTreArrowMetadata,
91) -> Result<AnnotatedArrowTable> {
92 let (schema, warnings) = annotate_arrow_schema(table.schema().as_ref(), metadata);
93 let table = table.with_schema(schema)?;
94 Ok(AnnotatedArrowTable { table, warnings })
95}
96
97pub fn annotate_arrow_schema(
98 schema: &Schema,
99 metadata: &AhriTreArrowMetadata,
100) -> (Schema, Vec<String>) {
101 let mut schema_metadata = without_ahri_tre_metadata(schema.metadata());
102 for (key, value) in &metadata.schema {
103 schema_metadata.insert(key.clone(), value.clone());
104 }
105
106 let field_metadata_index = FieldMetadataIndex::new(&metadata.fields);
107
108 let mut warnings = Vec::new();
109 let mut matched_fields = BTreeSet::new();
110 let fields = schema
111 .fields()
112 .iter()
113 .map(|field| {
114 let field = field.as_ref();
115 let mut merged = without_ahri_tre_metadata(field.metadata());
116 if let Some((field_index, tre_metadata)) =
117 field_metadata_index.find(field.name(), &matched_fields, &mut warnings)
118 {
119 for (key, value) in tre_metadata.iter() {
120 merged.insert(key.clone(), value.clone());
121 }
122 matched_fields.insert(field_index);
123 } else {
124 warnings.push(format!(
125 "Arrow field {} has no registered AHRI_TRE variable metadata",
126 field.name()
127 ));
128 }
129 clone_field_with_metadata(field, merged)
130 })
131 .collect::<Vec<_>>();
132
133 for (field_index, field_metadata) in metadata.fields.iter().enumerate() {
134 if !matched_fields.contains(&field_index) {
135 warnings.push(format!(
136 "registered AHRI_TRE variable {} is not present in the Arrow schema",
137 field_metadata.field_name
138 ));
139 }
140 }
141
142 (Schema::new_with_metadata(fields, schema_metadata), warnings)
143}
144
145struct FieldMetadataIndex<'a> {
146 fields: &'a [AhriTreFieldMetadata],
147 exact: BTreeMap<&'a str, Vec<usize>>,
148 canonical: BTreeMap<String, Vec<usize>>,
149}
150
151impl<'a> FieldMetadataIndex<'a> {
152 fn new(fields: &'a [AhriTreFieldMetadata]) -> Self {
153 let mut exact = BTreeMap::new();
154 let mut canonical = BTreeMap::new();
155 for (index, field_metadata) in fields.iter().enumerate() {
156 exact
157 .entry(field_metadata.field_name.as_str())
158 .or_insert_with(Vec::new)
159 .push(index);
160 canonical
161 .entry(canonical_field_name(&field_metadata.field_name))
162 .or_insert_with(Vec::new)
163 .push(index);
164 }
165 Self {
166 fields,
167 exact,
168 canonical,
169 }
170 }
171
172 fn find(
173 &self,
174 arrow_field_name: &str,
175 matched_fields: &BTreeSet<usize>,
176 warnings: &mut Vec<String>,
177 ) -> Option<(usize, &'a BTreeMap<String, String>)> {
178 if let Some(indexes) = self.exact.get(arrow_field_name) {
179 return self.unique_match(
180 arrow_field_name,
181 indexes,
182 matched_fields,
183 "exact field name",
184 warnings,
185 );
186 }
187
188 let canonical_name = canonical_field_name(arrow_field_name);
189 let indexes = self.canonical.get(&canonical_name)?;
190 self.unique_match(
191 arrow_field_name,
192 indexes,
193 matched_fields,
194 "canonical field name",
195 warnings,
196 )
197 }
198
199 fn unique_match(
200 &self,
201 arrow_field_name: &str,
202 indexes: &[usize],
203 matched_fields: &BTreeSet<usize>,
204 match_kind: &str,
205 warnings: &mut Vec<String>,
206 ) -> Option<(usize, &'a BTreeMap<String, String>)> {
207 let unmatched = indexes
208 .iter()
209 .copied()
210 .filter(|index| !matched_fields.contains(index))
211 .collect::<Vec<_>>();
212 match unmatched.as_slice() {
213 [index] => Some((*index, &self.fields[*index].metadata)),
214 [] => None,
215 _ => {
216 warnings.push(format!(
217 "Arrow field {arrow_field_name} has ambiguous registered AHRI_TRE variable metadata by {match_kind}"
218 ));
219 None
220 }
221 }
222 }
223}
224
225fn canonical_field_name(field_name: &str) -> String {
226 field_name.to_ascii_lowercase()
227}
228
229pub fn extract_field_metadata_proposals(schema: &Schema) -> Vec<FieldMetadataProposal> {
230 schema
231 .fields()
232 .iter()
233 .map(|field| extract_field_metadata_proposal(field.as_ref()))
234 .collect()
235}
236
237pub fn extract_field_metadata_proposal(field: &Field) -> FieldMetadataProposal {
238 let metadata = field.metadata();
239 let description = metadata
240 .get(keys::DESCRIPTION)
241 .cloned()
242 .or_else(|| conventional_description(metadata));
243 let mut warnings = Vec::new();
244 let vocabulary_items = extract_vocabulary_items(metadata).unwrap_or_else(|metadata_warnings| {
245 warnings.extend(metadata_warnings);
246 Vec::new()
247 });
248
249 FieldMetadataProposal {
250 field_name: field.name().clone(),
251 variable_name: metadata.get(keys::VARIABLE_NAME).cloned(),
252 value_type_id: metadata
253 .get(keys::VALUE_TYPE_ID)
254 .and_then(|value| value.trim().parse::<i64>().ok()),
255 value_type: metadata.get(keys::VALUE_TYPE).cloned(),
256 description,
257 vocabulary_name: metadata.get(keys::VOCABULARY_NAME).cloned(),
258 vocabulary_description: metadata.get(keys::VOCABULARY_DESCRIPTION).cloned(),
259 vocabulary_items,
260 warnings,
261 }
262}
263
264pub fn extract_json_field_metadata_proposals(text: &str) -> JsonMetadataExtraction {
265 let Ok(root) = serde_json::from_str::<Value>(text) else {
266 return JsonMetadataExtraction::default();
267 };
268 let Value::Object(root) = root else {
269 return JsonMetadataExtraction::default();
270 };
271
272 let common = extract_common_json_metadata(&root);
273 let ahri_tre = extract_ahri_tre_json_metadata(&root);
274 merge_json_metadata(common, ahri_tre)
275}
276
277fn extract_ahri_tre_json_metadata(root: &serde_json::Map<String, Value>) -> JsonMetadataExtraction {
278 let Some(contract) = root
279 .get("ahri_tre")
280 .or_else(|| root.get("ahri_tre_metadata"))
281 else {
282 return JsonMetadataExtraction::default();
283 };
284
285 let Some(fields) = contract
286 .get("fields")
287 .or_else(|| contract.get("variables"))
288 .and_then(Value::as_array)
289 else {
290 return JsonMetadataExtraction {
291 proposals: Vec::new(),
292 warnings: vec![
293 "ignored AHRI_TRE JSON metadata: expected ahri_tre.fields array".to_string(),
294 ],
295 };
296 };
297
298 let mut extraction = JsonMetadataExtraction::default();
299 for field in fields {
300 let Value::Object(field) = field else {
301 extraction
302 .warnings
303 .push("ignored AHRI_TRE JSON metadata field: expected object".to_string());
304 continue;
305 };
306 let Some(field_name) = field
307 .get("field_name")
308 .or_else(|| field.get("name"))
309 .and_then(Value::as_str)
310 .filter(|value| !value.trim().is_empty())
311 else {
312 extraction
313 .warnings
314 .push("ignored AHRI_TRE JSON metadata field: missing field_name".to_string());
315 continue;
316 };
317 let mut proposal = empty_field_metadata_proposal(field_name);
318 proposal.variable_name = field
319 .get("variable_name")
320 .or_else(|| field.get("tre_variable_name"))
321 .and_then(json_value_string_ref);
322 proposal.value_type_id = field
323 .get("value_type_id")
324 .and_then(Value::as_i64)
325 .or_else(|| {
326 field
327 .get("value_type_id")
328 .and_then(Value::as_str)
329 .and_then(|value| value.trim().parse::<i64>().ok())
330 });
331 proposal.value_type = field
332 .get("value_type")
333 .or_else(|| field.get("type"))
334 .and_then(json_value_string_ref);
335 proposal.description = field
336 .get("description")
337 .or_else(|| field.get("label"))
338 .and_then(json_value_string_ref);
339 proposal.vocabulary_name = field.get("vocabulary_name").and_then(json_value_string_ref);
340 proposal.vocabulary_description = field
341 .get("vocabulary_description")
342 .and_then(json_value_string_ref);
343 if let Some(value) = field
344 .get("vocabulary_items")
345 .or_else(|| field.get("categories"))
346 .or_else(|| field.get("codes"))
347 {
348 match parse_vocabulary_items_value("AHRI_TRE JSON vocabulary_items", value.clone()) {
349 Ok(items) => proposal.vocabulary_items = items,
350 Err(warnings) => proposal.warnings.extend(warnings),
351 }
352 }
353 extraction.proposals.push(proposal);
354 }
355 extraction
356}
357
358fn extract_common_json_metadata(root: &serde_json::Map<String, Value>) -> JsonMetadataExtraction {
359 if let Some(extraction) = extract_json_schema_metadata(root)
360 && (!extraction.proposals.is_empty() || !extraction.warnings.is_empty())
361 {
362 return extraction;
363 }
364 JsonMetadataExtraction::default()
365}
366
367fn extract_json_schema_metadata(
368 root: &serde_json::Map<String, Value>,
369) -> Option<JsonMetadataExtraction> {
370 let properties = root.get("properties")?.as_object()?;
371 let mut extraction = JsonMetadataExtraction::default();
372 for (field_name, descriptor) in properties {
373 let Value::Object(descriptor) = descriptor else {
374 extraction.warnings.push(format!(
375 "ignored JSON Schema property {field_name}: expected object descriptor"
376 ));
377 continue;
378 };
379 let mut proposal = empty_field_metadata_proposal(field_name);
380 proposal.value_type = descriptor
381 .get("type")
382 .and_then(json_schema_type_name)
383 .or_else(|| {
384 descriptor
385 .get("format")
386 .and_then(Value::as_str)
387 .and_then(json_schema_format_type_name)
388 });
389 proposal.description = descriptor
390 .get("description")
391 .or_else(|| descriptor.get("title"))
392 .and_then(json_value_string_ref);
393
394 if let Some(items) = json_schema_one_of_vocabulary_items(descriptor) {
395 proposal.vocabulary_items = items;
396 proposal.value_type = Some("enumeration".to_string());
397 } else if let Some(items) = json_schema_enum_names_vocabulary_items(descriptor) {
398 proposal.vocabulary_items = items;
399 proposal.value_type = Some("enumeration".to_string());
400 } else if descriptor.get("enum").is_some() {
401 proposal.warnings.push(format!(
402 "ignored JSON Schema property {field_name} enum: expected enumNames or oneOf const labels for categorical metadata"
403 ));
404 }
405 extraction.proposals.push(proposal);
406 }
407 Some(extraction)
408}
409
410fn merge_json_metadata(
411 common: JsonMetadataExtraction,
412 ahri_tre: JsonMetadataExtraction,
413) -> JsonMetadataExtraction {
414 let mut proposals_by_field = BTreeMap::new();
415 for proposal in common.proposals {
416 proposals_by_field.insert(proposal.field_name.clone(), proposal);
417 }
418 for ahri_proposal in ahri_tre.proposals {
419 proposals_by_field
420 .entry(ahri_proposal.field_name.clone())
421 .and_modify(|common_proposal| {
422 overlay_ahri_tre_proposal(common_proposal, &ahri_proposal)
423 })
424 .or_insert(ahri_proposal);
425 }
426 let mut warnings = common.warnings;
427 warnings.extend(ahri_tre.warnings);
428 JsonMetadataExtraction {
429 proposals: proposals_by_field.into_values().collect(),
430 warnings,
431 }
432}
433
434fn overlay_ahri_tre_proposal(common: &mut FieldMetadataProposal, ahri_tre: &FieldMetadataProposal) {
435 if ahri_tre.variable_name.is_some() {
436 common.variable_name = ahri_tre.variable_name.clone();
437 }
438 if ahri_tre.value_type_id.is_some() {
439 common.value_type_id = ahri_tre.value_type_id;
440 }
441 if ahri_tre.value_type.is_some() {
442 common.value_type = ahri_tre.value_type.clone();
443 }
444 if ahri_tre.description.is_some() {
445 common.description = ahri_tre.description.clone();
446 }
447 if ahri_tre.vocabulary_name.is_some() {
448 common.vocabulary_name = ahri_tre.vocabulary_name.clone();
449 }
450 if ahri_tre.vocabulary_description.is_some() {
451 common.vocabulary_description = ahri_tre.vocabulary_description.clone();
452 }
453 if !ahri_tre.vocabulary_items.is_empty() {
454 common.vocabulary_items = ahri_tre.vocabulary_items.clone();
455 }
456 common.warnings.extend(ahri_tre.warnings.clone());
457}
458
459fn empty_field_metadata_proposal(field_name: &str) -> FieldMetadataProposal {
460 FieldMetadataProposal {
461 field_name: field_name.to_string(),
462 variable_name: None,
463 value_type_id: None,
464 value_type: None,
465 description: None,
466 vocabulary_name: None,
467 vocabulary_description: None,
468 vocabulary_items: Vec::new(),
469 warnings: Vec::new(),
470 }
471}
472
473fn without_ahri_tre_metadata(metadata: &HashMap<String, String>) -> HashMap<String, String> {
474 metadata
475 .iter()
476 .filter(|(key, _)| !key.starts_with(AHRI_TRE_METADATA_PREFIX))
477 .map(|(key, value)| (key.clone(), value.clone()))
478 .collect()
479}
480
481fn clone_field_with_metadata(field: &Field, metadata: HashMap<String, String>) -> Field {
482 field.clone().with_metadata(metadata)
483}
484
485fn conventional_description(metadata: &HashMap<String, String>) -> Option<String> {
486 CONVENTIONAL_DESCRIPTION_KEYS
487 .iter()
488 .find_map(|key| metadata.get(*key).cloned())
489}
490
491fn extract_vocabulary_items(
492 metadata: &HashMap<String, String>,
493) -> std::result::Result<Vec<VocabularyItemProposal>, Vec<String>> {
494 if let Some(value) = metadata.get(keys::VOCABULARY_ITEMS) {
495 return parse_vocabulary_items(keys::VOCABULARY_ITEMS, value);
496 }
497
498 for key in CONVENTIONAL_VOCABULARY_KEYS {
499 if let Some(value) = metadata.get(key) {
500 return parse_vocabulary_items(key, value);
501 }
502 }
503 Ok(Vec::new())
504}
505
506fn parse_vocabulary_items(
507 key: &str,
508 value: &str,
509) -> std::result::Result<Vec<VocabularyItemProposal>, Vec<String>> {
510 let parsed = serde_json::from_str::<Value>(value).map_err(|error| {
511 vec![format!(
512 "ignored metadata key {key}: expected JSON code mapping ({error})"
513 )]
514 })?;
515 parse_vocabulary_items_value(key, parsed)
516}
517
518fn parse_vocabulary_items_value(
519 key: &str,
520 parsed: Value,
521) -> std::result::Result<Vec<VocabularyItemProposal>, Vec<String>> {
522 let items = match parsed {
523 Value::Object(map) => map
524 .into_iter()
525 .filter_map(|(code, value)| {
526 description_from_json_value(value).map(|description| VocabularyItemProposal {
527 code,
528 description: Some(description),
529 })
530 })
531 .collect::<Vec<_>>(),
532 Value::Array(values) => {
533 let mut items = Vec::with_capacity(values.len());
534 for value in values {
535 let Value::Object(mut object) = value else {
536 return Err(vec![format!(
537 "ignored metadata key {key}: array items must be objects with code and label or description"
538 )]);
539 };
540 let Some(code) = object.remove("code").and_then(json_string) else {
541 return Err(vec![format!(
542 "ignored metadata key {key}: vocabulary item is missing code"
543 )]);
544 };
545 let description = object
546 .remove("label")
547 .or_else(|| object.remove("description"))
548 .and_then(description_from_json_value);
549 items.push(VocabularyItemProposal { code, description });
550 }
551 items
552 }
553 _ => {
554 return Err(vec![format!(
555 "ignored metadata key {key}: expected JSON object or array"
556 )]);
557 }
558 };
559 if items.is_empty() {
560 Err(vec![format!(
561 "ignored metadata key {key}: no vocabulary items found"
562 )])
563 } else {
564 Ok(items)
565 }
566}
567
568fn json_schema_type_name(value: &Value) -> Option<String> {
569 match value {
570 Value::String(value) => match value.as_str() {
571 "integer" => Some("integer".to_string()),
572 "number" => Some("decimal".to_string()),
573 "string" => Some("string".to_string()),
574 "boolean" => Some("boolean".to_string()),
575 _ => None,
576 },
577 Value::Array(values) => values.iter().find_map(json_schema_type_name),
578 _ => None,
579 }
580}
581
582fn json_schema_format_type_name(value: &str) -> Option<String> {
583 match value {
584 "date" => Some("date".to_string()),
585 "date-time" => Some("datetime".to_string()),
586 "time" => Some("time".to_string()),
587 _ => None,
588 }
589}
590
591fn json_schema_enum_names_vocabulary_items(
592 descriptor: &serde_json::Map<String, Value>,
593) -> Option<Vec<VocabularyItemProposal>> {
594 let enums = descriptor.get("enum")?.as_array()?;
595 let names = descriptor
596 .get("enumNames")
597 .or_else(|| descriptor.get("enumDescriptions"))?
598 .as_array()?;
599 if enums.len() != names.len() || enums.is_empty() {
600 return None;
601 }
602 let mut items = Vec::with_capacity(enums.len());
603 for (code, description) in enums.iter().zip(names) {
604 items.push(VocabularyItemProposal {
605 code: json_value_string_ref(code)?,
606 description: Some(json_value_string_ref(description)?),
607 });
608 }
609 Some(items)
610}
611
612fn json_schema_one_of_vocabulary_items(
613 descriptor: &serde_json::Map<String, Value>,
614) -> Option<Vec<VocabularyItemProposal>> {
615 let one_of = descriptor
616 .get("oneOf")
617 .or_else(|| descriptor.get("anyOf"))?
618 .as_array()?;
619 if one_of.is_empty() {
620 return None;
621 }
622 let mut items = Vec::with_capacity(one_of.len());
623 for item in one_of {
624 let item = item.as_object()?;
625 let code = item
626 .get("const")
627 .or_else(|| item.get("enum").and_then(|value| value.as_array()?.first()))
628 .and_then(json_value_string_ref)?;
629 let description = item
630 .get("title")
631 .or_else(|| item.get("description"))
632 .and_then(json_value_string_ref)?;
633 items.push(VocabularyItemProposal {
634 code,
635 description: Some(description),
636 });
637 }
638 Some(items)
639}
640
641fn description_from_json_value(value: Value) -> Option<String> {
642 match value {
643 Value::String(value) => Some(value),
644 Value::Object(mut object) => object
645 .remove("label")
646 .or_else(|| object.remove("description"))
647 .and_then(json_string),
648 _ => None,
649 }
650}
651
652fn json_string(value: Value) -> Option<String> {
653 match value {
654 Value::String(value) => Some(value),
655 Value::Number(value) => Some(value.to_string()),
656 Value::Bool(value) => Some(value.to_string()),
657 _ => None,
658 }
659}
660
661fn json_value_string_ref(value: &Value) -> Option<String> {
662 match value {
663 Value::String(value) => Some(value.clone()),
664 Value::Number(value) => Some(value.to_string()),
665 Value::Bool(value) => Some(value.to_string()),
666 _ => None,
667 }
668}
669
670pub fn extract_json_file_metadata_proposals(
673 path: &std::path::Path,
674) -> std::io::Result<JsonMetadataExtraction> {
675 use serde::Deserializer;
676 use serde::de::{MapAccess, Visitor};
677 use std::{
678 cell::Cell,
679 io::{BufReader, Read},
680 rc::Rc,
681 };
682 const LIMIT: usize = crate::streaming::MAX_BATCH_BYTES;
683 struct BudgetReader<R> {
684 inner: R,
685 remaining: Rc<Cell<Option<usize>>>,
686 }
687 impl<R: Read> Read for BudgetReader<R> {
688 fn read(&mut self, out: &mut [u8]) -> std::io::Result<usize> {
689 if let Some(left) = self.remaining.get() {
690 if left == 0 {
691 return Err(std::io::Error::other("JSON metadata exceeds its limit"));
692 }
693 let len = out.len().min(left);
694 let count = self.inner.read(&mut out[..len])?;
695 self.remaining.set(Some(left - count));
696 Ok(count)
697 } else {
698 self.inner.read(out)
699 }
700 }
701 }
702 struct Metadata(Rc<Cell<Option<usize>>>);
703 impl<'de> Visitor<'de> for Metadata {
704 type Value = serde_json::Map<String, Value>;
705 fn expecting(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
706 f.write_str("JSON object")
707 }
708 fn visit_map<M: MapAccess<'de>>(
709 self,
710 mut map: M,
711 ) -> std::result::Result<Self::Value, M::Error> {
712 let mut metadata = serde_json::Map::new();
713 let mut retained = 0usize;
714 loop {
715 self.0.set(Some(LIMIT - retained));
716 let Some(key) = map.next_key::<String>()? else {
717 break;
718 };
719 if matches!(
720 key.as_str(),
721 "properties" | "ahri_tre" | "ahri_tre_metadata"
722 ) {
723 let value = map.next_value::<Value>()?;
724 retained = LIMIT - self.0.get().unwrap_or(0);
725 metadata.insert(key, value);
726 } else {
727 self.0.set(None);
728 map.next_value::<serde::de::IgnoredAny>()?;
729 }
730 }
731 Ok(metadata)
732 }
733 }
734 let remaining = Rc::new(Cell::new(Some(LIMIT)));
735 let reader = BudgetReader {
736 inner: BufReader::new(std::fs::File::open(path)?),
737 remaining: remaining.clone(),
738 };
739 let mut deserializer = serde_json::Deserializer::from_reader(reader);
740 let root = match deserializer.deserialize_map(Metadata(remaining)) {
741 Ok(root) => root,
742 Err(error) if error.is_io() => {
743 return Err(std::io::Error::other("JSON metadata exceeds its limit"));
744 }
745 Err(_) => return Ok(JsonMetadataExtraction::default()),
746 };
747 Ok(merge_json_metadata(
748 extract_common_json_metadata(&root),
749 extract_ahri_tre_json_metadata(&root),
750 ))
751}
752
753#[cfg(test)]
754mod tests {
755 use super::*;
756 use arrow_schema::DataType;
757
758 #[test]
759 fn ahri_tre_description_wins_over_conventional_description() {
760 let field = Field::new("sex", DataType::Utf8, true).with_metadata(
761 [
762 (
763 "description".to_string(),
764 "Conventional description".to_string(),
765 ),
766 (
767 keys::DESCRIPTION.to_string(),
768 "AHRI TRE description".to_string(),
769 ),
770 ]
771 .into_iter()
772 .collect(),
773 );
774
775 let proposal = extract_field_metadata_proposal(&field);
776
777 assert_eq!(
778 proposal.description.as_deref(),
779 Some("AHRI TRE description")
780 );
781 }
782
783 #[test]
784 fn annotation_matches_registered_metadata_by_canonical_field_name() {
785 let schema = Schema::new(vec![Field::new("HouseholdId", DataType::Utf8, true)]);
786 let metadata = AhriTreArrowMetadata {
787 schema: BTreeMap::new(),
788 fields: vec![AhriTreFieldMetadata {
789 field_name: "householdid".to_string(),
790 metadata: [
791 (keys::VARIABLE_NAME.to_string(), "householdid".to_string()),
792 (keys::VALUE_TYPE_ID.to_string(), "7".to_string()),
793 ]
794 .into_iter()
795 .collect(),
796 }],
797 };
798
799 let (annotated_schema, warnings) = annotate_arrow_schema(&schema, &metadata);
800
801 assert!(warnings.is_empty(), "{warnings:?}");
802 let field = annotated_schema.field_with_name("HouseholdId").unwrap();
803 assert_eq!(
804 field.metadata()[keys::VARIABLE_NAME],
805 "householdid",
806 "canonical variable metadata should attach to the original Arrow field name"
807 );
808 }
809
810 #[test]
811 fn annotation_warns_instead_of_guessing_ambiguous_canonical_metadata() {
812 let schema = Schema::new(vec![Field::new("HouseholdId", DataType::Utf8, true)]);
813 let metadata = AhriTreArrowMetadata {
814 schema: BTreeMap::new(),
815 fields: vec![
816 AhriTreFieldMetadata {
817 field_name: "householdid".to_string(),
818 metadata: [(keys::VARIABLE_NAME.to_string(), "householdid".to_string())]
819 .into_iter()
820 .collect(),
821 },
822 AhriTreFieldMetadata {
823 field_name: "HOUSEHOLDID".to_string(),
824 metadata: [(keys::VARIABLE_NAME.to_string(), "HOUSEHOLDID".to_string())]
825 .into_iter()
826 .collect(),
827 },
828 ],
829 };
830
831 let (annotated_schema, warnings) = annotate_arrow_schema(&schema, &metadata);
832
833 assert!(
834 warnings
835 .iter()
836 .any(|warning| warning.contains("ambiguous registered AHRI_TRE variable metadata")),
837 "{warnings:?}"
838 );
839 assert!(
840 !annotated_schema
841 .field_with_name("HouseholdId")
842 .unwrap()
843 .metadata()
844 .contains_key(keys::VARIABLE_NAME)
845 );
846 }
847
848 #[test]
849 fn conventional_description_keys_are_recognized_but_unrelated_metadata_is_ignored() {
850 let field = Field::new("age", DataType::Int64, true).with_metadata(
851 [
852 ("title".to_string(), "Age at visit".to_string()),
853 ("owner".to_string(), "not TRE semantics".to_string()),
854 ]
855 .into_iter()
856 .collect(),
857 );
858
859 let proposal = extract_field_metadata_proposal(&field);
860
861 assert_eq!(proposal.description.as_deref(), Some("Age at visit"));
862 assert!(proposal.variable_name.is_none());
863 assert!(proposal.vocabulary_items.is_empty());
864 }
865
866 #[test]
867 fn ahri_tre_vocabulary_items_preserve_codes_and_order() {
868 let field = Field::new("status", DataType::Utf8, true).with_metadata(
869 [(
870 keys::VOCABULARY_ITEMS.to_string(),
871 r#"[{"code":"A","label":"Alive"},{"code":"D","description":"Dead"}]"#.to_string(),
872 )]
873 .into_iter()
874 .collect(),
875 );
876
877 let proposal = extract_field_metadata_proposal(&field);
878
879 assert_eq!(
880 proposal.vocabulary_items,
881 vec![
882 VocabularyItemProposal {
883 code: "A".to_string(),
884 description: Some("Alive".to_string()),
885 },
886 VocabularyItemProposal {
887 code: "D".to_string(),
888 description: Some("Dead".to_string()),
889 },
890 ]
891 );
892 assert!(proposal.warnings.is_empty());
893 }
894
895 #[test]
896 fn conventional_code_label_mapping_creates_vocabulary_items() {
897 let field = Field::new("sex", DataType::Utf8, true).with_metadata(
898 [(
899 "value_labels".to_string(),
900 r#"{"F":"Female","M":"Male"}"#.to_string(),
901 )]
902 .into_iter()
903 .collect(),
904 );
905
906 let proposal = extract_field_metadata_proposal(&field);
907
908 assert_eq!(
909 proposal.vocabulary_items,
910 vec![
911 VocabularyItemProposal {
912 code: "F".to_string(),
913 description: Some("Female".to_string()),
914 },
915 VocabularyItemProposal {
916 code: "M".to_string(),
917 description: Some("Male".to_string()),
918 },
919 ]
920 );
921 }
922
923 #[test]
924 fn ahri_tre_json_metadata_extracts_field_contracts() {
925 let extraction = extract_json_field_metadata_proposals(
926 r#"{
927 "ahri_tre": {
928 "fields": [{
929 "field_name": "status",
930 "variable_name": "status_code",
931 "value_type": "enumeration",
932 "description": "AHRI TRE status",
933 "vocabulary_name": "status_vocab",
934 "vocabulary_description": "Status code list",
935 "vocabulary_items": [
936 {"code": "A", "label": "Alive"},
937 {"code": "D", "description": "Dead"}
938 ]
939 }]
940 }
941 }"#,
942 );
943
944 assert!(extraction.warnings.is_empty());
945 assert_eq!(extraction.proposals.len(), 1);
946 let proposal = &extraction.proposals[0];
947 assert_eq!(proposal.field_name, "status");
948 assert_eq!(proposal.variable_name.as_deref(), Some("status_code"));
949 assert_eq!(proposal.value_type.as_deref(), Some("enumeration"));
950 assert_eq!(proposal.description.as_deref(), Some("AHRI TRE status"));
951 assert_eq!(
952 proposal.vocabulary_items,
953 vec![
954 VocabularyItemProposal {
955 code: "A".to_string(),
956 description: Some("Alive".to_string()),
957 },
958 VocabularyItemProposal {
959 code: "D".to_string(),
960 description: Some("Dead".to_string()),
961 },
962 ]
963 );
964 }
965
966 #[test]
967 fn json_schema_metadata_extracts_descriptions_types_and_one_of_vocabularies() {
968 let extraction = extract_json_field_metadata_proposals(
969 r#"{
970 "$schema": "https://json-schema.org/draft/2020-12/schema",
971 "type": "object",
972 "properties": {
973 "age": {"type": "integer", "description": "Age in years"},
974 "sex": {
975 "type": "string",
976 "title": "Sex at baseline",
977 "oneOf": [
978 {"const": "F", "title": "Female"},
979 {"const": "M", "description": "Male"}
980 ]
981 }
982 }
983 }"#,
984 );
985
986 assert!(extraction.warnings.is_empty());
987 let age = extraction
988 .proposals
989 .iter()
990 .find(|proposal| proposal.field_name == "age")
991 .expect("age proposal should exist");
992 assert_eq!(age.value_type.as_deref(), Some("integer"));
993 assert_eq!(age.description.as_deref(), Some("Age in years"));
994 let sex = extraction
995 .proposals
996 .iter()
997 .find(|proposal| proposal.field_name == "sex")
998 .expect("sex proposal should exist");
999 assert_eq!(sex.value_type.as_deref(), Some("enumeration"));
1000 assert_eq!(sex.description.as_deref(), Some("Sex at baseline"));
1001 assert_eq!(
1002 sex.vocabulary_items,
1003 vec![
1004 VocabularyItemProposal {
1005 code: "F".to_string(),
1006 description: Some("Female".to_string()),
1007 },
1008 VocabularyItemProposal {
1009 code: "M".to_string(),
1010 description: Some("Male".to_string()),
1011 },
1012 ]
1013 );
1014 }
1015
1016 #[test]
1017 fn ahri_tre_json_metadata_overrides_common_schema_for_same_field() {
1018 let extraction = extract_json_field_metadata_proposals(
1019 r#"{
1020 "properties": {
1021 "status": {
1022 "type": "string",
1023 "description": "Common status",
1024 "enum": ["A", "D"],
1025 "enumNames": ["Active", "Dormant"]
1026 }
1027 },
1028 "ahri_tre": {
1029 "fields": [{
1030 "field_name": "status",
1031 "variable_name": "status_code",
1032 "description": "AHRI TRE status",
1033 "vocabulary_items": {"A": "Alive", "D": "Dead"}
1034 }]
1035 }
1036 }"#,
1037 );
1038
1039 assert_eq!(extraction.proposals.len(), 1);
1040 let proposal = &extraction.proposals[0];
1041 assert_eq!(proposal.variable_name.as_deref(), Some("status_code"));
1042 assert_eq!(proposal.description.as_deref(), Some("AHRI TRE status"));
1043 assert_eq!(
1044 proposal
1045 .vocabulary_items
1046 .iter()
1047 .map(|item| (item.code.as_str(), item.description.as_deref()))
1048 .collect::<Vec<_>>(),
1049 vec![("A", Some("Alive")), ("D", Some("Dead"))]
1050 );
1051 }
1052
1053 #[test]
1054 fn ordinary_json_data_is_not_interpreted_as_metadata() {
1055 let extraction = extract_json_field_metadata_proposals(
1056 r#"[{"status":"A","age":34},{"status":"D","age":41}]"#,
1057 );
1058
1059 assert!(extraction.proposals.is_empty());
1060 assert!(extraction.warnings.is_empty());
1061 }
1062
1063 #[test]
1064 fn ambiguous_json_schema_enum_is_reported_without_vocab_semantics() {
1065 let extraction = extract_json_field_metadata_proposals(
1066 r#"{
1067 "properties": {
1068 "status": {
1069 "type": "string",
1070 "description": "Status",
1071 "enum": ["A", "D"]
1072 }
1073 }
1074 }"#,
1075 );
1076
1077 assert_eq!(extraction.proposals.len(), 1);
1078 let proposal = &extraction.proposals[0];
1079 assert!(proposal.vocabulary_items.is_empty());
1080 assert!(
1081 proposal
1082 .warnings
1083 .iter()
1084 .any(|warning| warning.contains("expected enumNames or oneOf const labels")),
1085 "{:?}",
1086 proposal.warnings
1087 );
1088 }
1089}